Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
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
14 changes: 14 additions & 0 deletions Makefile

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.

same comment as above, don't forget to remove. And no need for this in future if you don't use replace

Original file line number Diff line number Diff line change
Expand Up @@ -204,9 +204,17 @@ $(BUF): | $(LOCALBIN)

GO_API_VER = $(shell go list -m -f '{{.Version}}' go.temporal.io/api \
|| (echo "failed to fetch version for go.temporal.io/api" >&2))
GO_API_REPLACED = $(shell go list -m -f '{{if .Replace}}true{{end}}' go.temporal.io/api \
|| (echo "failed to resolve replacement for go.temporal.io/api" >&2))
PROTOGEN := $(LOCALBIN)/protogen-$(GO_API_VER)
ifeq ($(GO_API_REPLACED),true)
$(PROTOGEN): | $(LOCALBIN)
@printf $(COLOR) "Building protogen through the replaced go.temporal.io/api module..."
@go build -o $@ go.temporal.io/api/cmd/protogen
else
$(PROTOGEN): | $(LOCALBIN)
$(call go-install-tool,$(PROTOGEN),go.temporal.io/api/cmd/protogen,$(GO_API_VER))
endif

ACTIONLINT_VER := v1.7.7
ACTIONLINT := $(LOCALBIN)/actionlint-$(ACTIONLINT_VER)
Expand Down Expand Up @@ -285,8 +293,14 @@ $(STAMPDIR)/protoc-gen-go-grpc-$(PROTOC_GEN_GO_GRPC_VER): | $(STAMPDIR) $(LOCALB
$(PROTOC_GEN_GO_GRPC): $(STAMPDIR)/protoc-gen-go-grpc-$(PROTOC_GEN_GO_GRPC_VER)

PROTOC_GEN_GO_HELPERS := $(LOCALBIN)/protoc-gen-go-helpers-$(GO_API_VER)
ifeq ($(GO_API_REPLACED),true)
$(STAMPDIR)/protoc-gen-go-helpers-$(GO_API_VER): | $(STAMPDIR) $(LOCALBIN)
@printf $(COLOR) "Building protoc-gen-go-helpers through the replaced go.temporal.io/api module..."
@go build -o $(PROTOC_GEN_GO_HELPERS) go.temporal.io/api/cmd/protoc-gen-go-helpers
else
$(STAMPDIR)/protoc-gen-go-helpers-$(GO_API_VER): | $(STAMPDIR) $(LOCALBIN)
$(call go-install-tool,$(PROTOC_GEN_GO_HELPERS),go.temporal.io/api/cmd/protoc-gen-go-helpers,$(GO_API_VER))
endif
@touch $@
$(PROTOC_GEN_GO_HELPERS): $(STAMPDIR)/protoc-gen-go-helpers-$(GO_API_VER)

Expand Down
8 changes: 8 additions & 0 deletions chasm/lib/activity/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,12 @@ var (
`Allows non-zero start_delay on StartActivityExecution requests.`,
)

EnableEagerStart = dynamicconfig.NewNamespaceBoolSetting(
"activity.enableEagerStart",
true,
`Allows the first standalone activity task to be returned directly by StartActivityExecution.`,
)

EnableCallbacks = dynamicconfig.NewNamespaceBoolSetting(
"activity.enableCallbacks",
false,
Expand Down Expand Up @@ -63,6 +69,7 @@ type Config struct {
EnableCallbacks dynamicconfig.BoolPropertyFnWithNamespaceFilter
EnabledCallbackKinds dynamicconfig.TypedPropertyFnWithNamespaceFilter[[]callbacks.Kind]
Enabled dynamicconfig.BoolPropertyFnWithNamespaceFilter
EnableEagerStart dynamicconfig.BoolPropertyFnWithNamespaceFilter
EnableStandaloneActivityOperatorCommands dynamicconfig.BoolPropertyFnWithNamespaceFilter
LongPollBuffer dynamicconfig.DurationPropertyFnWithNamespaceFilter
LongPollTimeout dynamicconfig.DurationPropertyFnWithNamespaceFilter
Expand All @@ -85,6 +92,7 @@ func ConfigProvider(dc *dynamicconfig.Collection) *Config {
EnableCallbacks: EnableCallbacks.Get(dc),
EnabledCallbackKinds: EnabledCallbackKinds.Get(dc),
Enabled: Enabled.Get(dc),
EnableEagerStart: EnableEagerStart.Get(dc),
EnableStandaloneActivityOperatorCommands: EnableStandaloneActivityOperatorCommands.Get(dc),
LongPollBuffer: LongPollBuffer.Get(dc),
LongPollTimeout: LongPollTimeout.Get(dc),
Expand Down
24 changes: 23 additions & 1 deletion chasm/lib/activity/frontend.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(

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.

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(
Expand Down
81 changes: 81 additions & 0 deletions chasm/lib/activity/frontend_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,8 @@ import (
"go.temporal.io/api/workflowservice/v1"
"go.temporal.io/server/common/dynamicconfig"
"go.temporal.io/server/common/log"
"go.temporal.io/server/common/metrics"
"go.temporal.io/server/common/metrics/metricstest"
"go.temporal.io/server/common/namespace"
"google.golang.org/protobuf/types/known/durationpb"
)
Expand Down Expand Up @@ -215,3 +217,82 @@ func TestRequestIdStableAcrossRetries(t *testing.T) {
})
})
}

func TestEagerStartFallback(t *testing.T) {
newHandler := func(eagerEnabled bool, metricsHandler metrics.Handler) *frontendHandler {
return &frontendHandler{
config: &Config{
BlobSizeLimitError: defaultBlobSizeLimitError,
BlobSizeLimitWarn: defaultBlobSizeLimitWarn,
DefaultActivityRetryPolicy: getDefaultRetrySettings,
EnableEagerStart: dynamicconfig.GetBoolPropertyFnFilteredByNamespace(eagerEnabled),
MaxIDLengthLimit: func() int { return defaultMaxIDLengthLimit },
MaxUserMetadataDetailsSize: defaultMaxUserMetadataDetailsSize,
MaxUserMetadataSummarySize: defaultMaxUserMetadataSummarySize,
StartDelayEnabled: dynamicconfig.GetBoolPropertyFnFilteredByNamespace(true),
},
linkValidator: newLinkValidator(
defaultMaxLinksPerRequest,
func(string) int { return 2000 },
defaultLinkMaxSize,
),
logger: log.NewNoopLogger(),
metricsHandler: metricsHandler,
}
}

newRequest := func() *workflowservice.StartActivityExecutionRequest {
return &workflowservice.StartActivityExecutionRequest{
Namespace: "test-namespace",
RequestEagerExecution: true,
ActivityId: "test-activity",
ActivityType: &commonpb.ActivityType{Name: "test-type"},
TaskQueue: &taskqueuepb.TaskQueue{Name: "test-queue"},
StartToCloseTimeout: durationpb.New(time.Minute),
}
}

assertMetric := func(t *testing.T, capture *metricstest.Capture, metricName string, reason string) {
t.Helper()
recordings := capture.SnapshotMetric(metricName)
require.Len(t, recordings, 1)
if reason != "" {
require.Equal(t, reason, recordings[0].Tags["reason"])
}
}

t.Run("enabled", func(t *testing.T) {
metricsHandler := metricstest.NewCaptureHandler()
capture := metricsHandler.StartCapture()
defer metricsHandler.StopCapture(capture)

req, err := newHandler(true, metricsHandler).validateAndPopulateStartRequest(context.Background(), newRequest(), "test-namespace-id")
require.NoError(t, err)
require.True(t, req.GetRequestEagerExecution())
assertMetric(t, capture, metrics.StandaloneActivityEagerStartAcceptedCounter.Name(), "")
})

t.Run("namespace disabled", func(t *testing.T) {
metricsHandler := metricstest.NewCaptureHandler()
capture := metricsHandler.StartCapture()
defer metricsHandler.StopCapture(capture)

req, err := newHandler(false, metricsHandler).validateAndPopulateStartRequest(context.Background(), newRequest(), "test-namespace-id")
require.NoError(t, err)
require.False(t, req.GetRequestEagerExecution())
assertMetric(t, capture, metrics.StandaloneActivityEagerStartDeniedCounter.Name(), "dynamic_config_disabled")
})

t.Run("start delay", func(t *testing.T) {
metricsHandler := metricstest.NewCaptureHandler()
capture := metricsHandler.StartCapture()
defer metricsHandler.StopCapture(capture)

req := newRequest()
req.StartDelay = durationpb.New(time.Minute)
req, err := newHandler(true, metricsHandler).validateAndPopulateStartRequest(context.Background(), req, "test-namespace-id")
require.NoError(t, err)
require.False(t, req.GetRequestEagerExecution())
assertMetric(t, capture, metrics.StandaloneActivityEagerStartDeniedCounter.Name(), "start_delay")
})
}
32 changes: 28 additions & 4 deletions chasm/lib/activity/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down Expand Up @@ -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(

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.

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()
Expand Down Expand Up @@ -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{
Expand All @@ -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
}
Expand Down
4 changes: 4 additions & 0 deletions chasm/lib/activity/model/model.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,10 @@ type Outcome struct {
// Initial is the state of a newly created Activity.
func Initial(cfg Config) AbstractState {
s := AbstractState{Status: Scheduled, AttemptCount: 1}
if cfg.InitialAttemptStarted {
s.Status = Started
return s
}
if cfg.HasStartDelay {
s.Dispatchability = StartDelayPending
}
Expand Down
26 changes: 25 additions & 1 deletion chasm/lib/activity/model/model_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,27 @@ func TestInitial(t *testing.T) {
require.Equal(t, AbstractState{Status: Scheduled, AttemptCount: 1}, Initial(Config{HasScheduleToClose: true}))
}

func TestEagerInitialAttempt(t *testing.T) {
cfg := Config{
InitialAttemptStarted: true,
HasScheduleToClose: true,
HasScheduleToStart: true,
HasHeartbeat: true,
}
initial := Initial(cfg)
require.Equal(t, AbstractState{Status: Started, AttemptCount: 1}, initial)
require.False(t, FindsTask(initial), "an eager task is already started and must not be dispatched again")
require.False(t, Possible(cfg, initial, ScheduleToStartElapsesType), "an eager first attempt has no schedule-to-start window")
require.True(t, Possible(cfg, initial, ScheduleToCloseElapsesType), "schedule-to-close still protects an eager first attempt")
require.True(t, Possible(cfg, initial, StartToCloseElapsesType), "start-to-close still protects an eager first attempt")
require.True(t, Possible(cfg, initial, HeartbeatElapsesType), "heartbeat still protects an eager first attempt")

retry := Transition(cfg, initial, FailRetryably).Next
require.Equal(t, AbstractState{Status: Scheduled, AttemptCount: 2, Dispatchability: BackoffPending}, retry)
retry = Transition(cfg, retry, BackoffElapses).Next
require.True(t, FindsTask(retry), "an eager first attempt's retry must return to normal dispatch")
}

func TestPollFromScheduledStarts(t *testing.T) {
out := Transition(Config{}, Initial(Config{}), Event{Type: PollType})
require.Equal(t, NoError, out.Reject)
Expand Down Expand Up @@ -40,7 +61,10 @@ func TestPauseWhileStartedIsPauseRequested(t *testing.T) {
// backedOffRetry returns a Scheduled state with a pending retry backoff (attempt 2), reached the way
// a worker would: poll the first attempt, then fail it retryably.
func backedOffRetry(t require.TestingT, cfg Config) AbstractState {
started := Transition(cfg, Initial(cfg), Event{Type: PollType}).Next
started := Initial(cfg)
if started.Status != Started {
started = Transition(cfg, started, Event{Type: PollType}).Next
}
s := Transition(cfg, started, Event{Type: RespondFailedType, Failure: &Failure{}}).Next
require.Equal(t, Scheduled, s.Status, "a retryable failure must schedule a retry")
require.Equal(t, BackoffPending, s.Dispatchability, "a retry must wait for its backoff")
Expand Down
12 changes: 7 additions & 5 deletions chasm/lib/activity/model/vocabulary.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,11 +43,13 @@ type AbstractState struct {
// Config is what the model needs to know about an activity's configuration: which options are set,
// and what their durations imply regarding retries.
type Config struct {
HasScheduleToClose bool
HasScheduleToStart bool
HasHeartbeat bool
HasStartDelay bool
MaxAttempts int32 // 0 = unlimited
// InitialAttemptStarted means the first attempt was delivered eagerly and is already running.
InitialAttemptStarted bool
HasScheduleToClose bool
HasScheduleToStart bool
HasHeartbeat bool
HasStartDelay bool
MaxAttempts int32 // 0 = unlimited

// NonRetryableTimeouts are the timeout elapses whose failure the retry policy refuses to retry,
// so that the timeout closes the activity instead of scheduling another attempt.
Expand Down
63 changes: 63 additions & 0 deletions chasm/lib/activity/responses.go
Original file line number Diff line number Diff line change
Expand Up @@ -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},

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.

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.
Expand Down
Loading
Loading