iop/packages/go/agenttask/reconcile.go

275 lines
7.4 KiB
Go

package agenttask
import (
"context"
"errors"
"fmt"
"sort"
"sync"
)
func (m *Manager) Reconcile(ctx context.Context) error {
m.reconcileMu.Lock()
defer m.reconcileMu.Unlock()
active, err := m.observeWorkflows(ctx)
if err != nil {
return err
}
claimed := make([]ProjectID, 0, len(active))
for _, projectID := range active {
ok, claimErr := m.claimProject(ctx, projectID)
if claimErr != nil {
return claimErr
}
if !ok {
m.emit(ctx, Event{
Type: EventBlocked,
ProjectID: projectID,
Detail: string(BlockerDuplicateProjectLease),
})
continue
}
claimed = append(claimed, projectID)
}
defer func() {
for _, projectID := range claimed {
m.releaseProject(context.WithoutCancel(ctx), projectID)
}
}()
if len(claimed) == 0 {
return nil
}
projectContexts := make(map[ProjectID]context.Context, len(claimed))
projectCleanups := make([]func(), 0, len(claimed))
for _, projectID := range claimed {
projectCtx, cleanup := m.beginProjectRun(ctx, projectID)
projectContexts[projectID] = projectCtx
projectCleanups = append(projectCleanups, cleanup)
}
defer func() {
for _, cleanup := range projectCleanups {
cleanup()
}
}()
var reconcileErrors []error
for round := 0; round < 10_000; round++ {
if err := m.refreshDependencies(ctx, claimed); err != nil {
return err
}
candidates, err := m.runnableWorks(ctx, claimed)
if err != nil {
return err
}
if len(candidates) > 0 {
var wait sync.WaitGroup
errs := make(chan error, len(candidates))
for _, candidate := range candidates {
candidate := candidate
wait.Add(1)
go func() {
defer wait.Done()
projectCtx := projectContexts[candidate.ProjectID]
if runErr := m.runWork(projectCtx, candidate.ProjectID, candidate.WorkUnitID); runErr != nil {
errs <- runErr
}
}()
}
wait.Wait()
close(errs)
for runErr := range errs {
reconcileErrors = append(reconcileErrors, runErr)
}
}
integrated, integrationErr := m.integratePending(ctx, claimed)
if integrationErr != nil {
reconcileErrors = append(reconcileErrors, integrationErr)
}
if len(candidates) == 0 && !integrated {
break
}
if ctx.Err() != nil {
return errors.Join(append(reconcileErrors, ctx.Err())...)
}
}
if err := m.refreshProjectStatuses(ctx, claimed); err != nil {
return err
}
return errors.Join(reconcileErrors...)
}
type runnableWork struct {
ProjectID ProjectID
WorkUnitID WorkUnitID
Ordinal DispatchOrdinal
}
func (m *Manager) refreshDependencies(ctx context.Context, active []ProjectID) error {
activeSet := make(map[ProjectID]struct{}, len(active))
for _, id := range active {
activeSet[id] = struct{}{}
}
return m.mutate(ctx, func(state *ManagerState) error {
projectIDs := make([]ProjectID, 0, len(activeSet))
for id := range activeSet {
projectIDs = append(projectIDs, id)
}
sort.Slice(projectIDs, func(left, right int) bool {
return projectIDs[left] < projectIDs[right]
})
for _, projectID := range projectIDs {
project := state.Projects[projectID]
if project.Status != ProjectStatusRunning || project.Workflow == nil {
continue
}
workIDs := make([]WorkUnitID, 0, len(project.Works))
for id := range project.Works {
workIDs = append(workIDs, id)
}
sort.Slice(workIDs, func(left, right int) bool {
return workIDs[left] < workIDs[right]
})
for _, workID := range workIDs {
work := project.Works[workID]
if work.State != WorkStateObserved {
continue
}
dependency := evaluateDependencies(work.Unit, *project.Workflow, project.Works)
switch dependency.Status {
case dependencyReady:
if err := transitionWork(&work, WorkStateReady); err != nil {
return err
}
if work.DispatchOrdinal == 0 {
state.NextOrdinal++
work.DispatchOrdinal = state.NextOrdinal
}
if work.Attempt == 0 {
work.Attempt = 1
work.AttemptID = attemptID(work.Unit.ID, work.Attempt)
}
work.Blocker = nil
work.UpdatedAt = m.clock.Now()
m.emit(ctx, Event{
Type: EventDependencyReady,
ProjectID: projectID,
WorkspaceID: project.WorkspaceID,
WorkUnitID: workID,
CommandID: project.Intent.CommandID,
WorkflowRevision: project.Intent.WorkflowRevision,
AttemptID: work.AttemptID,
Ordinal: work.DispatchOrdinal,
State: work.State,
WriteSetKind: work.Unit.WriteSetKind,
IsolationMode: work.Unit.IsolationMode,
})
case dependencyMissing:
blockWorkDependency(&work, BlockerDependencyMissing, dependency.Ref)
case dependencyAmbiguous:
blockWorkDependency(&work, BlockerDependencyAmbiguous, dependency.Ref)
case dependencyBlocked:
blockWorkDependency(&work, BlockerDependencyBlocked, dependency.Ref)
case dependencyWaiting:
continue
}
project.Works[workID] = work
}
state.Projects[projectID] = project
}
return nil
})
}
func blockWorkDependency(work *WorkRecord, code BlockerCode, ref string) {
_ = transitionWork(work, WorkStateBlocked)
work.Blocker = &Blocker{
Code: code,
Message: fmt.Sprintf("explicit predecessor %q is %s", ref, code),
}
}
func (m *Manager) runnableWorks(
ctx context.Context,
active []ProjectID,
) ([]runnableWork, error) {
state, err := m.load(ctx)
if err != nil {
return nil, err
}
var candidates []runnableWork
for _, projectID := range active {
project := state.Projects[projectID]
if project.Status != ProjectStatusRunning {
continue
}
for workID, work := range project.Works {
if work.State != WorkStateReady && work.State != WorkStateReviewing {
continue
}
candidates = append(candidates, runnableWork{
ProjectID: projectID, WorkUnitID: workID, Ordinal: work.DispatchOrdinal,
})
}
}
sort.Slice(candidates, func(left, right int) bool {
if candidates[left].Ordinal != candidates[right].Ordinal {
return candidates[left].Ordinal < candidates[right].Ordinal
}
if candidates[left].ProjectID != candidates[right].ProjectID {
return candidates[left].ProjectID < candidates[right].ProjectID
}
return candidates[left].WorkUnitID < candidates[right].WorkUnitID
})
return candidates, nil
}
func (m *Manager) refreshProjectStatuses(ctx context.Context, active []ProjectID) error {
return m.mutate(ctx, func(state *ManagerState) error {
for _, projectID := range active {
project := state.Projects[projectID]
if project.Status == ProjectStatusStopped {
continue
}
selected := 0
completed := 0
activeWork := 0
for _, work := range project.Works {
selected++
switch {
case work.State == WorkStateCompleted:
completed++
case !work.State.Terminal():
activeWork++
}
}
switch {
case selected > 0 && completed == selected:
project.Status = ProjectStatusCompleted
project.Blocker = nil
case activeWork > 0:
project.Status = ProjectStatusRunning
default:
project.Status = ProjectStatusBlocked
}
project.UpdatedAt = m.clock.Now()
state.Projects[projectID] = project
if project.Status == ProjectStatusCompleted {
var cmdID CommandID
var wfRev WorkflowRevision
if project.Intent != nil {
cmdID = project.Intent.CommandID
wfRev = project.Intent.WorkflowRevision
}
m.emit(ctx, Event{
Type: EventCompleted,
ProjectID: projectID,
WorkspaceID: project.WorkspaceID,
CommandID: cmdID,
WorkflowRevision: wfRev,
})
}
}
return nil
})
}