This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / mill / mill.go
28 kB 1119 lines
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}