go.temporal.io/server/chasm/nexus_completion.go

83 LOC · 30 covered · 53 uncovered · 8 ranges · 43 concepts · 4 introducers · 17 tests

File neighbourhood

The centred file is linked to every concept that introduces one of its ranges, every test that runs code from the file, and the gray connector concepts standing between those tests and the file's own introducer concepts. Undirected links join concepts to every file where they introduce source and concepts to the tests they introduce; arrows show specialization between the displayed concepts and bridge only concepts omitted from this view. Concept colors match the source ranges below; connector concepts have no source color and are shown in gray.

Focused file, its introducer and connector concepts, their introduced files, and tests that run code from the file

In the embedded map, ordinary wheel input scrolls the page; use the visible controls to zoom and drag to pan. Open the full-screen map for canvas navigation: wheel pans, Ctrl/Command plus wheel zooms, and arrow keys pan when this region is focused. On touch screens, open the full-screen map to pan or pinch. If JavaScript or WebGL is unavailable, use the related-file, concept, and source links on this page.

Focused file, its introducer and connector concepts, their introduced files, and tests that run code from the filego.temporal.io/server/chasm/lib/scheduler/invoker_tasks.go · 807 LOCscheduler/invoker_tasks.…go.temporal.io/server/chasm/lib/scheduler/scheduler.go · 1081 LOCscheduler/scheduler.gogo.temporal.io/server/common/testing/mockapi/workflowservicemock/v1/service_grpc.pb.mock.go · 2483 LOCv1/service_grpc.pb.mock.…TestExecuteTask_RetryableFailure · 0 introduced LOCTestExecuteTask_Retryabl…TestExecuteTask_BackoffUsesFrameworkClock · 0 introduced LOCTestExecuteTask_BackoffU…invoker_tasks.go ×5 · 13 introduced LOCinvoker_tasks.go ×5TestExecuteTask_DistinctRequestsCanReuseCompletedWorkflowID · 0 introduced LOCTestExecuteTask_Distinct…invoker_tasks.go ×1 · 1 introduced LOCinvoker_tasks.go ×1TestExecuteTask_Basic · 0 introduced LOCTestExecuteTask_Basicinvoker_tasks.go ×1 · 2 introduced LOCinvoker_tasks.go ×1invoker_tasks.go ×1 · 9 introduced LOCinvoker_tasks.go ×1TestExecuteTask_AlreadyStarted · 0 introduced LOCTestExecuteTask_AlreadyS…invoker_tasks.go ×3 · 14 introduced LOCinvoker_tasks.go ×3invoker_tasks.go ×1 · 2 introduced LOCinvoker_tasks.go ×1invoker_tasks.go ×7 · 89 introduced LOCinvoker_tasks.go ×7invocable_internal.go ×2 · 7 introduced LOCinvocable_internal.go ×2invocable_internal.go ×3 · 15 introduced LOCinvocable_internal.go ×3invocable_internal.go ×2 · 3 introduced LOCinvocable_internal.go ×2invocable_internal.go ×4 · 10 introduced LOCinvocable_internal.go ×4success-with-successful-operation · 0 introduced LOCsuccess-with-successful-…invocable_internal.go ×2 · 5 introduced LOCinvocable_internal.go ×2invocable_internal.go ×2 · 12 introduced LOCinvocable_internal.go ×2invocable_internal.go ×1 · 2 introduced LOCinvocable_internal.go ×1invocable_internal.go ×4 · 14 introduced LOCinvocable_internal.go ×4invocable_internal.go ×2 · 5 introduced LOCinvocable_internal.go ×2invocable_internal.go ×1 · 2 introduced LOCinvocable_internal.go ×1tasks.go ×1 · 3 introduced LOCtasks.go ×1invocable_internal.go ×2 · 7 introduced LOCinvocable_internal.go ×2invocable_internal.go ×4 · 21 introduced LOCinvocable_internal.go ×4chasm_invocation.go ×2 · 4 introduced LOCchasm_invocation.go ×2chasm_invocation.go ×3 · 15 introduced LOCchasm_invocation.go ×3chasm_invocation.go ×2 · 3 introduced LOCchasm_invocation.go ×2chasm_invocation.go ×4 · 10 introduced LOCchasm_invocation.go ×4success-with-successful-operation · 0 introduced LOCsuccess-with-successful-…chasm_invocation.go ×1 · 4 introduced LOCchasm_invocation.go ×1chasm_invocation.go ×2 · 12 introduced LOCchasm_invocation.go ×2chasm_invocation.go ×5 · 17 introduced LOCchasm_invocation.go ×5chasm_invocation.go ×1 · 2 introduced LOCchasm_invocation.go ×1chasm_invocation.go ×1 · 5 introduced LOCchasm_invocation.go ×1chasm_invocation.go ×4 · 21 introduced LOCchasm_invocation.go ×4request_response.pb.go ×1 · 5 introduced LOCrequest_response.pb.go ×…request_response.pb.go ×1 · 5 introduced LOCrequest_response.pb.go ×…service_grpc.pb.mock.go ×3 · 13 introduced LOCservice_grpc.pb.mock.go …nexus_completion.go ×2 · 4 introduced LOCnexus_completion.go ×2nexus_completion.go ×1 · 2 introduced LOCnexus_completion.go ×1nexus_completion.go ×1 · 3 introduced LOCnexus_completion.go ×1invalid-base64-header · introduced test · go.temporal.io/server/chasm/lib/callback/TestExecuteInvocationTaskChasm_Outcomes/invalid-base64-headerinvalid-base64-headerinvalid-protobuf-in-ref · introduced test · go.temporal.io/server/chasm/lib/callback/TestExecuteInvocationTaskChasm_Outcomes/invalid-protobuf-in-refinvalid-protobuf-in-refnon-retryable-rpc-error · introduced test · go.temporal.io/server/chasm/lib/callback/TestExecuteInvocationTaskChasm_Outcomes/non-retryable-rpc-errornon-retryable-rpc-errorretryable-rpc-error · introduced test · go.temporal.io/server/chasm/lib/callback/TestExecuteInvocationTaskChasm_Outcomes/retryable-rpc-errorretryable-rpc-errorsuccess-with-failed-operation · introduced test · go.temporal.io/server/chasm/lib/callback/TestExecuteInvocationTaskChasm_Outcomes/success-with-failed-operationsuccess-with-failed-oper…success-with-successful-operation · introduced test · go.temporal.io/server/chasm/lib/callback/TestExecuteInvocationTaskChasm_Outcomes/success-with-successful-operationsuccess-with-successful-…TestExecuteTask_AlreadyStarted · introduced test · go.temporal.io/server/chasm/lib/scheduler/TestExecuteTask_AlreadyStartedTestExecuteTask_AlreadyS…TestExecuteTask_BackoffUsesFrameworkClock · introduced test · go.temporal.io/server/chasm/lib/scheduler/TestExecuteTask_BackoffUsesFrameworkClockTestExecuteTask_BackoffU…TestExecuteTask_Basic · introduced test · go.temporal.io/server/chasm/lib/scheduler/TestExecuteTask_BasicTestExecuteTask_BasicTestExecuteTask_DistinctRequestsCanReuseCompletedWorkflowID · introduced test · go.temporal.io/server/chasm/lib/scheduler/TestExecuteTask_DistinctRequestsCanReuseCompletedWorkflowIDTestExecuteTask_Distinct…TestExecuteTask_ExceedsMaxActionsPerExecution · introduced test · go.temporal.io/server/chasm/lib/scheduler/TestExecuteTask_ExceedsMaxActionsPerExecutionTestExecuteTask_ExceedsM…TestExecuteTask_RetryableFailure · introduced test · go.temporal.io/server/chasm/lib/scheduler/TestExecuteTask_RetryableFailureTestExecuteTask_Retryabl…invalid-base64-header · introduced test · go.temporal.io/server/components/callbacks/TestProcessInvocationTaskChasm_Outcomes/invalid-base64-headerinvalid-base64-headernon-retryable-rpc-error · introduced test · go.temporal.io/server/components/callbacks/TestProcessInvocationTaskChasm_Outcomes/non-retryable-rpc-errornon-retryable-rpc-errorretryable-rpc-error · introduced test · go.temporal.io/server/components/callbacks/TestProcessInvocationTaskChasm_Outcomes/retryable-rpc-errorretryable-rpc-errorsuccess-with-failed-operation · introduced test · go.temporal.io/server/components/callbacks/TestProcessInvocationTaskChasm_Outcomes/success-with-failed-operationsuccess-with-failed-oper…success-with-successful-operation · introduced test · go.temporal.io/server/components/callbacks/TestProcessInvocationTaskChasm_Outcomes/success-with-successful-operationsuccess-with-successful-…Focused file · go.temporal.io/server/chasm/nexus_completion.go · 83 LOCchasm/nexus_completion.g…

Graph controls are ready.

Interactive rendering requires JavaScript and WebGL. Use the related-file, concept, and source links on this page while the interactive map is unavailable.

1 package chasm
2
3 import (
4 "encoding/base64"
5
6 commonpb "go.temporal.io/api/common/v1"
7 persistencespb "go.temporal.io/server/api/persistence/v1"
8 tokenspb "go.temporal.io/server/api/token/v1"
9 "google.golang.org/protobuf/proto"
10 )
11
12 // NexusCompletionHandlerURL is the user-visible URL for Nexus->CHASM callbacks.
13 const NexusCompletionHandlerURL = "temporal://internal"
14
15 // nexusCallbackTokenHeader is the callback header key carrying the completion token.
16 // NOTE: There's a constant defined for this in common/nexus but to avoid circular dependencies we
17 // redefine it here. nexus.Header lookups are case-insensitive, so this matches the canonical key.
18 const nexusCallbackTokenHeader = "temporal-callback-token"
19
20 // NexusCompletionHandler is implemented by CHASM components that want to handle Nexus operation completion callbacks.
21 type NexusCompletionHandler interface {
22 HandleNexusCompletion(ctx MutableContext, completion *persistencespb.ChasmNexusCompletion) error
23 }
24
25 // GenerateNexusCallback builds a Nexus completion callback targeting the CHASM component identified by
26 // serializedRef (obtained from Context.Ref). When encodeToken is true, the request ID is packed into
27 // the callback token (a NexusOperationCompletion envelope), so the completion is matched by a request
28 // ID that rides in the callback header and survives continue-as-new, rather than one read from mutable
29 // state. When encodeToken is false, the legacy format is emitted: the token is the bare base64-encoded
30 // ChasmComponentRef with no request ID. The caller chooses encodeToken (e.g. gated behind dynamic
31 // config) to keep the envelope format off the wire until the whole fleet can read it. Either format is
32 // always decodable by UnpackNexusCallbackToken.
33 > func GenerateNexusCallback(serializedRef []byte, requestID string, encodeToken bool) (*commonpb.Callback, error) { invoker_tasks.go ×7
34 > var token string
35 > if encodeToken {
36 > var err error
37 > token, err = packNexusCallbackToken(serializedRef, requestID)
38 > if err != nil {
39 return nil, err
40 }
41 } else {
42 // Legacy format: the token is the bare base64-encoded ChasmComponentRef.
43 token = base64.RawURLEncoding.EncodeToString(serializedRef)
44 }
45 > return &commonpb.Callback{ invoker_tasks.go ×7
46 > Variant: &commonpb.Callback_Nexus_{
47 > Nexus: &commonpb.Callback_Nexus{
48 > Url: NexusCompletionHandlerURL,
49 > Header: map[string]string{nexusCallbackTokenHeader: token},
50 > },
51 > },
52 > }, nil
53 }
54
55 // packNexusCallbackToken encodes a CHASM component ref and request ID into a callback token.
56 > func packNexusCallbackToken(componentRef []byte, requestID string) (string, error) { invoker_tasks.go ×7
57 > b, err := proto.Marshal(&tokenspb.NexusOperationCompletion{
58 > ComponentRef: componentRef,
59 > RequestId: requestID,
60 > })
61 > if err != nil {
62 return "", err
63 }
64 > return base64.RawURLEncoding.EncodeToString(b), nil invoker_tasks.go ×7
65 }
66
67 // UnpackNexusCallbackToken decodes a callback token produced by GenerateNexusCallback, returning the
68 // component ref and request ID. It accepts both token formats regardless of how the token was written:
69 // the NexusOperationCompletion envelope, and (for backward compatibility) the legacy bare base64-encoded
70 // ChasmComponentRef, in which case the request ID is empty.
71 > func UnpackNexusCallbackToken(encoded string) (componentRef []byte, requestID string, err error) { nexus_completion.go ×1
72 > raw, err := base64.RawURLEncoding.DecodeString(encoded)
73 > if err != nil {
74 > return nil, "", err nexus_completion.go ×1
75 > }
76 > completion := &tokenspb.NexusOperationCompletion{} nexus_completion.go ×2
77 > if proto.Unmarshal(raw, completion) == nil && len(completion.GetComponentRef()) > 0 &&
78 > proto.Unmarshal(completion.GetComponentRef(), &persistencespb.ChasmComponentRef{}) == nil {
79 return completion.GetComponentRef(), completion.GetRequestId(), nil
80 }
81 // Legacy format: the raw bytes are the ChasmComponentRef directly.
82 > return raw, "", nil nexus_completion.go ×2
83 }