diff --git a/src/web/reboot-restore-registry.ts b/src/web/reboot-restore-registry.ts index 0dfab538..94a7c63f 100644 --- a/src/web/reboot-restore-registry.ts +++ b/src/web/reboot-restore-registry.ts @@ -48,16 +48,18 @@ export class RebootRestoreRegistry { /** When the boot pass built the plan, in ms since the epoch. */ private builtAt = 0; /** - * Per owner, bumped by anything that invalidates that owner's entries while a - * restore is already holding them. A Dismiss arriving mid-restore must win: - * without this the route's `finally` would put its unspent entries back and - * resurrect the offer the user just cleared, with a fresh 24-hour life. + * Entries handed to a restore that has not finished, by session id, each + * remembering which caller is spending it. * - * Keyed by owner rather than global, because `clear()` is ownership-scoped. A - * single counter would let one user's Dismiss discard another user's unspent - * entries, and the plan is in-memory, so those offers would be gone for good. + * A taken entry is still part of the offer until its restore resolves it, so + * it has to stay reachable by everything that can invalidate an offer. Holding + * the entries themselves — rather than a counter to compare against later — + * means `clear()` filters them by the SAME `canAccess(entry.owner)` predicate + * it already applies to the plan. A counter cannot do that, because the caller + * spending an entry need not be its owner: an admin may restore another user's + * sessions, and then the spender and the owner are different keys. */ - private generations = new Map(); + private parked = new Map(); /** * Owners with a restore in flight, between its take and its last pane. * Keyed by owner so one user's restore does not turn another user's click into @@ -70,17 +72,8 @@ export class RebootRestoreRegistry { set(entries: readonly RebootRestoreEntry[]): void { this.entries = new Map(entries.map((entry) => [entry.sessionId, entry])); this.builtAt = entries.length > 0 ? Date.now() : 0; - this.bumpAll(); - } - - /** - * The generations of the owners of `entries`, for a caller that will hand some - * of them back later. Pass the result to {@link restore}. - */ - snapshotGenerations(entries: readonly RebootRestoreEntry[]): Map { - const snapshot = new Map(); - for (const entry of entries) snapshot.set(entry.owner, this.generations.get(entry.owner) ?? 0); - return snapshot; + // A fresh boot plan supersedes anything an in-flight restore still holds. + this.parked.clear(); } /** @@ -104,7 +97,11 @@ export class RebootRestoreRegistry { * * @param sessionIds The ids to spend, or undefined for every visible entry. */ - take(canAccess: (owner: string | undefined) => boolean, sessionIds?: readonly string[]): RebootRestoreEntry[] { + take( + canAccess: (owner: string | undefined) => boolean, + sessionIds: readonly string[] | undefined, + spender: string | undefined + ): RebootRestoreEntry[] { this.dropIfExpired(); const wanted = sessionIds ? new Set(sessionIds) : undefined; const taken: RebootRestoreEntry[] = []; @@ -112,6 +109,9 @@ export class RebootRestoreRegistry { if (wanted && !wanted.has(entry.sessionId)) continue; if (!canAccess(entry.owner)) continue; this.entries.delete(entry.sessionId); + // Parked rather than forgotten: until this restore resolves the entry, a + // dismiss still has to be able to reach and cancel it. + this.parked.set(entry.sessionId, { entry, spender }); taken.push(entry); } return taken; @@ -126,18 +126,19 @@ export class RebootRestoreRegistry { * hand is NOT put back, because that one cannot stop being true, and an entry * the banner keeps re-offering forever is noise only Dismiss can clear. */ - restore(entries: readonly RebootRestoreEntry[], generations?: ReadonlyMap): void { + releaseFlight(spender: string | undefined, keep: readonly RebootRestoreEntry[]): void { + const wanted = new Set(keep.map((entry) => entry.sessionId)); let added = 0; - for (const entry of entries) { - // A dismiss (or a fresh boot plan) for THIS entry's owner since the caller - // took it means it is no longer wanted back. Another owner's dismiss is - // none of this entry's business. - if (generations) { - const taken = generations.get(entry.owner); - if (taken !== undefined && taken !== (this.generations.get(entry.owner) ?? 0)) continue; + for (const [sessionId, held] of [...this.parked]) { + if (held.spender !== spender) continue; + this.parked.delete(sessionId); + // Still parked means nothing cancelled it while the restore ran. A dismiss, + // an expiry or a fresh boot plan removes it from `parked`, and then it does + // not come back however the restore ended. + if (wanted.has(sessionId)) { + this.entries.set(sessionId, held.entry); + added += 1; } - this.entries.set(entry.sessionId, entry); - added += 1; } if (added > 0 && this.builtAt === 0) this.builtAt = Date.now(); } @@ -146,16 +147,17 @@ export class RebootRestoreRegistry { clear(canAccess: (owner: string | undefined) => boolean): number { const removable = [...this.entries.values()].filter((entry) => canAccess(entry.owner)); for (const entry of removable) this.entries.delete(entry.sessionId); + // Entries a restore is holding are dismissed by the same rule, so a dismiss + // that lands mid-restore wins. Judged on the ENTRY's owner, exactly as above, + // rather than on who happens to be restoring it. + let parkedRemoved = 0; + for (const [sessionId, held] of [...this.parked]) { + if (!canAccess(held.entry.owner)) continue; + this.parked.delete(sessionId); + parkedRemoved += 1; + } if (this.entries.size === 0) this.builtAt = 0; - // A restore in flight for these owners must not put their entries back. The - // in-flight owners are the ones that matter and the ones the plan can no - // longer name: `take()` has already removed their entries, so a dismiss that - // lands mid-restore sees nothing of theirs to remove. The bump is limited to - // owners this caller could see, so it cannot reach anyone else's restore. - const invalidated = new Set(removable.map((entry) => entry.owner)); - for (const owner of this.spending) if (canAccess(owner)) invalidated.add(owner); - for (const owner of invalidated) this.bump(owner); - return removable.length; + return removable.length + parkedRemoved; } /** @@ -176,27 +178,16 @@ export class RebootRestoreRegistry { /** Test hook: forget everything, including the single-flight claim. */ reset(): void { this.entries.clear(); + this.parked.clear(); this.builtAt = 0; this.spending.clear(); - this.generations.clear(); - } - - private bump(owner: string | undefined): void { - this.generations.set(owner, (this.generations.get(owner) ?? 0) + 1); - } - - /** Invalidate every owner's in-flight returns, including owners not yet seen. */ - private bumpAll(): void { - for (const owner of new Set([...this.entries.values()].map((entry) => entry.owner))) this.bump(owner); - for (const owner of [...this.generations.keys()]) this.bump(owner); } private dropIfExpired(): void { if (this.builtAt > 0 && Date.now() - this.builtAt > PLAN_TTL_MS) { - // Bump before clearing, while the owners are still known: a restore that - // took entries just before the expiry must not hand them back afterwards - // and give an expired plan another full day of life. - this.bumpAll(); + // A restore that took entries just before the expiry must not hand them + // back afterwards and give an expired plan another full day of life. + this.parked.clear(); this.entries.clear(); this.builtAt = 0; } diff --git a/src/web/routes/reboot-restore-routes.ts b/src/web/routes/reboot-restore-routes.ts index da742cc4..3516ff19 100644 --- a/src/web/routes/reboot-restore-routes.ts +++ b/src/web/routes/reboot-restore-routes.ts @@ -92,8 +92,7 @@ export function registerRebootRestoreRoutes(app: FastifyInstance, ctx: RebootRes if (!rebootRestoreRegistry.beginSpending(owner)) { return reply.code(409).send(createErrorResponse(ApiErrorCode.CONFLICT, 'A reboot restore is already running')); } - const taken = rebootRestoreRegistry.take(canAccess, body.sessionIds); - const generations = rebootRestoreRegistry.snapshotGenerations(taken); + const taken = rebootRestoreRegistry.take(canAccess, body.sessionIds, owner); // Entries nothing built a pane for, returned to the plan on every exit path // including a throw. Without this a failure between here and the loop would // spend the offer and rebuild nothing, and the plan cannot be rebuilt. @@ -252,9 +251,9 @@ export function registerRebootRestoreRoutes(app: FastifyInstance, ctx: RebootRes } finally { // Anything that never became a pane goes back on offer, including after a // throw, so a transient failure costs a retry rather than the whole plan. - // Passing the generations makes a Dismiss that landed mid-restore win, for - // the owners it actually covered. - rebootRestoreRegistry.restore([...unspent], generations); + // Ends the flight: entries still parked for it come back if they are in + // `unspent`, and a Dismiss that unparked them meanwhile wins. + rebootRestoreRegistry.releaseFlight(owner, [...unspent]); rebootRestoreRegistry.endSpending(owner); } }); diff --git a/src/web/server.ts b/src/web/server.ts index 095b094d..2358b2ff 100644 --- a/src/web/server.ts +++ b/src/web/server.ts @@ -3010,26 +3010,42 @@ export class WebServer extends EventEmitter { if (!session) return; this.sessions.delete(sessionId); - // --- the inverse of setupSessionListeners(), in its order --- - const summaryTracker = this.runSummaryTrackers.get(sessionId); - if (summaryTracker) { - summaryTracker.stop(); - this.runSummaryTrackers.delete(sessionId); - } - // An fs.watch on the workspace (or on @fix_plan.md) that nothing else closes. - session.ralphTracker.stopWatchingFixPlan(); - // An FSWatcher on the workspace, likewise. - imageWatcher.unwatchSession(sessionId); + // --- the inverse of setupSessionListeners(), in reverse order --- + // Listeners first: while they are attached, one of them can still reach a + // tracker this is about to stop. const listeners = this.sessionListenerRefs.get(sessionId); if (listeners) { detachSessionListeners(session, listeners); this.sessionListenerRefs.delete(sessionId); } + // An FSWatcher on the workspace that nothing else closes. + imageWatcher.unwatchSession(sessionId); + // An fs.watch on the workspace (or on @fix_plan.md), likewise. + session.ralphTracker.stopWatchingFixPlan(); + const summaryTracker = this.runSummaryTrackers.get(sessionId); + if (summaryTracker) { + summaryTracker.stop(); + this.runSummaryTrackers.delete(sessionId); + } + + // --- what anything else may have attached to this id in the meantime --- + // A rebuild can fail AFTER startInteractive() resolved, and a restored + // workspace still carries Codeman's hooks, so the CLI can post a hook event + // within milliseconds. Each of these outlives the listeners and would + // otherwise meet the retry, which reuses the same session id by design. + this.stopTranscriptWatcher(sessionId); + attachmentRegistry.clearSession(sessionId); + sessionWaits.notifySignal(sessionId, 'exit'); + sessionWaits.cancelAll(sessionId); + approvalInbox.resolveForSession(sessionId, 'session_ended'); // --- the inverse of the construction itself --- this.sse.cleanupSessionBatches(sessionId); this.persistDeb.cancelKey(sessionId); fileStreamManager.closeSessionStreams(sessionId); + // `lastRecordedTokens` is deliberately NOT deleted: the `after-spawn` phase + // seeds it as the daily-usage baseline for these restored totals, and the + // retry reuses the id, so dropping it would count them as new usage. // The per-session custom-model config dir carries the endpoint's API key, and // `before-spawn` may already have written it. Nothing else would ever remove // it: the stale sweep only touches state.json. A retry rewrites it. @@ -3039,6 +3055,9 @@ export class WebServer extends EventEmitter { await session.stop(true); } catch (err) { console.warn(`[Server] stopping a partially built session failed: ${getErrorMessage(err)}`); + // `stop()` kills the mux session in its last block, after destroying its + // trackers, so a throw on the way there leaves the pane running. + await this.mux.killSession(sessionId).catch(() => {}); } try { await this.tabLayouts.sessionsRemoved([{ id: sessionId, owner: session.owner }]); diff --git a/test/discard-partially-built-session.test.ts b/test/discard-partially-built-session.test.ts index 7232e969..012009f4 100644 --- a/test/discard-partially-built-session.test.ts +++ b/test/discard-partially-built-session.test.ts @@ -18,10 +18,11 @@ * leaves that entry makes the next attempt wire nothing at all, and the user * gets a tab that never shows output. */ -import { mkdirSync, rmSync } from 'node:fs'; +import { mkdirSync } from 'node:fs'; import { homedir } from 'node:os'; import { join } from 'node:path'; -import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from 'vitest'; +import { safeRmHomeTree } from './mocks/test-helpers.js'; import { WebServer } from '../src/web/server.js'; import { Session } from '../src/session.js'; @@ -55,9 +56,11 @@ function buildSession(): Session { }); } -beforeEach(() => { +beforeAll(() => { mkdirSync(WORKSPACE, { recursive: true }); - // Test mode: no port is opened and no CLI is launched. + // Test mode: no port is opened and no CLI is launched. One server for the file, + // stopped at the end: the constructor registers handlers on the module-level + // image, subagent, team and workflow watchers, and only stop() removes them. server = new WebServer(0, false, true); internals = server as unknown as ServerInternals; mux = new TmuxManager(); @@ -65,7 +68,11 @@ beforeEach(() => { afterEach(async () => { await internals.discardPartiallyBuiltSession(SESSION_ID).catch(() => {}); - rmSync(WORKSPACE, { recursive: true, force: true }); +}); + +afterAll(async () => { + await server.stop().catch(() => {}); + safeRmHomeTree(WORKSPACE); }); describe('discarding a session whose pane never started', () => { @@ -85,25 +92,36 @@ describe('discarding a session whose pane never started', () => { await internals.setupSessionListeners(first); expect(internals.sessionListenerRefs.has(SESSION_ID)).toBe(true); + const firstRefs = internals.sessionListenerRefs.get(SESSION_ID); await internals.discardPartiallyBuiltSession(SESSION_ID); expect(internals.sessionListenerRefs.has(SESSION_ID)).toBe(false); // The retry reuses the id by design. `setupSessionListeners()` returns early - // while the refs are still there, so a session built now would run blind: - // no terminal output, no status updates, no exit broadcast. + // while the refs are still there, so a session built now would run blind: no + // terminal output, no status updates, no exit broadcast. Asserting a DIFFERENT + // refs object is what distinguishes wiring the retry from finding the corpse + // of the first attempt still in place. const retry = buildSession(); await internals.registerSessionWithLayout(retry); await internals.setupSessionListeners(retry); - expect(internals.sessionListenerRefs.has(SESSION_ID)).toBe(true); + const retryRefs = internals.sessionListenerRefs.get(SESSION_ID); + expect(retryRefs).toBeDefined(); + expect(retryRefs).not.toBe(firstRefs); }); it('stops the run-summary tracker, whose interval would otherwise keep firing', async () => { const session = buildSession(); await internals.registerSessionWithLayout(session); await internals.setupSessionListeners(session); - expect(internals.runSummaryTrackers.has(SESSION_ID)).toBe(true); + const tracker = internals.runSummaryTrackers.get(SESSION_ID) as { stop: () => void }; + expect(tracker).toBeDefined(); + // Dropping the map entry is not enough: the tracker arms a setInterval in its + // constructor, and only stop() clears it, so a discard that merely forgot the + // entry would leave the timer running for the life of the process. + const stopped = vi.spyOn(tracker, 'stop'); await internals.discardPartiallyBuiltSession(SESSION_ID); + expect(stopped).toHaveBeenCalled(); expect(internals.runSummaryTrackers.has(SESSION_ID)).toBe(false); }); diff --git a/test/reboot-restore.test.ts b/test/reboot-restore.test.ts index 76506358..aaf9a83e 100644 --- a/test/reboot-restore.test.ts +++ b/test/reboot-restore.test.ts @@ -198,15 +198,15 @@ describe('the plan the banner spends', () => { it('hands an entry to the first caller and nothing to the second', () => { const registry = new RebootRestoreRegistry(); registry.set([entryFor('a'), entryFor('b')]); - expect(registry.take(all).map((e) => e.sessionId)).toEqual(['a', 'b']); + expect(registry.take(all, undefined, undefined).map((e) => e.sessionId)).toEqual(['a', 'b']); // The double-click: two panes on one conversation is what this prevents. - expect(registry.take(all)).toEqual([]); + expect(registry.take(all, undefined, undefined)).toEqual([]); }); it('spends only the ids a caller asked for', () => { const registry = new RebootRestoreRegistry(); registry.set([entryFor('a'), entryFor('b')]); - expect(registry.take(all, ['b']).map((e) => e.sessionId)).toEqual(['b']); + expect(registry.take(all, ['b'], undefined).map((e) => e.sessionId)).toEqual(['b']); expect(registry.list(all).map((e) => e.sessionId)).toEqual(['a']); }); @@ -215,7 +215,7 @@ describe('the plan the banner spends', () => { registry.set([entryFor('mine', 'alice'), entryFor('theirs', 'bob')]); const asAlice = (owner: string | undefined) => owner === 'alice'; expect(registry.list(asAlice).map((e) => e.sessionId)).toEqual(['mine']); - expect(registry.take(asAlice).map((e) => e.sessionId)).toEqual(['mine']); + expect(registry.take(asAlice, undefined, 'alice').map((e) => e.sessionId)).toEqual(['mine']); // Bob's entry is still on offer for Bob. expect(registry.list(() => true).map((e) => e.sessionId)).toEqual(['theirs']); }); @@ -223,8 +223,8 @@ describe('the plan the banner spends', () => { it('puts back an entry that no pane was created for', () => { const registry = new RebootRestoreRegistry(); registry.set([entryFor('a')]); - const taken = registry.take(all); - registry.restore(taken); + const taken = registry.take(all, undefined, undefined); + registry.releaseFlight(undefined, taken); expect(registry.list(all).map((e) => e.sessionId)).toEqual(['a']); }); diff --git a/test/routes/reboot-restore-rebuild-failure.test.ts b/test/routes/reboot-restore-rebuild-failure.test.ts index 9ee93fb4..7a69e30d 100644 --- a/test/routes/reboot-restore-rebuild-failure.test.ts +++ b/test/routes/reboot-restore-rebuild-failure.test.ts @@ -183,12 +183,23 @@ describe('a rebuild that succeeds', () => { callOrder.push(`reapply:${phase}`); } ); + (ctx.setupSessionListeners as ReturnType).mockImplementation(async () => { + callOrder.push('setupSessionListeners'); + }); const app = await createHarness(ctx); await app.inject({ method: 'POST', url: '/api/reboot-restore/restore', payload: {} }); - // The custom-model environment has to reach the process; the token totals - // must not land on a session whose pane never started. - expect(callOrder).toEqual(['reapply:before-spawn', 'startInteractive', 'reapply:after-spawn']); + // `setupSessionListeners()` READS the image-watcher flag that `before-spawn` + // restores, so the phase has to precede it or the session comes back + // reporting the watcher as on with nothing watching. The custom-model + // environment has to reach the process, and the token totals must not land + // on a session whose pane never started. + expect(callOrder).toEqual([ + 'reapply:before-spawn', + 'setupSessionListeners', + 'startInteractive', + 'reapply:after-spawn', + ]); await app.close(); }); @@ -282,6 +293,32 @@ describe('a failure before any entry is considered', () => { }); describe('a dismiss that lands while a restore is running', () => { + it('wins when an admin is restoring the entries and their owner dismisses', async () => { + const theirs = offerEntry('theirs', 'bob'); + rebootRestoreRegistry.set([theirs]); + // An admin may spend another user's entries, so the caller doing the restore + // and the owner of what is being restored are different people. + const taken = rebootRestoreRegistry.take(() => true, undefined, 'admin'); + expect(taken.map((e) => e.sessionId)).toEqual(['theirs']); + + // Bob dismisses his own banner. Nothing of his is in the plan any more, and + // the restore is running under a different name than his. + rebootRestoreRegistry.clear((owner) => owner === 'bob'); + rebootRestoreRegistry.releaseFlight('admin', taken); + + expect(rebootRestoreRegistry.list(() => true)).toEqual([]); + }); + + it('wins when an admin dismisses everything mid-restore', async () => { + rebootRestoreRegistry.set([offerEntry('theirs', 'bob')]); + const taken = rebootRestoreRegistry.take(() => true, undefined, 'admin'); + + rebootRestoreRegistry.clear(() => true); + rebootRestoreRegistry.releaseFlight('admin', taken); + + expect(rebootRestoreRegistry.list(() => true)).toEqual([]); + }); + it('wins, rather than being undone when the route hands its entries back', async () => { rebootRestoreRegistry.set([offerEntry('a')]); const ctx = createMockRouteContext({ workspaceHooksEnabled: false }); @@ -304,15 +341,13 @@ describe('a dismiss that lands while a restore is running', () => { it('reaches an in-flight restore the dismisser can see, even once its entries are taken', async () => { const mine = offerEntry('mine', 'alice'); rebootRestoreRegistry.set([mine]); - expect(rebootRestoreRegistry.beginSpending('alice')).toBe(true); - const generations = rebootRestoreRegistry.snapshotGenerations([mine]); - const taken = rebootRestoreRegistry.take((owner) => owner === 'alice'); + const taken = rebootRestoreRegistry.take((owner) => owner === 'alice', undefined, 'alice'); + expect(taken).toHaveLength(1); - // The plan is empty now, so a dismiss has nothing of Alice's to remove; the - // invalidation has to come from her claimed flight. + // The plan is empty now, so the dismiss has nothing of Alice's left in the + // plan; it has to reach the entry the restore is holding. rebootRestoreRegistry.clear((owner) => owner === 'alice'); - rebootRestoreRegistry.restore(taken, generations); - rebootRestoreRegistry.endSpending('alice'); + rebootRestoreRegistry.releaseFlight('alice', taken); expect(rebootRestoreRegistry.list(() => true)).toEqual([]); }); @@ -322,13 +357,8 @@ describe('a dismiss that lands while a restore is running', () => { const theirs = offerEntry('theirs', 'bob'); rebootRestoreRegistry.set([mine, theirs]); - // Bob is mid-restore, holding his own entry. The claimed flight is what makes - // this the interesting case: a dismiss can no longer see Bob's entries in the - // plan, so the invalidation has to come from the in-flight set, filtered by - // what the dismissing user may access. - expect(rebootRestoreRegistry.beginSpending('bob')).toBe(true); - const bobsGenerations = rebootRestoreRegistry.snapshotGenerations([theirs]); - const bobsTaken = rebootRestoreRegistry.take((owner) => owner === 'bob'); + // Bob is mid-restore, holding his own entry. + const bobsTaken = rebootRestoreRegistry.take((owner) => owner === 'bob', undefined, 'bob'); expect(bobsTaken.map((e) => e.sessionId)).toEqual(['theirs']); // Alice dismisses her own banner meanwhile. @@ -336,8 +366,7 @@ describe('a dismiss that lands while a restore is running', () => { // Bob's restore finishes and hands his entry back. Alice's dismiss covered // her entries, not his, so his offer survives. - rebootRestoreRegistry.restore(bobsTaken, bobsGenerations); - rebootRestoreRegistry.endSpending('bob'); + rebootRestoreRegistry.releaseFlight('bob', bobsTaken); expect(rebootRestoreRegistry.list(() => true).map((e) => e.sessionId)).toEqual(['theirs']); }); });