Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377package 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 stickfunc (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 oncefunc (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 stuckfunc (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()}