diff --git a/docs/changes/unreleased/1733-session-slow-tests-control-their-clocks.md b/docs/changes/unreleased/1733-session-slow-tests-control-their-clocks.md new file mode 100644 index 000000000..d71a2309a --- /dev/null +++ b/docs/changes/unreleased/1733-session-slow-tests-control-their-clocks.md @@ -0,0 +1,9 @@ +--- +kind: internal +title: twelve slow session tests wait on acknowledgements instead of wall clocks +pr: 1733 +surface: [engine] +invalidates: + - "The slowest internal/session tests (the watch tickers, the team loop breaker, the slow phase listener, the young-bash steer grace, the stalled checker, the task-baseline attribution trio) took about 123 seconds between them, waiting on real time. They now drive the same real commands, queues and landings through acknowledgements and explicit clock advances, about 4 seconds in all. The two plandb CLI tests that run a real loop through bash (about 35 seconds each) are unchanged and remain the slowest in the package." + - "jobRegistry.watchTickWait, Agent.steerAfter, Config.teamWatchManual and Config.auditTimeout are new private test seams. Each is nil or false in production and falls through to the call it replaced, so no window, default or limit moved." +--- diff --git a/internal/session/callwindow.go b/internal/session/callwindow.go index 9d8155ebd..7f1b1d44b 100644 --- a/internal/session/callwindow.go +++ b/internal/session/callwindow.go @@ -51,6 +51,9 @@ type callWindow struct { // call, which is the one effort a person's pinned rung yields to (the // adapter's WithRequiredReasoningEffort states the rule). answer bool + // timeout is a test-only context-clock seam. A nil seam keeps the real + // deadline, so production windows and the information told stay unchanged. + timeout func(context.Context, time.Duration) (context.Context, context.CancelFunc) } // callWindowKey is the context key [callWindow] rides under. @@ -60,7 +63,11 @@ type callWindowKey struct{} // told. It is the one way this package opens a told window, and the checker's // calls are held to it by a law (callwindow_law_test.go). func openCallWindow(ctx context.Context, bound time.Duration, window callWindow) (context.Context, context.CancelFunc) { - ctx, cancel := context.WithTimeout(ctx, bound) + timeout := context.WithTimeout + if window.timeout != nil { + timeout = window.timeout + } + ctx, cancel := timeout(ctx, bound) return context.WithValue(ctx, callWindowKey{}, window), cancel } diff --git a/internal/session/controlled_waits_test.go b/internal/session/controlled_waits_test.go new file mode 100644 index 000000000..69892d0d2 --- /dev/null +++ b/internal/session/controlled_waits_test.go @@ -0,0 +1,146 @@ +package session + +import ( + "context" + "os" + "path/filepath" + "sync/atomic" + "testing" + "time" + + "github.com/Agent-Field/codeaf/internal/exec/bare" +) + +// awaitTestCompletion waits for a real acknowledgement. Its failure guard +// follows the suite's budget, so repository I/O on a loaded box cannot turn an +// otherwise correct landing into a race against an unrelated fixed window. +func awaitTestCompletion(t *testing.T, done <-chan struct{}, what string) { + t.Helper() + timer := time.NewTimer(awaitPatience(t)) + defer timer.Stop() + select { + case <-done: + case <-timer.C: + t.Fatalf("no acknowledgement of %s arrived within the suite's patience", what) + } +} + +// controlledWatch drives the real command, delta, notification and stopping +// loop. Only the wait between ticks is replaced, so a completed tick is a fact +// the test can read before changing the command's input. +type controlledWatch struct { + completed chan struct{} + next chan struct{} +} + +func controlWatch(agent *Agent) *controlledWatch { + clock := &controlledWatch{completed: make(chan struct{}, 1), next: make(chan struct{})} + agent.jobs.watchTickWait = func(ctx context.Context) { + select { + case clock.completed <- struct{}{}: + case <-ctx.Done(): + return + } + select { + case <-clock.next: + case <-ctx.Done(): + } + } + return clock +} + +func (clock *controlledWatch) completedTick(t *testing.T, watched *job) { + t.Helper() + select { + case <-clock.completed: + case <-watched.done: + case <-time.After(awaitPatience(t)): + t.Fatal("the watch did not complete its commanded tick") + } +} + +func (clock *controlledWatch) tick(t *testing.T, watched *job) { + t.Helper() + select { + case clock.next <- struct{}{}: + case <-watched.done: + t.Fatal("the watch stopped before its next commanded tick") + case <-time.After(awaitPatience(t)): + t.Fatal("the watch did not accept its next commanded tick") + } + clock.completedTick(t, watched) +} + +// heldSteerBash keeps a real command alive until the test releases it. Its +// lifetime follows the steer acknowledgement rather than a guessed sleep that +// could end before a loaded machine reaches the assertion. +func heldSteerBash(t *testing.T, ending string) (string, func()) { + t.Helper() + release := filepath.Join(t.TempDir(), "release-bash") + command := "while [ ! -f " + shellQuoted(release) + " ]; do sleep 0.01; done; echo " + shellQuoted(ending) + return command, func() { + if err := os.WriteFile(release, nil, 0o600); err != nil { + t.Fatalf("release the held bash: %v", err) + } + } +} + +// controlledSteer holds both clocks: a command stays young until advanced, +// and a scheduled callback runs only when fired. The real scheduling delay is +// retained so the test can still assert that the grace was armed at its bound. +type controlledSteer struct { + delay time.Duration + fire func() +} + +func controlSteer(agent *Agent) *controlledSteer { + clock := &controlledSteer{} + agent.mu.Lock() + agent.steerAge = func(*bare.BashCall) time.Duration { return 0 } + agent.steerAfter = func(delay time.Duration, fire func()) *time.Timer { + clock.delay, clock.fire = delay, fire + // The timer is only a stop handle; no wall-clock callback can race the + // command-age advance or the deliberate firing made by this test. + return time.NewTimer(24 * time.Hour) + } + agent.mu.Unlock() + return clock +} + +// controlledCheckCall exposes the share the product assigned and the expiry +// that a test fires only once its stalled provider call has actually started. +type controlledCheckCall struct { + bound time.Duration + expire func() +} + +// controlledCheckContext retains a told deadline without arming a real timer. +// Its explicit expiry reports DeadlineExceeded, while parent cancellation and +// cleanup keep their ordinary cancellation meaning. +type controlledCheckContext struct { + context.Context + until time.Time + expired atomic.Bool +} + +func (ctx *controlledCheckContext) Deadline() (time.Time, bool) { return ctx.until, true } + +func (ctx *controlledCheckContext) Err() error { + err := ctx.Context.Err() + if err != nil && ctx.expired.Load() { + return context.DeadlineExceeded + } + return err +} + +func controlCheck(calls chan<- controlledCheckCall) func(context.Context, time.Duration) (context.Context, context.CancelFunc) { + return func(parent context.Context, bound time.Duration) (context.Context, context.CancelFunc) { + ctx, cancel := context.WithCancel(parent) + opened := &controlledCheckContext{Context: ctx, until: time.Now().Add(bound)} + calls <- controlledCheckCall{bound: bound, expire: func() { + opened.expired.Store(true) + cancel() + }} + return opened, cancel + } +} diff --git a/internal/session/jobs.go b/internal/session/jobs.go index a680d9434..d823b3a48 100644 --- a/internal/session/jobs.go +++ b/internal/session/jobs.go @@ -559,6 +559,9 @@ type jobRegistry struct { retentionShut bool // The stage callback is a per-registry test seam; production leaves it nil. retentionStage func(string) + // watchTickWait lets a test acknowledge a completed tick and release the + // next one. Production leaves it nil and waits on the watch's real ticker. + watchTickWait func(context.Context) } func newJobRegistry(workspace string, place Place, notify func(string), watch ...func(string, string, bool)) *jobRegistry { diff --git a/internal/session/loop_speed_test.go b/internal/session/loop_speed_test.go index 5656a44a1..1bbd2502b 100644 --- a/internal/session/loop_speed_test.go +++ b/internal/session/loop_speed_test.go @@ -479,23 +479,33 @@ func TestAMarksReadingRidesBesideTheStepAndCutsItOnIndependentParts(t *testing.T // `phase.done()`, paid a whole Update-and-View cycle before the engine goroutine // got its call back. On every call, not only a cancelled one. // -// The listener here takes as long as every other auxiliary in this file, and the -// two figures the law is about are the same two. It passes because the news is -// left on a desk (sidecar.go) rather than carried. +// The listener stays blocked until the turn finishes, and the two figures the +// law is about are still asserted. It passes because the news is left on a +// desk (sidecar.go) rather than carried; joining that desk costs no timed wait +// once the test releases the listener's backlog. func TestASlowListenerNeverHoldsTheTurnThatIsTellingIt(t *testing.T) { // The listener is released the instant the test ends, so a desk still working // through a backlog never outlives the thing it was describing. release := make(chan struct{}) - t.Cleanup(func() { close(release) }) + var once sync.Once + free := func() { once.Do(func() { close(release) }) } + defer free() + t.Cleanup(free) + entered := make(chan struct{}, 1) var told int64 previous := OnPhaseNews(func(PhaseNews) { atomic.AddInt64(&told, 1) select { - case <-release: - case <-time.After(paceSlowReading): + case entered <- struct{}{}: + default: } + <-release + }) + t.Cleanup(func() { + free() + phaseDesk.settled() + OnPhaseNews(previous) }) - t.Cleanup(func() { OnPhaseNews(previous) }) rounds := 3 completer := &paceCompleter{ @@ -517,6 +527,10 @@ func TestASlowListenerNeverHoldsTheTurnThatIsTellingIt(t *testing.T) { "want at most %dms — a phase was posted down the turn's own stack", pace.StepGapMS, paceLawBudget.Milliseconds()) } + awaitTestCompletion(t, entered, "the blocked listener receiving a phase") + // The turn has finished while the listener is still blocked. Release its + // backlog before joining it, rather than paying a timeout for every phase. + free() phaseDesk.settled() // AND THE LISTENER REALLY WAS TOLD. A turn that posted nothing would pass this // law by saying nothing, which is the other way to break the status line. diff --git a/internal/session/session.go b/internal/session/session.go index f021f80d9..72bb398f3 100644 --- a/internal/session/session.go +++ b/internal/session/session.go @@ -1369,6 +1369,10 @@ type Config struct { // person reads when the window runs out without waiting five real minutes for // it. auditWindow time.Duration + // auditTimeout lets a test expire a checker's context after advancing the + // checking clock. Production leaves it nil and uses context.WithTimeout + // with the same window and call share. + auditTimeout func(context.Context, time.Duration) (context.Context, context.CancelFunc) // clock is THE AGENT'S ONE READING OF THE WORLD'S TIME, and it is UNEXPORTED // AND FOR TESTS ONLY ([Agent.now]). The product's answer is [time.Now]. @@ -1383,6 +1387,9 @@ type Config struct { // other. The other caller today is the trail that records a request's own // length (task_calltrail.go). clock func() time.Time + // teamWatchManual leaves traffic-clock advances to tests calling + // teamWatchTick. Production leaves it false and starts the real ticker. + teamWatchManual bool // AskConsent says somebody is watching this agent's events and will answer // an EventConsentRequest with [Agent.ResolveConsent]. @@ -3025,6 +3032,9 @@ type Agent struct { // steerAge is the foreground-command age seam used by steer tests. A nil // seam reads the process's real start through [bare.BashCall.RunningFor]. steerAge func(*bare.BashCall) time.Duration + // steerAfter lets a test hold and fire the scheduled second look itself. + // Production leaves it nil and uses time.AfterFunc at the same grace. + steerAfter func(time.Duration, func()) *time.Timer // ambient is periodic watch news that must wait for a TURN boundary. // // It is separate from steering because a step boundary is not a turn diff --git a/internal/session/steer_grace.go b/internal/session/steer_grace.go index 8a25aa342..56e1c9e77 100644 --- a/internal/session/steer_grace.go +++ b/internal/session/steer_grace.go @@ -132,7 +132,11 @@ func (a *Agent) armSteerGraceLocked() { } a.stopSteerGraceLocked() watch := &steerWatch{turn: a.turnSeq, calls: young} - watch.timer = time.AfterFunc(wait+steerGraceMargin, func() { a.steerGraceFired(watch) }) + after := time.AfterFunc + if a.steerAfter != nil { + after = a.steerAfter + } + watch.timer = after(wait+steerGraceMargin, func() { a.steerGraceFired(watch) }) a.steerGrace = watch } diff --git a/internal/session/steer_grace_test.go b/internal/session/steer_grace_test.go index b11344151..264155fe1 100644 --- a/internal/session/steer_grace_test.go +++ b/internal/session/steer_grace_test.go @@ -13,7 +13,6 @@ import ( "context" "strings" "sync" - "sync/atomic" "testing" "time" @@ -33,13 +32,14 @@ func advanceSteerAge(agent *Agent, age time.Duration) { // not at the command's own ending, and not at the background clock. func TestASteerLandsWhenAYoungBashCrossesTheGrace(t *testing.T) { t.Parallel() - var reached atomic.Int64 + command, release := heldSteerBash(t, "grace-bash-finished") + reached := make(chan struct{}) completer := &scriptedCompleter{steps: []step{ func(context.Context, []ai.Message) (*ai.Response, error) { - return toolResponse("grace-bash", "bash", `{"command":"sleep 7; echo grace-bash-finished"}`), nil + return bashCall("grace-bash", command)(context.Background(), nil) }, func(context.Context, []ai.Message) (*ai.Response, error) { - reached.Store(time.Now().UnixNano()) + close(reached) return textResponse("2026-09-05, and the suite is still running"), nil }, func(context.Context, []ai.Message) (*ai.Response, error) { @@ -53,21 +53,32 @@ func TestASteerLandsWhenAYoungBashCrossesTheGrace(t *testing.T) { // steer waiting for it. Twenty seconds stands in for livechat's thirty and // is longer than anything this test does. agent, _ := newTestAgent(t, completer, func(config *Config) { config.BashBackgroundAfterSeconds = 20 }) + clock := controlSteer(agent) turn := mustSubmit(t, agent, "run the suite") waitFor(t, "the foreground bash to start", func() bool { return len(agent.inFlightBash.snapshot()) == 1 }) calls := agent.inFlightBash.snapshot() - if age := calls[0].RunningFor(); age >= steerBashAge { + if age := agent.steerBashRunningFor(calls[0]); age >= steerBashAge { t.Fatalf("the bash was already %s old, so this is not the young case", age) } - sent := time.Now() steered := mustSteer(t, agent, "while that runs, what is today's date") - - waitFor(t, "the steer to reach the model", func() bool { return reached.Load() != 0 }) - waited := time.Unix(0, reached.Load()).Sub(sent) - t.Logf("the steer reached the model %s after it was sent", waited) - if bound := steerBashAge + 2*time.Second; waited > bound { - t.Fatalf("the steer reached the model %s after it was sent, want it inside %s", waited, bound) + agent.mu.Lock() + watch := agent.steerGrace + agent.mu.Unlock() + if watch == nil || clock.fire == nil { + t.Fatal("no second look was armed for the young bash") + } + if want := steerBashAge + steerGraceMargin; clock.delay != want { + t.Fatalf("the second look was scheduled after %s, want %s", clock.delay, want) + } + if agent.jobs.find(1) != nil { + t.Fatal("the young bash was adopted before its grace") } + // The command is held until its eventual exit is asked for below. Advancing + // its age and firing the actual scheduled callback proves the second look + // without racing the process's lifetime against a wall-clock window. + advanceSteerAge(agent, steerBashAge+steerGraceMargin) + clock.fire() + awaitTestCompletion(t, reached, "the delayed steer reaching the model") second := completer.request(1) if got, want := userLines(second), []string{"run the suite", "while that runs, what is today's date"}; !equalStrings(got, want) { @@ -85,6 +96,7 @@ func TestASteerLandsWhenAYoungBashCrossesTheGrace(t *testing.T) { } // ONE PROCESS, ADOPTED ONCE, AND ITS ENDING STILL ARRIVES. + release() waitFor(t, "the adopted job's exit note", func() bool { return notesContain(agent, "grace-bash-finished") }) waitFor(t, "the owed exit request", func() bool { return completer.requests() >= 3 }) if got := strings.Join(userLines(completer.request(2)), "\n"); !strings.Contains(got, "job 1 exited 0") { diff --git a/internal/session/task_audit.go b/internal/session/task_audit.go index 9b9da009c..41b40c4ed 100644 --- a/internal/session/task_audit.go +++ b/internal/session/task_audit.go @@ -1289,7 +1289,7 @@ func (a *Agent) auditOnce(ctx context.Context, node *TaskNode, tree taskTree, gr // through two thirty-second shares without a word (#941). Opened this way, // every request the checker makes under it carries the time it has left, and // the adapter's effort ladder sizes its thinking to fit (callwindow.go). - auditCtx, done := openCallWindow(ctx, bound, callWindow{}) + auditCtx, done := openCallWindow(ctx, bound, callWindow{timeout: a.config.auditTimeout}) defer done() fmt.Fprintf(log, "audit: verifying against the acceptance\n") diff --git a/internal/session/task_baseline_test.go b/internal/session/task_baseline_test.go index 55398cfe7..dcb93ac01 100644 --- a/internal/session/task_baseline_test.go +++ b/internal/session/task_baseline_test.go @@ -23,14 +23,14 @@ import ( // own checkout has the failure fixed without a commit, so only a detached read // of the branch can still find the committed red. func TestAReadingOfTheBaseSaysWhichChecksWereAlreadyFailing(t *testing.T) { - repo := newGoModuleRepo(t) - test := filepath.Join(repo, "base_test.go") - writeFile(t, test, "package taskaudit\n\nimport \"testing\"\n\nfunc TestBaseWasRed(t *testing.T) { t.Fatal(\"old red\") }\n") + repo := newTestRepo(t) + test := filepath.Join(repo, "check.sh") + writeFile(t, test, "printf 'old red\\n'; exit 1\n") mustGit(t, repo, "add", "-A") mustGit(t, repo, "-c", "user.name=t", "-c", "user.email=t@t", "commit", "-m", "old red") - writeFile(t, test, "package taskaudit\n\nimport \"testing\"\n\nfunc TestBaseWasRed(t *testing.T) {}\n") + writeFile(t, test, "exit 0\n") - command := "go test ./..." + command := "sh check.sh" photograph := (&Agent{}).baseChecksFor(context.Background(), taskTree{ dir: repo, root: repo, branch: "work", ground: repo, seal: baselineCommitForTest(t, repo), }, []string{command}) @@ -292,16 +292,17 @@ func proposeTaskWithAcceptance(title, brief, acceptance, ground string, checks [ // the real task door. The task writes an unrelated deliverable over a committed // failing suite, and the finished report carries the shared session sentence. func TestAPreExistingRedDoesNotStandBetweenTheWorkAndItsLanding(t *testing.T) { - repo := newGoModuleRepo(t) - writeFile(t, filepath.Join(repo, "old_red_test.go"), - "package taskaudit\n\nimport \"testing\"\n\nfunc TestOldRed(t *testing.T) { t.Fatal(\"already red\") }\n") + repo := newTestRepo(t) + // A real failing shell check proves attribution without rebuilding a + // disposable Go module; both detached and landing trees still run it. + writeFile(t, filepath.Join(repo, "check.sh"), "printf 'already red\n'; exit 1\n") mustGit(t, repo, "add", "-A") mustGit(t, repo, "-c", "user.name=t", "-c", "user.email=t@t", "commit", "-m", "old red") t.Setenv("HOME", t.TempDir()) completer := &routedCompleter{ parent: []step{ - proposeTaskWithAcceptance("Write the note", "write note.txt", "note.txt exists and `go test ./...` passes", repo, []string{"go test ./..."}), + proposeTaskWithAcceptance("Write the note", "write note.txt", "note.txt exists and `sh check.sh` passes", repo, []string{"sh check.sh"}), finalText("handed off"), }, child: []step{ @@ -321,38 +322,43 @@ func TestAPreExistingRedDoesNotStandBetweenTheWorkAndItsLanding(t *testing.T) { if node == nil { t.Fatalf("no node was admitted; events = %#v", events) } - waitDoneNode(t, node) + awaitTestCompletion(t, node.done, "the task over the pre-existing red to land") notice := node.notice() if notice.State != TaskDone { t.Fatalf("state = %q, report = %q, want done", notice.State, notice.Report) } - want := "1 check was already failing before this work; that does not show the requested result works: go test ./..." + want := "1 check was already failing before this work; that does not show the requested result works: sh check.sh" if !strings.Contains(notice.Report, want) { t.Fatalf("report = %q, want %q", notice.Report, want) } } // TestAGreenCheckTheWorkTurnedRedIsTheWorksOwn proves C6. The base is clean, -// the child's only file introduces a failing test, and the checker's refusal +// the child's only file turns the check red, and the checker's refusal // remains a refusal without borrowing the old-red sentence. func TestAGreenCheckTheWorkTurnedRedIsTheWorksOwn(t *testing.T) { - repo := newGoModuleRepo(t) + repo := newTestRepo(t) + // The base check passes until the worker creates the regression marker. + // Its failure uses the same test-name format the product extracts from Go. + writeFile(t, filepath.Join(repo, "check.sh"), + "if [ -f work.txt ]; then printf '%s\n' '--- FAIL: TestWorkBreaksGreen (0.00s)' 'FAIL'; exit 1; fi\n") + mustGit(t, repo, "add", "check.sh") + mustGit(t, repo, "-c", "user.name=t", "-c", "user.email=t@t", "commit", "-m", "green check") t.Setenv("HOME", t.TempDir()) completer := &routedCompleter{ parent: []step{ - proposeTaskWithAcceptance("Add the check", "write work_test.go", "work_test.go exists and `go test ./...` passes", repo, []string{"go test ./..."}), + proposeTaskWithAcceptance("Add the check", "write work.txt", "work.txt exists and `sh check.sh` passes", repo, []string{"sh check.sh"}), finalText("handed off"), }, child: []step{ - writeCall("write-red", "work_test.go", - "package taskaudit\n\nimport \"testing\"\n\nfunc TestWorkBreaksGreen(t *testing.T) { t.Fatal(\"new red\") }\n"), - finalText("Wrote work_test.go."), + writeCall("write-red", "work.txt", "the regression marker\n"), + finalText("Wrote work.txt."), }, audit: []step{ - bashCall("run-check", "go test ./..."), - verdictFromEvidence("FAIL", "REFUTED — go test ./... now fails: TestWorkBreaksGreen", "VERIFIED — go test ./... passes"), + bashCall("run-check", "sh check.sh"), + verdictFromEvidence("FAIL", "REFUTED — sh check.sh now fails: TestWorkBreaksGreen", "VERIFIED — sh check.sh passes"), }, } agent, _ := newTestAgent(t, completer, func(config *Config) { @@ -366,7 +372,7 @@ func TestAGreenCheckTheWorkTurnedRedIsTheWorksOwn(t *testing.T) { if node == nil { t.Fatalf("no node was admitted; events = %#v", events) } - waitDoneNode(t, node) + awaitTestCompletion(t, node.done, "the task's refusal of the newly red check") notice := node.notice() if notice.State != TaskFailed { diff --git a/internal/session/team_wakewatch.go b/internal/session/team_wakewatch.go index d5a198629..d217c34d1 100644 --- a/internal/session/team_wakewatch.go +++ b/internal/session/team_wakewatch.go @@ -155,6 +155,9 @@ func (a *Agent) watchTeamTraffic(profile string) { // before the loop, so one already past its bound closes on the start // rather than waiting out a tick (team_wrapup.go). a.teamWrapUpResume(profile, time.Now()) + if a.config.teamWatchManual { + return + } guard.Go("team traffic wake", func() { a.teamWatchLoop(profile) }) } diff --git a/internal/session/team_wakewatch_test.go b/internal/session/team_wakewatch_test.go index fcd14b6c8..7b59a0955 100644 --- a/internal/session/team_wakewatch_test.go +++ b/internal/session/team_wakewatch_test.go @@ -259,10 +259,25 @@ func TestTeamWakeFinishedAfterReplyStillDeduplesAfterManagerRestart(t *testing.T } func TestTeamWakeTheLoopBreakerCountsTenMemberReplies(t *testing.T) { - fastTeamWake(t) + // The real traffic reader is driven at explicit logical instants, so no + // timer can consume a reply before the test advances the settle clock. + manual := func(config *Config) { config.teamWatchManual = true } fixture := newWakingTeamFixture(t) managerCalls := oneAnswer(teamLoopRounds + 1) - manager := teamAgent(t, fixture, fixture.manager, managerCalls, nil) + manager := teamAgent(t, fixture, fixture.manager, managerCalls, manual) + wakes, stopWakes := manager.WatchWakes() + defer stopWakes() + now := time.Now() + tick := func() { + if !manager.teamWatchTick(fixture.profile, now) { + t.Fatal("the manager stopped watching its team") + } + now = now.Add(teamWakeSettle) + if !manager.teamWatchTick(fixture.profile, now) { + t.Fatal("the manager stopped watching at the settle boundary") + } + now = now.Add(teamWakeSettle) + } posted := make(chan int, teamLoopRounds+1) releases := make([]chan struct{}, teamLoopRounds+1) for i := range releases { @@ -288,7 +303,7 @@ func TestTeamWakeTheLoopBreakerCountsTenMemberReplies(t *testing.T) { } return textResponse("finished"), nil }} - web = teamAgent(t, fixture, fixture.web, rounds, nil) + web = teamAgent(t, fixture, fixture.web, rounds, manual) for i := 0; i <= teamLoopRounds; i++ { rounds.arm(i) events := mustSubmit(t, web, "answer the manager") @@ -297,18 +312,35 @@ func TestTeamWakeTheLoopBreakerCountsTenMemberReplies(t *testing.T) { if got != i { t.Fatalf("posted round %d before %d", got, i) } - case <-time.After(5 * time.Second): + case <-time.After(awaitPatience(t)): t.Fatal("the member did not post") } + tick() if i < teamLoopRounds { - waitRequests(t, managerCalls, i+1) - waitIdle(t, manager) + select { + case stream := <-wakes: + collect(t, stream) + case <-time.After(awaitPatience(t)): + t.Fatal("the member reply did not wake the manager") + } } else { - waitEvent(t, fixture, "the team has woken me 10 times") + // The manual tick writes its hold notice before returning, so its + // completion is the acknowledgement; no polling window is needed. + if got := trafficEvents(t, fixture, "the team has woken me 10 times"); len(got) == 0 { + t.Fatal("the eleventh reply did not leave the loop notice in the traffic") + } } close(releases[i]) collect(t, events) - time.Sleep(3*teamWakeSettle + 10*teamWatchEvery) + // The member's stream closes after its finished event is written. + // Read that event and cross the settle boundary again to prove it + // cannot duplicate the explicit reply's wake. + tick() + select { + case <-wakes: + t.Fatalf("round %d woke the manager again after its reply", i) + default: + } if got := managerCalls.requests(); got != min(i+1, teamLoopRounds) { t.Fatalf("after round %d the manager made %d calls", i, got) } diff --git a/internal/session/tools_watch.go b/internal/session/tools_watch.go index 14aa402a1..a5d6164cf 100644 --- a/internal/session/tools_watch.go +++ b/internal/session/tools_watch.go @@ -515,6 +515,10 @@ func (r *jobRegistry) runWatch(ctx context.Context, cancel context.CancelFunc, w if fired { break } + if r.watchTickWait != nil { + r.watchTickWait(ctx) + continue + } select { case <-ctx.Done(): case <-ticker.C: diff --git a/internal/session/tools_watch_test.go b/internal/session/tools_watch_test.go index 6a57c61b5..7df4fdd55 100644 --- a/internal/session/tools_watch_test.go +++ b/internal/session/tools_watch_test.go @@ -2,12 +2,10 @@ package session // Watch tests. // -// Every timing assertion here rides a REAL two-second timer — the tool's own -// floor — because the thing under test is a loop whose whole subject is the -// passage of time, and a fake clock would test the fake. The cost is paid once: -// the slow cases run in parallel with each other, so the file's wall time is -// about one interval plus change rather than the sum of them, and every wait is -// a polled condition with a deadline rather than a fixed sleep. +// Every case drives the real watch loop and shell command. Cases about delta +// selection and quiet streaks acknowledge each tick and release the next one +// through the registry's wait seam; timer-driven cases still cover the real +// two-second floor. A failure deadline buys patience rather than ordering. import ( "encoding/json" @@ -148,6 +146,7 @@ func TestWatchRejectsBadArguments(t *testing.T) { func TestWatchChangeIsBaselineSilenceThenDeltaOnly(t *testing.T) { t.Parallel() agent, workspace := jobsAgent(t) + clock := controlWatch(agent) feed(t, workspace, "app.log", "old one", "old two") text, isError := startWatchTool(t, agent, map[string]any{ @@ -157,16 +156,18 @@ func TestWatchChangeIsBaselineSilenceThenDeltaOnly(t *testing.T) { t.Fatalf("watch failed to start: %s", text) } id := watchID(t, agent) + watched := agent.jobs.find(id) // Two ticks over an unchanged file: the first is the baseline, the second // has nothing to say. Neither may speak. - waitTicks(t, agent, id, 2) + clock.completedTick(t, watched) + clock.tick(t, watched) if queued := sessionNotes(agent); len(queued) != 0 { t.Fatalf("an unchanged watch spoke: %v", queued) } feed(t, workspace, "app.log", "old one", "old two", "new three", "new four") - waitFor(t, "the delta note", func() bool { return len(sessionNotes(agent)) > 0 }) + clock.tick(t, watched) queued := sessionNotes(agent) if len(queued) != 1 { @@ -236,6 +237,7 @@ func TestWatchNoteCapsAtFortyLines(t *testing.T) { func TestWatchMatchDeliversOnlyMatchingLines(t *testing.T) { t.Parallel() agent, workspace := jobsAgent(t) + clock := controlWatch(agent) feed(t, workspace, "app.log", "INFO starting") text, isError := startWatchTool(t, agent, map[string]any{ @@ -246,17 +248,25 @@ func TestWatchMatchDeliversOnlyMatchingLines(t *testing.T) { t.Fatalf("watch failed to start: %s", text) } id := watchID(t, agent) - waitTicks(t, agent, id, 1) + watched := agent.jobs.find(id) + clock.completedTick(t, watched) + if queued := sessionNotes(agent); len(queued) != 0 { + t.Fatalf("the matching watch's baseline spoke: %v", queued) + } // New lines, none of them matching: silence. feed(t, workspace, "app.log", "INFO starting", "INFO listening") - waitTicks(t, agent, id, 3) + clock.tick(t, watched) + clock.tick(t, watched) if queued := sessionNotes(agent); len(queued) != 0 { t.Fatalf("a non-matching change spoke: %v", queued) } feed(t, workspace, "app.log", "INFO starting", "INFO listening", "ERROR disk full", "INFO retrying") - waitFor(t, "the match note", func() bool { return len(sessionNotes(agent)) > 0 }) + clock.tick(t, watched) + if queued := sessionNotes(agent); len(queued) != 1 { + t.Fatalf("want exactly one matching note, got %v", queued) + } note := sessionNotes(agent)[0] if !strings.HasPrefix(note, "1 line matching /ERROR/") { @@ -362,6 +372,7 @@ func TestWatchUntilDeliversFinalNoteAndStops(t *testing.T) { func TestWatchStopsAfterThreeIdenticalFailures(t *testing.T) { t.Parallel() agent, _ := jobsAgent(t) + clock := controlWatch(agent) text, isError := startWatchTool(t, agent, map[string]any{ "command": "exit 7", "every_seconds": watchMinEvery, "name": "broken", @@ -371,7 +382,14 @@ func TestWatchStopsAfterThreeIdenticalFailures(t *testing.T) { } id := watchID(t, agent) - waitFor(t, "the failure note", func() bool { return len(sessionNotes(agent)) > 0 }) + watched := agent.jobs.find(id) + clock.completedTick(t, watched) + for tick := 1; tick < watchFailLimit; tick++ { + if queued := sessionNotes(agent); len(queued) != 0 { + t.Fatalf("a watch stopped before its failure limit: %v", queued) + } + clock.tick(t, watched) + } queued := sessionNotes(agent) if len(queued) != 1 { t.Fatalf("a broken watch reported more than once: %v", queued) @@ -519,6 +537,7 @@ func TestCloseStopsWatches(t *testing.T) { func TestWatchQuietFiresWhenTheOutputStopsMoving(t *testing.T) { t.Parallel() agent, workspace := jobsAgent(t) + clock := controlWatch(agent) feed(t, workspace, "build.log", "compiling one") text, isError := startWatchTool(t, agent, map[string]any{ @@ -533,22 +552,27 @@ func TestWatchQuietFiresWhenTheOutputStopsMoving(t *testing.T) { t.Fatalf("the start line does not state the quiet terms: %q", text) } id := watchID(t, agent) + watched := agent.jobs.find(id) // WHILE THE OUTPUT MOVES, NOTHING IS SAID. Two changes across three ticks // keep resetting the run, and a quiet watch that spoke here would be // announcing the opposite of what it was asked to watch for. - waitTicks(t, agent, id, 1) + clock.completedTick(t, watched) feed(t, workspace, "build.log", "compiling one", "compiling two") - waitTicks(t, agent, id, 2) + clock.tick(t, watched) feed(t, workspace, "build.log", "compiling one", "compiling two", "linking") - waitTicks(t, agent, id, 3) + clock.tick(t, watched) if queued := sessionNotes(agent); len(queued) != 0 { t.Fatalf("a quiet watch spoke while the output was moving: %v", queued) } // Now leave it alone. Two unchanged ticks in a row are the terms, and the // note that follows is the last one this watch ever sends. - waitFor(t, "the quiet note", func() bool { return len(sessionNotes(agent)) > 0 }) + clock.tick(t, watched) + if queued := sessionNotes(agent); len(queued) != 0 { + t.Fatalf("a quiet watch fired after only one unchanged tick: %v", queued) + } + clock.tick(t, watched) queued := sessionNotes(agent) if len(queued) != 1 { t.Fatalf("want exactly one note, got %v", queued) @@ -560,8 +584,10 @@ func TestWatchQuietFiresWhenTheOutputStopsMoving(t *testing.T) { t.Fatalf("the final note does not carry the last output line: %q", queued[0]) } - waitFor(t, "the watch to stop", func() bool { return !agent.jobs.find(id).running() }) - waitSignal(t, agent.jobs.find(id).done, "the quiet watch loop to return") + waitSignal(t, watched.done, "the quiet watch loop to return") + if watched.running() { + t.Fatal("a quiet watch kept running after it fired") + } // It ends the way `until` ends: its returned loop has no ticker left. before := agent.jobs.find(id).tickCount() if after := agent.jobs.find(id).tickCount(); after != before { diff --git a/internal/session/unattendeddoor_test.go b/internal/session/unattendeddoor_test.go index dc048659c..7292c5c3a 100644 --- a/internal/session/unattendeddoor_test.go +++ b/internal/session/unattendeddoor_test.go @@ -26,6 +26,7 @@ import ( "os" "path/filepath" "strings" + "sync/atomic" "testing" "time" @@ -837,41 +838,53 @@ func TestACommitRefusedInAWritableTreeIsAboutTheWork(t *testing.T) { // Here the first call hangs, is cut at its share ([auditCallShare]), and the // second answers — with the landing saying which try it was. func TestAStalledCheckIsAbandonedAndTheSecondCallAnswers(t *testing.T) { - repo := newGoModuleRepo(t) + repo := newTestRepo(t) t.Setenv("HOME", t.TempDir()) + const window = auditDeadline + opened := make(chan controlledCheckCall, 2) + started := time.Now() + var elapsed atomic.Int64 completer := &routedCompleter{ parent: []step{ - proposeCall("Add the greeting", "write greet.go"), + proposeCall("Add the greeting", "write greet.txt"), finalText("handed off"), }, child: []step{ - writeCall("call-src", "greet.go", "package greet\n\nfunc Greet() string { return \"hi\" }\n"), - finalText("Wrote greet.go with the greeting."), + writeCall("call-src", "greet.txt", "hi\n"), + finalText("Wrote greet.txt with the greeting."), }, audit: []step{ // THE HUNG STREAM. It answers nothing and never refuses, which is // exactly what the wire did. func(ctx context.Context, _ []ai.Message) (*ai.Response, error) { + call := <-opened + if want := window / auditCallShare; call.bound != want { + t.Errorf("the first call was given %s, want its share %s", call.bound, want) + } + // Expire only after the stalled call starts, at its assigned share. + // The same logical clock leaves the rest of the window for retry. + elapsed.Store(int64(call.bound)) + call.expire() <-ctx.Done() return nil, ctx.Err() }, - verdict("VERIFIED — greet.go has the greeting the brief asked for"), + verdict("VERIFIED — greet.txt has the greeting the brief asked for"), }, } agent, _ := newTestAgent(t, completer, func(config *Config) { config.Workspace = repo config.AskConsent = false config.TaskAutoApproveSeconds = 0 - // The first call must time out, while the second still needs room for - // real repository setup when other package suites share the machine. - config.auditWindow = 10 * time.Second + config.auditWindow = window + config.clock = func() time.Time { return started.Add(time.Duration(elapsed.Load())) } + config.auditTimeout = controlCheck(opened) }) graph := agent.graph() collect(t, mustSubmit(t, agent, "add a greeting")) node := graph.node(1) - waitDoneNode(t, node) + awaitTestCompletion(t, node.done, "the retried task's landing") notice := node.notice() if notice.State != TaskDone { diff --git a/internal/session/watchwake_test.go b/internal/session/watchwake_test.go index e0ee01d51..c325c520f 100644 --- a/internal/session/watchwake_test.go +++ b/internal/session/watchwake_test.go @@ -15,9 +15,10 @@ package session // is telemetry with a complete log behind it and must stay quiet, and a turn per // delta would turn a quiet observer into an autonomous conversation. // -// Every case here drives the REAL timer loop rather than calling the lane's door, -// for tools_watch_test.go's reason: the thing under test is which news a tick is, -// and a fixture that decided that for itself would be testing the fixture. +// Every case here drives the real watch loop rather than calling the lane's +// door. The tick-silence case releases acknowledged ticks through the wait seam; +// the other cases retain the real timer. The loop still decides which news a +// tick is, because a fixture that chose the lane would be testing the fixture. import ( "context" @@ -138,10 +139,14 @@ func TestWatchTicksNeverStartATurn(t *testing.T) { }, }} agent, workspace := newTestAgent(t, completer, nil) + clock := controlWatch(agent) - firingWatch(t, agent, workspace, "app.log", "this never appears", "starting up") + id := firingWatch(t, agent, workspace, "app.log", "this never appears", "starting up") + watched := agent.jobs.find(id) + clock.completedTick(t, watched) for round, line := range []string{"one", "two", "three"} { feed(t, workspace, "app.log", "starting up", line) + clock.tick(t, watched) want := round + 1 waitFor(t, fmt.Sprintf("tick %d to reach the ambient queue", want), func() bool { return len(ambientQueue(agent)) == want @@ -154,6 +159,12 @@ func TestWatchTicksNeverStartATurn(t *testing.T) { if queued := steeringQueue(agent); len(queued) != 0 { t.Fatalf("a tick reached the owed lane: %v", queued) } + agent.mu.Lock() + running := agent.running + agent.mu.Unlock() + if running { + t.Fatal("ordinary watch ticks started a turn before its model was scheduled") + } } // ── and the firing wakes ────────────────────────────────────────────────────