Files
oh-my-pi/packages/coding-agent/test/async-yield-queue.test.ts
T
can1357 5ff277349c refactor(coding-agent): consolidated tool surface onto xd:// devices and hub
- Added the `xd://` virtual device protocol (`internal-urls/xd-protocol.ts`, `tools/xdev.ts`): tools declaring `loadMode: "discoverable"` are unmounted from the request tools array and driven via `read xd://` (list/docs+schema) and `write xd://<tool>` (execute), gated by the `tools.xdev` setting (default on) and inlined into the system prompt.
- Merged the `irc`, `job`, and `launch` tools into a single `hub` tool (`tools/hub/`, `async/job-manager.ts`): messaging keeps `send`/`inbox`/`list`, job control maps to `wait`/`cancel`/`jobs`, process supervision keeps `start`/`logs`/`stop`/`restart`/`describe` with `ps`, and the unified `wait` races background jobs against peer messages; SDK `IrcTool`/`JobTool`/`LaunchTool` are replaced by `HubTool`.
- Removed the hidden `resolve` tool in favor of the `xd://resolve`/`xd://reject`/`xd://propose` resolution devices, auto-including `write` whenever a deferrable tool or plan mode is present.
- Removed the BM25 tool-discovery system: the `search_tool_bm25` tool, the `tool-discovery` module, the `tools.discoveryMode`/`mcp.discoveryMode`/`mcp.discoveryDefaultServers`/`tools.essentialOverride` settings, per-tool MCP selection, and the `mcp_tool_selection` message type.
- Unified tool presentation on `ToolLoadMode` (`essential`|`discoverable`), replacing the custom-tool `xdev?: boolean` opt-out; custom, extension, MCP, RPC host, image-generation, and TTS tools now default to `discoverable`, and added a `satisfies` predicate to `SoftToolRequirement`.
- Removed the standalone `ssh` command tool and `ssh/ssh-executor` (the `ssh://` read/write/search protocol stays), and made `--tools` address hidden built-ins.
- Updated collab-web to render `xd://` dispatches and `hub` op families, dropped the `search_tool_bm25`/`ssh`/`report-finding` renderers, refreshed tool docs and prompts, and migrated the affected tests and changelogs.
2026-07-15 15:16:29 +02:00

175 lines
5.4 KiB
TypeScript

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 { type CoordinationDetails, HubTool } from "../src/tools/hub";
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(asyncJobManager?: AsyncJobManager): ToolSession {
return {
cwd: process.cwd(),
hasUI: false,
settings: {
get: (key: string) => (key === "async.pollWaitDuration" ? "5s" : undefined),
},
getSessionFile: () => null,
getSessionSpawns: () => null,
getAgentId: () => null,
asyncJobManager,
} 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 HubTool(createToolSession(harness.manager));
const result = await tool.execute("tool-call", { op: "wait", ids: [jobId] });
expect((result.details as CoordinationDetails)?.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]);
});
});