mirror of
https://github.com/wu736139669/hapi.git
synced 2026-10-06 18:39:47 +00:00
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:
@@ -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);
|
||||
|
||||
@@ -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> {
|
||||
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user