diff --git a/cli/src/codex/shared/root.test.ts b/cli/src/codex/shared/root.test.ts index 3f6dcb33..4f996a23 100644 --- a/cli/src/codex/shared/root.test.ts +++ b/cli/src/codex/shared/root.test.ts @@ -15,6 +15,7 @@ vi.mock('../codexAppServerClient', () => ({ thread = { id: 'thread', turns: [] as NativeTurn[] }; settings: Record = { 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 = {}) { + 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>, status = 'completed') { @@ -104,6 +108,13 @@ async function completePlan(f: Awaited>, 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); diff --git a/cli/src/codex/shared/root.ts b/cli/src/codex/shared/root.ts index 47009cf0..3e9d6972 100644 --- a/cli/src/codex/shared/root.ts +++ b/cli/src/codex/shared/root.ts @@ -180,7 +180,10 @@ export class SharedCodexRoot { 'shell_environment_policy.set.HAPI_SESSION_ID': this.session.sessionId } }; } - async bind(threadId: string, response: Record, subscribe: boolean): Promise { + /** 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, subscribe: boolean): Promise> { 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): Promise { await this.projection.history(response.thread); await this.refresh(); await this.refreshChildren(true); } async activate(options: SharedLaunchOptions = {}): Promise { diff --git a/cli/src/codex/shared/runtime.ts b/cli/src/codex/shared/runtime.ts index 3a054ce2..987f1f24 100644 --- a/cli/src/codex/shared/runtime.ts +++ b/cli/src/codex/shared/runtime.ts @@ -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, parent?: SharedCodexRoot, initialOptions?: SharedLaunchOptions): Promise => operation(async () => { const root = await prepare(string(params.cwd) ?? launch.cwd, undefined, parent);