Skip to content

Commit c821bc3

Browse files
committed
Preserve live steer drain order through Agent.deliver
Ingest used to start before enqueue, so two steers at one boundary could reverse if the second mention resolved first.
1 parent ff577d9 commit c821bc3

4 files changed

Lines changed: 225 additions & 17 deletions

File tree

src/tui/queued-delivery-hop.test.ts

Lines changed: 98 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,8 @@ import { attachSessionBridge, type SessionBridge } from "./runtime-bridge";
77
import { createLiveSessionPort } from "./live-session-port";
88
import { createAppShell } from "./shell";
99
import { withTestRenderer } from "./harness";
10-
import { routeQueuedDelivery } from "./queued-delivery.js";
10+
import { createLiveSteerDeliver, routeQueuedDelivery } from "./queued-delivery.js";
11+
import { createSessionOperationQueue } from "./session-operation-queue.js";
1112
import { badgeCount } from "./session-queue";
1213

1314
function lastHopPort(bridgeRef: { current: SessionBridge | undefined }) {
@@ -59,6 +60,102 @@ describe("queued delivery last hop", () => {
5960
);
6061
});
6162

63+
test("two live steers at one tool.boundary keep drain order through Agent.deliver", async () => {
64+
await withTestRenderer(
65+
async (h) => {
66+
const shell = createAppShell(h.renderer, {
67+
terminal: { columns: 80, rows: 24 },
68+
wireKeys: false,
69+
run: "busy",
70+
});
71+
const bridgeRef: { current: SessionBridge | undefined } = { current: undefined };
72+
const sends: string[] = [];
73+
const delivered: string[] = [];
74+
const { enqueue, awaitTail } = createSessionOperationQueue();
75+
let resolveSlow!: () => void;
76+
const slow = new Promise<void>((resolve) => {
77+
resolveSlow = resolve;
78+
});
79+
const port = createLiveSessionPort({
80+
send: (text) => {
81+
sends.push(text);
82+
},
83+
interrupt: () => {},
84+
deliver: routeQueuedDelivery({
85+
send: (text) => {
86+
sends.push(text);
87+
},
88+
deliverSteer: createLiveSteerDeliver({
89+
enqueue,
90+
ingest: async (text) => {
91+
if (text.includes("@mention")) await slow;
92+
return { text, attachments: [] };
93+
},
94+
deliver: (text) => {
95+
delivered.push(text);
96+
},
97+
captureGeneration: () => () => true,
98+
onFailure: (err) => {
99+
throw err;
100+
},
101+
}),
102+
parentCycleLive: () => bridgeRef.current?.parentCycleLive === true,
103+
}),
104+
});
105+
const bridge = attachSessionBridge(shell, port);
106+
bridgeRef.current = bridge;
107+
try {
108+
bridge.submit("@mention first", "steer");
109+
bridge.submit("plain second", "steer");
110+
expect(badgeCount(shell.session)).toBe(2);
111+
bridge.handle({ type: "tool.boundary" });
112+
await Promise.resolve();
113+
await Promise.resolve();
114+
expect(delivered).toEqual([]);
115+
resolveSlow();
116+
await awaitTail();
117+
expect(delivered).toEqual(["@mention first", "plain second"]);
118+
expect(sends).toEqual([]);
119+
} finally {
120+
bridge.dispose();
121+
shell.dispose();
122+
}
123+
},
124+
{ width: 80, height: 24 },
125+
);
126+
});
127+
128+
test("inference.done with outstanding tools last-hops to deliverSteer, not send", async () => {
129+
await withTestRenderer(
130+
async (h) => {
131+
const shell = createAppShell(h.renderer, {
132+
terminal: { columns: 80, rows: 24 },
133+
wireKeys: false,
134+
run: "busy",
135+
});
136+
const bridgeRef: { current: SessionBridge | undefined } = { current: undefined };
137+
const { port, sends, steers } = lastHopPort(bridgeRef);
138+
const bridge = attachSessionBridge(shell, port);
139+
bridgeRef.current = bridge;
140+
try {
141+
bridge.submit("asap", "steer");
142+
expect(badgeCount(shell.session)).toBe(1);
143+
bridge.handle({
144+
type: "tool.start",
145+
data: { call: { id: "c1", name: "run_shell" } },
146+
});
147+
bridge.handle({ type: "inference.done", data: {} });
148+
expect(steers).toEqual(["asap"]);
149+
expect(sends).toEqual([]);
150+
} finally {
151+
bridge.dispose();
152+
shell.dispose();
153+
}
154+
},
155+
{ width: 80, height: 24 },
156+
);
157+
});
158+
62159
test("text-only settle leftover steer last-hops to send", async () => {
63160
await withTestRenderer(
64161
async (h) => {

src/tui/queued-delivery.test.ts

Lines changed: 72 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,11 @@
11
import { describe, expect, test } from "bun:test";
22
import type { PendingImageAttachment } from "./image-attachments.js";
3-
import { createDeliveryGeneration, routeQueuedDelivery } from "./queued-delivery.js";
3+
import {
4+
createDeliveryGeneration,
5+
createLiveSteerDeliver,
6+
routeQueuedDelivery,
7+
} from "./queued-delivery.js";
8+
import { createSessionOperationQueue } from "./session-operation-queue.js";
49

510
const image: PendingImageAttachment = {
611
id: "img-1",
@@ -106,3 +111,69 @@ describe("createDeliveryGeneration", () => {
106111
expect(second()).toBe(false);
107112
});
108113
});
114+
115+
describe("createLiveSteerDeliver", () => {
116+
test("slow first ingest does not let a later steer deliver first", async () => {
117+
const delivered: string[] = [];
118+
const { enqueue, awaitTail } = createSessionOperationQueue();
119+
let resolveSlow!: () => void;
120+
const slow = new Promise<void>((resolve) => {
121+
resolveSlow = resolve;
122+
});
123+
const deliverSteer = createLiveSteerDeliver({
124+
enqueue,
125+
ingest: async (text) => {
126+
if (text.startsWith("@mention")) await slow;
127+
return { text, attachments: [] };
128+
},
129+
deliver: (text) => {
130+
delivered.push(text);
131+
},
132+
captureGeneration: () => () => true,
133+
onFailure: (err) => {
134+
throw err;
135+
},
136+
});
137+
138+
deliverSteer("@mention A");
139+
deliverSteer("plain B");
140+
await Promise.resolve();
141+
await Promise.resolve();
142+
expect(delivered).toEqual([]);
143+
144+
resolveSlow();
145+
await awaitTail();
146+
expect(delivered).toEqual(["@mention A", "plain B"]);
147+
});
148+
149+
test("generation bump during ingest drops both in-flight live steers", async () => {
150+
const delivered: string[] = [];
151+
const { enqueue, awaitTail } = createSessionOperationQueue();
152+
const generation = createDeliveryGeneration();
153+
let resolveSlow!: () => void;
154+
const slow = new Promise<void>((resolve) => {
155+
resolveSlow = resolve;
156+
});
157+
const deliverSteer = createLiveSteerDeliver({
158+
enqueue,
159+
ingest: async (text) => {
160+
if (text === "A") await slow;
161+
return { text, attachments: [] };
162+
},
163+
deliver: (text) => {
164+
delivered.push(text);
165+
},
166+
captureGeneration: generation.capture,
167+
onFailure: (err) => {
168+
throw err;
169+
},
170+
});
171+
172+
deliverSteer("A");
173+
deliverSteer("B");
174+
generation.bump();
175+
resolveSlow();
176+
await awaitTail();
177+
expect(delivered).toEqual([]);
178+
});
179+
});

src/tui/queued-delivery.ts

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,3 +44,43 @@ export function createDeliveryGeneration() {
4444
},
4545
};
4646
}
47+
48+
export interface IngestedSteer {
49+
readonly text: string;
50+
readonly attachments: readonly PendingImageAttachment[];
51+
}
52+
53+
export interface CreateLiveSteerDeliverArgs {
54+
/**
55+
* FIFO session queue. Ingest must run on this queue — not in a
56+
* fire-and-forget IIFE — so two steers at one boundary cannot reverse
57+
* if the second ingest finishes first.
58+
*/
59+
enqueue: (op: () => Promise<void>) => Promise<void>;
60+
ingest: (text: string, attachments: readonly PendingImageAttachment[]) => Promise<IngestedSteer>;
61+
/** Agent.deliver (or the sessionOps enqueue that wraps it). */
62+
deliver: (text: string, attachments: readonly PendingImageAttachment[]) => void;
63+
captureGeneration: () => () => boolean;
64+
onFailure: (err: unknown) => void;
65+
}
66+
67+
/**
68+
* Live inject: enqueue ingest, then deliver, in drain order. Previously
69+
* each item started ingest immediately, so Agent.deliver could reverse.
70+
*/
71+
export function createLiveSteerDeliver(
72+
args: CreateLiveSteerDeliverArgs,
73+
): (text: string, attachments?: readonly PendingImageAttachment[]) => void {
74+
return (text, attachments) => {
75+
const stillCurrent = args.captureGeneration();
76+
const pending = attachments ?? [];
77+
void args
78+
.enqueue(async () => {
79+
if (!stillCurrent()) return;
80+
const ingested = await args.ingest(text, pending);
81+
if (!stillCurrent()) return;
82+
args.deliver(ingested.text, ingested.attachments);
83+
})
84+
.catch(args.onFailure);
85+
};
86+
}

src/tui/runner.ts

Lines changed: 15 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -185,7 +185,11 @@ import { setActiveWebProviderBrand } from "./tool-formatter.js";
185185
import { consumeStream } from "../session/stream-consumer.js";
186186
import { createCycleTextRecorder } from "../session/stream-journal.js";
187187
import { mountRunnerHost } from "./runner-host.js";
188-
import { createDeliveryGeneration, routeQueuedDelivery } from "./queued-delivery.js";
188+
import {
189+
createDeliveryGeneration,
190+
createLiveSteerDeliver,
191+
routeQueuedDelivery,
192+
} from "./queued-delivery.js";
189193
import { createRuntimeShutdown } from "./runtime-shutdown.js";
190194
import {
191195
applyFocus,
@@ -2294,20 +2298,16 @@ export async function runTUI(initialConfig: Config): Promise<number> {
22942298
deliver: routeQueuedDelivery({
22952299
send,
22962300
parentCycleLive: () => host.bridge.parentCycleLive,
2297-
deliverSteer: (text, attachments) => {
2298-
const stillCurrent = deliveryGeneration.capture();
2299-
void (async () => {
2300-
if (!stillCurrent()) return;
2301-
const ingested = await ingestOperatorPrompt(
2302-
text,
2303-
config.cwd,
2304-
imageAttachmentFromPath,
2305-
attachments ?? [],
2306-
);
2307-
if (!stillCurrent()) return;
2308-
agentProxy.deliver(userInboundMessage(ingested.text, ingested.attachments));
2309-
})().catch(handleSendFailure);
2310-
},
2301+
deliverSteer: createLiveSteerDeliver({
2302+
enqueue: sessionOps.enqueue,
2303+
ingest: (text, pending) =>
2304+
ingestOperatorPrompt(text, config.cwd, imageAttachmentFromPath, pending),
2305+
deliver: (text, pending) => {
2306+
agentProxy.deliver(userInboundMessage(text, pending));
2307+
},
2308+
captureGeneration: deliveryGeneration.capture,
2309+
onFailure: handleSendFailure,
2310+
}),
23112311
}),
23122312
// Consent by proceeding requires the disclosure to be on screen before the
23132313
// first prompt activates the held telemetry instance: the landing shows it,

0 commit comments

Comments
 (0)