package mill import ( "context" "encoding/json" "errors" "fmt" "io" "log/slog" "maps" "os" "path/filepath" "slices" "sort" "strings" "sync" "time" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" "tangled.org/core/notifier" "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/models" "tangled.org/core/spindle/observability" "tangled.org/core/spindle/quota" "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/storage" "tangled.org/core/tid" millproto "tangled.org/core/spindle/mill/proto" millv1 "tangled.org/core/spindle/mill/proto/gen" ) const ( defaultReconnectGrace = 45 * time.Second defaultJobTimeout = 24 * time.Hour defaultBidTimeout = 5 * time.Second defaultTopK = 3 defaultMaxPending = 100 defaultCancelAckTimeout = 10 * time.Second defaultCancelTeardownTimeout = 10 * time.Minute defaultCacheDeleteTimeout = 30 * time.Second ) var errProtocolViolation = errors.New("executor protocol violation") func protoErrf(format string, args ...any) error { return fmt.Errorf("%w: %s", errProtocolViolation, fmt.Sprintf(format, args...)) } type Config struct { // mill appends live-tailed executor lines here so logview can follow running remote jobs LogDir string MaxPending int ReconnectGrace time.Duration CacheStoreID string CacheMaxBytesPerOwner int64 JobTimeout time.Duration BidTimeout time.Duration TopK int CancelAckTimeout time.Duration // teardown can take minutes after the cancel is acknowledged CancelTeardownTimeout time.Duration // bounds time spent deleting an object from the shared cache CacheDeleteTimeout time.Duration // how many protocol-violating session deaths in a row quarantine the node QuarantineStrikes int } type cacheDeletion struct { entries []db.CacheEntry objects []string } type Mill struct { l *slog.Logger cfg Config db *db.DB n *notifier.Notifier metrics *observability.Metrics qm *quota.Manager cache storage.Storage mu sync.Mutex quotaLifecycleMu sync.Mutex sessions map[string]*millSession leases map[string]*RemoteLease reservations map[string]*RemoteLease nodeSeqno map[string]uint64 quotaLeases map[string]*RemoteLease protoStrikes map[string]int pending int changeCh chan struct{} cacheDeleteOnce sync.Once cacheDeleteMu sync.Mutex cacheDeleteQueue []cacheDeletion cacheDeleteWake chan struct{} leaseSeq uint64 } func New(l *slog.Logger, cfg Config) *Mill { if cfg.ReconnectGrace <= 0 { cfg.ReconnectGrace = defaultReconnectGrace } if cfg.JobTimeout <= 0 { cfg.JobTimeout = defaultJobTimeout } if cfg.BidTimeout <= 0 { cfg.BidTimeout = defaultBidTimeout } if cfg.TopK <= 0 { cfg.TopK = defaultTopK } if cfg.MaxPending <= 0 { cfg.MaxPending = defaultMaxPending } if cfg.CancelAckTimeout <= 0 { cfg.CancelAckTimeout = defaultCancelAckTimeout } if cfg.CancelTeardownTimeout <= 0 { cfg.CancelTeardownTimeout = defaultCancelTeardownTimeout } if cfg.CacheDeleteTimeout <= 0 { cfg.CacheDeleteTimeout = defaultCacheDeleteTimeout } m := &Mill{ l: l, cfg: cfg, sessions: make(map[string]*millSession), leases: make(map[string]*RemoteLease), reservations: make(map[string]*RemoteLease), nodeSeqno: make(map[string]uint64), changeCh: make(chan struct{}), quotaLeases: make(map[string]*RemoteLease), } m.startCacheDeletionWorker() return m } func (m *Mill) Attach(d *db.DB, n *notifier.Notifier, qm *quota.Manager) { m.mu.Lock() m.db = d m.n = n m.qm = qm m.mu.Unlock() } func (m *Mill) AttachCache(store storage.Storage) { m.mu.Lock() m.cache = store m.mu.Unlock() m.ensureCacheSentinel(store) } func (m *Mill) ensureCacheSentinel(store storage.Storage) { if store == nil { return } rc, err := store.Get(context.Background(), millproto.CacheSentinelKey) if err == nil { defer rc.Close() value, readErr := io.ReadAll(io.LimitReader(rc, int64(len(millproto.CacheSentinelValue)+1))) if readErr != nil || string(value) != millproto.CacheSentinelValue { m.l.Warn("cache store sentinel has unexpected contents", "key", millproto.CacheSentinelKey) } return } if !errors.Is(err, storage.ErrNotExist) { m.l.Warn("failed to read cache store sentinel", "key", millproto.CacheSentinelKey, "err", err) return } if err := store.Put(context.Background(), millproto.CacheSentinelKey, strings.NewReader(millproto.CacheSentinelValue)); err != nil { m.l.Warn("failed to initialize cache store sentinel", "key", millproto.CacheSentinelKey, "err", err) } } func (m *Mill) nextLeaseID() string { m.mu.Lock() m.leaseSeq++ seq := m.leaseSeq m.mu.Unlock() return fmt.Sprintf("%s-%d", tid.TID(), seq) } func (m *Mill) dropReservation(id string) { m.mu.Lock() delete(m.reservations, id) m.mu.Unlock() } func (m *Mill) notifyChange() { m.mu.Lock() m.notifyChangeLocked() m.mu.Unlock() } func (m *Mill) notifyChangeLocked() { close(m.changeCh) m.changeCh = make(chan struct{}) } func (m *Mill) currentChangeCh() <-chan struct{} { m.mu.Lock() defer m.mu.Unlock() return m.changeCh } func (m *Mill) attachSession(sess *millSession) (uint64, bool) { // load db state before taking the mill lock var dbCursor uint64 haveCursor := false if m.db != nil { if cur, err := m.db.GetExecutorCursor(sess.nodeID, sess.epoch); err == nil { dbCursor, haveCursor = cur, true } } m.mu.Lock() defer m.mu.Unlock() if old := m.sessions[sess.nodeID]; old != nil { if old.live(m.cfg.ReconnectGrace) { return 0, false } if old.recoveryExpired { m.l.Warn("rejecting executor reconnect while lease failure retries", "node", sess.nodeID) return 0, false } if old.graceTimer != nil { old.graceTimer.Stop() } old.close() if old.recovering { m.armSessionRecoveryLocked(sess, old.graceDeadline) } else { m.armSessionRecoveryLocked(sess, time.Now().Add(m.cfg.ReconnectGrace)) } if old.disconnected { m.metrics.RecordMillReconnect() m.l.Info("executor reconnected", "node", sess.nodeID) } else { m.l.Warn("replacing silent executor session", "node", sess.nodeID) } } m.sessions[sess.nodeID] = sess // wakes commit retries waiting out reconnect grace m.notifyChangeLocked() key := sess.nodeID + "/" + sess.epoch if haveCursor { m.nodeSeqno[key] = dbCursor } return m.nodeSeqno[key], true } func (m *Mill) touchSession(sess *millSession) bool { m.mu.Lock() defer m.mu.Unlock() if m.sessions[sess.nodeID] != sess || sess.disconnected { return false } sess.lastSeen = time.Now() return true } func (m *Mill) armSessionRecoveryLocked(sess *millSession, deadline time.Time) { sess.recovering = true sess.recoveryReady = false sess.graceDeadline = deadline delay := time.Until(deadline) if delay < 0 { delay = 0 } sess.graceTimer = time.AfterFunc(delay, func() { m.failLeasesAfterGrace(sess) }) } func (m *Mill) extendGraceLocked(sess *millSession) { if !sess.recovering || sess.recoveryExpired { return } deadline := time.Now().Add(m.cfg.ReconnectGrace) if deadline.After(sess.graceDeadline) { sess.graceDeadline = deadline if sess.graceTimer != nil { sess.graceTimer.Stop() } delay := time.Until(deadline) sess.graceTimer = time.AfterFunc(delay, func() { m.failLeasesAfterGrace(sess) }) } } func (m *Mill) maybeFinishSessionRecoveryLocked(sess *millSession) { if m.sessions[sess.nodeID] != sess || sess.disconnected || !sess.recovering || sess.recoveryExpired || !sess.recoveryReady || len(sess.pendingCancels) != 0 { return } if sess.graceTimer != nil { sess.graceTimer.Stop() } sess.recovering = false sess.graceDeadline = time.Time{} sess.graceTimer = nil } func (m *Mill) detachSession(sess *millSession) { m.mu.Lock() if m.sessions[sess.nodeID] != sess { // already replaced by a reconnect m.mu.Unlock() sess.close() return } sess.disconnected = true if !sess.recovering { m.armSessionRecoveryLocked(sess, time.Now().Add(m.cfg.ReconnectGrace)) } m.mu.Unlock() sess.close() m.l.Warn("executor session lost; entering reconnect grace", "node", sess.nodeID, "grace", m.cfg.ReconnectGrace) m.notifyChange() } func (m *Mill) noteSessionError(sess *millSession, err error) { if errors.Is(err, errProtocolViolation) { m.l.Warn("executor session rejected on protocol violation", "node", sess.nodeID, "err", err) } } func (m *Mill) sessionReady(sess *millSession) { for _, lease := range m.cancelledLeasesForNode(sess.nodeID) { _, reason := lease.cancelRequested() m.sendCancel(sess, lease, reason) } m.notifyChange() } func (m *Mill) hasLiveExecutor(nodeID string) bool { m.mu.Lock() defer m.mu.Unlock() sess := m.sessions[nodeID] return sess != nil && sess.live(m.cfg.ReconnectGrace) } func (m *Mill) cancelledLeasesForNode(nodeID string) []*RemoteLease { m.mu.Lock() var candidates []*RemoteLease for _, lease := range m.leases { if lease.nodeID == nodeID { candidates = append(candidates, lease) } } m.mu.Unlock() var leases []*RemoteLease for _, lease := range candidates { if lease.cancelPending() { leases = append(leases, lease) } } return leases } // reconnect recovery did not complete before grace expired. fails every // lease the node held, and reschedules itself if any failure didn't stick func (m *Mill) failLeasesAfterGrace(sess *millSession) { m.failLeasesAfterGraceAttempt(sess, false) } func (m *Mill) failLeasesAfterGraceAttempt(sess *millSession, retry bool) { m.mu.Lock() if m.sessions[sess.nodeID] != sess || !sess.recovering || (sess.recoveryExpired && !retry) { m.mu.Unlock() return } sess.recoveryExpired = true var dead []*RemoteLease for _, lease := range m.leases { if lease.nodeID == sess.nodeID { dead = append(dead, lease) } } m.mu.Unlock() m.l.Warn("executor declared dead; failing its in-flight jobs", "node", sess.nodeID, "jobs", len(dead)) deadReason := "executor lost" success := true for _, lease := range dead { switch { case lease.orphaned: // restored leases have no RunStep waiting, fail them straight // into the event stream if err := m.finishOrphan(lease, string(models.StatusKindFailed), &deadReason, nil); err != nil { m.l.Error("finish orphan after executor loss failed, will retry", "lease", lease.id, "err", err) success = false } default: // live leases have a blocked RunStep, keep a pending cancellation // as the terminal reason status := string(models.StatusKindFailed) reason := deadReason if cancelled, cancelReason := lease.cancelRequested(); cancelled { status = string(models.StatusKindCancelled) reason = cancelReason } if err := m.finishLiveLease(lease, status, reason); err != nil { m.l.Error("finish lease after executor loss failed, will retry", "lease", lease.id, "err", err) success = false } } } m.mu.Lock() if !success { // something didn't finish cleanly, try again shortly. finished // leases skip themselves on the next pass sess.graceTimer = time.AfterFunc(5*time.Second, func() { m.failLeasesAfterGraceAttempt(sess, true) }) m.mu.Unlock() return } // everything failed cleanly. forget the node and wake placement current := m.sessions[sess.nodeID] == sess if current { delete(m.sessions, sess.nodeID) } m.mu.Unlock() if current { sess.close() } m.notifyChange() m.sweepUnclaimedOrphans() } func (m *Mill) place(ctx context.Context, engineName string, wid models.WorkflowId, wf *models.Workflow) (engine.WorkflowSlot, error) { placementStart := time.Now() placementResult := "error" defer func() { m.metrics.RecordMillPlacementResult(ctx, engineName, placementResult, time.Since(placementStart)) }() m.mu.Lock() if m.cfg.MaxPending > 0 && m.pending >= m.cfg.MaxPending { max := m.cfg.MaxPending cur := m.pending m.mu.Unlock() placementResult = "rejected" m.metrics.RecordMillPlacementAdmission(false, "max_pending_reached") return nil, fmt.Errorf("%w: mill has %d pending jobs (max %d)", engine.ErrNoWorkflowSlots, cur, max) } m.metrics.RecordMillPlacementAdmission(true, "allowed") m.pending++ m.mu.Unlock() defer func() { m.mu.Lock() m.pending-- m.mu.Unlock() }() for { if err := ctx.Err(); err != nil { placementResult = "cancelled" if errors.Is(err, context.DeadlineExceeded) { placementResult = "timeout" } return nil, err } // capture the channel before bidding // a change during the bid makes the next wait return immediately ch := m.currentChangeCh() lease, err := m.bid(ctx, engineName, wid, wf) if err != nil { return nil, err } if lease != nil { wf.CacheNamespace = lease.cacheNamespace slot, retry, err := m.admit(ctx, engineName, wid, wf, lease) if err != nil { var failure *engine.WorkflowFailure if errors.As(err, &failure) && failure.Class == engine.FailureClassPolicy { placementResult = "rejected" } return nil, err } if slot != nil { placementResult = "success" return slot, nil } if retry { continue } } // no executor available. wait for a change or ctx select { case <-ctx.Done(): placementResult = "cancelled" if errors.Is(ctx.Err(), context.DeadlineExceeded) { placementResult = "timeout" } return nil, ctx.Err() case <-ch: } } } func (m *Mill) admit(ctx context.Context, engineName string, wid models.WorkflowId, wf *models.Workflow, lease *RemoteLease) (engine.WorkflowSlot, bool, error) { req, ok := m.quotaRequest(wid, wf, lease.resources) if !ok { return m.publish(wid, wf, lease, nil) } qlease, res, err := m.qm.TryAcquire(ctx, req) switch { case err != nil: m.releaseSeat(lease) return nil, false, fmt.Errorf("acquire workflow quota: %w", err) case res.Allowed: return m.publish(wid, wf, lease, qlease) case !res.Temporary: m.releaseSeat(lease) return nil, false, engine.ClassifiedFailure( engine.FailureClassPolicy, engine.FailureReasonQuotaDenied, fmt.Errorf("%w: workflow quota denied: %s", engine.ErrWorkflowFailed, res.Reason), ) } // don't hold fleet capacity while waiting for quota m.releaseSeat(lease) qlease, err = m.qm.Acquire(ctx, req) if err != nil { if errors.Is(err, quota.ErrDenied) { err = engine.ClassifiedFailure(engine.FailureClassPolicy, engine.FailureReasonQuotaDenied, err) } return nil, false, fmt.Errorf("wait for workflow quota: %w", err) } // a re-bid must match the reserved resources next, err := m.bid(ctx, engineName, wid, wf) if err != nil { _ = m.releaseAbandonedQuota(qlease) return nil, false, err } if next == nil { _ = m.releaseAbandonedQuota(qlease) return nil, false, nil } if !maps.Equal(next.resources, req.Resources) { m.l.Warn("re-bid reported different resources than the quota reserved, retrying placement", "lease", next.id, "node", next.nodeID, "engine", engineName) m.releaseSeat(next) _ = m.releaseAbandonedQuota(qlease) return nil, true, nil } return m.publish(wid, wf, next, qlease) } func (m *Mill) releaseAbandonedQuota(qlease quota.Lease) error { if qlease == nil { return nil } err := qlease.ReleaseWithError() if err != nil { dummy := &RemoteLease{ id: "abandoned-" + qlease.ID(), quotaID: qlease.ID(), state: leaseDone, cleanedUp: false, } dummy.finishMu.Lock() m.scheduleCleanupLocked(dummy) dummy.finishMu.Unlock() m.l.Error("abandoned quota release failed, retry scheduled", "quota_id", qlease.ID(), "err", err) } return err } // teardown releases the attached quota exactly once func (m *Mill) publish(wid models.WorkflowId, wf *models.Workflow, lease *RemoteLease, qlease quota.Lease) (engine.WorkflowSlot, bool, error) { lease.wid = wid if id, ok := m.quotaIdentity(wf); ok { lease.ownerDID = id.OwnerDID lease.repoDID = id.RepoDID } if qlease != nil { lease.quotaLease = qlease lease.quotaID = qlease.ID() } if err := m.persistLease(lease, leaseRowReserved); err != nil { m.releaseSeat(lease) // release an unrecorded reservation releaseErr := m.releaseQuotaLocked(lease) if releaseErr != nil { lease.finishMu.Lock() m.scheduleCleanupLocked(lease) lease.finishMu.Unlock() } return nil, false, errors.Join(fmt.Errorf("persist reserved mill lease: %w", err), releaseErr) } m.mu.Lock() delete(m.reservations, lease.id) m.leases[lease.id] = lease if st, ok := wf.Data.(*millWorkflowState); ok && st != nil { st.Lease = lease } m.mu.Unlock() return &millSlot{fleet: m, lease: lease}, false, nil } func (m *Mill) releaseSeat(lease *RemoteLease) { m.dropReservation(lease.id) m.releaseRemote(lease) } func (m *Mill) quotaRequest(wid models.WorkflowId, wf *models.Workflow, resources quota.Resources) (quota.ReserveRequest, bool) { if m.qm == nil || len(resources) == 0 { return quota.ReserveRequest{}, false } id, ok := m.quotaIdentity(wf) if !ok { m.l.Warn("placing workflow without quota because the pipeline has no repo identity", "workflow", wid.String()) return quota.ReserveRequest{}, false } resID := quota.WorkflowReservationID(wf.RunID, id.OwnerDID, id.RepoDID, wid.Knot, wid.Rkey, wid.Name) return quota.ReserveRequest{ ID: resID, Kind: quota.KindWorkflow, Key: resID, Identity: id, Resources: resources, }, true } func (m *Mill) quotaIdentity(wf *models.Workflow) (quota.Identity, bool) { st, ok := wf.Data.(*millWorkflowState) if !ok || st == nil || st.RawPipeline.TriggerMetadata == nil { return quota.Identity{}, false } repo := st.RawPipeline.TriggerMetadata.Repo if repo == nil { return quota.Identity{}, false } id := quota.Identity{OwnerDID: repo.Did, RepoDID: repo.Did} if repo.RepoDid != nil && *repo.RepoDid != "" { id.RepoDID = *repo.RepoDid } if !strings.HasPrefix(id.OwnerDID, "did:") || !strings.HasPrefix(id.RepoDID, "did:") { return quota.Identity{}, false } return id, true } func (m *Mill) bid(ctx context.Context, engineName string, wid models.WorkflowId, wf *models.Workflow) (*RemoteLease, error) { rawPipeline, rawWorkflow, err := marshalJob(wf) if err != nil { return nil, err } requiredLabels := requiredLabels(wf) candidates := m.rankCandidates(engineName, requiredLabels, len(wf.Caches) > 0) if len(candidates) == 0 { return nil, nil } type bidResult struct { sess *millSession lease *RemoteLease rank int incompatible bool reason string failureClass string failureReason string } limit := m.cfg.TopK if limit <= 0 { limit = len(candidates) } if len(candidates) > limit { candidates = candidates[:limit] } results := make(chan bidResult, limit) ask := func(rank int, sess *millSession) { bidCtx, cancel := context.WithTimeout(ctx, m.cfg.BidTimeout) defer cancel() leaseID := m.nextLeaseID() lease := newLease(leaseID, sess.nodeID, sess.epoch, engineName) lease.cacheNamespace = sess.cacheNamespace m.mu.Lock() m.reservations[leaseID] = lease m.mu.Unlock() bidCtx, span := observability.Tracer().Start(bidCtx, "mill.placement.bid") if span.IsRecording() { var ownerDID, repoDID string if wf != nil { ownerDID = wf.OwnerDID repoDID = wf.RepoDID } if (ownerDID == "" || repoDID == "") && wf != nil { if id, ok := m.quotaIdentity(wf); ok { ownerDID = id.OwnerDID repoDID = id.RepoDID } } attrs := []attribute.KeyValue{ attribute.String(observability.WorkflowIDKey, wid.String()), attribute.String(observability.LeaseIDKey, leaseID), attribute.String(observability.ExecutorNodeIDKey, sess.nodeID), attribute.String(observability.PipelineIDKey, wid.PipelineId.AtUri().String()), } if ownerDID != "" { attrs = append(attrs, attribute.String(observability.OwnerDIDKey, ownerDID)) } if repoDID != "" { attrs = append(attrs, attribute.String(observability.RepoDIDKey, repoDID)) } span.SetAttributes(attrs...) } accepted := false defer func() { if accepted { span.SetStatus(codes.Ok, "accepted") } else { span.SetStatus(codes.Error, "placement bid failed") } span.End() }() traceparent, tracestate := observability.InjectToTraceparentAndTracestate(bidCtx) msg := &millproto.Message{ReserveSeat: &millv1.ReserveSeat{ LeaseId: leaseID, TargetEngine: engineName, RawPipelineJson: rawPipeline, RawWorkflowJson: rawWorkflow, Knot: wid.Knot, Rkey: wid.Rkey, TtlSeconds: uint32(m.cfg.ReconnectGrace / time.Second), Traceparent: traceparent, Tracestate: tracestate, RepoDid: wf.RepoDID, }} resp, err := sess.request(bidCtx, leaseID, msg) if err != nil { m.dropReservation(leaseID) results <- bidResult{} return } rr := resp.GetReserveResult() if rr == nil { m.dropReservation(leaseID) results <- bidResult{} return } if !rr.GetAccepted() { m.dropReservation(leaseID) if rr.GetRejectClass() == millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE { results <- bidResult{ sess: sess, rank: rank, incompatible: true, reason: rr.GetRejectReason(), failureClass: rr.GetFailureClass(), failureReason: rr.GetFailureReason(), } return } results <- bidResult{} return } // invalid resource claims are protocol errors if err := quota.ValidateResources(rr.GetQuotaResources()); err != nil { m.l.WarnContext(bidCtx, "executor reported invalid workflow resources", "node", sess.nodeID, "lease", leaseID, "err", err) m.dropReservation(leaseID) m.releaseRemote(lease) results <- bidResult{} return } lease.resources = maps.Clone(rr.GetQuotaResources()) lease.millRecordsTerminalMetrics = rr.GetSupportsMillTerminalMetrics() accepted = true results <- bidResult{sess: sess, lease: lease, rank: rank} } next := 0 inFlight := 0 for next < len(candidates) && inFlight < limit { inFlight++ go ask(next, candidates[next]) next++ } var winner *bidResult var losers []*RemoteLease var incompatible []string var incompatibleClass, incompatibleReason string attributionConsistent := true attributionSet := false // any soft reject (transient or timeout) means the fleet was just // busy, so an all-incompatible outcome isn't a hard placement error softReject := false for inFlight > 0 { r := <-results inFlight-- // incompatible rejects get reported to the user, any other failure // just means the fleet is busy if r.incompatible { if r.reason != "" { incompatible = append(incompatible, r.reason) } if r.failureClass == "" || r.failureReason == "" { attributionConsistent = false } else if !attributionSet { incompatibleClass = r.failureClass incompatibleReason = r.failureReason attributionSet = true } else if incompatibleClass != r.failureClass || incompatibleReason != r.failureReason { attributionConsistent = false } } else if r.lease == nil { softReject = true } // a failed bid means ask the next candidate, unless someone won if r.lease == nil { for winner == nil && next < len(candidates) && inFlight < limit { inFlight++ go ask(next, candidates[next]) next++ } continue } // if this bid is worse than the winner it goes to the losers pile if winner != nil && r.rank >= winner.rank { losers = append(losers, r.lease) continue } // otherwise it's the new best and the old winner joins the losers if winner != nil { losers = append(losers, winner.lease) } winner = &r } // let the losers go so they free their seats right away for _, l := range losers { m.dropReservation(l.id) m.releaseRemote(l) } if winner == nil { if len(incompatible) > 0 && !softReject { err := fmt.Errorf("no compatible executor for %s: %s", engineName, strings.Join(incompatible, "; ")) if attributionSet && attributionConsistent { return nil, engine.ClassifiedFailure( engine.FailureClass(incompatibleClass), engine.FailureReason(incompatibleReason), err, ) } return nil, err } return nil, nil } return winner.lease, nil } // ranks nodes that are least busy first. if a resource is used a lot // then that node will lose to one that is more even across the board. func (m *Mill) rankCandidates(engineName string, requiredLabels []string, cacheRequired bool) []*millSession { m.mu.Lock() defer m.mu.Unlock() type ranked struct { sess *millSession worst float64 sum float64 } var rs []ranked for _, s := range m.sessions { // only live, reporting sessions can take work if s.disconnected || s.recovering || len(s.pendingCancels) != 0 { continue } if s.snapshot == nil { continue } if cacheRequired && (m.cfg.CacheStoreID == "" || s.cacheStoreID != m.cfg.CacheStoreID || s.legacyProtocol) { continue } // the engine has to exist and have room right now ea, ok := s.snapshot.GetEngines()[engineName] if !ok || !ea.GetAvailable() { continue } // and satisfy the wf's label requirements if !hasLabels(s.labels, requiredLabels) { continue } worst, sum := loadScore(ea.GetLoad()) rs = append(rs, ranked{sess: s, worst: worst, sum: sum}) } slices.SortStableFunc(rs, func(a, b ranked) int { if a.worst < b.worst { return -1 } if a.worst > b.worst { return 1 } if a.sum < b.sum { return -1 } if a.sum > b.sum { return 1 } return 0 }) out := make([]*millSession, len(rs)) for i := range rs { out[i] = rs[i].sess } return out } func loadScore(load map[string]float64) (worst, sum float64) { for _, v := range load { if v > worst { worst = v } sum += v } return worst, sum } func requiredLabels(wf *models.Workflow) []string { st, ok := wf.Data.(*millWorkflowState) if !ok || st == nil { return nil } return st.RawWorkflow.RunsOn } func hasLabels(labels []string, required []string) bool { for _, want := range required { if !slices.Contains(labels, want) { return false } } return true } func (m *Mill) commitAndWait(ctx context.Context, wf *models.Workflow, unlocked []secrets.UnlockedSecret) (err error) { st, ok := wf.Data.(*millWorkflowState) if !ok || st == nil || st.Lease == nil { return fmt.Errorf("mill workflow state missing lease") } lease := st.Lease ctx, span := observability.Tracer().Start(ctx, "mill.commit") if span.IsRecording() { attrs := []attribute.KeyValue{ attribute.String(observability.WorkflowIDKey, lease.wid.String()), attribute.String(observability.LeaseIDKey, lease.id), attribute.String(observability.ExecutorNodeIDKey, lease.nodeID), attribute.String(observability.PipelineIDKey, lease.wid.PipelineId.AtUri().String()), } if lease.ownerDID != "" { attrs = append(attrs, attribute.String(observability.OwnerDIDKey, lease.ownerDID)) } if lease.repoDID != "" { attrs = append(attrs, attribute.String(observability.RepoDIDKey, lease.repoDID)) } span.SetAttributes(attrs...) } defer func() { if err != nil { span.SetStatus(codes.Error, "commit failed") } else { span.SetStatus(codes.Ok, "success") } span.End() }() capabilities, err := cacheCapabilities(wf.RepoDID, wf.CacheBindings) if err != nil { return err } if m.db != nil { if err := m.db.SaveMillCacheCapabilities(lease.id, capabilities); err != nil { return fmt.Errorf("persist mill cache capabilities: %w", err) } } pbSecrets := make([]*millv1.Secret, len(unlocked)) for i, s := range unlocked { pbSecrets[i] = &millv1.Secret{Key: s.Key, Value: s.Value} } traceparent, tracestate := observability.InjectToTraceparentAndTracestate(ctx) commit := &millproto.Message{CommitLease: &millv1.CommitLease{ LeaseId: lease.id, Secrets: pbSecrets, Traceparent: traceparent, Tracestate: tracestate, MillRecordsTerminalMetrics: lease.millRecordsTerminalMetrics, CacheBindings: cacheBindingsToProto(wf.CacheBindings), }} // commit retries ride reconnects, a reservation outlives one // disconnect. lost session or slow executor just means wait and retry, // only job timeout or a dead lease stops the loop for { if res, ok := pollTerminal(lease); ok { return terminalError(res.Status) } if !lease.markCommitting() { return engine.ErrWorkflowFailed } sess := m.sessionForNode(lease.nodeID) if sess == nil { if done, err := m.waitCommitRetry(ctx, lease); done || err != nil { return err } continue } reqCtx, cancel := context.WithTimeout(ctx, m.cfg.BidTimeout) resp, err := sess.request(reqCtx, lease.id, commit) cancel() if err != nil { switch { case errors.Is(err, errSessionClosed): // session died mid-request. wait out the grace, then retry on // the new one if done, err := m.waitCommitRetry(ctx, lease); done || err != nil { return err } continue case errors.Is(err, context.DeadlineExceeded) && ctx.Err() == nil: // executor didn't answer in time, but its seat is still held so // retrying is safe continue case errors.Is(err, context.DeadlineExceeded): // the job ctx itself ran out, a real timeout return engine.ErrTimedOut case errors.Is(err, context.Canceled): return err default: m.l.Warn("commit lease send failed; waiting for reconnect", "lease", lease.id, "node", lease.nodeID, "err", err) if done, err := m.waitCommitRetry(ctx, lease); done || err != nil { return err } continue } } if resp.GetCommitted() == nil { return engine.ErrWorkflowFailed } lease.markRunning() if err := m.persistLease(lease, leaseRowRunning); err != nil { m.l.Error("persist running mill lease", "lease", lease.id, "err", err) } if cancelled, reason := lease.cancelRequested(); cancelled { m.sendCancel(sess, lease, reason) } break } select { case res := <-lease.terminal: return terminalError(res.Status) case <-ctx.Done(): if ctx.Err() == context.DeadlineExceeded { return engine.ErrTimedOut } return ctx.Err() } } func (m *Mill) waitCommitRetry(ctx context.Context, lease *RemoteLease) (bool, error) { // grab the channel before checking for a live session again // a reconnect will still close the channel if it happens in between ch := m.currentChangeCh() if m.sessionForNode(lease.nodeID) != nil { return false, nil } select { case res := <-lease.terminal: return true, terminalError(res.Status) case <-ctx.Done(): if ctx.Err() == context.DeadlineExceeded { return true, engine.ErrTimedOut } return true, ctx.Err() case <-ch: return false, nil } } func terminalError(status millv1.TerminalStatus) error { switch status { case millv1.TerminalStatus_TERMINAL_STATUS_SUCCESS: return nil case millv1.TerminalStatus_TERMINAL_STATUS_TIMEOUT: return engine.ErrTimedOut case millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED: return engine.ErrWorkflowCanceled default: return engine.ErrWorkflowFailed } } func pollTerminal(lease *RemoteLease) (*millv1.AttemptResult, bool) { select { case res := <-lease.terminal: return res, true default: return nil, false } } func (m *Mill) destroy(wid models.WorkflowId) { m.mu.Lock() var lease *RemoteLease for _, l := range m.leases { if l.wid == wid { lease = l break } } m.mu.Unlock() if lease == nil { return } reason := "workflow destroyed" switch lease.requestCancel(reason) { case cancelLocal: if sess := m.sessionForNode(lease.nodeID); sess != nil { _ = sess.send(&millproto.Message{ReleaseLease: &millv1.ReleaseLease{LeaseId: lease.id}}) } lease.deliverCancelled(reason) case cancelRemote: if sess := m.sessionForNode(lease.nodeID); sess != nil { m.sendCancel(sess, lease, reason) } } } func (m *Mill) releaseSlot(s *millSlot) { lease := s.lease state, cancelled := lease.releaseState() if state == leaseReserved { m.releaseRemote(lease) } else if cancelled && state != leaseDone { return } if err := m.cleanupLease(lease); err != nil { m.l.Error("releaseSlot cleanupLease failed", "lease", lease.id, "err", err) } } func (m *Mill) cleanupLease(lease *RemoteLease) error { lease.finishMu.Lock() defer lease.finishMu.Unlock() return m.cleanupLeaseLocked(lease) } func (m *Mill) cleanupLeaseLocked(lease *RemoteLease) error { if lease.cleanedUp { return nil } if m.db != nil { // fence unreported save uploads before dropping the lease capabilities; // a late object is then covered by a durable deletion tombstone pendingKeys, err := m.db.DiscardPendingCacheEntriesForLease(lease.id, time.Now()) if err != nil { m.scheduleCleanupLocked(lease) return err } if len(pendingKeys) > 0 { objects := make(map[string]struct{}, len(pendingKeys)) for _, key := range pendingKeys { objects[key] = struct{}{} } m.enqueueCacheDeletions(nil, objects) } // delete the row before releasing the reservation if err := m.db.DeleteMillLease(lease.id); err != nil { m.scheduleCleanupLocked(lease) return err } } m.mu.Lock() delete(m.leases, lease.id) m.mu.Unlock() if err := m.releaseQuotaLocked(lease); err != nil { m.scheduleCleanupLocked(lease) return err } lease.cleanedUp = true m.notifyChange() return nil } func (m *Mill) releaseQuotaLocked(lease *RemoteLease) error { if m.qm == nil { return nil } lease.mu.Lock() workflowID := lease.quotaID if workflowID == "" && lease.quotaLease != nil { workflowID = lease.quotaLease.ID() } remoteQuotas := maps.Clone(lease.remoteQuotas) lease.mu.Unlock() ctx := context.Background() var releaseErr error if workflowID != "" { if err := m.qm.Release(ctx, workflowID); err != nil { releaseErr = errors.Join(releaseErr, fmt.Errorf("release workflow quota: %w", err)) } else { lease.mu.Lock() if lease.quotaID == workflowID { lease.quotaLease = nil lease.quotaID = "" } lease.mu.Unlock() } } for id, rq := range remoteQuotas { if rq.committing { if err := m.qm.Commit(ctx, id); err != nil { releaseErr = errors.Join(releaseErr, fmt.Errorf("commit remote quota %s: %w", id, err)) continue } } else { if err := m.qm.Release(ctx, id); err != nil { releaseErr = errors.Join(releaseErr, fmt.Errorf("release remote quota %s: %w", id, err)) continue } } lease.mu.Lock() delete(lease.remoteQuotas, id) lease.mu.Unlock() m.quotaLifecycleMu.Lock() if m.quotaLeases[id] == lease { delete(m.quotaLeases, id) } m.quotaLifecycleMu.Unlock() } return releaseErr } func (m *Mill) scheduleCleanupLocked(lease *RemoteLease) { if lease.cleanedUp || lease.cleanupRetry { return } lease.cleanupRetry = true time.AfterFunc(5*time.Second, func() { lease.finishMu.Lock() lease.cleanupRetry = false err := m.cleanupLeaseLocked(lease) lease.finishMu.Unlock() if err != nil { m.l.Error("retry lease cleanup failed", "lease", lease.id, "err", err) } }) } func (m *Mill) releaseRemote(lease *RemoteLease) { lease.setState(leaseDone) if sess := m.sessionForNode(lease.nodeID); sess != nil { _ = sess.send(&millproto.Message{ReleaseLease: &millv1.ReleaseLease{LeaseId: lease.id}}) } } func (m *Mill) sendCancel(sess *millSession, lease *RemoteLease, reason string) { lease.finishMu.Lock() m.mu.Lock() if m.sessions[lease.nodeID] != sess || sess.disconnected || m.leases[lease.id] != lease || lease.getState() == leaseDone { m.mu.Unlock() lease.finishMu.Unlock() return } sess.pendingCancels[lease.id] = struct{}{} m.mu.Unlock() lease.finishMu.Unlock() if err := sess.send(&millproto.Message{CancelAttempt: &millv1.CancelAttempt{ LeaseId: lease.id, Reason: reason, }}); err != nil { m.settleSessionCancel(sess, lease.id) sess.close() return } time.AfterFunc(m.cfg.CancelAckTimeout, func() { m.checkCancelAck(sess, lease) }) } func (m *Mill) sendUntrackedCancel(sess *millSession, leaseID, reason string) error { m.mu.Lock() if m.sessions[sess.nodeID] != sess || sess.disconnected { m.mu.Unlock() return errSessionClosed } sess.pendingCancels[leaseID] = struct{}{} m.mu.Unlock() if err := sess.send(&millproto.Message{CancelAttempt: &millv1.CancelAttempt{ LeaseId: leaseID, Reason: reason, }}); err != nil { m.settleSessionCancel(sess, leaseID) sess.close() return err } time.AfterFunc(m.cfg.CancelAckTimeout, func() { m.checkUntrackedCancelAck(sess, leaseID) }) return nil } func (m *Mill) checkUntrackedCancelAck(sess *millSession, leaseID string) { m.mu.Lock() _, pending := sess.pendingCancels[leaseID] isCurrent := m.sessions[sess.nodeID] == sess && !sess.disconnected m.mu.Unlock() if !pending || !isCurrent { return } m.l.Warn("node did not acknowledge unknown lease cancel within deadline; closing session", "node", sess.nodeID, "lease", leaseID, "deadline", m.cfg.CancelAckTimeout) sess.close() } func (m *Mill) sessionForNode(nodeID string) *millSession { m.mu.Lock() defer m.mu.Unlock() sess := m.sessions[nodeID] if sess == nil || sess.disconnected { return nil } return sess } func (m *Mill) onSnapshot(sess *millSession, snap *millv1.NodeSnapshot) error { m.mu.Lock() if sess.snapshot != nil && snap.Seqno <= sess.snapshot.Seqno { m.mu.Unlock() return protoErrf("snapshot seqno regression. Got %d, last seen %d", snap.Seqno, sess.snapshot.Seqno) } sess.snapshot = snap m.mu.Unlock() if err := m.reconcileLeases(sess, snap.GetActiveLeaseIds()); err != nil { return err } m.mu.Lock() for _, id := range snap.GetActiveLeaseIds() { if lease := m.leases[id]; lease != nil && lease.nodeID == sess.nodeID && lease.epoch == sess.epoch { lease.claimed = true } } sess.recoveryReady = true m.maybeFinishSessionRecoveryLocked(sess) m.mu.Unlock() m.notifyChange() return nil } func (m *Mill) onEventBatch(sess *millSession, batch *millv1.EventBatch) error { if batch == nil { return nil } if batch.Epoch != sess.epoch { return protoErrf("batch epoch %q does not match session %q", batch.Epoch, sess.epoch) } m.mu.Lock() currentKey := sess.nodeID + "/" + sess.epoch current := m.nodeSeqno[currentKey] m.mu.Unlock() expected := current + 1 var newEntries []*millv1.Event for _, entry := range batch.Events { // reconnect replays old seqnos, drop those if entry.Seqno <= current { continue } // gap means executor lost rows. so dont apply a partial batch if entry.Seqno != expected { return protoErrf("gap in stream seqnos. Expected %d, got %d", expected, entry.Seqno) } newEntries = append(newEntries, entry) expected++ } // all replays, still ack so the executor can trim its outbox if len(newEntries) == 0 { return m.sendAck(sess, current) } type pendingTerminal struct { lease *RemoteLease ar *millv1.AttemptResult } type startupObservation struct { engine string delay time.Duration } leaseSet := make(map[*RemoteLease]struct{}) for _, entry := range newEntries { m.mu.Lock() lease := m.leases[entry.LeaseId] m.mu.Unlock() if lease == nil || lease.nodeID != sess.nodeID { continue } if lease.epoch != "" && lease.epoch != sess.epoch { return protoErrf("lease %q epoch %q does not match session %q", lease.id, lease.epoch, sess.epoch) } leaseSet[lease] = struct{}{} } lockedLeases := make([]*RemoteLease, 0, len(leaseSet)) for lease := range leaseSet { lockedLeases = append(lockedLeases, lease) } sort.Slice(lockedLeases, func(i, j int) bool { return lockedLeases[i].id < lockedLeases[j].id }) for _, lease := range lockedLeases { lease.finishMu.Lock() } unlockLeases := func() { for i := len(lockedLeases) - 1; i >= 0; i-- { lockedLeases[i].finishMu.Unlock() } } // precompute storage sizes so object-storage stat does not run inside the write tx precomputedSizes := make(map[string]int64) for _, entry := range newEntries { if u := entry.GetCacheUpdate(); u != nil && u.GetAction() == millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_STORED { var storageKey string if m.db != nil { _ = m.db.QueryRow(`select storage_key from mill_cache_capabilities where lease_id = ? and action = 'save' and cache_id = ?`, entry.LeaseId, u.GetId()).Scan(&storageKey) } if storageKey != "" { if sz, err := m.cacheObjectSize(sess.ctx, storageKey); err == nil { precomputedSizes[storageKey] = sz } } } } type batchApply struct { terminalCancelIDs []string pendingTerminals []pendingTerminal artifactLeases []*RemoteLease startupObservations []startupObservation highestSeqno uint64 cacheDeletes map[string]db.CacheEntry cacheObjectDeletes map[string]struct{} } applyFunc := func(tx *db.EventBatchTx) (batchApply, error) { var r batchApply r.highestSeqno = current r.cacheDeletes = make(map[string]db.CacheEntry) r.cacheObjectDeletes = make(map[string]struct{}) finishedInBatch := make(map[string]struct{}) for _, entry := range newEntries { if entry.GetAttemptResult() != nil { r.terminalCancelIDs = append(r.terminalCancelIDs, entry.LeaseId) } m.mu.Lock() lease := m.leases[entry.LeaseId] m.mu.Unlock() // events for leases this node doesn't own are skipped but still // count as processed if lease == nil || lease.nodeID != sess.nodeID { r.highestSeqno = entry.Seqno continue } // events arriving for a different epoch are invalid (different session) if lease.epoch != "" && lease.epoch != sess.epoch { return batchApply{}, protoErrf("lease %q epoch %q does not match session %q", lease.id, lease.epoch, sess.epoch) } lease.mu.Lock() state := lease.state lease.mu.Unlock() // done leases can replay terminals on reconnect, skip them if state == leaseDone { r.highestSeqno = entry.Seqno continue } if _, finished := finishedInBatch[lease.id]; finished { return batchApply{}, protoErrf("stream entry follows terminal for lease %q", lease.id) } switch { case entry.GetStatusEvent() != nil: ev := entry.GetStatusEvent() statusStr := string(models.StatusKindRunning) var errMsg *string if e := ev.GetError(); e != "" { errMsg = &e } var exitCode *int64 if c := ev.GetExitCode(); c != 0 { exitCode = &c } pipelineAtUri := string(lease.wid.PipelineId.AtUri()) if tx != nil { alreadyRunning, err := tx.HasWorkflowStatus(context.Background(), pipelineAtUri, lease.wid.Name, statusStr) if err != nil { return batchApply{}, err } if err := tx.InsertStatusEvent(pipelineAtUri, lease.wid.Name, statusStr, errMsg, exitCode); err != nil { return batchApply{}, err } if !alreadyRunning { delay, ok, err := tx.WorkflowStartupDelay(context.Background(), pipelineAtUri, lease.wid.Name) if err != nil { return batchApply{}, err } if ok { r.startupObservations = append(r.startupObservations, startupObservation{ engine: lease.engine, delay: delay, }) } } } case entry.GetCacheUpdate() != nil: update := entry.GetCacheUpdate() if tx != nil { at := time.Now() capabilityAction := "" switch update.GetAction() { case millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_USED, millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_MISSING: capabilityAction = "restore" case millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_STORED, millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_DISCARDED: capabilityAction = "save" default: return batchApply{}, protoErrf("unsupported cache update action %v", update.GetAction()) } ref, authorized, err := tx.ConsumeMillCacheCapability( context.Background(), lease.id, capabilityAction, update.GetId(), ) if err != nil { return batchApply{}, err } if !authorized { return batchApply{}, protoErrf( "cache update %q was not planned for lease %q", update.GetId(), lease.id, ) } switch update.GetAction() { case millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_USED: if err := tx.TouchCacheEntry(context.Background(), update.GetId(), at); err != nil { return batchApply{}, err } case millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_STORED: actualSize, ok := precomputedSizes[ref] if !ok { var err error actualSize, err = m.cacheObjectSize(context.Background(), ref) if err != nil { return batchApply{}, fmt.Errorf("stat stored cache object %q: %w", ref, err) } } superseded, ready, err := tx.MarkCacheEntryReadyWithChecksum( context.Background(), update.GetId(), actualSize, update.GetChecksum(), m.cfg.CacheMaxBytesPerOwner, at, ) if err != nil { return batchApply{}, err } for _, old := range superseded { r.cacheDeletes[old.ID] = old } if !ready { if err := tx.QueueCacheObjectDeletion(context.Background(), ref, at); err != nil { return batchApply{}, err } r.cacheObjectDeletes[ref] = struct{}{} } case millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_DISCARDED, millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_MISSING: // A missing restore object must not let one executor delete a // ready generation that another executor published. Both // actions can only discard entries still owned by this save. pendingOnly := true discarded, err := tx.DiscardCacheEntry( context.Background(), update.GetId(), pendingOnly, ) if err != nil { return batchApply{}, err } if discarded != nil { r.cacheDeletes[discarded.ID] = *discarded } } } case entry.GetAttemptResult() != nil: ar := entry.GetAttemptResult() statusStr := "success" switch ar.Status { case millv1.TerminalStatus_TERMINAL_STATUS_SUCCESS: statusStr = "success" case millv1.TerminalStatus_TERMINAL_STATUS_FAILED: statusStr = "failed" case millv1.TerminalStatus_TERMINAL_STATUS_TIMEOUT: statusStr = "timeout" case millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED: statusStr = "cancelled" default: return batchApply{}, protoErrf("unsupported terminal status %v", ar.Status) } var errMsg *string if e := ar.GetError(); e != "" { errMsg = &e } var exitCode *int64 if c := ar.GetExitCode(); c != 0 { exitCode = &c } finishedInBatch[lease.id] = struct{}{} pipelineAtUri := string(lease.wid.PipelineId.AtUri()) if tx != nil { if err := tx.InsertStatusEvent(pipelineAtUri, lease.wid.Name, statusStr, errMsg, exitCode); err != nil { return batchApply{}, err } if err := tx.DeleteLease(lease.id); err != nil { return batchApply{}, err } if a := ar.GetLogArtifact(); a != nil { if a.GetRef() == "" { return batchApply{}, protoErrf("empty log artifact ref") } // the ref names an object the mill deletes on wipe; // pin it to the derived log key so an executor cannot // point cleanup at someone else's objects if a.GetRef() != "logs/"+lease.id+".log" { return batchApply{}, protoErrf("log artifact ref %q does not match lease %s", a.GetRef(), lease.id) } if !strings.HasPrefix(a.GetHash(), "sha256:") { return batchApply{}, protoErrf("invalid log artifact hash %q", a.GetHash()) } if lease.wid.Knot == "" || lease.wid.Rkey == "" || lease.wid.Name == "" { m.l.Error("skipping log artifact with incomplete workflow identity", "lease", lease.id, "wid", lease.wid) } else { if tx != nil { if err := tx.InsertArtifactRef(lease.id, lease.repoDID, lease.wid, a.GetRef(), a.GetHash()); err != nil { return batchApply{}, err } } r.artifactLeases = append(r.artifactLeases, lease) } } } r.pendingTerminals = append(r.pendingTerminals, pendingTerminal{ lease: lease, ar: ar, }) } r.highestSeqno = entry.Seqno } if tx != nil { if err := tx.AdvanceCursor(sess.nodeID, sess.epoch, r.highestSeqno); err != nil { return batchApply{}, err } } return r, nil } var applied batchApply var err error if m.db != nil { applied, err = db.ApplyEventBatchCtx(sess.ctx, m.db, m.n, applyFunc) } else { applied, err = applyFunc(nil) } if err != nil { unlockLeases() return err } for _, observation := range applied.startupObservations { m.metrics.RecordWorkflowStartupDelay(context.Background(), observation.engine, "mill", observation.delay) } m.enqueueCacheDeletions(applied.cacheDeletes, applied.cacheObjectDeletes) if m.cfg.LogDir != "" { for _, lease := range applied.artifactLeases { path := models.LogFilePath(m.cfg.LogDir, lease.wid) if err := os.Remove(path); err != nil && !errors.Is(err, os.ErrNotExist) { m.l.Warn("failed to remove live log file after artifact recorded", "path", path, "err", err) } } } m.mu.Lock() if applied.highestSeqno > m.nodeSeqno[currentKey] { m.nodeSeqno[currentKey] = applied.highestSeqno } if sess.recovering && !sess.recoveryExpired { m.extendGraceLocked(sess) } m.maybeFinishSessionRecoveryLocked(sess) m.mu.Unlock() for _, leaseID := range applied.terminalCancelIDs { m.settleSessionCancel(sess, leaseID) } for _, pt := range applied.pendingTerminals { if pt.lease.millRecordsTerminalMetrics && pt.ar.GetMillRecordsTerminalMetrics() && pt.ar.GetFailureClass() != "" && pt.ar.GetFailureReason() != "" { m.metrics.RecordWorkflowTerminal( pt.lease.engine, terminalMetricResult(pt.ar.GetStatus()), pt.ar.GetFailureClass(), pt.ar.GetFailureReason(), ) } // orphans have no waiting RunStep, just mark and clean up if pt.lease.orphaned { pt.lease.markDone() _ = m.cleanupLeaseLocked(pt.lease) continue } pt.lease.deliverTerminal(pt.ar) if pt.lease.cleanupReady() { _ = m.cleanupLeaseLocked(pt.lease) } } unlockLeases() return m.sendAck(sess, applied.highestSeqno) } func terminalMetricResult(status millv1.TerminalStatus) string { switch status { case millv1.TerminalStatus_TERMINAL_STATUS_SUCCESS: return "success" case millv1.TerminalStatus_TERMINAL_STATUS_TIMEOUT: return "timeout" case millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED: return "cancelled" default: return "failure" } } func (m *Mill) cacheObjectSize(ctx context.Context, key string) (int64, error) { m.mu.Lock() store := m.cache m.mu.Unlock() if store == nil { return 0, errors.New("cache store is not configured") } if statter, ok := store.(storage.StatStorage); ok { return statter.Stat(ctx, key) } r, err := store.Get(ctx, key) if err != nil { return 0, err } defer r.Close() n, err := io.Copy(io.Discard, r) if err != nil { return 0, err } return n, nil } func (m *Mill) startCacheDeletionWorker() { m.cacheDeleteOnce.Do(func() { m.cacheDeleteWake = make(chan struct{}, 1) go m.cacheDeletionLoop() }) } func (m *Mill) enqueueCacheDeletions(entries map[string]db.CacheEntry, objects map[string]struct{}) { if len(entries) == 0 && len(objects) == 0 { return } m.startCacheDeletionWorker() deletion := cacheDeletion{ entries: make([]db.CacheEntry, 0, len(entries)), objects: make([]string, 0, len(objects)), } for _, entry := range entries { deletion.entries = append(deletion.entries, entry) } for key := range objects { deletion.objects = append(deletion.objects, key) } m.cacheDeleteMu.Lock() m.cacheDeleteQueue = append(m.cacheDeleteQueue, deletion) m.cacheDeleteMu.Unlock() select { case m.cacheDeleteWake <- struct{}{}: default: } } func (m *Mill) cacheDeletionLoop() { for range m.cacheDeleteWake { for { m.cacheDeleteMu.Lock() if len(m.cacheDeleteQueue) == 0 { m.cacheDeleteMu.Unlock() break } deletion := m.cacheDeleteQueue[0] m.cacheDeleteQueue[0] = cacheDeletion{} m.cacheDeleteQueue = m.cacheDeleteQueue[1:] m.cacheDeleteMu.Unlock() m.deleteCacheEntries(deletion.entries) m.deleteQueuedCacheObjects(deletion.objects) } } } func (m *Mill) cacheDeleteContext() (context.Context, context.CancelFunc) { timeout := m.cfg.CacheDeleteTimeout if timeout <= 0 { timeout = defaultCacheDeleteTimeout } return context.WithTimeout(context.Background(), timeout) } func (m *Mill) deleteCacheEntries(entries []db.CacheEntry) { if len(entries) == 0 { return } m.mu.Lock() cache := m.cache database := m.db m.mu.Unlock() if cache == nil { m.l.Warn("cache objects need deletion but the mill has no cache store", "count", len(entries)) return } for _, entry := range entries { ctx, cancel := m.cacheDeleteContext() err := cache.Delete(ctx, entry.StorageKey) cancel() if err != nil { m.l.Warn("delete cache object", "id", entry.ID, "ref", entry.StorageKey, "err", err) continue } if database != nil { ctx, cancel = m.cacheDeleteContext() err = database.DeleteCacheEntry(ctx, entry.ID) cancel() if err != nil { m.l.Warn("delete cache metadata", "id", entry.ID, "err", err) } } } } func (m *Mill) deleteQueuedCacheObjects(keys []string) { if len(keys) == 0 { return } m.mu.Lock() cache := m.cache database := m.db m.mu.Unlock() if cache == nil { m.l.Warn("cache objects need deletion but the mill has no cache store", "count", len(keys)) return } for _, key := range keys { ctx, cancel := m.cacheDeleteContext() err := cache.Delete(ctx, key) cancel() if err != nil { m.l.Warn("delete unindexed cache object", "ref", key, "err", err) continue } if database != nil { ctx, cancel = m.cacheDeleteContext() err = database.CompleteCacheObjectDeletion(ctx, key) cancel() if err != nil { m.l.Warn("complete cache object deletion", "ref", key, "err", err) } } } } func (m *Mill) sendAck(sess *millSession, seqno uint64) error { msg := &millproto.Message{Ack: &millv1.Ack{ Epoch: sess.epoch, UpToSeqno: seqno, }} if err := sess.send(msg); err != nil { return fmt.Errorf("send ack message: %w", err) } return nil } func (m *Mill) onLiveLog(sess *millSession, ll *millv1.LiveLog) error { if ll == nil || ll.GetLeaseId() == "" { return nil } m.mu.Lock() lease := m.leases[ll.GetLeaseId()] m.mu.Unlock() if lease == nil || lease.nodeID != sess.nodeID { return nil } raw := ll.GetRawJson() if m.cfg.LogDir == "" || len(raw) == 0 { if m.n != nil { m.n.NotifyAll() } return nil } lease.mu.Lock() isDone := (lease.state == leaseDone) lease.mu.Unlock() if isDone { if m.n != nil { m.n.NotifyAll() } return nil } logPath := models.LogFilePath(m.cfg.LogDir, lease.wid) if err := os.MkdirAll(filepath.Dir(logPath), 0755); err != nil { m.l.Warn("failed to create log dir", "path", filepath.Dir(logPath), "err", err) if m.n != nil { m.n.NotifyAll() } return nil } f, err := os.OpenFile(logPath, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0600) if err != nil { m.l.Warn("failed to open log file", "path", logPath, "err", err) if m.n != nil { m.n.NotifyAll() } return nil } if _, err := f.Write(raw); err != nil { m.l.Warn("failed to write log file", "path", logPath, "err", err) } _ = f.Close() if m.n != nil { m.n.NotifyAll() } return nil } func (m *Mill) settleSessionCancel(sess *millSession, leaseID string) { m.mu.Lock() _, pending := sess.pendingCancels[leaseID] delete(sess.pendingCancels, leaseID) m.maybeFinishSessionRecoveryLocked(sess) m.mu.Unlock() if pending { m.notifyChange() } } func (m *Mill) onCancelAck(sess *millSession, ca *millv1.CancelAck) { leaseID := ca.GetLeaseId() m.settleSessionCancel(sess, leaseID) m.mu.Lock() lease := m.leases[leaseID] m.mu.Unlock() if lease == nil || lease.nodeID != sess.nodeID { return } if !lease.markCancelAcked() { return } time.AfterFunc(m.cfg.CancelTeardownTimeout, func() { m.checkCancelTeardown(lease) }) } func (m *Mill) checkCancelAck(sess *millSession, lease *RemoteLease) { lease.mu.Lock() isDone := (lease.state == leaseDone) isAcked := lease.cancelAcked lease.mu.Unlock() if isDone || isAcked { m.settleSessionCancel(sess, lease.id) return } m.mu.Lock() isCurrent := m.sessions[lease.nodeID] == sess && !sess.disconnected m.mu.Unlock() if !isCurrent { return } m.l.Warn("node did not acknowledge cancel request within deadline; closing session", "node", lease.nodeID, "lease", lease.id, "deadline", m.cfg.CancelAckTimeout) sess.close() } // don't fail the node's other jobs because one teardown got stuck func (m *Mill) checkCancelTeardown(lease *RemoteLease) { if lease.getState() == leaseDone { return } _, cancelReason := lease.cancelRequested() reason := fmt.Sprintf("%s (node teardown exceeded %s)", cancelReason, m.cfg.CancelTeardownTimeout) m.l.Warn("acked cancel exceeded teardown deadline; finishing lease without the node's terminal", "node", lease.nodeID, "lease", lease.id, "deadline", m.cfg.CancelTeardownTimeout) status := string(models.StatusKindCancelled) if lease.orphaned { if err := m.finishOrphan(lease, status, &reason, nil); err != nil { m.l.Error("finish orphan after cancel teardown deadline failed", "lease", lease.id, "err", err) } return } if err := m.finishLiveLease(lease, status, reason); err != nil { m.l.Error("finish lease after cancel teardown deadline failed", "lease", lease.id, "err", err) } } func marshalJob(wf *models.Workflow) (pipeline string, workflow string, err error) { st, ok := wf.Data.(*millWorkflowState) if !ok || st == nil { return "", "", fmt.Errorf("mill workflow state missing") } p, err := json.Marshal(st.RawPipeline) if err != nil { return "", "", fmt.Errorf("marshal pipeline: %w", err) } w, err := json.Marshal(st.RawWorkflow) if err != nil { return "", "", fmt.Errorf("marshal workflow: %w", err) } return string(p), string(w), nil } func (m *Mill) RegisterMetrics(metrics *observability.Metrics) { if metrics == nil { return } m.metrics = metrics metrics.RegisterMillGauges( func() float64 { m.mu.Lock() defer m.mu.Unlock() return float64(m.pending) }, func() float64 { return float64(m.cfg.MaxPending) }, func() float64 { m.mu.Lock() defer m.mu.Unlock() return float64(len(m.leases)) }, func() float64 { m.mu.Lock() defer m.mu.Unlock() return float64(len(m.reservations)) }, func() float64 { m.mu.Lock() defer m.mu.Unlock() active := 0 for _, sess := range m.sessions { if !sess.disconnected { active++ } } return float64(active) }, func() float64 { m.mu.Lock() defer m.mu.Unlock() disconnected := 0 for _, sess := range m.sessions { if sess.disconnected { disconnected++ } } return float64(disconnected) }, ) } func (m *Mill) authorizeAndLockQuota(s *millSession, reservationID, requestedLeaseID string) (*RemoteLease, error) { // snapshot the current owner before locking its teardown path m.quotaLifecycleMu.Lock() lease, ok := m.quotaLeases[reservationID] m.quotaLifecycleMu.Unlock() if !ok || lease == nil { return nil, fmt.Errorf("quota reservation %q is not attached to a live lease", reservationID) } // keep teardown from changing the lease during authorization lease.finishMu.Lock() // serialize the final ownership check with quota phase transitions m.quotaLifecycleMu.Lock() if m.quotaLeases[reservationID] != lease { m.quotaLifecycleMu.Unlock() lease.finishMu.Unlock() return nil, fmt.Errorf("quota reservation %q is not attached to a live lease", reservationID) } if lease.id != requestedLeaseID { m.quotaLifecycleMu.Unlock() lease.finishMu.Unlock() return nil, fmt.Errorf("quota reservation %q belongs to another workflow", reservationID) } if lease.nodeID != s.nodeID { m.quotaLifecycleMu.Unlock() lease.finishMu.Unlock() return nil, fmt.Errorf("quota reservation %q belongs to another executor", reservationID) } if lease.epoch != s.epoch { m.quotaLifecycleMu.Unlock() lease.finishMu.Unlock() return nil, fmt.Errorf("quota reservation %q has stale epoch %q (current: %q)", reservationID, lease.epoch, s.epoch) } lease.mu.Lock() rq, hasQuota := lease.remoteQuotas[reservationID] active := hasQuota && (lease.state != leaseDone || rq.committing) lease.mu.Unlock() if !active { m.quotaLifecycleMu.Unlock() lease.finishMu.Unlock() return nil, fmt.Errorf("quota reservation %q is not active", reservationID) } // the caller releases both locks after its phase transition return lease, nil } func (m *Mill) handleQuotaRequest(s *millSession, req *millv1.QuotaRequest) { var ( res quota.Reservation err error ) switch req.GetOperation() { case millv1.QuotaOperation_QUOTA_OPERATION_RESERVE: res, err = m.reserveRemoteQuota(s, req) case millv1.QuotaOperation_QUOTA_OPERATION_BEGIN_COMMIT, millv1.QuotaOperation_QUOTA_OPERATION_COMMIT, millv1.QuotaOperation_QUOTA_OPERATION_RELEASE: err = m.transitionRemoteQuota(s, req) default: err = fmt.Errorf("unsupported quota operation %q", req.GetOperation()) } resp := &millv1.QuotaResponse{ RequestId: req.GetRequestId(), ReservationId: res.ID, Allowed: res.Allowed, Temporary: res.Temporary, Reason: res.Reason, Resource: string(res.Resource), } if err != nil { resp.Error = err.Error() } _ = s.send(&millproto.Message{QuotaResp: resp}) } func (m *Mill) reserveRemoteQuota(s *millSession, req *millv1.QuotaRequest) (quota.Reservation, error) { if m.qm == nil { return quota.Reservation{}, errors.New("quota manager not attached") } if req.GetReservationId() != "" { return quota.Reservation{}, errors.New("reserve operation cannot name a reservation") } kind, resources, err := remoteQuotaPolicy(req) if err != nil { return quota.Reservation{}, err } m.mu.Lock() lease := m.leases[req.GetLeaseId()] m.mu.Unlock() if lease == nil || lease.nodeID != s.nodeID || lease.epoch != s.epoch { return quota.Reservation{}, fmt.Errorf("no live lease %q on this node with matching epoch", req.GetLeaseId()) } lease.mu.Lock() id := quota.Identity{OwnerDID: lease.ownerDID, RepoDID: lease.repoDID} lease.mu.Unlock() if id.OwnerDID == "" || id.RepoDID == "" { return quota.Reservation{}, errors.New("lease has no charged subject") } lease.finishMu.Lock() defer lease.finishMu.Unlock() if lease.cleanedUp || lease.getState() == leaseDone { return quota.Reservation{}, fmt.Errorf("lease %q ended before quota reservation", lease.id) } m.quotaLifecycleMu.Lock() defer m.quotaLifecycleMu.Unlock() reservationKey := string(kind) + "\x00" + req.GetKey() lease.mu.Lock() var existingID string mismatch := false for resID, rq := range lease.remoteQuotas { if rq.key == reservationKey { existingID = resID mismatch = !maps.Equal(rq.resources, resources) break } } lease.mu.Unlock() if existingID != "" { // a repeat is only idempotent for the exact same charge; anything // else is an executor trying to inflate an existing reservation if mismatch { return quota.Reservation{}, fmt.Errorf("cache key %q already reserved with different resources", req.GetKey()) } return quota.Reservation{ ID: existingID, Allowed: true, Reason: quota.ReasonWithinLimit, }, nil } qlease, res, err := m.qm.TryAcquire(context.Background(), quota.ReserveRequest{ Kind: kind, Key: req.GetKey(), Identity: id, Resources: resources, }) if err != nil { return quota.Reservation{}, err } if qlease == nil || res.ID == "" { return res, nil } existingLease := m.quotaLeases[res.ID] if existingLease != nil && existingLease != lease { return quota.Reservation{}, fmt.Errorf("quota reservation %q is already owned by another workflow", res.ID) } m.quotaLeases[res.ID] = lease lease.mu.Lock() rq := lease.remoteQuotas[res.ID] rq.key = reservationKey rq.resources = resources lease.remoteQuotas[res.ID] = rq lease.mu.Unlock() return res, nil } func remoteQuotaPolicy(req *millv1.QuotaRequest) (quota.Kind, quota.Resources, error) { kind := quota.Kind(req.GetKind()) if kind != quota.KindNixCache && kind != quota.KindGenericCache { return "", nil, fmt.Errorf("unsupported remote quota kind %q", req.GetKind()) } resources := req.GetResources() if len(resources) == 0 { return "", nil, errors.New("empty remote quota resources") } if err := quota.ValidateResources(resources); err != nil { return "", nil, fmt.Errorf("invalid remote quota resources: %w", err) } cacheBytes, ok := resources[quota.ResourceCacheStorageBytes] if !ok || cacheBytes <= 0 || len(resources) != 1 { return "", nil, errors.New("remote cache quota requires one positive cache_storage_bytes resource") } return kind, maps.Clone(resources), nil } func (m *Mill) transitionRemoteQuota(s *millSession, req *millv1.QuotaRequest) error { reservationID := req.GetReservationId() if reservationID == "" { return errors.New("quota transition requires a reservation") } if req.GetKind() != "" || req.GetKey() != "" || req.GetResources() != nil { return errors.New("quota transition cannot include reservation parameters") } // successful terminal transitions are idempotent if req.GetOperation() != millv1.QuotaOperation_QUOTA_OPERATION_BEGIN_COMMIT { m.quotaLifecycleMu.Lock() _, ok := m.quotaLeases[reservationID] m.quotaLifecycleMu.Unlock() if !ok { return nil } } lease, err := m.authorizeAndLockQuota(s, reservationID, req.GetLeaseId()) if err != nil { return err } defer lease.finishMu.Unlock() defer m.quotaLifecycleMu.Unlock() if m.qm == nil { return errors.New("quota manager not attached") } switch req.GetOperation() { case millv1.QuotaOperation_QUOTA_OPERATION_BEGIN_COMMIT: if err := m.qm.BeginCommit(context.Background(), reservationID); err != nil { return err } lease.mu.Lock() rq := lease.remoteQuotas[reservationID] rq.committing = true lease.remoteQuotas[reservationID] = rq lease.mu.Unlock() case millv1.QuotaOperation_QUOTA_OPERATION_COMMIT: if err := m.qm.Commit(context.Background(), reservationID); err != nil { return err } m.forgetQuotaReservationLocked(reservationID) case millv1.QuotaOperation_QUOTA_OPERATION_RELEASE: if err := m.qm.Release(context.Background(), reservationID); err != nil { return err } m.forgetQuotaReservationLocked(reservationID) } return nil } func (m *Mill) forgetQuotaReservationLocked(reservationID string) { lease := m.quotaLeases[reservationID] delete(m.quotaLeases, reservationID) if lease == nil { return } lease.mu.Lock() delete(lease.remoteQuotas, reservationID) lease.mu.Unlock() }