This repository has no description
1package mill
2
3import (
4 "context"
5 "encoding/json"
6 "errors"
7 "fmt"
8 "log/slog"
9 "os"
10 "path/filepath"
11 "slices"
12 "strings"
13 "sync"
14 "time"
15
16 "tangled.org/core/notifier"
17 "tangled.org/core/spindle/db"
18 "tangled.org/core/spindle/engine"
19 "tangled.org/core/spindle/models"
20 "tangled.org/core/spindle/secrets"
21 "tangled.org/core/tid"
22
23 millproto "tangled.org/core/spindle/mill/proto"
24 millv1 "tangled.org/core/spindle/mill/proto/gen"
25)
26
27const (
28 defaultReconnectGrace = 45 * time.Second
29 defaultJobTimeout = 24 * time.Hour
30 defaultBidTimeout = 5 * time.Second
31 defaultTopK = 3
32 defaultMaxPending = 100
33)
34
35type Config struct {
36 // mill appends live-tailed executor lines here so logview can follow running remote jobs.
37 LogDir string
38 MaxPending int
39 ReconnectGrace time.Duration
40 JobTimeout time.Duration
41 BidTimeout time.Duration
42 TopK int
43 CancelTimeout time.Duration
44}
45
46type Mill struct {
47 l *slog.Logger
48 cfg Config
49
50 db *db.DB
51 n *notifier.Notifier
52
53 mu sync.Mutex
54 sessions map[string]*millSession
55 leases map[string]*RemoteLease
56 reservations map[string]*RemoteLease
57 nodeSeqno map[string]uint64
58 pending int
59 changeCh chan struct{} // closed and replaced to wake placement waiters
60
61 leaseSeq uint64
62}
63
64func New(l *slog.Logger, cfg Config) *Mill {
65 if cfg.ReconnectGrace <= 0 {
66 cfg.ReconnectGrace = defaultReconnectGrace
67 }
68 if cfg.JobTimeout <= 0 {
69 cfg.JobTimeout = defaultJobTimeout
70 }
71 if cfg.BidTimeout <= 0 {
72 cfg.BidTimeout = defaultBidTimeout
73 }
74 if cfg.TopK <= 0 {
75 cfg.TopK = defaultTopK
76 }
77 if cfg.MaxPending <= 0 {
78 cfg.MaxPending = defaultMaxPending
79 }
80 if cfg.CancelTimeout <= 0 {
81 cfg.CancelTimeout = 10 * time.Second
82 }
83 return &Mill{
84 l: l,
85 cfg: cfg,
86 sessions: make(map[string]*millSession),
87 leases: make(map[string]*RemoteLease),
88 reservations: make(map[string]*RemoteLease),
89 nodeSeqno: make(map[string]uint64),
90 changeCh: make(chan struct{}),
91 }
92}
93
94func (m *Mill) Attach(d *db.DB, n *notifier.Notifier) {
95 m.mu.Lock()
96 m.db = d
97 m.n = n
98 m.mu.Unlock()
99}
100
101func (m *Mill) nextLeaseID() string {
102 m.mu.Lock()
103 m.leaseSeq++
104 seq := m.leaseSeq
105 m.mu.Unlock()
106 return fmt.Sprintf("%s-%d", tid.TID(), seq)
107}
108
109func (m *Mill) dropReservation(id string) {
110 m.mu.Lock()
111 delete(m.reservations, id)
112 m.mu.Unlock()
113}
114
115func (m *Mill) notifyChange() {
116 m.mu.Lock()
117 m.notifyChangeLocked()
118 m.mu.Unlock()
119}
120
121func (m *Mill) notifyChangeLocked() {
122 close(m.changeCh)
123 m.changeCh = make(chan struct{})
124}
125
126func (m *Mill) currentChangeCh() <-chan struct{} {
127 m.mu.Lock()
128 defer m.mu.Unlock()
129 return m.changeCh
130}
131
132func (m *Mill) attachSession(sess *millSession) (uint64, bool) {
133 m.mu.Lock()
134 defer m.mu.Unlock()
135
136 if old := m.sessions[sess.nodeID]; old != nil {
137 if old.live(m.cfg.ReconnectGrace) {
138 // a second live session for the same identity is a hijack
139 // attempt, reject it
140 return 0, false
141 }
142 if old.graceTimer != nil {
143 old.graceTimer.Stop()
144 }
145 old.close()
146 if old.disconnected {
147 m.l.Info("executor reconnected", "node", sess.nodeID)
148 } else {
149 m.l.Warn("replacing silent executor session", "node", sess.nodeID)
150 }
151 }
152 m.sessions[sess.nodeID] = sess
153 // wakes commit retries waiting out reconnect grace
154 m.notifyChangeLocked()
155 return m.nodeSeqno[sess.nodeID+"/"+sess.epoch], true
156}
157
158func (m *Mill) touchSession(sess *millSession) bool {
159 m.mu.Lock()
160 defer m.mu.Unlock()
161 if m.sessions[sess.nodeID] != sess || sess.disconnected {
162 return false
163 }
164 sess.lastSeen = time.Now()
165 return true
166}
167
168func (m *Mill) detachSession(sess *millSession) {
169 m.mu.Lock()
170 if m.sessions[sess.nodeID] != sess {
171 // already replaced by a reconnect
172 m.mu.Unlock()
173 sess.close()
174 return
175 }
176 sess.disconnected = true
177 sess.graceTimer = time.AfterFunc(m.cfg.ReconnectGrace, func() { m.failLeasesAfterGrace(sess) })
178 m.mu.Unlock()
179
180 sess.close()
181 m.l.Warn("executor session lost; entering reconnect grace", "node", sess.nodeID, "grace", m.cfg.ReconnectGrace)
182 m.notifyChange()
183}
184
185func (m *Mill) sessionReady(sess *millSession) {
186 for _, lease := range m.cancelledLeasesForNode(sess.nodeID) {
187 _, reason := lease.cancelRequested()
188 m.sendCancel(sess, lease, reason)
189 }
190 m.notifyChange()
191}
192
193func (m *Mill) cancelledLeasesForNode(nodeID string) []*RemoteLease {
194 m.mu.Lock()
195 var candidates []*RemoteLease
196 for _, lease := range m.leases {
197 if lease.nodeID == nodeID {
198 candidates = append(candidates, lease)
199 }
200 }
201 m.mu.Unlock()
202
203 var leases []*RemoteLease
204 for _, lease := range candidates {
205 if cancelled, _ := lease.cancelRequested(); cancelled && lease.getState() != leaseDone {
206 leases = append(leases, lease)
207 }
208 }
209 return leases
210}
211
212// reconnect grace ran out without the executor coming back. fails every
213// lease the node held, and reschedules itself if any failure didn't stick
214func (m *Mill) failLeasesAfterGrace(sess *millSession) {
215 m.mu.Lock()
216 if m.sessions[sess.nodeID] != sess || !sess.disconnected {
217 // reconnected in the meantime
218 m.mu.Unlock()
219 return
220 }
221 var dead []*RemoteLease
222 for _, lease := range m.leases {
223 if lease.nodeID == sess.nodeID {
224 dead = append(dead, lease)
225 }
226 }
227 m.mu.Unlock()
228
229 m.l.Warn("executor declared dead; failing its in-flight jobs", "node", sess.nodeID, "jobs", len(dead))
230 deadReason := "executor lost"
231 success := true
232 for _, lease := range dead {
233 switch {
234 case lease.orphaned:
235 // restored leases have no RunStep waiting, fail them straight
236 // into the event stream
237 if err := m.finishOrphan(lease, string(models.StatusKindFailed), &deadReason, nil); err != nil {
238 m.l.Error("finish orphan after executor loss failed, will retry", "lease", lease.id, "err", err)
239 success = false
240 }
241 default:
242 // live leases have a blocked RunStep, keep a pending cancellation
243 // as the terminal reason
244 status := string(models.StatusKindFailed)
245 reason := deadReason
246 if cancelled, cancelReason := lease.cancelRequested(); cancelled {
247 status = string(models.StatusKindCancelled)
248 reason = cancelReason
249 }
250 if err := m.finishLiveLease(lease, status, reason); err != nil {
251 m.l.Error("finish lease after executor loss failed, will retry", "lease", lease.id, "err", err)
252 success = false
253 }
254 }
255 }
256
257 m.mu.Lock()
258 if !success {
259 // something didn't finish cleanly, try again shortly. finished
260 // leases skip themselves on the next pass
261 sess.graceTimer = time.AfterFunc(5*time.Second, func() { m.failLeasesAfterGrace(sess) })
262 m.mu.Unlock()
263 return
264 }
265
266 // everything failed cleanly. forget the node and wake placement
267 if m.sessions[sess.nodeID] == sess {
268 delete(m.sessions, sess.nodeID)
269 }
270 m.mu.Unlock()
271 m.notifyChange()
272 m.sweepUnclaimedOrphans()
273}
274
275func (m *Mill) place(ctx context.Context, engineName string, wid models.WorkflowId, wf *models.Workflow) (engine.WorkflowSlot, error) {
276 m.mu.Lock()
277 if m.cfg.MaxPending > 0 && m.pending >= m.cfg.MaxPending {
278 max := m.cfg.MaxPending
279 cur := m.pending
280 m.mu.Unlock()
281 return nil, fmt.Errorf("%w: mill has %d pending jobs (max %d)", engine.ErrNoWorkflowSlots, cur, max)
282 }
283 m.pending++
284 m.mu.Unlock()
285 defer func() {
286 m.mu.Lock()
287 m.pending--
288 m.mu.Unlock()
289 }()
290
291 for {
292 if err := ctx.Err(); err != nil {
293 return nil, err
294 }
295
296 // grab the channel before bidding. a change mid-bid closes it, so
297 // the wait below re-bids right away
298 ch := m.currentChangeCh()
299
300 lease, err := m.bid(ctx, engineName, wid, wf)
301 if err != nil {
302 return nil, err
303 }
304 if lease != nil {
305 lease.wid = wid
306 if err := m.persistLease(lease, leaseRowReserved); err != nil {
307 m.releaseRemote(lease)
308 m.mu.Lock()
309 delete(m.reservations, lease.id)
310 m.mu.Unlock()
311 return nil, fmt.Errorf("persist reserved mill lease: %w", err)
312 }
313 m.mu.Lock()
314 delete(m.reservations, lease.id)
315 m.leases[lease.id] = lease
316 if st, ok := wf.Data.(*millWorkflowState); ok && st != nil {
317 st.Lease = lease
318 }
319 m.mu.Unlock()
320 return &millSlot{fleet: m, lease: lease}, nil
321 }
322
323 // no executor available. wait for a change or ctx
324 select {
325 case <-ctx.Done():
326 return nil, ctx.Err()
327 case <-ch:
328 }
329 }
330}
331
332func (m *Mill) bid(ctx context.Context, engineName string, wid models.WorkflowId, wf *models.Workflow) (*RemoteLease, error) {
333 rawPipeline, rawWorkflow, err := marshalJob(wf)
334 if err != nil {
335 return nil, err
336 }
337
338 requiredLabels := requiredLabels(wf)
339 candidates := m.rankCandidates(engineName, requiredLabels)
340 if len(candidates) == 0 {
341 return nil, nil
342 }
343
344 type bidResult struct {
345 sess *millSession
346 lease *RemoteLease
347 rank int
348 incompatible bool
349 reason string
350 }
351 limit := m.cfg.TopK
352 if limit <= 0 {
353 limit = len(candidates)
354 }
355 if len(candidates) > limit {
356 candidates = candidates[:limit]
357 }
358 results := make(chan bidResult, limit)
359 ask := func(rank int, sess *millSession) {
360 bidCtx, cancel := context.WithTimeout(ctx, m.cfg.BidTimeout)
361 defer cancel()
362 leaseID := m.nextLeaseID()
363 lease := newLease(leaseID, sess.nodeID, sess.epoch, engineName)
364 m.mu.Lock()
365 m.reservations[leaseID] = lease
366 m.mu.Unlock()
367 msg := &millproto.Message{ReserveSeat: &millv1.ReserveSeat{
368 LeaseId: leaseID,
369 TargetEngine: engineName,
370 RawPipelineJson: rawPipeline,
371 RawWorkflowJson: rawWorkflow,
372 Knot: wid.Knot,
373 Rkey: wid.Rkey,
374 TtlSeconds: uint32(m.cfg.ReconnectGrace / time.Second),
375 }}
376 resp, err := sess.request(bidCtx, leaseID, msg)
377 if err != nil {
378 m.dropReservation(leaseID)
379 results <- bidResult{}
380 return
381 }
382 rr := resp.GetReserveResult()
383 if rr == nil {
384 m.dropReservation(leaseID)
385 results <- bidResult{}
386 return
387 }
388 if !rr.GetAccepted() {
389 m.dropReservation(leaseID)
390 if rr.GetRejectClass() == millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE {
391 results <- bidResult{sess: sess, rank: rank, incompatible: true, reason: rr.GetRejectReason()}
392 return
393 }
394 results <- bidResult{}
395 return
396 }
397 results <- bidResult{sess: sess, lease: lease, rank: rank}
398 }
399 next := 0
400 inFlight := 0
401 for next < len(candidates) && inFlight < limit {
402 inFlight++
403 go ask(next, candidates[next])
404 next++
405 }
406
407 var winner *bidResult
408 var losers []*RemoteLease
409 var incompatible []string
410 // any soft reject (transient or timeout) means the fleet was just
411 // busy, so an all-incompatible outcome isn't a hard placement error
412 softReject := false
413 for inFlight > 0 {
414 r := <-results
415 inFlight--
416 // incompatible rejects get reported to the user, any other failure
417 // just means the fleet is busy
418 if r.incompatible {
419 if r.reason != "" {
420 incompatible = append(incompatible, r.reason)
421 }
422 } else if r.lease == nil {
423 softReject = true
424 }
425 // a failed bid means ask the next candidate, unless someone won
426 if r.lease == nil {
427 for winner == nil && next < len(candidates) && inFlight < limit {
428 inFlight++
429 go ask(next, candidates[next])
430 next++
431 }
432 continue
433 }
434 // if this bid is worse than the winner it goes to the losers pile
435 if winner != nil && r.rank >= winner.rank {
436 losers = append(losers, r.lease)
437 continue
438 }
439 // otherwise it's the new best and the old winner joins the losers
440 if winner != nil {
441 losers = append(losers, winner.lease)
442 }
443 winner = &r
444 }
445
446 // let the losers go so they free their seats right away
447 for _, l := range losers {
448 m.dropReservation(l.id)
449 m.releaseRemote(l)
450 }
451
452 if winner == nil {
453 if len(incompatible) > 0 && !softReject {
454 return nil, fmt.Errorf("no compatible executor for %s: %s", engineName, strings.Join(incompatible, "; "))
455 }
456 return nil, nil
457 }
458 return winner.lease, nil
459}
460
461// ranks nodes that are least busy first. if a resource is used a lot
462// then that node will lose to one that is more even across the board.
463func (m *Mill) rankCandidates(engineName string, requiredLabels []string) []*millSession {
464 m.mu.Lock()
465 defer m.mu.Unlock()
466
467 type ranked struct {
468 sess *millSession
469 worst float64
470 sum float64
471 }
472 var rs []ranked
473 for _, s := range m.sessions {
474 // only live, reporting sessions can take work
475 if s.disconnected {
476 continue
477 }
478 if s.snapshot == nil {
479 continue
480 }
481 // the engine has to exist and have room right now
482 ea, ok := s.snapshot.GetEngines()[engineName]
483 if !ok || !ea.GetAvailable() {
484 continue
485 }
486 // and satisfy the wf's label requirements
487 if !hasLabels(s.labels, requiredLabels) {
488 continue
489 }
490 worst, sum := loadScore(ea.GetLoad())
491 rs = append(rs, ranked{sess: s, worst: worst, sum: sum})
492 }
493 slices.SortStableFunc(rs, func(a, b ranked) int {
494 if a.worst < b.worst {
495 return -1
496 }
497 if a.worst > b.worst {
498 return 1
499 }
500 if a.sum < b.sum {
501 return -1
502 }
503 if a.sum > b.sum {
504 return 1
505 }
506 return 0
507 })
508
509 out := make([]*millSession, len(rs))
510 for i := range rs {
511 out[i] = rs[i].sess
512 }
513 return out
514}
515
516func loadScore(load map[string]float64) (worst, sum float64) {
517 for _, v := range load {
518 if v > worst {
519 worst = v
520 }
521 sum += v
522 }
523 return worst, sum
524}
525
526func requiredLabels(wf *models.Workflow) []string {
527 st, ok := wf.Data.(*millWorkflowState)
528 if !ok || st == nil {
529 return nil
530 }
531 return st.RawWorkflow.RunsOn
532}
533
534func hasLabels(labels []string, required []string) bool {
535 for _, want := range required {
536 if !slices.Contains(labels, want) {
537 return false
538 }
539 }
540 return true
541}
542
543func (m *Mill) commitAndWait(ctx context.Context, wf *models.Workflow, unlocked []secrets.UnlockedSecret) error {
544 st, ok := wf.Data.(*millWorkflowState)
545 if !ok || st == nil || st.Lease == nil {
546 return fmt.Errorf("mill workflow state missing lease")
547 }
548 lease := st.Lease
549
550 pbSecrets := make([]*millv1.Secret, len(unlocked))
551 for i, s := range unlocked {
552 pbSecrets[i] = &millv1.Secret{Key: s.Key, Value: s.Value}
553 }
554
555 commit := &millproto.Message{CommitLease: &millv1.CommitLease{
556 LeaseId: lease.id,
557 Secrets: pbSecrets,
558 }}
559
560 // commit retries ride reconnects, a reservation outlives one
561 // disconnect. lost session or slow executor just means wait and retry,
562 // only job timeout or a dead lease stops the loop
563 for {
564 if res, ok := pollTerminal(lease); ok {
565 return terminalError(res.Status)
566 }
567 if !lease.markCommitting() {
568 return engine.ErrWorkflowFailed
569 }
570
571 sess := m.sessionForNode(lease.nodeID)
572 if sess == nil {
573 if done, err := m.waitCommitRetry(ctx, lease); done || err != nil {
574 return err
575 }
576 continue
577 }
578
579 reqCtx, cancel := context.WithTimeout(ctx, m.cfg.BidTimeout)
580 resp, err := sess.request(reqCtx, lease.id, commit)
581 cancel()
582 if err != nil {
583 switch {
584 case errors.Is(err, errSessionClosed):
585 // session died mid-request. wait out the grace, then retry on
586 // the new one
587 if done, err := m.waitCommitRetry(ctx, lease); done || err != nil {
588 return err
589 }
590 continue
591 case errors.Is(err, context.DeadlineExceeded) && ctx.Err() == nil:
592 // executor didn't answer in time, but its seat is still held so
593 // retrying is safe
594 continue
595 case errors.Is(err, context.DeadlineExceeded):
596 // the job ctx itself ran out, a real timeout
597 return engine.ErrTimedOut
598 case errors.Is(err, context.Canceled):
599 return err
600 default:
601 m.l.Warn("commit lease send failed; waiting for reconnect", "lease", lease.id, "node", lease.nodeID, "err", err)
602 if done, err := m.waitCommitRetry(ctx, lease); done || err != nil {
603 return err
604 }
605 continue
606 }
607 }
608 if resp.GetCommitted() == nil {
609 return engine.ErrWorkflowFailed
610 }
611 lease.markRunning()
612 if err := m.persistLease(lease, leaseRowRunning); err != nil {
613 m.l.Error("persist running mill lease", "lease", lease.id, "err", err)
614 }
615 if cancelled, reason := lease.cancelRequested(); cancelled {
616 m.sendCancel(sess, lease, reason)
617 }
618 break
619 }
620
621 select {
622 case res := <-lease.terminal:
623 return terminalError(res.Status)
624 case <-ctx.Done():
625 if ctx.Err() == context.DeadlineExceeded {
626 return engine.ErrTimedOut
627 }
628 return ctx.Err()
629 }
630}
631
632func (m *Mill) waitCommitRetry(ctx context.Context, lease *RemoteLease) (bool, error) {
633 // grab the channel before checking for a live session again
634 // a reconnect will still close the channel if it happens in between
635 ch := m.currentChangeCh()
636 if m.sessionForNode(lease.nodeID) != nil {
637 return false, nil
638 }
639 select {
640 case res := <-lease.terminal:
641 return true, terminalError(res.Status)
642 case <-ctx.Done():
643 if ctx.Err() == context.DeadlineExceeded {
644 return true, engine.ErrTimedOut
645 }
646 return true, ctx.Err()
647 case <-ch:
648 return false, nil
649 }
650}
651
652func terminalError(status millv1.TerminalStatus) error {
653 switch status {
654 case millv1.TerminalStatus_SUCCESS:
655 return nil
656 case millv1.TerminalStatus_TIMEOUT:
657 return engine.ErrTimedOut
658 case millv1.TerminalStatus_CANCELLED:
659 return engine.ErrWorkflowCanceled
660 default:
661 return engine.ErrWorkflowFailed
662 }
663}
664
665func pollTerminal(lease *RemoteLease) (*millv1.AttemptResult, bool) {
666 select {
667 case res := <-lease.terminal:
668 return res, true
669 default:
670 return nil, false
671 }
672}
673
674func (m *Mill) destroy(wid models.WorkflowId) {
675 m.mu.Lock()
676 var lease *RemoteLease
677 for _, l := range m.leases {
678 if l.wid == wid {
679 lease = l
680 break
681 }
682 }
683 m.mu.Unlock()
684 if lease == nil {
685 return
686 }
687 reason := "workflow destroyed"
688 switch lease.requestCancel(reason) {
689 case cancelLocal:
690 if sess := m.sessionForNode(lease.nodeID); sess != nil {
691 _ = sess.send(&millproto.Message{ReleaseLease: &millv1.ReleaseLease{LeaseId: lease.id}})
692 }
693 lease.deliverCancelled(reason)
694 case cancelRemote:
695 if sess := m.sessionForNode(lease.nodeID); sess != nil {
696 m.sendCancel(sess, lease, reason)
697 }
698 }
699}
700
701func (m *Mill) releaseSlot(s *millSlot) {
702 lease := s.lease
703 state, cancelled := lease.releaseState()
704 if state == leaseReserved {
705 m.releaseRemote(lease)
706 } else if cancelled && state != leaseDone {
707 return
708 }
709 if err := m.cleanupLease(lease); err != nil {
710 m.l.Error("releaseSlot cleanupLease failed", "lease", lease.id, "err", err)
711 }
712}
713func (m *Mill) cleanupLease(lease *RemoteLease) error {
714 lease.finishMu.Lock()
715 defer lease.finishMu.Unlock()
716 return m.cleanupLeaseLocked(lease)
717}
718
719func (m *Mill) cleanupLeaseLocked(lease *RemoteLease) error {
720 if lease.cleanedUp {
721 return nil
722 }
723 if m.db != nil {
724 if err := m.db.DeleteMillLease(lease.id); err != nil {
725 m.scheduleCleanupLocked(lease)
726 return err
727 }
728 }
729 m.mu.Lock()
730 delete(m.leases, lease.id)
731 m.mu.Unlock()
732 lease.cleanedUp = true
733 m.notifyChange()
734 return nil
735}
736
737func (m *Mill) scheduleCleanupLocked(lease *RemoteLease) {
738 if lease.cleanedUp || lease.cleanupRetry {
739 return
740 }
741 lease.cleanupRetry = true
742 time.AfterFunc(5*time.Second, func() {
743 lease.finishMu.Lock()
744 lease.cleanupRetry = false
745 err := m.cleanupLeaseLocked(lease)
746 lease.finishMu.Unlock()
747 if err != nil {
748 m.l.Error("retry lease cleanup failed", "lease", lease.id, "err", err)
749 }
750 })
751}
752
753func (m *Mill) releaseRemote(lease *RemoteLease) {
754 lease.setState(leaseDone)
755 if sess := m.sessionForNode(lease.nodeID); sess != nil {
756 _ = sess.send(&millproto.Message{ReleaseLease: &millv1.ReleaseLease{LeaseId: lease.id}})
757 }
758}
759
760func (m *Mill) sendCancel(sess *millSession, lease *RemoteLease, reason string) {
761 if err := sess.send(&millproto.Message{CancelAttempt: &millv1.CancelAttempt{
762 LeaseId: lease.id,
763 Reason: reason,
764 }}); err != nil {
765 return
766 }
767 time.AfterFunc(m.cfg.CancelTimeout, func() { m.checkCancelDeadline(lease) })
768}
769
770func (m *Mill) sessionForNode(nodeID string) *millSession {
771 m.mu.Lock()
772 defer m.mu.Unlock()
773 sess := m.sessions[nodeID]
774 if sess == nil || sess.disconnected {
775 return nil
776 }
777 return sess
778}
779
780func (m *Mill) onSnapshot(sess *millSession, snap *millv1.NodeSnapshot) error {
781 m.mu.Lock()
782 if sess.snapshot != nil && snap.Seqno <= sess.snapshot.Seqno {
783 m.mu.Unlock()
784 return fmt.Errorf("protocol error: snapshot seqno regression. Got %d, last seen %d", snap.Seqno, sess.snapshot.Seqno)
785 }
786 sess.snapshot = snap
787 m.mu.Unlock()
788
789 if err := m.reconcileLeases(sess, snap.GetActiveLeaseIds()); err != nil {
790 return err
791 }
792 m.mu.Lock()
793 for _, id := range snap.GetActiveLeaseIds() {
794 if lease := m.leases[id]; lease != nil && lease.nodeID == sess.nodeID && lease.epoch == sess.epoch {
795 lease.claimed = true
796 }
797 }
798 m.mu.Unlock()
799 m.notifyChange()
800 return nil
801}
802
803func (m *Mill) onEventBatch(sess *millSession, batch *millv1.EventBatch) error {
804 if batch == nil {
805 return nil
806 }
807
808 if batch.Epoch != sess.epoch {
809 return fmt.Errorf("protocol error: batch epoch %q does not match session %q", batch.Epoch, sess.epoch)
810 }
811
812 m.mu.Lock()
813 currentKey := sess.nodeID + "/" + sess.epoch
814 current := m.nodeSeqno[currentKey]
815 m.mu.Unlock()
816
817 expected := current + 1
818 var newEntries []*millv1.Event
819 for _, entry := range batch.Events {
820 // reconnect replays old seqnos, drop those
821 if entry.Seqno <= current {
822 continue
823 }
824 // gap means executor lost rows. so dont apply a partial batch
825 if entry.Seqno != expected {
826 return fmt.Errorf("protocol error: gap in stream seqnos. Expected %d, got %d", expected, entry.Seqno)
827 }
828 newEntries = append(newEntries, entry)
829 expected++
830 }
831
832 // all replays, still ack so the executor can trim its outbox
833 if len(newEntries) == 0 {
834 return m.sendAck(sess, current)
835 }
836
837 type pendingTerminal struct {
838 lease *RemoteLease
839 ar *millv1.AttemptResult
840 }
841 var pendingTerminals []pendingTerminal
842 var artifactLeases []*RemoteLease
843 finishedInBatch := make(map[string]struct{})
844
845 var highestSeqno uint64 = current
846
847 applyFunc := func(tx *db.EventBatchTx) error {
848 for _, entry := range newEntries {
849 m.mu.Lock()
850 lease := m.leases[entry.LeaseId]
851 m.mu.Unlock()
852
853 // events for leases this node doesn't own are skipped but still
854 // count as processed
855 if lease == nil || lease.nodeID != sess.nodeID {
856 highestSeqno = entry.Seqno
857 continue
858 }
859
860 // events arriving for a different epoch are invalid (different session)
861 if lease.epoch != "" && lease.epoch != sess.epoch {
862 return fmt.Errorf("protocol error: lease %q epoch %q does not match session %q", lease.id, lease.epoch, sess.epoch)
863 }
864
865 lease.mu.Lock()
866 state := lease.state
867 lease.mu.Unlock()
868
869 // done leases can replay terminals on reconnect, skip them
870 if state == leaseDone {
871 highestSeqno = entry.Seqno
872 continue
873 }
874
875 if _, finished := finishedInBatch[lease.id]; finished {
876 return fmt.Errorf("protocol error: stream entry follows terminal for lease %q", lease.id)
877 }
878
879 switch {
880 case entry.GetStatusEvent() != nil:
881 ev := entry.GetStatusEvent()
882 statusStr := string(models.StatusKindRunning)
883 var errMsg *string
884 if e := ev.GetError(); e != "" {
885 errMsg = &e
886 }
887 var exitCode *int64
888 if c := ev.GetExitCode(); c != 0 {
889 exitCode = &c
890 }
891 pipelineAtUri := string(lease.wid.PipelineId.AtUri())
892 if tx != nil {
893 if err := tx.InsertStatusEvent(pipelineAtUri, lease.wid.Name, statusStr, errMsg, exitCode); err != nil {
894 return err
895 }
896 }
897
898 case entry.GetAttemptResult() != nil:
899 ar := entry.GetAttemptResult()
900 statusStr := "success"
901 switch ar.Status {
902 case millv1.TerminalStatus_SUCCESS:
903 statusStr = "success"
904 case millv1.TerminalStatus_FAILED:
905 statusStr = "failed"
906 case millv1.TerminalStatus_TIMEOUT:
907 statusStr = "timeout"
908 case millv1.TerminalStatus_CANCELLED:
909 statusStr = "cancelled"
910 default:
911 return fmt.Errorf("protocol error: unsupported terminal status %v", ar.Status)
912 }
913 var errMsg *string
914 if e := ar.GetError(); e != "" {
915 errMsg = &e
916 }
917 var exitCode *int64
918 if c := ar.GetExitCode(); c != 0 {
919 exitCode = &c
920 }
921 finishedInBatch[lease.id] = struct{}{}
922 pipelineAtUri := string(lease.wid.PipelineId.AtUri())
923 if tx != nil {
924 if err := tx.InsertStatusEvent(pipelineAtUri, lease.wid.Name, statusStr, errMsg, exitCode); err != nil {
925 return err
926 }
927 if err := tx.DeleteLease(lease.id); err != nil {
928 return err
929 }
930 if a := ar.GetLogArtifact(); a != nil {
931 if a.GetRef() == "" {
932 return fmt.Errorf("protocol error: empty log artifact ref")
933 }
934 if !strings.HasPrefix(a.GetHash(), "sha256:") {
935 return fmt.Errorf("protocol error: invalid log artifact hash %q", a.GetHash())
936 }
937 if tx != nil {
938 if err := tx.InsertArtifactRef(lease.id, lease.wid.Name, a.GetRef(), a.GetHash()); err != nil {
939 return err
940 }
941 }
942 artifactLeases = append(artifactLeases, lease)
943 }
944 }
945 pendingTerminals = append(pendingTerminals, pendingTerminal{
946 lease: lease,
947 ar: ar,
948 })
949 }
950
951 highestSeqno = entry.Seqno
952 }
953
954 if tx != nil {
955 return tx.AdvanceCursor(sess.nodeID, sess.epoch, highestSeqno)
956 }
957 return nil
958 }
959
960 var err error
961 if m.db != nil {
962 err = m.db.ApplyEventBatch(m.n, applyFunc)
963 } else {
964 err = applyFunc(nil)
965 }
966
967 if err != nil {
968 return err
969 }
970
971 if m.cfg.LogDir != "" {
972 for _, lease := range artifactLeases {
973 path := models.LogFilePath(m.cfg.LogDir, lease.wid)
974 if err := os.Remove(path); err != nil && !errors.Is(err, os.ErrNotExist) {
975 m.l.Warn("failed to remove live log file after artifact recorded", "path", path, "err", err)
976 }
977 }
978 }
979
980 m.mu.Lock()
981 if highestSeqno > m.nodeSeqno[currentKey] {
982 m.nodeSeqno[currentKey] = highestSeqno
983 }
984 m.mu.Unlock()
985
986 for _, pt := range pendingTerminals {
987 // orphans have no waiting RunStep, just mark and clean up
988 if pt.lease.orphaned {
989 pt.lease.markDone()
990 _ = m.cleanupLease(pt.lease)
991 continue
992 }
993 pt.lease.deliverTerminal(pt.ar)
994 if pt.lease.cleanupReady() {
995 _ = m.cleanupLease(pt.lease)
996 }
997 }
998
999 return m.sendAck(sess, highestSeqno)
1000}
1001
1002func (m *Mill) sendAck(sess *millSession, seqno uint64) error {
1003 msg := &millproto.Message{Ack: &millv1.Ack{
1004 Epoch: sess.epoch,
1005 UpToSeqno: seqno,
1006 }}
1007 if err := sess.send(msg); err != nil {
1008 return fmt.Errorf("send ack message: %w", err)
1009 }
1010 return nil
1011}
1012func (m *Mill) onLiveLog(sess *millSession, ll *millv1.LiveLog) error {
1013 if ll == nil || ll.GetLeaseId() == "" {
1014 return nil
1015 }
1016 m.mu.Lock()
1017 lease := m.leases[ll.GetLeaseId()]
1018 m.mu.Unlock()
1019 if lease == nil || lease.nodeID != sess.nodeID {
1020 return nil
1021 }
1022 raw := ll.GetRawJson()
1023 if m.cfg.LogDir == "" || len(raw) == 0 {
1024 if m.n != nil {
1025 m.n.NotifyAll()
1026 }
1027 return nil
1028 }
1029 lease.mu.Lock()
1030 isDone := (lease.state == leaseDone)
1031 lease.mu.Unlock()
1032 if isDone {
1033 if m.n != nil {
1034 m.n.NotifyAll()
1035 }
1036 return nil
1037 }
1038
1039 logPath := models.LogFilePath(m.cfg.LogDir, lease.wid)
1040 if err := os.MkdirAll(filepath.Dir(logPath), 0755); err != nil {
1041 m.l.Warn("failed to create log dir", "path", filepath.Dir(logPath), "err", err)
1042 if m.n != nil {
1043 m.n.NotifyAll()
1044 }
1045 return nil
1046 }
1047 f, err := os.OpenFile(logPath, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0600)
1048 if err != nil {
1049 m.l.Warn("failed to open log file", "path", logPath, "err", err)
1050 if m.n != nil {
1051 m.n.NotifyAll()
1052 }
1053 return nil
1054 }
1055 if _, err := f.Write(raw); err != nil {
1056 m.l.Warn("failed to write log file", "path", logPath, "err", err)
1057 }
1058 _ = f.Close()
1059
1060 if m.n != nil {
1061 m.n.NotifyAll()
1062 }
1063 return nil
1064}
1065
1066func (m *Mill) onCancelAck(sess *millSession, ca *millv1.CancelAck) {
1067 m.mu.Lock()
1068 lease := m.leases[ca.GetLeaseId()]
1069 m.mu.Unlock()
1070 if lease == nil || lease.nodeID != sess.nodeID {
1071 return
1072 }
1073 lease.mu.Lock()
1074 lease.cancelAcked = true
1075 lease.mu.Unlock()
1076}
1077
1078func (m *Mill) checkCancelDeadline(lease *RemoteLease) {
1079 lease.mu.Lock()
1080 isDone := (lease.state == leaseDone)
1081 isAcked := lease.cancelAcked
1082 lease.mu.Unlock()
1083
1084 if isDone && isAcked {
1085 return
1086 }
1087
1088 m.l.Warn("node failed to comply with cancel request within deadline, quarantining", "node", lease.nodeID, "lease", lease.id, "done", isDone, "acked", isAcked)
1089
1090 reason := fmt.Sprintf("cancel noncompliance for lease %s (done: %t, acked: %t)", lease.id, isDone, isAcked)
1091 if m.db != nil {
1092 if err := m.db.QuarantineExecutor(lease.nodeID, reason); err != nil {
1093 m.l.Error("failed to quarantine executor", "node", lease.nodeID, "err", err)
1094 }
1095 }
1096
1097 m.mu.Lock()
1098 sess := m.sessions[lease.nodeID]
1099 m.mu.Unlock()
1100 if sess != nil {
1101 sess.close()
1102 }
1103}
1104
1105func marshalJob(wf *models.Workflow) (pipeline string, workflow string, err error) {
1106 st, ok := wf.Data.(*millWorkflowState)
1107 if !ok || st == nil {
1108 return "", "", fmt.Errorf("mill workflow state missing")
1109 }
1110 p, err := json.Marshal(st.RawPipeline)
1111 if err != nil {
1112 return "", "", fmt.Errorf("marshal pipeline: %w", err)
1113 }
1114 w, err := json.Marshal(st.RawWorkflow)
1115 if err != nil {
1116 return "", "", fmt.Errorf("marshal workflow: %w", err)
1117 }
1118 return string(p), string(w), nil
1119}