From f443fc9b229a737d7ee9bdf466a1ea95366a81a3 Mon Sep 17 00:00:00 2001 From: Sam Mathis Date: Fri, 2 Oct 2026 11:12:45 -0700 Subject: [PATCH 1/2] Includes input/output for nexus op describe JSON --- internal/temporalcli/commands.link_test.go | 11 +- .../temporalcli/commands.nexus_operation.go | 53 +++++---- .../commands.nexus_operation_test.go | 108 ++++++++++++++++-- 3 files changed, 131 insertions(+), 41 deletions(-) diff --git a/internal/temporalcli/commands.link_test.go b/internal/temporalcli/commands.link_test.go index d0dc74d3c..f14983948 100644 --- a/internal/temporalcli/commands.link_test.go +++ b/internal/temporalcli/commands.link_test.go @@ -12,7 +12,6 @@ import ( enumspb "go.temporal.io/api/enums/v1" nexuspb "go.temporal.io/api/nexus/v1" "go.temporal.io/api/workflowservice/v1" - "go.temporal.io/sdk/client" ) func workflowEventLink() *commonpb.Link { @@ -161,13 +160,13 @@ func TestPrintNexusOperationDescription_Links(t *testing.T) { var buf bytes.Buffer cctx := &CommandContext{Printer: &printer.Printer{Output: &buf}} - desc := &client.NexusOperationExecutionDescription{ - RawInfo: &nexuspb.NexusOperationExecutionInfo{ - Links: []*commonpb.Link{workflowEventLink()}, + desc := &workflowservice.DescribeNexusOperationExecutionResponse{ + Info: &nexuspb.NexusOperationExecutionInfo{ + OperationId: "op-1", + RunId: "run-1", + Links: []*commonpb.Link{workflowEventLink()}, }, } - desc.OperationID = "op-1" - desc.OperationRunID = "run-1" require.NoError(t, printNexusOperationDescription(cctx, desc)) out := buf.String() diff --git a/internal/temporalcli/commands.nexus_operation.go b/internal/temporalcli/commands.nexus_operation.go index c97a39703..0445f61a5 100644 --- a/internal/temporalcli/commands.nexus_operation.go +++ b/internal/temporalcli/commands.nexus_operation.go @@ -251,24 +251,29 @@ func (c *TemporalNexusOperationDescribeCommand) run(cctx *CommandContext, args [ } defer cl.Close() - handle := cl.GetNexusOperationHandle(client.GetNexusOperationHandleOptions{ - OperationID: c.OperationId, - RunID: c.RunId, + desc, err := cl.WorkflowService().DescribeNexusOperationExecution(cctx, &workflowservice.DescribeNexusOperationExecutionRequest{ + Namespace: c.Parent.Parent.Namespace, + OperationId: c.OperationId, + RunId: c.RunId, + IncludeInput: cctx.JSONOutput, + IncludeOutcome: cctx.JSONOutput, }) - - desc, err := handle.Describe(cctx, client.DescribeNexusOperationOptions{}) if err != nil { return fmt.Errorf("failed describing nexus operation: %w", err) } if c.Raw || cctx.JSONOutput { - return cctx.Printer.PrintStructured(desc.RawInfo, printer.StructuredOptions{}) + return cctx.Printer.PrintStructured(desc, printer.StructuredOptions{}) } return printNexusOperationDescription(cctx, desc) } -func printNexusOperationDescription(cctx *CommandContext, desc *client.NexusOperationExecutionDescription) error { - summary, _ := desc.GetSummary() +func printNexusOperationDescription(cctx *CommandContext, desc *workflowservice.DescribeNexusOperationExecutionResponse) error { + info := desc.GetInfo() + var summary string + if payload := info.GetUserMetadata().GetSummary(); payload != nil { + _ = DataConverterWithRawValue.FromPayload(payload, &summary) + } d := struct { OperationId string RunId string @@ -287,27 +292,27 @@ func printNexusOperationDescription(cctx *CommandContext, desc *client.NexusOper Identity string `cli:",cardOmitEmpty"` Summary string `cli:",cardOmitEmpty"` }{ - OperationId: desc.OperationID, - RunId: desc.OperationRunID, - Endpoint: desc.Endpoint, - Service: desc.Service, - Operation: desc.Operation, - Status: desc.Status.String(), - State: desc.State.String(), - Attempt: desc.Attempt, - ScheduleToCloseTimeout: desc.ScheduleToCloseTimeout, - ScheduledTime: desc.ScheduledTime, - CloseTime: desc.CloseTime, - ExpirationTime: desc.ExpirationTime, - BlockedReason: desc.BlockedReason, - OperationToken: desc.OperationToken, - Identity: desc.Identity, + OperationId: info.GetOperationId(), + RunId: info.GetRunId(), + Endpoint: info.GetEndpoint(), + Service: info.GetService(), + Operation: info.GetOperation(), + Status: info.GetStatus().String(), + State: info.GetState().String(), + Attempt: info.GetAttempt(), + ScheduleToCloseTimeout: info.GetScheduleToCloseTimeout().AsDuration(), + ScheduledTime: info.GetScheduleTime().AsTime(), + CloseTime: info.GetCloseTime().AsTime(), + ExpirationTime: info.GetExpirationTime().AsTime(), + BlockedReason: info.GetBlockedReason(), + OperationToken: info.GetOperationToken(), + Identity: info.GetIdentity(), Summary: summary, } if err := cctx.Printer.PrintStructured(d, printer.StructuredOptions{}); err != nil { return err } - return printLinks(cctx, desc.RawInfo.GetLinks()) + return printLinks(cctx, info.GetLinks()) } func (c *TemporalNexusOperationCancelCommand) run(cctx *CommandContext, args []string) error { diff --git a/internal/temporalcli/commands.nexus_operation_test.go b/internal/temporalcli/commands.nexus_operation_test.go index 71fc7e2db..e8aa7bc93 100644 --- a/internal/temporalcli/commands.nexus_operation_test.go +++ b/internal/temporalcli/commands.nexus_operation_test.go @@ -6,6 +6,7 @@ import ( "errors" "fmt" "strings" + "sync" "testing" "time" @@ -15,9 +16,11 @@ import ( nexuspb "go.temporal.io/api/nexus/v1" "go.temporal.io/api/operatorservice/v1" "go.temporal.io/api/serviceerror" + "go.temporal.io/api/workflowservice/v1" "go.temporal.io/sdk/client" "go.temporal.io/sdk/temporalnexus" "go.temporal.io/sdk/workflow" + "google.golang.org/grpc" ) func (s *SharedServerSuite) setupNexusEndpointAndWorker(t *testing.T) (string, *DevWorker) { @@ -186,18 +189,101 @@ func (s *SharedServerSuite) TestNexusOperationDescribe_JSON() { ) s.NoError(res.Err) - s.Eventually(func() bool { - res = s.Execute( - "nexus", "operation", "describe", - "--address", s.Address(), - "--operation-id", opID, - "--output", "json", - ) - return res.Err == nil - }, 30*time.Second, 500*time.Millisecond) + handle := s.Client.GetNexusOperationHandle(client.GetNexusOperationHandleOptions{OperationID: opID}) + var result string + require.NoError(s.T(), handle.Get(s.Context, &result)) + s.Equal("got: hello", result) + + var requestLock sync.Mutex + var describeRequest *workflowservice.DescribeNexusOperationExecutionRequest + s.CommandHarness.Options.AdditionalClientGRPCDialOptions = append( + s.CommandHarness.Options.AdditionalClientGRPCDialOptions, + grpc.WithChainUnaryInterceptor(func( + ctx context.Context, + method string, req, reply any, + cc *grpc.ClientConn, invoker grpc.UnaryInvoker, opts ...grpc.CallOption, + ) error { + if r, ok := req.(*workflowservice.DescribeNexusOperationExecutionRequest); ok { + requestLock.Lock() + describeRequest = r + requestLock.Unlock() + } + return invoker(ctx, method, req, reply, cc, opts...) + }), + ) - s.NoError(res.Err) - s.Contains(res.Stdout.String(), opID) + for _, format := range []string{"json", "jsonl", "text", "raw"} { + s.T().Run(format, func(t *testing.T) { + args := []string{ + "nexus", "operation", "describe", + "--address", s.Address(), + "--operation-id", opID, + } + if format == "raw" { + args = append(args, "--raw") + } else { + args = append(args, "--output", format) + } + res := s.Execute(args...) + require.NoError(t, res.Err) + requestLock.Lock() + req := describeRequest + describeRequest = nil + requestLock.Unlock() + require.NotNil(t, req) + includePayloads := format == "json" || format == "jsonl" + require.Equal(t, includePayloads, req.GetIncludeInput()) + require.Equal(t, includePayloads, req.GetIncludeOutcome()) + if includePayloads { + var out map[string]any + require.NoError(t, json.Unmarshal(res.Stdout.Bytes(), &out)) + info, ok := out["info"].(map[string]any) + require.True(t, ok) + require.Equal(t, opID, info["operationId"]) + require.Equal(t, "hello", out["input"]) + require.Equal(t, "got: hello", out["result"]) + } else { + require.Contains(t, res.Stdout.String(), opID) + } + }) + } +} + +func (s *SharedServerSuite) TestNexusOperationDescribe_JSONFailure() { + endpointName, w := s.setupNexusEndpointWithWorkflow(s.T(), func(ctx workflow.Context, input string) (string, error) { + return "", errors.New("describe failure") + }, "describe-failure-") + defer w.Stop() + + opID := "desc-failed-op-" + uuid.NewString()[:8] + res := s.Execute( + "nexus", "operation", "start", + "--address", s.Address(), + "--endpoint", endpointName, + "--service", "test-service", + "--operation", "test-op", + "--operation-id", opID, + "--input", `"hello"`, + ) + require.NoError(s.T(), res.Err) + + handle := s.Client.GetNexusOperationHandle(client.GetNexusOperationHandleOptions{OperationID: opID}) + var result string + require.ErrorContains(s.T(), handle.Get(s.Context, &result), "describe failure") + + res = s.Execute( + "nexus", "operation", "describe", + "--address", s.Address(), + "--operation-id", opID, + "--output", "json", + ) + require.NoError(s.T(), res.Err) + var out map[string]any + require.NoError(s.T(), json.Unmarshal(res.Stdout.Bytes(), &out)) + require.Equal(s.T(), "hello", out["input"]) + require.NotNil(s.T(), out["failure"]) + require.NotContains(s.T(), out, "result") + require.Contains(s.T(), res.Stdout.String(), "describe failure") } func (s *SharedServerSuite) TestNexusOperationCancel() { From a3a7e4c82047136acd965dec3ddadcfb31819f08 Mon Sep 17 00:00:00 2001 From: Sam Mathis Date: Fri, 2 Oct 2026 14:15:04 -0700 Subject: [PATCH 2/2] uses timestampToTime for SANO timestamps --- internal/temporalcli/commands.link_test.go | 3 +++ internal/temporalcli/commands.nexus_operation.go | 6 +++--- 2 files changed, 6 insertions(+), 3 deletions(-) diff --git a/internal/temporalcli/commands.link_test.go b/internal/temporalcli/commands.link_test.go index f14983948..2c0fa738f 100644 --- a/internal/temporalcli/commands.link_test.go +++ b/internal/temporalcli/commands.link_test.go @@ -171,6 +171,9 @@ func TestPrintNexusOperationDescription_Links(t *testing.T) { require.NoError(t, printNexusOperationDescription(cctx, desc)) out := buf.String() require.Contains(t, out, "op-1") + require.NotContains(t, out, "ScheduledTime") + require.NotContains(t, out, "CloseTime") + require.NotContains(t, out, "ExpirationTime") require.Contains(t, out, "Links: 1") require.Contains(t, out, "temporal:///namespaces/ns/workflows/wf-id/run-id/history") } diff --git a/internal/temporalcli/commands.nexus_operation.go b/internal/temporalcli/commands.nexus_operation.go index 0445f61a5..46608f80a 100644 --- a/internal/temporalcli/commands.nexus_operation.go +++ b/internal/temporalcli/commands.nexus_operation.go @@ -301,9 +301,9 @@ func printNexusOperationDescription(cctx *CommandContext, desc *workflowservice. State: info.GetState().String(), Attempt: info.GetAttempt(), ScheduleToCloseTimeout: info.GetScheduleToCloseTimeout().AsDuration(), - ScheduledTime: info.GetScheduleTime().AsTime(), - CloseTime: info.GetCloseTime().AsTime(), - ExpirationTime: info.GetExpirationTime().AsTime(), + ScheduledTime: timestampToTime(info.GetScheduleTime()), + CloseTime: timestampToTime(info.GetCloseTime()), + ExpirationTime: timestampToTime(info.GetExpirationTime()), BlockedReason: info.GetBlockedReason(), OperationToken: info.GetOperationToken(), Identity: info.GetIdentity(),