Skip to content

Restructure workflows for replication simulation #6936

New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

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
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ domains:
operations:
- op: start_workflow
at: 0s
workflowType: test-workflow
workflowType: timer-activity-loop-workflow
workflowID: wf1
cluster: cluster0
domain: test-domain-aa
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ domains:
operations:
- op: start_workflow
at: 0s
workflowType: test-workflow
workflowType: timer-activity-loop-workflow
workflowID: wf1
cluster: cluster0
domain: test-domain
Expand Down
5 changes: 3 additions & 2 deletions simulation/replication/worker/cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ import (

"github.com/uber/cadence/common"
simTypes "github.com/uber/cadence/simulation/replication/types"
"github.com/uber/cadence/simulation/replication/workflows"
)

var (
Expand Down Expand Up @@ -142,11 +143,11 @@ func main() {
workerOptions,
)

for name, wf := range workflows {
for name, wf := range workflows.Workflows {
w.RegisterWorkflowWithOptions(wf, workflow.RegisterOptions{Name: name})
}

for name, act := range activities {
for name, act := range workflows.Activities {
w.RegisterActivityWithOptions(act, activity.RegisterOptions{Name: name})
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@
// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
// THE SOFTWARE.

package main
package timeractivityloop

import (
"context"
Expand All @@ -31,33 +31,28 @@ import (
"github.com/uber/cadence/simulation/replication/types"
)

var (
workflows = map[string]any{"test-workflow": TestWorkflow}
activities = map[string]any{"test-activity": TestActivity}
)

func TestWorkflow(ctx workflow.Context, input types.WorkflowInput) (types.WorkflowOutput, error) {
func Workflow(ctx workflow.Context, input types.WorkflowInput) (types.WorkflowOutput, error) {
logger := workflow.GetLogger(ctx)
logger.Sugar().Infof("testWorkflow started with input: %+v", input)
logger.Sugar().Infof("timer-activity-loop-workflow started with input: %+v", input)

endTime := workflow.Now(ctx).Add(input.Duration)
count := 0
for {
logger.Sugar().Infof("testWorkflow iteration %d", count)
logger.Sugar().Infof("timer-activity-loop-workflow iteration %d", count)
selector := workflow.NewSelector(ctx)
activityFuture := workflow.ExecuteActivity(workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
TaskList: types.TasklistName,
ScheduleToStartTimeout: 10 * time.Second,
StartToCloseTimeout: 10 * time.Second,
}), TestActivity, "World")
}), FormatStringActivity, "World")
selector.AddFuture(activityFuture, func(f workflow.Future) {
logger.Info("testWorkflow completed activity")
logger.Info("timer-activity-loop-workflow completed activity")
})

// use timer future to send notification email if processing takes too long
timerFuture := workflow.NewTimer(ctx, types.TimerInterval)
selector.AddFuture(timerFuture, func(f workflow.Future) {
logger.Info("testWorkflow timer fired")
logger.Info("timer-activity-loop-workflow timer fired")
})

// wait for both activity and timer to complete
Expand All @@ -67,19 +62,19 @@ func TestWorkflow(ctx workflow.Context, input types.WorkflowInput) (types.Workfl

now := workflow.Now(ctx)
if now.Before(endTime) {
logger.Sugar().Infof("testWorkflow will continue iteration because [now %v] < [endTime %v]", now, endTime)
logger.Sugar().Infof("timer-activity-loop-workflow will continue iteration because [now %v] < [endTime %v]", now, endTime)
} else {
logger.Sugar().Infof("testWorkflow will exit because [now %v] >= [endTime %v]", now, endTime)
logger.Sugar().Infof("timer-activity-loop-workflow will exit because [now %v] >= [endTime %v]", now, endTime)
break
}
}

logger.Info("testWorkflow completed")
logger.Info("timer-activity-loop-workflow completed")
return types.WorkflowOutput{Count: count}, nil
}

func TestActivity(ctx context.Context, input string) (string, error) {
func FormatStringActivity(ctx context.Context, input string) (string, error) {
logger := activity.GetLogger(ctx)
logger.Info("testActivity started")
logger.Info("timer-activity-loop-workflow format-string-activity started")
return fmt.Sprintf("Hello, %s!", input), nil
}
15 changes: 15 additions & 0 deletions simulation/replication/workflows/workflows.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
package workflows

import (
timeractivityloop "github.com/uber/cadence/simulation/replication/workflows/timer_activity_loop"
)

// Add workflows and activities to this map to register them with the worker.
var (
Workflows = map[string]any{
"timer-activity-loop-workflow": timeractivityloop.Workflow,
}
Activities = map[string]any{
"timer-activity-loop-format-string-activity": timeractivityloop.FormatStringActivity,
}
)