Archived
feat: add hierarchical session tree navigator
This commit is contained in:
@@ -19,12 +19,13 @@ import {
|
||||
type ModelRuntime,
|
||||
type ResourceDiagnostic,
|
||||
} from "@earendil-works/pi-coding-agent";
|
||||
import type { ClientArchiveSessionsResponse, ClientCommand, ClientCommandResult, ClientMessagePage, ClientSession, ClientSessionCleanupExecuteResponse, ClientSessionCleanupPreviewResponse, ClientSessionModel, ClientSessionStatus, ClientThinkingLevel, SessionStreamSnapshot, SessionUiEvent } from "../types.js";
|
||||
import type { ClientArchiveSessionsResponse, ClientCommand, ClientCommandResult, ClientMessagePage, ClientSession, ClientSessionCleanupExecuteResponse, ClientSessionCleanupPreviewResponse, ClientSessionModel, ClientSessionStatus, ClientSessionTreeNavigateRequest, ClientSessionTreeNavigateResult, ClientThinkingLevel, SessionStreamSnapshot, SessionUiEvent } from "../types.js";
|
||||
import { projectBrowserMessage } from "../browserMessageProjection.js";
|
||||
import { pageMessagesAtSafeBoundary } from "./messagePaging.js";
|
||||
import type { SessionEventHub } from "../realtime/sessionEventHub.js";
|
||||
import { BUILTIN_COMMANDS } from "./builtinCommands.js";
|
||||
import { SessionCommandService } from "./sessionCommandService.js";
|
||||
import { projectSessionTree, type ProjectableSessionTreeNode } from "./sessionTreeProjection.js";
|
||||
import { SessionArchiveStore, type ArchivedSessionRecord, type ArchiveSessionInput } from "./sessionArchiveStore.js";
|
||||
import { findArchiveCandidateByIdOrPrefix, planSessionArchiveTree, type SessionArchiveTreeCandidate } from "./sessionArchiveTree.js";
|
||||
import type { ActiveSession } from "./sessionRuntimeStore.js";
|
||||
@@ -32,6 +33,7 @@ import { deterministicSessionName, fallbackSessionName, generateShortSessionName
|
||||
import { computeEditPreview, type EditPreviewResult } from "./editPreview.js";
|
||||
import { attachmentsToInlineImages, saveAttachmentsToWorkspace } from "./attachmentService.js";
|
||||
import { parsePromptAttachments } from "../../shared/promptAttachments.js";
|
||||
import { SESSION_TREE_CUSTOM_INSTRUCTIONS_MAX_LENGTH } from "../../shared/apiTypes.js";
|
||||
import type {
|
||||
SavedPromptAttachment,
|
||||
SessionBulkArchiveResponse,
|
||||
@@ -106,6 +108,51 @@ interface QueuedPrompt {
|
||||
echoUserMessage?: boolean;
|
||||
}
|
||||
|
||||
interface DeferredSubsessionNotification {
|
||||
parentId: string;
|
||||
childId: string;
|
||||
text: string;
|
||||
}
|
||||
|
||||
interface TreeExclusiveOperationTarget {
|
||||
sessionId: string;
|
||||
session?: PiAgentSession;
|
||||
runtime?: PiSessionRuntime;
|
||||
}
|
||||
|
||||
type PiTreeNavigationOptions =
|
||||
| { summarize: false }
|
||||
| { summarize: true; customInstructions?: string };
|
||||
|
||||
function sessionTreeNavigationOptions(request: ClientSessionTreeNavigateRequest): PiTreeNavigationOptions {
|
||||
switch (request.summary.mode) {
|
||||
case "none":
|
||||
return { summarize: false };
|
||||
case "default":
|
||||
return { summarize: true };
|
||||
case "custom": {
|
||||
const customInstructions = request.summary.instructions.trim();
|
||||
if (customInstructions === "") throw new Error("Custom branch-summary instructions are required");
|
||||
if (customInstructions.length > SESSION_TREE_CUSTOM_INSTRUCTIONS_MAX_LENGTH) {
|
||||
throw new Error(`Custom branch-summary instructions must be at most ${String(SESSION_TREE_CUSTOM_INSTRUCTIONS_MAX_LENGTH)} characters`);
|
||||
}
|
||||
return { summarize: true, customInstructions };
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
function decrementWeakCount<Key extends object>(counts: WeakMap<Key, number>, key: Key): void {
|
||||
const remaining = (counts.get(key) ?? 1) - 1;
|
||||
if (remaining <= 0) counts.delete(key);
|
||||
else counts.set(key, remaining);
|
||||
}
|
||||
|
||||
function decrementMapCount<Key>(counts: Map<Key, number>, key: Key): void {
|
||||
const remaining = (counts.get(key) ?? 1) - 1;
|
||||
if (remaining <= 0) counts.delete(key);
|
||||
else counts.set(key, remaining);
|
||||
}
|
||||
|
||||
interface TrackedSubsessionLink {
|
||||
parentSessionId: string;
|
||||
childSessionId: string;
|
||||
@@ -197,6 +244,7 @@ export interface PiSessionManager {
|
||||
getSessionFile(): string | undefined;
|
||||
getBranch(): unknown[];
|
||||
getEntries?(): readonly unknown[];
|
||||
getTree?(): readonly ProjectableSessionTreeNode[];
|
||||
getLeafId(): string | null;
|
||||
getHeader?(): { parentSession?: string } | null | undefined;
|
||||
appendCustomEntry?(customType: string, data?: unknown): string;
|
||||
@@ -276,6 +324,8 @@ export interface PiAgentSession {
|
||||
prompt(text: string, options?: { streamingBehavior?: "steer" | "followUp"; images?: ImageContent[] }): Promise<void>;
|
||||
sendCustomMessage(message: { customType: string; content: string; display: boolean; details?: unknown }, options?: { triggerTurn?: boolean; deliverAs?: "steer" | "followUp" | "nextTurn" }): Promise<void>;
|
||||
executeBash(command: string, onChunk?: (chunk: string) => void, options?: { excludeFromContext?: boolean }): Promise<{ output: string; exitCode: number | undefined; cancelled: boolean; truncated: boolean; fullOutputPath?: string }>;
|
||||
navigateTree?(targetId: string, options?: { summarize?: boolean; customInstructions?: string }): Promise<{ editorText?: string; cancelled: boolean; aborted?: boolean; summaryEntry?: unknown }>;
|
||||
abortBranchSummary?(): void;
|
||||
abort(): Promise<void>;
|
||||
clearQueue(): { steering: string[]; followUp: string[] };
|
||||
getSteeringMessages(): readonly string[];
|
||||
@@ -605,6 +655,15 @@ export class PiSessionService implements SessionRouteService {
|
||||
private readonly activities = new Map<string, { phase: "active" | "idle" | "error"; label: string; detail?: string; at: string }>();
|
||||
private readonly heartbeat: NodeJS.Timeout;
|
||||
private readonly commandService: SessionCommandService<PiAgentSession>;
|
||||
/** Runtime-identity gate held while Pi may await abandoned-branch summarization. */
|
||||
private readonly treeNavigations = new WeakSet<PiAgentSession>();
|
||||
/** Counts async operations that may append an entry before they settle. */
|
||||
private readonly sessionEntryMutationCounts = new WeakMap<PiAgentSession, number>();
|
||||
/** Runtime/session-identity reservations for operations that must not overlap tree navigation. */
|
||||
private readonly treeExclusiveRuntimeOperationCounts = new WeakMap<PiSessionRuntime, number>();
|
||||
private readonly treeExclusiveSessionOperationCounts = new Map<string, number>();
|
||||
private readonly deferredSubsessionNotifications = new WeakMap<PiAgentSession, DeferredSubsessionNotification[]>();
|
||||
private readonly deferredGeneratedSessionNames = new WeakMap<PiAgentSession, string>();
|
||||
private readonly compactionPromptQueues = new Map<string, QueuedPrompt[]>();
|
||||
private readonly compactionDrainTimers = new Map<string, NodeJS.Timeout>();
|
||||
private readonly authLossWarnings = new Set<string>();
|
||||
@@ -667,14 +726,27 @@ export class PiSessionService implements SessionRouteService {
|
||||
events,
|
||||
{
|
||||
onCompactionStart: (session) => {
|
||||
this.beginSessionEntryMutation(session, "compact the session");
|
||||
this.publishActivity(session, "compacting", "active");
|
||||
this.publishStatus(session);
|
||||
},
|
||||
onCompactionEnd: (session, result, detail) => {
|
||||
this.endSessionEntryMutation(session);
|
||||
this.publishActivity(session, result === "success" ? "compaction complete" : "compaction failed", result === "success" ? "idle" : "error", detail);
|
||||
this.publishStatus(session);
|
||||
},
|
||||
reloadSession: (session) => this.reloadSessionRuntime(session),
|
||||
getSessionTree: (session) => {
|
||||
if (typeof session.sessionManager.getTree !== "function" || typeof session.navigateTree !== "function") return undefined;
|
||||
return projectSessionTree(session.sessionManager.getTree(), session.sessionManager.getLeafId());
|
||||
},
|
||||
hasActiveWork: (session) => this.hasActiveWork(session),
|
||||
isTreeNavigationActive: (session) => this.treeNavigations.has(session),
|
||||
runSessionReplacement: (session, operation) => this.runTreeExclusiveOperation(
|
||||
[{ sessionId: session.sessionId, session }],
|
||||
"Stop current session activity before replacing the session",
|
||||
operation,
|
||||
),
|
||||
},
|
||||
{ listSessionNames: (cwd) => this.listSessionNames(cwd) },
|
||||
);
|
||||
@@ -789,7 +861,7 @@ export class PiSessionService implements SessionRouteService {
|
||||
active.runtime.setRebindSession(undefined);
|
||||
this.workspaceActivity?.removeSession(active.runtime.session.sessionId, active.runtime.session.sessionManager.getCwd());
|
||||
try {
|
||||
await active.runtime.session.abort();
|
||||
await this.abortSessionOperations(active.runtime.session);
|
||||
} finally {
|
||||
await active.runtime.dispose();
|
||||
}
|
||||
@@ -1197,19 +1269,33 @@ export class PiSessionService implements SessionRouteService {
|
||||
private async notifyParentOfSubsession(parentId: string, childId: string, text: string): Promise<void> {
|
||||
try {
|
||||
const session = await this.getOrOpenParentForSubsession(parentId, childId);
|
||||
await session.sendCustomMessage(
|
||||
{ customType: SUBSESSION_NOTIFICATION_CUSTOM_TYPE, content: text, display: true, details: { sessionId: childId } },
|
||||
{ triggerTurn: true, deliverAs: "followUp" },
|
||||
);
|
||||
this.publishStatus(session);
|
||||
if (this.treeNavigations.has(session)) {
|
||||
const pending = this.deferredSubsessionNotifications.get(session) ?? [];
|
||||
pending.push({ parentId, childId, text });
|
||||
this.deferredSubsessionNotifications.set(session, pending);
|
||||
return;
|
||||
}
|
||||
await this.deliverSubsessionNotification(session, { parentId, childId, text });
|
||||
} catch (error: unknown) {
|
||||
this.logger.info(
|
||||
{ parentSessionId: parentId, sessionId: childId, error: error instanceof Error ? error.message : String(error) },
|
||||
"failed to notify parent of subsession completion",
|
||||
);
|
||||
this.logSubsessionNotificationFailure(parentId, childId, error);
|
||||
}
|
||||
}
|
||||
|
||||
private async deliverSubsessionNotification(session: PiAgentSession, notification: DeferredSubsessionNotification): Promise<void> {
|
||||
await this.runSessionEntryMutation(session, "deliver a subsession notification", () => session.sendCustomMessage(
|
||||
{ customType: SUBSESSION_NOTIFICATION_CUSTOM_TYPE, content: notification.text, display: true, details: { sessionId: notification.childId } },
|
||||
{ triggerTurn: true, deliverAs: "followUp" },
|
||||
));
|
||||
this.publishStatus(session);
|
||||
}
|
||||
|
||||
private logSubsessionNotificationFailure(parentId: string, childId: string, error: unknown): void {
|
||||
this.logger.info(
|
||||
{ parentSessionId: parentId, sessionId: childId, error: error instanceof Error ? error.message : String(error) },
|
||||
"failed to notify parent of subsession completion",
|
||||
);
|
||||
}
|
||||
|
||||
async messages(ref: PiSessionLookup, page?: { before?: number; limit?: number }): Promise<unknown[] | ClientMessagePage> {
|
||||
const session = await this.getOrOpen(ref);
|
||||
return pageMessagesAtSafeBoundary(historyMessages(session), page);
|
||||
@@ -1251,14 +1337,16 @@ export class PiSessionService implements SessionRouteService {
|
||||
async setModel(ref: PiSessionLookup, provider: string, modelId: string): Promise<ClientSessionStatus> {
|
||||
await this.assertWritable(ref);
|
||||
const session = await this.getOrOpen(ref);
|
||||
this.assertTreeNavigationInactive(session, "change models");
|
||||
await session.modelRuntime.reloadConfig();
|
||||
this.assertTreeNavigationInactive(session, "change models");
|
||||
const candidates = session.scopedModels.length > 0
|
||||
? session.scopedModels.map((scoped) => scoped.model)
|
||||
: session.modelRuntime.getAvailableSnapshot();
|
||||
const model = candidates.find((candidate) => candidate.provider === provider && candidate.id === modelId)
|
||||
?? session.modelRuntime.getModel(provider, modelId);
|
||||
if (model === undefined) throw new Error(`Model not found: ${provider}/${modelId}`);
|
||||
await session.setModel(model);
|
||||
await this.runSessionEntryMutation(session, "change models", () => session.setModel(model));
|
||||
this.publishActivity(session, `model: ${model.id}`, "idle", model.provider);
|
||||
this.publishStatus(session);
|
||||
return this.statusFromSession(session);
|
||||
@@ -1267,7 +1355,7 @@ export class PiSessionService implements SessionRouteService {
|
||||
async cycleModel(ref: PiSessionLookup, direction: "forward" | "backward"): Promise<ClientSessionStatus> {
|
||||
await this.assertWritable(ref);
|
||||
const session = await this.getOrOpen(ref);
|
||||
const result = await session.cycleModel(direction);
|
||||
const result = await this.runSessionEntryMutation(session, "change models", () => session.cycleModel(direction));
|
||||
if (result === undefined) throw new Error(session.scopedModels.length > 0 ? "Only one model in scope" : "Only one model available");
|
||||
this.publishActivity(session, `model: ${result.model.id}`, "idle", result.model.provider);
|
||||
this.publishStatus(session);
|
||||
@@ -1282,6 +1370,7 @@ export class PiSessionService implements SessionRouteService {
|
||||
async setThinkingLevel(ref: PiSessionLookup, level: string): Promise<ClientSessionStatus> {
|
||||
await this.assertWritable(ref);
|
||||
const session = await this.getOrOpen(ref);
|
||||
this.assertTreeNavigationInactive(session, "change the thinking level");
|
||||
// pi owns the valid set; validate against the session's live levels rather
|
||||
// than a hardcoded union so this stays correct if pi changes the set.
|
||||
const available = session.getAvailableThinkingLevels();
|
||||
@@ -1296,6 +1385,7 @@ export class PiSessionService implements SessionRouteService {
|
||||
async cycleThinkingLevel(ref: PiSessionLookup): Promise<ClientSessionStatus> {
|
||||
await this.assertWritable(ref);
|
||||
const session = await this.getOrOpen(ref);
|
||||
this.assertTreeNavigationInactive(session, "change the thinking level");
|
||||
const level = session.cycleThinkingLevel();
|
||||
if (level === undefined) throw new Error("Current model does not support thinking");
|
||||
this.publishActivity(session, `thinking: ${level}`, "idle");
|
||||
@@ -1330,6 +1420,7 @@ export class PiSessionService implements SessionRouteService {
|
||||
const images = (await attachmentsToInlineImages(parsedAttachments)).map((entry) => entry.image);
|
||||
await this.assertWritable(ref);
|
||||
const session = await this.getOrOpen(ref);
|
||||
this.assertTreeNavigationInactive(session, "send a prompt");
|
||||
this.maybeGenerateSessionName(session, promptText);
|
||||
const isQueued = session.isStreaming || session.isCompacting;
|
||||
const behavior = isQueued ? requestedBehavior ?? "followUp" : undefined;
|
||||
@@ -1349,7 +1440,7 @@ export class PiSessionService implements SessionRouteService {
|
||||
this.publishActivity(session, behavior === "steer" ? "steering queued" : behavior === "followUp" ? "message queued" : "prompt accepted", "active");
|
||||
if (behavior === undefined && echoUserMessage) this.events.publish(session.sessionId, { type: "message.append", message: userMessage(text, images) });
|
||||
const promptOptions = buildPromptOptions(behavior, images);
|
||||
const promptPromise = session.prompt(text, promptOptions).catch((error: unknown) => {
|
||||
const promptPromise = this.runSessionEntryMutation(session, "send a prompt", () => session.prompt(text, promptOptions)).catch((error: unknown) => {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
this.publishActivity(session, "error", "error", message);
|
||||
this.events.publish(session.sessionId, { type: "session.error", message });
|
||||
@@ -1378,6 +1469,7 @@ export class PiSessionService implements SessionRouteService {
|
||||
await this.assertWritable(ref);
|
||||
const active = await this.getActive(ref);
|
||||
const { session } = active.runtime;
|
||||
this.assertTreeNavigationInactive(session, "run a shell command");
|
||||
const isExcluded = text.startsWith("!!");
|
||||
const command = (isExcluded ? text.slice(2) : text.slice(1)).trim();
|
||||
if (!command) throw new Error("Usage: !<shell command>");
|
||||
@@ -1385,11 +1477,11 @@ export class PiSessionService implements SessionRouteService {
|
||||
|
||||
this.publishActivity(session, "running bash", "active", command);
|
||||
this.events.publish(session.sessionId, { type: "shell.start", command, excludeFromContext: isExcluded });
|
||||
void session.executeBash(command, (chunk) => {
|
||||
void this.runSessionEntryMutation(session, "run a shell command", () => session.executeBash(command, (chunk) => {
|
||||
this.events.publish(session.sessionId, { type: "shell.chunk", chunk });
|
||||
this.publishActivity(session, "running bash", "active", command);
|
||||
this.publishStatus(session);
|
||||
}, { excludeFromContext: isExcluded }).then((result) => {
|
||||
}, { excludeFromContext: isExcluded })).then((result) => {
|
||||
this.events.publish(session.sessionId, {
|
||||
type: "shell.end",
|
||||
output: result.output,
|
||||
@@ -1421,43 +1513,105 @@ export class PiSessionService implements SessionRouteService {
|
||||
return this.commandService.respond(active.runtime.session.sessionId, requestId, value);
|
||||
}
|
||||
|
||||
async navigateTree(ref: PiSessionLookup, request: ClientSessionTreeNavigateRequest): Promise<ClientSessionTreeNavigateResult> {
|
||||
if (request.targetId.trim() === "") throw new Error("Session tree target is required");
|
||||
if (this.isTreeExclusiveSessionIdentityActive(sessionIdFromLookup(ref))) {
|
||||
throw new Error("Stop current session activity before navigating the session tree");
|
||||
}
|
||||
await this.assertWritable(ref);
|
||||
const options = sessionTreeNavigationOptions(request);
|
||||
const session = await this.getOrOpen(ref);
|
||||
if (typeof session.navigateTree !== "function") throw new Error("Session tree navigation is not supported by this Pi runtime");
|
||||
if (this.hasActiveWork(session)) throw new Error("Stop current session activity before navigating the session tree");
|
||||
|
||||
// Acquire synchronously after the active-work check. No leaf-producing work
|
||||
// may enter this runtime until Pi's potentially asynchronous summary settles.
|
||||
this.treeNavigations.add(session);
|
||||
try {
|
||||
if (session.sessionManager.getLeafId() !== request.expectedLeafId) {
|
||||
throw new Error("The session changed since /tree was opened. Reopen /tree and try again.");
|
||||
}
|
||||
|
||||
this.publishActivity(session, options.summarize ? "summarizing branch" : "navigating session tree", "active");
|
||||
this.publishStatus(session);
|
||||
const result = await session.navigateTree(request.targetId, options);
|
||||
if (result.cancelled) {
|
||||
if (this.isCurrentActiveSession(session)) {
|
||||
this.publishActivity(session, result.aborted === true ? "branch summary aborted" : "tree navigation cancelled", "idle");
|
||||
}
|
||||
return { cancelled: true, ...(result.aborted === undefined ? {} : { aborted: result.aborted }) };
|
||||
}
|
||||
|
||||
if (this.isCurrentActiveSession(session)) this.publishActivity(session, "session tree navigated", "idle");
|
||||
return { cancelled: false, ...(result.editorText === undefined ? {} : { editorText: result.editorText }) };
|
||||
} catch (error: unknown) {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
if (this.isCurrentActiveSession(session)) {
|
||||
this.publishActivity(session, "tree navigation failed", "error", message);
|
||||
this.events.publish(session.sessionId, { type: "session.error", message });
|
||||
}
|
||||
throw error;
|
||||
} finally {
|
||||
this.treeNavigations.delete(session);
|
||||
if (this.isCurrentActiveSession(session)) {
|
||||
this.flushDeferredTreeNavigationWork(session);
|
||||
this.publishStatus(session);
|
||||
} else {
|
||||
this.deferredGeneratedSessionNames.delete(session);
|
||||
this.deferredSubsessionNotifications.delete(session);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private async reloadSessionRuntime(session: PiAgentSession): Promise<void> {
|
||||
if (this.hasActiveWork(session)) throw new Error("Stop current session activity before reloading");
|
||||
this.publishActivity(session, "reloading resources", "active");
|
||||
const priorGeneration = this.notificationGenerationBySession.get(session);
|
||||
let candidateGeneration: SessionNotificationGeneration | undefined;
|
||||
try {
|
||||
await session.reload(priorGeneration === undefined ? undefined : {
|
||||
beforeSessionStart: () => {
|
||||
candidateGeneration = this.notificationStore.beginReplacement(priorGeneration, notificationIdentityForSession(session));
|
||||
this.notificationGenerationBySession.set(session, candidateGeneration);
|
||||
this.replaceSessionNotificationContext(session, candidateGeneration);
|
||||
},
|
||||
});
|
||||
if (candidateGeneration !== undefined) {
|
||||
this.publishNotificationMutations(this.notificationStore.commitReplacement(candidateGeneration));
|
||||
}
|
||||
this.publishActivity(session, "resources reloaded", "idle");
|
||||
this.publishStatus(session);
|
||||
} catch (error: unknown) {
|
||||
if (candidateGeneration !== undefined) {
|
||||
this.publishNotificationMutations(this.notificationStore.abortReplacement(candidateGeneration, "candidate"));
|
||||
this.notificationGenerationBySession.set(session, candidateGeneration);
|
||||
}
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
this.publishActivity(session, "reload failed", "error", message);
|
||||
this.events.publish(session.sessionId, { type: "session.error", message });
|
||||
this.publishStatus(session);
|
||||
throw error;
|
||||
}
|
||||
await this.runTreeExclusiveOperation(
|
||||
[{ sessionId: session.sessionId, session }],
|
||||
"Stop current session activity before reloading",
|
||||
async () => {
|
||||
this.publishActivity(session, "reloading resources", "active");
|
||||
const priorGeneration = this.notificationGenerationBySession.get(session);
|
||||
let candidateGeneration: SessionNotificationGeneration | undefined;
|
||||
try {
|
||||
await session.reload(priorGeneration === undefined ? undefined : {
|
||||
beforeSessionStart: () => {
|
||||
candidateGeneration = this.notificationStore.beginReplacement(priorGeneration, notificationIdentityForSession(session));
|
||||
this.notificationGenerationBySession.set(session, candidateGeneration);
|
||||
this.replaceSessionNotificationContext(session, candidateGeneration);
|
||||
},
|
||||
});
|
||||
if (candidateGeneration !== undefined) {
|
||||
this.publishNotificationMutations(this.notificationStore.commitReplacement(candidateGeneration));
|
||||
}
|
||||
this.publishActivity(session, "resources reloaded", "idle");
|
||||
this.publishStatus(session);
|
||||
} catch (error: unknown) {
|
||||
if (candidateGeneration !== undefined) {
|
||||
this.publishNotificationMutations(this.notificationStore.abortReplacement(candidateGeneration, "candidate"));
|
||||
this.notificationGenerationBySession.set(session, candidateGeneration);
|
||||
}
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
this.publishActivity(session, "reload failed", "error", message);
|
||||
this.events.publish(session.sessionId, { type: "session.error", message });
|
||||
this.publishStatus(session);
|
||||
throw error;
|
||||
}
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
async archive(ref: PiSessionLookup): Promise<void> {
|
||||
const session = await this.getOrOpen(ref);
|
||||
if (this.hasActiveWork(session)) throw new Error("Stop current session activity before archiving");
|
||||
const archiveInput = await this.archiveInputForSession(session);
|
||||
await this.closeActive(session.sessionId, { kind: "clear", reason: "archive" });
|
||||
await this.archiveStore.archive(archiveInput);
|
||||
await this.runTreeExclusiveOperation(
|
||||
[{ sessionId: session.sessionId, session }],
|
||||
"Stop current session activity before archiving",
|
||||
async () => {
|
||||
const archiveInput = await this.archiveInputForSession(session);
|
||||
await this.closeActive(session.sessionId, { kind: "clear", reason: "archive" });
|
||||
await this.archiveStore.archive(archiveInput);
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
async archiveMany(refs: readonly SessionBulkMutationRef[]): Promise<SessionBulkArchiveResponse> {
|
||||
@@ -1499,23 +1653,42 @@ export class PiSessionService implements SessionRouteService {
|
||||
}
|
||||
}
|
||||
|
||||
const readyInputs: ArchiveSessionInput[] = [];
|
||||
const readyPlanItems: { input: ArchiveSessionInput; active?: ActiveSession<PiSessionRuntime> }[] = [];
|
||||
for (const item of planItems) {
|
||||
try {
|
||||
await this.closeActive(item.input.sessionId, { kind: "clear", reason: "archive" });
|
||||
readyInputs.push(item.input);
|
||||
} catch (error: unknown) {
|
||||
failures.push({ sessionId: item.input.sessionId, error: errorMessage(error) });
|
||||
const active = this.activeForLookup({ id: item.input.sessionId, cwd: item.input.cwd });
|
||||
if (active !== undefined && this.hasActiveWork(active.runtime.session)) {
|
||||
failures.push({ sessionId: item.input.sessionId, error: "Stop current session activity before archiving" });
|
||||
continue;
|
||||
}
|
||||
readyPlanItems.push(active === undefined ? item : { ...item, active });
|
||||
}
|
||||
|
||||
const readyInputs: ArchiveSessionInput[] = [];
|
||||
const archivedSessionIds = [...alreadyArchivedSessionIds];
|
||||
try {
|
||||
const archived = await this.archiveStoreArchiveMany(readyInputs);
|
||||
archivedSessionIds.push(...archived.map((record) => record.sessionId));
|
||||
} catch (error: unknown) {
|
||||
for (const input of readyInputs) failures.push({ sessionId: input.sessionId, error: errorMessage(error) });
|
||||
}
|
||||
await this.runTreeExclusiveOperation(
|
||||
readyPlanItems.map(({ input, active }) => ({
|
||||
sessionId: input.sessionId,
|
||||
...(active === undefined ? {} : { session: active.runtime.session, runtime: active.runtime }),
|
||||
})),
|
||||
"Stop current session activity before archiving",
|
||||
async () => {
|
||||
for (const item of readyPlanItems) {
|
||||
try {
|
||||
await this.closeActive(item.input.sessionId, { kind: "clear", reason: "archive" });
|
||||
readyInputs.push(item.input);
|
||||
} catch (error: unknown) {
|
||||
failures.push({ sessionId: item.input.sessionId, error: errorMessage(error) });
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
const archived = await this.archiveStoreArchiveMany(readyInputs);
|
||||
archivedSessionIds.push(...archived.map((record) => record.sessionId));
|
||||
} catch (error: unknown) {
|
||||
for (const input of readyInputs) failures.push({ sessionId: input.sessionId, error: errorMessage(error) });
|
||||
}
|
||||
},
|
||||
);
|
||||
|
||||
return {
|
||||
archived: true,
|
||||
@@ -1533,12 +1706,21 @@ export class PiSessionService implements SessionRouteService {
|
||||
const busy = plan.targets.map((target) => target.activeSession).find((target) => target !== undefined && this.hasActiveWork(target));
|
||||
if (busy !== undefined) throw new Error(`Stop current session activity before archiving ${sessionDisplayName(busy)}`);
|
||||
|
||||
for (const target of plan.targets) {
|
||||
if (target.archived) this.publishNotificationMutations(this.notificationStore.clearSession(target.id, "archive"));
|
||||
}
|
||||
const archiveInputs = plan.unarchivedTargets.map((target) => archiveInputFromCandidate(target));
|
||||
for (const input of archiveInputs) await this.closeActive(input.sessionId, { kind: "clear", reason: "archive" });
|
||||
await this.archiveStoreArchiveMany(archiveInputs);
|
||||
await this.runTreeExclusiveOperation(
|
||||
plan.unarchivedTargets.map((target) => ({
|
||||
sessionId: target.id,
|
||||
...(target.activeSession === undefined ? {} : { session: target.activeSession }),
|
||||
})),
|
||||
`Stop current session activity before archiving ${sessionDisplayName(session)}`,
|
||||
async () => {
|
||||
for (const target of plan.targets) {
|
||||
if (target.archived) this.publishNotificationMutations(this.notificationStore.clearSession(target.id, "archive"));
|
||||
}
|
||||
for (const input of archiveInputs) await this.closeActive(input.sessionId, { kind: "clear", reason: "archive" });
|
||||
await this.archiveStoreArchiveMany(archiveInputs);
|
||||
},
|
||||
);
|
||||
|
||||
return {
|
||||
archived: true,
|
||||
@@ -1625,28 +1807,35 @@ export class PiSessionService implements SessionRouteService {
|
||||
const session = await this.getOrOpen(ref);
|
||||
if (this.hasActiveWork(session)) throw new Error("Stop current session activity before reloading");
|
||||
|
||||
const priorGeneration = this.notificationGenerationBySession.get(session);
|
||||
const { sessionId, cwd } = notificationIdentityForSession(session);
|
||||
let candidateGeneration: SessionNotificationGeneration | undefined;
|
||||
try {
|
||||
await this.closeActive(
|
||||
sessionId,
|
||||
priorGeneration === undefined ? CLEAR_RUNTIME_NOTIFICATIONS : DEFER_RUNTIME_NOTIFICATIONS,
|
||||
);
|
||||
candidateGeneration = priorGeneration === undefined
|
||||
? undefined
|
||||
: this.notificationStore.beginReplacement(priorGeneration, { sessionId, cwd });
|
||||
const reopened = await this.getActive(ref, candidateGeneration === undefined ? {} : { notificationGeneration: candidateGeneration });
|
||||
if (candidateGeneration !== undefined) {
|
||||
this.publishNotificationMutations(this.notificationStore.commitReplacement(candidateGeneration));
|
||||
}
|
||||
this.publishStatus(reopened.runtime.session);
|
||||
} catch (error: unknown) {
|
||||
if (candidateGeneration !== undefined) {
|
||||
this.publishNotificationMutations(this.notificationStore.abortReplacement(candidateGeneration));
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
const reopenedSession = await this.runTreeExclusiveOperation(
|
||||
[{ sessionId: session.sessionId, session }],
|
||||
"Stop current session activity before reloading",
|
||||
async () => {
|
||||
const priorGeneration = this.notificationGenerationBySession.get(session);
|
||||
const { sessionId, cwd } = notificationIdentityForSession(session);
|
||||
let candidateGeneration: SessionNotificationGeneration | undefined;
|
||||
try {
|
||||
await this.closeActive(
|
||||
sessionId,
|
||||
priorGeneration === undefined ? CLEAR_RUNTIME_NOTIFICATIONS : DEFER_RUNTIME_NOTIFICATIONS,
|
||||
);
|
||||
candidateGeneration = priorGeneration === undefined
|
||||
? undefined
|
||||
: this.notificationStore.beginReplacement(priorGeneration, { sessionId, cwd });
|
||||
const reopened = await this.getActive(ref, candidateGeneration === undefined ? {} : { notificationGeneration: candidateGeneration });
|
||||
if (candidateGeneration !== undefined) {
|
||||
this.publishNotificationMutations(this.notificationStore.commitReplacement(candidateGeneration));
|
||||
}
|
||||
return reopened.runtime.session;
|
||||
} catch (error: unknown) {
|
||||
if (candidateGeneration !== undefined) {
|
||||
this.publishNotificationMutations(this.notificationStore.abortReplacement(candidateGeneration));
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
},
|
||||
);
|
||||
this.publishStatus(reopenedSession);
|
||||
}
|
||||
|
||||
async detachParent(ref: PiSessionLookup): Promise<void> {
|
||||
@@ -1680,9 +1869,16 @@ export class PiSessionService implements SessionRouteService {
|
||||
const sessionId = active.runtime.session.sessionId;
|
||||
this.clearCompactionPromptQueue(sessionId);
|
||||
clearSessionQueue(active.runtime.session);
|
||||
await active.runtime.session.abort();
|
||||
this.publishActivity(active.runtime.session, "stopped", "idle");
|
||||
this.publishStatus(active.runtime.session);
|
||||
try {
|
||||
await this.abortSessionOperations(active.runtime.session);
|
||||
this.publishActivity(active.runtime.session, "stopped", "idle");
|
||||
} catch (error: unknown) {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
this.publishActivity(active.runtime.session, "stop failed", "error", message);
|
||||
throw error;
|
||||
} finally {
|
||||
this.publishStatus(active.runtime.session);
|
||||
}
|
||||
}
|
||||
|
||||
async stop(ref: PiSessionLookup): Promise<void> {
|
||||
@@ -1906,12 +2102,33 @@ export class PiSessionService implements SessionRouteService {
|
||||
active.unsubscribe();
|
||||
active.runtime.setRebindSession(undefined);
|
||||
try {
|
||||
await active.runtime.session.abort();
|
||||
await this.abortSessionOperations(active.runtime.session);
|
||||
} finally {
|
||||
await active.runtime.dispose();
|
||||
}
|
||||
}
|
||||
|
||||
private async abortSessionOperations(session: PiAgentSession): Promise<void> {
|
||||
let branchSummaryAbortFailed = false;
|
||||
let branchSummaryAbortError: unknown;
|
||||
try {
|
||||
session.abortBranchSummary?.();
|
||||
} catch (error: unknown) {
|
||||
branchSummaryAbortFailed = true;
|
||||
branchSummaryAbortError = error;
|
||||
}
|
||||
|
||||
try {
|
||||
await session.abort();
|
||||
} catch (abortError: unknown) {
|
||||
if (branchSummaryAbortFailed) {
|
||||
throw new AggregateError([branchSummaryAbortError, abortError], "Failed to abort session operations", { cause: abortError });
|
||||
}
|
||||
throw abortError;
|
||||
}
|
||||
if (branchSummaryAbortFailed) throw branchSummaryAbortError;
|
||||
}
|
||||
|
||||
private async assertWritable(ref: PiSessionLookup): Promise<void> {
|
||||
if (await this.getArchived(ref) !== undefined) throw new Error("Archived sessions are read-only. Restore the session to continue.");
|
||||
}
|
||||
@@ -1979,6 +2196,10 @@ export class PiSessionService implements SessionRouteService {
|
||||
return archived;
|
||||
}
|
||||
|
||||
private isCurrentActiveSession(session: PiAgentSession): boolean {
|
||||
return this.active.get(session.sessionId)?.runtime.session === session;
|
||||
}
|
||||
|
||||
private activeForLookup(ref: PiSessionLookup): ActiveSession<PiSessionRuntime> | undefined {
|
||||
const sessionId = sessionIdFromLookup(ref);
|
||||
const exact = this.active.get(sessionId);
|
||||
@@ -2260,10 +2481,37 @@ export class PiSessionService implements SessionRouteService {
|
||||
|
||||
private applyGeneratedSessionName(session: PiAgentSession, name: string | undefined): void {
|
||||
if (name === undefined || session.sessionName !== undefined) return;
|
||||
if (this.treeNavigations.has(session)) {
|
||||
this.deferredGeneratedSessionNames.set(session, name);
|
||||
return;
|
||||
}
|
||||
session.setSessionName(name);
|
||||
this.publishSessionName(session);
|
||||
}
|
||||
|
||||
private flushDeferredTreeNavigationWork(session: PiAgentSession): void {
|
||||
const generatedName = this.deferredGeneratedSessionNames.get(session);
|
||||
this.deferredGeneratedSessionNames.delete(session);
|
||||
if (generatedName !== undefined) {
|
||||
try {
|
||||
this.applyGeneratedSessionName(session, generatedName);
|
||||
} catch (error: unknown) {
|
||||
this.logger.info(
|
||||
{ sessionId: session.sessionId, error: error instanceof Error ? error.message : String(error) },
|
||||
"failed to apply deferred session name",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const notifications = this.deferredSubsessionNotifications.get(session) ?? [];
|
||||
this.deferredSubsessionNotifications.delete(session);
|
||||
for (const notification of notifications) {
|
||||
void this.deliverSubsessionNotification(session, notification).catch((error: unknown) => {
|
||||
this.logSubsessionNotificationFailure(notification.parentId, notification.childId, error);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
applyAuthChange(change: AuthChange = {}): void {
|
||||
// ModelRuntime.login()/logout() refresh the shared runtime before AuthService
|
||||
// emits the change, so no refresh is needed here. Keeping this synchronous
|
||||
@@ -2329,6 +2577,8 @@ export class PiSessionService implements SessionRouteService {
|
||||
}
|
||||
|
||||
private activityLabelFromStatus(session: PiAgentSession): string {
|
||||
if (this.treeNavigations.has(session)) return "navigating session tree";
|
||||
if (this.isSessionEntryMutationActive(session)) return "updating session";
|
||||
if (session.isCompacting) return "compacting";
|
||||
if (session.isBashRunning) return "running bash";
|
||||
if (session.isStreaming) return "agent running";
|
||||
@@ -2337,7 +2587,85 @@ export class PiSessionService implements SessionRouteService {
|
||||
}
|
||||
|
||||
private hasActiveWork(session: PiAgentSession): boolean {
|
||||
return sessionHasActiveWork(session, this.compactionQueuedMessages(session.sessionId).length);
|
||||
return this.treeNavigations.has(session)
|
||||
|| this.isSessionEntryMutationActive(session)
|
||||
|| this.isTreeExclusiveOperationActive(session)
|
||||
|| sessionHasActiveWork(session, this.compactionQueuedMessages(session.sessionId).length);
|
||||
}
|
||||
|
||||
private async runTreeExclusiveOperation<T>(
|
||||
targets: readonly TreeExclusiveOperationTarget[],
|
||||
activeError: string,
|
||||
operation: () => Promise<T>,
|
||||
): Promise<T> {
|
||||
const sessionIds = new Set<string>();
|
||||
const runtimes = new Set<PiSessionRuntime>();
|
||||
for (const target of targets) {
|
||||
const runtime = target.runtime ?? (target.session === undefined ? undefined : this.activeRuntimeForSession(target.session));
|
||||
const session = target.session ?? runtime?.session;
|
||||
if (session !== undefined && this.hasActiveWork(session)) throw new Error(activeError);
|
||||
sessionIds.add(target.sessionId);
|
||||
if (runtime !== undefined) runtimes.add(runtime);
|
||||
}
|
||||
|
||||
for (const sessionId of sessionIds) {
|
||||
this.treeExclusiveSessionOperationCounts.set(sessionId, (this.treeExclusiveSessionOperationCounts.get(sessionId) ?? 0) + 1);
|
||||
}
|
||||
for (const runtime of runtimes) {
|
||||
this.treeExclusiveRuntimeOperationCounts.set(runtime, (this.treeExclusiveRuntimeOperationCounts.get(runtime) ?? 0) + 1);
|
||||
}
|
||||
|
||||
try {
|
||||
return await operation();
|
||||
} finally {
|
||||
for (const runtime of runtimes) decrementWeakCount(this.treeExclusiveRuntimeOperationCounts, runtime);
|
||||
for (const sessionId of sessionIds) decrementMapCount(this.treeExclusiveSessionOperationCounts, sessionId);
|
||||
}
|
||||
}
|
||||
|
||||
private isTreeExclusiveSessionIdentityActive(sessionId: string): boolean {
|
||||
return (this.treeExclusiveSessionOperationCounts.get(sessionId) ?? 0) > 0;
|
||||
}
|
||||
|
||||
private isTreeExclusiveOperationActive(session: PiAgentSession): boolean {
|
||||
if (this.isTreeExclusiveSessionIdentityActive(session.sessionId)) return true;
|
||||
const runtime = this.activeRuntimeForSession(session);
|
||||
return runtime !== undefined && (this.treeExclusiveRuntimeOperationCounts.get(runtime) ?? 0) > 0;
|
||||
}
|
||||
|
||||
private activeRuntimeForSession(session: PiAgentSession): PiSessionRuntime | undefined {
|
||||
for (const active of new Set(this.active.values())) {
|
||||
if (active.runtime.session === session) return active.runtime;
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
private assertTreeNavigationInactive(session: PiAgentSession, action: string): void {
|
||||
if (this.treeNavigations.has(session)) throw new Error(`Cannot ${action} while session tree navigation is active`);
|
||||
}
|
||||
|
||||
private async runSessionEntryMutation<T>(session: PiAgentSession, action: string, operation: () => Promise<T>): Promise<T> {
|
||||
this.beginSessionEntryMutation(session, action);
|
||||
try {
|
||||
return await operation();
|
||||
} finally {
|
||||
this.endSessionEntryMutation(session);
|
||||
}
|
||||
}
|
||||
|
||||
private beginSessionEntryMutation(session: PiAgentSession, action: string): void {
|
||||
this.assertTreeNavigationInactive(session, action);
|
||||
this.sessionEntryMutationCounts.set(session, (this.sessionEntryMutationCounts.get(session) ?? 0) + 1);
|
||||
}
|
||||
|
||||
private endSessionEntryMutation(session: PiAgentSession): void {
|
||||
const remaining = (this.sessionEntryMutationCounts.get(session) ?? 1) - 1;
|
||||
if (remaining <= 0) this.sessionEntryMutationCounts.delete(session);
|
||||
else this.sessionEntryMutationCounts.set(session, remaining);
|
||||
}
|
||||
|
||||
private isSessionEntryMutationActive(session: PiAgentSession): boolean {
|
||||
return (this.sessionEntryMutationCounts.get(session) ?? 0) > 0;
|
||||
}
|
||||
|
||||
private publishActivityForEvent(session: PiAgentSession, event: unknown): void {
|
||||
|
||||
Reference in New Issue
Block a user