fix(coding-agent-turn-interrupt/queue-ux): resolved steering abort state

- Replaced queued-message interrupt flow with session abort calls on empty submit and escape.
- Removed interrupting state and notifyInterrupting teardown paths from abort handling.
- Updated AgentSession queue operations to use shared steering and follow-up queue views.
- Propagated isAborting through session state and collab payloads to suppress late updates.
This commit is contained in:
can1357
2026-06-13 17:19:57 +02:00
parent 3def4a6979
commit 705750453d
23 changed files with 440 additions and 1304 deletions
+21 -10
View File
@@ -706,6 +706,11 @@ export class Agent {
this.#state.messages = ms.slice();
}
replaceQueues(steering: AgentMessage[], followUp: AgentMessage[]) {
this.#steeringQueue = steering.slice();
this.#followUpQueue = followUp.slice();
}
appendMessage(m: AgentMessage) {
this.#state.messages.push(m);
}
@@ -751,16 +756,22 @@ export class Agent {
return this.#steeringQueue.length > 0 || this.#followUpQueue.length > 0;
}
/**
* Drain queued messages for an immediate resume: all steering messages if any
* (honoring steeringMode), otherwise follow-ups. Mirrors continue()'s dequeue
* precedence so an empty-Enter flush and the natural turn-end resume agree on
* what to deliver next.
*/
takeQueuedMessages(): AgentMessage[] {
const steering = this.#dequeueSteeringMessages();
if (steering.length > 0) return steering;
return this.#dequeueFollowUpMessages();
/** Non-consuming view of the pending steering queue (insertion order, newest
* last). The session layer derives its queued-message display/count from
* this live view instead of a mirror, so the agent-core queue stays the
* single source of truth. */
peekSteeringQueue(): readonly AgentMessage[] {
return this.#steeringQueue;
}
/** Non-consuming view of the pending follow-up queue. See
* {@link peekSteeringQueue}. */
peekFollowUpQueue(): readonly AgentMessage[] {
return this.#followUpQueue;
}
get isAborting(): boolean {
return this.#abortController?.signal.aborted === true && this.#state.isStreaming;
}
#dequeueSteeringMessages(): AgentMessage[] {
+3 -7
View File
@@ -365,13 +365,8 @@ export class CollabHost {
const name = peer.name;
const content: string | (TextContent | ImageContent)[] =
images && images.length > 0 ? [{ type: "text", text }, ...images] : text;
const details: CollabPromptDetails & { __pendingDisplayTag?: string } = { from: name };
const details: CollabPromptDetails = { from: name };
if (this.#ctx.session.isStreaming) {
// Mid-turn guest prompts are steered: register the pending-display twin
// so queuedMessageCount reflects the queued steer (host pending bar +
// guests' "queued ×N" badge). The tag dequeues the entry when the agent
// consumes the message (mirrors the skill-prompt path).
details.__pendingDisplayTag = this.#ctx.session.enqueueCustomMessageDisplay(text, "steer");
this.#ctx.updatePendingMessagesDisplay();
this.#ctx.ui.requestRender();
this.#scheduleStateBroadcast();
@@ -385,7 +380,7 @@ export class CollabHost {
details,
attribution: "user",
},
{ streamingBehavior: "steer" },
{ streamingBehavior: "steer", queueChipText: text },
)
.catch(err => {
logger.warn("collab guest prompt failed", { error: String(err) });
@@ -422,6 +417,7 @@ export class CollabHost {
const breakdown = this.#ctx.statusLine.getCachedContextBreakdown();
return {
isStreaming: session.isStreaming,
isAborting: session.isAborting,
queuedMessageCount: session.queuedMessageCount,
sessionName: session.sessionName,
cwd: this.#ctx.sessionManager.getCwd(),
@@ -39,10 +39,10 @@ import type { SessionMessageEntry } from "../../session/session-manager";
import { parseSessionEntries } from "../../session/session-manager";
import { createIrcMessageCard } from "../../tools/irc";
import { replaceTabs, TRUNCATE_LENGTHS, truncateToWidth } from "../../tools/render-utils";
import { hasVisibleThinking } from "../../utils/thinking-display";
import type { ObservableSession, SessionObserverRegistry } from "../session-observer-registry";
import { getEditorTheme, theme } from "../theme/theme";
import { matchesSelectDown, matchesSelectUp } from "../utils/keybinding-matchers";
import { hasVisibleThinking } from "../../utils/thinking-display";
import { AssistantMessageComponent } from "./assistant-message";
import { BashExecutionComponent } from "./bash-execution";
import { BranchSummaryMessageComponent } from "./branch-summary-message";
@@ -371,9 +371,9 @@ export class AssistantMessageComponent extends Container {
}
// Add spacing only when another visible assistant content block follows.
// This avoids a superfluous blank line before separately-rendered tool execution blocks.
const hasVisibleContentAfter = message.content.slice(i + 1).some(
c => (c.type === "text" && c.text.trim()) || (c.type === "thinking" && hasVisibleThinking(c)),
);
const hasVisibleContentAfter = message.content
.slice(i + 1)
.some(c => (c.type === "text" && c.text.trim()) || (c.type === "thinking" && hasVisibleThinking(c)));
// Thinking traces in thinkingText color, italic
const md = new Markdown(thinkingText, 1, 0, getMarkdownTheme(), {
@@ -17,10 +17,10 @@ import { getSymbolTheme, theme } from "../../modes/theme/theme";
import type { InteractiveModeContext, TodoPhase } from "../../modes/types";
import type { PlanApprovalDetails } from "../../plan-mode/approved-plan";
import type { AgentSessionEvent } from "../../session/agent-session";
import { isSilentAbort, readPendingDisplayTag, resolveAbortLabel } from "../../session/messages";
import { isSilentAbort, readQueueChipText, resolveAbortLabel } from "../../session/messages";
import type { ResolveToolDetails } from "../../tools/resolve";
import { interruptHint } from "../shared";
import { hasVisibleThinking } from "../../utils/thinking-display";
import { interruptHint } from "../shared";
import { StreamingRevealController } from "./streaming-reveal";
import { ToolArgsRevealController } from "./tool-args-reveal";
@@ -38,16 +38,6 @@ const IRC_MESSAGE_VISIBLE_TTL_MS = 10_000;
*/
const MAX_LIVE_IRC_CARDS = 4;
/**
* Loader label shown the instant a user interrupt (Esc) is requested, kept until
* the agent turn fully unwinds. Esc fires the abort synchronously, but the loop
* only stops the spinner at `agent_end`, which it cannot reach until every
* in-flight tool settles its abort in `executeToolCalls` (`Promise.allSettled`).
* Swapping the steady "Working…" for this acknowledges the keypress instead of
* reading as an ignored Esc for the seconds a slow tool takes to tear down.
*/
export const INTERRUPTING_WORKING_MESSAGE = "Interrupting…";
type AgentSessionEventHandlers = {
[E in AgentSessionEventKind]: (event: Extract<AgentSessionEvent, { type: E }>) => Promise<void>;
};
@@ -64,8 +54,6 @@ export class EventController {
#renderedCustomMessages = new Set<string>();
#lastIntent: string | undefined = undefined;
#backgroundToolCallIds = new Set<string>();
#agentTurnActive = false;
#interrupting = false;
#readToolCallArgs = new Map<string, Record<string, unknown>>();
#readToolCallAssistantComponents = new Map<string, AssistantMessageComponent>();
#lastAssistantComponent: AssistantMessageComponent | undefined = undefined;
@@ -75,7 +63,6 @@ export class EventController {
// #handleMessageEnd / #handleAgentStart).
#pinnedErrorComponent: AssistantMessageComponent | undefined = undefined;
#idleCompactionTimer?: NodeJS.Timeout;
#deferredInterruptedEndTimer?: NodeJS.Timeout;
#ircExpiryTimers = new Map<string, NodeJS.Timeout>();
// Insertion-ordered IRC cards not yet retired; values are the transcript
// components each card contributed (see #retireIrcCard for the guard).
@@ -137,7 +124,6 @@ export class EventController {
this.#streamingReveal.stop();
this.#toolArgsReveal.stop();
this.#cancelIdleCompaction();
this.#cancelDeferredInterruptedEnd();
for (const timer of this.#ircExpiryTimers.values()) {
clearTimeout(timer);
}
@@ -196,7 +182,7 @@ export class EventController {
return true;
}
#updateWorkingMessageFromIntent(intent: unknown): void {
if (this.#interrupting) return;
if (this.ctx.session.isAborting) return;
// Streamed JSON can deliver non-string `_i` (object, number, boolean) before
// schema validation; `?.` only guards null/undefined, so guard the type too.
if (typeof intent !== "string") return;
@@ -206,24 +192,6 @@ export class EventController {
this.ctx.setWorkingMessage(`${trimmed}${interruptHint()}`);
}
/**
* Acknowledge a user interrupt immediately: keep/recreate the loader, switch
* it to `INTERRUPTING_WORKING_MESSAGE`, and freeze intent-driven
* working-message updates for the rest of the turn so a late
* `tool_execution_start` intent cannot repaint a "Working…/<intent>" line
* over the acknowledgment. Reset at the next `agent_start`. No-op outside an
* active turn or if already set.
*/
notifyInterrupting(): void {
if ((!this.#agentTurnActive && !this.ctx.session.isStreaming) || this.#interrupting) return;
this.#agentTurnActive = true;
this.#interrupting = true;
this.#cancelDeferredInterruptedEnd();
this.ctx.ensureLoadingAnimation();
this.ctx.setWorkingMessage(INTERRUPTING_WORKING_MESSAGE);
this.ctx.ui.requestRender();
}
subscribeToAgent(): void {
this.ctx.unsubscribe = this.ctx.session.subscribe(async (event: AgentSessionEvent) => {
await this.handleEvent(event);
@@ -241,14 +209,11 @@ export class EventController {
this.#renderedCustomMessages.clear();
this.#lastIntent = undefined;
this.#backgroundToolCallIds.clear();
this.#agentTurnActive = false;
this.#interrupting = false;
this.#readToolCallArgs.clear();
this.#readToolCallAssistantComponents.clear();
this.#lastAssistantComponent = undefined;
this.#pinnedErrorComponent = undefined;
this.#cancelIdleCompaction();
this.#cancelDeferredInterruptedEnd();
for (const timer of this.#ircExpiryTimers.values()) {
clearTimeout(timer);
}
@@ -273,9 +238,6 @@ export class EventController {
}
async #handleAgentStart(_event: Extract<AgentSessionEvent, { type: "agent_start" }>): Promise<void> {
this.#agentTurnActive = true;
this.#interrupting = false;
this.#cancelDeferredInterruptedEnd();
this.#lastIntent = undefined;
this.#readToolCallArgs.clear();
this.#readToolCallAssistantComponents.clear();
@@ -305,15 +267,10 @@ export class EventController {
this.#renderedCustomMessages.add(signature);
this.#resetReadGroup();
this.ctx.addMessageToChat(event.message);
// Tag-keyed pending-bar refresh: when AgentSession.#handleAgentEvent
// spliced this dequeued custom message out of #steeringMessages /
// #followUpMessages (it ran before this emit), the array state is
// already correct — pendingMessagesContainer just needs to be
// re-rendered to match. Gated on tag presence so non-queued customs
// (ttsr-injection, irc:*, async-result, hookMessage) skip the
// rebuild; their dispatch path never registered a pending chip.
// Mirrors the user-role refresh at the bottom of this function.
if (event.message.role === "custom" && readPendingDisplayTag(event.message.details)) {
// Queued custom-message chips are derived from the agent queue; refresh the
// pending bar when the queued custom is consumed so the chip disappears
// immediately.
if (event.message.role === "custom" && readQueueChipText(event.message.details)) {
this.ctx.updatePendingMessagesDisplay();
}
this.ctx.ui.requestRender();
@@ -833,35 +790,17 @@ export class EventController {
// this event belongs to a turn that has already been replaced. The session
// dispatches to listeners fire-and-forget across an async extension-emit hop
// (#emitSessionEvent), so an interrupted turn's agent_end can land AFTER the
// resumed turn's agent_start — e.g. empty-Enter interrupt-and-flush of a
// queued steer (interruptAndFlushQueuedMessages) or any post-turn
// agent.continue(). Running the turn-end teardown now would stop the loader
// the live turn just created, leaving "Working…" gone while the agent keeps
// running. The live turn owns the loader and finalizes it at its own
// agent_end (isStreaming === false by then). Mirrors the collab guest's
// !isStreaming loader reconciler.
// resumed turn's agent_start (e.g. any post-turn agent.continue()). Running
// the turn-end teardown now would stop the loader the live turn just created,
// leaving "Working…" gone while the agent keeps running. The live turn owns
// the loader and finalizes it at its own agent_end (isStreaming === false by
// then). Mirrors the collab guest's !isStreaming loader reconciler.
if (this.ctx.session.isStreaming) return;
// Empty-Enter interrupt-and-flush has one more race: the interrupted turn's
// final agent_end can arrive in the short gap after abort teardown but before
// the queued continuation calls agent.continue() (and flips isStreaming). Keep
// the "Interrupting…" loader alive while AgentSession says a queued resume is
// being armed; the deferred reconciler tears it down if no new stream starts.
if (
this.#interrupting &&
(this.ctx.session.isResumingQueuedMessages || this.ctx.session.queuedMessageCount > 0)
) {
this.#deferInterruptedEndTeardown();
return;
}
this.#cancelDeferredInterruptedEnd();
await this.#finishAgentEnd();
}
async #finishAgentEnd(): Promise<void> {
this.#agentTurnActive = false;
this.#interrupting = false;
this.#streamingReveal.stop();
this.#toolArgsReveal.flushAll();
if (this.ctx.loadingAnimation) {
@@ -1063,27 +1002,6 @@ export class EventController {
await this.ctx.reloadTodos();
}
#cancelDeferredInterruptedEnd(): void {
if (this.#deferredInterruptedEndTimer) {
clearTimeout(this.#deferredInterruptedEndTimer);
this.#deferredInterruptedEndTimer = undefined;
}
}
#deferInterruptedEndTeardown(): void {
if (this.#deferredInterruptedEndTimer) return;
this.#deferredInterruptedEndTimer = setTimeout(() => {
this.#deferredInterruptedEndTimer = undefined;
if (this.ctx.session.isStreaming) return;
if (this.ctx.session.isResumingQueuedMessages) {
this.#deferInterruptedEndTeardown();
return;
}
void this.#finishAgentEnd();
}, 16);
this.#deferredInterruptedEndTimer.unref?.();
}
#cancelIdleCompaction(): void {
if (this.#idleCompactionTimer) {
clearTimeout(this.#idleCompactionTimer);
@@ -152,7 +152,6 @@ export class InputController {
if (this.ctx.loopModeEnabled) {
this.ctx.pauseLoop();
if (this.ctx.session.isStreaming) {
this.ctx.notifyInterrupting();
void this.ctx.session.abort({ reason: USER_INTERRUPT_LABEL });
} else {
this.ctx.cancelPendingSubmission();
@@ -182,7 +181,6 @@ export class InputController {
// session is never streaming, so the native abort path below would
// no-op.
if (this.ctx.collabGuest.state?.isStreaming || this.ctx.loadingAnimation) {
if (!this.ctx.collabGuest.readOnly) this.ctx.notifyInterrupting();
this.ctx.collabGuest.sendAbort();
}
return;
@@ -205,7 +203,6 @@ export class InputController {
this.ctx.isPythonMode = false;
this.ctx.updateEditorBorderColor();
} else if (this.ctx.session.isStreaming) {
this.ctx.notifyInterrupting();
void this.ctx.session.abort({ reason: USER_INTERRUPT_LABEL });
} else if (this.ctx.editor.getText().trim()) {
// Esc with typed text clears the draft instead of (or before) any double-Esc action
@@ -405,31 +402,16 @@ export class InputController {
return;
}
// Empty submit while streaming with queued steering: interrupt now and
// immediately resume so the visible `Steer:` entry is sent without
// waiting for the current tool/model boundary.
// Empty submit while streaming with queued messages: abort the active
// turn and let the post-unwind drain deliver the agent-core queue.
if (!text && this.ctx.session.isStreaming) {
const queuedMessages = this.ctx.session.getQueuedMessages();
if (queuedMessages.steering.length > 0) {
this.ctx.notifyInterrupting();
try {
await this.ctx.session.interruptAndFlushQueuedMessages({ reason: USER_INTERRUPT_LABEL });
} catch (error) {
logger.warn("queued steer interrupt failed", {
error: error instanceof Error ? error.message : String(error),
});
this.ctx.showError(error instanceof Error ? error.message : String(error));
}
if (this.ctx.session.queuedMessageCount > 0) {
const aborting = this.ctx.session.abort({ reason: USER_INTERRUPT_LABEL });
await aborting;
this.ctx.updatePendingMessagesDisplay();
this.ctx.ui.requestRender();
return;
}
if (this.ctx.session.queuedMessageCount > 0) {
// Preserve the existing empty-submit flush for non-steer queues.
this.ctx.notifyInterrupting();
await this.ctx.session.abort({ reason: USER_INTERRUPT_LABEL });
return;
}
return;
}
if (!text) return;
@@ -701,16 +683,9 @@ export class InputController {
async #submitToFocusedSession(text: string, streamingBehavior: "steer" | "followUp"): Promise<void> {
const target = this.ctx.viewSession;
if (!text) {
// Mirror the empty-submit steer flush against the focused session.
if (target.isStreaming && target.getQueuedMessages().steering.length > 0) {
try {
await target.interruptAndFlushQueuedMessages({ reason: USER_INTERRUPT_LABEL });
} catch (error) {
logger.warn("focused queued steer interrupt failed", {
error: error instanceof Error ? error.message : String(error),
});
this.ctx.showError(error instanceof Error ? error.message : String(error));
}
if (target.isStreaming && target.queuedMessageCount > 0) {
const aborting = target.abort({ reason: USER_INTERRUPT_LABEL });
await aborting;
this.ctx.updatePendingMessagesDisplay();
this.ctx.ui.requestRender();
}
@@ -855,15 +830,6 @@ export class InputController {
args: args || undefined,
lineCount: body ? body.split("\n").length : 0,
};
// When the agent is streaming, register the compact slash-form text as
// the pending-display twin BEFORE dispatching the CustomMessage. The
// returned tag is embedded in details so AgentSession.#handleAgentEvent
// can remove the matching display entry when the agent consumes this
// message (mirrors the user-message dequeue path).
if (this.ctx.session.isStreaming) {
const tag = this.ctx.session.enqueueCustomMessageDisplay(text, streamingBehavior);
details.__pendingDisplayTag = tag;
}
await this.ctx.session.promptCustomMessage(
{
customType: SKILL_PROMPT_MESSAGE_TYPE,
@@ -872,7 +838,7 @@ export class InputController {
details,
attribution: "user",
},
{ streamingBehavior },
{ streamingBehavior, queueChipText: text },
);
if (this.ctx.session.isStreaming) {
this.ctx.updatePendingMessagesDisplay();
@@ -962,7 +928,7 @@ export class InputController {
if (allQueued.length === 0) {
this.ctx.updatePendingMessagesDisplay();
if (options?.abort) {
this.ctx.session.abort({ reason: USER_INTERRUPT_LABEL });
void this.ctx.session.abort({ reason: USER_INTERRUPT_LABEL });
}
return 0;
}
@@ -981,7 +947,7 @@ export class InputController {
}
this.ctx.updatePendingMessagesDisplay();
if (options?.abort) {
this.ctx.session.abort({ reason: USER_INTERRUPT_LABEL });
void this.ctx.session.abort({ reason: USER_INTERRUPT_LABEL });
}
return allQueued.length;
}
@@ -1,8 +1,8 @@
import type { AssistantMessage } from "@oh-my-pi/pi-ai";
import { getSegmenter } from "@oh-my-pi/pi-tui";
import { LRUCache } from "lru-cache/raw";
import type { AssistantMessageComponent } from "../components/assistant-message";
import { hasVisibleThinking } from "../../utils/thinking-display";
import type { AssistantMessageComponent } from "../components/assistant-message";
export const STREAMING_REVEAL_FRAME_MS = 1000 / 30;
export const MIN_STEP = 3;
@@ -3023,10 +3023,6 @@ export class InteractiveMode implements InteractiveModeContext {
this.setWorkingMessage(message);
}
notifyInterrupting(): void {
this.#eventController.notifyInterrupting();
}
showNewVersionNotification(newVersion: string): void {
this.#uiHelpers.showNewVersionNotification(newVersion);
}
-3
View File
@@ -209,9 +209,6 @@ export interface InteractiveModeContext {
flushPendingModelSwitch(): Promise<void>;
setWorkingMessage(message?: string): void;
applyPendingWorkingMessage(): void;
/** Acknowledge a user interrupt (Esc) by switching the loader to an
* "Interrupting…" label until the agent turn unwinds. */
notifyInterrupting(): void;
ensureLoadingAnimation(): void;
startPendingSubmission(input: {
text: string;
+106 -267
View File
@@ -253,7 +253,7 @@ import {
type CustomMessage,
convertToLlm,
type PythonExecutionMessage,
readPendingDisplayTag,
readQueueChipText,
SILENT_ABORT_MARKER,
SKILL_PROMPT_MESSAGE_TYPE,
stripImagesFromMessage,
@@ -862,18 +862,42 @@ function extractPermissionLocations(
// AgentSession Class
// ============================================================================
/** Internal record stored in the steering/followUp display queues. The optional
* `tag` is set only by `enqueueCustomMessageDisplay` (used for skill-prompt
* custom messages queued during streaming) and is matched by the custom-role
* `message_start` dequeue branch; user-message pushes leave it undefined and
* rely on the existing text-equality match. `images` carries the original
* (pre-normalization) image blocks so queue restoration (Esc / Alt+Up) can
* hand them back to the editor instead of dropping them. */
type QueuedDisplayEntry = { text: string; tag?: string; images?: ImageContent[] };
/** Entry returned by {@link AgentSession.clearQueue} / {@link AgentSession.popLastQueuedMessage}. */
export type RestoredQueuedMessage = { text: string; images?: ImageContent[] };
function queuedTextContent(message: AgentMessage): string | undefined {
if (!("content" in message)) return undefined;
const content = message.content;
if (typeof content === "string") return content;
return content.find((part): part is TextContent => part.type === "text")?.text;
}
function queuedImageContent(message: AgentMessage): ImageContent[] | undefined {
if (!("content" in message) || typeof message.content === "string") return undefined;
const images = message.content.filter(
(part): part is ImageContent =>
part.type === "image" && typeof part.data === "string" && typeof part.mimeType === "string",
);
return images.length > 0 ? images : undefined;
}
function isDisplayableQueuedMessage(message: AgentMessage): boolean {
return !(message.role === "custom" && message.display === false);
}
function queueChipText(message: AgentMessage): string {
if (message.role === "custom") {
return readQueueChipText(message.details) ?? queuedTextContent(message) ?? "";
}
const text = queuedTextContent(message) ?? "";
if (text) return text;
return queuedImageContent(message) ? "[Image]" : "";
}
function toRestoredQueuedMessage(message: AgentMessage): RestoredQueuedMessage {
return { text: queueChipText(message), images: queuedImageContent(message) };
}
export class AgentSession {
readonly agent: Agent;
readonly sessionManager: SessionManager;
@@ -904,15 +928,6 @@ export class AgentSession {
#eventListeners: AgentSessionEventListener[] = [];
#commandMetadataChangedListeners: CommandMetadataChangedListener[] = [];
/** Tracks pending steering messages for UI display. Removed when delivered.
* Entry shape: `{ text }` for plain-text steers (user-message dequeue
* matches by `.text`); `{ text, tag }` for queued custom messages (skill
* invocations dispatched while streaming) — the custom-role dequeue
* matches by `.tag` so duplicate-args queued skills cannot collide. */
#steeringMessages: QueuedDisplayEntry[] = [];
/** Tracks pending follow-up messages for UI display. Removed when delivered.
* See `#steeringMessages` for entry shape. */
#followUpMessages: QueuedDisplayEntry[] = [];
/** Messages queued to be included with the next user prompt as context ("asides"). */
#pendingNextTurnMessages: CustomMessage[] = [];
#scheduledHiddenNextTurnGeneration: number | undefined = undefined;
@@ -1058,10 +1073,6 @@ export class AgentSession {
* without producing an aborted message_end). */
#planCompactAbortPending = false;
/** Monotonic counter for `enqueueCustomMessageDisplay` tag generation;
* combined with `Date.now()` so tags stay unique even across rapid
* same-tick enqueues. */
#customDisplayTagCounter = 0;
#postPromptTasks = new Set<Promise<unknown>>();
#postPromptTasksPromise: Promise<void> | undefined = undefined;
#postPromptTasksResolve: (() => void) | undefined = undefined;
@@ -1083,8 +1094,6 @@ export class AgentSession {
// has decremented #promptInFlightCount, hitting AgentBusyError. Flushed from
// both #endInFlight (normal) and #resetInFlight (abort).
#pendingAgentEndEmit: AgentSessionEvent | undefined;
#resumingQueuedMessages = false;
#queuedFlushInterrupt: Promise<void> | undefined;
#obfuscator: SecretObfuscator | undefined;
#checkpointState: CheckpointState | undefined = undefined;
#pendingRewindReport: string | undefined = undefined;
@@ -1473,28 +1482,6 @@ export class AgentSession {
this.#planCompactAbortPending = false;
}
/** Register a compact display string for a custom message that the caller is
* about to dispatch via `promptCustomMessage` / `sendCustomMessage`.
* Returns a stable tag the caller MUST embed in
* `CustomMessage.details.__pendingDisplayTag` so the agent-side
* `message_start` handler can remove the matching display entry when the
* queued message is consumed.
*
* Does NOT push to the agent's steering/followUp queue — that happens
* separately inside `sendCustomMessage`. */
enqueueCustomMessageDisplay(text: string, mode: "steer" | "followUp"): string {
const tag = `omp-cmd-${Date.now()}-${++this.#customDisplayTagCounter}`;
const displayText = text.trim();
if (!displayText) return tag;
const entry: QueuedDisplayEntry = { text: displayText, tag };
if (mode === "steer") {
this.#steeringMessages.push(entry);
} else {
this.#followUpMessages.push(entry);
}
return tag;
}
getAsyncJobSnapshot(options?: { recentLimit?: number }): AsyncJobSnapshot | null {
const manager = this.#asyncJobManager;
if (!manager) return null;
@@ -1616,45 +1603,6 @@ export class AgentSession {
/** Internal handler for agent events - shared by subscribe and reconnect */
#handleAgentEvent = async (event: AgentEvent): Promise<void> => {
// When a user message starts, check if it's from either queue and remove it BEFORE emitting
// This ensures the UI sees the updated queue state
if (event.type === "message_start" && event.message.role === "user") {
const messageText = this.#getUserMessageText(event.message);
if (messageText) {
// Check steering queue first (match by .text on tagged records)
const steeringIndex = this.#steeringMessages.findIndex(e => e.text === messageText);
if (steeringIndex !== -1) {
this.#steeringMessages.splice(steeringIndex, 1);
} else {
// Check follow-up queue
const followUpIndex = this.#followUpMessages.findIndex(e => e.text === messageText);
if (followUpIndex !== -1) {
this.#followUpMessages.splice(followUpIndex, 1);
}
}
}
}
// Tag-based dequeue for custom messages (skills queued via promptCustomMessage).
// The InputController attached a stable tag via CustomMessage.details when it
// registered the display chip; pull it back here to remove the matching entry
// from the pending bar atomically with the agent's queue consumption. Match by
// tag (not text) — two queued skills with identical args cannot collide.
if (event.type === "message_start" && event.message.role === "custom") {
const tag = readPendingDisplayTag(event.message.details);
if (tag) {
const steerIdx = this.#steeringMessages.findIndex(e => e.tag === tag);
if (steerIdx !== -1) {
this.#steeringMessages.splice(steerIdx, 1);
} else {
const followUpIdx = this.#followUpMessages.findIndex(e => e.tag === tag);
if (followUpIdx !== -1) {
this.#followUpMessages.splice(followUpIdx, 1);
}
}
}
}
// Plan-mode → compaction transition: stamp `SILENT_ABORT_MARKER` on the
// persisted message BEFORE the obfuscator's display-side copy below.
// Invariant (must hold across refactors): this branch precedes the
@@ -2622,18 +2570,6 @@ export class AgentSession {
return Array.from(candidates);
}
/** Extract text content from a message */
#getUserMessageText(message: Message): string {
if (message.role !== "user") return "";
const content = message.content;
if (typeof content === "string") return content;
const textBlocks = content.filter(c => c.type === "text");
const text = textBlocks.map(c => (c as TextContent).text).join("");
if (text.length > 0) return text;
const hasImages = content.some(c => c.type === "image");
return hasImages ? "[Image]" : "";
}
/** Find the last assistant message in agent state (including aborted ones) */
#findLastAssistantMessage(): AssistantMessage | undefined {
const messages = this.agent.state.messages;
@@ -3342,9 +3278,8 @@ export class AgentSession {
return this.agent.state.isStreaming || this.#promptInFlightCount > 0;
}
/** True while an explicit interrupt-and-flush is between aborting the old turn and starting the queued continuation. */
get isResumingQueuedMessages(): boolean {
return this.#resumingQueuedMessages;
get isAborting(): boolean {
return this.agent.isAborting;
}
/** Wait until streaming and deferred recovery work are fully settled. */
@@ -4657,9 +4592,9 @@ export class AgentSession {
throw new AgentBusyError();
}
if (options.streamingBehavior === "followUp") {
await this.#queueFollowUp(expandedText, options?.images);
await this.#queueUserMessage(expandedText, options?.images, "followUp");
} else {
await this.#queueSteer(expandedText, options?.images);
await this.#queueUserMessage(expandedText, options?.images, "steer");
}
// Steer/follow-up the keyword notices alongside the queued user message.
for (const notice of keywordNotices) {
@@ -4710,7 +4645,7 @@ export class AgentSession {
async promptCustomMessage<T = unknown>(
message: Pick<CustomMessage<T>, "customType" | "content" | "display" | "details" | "attribution">,
options?: Pick<PromptOptions, "streamingBehavior" | "toolChoice">,
options?: Pick<PromptOptions, "streamingBehavior" | "toolChoice"> & { queueChipText?: string },
): Promise<void> {
const textContent =
typeof message.content === "string"
@@ -4734,7 +4669,10 @@ export class AgentSession {
if (!options?.streamingBehavior) {
throw new AgentBusyError();
}
await this.sendCustomMessage(message, { deliverAs: options.streamingBehavior });
await this.sendCustomMessage(message, {
deliverAs: options.streamingBehavior,
queueChipText: options.queueChipText,
});
for (const notice of keywordNotices) {
await this.sendCustomMessage(notice, { deliverAs: options.streamingBehavior });
}
@@ -5071,7 +5009,7 @@ export class AgentSession {
}
const expandedText = expandPromptTemplate(text, [...this.#promptTemplates]);
await this.#queueSteer(expandedText, images);
await this.#queueUserMessage(expandedText, images, "steer");
}
/**
@@ -5083,73 +5021,47 @@ export class AgentSession {
}
const expandedText = expandPromptTemplate(text, [...this.#promptTemplates]);
await this.#queueFollowUp(expandedText, images);
await this.#queueUserMessage(expandedText, images, "followUp");
}
/**
* Internal: Queue a steering message (already expanded, no extension command check).
*/
async #queueSteer(text: string, images?: ImageContent[]): Promise<void> {
async #queueUserMessage(
text: string,
images: ImageContent[] | undefined,
mode: "steer" | "followUp",
): Promise<void> {
const normalizedImages = await this.#normalizeImagesForModel(images);
const displayText = text || (images && images.length > 0 ? "[Image]" : "");
this.#steeringMessages.push({ text: displayText, images });
const content: (TextContent | ImageContent)[] = [{ type: "text", text }];
if (normalizedImages && normalizedImages.length > 0) {
if (normalizedImages?.length) {
content.push(...normalizedImages);
}
this.agent.steer({
role: "user",
content,
steering: true,
attribution: "user",
timestamp: Date.now(),
});
// A steer can land on an idle session: the caller checked isStreaming
// before the (potentially slow) image normalization above, so the turn
// may have ended in between. Without a drain the message would strand in
// the queue until the next manual prompt — schedule an immediate continue,
// mirroring #queueFollowUp's idle-path delivery.
if (this.#canAutoContinueForFollowUp()) {
this.#scheduleAgentContinue({
shouldContinue: () => this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages(),
if (mode === "followUp") {
this.agent.followUp({
role: "user",
content,
attribution: "user",
timestamp: Date.now(),
});
} else {
this.agent.steer({
role: "user",
content,
steering: true,
attribution: "user",
timestamp: Date.now(),
});
}
this.#scheduleIdleQueueDrain();
}
#scheduleIdleQueueDrain(): void {
if (!this.#canAutoContinueForFollowUp()) return;
this.#scheduleAgentContinue({
shouldContinue: () => this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages(),
});
}
/**
* Internal: Queue a follow-up message (already expanded, no extension command check).
*/
async #queueFollowUp(text: string, images?: ImageContent[]): Promise<void> {
const normalizedImages = await this.#normalizeImagesForModel(images);
const displayText = text || (images && images.length > 0 ? "[Image]" : "");
this.#followUpMessages.push({ text: displayText, images });
const content: (TextContent | ImageContent)[] = [{ type: "text", text }];
if (normalizedImages && normalizedImages.length > 0) {
content.push(...normalizedImages);
}
this.agent.followUp({
role: "user",
content,
attribution: "user",
timestamp: Date.now(),
});
// When fully idle AND the session is in a resumable assistant-ended state,
// schedule an immediate continue so the queued follow-up is delivered
// without waiting for the next user turn. We gate on isStreaming (model
// actively producing), isRetrying (auto-retry backoff is sleeping between
// attempts, #retryPromise set), and the last message being assistant —
// agent.continue() only dequeues follow-ups from an assistant-ended state;
// resuming from user/toolResult state runs an extra model call on the
// stale prompt before draining the queue.
if (this.#canAutoContinueForFollowUp()) {
this.#scheduleAgentContinue({
shouldContinue: () => this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages(),
});
}
}
/**
* Gate for idle-path follow-up auto-continue. See `#queueFollowUp` for rationale.
* Gate for idle-path queued-message auto-continue. See `#scheduleIdleQueueDrain` for rationale.
*/
#canAutoContinueForFollowUp(): boolean {
if (this.isStreaming) return false;
@@ -5258,14 +5170,24 @@ export class AgentSession {
*/
async sendCustomMessage<T = unknown>(
message: Pick<CustomMessage<T>, "customType" | "content" | "display" | "details" | "attribution">,
options?: { triggerTurn?: boolean; deliverAs?: "steer" | "followUp" | "nextTurn" },
options?: { triggerTurn?: boolean; deliverAs?: "steer" | "followUp" | "nextTurn"; queueChipText?: string },
): Promise<void> {
const details =
options?.queueChipText && options.deliverAs !== "nextTurn"
? ({
...((message.details && typeof message.details === "object" ? message.details : {}) as Record<
string,
unknown
>),
__queueChipText: options.queueChipText,
} as T)
: message.details;
const appMessage: CustomMessage<T> = {
role: "custom",
customType: message.customType,
content: message.content,
display: message.display,
details: message.details,
details,
attribution: message.attribution ?? "agent",
timestamp: Date.now(),
};
@@ -5281,14 +5203,7 @@ export class AgentSession {
} else {
this.agent.steer(normalizedAppMessage);
}
// The isStreaming check above can be stale: image normalization is
// awaited, so the turn may have ended in between, leaving the message
// queued on an idle agent. Mirror #queueSteer's idle-path delivery.
if (this.#canAutoContinueForFollowUp()) {
this.#scheduleAgentContinue({
shouldContinue: () => this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages(),
});
}
this.#scheduleIdleQueueDrain();
return;
}
@@ -5363,11 +5278,11 @@ export class AgentSession {
}
if (options?.deliverAs === "followUp") {
await this.#queueFollowUp(text, images);
await this.#queueUserMessage(text, images, "followUp");
return;
}
if (options?.deliverAs === "steer") {
await this.#queueSteer(text, images);
await this.#queueUserMessage(text, images, "steer");
return;
}
@@ -5378,57 +5293,37 @@ export class AgentSession {
});
}
/**
* Clear queued messages and return them (text plus any attached images).
* Useful for restoring to editor when user aborts. The internal entry
* arrays are handed out as-is — a `tag` (if any) is inert once the record
* leaves the queue.
*/
/** Clear queued messages and return them (text plus any attached images). */
clearQueue(): { steering: RestoredQueuedMessage[]; followUp: RestoredQueuedMessage[] } {
const steering = this.#steeringMessages;
const followUp = this.#followUpMessages;
this.#steeringMessages = [];
this.#followUpMessages = [];
const steering = this.agent.peekSteeringQueue().map(toRestoredQueuedMessage);
const followUp = this.agent.peekFollowUpQueue().map(toRestoredQueuedMessage);
this.agent.clearAllQueues();
return { steering, followUp };
}
/** Number of pending messages (includes steering, follow-up, and next-turn messages) */
/** Number of pending displayable messages (includes steering, follow-up, and next-turn messages) */
get queuedMessageCount(): number {
return this.#steeringMessages.length + this.#followUpMessages.length + this.#pendingNextTurnMessages.length;
return (
this.agent.peekSteeringQueue().filter(isDisplayableQueuedMessage).length +
this.agent.peekFollowUpQueue().filter(isDisplayableQueuedMessage).length +
this.#pendingNextTurnMessages.length
);
}
/** Get pending messages (read-only). Returns the public text-only view;
* internal `{text, tag?}` records are mapped to `.text` so callers
* (`updatePendingMessagesDisplay`, `restoreQueuedMessagesToEditor`) see
* the unchanged historical shape. */
getQueuedMessages(): { steering: readonly string[]; followUp: readonly string[] } {
return {
steering: this.#steeringMessages.map(e => e.text),
followUp: this.#followUpMessages.map(e => e.text),
steering: this.agent.peekSteeringQueue().filter(isDisplayableQueuedMessage).map(queueChipText),
followUp: this.agent.peekFollowUpQueue().filter(isDisplayableQueuedMessage).map(queueChipText),
};
}
/**
* Pop the last queued message (steering first, then follow-up).
* Used by dequeue keybinding to restore messages to editor one at a time.
* Returns the popped entry's text and images; the tag (if any) dies with
* the record — no orphan state can outlive the queue entry.
*/
popLastQueuedMessage(): RestoredQueuedMessage | undefined {
// Pop from steering first (LIFO)
if (this.#steeringMessages.length > 0) {
const entry = this.#steeringMessages.pop();
this.agent.popLastSteer();
return entry;
}
// Then from follow-up
if (this.#followUpMessages.length > 0) {
const entry = this.#followUpMessages.pop();
this.agent.popLastFollowUp();
return entry;
}
return undefined;
const message = this.agent.popLastSteer() ?? this.agent.popLastFollowUp();
return message ? toRestoredQueuedMessage(message) : undefined;
}
get skillsSettings(): SkillsSettings | undefined {
@@ -5510,56 +5405,6 @@ export class AgentSession {
}
}
/**
* Abort active work, then immediately resume the agent so queued steer/follow-up
* messages drain instead of waiting for another natural turn boundary.
*
* The drained queue is re-run via `agent.prompt()`, which appends and runs
* regardless of the trailing message role. Earlier this called
* `agent.continue()`, which only dequeues steers after an assistant message
* and throws "No messages to continue from" on an empty context — both states
* a flushed empty-Enter can land in, which stranded the steer in the queue and
* surfaced that error in the TUI.
*
* Concurrent calls (e.g. a double empty-Enter) coalesce: the first owns the
* abort→resume handoff while the rest await it and return. The guard clears
* once the resumed turn has started (not at turn end), so a genuinely-later
* flush of a freshly queued steer still works.
*/
async interruptAndFlushQueuedMessages(options?: { reason?: string }): Promise<void> {
const inFlight = this.#queuedFlushInterrupt;
if (inFlight) {
await inFlight;
return;
}
if (!this.agent.hasQueuedMessages()) return;
const { promise: handoff, resolve: settleHandoff } = Promise.withResolvers<void>();
this.#queuedFlushInterrupt = handoff;
this.#resumingQueuedMessages = true;
let turn: Promise<void> | undefined;
try {
await this.abort({ reason: options?.reason });
if (this.isCompacting || this.isGeneratingHandoff) return;
await this.#maybeRestoreRetryFallbackPrimary();
// A turn slipped in while we were settling (e.g. a fresh user prompt):
// leave the queue for its next steering boundary rather than draining
// and double-running, which would throw AgentBusyError. Checking
// isStreaming and prompting below is synchronous, so no turn can start
// between the drain and the resume.
if (this.agent.state.isStreaming) return;
const queued = this.agent.takeQueuedMessages();
if (queued.length === 0) return;
this.#resumingQueuedMessages = false;
turn = this.agent.prompt(queued);
} finally {
this.#resumingQueuedMessages = false;
this.#queuedFlushInterrupt = undefined;
settleHandoff();
}
await turn;
}
/**
* Start a new session, optionally with initial messages and parent tracking.
* Clears all messages and starts a new session.
@@ -5610,8 +5455,6 @@ export class AgentSession {
this.#rekeyMnemopiMemoryForCurrentSessionId();
this.#resetHindsightConversationTrackingIfHindsight();
this.#resetMnemopiConversationTrackingIfMnemopi();
this.#steeringMessages = [];
this.#followUpMessages = [];
this.#pendingNextTurnMessages = [];
this.#scheduledHiddenNextTurnGeneration = undefined;
@@ -6756,8 +6599,6 @@ export class AgentSession {
this.#rekeyMnemopiMemoryForCurrentSessionId();
this.#resetHindsightConversationTrackingIfHindsight();
this.#resetMnemopiConversationTrackingIfMnemopi();
this.#steeringMessages = [];
this.#followUpMessages = [];
this.#pendingNextTurnMessages = [];
this.#scheduledHiddenNextTurnGeneration = undefined;
this.#todoReminderCount = 0;
@@ -9655,8 +9496,8 @@ export class AgentSession {
// the existing message objects is sufficient and avoids structured-clone failures for
// extension/custom metadata that is valid to persist but not cloneable.
const previousAgentMessages = [...this.agent.state.messages];
const previousSteeringMessages = [...this.#steeringMessages];
const previousFollowUpMessages = [...this.#followUpMessages];
const previousSteeringMessages = [...this.agent.peekSteeringQueue()];
const previousFollowUpMessages = [...this.agent.peekFollowUpQueue()];
const previousPendingNextTurnMessages = [...this.#pendingNextTurnMessages];
const previousScheduledHiddenNextTurnGeneration = this.#scheduledHiddenNextTurnGeneration;
const previousModel = this.model;
@@ -9673,8 +9514,7 @@ export class AgentSession {
? this.#getSessionDefaultSelectedMCPToolNames(previousSessionFile)
: undefined;
this.#steeringMessages = [];
this.#followUpMessages = [];
this.agent.clearAllQueues();
this.#pendingNextTurnMessages = [];
this.#scheduledHiddenNextTurnGeneration = undefined;
@@ -9815,8 +9655,7 @@ export class AgentSession {
this.#baseSystemPrompt = previousBaseSystemPrompt;
this.agent.setSystemPrompt(previousSystemPrompt);
this.agent.replaceMessages(previousAgentMessages);
this.#steeringMessages = previousSteeringMessages;
this.#followUpMessages = previousFollowUpMessages;
this.agent.replaceQueues(previousSteeringMessages, previousFollowUpMessages);
this.#pendingNextTurnMessages = previousPendingNextTurnMessages;
this.#scheduledHiddenNextTurnGeneration = previousScheduledHiddenNextTurnGeneration;
if (previousModel) {
+9 -11
View File
@@ -40,13 +40,11 @@ export interface SkillPromptDetails {
path: string;
args?: string;
lineCount: number;
/** Internal: tag used by AgentSession to remove the pending-display chip
* from `#steeringMessages` / `#followUpMessages` when the agent consumes
* this message. Not surfaced to renderers; the `__` prefix signals
* "private". Optional — non-streaming skill prompts never set it. Stripped
* from persisted `details` by `SessionManager.appendCustomMessageEntry`
* via the `INTERNAL_DETAILS_FIELDS` allowlist below. */
__pendingDisplayTag?: string;
/** Internal: compact label shown for a queued custom message. Optional —
* non-streaming skill prompts never set it. Stripped from persisted
* `details` by `SessionManager.appendCustomMessageEntry` via the
* `INTERNAL_DETAILS_FIELDS` allowlist below. */
__queueChipText?: string;
}
/** Sentinel value for `AssistantMessage.errorMessage` indicating that the abort
@@ -104,12 +102,12 @@ export function resolveAbortLabel(errorMessage: string | undefined, retryAttempt
return "Operation aborted";
}
/** Extract the optional `__pendingDisplayTag` field from a CustomMessage's
/** Extract the optional `__queueChipText` field from a CustomMessage's
* `details` blob. Safe over `unknown`; returns undefined when the field is
* absent or non-string. */
export function readPendingDisplayTag(details: unknown): string | undefined {
export function readQueueChipText(details: unknown): string | undefined {
if (typeof details !== "object" || details === null) return undefined;
const candidate = (details as { __pendingDisplayTag?: unknown }).__pendingDisplayTag;
const candidate = (details as { __queueChipText?: unknown }).__queueChipText;
return typeof candidate === "string" ? candidate : undefined;
}
@@ -118,7 +116,7 @@ export function readPendingDisplayTag(details: unknown): string | undefined {
* the CustomMessageEntry to disk. Scoped intentionally narrow: only fields
* declared here are stripped. Adding a new entry is a deliberate, reviewed
* change — unrelated future payload fields are never silently dropped. */
export const INTERNAL_DETAILS_FIELDS = ["__pendingDisplayTag"] as const;
export const INTERNAL_DETAILS_FIELDS = ["__queueChipText"] as const;
/** Return a `details` copy with every key in `INTERNAL_DETAILS_FIELDS`
* removed. Returns the input unchanged when there is nothing to strip
@@ -5,6 +5,7 @@ import type { AgentMessage, ThinkingLevel } from "@oh-my-pi/pi-agent-core";
import { INTENT_FIELD } from "@oh-my-pi/pi-agent-core";
import type { AssistantMessage, Model } from "@oh-my-pi/pi-ai";
import { isZodSchema, zodToWireSchema } from "@oh-my-pi/pi-ai/utils/schema";
import { getVisibleThinkingText } from "../utils/thinking-display";
import {
type BashExecutionMessage,
type BranchSummaryMessage,
@@ -16,7 +17,6 @@ import {
type PythonExecutionMessage,
pythonExecutionToText,
} from "./messages";
import { getVisibleThinkingText } from "../utils/thinking-display";
/** Minimal tool shape for dump output (matches AgentTool fields used by formatSessionDumpText). */
export interface SessionDumpToolInfo {
@@ -2623,11 +2623,7 @@ export class SessionManager {
// Hot path: writer is open and all entries have been written via writeSync.
// Just fsync the fd — the data is already in the kernel page cache.
if (
this.#persistWriter?.isOpen() &&
this.#flushed &&
!this.#needsFullRewriteOnNextPersist
) {
if (this.#persistWriter?.isOpen() && this.#flushed && !this.#needsFullRewriteOnNextPersist) {
this.#persistWriter.fsyncSync();
return;
}
@@ -154,338 +154,6 @@ describe("AgentSession concurrent prompt guard", () => {
await firstPrompt.catch(() => {});
});
it("interrupts active work and immediately sends queued steering messages", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const callMessages: Message[][] = [];
const agent = new Agent({
getApiKey: () => "test-key",
initialState: {
model,
systemPrompt: ["Test"],
tools: [],
},
convertToLlm,
streamFn: (_model, context, options) => {
const callIndex = callMessages.length;
callMessages.push([...context.messages]);
const stream = new AssistantMessageEventStream();
queueMicrotask(() => {
stream.push({ type: "start", partial: createAssistantMessage("") });
if (callIndex > 0) {
stream.push({ type: "done", reason: "stop", message: createAssistantMessage("Handled steer") });
}
});
options?.signal?.addEventListener(
"abort",
() => {
stream.push({
type: "error",
reason: "aborted",
error: createAssistantMessage("Interrupted"),
});
},
{ once: true },
);
return stream;
},
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-interrupt-flush.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-interrupt-flush.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({
agent,
sessionManager,
settings,
modelRegistry,
});
const firstPrompt = session.prompt("First message").catch(() => {});
await waitFor(() => session.isStreaming && callMessages.length === 1);
await session.steer("Send this now");
expect(session.getQueuedMessages().steering).toEqual(["Send this now"]);
await session.interruptAndFlushQueuedMessages({ reason: "Interrupted by user" });
await firstPrompt;
expect(callMessages).toHaveLength(2);
expect(
callMessages[1]?.some(message => {
if (typeof message.content === "string") {
return message.content.includes("Send this now");
}
return message.content.some(content => content.type === "text" && content.text.includes("Send this now"));
}),
).toBe(true);
expect(session.getQueuedMessages()).toEqual({ steering: [], followUp: [] });
});
it("coalesces repeated interrupt-and-flush requests for one queued steer", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const callMessages: Message[][] = [];
const agent = new Agent({
getApiKey: () => "test-key",
initialState: {
model,
systemPrompt: ["Test"],
tools: [],
},
convertToLlm,
streamFn: (_model, context, options) => {
const callIndex = callMessages.length;
callMessages.push([...context.messages]);
const stream = new AssistantMessageEventStream();
queueMicrotask(() => {
stream.push({ type: "start", partial: createAssistantMessage("") });
if (callIndex > 0) {
stream.push({ type: "done", reason: "stop", message: createAssistantMessage("Handled steer") });
}
});
options?.signal?.addEventListener(
"abort",
() => {
stream.push({
type: "error",
reason: "aborted",
error: createAssistantMessage("Interrupted"),
});
},
{ once: true },
);
return stream;
},
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-interrupt-flush-repeat.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-interrupt-flush-repeat.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({
agent,
sessionManager,
settings,
modelRegistry,
});
const firstPrompt = session.prompt("First message").catch(() => {});
await waitFor(() => session.isStreaming && callMessages.length === 1);
await session.steer("Send this once");
expect(session.getQueuedMessages().steering).toEqual(["Send this once"]);
const flushes = Array.from({ length: 6 }, () =>
session.interruptAndFlushQueuedMessages({ reason: "Interrupted by user" }),
);
await expect(Promise.all(flushes)).resolves.toEqual([
undefined,
undefined,
undefined,
undefined,
undefined,
undefined,
]);
await firstPrompt;
expect(callMessages).toHaveLength(2);
const resumedCall = callMessages[1];
expect(
resumedCall?.some(message => {
if (typeof message.content === "string") {
return message.content.includes("Send this once");
}
return message.content.some(content => content.type === "text" && content.text.includes("Send this once"));
}),
).toBe(true);
expect(session.getQueuedMessages()).toEqual({ steering: [], followUp: [] });
});
it("does not resume when the queued steer is cleared during interrupt-and-flush", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const callMessages: Message[][] = [];
const agent = new Agent({
getApiKey: () => "test-key",
initialState: {
model,
systemPrompt: ["Test"],
tools: [],
},
convertToLlm,
streamFn: (_model, context, options) => {
callMessages.push([...context.messages]);
const stream = new AssistantMessageEventStream();
queueMicrotask(() => {
stream.push({ type: "start", partial: createAssistantMessage("") });
});
options?.signal?.addEventListener(
"abort",
() => {
stream.push({
type: "error",
reason: "aborted",
error: createAssistantMessage("Interrupted"),
});
},
{ once: true },
);
return stream;
},
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-interrupt-flush-clear.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-interrupt-flush-clear.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({
agent,
sessionManager,
settings,
modelRegistry,
});
const firstPrompt = session.prompt("First message").catch(() => {});
await waitFor(() => session.isStreaming && callMessages.length === 1);
await session.steer("Restore this to the editor instead");
const abort = session.abort.bind(session);
vi.spyOn(session, "abort").mockImplementation(async options => {
await abort(options);
session.clearQueue();
});
await expect(session.interruptAndFlushQueuedMessages({ reason: "Interrupted by user" })).resolves.toBeUndefined();
await firstPrompt;
expect(callMessages).toHaveLength(1);
expect(session.getQueuedMessages()).toEqual({ steering: [], followUp: [] });
});
it("delivers queued steering after interrupting mid-tool execution (queue survives external abort)", async () => {
// Regression: pressing Enter with a queued steer while a tool was running
// aborted the run, but the post-abort steering poll inside executeToolCalls
// drained the queue into the dying run — the message landed in history and
// interruptAndFlushQueuedMessages saw an empty queue, so it never resumed.
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const callMessages: Message[][] = [];
let toolStarted = false;
const blockingTool: AgentTool = {
name: "mock_blocker",
label: "Mock Blocker",
description: "Blocks until aborted",
parameters: z.object({}),
execute: async (_id, _args, signal) => {
toolStarted = true;
const { promise, resolve } = Promise.withResolvers<void>();
if (signal?.aborted) resolve();
else signal?.addEventListener("abort", () => resolve(), { once: true });
await promise;
return { content: [{ type: "text" as const, text: "tool aborted" }] };
},
};
const toolCallContent: ToolCall = {
type: "toolCall",
id: "call_steer_flush_001",
name: "mock_blocker",
arguments: {},
};
const agent = new Agent({
getApiKey: () => "test-key",
initialState: { model, systemPrompt: ["Test"], tools: [blockingTool] },
convertToLlm,
streamFn: (_model, context, options) => {
const callIndex = callMessages.length;
callMessages.push([...context.messages]);
const stream = new AssistantMessageEventStream();
const signal = options?.signal;
queueMicrotask(() => {
if (signal?.aborted) {
// Post-abort model call inside the dying run.
stream.push({ type: "error", reason: "aborted", error: createAssistantMessage("Interrupted") });
return;
}
if (callIndex === 0) {
const partial: AssistantMessage = {
...createAssistantMessage(""),
content: [toolCallContent],
stopReason: "toolUse",
};
stream.push({ type: "start", partial });
stream.push({ type: "toolcall_start", contentIndex: 0, partial });
stream.push({ type: "toolcall_end", contentIndex: 0, toolCall: toolCallContent, partial });
stream.push({ type: "done", reason: "toolUse", message: partial });
return;
}
const done = createAssistantMessage("Handled steer");
stream.push({ type: "start", partial: done });
stream.push({ type: "done", reason: "stop", message: done });
signal?.addEventListener(
"abort",
() => {
stream.push({ type: "error", reason: "aborted", error: createAssistantMessage("Interrupted") });
},
{ once: true },
);
});
return stream;
},
});
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-steer-tool-abort.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-steer-tool-abort.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
session = new AgentSession({
agent,
sessionManager,
settings,
modelRegistry,
});
const firstPrompt = session.prompt("First message").catch(() => {});
await waitFor(() => toolStarted);
await session.steer("Send this now");
expect(session.getQueuedMessages().steering).toEqual(["Send this now"]);
await session.interruptAndFlushQueuedMessages({ reason: "Interrupted by user" });
await firstPrompt;
// The resumed run's model call must carry the steer message.
const lastCall = callMessages[callMessages.length - 1];
expect(
lastCall?.some(message => {
if (typeof message.content === "string") return message.content.includes("Send this now");
return message.content.some(content => content.type === "text" && content.text.includes("Send this now"));
}),
).toBe(true);
// The agent actually resumed and produced a response after the steer.
const lastAssistant = [...agent.state.messages]
.reverse()
.find((m): m is AssistantMessage => m.role === "assistant");
expect(lastAssistant?.content).toEqual([{ type: "text", text: "Handled steer" }]);
expect(session.getQueuedMessages()).toEqual({ steering: [], followUp: [] });
});
it("should allow followUp() while streaming", async () => {
await createSession();
@@ -1,9 +1,8 @@
/**
* Contract: a guest prompt that arrives while the host agent is streaming is
* steered AND becomes visible as a queued message — the host registers the
* pending-display twin (so `queuedMessageCount` covers it, feeding the web
* composer's "queued ×N" badge and the TUI pending bar) and tags the custom
* message so the entry is dequeued when the agent consumes it.
* steered AND becomes visible as a queued message. The host sends the guest text
* as `queueChipText` on the queued custom message; the session derives
* `queuedMessageCount` from the agent-core queue for host and guest UI state.
*/
import { afterEach, describe, expect, it } from "bun:test";
import { importRoomKey } from "@oh-my-pi/pi-coding-agent/collab/crypto";
@@ -76,20 +75,19 @@ function startTestRelay(): { url: string; stop(): void } {
}
interface CapturedPrompt {
details?: { from?: string; __pendingDisplayTag?: string };
details?: { from?: string };
options?: { streamingBehavior?: "steer"; queueChipText?: string };
}
interface StreamingHostHarness {
ctx: InteractiveModeContext;
prompts: CapturedPrompt[];
enqueued: { text: string; mode: string; tag: string }[];
nextPrompt(): Promise<CapturedPrompt>;
}
/** Context double for a host whose agent is mid-turn (isStreaming === true). */
function makeStreamingHostContext(): StreamingHostHarness {
const prompts: CapturedPrompt[] = [];
const enqueued: { text: string; mode: string; tag: string }[] = [];
const promptWaiters: ((prompt: CapturedPrompt) => void)[] = [];
const ctx = {
settings: { get: () => "" },
@@ -105,20 +103,16 @@ function makeStreamingHostContext(): StreamingHostHarness {
session: {
isStreaming: true,
get queuedMessageCount(): number {
return enqueued.length;
return prompts.filter(prompt => prompt.options?.queueChipText).length;
},
isAborting: false,
sessionName: "test",
model: undefined,
thinkingLevel: undefined,
subscribe: () => () => {},
emitNotice: () => {},
enqueueCustomMessageDisplay: (text: string, mode: string) => {
const tag = `tag-${enqueued.length + 1}`;
enqueued.push({ text, mode, tag });
return tag;
},
promptCustomMessage: (message: CapturedPrompt) => {
const captured: CapturedPrompt = { details: message.details };
promptCustomMessage: (message: CapturedPrompt, options?: CapturedPrompt["options"]) => {
const captured: CapturedPrompt = { details: message.details, options };
prompts.push(captured);
for (const waiter of promptWaiters.splice(0)) waiter(captured);
return Promise.resolve();
@@ -140,7 +134,7 @@ function makeStreamingHostContext(): StreamingHostHarness {
promptWaiters.push(resolve);
return promise;
};
return { ctx, prompts, enqueued, nextPrompt };
return { ctx, prompts, nextPrompt };
}
interface TestGuest {
@@ -197,10 +191,8 @@ describe("collab mid-turn guest prompts", () => {
guest.socket.send({ t: "prompt", text: "steer the host" });
const prompt = await prompted;
// Display twin registered before dispatch, tag forwarded in details so
// the session dequeues the entry when the agent consumes the message.
expect(harness.enqueued).toEqual([{ text: "steer the host", mode: "steer", tag: "tag-1" }]);
expect(prompt.details).toEqual({ from: "writer", __pendingDisplayTag: "tag-1" });
expect(prompt.details).toEqual({ from: "writer" });
expect(prompt.options).toEqual({ streamingBehavior: "steer", queueChipText: "steer the host" });
// The queued steer must reach guests through state.queuedMessageCount —
// that field drives the web composer's "queued ×N" badge.
@@ -57,20 +57,21 @@ function createContext(): {
abortEval: Spy;
addMessageToChat: Spy;
cancelPendingSubmission: Spy;
clearEditor: Spy;
clearQueue: Spy;
flushSync: Spy;
getQueuedMessages: Spy;
ensureLoadingAnimation: Spy;
handleBtwCommand: Spy;
interruptAndFlushQueuedMessages: Spy;
handleBtwEscape: Spy;
handleOmfgEscape: Spy;
hasActiveBtw: Spy;
hasActiveOmfg: Spy;
notifyInterrupting: Spy;
onInputCallback: Spy;
prompt: Spy;
requestRender: Spy;
resetDisplay: Spy;
shutdown: Spy;
startPendingSubmission: StartPendingSubmissionSpy;
updatePendingMessagesDisplay: Spy;
};
@@ -84,7 +85,6 @@ function createContext(): {
const cancelPendingSubmission = vi.fn(() => false);
const clearQueue = vi.fn(() => ({ steering: [], followUp: [] }));
const getQueuedMessages = vi.fn(() => ({ steering: [], followUp: [] }));
const interruptAndFlushQueuedMessages = vi.fn(async () => {});
const onInputCallback = vi.fn();
const requestRender = vi.fn();
const resetDisplay = vi.fn();
@@ -94,7 +94,6 @@ function createContext(): {
const hasActiveBtw = vi.fn(() => false);
const handleOmfgEscape = vi.fn(() => true);
const hasActiveOmfg = vi.fn(() => false);
const notifyInterrupting = vi.fn();
const updatePendingMessagesDisplay = vi.fn();
const prompt = vi.fn();
const startPendingSubmission = vi.fn(
@@ -155,7 +154,6 @@ function createContext(): {
abortEval,
clearQueue,
getQueuedMessages,
interruptAndFlushQueuedMessages,
prompt,
} as unknown as InteractiveModeContext["session"],
viewSession: {
@@ -168,6 +166,7 @@ function createContext(): {
} as unknown as InteractiveModeContext["viewSession"],
sessionManager: {
getSessionName: () => "existing session",
flushSync: vi.fn(),
} as unknown as InteractiveModeContext["sessionManager"],
keybindings: {
getKeys: () => [],
@@ -182,7 +181,6 @@ function createContext(): {
addMessageToChat,
cancelPendingSubmission,
ensureLoadingAnimation,
notifyInterrupting,
finishPendingSubmission: vi.fn(),
flushPendingBashComponents: vi.fn(),
markPendingSubmissionStarted: vi.fn(() => true),
@@ -203,6 +201,8 @@ function createContext(): {
showTreeSelector: vi.fn(),
showUserMessageSelector: vi.fn(),
showSessionSelector: vi.fn(),
shutdown: vi.fn(async () => {}),
clearEditor: vi.fn(),
} as unknown as InteractiveModeContext;
return {
@@ -215,19 +215,20 @@ function createContext(): {
addMessageToChat,
cancelPendingSubmission,
clearQueue,
clearEditor: ctx.clearEditor as Spy,
getQueuedMessages,
ensureLoadingAnimation,
interruptAndFlushQueuedMessages,
flushSync: ctx.sessionManager.flushSync as Spy,
handleBtwCommand,
handleBtwEscape,
hasActiveBtw,
handleOmfgEscape,
hasActiveOmfg,
notifyInterrupting,
onInputCallback,
prompt,
requestRender,
resetDisplay,
shutdown: ctx.shutdown as Spy,
startPendingSubmission,
updatePendingMessagesDisplay,
},
@@ -268,25 +269,24 @@ describe("InputController escape behavior", () => {
expect(spies.abort).not.toHaveBeenCalled();
});
it("acknowledges empty-submit interrupt before flushing queued steering", async () => {
it("empty-submit with a queued message aborts the active stream and refreshes pending display", async () => {
const { ctx, editor, spies } = createContext();
(ctx.session as { isStreaming: boolean; queuedMessageCount: number }).isStreaming = true;
(ctx.session as { isStreaming: boolean; queuedMessageCount: number }).queuedMessageCount = 1;
spies.getQueuedMessages.mockReturnValue({ steering: ["queued steer"], followUp: [] });
const order: string[] = [];
spies.notifyInterrupting.mockImplementation(() => {
order.push("notify");
spies.abort.mockImplementation(async () => {
order.push("abort");
});
spies.interruptAndFlushQueuedMessages.mockImplementation(async () => {
order.push("flush");
spies.updatePendingMessagesDisplay.mockImplementation(() => {
order.push("refresh");
});
const controller = new InputController(ctx);
controller.setupEditorSubmitHandler();
await editor.onSubmit?.("");
expect(order).toEqual(["notify", "flush"]);
expect(spies.interruptAndFlushQueuedMessages).toHaveBeenCalledWith({ reason: USER_INTERRUPT_LABEL });
expect(order).toEqual(["abort", "refresh"]);
expect(spies.abort).toHaveBeenCalledWith({ reason: USER_INTERRUPT_LABEL });
expect(spies.updatePendingMessagesDisplay).toHaveBeenCalledTimes(1);
expect(spies.requestRender).toHaveBeenCalledTimes(1);
});
@@ -480,7 +480,45 @@ describe("InputController escape behavior", () => {
editor.onEscape?.(); // clears text, must also reset the timer
editor.onEscape?.(); // empty again: should only re-arm, not trigger
expect(ctx.showTreeSelector).not.toHaveBeenCalled();
expect(ctx.showUserMessageSelector).not.toHaveBeenCalled();
});
});
describe("InputController Ctrl+C behavior", () => {
it("sync-flushes the session JSONL on first Ctrl+C (editor clear)", () => {
const { ctx, editor, spies } = createContext();
const controller = new InputController(ctx);
controller.setupKeyHandlers();
editor.onClear?.(); // first Ctrl+C: clears editor
expect(spies.clearEditor).toHaveBeenCalledTimes(1);
expect(spies.flushSync).toHaveBeenCalledTimes(1);
expect(spies.shutdown).not.toHaveBeenCalled();
});
it("sync-flushes the session JSONL on second Ctrl+C (shutdown)", () => {
const { ctx, editor, spies } = createContext();
const controller = new InputController(ctx);
controller.setupKeyHandlers();
editor.onClear?.(); // first Ctrl+C
editor.onClear?.(); // second Ctrl+C within 500ms → shutdown
expect(spies.shutdown).toHaveBeenCalledTimes(1);
// flushSync fires on both presses; the first-press flush is the
// guarantee that the JSONL is on disk even if the user closes the
// terminal before the second press.
expect(spies.flushSync).toHaveBeenCalledTimes(2);
});
it("does not flush when Ctrl+C is not pressed", () => {
const { ctx, editor, spies } = createContext();
const controller = new InputController(ctx);
controller.setupKeyHandlers();
editor.onEscape?.(); // Esc is a different handler
expect(spies.flushSync).not.toHaveBeenCalled();
});
});
@@ -55,8 +55,6 @@ async function createContext() {
const terminalWrite = vi.fn();
const prompt = vi.fn(async () => {});
const abort = vi.fn(async () => {});
const interruptAndFlushQueuedMessages = vi.fn(async () => {});
const getQueuedMessages = vi.fn(() => ({ steering: [] as string[], followUp: [] as string[] }));
const updatePendingMessagesDisplay = vi.fn();
const editor: FakeEditor = {
setText(text: string) {
@@ -96,9 +94,7 @@ async function createContext() {
extensionRunner: undefined,
prompt,
queuedMessageCount: 0,
getQueuedMessages,
abort,
interruptAndFlushQueuedMessages,
} as unknown as InteractiveModeContext["session"],
keybindings: {
getKeys(action: string) {
@@ -149,7 +145,6 @@ async function createContext() {
showModelSelector,
updateEditorBorderColor: vi.fn(),
hasActiveBtw: vi.fn(() => false),
notifyInterrupting: vi.fn(),
showError: vi.fn(),
} as unknown as InteractiveModeContext;
@@ -165,8 +160,6 @@ async function createContext() {
updatePendingMessagesDisplay,
requestRender,
abort,
interruptAndFlushQueuedMessages,
getQueuedMessages,
resetDisplay,
},
};
@@ -196,19 +189,17 @@ describe("InputController keybinding setup", () => {
expect(spies.resetDisplay).toHaveBeenCalledTimes(1);
});
it("empty Enter interrupts and sends a queued steering message", async () => {
it("empty Enter aborts the active stream when queued messages are pending", async () => {
const { InputController, ctx, editor, spies } = await createContext();
const session = ctx.session as unknown as { isStreaming: boolean; queuedMessageCount: number };
session.isStreaming = true;
session.queuedMessageCount = 1;
spies.getQueuedMessages.mockReturnValue({ steering: ["Send this now"], followUp: [] });
const controller = new InputController(ctx);
controller.setupEditorSubmitHandler();
await editor.onSubmit?.("");
expect(spies.interruptAndFlushQueuedMessages).toHaveBeenCalledWith({ reason: "Interrupted by user" });
expect(spies.abort).not.toHaveBeenCalled();
expect(spies.abort).toHaveBeenCalledWith({ reason: "Interrupted by user" });
expect(spies.updatePendingMessagesDisplay).toHaveBeenCalledTimes(1);
expect(spies.requestRender).toHaveBeenCalledTimes(1);
expect(spies.prompt).not.toHaveBeenCalled();
@@ -1,23 +1,11 @@
/**
* Phase 6 — E layer.
* Skill/custom queued-message display contracts.
*
* Tests the skill-queue + custom-role dequeue contract that ties together:
* - InputController.#invokeSkillCommand (tag generation when streaming);
* - AgentSession.enqueueCustomMessageDisplay + #handleAgentEvent's
* custom-role `message_start` dequeue;
* - UiHelpers.updatePendingMessagesDisplay (compact slash-form rendering);
* - InputController.restoreQueuedMessagesToEditor (recovery of the slash-form
* into the editor).
*
* Tests split into:
* - E1-E3: InputController-side tag generation, stubbed session;
* - E4-E7: Real AgentSession driving synthetic `message_start` events
* for the tag-based custom-role dequeue;
* - E8: real UiHelpers render against a queued-display entry;
* - E9: real InputController.restoreQueuedMessagesToEditor.
* Custom queued chips now ride on the queued AgentMessage itself via
* details.__queueChipText. The session derives pending display directly from
* the agent-core queue; there is no separate display mirror to splice.
*/
import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
import * as fs from "node:fs";
import { afterEach, beforeEach, describe, expect, it, type Mock, vi } from "bun:test";
import * as path from "node:path";
import { Agent } from "@oh-my-pi/pi-agent-core";
import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
@@ -35,27 +23,26 @@ import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manage
import { Container } from "@oh-my-pi/pi-tui";
import { TempDir } from "@oh-my-pi/pi-utils";
// ============================================================================
// Shared helpers
// ============================================================================
function writeSkillFile(dir: string, skillName: string, body: string): string {
const skillPath = path.join(dir, `${skillName}.md`);
fs.writeFileSync(skillPath, `---\nname: ${skillName}\n---\n${body}\n`);
return skillPath;
}
// ============================================================================
// E1-E3: InputController tag generation with a stubbed session.
// ============================================================================
type StubEditor = {
setText: (text: string) => void;
getText: () => string;
addToHistory: ReturnType<typeof vi.fn>;
addToHistory: Mock<(...args: unknown[]) => unknown>;
onSubmit?: (text: string) => Promise<void>;
};
type PromptCustomMessage = Mock<
(
message: { details: SkillPromptDetails },
options?: { streamingBehavior?: "steer" | "followUp"; queueChipText?: string },
) => Promise<void>
>;
async function writeSkillFile(dir: string, skillName: string, body: string): Promise<string> {
const skillPath = path.join(dir, `${skillName}.md`);
await Bun.write(skillPath, `---\nname: ${skillName}\n---\n${body}\n`);
return skillPath;
}
function createStubInputControllerContext(opts: { skillCommands: Map<string, string>; isStreaming: boolean }) {
let editorText = "";
const editor: StubEditor = {
@@ -67,10 +54,7 @@ function createStubInputControllerContext(opts: { skillCommands: Map<string, str
},
addToHistory: vi.fn(),
};
const enqueueCustomMessageDisplay = vi.fn((_text: string, _mode: "steer" | "followUp") => "sk-test-0");
// Annotate parameters so `mock.calls[N]` is typed as a tuple (not `[]`) and
// `message` carries required skill prompt details for assertion below.
const promptCustomMessage = vi.fn(async (_message: { details: SkillPromptDetails }, _options?: unknown) => {});
const promptCustomMessage: PromptCustomMessage = vi.fn(async () => {});
const prompt = vi.fn(async (_text: string, _options?: unknown) => {});
const handleGoalModeCommand = vi.fn(async (_rest?: string) => {});
const updatePendingMessagesDisplay = vi.fn();
@@ -87,7 +71,6 @@ function createStubInputControllerContext(opts: { skillCommands: Map<string, str
isBashRunning: false,
isEvalRunning: false,
extensionRunner: undefined,
enqueueCustomMessageDisplay,
prompt,
promptCustomMessage,
},
@@ -98,7 +81,6 @@ function createStubInputControllerContext(opts: { skillCommands: Map<string, str
handleGoalModeCommand,
goalModeEnabled: false,
updatePendingMessagesDisplay,
// Defaults that InputController touches on submit but don't matter here.
isBashMode: false,
isPythonMode: false,
pendingImages: [],
@@ -109,16 +91,24 @@ function createStubInputControllerContext(opts: { skillCommands: Map<string, str
withLocalSubmission: async (_text: string, fn: () => unknown) => fn(),
} as unknown as InteractiveModeContext;
return { ctx, editor, enqueueCustomMessageDisplay, prompt, promptCustomMessage, handleGoalModeCommand };
return {
ctx,
editor,
prompt,
promptCustomMessage,
handleGoalModeCommand,
updatePendingMessagesDisplay,
requestRender,
};
}
describe("InputController #invokeSkillCommand (E1-E3)", () => {
describe("InputController skill queue chip metadata", () => {
let tempDir: TempDir;
let skillCommands: Map<string, string>;
beforeEach(() => {
beforeEach(async () => {
tempDir = TempDir.createSync("@pi-skill-queue-stub-");
const skillPath = writeSkillFile(tempDir.path(), "test-skill", "Do the thing.");
const skillPath = await writeSkillFile(tempDir.path(), "test-skill", "Do the thing.");
skillCommands = new Map<string, string>([["skill:test-skill", skillPath]]);
});
@@ -127,60 +117,48 @@ describe("InputController #invokeSkillCommand (E1-E3)", () => {
vi.restoreAllMocks();
});
it("E1: streaming + steer -> enqueueCustomMessageDisplay called and details.__pendingDisplayTag set", async () => {
const { ctx, editor, enqueueCustomMessageDisplay, promptCustomMessage } = createStubInputControllerContext({
skillCommands,
isStreaming: true,
});
it("passes slash-form queueChipText for streaming skill steers", async () => {
const { ctx, editor, promptCustomMessage, updatePendingMessagesDisplay, requestRender } =
createStubInputControllerContext({ skillCommands, isStreaming: true });
const controller = new InputController(ctx);
controller.setupEditorSubmitHandler();
editor.setText("/skill:test-skill arg1 arg2");
await editor.onSubmit?.("/skill:test-skill arg1 arg2");
expect(enqueueCustomMessageDisplay).toHaveBeenCalledTimes(1);
expect(enqueueCustomMessageDisplay).toHaveBeenCalledWith("/skill:test-skill arg1 arg2", "steer");
expect(promptCustomMessage).toHaveBeenCalledTimes(1);
const firstCall = promptCustomMessage.mock.calls[0];
expect(firstCall).toBeDefined();
if (!firstCall) {
throw new Error("expected promptCustomMessage to be called");
}
const messageArg = firstCall[0];
expect(messageArg.details.__pendingDisplayTag).toBe("sk-test-0");
expect(promptCustomMessage.mock.calls[0]?.[1]).toEqual({
streamingBehavior: "steer",
queueChipText: "/skill:test-skill arg1 arg2",
});
expect(promptCustomMessage.mock.calls[0]?.[0].details.__queueChipText).toBeUndefined();
expect(updatePendingMessagesDisplay).toHaveBeenCalledTimes(1);
expect(requestRender).toHaveBeenCalledTimes(1);
});
it("E2: streaming + followUp -> enqueueCustomMessageDisplay called with mode 'followUp', tag embedded", async () => {
const { ctx, editor, enqueueCustomMessageDisplay, promptCustomMessage } = createStubInputControllerContext({
it("passes slash-form queueChipText for streaming skill follow-ups", async () => {
const { ctx, editor, promptCustomMessage } = createStubInputControllerContext({
skillCommands,
isStreaming: true,
});
const controller = new InputController(ctx);
editor.setText("/skill:test-skill arg1 arg2");
// `handleFollowUp` is the Ctrl+Enter dispatcher; it routes through the same
// `#invokeSkillCommand` helper with mode "followUp".
await controller.handleFollowUp();
expect(enqueueCustomMessageDisplay).toHaveBeenCalledWith("/skill:test-skill arg1 arg2", "followUp");
const firstCall = promptCustomMessage.mock.calls[0];
expect(firstCall).toBeDefined();
if (!firstCall) {
throw new Error("expected promptCustomMessage to be called");
}
const messageArg = firstCall[0];
expect(messageArg.details.__pendingDisplayTag).toBe("sk-test-0");
expect(promptCustomMessage.mock.calls[0]?.[1]).toEqual({
streamingBehavior: "followUp",
queueChipText: "/skill:test-skill arg1 arg2",
});
});
it("E2b: streaming follow-up applies builtin slash commands instead of queueing them", async () => {
it("streaming follow-up applies builtin slash commands instead of queueing them", async () => {
const { ctx, editor, prompt, handleGoalModeCommand } = createStubInputControllerContext({
skillCommands,
isStreaming: true,
});
const controller = new InputController(ctx);
editor.setText("/goal set Ship the release");
await controller.handleFollowUp();
@@ -189,32 +167,25 @@ describe("InputController #invokeSkillCommand (E1-E3)", () => {
expect(editor.getText()).toBe("");
});
it("E3: not streaming -> enqueueCustomMessageDisplay NOT called and tag absent", async () => {
const { ctx, editor, enqueueCustomMessageDisplay, promptCustomMessage } = createStubInputControllerContext({
it("idle skill prompt still leaves queueChipText out of persisted details", async () => {
const { ctx, editor, promptCustomMessage } = createStubInputControllerContext({
skillCommands,
isStreaming: false,
});
const controller = new InputController(ctx);
controller.setupEditorSubmitHandler();
editor.setText("/skill:test-skill arg1 arg2");
await editor.onSubmit?.("/skill:test-skill arg1 arg2");
expect(enqueueCustomMessageDisplay).not.toHaveBeenCalled();
const firstCall = promptCustomMessage.mock.calls[0];
expect(firstCall).toBeDefined();
if (!firstCall) {
throw new Error("expected promptCustomMessage to be called");
}
const messageArg = firstCall[0];
expect(messageArg.details.__pendingDisplayTag).toBeUndefined();
expect(promptCustomMessage.mock.calls[0]?.[1]).toEqual({
streamingBehavior: "steer",
queueChipText: "/skill:test-skill arg1 arg2",
});
expect(promptCustomMessage.mock.calls[0]?.[0].details.__queueChipText).toBeUndefined();
});
});
// ============================================================================
// E4-E7: Real AgentSession driving synthetic `message_start` events.
// ============================================================================
interface SessionFixture {
tempDir: TempDir;
authStorage: AuthStorage;
@@ -248,24 +219,24 @@ async function createRealSession(): Promise<SessionFixture> {
return { tempDir, authStorage, session };
}
/** Emit a `message_start` for a custom message whose `details` carries the supplied tag. */
function emitCustomMessageStart(session: AgentSession, content: string, tag?: string): void {
const details: { __pendingDisplayTag?: string } | undefined =
tag === undefined ? undefined : { __pendingDisplayTag: tag };
session.agent.emitExternalEvent({
type: "message_start",
message: {
role: "custom",
customType: SKILL_PROMPT_MESSAGE_TYPE,
content,
display: true,
details,
timestamp: Date.now(),
},
function queueCustomSteer(session: AgentSession, chip: string, content = "skill body"): void {
session.agent.steer({
role: "custom",
customType: SKILL_PROMPT_MESSAGE_TYPE,
content,
display: true,
details: {
name: "foo",
path: "/s.md",
args: "bar",
lineCount: 1,
__queueChipText: chip,
} satisfies SkillPromptDetails,
timestamp: Date.now(),
});
}
describe("AgentSession custom-role tag dequeue (E4-E7)", () => {
describe("AgentSession derived queued custom display", () => {
let fixture: SessionFixture | undefined;
afterEach(async () => {
@@ -278,95 +249,42 @@ describe("AgentSession custom-role tag dequeue (E4-E7)", () => {
vi.restoreAllMocks();
});
it("E4: message_start with role=custom + matching tag removes the tagged display entry", async () => {
it("derives queued custom chip text directly from the agent steering queue", async () => {
fixture = await createRealSession();
const { session } = fixture;
const tag = session.enqueueCustomMessageDisplay("/skill:foo bar", "steer");
expect(tag).not.toBe("");
expect(session.getQueuedMessages().steering).toEqual(["/skill:foo bar"]);
emitCustomMessageStart(session, "irrelevant content", tag);
await Promise.resolve();
await Promise.resolve();
queueCustomSteer(session, "/skill:foo bar");
expect(session.getQueuedMessages().steering).toEqual(["/skill:foo bar"]);
expect(session.queuedMessageCount).toBe(1);
});
it("excludes display-suppressed custom messages from pending chips and counts", async () => {
fixture = await createRealSession();
const { session } = fixture;
session.agent.steer({
role: "custom",
customType: "internal",
content: "hidden",
display: false,
details: { __queueChipText: "hidden" },
timestamp: Date.now(),
});
expect(session.getQueuedMessages().steering).toEqual([]);
// And internal queue counters reflect the empty steer/followUp arrays. The
// pending-next-turn store stays at zero too because this test never queued one.
expect(session.queuedMessageCount).toBe(0);
});
it("E5: message_start with role=custom but no tag is a no-op", async () => {
it("popLastQueuedMessage restores chip text and removes the core queue entry", async () => {
fixture = await createRealSession();
const { session } = fixture;
session.enqueueCustomMessageDisplay("/skill:foo bar", "steer");
const beforeCount = session.queuedMessageCount;
expect(beforeCount).toBe(1);
queueCustomSteer(session, "/skill:foo bar");
emitCustomMessageStart(session, "irrelevant content"); // no tag
await Promise.resolve();
await Promise.resolve();
expect(session.getQueuedMessages().steering).toEqual(["/skill:foo bar"]);
expect(session.queuedMessageCount).toBe(beforeCount);
});
it("E6: two queued skills with identical args text are dequeued independently by tag", async () => {
fixture = await createRealSession();
const { session } = fixture;
const tag1 = session.enqueueCustomMessageDisplay("/skill:foo bar", "steer");
const tag2 = session.enqueueCustomMessageDisplay("/skill:foo bar", "steer");
expect(tag1).not.toBe(tag2);
expect(session.getQueuedMessages().steering).toEqual(["/skill:foo bar", "/skill:foo bar"]);
// Consume the SECOND-enqueued tag. After dequeue, the SURVIVING entry must be
// the one that was added FIRST — proves the dequeue keys off `tag`, not off
// `indexOf(text)` (which would always have removed the first match).
emitCustomMessageStart(session, "any", tag2);
await Promise.resolve();
await Promise.resolve();
expect(session.getQueuedMessages().steering).toEqual(["/skill:foo bar"]);
// Now dequeue the first; nothing left.
emitCustomMessageStart(session, "any", tag1);
await Promise.resolve();
await Promise.resolve();
expect(session.getQueuedMessages().steering).toEqual([]);
});
it("E7: popLastQueuedMessage on a tagged entry leaves no orphan tag state", async () => {
fixture = await createRealSession();
const { session } = fixture;
const firstTag = session.enqueueCustomMessageDisplay("/skill:foo bar", "steer");
const popped = session.popLastQueuedMessage();
expect(popped?.text).toBe("/skill:foo bar");
expect(session.getQueuedMessages().steering).toEqual([]);
// Push a NEW tagged entry with the same text. Emitting `message_start` for the
// FIRST (popped) tag must be a no-op — the dequeue cannot reach into the new
// entry because the popped tag died with its record.
const secondTag = session.enqueueCustomMessageDisplay("/skill:foo bar", "steer");
expect(secondTag).not.toBe(firstTag);
emitCustomMessageStart(session, "any", firstTag);
await Promise.resolve();
await Promise.resolve();
expect(session.getQueuedMessages().steering).toEqual(["/skill:foo bar"]);
// Sanity: the second tag still works.
emitCustomMessageStart(session, "any", secondTag);
await Promise.resolve();
await Promise.resolve();
expect(session.popLastQueuedMessage()?.text).toBe("/skill:foo bar");
expect(session.getQueuedMessages().steering).toEqual([]);
});
});
// ============================================================================
// E8-E9: Real UiHelpers / InputController against a session populated through
// enqueueCustomMessageDisplay.
// ============================================================================
function createStubInteractiveModeContextForUiHelpers(session: AgentSession) {
let editorText = "";
const editor: StubEditor = {
@@ -399,19 +317,10 @@ function createStubInteractiveModeContextForUiHelpers(session: AgentSession) {
return { ctx, editor, pendingMessagesContainer };
}
describe("UiHelpers / InputController against the queued-display layer (E8-E9)", () => {
describe("UiHelpers / InputController against derived queued custom display", () => {
let fixture: SessionFixture | undefined;
beforeEach(async () => {
// E8 invokes the real `theme.fg(...)` codepath inside
// updatePendingMessagesDisplay; without an initialized theme module the
// global `theme` variable is undefined. Installs `dark` per-test —
// matches the established suite convention used by other test files
// (bash-execution-clamp.test.ts, bash-execution-sixel.test.ts) where
// `dark` is the agreed default for every test that needs a theme.
// No `afterEach` restore is required by that convention; the theme
// module exposes no reset API, and `dark` is the suite-wide assumed
// post-state.
const themeInstance = await getThemeByName("dark");
expect(themeInstance).toBeDefined();
setThemeInstance(themeInstance!);
@@ -427,60 +336,35 @@ describe("UiHelpers / InputController against the queued-display layer (E8-E9)",
vi.restoreAllMocks();
});
it("E8: updatePendingMessagesDisplay renders the compact slash form for queued skills", async () => {
it("renders the compact slash form for queued skills", async () => {
fixture = await createRealSession();
const { session } = fixture;
session.enqueueCustomMessageDisplay("/skill:test-skill arg1 arg2", "steer");
queueCustomSteer(session, "/skill:test-skill arg1 arg2");
const { ctx, pendingMessagesContainer } = createStubInteractiveModeContextForUiHelpers(session);
const uiHelpers = new UiHelpers(ctx);
uiHelpers.updatePendingMessagesDisplay();
// Render the container at a generous width and assert the compact slash-form
// chip appears verbatim. Matches the user-facing "Steer: /skill:..." format.
const rendered = pendingMessagesContainer.render(120).join("\n");
expect(rendered).toMatch(/Steer: \/skill:test-skill arg1 arg2/);
});
it("E9: restoreQueuedMessagesToEditor recovers the compact slash form into the editor and clears the queue", async () => {
it("restores the compact slash form into the editor and clears the queue", async () => {
fixture = await createRealSession();
const { session } = fixture;
session.enqueueCustomMessageDisplay("/skill:test-skill arg1 arg2", "steer");
queueCustomSteer(session, "/skill:test-skill arg1 arg2");
const { ctx, editor } = createStubInteractiveModeContextForUiHelpers(session);
const controller = new InputController(ctx);
const count = controller.restoreQueuedMessagesToEditor();
expect(count).toBe(1);
expect(editor.getText()).toBe("/skill:test-skill arg1 arg2");
// Queue cleared on both arrays.
const { steering, followUp } = session.getQueuedMessages();
expect(steering).toEqual([]);
expect(followUp).toEqual([]);
expect(session.getQueuedMessages()).toEqual({ steering: [], followUp: [] });
});
});
// ============================================================================
// E10: EventController refreshes the pending-messages bar on tagged custom
// dequeue.
//
// Regression guard for the Codex P2 review finding on PR #1043: the
// custom-role `message_start` branch in AgentSession.#handleAgentEvent spliced
// the matching entry out of #steeringMessages / #followUpMessages correctly,
// but EventController.#handleMessageStart only called updatePendingMessagesDisplay
// from the `role === "user"` branch. The custom branch — which is where queued
// /skill: invocations flow — never rebuilt `pendingMessagesContainer`, so the
// chip kept painting until an unrelated trigger fired a refresh.
//
// The fix: in EventController's custom branch, when the dequeued message
// carries the `__pendingDisplayTag` (proof it was queued via
// enqueueCustomMessageDisplay), call updatePendingMessagesDisplay() before
// requestRender(). E10 covers both gate branches:
// - positive: tagged custom -> refresh fires once
// - negative: untagged custom (ttsr-injection, irc:*, async-result, hookMessage)
// -> refresh NOT fired (over-refresh guard)
// ============================================================================
function createEventControllerFixtureForE10() {
function createEventControllerFixture() {
const updatePendingMessagesDisplay = vi.fn();
const addMessageToChat = vi.fn();
const requestRender = vi.fn();
@@ -503,20 +387,14 @@ function createEventControllerFixtureForE10() {
return { controller, updatePendingMessagesDisplay, addMessageToChat };
}
describe("EventController custom-role dequeue refresh (E10)", () => {
describe("EventController custom queued-message refresh", () => {
afterEach(() => {
vi.restoreAllMocks();
});
it("E10: message_start with role=custom refreshes pending bar ONLY when __pendingDisplayTag is present", async () => {
const { controller, updatePendingMessagesDisplay, addMessageToChat } = createEventControllerFixtureForE10();
// Positive case: tagged custom => refresh fires exactly once. The tag is the
// unambiguous signal "this message was queued via enqueueCustomMessageDisplay";
// AgentSession.#handleAgentEvent has already spliced the matching entry out of
// the display arrays (ran before this emit), so the rebuild repaints the now-
// correct queue state.
const taggedEvent: Extract<AgentSessionEvent, { type: "message_start" }> = {
it("refreshes the pending bar only for custom messages carrying __queueChipText", async () => {
const { controller, updatePendingMessagesDisplay, addMessageToChat } = createEventControllerFixture();
const queuedEvent: Extract<AgentSessionEvent, { type: "message_start" }> = {
type: "message_start",
message: {
role: "custom",
@@ -524,7 +402,7 @@ describe("EventController custom-role dequeue refresh (E10)", () => {
content: "first",
display: true,
details: {
__pendingDisplayTag: "sk-test-0",
__queueChipText: "/skill:foo bar",
name: "foo",
path: "/s.md",
args: "bar",
@@ -533,18 +411,11 @@ describe("EventController custom-role dequeue refresh (E10)", () => {
timestamp: Date.now(),
},
};
await controller.handleEvent(taggedEvent);
await controller.handleEvent(queuedEvent);
expect(updatePendingMessagesDisplay).toHaveBeenCalledTimes(1);
// Chat rendering still ran — refresh is additive, not a replacement for the
// chat path.
expect(addMessageToChat).toHaveBeenCalledTimes(1);
// Negative case: untagged custom => refresh NOT fired. Over-refresh guard.
// Non-queued customs (ttsr-injection, irc:*, async-result, hookMessage) never
// registered a pending chip, so rebuilding pendingMessagesContainer for them
// would be pure waste. Distinct timestamp avoids the #renderedCustomMessages
// signature-dedup early-return.
const untaggedEvent: Extract<AgentSessionEvent, { type: "message_start" }> = {
const unqueuedEvent: Extract<AgentSessionEvent, { type: "message_start" }> = {
type: "message_start",
message: {
role: "custom",
@@ -555,12 +426,9 @@ describe("EventController custom-role dequeue refresh (E10)", () => {
timestamp: Date.now() + 1,
},
};
await controller.handleEvent(untaggedEvent);
// Still exactly 1 — no additional call from the untagged path.
await controller.handleEvent(unqueuedEvent);
expect(updatePendingMessagesDisplay).toHaveBeenCalledTimes(1);
// Chat rendering still ran for the untagged custom (the chat-add path is
// unconditional inside the custom branch; only the pending-bar refresh is
// tag-gated).
expect(addMessageToChat).toHaveBeenCalledTimes(2);
});
});
@@ -1,142 +1,66 @@
/**
* Regression test for the empty-messages interrupt-and-flush bug.
*
* Scenario the user hit:
* 1. The agent is already streaming (a hidden turn: context promotion,
* auto-retry, auto-compaction, an extension-emitted turn, etc.).
* 2. The user types "hi" and presses Enter. Because isStreaming is true,
* the input controller routes the message to `session.steer("hi")`
* (see the streamingBehavior branch in `AgentSession.prompt`).
* Critically, the user message is NEVER appended to `#state.messages`
* in the agent — it lives only in the steering queue.
* 3. The user presses Enter again on an empty editor to inject the steer
* immediately. The input controller calls
* `session.interruptAndFlushQueuedMessages()`.
* 4. `interruptAndFlushQueuedMessages` aborts the streaming turn, then
* calls `agent.continue()` to drain the queued steer. `continue()` is:
*
* if (messages.length === 0) throw "No messages to continue from";
* if (last role === "assistant") { drainSteering(); runLoop(drained); }
* else runLoop(undefined);
*
* When `#state.messages` is empty, step 4 throws and:
* - `agent.continue()` propagates the throw out of
* `interruptAndFlushQueuedMessages`.
* - The input controller catches it and surfaces the error string
* ("Error: No messages to continue from") — the visible
* "Error:" line in the user's first screenshot.
* - The steering queue is never drained, so the session's
* `#steeringMessages` mirror is never cleared by the
* message_start handler. The "Steer: hi" chip stays visible
* forever (the second screenshot) until the user manually
* dequeues with Alt+Up or starts a new session.
*
* Contract this test defends:
* - `interruptAndFlushQueuedMessages` MUST deliver a queued steer even
* when the agent has no prior messages in `#state.messages`.
* - The flush MUST NOT reject with any error.
* - It MUST clear both the agent queue and the session mirror in a
* single atomic step, so the visible chip disappears the moment the
* user presses Enter, not after a downstream message_start event
* fires.
* - The delivered messages must contain the original steer text.
*/
import { describe, expect, it, vi } from "bun:test";
import { InputController } from "@oh-my-pi/pi-coding-agent/modes/controllers/input-controller";
import type { InteractiveModeContext } from "@oh-my-pi/pi-coding-agent/modes/types";
import { USER_INTERRUPT_LABEL } from "@oh-my-pi/pi-coding-agent/session/messages";
import { afterEach, beforeEach, describe, expect, it } from "bun:test";
import * as fs from "node:fs";
import * as os from "node:os";
import * as path from "node:path";
import { Agent } from "@oh-my-pi/pi-agent-core";
import type { Message } from "@oh-my-pi/pi-ai";
import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream";
import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
import { convertToLlm } from "@oh-my-pi/pi-coding-agent/session/messages";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { Snowflake } from "@oh-my-pi/pi-utils";
import { createAssistantMessage } from "./helpers/agent-session-setup";
describe("interrupt-and-flush must not strand a queued steer when the agent has no prior messages", () => {
let session: AgentSession;
let tempDir: string;
const authStorages: AuthStorage[] = [];
beforeEach(() => {
tempDir = path.join(os.tmpdir(), `pi-flush-empty-messages-${Snowflake.next()}`);
fs.mkdirSync(tempDir, { recursive: true });
});
afterEach(async () => {
if (session) {
await session.dispose();
}
for (const authStorage of authStorages.splice(0)) {
authStorage.close();
}
if (tempDir && fs.existsSync(tempDir)) {
fs.rmSync(tempDir, { recursive: true });
}
});
it("delivers a steer queued on a session whose #state.messages is empty", async () => {
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
const callMessages: Message[][] = [];
const agent = new Agent({
getApiKey: () => "test-key",
initialState: { model, systemPrompt: ["Test"], tools: [] },
convertToLlm,
streamFn(_model, context, _options) {
callMessages.push([...context.messages]);
const stream = new AssistantMessageEventStream();
queueMicrotask(() => {
stream.push({ type: "start", partial: createAssistantMessage("") });
stream.push({ type: "done", reason: "stop", message: createAssistantMessage("ack") });
});
return stream;
function createContext() {
let editorText = "";
const abort = vi.fn(async () => {});
const prompt = vi.fn(async () => {});
const updatePendingMessagesDisplay = vi.fn();
const requestRender = vi.fn();
const showError = vi.fn();
const ctx = {
editor: {
setText(text: string) {
editorText = text;
},
});
getText() {
return editorText;
},
addToHistory: vi.fn(),
},
ui: { requestRender },
session: {
isStreaming: true,
isCompacting: false,
isBashRunning: false,
isEvalRunning: false,
queuedMessageCount: 1,
extensionRunner: undefined,
abort,
prompt,
},
get viewSession() {
return (this as typeof ctx).session;
},
pendingImages: [],
pendingImageLinks: [],
compactionQueuedMessages: [],
locallySubmittedUserSignatures: new Set<string>(),
isBashMode: false,
isPythonMode: false,
loopModeEnabled: false,
updatePendingMessagesDisplay,
showError,
hasActiveBtw: () => false,
hasActiveOmfg: () => false,
} as unknown as InteractiveModeContext;
return { ctx, abort, prompt, updatePendingMessagesDisplay, requestRender, showError };
}
const sessionManager = SessionManager.inMemory();
const settings = Settings.isolated();
const authStorage = await AuthStorage.create(path.join(tempDir, "auth.db"));
authStorages.push(authStorage);
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models.yml"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
describe("empty submit with queued messages", () => {
it("aborts the active stream instead of eagerly prompting a drained queue", async () => {
const { ctx, abort, prompt, updatePendingMessagesDisplay, requestRender, showError } = createContext();
const controller = new InputController(ctx);
controller.setupEditorSubmitHandler();
session = new AgentSession({ agent, sessionManager, settings, modelRegistry });
await ctx.editor.onSubmit?.("");
// Reproduce the user flow: a hidden turn is already streaming, so
// the input controller routed the user's "hi" to session.steer
// instead of session.prompt. The session mirror and agent queue
// are populated; #state.messages stays empty.
await session.steer("hi");
expect(session.agent.hasQueuedMessages()).toBe(true);
expect(session.getQueuedMessages().steering).toEqual(["hi"]);
// Empty-Enter → interruptAndFlushQueuedMessages. Must NOT throw.
// Awaiting the flush is enough: it resolves only after the resumed
// turn's prompt() has run and the assistant message has been
// emitted. The call to callMessages.push() happens synchronously
// inside streamFn, so by the time the flush resolves the model
// call has been recorded.
await expect(
session.interruptAndFlushQueuedMessages({ reason: "Interrupted by user" }),
).resolves.toBeUndefined();
// The steer must be delivered as a fresh turn.
expect(callMessages.length).toBeGreaterThanOrEqual(1);
const delivered = callMessages[0]?.some((message: Message) => {
if (typeof message.content === "string") return message.content.includes("hi");
return message.content.some(c => c.type === "text" && c.text.includes("hi"));
});
expect(delivered).toBe(true);
// Both queues must clear atomically — no stranded chip.
expect(session.agent.hasQueuedMessages()).toBe(false);
expect(session.getQueuedMessages()).toEqual({ steering: [], followUp: [] });
expect(abort).toHaveBeenCalledWith({ reason: USER_INTERRUPT_LABEL });
expect(prompt).not.toHaveBeenCalled();
expect(showError).not.toHaveBeenCalled();
expect(updatePendingMessagesDisplay).toHaveBeenCalledTimes(1);
expect(requestRender).toHaveBeenCalledTimes(1);
});
});
@@ -1,22 +1,18 @@
import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "bun:test";
import { resetSettingsForTest, Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import {
EventController,
INTERRUPTING_WORKING_MESSAGE,
} from "@oh-my-pi/pi-coding-agent/modes/controllers/event-controller";
import { EventController } from "@oh-my-pi/pi-coding-agent/modes/controllers/event-controller";
import { initTheme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
import type { InteractiveModeContext } from "@oh-my-pi/pi-coding-agent/modes/types";
import type { AgentSessionEvent } from "@oh-my-pi/pi-coding-agent/session/agent-session";
function createContext() {
const order: string[] = [];
const setWorkingMessage = vi.fn((_message?: string) => {
order.push("setWorkingMessage");
});
const ensureLoadingAnimation = vi.fn(() => {
order.push("ensureLoadingAnimation");
});
const setWorkingMessage = vi.fn();
const ensureLoadingAnimation = vi.fn();
const pendingTools = new Map<string, unknown>();
const session = {
getToolByName: () => undefined,
isAborting: false,
};
const ctx = {
isInitialized: true,
settings: { get: () => false },
@@ -28,9 +24,9 @@ function createContext() {
clearPinnedError: vi.fn(),
ensureLoadingAnimation,
ui: { requestRender: vi.fn() },
session: { getToolByName: () => undefined },
session,
} as unknown as InteractiveModeContext;
return { ctx, pendingTools, setWorkingMessage, ensureLoadingAnimation, order };
return { ctx, pendingTools, setWorkingMessage, session };
}
const AGENT_START = { type: "agent_start" } as unknown as AgentSessionEvent;
@@ -48,7 +44,7 @@ function toolStartWithIntent(toolCallId: string, intent: string): AgentSessionEv
} as unknown as AgentSessionEvent;
}
describe("EventController user interrupt acknowledgement", () => {
describe("EventController aborted-turn working messages", () => {
beforeAll(async () => {
await initTheme(false);
});
@@ -63,38 +59,20 @@ describe("EventController user interrupt acknowledgement", () => {
resetSettingsForTest();
});
it("swaps the loader to the interrupting label once a turn is active", async () => {
const { ctx, setWorkingMessage, ensureLoadingAnimation, order } = createContext();
it("suppresses late intent-driven working-message updates while aborting", async () => {
const { ctx, pendingTools, setWorkingMessage, session } = createContext();
const controller = new EventController(ctx);
await controller.handleEvent(AGENT_START);
setWorkingMessage.mockClear();
ensureLoadingAnimation.mockClear();
order.length = 0;
session.isAborting = true;
controller.notifyInterrupting();
expect(order).toEqual(["ensureLoadingAnimation", "setWorkingMessage"]);
expect(ensureLoadingAnimation).toHaveBeenCalledTimes(1);
expect(setWorkingMessage).toHaveBeenCalledWith(INTERRUPTING_WORKING_MESSAGE);
});
it("freezes intent-driven working-message updates while interrupting", async () => {
const { ctx, pendingTools, setWorkingMessage } = createContext();
const controller = new EventController(ctx);
await controller.handleEvent(AGENT_START);
controller.notifyInterrupting();
setWorkingMessage.mockClear();
// A tool whose args already started streaming before the abort still emits a
// late tool_execution_start; without the freeze its intent would repaint the
// loader over the "Interrupting…" acknowledgement.
pendingTools.set("late-call", {});
await controller.handleEvent(toolStartWithIntent("late-call", "Reticulating splines"));
expect(setWorkingMessage).not.toHaveBeenCalled();
});
it("lets intent updates drive the loader when not interrupting", async () => {
it("lets intent updates drive the loader when not aborting", async () => {
const { ctx, pendingTools, setWorkingMessage } = createContext();
const controller = new EventController(ctx);
await controller.handleEvent(AGENT_START);
@@ -107,15 +85,16 @@ describe("EventController user interrupt acknowledgement", () => {
expect(setWorkingMessage.mock.calls[0]?.[0]).toContain("Searching files");
});
it("clears the interrupt freeze at the next agent_start", async () => {
const { ctx, pendingTools, setWorkingMessage } = createContext();
it("resumes intent updates once aborting clears", async () => {
const { ctx, pendingTools, setWorkingMessage, session } = createContext();
const controller = new EventController(ctx);
await controller.handleEvent(AGENT_START);
controller.notifyInterrupting();
session.isAborting = true;
// New turn: the freeze must lift so the next turn's intents render again.
await controller.handleEvent(AGENT_START);
pendingTools.set("late-call", {});
await controller.handleEvent(toolStartWithIntent("late-call", "Reticulating splines"));
setWorkingMessage.mockClear();
session.isAborting = false;
pendingTools.set("call-2", {});
await controller.handleEvent(toolStartWithIntent("call-2", "Editing module"));
@@ -123,14 +102,4 @@ describe("EventController user interrupt acknowledgement", () => {
expect(setWorkingMessage).toHaveBeenCalledTimes(1);
expect(setWorkingMessage.mock.calls[0]?.[0]).toContain("Editing module");
});
it("is a no-op before any turn starts", () => {
const { ctx, setWorkingMessage } = createContext();
const controller = new EventController(ctx);
controller.notifyInterrupting();
expect(setWorkingMessage).not.toHaveBeenCalled();
});
});
@@ -9,10 +9,10 @@ import { TERMINAL } from "@oh-my-pi/pi-tui";
/**
* Models the loader lifecycle InteractiveMode owns: `agent_start` creates the
* loader via `ensureLoadingAnimation`; `agent_end` stops and drops it. The
* streaming/resume getters are backed by mutable flags the tests drive directly.
* streaming getter is backed by mutable flags the tests drive directly.
*/
function createContext() {
const streamState = { isStreaming: false, isResumingQueuedMessages: false, queuedMessageCount: 0 };
const streamState = { isStreaming: false };
const loader = { stop: vi.fn() };
const ctx = {
isInitialized: true,
@@ -39,12 +39,6 @@ function createContext() {
get isStreaming() {
return streamState.isStreaming;
},
get isResumingQueuedMessages() {
return streamState.isResumingQueuedMessages;
},
get queuedMessageCount() {
return streamState.queuedMessageCount;
},
getToolByName: () => undefined,
},
} as unknown as InteractiveModeContext;
@@ -82,9 +76,9 @@ describe("EventController superseded agent_end", () => {
await controller.handleEvent(AGENT_START);
expect(ctx.loadingAnimation).toBeDefined();
// User interrupt-and-flush of a queued steer: the resumed turn's agent_start
// arrives and the agent is streaming again. The interrupted turn's agent_end
// is still in flight through the async event pipeline.
// User abort of a queued steer: the resumed turn's agent_start arrives and
// the agent is streaming again. The interrupted turn's agent_end is still in
// flight through the async event pipeline.
streamState.isStreaming = true;
await controller.handleEvent(AGENT_START);
@@ -98,32 +92,6 @@ describe("EventController superseded agent_end", () => {
expect(TERMINAL.sendNotification).not.toHaveBeenCalled();
});
it("keeps the loader alive while an interrupted queued resume is being armed", async () => {
vi.useFakeTimers();
const { ctx, streamState, loader } = createContext();
const controller = new EventController(ctx);
await controller.handleEvent(AGENT_START);
controller.notifyInterrupting();
streamState.isStreaming = false;
streamState.isResumingQueuedMessages = true;
streamState.queuedMessageCount = 1;
await controller.handleEvent(AGENT_END);
vi.advanceTimersByTime(16);
expect(loader.stop).not.toHaveBeenCalled();
expect(ctx.loadingAnimation).toBeDefined();
streamState.isStreaming = true;
streamState.isResumingQueuedMessages = false;
await controller.handleEvent(AGENT_START);
vi.advanceTimersByTime(16);
expect(loader.stop).not.toHaveBeenCalled();
expect(ctx.loadingAnimation).toBeDefined();
});
it("tears the loader down on the live turn's own final agent_end", async () => {
const { ctx, streamState, loader } = createContext();
const controller = new EventController(ctx);
@@ -27,7 +27,7 @@ function readPersistedCustomMessageEntry<T>(session: SessionManager, id: string)
}
describe("SessionManager.appendCustomMessageEntry (allowlist strip + persistence contract)", () => {
it("F1: strips __pendingDisplayTag from persisted details while preserving all other SkillPromptDetails fields", () => {
it("F1: strips __queueChipText from persisted details while preserving all other SkillPromptDetails fields", () => {
const session = SessionManager.inMemory();
const id = session.appendCustomMessageEntry<SkillPromptDetails>(
SKILL_TYPE,
@@ -38,7 +38,7 @@ describe("SessionManager.appendCustomMessageEntry (allowlist strip + persistence
path: "/s.md",
args: "bar",
lineCount: 10,
__pendingDisplayTag: "omp-cmd-1-0",
__queueChipText: "omp-cmd-1-0",
},
"user",
);
@@ -52,7 +52,7 @@ describe("SessionManager.appendCustomMessageEntry (allowlist strip + persistence
});
// Explicit absence assertion — defends against `toEqual` semantics drift
// where an `undefined`-valued key would still satisfy deep equality.
expect(Object.hasOwn(entry.details!, "__pendingDisplayTag")).toBe(false);
expect(Object.hasOwn(entry.details!, "__queueChipText")).toBe(false);
});
it("F2: persists details deep-equal to the input when no allowlisted field is present", () => {
+1
View File
@@ -232,6 +232,7 @@ export interface SessionState {
thinkingLevel?: string;
contextUsage?: ContextUsage;
participants: Participant[];
isAborting?: boolean;
}
export interface AgentSnapshot {