-
Notifications
You must be signed in to change notification settings - Fork 94
Expand file tree
/
Copy pathsignal.go
More file actions
46 lines (40 loc) · 1.41 KB
/
Copy pathsignal.go
File metadata and controls
46 lines (40 loc) · 1.41 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
package workflow
import (
"encoding/json"
"fmt"
"github.com/cschleiden/go-workflows/core"
"github.com/cschleiden/go-workflows/internal/contextvalue"
"github.com/cschleiden/go-workflows/internal/log"
"github.com/cschleiden/go-workflows/internal/signals"
"github.com/cschleiden/go-workflows/internal/sync"
"github.com/cschleiden/go-workflows/internal/workflowstate"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/trace"
)
// NewSignalChannel returns a new signal channel.
func NewSignalChannel[T any](ctx Context, name string) Channel[T] {
wfState := workflowstate.WorkflowState(ctx)
return workflowstate.GetSignalChannel[T](ctx, wfState, name)
}
// SignalWorkflow sends a signal to another running workflow instance.
func SignalWorkflow[T any](ctx Context, instanceID string, name string, arg T) Future[any] {
ctx, span := Tracer(ctx).Start(ctx, "SignalWorkflow",
trace.WithAttributes(
attribute.String(log.SignalNameKey, name),
),
)
defer span.End()
input, err := contextvalue.Converter(ctx).To(arg)
if err != nil {
f := sync.NewFuture[any]()
f.Set(nil, fmt.Errorf("converting signal input for workflow %s signal %s: %w", instanceID, name, err))
return f
}
var a *signals.Activities
return ExecuteActivity[any](ctx, ActivityOptions{
RetryOptions: RetryOptions{
MaxAttempts: 1,
},
Queue: core.QueueSystem,
}, a.DeliverWorkflowSignal, instanceID, name, json.RawMessage(input))
}