From f01d672a4c010cb1f3c8dfbf36ee3418a224d15a Mon Sep 17 00:00:00 2001 From: chaptersix <13949480+chaptersix@users.noreply.github.com> Date: Fri, 4 Sep 2026 22:16:09 -0500 Subject: [PATCH] Project schedule executions and activity scheduling metadata into visibility --- chasm/lib/activity/activity.go | 80 ++++++++++++----- chasm/lib/activity/activity_test.go | 29 ++++++ .../gen/activitypb/v1/activity_state.pb.go | 89 +++++++++++-------- .../activity/proto/v1/activity_state.proto | 4 + chasm/lib/scheduler/action_activity.go | 2 +- chasm/lib/scheduler/library.go | 3 + chasm/lib/scheduler/scheduler.go | 48 ++++++---- chasm/search_attribute.go | 1 + 8 files changed, 180 insertions(+), 76 deletions(-) diff --git a/chasm/lib/activity/activity.go b/chasm/lib/activity/activity.go index 9360fbe479a..dbe3dbf1eb2 100644 --- a/chasm/lib/activity/activity.go +++ b/chasm/lib/activity/activity.go @@ -3,6 +3,8 @@ package activity import ( "errors" "fmt" + "maps" + "time" "github.com/nexus-rpc/sdk-go/nexus" apiactivitypb "go.temporal.io/api/activity/v1" //nolint:importas @@ -23,7 +25,9 @@ import ( "go.temporal.io/server/common/metrics" commonnexus "go.temporal.io/server/common/nexus" "go.temporal.io/server/common/nexus/nexusrpc" + "go.temporal.io/server/common/payload" "go.temporal.io/server/common/retrypolicy" + "go.temporal.io/server/common/searchattribute/sadefs" serviceerrors "go.temporal.io/server/common/serviceerror" "go.temporal.io/server/service/history/consts" "google.golang.org/protobuf/types/known/timestamppb" @@ -160,35 +164,56 @@ func NewStandaloneActivity( ctx chasm.MutableContext, request *workflowservice.StartActivityExecutionRequest, ) (*Activity, error) { + searchAttributes := maps.Clone(request.GetSearchAttributes().GetIndexedFields()) + var scheduledByID string + if value := searchAttributes[sadefs.TemporalScheduledById]; value != nil { + if err := payload.Decode(value, &scheduledByID); err != nil { + return nil, fmt.Errorf("decode scheduled-by ID: %w", err) + } + delete(searchAttributes, sadefs.TemporalScheduledById) + } + var scheduledStartTime time.Time + if value := searchAttributes[sadefs.TemporalScheduledStartTime]; value != nil { + if err := payload.Decode(value, &scheduledStartTime); err != nil { + return nil, fmt.Errorf("decode scheduled start time: %w", err) + } + delete(searchAttributes, sadefs.TemporalScheduledStartTime) + } visibility := chasm.NewVisibilityWithData( ctx, - request.GetSearchAttributes().GetIndexedFields(), + searchAttributes, nil, ) - activity := &Activity{ - ActivityState: &activitypb.ActivityState{ - ActivityType: request.ActivityType, - TaskQueue: request.GetTaskQueue(), - ScheduleToCloseTimeout: request.GetScheduleToCloseTimeout(), - ScheduleToStartTimeout: request.GetScheduleToStartTimeout(), - StartToCloseTimeout: request.GetStartToCloseTimeout(), - HeartbeatTimeout: request.GetHeartbeatTimeout(), - RetryPolicy: request.GetRetryPolicy(), - Priority: request.Priority, - StartDelay: request.GetStartDelay(), - OriginalOptions: &apiactivitypb.ActivityOptions{ - TaskQueue: common.CloneProto(request.GetTaskQueue()), - ScheduleToCloseTimeout: common.CloneProto(request.GetScheduleToCloseTimeout()), - ScheduleToStartTimeout: common.CloneProto(request.GetScheduleToStartTimeout()), - StartToCloseTimeout: common.CloneProto(request.GetStartToCloseTimeout()), - HeartbeatTimeout: common.CloneProto(request.GetHeartbeatTimeout()), - RetryPolicy: common.CloneProto(request.GetRetryPolicy()), - Priority: common.CloneProto(request.GetPriority()), - StartDelay: common.CloneProto(request.GetStartDelay()), - }, + state := &activitypb.ActivityState{ + ActivityType: request.ActivityType, + TaskQueue: request.GetTaskQueue(), + ScheduleToCloseTimeout: request.GetScheduleToCloseTimeout(), + ScheduleToStartTimeout: request.GetScheduleToStartTimeout(), + StartToCloseTimeout: request.GetStartToCloseTimeout(), + HeartbeatTimeout: request.GetHeartbeatTimeout(), + RetryPolicy: request.GetRetryPolicy(), + Priority: request.Priority, + StartDelay: request.GetStartDelay(), + OriginalOptions: &apiactivitypb.ActivityOptions{ + TaskQueue: common.CloneProto(request.GetTaskQueue()), + ScheduleToCloseTimeout: common.CloneProto(request.GetScheduleToCloseTimeout()), + ScheduleToStartTimeout: common.CloneProto(request.GetScheduleToStartTimeout()), + StartToCloseTimeout: common.CloneProto(request.GetStartToCloseTimeout()), + HeartbeatTimeout: common.CloneProto(request.GetHeartbeatTimeout()), + RetryPolicy: common.CloneProto(request.GetRetryPolicy()), + Priority: common.CloneProto(request.GetPriority()), + StartDelay: common.CloneProto(request.GetStartDelay()), }, - LastAttempt: chasm.NewDataField(ctx, &activitypb.ActivityAttemptState{}), + } + state.ScheduledById = scheduledByID + if !scheduledStartTime.IsZero() { + state.ScheduledStartTime = timestamppb.New(scheduledStartTime) + } + + activity := &Activity{ + ActivityState: state, + LastAttempt: chasm.NewDataField(ctx, &activitypb.ActivityAttemptState{}), RequestData: chasm.NewDataField(ctx, &activitypb.ActivityRequestData{ Input: request.Input, Header: request.Header, @@ -700,10 +725,17 @@ func (a *Activity) validateActivityTaskToken( // SearchAttributes implements chasm.VisibilitySearchAttributesProvider interface. // Returns the current search attribute values for this activity execution. func (a *Activity) SearchAttributes(_ chasm.Context) []chasm.SearchAttributeKeyValue { - return []chasm.SearchAttributeKeyValue{ + attributes := []chasm.SearchAttributeKeyValue{ TypeSearchAttribute.Value(a.GetActivityType().GetName()), StatusSearchAttribute.Value(InternalStatusToAPIStatus(a.GetStatus()).String()), chasm.SearchAttributeTaskQueue.Value(a.GetTaskQueue().GetName()), chasm.SearchAttributeExecutionTime.Value(a.firstDispatchTime()), } + if a.GetScheduledById() != "" { + attributes = append(attributes, chasm.SearchAttributeTemporalScheduledByID.Value(a.GetScheduledById())) + } + if scheduledStartTime := a.GetScheduledStartTime(); scheduledStartTime != nil { + attributes = append(attributes, chasm.SearchAttributeTemporalScheduledStartTime.Value(scheduledStartTime.AsTime())) + } + return attributes } diff --git a/chasm/lib/activity/activity_test.go b/chasm/lib/activity/activity_test.go index df704edda0b..61f73fdcb41 100644 --- a/chasm/lib/activity/activity_test.go +++ b/chasm/lib/activity/activity_test.go @@ -72,6 +72,35 @@ func TestSearchAttributesIncludesExecutionTime(t *testing.T) { } } +func TestNewStandaloneActivity_ExtractsScheduledMetadata(t *testing.T) { + now := time.Date(2000, 1, 1, 0, 0, 0, 0, time.UTC) + ctx := &chasm.MockMutableContext{MockContext: chasm.MockContext{ + HandleNow: func(chasm.Component) time.Time { return now }, + }} + activity, err := NewStandaloneActivity(ctx, &workflowservice.StartActivityExecutionRequest{ + Namespace: "ns", + ActivityId: "act", + ActivityType: &commonpb.ActivityType{Name: "T"}, + TaskQueue: &taskqueuepb.TaskQueue{Name: "tq"}, + StartToCloseTimeout: durationpb.New(time.Minute), + SearchAttributes: &commonpb.SearchAttributes{IndexedFields: map[string]*commonpb.Payload{ + sadefs.TemporalScheduledById: chasm.SearchAttributeTemporalScheduledByID.Value("schedule-id").Value.MustEncode(), + sadefs.TemporalScheduledStartTime: chasm.SearchAttributeTemporalScheduledStartTime.Value(now). + Value.MustEncode(), + }}, + }) + require.NoError(t, err) + require.Equal(t, "schedule-id", activity.GetScheduledById()) + require.Equal(t, now, activity.GetScheduledStartTime().AsTime()) + + attributes := make(map[string]any) + for _, attribute := range activity.SearchAttributes(ctx) { + attributes[attribute.Field] = attribute.Value.Value() + } + require.Equal(t, "schedule-id", attributes[sadefs.TemporalScheduledById]) + require.Equal(t, now, attributes[sadefs.TemporalScheduledStartTime]) +} + func TestHandleStarted(t *testing.T) { testTime := time.Date(2000, 1, 1, 0, 0, 0, 0, time.UTC) testRequestID := "test-request-id" diff --git a/chasm/lib/activity/gen/activitypb/v1/activity_state.pb.go b/chasm/lib/activity/gen/activitypb/v1/activity_state.pb.go index 5e588d2b142..4be2de91815 100644 --- a/chasm/lib/activity/gen/activitypb/v1/activity_state.pb.go +++ b/chasm/lib/activity/gen/activitypb/v1/activity_state.pb.go @@ -242,8 +242,7 @@ type ActivityState struct { // retries will be attempted. Either this or `start_to_close_timeout` must be specified. // // (-- api-linter: core::0140::prepositions=disabled - // - // aip.dev/not-precedent: "to" is used to indicate interval. --) + // aip.dev/not-precedent: "to" is used to indicate interval. --) ScheduleToCloseTimeout *durationpb.Duration `protobuf:"bytes,3,opt,name=schedule_to_close_timeout,json=scheduleToCloseTimeout,proto3" json:"schedule_to_close_timeout,omitempty"` // Limits time an activity task can stay in a task queue before a worker picks it up. This // timeout is always non retryable, as all a retry would achieve is to put it back into the same @@ -251,16 +250,14 @@ type ActivityState struct { // specified. // // (-- api-linter: core::0140::prepositions=disabled - // - // aip.dev/not-precedent: "to" is used to indicate interval. --) + // aip.dev/not-precedent: "to" is used to indicate interval. --) ScheduleToStartTimeout *durationpb.Duration `protobuf:"bytes,4,opt,name=schedule_to_start_timeout,json=scheduleToStartTimeout,proto3" json:"schedule_to_start_timeout,omitempty"` // Maximum time an activity is allowed to execute after being picked up by a worker. This // timeout is always retryable. Either this or `schedule_to_close_timeout` must be // specified. // // (-- api-linter: core::0140::prepositions=disabled - // - // aip.dev/not-precedent: "to" is used to indicate interval. --) + // aip.dev/not-precedent: "to" is used to indicate interval. --) StartToCloseTimeout *durationpb.Duration `protobuf:"bytes,5,opt,name=start_to_close_timeout,json=startToCloseTimeout,proto3" json:"start_to_close_timeout,omitempty"` // Maximum permitted time between successful worker heartbeats. HeartbeatTimeout *durationpb.Duration `protobuf:"bytes,6,opt,name=heartbeat_timeout,json=heartbeatTimeout,proto3" json:"heartbeat_timeout,omitempty"` @@ -319,8 +316,11 @@ type ActivityState struct { LastResetRequestId string `protobuf:"bytes,22,opt,name=last_reset_request_id,json=lastResetRequestId,proto3" json:"last_reset_request_id,omitempty"` // Used to de-dupe update requests. LastUpdateOptionsRequestId string `protobuf:"bytes,23,opt,name=last_update_options_request_id,json=lastUpdateOptionsRequestId,proto3" json:"last_update_options_request_id,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + // Scheduling metadata is stored separately from namespace-defined attributes. + ScheduledById string `protobuf:"bytes,24,opt,name=scheduled_by_id,json=scheduledById,proto3" json:"scheduled_by_id,omitempty"` + ScheduledStartTime *timestamppb.Timestamp `protobuf:"bytes,25,opt,name=scheduled_start_time,json=scheduledStartTime,proto3" json:"scheduled_start_time,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *ActivityState) Reset() { @@ -514,6 +514,20 @@ func (x *ActivityState) GetLastUpdateOptionsRequestId() string { return "" } +func (x *ActivityState) GetScheduledById() string { + if x != nil { + return x.ScheduledById + } + return "" +} + +func (x *ActivityState) GetScheduledStartTime() *timestamppb.Timestamp { + if x != nil { + return x.ScheduledStartTime + } + return nil +} + type ActivityCancelState struct { state protoimpl.MessageState `protogen:"open.v1"` RequestId string `protobuf:"bytes,1,opt,name=request_id,json=requestId,proto3" json:"request_id,omitempty"` @@ -1246,7 +1260,7 @@ var File_temporal_server_chasm_lib_activity_proto_v1_activity_state_proto protor const file_temporal_server_chasm_lib_activity_proto_v1_activity_state_proto_rawDesc = "" + "\n" + - "@temporal/server/chasm/lib/activity/proto/v1/activity_state.proto\x12+temporal.server.chasm.lib.activity.proto.v1\x1a\x1egoogle/protobuf/duration.proto\x1a\x1fgoogle/protobuf/timestamp.proto\x1a&temporal/api/activity/v1/message.proto\x1a$temporal/api/common/v1/message.proto\x1a(temporal/api/deployment/v1/message.proto\x1a$temporal/api/enums/v1/workflow.proto\x1a%temporal/api/failure/v1/message.proto\x1a'temporal/api/sdk/v1/user_metadata.proto\x1a'temporal/api/taskqueue/v1/message.proto\"\xb9\r\n" + + "@temporal/server/chasm/lib/activity/proto/v1/activity_state.proto\x12+temporal.server.chasm.lib.activity.proto.v1\x1a\x1egoogle/protobuf/duration.proto\x1a\x1fgoogle/protobuf/timestamp.proto\x1a&temporal/api/activity/v1/message.proto\x1a$temporal/api/common/v1/message.proto\x1a(temporal/api/deployment/v1/message.proto\x1a$temporal/api/enums/v1/workflow.proto\x1a%temporal/api/failure/v1/message.proto\x1a'temporal/api/sdk/v1/user_metadata.proto\x1a'temporal/api/taskqueue/v1/message.proto\"\xaf\x0e\n" + "\rActivityState\x12I\n" + "\ractivity_type\x18\x01 \x01(\v2$.temporal.api.common.v1.ActivityTypeR\factivityType\x12C\n" + "\n" + @@ -1273,7 +1287,9 @@ const file_temporal_server_chasm_lib_activity_proto_v1_activity_state_proto_rawD "\x15reset_restore_options\x18\x14 \x01(\bR\x13resetRestoreOptions\x125\n" + "\x17last_unpause_request_id\x18\x15 \x01(\tR\x14lastUnpauseRequestId\x121\n" + "\x15last_reset_request_id\x18\x16 \x01(\tR\x12lastResetRequestId\x12B\n" + - "\x1elast_update_options_request_id\x18\x17 \x01(\tR\x1alastUpdateOptionsRequestId\"\xa7\x01\n" + + "\x1elast_update_options_request_id\x18\x17 \x01(\tR\x1alastUpdateOptionsRequestId\x12&\n" + + "\x0fscheduled_by_id\x18\x18 \x01(\tR\rscheduledById\x12L\n" + + "\x14scheduled_start_time\x18\x19 \x01(\v2\x1a.google.protobuf.TimestampR\x12scheduledStartTime\"\xa7\x01\n" + "\x13ActivityCancelState\x12\x1d\n" + "\n" + "request_id\x18\x01 \x01(\tR\trequestId\x12=\n" + @@ -1409,32 +1425,33 @@ var file_temporal_server_chasm_lib_activity_proto_v1_activity_state_proto_depIdx 19, // 13: temporal.server.chasm.lib.activity.proto.v1.ActivityState.original_options:type_name -> temporal.api.activity.v1.ActivityOptions 5, // 14: temporal.server.chasm.lib.activity.proto.v1.ActivityState.last_pause_state:type_name -> temporal.server.chasm.lib.activity.proto.v1.ActivityPauseState 17, // 15: temporal.server.chasm.lib.activity.proto.v1.ActivityState.first_attempt_started_time:type_name -> google.protobuf.Timestamp - 17, // 16: temporal.server.chasm.lib.activity.proto.v1.ActivityCancelState.request_time:type_name -> google.protobuf.Timestamp - 17, // 17: temporal.server.chasm.lib.activity.proto.v1.ActivityPauseState.pause_time:type_name -> google.protobuf.Timestamp - 15, // 18: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.current_retry_interval:type_name -> google.protobuf.Duration - 17, // 19: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.started_time:type_name -> google.protobuf.Timestamp - 17, // 20: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.complete_time:type_name -> google.protobuf.Timestamp - 10, // 21: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.last_failure_details:type_name -> temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.LastFailureDetails - 20, // 22: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.last_deployment_version:type_name -> temporal.api.deployment.v1.WorkerDeploymentVersion - 17, // 23: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.dispatch_time:type_name -> google.protobuf.Timestamp - 1, // 24: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.current_retry_interval_source:type_name -> temporal.server.chasm.lib.activity.proto.v1.ActivityRetryIntervalSource - 21, // 25: temporal.server.chasm.lib.activity.proto.v1.ActivityHeartbeatState.details:type_name -> temporal.api.common.v1.Payloads - 17, // 26: temporal.server.chasm.lib.activity.proto.v1.ActivityHeartbeatState.recorded_time:type_name -> google.protobuf.Timestamp - 21, // 27: temporal.server.chasm.lib.activity.proto.v1.ActivityRequestData.input:type_name -> temporal.api.common.v1.Payloads - 22, // 28: temporal.server.chasm.lib.activity.proto.v1.ActivityRequestData.header:type_name -> temporal.api.common.v1.Header - 23, // 29: temporal.server.chasm.lib.activity.proto.v1.ActivityRequestData.user_metadata:type_name -> temporal.api.sdk.v1.UserMetadata - 11, // 30: temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.successful:type_name -> temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.Successful - 12, // 31: temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.failed:type_name -> temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.Failed - 24, // 32: temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.retry_state:type_name -> temporal.api.enums.v1.RetryState - 17, // 33: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.LastFailureDetails.time:type_name -> google.protobuf.Timestamp - 25, // 34: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.LastFailureDetails.failure:type_name -> temporal.api.failure.v1.Failure - 21, // 35: temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.Successful.output:type_name -> temporal.api.common.v1.Payloads - 25, // 36: temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.Failed.failure:type_name -> temporal.api.failure.v1.Failure - 37, // [37:37] is the sub-list for method output_type - 37, // [37:37] is the sub-list for method input_type - 37, // [37:37] is the sub-list for extension type_name - 37, // [37:37] is the sub-list for extension extendee - 0, // [0:37] is the sub-list for field type_name + 17, // 16: temporal.server.chasm.lib.activity.proto.v1.ActivityState.scheduled_start_time:type_name -> google.protobuf.Timestamp + 17, // 17: temporal.server.chasm.lib.activity.proto.v1.ActivityCancelState.request_time:type_name -> google.protobuf.Timestamp + 17, // 18: temporal.server.chasm.lib.activity.proto.v1.ActivityPauseState.pause_time:type_name -> google.protobuf.Timestamp + 15, // 19: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.current_retry_interval:type_name -> google.protobuf.Duration + 17, // 20: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.started_time:type_name -> google.protobuf.Timestamp + 17, // 21: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.complete_time:type_name -> google.protobuf.Timestamp + 10, // 22: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.last_failure_details:type_name -> temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.LastFailureDetails + 20, // 23: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.last_deployment_version:type_name -> temporal.api.deployment.v1.WorkerDeploymentVersion + 17, // 24: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.dispatch_time:type_name -> google.protobuf.Timestamp + 1, // 25: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.current_retry_interval_source:type_name -> temporal.server.chasm.lib.activity.proto.v1.ActivityRetryIntervalSource + 21, // 26: temporal.server.chasm.lib.activity.proto.v1.ActivityHeartbeatState.details:type_name -> temporal.api.common.v1.Payloads + 17, // 27: temporal.server.chasm.lib.activity.proto.v1.ActivityHeartbeatState.recorded_time:type_name -> google.protobuf.Timestamp + 21, // 28: temporal.server.chasm.lib.activity.proto.v1.ActivityRequestData.input:type_name -> temporal.api.common.v1.Payloads + 22, // 29: temporal.server.chasm.lib.activity.proto.v1.ActivityRequestData.header:type_name -> temporal.api.common.v1.Header + 23, // 30: temporal.server.chasm.lib.activity.proto.v1.ActivityRequestData.user_metadata:type_name -> temporal.api.sdk.v1.UserMetadata + 11, // 31: temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.successful:type_name -> temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.Successful + 12, // 32: temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.failed:type_name -> temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.Failed + 24, // 33: temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.retry_state:type_name -> temporal.api.enums.v1.RetryState + 17, // 34: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.LastFailureDetails.time:type_name -> google.protobuf.Timestamp + 25, // 35: temporal.server.chasm.lib.activity.proto.v1.ActivityAttemptState.LastFailureDetails.failure:type_name -> temporal.api.failure.v1.Failure + 21, // 36: temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.Successful.output:type_name -> temporal.api.common.v1.Payloads + 25, // 37: temporal.server.chasm.lib.activity.proto.v1.ActivityOutcome.Failed.failure:type_name -> temporal.api.failure.v1.Failure + 38, // [38:38] is the sub-list for method output_type + 38, // [38:38] is the sub-list for method input_type + 38, // [38:38] is the sub-list for extension type_name + 38, // [38:38] is the sub-list for extension extendee + 0, // [0:38] is the sub-list for field type_name } func init() { file_temporal_server_chasm_lib_activity_proto_v1_activity_state_proto_init() } diff --git a/chasm/lib/activity/proto/v1/activity_state.proto b/chasm/lib/activity/proto/v1/activity_state.proto index fc3fc5af816..077400d768e 100644 --- a/chasm/lib/activity/proto/v1/activity_state.proto +++ b/chasm/lib/activity/proto/v1/activity_state.proto @@ -161,6 +161,10 @@ message ActivityState { // Used to de-dupe update requests. string last_update_options_request_id = 23; + + // Scheduling metadata is stored separately from namespace-defined attributes. + string scheduled_by_id = 24; + google.protobuf.Timestamp scheduled_start_time = 25; } message ActivityCancelState { diff --git a/chasm/lib/scheduler/action_activity.go b/chasm/lib/scheduler/action_activity.go index ffdedf4cea2..5644cba5ca2 100644 --- a/chasm/lib/scheduler/action_activity.go +++ b/chasm/lib/scheduler/action_activity.go @@ -47,7 +47,7 @@ func (activityAction) Start(ctx context.Context, clients actionClients, input ac ScheduleToCloseTimeout: spec.ScheduleToCloseTimeout, ScheduleToStartTimeout: spec.ScheduleToStartTimeout, StartToCloseTimeout: spec.StartToCloseTimeout, HeartbeatTimeout: spec.HeartbeatTimeout, RetryPolicy: spec.RetryPolicy, Input: spec.Input, IdReusePolicy: reuse, IdConflictPolicy: enumspb.ACTIVITY_ID_CONFLICT_POLICY_FAIL, - SearchAttributes: spec.SearchAttributes, Header: spec.Header, + SearchAttributes: scheduler.startActionSearchAttributes(start.GetNominalTime().AsTime()), Header: spec.Header, UserMetadata: spec.UserMetadata, Priority: spec.Priority, CompletionCallbacks: []*commonpb.Callback{input.Callback}, StartDelay: spec.StartDelay, }) if err != nil { diff --git a/chasm/lib/scheduler/library.go b/chasm/lib/scheduler/library.go index 9183283162d..4d832695941 100644 --- a/chasm/lib/scheduler/library.go +++ b/chasm/lib/scheduler/library.go @@ -68,6 +68,9 @@ func (l *Library) Components() []*chasm.RegistrableComponent { scheduleIdleCloseTimeSearchAttribute, scheduleRunningWorkflowCountSearchAttribute, scheduleBufferedStartsCountSearchAttribute, + scheduleActionKindSearchAttribute, + scheduleActionTypeSearchAttribute, + scheduleRunningExecutionCountSearchAttribute, ), // Exposes Tweakables to scheduler components via the CHASM context // (see tweakablesFromContext). diff --git a/chasm/lib/scheduler/scheduler.go b/chasm/lib/scheduler/scheduler.go index b30e9a89238..cdf24e611e8 100644 --- a/chasm/lib/scheduler/scheduler.go +++ b/chasm/lib/scheduler/scheduler.go @@ -82,13 +82,19 @@ const ScheduleNextActionTimeName = "ScheduleNextActionTime" const ScheduleIdleCloseTimeName = "ScheduleIdleCloseTime" const ScheduleRunningWorkflowCountName = "ScheduleRunningWorkflowCount" const ScheduleBufferedStartsCountName = "ScheduleBufferedStartsCount" +const ScheduleActionKindName = "ScheduleActionKind" +const ScheduleActionTypeName = "ScheduleActionType" +const ScheduleRunningExecutionCountName = "ScheduleRunningExecutionCount" var ( - executionStatusSearchAttribute = chasm.NewSearchAttributeKeyword("ExecutionStatus", chasm.SearchAttributeFieldLowCardinalityKeyword01) - scheduleNextActionTimeSearchAttribute = chasm.NewSearchAttributeDateTime(ScheduleNextActionTimeName, chasm.SearchAttributeFieldDateTime01) - scheduleIdleCloseTimeSearchAttribute = chasm.NewSearchAttributeDateTime(ScheduleIdleCloseTimeName, chasm.SearchAttributeFieldDateTime02) - scheduleRunningWorkflowCountSearchAttribute = chasm.NewSearchAttributeInt(ScheduleRunningWorkflowCountName, chasm.SearchAttributeFieldInt01) - scheduleBufferedStartsCountSearchAttribute = chasm.NewSearchAttributeInt(ScheduleBufferedStartsCountName, chasm.SearchAttributeFieldInt02) + executionStatusSearchAttribute = chasm.NewSearchAttributeKeyword("ExecutionStatus", chasm.SearchAttributeFieldLowCardinalityKeyword01) + scheduleNextActionTimeSearchAttribute = chasm.NewSearchAttributeDateTime(ScheduleNextActionTimeName, chasm.SearchAttributeFieldDateTime01) + scheduleIdleCloseTimeSearchAttribute = chasm.NewSearchAttributeDateTime(ScheduleIdleCloseTimeName, chasm.SearchAttributeFieldDateTime02) + scheduleRunningWorkflowCountSearchAttribute = chasm.NewSearchAttributeInt(ScheduleRunningWorkflowCountName, chasm.SearchAttributeFieldInt01) + scheduleBufferedStartsCountSearchAttribute = chasm.NewSearchAttributeInt(ScheduleBufferedStartsCountName, chasm.SearchAttributeFieldInt02) + scheduleActionKindSearchAttribute = chasm.NewSearchAttributeKeyword(ScheduleActionKindName, chasm.SearchAttributeFieldKeyword02) + scheduleActionTypeSearchAttribute = chasm.NewSearchAttributeKeyword(ScheduleActionTypeName, chasm.SearchAttributeFieldKeyword01) + scheduleRunningExecutionCountSearchAttribute = chasm.NewSearchAttributeInt(ScheduleRunningExecutionCountName, chasm.SearchAttributeFieldInt03) ) var initialSerializedConflictToken = serializeConflictToken(scheduler.InitialConflictToken) @@ -370,12 +376,18 @@ func (s *Scheduler) LifecycleState(ctx chasm.Context) chasm.LifecycleState { func (s *Scheduler) ContextMetadata(_ chasm.Context) map[string]string { md := make(map[string]string, 2) - if wfType := s.Schedule.GetAction().GetStartWorkflow().GetWorkflowType().GetName(); wfType != "" { - md[contextutil.MetadataKeyWorkflowType] = wfType + metadata := s.actionMetadata() + typeKey, queueKey := contextutil.MetadataKeyWorkflowType, contextutil.MetadataKeyWorkflowTaskQueue + if metadata.Kind == enumspb.EXECUTION_TYPE_ACTIVITY { + typeKey, queueKey = contextutil.MetadataKeyStandaloneActivityType, contextutil.MetadataKeyStandaloneActivityTaskQueue } - if tq := s.Schedule.GetAction().GetStartWorkflow().GetTaskQueue().GetName(); tq != "" { - md[contextutil.MetadataKeyWorkflowTaskQueue] = tq + if metadata.Type != "" { + md[typeKey] = metadata.Type } + if metadata.TaskQueue != "" { + md[queueKey] = metadata.TaskQueue + } + if len(md) == 0 { return nil } @@ -1136,6 +1148,8 @@ func (s *Scheduler) SearchAttributes(ctx chasm.Context) []chasm.SearchAttributeK out := []chasm.SearchAttributeKeyValue{ executionStatusSearchAttribute.Value(s.executionStatus()), chasm.SearchAttributeTemporalSchedulePaused.Value(s.Schedule.GetState().GetPaused()), + scheduleActionKindSearchAttribute.Value(s.actionMetadata().Kind.String()), + scheduleActionTypeSearchAttribute.Value(s.actionMetadata().Type), } if !s.Closed { if gen := s.Generator.Get(ctx); len(gen.FutureActionTimes) > 0 { @@ -1152,6 +1166,7 @@ func (s *Scheduler) SearchAttributes(ctx chasm.Context) []chasm.SearchAttributeK // Emitted even when zero so that exact and range queries both work. out = append(out, scheduleRunningWorkflowCountSearchAttribute.Value(runningWorkflowCount), + scheduleRunningExecutionCountSearchAttribute.Value(int64(len(invoker.runningExecutions()))), scheduleBufferedStartsCountSearchAttribute.Value(bufferedStartsCount), ) } @@ -1194,12 +1209,15 @@ func (s *Scheduler) ListInfo( recentActions = util.SliceTail(recentActions, recentActionCountForList) return &schedulepb.ScheduleListInfo{ - Spec: spec, - WorkflowType: s.Schedule.Action.GetStartWorkflow().GetWorkflowType(), - Notes: s.Schedule.State.Notes, - Paused: s.Schedule.State.Paused, - RecentActions: recentActions, - FutureActionTimes: generator.FutureActionTimes, + Spec: spec, + WorkflowType: s.Schedule.Action.GetStartWorkflow().GetWorkflowType(), + ActionKind: s.actionMetadata().Kind, + ActionType: s.actionMetadata().Type, + RunningExecutionCount: int64(len(s.Invoker.Get(ctx).runningExecutions())), + Notes: s.Schedule.State.Notes, + Paused: s.Schedule.State.Paused, + RecentActions: recentActions, + FutureActionTimes: generator.FutureActionTimes, } } diff --git a/chasm/search_attribute.go b/chasm/search_attribute.go index 0842dd22a0d..1cae27ea06f 100644 --- a/chasm/search_attribute.go +++ b/chasm/search_attribute.go @@ -41,6 +41,7 @@ var ( SearchAttributeFieldInt01 = newSearchAttributeFieldInt(1) SearchAttributeFieldInt02 = newSearchAttributeFieldInt(2) + SearchAttributeFieldInt03 = newSearchAttributeFieldInt(3) SearchAttributeFieldDouble01 = newSearchAttributeFieldDouble(1) SearchAttributeFieldDouble02 = newSearchAttributeFieldDouble(2)