@@ -28,6 +28,8 @@ import (
2828 "context"
2929 "fmt"
3030 "maps"
31+ "regexp"
32+ "strconv"
3133 "strings"
3234 "sync"
3335 "time"
@@ -52,8 +54,29 @@ const (
5254 k8sDeployTimeout = 80 * time .Second
5355 k8sTransitionTimeout = 80 * time .Second
5456 k8sWatchRetryDelay = 5 * time .Second
57+
58+ // k8sJitTaskTemplateName is the shared TaskTemplate used to run JIT/DPL
59+ // pipeline tasks, since their task class name is unique per generated
60+ // workflow and can't be pre-registered as its own TaskTemplate.
61+ k8sJitTaskTemplateName = "dpl"
5562)
5663
64+ // jitClassNameRe matches JIT-generated task class identifiers of the form
65+ // "jit-<40-hex-char-sha1>-<devicename>" (see configuration/template/dplutil.go)
66+ // and captures the devicename, e.g. "readout-proxy" or "Dispatcher".
67+ var jitClassNameRe = regexp .MustCompile (`^jit-[0-9a-f]{40}-(.+)$` )
68+
69+ func isJitClassName (name string ) bool {
70+ return jitClassNameRe .MatchString (name )
71+ }
72+
73+ // jitOnK8sEnabled reports whether ECS may route JIT/DPL task classes to the
74+ // K8s path. Missing/false means JIT keeps running through Mesos, same as
75+ // before this bridge existed.
76+ func jitOnK8sEnabled () bool {
77+ return viper .GetBool ("jitOnK8s" )
78+ }
79+
5780// k8sEnvRegistry maps ECS environment IDs to K8s custom Environment names.
5881type k8sEnvRegistry struct {
5982 mu sync.RWMutex
@@ -83,6 +106,27 @@ func (r *k8sEnvRegistry) delete(envId uid.ID) {
83106 r .mu .Unlock ()
84107}
85108
109+ // k8sPortAllocator hands out control ports for JIT tasks, one monotonically
110+ // increasing counter per node (hostNetwork is used, so ports must not collide
111+ // between JIT tasks scheduled to the same node)
112+ type k8sPortAllocator struct {
113+ portBase uint16
114+ next map [string ]uint16
115+ }
116+
117+ func newK8sPortAllocator (port uint16 ) * k8sPortAllocator {
118+ return & k8sPortAllocator {portBase : port , next : make (map [string ]uint16 )}
119+ }
120+
121+ func (allocator * k8sPortAllocator ) allocate (hostname string ) uint16 {
122+ port , ok := allocator .next [hostname ]
123+ if ! ok {
124+ port = allocator .portBase
125+ }
126+ allocator .next [hostname ] = port + 1
127+ return port
128+ }
129+
86130// newK8sClientFromViper creates a K8s client from viper config.
87131// Returns nil, nil if kubeNamespace is not configured (K8s disabled).
88132func newK8sClientFromViper () (* k8sclient.Client , error ) {
@@ -112,7 +156,7 @@ func (m *Manager) deployKubernetesTasks(ctx context.Context, envId uid.ID, descr
112156 return nil , err
113157 }
114158
115- nodeToRefs , err := m .buildK8sNodeTaskRefs (entries )
159+ nodeToRefs , err := m .buildK8sNodeTaskRefs (envId , entries )
116160 if err != nil {
117161 log .WithField ("partition" , envId ).WithError (err ).Error ("failed to build K8s node task refs" )
118162 return nil , err
@@ -187,8 +231,12 @@ func (m *Manager) createK8sTaskEntries(envId uid.ID, descriptors Descriptors) ([
187231 return entries , nil
188232}
189233
190- func (m * Manager ) buildK8sNodeTaskRefs (entries []k8sTaskEntry ) (map [string ][]v1alpha1.TaskReference , error ) {
234+ const jitK8sBasePortStr = "jitK8sBasePort"
235+
236+ func (m * Manager ) buildK8sNodeTaskRefs (envId uid.ID , entries []k8sTaskEntry ) (map [string ][]v1alpha1.TaskReference , error ) {
191237 nodeToRefs := make (map [string ][]v1alpha1.TaskReference )
238+
239+ portAllocator := newK8sPortAllocator (viper .GetUint16 (jitK8sBasePortStr ))
192240 for _ , e := range entries {
193241 t := e .task
194242 desc := e .desc
@@ -217,18 +265,50 @@ func (m *Manager) buildK8sNodeTaskRefs(entries []k8sTaskEntry) (map[string][]v1a
217265 }
218266 }
219267
220- argsCLI := make ([]string , 0 , len (cmd .Arguments ))
221- for _ , arg := range cmd .Arguments {
222- if strings .TrimSpace (arg ) != "" {
223- argsCLI = append (argsCLI , arg )
268+ refName := taskClass .Identifier .Name
269+ nameSuffix := ""
270+ isJit := false
271+ if match := jitClassNameRe .FindStringSubmatch (refName ); match != nil {
272+ isJit = true
273+ refName = k8sJitTaskTemplateName
274+ nameSuffix = match [1 ]
275+ }
276+
277+ var argsCLI []string
278+ if isJit {
279+ // Unlike readout/stfbuilder/stfsender (one fixed port per task type,
280+ // baked into their TaskTemplate), JIT devices are dynamic and several
281+ // can run on the same node, so each needs its own control port. K8s has
282+ // no Mesos-style resource offer to claim a verified-free port from, so
283+ // ECS tracks per-node allocation itself (mirrors scheduler.go's
284+ // Mesos-offer-based control port claim for FAIRMQ tasks).
285+ controlPort := portAllocator .allocate (t .hostname )
286+ cmd .Env = append (cmd .Env , fmt .Sprintf ("OCC_CONTROL_PORT=%d" , controlPort ))
287+ cmd .Arguments = append (cmd .Arguments , "--control-port" , strconv .FormatUint (uint64 (controlPort ), 10 ))
288+
289+ // The "dpl" TaskTemplate runs `bash -c <script>`, so the whole
290+ // generated pipeline (driver command + OCC flags, joined the same
291+ // way the Mesos executor does it) must be a single Args entry.
292+ value := ""
293+ if cmd .Value != nil {
294+ value = * cmd .Value
295+ }
296+ argsCLI = []string {strings .Join (append ([]string {value }, cmd .Arguments ... ), " " )}
297+ } else {
298+ argsCLI = make ([]string , 0 , len (cmd .Arguments ))
299+ for _ , arg := range cmd .Arguments {
300+ if strings .TrimSpace (arg ) != "" {
301+ argsCLI = append (argsCLI , arg )
302+ }
224303 }
225304 }
226305
227306 ref := v1alpha1.TaskReference {
228- Name : taskClass .Identifier .Name ,
229- TaskID : t .taskId ,
230- ArgsCLI : argsCLI ,
231- Env : cmdEnvToK8sEnvVars (cmd .Env ),
307+ Name : refName ,
308+ TaskID : t .taskId ,
309+ NameSuffix : nameSuffix ,
310+ ArgsCLI : argsCLI ,
311+ Env : cmdEnvToK8sEnvVars (cmd .Env ),
232312 }
233313 nodeToRefs [t .hostname ] = append (nodeToRefs [t .hostname ], ref )
234314 }
0 commit comments