Archived
refactor: add session route service seam
This commit is contained in:
@@ -13,7 +13,7 @@ import {
|
||||
type CreateAgentSessionRuntimeFactory,
|
||||
type EditToolDetails,
|
||||
} from "@earendil-works/pi-coding-agent";
|
||||
import type { ClientArchiveSessionsResponse, ClientCommand, ClientCommandResult, ClientMessagePage, ClientSession, ClientSessionCleanupExecuteResponse, ClientSessionCleanupPreviewResponse, ClientSessionModel, ClientSessionRef, ClientSessionStatus, ClientThinkingLevel, SessionUiEvent } from "../types.js";
|
||||
import type { ClientArchiveSessionsResponse, ClientCommand, ClientCommandResult, ClientMessagePage, ClientSession, ClientSessionCleanupExecuteResponse, ClientSessionCleanupPreviewResponse, ClientSessionModel, ClientSessionStatus, ClientThinkingLevel, SessionUiEvent } from "../types.js";
|
||||
import { pageMessagesAtSafeBoundary } from "./messagePaging.js";
|
||||
import type { SessionEventHub } from "../realtime/sessionEventHub.js";
|
||||
import { BUILTIN_COMMANDS } from "./builtinCommands.js";
|
||||
@@ -28,6 +28,7 @@ import { createPiSessionManagerGateway } from "./piSessionManagerGateway.js";
|
||||
import { attachmentsToInlineImages, saveAttachmentsToWorkspace } from "./attachmentService.js";
|
||||
import { parsePromptAttachments } from "../../shared/promptAttachments.js";
|
||||
import type { SavedPromptAttachment } from "../../shared/apiTypes.js";
|
||||
import type { SessionRouteLookup, SessionRouteRef, SessionRouteService } from "./sessionService.js";
|
||||
|
||||
import { cwdPathsEqual } from "../workingDirectory.js";
|
||||
import type { WorkspaceActivityService } from "../activity/workspaceActivityService.js";
|
||||
@@ -115,9 +116,8 @@ function parsePromptStreamingBehavior(value: unknown): QueuedPromptKind | undefi
|
||||
|
||||
type SessionArchiveRepository = Pick<SessionArchiveStore, "list" | "get" | "archive" | "restore" | "isArchived"> & { deleteArchived?: (sessionId: string) => Promise<void> };
|
||||
|
||||
export type PiSessionRef = ClientSessionRef;
|
||||
|
||||
type PiSessionLookup = string | PiSessionRef;
|
||||
export type PiSessionRef = SessionRouteRef;
|
||||
type PiSessionLookup = SessionRouteLookup;
|
||||
|
||||
export interface PiSessionListEntry {
|
||||
id: string;
|
||||
@@ -303,7 +303,7 @@ export interface PiSessionServiceDependencies {
|
||||
now?: () => Date;
|
||||
}
|
||||
|
||||
export class PiSessionService {
|
||||
export class PiSessionService implements SessionRouteService {
|
||||
private readonly active = new Map<string, ActiveSession<PiSessionRuntime>>();
|
||||
private readonly activities = new Map<string, { phase: "active" | "idle" | "error"; label: string; detail?: string; at: string }>();
|
||||
private readonly heartbeat: NodeJS.Timeout;
|
||||
|
||||
@@ -4,7 +4,8 @@ import fastifyWebsocket from "@fastify/websocket";
|
||||
import { afterEach, beforeEach, describe, expect, it } from "vitest";
|
||||
import type { SessionCleanupExecuteResponse, SessionCleanupPreviewResponse } from "../../shared/apiTypes.js";
|
||||
import { SessionEventHub } from "../realtime/sessionEventHub.js";
|
||||
import { PiSessionService, type PiSessionManagerGateway, type PiSessionRef } from "./piSessionService.js";
|
||||
import { PiSessionService, type PiSessionManagerGateway } from "./piSessionService.js";
|
||||
import type { SessionRouteLookup, SessionRouteService } from "./sessionService.js";
|
||||
import { registerSessionRoutes } from "./sessionRoutes.js";
|
||||
import type { NormalizedSessionCleanupRequest } from "./sessionCleanup.js";
|
||||
|
||||
@@ -39,7 +40,7 @@ describe("session routes", () => {
|
||||
const routeApp = Fastify({ logger: false });
|
||||
await routeApp.register(fastifyWebsocket);
|
||||
const eventHub = new SessionEventHub();
|
||||
const routeService = new CapturingRouteSessionService(eventHub);
|
||||
const routeService = new CapturingRouteSessionService();
|
||||
registerSessionRoutes(routeApp, routeService, eventHub);
|
||||
|
||||
try {
|
||||
@@ -59,7 +60,7 @@ describe("session routes", () => {
|
||||
const routeApp = Fastify({ logger: false });
|
||||
await routeApp.register(fastifyWebsocket);
|
||||
const eventHub = new SessionEventHub();
|
||||
const routeService = new CapturingRouteSessionService(eventHub);
|
||||
const routeService = new CapturingRouteSessionService();
|
||||
registerSessionRoutes(routeApp, routeService, eventHub);
|
||||
|
||||
const attachments = [{ kind: "image", mimeType: "image/png", data: "QUJD", name: "shot.png" }];
|
||||
@@ -81,7 +82,7 @@ describe("session routes", () => {
|
||||
const routeApp = Fastify({ logger: false });
|
||||
await routeApp.register(fastifyWebsocket);
|
||||
const eventHub = new SessionEventHub();
|
||||
const routeService = new CapturingRouteSessionService(eventHub);
|
||||
const routeService = new CapturingRouteSessionService();
|
||||
registerSessionRoutes(routeApp, routeService, eventHub);
|
||||
|
||||
try {
|
||||
@@ -104,7 +105,7 @@ describe("session routes", () => {
|
||||
const routeApp = Fastify({ logger: false });
|
||||
await routeApp.register(fastifyWebsocket);
|
||||
const eventHub = new SessionEventHub();
|
||||
const routeService = new CapturingRouteSessionService(eventHub);
|
||||
const routeService = new CapturingRouteSessionService();
|
||||
registerSessionRoutes(routeApp, routeService, eventHub);
|
||||
|
||||
try {
|
||||
@@ -124,7 +125,7 @@ describe("session routes", () => {
|
||||
const routeApp = Fastify({ logger: false });
|
||||
await routeApp.register(fastifyWebsocket);
|
||||
const eventHub = new SessionEventHub();
|
||||
const routeService = new CapturingRouteSessionService(eventHub);
|
||||
const routeService = new CapturingRouteSessionService();
|
||||
routeService.reloadError = new Error("Stop current session activity before reloading");
|
||||
registerSessionRoutes(routeApp, routeService, eventHub);
|
||||
|
||||
@@ -143,7 +144,7 @@ describe("session routes", () => {
|
||||
const routeApp = Fastify({ logger: false });
|
||||
await routeApp.register(fastifyWebsocket);
|
||||
const eventHub = new SessionEventHub();
|
||||
const routeService = new CapturingRouteSessionService(eventHub);
|
||||
const routeService = new CapturingRouteSessionService();
|
||||
registerSessionRoutes(routeApp, routeService, eventHub);
|
||||
|
||||
try {
|
||||
@@ -164,7 +165,7 @@ describe("session routes", () => {
|
||||
const routeApp = Fastify({ logger: false });
|
||||
await routeApp.register(fastifyWebsocket);
|
||||
const eventHub = new SessionEventHub();
|
||||
const routeService = new CapturingRouteSessionService(eventHub);
|
||||
const routeService = new CapturingRouteSessionService();
|
||||
registerSessionRoutes(routeApp, routeService, eventHub);
|
||||
|
||||
try {
|
||||
@@ -180,34 +181,38 @@ describe("session routes", () => {
|
||||
});
|
||||
});
|
||||
|
||||
class CapturingRouteSessionService extends PiSessionService {
|
||||
class CapturingRouteSessionService implements SessionRouteService {
|
||||
readonly calls: unknown[] = [];
|
||||
readonly reloadCalls: (string | PiSessionRef)[] = [];
|
||||
readonly reloadCalls: SessionRouteLookup[] = [];
|
||||
readonly cleanupPreviewCalls: NormalizedSessionCleanupRequest[] = [];
|
||||
readonly cleanupCalls: NormalizedSessionCleanupRequest[] = [];
|
||||
reloadError: Error | undefined;
|
||||
|
||||
constructor(eventHub: SessionEventHub) {
|
||||
super(eventHub, { sessionManager: new RejectingSessionManager(), heartbeatIntervalMs: 60_000 });
|
||||
}
|
||||
|
||||
override cleanupPreview(request: NormalizedSessionCleanupRequest): Promise<SessionCleanupPreviewResponse> {
|
||||
cleanupPreview(request: NormalizedSessionCleanupRequest): Promise<SessionCleanupPreviewResponse> {
|
||||
this.cleanupPreviewCalls.push(request);
|
||||
return Promise.resolve({ generatedAt: "2026-06-25T00:00:00.000Z", thresholds: request.thresholds, projects: [], totals: { archiveCount: 0, deleteCount: 0 } });
|
||||
}
|
||||
|
||||
override cleanup(request: NormalizedSessionCleanupRequest): Promise<SessionCleanupExecuteResponse> {
|
||||
cleanup(request: NormalizedSessionCleanupRequest): Promise<SessionCleanupExecuteResponse> {
|
||||
this.cleanupCalls.push(request);
|
||||
return Promise.resolve({ generatedAt: "2026-06-25T00:00:00.000Z", thresholds: request.thresholds, projects: [], totals: { archiveCount: 0, deleteCount: 0 }, archivedSessionIds: [], deletedSessionIds: [] });
|
||||
}
|
||||
|
||||
override reload(lookup: string | PiSessionRef): Promise<void> {
|
||||
reload(lookup: SessionRouteLookup): Promise<void> {
|
||||
this.reloadCalls.push(lookup);
|
||||
if (this.reloadError !== undefined) return Promise.reject(this.reloadError);
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
override status(lookup: string | PiSessionRef) {
|
||||
dispose(): Promise<void> {
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
list(): never { throw unusedRouteMethod("list"); }
|
||||
start(): never { throw unusedRouteMethod("start"); }
|
||||
messages(): Promise<unknown[]> { return Promise.resolve([]); }
|
||||
|
||||
status(lookup: SessionRouteLookup) {
|
||||
this.calls.push(lookup);
|
||||
return Promise.resolve({
|
||||
sessionId: sessionIdFromLookup(lookup),
|
||||
@@ -221,12 +226,20 @@ class CapturingRouteSessionService extends PiSessionService {
|
||||
});
|
||||
}
|
||||
|
||||
override prompt(lookup: string | PiSessionRef, text: unknown, _streamingBehavior?: unknown, attachments?: unknown): Promise<void> {
|
||||
availableModels(): Promise<[]> { return Promise.resolve([]); }
|
||||
setModel(): never { throw unusedRouteMethod("setModel"); }
|
||||
cycleModel(): never { throw unusedRouteMethod("cycleModel"); }
|
||||
availableThinkingLevels(): Promise<[]> { return Promise.resolve([]); }
|
||||
setThinkingLevel(): never { throw unusedRouteMethod("setThinkingLevel"); }
|
||||
cycleThinkingLevel(): never { throw unusedRouteMethod("cycleThinkingLevel"); }
|
||||
commands(): Promise<[]> { return Promise.resolve([]); }
|
||||
|
||||
prompt(lookup: SessionRouteLookup, text: unknown, _streamingBehavior?: unknown, attachments?: unknown): Promise<void> {
|
||||
this.calls.push(attachments === undefined ? { lookup, text } : { lookup, text, attachments });
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
override saveAttachments(_lookup: string | PiSessionRef, attachments: unknown, folder?: string) {
|
||||
saveAttachments(_lookup: SessionRouteLookup, attachments: unknown, folder?: string) {
|
||||
const list = Array.isArray(attachments) ? attachments : [];
|
||||
return Promise.resolve(list.map((attachment: { mimeType: string; data: string; name?: string }) => ({
|
||||
path: `${folder ?? ".pi-web/attachments"}/${attachment.name ?? "file.png"}`,
|
||||
@@ -234,6 +247,19 @@ class CapturingRouteSessionService extends PiSessionService {
|
||||
size: Buffer.from(attachment.data, "base64").byteLength,
|
||||
})));
|
||||
}
|
||||
|
||||
shell(): never { throw unusedRouteMethod("shell"); }
|
||||
runCommand(): never { throw unusedRouteMethod("runCommand"); }
|
||||
respondToCommand(): never { throw unusedRouteMethod("respondToCommand"); }
|
||||
abort(): never { throw unusedRouteMethod("abort"); }
|
||||
stop(): never { throw unusedRouteMethod("stop"); }
|
||||
archive(): never { throw unusedRouteMethod("archive"); }
|
||||
archiveTree(): never { throw unusedRouteMethod("archiveTree"); }
|
||||
restore(): never { throw unusedRouteMethod("restore"); }
|
||||
deleteArchived(): never { throw unusedRouteMethod("deleteArchived"); }
|
||||
|
||||
|
||||
detachParent(): never { throw unusedRouteMethod("detachParent"); }
|
||||
}
|
||||
|
||||
class RejectingSessionManager implements PiSessionManagerGateway {
|
||||
@@ -260,6 +286,10 @@ class RejectingSessionManager implements PiSessionManagerGateway {
|
||||
}
|
||||
}
|
||||
|
||||
function sessionIdFromLookup(lookup: string | PiSessionRef): string {
|
||||
function sessionIdFromLookup(lookup: SessionRouteLookup): string {
|
||||
return typeof lookup === "string" ? lookup : lookup.id;
|
||||
}
|
||||
|
||||
function unusedRouteMethod(name: string): Error {
|
||||
return new Error(`Route test did not expect ${name} to be called`);
|
||||
}
|
||||
|
||||
@@ -2,10 +2,10 @@ import type { FastifyInstance } from "fastify";
|
||||
import type { SessionCleanupRequest } from "../../shared/apiTypes.js";
|
||||
import { normalizeRequestCwd } from "../workingDirectory.js";
|
||||
import type { SessionEventHub } from "../realtime/sessionEventHub.js";
|
||||
import type { PiSessionRef, PiSessionService } from "./piSessionService.js";
|
||||
import type { SessionRouteLookup, SessionRouteService } from "./sessionService.js";
|
||||
import { normalizeSessionCleanupRequest } from "./sessionCleanup.js";
|
||||
|
||||
type SessionLookup = string | PiSessionRef;
|
||||
type SessionLookup = SessionRouteLookup;
|
||||
|
||||
interface SessionQuery {
|
||||
cwd?: string;
|
||||
@@ -29,7 +29,7 @@ interface AttachmentsRequestBody {
|
||||
folder?: unknown;
|
||||
}
|
||||
|
||||
export function registerSessionRoutes(app: FastifyInstance, sessions: PiSessionService, eventHub: SessionEventHub, prefix = ""): void {
|
||||
export function registerSessionRoutes(app: FastifyInstance, sessions: SessionRouteService, eventHub: SessionEventHub, prefix = ""): void {
|
||||
app.get<{ Querystring: SessionQuery }>(`${prefix}/sessions`, async (request, reply) => {
|
||||
if (request.query.cwd === undefined || request.query.cwd === "") return reply.code(400).send({ error: "cwd query parameter is required" });
|
||||
try {
|
||||
@@ -204,9 +204,9 @@ export function registerSessionRoutes(app: FastifyInstance, sessions: PiSessionS
|
||||
}
|
||||
});
|
||||
|
||||
app.post<{ Params: { sessionId: string }; Body: { cwd?: unknown } | undefined }>(`${prefix}/sessions/:sessionId/stop`, (request, reply) => {
|
||||
app.post<{ Params: { sessionId: string }; Body: { cwd?: unknown } | undefined }>(`${prefix}/sessions/:sessionId/stop`, async (request, reply) => {
|
||||
try {
|
||||
sessions.stop(sessionLookupFromBody(request.params.sessionId, optionalRecord(request.body)));
|
||||
await sessions.stop(sessionLookupFromBody(request.params.sessionId, optionalRecord(request.body)));
|
||||
return { stopped: true };
|
||||
} catch (error) {
|
||||
return reply.code(mutationErrorStatus(error)).send({ error: errorMessage(error) });
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
import type { SavedPromptAttachment } from "../../shared/apiTypes.js";
|
||||
import type {
|
||||
ClientArchiveSessionsResponse,
|
||||
ClientCommand,
|
||||
ClientCommandResult,
|
||||
ClientMessagePage,
|
||||
ClientSession,
|
||||
ClientSessionCleanupExecuteResponse,
|
||||
ClientSessionCleanupPreviewResponse,
|
||||
ClientSessionModel,
|
||||
ClientSessionRef,
|
||||
ClientSessionStatus,
|
||||
ClientThinkingLevel,
|
||||
} from "../types.js";
|
||||
import type { NormalizedSessionCleanupRequest } from "./sessionCleanup.js";
|
||||
|
||||
export type SessionRouteRef = ClientSessionRef;
|
||||
export type SessionRouteLookup = string | SessionRouteRef;
|
||||
|
||||
/**
|
||||
* Route-facing session contract for PI WEB's HTTP/WebSocket API.
|
||||
*
|
||||
* Keep this surface neutral: implementations may be backed by the native Pi SDK,
|
||||
* an out-of-process agent bridge, or another daemon. Pi-specific lifecycle hooks
|
||||
* such as auth-change handling and daemon shutdown stay on the concrete service.
|
||||
*/
|
||||
export interface SessionRouteService {
|
||||
list(cwd: string): Promise<ClientSession[]>;
|
||||
start(cwd: string): Promise<ClientSession>;
|
||||
messages(ref: SessionRouteLookup, page?: { before?: number; limit?: number }): Promise<unknown[] | ClientMessagePage>;
|
||||
status(ref: SessionRouteLookup): Promise<ClientSessionStatus>;
|
||||
availableModels(ref: SessionRouteLookup): Promise<ClientSessionModel[]>;
|
||||
setModel(ref: SessionRouteLookup, provider: string, modelId: string): Promise<ClientSessionStatus>;
|
||||
cycleModel(ref: SessionRouteLookup, direction: "forward" | "backward"): Promise<ClientSessionStatus>;
|
||||
availableThinkingLevels(ref: SessionRouteLookup): Promise<ClientThinkingLevel[]>;
|
||||
setThinkingLevel(ref: SessionRouteLookup, level: string): Promise<ClientSessionStatus>;
|
||||
cycleThinkingLevel(ref: SessionRouteLookup): Promise<ClientSessionStatus>;
|
||||
commands(ref: SessionRouteLookup): Promise<ClientCommand[]>;
|
||||
prompt(ref: SessionRouteLookup, text: unknown, streamingBehavior?: unknown, attachments?: unknown): Promise<void>;
|
||||
saveAttachments(ref: SessionRouteLookup, attachments: unknown, folder?: string): Promise<SavedPromptAttachment[]>;
|
||||
cleanupPreview(request: NormalizedSessionCleanupRequest): Promise<ClientSessionCleanupPreviewResponse>;
|
||||
cleanup(request: NormalizedSessionCleanupRequest): Promise<ClientSessionCleanupExecuteResponse>;
|
||||
shell(ref: SessionRouteLookup, text: string): Promise<void>;
|
||||
runCommand(ref: SessionRouteLookup, text: string): Promise<ClientCommandResult>;
|
||||
respondToCommand(ref: SessionRouteLookup, requestId: string, value: string): Promise<ClientCommandResult>;
|
||||
abort(ref: SessionRouteLookup): Promise<void>;
|
||||
stop(ref: SessionRouteLookup): void | Promise<void>;
|
||||
archive(ref: SessionRouteLookup): Promise<void>;
|
||||
archiveTree(ref: SessionRouteLookup): Promise<ClientArchiveSessionsResponse>;
|
||||
restore(ref: SessionRouteLookup): Promise<void>;
|
||||
deleteArchived(ref: SessionRouteLookup): Promise<void>;
|
||||
reload(ref: SessionRouteLookup): Promise<void>;
|
||||
detachParent(ref: SessionRouteLookup): Promise<void>;
|
||||
}
|
||||
Reference in New Issue
Block a user