feat(coding-agent): added coding-agent follow-up queue with onBeforeYield

- Added optional `onBeforeYield` configuration and `setOnBeforeYield` in Agent, executed before follow-up checks.
- Added `YieldQueue` to `AgentSession`, with setup/teardown and streaming/idle flush via `setOnBeforeYield`.
- Replaced immediate async-result follow-up dispatch with queued batch entries, including stale-state suppression.
- Added MCP follow-up queueing in SDK, deduplicating updates by `serverName` and `uri`.
- Added changelog entries for `onBeforeYield`, async-result batching, MCP dedupe, and `display.shimmer` modes.
- Added yield queue unit tests for streaming emission, debounced idle batches, stale filtering, and error isolation.
This commit is contained in:
can1357
2026-05-22 13:05:30 +09:00
committed by Can Bölük
parent e39b48c405
commit 796f963da9
13 changed files with 690 additions and 36 deletions
+3
View File
@@ -1,6 +1,9 @@
# Changelog
## [Unreleased]
### Added
- Added `onBeforeYield` hook support so user code can run right before the agent loop checks for follow-up messages
## [15.1.3] - 2026-05-17
### Added
+1
View File
@@ -589,6 +589,7 @@ async function runLoopBody(
}
// Agent would stop here. Check for follow-up messages.
await config.onBeforeYield?.();
const followUpMessages = (await config.getFollowUpMessages?.()) || [];
if (followUpMessages.length > 0) {
// Set as pending so inner loop processes them
+6
View File
@@ -290,6 +290,7 @@ export class Agent {
#onSseEvent?: SimpleStreamOptions["onSseEvent"];
#onAssistantMessageEvent?: (message: AssistantMessage, event: AssistantMessageEvent) => void;
#onHarmonyLeak?: (event: HarmonyAuditEvent) => void | Promise<void>;
#onBeforeYield?: () => Promise<void> | void;
#telemetry?: AgentLoopConfig["telemetry"];
/** Buffered Cursor tool results with text length at time of call (for correct ordering) */
@@ -559,6 +560,10 @@ export class Agent {
this.#onAssistantMessageEvent = fn;
}
setOnBeforeYield(fn: (() => Promise<void> | void) | undefined): void {
this.#onBeforeYield = fn;
}
emitExternalEvent(event: AgentEvent) {
switch (event.type) {
case "message_start":
@@ -934,6 +939,7 @@ export class Agent {
return this.#dequeueSteeringMessages();
},
getFollowUpMessages: async () => this.#dequeueFollowUpMessages(),
onBeforeYield: () => this.#onBeforeYield?.(),
telemetry: this.#telemetry,
};
+7
View File
@@ -122,6 +122,13 @@ export interface AgentLoopConfig extends SimpleStreamOptions {
* continues with another turn.
*/
getFollowUpMessages?: () => Promise<AgentMessage[]>;
/**
* Hook fired right before the loop would exit.
*
* Called when the agent has no more tool calls and no steering messages,
* immediately before polling follow-up messages.
*/
onBeforeYield?: () => Promise<void> | void;
/**
* Provides tool execution context, resolved per tool call.
+33
View File
@@ -919,4 +919,37 @@ describe("agentLoopContinue with AgentMessage", () => {
expect(JSON.stringify(toolEnd.result)).toContain("hook exploded");
}
});
it("runs onBeforeYield before polling follow-up messages", async () => {
const context: AgentContext = {
systemPrompt: ["You are helpful."],
messages: [],
tools: [],
};
const queuedFollowUps: AgentMessage[] = [];
let hookCalls = 0;
const mock = createMockModel({
responses: [{ content: ["first"] }, { content: ["second"] }],
});
const config: AgentLoopConfig = {
model: mock.model,
convertToLlm: identityConverter,
onBeforeYield: () => {
hookCalls++;
if (hookCalls === 1) {
queuedFollowUps.push(createUserMessage("follow-up"));
}
},
getFollowUpMessages: async () => queuedFollowUps.splice(0),
};
const stream = agentLoop([createUserMessage("initial")], context, config, undefined, mock.stream);
for await (const _ of stream) {
// drain
}
const messages = await stream.result();
expect(hookCalls).toBe(2);
expect(messages.map(message => message.role)).toEqual(["user", "assistant", "user", "assistant"]);
expect(messages[2]).toMatchObject({ role: "user", content: "follow-up" });
});
});
+5
View File
@@ -1,6 +1,7 @@
# Changelog
## [Unreleased]
### Breaking Changes
- Changed PR and task-isolation worktree directory layout to hash-based `~/.omp/wt/<identifier>-<path-hash>` style paths, replacing the previous nested encoded-repo layout
@@ -10,11 +11,15 @@
- Added `omp worktree` command (alias `wt`) to list and manage agent-managed worktrees under `~/.omp/wt`
- Added `omp worktree clear` to remove orphaned worktree directories, with `--all` to include live PR-checkouts, `--dry-run` for preview, and `--json` reporting
- Added machine-readable JSON output to `omp worktree list` for scripted inspection
- Added `display.shimmer` appearance setting with `classic`, `kitt` (Knight Rider K.I.T.T. scanner), and `disabled` modes
### Changed
- Changed background job completion follow-ups to batch multiple finished jobs into a single `async-result` message, showing each completed job and its result in one place
- Changed MCP notification follow-ups to combine multiple resource updates into a single consolidated message and suppress duplicate server/uri entries
- Updated PR checkout to reuse `hashPath`-based worktree roots when creating and scanning worktrees for cleanup
- Updated `worktree` cleanup logic to gracefully prune parent git metadata after removing worktree directories
- Reworked working-message shimmer animation for 60fps rendering: ANSI sequences are coalesced per same-tier run instead of emitted per code point, palettes compile once and cache per active theme, and the band position is now fractional so motion is smooth at any frame rate
### Fixed
@@ -113,21 +113,39 @@ export class UiHelpers {
type?: "bash" | "task";
label?: string;
durationMs?: number;
jobs?: Array<{
jobId?: string;
type?: "bash" | "task";
label?: string;
durationMs?: number;
}>;
}>
).details;
const jobId = details?.jobId ?? "unknown";
const typeLabel = details?.type ? `[${details.type}]` : "[job]";
const duration =
typeof details?.durationMs === "number" ? formatDuration(details.durationMs) : undefined;
const line = [
theme.fg("success", `${theme.status.success} Background job completed`),
theme.fg("dim", typeLabel),
theme.fg("accent", jobId),
duration ? theme.fg("dim", `(${duration})`) : undefined,
]
.filter(Boolean)
.join(" ");
this.ctx.chatContainer.addChild(new Text(line, 1, 0));
const jobs =
details?.jobs && details.jobs.length > 0
? details.jobs
: [
{
jobId: details?.jobId,
type: details?.type,
label: details?.label,
durationMs: details?.durationMs,
},
];
for (const job of jobs) {
const jobId = job.jobId ?? "unknown";
const typeLabel = job.type ? `[${job.type}]` : "[job]";
const duration = typeof job.durationMs === "number" ? formatDuration(job.durationMs) : undefined;
const line = [
theme.fg("success", `${theme.status.success} Background job completed`),
theme.fg("dim", typeLabel),
theme.fg("accent", jobId),
duration ? theme.fg("dim", `(${duration})`) : undefined,
]
.filter(Boolean)
.join(" ");
this.ctx.chatContainer.addChild(new Text(line, 1, 0));
}
break;
}
if (message.customType === SKILL_PROMPT_MESSAGE_TYPE) {
@@ -1,5 +1,8 @@
<system-notice>
Background job {{jobId}} has completed. Resume your work using the result below.
{{#if multiple}}{{jobs.length}} background jobs have completed. Resume your work using the results below.
{{result}}
{{else}}Background job {{jobs.[0].jobId}} has completed. Resume your work using the result below.
{{/if}}{{#each jobs}}{{#if @root.multiple}}── Job {{this.jobId}}{{#if this.label}} ({{this.label}}){{/if}} ──
{{/if}}{{this.result}}{{#unless @last}}
{{/unless}}{{/each}}
</system-notice>
+95 -21
View File
@@ -31,7 +31,7 @@ import {
Snowflake,
} from "@oh-my-pi/pi-utils";
import chalk from "chalk";
import { AsyncJobManager, isBackgroundJobSupportEnabled } from "./async";
import { type AsyncJob, AsyncJobManager, isBackgroundJobSupportEnabled } from "./async";
import { createAutoresearchExtension } from "./autoresearch";
import { loadCapability } from "./capability";
import { type Rule, ruleCapability, setActiveRules } from "./capability/rule";
@@ -101,7 +101,7 @@ import {
import { AgentSession } from "./session/agent-session";
import { resolveAuthBrokerConfig } from "./session/auth-broker-config";
import { AuthBrokerClient, AuthStorage, RemoteAuthCredentialStore } from "./session/auth-storage";
import { convertToLlm } from "./session/messages";
import { type CustomMessage, convertToLlm } from "./session/messages";
import { SessionManager } from "./session/session-manager";
import { closeAllConnections } from "./ssh/connection-manager";
import { unmountAll } from "./ssh/sshfs-mount";
@@ -152,6 +152,83 @@ import { EventBus } from "./utils/event-bus";
import { buildNamedToolChoice } from "./utils/tool-choice";
import { buildWorkspaceTree, type WorkspaceTree } from "./workspace-tree";
type AsyncResultEntry = {
jobId: string;
result: string;
job: AsyncJob | undefined;
durationMs: number | undefined;
};
type AsyncResultJobDetails = {
jobId: string;
type?: "bash" | "task";
label?: string;
durationMs?: number;
};
type AsyncResultDetails = {
jobs: AsyncResultJobDetails[];
};
type McpNotificationEntry = {
serverName: string;
uri: string;
};
function buildAsyncResultBatchMessage(entries: AsyncResultEntry[]): CustomMessage<AsyncResultDetails> | null {
if (entries.length === 0) return null;
const jobs = entries.map(entry => ({
jobId: entry.jobId,
result: entry.result,
type: entry.job?.type,
label: entry.job?.label,
durationMs: entry.durationMs,
}));
const details: AsyncResultDetails = {
jobs: jobs.map(job => ({
jobId: job.jobId,
type: job.type,
label: job.label,
durationMs: job.durationMs,
})),
};
return {
role: "custom",
customType: "async-result",
content: prompt.render(asyncResultTemplate, {
multiple: jobs.length > 1,
jobs,
}),
display: true,
attribution: "agent",
details,
timestamp: Date.now(),
};
}
function buildMcpNotificationBatchMessage(entries: McpNotificationEntry[]): AgentMessage | null {
const resources: McpNotificationEntry[] = [];
const seen = new Set<string>();
for (const entry of entries) {
const key = `${entry.serverName}\0${entry.uri}`;
if (seen.has(key)) continue;
seen.add(key);
resources.push(entry);
}
if (resources.length === 0) return null;
const lines = [`[MCP notification] ${resources.length} resource(s) updated:`];
for (const resource of resources) {
lines.push(`- server="${resource.serverName}" uri=${resource.uri}`);
}
lines.push('Use read(path="mcp://<uri>") to inspect if relevant.');
return {
role: "user",
content: [{ type: "text", text: lines.join("\n") }],
attribution: "agent",
timestamp: Date.now(),
};
}
// Types
export interface CreateAgentSessionOptions {
/** Working directory for project-local discovery. Default: getProjectDir() */
@@ -1035,23 +1112,13 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
const formattedResult = await formatAsyncResultForFollowUp(result);
if (asyncJobManager!.isDeliverySuppressed(jobId)) return;
const message = prompt.render(asyncResultTemplate, { jobId, result: formattedResult });
const durationMs = job ? Math.max(0, Date.now() - job.startTime) : undefined;
await session.sendCustomMessage(
{
customType: "async-result",
content: message,
display: true,
attribution: "agent",
details: {
jobId,
type: job?.type,
label: job?.label,
durationMs,
},
},
{ deliverAs: "followUp", triggerTurn: true },
);
session.yieldQueue.enqueue<AsyncResultEntry>("async-result", {
jobId,
result: formattedResult,
job,
durationMs,
});
},
})
: undefined;
@@ -1902,6 +1969,15 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
providerSessionId: options.providerSessionId,
});
hasSession = true;
if (asyncJobManager) {
session.yieldQueue.register<AsyncResultEntry>("async-result", {
isStale: entry => asyncJobManager.isDeliverySuppressed(entry.jobId),
build: buildAsyncResultBatchMessage,
});
}
session.yieldQueue.register<McpNotificationEntry>("mcp-notification", {
build: buildMcpNotificationBatchMessage,
});
// Attach the live session to the pre-registered ref so peers can route IRC
// messages here. Refresh sessionFile in case it was unavailable at pre-register
@@ -2036,9 +2112,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
notificationDebounceTimers.delete(key);
// Re-check: user may have disabled notifications during the debounce window
if (!settings.get("mcp.notifications")) return;
void session.followUp(
`[MCP notification] Server "${serverName}" reports resource \`${uri}\` was updated. Use read(path="mcp://${uri}") to inspect if relevant.`,
);
session.yieldQueue.enqueue<McpNotificationEntry>("mcp-notification", { serverName, uri });
}, debounceMs),
);
});
@@ -205,6 +205,7 @@ import type {
} from "./session-manager";
import { getLatestCompactionEntry } from "./session-manager";
import { ToolChoiceQueue } from "./tool-choice-queue";
import { YieldQueue } from "./yield-queue";
/** Session-specific events that extend the core AgentEvent */
export type AgentSessionEvent =
@@ -735,6 +736,7 @@ export class AgentSession {
readonly agent: Agent;
readonly sessionManager: SessionManager;
readonly settings: Settings;
readonly yieldQueue: YieldQueue;
#powerAssertion: MacOSPowerAssertion | undefined;
@@ -1031,6 +1033,24 @@ export class AgentSession {
};
this.agent.setProviderResponseInterceptor(this.#onResponse);
this.agent.setRawSseEventInterceptor(this.#onSseEvent);
this.yieldQueue = new YieldQueue({
isStreaming: () => this.isStreaming,
injectStreaming: message => this.agent.followUp(message),
injectIdle: async messages => {
const first = messages[0];
if (!first) return;
await this.agent.prompt(messages.length === 1 ? first : messages);
},
scheduleIdleFlush: run => {
this.#schedulePostPromptTask(
async () => {
await run();
},
{ delayMs: 1 },
);
},
});
this.agent.setOnBeforeYield(() => this.yieldQueue.flush("streaming"));
this.#convertToLlm = config.convertToLlm ?? convertToLlm;
this.#rebuildSystemPrompt = config.rebuildSystemPrompt;
this.#getMcpServerInstructions = config.getMcpServerInstructions;
@@ -2720,6 +2740,8 @@ export class AgentSession {
async dispose(): Promise<void> {
this.#isDisposed = true;
this.#pendingBackgroundExchanges = [];
this.yieldQueue.clear();
this.agent.setOnBeforeYield(undefined);
this.#evalExecutionDisposing = true;
try {
if (this.#extensionRunner?.hasHandlers("session_shutdown")) {
@@ -0,0 +1,155 @@
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { logger } from "@oh-my-pi/pi-utils";
export interface YieldDispatcher<P> {
/** Drop entries already delivered through another path. Called per-entry at flush time. */
isStale?(entry: P): boolean;
/** Produce one batched AgentMessage from non-stale entries. Return null to skip. */
build(survivors: P[]): AgentMessage | null;
}
export interface YieldQueueOptions {
isStreaming: () => boolean;
injectStreaming(msg: AgentMessage): void;
injectIdle(messages: AgentMessage[]): Promise<void>;
scheduleIdleFlush(run: () => Promise<void>): void;
}
type YieldFlushMode = "streaming" | "idle";
interface StoredDispatcher {
isStale?: (entry: unknown) => boolean;
build: (survivors: unknown[]) => AgentMessage | null;
}
function formatError(error: unknown): string {
return error instanceof Error ? error.message : String(error);
}
export class YieldQueue {
readonly #options: YieldQueueOptions;
readonly #dispatchers = new Map<string, StoredDispatcher>();
readonly #entries = new Map<string, unknown[]>();
#idleFlushPending = false;
constructor(options: YieldQueueOptions) {
this.#options = options;
}
register<P>(kind: string, dispatcher: YieldDispatcher<P>): () => void {
const stored: StoredDispatcher = {
...(dispatcher.isStale ? { isStale: entry => dispatcher.isStale?.(entry as P) ?? false } : {}),
build: survivors => dispatcher.build(survivors as P[]),
};
this.#dispatchers.set(kind, stored);
return () => {
if (this.#dispatchers.get(kind) !== stored) return;
this.#dispatchers.delete(kind);
this.#entries.delete(kind);
};
}
enqueue<P>(kind: string, entry: P): void {
if (!this.#dispatchers.has(kind)) {
logger.warn("Yield queue entry ignored for unregistered kind", { kind });
return;
}
let entries = this.#entries.get(kind);
if (!entries) {
entries = [];
this.#entries.set(kind, entries);
}
entries.push(entry);
if (!this.#options.isStreaming()) {
this.#scheduleIdleFlush();
}
}
has(kind?: string): boolean {
if (kind !== undefined) return (this.#entries.get(kind)?.length ?? 0) > 0;
for (const entries of this.#entries.values()) {
if (entries.length > 0) return true;
}
return false;
}
async flush(mode: YieldFlushMode): Promise<void> {
if (mode === "idle") {
this.#idleFlushPending = false;
}
const idleMessages: AgentMessage[] = [];
for (const [kind, dispatcher] of this.#dispatchers) {
const entries = this.#drain(kind);
if (entries.length === 0) continue;
const message = this.#build(kind, dispatcher, entries);
if (!message) continue;
if (mode === "streaming") {
try {
this.#options.injectStreaming(message);
} catch (error) {
logger.warn("Yield queue streaming dispatch failed", { kind, error: formatError(error) });
}
} else {
idleMessages.push(message);
}
}
if (mode === "idle" && idleMessages.length > 0) {
try {
await this.#options.injectIdle(idleMessages);
} catch (error) {
logger.warn("Yield queue idle dispatch failed", { error: formatError(error) });
}
}
}
clear(): void {
this.#entries.clear();
this.#idleFlushPending = false;
}
#scheduleIdleFlush(): void {
if (this.#idleFlushPending) return;
this.#idleFlushPending = true;
try {
this.#options.scheduleIdleFlush(async () => {
this.#idleFlushPending = false;
if (this.#options.isStreaming()) return;
await this.flush("idle");
});
} catch (error) {
this.#idleFlushPending = false;
logger.warn("Yield queue idle flush scheduling failed", { error: formatError(error) });
}
}
#drain(kind: string): unknown[] {
const entries = this.#entries.get(kind);
if (!entries || entries.length === 0) return [];
this.#entries.delete(kind);
return entries;
}
#build(kind: string, dispatcher: StoredDispatcher, entries: unknown[]): AgentMessage | null {
const survivors: unknown[] = [];
for (const entry of entries) {
if (dispatcher.isStale) {
let stale: boolean;
try {
stale = dispatcher.isStale(entry);
} catch (error) {
logger.warn("Yield queue stale check failed", { kind, error: formatError(error) });
continue;
}
if (stale) continue;
}
survivors.push(entry);
}
if (survivors.length === 0) return null;
try {
return dispatcher.build(survivors);
} catch (error) {
logger.warn("Yield queue build failed", { kind, error: formatError(error) });
return null;
}
}
}
@@ -0,0 +1,173 @@
import { afterEach, describe, expect, test } from "bun:test";
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { type AsyncJob, AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async";
import type { CustomMessage } from "@oh-my-pi/pi-coding-agent/session/messages";
import { YieldQueue } from "@oh-my-pi/pi-coding-agent/session/yield-queue";
import type { ToolSession } from "@oh-my-pi/pi-coding-agent/tools";
import { JobTool } from "@oh-my-pi/pi-coding-agent/tools/job";
type AsyncEntry = {
jobId: string;
result: string;
job: AsyncJob | undefined;
durationMs: number | undefined;
};
type AsyncDetails = {
jobs: Array<{
jobId: string;
type?: "bash" | "task";
label?: string;
durationMs?: number;
}>;
};
function buildAsyncMessage(entries: AsyncEntry[]): CustomMessage<AsyncDetails> | null {
if (entries.length === 0) return null;
return {
role: "custom",
customType: "async-result",
content: entries.map(entry => entry.result).join("\n"),
display: true,
attribution: "agent",
details: {
jobs: entries.map(entry => ({
jobId: entry.jobId,
type: entry.job?.type,
label: entry.job?.label,
durationMs: entry.durationMs,
})),
},
timestamp: 0,
};
}
function asyncDetails(message: AgentMessage): AsyncDetails {
if (message.role !== "custom") throw new Error(`Expected custom message, got ${message.role}`);
return (message as CustomMessage<AsyncDetails>).details ?? { jobs: [] };
}
function createToolSession(): ToolSession {
return {
cwd: process.cwd(),
hasUI: false,
settings: {
get: (key: string) => (key === "async.pollWaitDuration" ? "5s" : undefined),
},
getSessionFile: () => null,
getSessionSpawns: () => null,
getAgentId: () => null,
} as unknown as ToolSession;
}
function createHarness(initialStreaming: boolean) {
let streaming = initialStreaming;
const followUps: AgentMessage[] = [];
const prompts: AgentMessage[][] = [];
const scheduledFlushes: Array<() => Promise<void>> = [];
const queue = new YieldQueue({
isStreaming: () => streaming,
injectStreaming: message => {
followUps.push(message);
},
injectIdle: async messages => {
prompts.push(messages);
},
scheduleIdleFlush: run => {
scheduledFlushes.push(run);
},
});
let manager!: AsyncJobManager;
queue.register<AsyncEntry>("async-result", {
isStale: entry => manager.isDeliverySuppressed(entry.jobId),
build: buildAsyncMessage,
});
manager = new AsyncJobManager({
onJobComplete: (jobId, result, job) => {
if (manager.isDeliverySuppressed(jobId)) return;
queue.enqueue<AsyncEntry>("async-result", {
jobId,
result,
job,
durationMs: job ? Math.max(0, Date.now() - job.startTime) : undefined,
});
},
});
AsyncJobManager.setInstance(manager);
return {
manager,
queue,
followUps,
prompts,
scheduledFlushes,
setStreaming: (value: boolean) => {
streaming = value;
},
};
}
async function waitUntil(predicate: () => boolean, message: string): Promise<void> {
const deadline = Date.now() + 2_000;
while (!predicate()) {
if (Date.now() >= deadline) throw new Error(message);
await Bun.sleep(5);
}
}
afterEach(async () => {
const manager = AsyncJobManager.instance();
if (manager) {
await manager.dispose({ timeoutMs: 200 });
}
AsyncJobManager.resetForTests();
});
describe("async result yield queue delivery", () => {
test("job poll acknowledgement suppresses already staged completion", async () => {
const harness = createHarness(true);
const jobId = harness.manager.register("bash", "race job", async () => "inline result");
await harness.manager.waitForAll();
await waitUntil(() => harness.queue.has("async-result"), "Timed out waiting for staged async result");
const tool = new JobTool(createToolSession());
const result = await tool.execute("tool-call", { poll: [jobId] });
expect(result.details?.jobs.find(job => job.id === jobId)?.status).toBe("completed");
await harness.queue.flush("streaming");
expect(harness.followUps).toHaveLength(0);
});
test("multiple completions in one yield window become one follow-up", async () => {
const harness = createHarness(true);
const firstJobId = harness.manager.register("bash", "first", async () => "first result");
const secondJobId = harness.manager.register("task", "second", async () => "second result");
await harness.manager.waitForAll();
expect(await harness.manager.drainDeliveries({ timeoutMs: 2_000 })).toBe(true);
await harness.queue.flush("streaming");
expect(harness.followUps).toHaveLength(1);
const deliveredIds = asyncDetails(harness.followUps[0]!)
.jobs.map(job => job.jobId)
.sort();
expect(deliveredIds).toEqual([firstJobId, secondJobId].sort());
});
test("idle completion prompts once after scheduled idle flush", async () => {
const harness = createHarness(false);
const jobId = harness.manager.register("bash", "idle job", async () => "idle result");
await harness.manager.waitForAll();
expect(await harness.manager.drainDeliveries({ timeoutMs: 2_000 })).toBe(true);
expect(harness.scheduledFlushes).toHaveLength(1);
expect(harness.prompts).toHaveLength(0);
await harness.scheduledFlushes[0]!();
expect(harness.prompts).toHaveLength(1);
expect(harness.prompts[0]).toHaveLength(1);
expect(asyncDetails(harness.prompts[0]![0]!).jobs.map(job => job.jobId)).toEqual([jobId]);
});
});
@@ -0,0 +1,154 @@
import { describe, expect, test } from "bun:test";
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { YieldQueue } from "@oh-my-pi/pi-coding-agent/session/yield-queue";
type Entry = {
id: string;
stale?: boolean;
};
function userMessage(text: string): AgentMessage {
return {
role: "user",
content: [{ type: "text", text }],
timestamp: 0,
};
}
function messageText(message: AgentMessage): string {
if (!("content" in message) || !Array.isArray(message.content)) return "";
const block = message.content[0];
return block?.type === "text" ? block.text : "";
}
function createHarness(initialStreaming: boolean) {
let streaming = initialStreaming;
const streamingMessages: AgentMessage[] = [];
const idleBatches: AgentMessage[][] = [];
const scheduledFlushes: Array<() => Promise<void>> = [];
const queue = new YieldQueue({
isStreaming: () => streaming,
injectStreaming: message => {
streamingMessages.push(message);
},
injectIdle: async messages => {
idleBatches.push(messages);
},
scheduleIdleFlush: run => {
scheduledFlushes.push(run);
},
});
return {
queue,
streamingMessages,
idleBatches,
scheduledFlushes,
setStreaming: (value: boolean) => {
streaming = value;
},
};
}
describe("YieldQueue", () => {
test("enqueue while streaming defers until streaming flush", async () => {
const harness = createHarness(true);
harness.queue.register<Entry>("items", {
build: entries => userMessage(entries.map(entry => entry.id).join(",")),
});
harness.queue.enqueue("items", { id: "a" });
expect(harness.scheduledFlushes).toHaveLength(0);
expect(harness.streamingMessages).toHaveLength(0);
expect(harness.queue.has("items")).toBe(true);
await harness.queue.flush("streaming");
expect(harness.queue.has()).toBe(false);
expect(harness.streamingMessages.map(messageText)).toEqual(["a"]);
});
test("enqueue while idle schedules one debounced idle flush", async () => {
const harness = createHarness(false);
harness.queue.register<Entry>("items", {
build: entries => userMessage(entries.map(entry => entry.id).join(",")),
});
harness.queue.enqueue("items", { id: "a" });
harness.queue.enqueue("items", { id: "b" });
expect(harness.scheduledFlushes).toHaveLength(1);
expect(harness.idleBatches).toHaveLength(0);
await harness.scheduledFlushes[0]!();
expect(harness.idleBatches).toHaveLength(1);
expect(harness.idleBatches[0]?.map(messageText)).toEqual(["a,b"]);
});
test("isStale drops stale entries and keeps survivors", async () => {
const harness = createHarness(true);
let survivorIds: string[] = [];
harness.queue.register<Entry>("items", {
isStale: entry => entry.stale === true,
build: entries => {
survivorIds = entries.map(entry => entry.id);
return userMessage(survivorIds.join(","));
},
});
harness.queue.enqueue("items", { id: "old", stale: true });
harness.queue.enqueue("items", { id: "fresh" });
await harness.queue.flush("streaming");
expect(survivorIds).toEqual(["fresh"]);
expect(harness.streamingMessages.map(messageText)).toEqual(["fresh"]);
});
test("build returning null does not inject", async () => {
const harness = createHarness(true);
harness.queue.register<Entry>("items", {
build: () => null,
});
harness.queue.enqueue("items", { id: "a" });
await harness.queue.flush("streaming");
expect(harness.streamingMessages).toHaveLength(0);
expect(harness.idleBatches).toHaveLength(0);
});
test("one kind failing in build does not abort other kinds", async () => {
const harness = createHarness(true);
harness.queue.register<Entry>("bad", {
build: () => {
throw new Error("boom");
},
});
harness.queue.register<Entry>("good", {
build: entries => userMessage(entries.map(entry => entry.id).join(",")),
});
harness.queue.enqueue("bad", { id: "bad" });
harness.queue.enqueue("good", { id: "good" });
await harness.queue.flush("streaming");
expect(harness.streamingMessages.map(messageText)).toEqual(["good"]);
});
test("flush preserves registration order across kinds", async () => {
const harness = createHarness(true);
harness.queue.register<Entry>("second", {
build: entries => userMessage(entries.map(entry => entry.id).join(",")),
});
harness.queue.register<Entry>("first", {
build: entries => userMessage(entries.map(entry => entry.id).join(",")),
});
harness.queue.enqueue("first", { id: "first" });
harness.queue.enqueue("second", { id: "second" });
await harness.queue.flush("streaming");
expect(harness.streamingMessages.map(messageText)).toEqual(["second", "first"]);
});
});