From 0645eae63630d2ec468b11438cbeeceaacb70dec Mon Sep 17 00:00:00 2001 From: chaptersix <13949480+chaptersix@users.noreply.github.com> Date: Sat, 12 Sep 2026 20:08:37 -0500 Subject: [PATCH] Avoid duplicate codec calls for System Nexus history --- cliext/client.go | 9 +- cliext/go.mod | 9 +- cliext/go.sum | 24 +- go.mod | 8 +- go.sum | 16 +- internal/temporalcli/client.go | 19 +- internal/temporalcli/commands.system_nexus.go | 54 --- .../temporalcli/commands.system_nexus_test.go | 381 ++++++------------ .../temporalcli/commands.workflow_exec.go | 87 ++-- .../temporalcli/commands.workflow_view.go | 3 +- .../commands.workflow_view_test.go | 69 +++- internal/temporalcli/commands_test.go | 2 +- 12 files changed, 254 insertions(+), 427 deletions(-) delete mode 100644 internal/temporalcli/commands.system_nexus.go diff --git a/cliext/client.go b/cliext/client.go index 48eaa7508..b9f5a1954 100644 --- a/cliext/client.go +++ b/cliext/client.go @@ -29,11 +29,6 @@ type ClientOptionsBuilder struct { // Logger is the slog logger to use for the client. If set, it will be // wrapped with the SDK's structured logger adapter. Logger *slog.Logger - - // PayloadCodec is populated by Build when a remote payload codec is - // configured. Callers can use it to decode payloads outside the gRPC - // interceptor chain (e.g. payloads nested inside opaque proto bytes). - PayloadCodec converter.PayloadCodec } type oauthCredentials struct { @@ -269,7 +264,6 @@ func (b *ClientOptionsBuilder) Build(ctx context.Context) (client.Options, error } clientOpts.ConnectionOptions.DialOptions = append( clientOpts.ConnectionOptions.DialOptions, grpc.WithChainUnaryInterceptor(interceptor)) - b.PayloadCodec = payloadCodec } // Set connect timeout for GetSystemInfo if provided. @@ -294,8 +288,7 @@ func parseKeyValuePairs(pairs []string) (map[string]string, error) { } // newRemotePayloadCodec constructs a remote payload codec from the configured endpoint, -// auth, and headers. The returned codec can be used both inside a gRPC interceptor and -// to decode payloads nested inside opaque proto bytes (e.g. system Nexus operation inputs). +// auth, and headers. func newRemotePayloadCodec( namespace string, codecEndpoint string, diff --git a/cliext/go.mod b/cliext/go.mod index 2853e64ed..73c0c1422 100644 --- a/cliext/go.mod +++ b/cliext/go.mod @@ -10,7 +10,7 @@ require ( go.temporal.io/sdk v1.47.0 go.temporal.io/sdk/contrib/envconfig v1.0.2 golang.org/x/oauth2 v0.36.0 - google.golang.org/grpc v1.82.1 + google.golang.org/grpc v1.83.1 ) require ( @@ -27,15 +27,14 @@ require ( github.com/robfig/cron v1.2.0 // indirect github.com/rogpeppe/go-internal v1.14.1 // indirect github.com/stretchr/objx v0.5.3 // indirect - go.opentelemetry.io/otel v1.44.0 // indirect - go.temporal.io/api v1.63.5 // indirect + go.temporal.io/api v1.63.6-0.20260910200743-859a8e8d17c0 // indirect golang.org/x/net v0.57.0 // indirect golang.org/x/sync v0.22.0 // indirect golang.org/x/sys v0.47.0 // indirect golang.org/x/text v0.40.0 // indirect golang.org/x/time v0.15.0 // indirect - google.golang.org/genproto/googleapis/api v0.0.0-20260420184626-e10c466a9529 // indirect - google.golang.org/genproto/googleapis/rpc v0.0.0-20260420184626-e10c466a9529 // indirect + google.golang.org/genproto/googleapis/api v0.0.0-20260526163538-3dc84a4a5aaa // indirect + google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect google.golang.org/protobuf v1.36.11 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/cliext/go.sum b/cliext/go.sum index bcea6cd8e..edfe5e086 100644 --- a/cliext/go.sum +++ b/cliext/go.sum @@ -57,14 +57,14 @@ go.opentelemetry.io/otel v1.44.0 h1:JjwHmHpA4iZ3wBxluu2fbbE7j4kqlE8jXyAyPXH7HqU= go.opentelemetry.io/otel v1.44.0/go.mod h1:BMgjTHL9WPRlRjL2oZCBTL4whCGtXch2H4BhOPIAyYc= go.opentelemetry.io/otel/metric v1.44.0 h1:1w0gILTcHdr3YI+ixLyjemwrVnsMURbTZFrSYCdDdmc= go.opentelemetry.io/otel/metric v1.44.0/go.mod h1:8O7hanEPBNgEMmybD3s2VBKcgWOCsA6tzHBPODAiquo= -go.opentelemetry.io/otel/sdk v1.43.0 h1:pi5mE86i5rTeLXqoF/hhiBtUNcrAGHLKQdhg4h4V9Dg= -go.opentelemetry.io/otel/sdk v1.43.0/go.mod h1:P+IkVU3iWukmiit/Yf9AWvpyRDlUeBaRg6Y+C58QHzg= -go.opentelemetry.io/otel/sdk/metric v1.43.0 h1:S88dyqXjJkuBNLeMcVPRFXpRw2fuwdvfCGLEo89fDkw= -go.opentelemetry.io/otel/sdk/metric v1.43.0/go.mod h1:C/RJtwSEJ5hzTiUz5pXF1kILHStzb9zFlIEe85bhj6A= +go.opentelemetry.io/otel/sdk v1.44.0 h1:nHYwb9lK+fJPU/dnT6s7W7Z8itMWyqrnVfbheVYrZ58= +go.opentelemetry.io/otel/sdk v1.44.0/go.mod h1:Osuydd3Se74nqjAKxid74N5eC+jfEqfTegHRnq58oK0= +go.opentelemetry.io/otel/sdk/metric v1.44.0 h1:3LlKgI+VjbVsjNRFZJZAJ30WjXC5VkNRks6si09iEfI= +go.opentelemetry.io/otel/sdk/metric v1.44.0/go.mod h1:5B5pMARnXxKhltooO4xUuCBorl65a4EpnTalObqOigA= go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/6gtIk= go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE= -go.temporal.io/api v1.63.5 h1:c11+kPYHkXXL3UiShPdbMD+xtvqGsbTibUA9ypmiCa4= -go.temporal.io/api v1.63.5/go.mod h1:SrlW2JMwVlDP4nRWSNznUFqnSHd+YeMDS1BkYo63HCQ= +go.temporal.io/api v1.63.6-0.20260910200743-859a8e8d17c0 h1:S765rdHH0PFi7zFEmR1ThDteFKUnJw+nqY7ryA6b5UI= +go.temporal.io/api v1.63.6-0.20260910200743-859a8e8d17c0/go.mod h1:acM0I9WPuYg8W3Pd9jOZvEgi7mRUttUQ4+e7fowKVnM= go.temporal.io/sdk v1.47.0 h1:lZ39w1+uWSjHTL0F3mSc0t4XUnKX8CCcWxqSiLbaHnc= go.temporal.io/sdk v1.47.0/go.mod h1:ilKs0twgP4JpP8pfhIgZumnOEyBiYn6ZO/ta//NnKMU= go.temporal.io/sdk/contrib/envconfig v1.0.2 h1:MGHfsuPUtsf7X9M6WYn3zYJj/mWsuYHnA1uuiL0KEuE= @@ -118,12 +118,12 @@ golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8T golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E= -google.golang.org/genproto/googleapis/api v0.0.0-20260420184626-e10c466a9529 h1:zUWMZsvo/IJcD1t6MNCPO/azZTwz0TvwCBqr5aifoVY= -google.golang.org/genproto/googleapis/api v0.0.0-20260420184626-e10c466a9529/go.mod h1:a5OGAgyRr4lqco7AG9hQM9Fwh0N2ZV4grR0eXFEsXQg= -google.golang.org/genproto/googleapis/rpc v0.0.0-20260420184626-e10c466a9529 h1:XF8+t6QQiS0o9ArVan/HW8Q7cycNPGsJf6GA2nXxYAg= -google.golang.org/genproto/googleapis/rpc v0.0.0-20260420184626-e10c466a9529/go.mod h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8= -google.golang.org/grpc v1.82.1 h1:NnAxzGRA0677vCa4BUkOAnO5+FfQqVl9iUXeD0IqcGE= -google.golang.org/grpc v1.82.1/go.mod h1:yzTZ1TB1Z3SG+LIYaI+WiE8D5+PZ3ArnrSp8zF3+/ZA= +google.golang.org/genproto/googleapis/api v0.0.0-20260526163538-3dc84a4a5aaa h1:Kjn0N0tCrDgiAFW+lGO4JZ3ck44CehvJQMAwj9QF0G8= +google.golang.org/genproto/googleapis/api v0.0.0-20260526163538-3dc84a4a5aaa/go.mod h1:q4lMZS6kskjT5HvCPrnnypcDPVJqT/f4nfxmkE7gryY= +google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa h1:mZHHdPZl0dbGHCflZgAq/Q468DWVFcU2whhB2KAo8fk= +google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa/go.mod h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8= +google.golang.org/grpc v1.83.1 h1:HIO0+BEtBP6soyqvqC8sNUjZ7bTs+0hFQuFF+RAy++Y= +google.golang.org/grpc v1.83.1/go.mod h1:kDyl6SKsiHKt0uylY5gtn5cEjkrIOhQOGDgIc4JGwzQ= google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= diff --git a/go.mod b/go.mod index fd0ed6586..7b8ceec3d 100644 --- a/go.mod +++ b/go.mod @@ -10,17 +10,17 @@ require ( github.com/fatih/color v1.19.0 github.com/google/uuid v1.6.0 github.com/mattn/go-isatty v0.0.23 - github.com/nexus-rpc/sdk-go v0.6.0 + github.com/nexus-rpc/sdk-go v0.7.0 github.com/olekukonko/tablewriter v0.0.5 github.com/spf13/cobra v1.10.2 github.com/spf13/pflag v1.0.10 github.com/stretchr/testify v1.12.0 github.com/temporalio/cli/cliext v0.0.0 github.com/temporalio/ui-server/v2 v2.54.1 - go.temporal.io/api v1.63.5 - go.temporal.io/sdk v1.47.0 + go.temporal.io/api v1.63.6-0.20260910200743-859a8e8d17c0 + go.temporal.io/sdk v1.48.0 go.temporal.io/sdk/contrib/envconfig v1.0.2 - go.temporal.io/server v1.32.0 + go.temporal.io/server v1.32.0-163.3 golang.org/x/exp v0.0.0-20260611194520-c48552f49976 golang.org/x/mod v0.40.0 golang.org/x/term v0.45.0 diff --git a/go.sum b/go.sum index 91db7848b..ba853c130 100644 --- a/go.sum +++ b/go.sum @@ -340,8 +340,8 @@ github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOF github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls= github.com/nexus-rpc/nexus-proto-annotations v0.1.0 h1:2fELd+9sqUtNu6Fg//pw8YFsxOvp8vZ8hfP0nHhNI80= github.com/nexus-rpc/nexus-proto-annotations v0.1.0/go.mod h1:n3UjF1bPCW8llR8tHvbxJ+27yPWrhpo8w/Yg1IOuY0Y= -github.com/nexus-rpc/sdk-go v0.6.0 h1:QRgnP2zTbxEbiyWG/aXH8uSC5LV/Mg1fqb19jb4DBlo= -github.com/nexus-rpc/sdk-go v0.6.0/go.mod h1:FHdPfVQwRuJFZFTF0Y2GOAxCrbIBNrcPna9slkGKPYk= +github.com/nexus-rpc/sdk-go v0.7.0 h1:38NrfY5rLnZAiMMs2ZfCKI/CSDzdfJG+27iAgfA8bUI= +github.com/nexus-rpc/sdk-go v0.7.0/go.mod h1:FHdPfVQwRuJFZFTF0Y2GOAxCrbIBNrcPna9slkGKPYk= github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno= github.com/olekukonko/tablewriter v0.0.5 h1:P2Ga83D34wi1o9J6Wh1mRuqd4mF/x/lgBS7N7AbDhec= github.com/olekukonko/tablewriter v0.0.5/go.mod h1:hPp6KlRPjbx+hW8ykQs1w3UBbZlj6HuIJcUGPhkA7kY= @@ -481,16 +481,16 @@ go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/ go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE= go.opentelemetry.io/proto/otlp v1.10.0 h1:IQRWgT5srOCYfiWnpqUYz9CVmbO8bFmKcwYxpuCSL2g= go.opentelemetry.io/proto/otlp v1.10.0/go.mod h1:/CV4QoCR/S9yaPj8utp3lvQPoqMtxXdzn7ozvvozVqk= -go.temporal.io/api v1.63.5 h1:c11+kPYHkXXL3UiShPdbMD+xtvqGsbTibUA9ypmiCa4= -go.temporal.io/api v1.63.5/go.mod h1:SrlW2JMwVlDP4nRWSNznUFqnSHd+YeMDS1BkYo63HCQ= +go.temporal.io/api v1.63.6-0.20260910200743-859a8e8d17c0 h1:S765rdHH0PFi7zFEmR1ThDteFKUnJw+nqY7ryA6b5UI= +go.temporal.io/api v1.63.6-0.20260910200743-859a8e8d17c0/go.mod h1:acM0I9WPuYg8W3Pd9jOZvEgi7mRUttUQ4+e7fowKVnM= go.temporal.io/auto-scaled-workers v0.2.0-1.32.0.158.0 h1:l+Rj0cIHMC2VB/DC+axrLLnZ2ISTgUpUw8TMiP4B9F8= go.temporal.io/auto-scaled-workers v0.2.0-1.32.0.158.0/go.mod h1:ZGNY7kCU0EZpXx02D81/jGXHbUcFcVIzzguJ/f8pvUQ= -go.temporal.io/sdk v1.47.0 h1:lZ39w1+uWSjHTL0F3mSc0t4XUnKX8CCcWxqSiLbaHnc= -go.temporal.io/sdk v1.47.0/go.mod h1:ilKs0twgP4JpP8pfhIgZumnOEyBiYn6ZO/ta//NnKMU= +go.temporal.io/sdk v1.48.0 h1:WDctKDVuh0Z8Nf7euAyqs/EwcPg1JTIIq1Fut8Tq118= +go.temporal.io/sdk v1.48.0/go.mod h1:SHv3+fLzD0GGZAwf0xNSvu8UmO1nFgG9WBSYoowApIk= go.temporal.io/sdk/contrib/envconfig v1.0.2 h1:MGHfsuPUtsf7X9M6WYn3zYJj/mWsuYHnA1uuiL0KEuE= go.temporal.io/sdk/contrib/envconfig v1.0.2/go.mod h1:MuMiH7hksps2uXnmKuAWaP9P6WbkSDy62kl64t1VJVg= -go.temporal.io/server v1.32.0 h1:JQoqsVREaGc8vcuDNfuRlySiV3tVi81XUOVC9LeL/1Y= -go.temporal.io/server v1.32.0/go.mod h1:SuxEWp1bDjSB7kHtUjyaUJBh7Qjyv5wB8lCrS0VFSvw= +go.temporal.io/server v1.32.0-163.3 h1:sltlJuke9JAatSqxVnc0zSzQjvH3ZUrZ8ZHPd/LxFEY= +go.temporal.io/server v1.32.0-163.3/go.mod h1:5Z6vEEG4JaRS4PRPMCwrRxyXERZF6WsyTqGhSpP1F4I= go.uber.org/atomic v1.5.0/go.mod h1:sABNBOSYdrvTF6hTgEIbc7YasKWGhgEQZyfxyTvoXHQ= go.uber.org/atomic v1.7.0/go.mod h1:fEN4uk6kAWBTFdckzkM89CLk9XfWZrxpCo0nPH17wJc= go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= diff --git a/internal/temporalcli/client.go b/internal/temporalcli/client.go index 9126a2b0f..6f081f858 100644 --- a/internal/temporalcli/client.go +++ b/internal/temporalcli/client.go @@ -23,17 +23,8 @@ import ( // so often used by callers after this call to know the currently configured // namespace. func dialClient(cctx *CommandContext, c *cliext.ClientOptions) (client.Client, error) { - cl, _, err := dialClientWithCodec(cctx, c) - return cl, err -} - -// dialClientWithCodec is like [dialClient] but also returns the configured remote -// payload codec, or nil if no codec is configured. The codec is the same instance -// used by the gRPC interceptor; callers can use it to decode payloads nested inside -// opaque proto bytes (e.g. the request/response of a system Nexus operation). -func dialClientWithCodec(cctx *CommandContext, c *cliext.ClientOptions) (client.Client, converter.PayloadCodec, error) { if cctx.RootCommand == nil { - return nil, nil, fmt.Errorf("root command unexpectedly missing when dialing client") + return nil, fmt.Errorf("root command unexpectedly missing when dialing client") } // Set default identity if not provided @@ -61,12 +52,12 @@ func dialClientWithCodec(cctx *CommandContext, c *cliext.ClientOptions) (client. // original setup error instead of attaching a guessed address or profile. var pathErr *fs.PathError if errors.As(err, &pathErr) { - return nil, nil, newConnectError(&connectDiagnosis{ + return nil, newConnectError(&connectDiagnosis{ Cause: causeCertFileUnreadable, Detail: pathErr.Path, }, connectMeta{}, err) } - return nil, nil, err + return nil, err } // We do not put codec on data converter here, it is applied via @@ -97,14 +88,14 @@ func dialClientWithCodec(cctx *CommandContext, c *cliext.ClientOptions) (client. cl, err := client.DialContext(dialCtx, clientOpts) if err != nil { - return nil, nil, dialConnectError(cctx, dialCtx, clientOpts, err) + return nil, dialConnectError(cctx, dialCtx, clientOpts, err) } // Since this namespace value is used by many commands after this call, // we are mutating it to be the derived one c.Namespace = clientOpts.Namespace - return cl, builder.PayloadCodec, nil + return cl, nil } // dialConnectError enriches a client.DialContext failure with a staged diff --git a/internal/temporalcli/commands.system_nexus.go b/internal/temporalcli/commands.system_nexus.go deleted file mode 100644 index da352b3c0..000000000 --- a/internal/temporalcli/commands.system_nexus.go +++ /dev/null @@ -1,54 +0,0 @@ -package temporalcli - -import ( - "context" - - commonpb "go.temporal.io/api/common/v1" - "go.temporal.io/api/proxy" - "go.temporal.io/api/workflowservice/v1" - "go.temporal.io/api/workflowservice/v1/workflowservicenexus" - "go.temporal.io/sdk/converter" - "google.golang.org/protobuf/proto" -) - -// systemNexusOpKey identifies a system Nexus operation by its (endpoint, operation) pair. -type systemNexusOpKey struct { - Endpoint string - Operation string -} - -// systemNexusOpTypes maps a system Nexus operation to the proto request and response types -// whose bytes are serialized in NexusOperationScheduled.Input and NexusOperationCompleted.Result. -type systemNexusOpTypes struct { - // NewRequest returns a fresh, zero-valued instance of the request proto. - NewRequest func() proto.Message - // NewResponse returns a fresh, zero-valued instance of the response proto. - NewResponse func() proto.Message -} - -// systemNexusOps is the global registry of known system Nexus operations on the -// __temporal_system endpoint. Add new entries here as the server adds support for more -// system operations. The keys' Operation values must match what the server records in -// NexusOperationScheduledEventAttributes.Operation. -// NOTE seankane: Part 2 of the System Operations work is to code generate this map from the -// go.temporal.io/api/workflowservice/v1/workflowservicenexus package. -var systemNexusOps = map[systemNexusOpKey]systemNexusOpTypes{ - { - Endpoint: temporalSystemNexusEndpoint, - Operation: workflowservicenexus.TemporalAPIWorkflowserviceV1WorkflowService.SignalWithStartWorkflowExecution.Name(), - }: { - NewRequest: func() proto.Message { return &workflowservice.SignalWithStartWorkflowExecutionRequest{} }, - NewResponse: func() proto.Message { return &workflowservice.SignalWithStartWorkflowExecutionResponse{} }, - }, -} - -// decodePayloadsInProto walks a proto message and applies codec.Decode to every Payload -// found inside it (including nested messages). The message is mutated in place. -func decodePayloadsInProto(ctx context.Context, msg proto.Message, codec converter.PayloadCodec) error { - return proxy.VisitPayloads(ctx, msg, proxy.VisitPayloadsOptions{ - SkipSearchAttributes: true, - Visitor: func(_ *proxy.VisitPayloadsContext, payloads []*commonpb.Payload) ([]*commonpb.Payload, error) { - return codec.Decode(payloads) - }, - }) -} diff --git a/internal/temporalcli/commands.system_nexus_test.go b/internal/temporalcli/commands.system_nexus_test.go index ad160df2e..0639b874c 100644 --- a/internal/temporalcli/commands.system_nexus_test.go +++ b/internal/temporalcli/commands.system_nexus_test.go @@ -1,332 +1,219 @@ package temporalcli import ( + "bytes" "context" - "fmt" "testing" "github.com/stretchr/testify/require" commonpb "go.temporal.io/api/common/v1" enumspb "go.temporal.io/api/enums/v1" historypb "go.temporal.io/api/history/v1" + "go.temporal.io/api/proxy" "go.temporal.io/api/temporalproto" "go.temporal.io/api/workflowservice/v1" "google.golang.org/protobuf/proto" ) -// markingCodec is a test [converter.PayloadCodec] that prefixes every payload's -// data with "decoded:" on Decode and tracks how many payloads it saw. It is used -// to verify that the codec is actually invoked on payloads nested inside opaque -// system Nexus operation bytes. -type markingCodec struct { - decodeCalls int -} - -func (c *markingCodec) Encode(payloads []*commonpb.Payload) ([]*commonpb.Payload, error) { - return payloads, nil -} - -func (c *markingCodec) Decode(payloads []*commonpb.Payload) ([]*commonpb.Payload, error) { - out := make([]*commonpb.Payload, len(payloads)) - for i, p := range payloads { - c.decodeCalls++ - out[i] = &commonpb.Payload{ - Metadata: p.Metadata, - Data: append([]byte("decoded:"), p.Data...), - } - } - return out, nil -} - -// failingCodec always returns an error from Decode; used to verify error propagation. -type failingCodec struct{} - -func (failingCodec) Encode(payloads []*commonpb.Payload) ([]*commonpb.Payload, error) { - return payloads, nil -} - -func (failingCodec) Decode(_ []*commonpb.Payload) ([]*commonpb.Payload, error) { - return nil, fmt.Errorf("codec decode failure for testing") -} - -func signalWithStartRequestPayload(t *testing.T, req *workflowservice.SignalWithStartWorkflowExecutionRequest) *commonpb.Payload { +func markedSystemPayload(t *testing.T, msg proto.Message) *commonpb.Payload { t.Helper() - data, err := proto.Marshal(req) + data, err := proto.Marshal(msg) require.NoError(t, err) return &commonpb.Payload{ - Metadata: map[string][]byte{"encoding": []byte("binary/protobuf")}, - Data: data, - } -} - -func signalWithStartResponsePayload(t *testing.T, resp *workflowservice.SignalWithStartWorkflowExecutionResponse) *commonpb.Payload { - t.Helper() - data, err := proto.Marshal(resp) - require.NoError(t, err) - return &commonpb.Payload{ - Metadata: map[string][]byte{"encoding": []byte("binary/protobuf")}, - Data: data, + Metadata: map[string][]byte{ + "encoding": []byte("binary/protobuf"), + "messageType": []byte(msg.ProtoReflect().Descriptor().FullName()), + proxy.SystemPayloadMetadataKey: []byte("true"), + }, + Data: data, } } -func TestUnwrapAndInjectRequest_NilPayloadIsNoOp(t *testing.T) { - iter := &structuredHistoryIter{ctx: context.Background()} - fields := map[string]any{} - require.NoError(t, iter.unwrapAndInjectRequest( - temporalSystemNexusEndpoint, "SignalWithStartWorkflowExecution", - nil, fields, temporalproto.CustomJSONMarshalOptions{})) - require.Empty(t, fields, "nil payload should not inject anything") -} - -func TestUnwrapAndInjectRequest_UnknownOperationIsNoOp(t *testing.T) { - p := &commonpb.Payload{Data: []byte("ignored")} - iter := &structuredHistoryIter{ctx: context.Background()} - fields := map[string]any{} - require.NoError(t, iter.unwrapAndInjectRequest( - temporalSystemNexusEndpoint, "NotARealOperation", - p, fields, temporalproto.CustomJSONMarshalOptions{})) - require.Empty(t, fields, "unknown operation should be a no-op") -} - -func TestUnwrapAndInjectRequest_UnknownEndpointIsNoOp(t *testing.T) { - p := &commonpb.Payload{Data: []byte("ignored")} - iter := &structuredHistoryIter{ctx: context.Background()} +func TestUnwrapAndInject_NilPayloadIsNoOp(t *testing.T) { fields := map[string]any{} - require.NoError(t, iter.unwrapAndInjectRequest( - "some-user-endpoint", "SignalWithStartWorkflowExecution", - p, fields, temporalproto.CustomJSONMarshalOptions{})) - require.Empty(t, fields, "non-system endpoint should be a no-op even if operation name matches") -} - -func TestUnwrapAndInjectRequest_BadProtoBytesReturnsError(t *testing.T) { - p := &commonpb.Payload{Data: []byte{0xff, 0xff, 0xff}} - iter := &structuredHistoryIter{ctx: context.Background()} - fields := map[string]any{} - err := iter.unwrapAndInjectRequest( - temporalSystemNexusEndpoint, "SignalWithStartWorkflowExecution", - p, fields, temporalproto.CustomJSONMarshalOptions{}) - require.Error(t, err) - require.ErrorContains(t, err, "failed unmarshaling system nexus payload") -} - -func TestUnwrapAndInjectResponse_NilPayloadIsNoOp(t *testing.T) { - iter := &structuredHistoryIter{ctx: context.Background()} - fields := map[string]any{} - require.NoError(t, iter.unwrapAndInjectResponse( - temporalSystemNexusEndpoint, "SignalWithStartWorkflowExecution", - nil, fields, temporalproto.CustomJSONMarshalOptions{})) + require.NoError(t, (&structuredHistoryIter{}).unwrapAndInject( + nil, fields, "unwrappedInput", temporalproto.CustomJSONMarshalOptions{})) require.Empty(t, fields) } -func TestUnwrapAndInjectResponse_UnknownOperationIsNoOp(t *testing.T) { - p := &commonpb.Payload{Data: []byte("ignored")} - iter := &structuredHistoryIter{ctx: context.Background()} - fields := map[string]any{} - require.NoError(t, iter.unwrapAndInjectResponse( - temporalSystemNexusEndpoint, "NotARealOperation", - p, fields, temporalproto.CustomJSONMarshalOptions{})) - require.Empty(t, fields) -} - -func TestUnwrapAndInjectResponse_BadProtoBytesReturnsError(t *testing.T) { - p := &commonpb.Payload{Data: []byte{0xff, 0xff, 0xff}} - iter := &structuredHistoryIter{ctx: context.Background()} - fields := map[string]any{} - err := iter.unwrapAndInjectResponse( - temporalSystemNexusEndpoint, "SignalWithStartWorkflowExecution", - p, fields, temporalproto.CustomJSONMarshalOptions{}) - require.Error(t, err) - require.ErrorContains(t, err, "failed unmarshaling system nexus payload") -} +func TestUnwrapAndInject_RejectsInvalidEnvelope(t *testing.T) { + valid := markedSystemPayload(t, &workflowservice.SignalWithStartWorkflowExecutionRequest{}) + tests := []struct { + name string + mutate func(*commonpb.Payload) + wantErr string + }{ + { + name: "missing marker", + mutate: func(payload *commonpb.Payload) { + delete(payload.Metadata, proxy.SystemPayloadMetadataKey) + }, + wantErr: "missing the __temporal_system_payload marker", + }, + { + name: "wrong encoding", + mutate: func(payload *commonpb.Payload) { + payload.Metadata["encoding"] = []byte("json/protobuf") + }, + wantErr: "must be encoded as binary/protobuf", + }, + { + name: "missing message type", + mutate: func(payload *commonpb.Payload) { + delete(payload.Metadata, "messageType") + }, + wantErr: "missing messageType metadata", + }, + { + name: "unknown message type", + mutate: func(payload *commonpb.Payload) { + payload.Metadata["messageType"] = []byte("temporal.api.unknown.v1.Message") + }, + wantErr: "references unknown message type", + }, + { + name: "invalid protobuf", + mutate: func(payload *commonpb.Payload) { + payload.Data = []byte{0xff, 0xff, 0xff} + }, + wantErr: "failed unmarshaling system nexus payload", + }, + } -func TestUnwrapAndInjectRequest_DecodesAllNestedPayloads(t *testing.T) { - // The Input/SignalInput fields hold the user-supplied payloads, which the codec - // should be applied to. The outer payload bytes are raw proto and are not codec-encoded. - inner1 := &commonpb.Payload{Metadata: map[string][]byte{"encoding": []byte("binary/plain")}, Data: []byte("hello")} - inner2 := &commonpb.Payload{Metadata: map[string][]byte{"encoding": []byte("binary/plain")}, Data: []byte("world")} - signalInner := &commonpb.Payload{Metadata: map[string][]byte{"encoding": []byte("binary/plain")}, Data: []byte("signal-arg")} - req := &workflowservice.SignalWithStartWorkflowExecutionRequest{ - Namespace: "ns", - WorkflowId: "wf", - Input: &commonpb.Payloads{Payloads: []*commonpb.Payload{inner1, inner2}}, - SignalInput: &commonpb.Payloads{Payloads: []*commonpb.Payload{signalInner}}, + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + payload := proto.Clone(valid).(*commonpb.Payload) + tc.mutate(payload) + err := (&structuredHistoryIter{}).unwrapAndInject( + payload, map[string]any{}, "unwrappedInput", temporalproto.CustomJSONMarshalOptions{}) + require.ErrorContains(t, err, tc.wantErr) + }) } - p := signalWithStartRequestPayload(t, req) +} - codec := &markingCodec{} - iter := &structuredHistoryIter{ctx: context.Background(), codec: codec} +func TestUnwrapAndInject_UsesEnvelopeMessageType(t *testing.T) { + payload := markedSystemPayload(t, &workflowservice.SignalWithStartWorkflowExecutionRequest{ + Namespace: "ns", + WorkflowId: "wf", + SignalName: "signal", + }) fields := map[string]any{} - require.NoError(t, iter.unwrapAndInjectRequest( - temporalSystemNexusEndpoint, "SignalWithStartWorkflowExecution", - p, fields, temporalproto.CustomJSONMarshalOptions{})) - require.Equal(t, 3, codec.decodeCalls, "codec should have been invoked once per nested payload") + require.NoError(t, (&structuredHistoryIter{}).unwrapAndInject( + payload, fields, "unwrappedInput", temporalproto.CustomJSONMarshalOptions{})) unwrapped, ok := fields["unwrappedInput"].(map[string]any) require.True(t, ok) require.Equal(t, "ns", unwrapped["namespace"]) require.Equal(t, "wf", unwrapped["workflowId"]) + require.Equal(t, "signal", unwrapped["signalName"]) } -func TestUnwrapAndInjectRequest_CodecErrorPropagates(t *testing.T) { - req := &workflowservice.SignalWithStartWorkflowExecutionRequest{ - Input: &commonpb.Payloads{Payloads: []*commonpb.Payload{{Data: []byte("x")}}}, - } - p := signalWithStartRequestPayload(t, req) - iter := &structuredHistoryIter{ctx: context.Background(), codec: failingCodec{}} - fields := map[string]any{} - err := iter.unwrapAndInjectRequest( - temporalSystemNexusEndpoint, "SignalWithStartWorkflowExecution", - p, fields, temporalproto.CustomJSONMarshalOptions{}) - require.Error(t, err) - require.Equal(t, "failed decoding payloads in system nexus payload: codec decode failure for testing", err.Error()) -} - -func TestDecodePayloadsInProto_VisitsAllPayloads(t *testing.T) { +func TestSystemNexusPayload_DecodedByVisitorAndUnwrappedLocally(t *testing.T) { req := &workflowservice.SignalWithStartWorkflowExecutionRequest{ + WorkflowId: "wf", Input: &commonpb.Payloads{Payloads: []*commonpb.Payload{ - {Data: []byte("a")}, - {Data: []byte("b")}, - }}, - SignalInput: &commonpb.Payloads{Payloads: []*commonpb.Payload{ - {Data: []byte("c")}, + {Data: []byte("encoded:first")}, + {Data: []byte("encoded:second")}, }}, } - codec := &markingCodec{} - require.NoError(t, decodePayloadsInProto(context.Background(), req, codec)) - require.Equal(t, 3, codec.decodeCalls) - require.Equal(t, []byte("decoded:a"), req.Input.Payloads[0].Data) - require.Equal(t, []byte("decoded:b"), req.Input.Payloads[1].Data) - require.Equal(t, []byte("decoded:c"), req.SignalInput.Payloads[0].Data) + payload := markedSystemPayload(t, req) + attrs := &historypb.NexusOperationScheduledEventAttributes{Input: payload} + + decodeRequests := 0 + decodedPayloads := 0 + err := proxy.VisitPayloads(context.Background(), attrs, proxy.VisitPayloadsOptions{ + Visitor: func(_ *proxy.VisitPayloadsContext, payloads []*commonpb.Payload) ([]*commonpb.Payload, error) { + decodeRequests++ + decodedPayloads += len(payloads) + decoded := make([]*commonpb.Payload, len(payloads)) + for i, payload := range payloads { + require.True(t, bytes.HasPrefix(payload.Data, []byte("encoded:"))) + decoded[i] = proto.Clone(payload).(*commonpb.Payload) + decoded[i].Data = bytes.TrimPrefix(payload.Data, []byte("encoded:")) + } + return decoded, nil + }, + }) + require.NoError(t, err) + require.Equal(t, 1, decodeRequests) + require.Equal(t, 2, decodedPayloads) + + fields := map[string]any{} + require.NoError(t, (&structuredHistoryIter{}).unwrapAndInject( + attrs.Input, fields, "unwrappedInput", temporalproto.CustomJSONMarshalOptions{})) + require.Equal(t, 1, decodeRequests, "local rendering must not revisit the codec") + require.Contains(t, fields, "unwrappedInput") } -func TestInjectSystemNexusUnwrapped_ScheduledKnownOp(t *testing.T) { - req := &workflowservice.SignalWithStartWorkflowExecutionRequest{ - Namespace: "ns", - WorkflowId: "wf-xyz", - SignalName: "ping", - } +func TestInjectSystemNexusUnwrapped_ScheduledUsesMessageTypeForAnyOperation(t *testing.T) { event := &historypb.HistoryEvent{ EventId: 5, EventType: enumspb.EVENT_TYPE_NEXUS_OPERATION_SCHEDULED, Attributes: &historypb.HistoryEvent_NexusOperationScheduledEventAttributes{ NexusOperationScheduledEventAttributes: &historypb.NexusOperationScheduledEventAttributes{ Endpoint: temporalSystemNexusEndpoint, - Operation: "SignalWithStartWorkflowExecution", - Input: signalWithStartRequestPayload(t, req), + Operation: "FutureSystemOperation", + Input: markedSystemPayload(t, &workflowservice.SignalWithStartWorkflowExecutionRequest{ + WorkflowId: "wf", + }), }, }, } - iter := &structuredHistoryIter{ctx: context.Background()} fields := map[string]any{} - require.NoError(t, iter.injectSystemNexusUnwrapped(event, fields, temporalproto.CustomJSONMarshalOptions{})) - - unwrapped, ok := fields["unwrappedInput"].(map[string]any) - require.True(t, ok, "expected unwrappedInput map to be set") - require.Equal(t, "ns", unwrapped["namespace"]) - require.Equal(t, "wf-xyz", unwrapped["workflowId"]) - require.Equal(t, "ping", unwrapped["signalName"]) + require.NoError(t, (&structuredHistoryIter{}).injectSystemNexusUnwrapped( + event, fields, temporalproto.CustomJSONMarshalOptions{})) + require.Contains(t, fields, "unwrappedInput") } -func TestInjectSystemNexusUnwrapped_ScheduledUnknownEndpointSkipped(t *testing.T) { - req := &workflowservice.SignalWithStartWorkflowExecutionRequest{Namespace: "ns"} +func TestInjectSystemNexusUnwrapped_UnknownEndpointIsNoOp(t *testing.T) { event := &historypb.HistoryEvent{ EventType: enumspb.EVENT_TYPE_NEXUS_OPERATION_SCHEDULED, Attributes: &historypb.HistoryEvent_NexusOperationScheduledEventAttributes{ NexusOperationScheduledEventAttributes: &historypb.NexusOperationScheduledEventAttributes{ - Endpoint: "user-endpoint", - Operation: "SignalWithStartWorkflowExecution", - Input: signalWithStartRequestPayload(t, req), + Endpoint: "user-endpoint", + Input: markedSystemPayload(t, &workflowservice.SignalWithStartWorkflowExecutionRequest{}), }, }, } - iter := &structuredHistoryIter{ctx: context.Background()} fields := map[string]any{} - require.NoError(t, iter.injectSystemNexusUnwrapped(event, fields, temporalproto.CustomJSONMarshalOptions{})) - _, ok := fields["unwrappedInput"] - require.False(t, ok, "non-system endpoint must not produce unwrappedInput") + require.NoError(t, (&structuredHistoryIter{}).injectSystemNexusUnwrapped( + event, fields, temporalproto.CustomJSONMarshalOptions{})) + require.Empty(t, fields) } -func TestInjectSystemNexusUnwrapped_CompletedUsesPriorScheduled(t *testing.T) { - resp := &workflowservice.SignalWithStartWorkflowExecutionResponse{ - RunId: "run-abc", - Started: true, - } - completed := &historypb.HistoryEvent{ - EventId: 6, +func TestInjectSystemNexusUnwrapped_CompletedUsesPriorScheduledEvent(t *testing.T) { + event := &historypb.HistoryEvent{ EventType: enumspb.EVENT_TYPE_NEXUS_OPERATION_COMPLETED, Attributes: &historypb.HistoryEvent_NexusOperationCompletedEventAttributes{ NexusOperationCompletedEventAttributes: &historypb.NexusOperationCompletedEventAttributes{ ScheduledEventId: 5, - Result: signalWithStartResponsePayload(t, resp), + Result: markedSystemPayload(t, &workflowservice.SignalWithStartWorkflowExecutionResponse{ + RunId: "run-id", + }), }, }, } - iter := &structuredHistoryIter{ - ctx: context.Background(), - systemNexusOps: map[int64]string{5: "SignalWithStartWorkflowExecution"}, - } + iter := &structuredHistoryIter{systemNexusOps: map[int64]string{5: "FutureSystemOperation"}} fields := map[string]any{} - require.NoError(t, iter.injectSystemNexusUnwrapped(completed, fields, temporalproto.CustomJSONMarshalOptions{})) - + require.NoError(t, iter.injectSystemNexusUnwrapped( + event, fields, temporalproto.CustomJSONMarshalOptions{})) unwrapped, ok := fields["unwrappedResult"].(map[string]any) - require.True(t, ok, "expected unwrappedResult map to be set") - require.Equal(t, "run-abc", unwrapped["runId"]) - require.Equal(t, true, unwrapped["started"]) + require.True(t, ok) + require.Equal(t, "run-id", unwrapped["runId"]) } -func TestInjectSystemNexusUnwrapped_CompletedWithoutPriorScheduledSkipped(t *testing.T) { - completed := &historypb.HistoryEvent{ +func TestInjectSystemNexusUnwrapped_CompletedWithoutPriorScheduledIsNoOp(t *testing.T) { + event := &historypb.HistoryEvent{ EventType: enumspb.EVENT_TYPE_NEXUS_OPERATION_COMPLETED, Attributes: &historypb.HistoryEvent_NexusOperationCompletedEventAttributes{ NexusOperationCompletedEventAttributes: &historypb.NexusOperationCompletedEventAttributes{ ScheduledEventId: 5, - Result: &commonpb.Payload{Data: []byte("garbage")}, + Result: markedSystemPayload(t, &workflowservice.SignalWithStartWorkflowExecutionResponse{}), }, }, } - iter := &structuredHistoryIter{ctx: context.Background()} - fields := map[string]any{} - require.NoError(t, iter.injectSystemNexusUnwrapped(completed, fields, temporalproto.CustomJSONMarshalOptions{})) - _, ok := fields["unwrappedResult"] - require.False(t, ok, "no prior scheduled means we don't know the op, so no unwrap") -} - -func TestInjectSystemNexusUnwrapped_NonNexusEventNoOp(t *testing.T) { - event := &historypb.HistoryEvent{ - EventType: enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_STARTED, - Attributes: &historypb.HistoryEvent_WorkflowExecutionStartedEventAttributes{ - WorkflowExecutionStartedEventAttributes: &historypb.WorkflowExecutionStartedEventAttributes{}, - }, - } - iter := &structuredHistoryIter{ctx: context.Background()} fields := map[string]any{} - require.NoError(t, iter.injectSystemNexusUnwrapped(event, fields, temporalproto.CustomJSONMarshalOptions{})) + require.NoError(t, (&structuredHistoryIter{}).injectSystemNexusUnwrapped( + event, fields, temporalproto.CustomJSONMarshalOptions{})) require.Empty(t, fields) } - -func TestInjectSystemNexusUnwrapped_AppliesCodecToScheduledInput(t *testing.T) { - req := &workflowservice.SignalWithStartWorkflowExecutionRequest{ - WorkflowId: "wf", - Input: &commonpb.Payloads{Payloads: []*commonpb.Payload{ - {Data: []byte("inner-input")}, - }}, - SignalInput: &commonpb.Payloads{Payloads: []*commonpb.Payload{ - {Data: []byte("inner-signal")}, - }}, - } - event := &historypb.HistoryEvent{ - EventType: enumspb.EVENT_TYPE_NEXUS_OPERATION_SCHEDULED, - Attributes: &historypb.HistoryEvent_NexusOperationScheduledEventAttributes{ - NexusOperationScheduledEventAttributes: &historypb.NexusOperationScheduledEventAttributes{ - Endpoint: temporalSystemNexusEndpoint, - Operation: "SignalWithStartWorkflowExecution", - Input: signalWithStartRequestPayload(t, req), - }, - }, - } - codec := &markingCodec{} - iter := &structuredHistoryIter{ctx: context.Background(), codec: codec} - fields := map[string]any{} - require.NoError(t, iter.injectSystemNexusUnwrapped(event, fields, temporalproto.CustomJSONMarshalOptions{})) - require.Equal(t, 2, codec.decodeCalls, "codec should run on both nested payloads") -} diff --git a/internal/temporalcli/commands.workflow_exec.go b/internal/temporalcli/commands.workflow_exec.go index cef46fdf2..194858b51 100644 --- a/internal/temporalcli/commands.workflow_exec.go +++ b/internal/temporalcli/commands.workflow_exec.go @@ -21,13 +21,15 @@ import ( "go.temporal.io/api/enums/v1" enumspb "go.temporal.io/api/enums/v1" "go.temporal.io/api/history/v1" + "go.temporal.io/api/proxy" taskqueuepb "go.temporal.io/api/taskqueue/v1" "go.temporal.io/api/temporalproto" "go.temporal.io/api/workflowservice/v1" "go.temporal.io/sdk/client" - "go.temporal.io/sdk/converter" "go.temporal.io/sdk/temporal" "google.golang.org/protobuf/proto" + "google.golang.org/protobuf/reflect/protoreflect" + "google.golang.org/protobuf/reflect/protoregistry" "google.golang.org/protobuf/types/known/durationpb" ) @@ -42,7 +44,7 @@ func (c *TemporalWorkflowStartCommand) run(cctx *CommandContext, args []string) } func (c *TemporalWorkflowExecuteCommand) run(cctx *CommandContext, args []string) error { - cl, codec, err := dialClientWithCodec(cctx, &c.Parent.ClientOptions) + cl, err := dialClient(cctx, &c.Parent.ClientOptions) if err != nil { return err } @@ -65,7 +67,6 @@ func (c *TemporalWorkflowExecuteCommand) run(cctx *CommandContext, args []string runID: run.GetRunID(), includeDetails: c.Detailed, follow: true, - codec: codec, } if err := iter.print(cctx); err != nil && cctx.Err() == nil { return fmt.Errorf("displaying history failed: %w", err) @@ -793,11 +794,6 @@ type structuredHistoryIter struct { // maps NexusOperationScheduled eventId → operation name for __temporal_system endpoint events systemNexusOps map[int64]string - - // codec is the remote payload codec configured for this client, or nil if none. Used to - // decode payloads nested inside system Nexus operation request/response bytes so they - // can be rendered alongside the rest of the event fields. - codec converter.PayloadCodec } func (s *structuredHistoryIter) print(cctx *CommandContext) error { @@ -982,8 +978,8 @@ func (s *structuredHistoryIter) flattenFields( } } // For system Nexus operation events, deserialize the request/response payload bytes - // into the typed proto, decode any payloads nested inside via the codec, and merge - // the decoded view into the output under "unwrappedInput" / "unwrappedResult". + // into the typed proto and merge the view into the output under "unwrappedInput" / + // "unwrappedResult". The gRPC interceptor has already handled codec decoding. if err := s.injectSystemNexusUnwrapped(event, fieldsMap, opts); err != nil { return nil, err } @@ -997,9 +993,9 @@ func (s *structuredHistoryIter) flattenFields( } // injectSystemNexusUnwrapped, if the given event is a known system Nexus operation, -// deserializes the underlying request (on Scheduled) or response (on Completed) proto, -// decodes any payloads nested inside via the codec, and inserts the resulting JSON -// representation into fieldsMap under "unwrappedInput" / "unwrappedResult". +// deserializes the underlying request (on Scheduled) or response (on Completed) proto +// and inserts the resulting JSON representation into fieldsMap under "unwrappedInput" / +// "unwrappedResult". func (s *structuredHistoryIter) injectSystemNexusUnwrapped( event *history.HistoryEvent, fieldsMap map[string]any, @@ -1008,59 +1004,28 @@ func (s *structuredHistoryIter) injectSystemNexusUnwrapped( switch event.EventType { case enums.EVENT_TYPE_NEXUS_OPERATION_SCHEDULED: attr := event.GetNexusOperationScheduledEventAttributes() - if attr == nil { + if attr == nil || attr.GetEndpoint() != temporalSystemNexusEndpoint { return nil } - return s.unwrapAndInjectRequest(attr.GetEndpoint(), attr.GetOperation(), attr.GetInput(), fieldsMap, opts) + return s.unwrapAndInject(attr.GetInput(), fieldsMap, "unwrappedInput", opts) case enums.EVENT_TYPE_NEXUS_OPERATION_COMPLETED: attr := event.GetNexusOperationCompletedEventAttributes() if attr == nil { return nil } - op, ok := s.systemNexusOps[attr.GetScheduledEventId()] + _, ok := s.systemNexusOps[attr.GetScheduledEventId()] if !ok { return nil } - return s.unwrapAndInjectResponse(temporalSystemNexusEndpoint, op, attr.GetResult(), fieldsMap, opts) + return s.unwrapAndInject(attr.GetResult(), fieldsMap, "unwrappedResult", opts) } return nil } -// unwrapAndInjectRequest looks up the registered request proto for (endpoint, operation), -// then injects the decoded view under "unwrappedInput". No-op for unregistered ops. -func (s *structuredHistoryIter) unwrapAndInjectRequest( - endpoint, operation string, - payload *commonpb.Payload, - fieldsMap map[string]any, - opts temporalproto.CustomJSONMarshalOptions, -) error { - types, ok := systemNexusOps[systemNexusOpKey{Endpoint: endpoint, Operation: operation}] - if !ok { - return nil - } - return s.unwrapAndInject(types.NewRequest(), payload, fieldsMap, "unwrappedInput", opts) -} - -// unwrapAndInjectResponse looks up the registered response proto for (endpoint, operation), -// then injects the decoded view under "unwrappedResult". No-op for unregistered ops. -func (s *structuredHistoryIter) unwrapAndInjectResponse( - endpoint, operation string, - payload *commonpb.Payload, - fieldsMap map[string]any, - opts temporalproto.CustomJSONMarshalOptions, -) error { - types, ok := systemNexusOps[systemNexusOpKey{Endpoint: endpoint, Operation: operation}] - if !ok { - return nil - } - return s.unwrapAndInject(types.NewResponse(), payload, fieldsMap, "unwrappedResult", opts) -} - -// unwrapAndInject is the shared body: unmarshal payload bytes into the supplied proto, -// decode any payloads nested inside via the codec, marshal back to JSON, and inject the -// resulting map into fieldsMap[key]. A nil payload is a no-op. +// unwrapAndInject resolves the marked envelope's protobuf type, unmarshals it, and injects +// its JSON representation into fieldsMap[key]. When configured, the gRPC payload codec +// interceptor has already visited nested payloads. A nil payload is a no-op. func (s *structuredHistoryIter) unwrapAndInject( - msg proto.Message, payload *commonpb.Payload, fieldsMap map[string]any, key string, @@ -1069,14 +1034,24 @@ func (s *structuredHistoryIter) unwrapAndInject( if payload == nil { return nil } + if string(payload.GetMetadata()[proxy.SystemPayloadMetadataKey]) != "true" { + return fmt.Errorf("system nexus payload is missing the %s marker", proxy.SystemPayloadMetadataKey) + } + if encoding := string(payload.GetMetadata()["encoding"]); encoding != "binary/protobuf" { + return fmt.Errorf("system nexus payload must be encoded as binary/protobuf but got %q", encoding) + } + messageType := protoreflect.FullName(payload.GetMetadata()["messageType"]) + if messageType == "" { + return fmt.Errorf("system nexus payload is missing messageType metadata") + } + messageDescriptor, err := protoregistry.GlobalTypes.FindMessageByName(messageType) + if err != nil { + return fmt.Errorf("system nexus payload references unknown message type %q: %w", messageType, err) + } + msg := messageDescriptor.New().Interface() if err := proto.Unmarshal(payload.Data, msg); err != nil { return fmt.Errorf("failed unmarshaling system nexus payload: %w", err) } - if s.codec != nil { - if err := decodePayloadsInProto(s.ctx, msg, s.codec); err != nil { - return fmt.Errorf("failed decoding payloads in system nexus payload: %w", err) - } - } unwrappedJSON, err := opts.Marshal(msg) if err != nil { return fmt.Errorf("failed marshaling unwrapped system nexus payload: %w", err) diff --git a/internal/temporalcli/commands.workflow_view.go b/internal/temporalcli/commands.workflow_view.go index c44711a11..98af51c1e 100644 --- a/internal/temporalcli/commands.workflow_view.go +++ b/internal/temporalcli/commands.workflow_view.go @@ -536,7 +536,7 @@ func (c *TemporalWorkflowShowCommand) run(cctx *CommandContext, _ []string) erro } // Call describe - cl, codec, err := dialClientWithCodec(cctx, &c.Parent.ClientOptions) + cl, err := dialClient(cctx, &c.Parent.ClientOptions) if err != nil { return err } @@ -552,7 +552,6 @@ func (c *TemporalWorkflowShowCommand) run(cctx *CommandContext, _ []string) erro includeDetails: c.Detailed, follow: c.Follow, reverse: c.Reverse, - codec: codec, } if !cctx.JSONOutput { cctx.Printer.Println(color.MagentaString("Progress:")) diff --git a/internal/temporalcli/commands.workflow_view_test.go b/internal/temporalcli/commands.workflow_view_test.go index 0c3929305..f00138a62 100644 --- a/internal/temporalcli/commands.workflow_view_test.go +++ b/internal/temporalcli/commands.workflow_view_test.go @@ -5,9 +5,12 @@ import ( "encoding/base64" "encoding/json" "fmt" + "net/http" + "net/http/httptest" "strconv" "strings" "sync" + "sync/atomic" "time" "github.com/google/uuid" @@ -22,17 +25,19 @@ import ( historypb "go.temporal.io/api/history/v1" nexuspb "go.temporal.io/api/nexus/v1" "go.temporal.io/api/operatorservice/v1" + "go.temporal.io/api/proxy" taskqueuepb "go.temporal.io/api/taskqueue/v1" workflowpb "go.temporal.io/api/workflow/v1" "go.temporal.io/api/workflowservice/v1" "go.temporal.io/api/workflowservice/v1/workflowservicenexus" "go.temporal.io/sdk/client" + "go.temporal.io/sdk/converter" "go.temporal.io/sdk/temporal" "go.temporal.io/sdk/temporalnexus" "go.temporal.io/sdk/worker" "go.temporal.io/sdk/workflow" - "go.temporal.io/server/common/payloads" "google.golang.org/grpc" + "google.golang.org/protobuf/proto" ) func (s *SharedServerSuite) TestWorkflow_Describe_ActivityFailing() { @@ -1271,7 +1276,7 @@ const temporalSystemNexusEndpointName = "__temporal_system" // to complete, then completes the caller. Returns the caller workflow ID. The target // workflow is registered for cleanup via s.T().Cleanup. The SDK cannot be used here // because it refuses endpoints with the reserved "__temporal_" prefix. -func (s *SharedServerSuite) runSystemNexusSWSWorkflow(ctx context.Context) string { +func (s *SharedServerSuite) runSystemNexusSWSWorkflow(ctx context.Context, input *common.Payloads) string { callerTaskQueue := "cli-sys-nexus-caller-" + uuid.NewString() targetTaskQueue := "cli-sys-nexus-target-" + uuid.NewString() targetWorkflowID := "cli-sys-nexus-target-" + uuid.NewString() @@ -1294,6 +1299,25 @@ func (s *SharedServerSuite) runSystemNexusSWSWorkflow(ctx context.Context) strin s.NoError(err) s.Equal(callerWorkflowID, pollResp.WorkflowExecution.WorkflowId) s.Equal(startResp.RunId, pollResp.WorkflowExecution.RunId) + operationRequest := &workflowservice.SignalWithStartWorkflowExecutionRequest{ + WorkflowId: targetWorkflowID, + SignalName: "cli-test-signal", + WorkflowType: &common.WorkflowType{Name: "target-workflow"}, + TaskQueue: &taskqueuepb.TaskQueue{Name: targetTaskQueue}, + Input: input, + } + operationInputData, err := proto.Marshal(operationRequest) + s.NoError(err) + // Generated System Nexus API clients mark protobuf envelopes that can contain + // nested payloads so payload visitors unwrap the envelope before invoking codecs. + operationInput := &common.Payload{ + Metadata: map[string][]byte{ + "encoding": []byte("binary/protobuf"), + "messageType": []byte(operationRequest.ProtoReflect().Descriptor().FullName()), + proxy.SystemPayloadMetadataKey: []byte("true"), + }, + Data: operationInputData, + } _, err = s.Client.WorkflowService().RespondWorkflowTaskCompleted(ctx, &workflowservice.RespondWorkflowTaskCompletedRequest{ Identity: "cli-test", @@ -1306,12 +1330,7 @@ func (s *SharedServerSuite) runSystemNexusSWSWorkflow(ctx context.Context) strin Endpoint: temporalSystemNexusEndpointName, Service: workflowservicenexus.TemporalAPIWorkflowserviceV1WorkflowService.ServiceName, Operation: workflowservicenexus.TemporalAPIWorkflowserviceV1WorkflowService.SignalWithStartWorkflowExecution.Name(), - Input: payloads.MustEncodeSingle(&workflowservice.SignalWithStartWorkflowExecutionRequest{ - WorkflowId: targetWorkflowID, - SignalName: "cli-test-signal", - WorkflowType: &common.WorkflowType{Name: "target-workflow"}, - TaskQueue: &taskqueuepb.TaskQueue{Name: targetTaskQueue}, - }), + Input: operationInput, }, }, }, @@ -1355,28 +1374,46 @@ func (s *SharedServerSuite) runSystemNexusSWSWorkflow(ctx context.Context) strin return callerWorkflowID } -// TestWorkflow_Show_SystemNexusOperationTransformsTypeNames drives a SignalWithStart -// Nexus operation against the __temporal_system endpoint from inside a workflow, then -// verifies that `workflow show` (default table mode) renders the resulting history -// events using the operation-prefixed names (e.g. SignalWithStartWorkflowExecutionScheduled) -// instead of the generic NexusOperation* names. -func (s *SharedServerSuite) TestWorkflow_Show_SystemNexusOperationTransformsTypeNames() { +// TestWorkflow_Show_SystemNexusOperationWithCodec drives a SignalWithStart Nexus +// operation containing codec-encoded input against the __temporal_system endpoint, then +// verifies that detailed `workflow show` decodes the nested input through the HTTP codec +// exactly once and locally unwraps it for display. +func (s *SharedServerSuite) TestWorkflow_Show_SystemNexusOperationWithCodec() { ctx, cancel := context.WithTimeout(s.Context, 60*time.Second) defer cancel() - callerWorkflowID := s.runSystemNexusSWSWorkflow(ctx) + input, err := converter.GetDefaultDataConverter().ToPayloads(map[string]string{"message": "codec-e2e-value"}) + s.NoError(err) + encodedPayloads, err := (prefixingCodec{}).Encode(input.Payloads) + s.NoError(err) + callerWorkflowID := s.runSystemNexusSWSWorkflow(ctx, &common.Payloads{Payloads: encodedPayloads}) + + var decodeRequests atomic.Int64 + codecHandler := converter.NewPayloadCodecHTTPHandler(prefixingCodec{}) + codecServer := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path == "/decode" { + decodeRequests.Add(1) + } + codecHandler.ServeHTTP(w, r) + })) + defer codecServer.Close() res := s.Execute( "workflow", "show", "--address", s.Address(), "-w", callerWorkflowID, + "--detailed", + "--codec-endpoint", codecServer.URL, ) s.NoError(res.Err) out := res.Stdout.String() s.Contains(out, "SignalWithStartWorkflowExecutionScheduled", "expected transformed Scheduled name in show output") s.Contains(out, "SignalWithStartWorkflowExecutionCompleted", "expected transformed Completed name in show output") + s.Contains(out, "unwrappedInput.input[0]") + s.Contains(out, "codec-e2e-value") s.NotContains(out, "NexusOperationScheduled", "raw event type name should be replaced by the unwrapped form") s.NotContains(out, "NexusOperationCompleted", "raw event type name should be replaced by the unwrapped form") + s.EqualValues(1, decodeRequests.Load(), "the gRPC interceptor should make the only codec-server request") } // TestWorkflow_Show_JSONOutputDoesNotUnwrapSystemNexus pins down that `workflow show -o json` @@ -1387,7 +1424,7 @@ func (s *SharedServerSuite) TestWorkflow_Show_JSONOutputDoesNotUnwrapSystemNexus ctx, cancel := context.WithTimeout(s.Context, 60*time.Second) defer cancel() - callerWorkflowID := s.runSystemNexusSWSWorkflow(ctx) + callerWorkflowID := s.runSystemNexusSWSWorkflow(ctx, nil) res := s.Execute( "workflow", "show", diff --git a/internal/temporalcli/commands_test.go b/internal/temporalcli/commands_test.go index 6b621ef18..2de546c50 100644 --- a/internal/temporalcli/commands_test.go +++ b/internal/temporalcli/commands_test.go @@ -242,7 +242,7 @@ func (s *SharedServerSuite) SetupSuite() { // Disable DescribeTaskQueue cache. "frontend.activityAPIsEnabled": true, "history.enableChasm": true, - // Required by TestWorkflow_Show_SystemNexusOperationTransformsTypeNames + // Required by TestWorkflow_Show_SystemNexusOperationWithCodec // to schedule a SignalWithStartWorkflowExecution Nexus operation against // the __temporal_system endpoint from inside a workflow. "history.enableSignalWithStartFromWorkflow": true,