go.temporal.io/server/components/nexusoperations/completion.go

283 LOC · 135 covered · 148 uncovered · 36 ranges · 81 concepts · 10 introducers · 20 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/components/nexusoperations/metrics.go · 116 LOCnexusoperations/metrics.…go.temporal.io/server/service/history/hsm/tree.go · 821 LOChsm/tree.goworkflow_handler.go ×11 · 117 introduced LOCworkflow_handler.go ×11common.go ×1 · 11 introduced LOCcommon.go ×1request_response.pb.go ×6 · 246 introduced LOCrequest_response.pb.go ×…visibility_store.go ×17 · 171 introduced LOCvisibility_store.go ×17telemetry.go ×2 · 15 introduced LOCtelemetry.go ×2fx.go ×1 · 4 introduced LOCfx.go ×1collector.go ×7 · 33 introduced LOCcollector.go ×7data_store_factory.go ×29 · 703 introduced LOCdata_store_factory.go ×2…TestNewServer · 0 introduced LOCTestNewServerpri_matcher.go ×1 · 2 introduced LOCpri_matcher.go ×1timer_queue_active_task_executor.go ×1 · 3 introduced LOCtimer_queue_active_task_…TestNewServer · 0 introduced LOCTestNewServermetric_client.go ×3 · 18 introduced LOCmetric_client.go ×3request_response.pb.go ×12 · 137 introduced LOCrequest_response.pb.go ×…workflow_task_completed_handler.go ×9 · 75 introduced LOCworkflow_task_completed_…pri_forwarder.go ×2 · 4 introduced LOCpri_forwarder.go ×2request_response.pb.go ×1 · 7 introduced LOCrequest_response.pb.go ×…logger.go ×1 · 5 introduced LOClogger.go ×1connections.go ×1 · 9 introduced LOCconnections.go ×1persistence_rate_limited_clients.go ×2 · 17 introduced LOCpersistence_rate_limited…logger.go ×1 · 2 introduced LOClogger.go ×1metric_client.go ×2 · 7 introduced LOCmetric_client.go ×2pri_matcher.go ×1 · 4 introduced LOCpri_matcher.go ×1task_queue_partition_manager.go ×2 · 10 introduced LOCtask_queue_partition_man…connections.go ×2 · 12 introduced LOCconnections.go ×2workflow_handler.go ×8 · 144 introduced LOCworkflow_handler.go ×8pri_matcher.go ×1 · 2 introduced LOCpri_matcher.go ×1handler.go ×1 · 20 introduced LOChandler.go ×1pri_matcher.go ×8 · 55 introduced LOCpri_matcher.go ×8matching_engine.go ×1 · 8 introduced LOCmatching_engine.go ×1server.go ×1 · 3 introduced LOCserver.go ×1logger.go ×2 · 21 introduced LOClogger.go ×2queue_scheduled.go ×1 · 2 introduced LOCqueue_scheduled.go ×1handler.go ×25 · 728 introduced LOChandler.go ×25onebox.go ×75 · 1256 introduced LOConebox.go ×75metric_client_gen.go ×4 · 50 introduced LOCmetric_client_gen.go ×4pri_forwarder.go ×1 · 6 introduced LOCpri_forwarder.go ×1matching_service_server_gen.go ×1 · 2 introduced LOCmatching_service_server_…http_api_server.go ×23 · 195 introduced LOChttp_api_server.go ×23db.go ×1 · 2 introduced LOCdb.go ×1pri_task_writer.go ×3 · 14 introduced LOCpri_task_writer.go ×3namespace_handover.go ×3 · 8 introduced LOCnamespace_handover.go ×3task_queue_partition_manager.go ×2 · 4 introduced LOCtask_queue_partition_man…matching_engine.go ×3 · 5 introduced LOCmatching_engine.go ×3service_grpc.pb.go ×19 · 357 introduced LOCservice_grpc.pb.go ×19request_response.pb.go ×6 · 114 introduced LOCrequest_response.pb.go ×…reader.go ×2 · 8 introduced LOCreader.go ×2service_grpc.pb.go ×20 · 740 introduced LOCservice_grpc.pb.go ×20endpoint_registry.go ×2 · 26 introduced LOCendpoint_registry.go ×2server.go ×3 · 13 introduced LOCserver.go ×3scanner.go ×1 · 2 introduced LOCscanner.go ×1mask_internal_error.go ×1 · 2 introduced LOCmask_internal_error.go ×…matching_engine.go ×2 · 5 introduced LOCmatching_engine.go ×2service_resolver.go ×4 · 49 introduced LOCservice_resolver.go ×4lite_server.go ×25 · 303 introduced LOClite_server.go ×25fx.go ×44 · 705 introduced LOCfx.go ×44adaptive_pool.go ×1 · 3 introduced LOCadaptive_pool.go ×1queue_immediate.go ×1 · 1 introduced LOCqueue_immediate.go ×1fx.go ×1 · 1 introduced LOCfx.go ×1service.go ×8 · 413 introduced LOCservice.go ×8rpc.go ×1 · 13 introduced LOCrpc.go ×1queue_scheduled.go ×1 · 1 introduced LOCqueue_scheduled.go ×1fx.go ×1 · 2 introduced LOCfx.go ×1fx.go ×44 · 4693 introduced LOCfx.go ×44executors.go ×1 · 4 introduced LOCexecutors.go ×1sync_failed · 0 introduced LOCsync_failedexecutors.go ×2 · 9 introduced LOCexecutors.go ×2executors.go ×3 · 10 introduced LOCexecutors.go ×3executors.go ×1 · 3 introduced LOCexecutors.go ×1executors.go ×3 · 6 introduced LOCexecutors.go ×3executors.go ×3 · 9 introduced LOCexecutors.go ×3completion.go ×2 · 12 introduced LOCcompletion.go ×2completion.go ×1 · 2 introduced LOCcompletion.go ×1completion.go ×1 · 1 introduced LOCcompletion.go ×1completion.go ×2 · 4 introduced LOCcompletion.go ×2completion.go ×18 · 57 introduced LOCcompletion.go ×18completion.go ×3 · 17 introduced LOCcompletion.go ×3completion.go ×2 · 13 introduced LOCcompletion.go ×2completion.go ×2 · 19 introduced LOCcompletion.go ×2completion.go ×4 · 15 introduced LOCcompletion.go ×4completion.go ×1 · 3 introduced LOCcompletion.go ×1canceled · introduced test · go.temporal.io/server/components/nexusoperations/TestCompletionHandler_EmitsCallerMetrics/canceledcanceledfailed · introduced test · go.temporal.io/server/components/nexusoperations/TestCompletionHandler_EmitsCallerMetrics/failedfailedsucceeded · introduced test · go.temporal.io/server/components/nexusoperations/TestCompletionHandler_EmitsCallerMetrics/succeededsucceededsync_canceled · introduced test · go.temporal.io/server/components/nexusoperations/TestProcessInvocationTask/sync_canceledsync_canceledsync_failed · introduced test · go.temporal.io/server/components/nexusoperations/TestProcessInvocationTask/sync_failedsync_failedsync_start · introduced test · go.temporal.io/server/components/nexusoperations/TestProcessInvocationTask/sync_startsync_startoperation_error · introduced test · go.temporal.io/server/components/nexusoperations/TestProcessInvocationTask_SystemEndpoint/operation_erroroperation_errorsync_start · introduced test · go.temporal.io/server/components/nexusoperations/TestProcessInvocationTask_SystemEndpoint/sync_startsync_startTestNewServer · introduced test · go.temporal.io/server/temporal/TestNewServerTestNewServerTestNewServerWithJSONEncoding · introduced test · go.temporal.io/server/temporal/TestNewServerWithJSONEncodingTestNewServerWithJSONEnc…with_OTEL_Collector_running · introduced test · go.temporal.io/server/temporal/TestNewServerWithOTEL/with_OTEL_Collector_runningwith_OTEL_Collector_runn…without_OTEL_Collector_running · introduced test · go.temporal.io/server/temporal/TestNewServerWithOTEL/without_OTEL_Collector_runningwithout_OTEL_Collector_r…ExampleNewServer · introduced test · go.temporal.io/server/temporaltest/ExampleNewServerExampleNewServerTestBaseServerOptions · introduced test · go.temporal.io/server/temporaltest/TestBaseServerOptionsTestBaseServerOptionsTestClientWithCustomInterceptor · introduced test · go.temporal.io/server/temporaltest/TestClientWithCustomInterceptorTestClientWithCustomInte…TestDefaultWorkerOptions · introduced test · go.temporal.io/server/temporaltest/TestDefaultWorkerOptionsTestDefaultWorkerOptionsTestNewServer · introduced test · go.temporal.io/server/temporaltest/TestNewServerTestNewServerTestNewWorkerWithOptions · introduced test · go.temporal.io/server/temporaltest/TestNewWorkerWithOptionsTestNewWorkerWithOptionsTestSearchAttributeRegistration · introduced test · go.temporal.io/server/temporaltest/TestSearchAttributeRegistrationTestSearchAttributeRegis…TestWorkerServiceHealthCheck · introduced test · go.temporal.io/server/tests/testcore/TestFunctionalTestBaseSuite/TestWorkerServiceHealthCheckTestWorkerServiceHealthC…Focused file · go.temporal.io/server/components/nexusoperations/completion.go · 283 LOCnexusoperations/completi…

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 nexusoperations
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "time"
8
9 "github.com/nexus-rpc/sdk-go/nexus"
10 commonpb "go.temporal.io/api/common/v1"
11 enumspb "go.temporal.io/api/enums/v1"
12 failurepb "go.temporal.io/api/failure/v1"
13 historypb "go.temporal.io/api/history/v1"
14 "go.temporal.io/api/serviceerror"
15 "go.temporal.io/server/common/metrics"
16 commonnexus "go.temporal.io/server/common/nexus"
17 "go.temporal.io/server/service/history/hsm"
18 "google.golang.org/protobuf/types/known/timestamppb"
19 )
20
21 func handleSuccessfulOperationResult(
22 node *hsm.Node,
23 operation Operation,
24 result *commonpb.Payload,
25 links []*commonpb.Link,
26 > ) error { completion.go ×2
27 > eventID, err := hsm.EventIDFromToken(operation.ScheduledEventToken)
28 > if err != nil {
29 return err
30 }
31 > event := node.AddHistoryEvent(enumspb.EVENT_TYPE_NEXUS_OPERATION_COMPLETED, func(e *historypb.HistoryEvent) { completion.go ×2
32 > // We must assign to this property, linter doesn't like this.
33 > // nolint:revive
34 > e.Attributes = &historypb.HistoryEvent_NexusOperationCompletedEventAttributes{
35 > NexusOperationCompletedEventAttributes: &historypb.NexusOperationCompletedEventAttributes{
36 > ScheduledEventId: eventID,
37 > Result: result,
38 > RequestId: operation.RequestId,
39 > },
40 > }
41 > e.Links = links
42 > })
43 > return CompletedEventDefinition{}.Apply(node.Parent, event)
44 }
45
46 func handleOperationError(
47 node *hsm.Node,
48 operation Operation,
49 opErr *nexus.OperationError,
50 > ) error { completion.go ×4
51 > eventID, err := hsm.EventIDFromToken(operation.ScheduledEventToken)
52 > if err != nil {
53 return err
54 }
55 > var originalCause *failurepb.Failure completion.go ×4
56 > // Special marker for Temporal->Temporal calls to indicate that the original failure should be unwrapped.
57 > // Temporal uses a wrapper operation error with no additional information to transmit the OperationError over the network.
58 > // The meaningful information is in the operation error's cause.
59 > unwrapError := opErr.OriginalFailure.Metadata["unwrap-error"] == "true"
60 >
61 > if unwrapError && opErr.OriginalFailure.Cause != nil {
62 var err error
63 originalCause, err = commonnexus.NexusFailureToTemporalFailure(*opErr.OriginalFailure.Cause)
64 if err != nil {
65 return serviceerror.NewInvalidArgumentf("Malformed failure: %v", err)
66 }
67 > } else { completion.go ×4
68 > // Transform the OperationError to either ApplicationFailure or CanceledFailure based on the operation error state.
69 > originalCause, err = commonnexus.NexusFailureToTemporalFailure(*opErr.OriginalFailure)
70 > if err != nil {
71 return serviceerror.NewInvalidArgumentf("Malformed failure: %v", err)
72 }
73 }
74
75 > switch opErr.State { // nolint:exhaustive completion.go ×4
76 > case nexus.OperationStateFailed: completion.go ×2
77 > event := node.AddHistoryEvent(enumspb.EVENT_TYPE_NEXUS_OPERATION_FAILED, func(e *historypb.HistoryEvent) {
78 > // We must assign to this property, linter doesn't like this.
79 > // nolint:revive
80 > e.Attributes = &historypb.HistoryEvent_NexusOperationFailedEventAttributes{
81 > NexusOperationFailedEventAttributes: &historypb.NexusOperationFailedEventAttributes{
82 > Failure: createNexusOperationFailure(operation, eventID, originalCause),
83 > ScheduledEventId: eventID,
84 > RequestId: operation.RequestId,
85 > },
86 > }
87 > })
88
89 > return FailedEventDefinition{}.Apply(node.Parent, event) completion.go ×2
90 > case nexus.OperationStateCanceled: completion.go ×3
91 > if originalCause.GetCanceledFailureInfo() == nil {
92 > // Old SDKs may send an ApplicationFailure for canceled operation causes. completion.go ×2
93 > originalCause = &failurepb.Failure{
94 > Message: originalCause.GetMessage(),
95 > StackTrace: originalCause.GetStackTrace(),
96 > FailureInfo: &failurepb.Failure_CanceledFailureInfo{
97 > CanceledFailureInfo: &failurepb.CanceledFailureInfo{},
98 > },
99 > Cause: originalCause.GetCause(),
100 > }
101 > }
102 > event := node.AddHistoryEvent(enumspb.EVENT_TYPE_NEXUS_OPERATION_CANCELED, func(e *historypb.HistoryEvent) { completion.go ×3
103 > // We must assign to this property, linter doesn't like this.
104 > // nolint:revive
105 > e.Attributes = &historypb.HistoryEvent_NexusOperationCanceledEventAttributes{
106 > NexusOperationCanceledEventAttributes: &historypb.NexusOperationCanceledEventAttributes{
107 > Failure: createNexusOperationFailure(operation, eventID, originalCause),
108 > ScheduledEventId: eventID,
109 > RequestId: operation.RequestId,
110 > },
111 > }
112 > })
113
114 > return CanceledEventDefinition{}.Apply(node.Parent, event) completion.go ×3
115 default:
116 // Both the Nexus Client and CompletionHandler reject invalid states, but just in case, we return this as a
117 // transition error.
118 return fmt.Errorf("unexpected operation state: %v", opErr.State)
119 }
120 }
121
122 // fabricateStartedEventIfMissing adds a NEXUS_OPERATION_STARTED history event and transitions the
123 // operation to NEXUS_OPERATION_STATE_STARTED if it has not started yet. It is necessary if the
124 // completion is received before the start response. Callers that need to know whether a start was
125 // fabricated should check the operation's state before calling (see TransitionStarted.Possible).
126 func fabricateStartedEventIfMissing(
127 node *hsm.Node,
128 requestID string,
129 operationToken string,
130 startTime *timestamppb.Timestamp,
131 links []*commonpb.Link,
132 > ) error { completion.go ×18
133 > operation, err := hsm.MachineData[Operation](node)
134 > if err != nil {
135 return err
136 }
137
138 // The operation was already started, ignore.
139 > if !TransitionStarted.Possible(operation) { completion.go ×18
140 return nil
141 }
142
143 > eventID, err := hsm.EventIDFromToken(operation.ScheduledEventToken) completion.go ×18
144 > if err != nil {
145 return err
146 }
147
148 > event := node.AddHistoryEvent(enumspb.EVENT_TYPE_NEXUS_OPERATION_STARTED, func(e *historypb.HistoryEvent) { completion.go ×18
149 > e.Attributes = &historypb.HistoryEvent_NexusOperationStartedEventAttributes{
150 > NexusOperationStartedEventAttributes: &historypb.NexusOperationStartedEventAttributes{
151 > ScheduledEventId: eventID,
152 > OperationToken: operationToken,
153 > // TODO(bergundy): Remove this fallback after the 1.27 release.
154 > OperationId: operationToken,
155 > RequestId: requestID,
156 > },
157 > }
158 > e.Links = links
159 > if startTime != nil {
160 e.EventTime = startTime
161 }
162 })
163 > return (StartedEventDefinition{}).Apply(node.Parent, event) completion.go ×18
164 }
165
166 // CompletionHandler resolves async Nexus operation completions delivered to the history service
167 // and emits the caller-side terminal metrics. It is provided via fx and injected into the history
168 // handler so the metrics handler and tag config are sourced from the dependency graph rather than
169 // threaded through the call site.
170 type CompletionHandler struct {
171 metricsHandler metrics.Handler
172 config *Config
173 }
174
175 // NewCompletionHandler returns a CompletionHandler. Wired via fx; see Module.
176 > func NewCompletionHandler(metricsHandler metrics.Handler, config *Config) *CompletionHandler { completion.go ×1
177 > return &CompletionHandler{metricsHandler: metricsHandler, config: config}
178 > }
179
180 // Handle resolves an async Nexus operation completion.
181 func (h *CompletionHandler) Handle(
182 ctx context.Context,
183 env hsm.Environment,
184 ref hsm.Ref,
185 requestID string,
186 operationToken string,
187 startTime *timestamppb.Timestamp,
188 links []*commonpb.Link,
189 result *commonpb.Payload,
190 opFailedError *nexus.OperationError,
191 > ) error { completion.go ×18
192 > // The initial version of the completion token did not include a request ID.
193 > // Only retry Access without a run ID if the request ID is not empty.
194 > isRetryableNotFoundErr := requestID != ""
195 > // emitMetrics is populated inside the Access closure and invoked only after the write
196 > // transaction commits successfully, so a failed commit (which retries the completion) does not
197 > // double-count the metric. See operationMetricsHandler's doc comment.
198 > var emitMetrics func()
199 > err := env.Access(ctx, ref, hsm.AccessWrite, func(node *hsm.Node) error {
200 > if err := node.CheckRunning(); err != nil {
201 return err
202 }
203 > operation, err := hsm.MachineData[Operation](node) completion.go ×18
204 > if err != nil {
205 return err
206 }
207 // If the operation has not started yet, this completion arrived before the start response and
208 // fabricateStartedEventIfMissing will transition it to started below; the executor never emitted
209 // schedule-to-start in that case, so we emit it here.
210 > fabricatedStart := TransitionStarted.Possible(operation) completion.go ×18
211 > if err := fabricateStartedEventIfMissing(node, requestID, operationToken, startTime, links); err != nil {
212 return err
213 }
214 > if requestID != "" && operation.RequestId != requestID { completion.go ×18
215 isRetryableNotFoundErr = false
216 return serviceerror.NewNotFound("operation not found")
217 }
218
219 > if opFailedError != nil { completion.go ×18
220 > err = handleOperationError(node, operation, opFailedError) completion.go ×1
221 > } else { completion.go ×18
222 > err = handleSuccessfulOperationResult(node, operation, result, nil) completion.go ×2
223 > }
224 // TODO(bergundy): Remove this once the operation auto-deletes itself from the tree on completion with state
225 // based replication.
226 > if errors.Is(err, hsm.ErrInvalidTransition) { completion.go ×18
227 isRetryableNotFoundErr = false
228 return serviceerror.NewNotFound("operation not found")
229 }
230 > if err != nil { completion.go ×18
231 return err
232 }
233 // fabricatedStart means the executor never emitted schedule-to-start, so we emit it here.
234 > emitScheduleToStart := fabricatedStart && operation.StartedTime != nil completion.go ×18
235 > emitMetrics = h.deferredCompletionMetric(operation, node.NamespaceName(), node.WorkflowTypeName(), opFailedError, emitScheduleToStart, env.Now())
236 > return nil
237 })
238 > if errors.As(err, new(*serviceerror.NotFound)) && isRetryableNotFoundErr && ref.WorkflowKey.RunID != "" { completion.go ×18
239 // Try again without a run ID in case the original run was reset.
240 ref.WorkflowKey.RunID = ""
241 // VersionedTransition is for a specific run. After reset, the TransitionCount will
242 // start from 1 again. Reset the TransitionCount to 0 here to fallback to old ref
243 // validation logic.
244 ref.StateMachineRef.MutableStateVersionedTransition = nil
245 ref.StateMachineRef.MachineInitialVersionedTransition.TransitionCount = 0
246 ref.StateMachineRef.MachineLastUpdateVersionedTransition.TransitionCount = 0
247 return h.Handle(ctx, env, ref, requestID, operationToken, startTime, links, result, opFailedError)
248 }
249 > if err != nil { completion.go ×18
250 return err
251 }
252 > if emitMetrics != nil { completion.go ×18
253 > emitMetrics()
254 > }
255 > return nil
256 }
257
258 // deferredCompletionMetric builds the post-commit caller-side emit for a resolved async completion:
259 // the terminal outcome (the canceled vs failed split mirrors handleOperationError so the metric
260 // matches the recorded transition) plus, when the start was fabricated here, schedule-to-start.
261 func (h *CompletionHandler) deferredCompletionMetric(
262 operation Operation,
263 namespaceName, workflowType string,
264 opFailedError *nexus.OperationError,
265 emitScheduleToStart bool,
266 closeTime time.Time,
267 > ) func() { completion.go ×18
268 > metricsHandler := h.metricsHandler
269 > metricTagConfig := h.config.ResolvedMetricTagConfig()
270 > return func() {
271 > if emitScheduleToStart {
272 > emitScheduleToStartLatency(metricsHandler, metricTagConfig, operation, namespaceName, workflowType, operation.StartedTime.AsTime())
273 > }
274 > switch {
275 > case opFailedError == nil: completion.go ×2
276 > emitOperationSucceeded(metricsHandler, metricTagConfig, operation, namespaceName, workflowType, closeTime)
277 > case opFailedError.State == nexus.OperationStateCanceled: completion.go ×2
278 > emitOperationCanceled(metricsHandler, metricTagConfig, operation, namespaceName, workflowType, closeTime)
279 > default: completion.go ×1
280 > emitOperationFailed(metricsHandler, metricTagConfig, operation, namespaceName, workflowType, closeTime)
281 }
282 }
283 }