From a1f749cdb6e185270a955e77848b364a2c3c68bb Mon Sep 17 00:00:00 2001 From: Federico Jaramillo Martinez Date: Tue, 14 Jul 2026 00:23:44 +0200 Subject: [PATCH] feat: add server session queue clearing --- .changeset/clear-session-message-queue.md | 5 + src/client/src/api/clients.test.ts | 30 +++ src/client/src/api/clients.ts | 1 + .../src/api/federatedRouteContract.test.ts | 1 + src/client/src/components/ChatView.test.ts | 100 +++++++++- src/client/src/components/ChatView.ts | 20 +- .../components/PiWebApp.clearQueue.test.ts | 151 +++++++++++++++ src/client/src/components/PiWebApp.ts | 23 ++- src/client/src/components/shared.ts | 9 +- .../sessionController.clearQueue.test.ts | 179 ++++++++++++++++++ .../src/controllers/sessionController.ts | 24 ++- src/server/app.remoteProxy.test.ts | 18 ++ .../sessiond/sessionProxyRoutes.test.ts | 11 ++ .../piSessionService.promptQueue.test.ts | 76 ++++++++ src/server/sessions/piSessionService.ts | 9 + src/server/sessions/sessionRoutes.test.ts | 68 ++++++- src/server/sessions/sessionRoutes.ts | 8 + src/shared/apiTypes.ts | 1 + src/shared/capabilities.test.ts | 20 ++ src/shared/capabilities.ts | 3 + src/shared/federatedRoutes.ts | 1 + 21 files changed, 735 insertions(+), 23 deletions(-) create mode 100644 .changeset/clear-session-message-queue.md create mode 100644 src/client/src/components/PiWebApp.clearQueue.test.ts create mode 100644 src/client/src/controllers/sessionController.clearQueue.test.ts diff --git a/.changeset/clear-session-message-queue.md b/.changeset/clear-session-message-queue.md new file mode 100644 index 0000000..32eb429 --- /dev/null +++ b/.changeset/clear-session-message-queue.md @@ -0,0 +1,5 @@ +--- +"@jmfederico/pi-web": patch +--- + +Add a capability-aware Clear queue action that removes queued session messages, including prompts held during compaction, without stopping active work. diff --git a/src/client/src/api/clients.test.ts b/src/client/src/api/clients.test.ts index 6146bbc..67b6742 100644 --- a/src/client/src/api/clients.test.ts +++ b/src/client/src/api/clients.test.ts @@ -246,6 +246,36 @@ describe("session API compatibility", () => { expect(url).toBe("https://pi.example.test/api/machines/remote%20a/sessions/s%201/prompt"); expect(JSON.parse(requestBody(init))).toEqual({ cwd: "/repo", text: "hello" }); }); + + it("clears a session queue through an encoded machine route and parses the returned status", async () => { + const fetchMock = stubJsonFetch({ + sessionId: "s /?", + isStreaming: true, + isCompacting: false, + isBashRunning: false, + pendingMessageCount: 0, + tokens: { input: 3, output: 2, cacheRead: 1, cacheWrite: 0, total: 6 }, + cost: 0.25, + ignored: "not part of SessionStatus", + }); + + await expect(sessionsApi.clearQueue({ id: "s /?", cwd: "/repo with spaces" }, "remote /?")).resolves.toEqual({ + sessionId: "s /?", + isStreaming: true, + isCompacting: false, + isBashRunning: false, + pendingMessageCount: 0, + queuedMessages: [], + tokens: { input: 3, output: 2, cacheRead: 1, cacheWrite: 0, total: 6 }, + cost: 0.25, + }); + + expect(fetchMock).toHaveBeenCalledOnce(); + const [url, init] = fetchCall(fetchMock, 0); + expect(url).toBe("https://pi.example.test/api/machines/remote%20%2F%3F/sessions/s%20%2F%3F/queue/clear"); + expect(init?.method).toBe("POST"); + expect(JSON.parse(requestBody(init))).toEqual({ cwd: "/repo with spaces" }); + }); }); describe("machine-scoped file suggestion API", () => { diff --git a/src/client/src/api/clients.ts b/src/client/src/api/clients.ts index 35b9940..16ca7fb 100644 --- a/src/client/src/api/clients.ts +++ b/src/client/src/api/clients.ts @@ -209,6 +209,7 @@ export const sessionsApi = { deleteArchivedMany: (sessions: readonly SessionLookup[], machineId = "local") => request(`${machinePrefix(machineId)}/sessions/bulk/delete-archived`, parseSessionBulkDeleteArchivedResponse, { method: "POST", body: sessionBulkMutationBody(sessions) }), messages: (session: SessionLookup, options?: { limit?: number; before?: number }, machineId = "local") => request(messagePath(session, options, machineId), parseMessagePage), status: (session: SessionLookup, machineId = "local") => request(sessionQueryPath(session, "status", machineId), parseSessionStatus), + clearQueue: (session: SessionLookup, machineId = "local") => request(sessionPath(session, "queue/clear", machineId), parseSessionStatus, { method: "POST", body: sessionBody(session) }), models: (session: SessionLookup, machineId = "local") => request(sessionQueryPath(session, "models", machineId), parseModelSelectionResponse), setModel: (session: SessionLookup, provider: string, modelId: string, machineId = "local") => request(sessionPath(session, "model", machineId), parseSessionStatus, { method: "POST", body: sessionBody(session, { provider, modelId }) }), cycleModel: (session: SessionLookup, direction: "forward" | "backward", machineId = "local") => request(sessionPath(session, "model/cycle", machineId), parseSessionStatus, { method: "POST", body: sessionBody(session, { direction }) }), diff --git a/src/client/src/api/federatedRouteContract.test.ts b/src/client/src/api/federatedRouteContract.test.ts index 28a5257..52a53c2 100644 --- a/src/client/src/api/federatedRouteContract.test.ts +++ b/src/client/src/api/federatedRouteContract.test.ts @@ -64,6 +64,7 @@ describe("federated route contract", () => { ignoreParseFailure(sessionsApi.deleteArchivedMany([session], machineId)), ignoreParseFailure(sessionsApi.messages(session, { limit: 20, before: 10 }, machineId)), ignoreParseFailure(sessionsApi.status(session, machineId)), + ignoreParseFailure(sessionsApi.clearQueue(session, machineId)), ignoreParseFailure(sessionsApi.models(session, machineId)), ignoreParseFailure(sessionsApi.setModel(session, "openai", "gpt", machineId)), ignoreParseFailure(sessionsApi.cycleModel(session, "forward", machineId)), diff --git a/src/client/src/components/ChatView.test.ts b/src/client/src/components/ChatView.test.ts index df3b167..2282bf1 100644 --- a/src/client/src/components/ChatView.test.ts +++ b/src/client/src/components/ChatView.test.ts @@ -1,5 +1,6 @@ import type { TemplateResult } from "lit"; -import { describe, expect, it } from "vitest"; +import { describe, expect, it, vi } from "vitest"; +import type { QueuedSessionMessage, SessionStatus } from "../api"; import type { ChatLine } from "./shared"; import { ChatView, chatMessageMetadataLabel, chatQueuedMessageSections } from "./ChatView"; @@ -12,19 +13,61 @@ describe("chatQueuedMessageSections", () => { expect(sections).toEqual([ { + source: "client", heading: "Queued until session starts", detail: "Will send once the backend session is ready", messages: [{ kind: "followUp", text: "queued before start" }], }, { + source: "server", heading: "Queued messages", - detail: "1 pending · Stop clears the queue", + detail: "1 pending", messages: [{ kind: "steer", text: "server queued" }], }, ]); }); }); +describe("ChatView queued-message clear action", () => { + // Direct handler extraction keeps this node-environment test focused on the + // Clear queue template wiring without introducing a component-wide DOM shim. + it("renders an accessible server-queue action and invokes its callback", () => { + const view = new ChatView(); + const onClearServerQueue = vi.fn(); + view.status = queuedStatus([{ kind: "steer", text: "server queued" }]); + view.canClearServerQueue = true; + view.onClearServerQueue = onClearServerQueue; + + const rendered = renderQueuedMessages(view); + const markup = templateStaticMarkup(rendered); + + expect(markup).toContain('type="button"'); + expect(markup).toContain('title="Clear queued messages without stopping active work"'); + expect(markup).toContain(">Clear queue"); + templateEventHandler(rendered, "Clear queue")(new Event("click")); + expect(onClearServerQueue).toHaveBeenCalledOnce(); + }); + + it("hides the action when the selected runtime does not support clearing", () => { + const view = new ChatView(); + view.status = queuedStatus([{ kind: "followUp", text: "server queued" }]); + view.canClearServerQueue = false; + view.onClearServerQueue = vi.fn(); + + expect(templateStaticMarkup(renderQueuedMessages(view))).not.toContain("Clear queue"); + }); + + it("does not expose the server action for the separate client pending-start queue", () => { + const view = new ChatView(); + view.status = queuedStatus([]); + view.clientQueuedMessages = [{ kind: "followUp", text: "waiting for session start" }]; + view.canClearServerQueue = true; + view.onClearServerQueue = vi.fn(); + + expect(templateStaticMarkup(renderQueuedMessages(view))).not.toContain("Clear queue"); + }); +}); + describe("chatMessageMetadataLabel", () => { it("uses one full date and model label without a model prefix", () => { const timestamp = "2026-07-10T19:15:30.000Z"; @@ -103,10 +146,17 @@ interface GroupBodyRenderCall { startIndex: number; } +type RenderQueuedMessages = (this: ChatView) => TemplateResult; type RenderMessageGroup = (this: ChatView, messages: ChatLine[], startIndex: number, endIndex: number, defaultOpen: boolean) => TemplateResult; type RenderMessageGroupBody = (this: ChatView, messages: ChatLine[], startIndex: number) => TemplateResult; type TemplateEventHandler = (event: Event) => void; +function renderQueuedMessages(view: ChatView): TemplateResult { + const method: unknown = Reflect.get(view, "renderQueuedMessages"); + if (!isRenderQueuedMessages(method)) throw new Error("ChatView.renderQueuedMessages is not callable"); + return method.call(view); +} + function renderMessageGroup(view: ChatView, messages: ChatLine[], startIndex: number, endIndex: number, defaultOpen: boolean): TemplateResult { const method: unknown = Reflect.get(view, "renderMessageGroup"); if (!isRenderMessageGroup(method)) throw new Error("ChatView.renderMessageGroup is not callable"); @@ -125,6 +175,10 @@ function observeGroupBodyRenders(view: ChatView): GroupBodyRenderCall[] { return calls; } +function isRenderQueuedMessages(value: unknown): value is RenderQueuedMessages { + return typeof value === "function"; +} + function isRenderMessageGroup(value: unknown): value is RenderMessageGroup { return typeof value === "function"; } @@ -134,13 +188,30 @@ function isRenderMessageGroupBody(value: unknown): value is RenderMessageGroupBo } function templateEventHandler(template: TemplateResult, marker: string): TemplateEventHandler { - const strings = templateStrings(template); - const values = templateValues(template); - for (let index = 0; index < values.length; index += 1) { - const value = values[index]; - if (strings[index]?.includes(marker) === true && isTemplateEventHandler(value)) return value; + let handler: TemplateEventHandler | undefined; + visit(template); + if (handler === undefined) throw new Error(`Expected template event handler near ${marker}`); + return handler; + + function visit(value: unknown): void { + if (handler !== undefined) return; + if (Array.isArray(value)) { + for (const item of value) visit(item); + return; + } + if (!isTemplateResult(value)) return; + const strings = templateStrings(value); + const values = templateValues(value); + for (let index = 0; index < values.length; index += 1) { + const candidate = values[index]; + const isNearMarker = strings[index]?.includes(marker) === true || strings[index + 1]?.includes(marker) === true; + if (isNearMarker && isTemplateEventHandler(candidate)) { + handler = candidate; + return; + } + visit(candidate); + } } - throw new Error(`Expected template event handler after ${marker}`); } function isTemplateEventHandler(value: unknown): value is TemplateEventHandler { @@ -221,3 +292,16 @@ function isTemplateResult(value: unknown): value is TemplateResult { function isStringArray(value: unknown): value is string[] { return Array.isArray(value) && value.every((item: unknown) => typeof item === "string"); } + +function queuedStatus(queuedMessages: QueuedSessionMessage[]): SessionStatus { + return { + sessionId: "session-1", + isStreaming: true, + isCompacting: false, + isBashRunning: false, + pendingMessageCount: queuedMessages.length, + queuedMessages, + tokens: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + cost: 0, + }; +} diff --git a/src/client/src/components/ChatView.ts b/src/client/src/components/ChatView.ts index 8cc549e..a6e62da 100644 --- a/src/client/src/components/ChatView.ts +++ b/src/client/src/components/ChatView.ts @@ -39,6 +39,7 @@ function clampNumber(value: number, min: number, max: number): number { } export interface QueuedMessageSection { + source: "client" | "server"; heading: string; detail: string; messages: QueuedSessionMessage[]; @@ -46,8 +47,8 @@ export interface QueuedMessageSection { export function chatQueuedMessageSections(clientQueued: QueuedSessionMessage[], serverQueued: QueuedSessionMessage[]): QueuedMessageSection[] { return [ - clientQueued.length === 0 ? undefined : { heading: "Queued until session starts", detail: "Will send once the backend session is ready", messages: clientQueued }, - serverQueued.length === 0 ? undefined : { heading: "Queued messages", detail: `${String(serverQueued.length)} pending · Stop clears the queue`, messages: serverQueued }, + clientQueued.length === 0 ? undefined : { source: "client", heading: "Queued until session starts", detail: "Will send once the backend session is ready", messages: clientQueued }, + serverQueued.length === 0 ? undefined : { source: "server", heading: "Queued messages", detail: `${String(serverQueued.length)} pending`, messages: serverQueued }, ].filter((section): section is QueuedMessageSection => section !== undefined); } @@ -89,6 +90,8 @@ export class ChatView extends LitElement { @property({ attribute: false }) clientQueuedMessages: QueuedSessionMessage[] = []; @property({ attribute: false }) status?: SessionStatus; @property({ attribute: false }) activity?: SessionActivity; + @property({ type: Boolean }) canClearServerQueue = false; + @property({ attribute: false }) onClearServerQueue?: () => void; @property({ attribute: false }) onLoadMore?: () => void; @query(".chat") private chat?: HTMLDivElement; @state() private pinnedToBottom = true; @@ -123,6 +126,9 @@ export class ChatView extends LitElement { private readonly onPageHide = () => { this.saveScrollPosition(); }; + private readonly handleClearServerQueue = (): void => { + this.onClearServerQueue?.(); + }; override connectedCallback(): void { super.connectedCallback(); @@ -261,11 +267,17 @@ export class ChatView extends LitElement { } private renderQueuedMessageList(section: QueuedMessageSection) { + const canClear = section.source === "server" && this.canClearServerQueue && this.onClearServerQueue !== undefined; return html`