Repository navigation
Add eager start support for standalone activities #12349
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
28e9572
02b75cd
aa1e032
96da079
0803b98
dccc1b0
7ff222a
5735e5c
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -42,6 +42,11 @@ var ErrStandaloneActivityDisabled = serviceerror.NewUnimplemented("Standalone ac | |
|
|
||
| var ErrStandaloneActivityOperatorCommandsDisabled = serviceerror.NewUnimplemented("Standalone activity operator commands are disabled") | ||
|
|
||
| const ( | ||
| eagerStartDeniedReasonDynamicConfigDisabled metrics.ReasonString = "dynamic_config_disabled" | ||
| eagerStartDeniedReasonStartDelay metrics.ReasonString = "start_delay" | ||
| ) | ||
|
|
||
| type frontendHandler struct { | ||
| FrontendHandler | ||
| callbackValidator callbacks.Validator | ||
|
|
@@ -386,7 +391,24 @@ func (h *frontendHandler) validateAndPopulateStartRequest( | |
| if req.GetStartDelay().AsDuration() > 0 && !h.config.StartDelayEnabled(req.GetNamespace()) { | ||
| return nil, serviceerror.NewInvalidArgument("start_delay is not enabled for this namespace") | ||
| } | ||
| // TODO(saa): when eager start is supported, deny it if start delay > 0 (same as workflow behavior). | ||
| if req.GetRequestEagerExecution() { | ||
| metricsHandler := h.metricsHandler.WithTags( | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. activity_eager_execution currently uses WFA’s namespace/task-queue metric scope, while this new SAA scope uses namespace/operation. Reusing the same metric with different label keys can cause Prometheus to reject the series. We addressed this WFA/SAA parity problem previously in #11806 by sharing standard Activity metric-scope construction,by adding empty SAA labels when WFA had labels SAA could not yet populate. Could we follow that pattern here: extract the current WFA activity_eager_execution scope into a shared common/metrics helper, migrate WFA to use it, then use it from SAA after the eager task is successfully constructed? This would also let us reuse ActivityEagerExecutionCounter instead of introducing standalone_activity_eager_start_accepted. Please add a WFA/SAA label-parity test. |
||
| metrics.NamespaceTag(req.GetNamespace()), | ||
| metrics.OperationTag("StartActivityExecution"), | ||
| ) | ||
| switch { | ||
| case !h.config.EnableEagerStart(req.GetNamespace()): | ||
| metrics.StandaloneActivityEagerStartDeniedCounter.With(metricsHandler). | ||
| Record(1, metrics.ReasonTag(eagerStartDeniedReasonDynamicConfigDisabled)) | ||
| req.RequestEagerExecution = false | ||
| case req.GetStartDelay().AsDuration() > 0: | ||
| metrics.StandaloneActivityEagerStartDeniedCounter.With(metricsHandler). | ||
| Record(1, metrics.ReasonTag(eagerStartDeniedReasonStartDelay)) | ||
| req.RequestEagerExecution = false | ||
| default: | ||
| metrics.StandaloneActivityEagerStartAcceptedCounter.With(metricsHandler).Record(1) | ||
| } | ||
| } | ||
|
|
||
| opts := activityOptionsFromStartRequest(req) | ||
| err := ValidateAndNormalizeStandaloneActivity( | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -99,7 +99,14 @@ func (h *handler) StartActivityExecution(ctx context.Context, req *activitypb.St | |
| } | ||
| } | ||
|
|
||
| err = TransitionScheduled.Apply(newActivity, mutableContext, nil) | ||
| if request.GetRequestEagerExecution() { | ||
| err = TransitionEagerStarted.Apply(newActivity, mutableContext, eagerStartEvent{ | ||
| requestID: request.GetRequestId(), | ||
| identity: request.GetIdentity(), | ||
| }) | ||
| } else { | ||
| err = TransitionScheduled.Apply(newActivity, mutableContext, nil) | ||
| } | ||
| if err != nil { | ||
| return nil, err | ||
| } | ||
|
|
@@ -129,6 +136,23 @@ func (h *handler) StartActivityExecution(ctx context.Context, req *activitypb.St | |
| ) | ||
| } | ||
|
|
||
| var eagerTask *workflowservice.PollActivityTaskQueueResponse | ||
| if result.Created && frontendReq.GetRequestEagerExecution() { | ||
| eagerTask, err = chasm.ReadComponent( | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. we shouldn't be doing a read here as it's not atomic with the startexec. We should move population of the eager task into the StartExecution startFn so it's within the transaction |
||
| ctx, | ||
| result.ExecutionRef, | ||
| (*Activity).buildEagerActivityTask, | ||
| eagerActivityTaskRequest{ | ||
| namespaceID: req.GetNamespaceId(), | ||
| namespace: frontendReq.GetNamespace(), | ||
| requestID: frontendReq.GetRequestId(), | ||
| }, | ||
| ) | ||
| if err != nil { | ||
| return nil, err | ||
| } | ||
| } | ||
|
|
||
| // Apply on_conflict_options to an existing activity. | ||
| // TODO: Use chasm.UpdateWithStartExecution to avoid a second transaction once the engine supports BusinessIDConflictPolicyFail in the updateFn path. | ||
| cbs := frontendReq.GetCompletionCallbacks() | ||
|
|
@@ -165,8 +189,9 @@ func (h *handler) StartActivityExecution(ctx context.Context, req *activitypb.St | |
|
|
||
| return &activitypb.StartActivityExecutionResponse{ | ||
| FrontendResponse: &workflowservice.StartActivityExecutionResponse{ | ||
| RunId: result.ExecutionKey.RunID, | ||
| Started: result.Created, | ||
| RunId: result.ExecutionKey.RunID, | ||
| Started: result.Created, | ||
| EagerActivityTask: eagerTask, | ||
| Link: &commonpb.Link{ | ||
| Variant: &commonpb.Link_Activity_{ | ||
| Activity: &commonpb.Link_Activity{ | ||
|
|
@@ -176,7 +201,6 @@ func (h *handler) StartActivityExecution(ctx context.Context, req *activitypb.St | |
| }, | ||
| }, | ||
| }, | ||
| // EagerTask: TODO when supported, need to call the same code that would handle the HandleStarted API | ||
| }, | ||
| }, nil | ||
| } | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -10,11 +10,74 @@ import ( | |
| "go.temporal.io/api/workflowservice/v1" | ||
| "go.temporal.io/server/chasm" | ||
| "go.temporal.io/server/chasm/lib/activity/gen/activitypb/v1" | ||
| "go.temporal.io/server/common/tasktoken" | ||
| "google.golang.org/protobuf/proto" | ||
| "google.golang.org/protobuf/types/known/durationpb" | ||
| "google.golang.org/protobuf/types/known/timestamppb" | ||
| ) | ||
|
|
||
| type eagerActivityTaskRequest struct { | ||
| namespaceID string | ||
| namespace string | ||
| requestID string | ||
| } | ||
|
|
||
| func (a *Activity) buildEagerActivityTask( | ||
| ctx chasm.Context, | ||
| request eagerActivityTaskRequest, | ||
| ) (*workflowservice.PollActivityTaskQueueResponse, error) { | ||
| attempt := a.LastAttempt.Get(ctx) | ||
| if !a.hasAttemptInProgress() || attempt.GetCount() != 1 || attempt.GetStartRequestId() != request.requestID { | ||
| return nil, nil | ||
| } | ||
|
|
||
| componentRef, err := ctx.Ref(a) | ||
| if err != nil { | ||
| return nil, err | ||
| } | ||
| key := ctx.ExecutionKey() | ||
| token, err := tasktoken.NewSerializer().Serialize(tasktoken.NewActivityTaskToken( | ||
| request.namespaceID, | ||
| "", | ||
| key.RunID, | ||
| 0, | ||
| key.BusinessID, | ||
| a.GetActivityType().GetName(), | ||
| attempt.GetCount(), | ||
| nil, | ||
| 0, | ||
| 0, | ||
| componentRef, | ||
| attempt.GetStartedStamp(), | ||
| )) | ||
| if err != nil { | ||
| return nil, err | ||
| } | ||
|
|
||
| requestData := a.RequestData.Get(ctx) | ||
| lastHeartbeat, _ := a.LastHeartbeat.TryGet(ctx) | ||
| return &workflowservice.PollActivityTaskQueueResponse{ | ||
| TaskToken: token, | ||
| WorkflowNamespace: request.namespace, | ||
| WorkflowExecution: &commonpb.WorkflowExecution{RunId: key.RunID}, | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. no need to set this. SAA has no parent wf |
||
| ActivityType: a.GetActivityType(), | ||
| ActivityId: key.BusinessID, | ||
| Header: requestData.GetHeader(), | ||
| Input: requestData.GetInput(), | ||
| HeartbeatDetails: lastHeartbeat.GetDetails(), | ||
| ScheduledTime: a.GetScheduleTime(), | ||
| CurrentAttemptScheduledTime: a.dispatchTimeForAttempt(attempt), | ||
| StartedTime: attempt.GetStartedTime(), | ||
| Attempt: attempt.GetCount(), | ||
| ScheduleToCloseTimeout: a.GetScheduleToCloseTimeout(), | ||
| StartToCloseTimeout: a.GetStartToCloseTimeout(), | ||
| HeartbeatTimeout: a.GetHeartbeatTimeout(), | ||
| RetryPolicy: a.GetRetryPolicy(), | ||
| Priority: a.GetPriority(), | ||
| ActivityRunId: key.RunID, | ||
| }, nil | ||
| } | ||
|
|
||
| // Projection of activity state onto the API response protos. | ||
|
|
||
| // InternalStatusToAPIStatus converts internal activity execution status to API status. | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
same comment as above, don't forget to remove. And no need for this in future if you don't use replace