Skip to content
Draft
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
3 changes: 2 additions & 1 deletion internal/signals/activities.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package signals

import (
"context"
"encoding/json"
)

type Signaler interface {
Expand All @@ -12,6 +13,6 @@ type Activities struct {
Signaler Signaler
}

func (a *Activities) DeliverWorkflowSignal(ctx context.Context, instanceID, signalName string, arg interface{}) error {
func (a *Activities) DeliverWorkflowSignal(ctx context.Context, instanceID, signalName string, arg json.RawMessage) error {
return a.Signaler.SignalWorkflow(ctx, instanceID, signalName, arg)
}
49 changes: 49 additions & 0 deletions tester/tester_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -241,6 +241,55 @@ func Test_SignalSubWorkflow(t *testing.T) {
require.Equal(t, 42, wfR)
}

func Test_SignalWorkflow_LargeIntegerPrecision(t *testing.T) {
type signalPayload struct {
ID uint64
}

const (
receiverInstanceID = "receiver"
signalChannel = "test-signal"
largeIntegerID = uint64(123456789012345678)
)

testPayload := signalPayload{ID: largeIntegerID}

var received signalPayload

receiverWorkflow := func(ctx workflow.Context) error {
ch := workflow.NewSignalChannel[signalPayload](ctx, signalChannel)
payload, ok := ch.Receive(ctx)
if !ok {
return errors.New("signal channel unexpectedly closed")
}
received = payload

return nil
}

senderWorkflow := func(ctx workflow.Context) error {
sub := workflow.CreateSubWorkflowInstance[any](ctx, workflow.SubWorkflowOptions{
InstanceID: receiverInstanceID,
}, receiverWorkflow)

if _, err := workflow.SignalWorkflow(ctx, receiverInstanceID, signalChannel, testPayload).Get(ctx); err != nil {
return err
}

_, err := sub.Get(ctx)
return err
}

wt := NewWorkflowTester[any](senderWorkflow, WithTestTimeout(5*time.Second))
wt.Registry().RegisterWorkflow(receiverWorkflow)
wt.Execute(context.Background())

require.True(t, wt.WorkflowFinished())
_, err := wt.WorkflowResult()
require.NoError(t, err)
require.Equal(t, testPayload.ID, received.ID)
}

func workflowSubworkflowSignal(ctx workflow.Context) (int, error) {
sw := workflow.CreateSubWorkflowInstance[int](ctx, workflow.SubWorkflowOptions{
InstanceID: "subworkflow",
Expand Down
14 changes: 13 additions & 1 deletion workflow/signal.go
Original file line number Diff line number Diff line change
@@ -1,9 +1,14 @@
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"
Expand All @@ -24,11 +29,18 @@ func SignalWorkflow[T any](ctx Context, instanceID string, name string, arg T) F
)
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, arg)
}, a.DeliverWorkflowSignal, instanceID, name, json.RawMessage(input))
}