-
Notifications
You must be signed in to change notification settings - Fork 29
Expand file tree
/
Copy pathrequirement_selecting_module.go
More file actions
173 lines (146 loc) · 4.9 KB
/
Copy pathrequirement_selecting_module.go
File metadata and controls
173 lines (146 loc) · 4.9 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
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
package host
import (
"context"
"errors"
"fmt"
"sync"
"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
}
// lazyModule wraps a ModuleAndHandler so that Start is called at most once
// and Close only fires for modules that were actually started. The mutex
// serializes start/close so a concurrent Close cannot race past an in-flight
// Start (leaving a started module unclosed) and vice versa.
type lazyModule struct {
ModuleAndHandler
mu sync.Mutex
started bool
closed bool
}
func (l *lazyModule) ensureStarted() {
l.mu.Lock()
defer l.mu.Unlock()
if l.started || l.closed {
return
}
l.Module.Start()
l.started = true
}
func (l *lazyModule) ensureClosed() {
l.mu.Lock()
defer l.mu.Unlock()
if l.closed {
return
}
l.closed = true
if l.started {
l.Module.Close()
}
}
// NewRequirementSelectingModule creates a module that routes trigger executions
// based on subscription requirements. main is prepended as modules[0]; additional
// modules follow. Subscribe always runs on modules[0].
func NewRequirementSelectingModule(main ModuleAndHandler, additional []ModuleAndHandler) Module {
modules := make([]*lazyModule, 1+len(additional))
modules[0] = &lazyModule{ModuleAndHandler: main}
for i, a := range additional {
modules[1+i] = &lazyModule{ModuleAndHandler: a}
}
return &requirementSelectingModule{modules: modules}
}
type triggerInfo struct {
moduleIdx int
preHook bool
requirements *sdk.Requirements
}
type requirementSelectingModule struct {
modules []*lazyModule
// triggerID → triggerInfo
cache sync.Map
}
func (r *requirementSelectingModule) Start() {
r.modules[0].ensureStarted()
}
func (r *requirementSelectingModule) Close() {
for _, m := range r.modules {
m.ensureClosed()
}
}
func (r *requirementSelectingModule) IsLegacyDAG() bool {
return r.modules[0].IsLegacyDAG()
}
func (r *requirementSelectingModule) Execute(ctx context.Context, request *sdk.ExecuteRequest, handler ExecutionHelper) (*sdk.ExecutionResult, error) {
if request.GetTrigger() == nil {
return r.subscribe(ctx, request, handler)
}
return r.trigger(ctx, request, handler)
}
func (r *requirementSelectingModule) subscribe(ctx context.Context, request *sdk.ExecuteRequest, handler ExecutionHelper) (*sdk.ExecutionResult, error) {
result, err := r.modules[0].Execute(ctx, request, handler)
if err != nil {
return nil, err
}
for i, sub := range result.GetTriggerSubscriptions().GetSubscriptions() {
matched := false
for j, m := range r.modules {
if CheckRequirements(ctx, m.RequirementsHandler, sub.Requirements) {
m.ensureStarted()
r.cache.Store(uint64(i), triggerInfo{moduleIdx: j, requirements: sub.Requirements, preHook: sub.PreHook})
matched = true
break
}
}
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)
}
}
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 {
info := val.(triggerInfo)
m := r.modules[info.moduleIdx]
if info.preHook {
prehook := &sdk.ExecuteRequest{Request: &sdk.ExecuteRequest_PreHook{PreHook: trigger}}
preHookResult, err := r.modules[0].Execute(ctx, prehook, handler)
if err != nil {
return nil, fmt.Errorf("pre-hook execution failed: %w", err)
}
switch preHookResult.Result.(type) {
case *sdk.ExecutionResult_Error:
return preHookResult, nil
}
restrictions := preHookResult.GetRestrictions()
handler = NewRestrictedExecutionHelper(handler, restrictions)
if rem, ok := m.Module.(RestrictionAwareModule); ok {
rem.SetRestrictions(handler.GetWorkflowExecutionID(), restrictions)
}
}
if rem, ok := m.Module.(RequirementEnforcingModule); ok && info.requirements != nil {
rem.SetRequirements(handler.GetWorkflowExecutionID(), info.requirements)
}
return m.Execute(ctx, request, handler)
}
return nil, errors.New("cannot trigger before gathering subscriptions")
}
var _ Module = &requirementSelectingModule{}