453 lines
16 KiB
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())
|
|
}
|
|
})
|
|
}
|