Skip to content
Draft
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
80 changes: 56 additions & 24 deletions chasm/lib/activity/activity.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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"
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
}
29 changes: 29 additions & 0 deletions chasm/lib/activity/activity_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
89 changes: 53 additions & 36 deletions chasm/lib/activity/gen/activitypb/v1/activity_state.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading
Loading