fix(coding-agent): fixed session magic-keyword ordering and stranded queue handling issues
- Fixed magic-keyword notices to preserve ordering in agent-session processing. - Fixed stranded queue behavior during queued steer/skill delivery in session logic. - Updated agent-session and input-controller tests covering suppression, keywords, and queues. - Updated unreleased changelog notes describing the magic-keyword and queue fixes.
This commit is contained in:
@@ -1,9 +1,14 @@
|
||||
# Changelog
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed magic-keyword steering notices (`ultrathink-notice`, `orchestrate-notice`, `workflow-notice`) to be prepended before the related user message so they influence that same turn
|
||||
- Fixed dequeuing or popping queued user messages to remove their preceding hidden magic-keyword notice companions, preventing orphaned queued notices
|
||||
- Fixed queued user steers to auto-resume after interrupts even when the transcript tail is a preserved advisor card or other non-conversational custom message
|
||||
- Fixed queued user follow-up messages to remain queued after an interrupt and only run on explicit resume, even when an IRC wake leaves a provider-valid tail
|
||||
- Fixed stranded IRC asides to wake a response turn after interruption instead of remaining pending
|
||||
- Fixed accepted IRC asides to be flushed into the transcript during disposal instead of being discarded
|
||||
- Fixed interactive submissions made while the TUI had no active input waiter: they now start a real prompt directly, with steer fallback if a background turn races in, instead of queueing behind a non-resumable idle transcript and appearing to do nothing.
|
||||
- Fixed pressing Esc (or Alt+Up dequeue) while agent-authored messages were queued — advisor concern/blocker notes, hidden goal/plan/budget steers, IRC/extension asides — dumping their text into the user's editor. Editor restoration (`clearQueue()`), pending chips (`getQueuedMessages()`), and `popLastQueuedMessage()` now surface only genuinely user-authored queued messages (plain user turns and `attribution: "user"` custom messages like `/skill`). Plain Alt+Up dequeue leaves all other queued messages in place for the continuing stream; only the Esc interrupt path keeps just advisor cards (so abort's preservation still re-records them as visible advice) and drops other internal steers, so a user interrupt can't be silently undone by an auto-resume on leftover internal context. `queuedMessageCount` still reflects all actual queued work (advisor cards included) so `hasPendingMessages()`/RPC and the empty-submit abort gate stay accurate.
|
||||
- Fixed `omp --continue`/`-c` sometimes resuming into a subagent transcript instead of the interactive session. Subagent (and HTML-export) `SessionManager.open()` calls run in the parent's terminal and were clobbering the per-TTY `--continue` breadcrumb with their own artifact-dir session file; these headless opens now suppress the breadcrumb. `continueRecent()` also recovers already-poisoned breadcrumbs by resolving any session file inside a parent's artifacts dir (`<parent>/<agentId>.jsonl`) back up to the top-level session.
|
||||
@@ -11795,4 +11800,4 @@ Initial public release.
|
||||
|
||||
## [0.7.6] - 2025-11-13
|
||||
|
||||
Previous releases did not maintain a changelog.
|
||||
Previous releases did not maintain a changelog.
|
||||
@@ -968,6 +968,31 @@ function isUserQueuedMessage(message: AgentMessage): boolean {
|
||||
return message.role === "custom" && message.attribution === "user" && message.display !== false;
|
||||
}
|
||||
|
||||
/** Custom-message types of the hidden magic-keyword notices that `#createMagicKeywordNotices`
|
||||
* enqueues alongside a user prompt. Keep in sync with that method. */
|
||||
const MAGIC_KEYWORD_NOTICE_TYPES: ReadonlySet<string> = new Set([
|
||||
"ultrathink-notice",
|
||||
"orchestrate-notice",
|
||||
"workflow-notice",
|
||||
]);
|
||||
|
||||
/**
|
||||
* A hidden, user-attributed companion of a queued user prompt: the magic-keyword
|
||||
* notices (`ultrathink`/`orchestrate`/`workflow`) enqueued alongside the user
|
||||
* message. They are `attribution: "user"` but `display: false`, so they are not
|
||||
* editor-restorable; when the user pulls their prompt back out of the queue these
|
||||
* must leave with it rather than linger as stale, companion-less steering. Scoped to
|
||||
* the known notice types so an unrelated hidden user custom is never silently dropped.
|
||||
*/
|
||||
function isHiddenUserCompanion(message: AgentMessage): boolean {
|
||||
return (
|
||||
message.role === "custom" &&
|
||||
message.attribution === "user" &&
|
||||
message.display === false &&
|
||||
MAGIC_KEYWORD_NOTICE_TYPES.has(message.customType)
|
||||
);
|
||||
}
|
||||
|
||||
function queueChipText(message: AgentMessage): string {
|
||||
if (message.role === "custom") {
|
||||
return readQueueChipText(message.details) ?? queuedTextContent(message) ?? "";
|
||||
@@ -1274,6 +1299,57 @@ export class AgentSession {
|
||||
#drainStrandedQueuedMessages(): void {
|
||||
if (this.#abortInProgress) return;
|
||||
this.#scheduleQueuedMessageDrain();
|
||||
this.#resumeStrandedIrcAsides();
|
||||
}
|
||||
|
||||
/** IRC asides that arrive after the loop's final aside poll — or while an abort skipped that
|
||||
* poll — land in #pendingIrcAsides with no loop left to drain them; the queued-message drain's
|
||||
* gate (agent.hasQueuedMessages()) does not count them. Once idle, wake a turn so the agent
|
||||
* responds to the peer. Skip only when a queued steer/follow-up will itself drive a resume turn
|
||||
* whose aside poll already consumes these (no double-wake). */
|
||||
#resumeStrandedIrcAsides(): void {
|
||||
if (this.#isDisposed || this.isStreaming) return;
|
||||
if (this.#pendingIrcAsides.length === 0) return;
|
||||
if (this.#canAutoContinueForFollowUp() && this.agent.hasQueuedMessages()) return;
|
||||
const records = this.#pendingIrcAsides;
|
||||
this.#pendingIrcAsides = [];
|
||||
this.#wakeForIrc(records);
|
||||
}
|
||||
|
||||
/** Fire-and-forget wake turn for incoming IRC — idle delivery and stranded-aside resume both
|
||||
* route here. Wrapped in #beginInFlight/#endInFlight so the turn is tracked and its settle
|
||||
* re-drains anything that stranded during it. A user interrupt may have intentionally left a
|
||||
* follow-up queued behind an invalid tail (seam #5); the wake turn's loop would otherwise drain
|
||||
* it, so park the follow-up queue across the wake and restore it after. It stays queued post-wake
|
||||
* because #canAutoContinueForFollowUp suppresses follow-up auto-resume while a user interrupt is
|
||||
* in effect, even though the wake left a provider-valid tail. */
|
||||
#wakeForIrc(records: CustomMessage[]): void {
|
||||
// Park only a *blocked* follow-up (one a user interrupt is intentionally holding); an
|
||||
// already-resumable follow-up can ride the wake turn normally without reordering.
|
||||
const parkedFollowUps =
|
||||
this.agent.peekSteeringQueue().length === 0 &&
|
||||
this.agent.peekFollowUpQueue().length > 0 &&
|
||||
!this.#canAutoContinueForFollowUp()
|
||||
? [...this.agent.peekFollowUpQueue()]
|
||||
: [];
|
||||
if (parkedFollowUps.length > 0) {
|
||||
this.agent.replaceQueues([...this.agent.peekSteeringQueue()], []);
|
||||
}
|
||||
this.#beginInFlight();
|
||||
void this.agent
|
||||
.prompt(records)
|
||||
.catch(error => {
|
||||
logger.warn("IRC wake turn failed", { error: String(error) });
|
||||
})
|
||||
.finally(() => {
|
||||
if (parkedFollowUps.length > 0) {
|
||||
this.agent.replaceQueues(
|
||||
[...this.agent.peekSteeringQueue()],
|
||||
[...parkedFollowUps, ...this.agent.peekFollowUpQueue()],
|
||||
);
|
||||
}
|
||||
this.#endInFlight();
|
||||
});
|
||||
}
|
||||
|
||||
/** Remove advisor concern/blocker cards from the agent-core steer/follow-up
|
||||
@@ -3657,7 +3733,7 @@ export class AgentSession {
|
||||
*/
|
||||
beginDispose(): void {
|
||||
this.#isDisposed = true;
|
||||
this.#pendingIrcAsides = [];
|
||||
this.#flushPendingIrcAsides();
|
||||
this.yieldQueue.clear();
|
||||
this.agent.setAsideMessageProvider(undefined);
|
||||
this.#stopAdvisorRuntime();
|
||||
@@ -5135,15 +5211,16 @@ export class AgentSession {
|
||||
if (!options?.streamingBehavior) {
|
||||
throw new AgentBusyError();
|
||||
}
|
||||
// Steer/follow-up the keyword notices BEFORE the queued user message so the
|
||||
// model reads the steering notice ahead of the prompt it modifies.
|
||||
for (const notice of keywordNotices) {
|
||||
await this.sendCustomMessage(notice, { deliverAs: options.streamingBehavior });
|
||||
}
|
||||
if (options.streamingBehavior === "followUp") {
|
||||
await this.#queueUserMessage(expandedText, options?.images, "followUp");
|
||||
} else {
|
||||
await this.#queueUserMessage(expandedText, options?.images, "steer");
|
||||
}
|
||||
// Steer/follow-up the keyword notices alongside the queued user message.
|
||||
for (const notice of keywordNotices) {
|
||||
await this.sendCustomMessage(notice, { deliverAs: options.streamingBehavior });
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
@@ -5182,8 +5259,10 @@ export class AgentSession {
|
||||
await this.#promptWithMessage(message, expandedText, {
|
||||
...options,
|
||||
images: normalizedImages,
|
||||
prependMessages: preludeMessages.length > 0 ? preludeMessages : undefined,
|
||||
appendMessages: keywordNotices.length > 0 ? keywordNotices : undefined,
|
||||
prependMessages:
|
||||
preludeMessages.length > 0 || keywordNotices.length > 0
|
||||
? [...preludeMessages, ...keywordNotices]
|
||||
: undefined,
|
||||
});
|
||||
} finally {
|
||||
// Clean up residual eager-todo directive if the prompt never consumed it
|
||||
@@ -5222,13 +5301,13 @@ export class AgentSession {
|
||||
if (!options?.streamingBehavior) {
|
||||
throw new AgentBusyError();
|
||||
}
|
||||
for (const notice of keywordNotices) {
|
||||
await this.sendCustomMessage(notice, { 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 });
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -5244,7 +5323,7 @@ export class AgentSession {
|
||||
|
||||
await this.#promptWithMessage(customMessage, textContent, {
|
||||
...options,
|
||||
appendMessages: keywordNotices.length > 0 ? keywordNotices : undefined,
|
||||
prependMessages: keywordNotices.length > 0 ? keywordNotices : undefined,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -5253,7 +5332,6 @@ export class AgentSession {
|
||||
expandedText: string,
|
||||
options?: Pick<PromptOptions, "toolChoice" | "images" | "skipCompactionCheck"> & {
|
||||
prependMessages?: AgentMessage[];
|
||||
appendMessages?: AgentMessage[];
|
||||
skipPostPromptRecoveryWait?: boolean;
|
||||
},
|
||||
): Promise<void> {
|
||||
@@ -5320,12 +5398,6 @@ export class AgentSession {
|
||||
|
||||
messages.push(message);
|
||||
|
||||
// Inject the ultrathink notice (and any other per-turn appends) right after the
|
||||
// user message so the model reads it as part of the same turn.
|
||||
if (options?.appendMessages) {
|
||||
messages.push(...options.appendMessages);
|
||||
}
|
||||
|
||||
// Early bail-out: if a newer abort/prompt cycle started during setup,
|
||||
// return before mutating shared state (nextTurn messages, system prompt).
|
||||
if (this.#promptGeneration !== generation) {
|
||||
@@ -5647,12 +5719,24 @@ export class AgentSession {
|
||||
#canAutoContinueForFollowUp(): boolean {
|
||||
if (this.isStreaming) return false;
|
||||
if (this.isRetrying) return false;
|
||||
// A queued steer resumes from ANY tail: Agent.continue() runs #runLoop(undefined),
|
||||
// whose initial steering poll injects the steer before the first provider call, so the
|
||||
// request tail becomes the steer (valid) regardless of any injected custom / bashExecution
|
||||
// / pythonExecution record a user interrupt left as the literal transcript tail. This is
|
||||
// why a queued user steer stranded behind a preserved advisor card (or a flushed IRC aside
|
||||
// / eval execution record) still resumes — no tail-role enumeration needed.
|
||||
if (this.agent.peekSteeringQueue().length > 0) return true;
|
||||
// Follow-up-only auto-resume stays suppressed while a deliberate user interrupt is in effect
|
||||
// (#advisorAutoResumeSuppressed, cleared on the next user prompt): the user stopped, so their
|
||||
// queued follow-up waits for an explicit resume — even if an interleaving IRC wake turn has
|
||||
// since left a provider-valid tail.
|
||||
if (this.#advisorAutoResumeSuppressed) return false;
|
||||
// Follow-up-only resume has no steer to inject, so Agent.continue() continues from the
|
||||
// existing context tail — which must itself be a valid provider tail. An injected
|
||||
// non-conversational tail (advisor card → `developer`, bash/python execution) would make
|
||||
// the first model call invalid, so leave the follow-up queued for the next explicit resume.
|
||||
const messages = this.agent.state.messages;
|
||||
const last = messages[messages.length - 1];
|
||||
// A user interrupt during tool execution can leave the transcript ending
|
||||
// with the emitted tool result, not the aborted assistant message. Continuing
|
||||
// from that state is still resumable: Agent.continue() first polls queued
|
||||
// steering before making the next model call.
|
||||
return last?.role === "assistant" || last?.role === "toolResult";
|
||||
}
|
||||
|
||||
@@ -5901,7 +5985,9 @@ export class AgentSession {
|
||||
const followUpAll = this.agent.peekFollowUpQueue();
|
||||
const steering = steeringAll.filter(isUserQueuedMessage).map(toRestoredQueuedMessage);
|
||||
const followUp = followUpAll.filter(isUserQueuedMessage).map(toRestoredQueuedMessage);
|
||||
const keep: (m: AgentMessage) => boolean = options?.forInterrupt ? isAdvisorCard : m => !isUserQueuedMessage(m);
|
||||
const keep: (m: AgentMessage) => boolean = options?.forInterrupt
|
||||
? isAdvisorCard
|
||||
: m => !isUserQueuedMessage(m) && !isHiddenUserCompanion(m);
|
||||
this.agent.replaceQueues(steeringAll.filter(keep), followUpAll.filter(keep));
|
||||
return { steering, followUp };
|
||||
}
|
||||
@@ -5938,18 +6024,26 @@ export class AgentSession {
|
||||
}
|
||||
return -1;
|
||||
};
|
||||
// Notices queue immediately before their user message, so dropping the popped
|
||||
// prompt means also dropping the contiguous hidden-user companions right before
|
||||
// it — companions of other queued prompts stay put.
|
||||
const removeWithCompanions = (queue: readonly AgentMessage[], userIndex: number): AgentMessage[] => {
|
||||
let start = userIndex;
|
||||
while (start > 0 && isHiddenUserCompanion(queue[start - 1])) start--;
|
||||
const next = queue.slice();
|
||||
next.splice(start, userIndex - start + 1);
|
||||
return next;
|
||||
};
|
||||
const fromSteer = lastUserIndex(steering);
|
||||
if (fromSteer >= 0) {
|
||||
const nextSteer = steering.slice();
|
||||
const [removed] = nextSteer.splice(fromSteer, 1);
|
||||
this.agent.replaceQueues(nextSteer, followUp.slice());
|
||||
const removed = steering[fromSteer];
|
||||
this.agent.replaceQueues(removeWithCompanions(steering, fromSteer), followUp.slice());
|
||||
return toRestoredQueuedMessage(removed);
|
||||
}
|
||||
const fromFollowUp = lastUserIndex(followUp);
|
||||
if (fromFollowUp >= 0) {
|
||||
const nextFollowUp = followUp.slice();
|
||||
const [removed] = nextFollowUp.splice(fromFollowUp, 1);
|
||||
this.agent.replaceQueues(steering.slice(), nextFollowUp);
|
||||
const removed = followUp[fromFollowUp];
|
||||
this.agent.replaceQueues(steering.slice(), removeWithCompanions(followUp, fromFollowUp));
|
||||
return toRestoredQueuedMessage(removed);
|
||||
}
|
||||
return undefined;
|
||||
@@ -10242,11 +10336,8 @@ export class AgentSession {
|
||||
if (autoReply) void this.#runIrcAutoReply(msg);
|
||||
return "injected";
|
||||
}
|
||||
// Idle: same wake primitive the yield queue uses for async-result
|
||||
// delivery — prompt the agent directly so a real turn runs.
|
||||
this.agent.prompt(record).catch(error => {
|
||||
logger.warn("IRC wake turn failed", { from: msg.from, to: msg.to, error: String(error) });
|
||||
});
|
||||
// Idle: wake a real turn so the recipient responds (shared with the stranded-aside resume).
|
||||
this.#wakeForIrc([record]);
|
||||
return "woken";
|
||||
}
|
||||
|
||||
|
||||
@@ -4,13 +4,18 @@
|
||||
* so they re-enter context when the user resumes. Internal (non-user) aborts keep
|
||||
* the prior behavior — advisor advice stays in the auto-continue path.
|
||||
*
|
||||
* Three seams:
|
||||
* Five seams:
|
||||
* 1. A concern already steered into the agent queue when the user hits Esc is
|
||||
* pulled out of the post-abort auto-continue path and re-recorded as advice.
|
||||
* 2. A concern parked hidden (#pendingNextTurnMessages) by the suppressed
|
||||
* delivery while the turn is still tearing down is reclaimed once idle.
|
||||
* 3. A non-user abort does NOT suppress: a steered advisor card still drives the
|
||||
* auto-continue, so the gate is keyed to the user interrupt, not any abort.
|
||||
* 4. A user message queued (as a steer) before the interrupt is delivered on
|
||||
* resume even though the preserved advisor card is the trailing message.
|
||||
* 5. The same queued as a follow-up: continuing from the preserved advisor card
|
||||
* (which converts to `developer`) would send an invalid provider tail, so the
|
||||
* follow-up stays queued for the next explicit resume rather than auto-running.
|
||||
*/
|
||||
import { afterEach, beforeEach, describe, expect, it } from "bun:test";
|
||||
import * as fs from "node:fs";
|
||||
@@ -21,6 +26,7 @@ import { createMockModel, type MockModel, type MockResponse } from "@oh-my-pi/pi
|
||||
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 type { IrcMessage } from "@oh-my-pi/pi-coding-agent/irc/bus";
|
||||
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 { USER_INTERRUPT_LABEL } from "@oh-my-pi/pi-coding-agent/session/messages";
|
||||
@@ -102,6 +108,20 @@ describe("AgentSession advisor auto-resume suppression", () => {
|
||||
return message.role === "custom" && (message as { customType?: string }).customType === ADVISOR_TYPE;
|
||||
}
|
||||
|
||||
function userMessageText(messages: AgentMessage[]): string[] {
|
||||
const out: string[] = [];
|
||||
for (const message of messages) {
|
||||
if (message.role !== "user") continue;
|
||||
const content = message.content;
|
||||
if (typeof content === "string") {
|
||||
out.push(content);
|
||||
} else {
|
||||
for (const part of content) if (part.type === "text") out.push(part.text);
|
||||
}
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
function capturePersistedAdvice(sessionManager: SessionManager): string[] {
|
||||
const persisted: string[] = [];
|
||||
sessionManager.onEntryAppended = entry => {
|
||||
@@ -178,4 +198,119 @@ describe("AgentSession advisor auto-resume suppression", () => {
|
||||
expect(session.agent.peekSteeringQueue()).toEqual([]);
|
||||
expect(mock.calls.length).toBe(2);
|
||||
});
|
||||
|
||||
it("resumes a queued user steer stranded behind a preserved advisor card", async () => {
|
||||
// Reported bug: typing a message during a run (queued as a steer) then pressing
|
||||
// enter again (empty-submit interrupt) recorded the advisor card but stranded the
|
||||
// user message — nothing resumed the run. The preserved advisor card is the
|
||||
// trailing `custom` message, which #canAutoContinueForFollowUp must look past.
|
||||
const { session, sessionManager, mock, streamStarted } = await createParkedSession([
|
||||
{ content: ["resumed on the steer"] },
|
||||
]);
|
||||
const persisted = capturePersistedAdvice(sessionManager);
|
||||
|
||||
const running = session.prompt("do the thing");
|
||||
await streamStarted;
|
||||
|
||||
await session.sendCustomMessage(advisorCard("breaks the build"), { deliverAs: "steer", triggerTurn: true });
|
||||
await session.prompt("also rename the helper", { streamingBehavior: "steer" });
|
||||
|
||||
await session.abort({ reason: USER_INTERRUPT_LABEL });
|
||||
await session.waitForIdle();
|
||||
await running.catch(() => {});
|
||||
|
||||
// Advisor card preserved as a visible/persisted card AND the user steer delivered
|
||||
// in exactly one resume turn (no spurious extra call).
|
||||
expect(session.agent.state.messages.filter(isAdvisorCard)).toHaveLength(1);
|
||||
expect(persisted).toEqual(["breaks the build"]);
|
||||
expect(session.agent.peekSteeringQueue()).toEqual([]);
|
||||
expect(mock.calls.length).toBe(2);
|
||||
expect(userMessageText(session.agent.state.messages)).toContain("also rename the helper");
|
||||
});
|
||||
|
||||
it("leaves a queued user follow-up queued behind a preserved advisor card", async () => {
|
||||
// Same stranding, but the user message was queued as a follow-up (Ctrl+Enter).
|
||||
// Only steering resumes safely behind a preserved advisor card: agentLoopContinue
|
||||
// injects steering before the next model call, keeping the request tail valid. A
|
||||
// follow-up would instead resume by continuing from the advisor card (which converts
|
||||
// to `developer`) as the tail — a provider-invalid request — so it is NOT auto-run.
|
||||
// It stays queued for the next explicit user resume.
|
||||
const { session, sessionManager, mock, streamStarted } = await createParkedSession([
|
||||
{ content: ["must not run"] },
|
||||
]);
|
||||
const persisted = capturePersistedAdvice(sessionManager);
|
||||
|
||||
const running = session.prompt("do the thing");
|
||||
await streamStarted;
|
||||
|
||||
await session.sendCustomMessage(advisorCard("missing a guard"), { deliverAs: "steer", triggerTurn: true });
|
||||
await session.prompt("then add the test", { streamingBehavior: "followUp" });
|
||||
|
||||
await session.abort({ reason: USER_INTERRUPT_LABEL });
|
||||
await session.waitForIdle();
|
||||
await running.catch(() => {});
|
||||
|
||||
// Advisor preserved as a visible/persisted card; the follow-up stays queued and
|
||||
// drives no resume (only the original, aborted turn ever called the model).
|
||||
expect(session.agent.state.messages.filter(isAdvisorCard)).toHaveLength(1);
|
||||
expect(persisted).toEqual(["missing a guard"]);
|
||||
expect(userMessageText([...session.agent.peekFollowUpQueue()])).toContain("then add the test");
|
||||
expect(userMessageText(session.agent.state.messages)).not.toContain("then add the test");
|
||||
expect(mock.calls.length).toBe(1);
|
||||
});
|
||||
|
||||
it("wakes a turn for an IRC aside stranded across a user interrupt", async () => {
|
||||
const { session, mock, streamStarted } = await createParkedSession([{ content: ["replying to peer"] }]);
|
||||
const running = session.prompt("do the thing");
|
||||
await streamStarted;
|
||||
// IRC arrives mid-turn → queued as a non-interrupting aside.
|
||||
await session.deliverIrcMessage({ id: "m1", from: "peer", to: "me", body: "ping", ts: Date.now() } as IrcMessage);
|
||||
// The user interrupt skips the loop's final aside poll, stranding the aside with no loop to
|
||||
// drain it. The settle drain must wake a turn so the peer still gets a response.
|
||||
await session.abort({ reason: USER_INTERRUPT_LABEL });
|
||||
await session.waitForIdle();
|
||||
await running.catch(() => {});
|
||||
|
||||
const sawIrc = session.agent.state.messages.some(
|
||||
m => m.role === "custom" && (m as { customType?: string }).customType === "irc:incoming",
|
||||
);
|
||||
expect(sawIrc).toBe(true);
|
||||
expect(mock.calls.length).toBe(2);
|
||||
});
|
||||
|
||||
it("flushes an accepted IRC aside on dispose instead of dropping it", async () => {
|
||||
const { session, streamStarted } = await createParkedSession();
|
||||
const running = session.prompt("do the thing");
|
||||
await streamStarted;
|
||||
await session.deliverIrcMessage({ id: "m1", from: "peer", to: "me", body: "ping", ts: Date.now() } as IrcMessage);
|
||||
// Dispose mid-flight persists the accepted aside to the transcript rather than dropping it.
|
||||
session.beginDispose();
|
||||
const sawIrc = session.agent.state.messages.some(
|
||||
m => m.role === "custom" && (m as { customType?: string }).customType === "irc:incoming",
|
||||
);
|
||||
expect(sawIrc).toBe(true);
|
||||
running.catch(() => {});
|
||||
});
|
||||
|
||||
it("responds to a stranded IRC aside while keeping a blocked follow-up queued", async () => {
|
||||
const { session, mock, streamStarted } = await createParkedSession([{ content: ["replying to peer"] }]);
|
||||
const running = session.prompt("do the thing");
|
||||
await streamStarted;
|
||||
// The user queues a follow-up (Ctrl+Enter) and an IRC ping lands as an aside...
|
||||
await session.prompt("then add the test", { streamingBehavior: "followUp" });
|
||||
await session.deliverIrcMessage({ id: "m2", from: "peer", to: "me", body: "ping", ts: Date.now() } as IrcMessage);
|
||||
// ...then the user interrupts. The IRC must still get a response, but the user's queued
|
||||
// follow-up must NOT auto-run (seam #5) even though the IRC wake turn leaves a valid tail.
|
||||
await session.abort({ reason: USER_INTERRUPT_LABEL });
|
||||
await session.waitForIdle();
|
||||
await running.catch(() => {});
|
||||
|
||||
const sawIrc = session.agent.state.messages.some(
|
||||
m => m.role === "custom" && (m as { customType?: string }).customType === "irc:incoming",
|
||||
);
|
||||
expect(sawIrc).toBe(true);
|
||||
expect(userMessageText([...session.agent.peekFollowUpQueue()])).toContain("then add the test");
|
||||
expect(userMessageText(session.agent.state.messages)).not.toContain("then add the test");
|
||||
expect(mock.calls.length).toBe(2);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -117,4 +117,20 @@ describe("AgentSession magic keyword settings", () => {
|
||||
expect(session.thinkingLevel).toBe(Effort.Low);
|
||||
expect(session.autoResolvedThinkingLevel()).toBe(Effort.Low);
|
||||
});
|
||||
|
||||
it("queues the magic-keyword notice before the user message", async () => {
|
||||
const created = await createMagicKeywordSession(root);
|
||||
session = created.session;
|
||||
authStorage = created.authStorage;
|
||||
const promptSpy = vi.spyOn(session.agent, "prompt").mockResolvedValue(undefined);
|
||||
|
||||
await session.prompt("ultrathink do the thing");
|
||||
|
||||
const promptMessages = promptSpy.mock.calls[0]![0] as unknown as Array<{ role?: string; customType?: string }>;
|
||||
const noticeIdx = promptMessages.findIndex(m => m.customType === "ultrathink-notice");
|
||||
const userIdx = promptMessages.findIndex(m => m.role === "user");
|
||||
expect(noticeIdx).toBeGreaterThanOrEqual(0);
|
||||
expect(userIdx).toBeGreaterThanOrEqual(0);
|
||||
expect(noticeIdx).toBeLessThan(userIdx);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -205,4 +205,95 @@ describe("AgentSession queued steer delivery", () => {
|
||||
expect(session.agent.hasQueuedMessages()).toBe(false);
|
||||
expect(session.getQueuedMessages().steering).toEqual([]);
|
||||
});
|
||||
|
||||
it("dequeuing an ultrathink prompt mid-stream restores the text and drops its companion notice", async () => {
|
||||
const { session } = await createSession([{ content: ["host answer"] }]);
|
||||
let queuedShape: string[] | undefined;
|
||||
let clearedSteering: unknown;
|
||||
let hasQueuedAfterClear: boolean | undefined;
|
||||
let injected = false;
|
||||
session.agent.setOnBeforeYield(async () => {
|
||||
if (injected) return;
|
||||
injected = true;
|
||||
// Real path: a magic-keyword prompt steered mid-stream enqueues the hidden
|
||||
// notice immediately before the user message.
|
||||
await session.prompt("ultrathink fix it", { streamingBehavior: "steer" });
|
||||
queuedShape = session.agent.peekSteeringQueue().map(m => (m.role === "custom" ? m.customType : m.role));
|
||||
// Alt+Up restore mid-flight: only the user's text returns; the companion
|
||||
// notice must not be left orphaned in the queue.
|
||||
const cleared = session.clearQueue();
|
||||
clearedSteering = cleared.steering;
|
||||
hasQueuedAfterClear = session.agent.hasQueuedMessages();
|
||||
});
|
||||
|
||||
await session.prompt("hello");
|
||||
|
||||
expect(queuedShape).toEqual(["ultrathink-notice", "user"]);
|
||||
expect(clearedSteering).toEqual([{ text: "ultrathink fix it", images: undefined }]);
|
||||
expect(hasQueuedAfterClear).toBe(false);
|
||||
});
|
||||
|
||||
it("a fresh user prompt delivers queued steer and follow-up work", async () => {
|
||||
const { session } = await createSession([{ content: ["one"] }, { content: ["two"] }, { content: ["three"] }]);
|
||||
// Queue real pending work before the user's next send.
|
||||
session.agent.steer({
|
||||
role: "user",
|
||||
content: [{ type: "text", text: "queued steer" }],
|
||||
steering: true,
|
||||
attribution: "user",
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
session.agent.followUp({
|
||||
role: "user",
|
||||
content: [{ type: "text", text: "queued follow-up" }],
|
||||
attribution: "user",
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
expect(session.agent.hasQueuedMessages()).toBe(true);
|
||||
|
||||
await session.prompt("hello");
|
||||
await session.waitForIdle();
|
||||
|
||||
// Sending a fresh prompt is the opportunity to drain everything: the steer folds
|
||||
// into the new turn and the follow-up runs as its continuation — nothing stranded.
|
||||
const userTexts = session.agent.state.messages
|
||||
.filter(message => message.role === "user")
|
||||
.map(message =>
|
||||
typeof message.content === "string"
|
||||
? message.content
|
||||
: message.content
|
||||
.filter(part => part.type === "text")
|
||||
.map(part => part.text)
|
||||
.join(""),
|
||||
);
|
||||
expect(userTexts).toContain("hello");
|
||||
expect(userTexts).toContain("queued steer");
|
||||
expect(userTexts).toContain("queued follow-up");
|
||||
expect(session.agent.hasQueuedMessages()).toBe(false);
|
||||
});
|
||||
|
||||
it("resumes a queued steer left behind a non-advisor custom transcript tail", async () => {
|
||||
const { session } = await createSession([{ content: ["first answer"] }, { content: ["resumed"] }]);
|
||||
await session.prompt("first");
|
||||
// A non-advisor custom (e.g. a flushed irc:incoming aside) is the literal transcript tail.
|
||||
// A queued steer must resume regardless of tail role — Agent.continue injects it via the
|
||||
// initial steering poll — so the old advisor-only look-back can no longer strand it.
|
||||
const aside = {
|
||||
role: "custom" as const,
|
||||
customType: "irc:incoming",
|
||||
content: "peer pinged you",
|
||||
display: true,
|
||||
attribution: "agent" as const,
|
||||
timestamp: Date.now(),
|
||||
};
|
||||
session.agent.emitExternalEvent({ type: "message_start", message: aside });
|
||||
session.agent.emitExternalEvent({ type: "message_end", message: aside });
|
||||
|
||||
const delivered = nextUserMessage(session, "resume me");
|
||||
await session.steer("resume me");
|
||||
await delivered;
|
||||
await session.waitForIdle();
|
||||
|
||||
expect(session.agent.peekSteeringQueue()).toEqual([]);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -249,6 +249,30 @@ function queueAdvisorSteer(session: AgentSession, note = "consider X"): void {
|
||||
});
|
||||
}
|
||||
|
||||
/** Mirror a hidden magic-keyword companion notice (`display:false`, `attribution:"user"`). */
|
||||
function queueMagicCompanion(session: AgentSession, customType = "ultrathink-notice"): void {
|
||||
session.agent.steer({
|
||||
role: "custom",
|
||||
customType,
|
||||
content: "hidden notice",
|
||||
display: false,
|
||||
attribution: "user",
|
||||
details: {},
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
}
|
||||
|
||||
/** Mirror a steered user prompt (`AgentSession.#queueUserMessage(..., "steer")`). */
|
||||
function queueUserSteer(session: AgentSession, text: string): void {
|
||||
session.agent.steer({
|
||||
role: "user",
|
||||
content: [{ type: "text", text }],
|
||||
steering: true,
|
||||
attribution: "user",
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
}
|
||||
|
||||
describe("AgentSession derived queued custom display", () => {
|
||||
let fixture: SessionFixture | undefined;
|
||||
|
||||
@@ -379,6 +403,39 @@ describe("AgentSession derived queued custom display", () => {
|
||||
expect(remaining).toHaveLength(1);
|
||||
expect(remaining[0]).toMatchObject({ customType: "advisor" });
|
||||
});
|
||||
|
||||
it("clearQueue drops a queued magic-keyword companion with its dequeued user prompt", async () => {
|
||||
fixture = await createRealSession();
|
||||
const { session } = fixture;
|
||||
// Queue order mirrors prompt("ultrathink do X", { streamingBehavior: "steer" }):
|
||||
// the hidden companion notice queues right before the user message.
|
||||
queueMagicCompanion(session, "ultrathink-notice");
|
||||
queueUserSteer(session, "ultrathink do X");
|
||||
|
||||
// The companion is display:false, so only the user prompt is displayable work.
|
||||
expect(session.queuedMessageCount).toBe(1);
|
||||
|
||||
// Alt+Up bulk restore returns the user's text and leaves no orphaned companion.
|
||||
const cleared = session.clearQueue();
|
||||
expect(cleared.steering).toEqual([{ text: "ultrathink do X", images: undefined }]);
|
||||
expect(session.agent.hasQueuedMessages()).toBe(false);
|
||||
});
|
||||
|
||||
it("popLastQueuedMessage drops only the popped prompt's preceding companion", async () => {
|
||||
fixture = await createRealSession();
|
||||
const { session } = fixture;
|
||||
// [ultrathink-notice, "first", orchestrate-notice, "second"].
|
||||
queueMagicCompanion(session, "ultrathink-notice");
|
||||
queueUserSteer(session, "first");
|
||||
queueMagicCompanion(session, "orchestrate-notice");
|
||||
queueUserSteer(session, "second");
|
||||
|
||||
expect(session.popLastQueuedMessage()?.text).toBe("second");
|
||||
// Only the popped prompt's companion (orchestrate-notice) leaves; the earlier
|
||||
// prompt and its own companion stay intact.
|
||||
const remaining = session.agent.peekSteeringQueue();
|
||||
expect(remaining.map(m => (m.role === "custom" ? m.customType : m.role))).toEqual(["ultrathink-notice", "user"]);
|
||||
});
|
||||
});
|
||||
|
||||
function createStubInteractiveModeContextForUiHelpers(session: AgentSession) {
|
||||
|
||||
@@ -582,7 +582,9 @@ describe("IRC", () => {
|
||||
});
|
||||
expect(outcome).toBe("woken");
|
||||
expect(promptSpy).toHaveBeenCalledTimes(1);
|
||||
const prompted = promptSpy.mock.calls[0]?.[0] as unknown as CustomMessage;
|
||||
// The idle wake routes through #wakeForIrc, which batches records into one prompt —
|
||||
// even a lone incoming message is delivered as a one-element array.
|
||||
const prompted = (promptSpy.mock.calls[0]?.[0] as unknown as CustomMessage[])[0];
|
||||
expect(prompted).toMatchObject({ role: "custom", customType: "irc:incoming" });
|
||||
expect(prompted.details).toMatchObject({ id: "msg-1", from: "0-Peer", message: "wake up" });
|
||||
|
||||
|
||||
Reference in New Issue
Block a user