diff --git a/pkg/workflows/host/requirement_selecting_module.go b/pkg/workflows/host/requirement_selecting_module.go index 3ce23f7e49..93b7555ec5 100644 --- a/pkg/workflows/host/requirement_selecting_module.go +++ b/pkg/workflows/host/requirement_selecting_module.go @@ -9,6 +9,10 @@ import ( "github.com/smartcontractkit/chainlink-protos/cre/go/sdk" ) +// ErrRunnerUnavailable means a runner of the required type exists but cannot yet +// satisfy the requirements (e.g. a gated capability); callers may retry, not fail. +var ErrRunnerUnavailable = errors.New("no runner can currently satisfy the trigger requirements") + type ModuleAndHandler struct { Module RequirementsHandler @@ -109,6 +113,9 @@ func (r *requirementSelectingModule) subscribe(ctx context.Context, request *sdk } } if !matched { + if r.runnerTypeAvailable(sub.Requirements) { + return nil, fmt.Errorf("%w for trigger %d", ErrRunnerUnavailable, i) + } return nil, fmt.Errorf("cannot find a runner that can satisfy the requirements for trigger %d", i) } } @@ -116,6 +123,17 @@ func (r *requirementSelectingModule) subscribe(ctx context.Context, request *sdk return result, nil } +// runnerTypeAvailable reports whether any module handles req's requirement types, +// even if it can't currently satisfy them: a present-but-unavailable runner vs none. +func (r *requirementSelectingModule) runnerTypeAvailable(req *sdk.Requirements) bool { + for _, m := range r.modules { + if HandlesRequirements(m.RequirementsHandler, req) { + return true + } + } + return false +} + func (r *requirementSelectingModule) trigger(ctx context.Context, request *sdk.ExecuteRequest, handler ExecutionHelper) (*sdk.ExecutionResult, error) { trigger := request.GetTrigger() if val, cached := r.cache.Load(trigger.Id); cached { diff --git a/pkg/workflows/host/requirement_selecting_module_test.go b/pkg/workflows/host/requirement_selecting_module_test.go index f91f539046..83e96feb15 100644 --- a/pkg/workflows/host/requirement_selecting_module_test.go +++ b/pkg/workflows/host/requirement_selecting_module_test.go @@ -237,7 +237,7 @@ func TestRequirementSelectingModule_Execute(t *testing.T) { _, err := m.Execute(t.Context(), subscribeRequest(), nil) require.Error(t, err) - assert.Contains(t, err.Error(), "cannot find a runner that can satisfy the requirements") + assert.ErrorIs(t, err, ErrRunnerUnavailable) }) t.Run("subscribe skips non-matching and selects later additional", func(t *testing.T) { @@ -542,6 +542,7 @@ func TestRequirementSelectingModule_TriggerCache(t *testing.T) { _, err := m.Execute(t.Context(), subscribeRequest(), nil) require.Error(t, err) assert.Contains(t, err.Error(), "cannot find a runner") + assert.NotErrorIs(t, err, ErrRunnerUnavailable) }) } diff --git a/pkg/workflows/host/requirements_gen/requirements_helper.go.tmpl b/pkg/workflows/host/requirements_gen/requirements_helper.go.tmpl index da55611a42..f822cc849a 100644 --- a/pkg/workflows/host/requirements_gen/requirements_helper.go.tmpl +++ b/pkg/workflows/host/requirements_gen/requirements_helper.go.tmpl @@ -33,3 +33,24 @@ func CheckRequirements(ctx context.Context, handler RequirementsHandler, req *sd return true } + +// HandlesRequirements reports whether handler covers every non-nil field in req (the +// right runner type), regardless of whether those callbacks pass. Distinguishes a +// present-but-unavailable runner from a missing one. +func HandlesRequirements(handler RequirementsHandler, req *sdk.Requirements) bool { + if req == nil { + return true + } + + if len(req.ProtoReflect().GetUnknown()) != 0 { + return false + } + +{{range .Fields}} + if req.{{.Name}} != nil && handler.{{.Name}} == nil { + return false + } +{{end}} + + return true +} diff --git a/pkg/workflows/host/requirements_helper_gen.go b/pkg/workflows/host/requirements_helper_gen.go index d9e5140f34..73099b9c0b 100644 --- a/pkg/workflows/host/requirements_helper_gen.go +++ b/pkg/workflows/host/requirements_helper_gen.go @@ -35,3 +35,22 @@ func CheckRequirements(ctx context.Context, handler RequirementsHandler, req *sd return true } + +// HandlesRequirements reports whether handler covers every non-nil field in req (the +// right runner type), regardless of whether those callbacks pass. Distinguishes a +// present-but-unavailable runner from a missing one. +func HandlesRequirements(handler RequirementsHandler, req *sdk.Requirements) bool { + if req == nil { + return true + } + + if len(req.ProtoReflect().GetUnknown()) != 0 { + return false + } + + if req.Tee != nil && handler.Tee == nil { + return false + } + + return true +} diff --git a/pkg/workflows/host/requirements_helper_gen_test.go b/pkg/workflows/host/requirements_helper_gen_test.go index 68bb23b3d7..ed53178b25 100644 --- a/pkg/workflows/host/requirements_helper_gen_test.go +++ b/pkg/workflows/host/requirements_helper_gen_test.go @@ -46,3 +46,38 @@ func Test_CheckRequirements(t *testing.T) { assert.True(t, CheckRequirements(context.Background(), handler, req)) }) } + +func Test_HandlesRequirements(t *testing.T) { + t.Parallel() + t.Run("nil req always passes", func(t *testing.T) { + assert.True(t, HandlesRequirements(RequirementsHandler{}, nil)) + }) + + t.Run("no fields always passes", func(t *testing.T) { + assert.True(t, HandlesRequirements(RequirementsHandler{}, &sdk.Requirements{})) + }) + + t.Run("unknown proto fields", func(t *testing.T) { + // Encode a field number (99) unknown to Requirements so proto.Unmarshal + // preserves it as unknown bytes. + b := protowire.AppendTag(nil, 99, protowire.VarintType) + b = protowire.AppendVarint(b, 1) + req := &sdk.Requirements{} + require.NoError(t, proto.Unmarshal(b, req)) + + assert.False(t, HandlesRequirements(RequirementsHandler{}, req)) + }) + + t.Run("required field with nil handler returns false", func(t *testing.T) { + req := &sdk.Requirements{Tee: &sdk.Tee{}} + assert.False(t, HandlesRequirements(RequirementsHandler{}, req)) + }) + + t.Run("handler present passes regardless of callback result", func(t *testing.T) { + req := &sdk.Requirements{Tee: &sdk.Tee{}} + // HandlesRequirements only checks coverage, not whether the callback passes, + // so a handler that would fail CheckRequirements still handles the request. + handler := RequirementsHandler{Tee: func(context.Context, *sdk.Tee) bool { return false }} + assert.True(t, HandlesRequirements(handler, req)) + }) +}