package agenttask import ( "context" "errors" "fmt" "reflect" "sort" "strconv" "strings" "sync" "time" ) type ManagerConfig struct { OwnerID string LeaseDuration time.Duration MaxReworkAttempts uint32 MaxFailureAttempts uint32 StateWriteAttempts int } type Manager struct { config ManagerConfig clock Clock store StateStore workflow WorkflowAdapter selector Selector isolation IsolationBackend invoker ProviderInvoker recovery RecoveryInspector evidence WorkflowEvidence reviewer Reviewer integrator Integrator events EventSink scheduler *Scheduler // renewalTicks is a package-private test seam. Production leaves it nil // and uses a time.Ticker; tests can drive renewal without wall-clock waits. renewalTicks func(time.Duration) <-chan time.Time reconcileMu sync.Mutex deliveryMu sync.Mutex activeMu sync.Mutex activeRuns map[ProjectID]context.CancelFunc } type deliveryErrorContextKey struct{} type deliveryErrorCollector struct { mu sync.Mutex errs []error } func withDeliveryErrors(ctx context.Context) (context.Context, *deliveryErrorCollector) { collector := &deliveryErrorCollector{} return context.WithValue(ctx, deliveryErrorContextKey{}, collector), collector } func recordDeliveryError(ctx context.Context, err error) { if err == nil { return } collector, _ := ctx.Value(deliveryErrorContextKey{}).(*deliveryErrorCollector) if collector == nil { return } collector.mu.Lock() collector.errs = append(collector.errs, err) collector.mu.Unlock() } func (c *deliveryErrorCollector) Err() error { if c == nil { return nil } c.mu.Lock() defer c.mu.Unlock() return errors.Join(c.errs...) } // leaseSet owns the collection of active durable lease claims for one // reconciliation. A background supervisor renews each claim by CAS at a // bounded fraction of LeaseDuration; if any renewal cannot prove the token // still matches, the guarded context is cancelled so that no late result can // commit durable state under a stale owner. type leaseSet struct { mu sync.Mutex claims []*leaseClaim cancel context.CancelFunc manager *Manager wg sync.WaitGroup } type leaseSetContextKey struct{} type leaseFenceBypassContextKey struct{} // renewAll renews every tracked claim by CAS. It returns the first error so // the supervisor can cancel the guarded context and stop the loop. func (ls *leaseSet) renewAll(ctx context.Context) error { ls.mu.Lock() defer ls.mu.Unlock() for _, claim := range ls.claims { if err := ls.manager.renewOneLease(ctx, claim); err != nil { return err } } return nil } // Validate confirms that every tracked claim still matches the current // durable state. It returns ErrLeaseLost for the first mismatch so callers // can reject an external result that arrived after ownership was lost. func (ls *leaseSet) Validate(ctx context.Context) error { select { case <-ctx.Done(): return ctx.Err() default: } ls.mu.Lock() claims := make([]*leaseClaim, len(ls.claims)) copy(claims, ls.claims) ls.mu.Unlock() for _, claim := range claims { if err := ls.manager.validateOneLease(ctx, claim); err != nil { return err } } return nil } // Add registers a newly acquired claim with the set so the supervisor also // renews it and the fence validator also checks it. func (ls *leaseSet) Add(claim *leaseClaim) { if claim == nil { return } ls.mu.Lock() ls.claims = append(ls.claims, claim) ls.mu.Unlock() } // Remove prevents any later renewal pass from observing claim. It waits for // an in-progress renewal pass through the set lock before the caller performs // the exact-token release. func (ls *leaseSet) Remove(claim *leaseClaim) { if claim == nil { return } ls.mu.Lock() defer ls.mu.Unlock() for index, current := range ls.claims { if current == claim { ls.claims = append(ls.claims[:index], ls.claims[index+1:]...) return } } } func (ls *leaseSet) snapshotClaims() []*leaseClaim { ls.mu.Lock() defer ls.mu.Unlock() claims := make([]*leaseClaim, len(ls.claims)) copy(claims, ls.claims) return claims } // Close cancels the supervisor and the guarded context and waits for the // renewal goroutine to exit. It does not touch durable state; callers // release each lease by exact token after closing. func (ls *leaseSet) Close() { if ls == nil || ls.cancel == nil { return } ls.cancel() ls.wg.Wait() } // maintainLeases starts a reconciliation-owned lease supervisor. The returned // context is cancelled when any renewal fails, so that every mutation under // it is guaranteed to observe a token the supervisor still holds. func (m *Manager) maintainLeases(ctx context.Context, initial *leaseClaim) (context.Context, *leaseSet) { ctx, cancel := context.WithCancel(ctx) ls := &leaseSet{ claims: []*leaseClaim{initial}, cancel: cancel, manager: m, } ctx = context.WithValue(ctx, leaseSetContextKey{}, ls) ls.wg.Add(1) go ls.runRenewalLoop(ctx) return ctx, ls } // runRenewalLoop ticks at a bounded fraction of LeaseDuration and renews all // tracked claims. A single CAS miss cancels the guarded context and stops. func (ls *leaseSet) runRenewalLoop(ctx context.Context) { defer ls.wg.Done() interval := ls.manager.config.LeaseDuration / 3 if interval < time.Millisecond { interval = time.Millisecond } if ticks := ls.manager.renewalTicks; ticks != nil { for { select { case <-ctx.Done(): return case <-ticks(interval): if err := ls.renewAll(unfencedLeaseContext(ctx)); err != nil { ls.cancel() return } } } } ticker := time.NewTicker(interval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: if err := ls.renewAll(unfencedLeaseContext(ctx)); err != nil { ls.cancel() return } } } } // Wait blocks until the renewal supervisor has fully stopped. func (ls *leaseSet) Wait() { ls.wg.Wait() } // renewOneLease renews a single claim by CAS. It returns ErrLeaseLost when // the current state no longer holds the tracked token. func (m *Manager) renewOneLease(ctx context.Context, claim *leaseClaim) error { switch claim.scope { case "device": return m.renewDevice(ctx, claim.token) case "workspace": return m.renewWorkspace(ctx, claim.token, claim.subject) case "project": return m.renewProject(ctx, claim.token, claim.subject) case "integration": return m.renewIntegration(ctx, claim.token, claim.subject) default: return fmt.Errorf("agenttask: unknown lease scope %q", claim.scope) } } // validateOneLease checks that a single claim still matches current state. func (m *Manager) validateOneLease(ctx context.Context, claim *leaseClaim) error { switch claim.scope { case "device": return m.validateDevice(ctx, claim.token) case "workspace": return m.validateWorkspace(ctx, claim.token, claim.subject) case "project": return m.validateProject(ctx, claim.token, claim.subject) case "integration": return m.validateIntegration(ctx, claim.token, claim.subject) default: return fmt.Errorf("agenttask: unknown lease scope %q", claim.scope) } } func NewManager( config ManagerConfig, clock Clock, store StateStore, workflow WorkflowAdapter, selector Selector, isolation IsolationBackend, invoker ProviderInvoker, recovery RecoveryInspector, evidence WorkflowEvidence, reviewer Reviewer, integrator Integrator, events EventSink, ) (*Manager, error) { if err := validateIdentity("manager_owner", config.OwnerID); err != nil { return nil, err } if store == nil || workflow == nil || selector == nil || isolation == nil || invoker == nil || recovery == nil || evidence == nil || reviewer == nil || integrator == nil { return nil, errors.New("agenttask: every execution port is required; unsafe fallback is disabled") } if clock == nil { clock = systemClock{} } if events == nil { events = nopEventSink{} } if config.LeaseDuration <= 0 { config.LeaseDuration = 30 * time.Second } if config.MaxReworkAttempts == 0 { config.MaxReworkAttempts = 3 } if config.MaxFailureAttempts == 0 { config.MaxFailureAttempts = 10 } if config.StateWriteAttempts <= 0 { config.StateWriteAttempts = 32 } return &Manager{ config: config, clock: clock, store: store, workflow: workflow, selector: selector, isolation: isolation, invoker: invoker, recovery: recovery, evidence: evidence, reviewer: reviewer, integrator: integrator, events: events, scheduler: NewScheduler(), activeRuns: make(map[ProjectID]context.CancelFunc), }, nil } func (m *Manager) StartProject(ctx context.Context, req StartRequest) error { if err := validateStartRequest(req); err != nil { return err } ctx, deliveryErrors := withDeliveryErrors(ctx) if err := m.flushPendingEvents(ctx); err != nil { recordDeliveryError(ctx, err) return deliveryErrors.Err() } autoResume := true if req.AutoResumeInterrupted != nil { autoResume = *req.AutoResumeInterrupted } intent := StartIntent{ CommandID: req.CommandID, ProjectID: req.ProjectID, WorkspaceID: req.WorkspaceID, MilestoneID: req.MilestoneID, WorkflowRevision: req.WorkflowRevision, ConfigRevision: req.ConfigRevision, GrantRevision: req.GrantRevision, AutoResumeInterrupted: autoResume, StartedAt: m.clock.Now(), } err := m.mutate(ctx, func(state *ManagerState) error { if previous, ok := state.Commands[req.CommandID]; ok { if sameCommandIntent(previous.Intent, intent) { return nil } return fmt.Errorf("agenttask: command %q was already used with different immutable input", req.CommandID) } project := state.Projects[req.ProjectID] if project.ProjectID != "" && project.WorkspaceID != "" && project.WorkspaceID != req.WorkspaceID { return fmt.Errorf("agenttask: project %q workspace identity changed", req.ProjectID) } project.ProjectID = req.ProjectID project.WorkspaceID = req.WorkspaceID project.Status = ProjectStatusStarted project.Intent = &intent project.Blocker = nil project.UpdatedAt = m.clock.Now() if project.Works == nil { project.Works = make(map[WorkUnitID]WorkRecord) } state.Commands[req.CommandID] = CommandRecord{Intent: intent} state.Projects[req.ProjectID] = project return nil }) if err == nil { m.emit(ctx, Event{ Type: EventManualStart, ProjectID: req.ProjectID, WorkspaceID: req.WorkspaceID, CommandID: req.CommandID, WorkflowRevision: req.WorkflowRevision, Detail: string(req.MilestoneID), }) } return errors.Join(err, deliveryErrors.Err()) } func (m *Manager) StopProject(ctx context.Context, projectID ProjectID) error { if err := validateIdentity("project", string(projectID)); err != nil { return err } ctx, deliveryErrors := withDeliveryErrors(ctx) if err := m.flushPendingEvents(ctx); err != nil { recordDeliveryError(ctx, err) return deliveryErrors.Err() } m.activeMu.Lock() cancel := m.activeRuns[projectID] m.activeMu.Unlock() hadActiveRun := cancel != nil if cancel != nil { cancel() } var commandID CommandID var workflowRev WorkflowRevision err := m.mutate(ctx, func(state *ManagerState) error { project, ok := state.Projects[projectID] if !ok { return fmt.Errorf("agenttask: project %q is not registered in manager state", projectID) } if project.Intent != nil { commandID = project.Intent.CommandID workflowRev = project.Intent.WorkflowRevision } project.Status = ProjectStatusStopped liveLease := project.Lease != nil && project.Lease.ExpiresAt.After(m.clock.Now()) if !hadActiveRun && !liveLease { project.Lease = nil if lease, ok := state.WorkspaceLeases[project.WorkspaceID]; ok && lease.OwnerID == m.config.OwnerID { delete(state.WorkspaceLeases, project.WorkspaceID) } if lease, ok := state.IntegrationLeases[project.WorkspaceID]; ok && lease.OwnerID == m.config.OwnerID { delete(state.IntegrationLeases, project.WorkspaceID) } } project.UpdatedAt = m.clock.Now() for id, work := range project.Works { if !work.State.Terminal() { preStop := work.State if preStop != WorkStateStopped { work.ResumeStage = preStop } if err := transitionWork(&work, WorkStateStopped); err != nil { return err } work.UpdatedAt = m.clock.Now() project.Works[id] = work } } state.Projects[projectID] = project return nil }) if err == nil { m.emit(ctx, Event{ Type: EventStopped, ProjectID: projectID, CommandID: commandID, WorkflowRevision: workflowRev, }) } return errors.Join(err, deliveryErrors.Err()) } func mutateDecision[T any]( m *Manager, ctx context.Context, change func(*ManagerState) (T, error), ) (T, error) { var zero T var conflict error for range m.config.StateWriteAttempts { state, revision, err := m.store.Load(ctx) if err != nil { return zero, err } if err := validateManagerState(state); err != nil { return zero, err } if err := validateContextClaims(ctx, state); err != nil { return zero, err } next := cloneState(state) if next.SchemaVersion != currentSchemaVersion { return zero, fmt.Errorf("agenttask: unsupported state schema %d", next.SchemaVersion) } result, err := change(&next) if err != nil { return zero, err } if err := validateManagerState(next); err != nil { return zero, err } _, err = m.store.CompareAndSwap(ctx, revision, next) if errors.Is(err, ErrRevisionConflict) { conflict = err continue } if err != nil { return zero, err } return result, nil } return zero, fmt.Errorf("agenttask: state CAS retry budget exhausted: %w", conflict) } func validateContextClaims(ctx context.Context, state ManagerState) error { if ctx.Value(leaseFenceBypassContextKey{}) != nil { return nil } leases, _ := ctx.Value(leaseSetContextKey{}).(*leaseSet) if leases == nil { return nil } for _, claim := range leases.snapshotClaims() { if err := claimMatchesState(claim, state); err != nil { leases.cancel() return err } } return nil } func claimMatchesState(claim *leaseClaim, state ManagerState) error { if claim == nil { return fmt.Errorf("%w: missing lease claim", ErrLeaseLost) } matched := false switch claim.scope { case "device": matched = state.DeviceLease != nil && state.DeviceLease.OwnerID == claim.owner && state.DeviceLease.Token == claim.token case "project": project, ok := state.Projects[ProjectID(claim.subject)] matched = ok && project.Lease != nil && project.Lease.OwnerID == claim.owner && project.Lease.Token == claim.token case "workspace": lease, ok := state.WorkspaceLeases[WorkspaceID(claim.subject)] matched = ok && lease.OwnerID == claim.owner && lease.Token == claim.token case "integration": lease, ok := state.IntegrationLeases[WorkspaceID(claim.subject)] matched = ok && lease.OwnerID == claim.owner && lease.Token == claim.token default: return fmt.Errorf("%w: unknown lease scope %q", ErrLeaseLost, claim.scope) } if !matched { return fmt.Errorf("%w: %s lease %q no longer matches its exact token", ErrLeaseLost, claim.scope, claim.subject) } return nil } func unfencedLeaseContext(ctx context.Context) context.Context { return context.WithValue(context.WithoutCancel(ctx), leaseFenceBypassContextKey{}, true) } func (m *Manager) mutate( ctx context.Context, change func(*ManagerState) error, ) error { _, err := mutateDecision(m, ctx, func(state *ManagerState) (struct{}, error) { return struct{}{}, change(state) }) return err } func (m *Manager) load(ctx context.Context) (ManagerState, error) { state, _, err := m.store.Load(ctx) if err != nil { return ManagerState{}, err } if err := validateManagerState(state); err != nil { return ManagerState{}, err } state = cloneState(state) if state.SchemaVersion != currentSchemaVersion { return ManagerState{}, fmt.Errorf("agenttask: unsupported state schema %d", state.SchemaVersion) } return state, nil } func durableIdentity(domain string, components ...string) string { var sb strings.Builder sb.WriteString(domain) for _, c := range components { sb.WriteString("/") sb.WriteString(strconv.Itoa(len(c))) sb.WriteString(":") sb.WriteString(c) } return sb.String() } func normalizeEventIdentity(event Event) Event { if event.EventID == "" { event.EventID = durableIdentity( "event-v1", string(event.Type), string(event.ProjectID), string(event.WorkspaceID), string(event.WorkUnitID), string(event.CommandID), string(event.WorkflowRevision), string(event.AttemptID), strconv.FormatUint(uint64(event.Ordinal), 10), string(event.ChangeSetID), event.ChangeSetRevision, strconv.FormatUint(uint64(event.IntegrationAttempt), 10), string(event.State), event.ProviderID, event.ProfileID, string(event.WriteSetKind), string(event.IsolationMode), event.Detail, ) } return event } func sameLogicalEvent(left, right Event) bool { left.Timestamp = time.Time{} right.Timestamp = time.Time{} return reflect.DeepEqual(left, right) } func (m *Manager) emit(ctx context.Context, event Event) { // Direct construction is retained only for package-level identity tests. // NewManager requires a durable store for every production manager. if m.store == nil { event = normalizeEventIdentity(event) if event.Timestamp.IsZero() { event.Timestamp = m.clock.Now() } if m.events != nil { recordDeliveryError(ctx, m.events.Emit(ctx, event)) } return } _, err := m.enqueueEvent(ctx, event) if err == nil { err = m.flushPendingEvents(ctx) } recordDeliveryError(ctx, err) } func (m *Manager) enqueueEvent(ctx context.Context, event Event) (EventDelivery, error) { event = normalizeEventIdentity(event) var conflict error for range m.config.StateWriteAttempts { state, revision, err := m.store.Load(ctx) if err != nil { return EventDelivery{}, err } if err := validateManagerState(state); err != nil { return EventDelivery{}, err } if err := validateContextClaims(ctx, state); err != nil { return EventDelivery{}, err } if pending, ok := state.PendingEvents[event.EventID]; ok { if !sameLogicalEvent(pending.Event, event) { return EventDelivery{}, fmt.Errorf( "agenttask: pending event %q was reused with conflicting logical content", event.EventID, ) } return cloneEventDelivery(pending), nil } project, ok := state.Projects[event.ProjectID] if !ok { return EventDelivery{}, fmt.Errorf( "agenttask: event %q project %q is absent from committed state", event.EventID, event.ProjectID, ) } if event.Timestamp.IsZero() { event.Timestamp = m.clock.Now() } project = cloneProject(project) delivery := EventDelivery{ Event: event, EvidenceRevision: revision, Project: &project, } if event.WorkUnitID != "" { work, ok := project.Works[event.WorkUnitID] if !ok { return EventDelivery{}, fmt.Errorf( "agenttask: event %q work %q is absent from committed state", event.EventID, event.WorkUnitID, ) } work = cloneWork(work) delivery.Work = &work } if err := validateEventDelivery(event.EventID, delivery); err != nil { return EventDelivery{}, err } next := cloneState(state) next.PendingEvents[event.EventID] = delivery if _, err := m.store.CompareAndSwap(ctx, revision, next); err != nil { if errors.Is(err, ErrRevisionConflict) { conflict = err continue } return EventDelivery{}, err } return cloneEventDelivery(delivery), nil } return EventDelivery{}, fmt.Errorf( "agenttask: pending event CAS retry budget exhausted: %w", conflict, ) } func (m *Manager) flushPendingEvents(ctx context.Context) error { if m.store == nil { return nil } m.deliveryMu.Lock() defer m.deliveryMu.Unlock() state, _, err := m.store.Load(ctx) if err != nil { return err } if err := validateManagerState(state); err != nil { return err } eventIDs := make([]string, 0, len(state.PendingEvents)) for eventID := range state.PendingEvents { eventIDs = append(eventIDs, eventID) } sort.Strings(eventIDs) for _, eventID := range eventIDs { if err := m.flushEventDeliveryLocked(ctx, eventID); err != nil { return fmt.Errorf("agenttask: deliver pending event %q: %w", eventID, err) } } return nil } func (m *Manager) flushEventDeliveryLocked(ctx context.Context, eventID string) error { state, _, err := m.store.Load(ctx) if err != nil { return err } if err := validateManagerState(state); err != nil { return err } if err := validateContextClaims(ctx, state); err != nil { return err } delivery, ok := state.PendingEvents[eventID] if !ok { return nil } if err := m.events.Emit(ctx, delivery.Event); err != nil { return err } var conflict error for range m.config.StateWriteAttempts { state, revision, err := m.store.Load(ctx) if err != nil { return err } if err := validateManagerState(state); err != nil { return err } if err := validateContextClaims(ctx, state); err != nil { return err } current, ok := state.PendingEvents[eventID] if !ok { return nil } if !reflect.DeepEqual(current, delivery) { return fmt.Errorf( "agenttask: pending event %q changed before acknowledgement", eventID, ) } next := cloneState(state) delete(next.PendingEvents, eventID) if _, err := m.store.CompareAndSwap(ctx, revision, next); err != nil { if errors.Is(err, ErrRevisionConflict) { conflict = err continue } return err } return nil } return fmt.Errorf( "agenttask: pending event acknowledgement CAS retry budget exhausted: %w", conflict, ) } func sameCommandIntent(left, right StartIntent) bool { left.StartedAt = time.Time{} right.StartedAt = time.Time{} return reflect.DeepEqual(left, right) } func (m *Manager) beginProjectRun( parent context.Context, projectID ProjectID, ) (context.Context, func()) { ctx, cancel := context.WithCancel(parent) monitorDone := make(chan struct{}) go m.monitorProjectRun(ctx, projectID, cancel, monitorDone) m.activeMu.Lock() if previous := m.activeRuns[projectID]; previous != nil { previous() } m.activeRuns[projectID] = cancel m.activeMu.Unlock() return ctx, func() { cancel() <-monitorDone m.activeMu.Lock() delete(m.activeRuns, projectID) m.activeMu.Unlock() } } func (m *Manager) monitorProjectRun( ctx context.Context, projectID ProjectID, cancel context.CancelFunc, done chan<- struct{}, ) { defer close(done) interval := m.config.LeaseDuration / 10 if interval < 25*time.Millisecond { interval = 25 * time.Millisecond } if interval > 250*time.Millisecond { interval = 250 * time.Millisecond } ticker := time.NewTicker(interval) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: state, err := m.load(ctx) if err != nil { cancel() return } project, ok := state.Projects[projectID] if !ok || project.Status == ProjectStatusStopped { cancel() return } } } } func attemptID(workID WorkUnitID, attempt uint32) AttemptID { return AttemptID(durableIdentity("attempt-v1", string(workID), strconv.FormatUint(uint64(attempt), 10))) } func dispatchKey(projectID ProjectID, workID WorkUnitID, attempt uint32) string { return durableIdentity("dispatch-v1", string(projectID), string(workID), strconv.FormatUint(uint64(attempt), 10)) } func reviewKey(projectID ProjectID, workID WorkUnitID, attempt uint32, artifact ArtifactID) string { return durableIdentity("review-v1", string(projectID), string(workID), strconv.FormatUint(uint64(attempt), 10), string(artifact)) } func integrationKey( projectID ProjectID, workID WorkUnitID, changeSet ChangeSetIdentity, attempt IntegrationAttempt, ) string { return durableIdentity( "integrate-v1", string(projectID), string(workID), string(changeSet.ID), string(changeSet.Revision), strconv.FormatUint(uint64(attempt), 10), ) }