iop/packages/go/agenttask/manager_test.go

483 lines
16 KiB
Go

package agenttask
import (
"context"
"errors"
"sync"
"testing"
"time"
)
func TestNoUnselectedStart(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("ready", WriteSetUnknown))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 2)
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
if harness.invoker.callCount() != 0 {
t.Fatalf("unselected ready project invoked provider %d times", harness.invoker.callCount())
}
project := harness.store.snapshot().Projects["project"]
if project.Status != ProjectStatusObserved || project.Intent != nil {
t.Fatalf("unselected project = %#v", project)
}
}
func TestManualStartFullProgression(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("work", WriteSetDisjoint))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 2)
harness.start("project", "workspace", nil)
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
state := harness.store.snapshot()
project := state.Projects["project"]
work := project.Works["work"]
if project.Status != ProjectStatusCompleted || work.State != WorkStateCompleted {
t.Fatalf("project/work = %s/%s, want completed/completed", project.Status, work.State)
}
if harness.invoker.callCount() != 1 || harness.reviewer.callCount() != 1 ||
harness.integrator.callCount() != 1 {
t.Fatalf(
"calls invoke=%d review=%d integrate=%d",
harness.invoker.callCount(), harness.reviewer.callCount(), harness.integrator.callCount(),
)
}
}
func TestInterruptedResumeDefaultsOn(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 1)
harness.start("project", "workspace", nil)
harness.store.edit(func(state *ManagerState) {
project := state.Projects["project"]
project.Status = ProjectStatusRunning
state.Projects["project"] = project
})
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
if harness.store.snapshot().Projects["project"].Status != ProjectStatusCompleted {
t.Fatalf("default interrupted resume did not complete")
}
if harness.invoker.callCount() != 1 {
t.Fatalf("default resume invocations = %d, want 1", harness.invoker.callCount())
}
}
func TestInterruptedResumeOverrideFalseStops(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 1)
disabled := false
harness.start("project", "workspace", &disabled)
harness.store.edit(func(state *ManagerState) {
project := state.Projects["project"]
project.Status = ProjectStatusRunning
state.Projects["project"] = project
})
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
if harness.store.snapshot().Projects["project"].Status != ProjectStatusStopped {
t.Fatalf("override-off interrupted project was not stopped")
}
if harness.invoker.callCount() != 0 {
t.Fatalf("override-off resume invoked provider %d times", harness.invoker.callCount())
}
}
func TestProjectIsolationKeepsIndependentProjectRunning(t *testing.T) {
snapshots := map[ProjectID]ProjectWorkflowSnapshot{
"broken": testSnapshot("broken", "workspace-broken", testUnit("broken-work", WriteSetUnknown)),
"healthy": testSnapshot("healthy", "workspace-healthy", testUnit("healthy-work", WriteSetUnknown)),
}
harness := newHarness(t, snapshots, 2)
harness.workflow.errors["broken"] = errors.New("fixture workflow parse failure")
harness.start("broken", "workspace-broken", nil)
harness.start("healthy", "workspace-healthy", nil)
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
state := harness.store.snapshot()
if state.Projects["broken"].Status != ProjectStatusBlocked {
t.Fatalf("broken project status = %s", state.Projects["broken"].Status)
}
if state.Projects["healthy"].Status != ProjectStatusCompleted {
t.Fatalf("healthy project status = %s", state.Projects["healthy"].Status)
}
if harness.invoker.callCount() != 1 {
t.Fatalf("provider calls = %d, want healthy project only", harness.invoker.callCount())
}
}
func TestManualStartDuplicateManagerLeasePreventsConcurrentInvocation(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 1)
harness.start("project", "workspace", nil)
harness.store.edit(func(state *ManagerState) {
project := state.Projects["project"]
project.Status = ProjectStatusRunning
project.Lease = &LeaseRecord{
OwnerID: "other-manager", Token: "live",
ExpiresAt: fixedClock{now: harness.manager.clock.Now()}.Now().Add(harness.manager.config.LeaseDuration),
}
state.Projects["project"] = project
})
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
if harness.invoker.callCount() != 0 {
t.Fatalf("live duplicate manager lease invoked provider")
}
}
func TestManualStartStopCancelsInvocationAndReleasesCapacity(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 1)
harness.invoker.delays["work"] = time.Second
harness.start("project", "workspace", nil)
done := make(chan error, 1)
go func() {
done <- harness.manager.Reconcile(context.Background())
}()
deadline := time.Now().Add(time.Second)
for harness.invoker.activeCount() == 0 && time.Now().Before(deadline) {
time.Sleep(time.Millisecond)
}
if harness.invoker.activeCount() == 0 {
t.Fatal("provider invocation did not start")
}
if err := harness.manager.StopProject(context.Background(), "project"); err != nil {
t.Fatalf("StopProject: %v", err)
}
select {
case err := <-done:
if !errors.Is(err, context.Canceled) {
t.Fatalf("Reconcile error = %v, want project cancellation", err)
}
case <-time.After(time.Second):
t.Fatal("Reconcile did not stop cancelled invocation")
}
state := harness.store.snapshot()
if state.Projects["project"].Status != ProjectStatusStopped {
t.Fatalf("project status = %s, want stopped", state.Projects["project"].Status)
}
if harness.manager.scheduler.Active("provider\x00profile") != 0 {
t.Fatal("provider scheduler capacity leaked after stop")
}
}
func TestStopProjectPersistsUntilExplicitRestart(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 1)
harness.start("project", "workspace", nil)
if err := harness.manager.StopProject(context.Background(), "project"); err != nil {
t.Fatalf("StopProject: %v", err)
}
for i := 0; i < 5; i++ {
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile %d: %v", i, err)
}
}
if harness.invoker.callCount() != 0 {
t.Fatalf("stopped project was executed %d times, want 0", harness.invoker.callCount())
}
state := harness.store.snapshot()
if state.Projects["project"].Status != ProjectStatusStopped {
t.Fatalf("project status = %s, want stopped", state.Projects["project"].Status)
}
req := StartRequest{
CommandID: "start-2-project",
ProjectID: "project",
WorkspaceID: "workspace",
MilestoneID: "milestone",
WorkflowRevision: "workflow-r1",
ConfigRevision: "config-r1",
GrantRevision: "grant-r1",
}
if err := harness.manager.StartProject(context.Background(), req); err != nil {
t.Fatalf("StartProject restart: %v", err)
}
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile after restart: %v", err)
}
state = harness.store.snapshot()
if state.Projects["project"].Status != ProjectStatusCompleted {
t.Fatalf("restarted project status = %s, want completed", state.Projects["project"].Status)
}
if harness.invoker.callCount() != 1 {
t.Fatalf("restarted project provider calls = %d, want 1", harness.invoker.callCount())
}
}
func TestStoppedWorkResumesFromDurableStage(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 1)
harness.start("project", "workspace", nil)
harness.store.edit(func(state *ManagerState) {
project := state.Projects["project"]
project.Status = ProjectStatusRunning
work := project.Works["work"]
work.Unit = testUnit("work", WriteSetUnknown)
work.State = WorkStateReviewing
work.ResumeStage = WorkStateReviewing
work.Attempt = 1
work.AttemptID = attemptID("work", 1)
work.Submission = &Submission{
ProjectID: "project",
WorkUnitID: "work",
AttemptID: attemptID("work", 1),
ArtifactID: "art-1",
Ready: true,
}
project.Works["work"] = work
state.Projects["project"] = project
})
if err := harness.manager.StopProject(context.Background(), "project"); err != nil {
t.Fatalf("StopProject: %v", err)
}
state := harness.store.snapshot()
work := state.Projects["project"].Works["work"]
if work.State != WorkStateStopped || work.ResumeStage != WorkStateReviewing {
t.Fatalf("stopped work state/resumeStage = %s/%s, want stopped/reviewing", work.State, work.ResumeStage)
}
req := StartRequest{
CommandID: "start-2-project",
ProjectID: "project",
WorkspaceID: "workspace",
MilestoneID: "milestone",
WorkflowRevision: "workflow-r1",
ConfigRevision: "config-r1",
GrantRevision: "grant-r1",
}
if err := harness.manager.StartProject(context.Background(), req); err != nil {
t.Fatalf("StartProject: %v", err)
}
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
state = harness.store.snapshot()
work = state.Projects["project"].Works["work"]
if work.ResumeStage != "" {
t.Fatalf("recovered work still has resumeStage = %s", work.ResumeStage)
}
if work.State != WorkStateCompleted {
t.Fatalf("recovered work state = %s, want completed", work.State)
}
}
func TestAutoResumeOverrideFalseRequiresExplicitRestart(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 1)
disabled := false
harness.start("project", "workspace", &disabled)
harness.store.edit(func(state *ManagerState) {
project := state.Projects["project"]
project.Status = ProjectStatusRunning
state.Projects["project"] = project
})
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
state := harness.store.snapshot()
if state.Projects["project"].Status != ProjectStatusStopped {
t.Fatalf("project status = %s, want stopped", state.Projects["project"].Status)
}
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile 2: %v", err)
}
if harness.invoker.callCount() != 0 {
t.Fatalf("provider calls = %d, want 0", harness.invoker.callCount())
}
req := StartRequest{
CommandID: "start-2-project",
ProjectID: "project",
WorkspaceID: "workspace",
MilestoneID: "milestone",
WorkflowRevision: "workflow-r1",
ConfigRevision: "config-r1",
GrantRevision: "grant-r1",
}
if err := harness.manager.StartProject(context.Background(), req); err != nil {
t.Fatalf("StartProject: %v", err)
}
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile 3: %v", err)
}
if harness.store.snapshot().Projects["project"].Status != ProjectStatusCompleted {
t.Fatalf("explicitly restarted project did not complete")
}
}
func TestClaimProjectCASConflictReturnsCommittedDecision(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 1)
harness.start("project", "workspace", nil)
store := harness.store
wrappedStore := &casConflictStore{
store: store,
onConflict: func() {
store.edit(func(state *ManagerState) {
project := state.Projects["project"]
project.Status = ProjectStatusStopped
state.Projects["project"] = project
})
},
}
harness.manager.store = wrappedStore
claimed, err := harness.manager.claimProject(context.Background(), "project")
if err != nil {
t.Fatalf("claimProject: %v", err)
}
if claimed {
t.Fatalf("claimProject returned true for stopped project after CAS conflict")
}
}
func TestClaimIntegrationCASConflictReturnsCommittedDecision(t *testing.T) {
harness := newHarness(t, nil, 1)
store := harness.store
wrappedStore := &casConflictStore{
store: store,
onConflict: func() {
store.edit(func(state *ManagerState) {
state.IntegrationLeases["workspace"] = LeaseRecord{
OwnerID: "other-owner",
Token: "foreign",
ExpiresAt: time.Now().Add(time.Hour),
}
})
},
}
harness.manager.store = wrappedStore
claimed, err := harness.manager.claimIntegration(context.Background(), "workspace")
if err != nil {
t.Fatalf("claimIntegration: %v", err)
}
if claimed {
t.Fatalf("claimIntegration returned true when foreign lease committed during CAS conflict")
}
}
func TestWorkflowActivationCASConflictUsesCommittedState(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 1)
harness.start("project", "workspace", nil)
store := harness.store
wrappedStore := &casConflictStore{
store: store,
onConflict: func() {
store.edit(func(state *ManagerState) {
project := state.Projects["project"]
project.Status = ProjectStatusStopped
state.Projects["project"] = project
})
},
}
harness.manager.store = wrappedStore
active, err := harness.manager.observeWorkflows(context.Background())
if err != nil {
t.Fatalf("observeWorkflows: %v", err)
}
if len(active) != 0 {
t.Fatalf("observeWorkflows returned active projects = %v, want empty after CAS conflict to stopped", active)
}
}
func TestWorkflowActivationCASConflictUsesCommittedEventDecision(t *testing.T) {
snapshot := testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown))
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{"project": snapshot}, 1)
harness.start("project", "workspace", nil)
harness.store.edit(func(state *ManagerState) {
project := state.Projects["project"]
project.Status = ProjectStatusRunning
state.Projects["project"] = project
})
store := harness.store
wrappedStore := &casConflictStore{
store: store,
onConflict: func() {
store.edit(func(state *ManagerState) {
project := state.Projects["project"]
project.Status = ProjectStatusStarted
state.Projects["project"] = project
})
},
}
harness.manager.store = wrappedStore
active, err := harness.manager.observeWorkflows(context.Background())
if err != nil {
t.Fatalf("observeWorkflows: %v", err)
}
if len(active) != 1 || active[0] != "project" {
t.Fatalf("observeWorkflows active = %v, want [project]", active)
}
eventsEmitted := harness.events.snapshot()
var autoResumeCount int
var observedEvent *Event
for _, e := range eventsEmitted {
if e.Type == EventAutoResume {
autoResumeCount++
}
if e.Type == EventObserved {
eCopy := e
observedEvent = &eCopy
}
}
if autoResumeCount != 0 {
t.Fatalf("auto-resume event count = %d, want 0 on explicit-start winner", autoResumeCount)
}
if observedEvent == nil || observedEvent.CommandID == "" || observedEvent.WorkflowRevision == "" {
t.Fatalf("observed event missing committed identity: %#v", observedEvent)
}
}
type casConflictStore struct {
store *memoryStore
mu sync.Mutex
injected bool
onConflict func()
}
func (s *casConflictStore) Load(ctx context.Context) (ManagerState, StateRevision, error) {
return s.store.Load(ctx)
}
func (s *casConflictStore) CompareAndSwap(ctx context.Context, expected StateRevision, next ManagerState) (StateRevision, error) {
s.mu.Lock()
if !s.injected {
s.injected = true
if s.onConflict != nil {
s.onConflict()
}
s.mu.Unlock()
return "", ErrRevisionConflict
}
s.mu.Unlock()
return s.store.CompareAndSwap(ctx, expected, next)
}