diff --git a/internal/devserver/server.go b/internal/devserver/server.go index 09d47854f..5ef689a7e 100644 --- a/internal/devserver/server.go +++ b/internal/devserver/server.go @@ -40,6 +40,7 @@ import ( uiserveroptions "github.com/temporalio/ui-server/v2/server/server_options" "go.temporal.io/api/enums/v1" "go.temporal.io/server/chasm/lib/activity" + chasmcallback "go.temporal.io/server/chasm/lib/callback" "go.temporal.io/server/common/authorization" "go.temporal.io/server/common/cluster" "go.temporal.io/server/common/config" @@ -250,6 +251,14 @@ func (s *StartOptions) buildServerOptions() ([]temporal.ServerOption, *slog.Leve dynConf[activity.EnableStandaloneActivityOperatorCommands.Key()] = true dynConf[dynamicconfig.FrontendEnableBatchOperationsForStandaloneActivities.Key()] = true + // The server calls no callback address until one is allowed. A local + // receiver is the usual target of a notification channel's callback on + // a dev server, and it rarely serves TLS. + dynConf[chasmcallback.AllowedAddresses.Key()] = []any{ + map[string]any{"Pattern": "127.0.0.1:*", "AllowInsecure": true}, + map[string]any{"Pattern": "localhost:*", "AllowInsecure": true}, + } + // Dynamic config if set for k, v := range s.DynamicConfigValues { dynConf[dynamicconfig.MakeKey(k)] = v diff --git a/internal/temporalcli/commands.channel.go b/internal/temporalcli/commands.channel.go new file mode 100644 index 000000000..35919f214 --- /dev/null +++ b/internal/temporalcli/commands.channel.go @@ -0,0 +1,649 @@ +package temporalcli + +import ( + "context" + "encoding/base64" + "encoding/json" + "errors" + "fmt" + "reflect" + "sort" + "strconv" + "strings" + "time" + "unicode" + "unicode/utf8" + + "github.com/fatih/color" + "github.com/google/uuid" + "github.com/temporalio/cli/internal/printer" + commonpb "go.temporal.io/api/common/v1" + enumspb "go.temporal.io/api/enums/v1" + notificationpb "go.temporal.io/api/notification/v1" + "go.temporal.io/api/serviceerror" + workflowpb "go.temporal.io/api/workflow/v1" + "go.temporal.io/api/workflowservice/v1" + "go.temporal.io/sdk/client" + "google.golang.org/protobuf/types/known/durationpb" +) + +// channelPollGrace is how long past the requested wait a poll may take before +// the CLI gives up on it, so a Service answering at its own deadline is not cut +// off. +const channelPollGrace = 10 * time.Second + +// channelTarget is the channel a command names: an independent channel by its +// name alone, or a channel linked to an execution by that owner and the name. +type channelTarget struct { + name string + execution *commonpb.Execution +} + +func (o *ChannelOptions) target() (channelTarget, error) { + if o.WorkflowId != "" && o.ActivityId != "" { + return channelTarget{}, fmt.Errorf( + "--workflow-id and --activity-id name different owners, set one of them") + } + if o.RunId != "" && o.WorkflowId == "" && o.ActivityId == "" { + return channelTarget{}, fmt.Errorf("--run-id requires --workflow-id or --activity-id") + } + t := channelTarget{name: o.Channel} + switch { + case o.WorkflowId != "": + t.execution = &commonpb.Execution{ + Type: enumspb.EXECUTION_TYPE_WORKFLOW, + BusinessId: o.WorkflowId, + RunId: o.RunId, + } + case o.ActivityId != "": + t.execution = &commonpb.Execution{ + Type: enumspb.EXECUTION_TYPE_ACTIVITY, + BusinessId: o.ActivityId, + RunId: o.RunId, + } + } + return t, nil +} + +// executionKindText is the owner's kind the way the flags name it, so output +// reads back as the flag that reaches the same owner. +func executionKindText(e *commonpb.Execution) string { + switch e.GetType() { + case enumspb.EXECUTION_TYPE_ACTIVITY: + return "activity" + case enumspb.EXECUTION_TYPE_NEXUS_OPERATION: + return "nexus operation" + default: + return "workflow" + } +} + +// executionText names an execution for a card or a message: its kind and ID, +// then the run when one was given. +func executionText(e *commonpb.Execution) string { + if e == nil { + return "" + } + text := executionKindText(e) + " " + e.GetBusinessId() + if e.GetRunId() != "" { + text += " (run " + e.GetRunId() + ")" + } + return text +} + +func (t channelTarget) String() string { + if t.execution == nil { + return fmt.Sprintf("channel %q", t.name) + } + return fmt.Sprintf("channel %q of %s %q", t.name, executionKindText(t.execution), + t.execution.GetBusinessId()) +} + +// label is the target for a line of normal output, where quotes would be noise. +func (t channelTarget) label() string { + if t.execution == nil { + return "channel " + t.name + } + return "channel " + t.name + " of " + executionKindText(t.execution) + " " + + t.execution.GetBusinessId() +} + +// plainChannelRefusal turns the Service's refusals into what to do next. Every +// other error comes back as it was. +func plainChannelRefusal(err error, t channelTarget) error { + var notFound *serviceerror.NotFound + if errors.As(err, ¬Found) { + plain := fmt.Sprintf("there is no channel %q. A channel exists once a writer "+ + "notifies it or a listener registers on it, and goes away after a while "+ + "with neither.", t.name) + if t.execution != nil { + plain = fmt.Sprintf("%s %q has no running execution to reach. A linked "+ + "channel lives only while its owner runs.", + executionKindText(t.execution), t.execution.GetBusinessId()) + } + return &refusalError{plain: plain, detail: notFound.Message, cause: err} + } + var exhausted *serviceerror.ResourceExhausted + if !errors.As(err, &exhausted) { + return err + } + var plain string + switch exhausted.Cause { + case enumspb.RESOURCE_EXHAUSTED_CAUSE_CONCURRENT_LIMIT: + plain = "the channel already has as many listeners as it may hold, so this one " + + "was not registered. Remove a listener it no longer needs, or use another channel." + case enumspb.RESOURCE_EXHAUSTED_CAUSE_RPS_LIMIT: + plain = "the namespace is notifying channels faster than it may, so this " + + "notification was not sent. Retry after a moment, or notify less often." + default: + plain = "the Service is over one of its limits and refused the call. Retry after " + + "a moment." + } + return &refusalError{plain: plain, detail: exhausted.Message, cause: err} +} + +// parseChannelMetadata reads KEY=VALUE pairs whose values are JSON, carried as +// JSON payloads the way --input carries a Signal's arguments. +func parseChannelMetadata(pairs []string) (map[string]*commonpb.Payload, error) { + if len(pairs) == 0 { + return nil, nil + } + out := make(map[string]*commonpb.Payload, len(pairs)) + for _, pair := range pairs { + key, value, ok := strings.Cut(pair, "=") + if !ok || key == "" { + return nil, fmt.Errorf("--metadata %q must be KEY=VALUE", pair) + } + if _, dup := out[key]; dup { + return nil, fmt.Errorf("--metadata key %q is given more than once", key) + } + if !json.Valid([]byte(value)) { + return nil, fmt.Errorf("--metadata value for %q is not valid JSON", key) + } + out[key] = &commonpb.Payload{ + Metadata: map[string][]byte{"encoding": []byte("json/plain")}, + Data: []byte(value), + } + } + return out, nil +} + +// positionText shows a position as the text it most often is, and as base64 +// when it holds bytes a terminal would mangle. +func positionText(p []byte) string { + if utf8.Valid(p) && strings.IndexFunc(string(p), func(r rune) bool { + return !unicode.IsPrint(r) + }) < 0 { + return string(p) + } + return "base64:" + base64.StdEncoding.EncodeToString(p) +} + +// metadataText puts a notification's metadata on one line, keys sorted so the +// same notification always reads the same. +func metadataText(md map[string]*commonpb.Payload) (string, error) { + keys := make([]string, 0, len(md)) + for k := range md { + keys = append(keys, k) + } + sort.Strings(keys) + parts := make([]string, 0, len(keys)) + for _, k := range keys { + v, err := payloadText(md[k]) + if err != nil { + return "", err + } + parts = append(parts, k+"="+v) + } + return strings.Join(parts, " "), nil +} + +type channelNotificationRow struct { + Counter int64 + Position string + Metadata string +} + +func notificationRow(n *notificationpb.Notification) (channelNotificationRow, error) { + md, err := metadataText(n.GetMetadata()) + if err != nil { + return channelNotificationRow{}, err + } + return channelNotificationRow{ + Counter: n.GetCounter(), + Position: positionText(n.GetPosition()), + Metadata: md, + }, nil +} + +// channelKindText reads a server that predates the linked kind, and so leaves +// the kind unset, as the independent channel it serves. +func channelKindText(kind notificationpb.ChannelKind) string { + if kind == notificationpb.CHANNEL_KIND_UNSPECIFIED { + kind = notificationpb.CHANNEL_KIND_INDEPENDENT + } + return kind.String() +} + +// printChannelSubscriptions lists the channels a workflow stands on, as its +// description reports them. A workflow that stands on none prints nothing, so +// the section is absent on servers that predate it as well. +func printChannelSubscriptions( + cctx *CommandContext, subs []*workflowpb.ChannelSubscriptionInfo, +) error { + if len(subs) == 0 { + return nil + } + rows := make([]struct { + Channel string + Kind string + LastCounter int64 + PendingCounter string + ScheduledCounter int64 + Listeners int32 + Retained int32 + }, len(subs)) + for i, s := range subs { + rows[i].Channel = s.GetChannel() + rows[i].Kind = channelKindText(s.GetKind()) + rows[i].LastCounter = s.GetLastCounter() + if pending := s.GetPendingNotification(); pending != nil { + rows[i].PendingCounter = strconv.FormatInt(pending.GetCounter(), 10) + } + rows[i].ScheduledCounter = s.GetScheduledCounter() + rows[i].Listeners = s.GetListenerCount() + rows[i].Retained = s.GetRetainedCount() + } + cctx.Printer.Println() + cctx.Printer.Println(color.MagentaString("Notification Channels: %v", len(subs))) + cctx.Printer.Println() + return cctx.Printer.PrintStructured(rows, printer.StructuredOptions{ + Table: &printer.TableOptions{}, + }) +} + +func (c *TemporalChannelNotifyCommand) run(cctx *CommandContext, _ []string) error { + target, err := c.target() + if err != nil { + return err + } + if c.Counter <= 0 { + return fmt.Errorf("--counter must be greater than zero") + } + metadata, err := parseChannelMetadata(c.Metadata) + if err != nil { + return err + } + cl, err := dialClient(cctx, &c.Parent.ClientOptions) + if err != nil { + return err + } + defer cl.Close() + + resp, err := cl.WorkflowService().NotifyChannel(cctx, &workflowservice.NotifyChannelRequest{ + Namespace: c.Parent.Namespace, + Notification: ¬ificationpb.Notification{ + Channel: c.Channel, + Position: []byte(c.Position), + Counter: int64(c.Counter), + Metadata: metadata, + }, + Identity: c.Parent.Identity, + RequestId: uuid.NewString(), + Execution: target.execution, + }) + if err != nil { + return fmt.Errorf("failed notifying %v: %w", target, plainChannelRefusal(err, target)) + } + if cctx.JSONOutput { + return cctx.Printer.PrintStructured(resp, printer.StructuredOptions{}) + } + cctx.Printer.Printlnf("Notified %v. Listeners reached: %v", + target.label(), resp.GetListenerCount()) + return nil +} + +func (c *TemporalChannelDescribeCommand) run(cctx *CommandContext, _ []string) error { + target, err := c.target() + if err != nil { + return err + } + cl, err := dialClient(cctx, &c.Parent.ClientOptions) + if err != nil { + return err + } + defer cl.Close() + + resp, err := cl.WorkflowService().DescribeChannel(cctx, + &workflowservice.DescribeChannelRequest{ + Namespace: c.Parent.Namespace, + Channel: c.Channel, + Execution: target.execution, + }) + if err != nil { + return fmt.Errorf("failed describing %v: %w", target, plainChannelRefusal(err, target)) + } + if cctx.JSONOutput { + return cctx.Printer.PrintStructured(resp, printer.StructuredOptions{}) + } + + info := struct { + Channel string + Kind string + LinkedTo string `cli:",cardOmitEmpty"` + RetainedCount int32 + LatestCounter int64 `cli:",cardOmitEmpty"` + LatestPosition string `cli:",cardOmitEmpty"` + LatestMetadata string `cli:",cardOmitEmpty"` + }{ + Channel: c.Channel, + Kind: channelKindText(resp.GetKind()), + LinkedTo: executionText(resp.GetLinkedTo()), + RetainedCount: resp.GetRetainedCount(), + } + if latest := resp.GetLatest(); latest != nil { + row, err := notificationRow(latest) + if err != nil { + return err + } + info.LatestCounter, info.LatestPosition, info.LatestMetadata = + row.Counter, row.Position, row.Metadata + } + cctx.Printer.Println(color.MagentaString("Channel:")) + if err := cctx.Printer.PrintStructured(info, printer.StructuredOptions{}); err != nil { + return err + } + + type workflowRow struct { + ListenerId string + WorkflowId string + RunId string + Registered time.Time + } + // Headers stay out of the table because they often carry credentials. + type callbackRow struct { + ListenerId string + Url string + Registered time.Time + } + var workflows []workflowRow + var callbacks []callbackRow + for _, l := range resp.GetListeners() { + // A linked channel's owner never registered, so it has no time; the + // conversion would turn that into the Unix epoch. + var registered time.Time + if l.GetRegisteredTime() != nil { + registered = l.GetRegisteredTime().AsTime() + } + if wf := l.GetWorkflow(); wf != nil { + workflows = append(workflows, workflowRow{ + ListenerId: l.GetListenerId(), + WorkflowId: wf.GetWorkflowId(), + RunId: wf.GetRunId(), + Registered: registered, + }) + } else if cb := l.GetCallback(); cb != nil { + callbacks = append(callbacks, callbackRow{ + ListenerId: l.GetListenerId(), + Url: cb.GetNexus().GetUrl(), + Registered: registered, + }) + } + } + sort.Slice(workflows, func(i, j int) bool { + return workflows[i].ListenerId < workflows[j].ListenerId + }) + sort.Slice(callbacks, func(i, j int) bool { + return callbacks[i].ListenerId < callbacks[j].ListenerId + }) + cctx.Printer.Println() + cctx.Printer.Println(color.MagentaString("Workflow listeners: %v", len(workflows))) + if len(workflows) > 0 { + if err := cctx.Printer.PrintStructured(workflows, printer.StructuredOptions{ + Table: &printer.TableOptions{}, + }); err != nil { + return err + } + } + cctx.Printer.Println() + cctx.Printer.Println(color.MagentaString("Callback listeners: %v", len(callbacks))) + if len(callbacks) > 0 { + return cctx.Printer.PrintStructured(callbacks, printer.StructuredOptions{ + Table: &printer.TableOptions{}, + }) + } + return nil +} + +func (c *TemporalChannelListenerAddCommand) run(cctx *CommandContext, _ []string) error { + target, err := c.target() + if err != nil { + return err + } + if c.CallbackUrl == "" { + return fmt.Errorf("--callback-url cannot be empty") + } + headers, err := stringKeysValues(c.Header) + if err != nil { + return fmt.Errorf("invalid --header: %w", err) + } + clientOpts := &c.Parent.Parent.ClientOptions + cl, err := dialClient(cctx, clientOpts) + if err != nil { + return err + } + defer cl.Close() + + resp, err := cl.WorkflowService().RegisterChannelListener(cctx, + &workflowservice.RegisterChannelListenerRequest{ + Namespace: clientOpts.Namespace, + Channel: c.Channel, + Callback: &commonpb.Callback{ + Variant: &commonpb.Callback_Nexus_{Nexus: &commonpb.Callback_Nexus{ + Url: c.CallbackUrl, + Header: headers, + }}, + }, + RequestId: uuid.NewString(), + Identity: clientOpts.Identity, + Execution: target.execution, + }) + if err != nil { + return fmt.Errorf("failed adding a listener to %v: %w", + target, plainChannelRefusal(err, target)) + } + if cctx.JSONOutput { + return cctx.Printer.PrintStructured(resp, printer.StructuredOptions{}) + } + cctx.Printer.Printlnf("Added listener %v to %v", resp.GetListenerId(), target.label()) + return nil +} + +func (c *TemporalChannelListenerRemoveCommand) run(cctx *CommandContext, _ []string) error { + target, err := c.target() + if err != nil { + return err + } + clientOpts := &c.Parent.Parent.ClientOptions + cl, err := dialClient(cctx, clientOpts) + if err != nil { + return err + } + defer cl.Close() + + _, err = cl.WorkflowService().UnregisterChannelListener(cctx, + &workflowservice.UnregisterChannelListenerRequest{ + Namespace: clientOpts.Namespace, + Channel: c.Channel, + ListenerId: c.ListenerId, + Identity: clientOpts.Identity, + Execution: target.execution, + }) + if err != nil { + return fmt.Errorf("failed removing listener %q from %v: %w", + c.ListenerId, target, plainChannelRefusal(err, target)) + } + cctx.Printer.Printlnf("Removed listener %v from %v", c.ListenerId, target.label()) + return nil +} + +// channelPoller hands out a channel's notifications one at a time. Without +// follow it makes one poll; with follow it polls again from the highest +// counter it has seen until the command is interrupted. +type channelPoller struct { + ctx context.Context + cl client.Client + namespace string + target channelTarget + after int64 + wait time.Duration + max int32 + follow bool + buf []*notificationpb.Notification + polled bool +} + +func (p *channelPoller) next() (*notificationpb.Notification, error) { + for len(p.buf) == 0 && (p.follow || !p.polled) { + if err := p.poll(); err != nil { + // An interrupted command is the poller leaving, not a failure. + if p.ctx.Err() != nil { + return nil, nil + } + return nil, fmt.Errorf("failed polling %v: %w", + p.target, plainChannelRefusal(err, p.target)) + } + } + if len(p.buf) == 0 { + return nil, nil + } + n := p.buf[0] + p.buf = p.buf[1:] + return n, nil +} + +func (p *channelPoller) poll() error { + ctx, cancel := context.WithTimeout(p.ctx, p.wait+channelPollGrace) + defer cancel() + resp, err := p.cl.WorkflowService().PollChannel(ctx, &workflowservice.PollChannelRequest{ + Namespace: p.namespace, + Channel: p.target.name, + AfterCounter: p.after, + Wait: durationpb.New(p.wait), + MaxNotifications: p.max, + Execution: p.target.execution, + }) + p.polled = true + if err != nil { + // A poll the Service held past its deadline is an empty answer, the same + // as one it returned with nothing. + if ctx.Err() != nil && p.ctx.Err() == nil { + return nil + } + return err + } + for _, n := range resp.GetNotifications() { + p.after = max(p.after, n.GetCounter()) + } + p.buf = append(p.buf, resp.GetNotifications()...) + return nil +} + +// channelRowIter adapts the poller to the printer's streaming table. +type channelRowIter struct{ poller *channelPoller } + +func (i *channelRowIter) Next() (any, error) { + n, err := i.poller.next() + if n == nil || err != nil { + return nil, err + } + return notificationRow(n) +} + +func (c *TemporalChannelPollCommand) run(cctx *CommandContext, _ []string) error { + target, err := c.target() + if err != nil { + return err + } + if c.AfterCounter < 0 { + return fmt.Errorf("--after-counter cannot be negative") + } + if c.Max < 0 { + return fmt.Errorf("--max cannot be negative") + } + if c.Wait.Duration() <= 0 { + return fmt.Errorf("--wait must be positive") + } + cl, err := dialClient(cctx, &c.Parent.ClientOptions) + if err != nil { + return err + } + defer cl.Close() + + poller := &channelPoller{ + ctx: cctx, + cl: cl, + namespace: c.Parent.Namespace, + target: target, + after: int64(c.AfterCounter), + wait: c.Wait.Duration(), + max: int32(c.Max), + follow: c.Follow, + } + if cctx.JSONOutput { + // This is a listing command subject to json vs jsonl rules + cctx.Printer.StartList() + defer cctx.Printer.EndList() + for { + n, err := poller.next() + if err != nil { + return err + } + if n == nil { + return nil + } + if err := cctx.Printer.PrintStructured(n, printer.StructuredOptions{}); err != nil { + return err + } + } + } + if !c.Follow { + var rows []channelNotificationRow + for { + n, err := poller.next() + if err != nil { + return err + } + if n == nil { + break + } + row, err := notificationRow(n) + if err != nil { + return err + } + rows = append(rows, row) + } + if len(rows) == 0 { + cctx.Printer.Printlnf("No notifications on %v after counter %v", + target.label(), c.AfterCounter) + return nil + } + return cctx.Printer.PrintStructured(rows, printer.StructuredOptions{ + Table: &printer.TableOptions{}, + }) + } + return cctx.Printer.PrintStructuredTableIter( + reflect.TypeOf(channelNotificationRow{}), + &channelRowIter{poller: poller}, + printer.StructuredOptions{ + Table: &printer.TableOptions{ + // Streaming rows cannot be measured first, so the columns + // before the metadata take fixed widths. + FieldWidths: map[string]int{ + "Counter": 12, + "Position": 24, + }, + }, + }, + ) +} diff --git a/internal/temporalcli/commands.channel_test.go b/internal/temporalcli/commands.channel_test.go new file mode 100644 index 000000000..800708649 --- /dev/null +++ b/internal/temporalcli/commands.channel_test.go @@ -0,0 +1,1079 @@ +package temporalcli_test + +import ( + "context" + "encoding/json" + "io" + "net" + "net/http" + "net/http/httptest" + "strings" + "sync" + "testing" + "time" + + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/temporalio/cli/internal/temporalcli" + commonpb "go.temporal.io/api/common/v1" + enumspb "go.temporal.io/api/enums/v1" + notificationpb "go.temporal.io/api/notification/v1" + "go.temporal.io/api/serviceerror" + workflowpb "go.temporal.io/api/workflow/v1" + "go.temporal.io/api/workflowservice/v1" + "go.temporal.io/sdk/client" + "go.temporal.io/sdk/workflow" + "google.golang.org/grpc" + "google.golang.org/protobuf/types/known/timestamppb" +) + +// fakeChannelService answers the channel calls from canned responses and +// records what the CLI sent, so the commands are checked without a server +// that carries channels. +type fakeChannelService struct { + workflowservice.UnimplementedWorkflowServiceServer + + mu sync.Mutex + notifies []*workflowservice.NotifyChannelRequest + registers []*workflowservice.RegisterChannelListenerRequest + unregister []*workflowservice.UnregisterChannelListenerRequest + polls []*workflowservice.PollChannelRequest + describes []*workflowservice.DescribeChannelRequest + + err error + describe *workflowservice.DescribeChannelResponse + // pollAnswers are handed out in order; past the end a poll waits for the + // caller to give up, the way an idle channel does. + pollAnswers [][]*notificationpb.Notification + // onPoll runs with the number of polls seen so far. + onPoll func(n int) + // describeWorkflow answers the workflow description with its channels. + describeWorkflow *workflowservice.DescribeWorkflowExecutionResponse +} + +func (f *fakeChannelService) DescribeWorkflowExecution( + context.Context, *workflowservice.DescribeWorkflowExecutionRequest, +) (*workflowservice.DescribeWorkflowExecutionResponse, error) { + f.mu.Lock() + defer f.mu.Unlock() + if f.err != nil { + return nil, f.err + } + return f.describeWorkflow, nil +} + +func (f *fakeChannelService) GetSystemInfo( + context.Context, *workflowservice.GetSystemInfoRequest, +) (*workflowservice.GetSystemInfoResponse, error) { + return &workflowservice.GetSystemInfoResponse{}, nil +} + +func (f *fakeChannelService) NotifyChannel( + _ context.Context, req *workflowservice.NotifyChannelRequest, +) (*workflowservice.NotifyChannelResponse, error) { + f.mu.Lock() + defer f.mu.Unlock() + f.notifies = append(f.notifies, req) + if f.err != nil { + return nil, f.err + } + return &workflowservice.NotifyChannelResponse{ListenerCount: 3}, nil +} + +func (f *fakeChannelService) RegisterChannelListener( + _ context.Context, req *workflowservice.RegisterChannelListenerRequest, +) (*workflowservice.RegisterChannelListenerResponse, error) { + f.mu.Lock() + defer f.mu.Unlock() + f.registers = append(f.registers, req) + if f.err != nil { + return nil, f.err + } + return &workflowservice.RegisterChannelListenerResponse{ListenerId: "listener-7"}, nil +} + +func (f *fakeChannelService) UnregisterChannelListener( + _ context.Context, req *workflowservice.UnregisterChannelListenerRequest, +) (*workflowservice.UnregisterChannelListenerResponse, error) { + f.mu.Lock() + defer f.mu.Unlock() + f.unregister = append(f.unregister, req) + if f.err != nil { + return nil, f.err + } + return &workflowservice.UnregisterChannelListenerResponse{}, nil +} + +func (f *fakeChannelService) DescribeChannel( + _ context.Context, req *workflowservice.DescribeChannelRequest, +) (*workflowservice.DescribeChannelResponse, error) { + f.mu.Lock() + f.describes = append(f.describes, req) + f.mu.Unlock() + if f.err != nil { + return nil, f.err + } + return f.describe, nil +} + +func (f *fakeChannelService) PollChannel( + ctx context.Context, req *workflowservice.PollChannelRequest, +) (*workflowservice.PollChannelResponse, error) { + f.mu.Lock() + f.polls = append(f.polls, req) + n := len(f.polls) + var answer []*notificationpb.Notification + idle := n > len(f.pollAnswers) + if !idle { + answer = f.pollAnswers[n-1] + } + onPoll, err := f.onPoll, f.err + f.mu.Unlock() + if onPoll != nil { + onPoll(n) + } + if err != nil { + return nil, err + } + if idle { + <-ctx.Done() + return nil, ctx.Err() + } + return &workflowservice.PollChannelResponse{Notifications: answer}, nil +} + +func startFakeChannelService(t *testing.T, f *fakeChannelService) string { + ln, err := net.Listen("tcp", "127.0.0.1:0") + require.NoError(t, err) + // The Service sends its errors as gRPC statuses; a bare server would send + // them as unknown errors and lose their kind. + srv := grpc.NewServer(grpc.UnaryInterceptor(func( + ctx context.Context, req any, _ *grpc.UnaryServerInfo, handler grpc.UnaryHandler, + ) (any, error) { + resp, err := handler(ctx, req) + if err != nil { + return nil, serviceerror.ToStatus(err).Err() + } + return resp, nil + })) + workflowservice.RegisterWorkflowServiceServer(srv, f) + go func() { _ = srv.Serve(ln) }() + t.Cleanup(srv.Stop) + return ln.Addr().String() +} + +func jsonPayload(data string) *commonpb.Payload { + return &commonpb.Payload{ + Metadata: map[string][]byte{"encoding": []byte("json/plain")}, + Data: []byte(data), + } +} + +func TestChannel_ArgumentValidation(t *testing.T) { + f := &fakeChannelService{} + addr := startFakeChannelService(t, f) + h := NewCommandHarness(t) + + for _, tc := range []struct { + args []string + want string + }{ + {[]string{"notify", "--position", "1", "--counter", "1"}, `"channel" not set`}, + {[]string{"notify", "-c", "ch", "--counter", "1"}, `"position" not set`}, + {[]string{"notify", "-c", "ch", "--position", "1", "--counter", "0"}, + "--counter must be greater than zero"}, + {[]string{"notify", "-c", "ch", "--position", "1", "--counter", "1", + "--metadata", "topic"}, "must be KEY=VALUE"}, + {[]string{"notify", "-c", "ch", "--position", "1", "--counter", "1", + "--metadata", "topic=scores"}, "not valid JSON"}, + {[]string{"notify", "-c", "ch", "--position", "1", "--counter", "1", + "--metadata", `a="x"`, "--metadata", `a="y"`}, "more than once"}, + {[]string{"listener", "add", "-c", "ch"}, `"callback-url" not set`}, + {[]string{"listener", "add", "-c", "ch", "--callback-url", "http://x", + "--header", "nope"}, "invalid --header"}, + {[]string{"listener", "remove", "-c", "ch"}, `"listener-id" not set`}, + {[]string{"poll", "-c", "ch", "--after-counter", "-1"}, "cannot be negative"}, + {[]string{"poll", "-c", "ch", "--max", "-1"}, "cannot be negative"}, + {[]string{"poll", "-c", "ch", "--wait", "0s"}, "--wait must be positive"}, + {[]string{"describe"}, `"channel" not set`}, + {[]string{"describe", "-c", "ch", "--run-id", "r1"}, + "--run-id requires --workflow-id or --activity-id"}, + {[]string{"notify", "-c", "ch", "--position", "1", "--counter", "1", "-r", "r1"}, + "--run-id requires --workflow-id or --activity-id"}, + {[]string{"poll", "-c", "ch", "-r", "r1"}, "--run-id requires --workflow-id or --activity-id"}, + {[]string{"listener", "add", "-c", "ch", "--callback-url", "http://x", "-r", "r1"}, + "--run-id requires --workflow-id or --activity-id"}, + {[]string{"listener", "remove", "-c", "ch", "--listener-id", "l", "-r", "r1"}, + "--run-id requires --workflow-id or --activity-id"}, + {[]string{"describe", "-c", "ch", "-w", "wf", "--activity-id", "act"}, + "--workflow-id and --activity-id name different owners"}, + {[]string{"notify", "-c", "ch", "--position", "1", "--counter", "1", "-w", "wf", + "--activity-id", "act"}, "--workflow-id and --activity-id name different owners"}, + {[]string{"poll", "-c", "ch", "-w", "wf", "--activity-id", "act"}, + "--workflow-id and --activity-id name different owners"}, + } { + res := h.Execute(append(append([]string{"channel"}, tc.args...), "--address", addr)...) + require.Error(t, res.Err, tc.args) + assert.Contains(t, res.Err.Error(), tc.want, tc.args) + } + // None of the refused commands reached the Service. + assert.Empty(t, f.notifies) + assert.Empty(t, f.registers) + assert.Empty(t, f.unregister) + assert.Empty(t, f.polls) + assert.Empty(t, f.describes) +} + +func TestChannel_Notify(t *testing.T) { + f := &fakeChannelService{} + addr := startFakeChannelService(t, f) + h := NewCommandHarness(t) + + res := h.Execute( + "channel", "notify", "--address", addr, "--namespace", "ns1", + "-c", "stream/scores", "--position", "42", "--counter", "42", + "--metadata", `topic="scores"`, "--metadata", `batch={"n": 2}`, + ) + require.NoError(t, res.Err) + h.ContainsOnSameLine(res.Stdout.String(), "Notified channel stream/scores", "3") + + require.Len(t, f.notifies, 1) + req := f.notifies[0] + assert.Equal(t, "ns1", req.GetNamespace()) + assert.NotEmpty(t, req.GetIdentity()) + assert.NotEmpty(t, req.GetRequestId()) + n := req.GetNotification() + assert.Equal(t, "stream/scores", n.GetChannel()) + assert.Equal(t, []byte("42"), n.GetPosition()) + assert.Equal(t, int64(42), n.GetCounter()) + require.Len(t, n.GetMetadata(), 2) + assert.Equal(t, "json/plain", string(n.GetMetadata()["topic"].GetMetadata()["encoding"])) + assert.Equal(t, `"scores"`, string(n.GetMetadata()["topic"].GetData())) + assert.Equal(t, `{"n": 2}`, string(n.GetMetadata()["batch"].GetData())) + + // A retried command is a new request, so the Service does not fold it + // with the first. + res = h.Execute( + "channel", "notify", "--address", addr, "-c", "ch", "--position", "p", + "--counter", "1", "-o", "json", + ) + require.NoError(t, res.Err) + var out workflowservice.NotifyChannelResponse + require.NoError(t, temporalcli.UnmarshalProtoJSONWithOptions(res.Stdout.Bytes(), &out, true)) + assert.Equal(t, int32(3), out.GetListenerCount()) + require.Len(t, f.notifies, 2) + assert.NotEqual(t, f.notifies[0].GetRequestId(), f.notifies[1].GetRequestId()) +} + +func TestChannel_Describe(t *testing.T) { + registered := time.Date(2026, 10, 1, 12, 0, 0, 0, time.UTC) + f := &fakeChannelService{describe: &workflowservice.DescribeChannelResponse{ + Listeners: []*notificationpb.ChannelListener{ + { + ListenerId: "l-callback", + Listener: ¬ificationpb.ChannelListener_Callback{Callback: &commonpb.Callback{ + Variant: &commonpb.Callback_Nexus_{Nexus: &commonpb.Callback_Nexus{ + Url: "https://example.com/notify", + Header: map[string]string{"Authorization": "Bearer secret"}, + }}, + }}, + RegisteredTime: timestamppb.New(registered), + }, + { + ListenerId: "l-workflow", + Listener: ¬ificationpb.ChannelListener_Workflow{ + Workflow: ¬ificationpb.WorkflowListener{WorkflowId: "wf-1", RunId: "run-1"}, + }, + RegisteredTime: timestamppb.New(registered), + }, + }, + Latest: ¬ificationpb.Notification{ + Channel: "ch", + Position: []byte{0x00, 0xff}, + Counter: 9, + Metadata: map[string]*commonpb.Payload{ + "topic": jsonPayload(`"scores"`), + "a": jsonPayload(`1`), + }, + }, + RetainedCount: 4, + }} + addr := startFakeChannelService(t, f) + h := NewCommandHarness(t) + + res := h.Execute("channel", "describe", "--address", addr, "-c", "ch") + require.NoError(t, res.Err) + out := res.Stdout.String() + h.ContainsOnSameLine(out, "Channel", "ch") + h.ContainsOnSameLine(out, "RetainedCount", "4") + h.ContainsOnSameLine(out, "LatestCounter", "9") + // Bytes a terminal would mangle print as base64. + h.ContainsOnSameLine(out, "LatestPosition", "base64:AP8=") + h.ContainsOnSameLine(out, "LatestMetadata", `a=1 topic="scores"`) + h.ContainsOnSameLine(out, "Workflow listeners: 1") + h.ContainsOnSameLine(out, "l-workflow", "wf-1", "run-1") + h.ContainsOnSameLine(out, "Callback listeners: 1") + h.ContainsOnSameLine(out, "l-callback", "https://example.com/notify") + assert.NotContains(t, out, "secret", "headers stay out of the text output") + + res = h.Execute("channel", "describe", "--address", addr, "-c", "ch", "-o", "json") + require.NoError(t, res.Err) + var described workflowservice.DescribeChannelResponse + require.NoError(t, temporalcli.UnmarshalProtoJSONWithOptions(res.Stdout.Bytes(), &described, true)) + assert.Equal(t, int32(4), described.GetRetainedCount()) + assert.Len(t, described.GetListeners(), 2) + assert.Equal(t, []byte{0x00, 0xff}, described.GetLatest().GetPosition()) + + // A channel that has seen nothing yet still describes. + f.describe = &workflowservice.DescribeChannelResponse{} + res = h.Execute("channel", "describe", "--address", addr, "-c", "ch") + require.NoError(t, res.Err) + assert.NotContains(t, res.Stdout.String(), "LatestCounter") + h.ContainsOnSameLine(res.Stdout.String(), "Workflow listeners: 0") +} + +func TestChannel_Listener(t *testing.T) { + f := &fakeChannelService{} + addr := startFakeChannelService(t, f) + h := NewCommandHarness(t) + + res := h.Execute( + "channel", "listener", "add", "--address", addr, "-c", "ch", + "--callback-url", "https://example.com/notify", + "--header", "Authorization=Bearer t=1", "--header", "X-Team=ai", + ) + require.NoError(t, res.Err) + h.ContainsOnSameLine(res.Stdout.String(), "Added listener listener-7 to channel ch") + require.Len(t, f.registers, 1) + reg := f.registers[0] + assert.Equal(t, "ch", reg.GetChannel()) + assert.NotEmpty(t, reg.GetRequestId()) + assert.NotEmpty(t, reg.GetIdentity()) + nexus := reg.GetCallback().GetNexus() + assert.Equal(t, "https://example.com/notify", nexus.GetUrl()) + assert.Equal(t, map[string]string{ + "Authorization": "Bearer t=1", + "X-Team": "ai", + }, nexus.GetHeader()) + + res = h.Execute( + "channel", "listener", "remove", "--address", addr, "-c", "ch", + "--listener-id", "listener-7", + ) + require.NoError(t, res.Err) + h.ContainsOnSameLine(res.Stdout.String(), "Removed listener listener-7 from channel ch") + require.Len(t, f.unregister, 1) + assert.Equal(t, "listener-7", f.unregister[0].GetListenerId()) + assert.Equal(t, "ch", f.unregister[0].GetChannel()) +} + +func TestChannel_Poll(t *testing.T) { + f := &fakeChannelService{pollAnswers: [][]*notificationpb.Notification{{ + {Channel: "ch", Position: []byte("offset-4"), Counter: 4, + Metadata: map[string]*commonpb.Payload{"topic": jsonPayload(`"scores"`)}}, + {Channel: "ch", Position: []byte("offset-6"), Counter: 6}, + }}} + addr := startFakeChannelService(t, f) + h := NewCommandHarness(t) + + res := h.Execute( + "channel", "poll", "--address", addr, "-c", "ch", + "--after-counter", "3", "--wait", "5s", "--max", "10", + ) + require.NoError(t, res.Err) + out := res.Stdout.String() + h.ContainsOnSameLine(out, "Counter", "Position", "Metadata") + h.ContainsOnSameLine(out, "4", "offset-4", `topic="scores"`) + h.ContainsOnSameLine(out, "6", "offset-6") + require.Len(t, f.polls, 1) + assert.Equal(t, int64(3), f.polls[0].GetAfterCounter()) + assert.Equal(t, 5*time.Second, f.polls[0].GetWait().AsDuration()) + assert.Equal(t, int32(10), f.polls[0].GetMaxNotifications()) + + // Nothing retained and nothing arriving within the wait reads as such. + f.pollAnswers = append(f.pollAnswers, nil) + res = h.Execute("channel", "poll", "--address", addr, "-c", "ch", "--after-counter", "6") + require.NoError(t, res.Err) + h.ContainsOnSameLine(res.Stdout.String(), "No notifications on channel ch after counter 6") +} + +func TestChannel_PollFollow(t *testing.T) { + h := NewCommandHarness(t) + f := &fakeChannelService{ + pollAnswers: [][]*notificationpb.Notification{ + {{Channel: "ch", Position: []byte("a"), Counter: 2}}, + {}, + { + {Channel: "ch", Position: []byte("b"), Counter: 5}, + {Channel: "ch", Position: []byte("c"), Counter: 3}, + }, + }, + // The fourth poll is the command waiting on an idle channel; stop it + // there the way an interrupt would. + onPoll: func(n int) { + if n == 4 { + h.CancelContext() + } + }, + } + addr := startFakeChannelService(t, f) + + res := h.Execute("channel", "poll", "--address", addr, "-c", "ch", "--follow", "-o", "jsonl") + require.NoError(t, res.Err) + var counters []int64 + for _, raw := range decodeJSONValues(t, res.Stdout.String()) { + var n notificationpb.Notification + require.NoError(t, temporalcli.UnmarshalProtoJSONWithOptions(raw, &n, true)) + counters = append(counters, n.GetCounter()) + } + assert.Equal(t, []int64{2, 5, 3}, counters) + + // Each poll starts after the highest counter seen so far, and an empty + // answer keeps the cursor where it was. + f.mu.Lock() + defer f.mu.Unlock() + require.Len(t, f.polls, 4) + var after []int64 + for _, p := range f.polls { + after = append(after, p.GetAfterCounter()) + } + assert.Equal(t, []int64{0, 2, 2, 5}, after) +} + +func TestChannel_PlainRefusals(t *testing.T) { + f := &fakeChannelService{} + addr := startFakeChannelService(t, f) + h := NewCommandHarness(t) + + f.err = serviceerror.NewNotFound("channel not found") + res := h.Execute("channel", "describe", "--address", addr, "-c", "missing") + require.Error(t, res.Err) + assert.Contains(t, res.Err.Error(), `there is no channel "missing"`) + assert.Contains(t, res.Err.Error(), "The server said: channel not found") + res = h.Execute("channel", "poll", "--address", addr, "-c", "missing") + require.Error(t, res.Err) + assert.Contains(t, res.Err.Error(), `there is no channel "missing"`) + + f.err = serviceerror.NewResourceExhaustedf(enumspb.RESOURCE_EXHAUSTED_CAUSE_CONCURRENT_LIMIT, + "channel already has 1000 listeners, which is the limit") + res = h.Execute( + "channel", "listener", "add", "--address", addr, "-c", "ch", "--callback-url", "http://x", + ) + require.Error(t, res.Err) + assert.Contains(t, res.Err.Error(), "as many listeners as it may hold") + assert.Contains(t, res.Err.Error(), "1000 listeners") + + f.err = &serviceerror.ResourceExhausted{ + Cause: enumspb.RESOURCE_EXHAUSTED_CAUSE_RPS_LIMIT, + Scope: enumspb.RESOURCE_EXHAUSTED_SCOPE_NAMESPACE, + Message: "namespace is over its limit", + } + res = h.Execute( + "channel", "notify", "--address", addr, "-c", "ch", "--position", "1", "--counter", "1", + ) + require.Error(t, res.Err) + assert.Contains(t, res.Err.Error(), "notifying channels faster than it may") + + // Errors the CLI has nothing to add to pass through as they came. + f.err = serviceerror.NewInvalidArgument("notification position exceeds 1024 bytes") + res = h.Execute( + "channel", "notify", "--address", addr, "-c", "ch", "--position", "1", "--counter", "1", + ) + require.Error(t, res.Err) + assert.Contains(t, res.Err.Error(), "position exceeds 1024 bytes") + assert.NotContains(t, res.Err.Error(), "The server said") +} + +// decodeJSONValues splits output holding one JSON value after another. +func decodeJSONValues(t *testing.T, s string) []json.RawMessage { + var out []json.RawMessage + dec := json.NewDecoder(strings.NewReader(s)) + for dec.More() { + var raw json.RawMessage + require.NoError(t, dec.Decode(&raw)) + out = append(out, raw) + } + return out +} + +func (s *SharedServerSuite) TestChannel_NotifyWithoutListenersIsRetained() { + ch := "channel-" + uuid.NewString() + res := s.Execute( + "channel", "notify", "--address", s.Address(), "-c", ch, + "--position", "offset-1", "--counter", "1", "--metadata", `topic="scores"`, + ) + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "Listeners reached", "0") + + res = s.Execute("channel", "describe", "--address", s.Address(), "-c", ch) + s.NoError(res.Err) + out := res.Stdout.String() + s.ContainsOnSameLine(out, "RetainedCount", "1") + s.ContainsOnSameLine(out, "LatestCounter", "1") + s.ContainsOnSameLine(out, "LatestPosition", "offset-1") + s.ContainsOnSameLine(out, "LatestMetadata", `topic="scores"`) + s.ContainsOnSameLine(out, "Callback listeners: 0") +} + +func (s *SharedServerSuite) TestChannel_StaleCounterIsFolded() { + ch := "channel-" + uuid.NewString() + notify := func(position, counter string) { + res := s.Execute( + "channel", "notify", "--address", s.Address(), "-c", ch, + "--position", position, "--counter", counter, + ) + s.NoError(res.Err) + } + notify("p5", "5") + // An equal or lower counter is folded into the one already kept. + notify("p5-again", "5") + notify("p3", "3") + res := s.Execute("channel", "describe", "--address", s.Address(), "-c", ch) + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "RetainedCount", "1") + s.ContainsOnSameLine(res.Stdout.String(), "LatestPosition", "p5") + + notify("p7", "7") + notify("p9", "9") + res = s.Execute( + "channel", "poll", "--address", s.Address(), "-c", ch, + "--after-counter", "5", "--wait", "2s", "-o", "jsonl", + ) + s.NoError(res.Err) + var counters []int64 + for _, raw := range decodeJSONValues(s.T(), res.Stdout.String()) { + var n notificationpb.Notification + s.NoError(temporalcli.UnmarshalProtoJSONWithOptions(raw, &n, true)) + counters = append(counters, n.GetCounter()) + } + s.Equal([]int64{7, 9}, counters) +} + +func (s *SharedServerSuite) TestChannel_PollWaitsForTheNextNotification() { + ch := "channel-" + uuid.NewString() + res := s.Execute( + "channel", "notify", "--address", s.Address(), "-c", ch, "--position", "p1", "--counter", "1", + ) + s.NoError(res.Err) + + done := make(chan *CommandResult, 1) + started := time.Now() + go func() { + done <- s.Execute( + "channel", "poll", "--address", s.Address(), "-c", ch, + "--after-counter", "1", "--wait", "20s", "-o", "json", + ) + }() + // The poll is parked on the channel by the time this lands. + time.Sleep(time.Second) + res = s.Execute( + "channel", "notify", "--address", s.Address(), "-c", ch, "--position", "p2", "--counter", "2", + ) + s.NoError(res.Err) + + select { + case res = <-done: + case <-time.After(60 * time.Second): + s.Fail("poll did not return after the notification") + } + s.NoError(res.Err) + s.Less(time.Since(started), 15*time.Second, "the poll answered on arrival, not at its wait") + raw := decodeJSONValues(s.T(), res.Stdout.String()) + s.Len(raw, 1) + var list []json.RawMessage + s.NoError(json.Unmarshal(raw[0], &list)) + s.Len(list, 1) + var n notificationpb.Notification + s.NoError(temporalcli.UnmarshalProtoJSONWithOptions(list[0], &n, true)) + s.Equal(int64(2), n.GetCounter()) + s.Equal([]byte("p2"), n.GetPosition()) +} + +func (s *SharedServerSuite) TestChannel_CallbackListenerRoundTrip() { + type delivery struct { + auth string + body []byte + } + deliveries := make(chan delivery, 8) + receiver := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(r.Body) + deliveries <- delivery{auth: r.Header.Get("Authorization"), body: body} + })) + defer receiver.Close() + + ch := "channel-" + uuid.NewString() + res := s.Execute( + "channel", "listener", "add", "--address", s.Address(), "-c", ch, + "--callback-url", receiver.URL+"/notify", "--header", "Authorization=Bearer t1", + "-o", "json", + ) + s.NoError(res.Err) + var added workflowservice.RegisterChannelListenerResponse + s.NoError(temporalcli.UnmarshalProtoJSONWithOptions(res.Stdout.Bytes(), &added, true)) + s.NotEmpty(added.GetListenerId()) + + res = s.Execute("channel", "describe", "--address", s.Address(), "-c", ch) + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "Callback listeners: 1") + s.ContainsOnSameLine(res.Stdout.String(), added.GetListenerId(), receiver.URL+"/notify") + + res = s.Execute( + "channel", "notify", "--address", s.Address(), "-c", ch, "--position", "p1", "--counter", "1", + ) + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "Listeners reached", "1") + select { + case d := <-deliveries: + s.Equal("Bearer t1", d.auth) + s.Contains(string(d.body), ch) + case <-time.After(30 * time.Second): + s.Fail("the callback was not called") + } + + res = s.Execute( + "channel", "listener", "remove", "--address", s.Address(), "-c", ch, + "--listener-id", added.GetListenerId(), + ) + s.NoError(res.Err) + res = s.Execute("channel", "describe", "--address", s.Address(), "-c", ch) + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "Callback listeners: 0") +} + +func (s *SharedServerSuite) TestChannel_UnknownChannel() { + res := s.Execute("channel", "describe", "--address", s.Address(), "-c", "never-"+uuid.NewString()) + s.ErrorContains(res.Err, "there is no channel") +} + +func workflowExecution(id, run string) *commonpb.Execution { + return &commonpb.Execution{ + Type: enumspb.EXECUTION_TYPE_WORKFLOW, BusinessId: id, RunId: run, + } +} + +func activityExecution(id, run string) *commonpb.Execution { + return &commonpb.Execution{ + Type: enumspb.EXECUTION_TYPE_ACTIVITY, BusinessId: id, RunId: run, + } +} + +// linkedChannelRequests checks that every call named the owner with the type +// its flag stands for, and the run only where it was given. +func linkedChannelRequests(t *testing.T, f *fakeChannelService, want *commonpb.Execution) { + t.Helper() + f.mu.Lock() + defer f.mu.Unlock() + withRun := []*commonpb.Execution{f.notifies[0].GetExecution()} + withoutRun := []*commonpb.Execution{ + f.describes[0].GetExecution(), f.polls[0].GetExecution(), + f.registers[0].GetExecution(), f.unregister[0].GetExecution(), + } + for _, got := range append(withRun, withoutRun...) { + assert.Equal(t, want.GetType(), got.GetType()) + assert.Equal(t, want.GetBusinessId(), got.GetBusinessId()) + } + assert.Equal(t, want.GetRunId(), withRun[0].GetRunId()) + for _, got := range withoutRun { + assert.Empty(t, got.GetRunId()) + } +} + +func TestChannel_LinkedToWorkflow(t *testing.T) { + f := &fakeChannelService{ + describe: &workflowservice.DescribeChannelResponse{ + Kind: notificationpb.CHANNEL_KIND_LINKED, + LinkedTo: workflowExecution("wf-1", "run-9"), + RetainedCount: 1, + Latest: ¬ificationpb.Notification{Channel: "ch", Position: []byte("p1"), Counter: 1}, + Listeners: []*notificationpb.ChannelListener{{ + ListenerId: "wf-1", + Listener: ¬ificationpb.ChannelListener_Workflow{ + Workflow: ¬ificationpb.WorkflowListener{WorkflowId: "wf-1", RunId: "run-9"}, + }, + }}, + }, + pollAnswers: [][]*notificationpb.Notification{{{Channel: "ch", Counter: 1}}}, + } + addr := startFakeChannelService(t, f) + h := NewCommandHarness(t) + linked := []string{"--address", addr, "-c", "ch", "--workflow-id", "wf-1"} + + res := h.Execute(append([]string{"channel", "notify", "--position", "p1", "--counter", "1", + "--run-id", "run-9"}, linked...)...) + require.NoError(t, res.Err) + h.ContainsOnSameLine(res.Stdout.String(), "Notified channel ch of workflow wf-1") + res = h.Execute(append([]string{"channel", "describe"}, linked...)...) + require.NoError(t, res.Err) + out := res.Stdout.String() + h.ContainsOnSameLine(out, "Kind", "Linked") + h.ContainsOnSameLine(out, "LinkedTo", "workflow wf-1 (run run-9)") + h.ContainsOnSameLine(out, "RetainedCount", "1") + // The owner listens without registering, so its row has no time. + h.ContainsOnSameLine(out, "Workflow listeners: 1") + assert.NotContains(t, out, "ago") + res = h.Execute(append([]string{"channel", "poll"}, linked...)...) + require.NoError(t, res.Err) + res = h.Execute(append([]string{"channel", "listener", "add", "--callback-url", + "http://x"}, linked...)...) + require.NoError(t, res.Err) + h.ContainsOnSameLine(res.Stdout.String(), "to channel ch of workflow wf-1") + res = h.Execute(append([]string{"channel", "listener", "remove", "--listener-id", + "listener-7"}, linked...)...) + require.NoError(t, res.Err) + + linkedChannelRequests(t, f, workflowExecution("wf-1", "run-9")) +} + +func TestChannel_LinkedToActivity(t *testing.T) { + f := &fakeChannelService{ + describe: &workflowservice.DescribeChannelResponse{ + Kind: notificationpb.CHANNEL_KIND_LINKED, + LinkedTo: activityExecution("act-1", ""), + RetainedCount: 2, + Latest: ¬ificationpb.Notification{Channel: "ch", Position: []byte("p2"), Counter: 2}, + }, + pollAnswers: [][]*notificationpb.Notification{{{Channel: "ch", Counter: 2}}}, + } + addr := startFakeChannelService(t, f) + h := NewCommandHarness(t) + linked := []string{"--address", addr, "-c", "ch", "--activity-id", "act-1"} + + res := h.Execute(append([]string{"channel", "notify", "--position", "p2", "--counter", "2", + "--run-id", "run-3"}, linked...)...) + require.NoError(t, res.Err) + h.ContainsOnSameLine(res.Stdout.String(), "Notified channel ch of activity act-1") + res = h.Execute(append([]string{"channel", "describe"}, linked...)...) + require.NoError(t, res.Err) + out := res.Stdout.String() + h.ContainsOnSameLine(out, "Kind", "Linked") + // The owner line carries no run when the Service names none. + h.ContainsOnSameLine(out, "LinkedTo", "activity act-1") + assert.NotContains(t, out, "run ") + h.ContainsOnSameLine(out, "RetainedCount", "2") + res = h.Execute(append([]string{"channel", "poll"}, linked...)...) + require.NoError(t, res.Err) + res = h.Execute(append([]string{"channel", "listener", "add", "--callback-url", + "http://x"}, linked...)...) + require.NoError(t, res.Err) + h.ContainsOnSameLine(res.Stdout.String(), "to channel ch of activity act-1") + res = h.Execute(append([]string{"channel", "listener", "remove", "--listener-id", + "listener-7"}, linked...)...) + require.NoError(t, res.Err) + + linkedChannelRequests(t, f, activityExecution("act-1", "run-3")) + + // The JSON output carries the owner as the Service sent it. + res = h.Execute(append([]string{"channel", "describe", "-o", "json"}, linked...)...) + require.NoError(t, res.Err) + var described workflowservice.DescribeChannelResponse + require.NoError(t, temporalcli.UnmarshalProtoJSONWithOptions(res.Stdout.Bytes(), &described, true)) + assert.Equal(t, enumspb.EXECUTION_TYPE_ACTIVITY, described.GetLinkedTo().GetType()) + assert.Equal(t, "act-1", described.GetLinkedTo().GetBusinessId()) +} + +func TestChannel_IndependentByDefault(t *testing.T) { + f := &fakeChannelService{describe: &workflowservice.DescribeChannelResponse{}} + addr := startFakeChannelService(t, f) + h := NewCommandHarness(t) + + res := h.Execute("channel", "notify", "--address", addr, "-c", "ch", + "--position", "p", "--counter", "1") + require.NoError(t, res.Err) + // A server that predates the linked kind leaves the kind unset. + res = h.Execute("channel", "describe", "--address", addr, "-c", "ch") + require.NoError(t, res.Err) + h.ContainsOnSameLine(res.Stdout.String(), "Kind", "Independent") + assert.NotContains(t, res.Stdout.String(), "LinkedTo") + assert.Nil(t, f.notifies[0].GetExecution()) + assert.Nil(t, f.describes[0].GetExecution()) +} + +func TestChannel_LinkedNotFound(t *testing.T) { + f := &fakeChannelService{err: serviceerror.NewNotFound("workflow execution already completed")} + addr := startFakeChannelService(t, f) + h := NewCommandHarness(t) + + res := h.Execute("channel", "notify", "--address", addr, "-c", "ch", "-w", "gone", + "--position", "p", "--counter", "1") + require.Error(t, res.Err) + assert.Contains(t, res.Err.Error(), `workflow "gone" has no running execution`) + assert.Contains(t, res.Err.Error(), `channel "ch" of workflow "gone"`) + + f.err = serviceerror.NewNotFound("activity execution already completed") + res = h.Execute("channel", "describe", "--address", addr, "-c", "ch", "--activity-id", "gone") + require.Error(t, res.Err) + assert.Contains(t, res.Err.Error(), `activity "gone" has no running execution`) + assert.Contains(t, res.Err.Error(), `channel "ch" of activity "gone"`) +} + +// runningWorkflowDescription is a description with what the text output +// dereferences, for a run that is still open so no close event is fetched. +func runningWorkflowDescription( + subs ...*workflowpb.ChannelSubscriptionInfo, +) *workflowservice.DescribeWorkflowExecutionResponse { + return &workflowservice.DescribeWorkflowExecutionResponse{ + WorkflowExecutionInfo: &workflowpb.WorkflowExecutionInfo{ + Execution: &commonpb.WorkflowExecution{WorkflowId: "wf-1", RunId: "run-1"}, + Type: &commonpb.WorkflowType{Name: "DevWorkflow"}, + Status: enumspb.WORKFLOW_EXECUTION_STATUS_RUNNING, + StartTime: timestamppb.New(time.Date(2026, 10, 2, 9, 0, 0, 0, time.UTC)), + }, + ChannelSubscriptions: subs, + } +} + +func TestChannel_WorkflowDescribeListsChannels(t *testing.T) { + f := &fakeChannelService{describeWorkflow: runningWorkflowDescription( + &workflowpb.ChannelSubscriptionInfo{ + Channel: "demo", + Kind: notificationpb.CHANNEL_KIND_INDEPENDENT, + SubscribedEventId: 5, + LastCounter: 4, + PendingNotification: ¬ificationpb.Notification{Channel: "demo", Counter: 5}, + }, + &workflowpb.ChannelSubscriptionInfo{ + Channel: "stream/scores", + Kind: notificationpb.CHANNEL_KIND_LINKED, + LastCounter: 3, + ScheduledCounter: 0, + ListenerCount: 2, + RetainedCount: 16, + AcceptedCount: 3, + }, + )} + addr := startFakeChannelService(t, f) + h := NewCommandHarness(t) + + res := h.Execute("workflow", "describe", "--address", addr, "-w", "wf-1") + require.NoError(t, res.Err) + out := res.Stdout.String() + h.ContainsOnSameLine(out, "Notification Channels: 2") + h.ContainsOnSameLine(out, "Channel", "Kind", "LastCounter", "PendingCounter", + "ScheduledCounter", "Listeners", "Retained") + // The pending column is blank, not zero, when nothing is pending. + assert.Equal(t, []string{"demo", "Independent", "4", "5", "0", "0", "0"}, + strings.Fields(lineContaining(t, out, "demo"))) + assert.Equal(t, []string{"stream/scores", "Linked", "3", "0", "2", "16"}, + strings.Fields(lineContaining(t, out, "stream/scores"))) + + res = h.Execute("workflow", "describe", "--address", addr, "-w", "wf-1", "-o", "json") + require.NoError(t, res.Err) + var described workflowservice.DescribeWorkflowExecutionResponse + require.NoError(t, temporalcli.UnmarshalProtoJSONWithOptions(res.Stdout.Bytes(), &described, true)) + require.Len(t, described.GetChannelSubscriptions(), 2) + assert.Equal(t, int64(5), described.GetChannelSubscriptions()[0].GetSubscribedEventId()) + assert.Equal(t, int64(5), + described.GetChannelSubscriptions()[0].GetPendingNotification().GetCounter()) + linked := described.GetChannelSubscriptions()[1] + assert.Equal(t, notificationpb.CHANNEL_KIND_LINKED, linked.GetKind()) + assert.Equal(t, int64(3), linked.GetAcceptedCount()) + + // A workflow standing on no channel has no section at all. + f.describeWorkflow = runningWorkflowDescription() + res = h.Execute("workflow", "describe", "--address", addr, "-w", "wf-1") + require.NoError(t, res.Err) + assert.NotContains(t, res.Stdout.String(), "Notification Channels") +} + +func lineContaining(t *testing.T, text, piece string) string { + for _, line := range strings.Split(text, "\n") { + if strings.Contains(line, piece) { + return line + } + } + require.Failf(t, "line not found", "no line contains %q", piece) + return "" +} + +func (s *SharedServerSuite) TestChannel_WorkflowDescribeListsLinked() { + s.Worker().OnDevWorkflow(func(ctx workflow.Context, a any) (any, error) { + workflow.GetSignalChannel(ctx, "finish").Receive(ctx, nil) + return nil, nil + }) + run, err := s.Client.ExecuteWorkflow( + s.Context, + client.StartWorkflowOptions{TaskQueue: s.Worker().Options.TaskQueue}, + DevWorkflow, + "ignored", + ) + s.NoError(err) + describe := func(args ...string) *CommandResult { + return s.Execute(append([]string{"workflow", "describe", "--address", s.Address(), + "-w", run.GetID()}, args...)...) + } + + // A fresh workflow stands on no channel, so the section is absent. + res := describe() + s.NoError(res.Err) + s.NotContains(res.Stdout.String(), "Notification Channels") + + ch := "channel-" + uuid.NewString() + res = s.Execute("channel", "notify", "--address", s.Address(), "-c", ch, + "--workflow-id", run.GetID(), "--position", "p1", "--counter", "1") + s.NoError(res.Err) + + res = describe() + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "Notification Channels: 1") + s.ContainsOnSameLine(res.Stdout.String(), ch, "Linked") + + res = describe("-o", "json") + s.NoError(res.Err) + var described workflowservice.DescribeWorkflowExecutionResponse + s.NoError(temporalcli.UnmarshalProtoJSONWithOptions(res.Stdout.Bytes(), &described, true)) + s.Len(described.GetChannelSubscriptions(), 1) + linked := described.GetChannelSubscriptions()[0] + s.Equal(ch, linked.GetChannel()) + s.Equal(notificationpb.CHANNEL_KIND_LINKED, linked.GetKind()) + s.Equal(int32(1), linked.GetRetainedCount()) + s.Equal(int64(1), linked.GetAcceptedCount()) + + s.NoError(s.Client.SignalWorkflow(s.Context, run.GetID(), "", "finish", nil)) + s.NoError(run.Get(s.Context, nil)) +} + +func (s *SharedServerSuite) TestChannel_LinkedToWorkflow() { + s.Worker().OnDevWorkflow(func(ctx workflow.Context, a any) (any, error) { + workflow.GetSignalChannel(ctx, "finish").Receive(ctx, nil) + return nil, nil + }) + run, err := s.Client.ExecuteWorkflow( + s.Context, + client.StartWorkflowOptions{TaskQueue: s.Worker().Options.TaskQueue}, + DevWorkflow, + "ignored", + ) + s.NoError(err) + ch := "channel-" + uuid.NewString() + linked := []string{"--address", s.Address(), "-c", ch, "--workflow-id", run.GetID()} + channel := func(args ...string) *CommandResult { + return s.Execute(append(append([]string{"channel"}, args...), linked...)...) + } + + // The owner listens by construction, so the first notify reaches it. + res := channel("notify", "--position", "p1", "--counter", "1") + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "Listeners reached", "1") + + res = channel("describe") + s.NoError(res.Err) + out := res.Stdout.String() + s.ContainsOnSameLine(out, "Kind", "Linked") + s.ContainsOnSameLine(out, "LinkedTo", + "workflow "+run.GetID()+" (run "+run.GetRunID()+")") + s.ContainsOnSameLine(out, "RetainedCount", "1") + s.ContainsOnSameLine(out, "LatestPosition", "p1") + s.ContainsOnSameLine(out, "Workflow listeners: 1") + + res = channel("poll", "--after-counter", "0", "--wait", "5s", "-o", "jsonl") + s.NoError(res.Err) + raw := decodeJSONValues(s.T(), res.Stdout.String()) + s.Len(raw, 1) + var n notificationpb.Notification + s.NoError(temporalcli.UnmarshalProtoJSONWithOptions(raw[0], &n, true)) + s.Equal(int64(1), n.GetCounter()) + s.Equal(enumspb.EXECUTION_TYPE_WORKFLOW, n.GetLinkedTo().GetType()) + s.Equal(run.GetID(), n.GetLinkedTo().GetBusinessId()) + + res = channel("listener", "add", "--callback-url", "http://127.0.0.1:1/notify", "-o", "json") + s.NoError(res.Err) + var added workflowservice.RegisterChannelListenerResponse + s.NoError(temporalcli.UnmarshalProtoJSONWithOptions(res.Stdout.Bytes(), &added, true)) + res = channel("describe") + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "Callback listeners: 1") + s.ContainsOnSameLine(res.Stdout.String(), added.GetListenerId(), "http://127.0.0.1:1/notify") + res = channel("listener", "remove", "--listener-id", added.GetListenerId()) + s.NoError(res.Err) + res = channel("describe") + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "Callback listeners: 0") + + // A name nobody has used still exists on a running workflow. + res = s.Execute("channel", "describe", "--address", s.Address(), + "-c", "untouched-"+uuid.NewString(), "--workflow-id", run.GetID()) + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "Kind", "Linked") + s.ContainsOnSameLine(res.Stdout.String(), "RetainedCount", "0") + s.ContainsOnSameLine(res.Stdout.String(), "Workflow listeners: 0") + s.ContainsOnSameLine(res.Stdout.String(), "Callback listeners: 0") + + // The independent channel of the same name is a different channel. + res = s.Execute("channel", "describe", "--address", s.Address(), "-c", ch) + s.ErrorContains(res.Err, "there is no channel") + + res = s.Execute("channel", "describe", "--address", s.Address(), "-c", ch, "--run-id", "r") + s.ErrorContains(res.Err, "--run-id requires --workflow-id or --activity-id") + + s.NoError(s.Client.SignalWorkflow(s.Context, run.GetID(), "", "finish", nil)) + s.NoError(run.Get(s.Context, nil)) +} + +func (s *SharedServerSuite) TestChannel_LinkedToActivity() { + // The activity stays running until the test is done with its channels. + release := make(chan struct{}) + s.Worker().OnDevActivity(func(ctx context.Context, a any) (any, error) { + select { + case <-release: + case <-ctx.Done(): + } + return nil, nil + }) + defer close(release) + activityID := "activity-" + uuid.NewString() + res := s.Execute("activity", "start", "--address", s.Address(), "-o", "json", + "--activity-id", activityID, "--type", "DevActivity", + "--task-queue", s.Worker().Options.TaskQueue, "--start-to-close-timeout", "1m") + s.NoError(res.Err) + var started map[string]any + s.NoError(json.Unmarshal(res.Stdout.Bytes(), &started)) + runID, _ := started["runId"].(string) + s.NotEmpty(runID) + + ch := "channel-" + uuid.NewString() + linked := []string{"--address", s.Address(), "-c", ch, "--activity-id", activityID} + channel := func(args ...string) *CommandResult { + return s.Execute(append(append([]string{"channel"}, args...), linked...)...) + } + + // The activity does not listen on its own channels, so the notify reaches + // nobody and is retained for pollers. + res = channel("notify", "--position", "p1", "--counter", "1") + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "Listeners reached", "0") + + res = channel("describe") + s.NoError(res.Err) + out := res.Stdout.String() + s.ContainsOnSameLine(out, "Kind", "Linked") + s.ContainsOnSameLine(out, "LinkedTo", "activity "+activityID+" (run "+runID+")") + s.ContainsOnSameLine(out, "RetainedCount", "1") + s.ContainsOnSameLine(out, "LatestPosition", "p1") + s.ContainsOnSameLine(out, "Workflow listeners: 0") + + res = channel("poll", "--after-counter", "0", "--wait", "5s", "-o", "jsonl") + s.NoError(res.Err) + raw := decodeJSONValues(s.T(), res.Stdout.String()) + s.Len(raw, 1) + var n notificationpb.Notification + s.NoError(temporalcli.UnmarshalProtoJSONWithOptions(raw[0], &n, true)) + s.Equal(int64(1), n.GetCounter()) + s.Equal(enumspb.EXECUTION_TYPE_ACTIVITY, n.GetLinkedTo().GetType()) + s.Equal(activityID, n.GetLinkedTo().GetBusinessId()) + + res = channel("listener", "add", "--callback-url", "http://127.0.0.1:1/notify", "-o", "json") + s.NoError(res.Err) + var added workflowservice.RegisterChannelListenerResponse + s.NoError(temporalcli.UnmarshalProtoJSONWithOptions(res.Stdout.Bytes(), &added, true)) + res = channel("describe") + s.NoError(res.Err) + s.ContainsOnSameLine(res.Stdout.String(), "Callback listeners: 1") + res = channel("listener", "remove", "--listener-id", added.GetListenerId()) + s.NoError(res.Err) + + // An activity nobody started has no channels to reach. + res = s.Execute("channel", "describe", "--address", s.Address(), "-c", ch, + "--activity-id", "never-"+uuid.NewString()) + s.ErrorContains(res.Err, "has no running execution to reach") +} diff --git a/internal/temporalcli/commands.gen.go b/internal/temporalcli/commands.gen.go index a48152714..ed56d7cc1 100644 --- a/internal/temporalcli/commands.gen.go +++ b/internal/temporalcli/commands.gen.go @@ -97,6 +97,23 @@ func (v *WorkflowReferenceOptions) BuildFlags(f *pflag.FlagSet) { f.StringVarP(&v.RunId, "run-id", "r", "", "Run ID.") } +type ChannelOptions struct { + Channel string + WorkflowId string + ActivityId string + RunId string + FlagSet *pflag.FlagSet +} + +func (v *ChannelOptions) BuildFlags(f *pflag.FlagSet) { + v.FlagSet = f + f.StringVarP(&v.Channel, "channel", "c", "", "Name of the notification channel. Required.") + _ = cobra.MarkFlagRequired(f, "channel") + f.StringVarP(&v.WorkflowId, "workflow-id", "w", "", "Workflow ID of the Workflow Execution the channel is linked to. Without an owner the command uses the independent channel of that name. Cannot use with --activity-id.") + f.StringVar(&v.ActivityId, "activity-id", "", "Activity ID of the standalone Activity the channel is linked to. Cannot use with --workflow-id.") + f.StringVarP(&v.RunId, "run-id", "r", "", "Run ID of the linked execution. Defaults to the current run. Requires --workflow-id or --activity-id.") +} + type DeploymentNameOptions struct { Name string FlagSet *pflag.FlagSet @@ -530,6 +547,7 @@ func NewTemporalCommand(cctx *CommandContext) *TemporalCommand { s.Command.Args = cobra.NoArgs s.Command.AddCommand(&NewTemporalActivityCommand(cctx, &s).Command) s.Command.AddCommand(&NewTemporalBatchCommand(cctx, &s).Command) + s.Command.AddCommand(&NewTemporalChannelCommand(cctx, &s).Command) s.Command.AddCommand(&NewTemporalConfigCommand(cctx, &s).Command) s.Command.AddCommand(&NewTemporalEnvCommand(cctx, &s).Command) s.Command.AddCommand(&NewTemporalNexusCommand(cctx, &s).Command) @@ -1191,6 +1209,213 @@ func NewTemporalBatchTerminateCommand(cctx *CommandContext, parent *TemporalBatc return &s } +type TemporalChannelCommand struct { + Parent *TemporalCommand + Command cobra.Command + cliext.ClientOptions +} + +func NewTemporalChannelCommand(cctx *CommandContext, parent *TemporalCommand) *TemporalChannelCommand { + var s TemporalChannelCommand + s.Parent = parent + s.Command.Use = "channel" + s.Command.Short = "Notify and listen on notification channels" + if hasHighlighting { + s.Command.Long = "A notification channel is a name a writer and its listeners agree on.\nThe writer notifies the channel when a source it writes moves, such as\na Stream gaining records, and never learns who listens. A Workflow\nlistener gets a Workflow Task, a callback listener gets an HTTP call,\nand a client long-polls:\n\n\x1b[1mtemporal channel [command] [options]\x1b[0m\n\nFor example:\n\n\x1b[1mtemporal channel poll \\\n --channel YourChannel \\\n --follow\x1b[0m\n\nA notification tells listeners where the source stands. It carries no\ndata; the listener reads the source itself.\n\nAn independent channel exists on its own, and any number of Workflows,\ncallbacks and clients listen on it. A channel linked to a Workflow\nExecution lives with that Workflow, which listens on it without\nsubscribing. A channel linked to a standalone Activity lives with that\nActivity the same way. Name the owner with \x1b[1m--workflow-id\x1b[0m or\n\x1b[1m--activity-id\x1b[0m on any channel command:\n\n\x1b[1mtemporal channel notify \\\n --channel YourChannel \\\n --workflow-id YourWorkflowId \\\n --position 42 \\\n --counter 42\x1b[0m" + } else { + s.Command.Long = "A notification channel is a name a writer and its listeners agree on.\nThe writer notifies the channel when a source it writes moves, such as\na Stream gaining records, and never learns who listens. A Workflow\nlistener gets a Workflow Task, a callback listener gets an HTTP call,\nand a client long-polls:\n\n```\ntemporal channel [command] [options]\n```\n\nFor example:\n\n```\ntemporal channel poll \\\n --channel YourChannel \\\n --follow\n```\n\nA notification tells listeners where the source stands. It carries no\ndata; the listener reads the source itself.\n\nAn independent channel exists on its own, and any number of Workflows,\ncallbacks and clients listen on it. A channel linked to a Workflow\nExecution lives with that Workflow, which listens on it without\nsubscribing. A channel linked to a standalone Activity lives with that\nActivity the same way. Name the owner with `--workflow-id` or\n`--activity-id` on any channel command:\n\n```\ntemporal channel notify \\\n --channel YourChannel \\\n --workflow-id YourWorkflowId \\\n --position 42 \\\n --counter 42\n```" + } + s.Command.Args = cobra.NoArgs + s.Command.AddCommand(&NewTemporalChannelDescribeCommand(cctx, &s).Command) + s.Command.AddCommand(&NewTemporalChannelListenerCommand(cctx, &s).Command) + s.Command.AddCommand(&NewTemporalChannelNotifyCommand(cctx, &s).Command) + s.Command.AddCommand(&NewTemporalChannelPollCommand(cctx, &s).Command) + s.ClientOptions.BuildFlags(s.Command.PersistentFlags()) + s.ClientOptions.HideFlags() + return &s +} + +type TemporalChannelDescribeCommand struct { + Parent *TemporalChannelCommand + Command cobra.Command + ChannelOptions +} + +func NewTemporalChannelDescribeCommand(cctx *CommandContext, parent *TemporalChannelCommand) *TemporalChannelDescribeCommand { + var s TemporalChannelDescribeCommand + s.Parent = parent + s.Command.DisableFlagsInUseLine = true + s.Command.Use = "describe [flags]" + s.Command.Short = "Show a channel's listeners and latest notification" + if hasHighlighting { + s.Command.Long = "Show who listens on a channel, its latest notification, and how many\nnotifications it keeps for pollers:\n\n\x1b[1mtemporal channel describe \\\n --channel YourChannel\x1b[0m\n\nAn independent channel exists once a writer notifies it or a listener\nregisters on it, and goes away after a while without listeners or\nactivity. A channel linked to a running Workflow Execution or standalone\nActivity always exists, and its description names that owner:\n\n\x1b[1mtemporal channel describe \\\n --channel YourChannel \\\n --workflow-id YourWorkflowId\x1b[0m" + } else { + s.Command.Long = "Show who listens on a channel, its latest notification, and how many\nnotifications it keeps for pollers:\n\n```\ntemporal channel describe \\\n --channel YourChannel\n```\n\nAn independent channel exists once a writer notifies it or a listener\nregisters on it, and goes away after a while without listeners or\nactivity. A channel linked to a running Workflow Execution or standalone\nActivity always exists, and its description names that owner:\n\n```\ntemporal channel describe \\\n --channel YourChannel \\\n --workflow-id YourWorkflowId\n```" + } + s.Command.Args = cobra.NoArgs + s.ChannelOptions.BuildFlags(s.Command.Flags()) + s.Command.Run = func(c *cobra.Command, args []string) { + if err := s.run(cctx, args); err != nil { + cctx.Options.Fail(err) + } + } + return &s +} + +type TemporalChannelListenerCommand struct { + Parent *TemporalChannelCommand + Command cobra.Command +} + +func NewTemporalChannelListenerCommand(cctx *CommandContext, parent *TemporalChannelCommand) *TemporalChannelListenerCommand { + var s TemporalChannelListenerCommand + s.Parent = parent + s.Command.Use = "listener" + s.Command.Short = "Add or remove a channel's callback listeners" + if hasHighlighting { + s.Command.Long = "Register a callback the Service calls with every notification on a\nchannel, or remove one:\n\n\x1b[1mtemporal channel listener [command] [options]\x1b[0m\n\nA Workflow listens on a channel from its own code instead, and stops\nlistening when its run ends." + } else { + s.Command.Long = "Register a callback the Service calls with every notification on a\nchannel, or remove one:\n\n```\ntemporal channel listener [command] [options]\n```\n\nA Workflow listens on a channel from its own code instead, and stops\nlistening when its run ends." + } + s.Command.Args = cobra.NoArgs + s.Command.AddCommand(&NewTemporalChannelListenerAddCommand(cctx, &s).Command) + s.Command.AddCommand(&NewTemporalChannelListenerRemoveCommand(cctx, &s).Command) + return &s +} + +type TemporalChannelListenerAddCommand struct { + Parent *TemporalChannelListenerCommand + Command cobra.Command + ChannelOptions + CallbackUrl string + Header []string +} + +func NewTemporalChannelListenerAddCommand(cctx *CommandContext, parent *TemporalChannelListenerCommand) *TemporalChannelListenerAddCommand { + var s TemporalChannelListenerAddCommand + s.Parent = parent + s.Command.DisableFlagsInUseLine = true + s.Command.Use = "add [flags]" + s.Command.Short = "Register a callback listener" + if hasHighlighting { + s.Command.Long = "Register a callback on a channel. The Service calls the URL with every\nnotification the channel gets from now on:\n\n\x1b[1mtemporal channel listener add \\\n --channel YourChannel \\\n --callback-url https://example.com/notify \\\n --header \"Authorization=Bearer YourToken\"\x1b[0m\n\nThe command prints the listener ID that \x1b[1mtemporal channel listener\nremove\x1b[0m takes.\n\nThe Service only calls addresses its \x1b[1mcallback.allowedAddresses\x1b[0m\ndynamic configuration value allows. The development server allows\n\x1b[1m127.0.0.1\x1b[0m and \x1b[1mlocalhost\x1b[0m on any port over plain HTTP. Another\nService needs the address added, for example:\n\n\x1b[1m--dynamic-config-value \\\n 'callback.allowedAddresses=[{\"Pattern\":\"*.example.com:443\"}]'\x1b[0m" + } else { + s.Command.Long = "Register a callback on a channel. The Service calls the URL with every\nnotification the channel gets from now on:\n\n```\ntemporal channel listener add \\\n --channel YourChannel \\\n --callback-url https://example.com/notify \\\n --header \"Authorization=Bearer YourToken\"\n```\n\nThe command prints the listener ID that `temporal channel listener\nremove` takes.\n\nThe Service only calls addresses its `callback.allowedAddresses`\ndynamic configuration value allows. The development server allows\n`127.0.0.1` and `localhost` on any port over plain HTTP. Another\nService needs the address added, for example:\n\n```\n--dynamic-config-value \\\n 'callback.allowedAddresses=[{\"Pattern\":\"*.example.com:443\"}]'\n```" + } + s.Command.Args = cobra.NoArgs + s.Command.Flags().StringVar(&s.CallbackUrl, "callback-url", "", "URL the Service calls with each notification. Required.") + _ = cobra.MarkFlagRequired(s.Command.Flags(), "callback-url") + s.Command.Flags().StringArrayVar(&s.Header, "header", nil, "Header sent with every call, in `KEY=VALUE` format. Can be passed multiple times.") + s.ChannelOptions.BuildFlags(s.Command.Flags()) + s.Command.Run = func(c *cobra.Command, args []string) { + if err := s.run(cctx, args); err != nil { + cctx.Options.Fail(err) + } + } + return &s +} + +type TemporalChannelListenerRemoveCommand struct { + Parent *TemporalChannelListenerCommand + Command cobra.Command + ChannelOptions + ListenerId string +} + +func NewTemporalChannelListenerRemoveCommand(cctx *CommandContext, parent *TemporalChannelListenerCommand) *TemporalChannelListenerRemoveCommand { + var s TemporalChannelListenerRemoveCommand + s.Parent = parent + s.Command.DisableFlagsInUseLine = true + s.Command.Use = "remove [flags]" + s.Command.Short = "Remove a callback listener" + if hasHighlighting { + s.Command.Long = "Stop calling a callback listener:\n\n\x1b[1mtemporal channel listener remove \\\n --channel YourChannel \\\n --listener-id YourListenerId\x1b[0m\n\nFind the listener ID with \x1b[1mtemporal channel describe\x1b[0m." + } else { + s.Command.Long = "Stop calling a callback listener:\n\n```\ntemporal channel listener remove \\\n --channel YourChannel \\\n --listener-id YourListenerId\n```\n\nFind the listener ID with `temporal channel describe`." + } + s.Command.Args = cobra.NoArgs + s.Command.Flags().StringVar(&s.ListenerId, "listener-id", "", "ID of the listener to remove. Required.") + _ = cobra.MarkFlagRequired(s.Command.Flags(), "listener-id") + s.ChannelOptions.BuildFlags(s.Command.Flags()) + s.Command.Run = func(c *cobra.Command, args []string) { + if err := s.run(cctx, args); err != nil { + cctx.Options.Fail(err) + } + } + return &s +} + +type TemporalChannelNotifyCommand struct { + Parent *TemporalChannelCommand + Command cobra.Command + ChannelOptions + Position string + Counter int + Metadata []string +} + +func NewTemporalChannelNotifyCommand(cctx *CommandContext, parent *TemporalChannelCommand) *TemporalChannelNotifyCommand { + var s TemporalChannelNotifyCommand + s.Parent = parent + s.Command.DisableFlagsInUseLine = true + s.Command.Use = "notify [flags]" + s.Command.Short = "Notify a channel's listeners" + if hasHighlighting { + s.Command.Long = "Tell every listener of a channel that its source moved. The command\nprints how many listeners the notification reached:\n\n\x1b[1mtemporal channel notify \\\n --channel YourChannel \\\n --position 42 \\\n --counter 42 \\\n --metadata 'topic=\"scores\"'\x1b[0m\n\nThe position is where the source stands now, in the writer's own\nterms. The counter orders notifications on the channel: when several\nwait for a listener at once, it gets the one with the highest counter.\nMetadata values are JSON." + } else { + s.Command.Long = "Tell every listener of a channel that its source moved. The command\nprints how many listeners the notification reached:\n\n```\ntemporal channel notify \\\n --channel YourChannel \\\n --position 42 \\\n --counter 42 \\\n --metadata 'topic=\"scores\"'\n```\n\nThe position is where the source stands now, in the writer's own\nterms. The counter orders notifications on the channel: when several\nwait for a listener at once, it gets the one with the highest counter.\nMetadata values are JSON." + } + s.Command.Args = cobra.NoArgs + s.Command.Flags().StringVar(&s.Position, "position", "", "Where the source stands after the write. Sent as text; the Service does not read it. Required.") + _ = cobra.MarkFlagRequired(s.Command.Flags(), "position") + s.Command.Flags().IntVar(&s.Counter, "counter", 0, "Order of this notification among the channel's notifications. Must be greater than zero. Required.") + _ = cobra.MarkFlagRequired(s.Command.Flags(), "counter") + s.Command.Flags().StringArrayVar(&s.Metadata, "metadata", nil, "Detail for listeners, in `KEY=VALUE` format. Values must be JSON. Can be passed multiple times.") + s.ChannelOptions.BuildFlags(s.Command.Flags()) + s.Command.Run = func(c *cobra.Command, args []string) { + if err := s.run(cctx, args); err != nil { + cctx.Options.Fail(err) + } + } + return &s +} + +type TemporalChannelPollCommand struct { + Parent *TemporalChannelCommand + Command cobra.Command + ChannelOptions + AfterCounter int + Wait cliext.FlagDuration + Max int + Follow bool +} + +func NewTemporalChannelPollCommand(cctx *CommandContext, parent *TemporalChannelCommand) *TemporalChannelPollCommand { + var s TemporalChannelPollCommand + s.Parent = parent + s.Command.DisableFlagsInUseLine = true + s.Command.Use = "poll [flags]" + s.Command.Short = "Wait for a channel's notifications" + if hasHighlighting { + s.Command.Long = "Show the notifications a channel keeps, waiting for one when there is\nnone yet:\n\n\x1b[1mtemporal channel poll \\\n --channel YourChannel\x1b[0m\n\nKeep polling and print notifications as they arrive:\n\n\x1b[1mtemporal channel poll \\\n --channel YourChannel \\\n --after-counter 10 \\\n --follow\x1b[0m\n\nEach notification prints with its counter, position and metadata. With\n\x1b[1m--output json\x1b[0m each notification is one JSON object." + } else { + s.Command.Long = "Show the notifications a channel keeps, waiting for one when there is\nnone yet:\n\n```\ntemporal channel poll \\\n --channel YourChannel\n```\n\nKeep polling and print notifications as they arrive:\n\n```\ntemporal channel poll \\\n --channel YourChannel \\\n --after-counter 10 \\\n --follow\n```\n\nEach notification prints with its counter, position and metadata. With\n`--output json` each notification is one JSON object." + } + s.Command.Args = cobra.NoArgs + s.Command.Flags().IntVar(&s.AfterCounter, "after-counter", 0, "Show only notifications with a counter above this one.") + s.Wait = cliext.MustParseFlagDuration("30s") + s.Command.Flags().Var(&s.Wait, "wait", "How long one poll waits for a notification when there is none. The Service may answer sooner.") + s.Command.Flags().IntVar(&s.Max, "max", 0, "Most notifications one poll returns. Default is zero (the Service's limit).") + s.Command.Flags().BoolVarP(&s.Follow, "follow", "f", false, "Keep polling from the last counter seen until interrupted.") + s.ChannelOptions.BuildFlags(s.Command.Flags()) + s.Command.Run = func(c *cobra.Command, args []string) { + if err := s.run(cctx, args); err != nil { + cctx.Options.Fail(err) + } + } + return &s +} + type TemporalConfigCommand struct { Parent *TemporalCommand Command cobra.Command @@ -2848,9 +3073,9 @@ func NewTemporalServerStartDevCommand(cctx *CommandContext, parent *TemporalServ s.Command.Use = "start-dev [flags]" s.Command.Short = "Start Temporal development server" if hasHighlighting { - s.Command.Long = "Run a development Temporal Server on your local system.\n\n\x1b[1m+------------------------------------------------------------------------+\n| WARNING: The development server is not intended for production use. |\n| It skips certain HTTP security checks to make local use simpler. |\n| |\n| For production use, see: |\n| https://docs.temporal.io/production-deployment |\n+------------------------------------------------------------------------+\x1b[0m\n\nView the Web UI for the default configuration at: http://localhost:8233\n\n\x1b[1mtemporal server start-dev\x1b[0m\n\nAdd persistence for Workflow Executions across runs:\n\n\x1b[1mtemporal server start-dev \\\n --db-filename path-to-your-local-persistent-store\x1b[0m\n\nSet the port from the front-end gRPC Service (7233 default):\n\n\x1b[1mtemporal server start-dev \\\n --port 7000\x1b[0m\n\nUse a custom port for the Web UI. The default is the gRPC port (7233 default)\nplus 1000 (8233):\n\n\x1b[1mtemporal server start-dev \\\n --ui-port 3000\x1b[0m" + s.Command.Long = "Run a development Temporal Server on your local system.\n\n\x1b[1m+------------------------------------------------------------------------+\n| WARNING: The development server is not intended for production use. |\n| It skips certain HTTP security checks to make local use simpler. |\n| |\n| For production use, see: |\n| https://docs.temporal.io/production-deployment |\n+------------------------------------------------------------------------+\x1b[0m\n\nView the Web UI for the default configuration at: http://localhost:8233\n\n\x1b[1mtemporal server start-dev\x1b[0m\n\nAdd persistence for Workflow Executions across runs:\n\n\x1b[1mtemporal server start-dev \\\n --db-filename path-to-your-local-persistent-store\x1b[0m\n\nSet the port from the front-end gRPC Service (7233 default):\n\n\x1b[1mtemporal server start-dev \\\n --port 7000\x1b[0m\n\nUse a custom port for the Web UI. The default is the gRPC port (7233 default)\nplus 1000 (8233):\n\n\x1b[1mtemporal server start-dev \\\n --ui-port 3000\x1b[0m\n\nCallbacks to \x1b[1m127.0.0.1\x1b[0m and \x1b[1mlocalhost\x1b[0m on any port are allowed over\nplain HTTP, so a local receiver can listen on a notification channel.\nReplace the list with a dynamic configuration value:\n\n\x1b[1mtemporal server start-dev \\\n --dynamic-config-value 'callback.allowedAddresses=[]'\x1b[0m" } else { - s.Command.Long = "Run a development Temporal Server on your local system.\n\n```\n+------------------------------------------------------------------------+\n| WARNING: The development server is not intended for production use. |\n| It skips certain HTTP security checks to make local use simpler. |\n| |\n| For production use, see: |\n| https://docs.temporal.io/production-deployment |\n+------------------------------------------------------------------------+\n```\n\nView the Web UI for the default configuration at: http://localhost:8233\n\n```\ntemporal server start-dev\n```\n\nAdd persistence for Workflow Executions across runs:\n\n```\ntemporal server start-dev \\\n --db-filename path-to-your-local-persistent-store\n```\n\nSet the port from the front-end gRPC Service (7233 default):\n\n```\ntemporal server start-dev \\\n --port 7000\n```\n\nUse a custom port for the Web UI. The default is the gRPC port (7233 default)\nplus 1000 (8233):\n\n```\ntemporal server start-dev \\\n --ui-port 3000\n```" + s.Command.Long = "Run a development Temporal Server on your local system.\n\n```\n+------------------------------------------------------------------------+\n| WARNING: The development server is not intended for production use. |\n| It skips certain HTTP security checks to make local use simpler. |\n| |\n| For production use, see: |\n| https://docs.temporal.io/production-deployment |\n+------------------------------------------------------------------------+\n```\n\nView the Web UI for the default configuration at: http://localhost:8233\n\n```\ntemporal server start-dev\n```\n\nAdd persistence for Workflow Executions across runs:\n\n```\ntemporal server start-dev \\\n --db-filename path-to-your-local-persistent-store\n```\n\nSet the port from the front-end gRPC Service (7233 default):\n\n```\ntemporal server start-dev \\\n --port 7000\n```\n\nUse a custom port for the Web UI. The default is the gRPC port (7233 default)\nplus 1000 (8233):\n\n```\ntemporal server start-dev \\\n --ui-port 3000\n```\n\nCallbacks to `127.0.0.1` and `localhost` on any port are allowed over\nplain HTTP, so a local receiver can listen on a notification channel.\nReplace the list with a dynamic configuration value:\n\n```\ntemporal server start-dev \\\n --dynamic-config-value 'callback.allowedAddresses=[]'\n```" } s.Command.Args = cobra.NoArgs s.Command.Flags().StringVarP(&s.DbFilename, "db-filename", "f", "", "Path to file for persistent Temporal state store. By default, Workflow Executions are lost when the server process dies.") diff --git a/internal/temporalcli/commands.workflow_view.go b/internal/temporalcli/commands.workflow_view.go index c44711a11..ca767934d 100644 --- a/internal/temporalcli/commands.workflow_view.go +++ b/internal/temporalcli/commands.workflow_view.go @@ -252,6 +252,9 @@ func (c *TemporalWorkflowDescribeCommand) run(cctx *CommandContext, args []strin if err := printCallbacks(cctx, resp.Callbacks); err != nil { return err } + if err := printChannelSubscriptions(cctx, resp.ChannelSubscriptions); err != nil { + return err + } if running { cctx.Printer.Println() diff --git a/internal/temporalcli/commands.yaml b/internal/temporalcli/commands.yaml index 3a94ce74b..18b1020fe 100644 --- a/internal/temporalcli/commands.yaml +++ b/internal/temporalcli/commands.yaml @@ -950,6 +950,244 @@ commands: description: Reason for terminating the batch job. required: true + - name: temporal channel + summary: Notify and listen on notification channels + description: | + A notification channel is a name a writer and its listeners agree on. + The writer notifies the channel when a source it writes moves, such as + a Stream gaining records, and never learns who listens. A Workflow + listener gets a Workflow Task, a callback listener gets an HTTP call, + and a client long-polls: + + ``` + temporal channel [command] [options] + ``` + + For example: + + ``` + temporal channel poll \ + --channel YourChannel \ + --follow + ``` + + A notification tells listeners where the source stands. It carries no + data; the listener reads the source itself. + + An independent channel exists on its own, and any number of Workflows, + callbacks and clients listen on it. A channel linked to a Workflow + Execution lives with that Workflow, which listens on it without + subscribing. A channel linked to a standalone Activity lives with that + Activity the same way. Name the owner with `--workflow-id` or + `--activity-id` on any channel command: + + ``` + temporal channel notify \ + --channel YourChannel \ + --workflow-id YourWorkflowId \ + --position 42 \ + --counter 42 + ``` + option-sets: + - client + docs: + description-header: >- + Temporal Channel commands notify a notification channel, poll it, + describe it, and add or remove its callback listeners. + keywords: + - channel + - channel describe + - channel listener add + - channel listener remove + - channel notify + - channel poll + - cli reference + - command-line-interface-cli + - notification channel + - temporal cli + tags: + - Temporal CLI + - Channels + + - name: temporal channel describe + summary: Show a channel's listeners and latest notification + description: | + Show who listens on a channel, its latest notification, and how many + notifications it keeps for pollers: + + ``` + temporal channel describe \ + --channel YourChannel + ``` + + An independent channel exists once a writer notifies it or a listener + registers on it, and goes away after a while without listeners or + activity. A channel linked to a running Workflow Execution or standalone + Activity always exists, and its description names that owner: + + ``` + temporal channel describe \ + --channel YourChannel \ + --workflow-id YourWorkflowId + ``` + option-sets: + - channel + + - name: temporal channel listener + summary: Add or remove a channel's callback listeners + description: | + Register a callback the Service calls with every notification on a + channel, or remove one: + + ``` + temporal channel listener [command] [options] + ``` + + A Workflow listens on a channel from its own code instead, and stops + listening when its run ends. + + - name: temporal channel listener add + summary: Register a callback listener + description: | + Register a callback on a channel. The Service calls the URL with every + notification the channel gets from now on: + + ``` + temporal channel listener add \ + --channel YourChannel \ + --callback-url https://example.com/notify \ + --header "Authorization=Bearer YourToken" + ``` + + The command prints the listener ID that `temporal channel listener + remove` takes. + + The Service only calls addresses its `callback.allowedAddresses` + dynamic configuration value allows. The development server allows + `127.0.0.1` and `localhost` on any port over plain HTTP. Another + Service needs the address added, for example: + + ``` + --dynamic-config-value \ + 'callback.allowedAddresses=[{"Pattern":"*.example.com:443"}]' + ``` + options: + - name: callback-url + type: string + description: URL the Service calls with each notification. + required: true + - name: header + type: string[] + description: | + Header sent with every call, in `KEY=VALUE` format. + Can be passed multiple times. + option-sets: + - channel + + - name: temporal channel listener remove + summary: Remove a callback listener + description: | + Stop calling a callback listener: + + ``` + temporal channel listener remove \ + --channel YourChannel \ + --listener-id YourListenerId + ``` + + Find the listener ID with `temporal channel describe`. + options: + - name: listener-id + type: string + description: ID of the listener to remove. + required: true + option-sets: + - channel + + - name: temporal channel notify + summary: Notify a channel's listeners + description: | + Tell every listener of a channel that its source moved. The command + prints how many listeners the notification reached: + + ``` + temporal channel notify \ + --channel YourChannel \ + --position 42 \ + --counter 42 \ + --metadata 'topic="scores"' + ``` + + The position is where the source stands now, in the writer's own + terms. The counter orders notifications on the channel: when several + wait for a listener at once, it gets the one with the highest counter. + Metadata values are JSON. + options: + - name: position + type: string + description: | + Where the source stands after the write. + Sent as text; the Service does not read it. + required: true + - name: counter + type: int + description: | + Order of this notification among the channel's notifications. + Must be greater than zero. + required: true + - name: metadata + type: string[] + description: | + Detail for listeners, in `KEY=VALUE` format. + Values must be JSON. + Can be passed multiple times. + option-sets: + - channel + + - name: temporal channel poll + summary: Wait for a channel's notifications + description: | + Show the notifications a channel keeps, waiting for one when there is + none yet: + + ``` + temporal channel poll \ + --channel YourChannel + ``` + + Keep polling and print notifications as they arrive: + + ``` + temporal channel poll \ + --channel YourChannel \ + --after-counter 10 \ + --follow + ``` + + Each notification prints with its counter, position and metadata. With + `--output json` each notification is one JSON object. + options: + - name: after-counter + type: int + description: Show only notifications with a counter above this one. + - name: wait + type: duration + description: | + How long one poll waits for a notification when there is none. + The Service may answer sooner. + default: 30s + - name: max + type: int + description: | + Most notifications one poll returns. + Default is zero (the Service's limit). + - name: follow + short: f + type: bool + description: Keep polling from the last counter seen until interrupted. + option-sets: + - channel + - name: temporal config summary: Manage config files (EXPERIMENTAL) description: | @@ -3380,6 +3618,15 @@ commands: temporal server start-dev \ --ui-port 3000 ``` + + Callbacks to `127.0.0.1` and `localhost` on any port are allowed over + plain HTTP, so a local receiver can listen on a notification channel. + Replace the list with a dynamic configuration value: + + ``` + temporal server start-dev \ + --dynamic-config-value 'callback.allowedAddresses=[]' + ``` options: - name: db-filename short: f @@ -5392,6 +5639,32 @@ option-sets: short: r description: Run ID. + - name: channel + options: + - name: channel + type: string + short: c + description: Name of the notification channel. + required: true + - name: workflow-id + type: string + short: w + description: | + Workflow ID of the Workflow Execution the channel is linked to. + Without an owner the command uses the independent channel of that name. + Cannot use with --activity-id. + - name: activity-id + type: string + description: | + Activity ID of the standalone Activity the channel is linked to. + Cannot use with --workflow-id. + - name: run-id + type: string + short: r + description: | + Run ID of the linked execution. + Defaults to the current run. Requires --workflow-id or --activity-id. + - name: deployment-name options: - name: name diff --git a/internal/temporalcli/payload.go b/internal/temporalcli/payload.go index a2f3bc9d8..5875ab458 100644 --- a/internal/temporalcli/payload.go +++ b/internal/temporalcli/payload.go @@ -7,6 +7,7 @@ import ( "strings" "go.temporal.io/api/common/v1" + "go.temporal.io/api/temporalproto" ) // CreatePayloads creates API Payload objects from given data and metadata slices. @@ -40,3 +41,18 @@ func CreatePayloads(data [][]byte, metadata map[string][][]byte, isBase64 bool) } return ret, nil } + +// payloadText renders a payload for a table cell the way the printer renders +// payloads elsewhere: shorthand, on one line. +func payloadText(p *common.Payload) (string, error) { + if p == nil { + return "", nil + } + b, err := temporalproto.CustomJSONMarshalOptions{ + Metadata: map[string]any{common.EnablePayloadShorthandMetadataKey: true}, + }.Marshal(p) + if err != nil { + return "", fmt.Errorf("failed rendering payload: %w", err) + } + return string(b), nil +} diff --git a/internal/temporalcli/refusal.go b/internal/temporalcli/refusal.go new file mode 100644 index 000000000..a086a374b --- /dev/null +++ b/internal/temporalcli/refusal.go @@ -0,0 +1,18 @@ +package temporalcli + +// refusalError is a server refusal read into plain language, so the message +// says what to do next instead of repeating the server's wording. +type refusalError struct { + plain string + detail string + cause error +} + +func (e *refusalError) Error() string { + if e.detail == "" { + return e.plain + } + return e.plain + " The server said: " + e.detail +} + +func (e *refusalError) Unwrap() error { return e.cause }