| 1 | package control |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "errors" |
| 6 | "os" |
| 7 | "path/filepath" |
| 8 | "sync/atomic" |
| 9 | "testing" |
| 10 | "time" |
| 11 | |
| 12 | "reasonix/internal/event" |
| 13 | "reasonix/internal/turnevent" |
| 14 | ) |
| 15 | |
| 16 | type turnEventGateRunner struct { |
| 17 | started chan struct{} |
| 18 | release chan struct{} |
| 19 | calls atomic.Int32 |
| 20 | } |
| 21 | |
| 22 | func (r *turnEventGateRunner) Run(context.Context, string) error { |
| 23 | r.calls.Add(1) |
| 24 | close(r.started) |
| 25 | <-r.release |
| 26 | return nil |
| 27 | } |
| 28 | |
| 29 | func TestTurnAdmissionIsDurableBeforeRunnerStarts(t *testing.T) { |
| 30 | dir := t.TempDir() |
| 31 | runner := &turnEventGateRunner{started: make(chan struct{}), release: make(chan struct{})} |
| 32 | done := make(chan event.Event, 1) |
| 33 | c := newOwnedTestController(t, Options{ |
| 34 | Runner: runner, |
| 35 | Sink: event.FuncSink(func(e event.Event) { |
| 36 | if e.Kind == event.TurnDone { |
| 37 | done <- e |
| 38 | } |
| 39 | }), |
| 40 | SessionDir: dir, SessionPath: filepath.Join(dir, "session.jsonl"), |
| 41 | }) |
| 42 | t.Cleanup(c.Close) |
| 43 | |
| 44 | c.Submit("run") |
| 45 | select { |
| 46 | case <-runner.started: |
| 47 | case <-time.After(5 * time.Second): |
| 48 | t.Fatal("runner did not start") |
| 49 | } |
| 50 | records, err := c.TurnEventsAfter(0) |
| 51 | if err != nil { |
| 52 | t.Fatalf("TurnEventsAfter: %v", err) |
| 53 | } |
| 54 | if len(records) < 2 || records[0].Status != event.TurnQueued || records[1].Kind != "turn_started" || records[1].Status != event.TurnInProgress { |
| 55 | t.Fatalf("admission prefix = %+v, want queued then durable in_progress start", records) |
| 56 | } |
| 57 | close(runner.release) |
| 58 | select { |
| 59 | case <-done: |
| 60 | case <-time.After(5 * time.Second): |
| 61 | t.Fatal("turn did not finish") |
| 62 | } |
| 63 | records, err = c.TurnEventsAfter(0) |
| 64 | if err != nil { |
| 65 | t.Fatalf("TurnEventsAfter terminal: %v", err) |
| 66 | } |
| 67 | started := 0 |
| 68 | for _, record := range records { |
| 69 | if record.Kind == "turn_started" { |
| 70 | started++ |
| 71 | } |
| 72 | } |
| 73 | if started != 1 { |
| 74 | t.Fatalf("turn_started records = %d, want exactly one", started) |
| 75 | } |
| 76 | } |
| 77 | |
| 78 | func TestTurnAdmissionLedgerFailureDoesNotRunProvider(t *testing.T) { |
| 79 | dir := t.TempDir() |
| 80 | blockedParent := filepath.Join(dir, "not-a-directory") |
| 81 | if err := os.WriteFile(blockedParent, []byte("block"), 0o600); err != nil { |
| 82 | t.Fatalf("write blocker: %v", err) |
| 83 | } |
| 84 | runner := &turnEventGateRunner{started: make(chan struct{}), release: make(chan struct{})} |
| 85 | done := make(chan event.Event, 1) |
| 86 | c := newOwnedTestController(t, Options{ |
| 87 | Runner: runner, |
| 88 | Sink: event.FuncSink(func(e event.Event) { |
| 89 | if e.Kind == event.TurnDone { |
| 90 | done <- e |
| 91 | } |
| 92 | }), |
| 93 | SessionDir: dir, SessionPath: filepath.Join(blockedParent, "session.jsonl"), |
| 94 | }) |
| 95 | t.Cleanup(c.Close) |
| 96 | |
| 97 | c.Submit("must not reach provider") |
| 98 | select { |
| 99 | case terminal := <-done: |
| 100 | if !errors.Is(terminal.Err, turnevent.ErrTurnLedgerUnavailable) { |
| 101 | t.Fatalf("terminal error = %v, want explicit ledger admission failure", terminal.Err) |
| 102 | } |
| 103 | case <-time.After(5 * time.Second): |
| 104 | t.Fatal("failed admission did not terminate") |
| 105 | } |
| 106 | if got := runner.calls.Load(); got != 0 { |
| 107 | t.Fatalf("runner calls = %d, want provider side effects blocked", got) |
| 108 | } |
| 109 | } |
| 110 | |
| 111 | func TestAsyncStreamLedgerFailureCancelsTurnWithoutPublishingChunk(t *testing.T) { |
| 112 | root := filepath.Join(t.TempDir(), "session-dir") |
| 113 | if err := os.MkdirAll(root, 0o700); err != nil { |
| 114 | t.Fatal(err) |
| 115 | } |
| 116 | started := make(chan struct{}) |
| 117 | release := make(chan struct{}) |
| 118 | cancelled := make(chan struct{}) |
| 119 | done := make(chan event.Event, 1) |
| 120 | var publishedText atomic.Int32 |
| 121 | c := newOwnedTestController(t, Options{ |
| 122 | Sink: event.FuncSink(func(e event.Event) { |
| 123 | if e.Kind == event.Text { |
| 124 | publishedText.Add(1) |
| 125 | } |
| 126 | if e.Kind == event.TurnDone { |
| 127 | done <- e |
| 128 | } |
| 129 | }), |
| 130 | SessionDir: root, SessionPath: filepath.Join(root, "session.jsonl"), |
| 131 | }) |
| 132 | t.Cleanup(c.Close) |
| 133 | |
| 134 | c.runGuarded(func(ctx context.Context) error { |
| 135 | close(started) |
| 136 | <-release |
| 137 | c.sink.Emit(event.Event{Kind: event.ToolResult, Tool: event.Tool{ID: "failed-store", Name: "probe", Output: "must stay behind the v3 log"}}) |
| 138 | <-ctx.Done() |
| 139 | close(cancelled) |
| 140 | return ctx.Err() |
| 141 | }) |
| 142 | <-started |
| 143 | v3 := c.sessionEventStore() |
| 144 | if v3 == nil { |
| 145 | t.Fatal("controller did not open a v3 event store") |
| 146 | } |
| 147 | if err := v3.Close(context.Background()); err != nil { |
| 148 | t.Fatalf("close active v3 handle: %v", err) |
| 149 | } |
| 150 | close(release) |
| 151 | |
| 152 | select { |
| 153 | case <-cancelled: |
| 154 | case <-time.After(5 * time.Second): |
| 155 | t.Fatal("asynchronous stream persistence failure did not cancel the turn") |
| 156 | } |
| 157 | select { |
| 158 | case terminal := <-done: |
| 159 | if terminal.Status != event.TurnFailed || terminal.Err == nil { |
| 160 | t.Fatalf("control-plane terminal = %+v, want explicit storage failure", terminal) |
| 161 | } |
| 162 | case <-time.After(5 * time.Second): |
| 163 | t.Fatal("storage failure did not release frontend running state") |
| 164 | } |
| 165 | if got := publishedText.Load(); got != 0 { |
| 166 | t.Fatalf("published text chunks = %d, want none before durable append", got) |
| 167 | } |
| 168 | if err := c.turnEventLedgerError(); err == nil { |
| 169 | t.Fatal("controller accepted new work after the ledger was poisoned") |
| 170 | } |
| 171 | } |
| 172 |