iop/packages/go/agenttask/manager_integration_test.go

453 lines
16 KiB
Go

package agenttask
import (
"context"
"errors"
"fmt"
"strings"
"sync"
"testing"
"time"
)
var errInjectedEventSink = errors.New("injected event sink failure")
type observedEventDelivery struct {
Event Event
Delivery EventDelivery
}
type stateObservingEventSink struct {
store *memoryStore
failType EventType
failOnce bool
mu sync.Mutex
observations []observedEventDelivery
}
func (s *stateObservingEventSink) Emit(_ context.Context, event Event) error {
state := s.store.snapshot()
delivery, ok := state.PendingEvents[event.EventID]
if !ok {
return fmt.Errorf("event %q reached sink without a pending delivery", event.EventID)
}
if !sameLogicalEvent(delivery.Event, event) ||
!delivery.Event.Timestamp.Equal(event.Timestamp) {
return fmt.Errorf("event %q differs from its pending delivery", event.EventID)
}
s.mu.Lock()
s.observations = append(s.observations, observedEventDelivery{
Event: event,
Delivery: cloneEventDelivery(delivery),
})
shouldFail := s.failOnce && event.Type == s.failType
if shouldFail {
s.failOnce = false
}
s.mu.Unlock()
if shouldFail {
return errInjectedEventSink
}
return nil
}
func (s *stateObservingEventSink) snapshot() []observedEventDelivery {
s.mu.Lock()
defer s.mu.Unlock()
out := make([]observedEventDelivery, len(s.observations))
for index, observation := range s.observations {
out[index] = observedEventDelivery{
Event: observation.Event,
Delivery: cloneEventDelivery(observation.Delivery),
}
}
return out
}
func findObservedDelivery(
t *testing.T,
observations []observedEventDelivery,
eventType EventType,
predicate func(EventDelivery) bool,
) EventDelivery {
t.Helper()
for _, observation := range observations {
if observation.Event.Type == eventType && predicate(observation.Delivery) {
return observation.Delivery
}
}
t.Fatalf("missing committed delivery for event type %q", eventType)
return EventDelivery{}
}
func TestManagerEventDeliveryUsesCommittedEvidence(t *testing.T) {
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{
"project": testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown)),
}, 1)
harness.reviewer.sequences["work"] = []ReviewVerdict{
ReviewVerdictWarn,
ReviewVerdictPass,
}
sink := &stateObservingEventSink{store: harness.store}
harness.manager.events = sink
harness.start("project", "workspace", nil)
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
observations := sink.snapshot()
dependency := findObservedDelivery(t, observations, EventDependencyReady, func(delivery EventDelivery) bool {
return delivery.Work != nil && delivery.Work.Unit.ID == "work"
})
if dependency.Work.State != WorkStateReady ||
dependency.Work.Attempt != 1 ||
dependency.Work.AttemptID == "" ||
dependency.Work.DispatchOrdinal == 0 {
t.Fatalf("dependency delivery used uncommitted work: %+v", dependency.Work)
}
review := findObservedDelivery(t, observations, EventReviewResult, func(delivery EventDelivery) bool {
return delivery.Work != nil &&
delivery.Work.Review != nil &&
delivery.Work.Review.Verdict == ReviewVerdictWarn
})
if review.Work.State != WorkStateReviewing ||
review.Work.Review.AttemptID != review.Work.AttemptID {
t.Fatalf("review delivery used stale evidence: %+v", review.Work)
}
followup := findObservedDelivery(t, observations, EventFollowup, func(delivery EventDelivery) bool {
return delivery.Work != nil && delivery.Work.Attempt == 2
})
if followup.Work.State != WorkStateReady ||
followup.Event.AttemptID != followup.Work.AttemptID {
t.Fatalf("follow-up delivery did not preserve the next attempt: %+v", followup)
}
terminal := findObservedDelivery(t, observations, EventIntegrationResult, func(delivery EventDelivery) bool {
return delivery.Work != nil && delivery.Work.State == WorkStateCompleted
})
if terminal.Work.Integration == nil ||
terminal.Work.Integration.Outcome != IntegrationOutcomeIntegrated ||
!terminal.Work.CompletionVerified {
t.Fatalf("terminal delivery used incomplete integration evidence: %+v", terminal.Work)
}
if pending := harness.store.snapshot().PendingEvents; len(pending) != 0 {
t.Fatalf("successful deliveries remain pending: %+v", pending)
}
}
func TestManagerEventDeliveryRecoversSinkFailure(t *testing.T) {
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{
"project": testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown)),
}, 1)
clock := &advancingClock{now: time.Date(2026, 7, 30, 19, 0, 0, 0, time.UTC)}
harness.manager.clock = clock
harness.start("project", "workspace", nil)
failing := &stateObservingEventSink{
store: harness.store,
failType: EventDependencyReady,
failOnce: true,
}
harness.manager.events = failing
err := harness.manager.Reconcile(context.Background())
if !errors.Is(err, errInjectedEventSink) {
t.Fatalf("Reconcile error = %v, want injected sink failure", err)
}
state := harness.store.snapshot()
if len(state.PendingEvents) != 1 {
t.Fatalf("pending deliveries = %d, want 1", len(state.PendingEvents))
}
var pending EventDelivery
for _, delivery := range state.PendingEvents {
pending = delivery
}
if pending.Event.Type != EventDependencyReady ||
pending.Work == nil ||
pending.Work.State != WorkStateReady {
t.Fatalf("failed delivery did not retain exact committed evidence: %+v", pending)
}
pendingID := pending.Event.EventID
pendingTime := pending.Event.Timestamp
clock.Advance(2 * time.Hour)
recoveredSink := &stateObservingEventSink{store: harness.store}
restarted, newErr := NewManager(
harness.manager.config,
clock,
harness.store,
harness.workflow,
harness.manager.selector,
harness.isolation,
harness.invoker,
harness.recovery,
harness.evidence,
harness.reviewer,
harness.integrator,
recoveredSink,
)
if newErr != nil {
t.Fatalf("restart NewManager: %v", newErr)
}
if err := restarted.Reconcile(context.Background()); err != nil {
t.Fatalf("restart Reconcile: %v", err)
}
delivered := 0
for _, observation := range recoveredSink.snapshot() {
if observation.Event.EventID != pendingID {
continue
}
delivered++
if !observation.Event.Timestamp.Equal(pendingTime) {
t.Fatalf(
"recovered timestamp = %s, want %s",
observation.Event.Timestamp,
pendingTime,
)
}
}
if delivered != 1 {
t.Fatalf("recovered pending delivery count = %d, want 1", delivered)
}
if pending := harness.store.snapshot().PendingEvents; len(pending) != 0 {
t.Fatalf("recovered deliveries remain pending: %+v", pending)
}
}
func TestManagerS03S16MultiProjectManualResumeAndParallelTrace(t *testing.T) {
projectA := testSnapshot(
"project-a", "shared-workspace",
testUnit("a-disjoint", WriteSetDisjoint),
testUnit("a-overlap", WriteSetOverlap),
)
projectB := testSnapshot(
"project-b", "workspace-b",
testUnit("b-unknown", WriteSetUnknown),
)
unselected := testSnapshot(
"unselected", "workspace-unselected",
testUnit("must-not-run", WriteSetUnknown),
)
harness := newHarness(
t,
map[ProjectID]ProjectWorkflowSnapshot{
"project-a": projectA, "project-b": projectB, "unselected": unselected,
},
3,
)
for _, workID := range []WorkUnitID{"a-disjoint", "a-overlap", "b-unknown"} {
harness.invoker.delays[workID] = 20 * time.Millisecond
}
harness.start("project-a", "shared-workspace", nil)
harness.start("project-b", "workspace-b", nil)
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
state := harness.store.snapshot()
if state.Projects["project-a"].Status != ProjectStatusCompleted ||
state.Projects["project-b"].Status != ProjectStatusCompleted {
t.Fatalf(
"manual projects = %s/%s",
state.Projects["project-a"].Status, state.Projects["project-b"].Status,
)
}
if state.Projects["unselected"].Status != ProjectStatusObserved {
t.Fatalf("unselected project status = %s", state.Projects["unselected"].Status)
}
if harness.invoker.callCount() != 3 {
t.Fatalf("invocations = %d, want selected work only", harness.invoker.callCount())
}
if harness.invoker.maxConcurrency() < 2 {
t.Fatalf("max concurrency = %d, want isolated parallel dispatch", harness.invoker.maxConcurrency())
}
trace := make([]string, 0)
for _, event := range harness.events.snapshot() {
if event.Type != EventDispatchStarted && event.Type != EventIntegrationResult {
continue
}
trace = append(trace, fmt.Sprintf(
"%s:%s:%s:%d:%s:%s",
event.Type, event.ProjectID, event.WorkUnitID, event.Ordinal,
event.WriteSetKind, event.IsolationMode,
))
if event.WorkUnitID == "must-not-run" {
t.Fatalf("unselected work appeared in execution trace: %v", trace)
}
}
joined := strings.Join(trace, "\n")
for _, expected := range []string{"a-disjoint", "a-overlap", "b-unknown"} {
if !strings.Contains(joined, expected) {
t.Fatalf("trace missing %s:\n%s", expected, joined)
}
}
t.Logf(
"S03 trace: unselected_invocations=0 manual_projects=2 terminal=2; "+
"S16 trace: explicit_dependency_only=true isolated_parallel_max=%d integration_ordinals=%v\n%s",
harness.invoker.maxConcurrency(), harness.integrator.ordinals(), joined,
)
}
func TestManagerWorkflowEvidenceGateAndPiRepair(t *testing.T) {
t.Run("complete evidence invokes review", func(t *testing.T) {
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{
"project": testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown)),
}, 1)
harness.start("project", "workspace", nil)
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
if harness.reviewer.callCount() != 1 {
t.Fatalf("review calls = %d, want 1", harness.reviewer.callCount())
}
})
for _, test := range []struct {
name string
configure func(*managerHarness)
wantCode BlockerCode
wantRepair int
}{
{
name: "placeholder from another provider never repairs or reviews",
configure: func(harness *managerHarness) {
harness.evidence.observeFunc = func(request WorkflowEvidenceRequest, _ int) (ArtifactEvidence, error) {
return ArtifactEvidence{
Active: true, Identity: artifactIdentity(request.Project, request.Work, request.Submission),
Completeness: ArtifactPlaceholder,
}, nil
}
},
wantCode: BlockerEvidenceRepairDenied,
},
{
name: "wrong active artifact identity never invokes review",
configure: func(harness *managerHarness) {
harness.evidence.observeFunc = func(request WorkflowEvidenceRequest, _ int) (ArtifactEvidence, error) {
identity := artifactIdentity(request.Project, request.Work, request.Submission)
identity.AttemptID = "stale-attempt"
return ArtifactEvidence{Active: true, Identity: identity, Completeness: ArtifactComplete}, nil
}
},
wantCode: BlockerArtifactMismatch,
},
{
name: "Pi repair rejects a stale native locator",
configure: func(harness *managerHarness) {
harness.manager.selector = fakeSelector{capacity: 1, providerID: "pi"}
harness.invoker.locators["work"] = []LocatorRecord{{
Kind: LocatorSession, Opaque: "native-session", Revision: "session-r1",
ProjectID: "project", WorkspaceID: "workspace", WorkUnitID: "work", AttemptID: attemptID("work", 1),
}}
harness.evidence.observeFunc = func(request WorkflowEvidenceRequest, _ int) (ArtifactEvidence, error) {
identity := artifactIdentity(request.Project, request.Work, request.Submission)
return ArtifactEvidence{
Active: true, Identity: identity, Completeness: ArtifactPlaceholder,
RepairIntent: &EvidenceRepairIntent{
Identity: identity, DispatchOrdinal: request.Work.DispatchOrdinal,
NativeLocator: LocatorRecord{
Kind: LocatorSession, Opaque: "stale-session", Revision: "session-r1",
ProjectID: "project", WorkspaceID: "workspace", WorkUnitID: "work", AttemptID: attemptID("work", 1),
},
},
}, nil
}
},
wantCode: BlockerEvidenceRepairDenied,
},
} {
t.Run(test.name, func(t *testing.T) {
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{
"project": testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown)),
}, 1)
test.configure(harness)
harness.start("project", "workspace", nil)
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
work := harness.store.snapshot().Projects["project"].Works["work"]
if work.Blocker == nil || work.Blocker.Code != test.wantCode {
t.Fatalf("work blocker = %#v, want %q", work.Blocker, test.wantCode)
}
if harness.reviewer.callCount() != 0 {
t.Fatalf("review calls = %d, want 0", harness.reviewer.callCount())
}
_, repairs := harness.evidence.counts()
if repairs != test.wantRepair {
t.Fatalf("repair calls = %d, want %d", repairs, test.wantRepair)
}
})
}
t.Run("Pi repair rematches before review", func(t *testing.T) {
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{
"project": testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown)),
}, 1)
harness.manager.selector = fakeSelector{capacity: 1, providerID: "pi"}
harness.invoker.locators["work"] = []LocatorRecord{{
Kind: LocatorSession, Opaque: "native-session", Revision: "session-r1",
ProjectID: "project", WorkspaceID: "workspace", WorkUnitID: "work", AttemptID: attemptID("work", 1),
}}
harness.evidence.observeFunc = func(request WorkflowEvidenceRequest, call int) (ArtifactEvidence, error) {
identity := artifactIdentity(request.Project, request.Work, request.Submission)
if call == 0 {
return ArtifactEvidence{
Active: true, Identity: identity, Completeness: ArtifactPlaceholder,
RepairIntent: &EvidenceRepairIntent{
Identity: identity, DispatchOrdinal: request.Work.DispatchOrdinal,
NativeLocator: request.Work.Locators[LocatorSession],
},
}, nil
}
return ArtifactEvidence{Active: true, Identity: identity, Completeness: ArtifactComplete}, nil
}
harness.start("project", "workspace", nil)
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
observations, repairs := harness.evidence.counts()
if observations != 2 || repairs != 1 {
t.Fatalf("observe/repair calls = %d/%d, want 2/1", observations, repairs)
}
if harness.reviewer.callCount() != 1 {
t.Fatalf("review calls = %d, want 1 after fresh rematch", harness.reviewer.callCount())
}
})
t.Run("Pi repair without a fresh completed match never invokes review", func(t *testing.T) {
harness := newHarness(t, map[ProjectID]ProjectWorkflowSnapshot{
"project": testSnapshot("project", "workspace", testUnit("work", WriteSetUnknown)),
}, 1)
harness.manager.selector = fakeSelector{capacity: 1, providerID: "pi"}
harness.invoker.locators["work"] = []LocatorRecord{{
Kind: LocatorSession, Opaque: "native-session", Revision: "session-r1",
ProjectID: "project", WorkspaceID: "workspace", WorkUnitID: "work", AttemptID: attemptID("work", 1),
}}
harness.evidence.observeFunc = func(request WorkflowEvidenceRequest, call int) (ArtifactEvidence, error) {
identity := artifactIdentity(request.Project, request.Work, request.Submission)
if call == 0 {
return ArtifactEvidence{
Active: true, Identity: identity, Completeness: ArtifactPlaceholder,
RepairIntent: &EvidenceRepairIntent{
Identity: identity, DispatchOrdinal: request.Work.DispatchOrdinal,
NativeLocator: request.Work.Locators[LocatorSession],
},
}, nil
}
return ArtifactEvidence{Active: true, Identity: identity, Completeness: ArtifactPlaceholder}, nil
}
harness.start("project", "workspace", nil)
if err := harness.manager.Reconcile(context.Background()); err != nil {
t.Fatalf("Reconcile: %v", err)
}
work := harness.store.snapshot().Projects["project"].Works["work"]
if work.Blocker == nil || work.Blocker.Code != BlockerSubmissionIncomplete {
t.Fatalf("work blocker = %#v, want %q", work.Blocker, BlockerSubmissionIncomplete)
}
observations, repairs := harness.evidence.counts()
if observations != 2 || repairs != 1 {
t.Fatalf("observe/repair calls = %d/%d, want 2/1", observations, repairs)
}
if harness.reviewer.callCount() != 0 {
t.Fatalf("review calls = %d, want 0 without a fresh completed match", harness.reviewer.callCount())
}
})
}