Archived
Improve queued message handling
This commit is contained in:
@@ -33,7 +33,7 @@ function asRuntimeFactoryResult(runtime: AgentSessionRuntime): RuntimeFactoryRes
|
||||
|
||||
function fakeRuntime(sessionId = "session-1") {
|
||||
const promptCalls: { text: string; options: unknown }[] = [];
|
||||
const calls = { abort: 0, dispose: 0, prompt: promptCalls };
|
||||
const calls = { abort: 0, clearQueue: 0, dispose: 0, prompt: promptCalls };
|
||||
const session = {
|
||||
sessionId,
|
||||
sessionFile: `/tmp/${sessionId}.jsonl`,
|
||||
@@ -57,6 +57,12 @@ function fakeRuntime(sessionId = "session-1") {
|
||||
calls.abort += 1;
|
||||
return Promise.resolve();
|
||||
},
|
||||
clearQueue: () => {
|
||||
calls.clearQueue += 1;
|
||||
return { steering: [], followUp: [] };
|
||||
},
|
||||
getSteeringMessages: () => [],
|
||||
getFollowUpMessages: () => [],
|
||||
} as unknown as AgentSession;
|
||||
const runtime = {
|
||||
session,
|
||||
@@ -145,4 +151,117 @@ describe("PiSessionService", () => {
|
||||
expect(fake.calls.prompt).toEqual([{ text: "Build the thing", options: undefined }]);
|
||||
await service.dispose();
|
||||
});
|
||||
|
||||
it("includes queued message details in session status", async () => {
|
||||
const fake = fakeRuntime("status-session");
|
||||
(fake.runtime.session as unknown as { pendingMessageCount: number; getSteeringMessages: () => string[]; getFollowUpMessages: () => string[] }).pendingMessageCount = 2;
|
||||
(fake.runtime.session as unknown as { getSteeringMessages: () => string[] }).getSteeringMessages = () => ["adjust this turn"];
|
||||
(fake.runtime.session as unknown as { getFollowUpMessages: () => string[] }).getFollowUpMessages = () => ["then do this"];
|
||||
const service = new PiSessionService(new CapturingSessionEventHub(), {
|
||||
createRuntime: () => Promise.resolve(asRuntimeFactoryResult(fake.runtime)),
|
||||
createAgentRuntime: () => Promise.resolve(fake.runtime),
|
||||
sessionManager: {
|
||||
create: () => fakeSessionManager(),
|
||||
list: () => Promise.resolve([]),
|
||||
listAll: () => Promise.resolve([{ id: "status-session", path: "/sessions/status-session.jsonl", cwd: "/workspace", created: new Date("2026-01-01T00:00:00.000Z"), modified: new Date("2026-01-01T00:01:00.000Z"), messageCount: 0, firstMessage: "", allMessagesText: "" }]),
|
||||
open: () => fakeSessionManager(),
|
||||
},
|
||||
heartbeatIntervalMs: 60_000,
|
||||
});
|
||||
|
||||
await expect(service.status("status-session")).resolves.toMatchObject({
|
||||
pendingMessageCount: 2,
|
||||
queuedMessages: [{ kind: "steer", text: "adjust this turn" }, { kind: "followUp", text: "then do this" }],
|
||||
});
|
||||
await service.dispose();
|
||||
});
|
||||
|
||||
it("does not enqueue duplicate queued message text", async () => {
|
||||
const fake = fakeRuntime("dedupe-session");
|
||||
(fake.runtime.session as unknown as { isStreaming: boolean; pendingMessageCount: number; getFollowUpMessages: () => string[] }).isStreaming = true;
|
||||
(fake.runtime.session as unknown as { pendingMessageCount: number }).pendingMessageCount = 1;
|
||||
(fake.runtime.session as unknown as { getFollowUpMessages: () => string[] }).getFollowUpMessages = () => ["already queued"];
|
||||
const service = new PiSessionService(new CapturingSessionEventHub(), {
|
||||
createRuntime: () => Promise.resolve(asRuntimeFactoryResult(fake.runtime)),
|
||||
createAgentRuntime: () => Promise.resolve(fake.runtime),
|
||||
sessionManager: {
|
||||
create: () => fakeSessionManager(),
|
||||
list: () => Promise.resolve([]),
|
||||
listAll: () => Promise.resolve([{ id: "dedupe-session", path: "/sessions/dedupe-session.jsonl", cwd: "/workspace", created: new Date("2026-01-01T00:00:00.000Z"), modified: new Date("2026-01-01T00:01:00.000Z"), messageCount: 0, firstMessage: "", allMessagesText: "" }]),
|
||||
open: () => fakeSessionManager(),
|
||||
},
|
||||
heartbeatIntervalMs: 60_000,
|
||||
});
|
||||
|
||||
await service.prompt("dedupe-session", "already queued", "followUp");
|
||||
|
||||
expect(fake.calls.prompt).toEqual([]);
|
||||
await service.dispose();
|
||||
});
|
||||
|
||||
it("does not append queued prompts to the transcript before delivery", async () => {
|
||||
const hub = new CapturingSessionEventHub();
|
||||
const fake = fakeRuntime("queued-session");
|
||||
(fake.runtime.session as unknown as { isStreaming: boolean }).isStreaming = true;
|
||||
const service = new PiSessionService(hub, {
|
||||
createRuntime: () => Promise.resolve(asRuntimeFactoryResult(fake.runtime)),
|
||||
createAgentRuntime: () => Promise.resolve(fake.runtime),
|
||||
sessionManager: {
|
||||
create: () => fakeSessionManager(),
|
||||
list: () => Promise.resolve([]),
|
||||
listAll: () => Promise.resolve([{ id: "queued-session", path: "/sessions/queued-session.jsonl", cwd: "/workspace", created: new Date("2026-01-01T00:00:00.000Z"), modified: new Date("2026-01-01T00:01:00.000Z"), messageCount: 0, firstMessage: "", allMessagesText: "" }]),
|
||||
open: () => fakeSessionManager(),
|
||||
},
|
||||
heartbeatIntervalMs: 60_000,
|
||||
});
|
||||
|
||||
await service.prompt("queued-session", "Wait for the current turn", "followUp");
|
||||
|
||||
expect(fake.calls.prompt).toEqual([{ text: "Wait for the current turn", options: { streamingBehavior: "followUp" } }]);
|
||||
expect(hub.sessionEvents.some(({ event }) => event.type === "message.append")).toBe(false);
|
||||
await service.dispose();
|
||||
});
|
||||
|
||||
it("clears queued messages when aborting active work", async () => {
|
||||
const fake = fakeRuntime("abort-session");
|
||||
const service = new PiSessionService(new CapturingSessionEventHub(), {
|
||||
createRuntime: () => Promise.resolve(asRuntimeFactoryResult(fake.runtime)),
|
||||
createAgentRuntime: () => Promise.resolve(fake.runtime),
|
||||
sessionManager: {
|
||||
create: () => fakeSessionManager(),
|
||||
list: () => Promise.resolve([]),
|
||||
listAll: () => Promise.resolve([{ id: "abort-session", path: "/sessions/abort-session.jsonl", cwd: "/workspace", created: new Date("2026-01-01T00:00:00.000Z"), modified: new Date("2026-01-01T00:01:00.000Z"), messageCount: 0, firstMessage: "", allMessagesText: "" }]),
|
||||
open: () => fakeSessionManager(),
|
||||
},
|
||||
heartbeatIntervalMs: 60_000,
|
||||
});
|
||||
|
||||
await service.status("abort-session");
|
||||
await service.abort("abort-session");
|
||||
|
||||
expect(fake.calls.clearQueue).toBe(1);
|
||||
expect(fake.calls.abort).toBe(1);
|
||||
await service.dispose();
|
||||
});
|
||||
|
||||
it("clears queued messages when stopping a session runtime", async () => {
|
||||
const fake = fakeRuntime("stop-session");
|
||||
const service = new PiSessionService(new CapturingSessionEventHub(), {
|
||||
createRuntime: () => Promise.resolve(asRuntimeFactoryResult(fake.runtime)),
|
||||
createAgentRuntime: () => Promise.resolve(fake.runtime),
|
||||
sessionManager: {
|
||||
create: () => fakeSessionManager(),
|
||||
list: () => Promise.resolve([]),
|
||||
listAll: () => Promise.resolve([{ id: "stop-session", path: "/sessions/stop-session.jsonl", cwd: "/workspace", created: new Date("2026-01-01T00:00:00.000Z"), modified: new Date("2026-01-01T00:01:00.000Z"), messageCount: 0, firstMessage: "", allMessagesText: "" }]),
|
||||
open: () => fakeSessionManager(),
|
||||
},
|
||||
heartbeatIntervalMs: 60_000,
|
||||
});
|
||||
|
||||
await service.status("stop-session");
|
||||
service.stop("stop-session");
|
||||
|
||||
expect(fake.calls.clearQueue).toBe(1);
|
||||
await service.dispose();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -170,9 +170,15 @@ export class PiSessionService {
|
||||
await this.assertWritable(sessionId);
|
||||
const session = await this.getOrOpen(sessionId);
|
||||
this.maybeGenerateSessionName(session, text);
|
||||
const behavior = session.isStreaming || session.isCompacting ? streamingBehavior ?? "followUp" : undefined;
|
||||
const isQueued = session.isStreaming || session.isCompacting;
|
||||
const behavior = isQueued ? streamingBehavior ?? "followUp" : undefined;
|
||||
if (isQueued && hasQueuedMessageText(session, text)) {
|
||||
this.publishActivity(session, "duplicate queued message ignored", "active");
|
||||
this.publishStatus(session);
|
||||
return;
|
||||
}
|
||||
this.publishActivity(session, session.isCompacting ? "message queued during compaction" : behavior === "steer" ? "steering queued" : behavior === "followUp" ? "message queued" : "prompt accepted", "active");
|
||||
this.events.publish(sessionId, { type: "message.append", message: userTextMessage(text) });
|
||||
if (!isQueued) this.events.publish(sessionId, { type: "message.append", message: userTextMessage(text) });
|
||||
void session.prompt(text, behavior === undefined ? undefined : { streamingBehavior: behavior }).catch((error: unknown) => {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
this.publishActivity(session, "error", "error", message);
|
||||
@@ -246,6 +252,7 @@ export class PiSessionService {
|
||||
async abort(sessionId: string): Promise<void> {
|
||||
const active = this.active.get(sessionId);
|
||||
if (!active) return;
|
||||
clearSessionQueue(active.runtime.session);
|
||||
await active.runtime.session.abort();
|
||||
this.publishActivity(active.runtime.session, "stopped", "idle");
|
||||
this.publishStatus(active.runtime.session);
|
||||
@@ -254,6 +261,7 @@ export class PiSessionService {
|
||||
stop(sessionId: string): void {
|
||||
const active = this.active.get(sessionId);
|
||||
if (!active) return;
|
||||
clearSessionQueue(active.runtime.session);
|
||||
active.unsubscribe();
|
||||
void active.runtime.session.abort().finally(() => active.runtime.dispose());
|
||||
this.active.delete(sessionId);
|
||||
@@ -416,6 +424,7 @@ export class PiSessionService {
|
||||
isCompacting: session.isCompacting,
|
||||
isBashRunning: session.isBashRunning,
|
||||
pendingMessageCount: session.pendingMessageCount,
|
||||
queuedMessages: queuedMessagesFromSession(session),
|
||||
tokens: stats.tokens,
|
||||
cost: stats.cost,
|
||||
...(contextUsage === undefined ? {} : { contextUsage }),
|
||||
@@ -435,6 +444,23 @@ async function clearParentSession(sessionFile: string): Promise<void> {
|
||||
await writeFile(sessionFile, `${JSON.stringify(header)}${rest}`, "utf8");
|
||||
}
|
||||
|
||||
function clearSessionQueue(session: AgentSession): void {
|
||||
const candidate = session as AgentSession & { clearQueue?: () => unknown };
|
||||
candidate.clearQueue?.();
|
||||
}
|
||||
|
||||
function hasQueuedMessageText(session: AgentSession, text: string): boolean {
|
||||
return queuedMessagesFromSession(session).some((message) => message.text === text);
|
||||
}
|
||||
|
||||
function queuedMessagesFromSession(session: AgentSession): { kind: "steer" | "followUp"; text: string }[] {
|
||||
const candidate = session as AgentSession & { getSteeringMessages?: () => string[]; getFollowUpMessages?: () => string[] };
|
||||
return [
|
||||
...(candidate.getSteeringMessages?.() ?? []).map((text) => ({ kind: "steer" as const, text })),
|
||||
...(candidate.getFollowUpMessages?.() ?? []).map((text) => ({ kind: "followUp" as const, text })),
|
||||
];
|
||||
}
|
||||
|
||||
function userTextMessage(text: string): { role: "user"; content: string } {
|
||||
return { role: "user", content: text };
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user