Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 1 addition & 12 deletions chasm/lib/activity/activity.go
Original file line number Diff line number Diff line change
Expand Up @@ -1071,14 +1071,6 @@ func (a *Activity) unpause(
dispatchTime := a.unpauseDispatchTime(ctx, event)
attempt.DispatchTime = timestamppb.New(dispatchTime)

if event.req.GetResetAttempts() {
attempt.Count = 1
attempt.CurrentRetryInterval = nil
attempt.CurrentRetryIntervalSource = activitypb.ACTIVITY_RETRY_INTERVAL_SOURCE_UNSPECIFIED
}
if event.req.GetResetHeartbeat() {
a.LastHeartbeat = chasm.NewDataField(ctx, &activitypb.ActivityHeartbeatState{})
}
attempt.Stamp++
if timeout := a.GetScheduleToStartTimeout().AsDuration(); timeout > 0 {
ctx.AddTask(
Expand All @@ -1099,10 +1091,7 @@ func (a *Activity) unpauseDispatchTime(ctx chasm.MutableContext, event unpauseEv
unpauseTime = unpauseTime.Add(time.Duration(rand.Int63n(int64(jitter)))) //nolint:gosec
}
dispatchTime := a.dispatchTimeRespectingStartDelay(unpauseTime)
var retryDispatchTime *timestamppb.Timestamp
if !event.req.GetResetAttempts() {
retryDispatchTime = dispatchTimeForRetry(a.LastAttempt.Get(ctx))
}
retryDispatchTime := dispatchTimeForRetry(a.LastAttempt.Get(ctx))
if retryDispatchTime != nil && retryDispatchTime.AsTime().After(dispatchTime) {
return retryDispatchTime.AsTime()
}
Expand Down
8 changes: 3 additions & 5 deletions chasm/lib/activity/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -422,11 +422,9 @@ func (h *handler) UnpauseActivityExecution(ctx context.Context, req *activitypb.
WorkflowId: frontendReq.GetWorkflowId(),
RunId: frontendReq.GetRunId(),
},
Activity: &workflowservice.UnpauseActivityRequest_Id{Id: frontendReq.GetActivityId()},
Jitter: frontendReq.GetJitter(),
ResetAttempts: frontendReq.GetResetAttempts(),
ResetHeartbeat: frontendReq.GetResetHeartbeat(),
Identity: frontendReq.GetIdentity(),
Activity: &workflowservice.UnpauseActivityRequest_Id{Id: frontendReq.GetActivityId()},
Jitter: frontendReq.GetJitter(),
Identity: frontendReq.GetIdentity(),
Comment on lines +425 to +427

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not introduced here, but this change makes it the whole story for the API, so worth confirming it's intended: the two backends of UnpauseActivityExecution now disagree on retry-backoff timing.

  • Workflow activity (this branch): workflow.UnpauseActivity regenerates the retry task at now(+jitter) whenever the activity is SCHEDULED (service/history/workflow/activity.go:401-410), independent of the reset flags — so unpause drops the remaining backoff and dispatches immediately.
  • Standalone activity: unpauseDispatchTime now always consults dispatchTimeForRetry, so unpause waits out the remaining backoff. tests/activity_standalone_test.go:11095 asserts exactly that ("unpause must honor the remaining retry backoff").

Concretely: pause an activity sitting in a 30s retry backoff, then unpause. A workflow activity runs right away; a standalone activity waits for the original retry deadline. Previously reset_attempts=true gave the standalone path the workflow path's timing; with the flag gone, the only way to skip the backoff on a standalone activity is ResetActivityExecution, which also resets the attempt count and clears heartbeat details.

If the intent is that unpause is purely "resume", the legacy workflow path is the odd one out and may deserve a follow-up issue.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is something to fix in the workflow activity parity drive post saa release

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is intended: we will bring WFA into parity with SAA later. Neither have this API at GA.

},
})
if err != nil {
Expand Down
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ require (
go.opentelemetry.io/otel/sdk v1.43.0
go.opentelemetry.io/otel/sdk/metric v1.43.0
go.opentelemetry.io/otel/trace v1.44.0
go.temporal.io/api v1.63.5-0.20260731164320-beafdd4b0d8f
go.temporal.io/api v1.63.5-0.20260803183639-0e1e8c485f37
go.temporal.io/auto-scaled-workers v0.0.0-20260706201056-4320b34799ee
go.temporal.io/sdk v1.44.0
go.uber.org/fx v1.24.0
Expand Down
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -479,8 +479,8 @@ go.opentelemetry.io/proto/slim/otlp/collector/profiles/v1development v0.3.0 h1:R
go.opentelemetry.io/proto/slim/otlp/collector/profiles/v1development v0.3.0/go.mod h1:I89cynRj8y+383o7tEQVg2SVA6SRgDVIouWPUVXjx0U=
go.opentelemetry.io/proto/slim/otlp/profiles/v1development v0.3.0 h1:CQvJSldHRUN6Z8jsUeYv8J0lXRvygALXIzsmAeCcZE0=
go.opentelemetry.io/proto/slim/otlp/profiles/v1development v0.3.0/go.mod h1:xSQ+mEfJe/GjK1LXEyVOoSI1N9JV9ZI923X5kup43W4=
go.temporal.io/api v1.63.5-0.20260731164320-beafdd4b0d8f h1:n/cM2e930fSKURv5NSlvx0LaXEpyQY06Uo5MizOKYhU=
go.temporal.io/api v1.63.5-0.20260731164320-beafdd4b0d8f/go.mod h1:SrlW2JMwVlDP4nRWSNznUFqnSHd+YeMDS1BkYo63HCQ=
go.temporal.io/api v1.63.5-0.20260803183639-0e1e8c485f37 h1:ZDICI5Hxc97YsjpE6/WREWkUqQ7qq+ntTH+IbOuTxNw=
go.temporal.io/api v1.63.5-0.20260803183639-0e1e8c485f37/go.mod h1:SrlW2JMwVlDP4nRWSNznUFqnSHd+YeMDS1BkYo63HCQ=
go.temporal.io/auto-scaled-workers v0.0.0-20260706201056-4320b34799ee h1:y6A65Iml06cR3CpxW2Zn8FQLjniPDTnN4jtr69UZxXI=
go.temporal.io/auto-scaled-workers v0.0.0-20260706201056-4320b34799ee/go.mod h1:hhHijO9XRPIkAflLJJHix61M9FzbRPqk8fSydkcLkqw=
go.temporal.io/sdk v1.44.0 h1:suitPDukX74rW3/N1FqvEbZTZVJJsxMKhv0KMa/j7pU=
Expand Down
46 changes: 30 additions & 16 deletions tests/activity_api_pause_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,10 @@ import (
type activityPauseAPI struct {
name string
pause func(ctx context.Context, s *testcore.TestEnv, wfID, actID, identity, reason, requestID string) error
unpause func(ctx context.Context, s *testcore.TestEnv, wfID, actID, identity string, resetAttempts bool) error
unpause func(ctx context.Context, s *testcore.TestEnv, wfID, actID, identity string) error
// unpauseResettingAttempts is nil on an API with no reset_attempts flag. Only the deprecated
// UnpauseActivity has one; UnpauseActivityExecution deliberately does not.
Comment on lines 39 to +42

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Now that only one of the two APIs supports reset_attempts, the table gains a second func field, a nil check, a t.Skip, and a comment — all to express "this subtest only applies to the legacy API". Hoisting TestActivityPauseApi_WithReset out of the for _, api := range pauseAPIs() loop and having it call PauseActivity/UnpauseActivity directly would delete all four. The subtest body already builds its own env, workflow func and activity func; the only things it takes from api are pause and the unpause adapter, so nothing is shared that would need to be duplicated.

That also avoids the slightly odd shape where a subtest inside the UnpauseActivityExecution group exists only to be skipped.

Minor, on the comment itself: deliberately does not refers to the decision rather than the behavior. If you keep the field, something like // unpauseResettingAttempts is nil for APIs without a reset_attempts flag on unpause. says the same thing without it.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Declined. That would cause a large diff to the WFA test, and this PR is about SAA not WFA.

unpauseResettingAttempts func(ctx context.Context, s *testcore.TestEnv, wfID, actID, identity string) error
}

func pauseAPIs() []activityPauseAPI {
Expand All @@ -55,13 +58,22 @@ func pauseAPIs() []activityPauseAPI {
})
return err
},
unpause: func(ctx context.Context, s *testcore.TestEnv, wfID, actID, identity string, resetAttempts bool) error {
unpause: func(ctx context.Context, s *testcore.TestEnv, wfID, actID, identity string) error {
_, err := s.FrontendClient().UnpauseActivity(ctx, &workflowservice.UnpauseActivityRequest{
Namespace: s.Namespace().String(),
Execution: &commonpb.WorkflowExecution{WorkflowId: wfID},
Activity: &workflowservice.UnpauseActivityRequest_Id{Id: actID},
Identity: identity,
})
return err
},
unpauseResettingAttempts: func(ctx context.Context, s *testcore.TestEnv, wfID, actID, identity string) error {
_, err := s.FrontendClient().UnpauseActivity(ctx, &workflowservice.UnpauseActivityRequest{
Namespace: s.Namespace().String(),
Execution: &commonpb.WorkflowExecution{WorkflowId: wfID},
Activity: &workflowservice.UnpauseActivityRequest_Id{Id: actID},
Identity: identity,
ResetAttempts: resetAttempts,
ResetAttempts: true,
})
return err
},
Expand All @@ -79,13 +91,12 @@ func pauseAPIs() []activityPauseAPI {
})
return err
},
unpause: func(ctx context.Context, s *testcore.TestEnv, wfID, actID, identity string, resetAttempts bool) error {
unpause: func(ctx context.Context, s *testcore.TestEnv, wfID, actID, identity string) error {
_, err := s.FrontendClient().UnpauseActivityExecution(ctx, &workflowservice.UnpauseActivityExecutionRequest{
Namespace: s.Namespace().String(),
WorkflowId: wfID,
ActivityId: actID,
Identity: identity,
ResetAttempts: resetAttempts,
Namespace: s.Namespace().String(),
WorkflowId: wfID,
ActivityId: actID,
Identity: identity,
})
return err
},
Expand Down Expand Up @@ -211,7 +222,7 @@ func TestActivityApiPauseClientTestSuite(t *testing.T) {
s.Equal(testReason, description.PendingActivities[0].PauseInfo.GetManual().Reason)

// unpause the activity
require.NoError(t, api.unpause(ctx, s, workflowRun.GetID(), "activity-id", "", false))
require.NoError(t, api.unpause(ctx, s, workflowRun.GetID(), "activity-id", ""))

var out string
err = workflowRun.Get(ctx, &out)
Expand Down Expand Up @@ -329,7 +340,7 @@ func TestActivityApiPauseClientTestSuite(t *testing.T) {
shouldSucceed.Store(true)

// unpause the activity
require.NoError(t, api.unpause(ctx, s, workflowRun.GetID(), "activity-id", "", false))
require.NoError(t, api.unpause(ctx, s, workflowRun.GetID(), "activity-id", ""))

// wait for activity to complete
await.Require(t.Context(), t, func(t *await.T) {
Expand Down Expand Up @@ -425,7 +436,7 @@ func TestActivityApiPauseClientTestSuite(t *testing.T) {
s.Equal(testReason, description.PendingActivities[0].PauseInfo.GetManual().Reason)

// unpause the activity
require.NoError(t, api.unpause(ctx, s, workflowRun.GetID(), "activity-id", "", false))
require.NoError(t, api.unpause(ctx, s, workflowRun.GetID(), "activity-id", ""))

// wait for activity to complete
await.Require(t.Context(), t, func(t *await.T) {
Expand Down Expand Up @@ -505,7 +516,7 @@ func TestActivityApiPauseClientTestSuite(t *testing.T) {
require.NoError(t, api.pause(ctx, s, workflowRun.GetID(), "activity-id", "", "", testRequestID))

// unpause the activity
require.NoError(t, api.unpause(ctx, s, workflowRun.GetID(), "activity-id", "", false))
require.NoError(t, api.unpause(ctx, s, workflowRun.GetID(), "activity-id", ""))

// wait for activity to complete. It should happen immediately since noWait is set
await.Require(t.Context(), t, func(t *await.T) {
Expand All @@ -520,6 +531,9 @@ func TestActivityApiPauseClientTestSuite(t *testing.T) {

t.Run("TestActivityPauseApi_WithReset", func(t *testing.T) {
// pause/unpause the activity with reset option and noWait flag
if api.unpauseResettingAttempts == nil {
t.Skip("this API has no reset_attempts flag on unpause; Reset is the operation that restarts attempts")
}
s := testcore.NewEnv(t)

initialRetryInterval := 1 * time.Second
Expand Down Expand Up @@ -599,7 +613,7 @@ func TestActivityApiPauseClientTestSuite(t *testing.T) {
activityWasReset = true

// unpause the activity with reset
require.NoError(t, api.unpause(ctx, s, workflowRun.GetID(), "activity-id", "", true))
require.NoError(t, api.unpauseResettingAttempts(ctx, s, workflowRun.GetID(), "activity-id", ""))

// wait for activity to be running
await.Require(t.Context(), t, func(t *await.T) {
Expand Down Expand Up @@ -739,7 +753,7 @@ func TestActivityApiPauseClientTestSuite(t *testing.T) {
s.Equal(testReason, description.PendingActivities[0].PauseInfo.GetManual().Reason)

// unpause the activity
require.NoError(t, api.unpause(ctx, s, workflowRun.GetID(), "activity-id", "", false))
require.NoError(t, api.unpause(ctx, s, workflowRun.GetID(), "activity-id", ""))

var out string
err = workflowRun.Get(ctx, &out)
Expand Down Expand Up @@ -920,7 +934,7 @@ func TestActivityApiPauseClientTestSuite(t *testing.T) {

// step 4: unpause
activityWasReset.Store(true)
require.NoError(t, api.unpause(ctx, s, wfID, "activity-id", "", false))
require.NoError(t, api.unpause(ctx, s, wfID, "activity-id", ""))

await.Require(t.Context(), t, func(c *await.T) {
desc, err := s.SdkClient().DescribeWorkflowExecution(ctx, workflowRun.GetID(), workflowRun.GetRunID())
Expand Down
Loading
Loading