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
124 changes: 124 additions & 0 deletions apps/desktop/electron/main/agent-events.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,124 @@
/**
* Process-memory fan-out hub for Agent events toward external control clients.
*
* The desktop renderer receives every `AgentEventEnvelope` through
* `IPC.event.agentMessage`; the MCP control plane, however, has been POST-only,
* so a remote client polling `pi_agent_status` / `pi_session_get` every couple
* of seconds was the only way to observe progress — one desktop RPC per session
* per tick, and every event between two polls is lost. This hub ingests the
* same envelopes the renderer sees, in Electron main, and lets the control
* server push them to a subscribed HTTP client instead of being polled.
*
* It mirrors `pending-asks.ts`: fed by the local sidecar, the native agent
* event path and the remote event bridge; process memory only, never persisted
* or logged; bounded subscriber count; deliveries are shallow clones so a
* consumer cannot mutate what the next consumer sees.
*
* Backpressure lives with the transport, not here: an SSE response that falls
* behind decides which envelopes to skip (the control server drops streaming
* deltas first) and closes the connection when it cannot keep up. The hub only
* guarantees unbounded subscriber growth cannot happen.
*/

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

/** Hard bound on concurrent subscribers; the SSE endpoint rejects beyond this. */
export const MAX_AGENT_EVENT_SUBSCRIBERS = 16;

export type AgentEventListener = (envelope: AgentEventEnvelope) => void;

export type AgentEventSubscription = {
readonly clientId: string;
/** Remove the listener. Safe to call more than once. */
unsubscribe: () => void;
};

export type AgentEventSubscribeOptions = {
/** Only deliver envelopes for these session ids; omit for every session. */
sessionIds?: readonly string[];
};

export type AgentEventHub = {
/**
* Feed one envelope to every matching subscriber. A listener that throws is
* isolated: it is removed (a broken SSE response must not break the emitter)
* and the remaining listeners still receive the event.
*/
ingest: (envelope: AgentEventEnvelope) => void;
subscribe: (
listener: AgentEventListener,
options?: AgentEventSubscribeOptions,
) => AgentEventSubscription | null;
unsubscribe: (clientId: string) => void;
subscriberCount: () => number;
};

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

type ListenerRecord = {
clientId: string;
listener: AgentEventListener;
sessionIds: ReadonlySet<string> | null;
};

let clientIdCounter = 0;

export function createAgentEventHub(): AgentEventHub {
const listeners = new Map<string, ListenerRecord>();

const ingest = (envelope: AgentEventEnvelope): void => {
if (!envelope || typeof envelope !== "object") return;
const sessionId = envelope.sessionId;
if (!isNonEmptyString(sessionId)) return;
if (!envelope.event || typeof envelope.event !== "object") return;
// Shallow clone: the emitter may mutate its envelope after we return (the
// runtime reuses message objects across stream frames), and each listener
// must observe its own copy of the top level.
const snapshot: AgentEventEnvelope = { ...envelope };
for (const record of [...listeners.values()]) {
if (record.sessionIds && !record.sessionIds.has(sessionId)) continue;
try {
record.listener(snapshot);
} catch {
// A dead subscriber must not break the emitter or the others.
listeners.delete(record.clientId);
}
}
};

const subscribe = (
listener: AgentEventListener,
options?: AgentEventSubscribeOptions,
): AgentEventSubscription | null => {
if (typeof listener !== "function") return null;
if (listeners.size >= MAX_AGENT_EVENT_SUBSCRIBERS) return null;
const sessionIds = options?.sessionIds;
const normalized =
sessionIds && sessionIds.length > 0
? new Set(sessionIds.filter(isNonEmptyString))
: null;
if (normalized !== null && normalized.size === 0) return null;
const clientId = `agent-events-${++clientIdCounter}`;
listeners.set(clientId, { clientId, listener, sessionIds: normalized });
return {
clientId,
unsubscribe: () => {
listeners.delete(clientId);
},
};
};

const unsubscribe = (clientId: string): void => {
const normalized = String(clientId ?? "").trim();
if (normalized) listeners.delete(normalized);
};


const subscriberCount = (): number => listeners.size;

return { ingest, subscribe, unsubscribe, subscriberCount };
}

/** Shared process-memory hub used by local, native and remote event paths. */
export const agentEventHub = createAgentEventHub();
4 changes: 2 additions & 2 deletions apps/desktop/electron/main/bootstrap/startup.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import {
type McpControlController,
type McpControlInvokeInput,
} from "../mcp-control";
import { agentEventHub } from "../agent-events";
import type { ModelsDevCatalog } from "../models-dev-catalog";
import type { AppUpdaterController } from "../updater";
import type { HostProcess } from "../host-process";
Expand All @@ -37,7 +38,6 @@ import {
ensureCrashDumpsDirectory,
reportPreviousCrashDumps,
} from "../crash-report";

type IpcInvoker = (
channel: string,
args?: readonly unknown[],
Expand Down Expand Up @@ -343,7 +343,7 @@ export function registerApplicationStartup(deps: StartupDependencies): void {
port: process.env.PI_DESKTOP_MCP_PORT
? Number(process.env.PI_DESKTOP_MCP_PORT)
: undefined,
controller: state.desktopControl ?? undefined,
eventHub: agentEventHub,
log: (level, message, data) => logger.app("runtime", level, message, { data }),
});
await state.mcpControl.start();
Expand Down
7 changes: 7 additions & 0 deletions apps/desktop/electron/main/ipc/agent-ipc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -741,6 +741,13 @@ export function registerAgentIpc({
return sidecar.call("agent.getStatus", { sessionId });
});

handle(IPC.invoke.agentGetStatuses, async (req: { sessionIds: string[] }) => {
if (!sidecar) throw new Error("sidecar unavailable");
return sidecar.call("agent.getStatuses", {
sessionIds: Array.isArray(req?.sessionIds) ? req.sessionIds : [],
});
});

// The Host-owned turn queue (D375 / D386). The renderer mirrors it; the
// headless module admits, orders, and drains it.
handle(IPC.invoke.agentQueuePush, async (req: AgentQueuePushRequest) => {
Expand Down
106 changes: 98 additions & 8 deletions apps/desktop/electron/main/mcp-control.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ import { isIP } from "node:net";
import { chmod, mkdir, readFile, writeFile } from "node:fs/promises";
import { join } from "node:path";
import { SESSION_COLLABORATION_OPERATIONS } from "./session-collaboration-control";

import { agentEventHub, type AgentEventHub } from "./agent-events";
/** A small JSON Schema subset used by MCP's tools/list response. */
export type McpJsonSchema = {
type?: "object" | "array" | "string" | "number" | "integer" | "boolean";
Expand Down Expand Up @@ -55,6 +55,7 @@ export type McpControlConnectionInfo = {
serverName: string;
protocol: "streamable-http";
url: string;
eventsUrl: string;
token: string;
pid: number;
startedAt?: string;
Expand Down Expand Up @@ -206,7 +207,7 @@ const CONTROL_OPERATION_SPECS: OperationSpec[] = [
spec("agentAbort", "agent/abort", "Abort an active Agent turn.", "write", ["request"]),
spec("agentStop", "agent/stop", "Request a graceful Agent stop.", "write", ["request"]),
spec("agentGetStatus", "agent/getStatus", "Read Agent runtime status.", "read", ["sessionId"]),
spec("sessionList", "session/list", "List durable sessions.", "read", []),
spec("agentGetStatuses", "agent/getStatuses", "Read Agent runtime status for multiple sessions in one desktop call.", "read", ["sessionIds"]),
spec("sessionCreate", "session/create", "Create a durable session.", "write", ["input"]),
spec("sessionFork", "session/fork", "Fork a session.", "write", ["input"]),
spec("sessionGet", "session/get", "Read a session and its transcript.", "read", ["input"]),
Expand Down Expand Up @@ -384,6 +385,13 @@ const CORE_TOOL_SPECS = [
"agent/getStatus",
(input) => [input.sessionId],
),
coreTool(
"pi_agent_status_batch",
"Read runtime status for multiple sessions in one desktop call.",
objectSchema({ sessionIds: { type: "array", items: stringSchema("Session id.") } }, ["sessionIds"]),
"agent/getStatuses",
(input) => [input],
),
coreTool(
"pi_agent_stop",
"Request a graceful stop at the next Agent turn boundary.",
Expand Down Expand Up @@ -723,6 +731,7 @@ export type McpControlServerOptions = {
invoke: IpcInvoke;
channels: Readonly<Record<string, string>>;
controller?: McpControlController;
eventHub?: AgentEventHub;
host?: string;
port?: number;
version?: string;
Expand Down Expand Up @@ -752,19 +761,21 @@ export class McpControlServer {
private readonly operations: McpControlOperation[];
private readonly operationById = new Map<string, McpControlOperation>();
private readonly toolsList: McpTool[];
private readonly sessions = new Set<string>();
private readonly serverName = MCP_SERVER_NAME;
private readonly sessions = new Set<string>();
private readonly eventHub: AgentEventHub;
private readonly eventConnections = new Map<ServerResponse, () => void>();
private server: ReturnType<typeof createServer> | null = null;
private token = "";
private port: number | null = null;
private startedAt: string | undefined;

constructor(options: McpControlServerOptions) {
this.dataDir = options.dataDir;
this.host = options.host ?? DEFAULT_HOST;
this.requestedPort = options.port ?? DEFAULT_PORT;
this.version = options.version ?? "1";
this.log = options.log ?? (() => undefined);
this.eventHub = options.eventHub ?? agentEventHub;
this.controller = options.controller ?? createMcpControlController({
channels: options.channels,
invoke: options.invoke,
Expand All @@ -786,6 +797,7 @@ export class McpControlServer {
active: this.isRunning,
serverName: this.serverName,
protocol: "streamable-http",
eventsUrl: `http://${hostname}:${this.port}/events`,
url: `http://${hostname}:${this.port}/mcp`,
token: this.token,
pid: process.pid,
Expand Down Expand Up @@ -860,6 +872,10 @@ export class McpControlServer {
const server = this.server;
this.server = null;
this.sessions.clear();
for (const [response, cleanup] of [...this.eventConnections.entries()]) {
cleanup();
if (!response.writableEnded) response.end();
}
if (server) {
if (typeof server.closeAllConnections === "function") {
server.closeAllConnections();
Expand Down Expand Up @@ -965,22 +981,26 @@ export class McpControlServer {
}
if (request.method === "OPTIONS") {
result.writeHead(204, {
Allow: "POST, DELETE, OPTIONS",
Allow: "GET, POST, DELETE, OPTIONS",
"Access-Control-Allow-Headers": "Authorization, Content-Type, Mcp-Session-Id, MCP-Protocol-Version, X-Pi-Desktop-Token",
"Access-Control-Allow-Methods": "POST, DELETE, OPTIONS",
"Access-Control-Allow-Methods": "GET, POST, DELETE, OPTIONS",
});
result.end();
return;
}
const pathname = new URL(request.url ?? "/", "http://127.0.0.1").pathname;
if (pathname !== "/mcp" && pathname !== "/mcp/") {
if (pathname !== "/mcp" && pathname !== "/mcp/" && pathname !== "/events" && pathname !== "/events/") {
this.sendHttp(result, 404, { error: "not found" });
return;
}
if (!this.isAuthorized(request)) {
this.sendHttp(result, 401, { error: "unauthorized" }, { "WWW-Authenticate": "Bearer" });
return;
}
if (pathname === "/events" || pathname === "/events/") {
await this.handleEvents(request, result);
return;
}
if (request.method === "DELETE") {
const sessionId = this.sessionId(request);
if (!sessionId || !this.sessions.has(sessionId)) {
Expand All @@ -993,7 +1013,7 @@ export class McpControlServer {
return;
}
if (request.method === "GET") {
this.sendHttp(result, 405, { error: "SSE stream not supported" }, { Allow: "POST, DELETE, OPTIONS" });
this.sendHttp(result, 405, { error: "method not allowed" }, { Allow: "POST, DELETE, OPTIONS" });
return;
}
if (request.method !== "POST") {
Expand Down Expand Up @@ -1123,6 +1143,76 @@ export class McpControlServer {
return { response: rpcError(id, -32601, `method not found: ${method}`) };
}

private async handleEvents(request: IncomingMessage, result: ServerResponse): Promise<void> {
if (request.method !== "GET") {
this.sendHttp(result, 405, { error: "method not allowed" }, { Allow: "GET, OPTIONS" });
return;
}
const url = new URL(request.url ?? "/events", "http://127.0.0.1");
const requestedIds = [
...url.searchParams.getAll("sessionId"),
...url.searchParams.getAll("sessionIds").flatMap((value) => value.split(",")),
]
.map((value) => value.trim())
.filter(Boolean);
const sessionIds = requestedIds.length > 0 ? [...new Set(requestedIds)].slice(0, 256) : undefined;
let subscription: ReturnType<AgentEventHub["subscribe"]>;
let heartbeat: ReturnType<typeof setInterval> | undefined;
let closed = false;
const cleanup = () => {
if (closed) return;
closed = true;
if (heartbeat) clearInterval(heartbeat);
subscription?.unsubscribe();
this.eventConnections.delete(result);
};
result.writeHead(200, {
"Content-Type": "text/event-stream; charset=utf-8",
"Cache-Control": "no-cache, no-store, must-revalidate",
Connection: "keep-alive",
"X-Accel-Buffering": "no",
});
result.write(`event: ready\ndata: ${JSON.stringify({ protocol: "pi-agent-events", version: 1 })}\n\n`);
subscription = this.eventHub.subscribe(
(envelope) => {
if (closed || result.destroyed || result.writableEnded) return;
try {
const ok = result.write(`event: agent\ndata: ${JSON.stringify(envelope)}\n\n`);
if (!ok) {
cleanup();
result.destroy();
}
} catch {
cleanup();
result.destroy();
}
},
sessionIds ? { sessionIds } : undefined,
);
if (!subscription) {
cleanup();
if (!result.writableEnded) result.end();
return;
}
heartbeat = setInterval(() => {
if (closed || result.destroyed || result.writableEnded) {
cleanup();
return;
}
try {
result.write(": heartbeat\n\n");
} catch {
cleanup();
result.destroy();
}
}, 15_000);
this.eventConnections.set(result, cleanup);
request.on("close", cleanup);
request.on("error", cleanup);
result.on("close", cleanup);
result.on("error", cleanup);
}

private buildTools(): McpTool[] {
const common = CORE_TOOL_SPECS.flatMap((entry) => {
const operation = this.operationById.get(entry.operationId);
Expand Down
2 changes: 2 additions & 0 deletions apps/desktop/electron/main/remote/remote-event-bridge.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
* connection layer (Stage 3) drives subscription and unsubscription.
*/
import { IPC } from "@pi-desktop/shared";
import { agentEventHub } from "../agent-events";
import type {
AgentEvent,
AgentEventEnvelope,
Expand Down Expand Up @@ -156,6 +157,7 @@ export function createRemoteEventBridge(options: RemoteEventBridgeOptions): Remo
...(envelope.parentToolCallId ? { parentToolCallId: envelope.parentToolCallId } : {}),
...(envelope.agentName ? { agentName: envelope.agentName } : {}),
};
agentEventHub.ingest(local);
emit(IPC.event.agentMessage, local);
};

Expand Down
Loading
Loading