Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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."
---
9 changes: 8 additions & 1 deletion internal/session/callwindow.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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
}

Expand Down
146 changes: 146 additions & 0 deletions internal/session/controlled_waits_test.go
Original file line number Diff line number Diff line change
@@ -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
}
}
3 changes: 3 additions & 0 deletions internal/session/jobs.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
28 changes: 21 additions & 7 deletions internal/session/loop_speed_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{
Expand All @@ -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.
Expand Down
10 changes: 10 additions & 0 deletions internal/session/session.go
Original file line number Diff line number Diff line change
Expand Up @@ -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].
Expand All @@ -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].
Expand Down Expand Up @@ -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
Expand Down
6 changes: 5 additions & 1 deletion internal/session/steer_grace.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down
36 changes: 24 additions & 12 deletions internal/session/steer_grace_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,6 @@ import (
"context"
"strings"
"sync"
"sync/atomic"
"testing"
"time"

Expand All @@ -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) {
Expand All @@ -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) {
Expand All @@ -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") {
Expand Down
2 changes: 1 addition & 1 deletion internal/session/task_audit.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down
Loading
Loading