iop/packages/go/agenttask/workflow.go

271 lines
8 KiB
Go

package agenttask
import (
"context"
"fmt"
"reflect"
"sort"
)
type workflowDecision struct {
Activated bool
AutoResume bool
CommandID CommandID
WorkflowRevision WorkflowRevision
}
func (m *Manager) observeWorkflows(ctx context.Context) ([]ProjectID, error) {
registered, err := m.workflow.RegisteredProjects(ctx)
if err != nil {
return nil, err
}
sort.Slice(registered, func(left, right int) bool {
return registered[left] < registered[right]
})
active := make([]ProjectID, 0, len(registered))
seen := make(map[ProjectID]struct{}, len(registered))
for _, projectID := range registered {
if _, duplicate := seen[projectID]; duplicate {
continue
}
seen[projectID] = struct{}{}
snapshot, snapshotErr := m.workflow.Snapshot(ctx, projectID)
if snapshotErr != nil {
m.blockProject(ctx, projectID, Blocker{
Code: BlockerWorkflowUnavailable,
Message: snapshotErr.Error(),
Retryable: true,
})
continue
}
if err := validateWorkflowSnapshot(snapshot); err != nil {
m.blockProject(ctx, projectID, Blocker{
Code: BlockerInvalidIdentity,
Message: err.Error(),
})
continue
}
if snapshot.ProjectID != projectID {
m.blockProject(ctx, projectID, Blocker{
Code: BlockerInvalidIdentity,
Message: "workflow adapter returned a different project identity",
})
continue
}
decision, err := mutateDecision(m, ctx, func(state *ManagerState) (workflowDecision, error) {
project := state.Projects[projectID]
if project.ProjectID == "" {
project = ProjectRecord{
ProjectID: projectID,
WorkspaceID: snapshot.WorkspaceID,
Status: ProjectStatusObserved,
Works: make(map[WorkUnitID]WorkRecord),
}
}
if project.WorkspaceID != snapshot.WorkspaceID {
project.Status = ProjectStatusBlocked
project.Blocker = &Blocker{
Code: BlockerInvalidIdentity,
Message: "registered workspace identity does not match workflow snapshot",
}
project.UpdatedAt = m.clock.Now()
state.Projects[projectID] = project
return workflowDecision{}, nil
}
project.Workflow = pointerWorkflow(snapshot)
project.UpdatedAt = m.clock.Now()
if project.Intent == nil {
project.Status = ProjectStatusObserved
project.Blocker = nil
state.Projects[projectID] = project
return workflowDecision{}, nil
}
if project.Intent.WorkflowRevision != snapshot.Revision {
project.Status = ProjectStatusBlocked
project.Blocker = &Blocker{
Code: BlockerWorkflowRevisionDrift,
Message: "manual start is pinned to a different workflow revision",
}
state.Projects[projectID] = project
return workflowDecision{
CommandID: project.Intent.CommandID,
WorkflowRevision: project.Intent.WorkflowRevision,
}, nil
}
if project.Status == ProjectStatusStopped {
state.Projects[projectID] = project
return workflowDecision{
CommandID: project.Intent.CommandID,
WorkflowRevision: project.Intent.WorkflowRevision,
}, nil
}
interrupted := project.Status == ProjectStatusRunning
explicitlyStarted := project.Status == ProjectStatusStarted
if interrupted && !project.Intent.AutoResumeInterrupted {
project.Status = ProjectStatusStopped
for id, work := range project.Works {
if !work.State.Terminal() {
preStop := work.State
if preStop != WorkStateStopped {
work.ResumeStage = preStop
}
work.State = WorkStateStopped
work.UpdatedAt = m.clock.Now()
project.Works[id] = work
}
}
state.Projects[projectID] = project
return workflowDecision{
CommandID: project.Intent.CommandID,
WorkflowRevision: project.Intent.WorkflowRevision,
}, nil
}
for _, unit := range snapshot.Units {
if unit.MilestoneID != project.Intent.MilestoneID {
continue
}
work, exists := project.Works[unit.ID]
if exists && !reflect.DeepEqual(work.Unit, unit) {
work.State = WorkStateBlocked
work.Blocker = &Blocker{
Code: BlockerWorkflowRevisionDrift,
Message: fmt.Sprintf("work unit %q changed under an immutable workflow revision", unit.ID),
}
work.UpdatedAt = m.clock.Now()
project.Works[unit.ID] = work
continue
}
if !exists {
work = WorkRecord{
Unit: cloneUnit(unit),
State: WorkStateObserved,
UpdatedAt: m.clock.Now(),
}
}
if unit.Completed {
work.State = WorkStateCompleted
work.CompletionVerified = true
}
if explicitlyStarted || work.State == WorkStateStopped || work.ResumeStage != "" {
recoverStoppedWork(&work)
} else if interrupted && project.Intent.AutoResumeInterrupted {
recoverInterruptedWork(&work)
}
project.Works[unit.ID] = work
}
project.Status = ProjectStatusRunning
project.Blocker = nil
project.UpdatedAt = m.clock.Now()
state.Projects[projectID] = project
return workflowDecision{
Activated: true,
AutoResume: interrupted && project.Intent.AutoResumeInterrupted,
CommandID: project.Intent.CommandID,
WorkflowRevision: project.Intent.WorkflowRevision,
}, nil
})
if err != nil {
return active, err
}
m.emit(ctx, Event{
Type: EventObserved,
ProjectID: projectID,
WorkspaceID: snapshot.WorkspaceID,
CommandID: decision.CommandID,
WorkflowRevision: snapshot.Revision,
Detail: string(snapshot.Revision),
})
if decision.Activated {
if decision.AutoResume {
m.emit(ctx, Event{
Type: EventAutoResume,
ProjectID: projectID,
WorkspaceID: snapshot.WorkspaceID,
CommandID: decision.CommandID,
WorkflowRevision: snapshot.Revision,
Detail: string(snapshot.Revision),
})
}
active = append(active, projectID)
}
}
return active, nil
}
func recoverStoppedWork(work *WorkRecord) {
target := work.ResumeStage
if target == "" {
target = work.State
}
switch target {
case WorkStatePreparing, WorkStateDispatching:
work.Attempt++
work.AttemptID = attemptID(work.Unit.ID, work.Attempt)
work.Target = nil
work.ContinuationTarget = nil
work.Isolation = nil
work.Submission = nil
work.Review = nil
work.ChangeSet = nil
work.Integration = nil
work.IntegrationAttempt = 0
work.Locators = make(map[LocatorKind]LocatorRecord)
work.Blocker = nil
work.CompletionVerified = false
work.State = WorkStateReady
case WorkStateReady, WorkStateStopped:
work.State = WorkStateReady
case WorkStateSubmitted, WorkStateReviewing:
work.State = WorkStateReviewing
case WorkStateIntegrating, WorkStatePendingIntegration:
work.State = WorkStatePendingIntegration
case WorkStateObserved:
work.State = WorkStateObserved
case WorkStateCompleted, WorkStateTerminalDeferred, WorkStateBlocked:
// stay terminal
default:
work.State = WorkStateReady
}
work.ResumeStage = ""
}
func recoverInterruptedWork(work *WorkRecord) {
switch work.State {
case WorkStatePreparing:
work.State = WorkStateReady
case WorkStateDispatching:
// Durable execution reconciliation decides whether a child is live,
// submitted, absent, or ambiguous before workflow activation.
case WorkStateSubmitted:
work.State = WorkStateReviewing
case WorkStateReviewing:
// Review is replayed with the same idempotency key.
case WorkStateIntegrating:
work.State = WorkStatePendingIntegration
}
}
func pointerWorkflow(snapshot ProjectWorkflowSnapshot) *ProjectWorkflowSnapshot {
value := cloneWorkflow(snapshot)
return &value
}
func (m *Manager) blockProject(ctx context.Context, projectID ProjectID, blocker Blocker) {
_ = m.mutate(ctx, func(state *ManagerState) error {
project := state.Projects[projectID]
if project.ProjectID == "" {
project.ProjectID = projectID
project.Works = make(map[WorkUnitID]WorkRecord)
}
project.Status = ProjectStatusBlocked
project.Blocker = &blocker
project.UpdatedAt = m.clock.Now()
state.Projects[projectID] = project
return nil
})
m.emit(ctx, Event{
Type: EventBlocked,
ProjectID: projectID,
Detail: string(blocker.Code),
})
}