package executor import ( "bytes" "context" "encoding/json" "errors" "fmt" "io" "log/slog" "maps" "net/http" "runtime" "strings" "sync" "time" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/gorilla/websocket" "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" "go.opentelemetry.io/otel/trace" "tangled.org/core/api/tangled" "tangled.org/core/netutil" "tangled.org/core/notifier" "tangled.org/core/spindle/artifactstore" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" millproto "tangled.org/core/spindle/mill/proto" millv1 "tangled.org/core/spindle/mill/proto/gen" "tangled.org/core/spindle/models" "tangled.org/core/spindle/observability" "tangled.org/core/spindle/quota" "tangled.org/core/spindle/storage" ) const ( dialBackoffMin = 1 * time.Second dialBackoffMax = 30 * time.Second defaultDrainOutboxGrace = 45 * time.Second snapshotEvery = 15 * time.Second defaultSeats = 4 ) type Executor struct { millURL string token string nodeID string seats int labels []string engines map[string]models.Engine db *db.DB n *notifier.Notifier cfg *config.Config l *slog.Logger writer artifactstore.Writer cache storage.Storage epoch string outboxBytes int64 maxOutboxBytes int64 // 10 MiB default outbox cap outboxIdleCh chan struct{} terminalSeqnos map[string]uint64 eventMu sync.Mutex sendMu sync.Mutex flushMu sync.Mutex sentSeqno uint64 connMu sync.Mutex enc messageEncoder sessionCancel context.CancelFunc mu sync.Mutex active map[string]*reservation cleanupPending int draining bool idleCh chan struct{} snapshotMu sync.Mutex nextSeqno uint64 quotaClient *QuotaClient lifecycleCtx context.Context jobsWG sync.WaitGroup } type reservation struct { leaseID string wid models.WorkflowId realEngine models.Engine slot engine.WorkflowSlot wf *models.Workflow pipeline *models.Pipeline vault *memVault committed bool cancelled bool cancel context.CancelCauseFunc ttlTimer *time.Timer stopTail func() runDone chan struct{} traceParent trace.SpanContext failureClass string failureReason string millRecordsTerminalMetrics bool } type messageEncoder interface { Encode(*millproto.Message) error } func New(cfg *config.Config, engines map[string]models.Engine, d *db.DB, n *notifier.Notifier, l *slog.Logger, writer artifactstore.Writer, cache storage.Storage) (*Executor, error) { seats := defaultSeats millURL := "" token := "" nodeID := "" var labels []string if cfg != nil { if cfg.Mill.Seats > 0 { seats = cfg.Mill.Seats } labels = normalizeLabels(cfg.Mill.Labels) millURL = cfg.Mill.URL token = cfg.Mill.SharedSecret nodeID = cfg.Server.Hostname } if d == nil || n == nil { return nil, fmt.Errorf("executor requires a database and notifier") } e := &Executor{ millURL: millURL, token: token, nodeID: nodeID, seats: seats, labels: labels, engines: engines, db: d, n: n, cfg: cfg, l: l.With("component", "mill.executor"), writer: writer, cache: cache, active: make(map[string]*reservation), maxOutboxBytes: 10 * 1024 * 1024, } if err := e.initOutbox(); err != nil { return nil, fmt.Errorf("initialize executor outbox: %w", err) } return e, nil } func (e *Executor) Connect(ctx context.Context) { e.lifecycleCtx = ctx recoveryDone := make(chan struct{}) go func() { defer close(recoveryDone) if err := e.recoverPendingArtifacts(ctx); err != nil && ctx.Err() == nil { e.l.Warn("recover pending artifacts failed", "err", err) } }() defer func() { <-recoveryDone }() sub := e.n.Subscribe() cursor, err := e.db.EventHighWater() if err != nil { e.n.Unsubscribe(sub) e.l.Error("establish event cursor failed", "err", err) return } e.drainEvents(&cursor) observerCtx, stopObserver := context.WithCancel(ctx) observerDone := make(chan struct{}) go func() { defer close(observerDone) e.observeLoop(observerCtx, sub, cursor) }() defer e.n.Unsubscribe(sub) backoff := dialBackoffMin for { if ctx.Err() != nil { break } started := time.Now() err := e.runSession(ctx) if ctx.Err() != nil { break } if time.Since(started) >= dialBackoffMax { backoff = dialBackoffMin } e.l.Warn("mill session ended; reconnecting", "err", err, "backoff", backoff) select { case <-ctx.Done(): break case <-time.After(backoff): } backoff = min(backoff*2, dialBackoffMax) } e.jobsWG.Wait() stopObserver() <-observerDone e.drainEvents(&cursor) } func (e *Executor) runSession(ctx context.Context) error { dev := e.cfg == nil || e.cfg.Server.Dev if _, err := netutil.EnforceWSSURL(e.millURL, dev); err != nil { return fmt.Errorf("mill url: %w", err) } header := http.Header{} if e.token != "" { header.Set("Authorization", "Bearer "+e.token) } conn, _, err := websocket.DefaultDialer.DialContext(ctx, e.millURL, header) if err != nil { return fmt.Errorf("dial mill: %w", err) } defer conn.Close() sessionCtx, cancelSession := context.WithCancel(ctx) defer cancelSession() stopClose := context.AfterFunc(sessionCtx, func() { _ = conn.Close() }) defer stopClose() stream := millproto.NewWSStream(conn) enc := millproto.NewEncoder(stream) dec := millproto.NewDecoder(stream) cacheStoreID := "" if e.cfg != nil && e.cfg.Cache.StoreID != "" && e.cacheStoreAvailable(ctx) { cacheStoreID = e.cfg.Cache.StoreID } hello := &millproto.Message{Hello: &millv1.Hello{ ProtocolVersion: millproto.ProtocolVersion, MinProtocolVersion: millproto.ProtocolMinVersion, MaxProtocolVersion: millproto.ProtocolMaxVersion, Arch: runtime.GOARCH, Labels: e.labels, Epoch: e.epoch, CacheStoreId: cacheStoreID, CacheNamespace: models.CacheNamespace(), }} if err := enc.Encode(hello); err != nil { return fmt.Errorf("send hello: %w", err) } resumeMsg, err := dec.Decode() if err != nil { return fmt.Errorf("read resume: %w", err) } resume := resumeMsg.GetResume() if resume == nil { return fmt.Errorf("expected resume, got something else") } if resume.GetEpoch() != e.epoch { return fmt.Errorf("resume epoch mismatch: got %q, want %q", resume.GetEpoch(), e.epoch) } e.connMu.Lock() e.sessionCancel = cancelSession e.enc = enc e.connMu.Unlock() if e.quotaClient != nil { e.quotaClient.SetSendFn(e.sendQuota) } defer func() { e.connMu.Lock() e.enc = nil e.sessionCancel = nil e.connMu.Unlock() if e.quotaClient != nil { e.quotaClient.OnDisconnect() } }() readErr := make(chan error, 1) go func() { for { msg, err := dec.Decode() if err != nil { readErr <- fmt.Errorf("read: %w", err) return } e.dispatch(sessionCtx, msg) } }() if err := e.replay(resume.GetAckSeqno()); err != nil { cancelSession() <-readErr return fmt.Errorf("replay failed: %w", err) } e.pushSnapshot() e.l.Info("connected to mill", "node", e.nodeID, "resumeFrom", resume.GetAckSeqno()) go e.snapshotLoop(sessionCtx, enc) return <-readErr } func (e *Executor) send(msg *millproto.Message) { e.connMu.Lock() enc := e.enc cancel := e.sessionCancel e.connMu.Unlock() if enc != nil { e.sendMu.Lock() err := enc.Encode(msg) e.sendMu.Unlock() if err != nil { e.l.Error("send failed, ending session", "err", err) if cancel != nil { cancel() } } } } func (e *Executor) cacheStoreAvailable(ctx context.Context) bool { if e.cache == nil { return false } rc, err := e.cache.Get(ctx, millproto.CacheSentinelKey) if err != nil { e.l.Warn("cache store is not shared with mill; disabling cache placement", "key", millproto.CacheSentinelKey, "err", err) return false } defer rc.Close() value, err := io.ReadAll(io.LimitReader(rc, int64(len(millproto.CacheSentinelValue)+1))) if err != nil || !bytes.Equal(value, []byte(millproto.CacheSentinelValue)) { e.l.Warn("cache store sentinel is invalid; disabling cache placement", "key", millproto.CacheSentinelKey) return false } return true } func (e *Executor) dispatch(ctx context.Context, msg *millproto.Message) { switch { case msg.GetReserveSeat() != nil: e.handleReserve(ctx, msg.GetReserveSeat()) case msg.GetCommitLease() != nil: e.handleCommit(ctx, msg.GetCommitLease()) case msg.GetReleaseLease() != nil: e.handleRelease(msg.GetReleaseLease().GetLeaseId()) case msg.GetCancelAttempt() != nil: e.handleCancel(msg.GetCancelAttempt().GetLeaseId()) case msg.GetAck() != nil: e.handleAck(msg.GetAck()) case msg.GetQuotaResp() != nil: if e.quotaClient != nil { e.quotaClient.HandleResponse(msg.GetQuotaResp().GetRequestId(), msg.GetQuotaResp()) } default: e.l.Warn("unhandled incoming message", "type", fmt.Sprintf("%T", msg)) } } func (e *Executor) SetQuotaClient(c *QuotaClient) { e.mu.Lock() defer e.mu.Unlock() e.quotaClient = c } func (e *Executor) sendQuota(msg *millproto.Message) error { e.connMu.Lock() enc := e.enc e.connMu.Unlock() if enc == nil { return errors.New("not connected to mill") } e.send(msg) return nil } func (e *Executor) sendReject(leaseID string, reason string, class millv1.RejectClass) { e.sendRejectWithAttribution(leaseID, reason, class, "", "") } func (e *Executor) sendRejectWithAttribution( leaseID, reason string, class millv1.RejectClass, failureClass, failureReason string, ) { e.send(&millproto.Message{ReserveResult: &millv1.ReserveResult{ LeaseId: leaseID, Accepted: false, RejectReason: reason, RejectClass: class, FailureClass: failureClass, FailureReason: failureReason, }}) } func (e *Executor) sendCommitted(leaseID string) { e.send(&millproto.Message{Committed: &millv1.Committed{LeaseId: leaseID}}) } func (e *Executor) sendCancelAck(leaseID string) { e.send(&millproto.Message{CancelAck: &millv1.CancelAck{LeaseId: leaseID}}) } func (e *Executor) releaseReservation(cleanup func()) { if cleanup != nil { cleanup() } e.pushSnapshot() } func (e *Executor) handleReserve(ctx context.Context, rs *millv1.ReserveSeat) { parentCtx := e.lifecycleCtx if parentCtx == nil { parentCtx = ctx } if parentCtx == nil { parentCtx = context.Background() } reserveCtx := observability.ExtractFromTraceparentAndTracestate(parentCtx, rs.GetTraceparent(), rs.GetTracestate()) reserveCtx, span := observability.Tracer().Start(reserveCtx, "executor.assignment") if span.IsRecording() { span.SetAttributes( attribute.String(observability.LeaseIDKey, rs.GetLeaseId()), attribute.String(observability.ExecutorNodeIDKey, e.nodeID), attribute.String(observability.PipelineIDKey, (&models.PipelineId{Knot: rs.GetKnot(), Rkey: rs.GetRkey()}).AtUri().String()), ) } accepted := false defer func() { if accepted { span.SetStatus(codes.Ok, "accepted") } else { span.SetStatus(codes.Error, "assignment rejected") } span.End() }() reject := func(reason string, class millv1.RejectClass) { e.sendReject(rs.GetLeaseId(), reason, class) } rejectError := func(prefix string, err error, class millv1.RejectClass) { failureClass, failureReason := engine.FailureAttribution("failure", err) e.sendRejectWithAttribution( rs.GetLeaseId(), prefix+err.Error(), class, failureClass, failureReason, ) } e.mu.Lock() draining := e.draining atCapacity := len(e.active) >= e.seatCapacity() e.mu.Unlock() if draining { reject("draining", millv1.RejectClass_REJECT_CLASS_TRANSIENT) return } if atCapacity { reject("no executor seats available", millv1.RejectClass_REJECT_CLASS_TRANSIENT) return } realEngine, ok := e.engines[rs.GetTargetEngine()] if !ok { reject("unknown engine "+rs.GetTargetEngine(), millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) return } slotter, ok := realEngine.(engine.WorkflowSlotter) if !ok { reject("engine does not support workflow slots", millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) return } var twf tangled.Pipeline_Workflow if err := json.Unmarshal([]byte(rs.GetRawWorkflowJson()), &twf); err != nil { reject("bad workflow json", millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) return } pipelineId := models.PipelineId{Knot: rs.GetKnot(), Rkey: rs.GetRkey()} wid := models.WorkflowId{PipelineId: pipelineId, Name: twf.Name} if span.IsRecording() { span.SetAttributes(attribute.String(observability.WorkflowIDKey, wid.String())) } var tpl tangled.Pipeline if err := json.Unmarshal([]byte(rs.GetRawPipelineJson()), &tpl); err != nil { reject("bad pipeline json", millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) return } if tpl.TriggerMetadata == nil { reject("pipeline missing trigger metadata", millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) return } repoDid, err := syntax.ParseDID(rs.GetRepoDid()) if err != nil { reject("bad repository did", millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) return } trustedSource := models.TrustedPipelineSource(tpl.TriggerMetadata, repoDid.String()) wf, err := realEngine.InitWorkflow(twf, tpl) if err != nil { rejectError("init workflow: ", err, millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) return } if e.quotaClient != nil { if binder, ok := realEngine.(engine.WorkflowQuotaStoreBinder); ok { if err := binder.BindWorkflowQuotaStore(wf, e.quotaClient.ForLease(rs.GetLeaseId())); err != nil { rejectError("bind workflow quota store: ", err, millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) return } } } wf.Engine = rs.GetTargetEngine() wf.RepoDID = rs.GetRepoDid() if validator, ok := realEngine.(engine.WorkflowPlacementValidator); ok { if err := validator.ValidateWorkflowPlacement(wf); err != nil { rejectError("validate workflow placement: ", err, millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) return } } if wf.Environment == nil { wf.Environment = make(map[string]string) } maps.Copy(wf.Environment, models.PipelineEnvVars(tpl.TriggerMetadata, pipelineId)) slot, err := slotter.AcquireWorkflowSlot(reserveCtx, wid, wf, engine.NoWait) if err != nil { class := millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE if errors.Is(err, engine.ErrNoWorkflowSlots) { class = millv1.RejectClass_REJECT_CLASS_TRANSIENT } reject(err.Error(), class) return } if span.IsRecording() && repoDid.String() != "" { span.SetAttributes(attribute.String(observability.RepoDIDKey, repoDid.String())) } res := &reservation{ leaseID: rs.GetLeaseId(), wid: wid, realEngine: realEngine, slot: slot, wf: wf, traceParent: trace.SpanContextFromContext(reserveCtx), pipeline: &models.Pipeline{ RepoDid: repoDid, TrustedSource: trustedSource, TriggerMetadata: tpl.TriggerMetadata, }, } e.snapshotMu.Lock() e.mu.Lock() if e.draining { e.mu.Unlock() e.snapshotMu.Unlock() slot.Release() reject("draining", millv1.RejectClass_REJECT_CLASS_TRANSIENT) return } if len(e.active) >= e.seatCapacity() { e.mu.Unlock() e.snapshotMu.Unlock() slot.Release() reject("no executor seats available", millv1.RejectClass_REJECT_CLASS_TRANSIENT) return } if len(e.active) == 0 && e.cleanupPending == 0 { e.idleCh = make(chan struct{}) } e.active[res.leaseID] = res res.ttlTimer = time.AfterFunc(ttlDuration(rs.GetTtlSeconds()), func() { e.expireReservation(res.leaseID) }) e.mu.Unlock() e.send(&millproto.Message{ReserveResult: &millv1.ReserveResult{ LeaseId: rs.GetLeaseId(), Accepted: true, QuotaResources: reportedResources(realEngine, wf), SupportsMillTerminalMetrics: true, }}) accepted = true e.pushSnapshotLocked() e.snapshotMu.Unlock() } func reportedResources(eng models.Engine, wf *models.Workflow) map[string]int64 { reporter, ok := eng.(engine.WorkflowQuotaReporter) if !ok { return nil } resources := reporter.QuotaResources(wf) if err := quota.ValidateResources(resources); err != nil { return nil } if len(resources) == 0 { return nil } return maps.Clone(resources) } func (e *Executor) handleCommit(ctx context.Context, cl *millv1.CommitLease) { e.mu.Lock() res := e.active[cl.GetLeaseId()] if res == nil { e.mu.Unlock() e.sendReject(cl.GetLeaseId(), "reservation missing or expired", millv1.RejectClass_REJECT_CLASS_TRANSIENT) return } if res.committed { e.mu.Unlock() e.sendCommitted(cl.GetLeaseId()) return } var bindings []models.CacheBinding if res.pipeline.TrustedSource { var err error bindings, err = cacheBindingsFromProto(res.pipeline.RepoDid.String(), res.wf.Caches, cl.GetCacheBindings()) if err != nil { e.mu.Unlock() e.sendReject(cl.GetLeaseId(), "invalid cache plan: "+err.Error(), millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE) return } } res.wf.CacheBindings = bindings res.committed = true res.millRecordsTerminalMetrics = cl.GetMillRecordsTerminalMetrics() if res.ttlTimer != nil { res.ttlTimer.Stop() } committedSecrets := cl.GetSecrets() if !res.pipeline.TrustedSource { committedSecrets = nil } res.vault = newMemVault(committedSecrets) parentCtx := e.lifecycleCtx if parentCtx == nil { parentCtx = ctx } if parentCtx == nil { parentCtx = context.Background() } if res.traceParent.IsValid() { parentCtx = trace.ContextWithSpanContext(parentCtx, res.traceParent) } commitCtx := observability.ExtractFromTraceparentAndTracestate(parentCtx, cl.GetTraceparent(), cl.GetTracestate()) runCtx, runSpan := observability.Tracer().Start(commitCtx, "executor.run") if runSpan.IsRecording() { attrs := []attribute.KeyValue{ attribute.String(observability.WorkflowIDKey, res.wid.String()), attribute.String(observability.LeaseIDKey, res.leaseID), attribute.String(observability.ExecutorNodeIDKey, e.nodeID), attribute.String(observability.PipelineIDKey, res.wid.PipelineId.AtUri().String()), } if res.pipeline.RepoDid.String() != "" { attrs = append(attrs, attribute.String(observability.RepoDIDKey, res.pipeline.RepoDid.String())) } if res.wf != nil && res.wf.OwnerDID != "" { attrs = append(attrs, attribute.String(observability.OwnerDIDKey, res.wf.OwnerDID)) } runSpan.SetAttributes(attrs...) } jobCtx, cancel := context.WithCancelCause(runCtx) if cl.GetMillRecordsTerminalMetrics() { jobCtx = observability.WithWorkflowTerminalObserver(jobCtx, func(_, _, failureClass, reason string) { e.mu.Lock() defer e.mu.Unlock() if e.active[res.leaseID] != res { return } res.failureClass = failureClass res.failureReason = reason }) } res.cancel = cancel res.runDone = make(chan struct{}) e.mu.Unlock() vault := res.vault re := newReservedEngine(res.realEngine, res.slot) res.pipeline.Workflows = map[models.Engine][]models.Workflow{re: {*res.wf}} cacheController := engine.NewPreplannedCacheController(func(ctx context.Context, update engine.CacheUpdate) error { return e.appendCacheUpdate(res.leaseID, update) }) e.startTail(res) e.jobsWG.Add(1) go func() { defer e.jobsWG.Done() defer close(res.runDone) defer runSpan.End() el := e.l.With( "lease_id", res.leaseID, "node_id", e.nodeID, ) engine.StartWorkflows(el, vault, e.cfg, nil, nil, e.db, e.n, e.cache, cacheController, jobCtx, res.pipeline, res.wid.PipelineId) }() e.sendCommitted(cl.GetLeaseId()) } func (e *Executor) handleRelease(leaseID string) { cleanup, ok := e.takeUncommittedReservation(leaseID, true) if !ok { return } e.releaseReservation(cleanup) } func (e *Executor) handleCancel(leaseID string) { e.mu.Lock() res := e.active[leaseID] if res == nil { e.mu.Unlock() if err := e.appendTerminal(leaseID, string(models.StatusKindCancelled), nil); err != nil { e.l.Error("persist cancelled reservation terminal", "lease", leaseID, "err", err) return } e.sendCancelAck(leaseID) return } res.cancelled = true cancel := res.cancel committed := res.committed var cleanup func() if !committed { cleanup = e.removeReservationLocked(res, true) } e.mu.Unlock() if !committed { if err := e.appendTerminal(leaseID, string(models.StatusKindCancelled), nil); err != nil { e.l.Error("persist cancelled reservation terminal", "lease", leaseID, "err", err) e.releaseReservation(cleanup) return } e.sendCancelAck(leaseID) e.releaseReservation(cleanup) return } e.sendCancelAck(leaseID) if cancel != nil { cancel(engine.ErrWorkflowCanceled) } } func (e *Executor) expireReservation(leaseID string) { cleanup, ok := e.takeUncommittedReservation(leaseID, true) if !ok { return } e.releaseReservation(cleanup) } func (e *Executor) takeUncommittedReservation(leaseID string, releaseSlot bool) (func(), bool) { e.mu.Lock() defer e.mu.Unlock() res := e.active[leaseID] if res == nil || res.committed { return nil, false } if res.ttlTimer != nil { res.ttlTimer.Stop() } return e.removeReservationLocked(res, releaseSlot), true } func (e *Executor) removeReservationLocked(res *reservation, releaseSlot bool) func() { delete(e.active, res.leaseID) e.cleanupPending++ slot := res.slot return func() { if releaseSlot && slot != nil { slot.Release() } e.mu.Lock() e.cleanupPending-- if len(e.active) == 0 && e.cleanupPending == 0 && e.idleCh != nil { close(e.idleCh) e.idleCh = nil } e.mu.Unlock() } } func (e *Executor) snapshotLoop(ctx context.Context, enc *millproto.Encoder) { ticker := time.NewTicker(snapshotEvery) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: e.pushSnapshot() } } } func (e *Executor) pushSnapshot() { e.snapshotMu.Lock() defer e.snapshotMu.Unlock() e.pushSnapshotLocked() } func (e *Executor) pushSnapshotLocked() { e.mu.Lock() activeLeases := make([]string, 0, len(e.active)) for leaseID := range e.active { activeLeases = append(activeLeases, leaseID) } draining := e.draining e.mu.Unlock() avail := make(map[string]*millv1.EngineAvailability) for name, eng := range e.engines { a := &millv1.EngineAvailability{Available: !draining} if getter, ok := eng.(interface{ Load() map[string]float64 }); ok { a.Load = getter.Load() } avail[name] = a } e.nextSeqno++ snap := &millproto.Message{ NodeSnapshot: &millv1.NodeSnapshot{ Seqno: e.nextSeqno, Engines: avail, ActiveLeaseIds: activeLeases, }, } e.send(snap) } func (e *Executor) Drain(ctx context.Context) error { e.mu.Lock() e.draining = true idleCh := e.idleCh if len(e.active) == 0 && e.cleanupPending == 0 { idleCh = make(chan struct{}) close(idleCh) } else if idleCh == nil { idleCh = make(chan struct{}) e.idleCh = idleCh } e.mu.Unlock() snapshotDone := make(chan struct{}) go func() { e.pushSnapshot() close(snapshotDone) }() if err := waitForDrainStage(ctx, snapshotDone); err != nil { return err } if err := waitForDrainStage(ctx, idleCh); err != nil { return err } jobsDone := make(chan struct{}) go func() { e.jobsWG.Wait() close(jobsDone) }() if err := waitForDrainStage(ctx, jobsDone); err != nil { return err } return e.drainOutbox(ctx) } func (e *Executor) drainOutbox(ctx context.Context) error { grace := defaultDrainOutboxGrace if e.cfg != nil && e.cfg.Mill.ReconnectGrace > 0 { grace = e.cfg.Mill.ReconnectGrace } timer := time.NewTimer(grace) defer timer.Stop() for { idle := e.outboxIdle() select { case <-idle: e.eventMu.Lock() outboxBytes := e.outboxBytes e.eventMu.Unlock() if outboxBytes == 0 { return nil } case <-timer.C: e.eventMu.Lock() outboxBytes := e.outboxBytes e.eventMu.Unlock() if outboxBytes > 0 { e.l.Warn("executor drain left durable outbox for replay", "bytes", outboxBytes, "grace", grace) } return nil case <-ctx.Done(): return ctx.Err() } } } func waitForDrainStage(ctx context.Context, done <-chan struct{}) error { select { case <-done: return nil case <-ctx.Done(): return ctx.Err() } } func (e *Executor) seatCapacity() int { if e.seats > 0 { return e.seats } return defaultSeats } func ttlDuration(secs uint32) time.Duration { if secs == 0 { return defaultReservationTTL } return time.Duration(secs) * time.Second } const defaultReservationTTL = 60 * time.Second func normalizeLabels(labels []string) []string { seen := make(map[string]struct{}, len(labels)) out := make([]string, 0, len(labels)) for _, label := range labels { label = strings.TrimSpace(label) if label == "" { continue } if _, ok := seen[label]; ok { continue } seen[label] = struct{}{} out = append(out, label) } return out } func (e *Executor) RegisterMetrics(metrics *observability.Metrics) { if metrics == nil { return } metrics.RegisterExecutorGauges( func() float64 { return float64(e.seatCapacity()) }, func() float64 { e.mu.Lock() defer e.mu.Unlock() reservations := 0 for _, r := range e.active { if !r.committed { reservations++ } } return float64(reservations) }, func() float64 { e.mu.Lock() defer e.mu.Unlock() jobs := 0 for _, r := range e.active { if r.committed { jobs++ } } return float64(jobs) }, func() float64 { e.eventMu.Lock() defer e.eventMu.Unlock() return float64(e.outboxBytes) }, ) }