From 796f963da9123e9923e672bcfc25f76d7163ff44 Mon Sep 17 00:00:00 2001 From: can1357 Date: Fri, 22 May 2026 13:05:30 +0900 Subject: [PATCH] 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. --- packages/agent/CHANGELOG.md | 3 + packages/agent/src/agent-loop.ts | 1 + packages/agent/src/agent.ts | 6 + packages/agent/src/types.ts | 7 + packages/agent/test/agent-loop.test.ts | 33 ++++ packages/coding-agent/CHANGELOG.md | 5 + .../src/modes/utils/ui-helpers.ts | 44 +++-- .../src/prompts/tools/async-result.md | 7 +- packages/coding-agent/src/sdk.ts | 116 +++++++++--- .../coding-agent/src/session/agent-session.ts | 22 +++ .../coding-agent/src/session/yield-queue.ts | 155 ++++++++++++++++ .../test/async-yield-queue.test.ts | 173 ++++++++++++++++++ .../test/session/yield-queue.test.ts | 154 ++++++++++++++++ 13 files changed, 690 insertions(+), 36 deletions(-) create mode 100644 packages/coding-agent/src/session/yield-queue.ts create mode 100644 packages/coding-agent/test/async-yield-queue.test.ts create mode 100644 packages/coding-agent/test/session/yield-queue.test.ts diff --git a/packages/agent/CHANGELOG.md b/packages/agent/CHANGELOG.md index f485dabcd..2cf068750 100644 --- a/packages/agent/CHANGELOG.md +++ b/packages/agent/CHANGELOG.md @@ -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 diff --git a/packages/agent/src/agent-loop.ts b/packages/agent/src/agent-loop.ts index 7dc409f2f..4a2626874 100644 --- a/packages/agent/src/agent-loop.ts +++ b/packages/agent/src/agent-loop.ts @@ -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 diff --git a/packages/agent/src/agent.ts b/packages/agent/src/agent.ts index c89c062c9..c5a9e45fc 100644 --- a/packages/agent/src/agent.ts +++ b/packages/agent/src/agent.ts @@ -290,6 +290,7 @@ export class Agent { #onSseEvent?: SimpleStreamOptions["onSseEvent"]; #onAssistantMessageEvent?: (message: AssistantMessage, event: AssistantMessageEvent) => void; #onHarmonyLeak?: (event: HarmonyAuditEvent) => void | Promise; + #onBeforeYield?: () => Promise | 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) | 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, }; diff --git a/packages/agent/src/types.ts b/packages/agent/src/types.ts index 26e11b468..0b2c3bcc9 100644 --- a/packages/agent/src/types.ts +++ b/packages/agent/src/types.ts @@ -122,6 +122,13 @@ export interface AgentLoopConfig extends SimpleStreamOptions { * continues with another turn. */ getFollowUpMessages?: () => Promise; + /** + * 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; /** * Provides tool execution context, resolved per tool call. diff --git a/packages/agent/test/agent-loop.test.ts b/packages/agent/test/agent-loop.test.ts index 44204ab86..dc4a97365 100644 --- a/packages/agent/test/agent-loop.test.ts +++ b/packages/agent/test/agent-loop.test.ts @@ -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" }); + }); }); diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 509863500..d206356a2 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] + ### Breaking Changes - Changed PR and task-isolation worktree directory layout to hash-based `~/.omp/wt/-` 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 diff --git a/packages/coding-agent/src/modes/utils/ui-helpers.ts b/packages/coding-agent/src/modes/utils/ui-helpers.ts index 1d4de40ad..ea2468f09 100644 --- a/packages/coding-agent/src/modes/utils/ui-helpers.ts +++ b/packages/coding-agent/src/modes/utils/ui-helpers.ts @@ -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) { diff --git a/packages/coding-agent/src/prompts/tools/async-result.md b/packages/coding-agent/src/prompts/tools/async-result.md index 370d8077f..ff3501758 100644 --- a/packages/coding-agent/src/prompts/tools/async-result.md +++ b/packages/coding-agent/src/prompts/tools/async-result.md @@ -1,5 +1,8 @@ -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}} diff --git a/packages/coding-agent/src/sdk.ts b/packages/coding-agent/src/sdk.ts index 3760549ef..a06ea56f1 100644 --- a/packages/coding-agent/src/sdk.ts +++ b/packages/coding-agent/src/sdk.ts @@ -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 | 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(); + 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://") 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("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("async-result", { + isStale: entry => asyncJobManager.isDeliverySuppressed(entry.jobId), + build: buildAsyncResultBatchMessage, + }); + } + session.yieldQueue.register("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("mcp-notification", { serverName, uri }); }, debounceMs), ); }); diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 6ee148338..571a84f81 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -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 { this.#isDisposed = true; this.#pendingBackgroundExchanges = []; + this.yieldQueue.clear(); + this.agent.setOnBeforeYield(undefined); this.#evalExecutionDisposing = true; try { if (this.#extensionRunner?.hasHandlers("session_shutdown")) { diff --git a/packages/coding-agent/src/session/yield-queue.ts b/packages/coding-agent/src/session/yield-queue.ts new file mode 100644 index 000000000..a329531af --- /dev/null +++ b/packages/coding-agent/src/session/yield-queue.ts @@ -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

{ + /** 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; + scheduleIdleFlush(run: () => Promise): 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(); + readonly #entries = new Map(); + #idleFlushPending = false; + + constructor(options: YieldQueueOptions) { + this.#options = options; + } + + register

(kind: string, dispatcher: YieldDispatcher

): () => 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

(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 { + 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; + } + } +} diff --git a/packages/coding-agent/test/async-yield-queue.test.ts b/packages/coding-agent/test/async-yield-queue.test.ts new file mode 100644 index 000000000..07984a619 --- /dev/null +++ b/packages/coding-agent/test/async-yield-queue.test.ts @@ -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 | 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).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> = []; + 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("async-result", { + isStale: entry => manager.isDeliverySuppressed(entry.jobId), + build: buildAsyncMessage, + }); + manager = new AsyncJobManager({ + onJobComplete: (jobId, result, job) => { + if (manager.isDeliverySuppressed(jobId)) return; + queue.enqueue("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 { + 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]); + }); +}); diff --git a/packages/coding-agent/test/session/yield-queue.test.ts b/packages/coding-agent/test/session/yield-queue.test.ts new file mode 100644 index 000000000..200ffcc79 --- /dev/null +++ b/packages/coding-agent/test/session/yield-queue.test.ts @@ -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> = []; + 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("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("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("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("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("bad", { + build: () => { + throw new Error("boom"); + }, + }); + harness.queue.register("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("second", { + build: entries => userMessage(entries.map(entry => entry.id).join(",")), + }); + harness.queue.register("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"]); + }); +});