fix(coding-agent): persisted vibe sessions across restarts
Vibe worker roster lived only in a process-local Map, so a resumed parent session started with an empty registry and vibe_send failed with "Unknown vibe session". Persist a versioned, parent-scoped lifecycle journal (spawn/turn/tombstone events), rehydrate validated idle workers through the persisted-subagent reviver on resume, and gate the flow with generation/CAS protection so stale finalizers cannot clobber a replacement worker. Killed transcripts stay readable but non-revivable; mode-exit commits tombstones atomically with the mode change and rolls back cleanly on storage failure. Ported from @mastertyko's fork branch fix/vibe-session-persistence. Fixes #5303
This commit is contained in:
@@ -33,7 +33,7 @@ import type { MCPManager } from "../mcp/manager";
|
||||
import type { MnemopiSessionState } from "../mnemopi/state";
|
||||
import subagentSystemPromptTemplate from "../prompts/system/subagent-system-prompt.md" with { type: "text" };
|
||||
import submitReminderTemplate from "../prompts/system/subagent-yield-reminder.md" with { type: "text" };
|
||||
import { AgentLifecycleManager } from "../registry/agent-lifecycle";
|
||||
import { AgentLifecycleManager, type AgentReviver } from "../registry/agent-lifecycle";
|
||||
import { AgentRegistry } from "../registry/agent-registry";
|
||||
import { type CreateAgentSessionOptions, createAgentSession, discoverAuthStorage } from "../sdk";
|
||||
import type { AgentSession, AgentSessionEvent } from "../session/agent-session";
|
||||
@@ -1945,9 +1945,11 @@ export async function finalizeSubagentLifecycle(args: {
|
||||
keepAlive: boolean;
|
||||
isolated: boolean;
|
||||
agentIdleTtlMs: number;
|
||||
reviveSession: (() => Promise<AgentSession>) | null;
|
||||
reviveSession: AgentReviver | null;
|
||||
}): Promise<void> {
|
||||
const registry = AgentRegistry.global();
|
||||
const ref = registry.get(args.id);
|
||||
const ownsRef = Boolean(ref && ref.session === args.session);
|
||||
const disposeSession = async (): Promise<void> => {
|
||||
try {
|
||||
await untilAborted(AbortSignal.timeout(5000), () => args.session.dispose());
|
||||
@@ -1962,7 +1964,7 @@ export async function finalizeSubagentLifecycle(args: {
|
||||
const resumableAbort =
|
||||
args.abortKind === "budget" && args.keepAlive && !args.isolated && args.reviveSession !== null;
|
||||
if (args.aborted && !resumableAbort) {
|
||||
registry.setStatus(args.id, "aborted");
|
||||
if (ref && ownsRef) registry.setStatus(args.id, "aborted", ref);
|
||||
await disposeSession();
|
||||
return;
|
||||
}
|
||||
@@ -1970,7 +1972,7 @@ export async function finalizeSubagentLifecycle(args: {
|
||||
if (!args.keepAlive) {
|
||||
// One-shot helper: dispose and unregister. No IRC, no revival.
|
||||
await disposeSession();
|
||||
registry.unregister(args.id);
|
||||
if (ref && ownsRef) registry.unregister(args.id, ref);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -1980,19 +1982,26 @@ export async function finalizeSubagentLifecycle(args: {
|
||||
// transcript stays reachable (history://), but ensureLive will throw.
|
||||
// Status must flip to "parked" before dispose so the sdk dispose
|
||||
// wrapper skips unregister.
|
||||
registry.setStatus(args.id, "parked");
|
||||
if (ref && ownsRef) registry.setStatus(args.id, "parked", ref);
|
||||
await disposeSession();
|
||||
registry.detachSession(args.id);
|
||||
if (ref && ownsRef) registry.detachSession(args.id, ref);
|
||||
return;
|
||||
}
|
||||
|
||||
// Keep-alive: finished and failed subagents both stay interrogable.
|
||||
// The lifecycle manager owns idle-TTL parking + revival from here on.
|
||||
registry.setStatus(args.id, "idle");
|
||||
AgentLifecycleManager.global().adopt(args.id, {
|
||||
idleTtlMs: args.agentIdleTtlMs,
|
||||
revive: args.reviveSession ?? undefined,
|
||||
});
|
||||
if (!ref || !ownsRef || !registry.setStatus(args.id, "idle", ref)) {
|
||||
await disposeSession();
|
||||
return;
|
||||
}
|
||||
AgentLifecycleManager.global().adopt(
|
||||
args.id,
|
||||
{
|
||||
idleTtlMs: args.agentIdleTtlMs,
|
||||
revive: args.reviveSession ?? undefined,
|
||||
},
|
||||
ref,
|
||||
);
|
||||
}
|
||||
|
||||
/** Options for {@link runSubagentFollowUpTurn}. */
|
||||
@@ -2240,7 +2249,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
});
|
||||
const progress = monitor.progress;
|
||||
let unsubscribe: (() => void) | null = null;
|
||||
let reviveSession: (() => Promise<AgentSession>) | null = null;
|
||||
let reviveSession: AgentReviver | null = null;
|
||||
// Adopted (kept-alive) subagents flip registry status from session events on
|
||||
// later turns: revive/wake → running, turn drained → idle. The subscription
|
||||
// intentionally survives this run; a disposed session emits nothing, so it
|
||||
@@ -2248,9 +2257,9 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
const installRegistryStatusSync = (target: AgentSession): void => {
|
||||
target.subscribe(event => {
|
||||
if (event.type === "agent_start") {
|
||||
AgentRegistry.global().setStatus(id, "running");
|
||||
AgentRegistry.global().setStatus(id, "running", target);
|
||||
} else if (event.type === "agent_end") {
|
||||
AgentRegistry.global().setStatus(id, "idle");
|
||||
AgentRegistry.global().setStatus(id, "idle", target);
|
||||
}
|
||||
});
|
||||
};
|
||||
@@ -2426,7 +2435,10 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
// the same JSONL file re-invokes createAgentSession with the exact options
|
||||
// of the original run (same agent id, tools, model, system prompt,
|
||||
// artifacts dir) — only the SessionManager differs.
|
||||
const buildSubagentSessionOptions = (sessionManagerForRun: SessionManager): CreateAgentSessionOptions => ({
|
||||
const buildSubagentSessionOptions = (
|
||||
sessionManagerForRun: SessionManager,
|
||||
expectedAgentRef: CreateAgentSessionOptions["expectedAgentRef"],
|
||||
): CreateAgentSessionOptions => ({
|
||||
cwd: worktree ?? cwd,
|
||||
authStorage,
|
||||
modelRegistry,
|
||||
@@ -2474,6 +2486,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
parentAgentId: options.parentAgentId,
|
||||
agentId: id,
|
||||
agentDisplayName: agent.name,
|
||||
expectedAgentRef,
|
||||
enableLsp: lspEnabled,
|
||||
skipPythonPreflight,
|
||||
enableMCP,
|
||||
@@ -2487,7 +2500,7 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
},
|
||||
});
|
||||
|
||||
const sessionPromise = createAgentSession(buildSubagentSessionOptions(sessionManager));
|
||||
const sessionPromise = createAgentSession(buildSubagentSessionOptions(sessionManager, null));
|
||||
let session: AgentSession;
|
||||
try {
|
||||
({ session } = await awaitAbortable(sessionPromise));
|
||||
@@ -2507,14 +2520,16 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
// the single-writer lock cleanly and restores the full message history
|
||||
// (createAgentSession → agent.replaceMessages). Isolated runs are not
|
||||
// resumable (worktree is merged + cleaned) and never get a reviver.
|
||||
reviveSession = async () => {
|
||||
reviveSession = async expectedAgentRef => {
|
||||
const reopened = await SessionManager.open(sessionFile, undefined, undefined, {
|
||||
suppressBreadcrumb: true,
|
||||
});
|
||||
if (options.parentArtifactManager) {
|
||||
reopened.adoptArtifactManager(options.parentArtifactManager);
|
||||
}
|
||||
const { session: revived } = await createAgentSession(buildSubagentSessionOptions(reopened));
|
||||
const { session: revived } = await createAgentSession(
|
||||
buildSubagentSessionOptions(reopened, expectedAgentRef),
|
||||
);
|
||||
installRegistryStatusSync(revived);
|
||||
return revived;
|
||||
};
|
||||
|
||||
@@ -18,10 +18,11 @@ import { ADVISOR_TRANSCRIPT_STEM } from "../advisor/transcript-recorder";
|
||||
*
|
||||
* The first allocation of a given name keeps the name as-is; subsequent
|
||||
* allocations of the same name get a `-2`, `-3`, … suffix. On resume, scans
|
||||
* existing output files so previously written outputs are never overwritten.
|
||||
* existing output and child-session files so prior state is never overwritten.
|
||||
*/
|
||||
export class AgentOutputManager {
|
||||
#initialized = false;
|
||||
#initializing: Promise<void> | undefined;
|
||||
/** Final ids already handed out, relative to this manager's scope. */
|
||||
readonly #taken = new Set<string>();
|
||||
readonly #getArtifactsDir: () => string | null;
|
||||
@@ -42,8 +43,12 @@ export class AgentOutputManager {
|
||||
*/
|
||||
async #ensureInitialized(): Promise<void> {
|
||||
if (this.#initialized) return;
|
||||
this.#initializing ??= this.#seedFromDisk();
|
||||
await this.#initializing;
|
||||
this.#initialized = true;
|
||||
}
|
||||
|
||||
async #seedFromDisk(): Promise<void> {
|
||||
const dir = this.#getArtifactsDir();
|
||||
if (!dir) return;
|
||||
|
||||
@@ -56,8 +61,9 @@ export class AgentOutputManager {
|
||||
|
||||
const prefix = this.#parentPrefix ? `${this.#parentPrefix}.` : "";
|
||||
for (const file of files) {
|
||||
if (!file.endsWith(".md")) continue;
|
||||
let rest = file.slice(0, -3); // drop ".md"
|
||||
const extension = file.endsWith(".jsonl") ? ".jsonl" : file.endsWith(".md") ? ".md" : undefined;
|
||||
if (!extension) continue;
|
||||
let rest = file.slice(0, -extension.length);
|
||||
if (prefix) {
|
||||
if (!rest.startsWith(prefix)) continue;
|
||||
rest = rest.slice(prefix.length);
|
||||
@@ -80,6 +86,22 @@ export class AgentOutputManager {
|
||||
return this.#parentPrefix ? `${this.#parentPrefix}.${candidate}` : candidate;
|
||||
}
|
||||
|
||||
/** Reserve final IDs discovered outside the output directory scan. */
|
||||
async reserve(ids: Iterable<string>): Promise<void> {
|
||||
await this.#ensureInitialized();
|
||||
const prefix = this.#parentPrefix ? `${this.#parentPrefix}.` : "";
|
||||
for (const id of ids) {
|
||||
let rest = id;
|
||||
if (prefix) {
|
||||
if (!rest.startsWith(prefix)) continue;
|
||||
rest = rest.slice(prefix.length);
|
||||
}
|
||||
const dot = rest.indexOf(".");
|
||||
const segment = dot === -1 ? rest : rest.slice(0, dot);
|
||||
if (segment) this.#taken.add(segment);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Allocate a unique ID.
|
||||
*
|
||||
|
||||
@@ -71,7 +71,7 @@ export function createPersistedSubagentReviverFactory(
|
||||
taskDepth++;
|
||||
parentId = registry.get(parentId)?.parentId;
|
||||
}
|
||||
return async () => {
|
||||
return async expectedRef => {
|
||||
// Re-open fresh on every revive: park closes the writer, so this takes
|
||||
// the single-writer lock cleanly and restores the full message history.
|
||||
const reopened = await SessionManager.open(sessionFile, undefined, undefined, {
|
||||
@@ -96,6 +96,7 @@ export function createPersistedSubagentReviverFactory(
|
||||
agentDisplayName: ref.displayName,
|
||||
parentTaskPrefix: ref.id,
|
||||
parentAgentId: ref.parentId,
|
||||
expectedAgentRef: expectedRef,
|
||||
taskDepth,
|
||||
toolNames: init.tools,
|
||||
outputSchema: init.outputSchema,
|
||||
@@ -119,8 +120,8 @@ export function createPersistedSubagentReviverFactory(
|
||||
// Without it the idle-TTL timer never clears on a turn and the lifecycle
|
||||
// could park the agent mid-run.
|
||||
session.subscribe(event => {
|
||||
if (event.type === "agent_start") registry.setStatus(ref.id, "running");
|
||||
else if (event.type === "agent_end") registry.setStatus(ref.id, "idle");
|
||||
if (event.type === "agent_start") registry.setStatus(ref.id, "running", session);
|
||||
else if (event.type === "agent_end") registry.setStatus(ref.id, "idle", session);
|
||||
});
|
||||
return session;
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user