Archived
feat: add server session queue clearing
This commit is contained in:
@@ -212,6 +212,82 @@ describe("PiSessionService prompt, queue, and auth warnings", () => {
|
||||
await service.dispose();
|
||||
});
|
||||
|
||||
it("clears runtime and compaction queues without interrupting active work", async () => {
|
||||
const steeringMessages = ["adjust this turn"];
|
||||
const followUpMessages = ["then do this"];
|
||||
const transcript = [{ role: "user", content: "keep this history" }];
|
||||
const hub = new CapturingSessionEventHub();
|
||||
const fake = fakeRuntime("clear-queue-session", {
|
||||
messages: transcript,
|
||||
isStreaming: true,
|
||||
isCompacting: true,
|
||||
pendingMessageCount: 2,
|
||||
getSteeringMessages: () => steeringMessages,
|
||||
getFollowUpMessages: () => followUpMessages,
|
||||
});
|
||||
const clearRuntimeQueue = vi.fn(() => {
|
||||
const cleared = { steering: [...steeringMessages], followUp: [...followUpMessages] };
|
||||
steeringMessages.length = 0;
|
||||
followUpMessages.length = 0;
|
||||
fake.session.pendingMessageCount = 0;
|
||||
return cleared;
|
||||
});
|
||||
fake.session.clearQueue = clearRuntimeQueue;
|
||||
const service = new PiSessionService(hub, {
|
||||
createAgentRuntime: runtimeCreator(fake.runtime),
|
||||
sessionManager: sessionGateway([sessionRecord("clear-queue-session")]),
|
||||
heartbeatIntervalMs: 60_000,
|
||||
});
|
||||
|
||||
await service.prompt(sessionRef("clear-queue-session"), "queued during compaction", "followUp");
|
||||
await expect(service.status(sessionRef("clear-queue-session"))).resolves.toMatchObject({
|
||||
isStreaming: true,
|
||||
isCompacting: true,
|
||||
pendingMessageCount: 3,
|
||||
queuedMessages: [
|
||||
{ kind: "steer", text: "adjust this turn" },
|
||||
{ kind: "followUp", text: "then do this" },
|
||||
{ kind: "followUp", text: "queued during compaction" },
|
||||
],
|
||||
});
|
||||
|
||||
const status = await service.clearQueue(sessionRef("clear-queue-session"));
|
||||
|
||||
expect(clearRuntimeQueue).toHaveBeenCalledOnce();
|
||||
expect(status).toMatchObject({
|
||||
isStreaming: true,
|
||||
isCompacting: true,
|
||||
pendingMessageCount: 0,
|
||||
queuedMessages: [],
|
||||
messageCount: 1,
|
||||
});
|
||||
expect(fake.session.messages).toBe(transcript);
|
||||
expect(fake.calls.prompt).toEqual([]);
|
||||
expect(fake.calls.abort).toBe(0);
|
||||
expect(fake.calls.dispose).toBe(0);
|
||||
const publishedStatuses = hub.sessionEvents.filter(({ event }) => event.type === "status.update");
|
||||
expect(publishedStatuses.at(-1)?.event).toEqual({ type: "status.update", status });
|
||||
await service.dispose();
|
||||
});
|
||||
|
||||
it("clears an already-empty queue idempotently", async () => {
|
||||
const fake = fakeRuntime("clear-empty-queue-session");
|
||||
const service = new PiSessionService(new CapturingSessionEventHub(), {
|
||||
createAgentRuntime: runtimeCreator(fake.runtime),
|
||||
sessionManager: sessionGateway([sessionRecord("clear-empty-queue-session")]),
|
||||
heartbeatIntervalMs: 60_000,
|
||||
});
|
||||
|
||||
const firstStatus = await service.clearQueue(sessionRef("clear-empty-queue-session"));
|
||||
const secondStatus = await service.clearQueue(sessionRef("clear-empty-queue-session"));
|
||||
|
||||
expect(fake.calls.clearQueue).toBe(2);
|
||||
expect(fake.calls.abort).toBe(0);
|
||||
expect(firstStatus).toMatchObject({ pendingMessageCount: 0, queuedMessages: [] });
|
||||
expect(secondStatus).toMatchObject({ pendingMessageCount: 0, queuedMessages: [] });
|
||||
await service.dispose();
|
||||
});
|
||||
|
||||
it("clears queued messages when aborting active work", async () => {
|
||||
const fake = fakeRuntime("abort-session");
|
||||
const service = new PiSessionService(new CapturingSessionEventHub(), {
|
||||
|
||||
@@ -1363,6 +1363,15 @@ export class PiSessionService {
|
||||
this.unregisterSubsession(session.sessionId);
|
||||
}
|
||||
|
||||
async clearQueue(ref: PiSessionLookup): Promise<ClientSessionStatus> {
|
||||
await this.assertWritable(ref);
|
||||
const session = await this.getOrOpen(ref);
|
||||
this.clearCompactionPromptQueue(session.sessionId);
|
||||
clearSessionQueue(session);
|
||||
this.publishStatus(session);
|
||||
return this.statusFromSession(session);
|
||||
}
|
||||
|
||||
async abort(ref: PiSessionLookup): Promise<void> {
|
||||
const active = this.activeForLookup(ref);
|
||||
if (active === undefined) return;
|
||||
|
||||
@@ -2,7 +2,7 @@ import { resolve } from "node:path";
|
||||
import Fastify, { type FastifyInstance } from "fastify";
|
||||
import fastifyWebsocket from "@fastify/websocket";
|
||||
import { afterEach, beforeEach, describe, expect, it } from "vitest";
|
||||
import type { MessagePage, SessionBulkArchiveResponse, SessionBulkDeleteArchivedResponse, SessionBulkMutationRef, SessionCleanupExecuteResponse, SessionCleanupPreviewResponse } from "../../shared/apiTypes.js";
|
||||
import type { MessagePage, SessionBulkArchiveResponse, SessionBulkDeleteArchivedResponse, SessionBulkMutationRef, SessionCleanupExecuteResponse, SessionCleanupPreviewResponse, SessionStatus } from "../../shared/apiTypes.js";
|
||||
import { SessionEventHub } from "../realtime/sessionEventHub.js";
|
||||
import { PiSessionService, type PiSessionManagerGateway, type PiSessionRef } from "./piSessionService.js";
|
||||
import { registerSessionRoutes } from "./sessionRoutes.js";
|
||||
@@ -165,6 +165,55 @@ describe("session routes", () => {
|
||||
}
|
||||
});
|
||||
|
||||
it("clears a session queue with workspace context and returns fresh status", async () => {
|
||||
const routeApp = Fastify({ logger: false });
|
||||
await routeApp.register(fastifyWebsocket);
|
||||
const eventHub = new SessionEventHub();
|
||||
const routeService = new CapturingRouteSessionService(eventHub);
|
||||
registerSessionRoutes(routeApp, routeService, eventHub);
|
||||
|
||||
try {
|
||||
const requestCwd = resolve("/repo");
|
||||
const response = await routeApp.inject({ method: "POST", url: "/sessions/session-1/queue/clear", payload: { cwd: requestCwd } });
|
||||
|
||||
expect(response.statusCode).toBe(200);
|
||||
expect(response.json()).toEqual({
|
||||
sessionId: "session-1",
|
||||
isStreaming: true,
|
||||
isCompacting: false,
|
||||
isBashRunning: false,
|
||||
pendingMessageCount: 0,
|
||||
queuedMessages: [],
|
||||
tokens: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
cost: 0,
|
||||
});
|
||||
expect(routeService.clearQueueCalls).toEqual([{ id: "session-1", cwd: requestCwd }]);
|
||||
} finally {
|
||||
await routeService.dispose();
|
||||
await routeApp.close();
|
||||
}
|
||||
});
|
||||
|
||||
it("maps archived queue-clear failures to a mutation error without requiring a body", async () => {
|
||||
const routeApp = Fastify({ logger: false });
|
||||
await routeApp.register(fastifyWebsocket);
|
||||
const eventHub = new SessionEventHub();
|
||||
const routeService = new CapturingRouteSessionService(eventHub);
|
||||
routeService.clearQueueError = new Error("Archived sessions are read-only. Restore the session to continue.");
|
||||
registerSessionRoutes(routeApp, routeService, eventHub);
|
||||
|
||||
try {
|
||||
const response = await routeApp.inject({ method: "POST", url: "/sessions/session-1/queue/clear" });
|
||||
|
||||
expect(response.statusCode).toBe(400);
|
||||
expect(response.json()).toEqual({ error: "Archived sessions are read-only. Restore the session to continue." });
|
||||
expect(routeService.clearQueueCalls).toEqual(["session-1"]);
|
||||
} finally {
|
||||
await routeService.dispose();
|
||||
await routeApp.close();
|
||||
}
|
||||
});
|
||||
|
||||
it("normalizes cleanup requests for preview and execute routes", async () => {
|
||||
const routeApp = Fastify({ logger: false });
|
||||
await routeApp.register(fastifyWebsocket);
|
||||
@@ -252,12 +301,14 @@ describe("session routes", () => {
|
||||
class CapturingRouteSessionService extends PiSessionService {
|
||||
readonly calls: unknown[] = [];
|
||||
readonly reloadCalls: (string | PiSessionRef)[] = [];
|
||||
readonly clearQueueCalls: (string | PiSessionRef)[] = [];
|
||||
messagesResponse: unknown[] | MessagePage = [];
|
||||
readonly cleanupPreviewCalls: NormalizedSessionCleanupRequest[] = [];
|
||||
readonly cleanupCalls: NormalizedSessionCleanupRequest[] = [];
|
||||
readonly bulkArchiveCalls: SessionBulkMutationRef[][] = [];
|
||||
readonly bulkDeleteCalls: SessionBulkMutationRef[][] = [];
|
||||
reloadError: Error | undefined;
|
||||
clearQueueError: Error | undefined;
|
||||
|
||||
constructor(eventHub: SessionEventHub) {
|
||||
super(eventHub, { sessionManager: new RejectingSessionManager(), heartbeatIntervalMs: 60_000 });
|
||||
@@ -289,6 +340,21 @@ class CapturingRouteSessionService extends PiSessionService {
|
||||
return Promise.resolve();
|
||||
}
|
||||
|
||||
override clearQueue(lookup: string | PiSessionRef): Promise<SessionStatus> {
|
||||
this.clearQueueCalls.push(lookup);
|
||||
if (this.clearQueueError !== undefined) return Promise.reject(this.clearQueueError);
|
||||
return Promise.resolve({
|
||||
sessionId: sessionIdFromLookup(lookup),
|
||||
isStreaming: true,
|
||||
isCompacting: false,
|
||||
isBashRunning: false,
|
||||
pendingMessageCount: 0,
|
||||
queuedMessages: [],
|
||||
tokens: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
cost: 0,
|
||||
});
|
||||
}
|
||||
|
||||
override messages(): Promise<unknown[] | MessagePage> {
|
||||
return Promise.resolve(this.messagesResponse);
|
||||
}
|
||||
|
||||
@@ -173,6 +173,14 @@ export function registerSessionRoutes(app: FastifyInstance, sessions: PiSessionS
|
||||
}
|
||||
});
|
||||
|
||||
app.post<{ Params: { sessionId: string }; Body: { cwd?: unknown } | undefined }>(`${prefix}/sessions/:sessionId/queue/clear`, async (request, reply) => {
|
||||
try {
|
||||
return await sessions.clearQueue(sessionLookupFromBody(request.params.sessionId, optionalRecord(request.body)));
|
||||
} catch (error) {
|
||||
return reply.code(mutationErrorStatus(error)).send({ error: errorMessage(error) });
|
||||
}
|
||||
});
|
||||
|
||||
app.post<{ Params: { sessionId: string }; Body: AttachmentsRequestBody | undefined }>(`${prefix}/sessions/:sessionId/attachments`, async (request, reply) => {
|
||||
try {
|
||||
const body = optionalRecord(request.body);
|
||||
|
||||
Reference in New Issue
Block a user