fix(codex): report session started before native history replay

The runner kills a runner-spawned child whose "session started" webhook
misses its bounded wait (15s by default), so a cold resume whose thread
has a large rollout always failed with "Session webhook timeout for PID
...": the codex shared runtime only reported after reconciling the full
native history (projection.history + refresh + refreshChildren), and
replaying tens of thousands of items took longer than that wait.

Measured on this machine before the fix: two ~800MB rollouts (36k/43k
lines) needed 17-20s to bind; one attempt was killed 2s before its late
webhook landed and was reaped as an orphan.

Split SharedCodexRoot.bind() into the local claim (queue, projection,
metadata, settings) and a separate syncHistory(); runSharedRuntime.bind()
now activates the session and reports it to the runner before the replay.
The session is controllable when the signal fires (controls registered),
and the replay continues in the background; re-sent history is idempotent
because projection uses stable message ids that the hub dedupes by
local_id.

Covered by root.test.ts: bind performs no native history reads, syncHistory
does the reconciliation.
This commit is contained in:
2026-09-29 20:18:51 +08:00
parent 8e1c4a80a1
commit d6cf54d724
3 changed files with 31 additions and 8 deletions
+14 -3
View File
@@ -15,6 +15,7 @@ vi.mock('../codexAppServerClient', () => ({
thread = { id: 'thread', turns: [] as NativeTurn[] };
settings: Record<string, unknown> = { model: 'mock', collaborationMode: { mode: 'default' } };
queue: Array<{ id: string; clientUserMessageId: unknown; input: unknown }> = [];
requests: string[] = [];
notify?: (method: string, params: unknown) => void;
abandoned?: () => void;
setNotificationHandler(handler: typeof this.notify) { this.notify = handler; }
@@ -25,6 +26,7 @@ vi.mock('../codexAppServerClient', () => ({
isInitialized() { return this.initialized; }
async disconnect() { this.initialized = false; }
async request(method: string, params: Record<string, unknown> = {}) {
this.requests.push(method);
if (method === 'thread/read' || method === 'thread/resume') return { ...this.settings, thread: structuredClone(this.thread) };
if (method === 'thread/list') return { data: [] };
if (method === 'thread/queue/list') return { data: this.queue };
@@ -53,7 +55,7 @@ afterEach(async () => {
finally { vi.useRealTimers(); }
});
async function fixture() {
async function fixture(options: { deferHistory?: boolean } = {}) {
const directory = await mkdtemp('/tmp/hapi-shared-root-');
let state: AgentState = { steeringActive: true };
let metadata: Metadata = { path: directory, host: 'test', flavor: 'codex' };
@@ -79,15 +81,17 @@ async function fixture() {
} satisfies RootHost);
cleanups.push(async () => { await root.close(false); await rm(directory, { recursive: true, force: true }); });
await root.prepare();
await root.bind('thread', { model: 'mock', thread: { turns: [] } }, false);
const response = await root.bind('thread', { model: 'mock', thread: { turns: [] } }, false);
if (!options.deferHistory) await root.syncHistory(response);
const native = root.client as unknown as {
initialized: boolean;
thread: { id: string; turns: NativeTurn[] };
queue: Array<{ id: string; clientUserMessageId: string; input: unknown }>;
requests: string[];
notify(method: string, params: unknown): void;
abandoned(): void;
};
return { root, native, rpc, send, state: () => state, updateState, reconnect: () => reconnect?.() };
return { root, native, rpc, send, state: () => state, updateState, reconnect: () => reconnect?.(), syncHistory: () => root.syncHistory(response) };
}
async function completePlan(f: Awaited<ReturnType<typeof fixture>>, status = 'completed') {
@@ -104,6 +108,13 @@ async function completePlan(f: Awaited<ReturnType<typeof fixture>>, status = 'co
}
describe('shared plan actions', () => {
it('keeps bind local and reconciles native history only in syncHistory', async () => {
const f = await fixture({ deferHistory: true });
expect(f.native.requests).not.toContain('thread/read');
await f.syncHistory();
expect(f.native.requests).toContain('thread/read');
});
it('preserves content while native turns, mode changes and disconnects withdraw controls', async () => {
const f = await fixture();
const id = await completePlan(f);
+9 -1
View File
@@ -180,7 +180,10 @@ export class SharedCodexRoot {
'shell_environment_policy.set.HAPI_SESSION_ID': this.session.sessionId
} };
}
async bind(threadId: string, response: Record<string, unknown>, subscribe: boolean): Promise<void> {
/** Local claim only: queue, projection, metadata and settings. Native
* history reconciliation is a separate step (syncHistory) so liveness
* signals never wait on an arbitrarily large thread replay. */
async bind(threadId: string, response: Record<string, unknown>, subscribe: boolean): Promise<Record<string, unknown>> {
if (this.threadId && this.threadId !== threadId) throw new Error('Cannot retarget a shared HAPI session');
this.threadId = threadId;
this.queue = new SharedCodexQueue(this.client, threadId, join(this.host.directory, `${this.session.sessionId}.queue.json`),
@@ -201,6 +204,11 @@ export class SharedCodexRoot {
}));
if (subscribe) response = record(await this.client.request('thread/resume', { threadId }));
this.acceptSettings(response); this.acceptSettings(this.host.settingsFor(threadId) ?? {});
return response;
}
/** Reconcile the full native history into the hub. Huge threads can take
* minutes; the HAPI session is already usable while this runs. */
async syncHistory(response: Record<string, unknown>): Promise<void> {
await this.projection.history(response.thread); await this.refresh(); await this.refreshChildren(true);
}
async activate(options: SharedLaunchOptions = {}): Promise<void> {
+8 -4
View File
@@ -201,15 +201,19 @@ export async function runSharedRuntime(options: SharedLaunchOptions, onReady?: (
await control.request('thread/metadata/update', { threadId, gitInfo: gitInfo(string(record(response.thread).cwd) ?? root.bootstrap.workingDirectory) });
assertRunning();
roots.set(threadId, root);
await root.bind(threadId, response, subscribe);
const effective = await root.bind(threadId, response, subscribe);
assertRunning();
// Cold-resumed threads predate the control connection's automatic
// new-thread subscription. Subscribe once without changing settings.
await control.request('thread/resume', { threadId });
await root.activate(initialOptions);
await root.session.flush();
await persist();
// The runner kills a child whose "session started" webhook misses its
// bounded wait. Report ownership as soon as the session is controllable;
// the native-history replay below can take minutes on large threads.
await notifyRunnerSessionStarted(root.session.sessionId, root.session.getMetadata() ?? root.bootstrap.metadata);
// Cold-resumed threads predate the control connection's automatic
// new-thread subscription. Subscribe once without changing settings.
await control.request('thread/resume', { threadId });
await root.syncHistory(effective);
};
const create = (method: 'thread/start' | 'thread/fork', params: Record<string, unknown>, parent?: SharedCodexRoot, initialOptions?: SharedLaunchOptions): Promise<SharedCodexRoot> => operation(async () => {
const root = await prepare(string(params.cwd) ?? launch.cwd, undefined, parent);