package tests import ( "bytes" "encoding/binary" "fmt" "strconv" "testing" "time" "github.com/google/uuid" commandpb "go.temporal.io/api/command/v1" commonpb "go.temporal.io/api/common/v1" enumspb "go.temporal.io/api/enums/v1" filterpb "go.temporal.io/api/filter/v1" historypb "go.temporal.io/api/history/v1" "go.temporal.io/api/serviceerror" taskqueuepb "go.temporal.io/api/taskqueue/v1" "go.temporal.io/api/workflowservice/v1" "go.temporal.io/server/common/convert" "go.temporal.io/server/common/dynamicconfig" "go.temporal.io/server/common/log/tag" "go.temporal.io/server/common/payload" "go.temporal.io/server/common/payloads" "go.temporal.io/server/common/rpc" "go.temporal.io/server/common/testing/parallelsuite" "go.temporal.io/server/service/history/consts" "go.temporal.io/server/tests/testcore" "google.golang.org/protobuf/types/known/durationpb" "google.golang.org/protobuf/types/known/timestamppb" ) type SignalWorkflowTestSuite struct { parallelsuite.Suite[*SignalWorkflowTestSuite] } func TestSignalWorkflowTestSuiteLegacy(t *testing.T) { parallelsuite.Run(t, &SignalWorkflowTestSuite{}, []testcore.TestOption{}) } func TestSignalWorkflowTestSuiteChasm(t *testing.T) { parallelsuite.Run( t, &SignalWorkflowTestSuite{}, []testcore.TestOption{ testcore.WithDynamicConfig(dynamicconfig.EnableChasm, true), testcore.WithDynamicConfig(dynamicconfig.EnableCHASMSignalBacklinks, true), }, ) } func (s *SignalWorkflowTestSuite) TestSignalWorkflow(opts []testcore.TestOption) { env := testcore.NewEnv(s.T(), opts...) id := "functional-signal-workflow-test" wt := "functional-signal-workflow-test-type" tl := "functional-signal-workflow-test-taskqueue" identity := "worker1" activityName := "activity_type1" workflowType := &commonpb.WorkflowType{Name: wt} taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL} // Send a signal to non-exist workflow header := &commonpb.Header{ Fields: map[string]*commonpb.Payload{"signal header key": payload.EncodeString("signal header value")}, } _, err0 := env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, RunId: uuid.NewString(), }, SignalName: "failed signal", Input: nil, Identity: identity, Header: header, }) s.Error(err0) s.IsType(&serviceerror.NotFound{}, err0) // Start workflow execution request := &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), Identity: identity, } we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request) s.NoError(err0) env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId)) // workflow logic workflowComplete := false activityScheduled := false activityData := int32(1) var signalEvent *historypb.HistoryEvent wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { if !activityScheduled { activityScheduled = true buf := new(bytes.Buffer) s.NoError(binary.Write(buf, binary.LittleEndian, activityData)) return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK, Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{ ActivityId: strconv.Itoa(1), ActivityType: &commonpb.ActivityType{Name: activityName}, TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, Input: payloads.EncodeBytes(buf.Bytes()), ScheduleToCloseTimeout: durationpb.New(100 * time.Second), ScheduleToStartTimeout: durationpb.New(2 * time.Second), StartToCloseTimeout: durationpb.New(50 * time.Second), HeartbeatTimeout: durationpb.New(5 * time.Second), }}, }}, nil } else if task.PreviousStartedEventId > 0 { for _, event := range task.History.Events[task.PreviousStartedEventId:] { if event.GetEventType() == enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_SIGNALED { signalEvent = event return []*commandpb.Command{}, nil } } } workflowComplete = true return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{ Result: payloads.EncodeString("Done"), }}, }}, nil } // activity handler atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) { return payloads.EncodeString("Activity Result"), false, nil } poller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.Namespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: wtHandler, ActivityTaskHandler: atHandler, Logger: env.Logger, T: s.T(), } // Make first command to schedule activity _, err := poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) // Send first signal using RunID signalName := "my signal" signalInput := payloads.EncodeString("my signal input") _, err = env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we.RunId, }, SignalName: signalName, Input: signalInput, Identity: identity, Header: header, }) s.NoError(err) // Process signal in workflow _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) s.False(workflowComplete) s.NotNil(signalEvent) s.Equal(signalName, signalEvent.GetWorkflowExecutionSignaledEventAttributes().SignalName) s.ProtoEqual(signalInput, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Input) s.Equal(identity, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Identity) s.ProtoEqual(header, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Header) // Send another signal without RunID signalName = "another signal" signalInput = payloads.EncodeString("another signal input") _, err = env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, }, SignalName: signalName, Input: signalInput, Identity: identity, }) s.NoError(err) // Process signal in workflow _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) s.False(workflowComplete) s.NotNil(signalEvent) s.Equal(signalName, signalEvent.GetWorkflowExecutionSignaledEventAttributes().SignalName) s.ProtoEqual(signalInput, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Input) s.Equal(identity, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Identity) // Terminate workflow execution _, err = env.FrontendClient().TerminateWorkflowExecution(s.Context(), &workflowservice.TerminateWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, }, Reason: "test signal", Details: nil, Identity: identity, }) s.NoError(err) // Send signal to terminated workflow _, err = env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we.RunId, }, SignalName: "failed signal 1", Input: nil, Identity: identity, }) s.Error(err) s.IsType(&serviceerror.NotFound{}, err) } func (s *SignalWorkflowTestSuite) TestSignalWorkflow_DuplicateRequest(opts []testcore.TestOption) { env := testcore.NewEnv(s.T(), opts...) id := "functional-signal-workflow-test-duplicate" wt := "functional-signal-workflow-test-duplicate-type" tl := "functional-signal-workflow-test-duplicate-taskqueue" identity := "worker1" activityName := "activity_type1" workflowType := &commonpb.WorkflowType{Name: wt} taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL} // Start workflow execution request := &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), Identity: identity, } we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request) s.NoError(err0) env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId)) // workflow logic workflowComplete := false activityScheduled := false activityData := int32(1) var signalEvent *historypb.HistoryEvent numOfSignaledEvent := 0 wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { if !activityScheduled { activityScheduled = true buf := new(bytes.Buffer) s.NoError(binary.Write(buf, binary.LittleEndian, activityData)) return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK, Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{ ActivityId: strconv.Itoa(1), ActivityType: &commonpb.ActivityType{Name: activityName}, TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, Input: payloads.EncodeBytes(buf.Bytes()), ScheduleToCloseTimeout: durationpb.New(100 * time.Second), ScheduleToStartTimeout: durationpb.New(2 * time.Second), StartToCloseTimeout: durationpb.New(50 * time.Second), HeartbeatTimeout: durationpb.New(5 * time.Second), }}, }}, nil } else if task.PreviousStartedEventId > 0 { numOfSignaledEvent = 0 for _, event := range task.History.Events[task.PreviousStartedEventId:] { if event.GetEventType() == enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_SIGNALED { signalEvent = event numOfSignaledEvent++ } } return []*commandpb.Command{}, nil } workflowComplete = true return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{ Result: payloads.EncodeString("Done"), }}, }}, nil } // activity handler atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) { return payloads.EncodeString("Activity Result"), false, nil } poller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.Namespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: wtHandler, ActivityTaskHandler: atHandler, Logger: env.Logger, T: s.T(), } // Make first command to schedule activity _, err := poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) // Send first signal signalName := "my signal" signalInput := payloads.EncodeString("my signal input") requestID := uuid.NewString() signalReqest := &workflowservice.SignalWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we.RunId, }, SignalName: signalName, Input: signalInput, Identity: identity, RequestId: requestID, } _, err = env.FrontendClient().SignalWorkflowExecution(s.Context(), signalReqest) s.NoError(err) // Process signal in workflow _, err = poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) s.False(workflowComplete) s.NotNil(signalEvent) s.Equal(signalName, signalEvent.GetWorkflowExecutionSignaledEventAttributes().SignalName) s.ProtoEqual(signalInput, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Input) s.Equal(identity, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Identity) s.Equal(1, numOfSignaledEvent) // Send another signal with same request id _, err = env.FrontendClient().SignalWorkflowExecution(s.Context(), signalReqest) s.NoError(err) // Process signal in workflow _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) s.False(workflowComplete) s.NotNil(signalEvent) s.Equal(0, numOfSignaledEvent) } func (s *SignalWorkflowTestSuite) TestSignalExternalWorkflowCommand(opts []testcore.TestOption) { // Explicitly enable cross namespace commands for this test, // need a dedicated cluster to enable cross namespace commands opts = append( opts, testcore.WithDedicatedCluster(), testcore.WithDynamicConfig(dynamicconfig.EnableCrossNamespaceCommands, true), ) env := testcore.NewEnv(s.T(), opts...) id := "functional-signal-external-workflow-test" wt := "functional-signal-external-workflow-test-type" tl := "functional-signal-external-workflow-test-taskqueue" identity := "worker1" activityName := "activity_type1" workflowType := &commonpb.WorkflowType{Name: wt} taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL} request := &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), Identity: identity, } we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request) s.NoError(err0) env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId)) externalRequest := &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.ExternalNamespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), Identity: identity, } we2, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), externalRequest) s.NoError(err0) env.Logger.Info("StartWorkflowExecution on external Namespace", tag.WorkflowNamespace(env.ExternalNamespace().String()), tag.WorkflowRunID(we2.RunId)) activityCount := int32(1) activityCounter := int32(0) signalName := "my signal" signalInput := payloads.EncodeString("my signal input") signalHeader := &commonpb.Header{ Fields: map[string]*commonpb.Payload{"signal header key": payload.EncodeString("signal header value")}, } wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { if activityCounter < activityCount { activityCounter++ buf := new(bytes.Buffer) s.NoError(binary.Write(buf, binary.LittleEndian, activityCounter)) return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK, Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{ ActivityId: convert.Int32ToString(activityCounter), ActivityType: &commonpb.ActivityType{Name: activityName}, TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, Input: payloads.EncodeBytes(buf.Bytes()), ScheduleToCloseTimeout: durationpb.New(100 * time.Second), ScheduleToStartTimeout: durationpb.New(10 * time.Second), StartToCloseTimeout: durationpb.New(50 * time.Second), HeartbeatTimeout: durationpb.New(5 * time.Second), }}, }}, nil } return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SIGNAL_EXTERNAL_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_SignalExternalWorkflowExecutionCommandAttributes{SignalExternalWorkflowExecutionCommandAttributes: &commandpb.SignalExternalWorkflowExecutionCommandAttributes{ Namespace: env.ExternalNamespace().String(), Execution: &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we2.GetRunId(), }, SignalName: signalName, Input: signalInput, Header: signalHeader, }}, }}, nil } atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) { return payloads.EncodeString("Activity Result"), false, nil } poller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.Namespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: wtHandler, ActivityTaskHandler: atHandler, Logger: env.Logger, T: s.T(), } workflowComplete := false externalActivityCount := int32(1) externalActivityCounter := int32(0) var signalEvent *historypb.HistoryEvent externalWFTHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { if externalActivityCounter < externalActivityCount { externalActivityCounter++ buf := new(bytes.Buffer) s.NoError(binary.Write(buf, binary.LittleEndian, externalActivityCounter)) return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK, Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{ ActivityId: convert.Int32ToString(externalActivityCounter), ActivityType: &commonpb.ActivityType{Name: activityName}, TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, Input: payloads.EncodeBytes(buf.Bytes()), ScheduleToCloseTimeout: durationpb.New(100 * time.Second), ScheduleToStartTimeout: durationpb.New(10 * time.Second), StartToCloseTimeout: durationpb.New(50 * time.Second), HeartbeatTimeout: durationpb.New(5 * time.Second), }}, }}, nil } else if task.PreviousStartedEventId > 0 { for _, event := range task.History.Events[task.PreviousStartedEventId:] { if event.GetEventType() == enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_SIGNALED { signalEvent = event return []*commandpb.Command{}, nil } } } workflowComplete = true return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{ Result: payloads.EncodeString("Done"), }}, }}, nil } //nolint:staticcheck // SA1019 TaskPoller replacement needs to be done holistically. externalPoller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.ExternalNamespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: externalWFTHandler, ActivityTaskHandler: atHandler, Logger: env.Logger, T: s.T(), } // Start both current and external workflows to make some progress. _, err := poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) _, err = externalPoller.PollAndProcessWorkflowTask() env.Logger.Info("external PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) err = externalPoller.PollAndProcessActivityTask(false) env.Logger.Info("external PollAndProcessActivityTask", tag.Error(err)) s.NoError(err) // Signal the external workflow with this command. _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) // in source workflow var historyEvents []*historypb.HistoryEvent CheckHistoryLoopForSignalSent: for i := 1; i < 10; i++ { historyEvents = env.GetHistory(env.Namespace().String(), &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we.RunId, }) signalRequestedEvent := historyEvents[len(historyEvents)-2] if signalRequestedEvent.GetEventType() != enumspb.EVENT_TYPE_EXTERNAL_WORKFLOW_EXECUTION_SIGNALED { env.Logger.Info("Signal still not sent") time.Sleep(100 * time.Millisecond) //nolint:forbidigo continue CheckHistoryLoopForSignalSent } break } s.EqualHistoryEvents(fmt.Sprintf(` 1 WorkflowExecutionStarted 2 WorkflowTaskScheduled 3 WorkflowTaskStarted 4 WorkflowTaskCompleted 5 ActivityTaskScheduled 6 ActivityTaskTimedOut 7 WorkflowTaskScheduled 8 WorkflowTaskStarted 9 WorkflowTaskCompleted 10 SignalExternalWorkflowExecutionInitiated 11 ExternalWorkflowExecutionSignaled {"InitiatedEventId":10,"WorkflowExecution":{"RunId":"%s","WorkflowId":"%s"}} 12 WorkflowTaskScheduled`, we2.RunId, id), historyEvents) // Process signal in workflow for external workflow s.Await(func(s *SignalWorkflowTestSuite) { _, err = externalPoller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) }, 10*time.Second, 100*time.Millisecond) s.False(workflowComplete) s.NotNil(signalEvent) s.Equal(signalName, signalEvent.GetWorkflowExecutionSignaledEventAttributes().SignalName) s.ProtoEqual(signalInput, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Input) s.ProtoEqual(signalHeader, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Header) s.Equal("history-service", signalEvent.GetWorkflowExecutionSignaledEventAttributes().Identity) } func (s *SignalWorkflowTestSuite) TestSignalWorkflow_Cron_NoWorkflowTaskCreated(opts []testcore.TestOption) { env := testcore.NewEnv(s.T(), opts...) id := "functional-signal-workflow-test-cron" wt := "functional-signal-workflow-test-cron-type" tl := "functional-signal-workflow-test-cron-taskqueue" identity := "worker1" cronSpec := "@every 2s" workflowType := &commonpb.WorkflowType{Name: wt} taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL} // Start workflow execution request := &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), Identity: identity, CronSchedule: cronSpec, } now := time.Now().UTC() we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request) s.NoError(err0) env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId)) // Send first signal using RunID signalName := "my signal" signalInput := payloads.EncodeString("my signal input") _, err := env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we.RunId, }, SignalName: signalName, Input: signalInput, Identity: identity, }) s.NoError(err) // workflow logic var workflowTaskDelay time.Duration wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { workflowTaskDelay = time.Now().UTC().Sub(now) return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{ Result: payloads.EncodeString("Done"), }}, }}, nil } poller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.Namespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: wtHandler, Logger: env.Logger, T: s.T(), } // Make first command to schedule activity _, err = poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) s.Greater(workflowTaskDelay, time.Second*2) } func (s *SignalWorkflowTestSuite) TestSignalWorkflow_WorkflowCloseAttempted(opts []testcore.TestOption) { env := testcore.NewEnv(s.T(), opts...) id := "functional-signal-workflow-workflow-close-attempted-test" wt := "functional-signal-workflow-workflow-close-attempted-test-type" tl := "functional-signal-workflow-workflow-close-attempted-test-taskqueue" identity := "worker1" workflowType := &commonpb.WorkflowType{Name: wt} taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL} we, err := env.FrontendClient().StartWorkflowExecution(s.Context(), &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(3 * time.Second), Identity: identity, }) s.NoError(err) attemptCount := 1 wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { if attemptCount == 1 { _, err := env.FrontendClient().SignalWorkflowExecution(s.Context(), &workflowservice.SignalWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we.RunId, }, SignalName: "buffered-signal", Identity: identity, RequestId: uuid.NewString(), }) s.NoError(err) } if attemptCount == 2 { ctx, _ := rpc.NewContextWithTimeoutAndVersionHeaders(time.Second) _, err := env.FrontendClient().SignalWorkflowExecution(ctx, &workflowservice.SignalWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we.RunId, }, SignalName: "rejected-signal", Identity: identity, RequestId: uuid.NewString(), }) var resourceExhausted *serviceerror.ResourceExhausted s.ErrorAs(err, &resourceExhausted) s.Equal(consts.ErrWorkflowClosing.Cause, resourceExhausted.Cause) s.Equal(consts.ErrWorkflowClosing.Message, resourceExhausted.Message) } attemptCount++ return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{ Result: payloads.EncodeString("Done"), }}, }}, nil } poller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.Namespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: wtHandler, Logger: env.Logger, T: s.T(), } _, err = poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.Error(err) _, err = poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) } func (s *SignalWorkflowTestSuite) TestSignalExternalWorkflowCommand_WithoutRunID(opts []testcore.TestOption) { // Explicitly enable cross namespace commands for this test, // need a dedicated cluster to enable cross namespace commands opts = append( opts, testcore.WithDedicatedCluster(), testcore.WithDynamicConfig(dynamicconfig.EnableCrossNamespaceCommands, true), ) env := testcore.NewEnv(s.T(), opts...) id := "functional-signal-external-workflow-test-without-run-id" wt := "functional-signal-external-workflow-test-without-run-id-type" tl := "functional-signal-external-workflow-test-without-run-id-taskqueue" identity := "worker1" activityName := "activity_type1" workflowType := &commonpb.WorkflowType{Name: wt} taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL} request := &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), Identity: identity, } we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request) s.NoError(err0) env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId)) externalRequest := &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.ExternalNamespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), Identity: identity, } we2, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), externalRequest) s.NoError(err0) env.Logger.Info("StartWorkflowExecution on external Namespace", tag.WorkflowNamespace(env.ExternalNamespace().String()), tag.WorkflowRunID(we2.RunId)) activityCount := int32(1) activityCounter := int32(0) signalName := "my signal" signalInput := payloads.EncodeString("my signal input") wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { if activityCounter < activityCount { activityCounter++ buf := new(bytes.Buffer) s.NoError(binary.Write(buf, binary.LittleEndian, activityCounter)) return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK, Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{ ActivityId: convert.Int32ToString(activityCounter), ActivityType: &commonpb.ActivityType{Name: activityName}, TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, Input: payloads.EncodeBytes(buf.Bytes()), ScheduleToCloseTimeout: durationpb.New(100 * time.Second), ScheduleToStartTimeout: durationpb.New(10 * time.Second), StartToCloseTimeout: durationpb.New(50 * time.Second), HeartbeatTimeout: durationpb.New(5 * time.Second), }}, }}, nil } return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SIGNAL_EXTERNAL_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_SignalExternalWorkflowExecutionCommandAttributes{SignalExternalWorkflowExecutionCommandAttributes: &commandpb.SignalExternalWorkflowExecutionCommandAttributes{ Namespace: env.ExternalNamespace().String(), Execution: &commonpb.WorkflowExecution{ WorkflowId: id, // No RunID in command }, SignalName: signalName, Input: signalInput, }}, }}, nil } atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) { return payloads.EncodeString("Activity Result"), false, nil } poller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.Namespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: wtHandler, ActivityTaskHandler: atHandler, Logger: env.Logger, T: s.T(), } workflowComplete := false externalActivityCount := int32(1) externalActivityCounter := int32(0) var signalEvent *historypb.HistoryEvent externalWFTHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { if externalActivityCounter < externalActivityCount { externalActivityCounter++ buf := new(bytes.Buffer) s.NoError(binary.Write(buf, binary.LittleEndian, externalActivityCounter)) return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK, Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{ ActivityId: convert.Int32ToString(externalActivityCounter), ActivityType: &commonpb.ActivityType{Name: activityName}, TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, Input: payloads.EncodeBytes(buf.Bytes()), ScheduleToCloseTimeout: durationpb.New(100 * time.Second), ScheduleToStartTimeout: durationpb.New(10 * time.Second), StartToCloseTimeout: durationpb.New(50 * time.Second), HeartbeatTimeout: durationpb.New(5 * time.Second), }}, }}, nil } else if task.PreviousStartedEventId > 0 { for _, event := range task.History.Events[task.PreviousStartedEventId:] { if event.GetEventType() == enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_SIGNALED { signalEvent = event return []*commandpb.Command{}, nil } } } workflowComplete = true return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{ Result: payloads.EncodeString("Done"), }}, }}, nil } //nolint:staticcheck // SA1019 TaskPoller replacement needs to be done holistically. externalPoller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.ExternalNamespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: externalWFTHandler, ActivityTaskHandler: atHandler, Logger: env.Logger, T: s.T(), } // Start both current and external workflows to make some progress. _, err := poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) _, err = externalPoller.PollAndProcessWorkflowTask() env.Logger.Info("external PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) err = externalPoller.PollAndProcessActivityTask(false) env.Logger.Info("external PollAndProcessActivityTask", tag.Error(err)) s.NoError(err) // Signal the external workflow with this command. _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) // in source workflow var historyEvents []*historypb.HistoryEvent CheckHistoryLoopForSignalSent: for i := 1; i < 10; i++ { historyEvents = env.GetHistory(env.Namespace().String(), &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we.RunId, }) signalRequestedEvent := historyEvents[len(historyEvents)-2] if signalRequestedEvent.GetEventType() != enumspb.EVENT_TYPE_EXTERNAL_WORKFLOW_EXECUTION_SIGNALED { env.Logger.Info("Signal still not sent") time.Sleep(100 * time.Millisecond) //nolint:forbidigo continue CheckHistoryLoopForSignalSent } break } s.EqualHistoryEvents(fmt.Sprintf(` 1 WorkflowExecutionStarted 2 WorkflowTaskScheduled 3 WorkflowTaskStarted 4 WorkflowTaskCompleted 5 ActivityTaskScheduled 6 ActivityTaskTimedOut 7 WorkflowTaskScheduled 8 WorkflowTaskStarted 9 WorkflowTaskCompleted 10 SignalExternalWorkflowExecutionInitiated 11 ExternalWorkflowExecutionSignaled {"InitiatedEventId":10,"WorkflowExecution":{"RunId":"","WorkflowId":"%s"}} 12 WorkflowTaskScheduled`, id), historyEvents) // Process signal in workflow for external workflow _, err = externalPoller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) s.False(workflowComplete) s.NotNil(signalEvent) s.Equal(signalName, signalEvent.GetWorkflowExecutionSignaledEventAttributes().SignalName) s.ProtoEqual(signalInput, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Input) s.Equal("history-service", signalEvent.GetWorkflowExecutionSignaledEventAttributes().Identity) } func (s *SignalWorkflowTestSuite) TestSignalExternalWorkflowCommand_UnKnownTarget(opts []testcore.TestOption) { // Explicitly enable cross namespace commands for this test, // need a dedicated cluster to enable cross namespace commands opts = append( opts, testcore.WithDedicatedCluster(), testcore.WithDynamicConfig(dynamicconfig.EnableCrossNamespaceCommands, true), ) env := testcore.NewEnv(s.T(), opts...) id := "functional-signal-unknown-workflow-command-test" wt := "functional-signal-unknown-workflow-command-test-type" tl := "functional-signal-unknown-workflow-command-test-taskqueue" identity := "worker1" activityName := "activity_type1" workflowType := &commonpb.WorkflowType{Name: wt} taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL} request := &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), Identity: identity, } we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request) s.NoError(err0) env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId)) activityCount := int32(1) activityCounter := int32(0) signalName := "my signal" signalInput := payloads.EncodeString("my signal input") wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { if activityCounter < activityCount { activityCounter++ buf := new(bytes.Buffer) s.NoError(binary.Write(buf, binary.LittleEndian, activityCounter)) return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK, Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{ ActivityId: convert.Int32ToString(activityCounter), ActivityType: &commonpb.ActivityType{Name: activityName}, TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, Input: payloads.EncodeBytes(buf.Bytes()), ScheduleToCloseTimeout: durationpb.New(100 * time.Second), ScheduleToStartTimeout: durationpb.New(10 * time.Second), StartToCloseTimeout: durationpb.New(50 * time.Second), HeartbeatTimeout: durationpb.New(5 * time.Second), }}, }}, nil } return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SIGNAL_EXTERNAL_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_SignalExternalWorkflowExecutionCommandAttributes{SignalExternalWorkflowExecutionCommandAttributes: &commandpb.SignalExternalWorkflowExecutionCommandAttributes{ Namespace: env.ExternalNamespace().String(), Execution: &commonpb.WorkflowExecution{ WorkflowId: "workflow_not_exist", RunId: we.GetRunId(), }, SignalName: signalName, Input: signalInput, }}, }}, nil } atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) { return payloads.EncodeString("Activity Result"), false, nil } poller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.Namespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: wtHandler, ActivityTaskHandler: atHandler, Logger: env.Logger, T: s.T(), } // Start workflows to make some progress. _, err := poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) // Signal the external workflow with this command. _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) var historyEvents []*historypb.HistoryEvent CheckHistoryLoopForCancelSent: for i := 1; i < 10; i++ { historyEvents = env.GetHistory(env.Namespace().String(), &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we.RunId, }) signalFailedEvent := historyEvents[len(historyEvents)-2] if signalFailedEvent.GetEventType() != enumspb.EVENT_TYPE_SIGNAL_EXTERNAL_WORKFLOW_EXECUTION_FAILED { env.Logger.Info("Cancellaton not cancelled yet") time.Sleep(100 * time.Millisecond) //nolint:forbidigo continue CheckHistoryLoopForCancelSent } break } s.EqualHistoryEvents(fmt.Sprintf(` 1 WorkflowExecutionStarted 2 WorkflowTaskScheduled 3 WorkflowTaskStarted 4 WorkflowTaskCompleted 5 ActivityTaskScheduled 6 ActivityTaskTimedOut 7 WorkflowTaskScheduled 8 WorkflowTaskStarted 9 WorkflowTaskCompleted 10 SignalExternalWorkflowExecutionInitiated 11 SignalExternalWorkflowExecutionFailed {"InitiatedEventId":10,"WorkflowExecution":{"RunId":"%s","WorkflowId":"workflow_not_exist"}} 12 WorkflowTaskScheduled`, we.RunId), historyEvents) } func (s *SignalWorkflowTestSuite) TestSignalExternalWorkflowCommand_SignalSelf(opts []testcore.TestOption) { env := testcore.NewEnv(s.T(), opts...) id := "functional-signal-self-workflow-command-test" wt := "functional-signal-self-workflow-command-test-type" tl := "functional-signal-self-workflow-command-test-taskqueue" identity := "worker1" activityName := "activity_type1" workflowType := &commonpb.WorkflowType{Name: wt} taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL} request := &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), Identity: identity, } we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request) s.NoError(err0) env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId)) activityCount := int32(1) activityCounter := int32(0) signalName := "my signal" signalInput := payloads.EncodeString("my signal input") wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { if activityCounter < activityCount { activityCounter++ buf := new(bytes.Buffer) s.NoError(binary.Write(buf, binary.LittleEndian, activityCounter)) return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK, Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{ ActivityId: convert.Int32ToString(activityCounter), ActivityType: &commonpb.ActivityType{Name: activityName}, TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, Input: payloads.EncodeBytes(buf.Bytes()), ScheduleToCloseTimeout: durationpb.New(100 * time.Second), ScheduleToStartTimeout: durationpb.New(10 * time.Second), StartToCloseTimeout: durationpb.New(50 * time.Second), HeartbeatTimeout: durationpb.New(5 * time.Second), }}, }}, nil } return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SIGNAL_EXTERNAL_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_SignalExternalWorkflowExecutionCommandAttributes{SignalExternalWorkflowExecutionCommandAttributes: &commandpb.SignalExternalWorkflowExecutionCommandAttributes{ Namespace: env.Namespace().String(), Execution: &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we.GetRunId(), }, SignalName: signalName, Input: signalInput, }}, }}, nil } atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) { return payloads.EncodeString("Activity Result"), false, nil } poller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.Namespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: wtHandler, ActivityTaskHandler: atHandler, Logger: env.Logger, T: s.T(), } // Start workflows to make some progress. _, err := poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) // Signal the external workflow with this command. _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) var historyEvents []*historypb.HistoryEvent CheckHistoryLoopForCancelSent: for i := 1; i < 10; i++ { historyEvents = env.GetHistory(env.Namespace().String(), &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we.RunId, }) signalFailedEvent := historyEvents[len(historyEvents)-2] if signalFailedEvent.GetEventType() != enumspb.EVENT_TYPE_SIGNAL_EXTERNAL_WORKFLOW_EXECUTION_FAILED { env.Logger.Info("Cancellaton not cancelled yet") time.Sleep(100 * time.Millisecond) //nolint:forbidigo continue CheckHistoryLoopForCancelSent } break } s.EqualHistoryEvents(fmt.Sprintf(` 1 WorkflowExecutionStarted 2 WorkflowTaskScheduled 3 WorkflowTaskStarted 4 WorkflowTaskCompleted 5 ActivityTaskScheduled 6 ActivityTaskTimedOut 7 WorkflowTaskScheduled 8 WorkflowTaskStarted 9 WorkflowTaskCompleted 10 SignalExternalWorkflowExecutionInitiated 11 SignalExternalWorkflowExecutionFailed {"InitiatedEventId":10,"WorkflowExecution":{"RunId":"%s","WorkflowId":"%s"}} 12 WorkflowTaskScheduled`, we.RunId, id), historyEvents) } func (s *SignalWorkflowTestSuite) TestSignalWithStartWorkflow(opts []testcore.TestOption) { env := testcore.NewEnv(s.T(), opts...) id := "functional-signal-with-start-workflow-test" wt := "functional-signal-with-start-workflow-test-type" tl := "functional-signal-with-start-workflow-test-taskqueue" identity := "worker1" activityName := "activity_type1" workflowType := &commonpb.WorkflowType{Name: wt} taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL} header := &commonpb.Header{ Fields: map[string]*commonpb.Payload{"tracing": payload.EncodeString("sample data")}, } // Start a workflow request := &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), Identity: identity, } we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request) s.NoError(err0) env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId)) // workflow logic workflowComplete := false activityScheduled := false activityData := int32(1) newWorkflowStarted := false var signalEvent, startedEvent *historypb.HistoryEvent wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { if !activityScheduled { activityScheduled = true buf := new(bytes.Buffer) s.NoError(binary.Write(buf, binary.LittleEndian, activityData)) return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK, Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{ ActivityId: strconv.Itoa(1), ActivityType: &commonpb.ActivityType{Name: activityName}, TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, Input: payloads.EncodeBytes(buf.Bytes()), ScheduleToCloseTimeout: durationpb.New(100 * time.Second), ScheduleToStartTimeout: durationpb.New(2 * time.Second), StartToCloseTimeout: durationpb.New(50 * time.Second), HeartbeatTimeout: durationpb.New(5 * time.Second), }}, }}, nil } else if task.PreviousStartedEventId > 0 { for _, event := range task.History.Events[task.PreviousStartedEventId:] { if event.GetEventType() == enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_SIGNALED { signalEvent = event return []*commandpb.Command{}, nil } } } else if newWorkflowStarted { newWorkflowStarted = false signalEvent = nil startedEvent = nil for _, event := range task.History.Events[task.PreviousStartedEventId:] { if event.GetEventType() == enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_SIGNALED { signalEvent = event } if event.GetEventType() == enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_STARTED { startedEvent = event } } if signalEvent != nil && startedEvent != nil { return []*commandpb.Command{}, nil } } workflowComplete = true return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{ Result: payloads.EncodeString("Done"), }}, }}, nil } // activity handler atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) { return payloads.EncodeString("Activity Result"), false, nil } poller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.Namespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: wtHandler, ActivityTaskHandler: atHandler, Logger: env.Logger, T: s.T(), } // Make first command to schedule activity _, err := poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) // Send a signal signalName := "my signal" signalInput := payloads.EncodeString("my signal input") wfIDReusePolicy := enumspb.WORKFLOW_ID_REUSE_POLICY_ALLOW_DUPLICATE sRequest := &workflowservice.SignalWithStartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, Header: header, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), SignalName: signalName, SignalInput: signalInput, Identity: identity, WorkflowIdReusePolicy: wfIDReusePolicy, } resp, err := env.FrontendClient().SignalWithStartWorkflowExecution(s.Context(), sRequest) s.NoError(err) s.False(resp.Started) s.Equal(we.GetRunId(), resp.GetRunId()) // Process signal in workflow _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) s.False(workflowComplete) s.NotNil(signalEvent) s.Equal(signalName, signalEvent.GetWorkflowExecutionSignaledEventAttributes().SignalName) s.ProtoEqual(signalInput, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Input) s.Equal(identity, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Identity) // Terminate workflow execution _, err = env.FrontendClient().TerminateWorkflowExecution(s.Context(), &workflowservice.TerminateWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, }, Reason: "test signal", Details: nil, Identity: identity, }) s.NoError(err) // Send signal to terminated workflow signalName = "signal to terminate" signalInput = payloads.EncodeString("signal to terminate input") sRequest.SignalName = signalName sRequest.SignalInput = signalInput sRequest.WorkflowId = id resp, err = env.FrontendClient().SignalWithStartWorkflowExecution(s.Context(), sRequest) s.NoError(err) s.True(resp.Started) s.NotNil(resp.GetRunId()) s.NotEqual(we.GetRunId(), resp.GetRunId()) newWorkflowStarted = true // Process signal in workflow _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) s.False(workflowComplete) s.NotNil(signalEvent) s.Equal(signalName, signalEvent.GetWorkflowExecutionSignaledEventAttributes().SignalName) s.ProtoEqual(signalInput, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Input) s.Equal(identity, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Identity) s.NotNil(startedEvent) s.ProtoEqual(header, startedEvent.GetWorkflowExecutionStartedEventAttributes().Header) // Send signal to not existed workflow id = "functional-signal-with-start-workflow-test-non-exist" signalName = "signal to non exist" signalInput = payloads.EncodeString("signal to non exist input") sRequest.SignalName = signalName sRequest.SignalInput = signalInput sRequest.WorkflowId = id resp, err = env.FrontendClient().SignalWithStartWorkflowExecution(s.Context(), sRequest) s.NoError(err) s.NotNil(resp.GetRunId()) s.True(resp.Started) newWorkflowStarted = true // Process signal in workflow _, err = poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) s.False(workflowComplete) s.NotNil(signalEvent) s.Equal(signalName, signalEvent.GetWorkflowExecutionSignaledEventAttributes().SignalName) s.ProtoEqual(signalInput, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Input) s.Equal(identity, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Identity) listOpenRequest := &workflowservice.ListOpenWorkflowExecutionsRequest{ Namespace: env.Namespace().String(), MaximumPageSize: 100, StartTimeFilter: &filterpb.StartTimeFilter{ EarliestTime: nil, LatestTime: timestamppb.New(time.Now().UTC()), }, Filters: &workflowservice.ListOpenWorkflowExecutionsRequest_ExecutionFilter{ ExecutionFilter: &filterpb.WorkflowExecutionFilter{ WorkflowId: id, }, }, } // Assert visibility is correct s.Eventually( func() bool { listResp, err := env.FrontendClient().ListOpenWorkflowExecutions(s.Context(), listOpenRequest) s.NoError(err) return len(listResp.Executions) == 1 }, testcore.WaitForESToSettle, 100*time.Millisecond, ) // Terminate workflow execution and assert visibility is correct _, err = env.FrontendClient().TerminateWorkflowExecution(s.Context(), &workflowservice.TerminateWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, }, Reason: "kill workflow", Details: nil, Identity: identity, }) s.NoError(err) s.Eventually( func() bool { listResp, err := env.FrontendClient().ListOpenWorkflowExecutions(s.Context(), listOpenRequest) s.NoError(err) return len(listResp.Executions) == 0 }, testcore.WaitForESToSettle, 100*time.Millisecond, ) listClosedRequest := &workflowservice.ListClosedWorkflowExecutionsRequest{ Namespace: env.Namespace().String(), MaximumPageSize: 100, StartTimeFilter: &filterpb.StartTimeFilter{ EarliestTime: nil, LatestTime: timestamppb.New(time.Now().UTC()), }, Filters: &workflowservice.ListClosedWorkflowExecutionsRequest_ExecutionFilter{ExecutionFilter: &filterpb.WorkflowExecutionFilter{ WorkflowId: id, }}, } listClosedResp, err := env.FrontendClient().ListClosedWorkflowExecutions(s.Context(), listClosedRequest) s.NoError(err) s.Len(listClosedResp.Executions, 1) } func (s *SignalWorkflowTestSuite) TestSignalWithStartWorkflow_ResolveIDDeduplication(opts []testcore.TestOption) { env := testcore.NewEnv(s.T(), opts...) // setting this to 0 to be sure we are terminating the current workflow env.OverrideDynamicConfig(dynamicconfig.WorkflowIdReuseMinimalInterval, 0) id := "functional-signal-with-start-workflow-id-reuse-test" wt := "functional-signal-with-start-workflow-id-reuse-test-type" tl := "functional-signal-with-start-workflow-id-reuse-test-taskqueue" identity := "worker1" activityName := "activity_type1" workflowType := &commonpb.WorkflowType{Name: wt} taskQueue := &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL} // Start a workflow request := &workflowservice.StartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), Identity: identity, } we, err0 := env.FrontendClient().StartWorkflowExecution(s.Context(), request) s.NoError(err0) env.Logger.Info("StartWorkflowExecution", tag.WorkflowRunID(we.RunId)) workflowComplete := false activityCount := int32(1) activityCounter := int32(0) wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { if activityCounter < activityCount { activityCounter++ buf := new(bytes.Buffer) s.NoError(binary.Write(buf, binary.LittleEndian, activityCounter)) return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_SCHEDULE_ACTIVITY_TASK, Attributes: &commandpb.Command_ScheduleActivityTaskCommandAttributes{ScheduleActivityTaskCommandAttributes: &commandpb.ScheduleActivityTaskCommandAttributes{ ActivityId: convert.Int32ToString(activityCounter), ActivityType: &commonpb.ActivityType{Name: activityName}, TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, Input: payloads.EncodeBytes(buf.Bytes()), ScheduleToCloseTimeout: durationpb.New(100 * time.Second), ScheduleToStartTimeout: durationpb.New(10 * time.Second), StartToCloseTimeout: durationpb.New(50 * time.Second), HeartbeatTimeout: durationpb.New(5 * time.Second), }}, }}, nil } workflowComplete = true return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{ Result: payloads.EncodeString("Done"), }}, }}, nil } atHandler := func(task *workflowservice.PollActivityTaskQueueResponse) (*commonpb.Payloads, bool, error) { return payloads.EncodeString("Activity Result"), false, nil } poller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.Namespace().String(), TaskQueue: taskQueue, Identity: identity, WorkflowTaskHandler: wtHandler, ActivityTaskHandler: atHandler, Logger: env.Logger, T: s.T(), } // Start workflows, make some progress and complete workflow _, err := poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) _, err = poller.PollAndProcessWorkflowTask() env.Logger.Info("PollAndProcessWorkflowTask", tag.Error(err)) s.NoError(err) s.True(workflowComplete) // test WorkflowIdReusePolicy: RejectDuplicate signalName := "my signal" signalInput := payloads.EncodeString("my signal input") sRequest := &workflowservice.SignalWithStartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: workflowType, TaskQueue: taskQueue, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), SignalName: signalName, SignalInput: signalInput, Identity: identity, WorkflowIdReusePolicy: enumspb.WORKFLOW_ID_REUSE_POLICY_REJECT_DUPLICATE, } ctx, _ := rpc.NewContextWithTimeoutAndVersionHeaders(5 * time.Second) resp, err := env.FrontendClient().SignalWithStartWorkflowExecution(ctx, sRequest) s.Nil(resp) s.Error(err) s.Contains(err.Error(), "reject duplicate workflow Id") s.IsType(&serviceerror.WorkflowExecutionAlreadyStarted{}, err) // test WorkflowIdReusePolicy: AllowDuplicateFailedOnly sRequest.WorkflowIdReusePolicy = enumspb.WORKFLOW_ID_REUSE_POLICY_ALLOW_DUPLICATE_FAILED_ONLY ctx, _ = rpc.NewContextWithTimeoutAndVersionHeaders(5 * time.Second) resp, err = env.FrontendClient().SignalWithStartWorkflowExecution(ctx, sRequest) s.Nil(resp) s.Error(err) s.Contains(err.Error(), "allow duplicate workflow Id if last run failed") s.IsType(&serviceerror.WorkflowExecutionAlreadyStarted{}, err) // test WorkflowIdReusePolicy: AllowDuplicate sRequest.WorkflowIdReusePolicy = enumspb.WORKFLOW_ID_REUSE_POLICY_ALLOW_DUPLICATE ctx, _ = rpc.NewContextWithTimeoutAndVersionHeaders(5 * time.Second) resp, err = env.FrontendClient().SignalWithStartWorkflowExecution(ctx, sRequest) s.NoError(err) s.NotEmpty(resp.GetRunId()) s.True(resp.Started) // Terminate workflow execution _, err = env.FrontendClient().TerminateWorkflowExecution(s.Context(), &workflowservice.TerminateWorkflowExecutionRequest{ Namespace: env.Namespace().String(), WorkflowExecution: &commonpb.WorkflowExecution{ WorkflowId: id, }, Reason: "test WorkflowIdReusePolicyAllowDuplicateFailedOnly", Details: nil, Identity: identity, }) s.NoError(err) // test WorkflowIdReusePolicy: AllowDuplicateFailedOnly sRequest.WorkflowIdReusePolicy = enumspb.WORKFLOW_ID_REUSE_POLICY_ALLOW_DUPLICATE_FAILED_ONLY resp, err = env.FrontendClient().SignalWithStartWorkflowExecution(s.Context(), sRequest) s.NoError(err) s.NotEmpty(resp.GetRunId()) s.True(resp.Started) // test WorkflowIdReusePolicy: TerminateIfRunning (for backwards compatibility) prevRunID := resp.RunId sRequest.WorkflowIdReusePolicy = enumspb.WORKFLOW_ID_REUSE_POLICY_TERMINATE_IF_RUNNING resp, err = env.FrontendClient().SignalWithStartWorkflowExecution(s.Context(), sRequest) s.NoError(err) s.NotEmpty(resp.GetRunId()) s.NotEqual(prevRunID, resp.GetRunId()) s.True(resp.Started) descResp, err := env.FrontendClient().DescribeWorkflowExecution(s.Context(), &workflowservice.DescribeWorkflowExecutionRequest{ Namespace: env.Namespace().String(), Execution: &commonpb.WorkflowExecution{WorkflowId: id, RunId: prevRunID}, }) s.NoError(err) s.Equal(enumspb.WORKFLOW_EXECUTION_STATUS_TERMINATED, descResp.WorkflowExecutionInfo.Status) // test WorkflowIdConflictPolicy: TerminateExisting (replaced TerminateIfRunning) prevRunID = resp.RunId sRequest.WorkflowIdReusePolicy = enumspb.WORKFLOW_ID_REUSE_POLICY_ALLOW_DUPLICATE sRequest.WorkflowIdConflictPolicy = enumspb.WORKFLOW_ID_CONFLICT_POLICY_TERMINATE_EXISTING resp, err = env.FrontendClient().SignalWithStartWorkflowExecution(s.Context(), sRequest) s.NoError(err) s.NotEmpty(resp.GetRunId()) s.NotEqual(prevRunID, resp.GetRunId()) s.True(resp.Started) descResp, err = env.FrontendClient().DescribeWorkflowExecution(s.Context(), &workflowservice.DescribeWorkflowExecutionRequest{ Namespace: env.Namespace().String(), Execution: &commonpb.WorkflowExecution{WorkflowId: id, RunId: prevRunID}, }) s.NoError(err) s.Equal(enumspb.WORKFLOW_EXECUTION_STATUS_TERMINATED, descResp.WorkflowExecutionInfo.Status) descResp, err = env.FrontendClient().DescribeWorkflowExecution(s.Context(), &workflowservice.DescribeWorkflowExecutionRequest{ Namespace: env.Namespace().String(), Execution: &commonpb.WorkflowExecution{ WorkflowId: id, RunId: resp.GetRunId(), }, }) s.NoError(err) s.Equal(enumspb.WORKFLOW_EXECUTION_STATUS_RUNNING, descResp.WorkflowExecutionInfo.Status) } func (s *SignalWorkflowTestSuite) TestSignalWithStartWorkflow_StartDelay(opts []testcore.TestOption) { env := testcore.NewEnv(s.T(), opts...) id := "functional-signal-with-start-workflow-start-delay-test" wt := "functional-signal-with-start-workflow-start-delay-test-type" tl := "functional-signal-with-start-workflow-start-delay-test-taskqueue" stickyTq := "functional-signal-with-start-workflow-start-delay-test-sticky-taskqueue" identity := "worker1" startDelay := 3 * time.Second signalName := "my signal" signalInput := payloads.EncodeString("my signal input") sRequest := &workflowservice.SignalWithStartWorkflowExecutionRequest{ RequestId: uuid.NewString(), Namespace: env.Namespace().String(), WorkflowId: id, WorkflowType: &commonpb.WorkflowType{Name: wt}, TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, Input: nil, WorkflowRunTimeout: durationpb.New(100 * time.Second), WorkflowTaskTimeout: durationpb.New(1 * time.Second), SignalName: signalName, SignalInput: signalInput, Identity: identity, WorkflowStartDelay: durationpb.New(startDelay), } reqStartTime := time.Now() we0, startErr := env.FrontendClient().SignalWithStartWorkflowExecution(s.Context(), sRequest) s.NoError(startErr) var signalEvent *historypb.HistoryEvent delayEndTime := time.Now() wtHandler := func(task *workflowservice.PollWorkflowTaskQueueResponse) ([]*commandpb.Command, error) { delayEndTime = time.Now() for _, event := range task.History.Events[task.PreviousStartedEventId:] { if event.GetEventType() == enumspb.EVENT_TYPE_WORKFLOW_EXECUTION_SIGNALED { signalEvent = event } } return []*commandpb.Command{{ CommandType: enumspb.COMMAND_TYPE_COMPLETE_WORKFLOW_EXECUTION, Attributes: &commandpb.Command_CompleteWorkflowExecutionCommandAttributes{CompleteWorkflowExecutionCommandAttributes: &commandpb.CompleteWorkflowExecutionCommandAttributes{ Result: payloads.EncodeString("Done"), }}, }}, nil } poller := &testcore.TaskPoller{ Client: env.FrontendClient(), Namespace: env.Namespace().String(), TaskQueue: &taskqueuepb.TaskQueue{Name: tl, Kind: enumspb.TASK_QUEUE_KIND_NORMAL}, StickyTaskQueue: &taskqueuepb.TaskQueue{Name: stickyTq, Kind: enumspb.TASK_QUEUE_KIND_STICKY, NormalName: tl}, Identity: identity, WorkflowTaskHandler: wtHandler, Logger: env.Logger, T: s.T(), } _, pollErr := poller.PollAndProcessWorkflowTask(testcore.WithDumpHistory) s.NoError(pollErr) s.GreaterOrEqual(delayEndTime.Sub(reqStartTime), startDelay) s.NotNil(signalEvent) s.Equal(signalName, signalEvent.GetWorkflowExecutionSignaledEventAttributes().SignalName) s.ProtoEqual(signalInput, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Input) s.Equal(identity, signalEvent.GetWorkflowExecutionSignaledEventAttributes().Identity) descResp, descErr := env.FrontendClient().DescribeWorkflowExecution(s.Context(), &workflowservice.DescribeWorkflowExecutionRequest{ Namespace: env.Namespace().String(), Execution: &commonpb.WorkflowExecution{ WorkflowId: id, RunId: we0.RunId, }, }) s.NoError(descErr) s.Equal(enumspb.WORKFLOW_EXECUTION_STATUS_COMPLETED, descResp.WorkflowExecutionInfo.Status) }