package agenttask import ( "context" "errors" "fmt" "reflect" "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 activeMu sync.Mutex activeRuns map[ProjectID]context.CancelFunc } // 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 } 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 err } func (m *Manager) StopProject(ctx context.Context, projectID ProjectID) error { if err := validateIdentity("project", string(projectID)); err != nil { return err } m.activeMu.Lock() cancel := m.activeRuns[projectID] m.activeMu.Unlock() 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 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 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 (m *Manager) emit(ctx context.Context, 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.Detail, ) } event.Timestamp = m.clock.Now() _ = m.events.Emit(ctx, event) } 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) m.activeMu.Lock() if previous := m.activeRuns[projectID]; previous != nil { previous() } m.activeRuns[projectID] = cancel m.activeMu.Unlock() return ctx, func() { cancel() m.activeMu.Lock() delete(m.activeRuns, projectID) m.activeMu.Unlock() } } 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), ) }