Skip to content
Merged
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
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ require (
go.opentelemetry.io/otel/sdk v1.34.0
go.opentelemetry.io/otel/sdk/metric v1.34.0
go.opentelemetry.io/otel/trace v1.34.0
go.temporal.io/api v1.58.1-0.20251126231839-2fcd2247e106
go.temporal.io/api v1.58.1-0.20251128181858-703071215042
go.temporal.io/sdk v1.35.0
go.uber.org/fx v1.24.0
go.uber.org/mock v0.6.0
Expand Down
4 changes: 2 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -390,8 +390,8 @@ go.opentelemetry.io/otel/trace v1.34.0 h1:+ouXS2V8Rd4hp4580a8q23bg0azF2nI8cqLYnC
go.opentelemetry.io/otel/trace v1.34.0/go.mod h1:Svm7lSjQD7kG7KJ/MUHPVXSDGz2OX4h0M2jHBhmSfRE=
go.opentelemetry.io/proto/otlp v1.5.0 h1:xJvq7gMzB31/d406fB8U5CBdyQGw4P399D1aQWU/3i4=
go.opentelemetry.io/proto/otlp v1.5.0/go.mod h1:keN8WnHxOy8PG0rQZjJJ5A2ebUoafqWp0eVQ4yIXvJ4=
go.temporal.io/api v1.58.1-0.20251126231839-2fcd2247e106 h1:V2H8rfBDapmWpIsNDWZLOS95WIIWAgPXnG7gpNrWO5Y=
go.temporal.io/api v1.58.1-0.20251126231839-2fcd2247e106/go.mod h1:iaxoP/9OXMJcQkETTECfwYq4cw/bj4nwov8b3ZLVnXM=
go.temporal.io/api v1.58.1-0.20251128181858-703071215042 h1:44+nPe+rGhYUwA1oDi46rkXEYEVfoAxOmb0myvTm4Es=
go.temporal.io/api v1.58.1-0.20251128181858-703071215042/go.mod h1:iaxoP/9OXMJcQkETTECfwYq4cw/bj4nwov8b3ZLVnXM=
go.temporal.io/sdk v1.35.0 h1:lRNAQ5As9rLgYa7HBvnmKyzxLcdElTuoFJ0FXM/AsLQ=
go.temporal.io/sdk v1.35.0/go.mod h1:1q5MuLc2MEJ4lneZTHJzpVebW2oZnyxoIOWX3oFVebw=
go.uber.org/atomic v1.5.0/go.mod h1:sABNBOSYdrvTF6hTgEIbc7YasKWGhgEQZyfxyTvoXHQ=
Expand Down
6 changes: 6 additions & 0 deletions service/frontend/workflow_handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -5950,6 +5950,9 @@ func (wh *WorkflowHandler) UpdateWorkflowExecutionOptions(
if err != nil {
return nil, serviceerror.NewInvalidArgumentf("error parsing UpdateMask: %s", err.Error())
}
if err := priorities.Validate(opts.GetPriority()); err != nil {
return nil, err
}

namespaceID, err := wh.namespaceRegistry.GetNamespaceID(namespace.Name(request.GetNamespace()))
if err != nil {
Expand Down Expand Up @@ -5987,6 +5990,9 @@ func (wh *WorkflowHandler) UpdateActivityOptions(
if request.GetActivity() == nil {
return nil, errActivityIDOrTypeNotSet
}
if err := priorities.Validate(request.GetActivityOptions().GetPriority()); err != nil {
return nil, err
}

namespaceID, err := wh.namespaceRegistry.GetNamespaceID(namespace.Name(request.GetNamespace()))
if err != nil {
Expand Down
101 changes: 101 additions & 0 deletions service/frontend/workflow_handler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import (
"github.com/stretchr/testify/mock"
"github.com/stretchr/testify/require"
"github.com/stretchr/testify/suite"
activitypb "go.temporal.io/api/activity/v1"
batchpb "go.temporal.io/api/batch/v1"
commonpb "go.temporal.io/api/common/v1"
enumspb "go.temporal.io/api/enums/v1"
Expand Down Expand Up @@ -914,6 +915,31 @@ func (s *WorkflowHandlerSuite) TestStartWorkflowExecution_Failed_InvalidAggregat
s.ErrorContains(err, "cannot attach more than 10 links per request, got 11")
}

func (s *WorkflowHandlerSuite) TestStartWorkflowExecution_Priority() {
config := s.newConfig()
wh := s.getWorkflowHandler(config)

s.mockSearchAttributesMapperProvider.EXPECT().GetMapper(gomock.Any()).Return(nil, nil).AnyTimes()

request := &workflowservice.StartWorkflowExecutionRequest{
Namespace: s.testNamespace.String(),
WorkflowId: "workflow-id",
WorkflowType: &commonpb.WorkflowType{
Name: "workflow-type",
},
TaskQueue: &taskqueuepb.TaskQueue{
Name: "task-queue",
},
Priority: &commonpb.Priority{PriorityKey: -1},
}

_, err := wh.StartWorkflowExecution(context.Background(), request)
var invalidArg *serviceerror.InvalidArgument
s.ErrorAs(err, &invalidArg)
s.ErrorContains(err, "priority key can't be negative")
// NOTE: only testing a single validation scenario here; the priority validation has its own unit tests
}

func (s *WorkflowHandlerSuite) TestSignalWithStartWorkflowExecution_InvalidWorkflowIdConflictPolicy() {
config := s.newConfig()
wh := s.getWorkflowHandler(config)
Expand Down Expand Up @@ -1011,6 +1037,32 @@ func (s *WorkflowHandlerSuite) TestSignalWithStartWorkflowExecution_Failed_Inval
s.ErrorContains(err, "link exceeds allowed size of 4000")
}

func (s *WorkflowHandlerSuite) TestSignalWithStartWorkflowExecution_Priority() {
config := s.newConfig()
wh := s.getWorkflowHandler(config)

s.mockSearchAttributesMapperProvider.EXPECT().GetMapper(gomock.Any()).Return(nil, nil).AnyTimes()

request := &workflowservice.SignalWithStartWorkflowExecutionRequest{
Namespace: s.testNamespace.String(),
WorkflowId: "workflow-id",
WorkflowType: &commonpb.WorkflowType{
Name: "workflow-type",
},
TaskQueue: &taskqueuepb.TaskQueue{
Name: "task-queue",
},
SignalName: "signal-name",
Priority: &commonpb.Priority{PriorityKey: -1},
}

_, err := wh.SignalWithStartWorkflowExecution(context.Background(), request)
var invalidArg *serviceerror.InvalidArgument
s.ErrorAs(err, &invalidArg)
s.ErrorContains(err, "priority key can't be negative")
// NOTE: only testing a single validation scenario here; the priority validation has its own unit tests
}

func (s *WorkflowHandlerSuite) TestSignalWorkflowExecution_Failed_InvalidLinks() {
s.mockSearchAttributesMapperProvider.EXPECT().GetMapper(gomock.Any()).AnyTimes().Return(nil, nil)
config := s.newConfig()
Expand Down Expand Up @@ -3958,3 +4010,52 @@ func (s *WorkflowHandlerSuite) TestUpdateTaskQueueConfig_Validation() {
s.NotNil(resp)
})
}

func (s *WorkflowHandlerSuite) TestUpdateWorkflowExecutionOptions_Priority() {
config := s.newConfig()
wh := s.getWorkflowHandler(config)

request := &workflowservice.UpdateWorkflowExecutionOptionsRequest{
Namespace: s.testNamespace.String(),
WorkflowExecution: &commonpb.WorkflowExecution{
WorkflowId: "workflow-id",
RunId: "run-id",
},
WorkflowExecutionOptions: &workflowpb.WorkflowExecutionOptions{
Priority: &commonpb.Priority{PriorityKey: -1},
},
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"priority"}},
}

_, err := wh.UpdateWorkflowExecutionOptions(context.Background(), request)
var invalidArg *serviceerror.InvalidArgument
s.ErrorAs(err, &invalidArg)
s.ErrorContains(err, "priority key can't be negative")
// NOTE: only testing a single validation scenario here; the priority validation has its own unit tests
}

func (s *WorkflowHandlerSuite) TestUpdateActivityOptions_Priority() {
config := s.newConfig()
wh := s.getWorkflowHandler(config)

request := &workflowservice.UpdateActivityOptionsRequest{
Namespace: s.testNamespace.String(),
Execution: &commonpb.WorkflowExecution{
WorkflowId: "workflow-id",
RunId: "run-id",
},
Activity: &workflowservice.UpdateActivityOptionsRequest_Id{
Id: "activity-id",
},
ActivityOptions: &activitypb.ActivityOptions{
Priority: &commonpb.Priority{PriorityKey: -1},
},
UpdateMask: &fieldmaskpb.FieldMask{Paths: []string{"priority"}},
}

_, err := wh.UpdateActivityOptions(context.Background(), request)
var invalidArg *serviceerror.InvalidArgument
s.ErrorAs(err, &invalidArg)
s.ErrorContains(err, "priority key can't be negative")
// NOTE: only testing a single validation scenario here; the priority validation has its own unit tests
}
7 changes: 3 additions & 4 deletions service/history/api/recordactivitytaskstarted/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -165,12 +165,11 @@ func recordActivityTaskStarted(
}

if ai.Stamp != request.Stamp {
// activity has changes before task is started.
// ErrActivityStampMismatch is the error to indicate that requested activity has mismatched stamp
// This happens when the workflow task was rescheduled.
errorMessage := fmt.Sprintf(
"Activity task with this stamp not found. Id: %s,: type: %s, current stamp: %d",
"Activity task rejected; stamp has changed. Id: %s,: type: %s, current stamp: %d",
ai.ActivityId, ai.ActivityType.Name, ai.Stamp)
return nil, rejectCodeUndefined, serviceerror.NewNotFound(errorMessage)
return nil, rejectCodeUndefined, serviceerrors.NewObsoleteMatchingTask(errorMessage)

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 seems much more consistent and a better error type, too.

}

wfBehavior := mutableState.GetEffectiveVersioningBehavior()
Expand Down
3 changes: 2 additions & 1 deletion service/history/api/startworkflow/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -670,7 +670,8 @@ func (s *Starter) handleUseExistingWorkflowOnConflictOptions(
requestID,
completionCallbacks,
links,
"",
"", // identity
nil, // priority
)
return api.UpdateWorkflowWithoutWorkflowTask, err
},
Expand Down
42 changes: 42 additions & 0 deletions service/history/api/updateactivityoptions/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
"go.temporal.io/api/workflowservice/v1"
"go.temporal.io/server/api/historyservice/v1"
persistencespb "go.temporal.io/server/api/persistence/v1"
"go.temporal.io/server/common"
"go.temporal.io/server/common/definition"
"go.temporal.io/server/common/namespace"
"go.temporal.io/server/common/util"
Expand Down Expand Up @@ -148,6 +149,7 @@ func processActivityOptionsUpdate(
ScheduleToStartTimeout: ai.ScheduleToStartTimeout,
StartToCloseTimeout: ai.StartToCloseTimeout,
HeartbeatTimeout: ai.HeartbeatTimeout,
Priority: common.CloneProto(ai.Priority),
RetryPolicy: &commonpb.RetryPolicy{
BackoffCoefficient: ai.RetryBackoffCoefficient,
InitialInterval: ai.RetryInitialInterval,
Expand Down Expand Up @@ -202,10 +204,48 @@ func mergeActivityOptions(
mergeInto.HeartbeatTimeout = mergeFrom.HeartbeatTimeout
}

if _, ok := updateFields["priority"]; ok {
mergeInto.Priority = mergeFrom.Priority
}

if _, ok := updateFields["priority.priorityKey"]; ok {
if mergeFrom.Priority == nil {
return serviceerror.NewInvalidArgument("Priority is not provided")
}
if mergeInto.Priority == nil {
mergeInto.Priority = &commonpb.Priority{}
}
mergeInto.Priority.PriorityKey = mergeFrom.Priority.PriorityKey
}

if _, ok := updateFields["priority.fairnessKey"]; ok {
if mergeFrom.Priority == nil {
return serviceerror.NewInvalidArgument("Priority is not provided")
}
if mergeInto.Priority == nil {
mergeInto.Priority = &commonpb.Priority{}
}
mergeInto.Priority.FairnessKey = mergeFrom.Priority.FairnessKey
}

if _, ok := updateFields["priority.fairnessWeight"]; ok {
if mergeFrom.Priority == nil {
return serviceerror.NewInvalidArgument("Priority is not provided")
}
if mergeInto.Priority == nil {
mergeInto.Priority = &commonpb.Priority{}
}
mergeInto.Priority.FairnessWeight = mergeFrom.Priority.FairnessWeight
}

if mergeInto.RetryPolicy == nil {
mergeInto.RetryPolicy = &commonpb.RetryPolicy{}
}

if _, ok := updateFields["retryPolicy"]; ok {
mergeInto.RetryPolicy = mergeFrom.RetryPolicy
}
Comment on lines +245 to +247

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 seems like a reasonable addition to me.

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.

I guess so. I was originally a little confused at the semantics, but it seems reasonable, to allow updating a whole sub-object without listing all the fields


if _, ok := updateFields["retryPolicy.initialInterval"]; ok {
if mergeFrom.RetryPolicy == nil {
return serviceerror.NewInvalidArgument("RetryPolicy is not provided")
Expand Down Expand Up @@ -295,6 +335,7 @@ func updateActivityOptions(
activityInfo.ScheduleToStartTimeout = activityOptions.ScheduleToStartTimeout
activityInfo.StartToCloseTimeout = activityOptions.StartToCloseTimeout
activityInfo.HeartbeatTimeout = activityOptions.HeartbeatTimeout
activityInfo.Priority = activityOptions.Priority
activityInfo.RetryMaximumInterval = activityOptions.RetryPolicy.MaximumInterval
activityInfo.RetryBackoffCoefficient = activityOptions.RetryPolicy.BackoffCoefficient
activityInfo.RetryInitialInterval = activityOptions.RetryPolicy.InitialInterval
Expand Down Expand Up @@ -375,6 +416,7 @@ func restoreOriginalOptions(
ScheduleToStartTimeout: originalOptions.ScheduleToStartTimeout,
StartToCloseTimeout: originalOptions.StartToCloseTimeout,
HeartbeatTimeout: originalOptions.HeartbeatTimeout,
Priority: originalOptions.Priority,
RetryPolicy: originalOptions.RetryPolicy,
}

Expand Down
Loading
Loading