Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 14 additions & 1 deletion apps/desktop/electron/main/ipc/agent-ipc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import type { PersistenceOutbox } from "../persistence-outbox";
import type { ComposerCommandService } from "./composer-ipc";
import type { IpcRegistrar } from "./types";
import { withPromptEnhancementTimeout } from "../prompt-enhancement-timeout";
import { pendingAsksRegistry } from "../pending-asks";

export type AgentIpcDependencies = {
registrar: IpcRegistrar;
Expand Down Expand Up @@ -703,6 +704,7 @@ export function registerAgentIpc({
agentExtensions.cancelPrompts(req.sessionId);
cancelSessionTools(req.sessionId, "Session turn was aborted");
result = await sidecar.call("agent.abort", req);
pendingAsksRegistry.clearSession(req.sessionId);
} finally {
// A turn that already stopped owning the session is refused inside the
// finalizer, so the identity captured above is the only one used here.
Expand Down Expand Up @@ -794,16 +796,27 @@ export function registerAgentIpc({
return resolved;
});

handle(
IPC.invoke.askToolPending,
async (input: { sessionId?: unknown } = {}) => {
const sessionId =
typeof input.sessionId === "string" ? input.sessionId.trim() : "";
return pendingAsksRegistry.list(sessionId || undefined);
},
);

handle(IPC.invoke.askToolResolve, async (resolution: AskToolResolution) => {
if (!sidecar) throw new Error("sidecar unavailable");
const sessionId = String(resolution?.sessionId ?? "").trim();
const requestId = String(resolution?.requestId ?? "").trim();
if (!sessionId || !requestId) throw new Error("asktool resolution identity required");
return sidecar.call("asktool.resolve", {
const result = await sidecar.call("asktool.resolve", {
...resolution,
sessionId,
requestId,
});
pendingAsksRegistry.settle(sessionId, requestId);
return result;
});

handle(IPC.invoke.plansPending, async (input: { sessionId?: string } = {}) => {
Expand Down
2 changes: 2 additions & 0 deletions apps/desktop/electron/main/ipc/session-ipc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import type { PluginRuntime } from "../plugin-runtime";
import { readSessionCollaboration } from "../services/session-collaboration";
import { searchSessionsAcrossSources } from "../services/session-search";
import type { IpcRegistrar } from "./types";
import { pendingAsksRegistry } from "../pending-asks";

type RuntimeSession = {
id?: string;
Expand Down Expand Up @@ -324,6 +325,7 @@ export function registerSessionIpc({
}
if (!host) throw new Error("host unavailable");
const res = await host.call("session.delete", { id });
pendingAsksRegistry.clearSession(id);
await persistenceOutbox.dropSession(id);
// Drop the session's pi-agent so a later session with the same id (or a
// stale runtime) can't answer with this session's context.
Expand Down
8 changes: 8 additions & 0 deletions apps/desktop/electron/main/mcp-control.ts
Original file line number Diff line number Diff line change
Expand Up @@ -244,6 +244,7 @@ const CONTROL_OPERATION_SPECS: OperationSpec[] = [
spec("pullsList", "pulls/list", "List pull requests for the active workspace.", "read", []),
spec("scheduledList", "scheduled/list", "List scheduled tasks.", "read", []),
spec("toolResolvePermission", "tool/resolvePermission", "Resolve a pending tool permission request.", "dangerous", ["resolution"]),
spec("askToolPending", "agent/askTool/pending", "List pending Agent questions.", "read", ["input"]),
spec("askToolResolve", "agent/askTool/resolve", "Answer an Agent question.", "dangerous", ["resolution"]),
spec("plansPending", "plans/pending", "List pending Plan or Goal approvals.", "read", ["input"]),
spec("plansResolve", "plans/resolve", "Approve or reject a Plan or Goal checkpoint.", "dangerous", ["resolution"]),
Expand Down Expand Up @@ -405,6 +406,13 @@ const CORE_TOOL_SPECS = [
"agent/compact",
(input) => [input],
),
coreTool(
"pi_asktool_pending",
"List pending Agent questions.",
objectSchema({ sessionId: stringSchema("Optional session id filter.") }),
"agent/askTool/pending",
(input) => [input],
),
coreTool(
"pi_plans_pending",
"List pending Plan or Goal approvals.",
Expand Down
227 changes: 227 additions & 0 deletions apps/desktop/electron/main/pending-asks.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,227 @@
/**
* Process-memory registry of pending Agent asktool questions.
*
* The desktop renderer receives every `asktool_request` as an in-memory
* `AgentEventEnvelope`; the local MCP control plane can answer such a
* question (`agent/askTool/resolve`) but could not read it back: the request
* id lives only in the renderer's zustand store. This registry ingests the
* same envelopes in Electron main so a phone client can discover the
* requestId and the questions it needs for the resolve operation.
*
* Security bounds:
* - Process memory only. Never persisted to disk, log, notification, or
* transcript. State dies with the process, matching the runtime's own
* pending-request lifetime.
* - Bounded per session: when a session accumulates more than
* {@link MAX_PENDING_ASKS_PER_SESSION} unresolved requests, the oldest
* entry is pruned first. No unbounded growth.
* - `list()` returns structural clones so callers cannot mutate the
* registry's references.
*
* Lifecycle:
* - `asktool_request` adds (or dedupes by requestId).
* - `tool_end` for the same session removes the entry whose `toolCallId`
* matches (the asktool call has completed).
* - `agent_end` clears the session: the turn is over and every unresolved
* request is dead for that session.
*/

import type {
AgentEventEnvelope,
AskToolQuestion,
AskToolRequest,
} from "@pi-desktop/shared";

/** Hard per-session bound; pruned oldest-first. */
export const MAX_PENDING_ASKS_PER_SESSION = 20;

/** Hard global bound on session buckets; prevents stale remote hosts from growing memory forever. */
export const MAX_PENDING_ASK_SESSIONS = 1024;

/** One pending ask as exposed over the read operation. */
export type PendingAsk = {
requestId: string;
sessionId: string;
toolCallId: string;
questions: AskToolQuestion[];
receivedAt: number;
};

export type PendingAsksResult = {
kind: "pending" | "none";
asks: PendingAsk[];
};

export type PendingAsksRegistry = {
ingest: (envelope: AgentEventEnvelope) => void;
settle: (sessionId: string, requestId: string) => void;
clearSession: (sessionId: string) => void;
/**
* Drop every bucket whose session id starts with `prefix`. Used when a remote
* host goes away: its asks can never be answered afterwards, and the ids are
* namespaced (`remote:<hostKey>:<hostSessionId>`), so the prefix is exact.
* Returns the number of buckets removed.
*/
clearSessionsWithPrefix: (prefix: string) => number;
list: (sessionId?: string) => PendingAsksResult;
};

const isNonEmptyString = (value: unknown): value is string =>
typeof value === "string" && value.trim().length > 0;

/**
* Create the main-process pending-asks registry. Pure in-memory state with
* no timers, no persistence, and no external dependencies beyond the
* shared `AgentEventEnvelope` contract.
*/
export function createPendingAsksRegistry(): PendingAsksRegistry {
/** sessionId -> Array<PendingAsk>, ordered by receipt. */
const bySession = new Map<string, PendingAsk[]>();

/**
* Keep the bucket count bounded. A remote host that disconnects without
* resolving its asks (or a long-lived desktop accumulating dead sessions)
* must not grow process memory without limit; the per-session cap alone does
* not bound the number of sessions. The oldest bucket (Map insertion order)
* goes first, matching the per-session oldest-first policy.
*/
const evictOldestSessions = (): void => {
while (bySession.size > MAX_PENDING_ASK_SESSIONS) {
const oldest = bySession.keys().next();
if (oldest.done) return;
bySession.delete(oldest.value);
}
};

const ingest = (envelope: AgentEventEnvelope): void => {
if (!envelope || typeof envelope !== "object") return;
const sessionId = String(envelope.sessionId ?? "").trim();
const event = envelope.event;
if (!sessionId || !event) return;
if (event.type === "asktool_request") {
const request: AskToolRequest | undefined = event.request;
if (
!request ||
!isNonEmptyString(request.requestId) ||
!isNonEmptyString(request.toolCallId) ||
request.sessionId !== sessionId ||
!Array.isArray(request.questions)
) {
return;
}
const questions: AskToolQuestion[] = [];
for (const raw of request.questions as unknown[]) {
if (!raw || typeof raw !== "object") return;
const question = raw as Partial<AskToolQuestion>;
if (
!isNonEmptyString(question.question) ||
!Array.isArray(question.options) ||
!question.options.every((option) => typeof option === "string")
) {
return;
}
questions.push({
question: question.question,
options: [...question.options],
...(question.multiSelect === true ? { multiSelect: true } : {}),
});
}
if (questions.length === 0) return;
const existing = bySession.get(sessionId) ?? [];
if (existing.some((ask) => ask.requestId === request.requestId)) return;
const entry: PendingAsk = {
requestId: request.requestId,
sessionId,
toolCallId: request.toolCallId,
questions,
receivedAt: Number.isFinite(envelope.ts) ? envelope.ts : Date.now(),
};
const next = [...existing, entry];
while (next.length > MAX_PENDING_ASKS_PER_SESSION) {
next.shift();
}
// Re-insert so this session becomes the newest bucket: Map iteration order
// is insertion order, which is what `evictOldestSessions` relies on.
bySession.delete(sessionId);
bySession.set(sessionId, next);
evictOldestSessions();
return;
}
if (event.type === "tool_end") {
const toolCallId = event.toolCallId;
if (!isNonEmptyString(toolCallId)) return;
const existing = bySession.get(sessionId);
if (!existing?.length) return;
const remaining = existing.filter((ask) => ask.toolCallId !== toolCallId);
if (remaining.length === existing.length) return;
if (remaining.length === 0) bySession.delete(sessionId);
else bySession.set(sessionId, remaining);
return;
}
if (event.type === "agent_end") {
bySession.delete(sessionId);
}
};

const settle = (sessionId: string, requestId: string): void => {
const normalizedSessionId = String(sessionId ?? "").trim();
const normalizedRequestId = String(requestId ?? "").trim();
if (!normalizedSessionId || !normalizedRequestId) return;
const existing = bySession.get(normalizedSessionId);
if (!existing?.length) return;
const remaining = existing.filter(
(ask) => ask.requestId !== normalizedRequestId,
);
if (remaining.length === existing.length) return;
if (remaining.length === 0) bySession.delete(normalizedSessionId);
else bySession.set(normalizedSessionId, remaining);
};

const clearSession = (sessionId: string): void => {
const normalized = String(sessionId ?? "").trim();
if (normalized) bySession.delete(normalized);
};

const clearSessionsWithPrefix = (prefix: string): number => {
const normalized = String(prefix ?? "");
if (!normalized) return 0;
let removed = 0;
for (const sessionId of [...bySession.keys()]) {
if (!sessionId.startsWith(normalized)) continue;
bySession.delete(sessionId);
removed += 1;
}
return removed;
};

const list = (sessionId?: string): PendingAsksResult => {
const sources: PendingAsk[][] = [];
if (isNonEmptyString(sessionId)) {
const existing = bySession.get(sessionId.trim());
if (existing?.length) sources.push(existing);
} else {
for (const bucket of bySession.values()) {
if (bucket.length) sources.push(bucket);
}
}
const asks = sources.flat().sort((left, right) => {
if (left.receivedAt !== right.receivedAt) return left.receivedAt - right.receivedAt;
return left.requestId < right.requestId ? -1 : left.requestId > right.requestId ? 1 : 0;
});
return {
kind: asks.length > 0 ? "pending" : "none",
asks: asks.map((ask) => ({
...ask,
questions: ask.questions.map((question) => ({
...question,
options: [...question.options],
})),
})),
};
};

return { ingest, settle, clearSession, clearSessionsWithPrefix, list };
}

/** Shared process-memory registry used by local, native and remote event paths. */
export const pendingAsksRegistry = createPendingAsksRegistry();
7 changes: 7 additions & 0 deletions apps/desktop/electron/main/remote/backend-router.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,13 @@ export const ROUTE_LOCAL = Symbol("pi-desktop.route-local");
/** Namespaced-id prefix for sessions owned by a remote host. */
const REMOTE_PREFIX = "remote:";

/**
* The same prefix, exported for callers that must reason about a whole host's
* sessions at once (e.g. clearing that host's pending-ask buckets when it goes
* away) instead of one namespaced id at a time.
*/
export const REMOTE_SESSION_PREFIX = REMOTE_PREFIX;

/**
* Delimiter that embeds the renderer-visible remote session id inside a tool
* permission `requestId`. `toolResolvePermission` carries only `{requestId,
Expand Down
2 changes: 2 additions & 0 deletions apps/desktop/electron/main/remote/remote-backend.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ import {
sessionIdForCall,
} from "./backend-router.js";
import { racpSessionToSummary, snapshotToSessionDetail } from "./remote-transcript.js";
import { pendingAsksRegistry } from "../pending-asks";

/** The subset of `RacpClient` this backend needs; kept minimal for testing. */
export type RemoteRacpClient = {
Expand Down Expand Up @@ -306,6 +307,7 @@ export function createRemoteBackend(options: RemoteBackendOptions): RemoteBacken
answers: resolution.answers,
context: context(),
});
pendingAsksRegistry.settle(resolution.sessionId, resolution.requestId);
return { ok: true };
}
case IPC.invoke.plansResolve: {
Expand Down
19 changes: 14 additions & 5 deletions apps/desktop/electron/main/remote/remote-event-bridge.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ import type {
ToolPermissionRequest,
} from "@pi-desktop/shared";
import { makeRemoteApprovalRequestId, makeRemoteSessionId } from "./backend-router.js";
import { pendingAsksRegistry } from "../pending-asks";

/** A minimal shape of the session field carried by host-scope session events.
* Both the RACP `RacpSession` and the host's smaller `SessionSummary` extend
Expand Down Expand Up @@ -120,10 +121,10 @@ function toAskToolRequest(
return {
requestId: input.id,
sessionId: remoteSessionId,
// Ask-tool needs a toolCallId to attach the answer to; the RACP schema
// supplies it as `parentToolCallId` when the input came from a subagent,
// and leaves it undefined for the top-level agent.
toolCallId: input.parentToolCallId ?? "",
// The RACP response is keyed by `input.id`, not this field. Preserve the
// parent tool call when present; for a top-level ask use the input id as a
// stable non-empty lifecycle key for the main-process registry.
toolCallId: input.parentToolCallId ?? input.id,
questions: input.questions.map((question) => ({
question: question.question,
options: question.options,
Expand Down Expand Up @@ -156,6 +157,7 @@ export function createRemoteEventBridge(options: RemoteEventBridgeOptions): Remo
...(envelope.parentToolCallId ? { parentToolCallId: envelope.parentToolCallId } : {}),
...(envelope.agentName ? { agentName: envelope.agentName } : {}),
};
pendingAsksRegistry.ingest(local);
emit(IPC.event.agentMessage, local);
};

Expand All @@ -181,6 +183,7 @@ export function createRemoteEventBridge(options: RemoteEventBridgeOptions): Remo
}
// "session.archived": pass through to the lifecycle handler for router
// cleanup, then refresh the renderer's session list.
pendingAsksRegistry.clearSession(remoteSessionId);
onLifecycle?.({ kind: "session.archived", hostSessionId: session.id, remoteSessionId, session });
emit(IPC.event.sessionsChanged, { reason: "remote.session.archived" });
};
Expand Down Expand Up @@ -244,8 +247,14 @@ export function createRemoteEventBridge(options: RemoteEventBridgeOptions): Remo
});
return;
}
case "input.resolved": {
const payload = envelope.payload;
if (isRecord(payload) && typeof payload.inputId === "string") {
pendingAsksRegistry.settle(remoteSessionId, payload.inputId);
}
return;
}
case "approval.resolved":
case "input.resolved":
case "terminal.changed":
case "terminal.output":
case "resync.required":
Expand Down
Loading
Loading