Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290package mill
import ( "fmt" "maps" "slices" "time"
"tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/models"
millv1 "tangled.org/core/spindle/mill/proto/gen")
const ( leaseRowReserved = "reserved" leaseRowRunning = "running")
func (m *Mill) persistLease(lease *RemoteLease, state string) error { if m.db == nil { return nil } lease.mu.Lock() quotaID := lease.quotaID lease.mu.Unlock() return m.db.SaveMillLease(db.MillLease{ LeaseID: lease.id, NodeID: lease.nodeID, Epoch: lease.epoch, Engine: lease.engine, Knot: lease.wid.Knot, Rkey: lease.wid.Rkey, Workflow: lease.wid.Name, State: state, QuotaReservationID: quotaID, OwnerDID: lease.ownerDID, RepoDID: lease.repoDID, MillRecordsTerminalMetrics: lease.millRecordsTerminalMetrics, })}
// reservation ids the fleet still has live leases for. startup recovery// keeps these and reclaims every other rowfunc (m *Mill) LiveQuotaReservationIDs() []string { m.mu.Lock() leases := slices.Collect(maps.Values(m.leases)) m.mu.Unlock()
var ids []string for _, lease := range leases { lease.mu.Lock() if lease.quotaID != "" { ids = append(ids, lease.quotaID) } ids = append(ids, slices.Collect(maps.Keys(lease.remoteQuotas))...) lease.mu.Unlock() } return ids}
func (m *Mill) RestoreState() error { if m.db == nil { return nil } cursors, err := m.db.ListExecutorCursors() if err != nil { return err } rows, err := m.db.ListMillLeases() if err != nil { return err }
m.mu.Lock() for _, c := range cursors { m.nodeSeqno[c.NodeID+"/"+c.Epoch] = c.AckedSeqno } for _, r := range rows { lease := newLease(r.LeaseID, r.NodeID, r.Epoch, r.Engine) lease.wid = models.WorkflowId{ PipelineId: models.PipelineId{Knot: r.Knot, Rkey: r.Rkey}, Name: r.Workflow, } // the charge outlives the mill process. only the id survives, so // release goes through the manager by id rather than a held lease lease.quotaID = r.QuotaReservationID lease.ownerDID = r.OwnerDID lease.repoDID = r.RepoDID lease.millRecordsTerminalMetrics = r.MillRecordsTerminalMetrics // restored leases start as orphans, an executor must reclaim it via // its first snapshot, or the sweep will fail it lease.orphaned = true lease.claimed = false if r.State == leaseRowRunning { lease.state = leaseRunning } m.leases[r.LeaseID] = lease } restored := len(rows) m.mu.Unlock()
if restored > 0 { m.l.Info("restored mill leases from previous run", "leases", restored, "cursors", len(cursors)) // executors get a grace window to reconnect and claim their leases time.AfterFunc(m.cfg.ReconnectGrace, m.sweepUnclaimedOrphans) } return nil}func (m *Mill) sweepUnclaimedOrphans() { m.mu.Lock() var unclaimed []*RemoteLease for _, lease := range m.leases { if !lease.orphaned { continue } if lease.claimed { continue } if lease.getState() == leaseDone { continue } if sess := m.sessions[lease.nodeID]; sess == nil || sess.disconnected { unclaimed = append(unclaimed, lease) } } m.mu.Unlock()
retry := false reason := "executor did not reconnect after mill restart" for _, lease := range unclaimed { m.l.Warn("failing unclaimed restored lease", "lease", lease.id, "node", lease.nodeID) if err := m.finishOrphan(lease, string(models.StatusKindFailed), &reason, nil); err != nil { m.l.Error("finish unclaimed restored lease", "lease", lease.id, "err", err) retry = true } } if retry { time.AfterFunc(5*time.Second, m.sweepUnclaimedOrphans) }}
func (m *Mill) reconcileLeases(sess *millSession, activeLeaseIDs []string) error { active := make(map[string]struct{}, len(activeLeaseIDs)) for _, id := range activeLeaseIDs { active[id] = struct{}{} }
known := make(map[string]struct{}) m.mu.Lock() var gone []*RemoteLease for _, lease := range m.leases { if lease.nodeID == sess.nodeID { known[lease.id] = struct{}{} if lease.epoch != sess.epoch { gone = append(gone, lease) } else if _, ok := active[lease.id]; !ok { gone = append(gone, lease) } } } for _, lease := range m.reservations { if lease.nodeID == sess.nodeID && lease.epoch == sess.epoch { known[lease.id] = struct{}{} } } m.mu.Unlock() var unknown []string for id := range active { if _, ok := known[id]; !ok { unknown = append(unknown, id) } }
for _, lease := range gone { status := string(models.StatusKindFailed) reason := "executor no longer holds lease" if cancelled, cancelReason := lease.cancelRequested(); cancelled { status = string(models.StatusKindCancelled) reason = cancelReason } m.l.Warn("finishing reconciled lease", "lease", lease.id, "node", sess.nodeID, "leaseInc", lease.epoch, "sessInc", sess.epoch) if lease.orphaned { if err := m.finishOrphan(lease, status, &reason, nil); err != nil { return err } } else if err := m.finishLiveLease(lease, status, reason); err != nil { return err } } for _, id := range unknown { if err := m.sendUntrackedCancel(sess, id, "lease is not owned by this mill"); err != nil { return fmt.Errorf("cancel unknown executor lease %q: %w", id, err) } } return nil}
func (m *Mill) completeLeaseRow(lease *RemoteLease, status string, errMsg *string, exitCode *int64) error { if m.db == nil { return nil } return m.db.CompleteMillLease( lease.id, string(lease.wid.PipelineId.AtUri()), lease.wid.Name, status, errMsg, exitCode, m.n, )}
func (m *Mill) finishLiveLease(lease *RemoteLease, status, reason string) error { lease.finishMu.Lock() defer lease.finishMu.Unlock() if lease.getState() == leaseDone { return nil } if err := m.completeLeaseRow(lease, status, &reason, nil); err != nil { return err } lease.deliverTerminal(&millv1.AttemptResult{ Status: mapTerminalStatusString(status), Error: reason, }) m.recordSyntheticTerminal(lease, status) if err := m.cleanupLeaseLocked(lease); err != nil { return err } return nil}
func (m *Mill) finishOrphan(lease *RemoteLease, status string, errMsg *string, exitCode *int64) error { lease.finishMu.Lock() defer lease.finishMu.Unlock() if lease.getState() == leaseDone { return nil } if err := m.completeLeaseRow(lease, status, errMsg, exitCode); err != nil { return err } lease.markDone() m.recordSyntheticTerminal(lease, status) if err := m.cleanupLeaseLocked(lease); err != nil { return err } return nil}
func (m *Mill) recordSyntheticTerminal(lease *RemoteLease, status string) { if !lease.millRecordsTerminalMetrics { return } result := "failure" class := engine.FailureClassInfrastructure reason := engine.FailureReasonExecutorLost switch status { case string(models.StatusKindSuccess): result = "success" class = engine.FailureClassNone reason = engine.FailureReasonSuccess case string(models.StatusKindTimeout): result = "timeout" class = engine.FailureClassUser reason = engine.FailureReasonTimeout case string(models.StatusKindCancelled): result = "cancelled" class = engine.FailureClassUser reason = engine.FailureReasonCancelled } m.metrics.RecordWorkflowTerminal(lease.engine, result, string(class), string(reason))}
func mapTerminalStatusString(s string) millv1.TerminalStatus { switch s { case "success": return millv1.TerminalStatus_TERMINAL_STATUS_SUCCESS case "failed": return millv1.TerminalStatus_TERMINAL_STATUS_FAILED case "timeout": return millv1.TerminalStatus_TERMINAL_STATUS_TIMEOUT case "cancelled": return millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED default: return millv1.TerminalStatus_TERMINAL_STATUS_UNSPECIFIED }}