Skip to content

Async activity should respect the context cancellation #468

Description

@raymondsze
package main

import (
	"context"
	"fmt"
	"net/http"
	"os"
	"os/signal"
	"time"

	"github.com/cschleiden/go-workflows/backend"
	"github.com/cschleiden/go-workflows/backend/sqlite"
	"github.com/cschleiden/go-workflows/client"
	"github.com/cschleiden/go-workflows/diag"
	"github.com/cschleiden/go-workflows/registry"
	"github.com/cschleiden/go-workflows/worker"
	"github.com/cschleiden/go-workflows/workflow"
	"github.com/google/uuid"
	"github.com/huandu/go-clone"
	"google.golang.org/grpc/codes"
	"google.golang.org/grpc/status"
	"google.golang.org/protobuf/types/known/wrapperspb"
)

func Activity1(ctx context.Context, a *wrapperspb.Int64Value) (*wrapperspb.Int64Value, error) {
	time.Sleep(time.Second)
	fmt.Println("Activity 1 Executed With Params", a)
	return nil, workflow.NewPermanentError(status.Errorf(codes.InvalidArgument, "invalid argument"))
}

func Activity2(ctx context.Context, a *wrapperspb.Int64Value) (*wrapperspb.Int64Value, error) {
	time.Sleep(time.Second)
	fmt.Println("Activity 2 Executed With Params", a)
	return a, nil
}

func Workflow(ctx workflow.Context, input string) error {
	fmt.Println("Workflow Executing With Params", input)
	wg := workflow.NewWaitGroup()
	ctx, cancel := workflow.WithCancelCause(ctx)

	for i := range 20 {
		wg.Add(1)
		workflow.Go(ctx, func(ctx workflow.Context) {
			defer wg.Done()
			activityOptions := clone.Clone(workflow.DefaultActivityOptions).(workflow.ActivityOptions)
			activityOptions.Queue = "raymond.queues.activity1"
			activityOptions.RetryOptions.RetryTimeout = 60 * time.Second
			result, err := workflow.ExecuteActivity[*wrapperspb.Int64Value](ctx, activityOptions, "raymond.activities.activity1", wrapperspb.Int64(int64(i))).Get(ctx)
			if err != nil {
				fmt.Println("Workflow Get Error From Activity 1", err)
				cancel(err)
				return
			}
			fmt.Println("Workflow Get Result From Activity 1", result.GetValue())
		})
	}
	wg.Wait(ctx)
	fmt.Println("Workflow Activities 1 Done")

	if ctx.Err() != nil {
		return ctx.Err()
	}

	wg = workflow.NewWaitGroup()
	for i := range 20 {
		wg.Add(1)
		workflow.Go(ctx, func(ctx workflow.Context) {
			defer wg.Done()
			activityOptions := clone.Clone(workflow.DefaultActivityOptions).(workflow.ActivityOptions)
			activityOptions.Queue = "raymond.queues.activity2"
			result, err := workflow.ExecuteActivity[*wrapperspb.Int64Value](ctx, activityOptions, "raymond.activities.activity2", wrapperspb.Int64(int64(i))).Get(ctx)
			if err != nil {
				cancel(err)
				return
			}
			fmt.Println("Workflow Get Result From Activity 2", result.GetValue())
		})
	}
	wg.Wait(ctx)
	fmt.Println("Workflow Activities 2 Done")

	if ctx.Err() != nil {
		return ctx.Err()
	}

	return nil
}

func runWorker(ctx context.Context, mb backend.Backend) {
	workerOptions := clone.Clone(worker.DefaultOptions).(worker.Options)
	workerOptions.WorkflowQueues = []workflow.Queue{"raymond.queues.example"}
	workflowWorker := worker.New(mb, &workerOptions)
	workflowWorker.RegisterWorkflow(Workflow, registry.WithName("raymond.workflows.example"))

	activityWorker1Options := clone.Clone(worker.DefaultOptions).(worker.Options)
	activityWorker1Options.ActivityQueues = []workflow.Queue{"raymond.queues.activity1"}
	activityWorker1Options.MaxParallelActivityTasks = 5
	activityWorker1 := worker.New(mb, &activityWorker1Options)
	activityWorker1.RegisterActivity(Activity1, registry.WithName("raymond.activities.activity1"))

	activityWorker2Options := clone.Clone(worker.DefaultOptions).(worker.Options)
	activityWorker2Options.ActivityQueues = []workflow.Queue{"raymond.queues.activity2"}
	activityWorker2Options.MaxParallelActivityTasks = 1
	activityWorker2 := worker.New(mb, &activityWorker2Options)
	activityWorker2.RegisterActivity(Activity2, registry.WithName("raymond.activities.activity2"))

	if err := workflowWorker.Start(ctx); err != nil {
		panic("could not start worker")
	}

	if err := activityWorker1.Start(ctx); err != nil {
		panic("could not start activity worker 1")
	}

	if err := activityWorker2.Start(ctx); err != nil {
		panic("could not start activity worker 2")
	}
}

func main() {
	ctx := context.Background()
	b := sqlite.NewSqliteBackend("simple.sqlite")

	m := http.NewServeMux()
	m.Handle("/diag/", http.StripPrefix("/diag", diag.NewServeMux(b)))
	go http.ListenAndServe(":8000", m)

	go runWorker(ctx, b)

	c := client.New(b)

	instance, err := c.CreateWorkflowInstance(ctx, client.WorkflowInstanceOptions{
		InstanceID: uuid.NewString(),
		Queue:      "raymond.queues.example",
	}, "raymond.workflows.example", "input-for-workflow")
	if err != nil {
		panic(err)
	}

	if err := c.WaitForWorkflowInstance(ctx, instance, time.Hour); err != nil {
		panic(err)
	}

	result, err := client.GetWorkflowResult[*wrapperspb.Int64Value](ctx, c, instance, time.Hour)
	if err != nil {
		panic(err)
	}

	fmt.Println("Workflow Result", result)

	c2 := make(chan os.Signal, 1)
	signal.Notify(c2, os.Interrupt)
	<-c2
}

Output:

Workflow Executing With Params input-for-workflow
Activity 1 Executed With Params value:15
Activity 1 Executed With Params value:16
Activity 1 Executed With Params value:17
Activity 1 Executed With Params value:19
Activity 1 Executed With Params value:18
Workflow Executing With Params input-for-workflow
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Activity 1 Executed With Params 
Activity 1 Executed With Params value:1
Activity 1 Executed With Params value:2
Activity 1 Executed With Params value:3
Activity 1 Executed With Params value:4
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Activity 1 Executed With Params value:5
Activity 1 Executed With Params value:6
Activity 1 Executed With Params value:7
Activity 1 Executed With Params value:8
Activity 1 Executed With Params value:9
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Activity 1 Executed With Params value:10
Activity 1 Executed With Params value:11
Activity 1 Executed With Params value:13
Activity 1 Executed With Params value:12
Activity 1 Executed With Params value:14
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Activity 1 Executed With Params value:15
Activity 1 Executed With Params value:16
Activity 1 Executed With Params value:18
Activity 1 Executed With Params value:17
Activity 1 Executed With Params value:19
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Get Error From Activity 1 rpc error: code = InvalidArgument desc = invalid argument
Workflow Activities 1 Done
panic: context canceled

Expectation:
Activity 1 Executed With Params value:5 to Activity 1 Executed With Params value:19 should not be executed due to context cancellation.

Metadata

Metadata

Assignees

No one assigned

    Labels

    featureNew feature or request

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions