fix(coding-agent): tracked acp tool args by session for replay
- Tracked ACP tool-call inputs per session and replayed them via `toolArgsById`/`getToolArgs` plumbing. - Merged ACP tool execution end content from start and result events so command output replay preserves original args. - Scoped ACP async-job draining by session `ownerId` and `agentId` with in-flight tracking and permission-gated deferred turns. - Refactored compaction telemetry and async tests with per-test telemetry setup and asynchronous teardown resets.
This commit is contained in:
@@ -2,12 +2,12 @@
|
||||
* Tests for OpenTelemetry instrumentation around oneshot LLM calls:
|
||||
* compaction summaries, handoff document, branch summary.
|
||||
*
|
||||
* Mirrors the InMemorySpanExporter + AsyncLocalStorageContextManager setup
|
||||
* used by `otel.test.ts`. Spies on `completeSimple` to avoid real HTTP traffic
|
||||
* Uses a per-test InMemorySpanExporter and explicit tracer. Spies on
|
||||
* `completeSimple` to avoid real HTTP traffic
|
||||
* while exercising the chat-span lifecycle (`startChatSpan` →
|
||||
* `runInActiveSpan` → `finishChatSpan` / `failChatSpan`).
|
||||
*/
|
||||
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "bun:test";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
|
||||
import {
|
||||
type CompactionPreparation,
|
||||
compact,
|
||||
@@ -27,8 +27,7 @@ import {
|
||||
import type { AgentMessage } from "@oh-my-pi/pi-agent-core/types";
|
||||
import type { AssistantMessage, Model, Usage } from "@oh-my-pi/pi-ai";
|
||||
import * as ai from "@oh-my-pi/pi-ai";
|
||||
import { context, SpanStatusCode, trace } from "@opentelemetry/api";
|
||||
import { AsyncLocalStorageContextManager } from "@opentelemetry/context-async-hooks";
|
||||
import { SpanStatusCode } from "@opentelemetry/api";
|
||||
import {
|
||||
BasicTracerProvider,
|
||||
InMemorySpanExporter,
|
||||
@@ -49,28 +48,18 @@ const MODEL: Model = {
|
||||
maxTokens: 32_768,
|
||||
};
|
||||
|
||||
const exporter = new InMemorySpanExporter();
|
||||
let exporter: InMemorySpanExporter;
|
||||
let provider: BasicTracerProvider;
|
||||
let contextManager: AsyncLocalStorageContextManager;
|
||||
|
||||
beforeAll(() => {
|
||||
trace.disable();
|
||||
context.disable();
|
||||
contextManager = new AsyncLocalStorageContextManager().enable();
|
||||
context.setGlobalContextManager(contextManager);
|
||||
beforeEach(() => {
|
||||
exporter = new InMemorySpanExporter();
|
||||
provider = new BasicTracerProvider({ spanProcessors: [new SimpleSpanProcessor(exporter)] });
|
||||
trace.setGlobalTracerProvider(provider);
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
afterEach(async () => {
|
||||
exporter.reset();
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
afterAll(async () => {
|
||||
await provider.shutdown();
|
||||
context.disable();
|
||||
trace.disable();
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
function makeTelemetryConfig(): AgentTelemetryConfig {
|
||||
|
||||
@@ -1,12 +1,14 @@
|
||||
# Changelog
|
||||
|
||||
## [Unreleased]
|
||||
### Fixed
|
||||
|
||||
- Fixed ACP ordinary file-editing calls (`edit`, `write`, `ast_edit`) incorrectly requesting `session/request_permission` before every call, while keeping permission prompts for edit operations that delete or move files; permission requests now report the gated tool call as `pending` so clients can render the approval UI instead of returning `Permission request cancelled` without a visible prompt.
|
||||
|
||||
### Fixed
|
||||
|
||||
- Fixed ACP command and custom tool-call notifications to carry the original tool arguments in replayed and final updates, so command text is preserved and raw input is no longer wrapped
|
||||
- Fixed ACP async-job draining to be scoped by session owner so `getAsyncJobSnapshot` and `drainAsyncJobDeliveriesForAcp` no longer consume or expose jobs from other sessions
|
||||
- Fixed async job status reporting to include in-flight completions so queued/delivering indicators remain accurate while callbacks are still running
|
||||
- Fixed `deferAgentInitiatedTurns` handling during ACP async-job draining so background completion follow-up turns are delivered even when agent-initiated turns are deferred
|
||||
- Fixed ACP ordinary file-editing calls (`edit`, `write`, `ast_edit`) incorrectly requesting `session/request_permission` before every call, while keeping permission prompts for edit operations that delete or move files; permission requests now report the gated tool call as `pending` so clients can render the approval UI instead of returning `Permission request cancelled` without a visible prompt. ([#1134](https://github.com/can1357/oh-my-pi/pull/1134) by [@jiwangyihao](https://github.com/jiwangyihao))
|
||||
- Fixed the session tree selector to preserve a readable message column when deeply nested branch gutters would otherwise consume the viewport. ([#1144](https://github.com/can1357/oh-my-pi/issues/1144))
|
||||
|
||||
## [15.1.3] - 2026-05-17
|
||||
|
||||
@@ -37,6 +37,8 @@ interface AsyncJobDelivery {
|
||||
attempt: number;
|
||||
nextAttemptAt: number;
|
||||
lastError?: string;
|
||||
ownerId?: string;
|
||||
promise?: Promise<void>;
|
||||
}
|
||||
|
||||
export interface AsyncJobDeliveryState {
|
||||
@@ -82,6 +84,7 @@ export class AsyncJobManager {
|
||||
|
||||
readonly #jobs = new Map<string, AsyncJob>();
|
||||
readonly #deliveries: AsyncJobDelivery[] = [];
|
||||
readonly #inFlightDeliveries: AsyncJobDelivery[] = [];
|
||||
readonly #suppressedDeliveries = new Set<string>();
|
||||
readonly #watchedJobs = new Set<string>();
|
||||
readonly #evictionTimers = new Map<string, NodeJS.Timeout>();
|
||||
@@ -221,16 +224,17 @@ export class AsyncJobManager {
|
||||
|
||||
getDeliveryState(filter?: AsyncJobFilter): AsyncJobDeliveryState {
|
||||
const deliveries = this.#filterDeliveries(filter);
|
||||
const inFlightDeliveries = this.#filterInFlightDeliveries(filter);
|
||||
const nextRetryAt = deliveries.reduce<number | undefined>((next, delivery) => {
|
||||
if (next === undefined) return delivery.nextAttemptAt;
|
||||
return Math.min(next, delivery.nextAttemptAt);
|
||||
}, undefined);
|
||||
|
||||
return {
|
||||
queued: deliveries.length,
|
||||
delivering: this.#deliveryLoop !== undefined && deliveries.length > 0,
|
||||
queued: deliveries.length + inFlightDeliveries.length,
|
||||
delivering: inFlightDeliveries.length > 0 || (this.#deliveryLoop !== undefined && deliveries.length > 0),
|
||||
nextRetryAt,
|
||||
pendingJobIds: deliveries.map(delivery => delivery.jobId),
|
||||
pendingJobIds: deliveries.concat(inFlightDeliveries).map(delivery => delivery.jobId),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -303,6 +307,12 @@ export class AsyncJobManager {
|
||||
if (delivered) continue;
|
||||
return false;
|
||||
}
|
||||
const inFlightDeliveries = this.#filterInFlightDeliveries();
|
||||
if (inFlightDeliveries.length > 0 && this.#filterDeliveries().length === 0) {
|
||||
const delivered = await this.#waitForDeliveryPromise(inFlightDeliveries[0]?.promise, deadline);
|
||||
if (delivered) continue;
|
||||
return false;
|
||||
}
|
||||
|
||||
this.#ensureDeliveryLoop();
|
||||
const loop = this.#deliveryLoop;
|
||||
@@ -338,6 +348,7 @@ export class AsyncJobManager {
|
||||
this.#clearEvictionTimers();
|
||||
this.#jobs.clear();
|
||||
this.#deliveries.length = 0;
|
||||
this.#inFlightDeliveries.length = 0;
|
||||
this.#suppressedDeliveries.clear();
|
||||
this.#watchedJobs.clear();
|
||||
return drained;
|
||||
@@ -398,49 +409,51 @@ export class AsyncJobManager {
|
||||
|
||||
#filterDeliveries(filter?: AsyncJobFilter): AsyncJobDelivery[] {
|
||||
const ownerId = filter?.ownerId;
|
||||
if (!ownerId) return this.#deliveries;
|
||||
return this.#deliveries.filter(delivery => this.#jobs.get(delivery.jobId)?.ownerId === ownerId);
|
||||
if (!ownerId) return this.#deliveries.filter(delivery => !this.isDeliverySuppressed(delivery.jobId));
|
||||
return this.#deliveries.filter(
|
||||
delivery => delivery.ownerId === ownerId && !this.isDeliverySuppressed(delivery.jobId),
|
||||
);
|
||||
}
|
||||
|
||||
#filterInFlightDeliveries(filter?: AsyncJobFilter): AsyncJobDelivery[] {
|
||||
const ownerId = filter?.ownerId;
|
||||
if (!ownerId) return this.#inFlightDeliveries.filter(delivery => !this.isDeliverySuppressed(delivery.jobId));
|
||||
return this.#inFlightDeliveries.filter(
|
||||
delivery => delivery.ownerId === ownerId && !this.isDeliverySuppressed(delivery.jobId),
|
||||
);
|
||||
}
|
||||
|
||||
async #deliverNextFiltered(filter: AsyncJobFilter, deadline: number): Promise<boolean> {
|
||||
let selected: AsyncJobDelivery | undefined;
|
||||
for (const delivery of this.#deliveries) {
|
||||
if (this.#jobs.get(delivery.jobId)?.ownerId !== filter.ownerId) continue;
|
||||
if (this.isDeliverySuppressed(delivery.jobId)) continue;
|
||||
if (!selected || delivery.nextAttemptAt < selected.nextAttemptAt) {
|
||||
selected = delivery;
|
||||
while (true) {
|
||||
let selected: AsyncJobDelivery | undefined;
|
||||
for (const delivery of this.#deliveries) {
|
||||
if (delivery.ownerId !== filter.ownerId) continue;
|
||||
if (this.isDeliverySuppressed(delivery.jobId)) continue;
|
||||
if (!selected || delivery.nextAttemptAt < selected.nextAttemptAt) {
|
||||
selected = delivery;
|
||||
}
|
||||
}
|
||||
}
|
||||
if (!selected) return true;
|
||||
|
||||
const now = Date.now();
|
||||
if (selected.nextAttemptAt > now) {
|
||||
if (selected.nextAttemptAt > deadline) return false;
|
||||
await Bun.sleep(selected.nextAttemptAt - now);
|
||||
}
|
||||
|
||||
const index = this.#deliveries.indexOf(selected);
|
||||
if (index === -1) return true;
|
||||
|
||||
try {
|
||||
await this.#onJobComplete(selected.jobId, selected.text, this.#jobs.get(selected.jobId));
|
||||
this.#deliveries.splice(index, 1);
|
||||
} catch (error) {
|
||||
selected.attempt += 1;
|
||||
selected.lastError = error instanceof Error ? error.message : String(error);
|
||||
selected.nextAttemptAt = Date.now() + this.#getRetryDelay(selected.attempt);
|
||||
this.#deliveries.splice(index, 1);
|
||||
if (!this.isDeliverySuppressed(selected.jobId)) {
|
||||
this.#deliveries.push(selected);
|
||||
if (!selected) {
|
||||
const inFlight = this.#filterInFlightDeliveries(filter);
|
||||
if (inFlight.length === 0) return true;
|
||||
return this.#waitForDeliveryPromise(inFlight[0]?.promise, deadline);
|
||||
}
|
||||
logger.warn("Async job completion delivery failed", {
|
||||
jobId: selected.jobId,
|
||||
attempt: selected.attempt,
|
||||
nextRetryAt: selected.nextAttemptAt,
|
||||
error: selected.lastError,
|
||||
});
|
||||
|
||||
const now = Date.now();
|
||||
if (selected.nextAttemptAt > now) {
|
||||
if (selected.nextAttemptAt > deadline) return false;
|
||||
await Bun.sleep(selected.nextAttemptAt - now);
|
||||
continue;
|
||||
}
|
||||
|
||||
const index = this.#deliveries.indexOf(selected);
|
||||
if (index === -1) continue;
|
||||
this.#deliveries.splice(index, 1);
|
||||
if (this.isDeliverySuppressed(selected.jobId)) continue;
|
||||
|
||||
return this.#waitForDeliveryPromise(this.#deliverDelivery(selected), deadline);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
isDeliverySuppressed(jobId: string): boolean {
|
||||
@@ -457,6 +470,7 @@ export class AsyncJobManager {
|
||||
text,
|
||||
attempt: 0,
|
||||
nextAttemptAt: Date.now(),
|
||||
ownerId: this.#jobs.get(jobId)?.ownerId,
|
||||
});
|
||||
this.#ensureDeliveryLoop();
|
||||
}
|
||||
@@ -492,20 +506,25 @@ export class AsyncJobManager {
|
||||
if (this.#deliveries[0] !== delivery) {
|
||||
continue;
|
||||
}
|
||||
// Check again after sleep
|
||||
if (this.isDeliverySuppressed(delivery.jobId)) {
|
||||
this.#deliveries.shift();
|
||||
continue;
|
||||
}
|
||||
|
||||
this.#deliveries.shift();
|
||||
await this.#deliverDelivery(delivery);
|
||||
}
|
||||
}
|
||||
|
||||
#deliverDelivery(delivery: AsyncJobDelivery): Promise<void> {
|
||||
const promise = (async () => {
|
||||
this.#inFlightDeliveries.push(delivery);
|
||||
try {
|
||||
await this.#onJobComplete(delivery.jobId, delivery.text, this.#jobs.get(delivery.jobId));
|
||||
this.#deliveries.shift();
|
||||
} catch (error) {
|
||||
delivery.attempt += 1;
|
||||
delivery.lastError = error instanceof Error ? error.message : String(error);
|
||||
delivery.nextAttemptAt = Date.now() + this.#getRetryDelay(delivery.attempt);
|
||||
this.#deliveries.shift();
|
||||
if (!this.isDeliverySuppressed(delivery.jobId)) {
|
||||
this.#deliveries.push(delivery);
|
||||
}
|
||||
@@ -515,8 +534,32 @@ export class AsyncJobManager {
|
||||
nextRetryAt: delivery.nextAttemptAt,
|
||||
error: delivery.lastError,
|
||||
});
|
||||
} finally {
|
||||
const index = this.#inFlightDeliveries.indexOf(delivery);
|
||||
if (index !== -1) this.#inFlightDeliveries.splice(index, 1);
|
||||
if (this.#deliveries.length > 0) this.#ensureDeliveryLoop();
|
||||
}
|
||||
})();
|
||||
delivery.promise = promise;
|
||||
return promise;
|
||||
}
|
||||
|
||||
async #waitForDeliveryPromise(promise: Promise<void> | undefined, deadline: number): Promise<boolean> {
|
||||
if (!promise) return true;
|
||||
if (deadline === Number.POSITIVE_INFINITY) {
|
||||
await promise;
|
||||
return true;
|
||||
}
|
||||
const remainingMs = deadline - Date.now();
|
||||
if (remainingMs <= 0) return false;
|
||||
let timedOut = false;
|
||||
await Promise.race([
|
||||
promise,
|
||||
Bun.sleep(remainingMs).then(() => {
|
||||
timedOut = true;
|
||||
}),
|
||||
]);
|
||||
return !timedOut;
|
||||
}
|
||||
|
||||
#getRetryDelay(attempt: number): number {
|
||||
|
||||
@@ -56,6 +56,7 @@ import {
|
||||
discoverAuthStorage,
|
||||
} from "./sdk";
|
||||
import type { AgentSession } from "./session/agent-session";
|
||||
import type { AuthStorage } from "./session/auth-storage";
|
||||
import { resolveResumableSession, type SessionInfo, SessionManager } from "./session/session-manager";
|
||||
import { resolvePromptInput } from "./system-prompt";
|
||||
import type { LspStartupServerInfo } from "./tools";
|
||||
@@ -202,7 +203,7 @@ interface AcpSessionFactoryOptions {
|
||||
baseOptions: CreateAgentSessionOptions;
|
||||
settings: Settings;
|
||||
sessionDir?: string;
|
||||
authStorage: Awaited<ReturnType<typeof discoverAuthStorage>>;
|
||||
authStorage: AuthStorage;
|
||||
modelRegistry: ModelRegistry;
|
||||
parsedArgs: Pick<Args, "apiKey">;
|
||||
rawArgs: string[];
|
||||
@@ -213,6 +214,7 @@ function createAcpSessionFactory(args: AcpSessionFactoryOptions): AcpSessionFact
|
||||
return async cwd => {
|
||||
const nextSettings = await args.settings.cloneForCwd(cwd);
|
||||
const nextSessionManager = SessionManager.create(cwd, args.sessionDir);
|
||||
const agentId = `acp:${nextSessionManager.getSessionId()}`;
|
||||
const { session: nextSession } = await args.createSession({
|
||||
...args.baseOptions,
|
||||
cwd,
|
||||
@@ -220,6 +222,7 @@ function createAcpSessionFactory(args: AcpSessionFactoryOptions): AcpSessionFact
|
||||
settings: nextSettings,
|
||||
authStorage: args.authStorage,
|
||||
modelRegistry: args.modelRegistry,
|
||||
agentId,
|
||||
hasUI: false,
|
||||
});
|
||||
if (args.parsedArgs.apiKey && !args.baseOptions.model && nextSession.model) {
|
||||
|
||||
@@ -140,6 +140,7 @@ type ManagedSessionRecord = {
|
||||
promptQueue: PromptQueueState;
|
||||
liveMessageId: string | undefined;
|
||||
liveMessageProgress: { textEmitted: boolean; thoughtEmitted: boolean } | undefined;
|
||||
toolArgsById: Map<string, unknown>;
|
||||
extensionsConfigured: boolean;
|
||||
// Installed inside `#scheduleBootstrapUpdates` (post-race-guard); released
|
||||
// in `#disposeSessionRecord`. Lives independent of any prompt turn.
|
||||
@@ -975,6 +976,7 @@ export class AcpAgent implements Agent {
|
||||
promptQueue: { promise: Promise.resolve(), release: undefined },
|
||||
liveMessageId: undefined,
|
||||
liveMessageProgress: undefined,
|
||||
toolArgsById: new Map(),
|
||||
extensionsConfigured: false,
|
||||
lifetimeUnsubscribe: undefined,
|
||||
};
|
||||
@@ -1037,14 +1039,22 @@ export class AcpAgent implements Agent {
|
||||
return;
|
||||
}
|
||||
|
||||
if (event.type === "tool_execution_start" || event.type === "tool_execution_update") {
|
||||
record.toolArgsById.set(event.toolCallId, event.args);
|
||||
}
|
||||
|
||||
this.#prepareLiveAssistantMessage(record, event);
|
||||
for (const notification of mapAgentSessionEventToAcpSessionUpdates(event, record.session.sessionId, {
|
||||
getMessageId: message => this.#getLiveMessageId(record, message),
|
||||
getMessageProgress: message => this.#getLiveMessageProgress(record, message),
|
||||
getToolArgs: toolCallId => record.toolArgsById.get(toolCallId),
|
||||
cwd: record.session.sessionManager.getCwd(),
|
||||
})) {
|
||||
await this.#connection.sessionUpdate(notification);
|
||||
}
|
||||
if (event.type === "tool_execution_end") {
|
||||
record.toolArgsById.delete(event.toolCallId);
|
||||
}
|
||||
this.#clearLiveAssistantMessageAfterEvent(record, event);
|
||||
|
||||
if (event.type === "agent_end") {
|
||||
@@ -1613,12 +1623,14 @@ export class AcpAgent implements Agent {
|
||||
async #replaySessionHistory(record: ManagedSessionRecord): Promise<void> {
|
||||
const cwd = record.session.sessionManager.getCwd();
|
||||
const replayedToolCallIds = new Set<string>();
|
||||
const replayedToolCallArgs = new Map<string, unknown>();
|
||||
for (const message of record.session.sessionManager.buildSessionContext().messages as ReplayableMessage[]) {
|
||||
for (const notification of this.#messageToReplayNotifications(
|
||||
record.session.sessionId,
|
||||
message,
|
||||
cwd,
|
||||
replayedToolCallIds,
|
||||
replayedToolCallArgs,
|
||||
)) {
|
||||
await this.#connection.sessionUpdate(notification);
|
||||
}
|
||||
@@ -1630,9 +1642,10 @@ export class AcpAgent implements Agent {
|
||||
message: ReplayableMessage,
|
||||
cwd: string,
|
||||
replayedToolCallIds: Set<string>,
|
||||
replayedToolCallArgs: Map<string, unknown>,
|
||||
): SessionNotification[] {
|
||||
if (message.role === "assistant") {
|
||||
return this.#replayAssistantMessage(sessionId, message, cwd, replayedToolCallIds);
|
||||
return this.#replayAssistantMessage(sessionId, message, cwd, replayedToolCallIds, replayedToolCallArgs);
|
||||
}
|
||||
if (
|
||||
message.role === "user" ||
|
||||
@@ -1660,7 +1673,10 @@ export class AcpAgent implements Agent {
|
||||
toolCallId: message.toolCallId,
|
||||
toolName: message.toolName,
|
||||
},
|
||||
{ includeStart: !replayedToolCallIds.has(message.toolCallId) },
|
||||
{
|
||||
includeStart: !replayedToolCallIds.has(message.toolCallId),
|
||||
toolArgs: replayedToolCallArgs.get(message.toolCallId),
|
||||
},
|
||||
);
|
||||
}
|
||||
if (
|
||||
@@ -1683,6 +1699,7 @@ export class AcpAgent implements Agent {
|
||||
message: ReplayableMessage,
|
||||
cwd: string,
|
||||
replayedToolCallIds: Set<string>,
|
||||
replayedToolCallArgs: Map<string, unknown>,
|
||||
): SessionNotification[] {
|
||||
const notifications: SessionNotification[] = [];
|
||||
const messageId = crypto.randomUUID();
|
||||
@@ -1734,6 +1751,7 @@ export class AcpAgent implements Agent {
|
||||
});
|
||||
notifications.push({ sessionId, update });
|
||||
replayedToolCallIds.add(toolItem.id);
|
||||
replayedToolCallArgs.set(toolItem.id, args);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1755,7 +1773,7 @@ export class AcpAgent implements Agent {
|
||||
return normalizeReplayToolArguments(item.arguments).args;
|
||||
}
|
||||
if (item.type === "tool_use" && "input" in item) {
|
||||
return { input: item.input };
|
||||
return item.input;
|
||||
}
|
||||
return {};
|
||||
}
|
||||
@@ -1764,7 +1782,7 @@ export class AcpAgent implements Agent {
|
||||
sessionId: string,
|
||||
cwd: string,
|
||||
message: Required<Pick<ReplayableMessage, "toolCallId" | "toolName">> & ReplayableMessage,
|
||||
options: { includeStart?: boolean } = {},
|
||||
options: { includeStart?: boolean; toolArgs?: unknown } = {},
|
||||
): SessionNotification[] {
|
||||
const args = this.#buildReplayToolArgs(message.details);
|
||||
const startEvent: AgentSessionEvent = {
|
||||
@@ -1784,7 +1802,10 @@ export class AcpAgent implements Agent {
|
||||
errorMessage: message.errorMessage,
|
||||
},
|
||||
};
|
||||
const notifications = mapAgentSessionEventToAcpSessionUpdates(endEvent, sessionId, { cwd });
|
||||
const notifications = mapAgentSessionEventToAcpSessionUpdates(endEvent, sessionId, {
|
||||
cwd,
|
||||
getToolArgs: toolCallId => (toolCallId === message.toolCallId ? options.toolArgs : undefined),
|
||||
});
|
||||
if (options.includeStart === false) {
|
||||
return notifications;
|
||||
}
|
||||
|
||||
@@ -18,6 +18,7 @@ interface MessageProgress {
|
||||
interface AcpEventMapperOptions {
|
||||
getMessageId?: (message: unknown) => string | undefined;
|
||||
getMessageProgress?: (message: unknown) => MessageProgress | undefined;
|
||||
getToolArgs?: (toolCallId: string) => unknown;
|
||||
/**
|
||||
* Session cwd. Tool call locations sent to ACP clients must be absolute
|
||||
* (the editor host needs them to open or focus files). When provided,
|
||||
@@ -185,8 +186,11 @@ export function mapAgentSessionEventToAcpSessionUpdates(
|
||||
return [toSessionNotification(sessionId, update)];
|
||||
}
|
||||
case "tool_execution_end": {
|
||||
const diffContent = extractDiffToolCallContent(event.result);
|
||||
const content = [...diffContent, ...extractToolCallContent(event.result)];
|
||||
const resultContent = [...extractDiffToolCallContent(event.result), ...extractToolCallContent(event.result)];
|
||||
const content = mergeToolUpdateContent(
|
||||
buildToolStartContent(event.toolName, getToolExecutionEndArgs(event, options)),
|
||||
resultContent,
|
||||
);
|
||||
const update: SessionUpdate = {
|
||||
sessionUpdate: "tool_call_update",
|
||||
toolCallId: event.toolCallId,
|
||||
@@ -380,6 +384,7 @@ function extractTodoEntries(phases: unknown[]): Array<{ content: string; status:
|
||||
|
||||
function isTodoStatus(status: unknown): status is TodoStatus {
|
||||
return status === "pending" || status === "in_progress" || status === "completed" || status === "abandoned";
|
||||
}
|
||||
export function buildToolCallStartUpdate(input: {
|
||||
toolCallId: string;
|
||||
toolName: string;
|
||||
@@ -419,6 +424,16 @@ export function normalizeReplayToolArguments(value: unknown): { args: unknown }
|
||||
}
|
||||
}
|
||||
|
||||
function getToolExecutionEndArgs(
|
||||
event: Extract<AgentSessionEvent, { type: "tool_execution_end" }>,
|
||||
options: AcpEventMapperOptions,
|
||||
): unknown {
|
||||
if ("args" in event) {
|
||||
return (event as { args?: unknown }).args;
|
||||
}
|
||||
return options.getToolArgs?.(event.toolCallId);
|
||||
}
|
||||
|
||||
function buildToolStartContent(toolName: string, args: unknown): ToolCallContent[] {
|
||||
if (!isCommandToolName(toolName)) {
|
||||
return [];
|
||||
|
||||
@@ -732,6 +732,7 @@ export class AgentSession {
|
||||
#planReferenceSent = false;
|
||||
#planReferencePath = "local://PLAN.md";
|
||||
#clientBridge: ClientBridge | undefined;
|
||||
#allowAcpAgentInitiatedTurns = false;
|
||||
/** Per-session memory of allow_always / reject_always decisions for gated tools. */
|
||||
#acpPermissionDecisions: Map<string, "allow_always" | "reject_always"> = new Map();
|
||||
|
||||
@@ -2777,9 +2778,15 @@ export class AgentSession {
|
||||
const ownerFilter = this.#agentId ? { ownerId: this.#agentId } : undefined;
|
||||
const before = manager.getDeliveryState(ownerFilter);
|
||||
if (before.queued === 0 && !before.delivering) return false;
|
||||
const drained = await manager.drainDeliveries({ timeoutMs: options?.timeoutMs, filter: ownerFilter });
|
||||
const after = manager.getDeliveryState(ownerFilter);
|
||||
return drained && (before.queued !== after.queued || before.delivering !== after.delivering);
|
||||
const previousAllowAcpAgentInitiatedTurns = this.#allowAcpAgentInitiatedTurns;
|
||||
this.#allowAcpAgentInitiatedTurns = true;
|
||||
try {
|
||||
const drained = await manager.drainDeliveries({ timeoutMs: options?.timeoutMs, filter: ownerFilter });
|
||||
const after = manager.getDeliveryState(ownerFilter);
|
||||
return drained && (before.queued !== after.queued || before.delivering !== after.delivering);
|
||||
} finally {
|
||||
this.#allowAcpAgentInitiatedTurns = previousAllowAcpAgentInitiatedTurns;
|
||||
}
|
||||
}
|
||||
|
||||
/** Most recent assistant message in agent state. */
|
||||
@@ -4419,7 +4426,7 @@ export class AgentSession {
|
||||
|
||||
if (options?.deliverAs === "nextTurn") {
|
||||
if (options?.triggerTurn) {
|
||||
if (this.#clientBridge?.deferAgentInitiatedTurns) {
|
||||
if (this.#clientBridge?.deferAgentInitiatedTurns && !this.#allowAcpAgentInitiatedTurns) {
|
||||
this.#queueHiddenNextTurnMessage(appMessage, false);
|
||||
return;
|
||||
}
|
||||
@@ -4438,7 +4445,7 @@ export class AgentSession {
|
||||
}
|
||||
|
||||
if (options?.triggerTurn) {
|
||||
if (this.#clientBridge?.deferAgentInitiatedTurns) {
|
||||
if (this.#clientBridge?.deferAgentInitiatedTurns && !this.#allowAcpAgentInitiatedTurns) {
|
||||
this.#queueHiddenNextTurnMessage(appMessage, false);
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -888,7 +888,7 @@ describe("ACP agent", () => {
|
||||
expect.objectContaining({
|
||||
sessionUpdate: "tool_call",
|
||||
toolCallId: "toolu_custom",
|
||||
rawInput: { input: "raw custom payload" },
|
||||
rawInput: "raw custom payload",
|
||||
}),
|
||||
);
|
||||
|
||||
|
||||
@@ -1,13 +1,18 @@
|
||||
import { describe, expect, it } from "bun:test";
|
||||
import * as fs from "node:fs";
|
||||
import * as os from "node:os";
|
||||
import path from "node:path";
|
||||
import type { SessionNotification } from "@agentclientprotocol/sdk";
|
||||
import type { AgentSideConnection, SessionNotification } from "@agentclientprotocol/sdk";
|
||||
import { zSessionNotification } from "@agentclientprotocol/sdk/dist/schema/zod.gen.js";
|
||||
import type { Model } from "@oh-my-pi/pi-ai";
|
||||
import { AcpAgent } from "../src/modes/acp/acp-agent";
|
||||
import {
|
||||
buildToolCallStartUpdate,
|
||||
mapAgentSessionEventToAcpSessionUpdates,
|
||||
normalizeReplayToolArguments,
|
||||
} from "../src/modes/acp/acp-event-mapper";
|
||||
import type { AgentSessionEvent } from "../src/session/agent-session";
|
||||
import type { AgentSession, AgentSessionEvent } from "../src/session/agent-session";
|
||||
import { SessionManager } from "../src/session/session-manager";
|
||||
import { expectAcpStructure, expectAcpStructureRejects } from "./helpers/acp-schema";
|
||||
|
||||
function makeAssistantMessage(text: string) {
|
||||
@@ -41,6 +46,55 @@ function expectAcpNotifications(updates: SessionNotification[]): void {
|
||||
}
|
||||
}
|
||||
|
||||
const TEST_MODEL: Model = {
|
||||
id: "claude-sonnet-4-20250514",
|
||||
name: "Claude Sonnet",
|
||||
api: "anthropic-messages",
|
||||
provider: "anthropic",
|
||||
baseUrl: "https://example.invalid",
|
||||
reasoning: true,
|
||||
input: ["text", "image"],
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 },
|
||||
contextWindow: 200_000,
|
||||
maxTokens: 8_192,
|
||||
};
|
||||
|
||||
class ReplayTestSession {
|
||||
sessionManager: SessionManager;
|
||||
sessionId: string;
|
||||
model: Model | undefined = TEST_MODEL;
|
||||
thinkingLevel: string | undefined;
|
||||
customCommands: [] = [];
|
||||
skills: [] = [];
|
||||
extensionRunner = undefined;
|
||||
settings = { get: (_key: string) => false };
|
||||
|
||||
constructor(cwd: string, sessionDir?: string) {
|
||||
this.sessionManager = SessionManager.create(cwd, sessionDir);
|
||||
this.sessionId = this.sessionManager.getSessionId();
|
||||
}
|
||||
|
||||
getAvailableModels(): Model[] {
|
||||
return [TEST_MODEL];
|
||||
}
|
||||
|
||||
getAvailableThinkingLevels(): ReadonlyArray<string> {
|
||||
return [];
|
||||
}
|
||||
|
||||
getPlanModeState(): undefined {
|
||||
return undefined;
|
||||
}
|
||||
|
||||
setClientBridge(_bridge: unknown): void {}
|
||||
|
||||
subscribe(_listener: (event: AgentSessionEvent) => void): () => void {
|
||||
return () => {};
|
||||
}
|
||||
|
||||
async refreshMCPTools(_tools: unknown): Promise<void> {}
|
||||
}
|
||||
|
||||
describe("ACP event mapper", () => {
|
||||
it("attaches a stable messageId to live assistant chunks", () => {
|
||||
const assistantMessage = makeAssistantMessage("chunk");
|
||||
@@ -331,6 +385,37 @@ describe("ACP event mapper", () => {
|
||||
expect(update.content).toContainEqual({ type: "terminal", terminalId: "term-1" });
|
||||
});
|
||||
|
||||
it("preserves command text when a command tool final update replaces content", () => {
|
||||
const updates = mapAgentSessionEventToAcpSessionUpdates(
|
||||
{
|
||||
type: "tool_execution_end",
|
||||
toolCallId: "tc-terminal-final-command",
|
||||
toolName: "bash",
|
||||
isError: false,
|
||||
result: {
|
||||
content: [{ type: "text", text: "done" }],
|
||||
details: { terminalId: "term-1" },
|
||||
},
|
||||
} as AgentSessionEvent,
|
||||
"session-1",
|
||||
{
|
||||
getToolArgs: toolCallId =>
|
||||
toolCallId === "tc-terminal-final-command" ? { command: "npm run check" } : undefined,
|
||||
},
|
||||
);
|
||||
|
||||
expect(updates).toHaveLength(1);
|
||||
expectAcpNotifications(updates);
|
||||
const update = updates[0]!.update as {
|
||||
sessionUpdate: string;
|
||||
content?: Array<{ type: string; terminalId?: string; content?: { type: string; text?: string } }>;
|
||||
};
|
||||
expect(update.sessionUpdate).toBe("tool_call_update");
|
||||
expect(update.content).toContainEqual({ type: "content", content: { type: "text", text: "$ npm run check" } });
|
||||
expect(update.content).toContainEqual({ type: "content", content: { type: "text", text: "done" } });
|
||||
expect(update.content).toContainEqual({ type: "terminal", terminalId: "term-1" });
|
||||
});
|
||||
|
||||
it("keeps terminal content alongside readable error and message fields", () => {
|
||||
const errorUpdates = mapAgentSessionEventToAcpSessionUpdates(
|
||||
{
|
||||
@@ -494,6 +579,90 @@ describe("ACP event mapper", () => {
|
||||
}
|
||||
});
|
||||
|
||||
it("replays assistant tool_use input through the ACP dispatcher without wrapping", async () => {
|
||||
const root = await fs.promises.mkdtemp(path.join(os.tmpdir(), "omp-acp-replay-contract-"));
|
||||
const cwd = path.join(root, "cwd");
|
||||
const sessionDir = path.join(root, "sessions");
|
||||
const initialSessionDir = path.join(root, "initial-session");
|
||||
const updates: SessionNotification[] = [];
|
||||
const sessions: ReplayTestSession[] = [];
|
||||
const abortController = new AbortController();
|
||||
try {
|
||||
await fs.promises.mkdir(cwd, { recursive: true });
|
||||
const connection = {
|
||||
sessionUpdate: async (notification: SessionNotification) => {
|
||||
updates.push(notification);
|
||||
},
|
||||
signal: abortController.signal,
|
||||
closed: Promise.resolve(),
|
||||
} as unknown as AgentSideConnection;
|
||||
const agent = new AcpAgent(
|
||||
connection,
|
||||
async (sessionCwd: string) => {
|
||||
const session = new ReplayTestSession(sessionCwd, sessionDir);
|
||||
sessions.push(session);
|
||||
return session as unknown as AgentSession;
|
||||
},
|
||||
new ReplayTestSession(cwd, initialSessionDir) as unknown as AgentSession,
|
||||
);
|
||||
const created = await agent.newSession({ cwd, mcpServers: [] });
|
||||
const session = sessions[0]!;
|
||||
session.sessionManager.appendMessage({
|
||||
role: "assistant",
|
||||
content: [
|
||||
{
|
||||
type: "tool_use",
|
||||
id: "toolu_replay_input",
|
||||
name: "bash",
|
||||
input: { command: "echo hi" },
|
||||
},
|
||||
],
|
||||
usage: {
|
||||
input: 10,
|
||||
output: 5,
|
||||
cacheRead: 0,
|
||||
cacheWrite: 0,
|
||||
totalTokens: 15,
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
},
|
||||
api: "anthropic-messages",
|
||||
provider: "anthropic",
|
||||
model: "claude-sonnet-4-20250514",
|
||||
stopReason: "stop",
|
||||
timestamp: Date.now(),
|
||||
} as unknown as Parameters<SessionManager["appendMessage"]>[0]);
|
||||
session.sessionManager.appendMessage({
|
||||
role: "toolResult",
|
||||
toolCallId: "toolu_replay_input",
|
||||
toolName: "bash",
|
||||
content: [{ type: "text", text: "done" }],
|
||||
details: { terminalId: "term-replay" },
|
||||
isError: false,
|
||||
timestamp: Date.now(),
|
||||
});
|
||||
|
||||
updates.length = 0;
|
||||
await agent.loadSession({ sessionId: created.sessionId, cwd, mcpServers: [] });
|
||||
|
||||
expectAcpNotifications(updates);
|
||||
const toolCall = updates.find(update => update.update.sessionUpdate === "tool_call")?.update as
|
||||
| { rawInput?: unknown; content?: unknown }
|
||||
| undefined;
|
||||
const finalUpdate = updates.find(update => update.update.sessionUpdate === "tool_call_update")?.update as
|
||||
| { content?: unknown }
|
||||
| undefined;
|
||||
|
||||
expect(toolCall?.rawInput).toEqual({ command: "echo hi" });
|
||||
expect(toolCall?.rawInput).not.toEqual({ input: { command: "echo hi" } });
|
||||
expect(toolCall?.content).toEqual([{ type: "content", content: { type: "text", text: "$ echo hi" } }]);
|
||||
expect(finalUpdate?.content).toContainEqual({ type: "content", content: { type: "text", text: "$ echo hi" } });
|
||||
expect(finalUpdate?.content).toContainEqual({ type: "content", content: { type: "text", text: "done" } });
|
||||
expect(finalUpdate?.content).toContainEqual({ type: "terminal", terminalId: "term-replay" });
|
||||
} finally {
|
||||
abortController.abort();
|
||||
await fs.promises.rm(root, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
it("builds replayed bash tool calls from JSON string arguments", () => {
|
||||
const replayArgs = normalizeReplayToolArguments(JSON.stringify({ command: "npm test", cwd: "/repo" }));
|
||||
const update = buildToolCallStartUpdate({
|
||||
|
||||
@@ -709,16 +709,3 @@ it("read tool: requestPermission is never called for non-gated tools", async ()
|
||||
expect(permissionSpy).toHaveBeenCalledTimes(0);
|
||||
expect(readTool.executeCalls).toBe(1);
|
||||
});
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// 5. No bridge → original tool object identity preserved (no wrapping)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
it("no bridge: original tool object is returned unchanged", async () => {
|
||||
const bashTool = makeFakeTool("bash");
|
||||
session = await createSession([bashTool]); // no bridge
|
||||
|
||||
await session.setActiveToolsByName(["bash"]);
|
||||
const activeBash = session.agent.state.tools.find(t => t.name === "bash");
|
||||
expect(activeBash).toBe(bashTool);
|
||||
});
|
||||
|
||||
@@ -10,6 +10,7 @@ import { Agent, AgentBusyError, type AgentTool } from "@oh-my-pi/pi-agent-core";
|
||||
import { type AssistantMessage, getBundledModel, type Message, type ToolCall } from "@oh-my-pi/pi-ai";
|
||||
import { createMockModel } from "@oh-my-pi/pi-ai/providers/mock";
|
||||
import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream";
|
||||
import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async";
|
||||
import type { Rule } from "@oh-my-pi/pi-coding-agent/capability/rule";
|
||||
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
|
||||
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
||||
@@ -45,6 +46,7 @@ describe("AgentSession concurrent prompt guard", () => {
|
||||
fs.rmSync(tempDir, { recursive: true });
|
||||
}
|
||||
vi.restoreAllMocks();
|
||||
AsyncJobManager.resetForTests();
|
||||
});
|
||||
|
||||
async function createSession() {
|
||||
@@ -334,9 +336,9 @@ describe("AgentSession concurrent prompt guard", () => {
|
||||
|
||||
const sessionManager = SessionManager.inMemory();
|
||||
const settings = Settings.isolated();
|
||||
const authStorage = await AuthStorage.create(":memory:");
|
||||
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-acp-idle.db"));
|
||||
authStorages.push(authStorage);
|
||||
const modelRegistry = new ModelRegistry(authStorage);
|
||||
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-acp-idle.yml"));
|
||||
authStorage.setRuntimeApiKey("anthropic", "test-key");
|
||||
|
||||
session = new AgentSession({
|
||||
@@ -383,6 +385,161 @@ describe("AgentSession concurrent prompt guard", () => {
|
||||
}),
|
||||
).toBe(true);
|
||||
});
|
||||
|
||||
it("runs drained ACP async completions as owned follow-up turns despite deferred client turns", async () => {
|
||||
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
|
||||
const mock = createMockModel({ handler: () => ({ content: ["Done"] }) });
|
||||
const agent = new Agent({
|
||||
getApiKey: () => "test-key",
|
||||
initialState: {
|
||||
model,
|
||||
systemPrompt: ["Test"],
|
||||
tools: [],
|
||||
},
|
||||
convertToLlm,
|
||||
streamFn: mock.stream,
|
||||
});
|
||||
|
||||
const sessionManager = SessionManager.inMemory();
|
||||
const settings = Settings.isolated();
|
||||
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-acp-async.db"));
|
||||
authStorages.push(authStorage);
|
||||
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-acp-async.yml"));
|
||||
authStorage.setRuntimeApiKey("anthropic", "test-key");
|
||||
|
||||
const ownerId = "acp-session-a";
|
||||
const deliveryGate = Promise.withResolvers<void>();
|
||||
let deliveryStarted = false;
|
||||
const asyncJobManager = new AsyncJobManager({
|
||||
maxRunningJobs: 2,
|
||||
retentionMs: 1_000,
|
||||
onJobComplete: async () => {
|
||||
deliveryStarted = true;
|
||||
await deliveryGate.promise;
|
||||
await session.sendCustomMessage(
|
||||
{
|
||||
customType: "async-result",
|
||||
content: "Background result",
|
||||
display: true,
|
||||
attribution: "agent",
|
||||
},
|
||||
{ deliverAs: "followUp", triggerTurn: true },
|
||||
);
|
||||
},
|
||||
});
|
||||
AsyncJobManager.setInstance(asyncJobManager);
|
||||
|
||||
session = new AgentSession({
|
||||
agent,
|
||||
sessionManager,
|
||||
settings,
|
||||
modelRegistry,
|
||||
agentId: ownerId,
|
||||
ownedAsyncJobManager: asyncJobManager,
|
||||
});
|
||||
session.setClientBridge({
|
||||
capabilities: {},
|
||||
deferAgentInitiatedTurns: true,
|
||||
});
|
||||
|
||||
await session.prompt("First message");
|
||||
expect(session.isStreaming).toBe(false);
|
||||
const callsAfterFirstPrompt = mock.calls.length;
|
||||
|
||||
try {
|
||||
asyncJobManager.register("bash", "owned job", async () => "Background result", {
|
||||
id: "owned-job",
|
||||
ownerId,
|
||||
});
|
||||
await waitFor(() => deliveryStarted);
|
||||
|
||||
const drainedPromise = session.drainAsyncJobDeliveriesForAcp({ timeoutMs: 1_000 });
|
||||
await waitFor(() => asyncJobManager.getDeliveryState({ ownerId }).delivering);
|
||||
deliveryGate.resolve();
|
||||
|
||||
await expect(drainedPromise).resolves.toBe(true);
|
||||
await session.waitForIdle();
|
||||
|
||||
expect(mock.calls).toHaveLength(callsAfterFirstPrompt + 1);
|
||||
expect(
|
||||
mock.calls.at(-1)?.context.messages.some(message => {
|
||||
if (typeof message.content === "string") {
|
||||
return message.content.includes("Background result");
|
||||
}
|
||||
|
||||
return message.content.some(
|
||||
content => content.type === "text" && content.text.includes("Background result"),
|
||||
);
|
||||
}),
|
||||
).toBe(true);
|
||||
} finally {
|
||||
deliveryGate.resolve();
|
||||
}
|
||||
});
|
||||
|
||||
it("scopes ACP async job snapshots and drains to the owning session id", async () => {
|
||||
const model = getBundledModel("anthropic", "claude-sonnet-4-5")!;
|
||||
const authStorage = await AuthStorage.create(path.join(tempDir, "testauth-acp-scope.db"));
|
||||
authStorages.push(authStorage);
|
||||
const modelRegistry = new ModelRegistry(authStorage, path.join(tempDir, "models-acp-scope.yml"));
|
||||
authStorage.setRuntimeApiKey("anthropic", "test-key");
|
||||
const settings = Settings.isolated();
|
||||
const deliveryGate = Promise.withResolvers<void>();
|
||||
const delivered: string[] = [];
|
||||
const started = new Set<string>();
|
||||
const asyncJobManager = new AsyncJobManager({
|
||||
maxRunningJobs: 3,
|
||||
retentionMs: 1_000,
|
||||
onJobComplete: async jobId => {
|
||||
started.add(jobId);
|
||||
if (jobId === "job-a") {
|
||||
await deliveryGate.promise;
|
||||
}
|
||||
delivered.push(jobId);
|
||||
},
|
||||
});
|
||||
AsyncJobManager.setInstance(asyncJobManager);
|
||||
|
||||
const agentA = new Agent({
|
||||
getApiKey: () => "test-key",
|
||||
initialState: { model, systemPrompt: ["Test"], tools: [] },
|
||||
streamFn: createMockModel({ handler: () => ({ content: ["Done"] }) }).stream,
|
||||
});
|
||||
const agentB = new Agent({
|
||||
getApiKey: () => "test-key",
|
||||
initialState: { model, systemPrompt: ["Test"], tools: [] },
|
||||
streamFn: createMockModel({ handler: () => ({ content: ["Done"] }) }).stream,
|
||||
});
|
||||
const sessionB = new AgentSession({
|
||||
agent: agentB,
|
||||
sessionManager: SessionManager.inMemory(),
|
||||
settings,
|
||||
modelRegistry,
|
||||
agentId: "acp-session-b",
|
||||
});
|
||||
session = new AgentSession({
|
||||
agent: agentA,
|
||||
sessionManager: SessionManager.inMemory(),
|
||||
settings,
|
||||
modelRegistry,
|
||||
agentId: "acp-session-a",
|
||||
ownedAsyncJobManager: asyncJobManager,
|
||||
});
|
||||
|
||||
try {
|
||||
asyncJobManager.register("bash", "A", async () => "A", { id: "job-a", ownerId: "acp-session-a" });
|
||||
await waitFor(() => started.has("job-a"));
|
||||
asyncJobManager.register("bash", "B", async () => "B", { id: "job-b", ownerId: "acp-session-b" });
|
||||
await waitFor(() => asyncJobManager.getDeliveryState({ ownerId: "acp-session-b" }).queued > 0);
|
||||
|
||||
expect(sessionB.getAsyncJobSnapshot()?.delivery.pendingJobIds).not.toContain("job-a");
|
||||
await expect(sessionB.drainAsyncJobDeliveriesForAcp({ timeoutMs: 1_000 })).resolves.toBe(true);
|
||||
expect(delivered).toEqual(["job-b"]);
|
||||
} finally {
|
||||
deliveryGate.resolve();
|
||||
await sessionB.dispose();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
describe("AgentSession TTSR resume gate", () => {
|
||||
|
||||
@@ -230,35 +230,104 @@ describe("AsyncJobManager", () => {
|
||||
});
|
||||
|
||||
test("scoped delivery drain returns once matching owner deliveries finish", async () => {
|
||||
let firstOwnerAttempts = 0;
|
||||
const secondOwnerCompletions: Array<{ jobId: string; text: string }> = [];
|
||||
let mainJobId = "";
|
||||
let releaseMainDelivery = (): void => {};
|
||||
let notifyMainDeliveryStarted = (): void => {};
|
||||
const mainDeliveryStarted = new Promise<void>(resolve => {
|
||||
notifyMainDeliveryStarted = resolve;
|
||||
});
|
||||
const mainDeliveryReleased = new Promise<void>(resolve => {
|
||||
releaseMainDelivery = resolve;
|
||||
});
|
||||
const subagentCompletions: Array<{ jobId: string; text: string }> = [];
|
||||
const manager = new AsyncJobManager({
|
||||
onJobComplete: async (jobId, text, job) => {
|
||||
if (job?.ownerId === "0-Main") {
|
||||
firstOwnerAttempts++;
|
||||
throw new Error("first owner delivery retry");
|
||||
retentionMs: 0,
|
||||
onJobComplete: async (jobId, text) => {
|
||||
if (jobId === mainJobId) {
|
||||
notifyMainDeliveryStarted();
|
||||
await mainDeliveryReleased;
|
||||
return;
|
||||
}
|
||||
secondOwnerCompletions.push({ jobId, text });
|
||||
subagentCompletions.push({ jobId, text });
|
||||
},
|
||||
});
|
||||
|
||||
manager.register("task", "main job", async () => "main result", { ownerId: "0-Main" });
|
||||
mainJobId = manager.register("task", "main job", async () => "main result", { ownerId: "0-Main" });
|
||||
const targetJobId = manager.register("task", "subagent job", async () => "subagent result", {
|
||||
ownerId: "3-AuthLoader",
|
||||
});
|
||||
await manager.waitForAll();
|
||||
const firstAttemptDeadline = Date.now() + 2_000;
|
||||
while (firstOwnerAttempts === 0) {
|
||||
if (Date.now() >= firstAttemptDeadline) throw new Error("Timed out waiting for first owner delivery attempt");
|
||||
await Bun.sleep(5);
|
||||
}
|
||||
await mainDeliveryStarted;
|
||||
|
||||
expect(manager.hasPendingDeliveries({ ownerId: "0-Main" })).toBe(true);
|
||||
const drained = await manager.drainDeliveries({ timeoutMs: 50, filter: { ownerId: "3-AuthLoader" } });
|
||||
|
||||
expect(drained).toBe(true);
|
||||
expect(secondOwnerCompletions).toEqual([{ jobId: targetJobId, text: "subagent result" }]);
|
||||
expect(subagentCompletions).toEqual([{ jobId: targetJobId, text: "subagent result" }]);
|
||||
expect(manager.hasPendingDeliveries({ ownerId: "3-AuthLoader" })).toBe(false);
|
||||
expect(manager.hasPendingDeliveries({ ownerId: "0-Main" })).toBe(true);
|
||||
|
||||
expect(manager.acknowledgeDeliveries([mainJobId])).toBe(0);
|
||||
expect(manager.hasPendingDeliveries({ ownerId: "0-Main" })).toBe(false);
|
||||
releaseMainDelivery();
|
||||
await Bun.sleep(0);
|
||||
});
|
||||
|
||||
test("scoped delivery drain times out while a matching delivery callback is in flight", async () => {
|
||||
let mainJobId = "";
|
||||
let targetJobId = "";
|
||||
let releaseMainDelivery = (): void => {};
|
||||
let notifyMainDeliveryStarted = (): void => {};
|
||||
let releaseTargetDelivery = (): void => {};
|
||||
let notifyTargetDeliveryStarted = (): void => {};
|
||||
const mainDeliveryStarted = new Promise<void>(resolve => {
|
||||
notifyMainDeliveryStarted = resolve;
|
||||
});
|
||||
const mainDeliveryReleased = new Promise<void>(resolve => {
|
||||
releaseMainDelivery = resolve;
|
||||
});
|
||||
const targetDeliveryStarted = new Promise<void>(resolve => {
|
||||
notifyTargetDeliveryStarted = resolve;
|
||||
});
|
||||
const targetDeliveryReleased = new Promise<void>(resolve => {
|
||||
releaseTargetDelivery = resolve;
|
||||
});
|
||||
const completions: string[] = [];
|
||||
const manager = new AsyncJobManager({
|
||||
onJobComplete: async jobId => {
|
||||
if (jobId === mainJobId) {
|
||||
notifyMainDeliveryStarted();
|
||||
await mainDeliveryReleased;
|
||||
return;
|
||||
}
|
||||
if (jobId === targetJobId) {
|
||||
notifyTargetDeliveryStarted();
|
||||
await targetDeliveryReleased;
|
||||
completions.push(jobId);
|
||||
}
|
||||
},
|
||||
});
|
||||
|
||||
mainJobId = manager.register("task", "main job", async () => "main result", { ownerId: "0-Main" });
|
||||
targetJobId = manager.register("task", "subagent job", async () => "subagent result", {
|
||||
ownerId: "3-AuthLoader",
|
||||
});
|
||||
await manager.waitForAll();
|
||||
await mainDeliveryStarted;
|
||||
|
||||
const timedOut = await manager.drainDeliveries({ timeoutMs: 10, filter: { ownerId: "3-AuthLoader" } });
|
||||
await targetDeliveryStarted;
|
||||
|
||||
expect(timedOut).toBe(false);
|
||||
expect(manager.hasPendingDeliveries({ ownerId: "3-AuthLoader" })).toBe(true);
|
||||
expect(completions).toEqual([]);
|
||||
|
||||
releaseTargetDelivery();
|
||||
const drained = await manager.drainDeliveries({ timeoutMs: 200, filter: { ownerId: "3-AuthLoader" } });
|
||||
expect(drained).toBe(true);
|
||||
expect(completions).toEqual([targetJobId]);
|
||||
|
||||
releaseMainDelivery();
|
||||
expect(await manager.drainDeliveries({ timeoutMs: 200 })).toBe(true);
|
||||
});
|
||||
|
||||
test("cancelAll with ownerId only cancels matching jobs", async () => {
|
||||
|
||||
Reference in New Issue
Block a user