fix(coding-agent): trim rpc prompt completion scope
This commit is contained in:
@@ -70,6 +70,7 @@ Important edge behavior from runtime:
|
||||
- `prompt` and `abort_and_prompt` return immediate success, then may emit a later error response with the **same** id if async prompt scheduling fails.
|
||||
- `prompt` success responses may include `data.agentInvoked`. `false` means the prompt completed locally without an agent turn; `true` means an agent turn was scheduled; omitted means the host must use session events for completion.
|
||||
- `abort_and_prompt` does not currently emit `data.agentInvoked` or `prompt_result`; hosts should treat it as the legacy abort-then-schedule path and rely on session events or same-id scheduling errors.
|
||||
|
||||
## Command Schema (canonical)
|
||||
|
||||
`RpcCommand` is defined in `src/modes/rpc/rpc-types.ts`:
|
||||
|
||||
@@ -25,6 +25,8 @@
|
||||
- Added streaming speech vocalization: with `speech.enabled` on, the assistant speaks its reply through the speakers as it streams. Assistant text deltas are fed *directly into the engine's incremental text input* (Kokoro's `TextSplitterStream` via the worker) as they arrive, rather than pre-chunked in JS and synthesized one batch call per sentence — the engine owns sentence segmentation and emits one audio chunk per sentence. A single persistent player (`StreamingAudioPlayer`) drains those chunks **gaplessly** (raw 32-bit-float PCM piped to one `ffmpeg`→PulseAudio/ALSA process on Linux; interruptible per-file `afplay`/PowerShell `SoundPlayer` on macOS/Windows), replacing the spawn-a-player-per-sentence path that added latency and audible gaps. Overspeech is handled end to end: a new turn, a sent message, or an Esc/Ctrl+C interrupt stops playback **instantly** (the player process is killed rather than letting the current sentence finish); holding the push-to-talk key **ducks** the volume while you speak and restores it when you stop; and sequential utterances queue and drain in order instead of overlapping. `speech.mode` (`all` | `assistant` | `yield`, default `assistant`) picks what is spoken — `all` adds thinking, `yield` speaks only the final message at turn end — and `speech.voice` selects the Kokoro voice. `ask`-tool questions are spoken in every mode. Synthesis reuses the local Kokoro engine (`tts.localModel`) through a new streaming synthesis path (`TtsClient.synthesizeStream`) that pushes text in and streams audio chunks back over the worker protocol.
|
||||
- Added live (streaming) speech-to-text: with `stt.enabled` on, transcription now appears in the composer *as you speak* instead of all at once after you stop. The recorder streams raw 16 kHz mono PCM from sox/ffmpeg/arecord stdout to the warm STT worker, where an energy-based endpointer (no extra model) splits speech into segments at natural pauses; each finalized segment is committed into the editor while the in-progress segment shows a live volatile preview that refreshes in place and is kept out of the undo history. Works with both the default Parakeet (sherpa-onnx) and the Whisper (transformers.js) tiers. Recorders that cannot stream to a pipe (the Windows PowerShell mci fallback) transparently fall back to single-shot transcription.
|
||||
|
||||
- Added `skills.enableAgentsUser` and `skills.enableAgentsProject` settings (default on) so the canonical OMP-native `~/.agent[s]/skills` and project-walkup `.agent[s]/skills` are configurable independently from the third-party Claude/Codex/Pi toggles.
|
||||
|
||||
- Added RPC prompt lifecycle hints so hosts can distinguish scheduled agent turns from local-only slash commands via `data.agentInvoked` and `prompt_result`.
|
||||
|
||||
### Changed
|
||||
|
||||
@@ -685,13 +685,10 @@ export class AcpAgent implements Agent {
|
||||
}
|
||||
}
|
||||
|
||||
#trackExtensionUserMessage(record: ManagedSessionRecord, task: Promise<boolean>): void {
|
||||
const tracked = task.then(
|
||||
() => undefined,
|
||||
(error: unknown) => {
|
||||
logger.warn("ACP extension sendUserMessage failed", { error });
|
||||
},
|
||||
);
|
||||
#trackExtensionUserMessage(record: ManagedSessionRecord, task: Promise<void>): void {
|
||||
const tracked = task.catch((error: unknown) => {
|
||||
logger.warn("ACP extension sendUserMessage failed", { error });
|
||||
});
|
||||
record.extensionUserMessageTasks.add(tracked);
|
||||
void tracked.finally(() => {
|
||||
record.extensionUserMessageTasks.delete(tracked);
|
||||
|
||||
@@ -23,7 +23,6 @@ import type {
|
||||
RpcHostToolDefinition,
|
||||
RpcHostToolResult,
|
||||
RpcHostToolUpdate,
|
||||
RpcPromptResultFrame,
|
||||
RpcResponse,
|
||||
RpcSessionState,
|
||||
RpcSubagentEventFrame,
|
||||
@@ -40,12 +39,6 @@ type DistributiveOmit<T, K extends keyof T> = T extends unknown ? Omit<T, K> : n
|
||||
/** RpcCommand without the id field (for internal send) */
|
||||
type RpcCommandBody = DistributiveOmit<RpcCommand, "id">;
|
||||
|
||||
type PromptCompletion = {
|
||||
promise: Promise<void>;
|
||||
resolve: () => void;
|
||||
settled: boolean;
|
||||
};
|
||||
|
||||
export interface RpcClientOptions {
|
||||
/** Path to the CLI entry point (default: searches for dist/cli.js) */
|
||||
cliPath?: string;
|
||||
@@ -176,15 +169,6 @@ function isRpcAvailableCommandsUpdateFrame(value: unknown): value is RpcAvailabl
|
||||
return value.type === "available_commands_update" && Array.isArray(value.commands);
|
||||
}
|
||||
|
||||
function isRpcPromptResultFrame(value: unknown): value is RpcPromptResultFrame {
|
||||
if (!isRecord(value)) return false;
|
||||
return (
|
||||
value.type === "prompt_result" &&
|
||||
(value.id === undefined || typeof value.id === "string") &&
|
||||
typeof value.agentInvoked === "boolean"
|
||||
);
|
||||
}
|
||||
|
||||
function isRpcHostToolCallRequest(value: unknown): value is RpcHostToolCallRequest {
|
||||
if (!isRecord(value)) return false;
|
||||
return (
|
||||
@@ -234,9 +218,6 @@ export class RpcClient {
|
||||
#requestId = 0;
|
||||
#extensionUiListeners: Set<(req: RpcExtensionUIRequest) => void> = new Set();
|
||||
#abortController = new AbortController();
|
||||
#promptCompletions = new Map<string, PromptCompletion>();
|
||||
#lastPromptCompletion: PromptCompletion | undefined;
|
||||
#promptCompletionListeners = new Set<() => void>();
|
||||
|
||||
constructor(private options: RpcClientOptions = {}) {
|
||||
this.#customTools = [...(options.customTools ?? [])];
|
||||
@@ -436,7 +417,7 @@ export class RpcClient {
|
||||
* Use waitForIdle() to wait for completion.
|
||||
*/
|
||||
async prompt(message: string, images?: ImageContent[]): Promise<void> {
|
||||
await this.#sendPrompt(message, images);
|
||||
await this.#send({ type: "prompt", message, images });
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -757,26 +738,18 @@ export class RpcClient {
|
||||
|
||||
/**
|
||||
* Wait for agent to become idle (no streaming).
|
||||
* Resolves on `agent_end` or a local-only `prompt_result`.
|
||||
* Resolves when agent_end event is received.
|
||||
*/
|
||||
waitForIdle(timeout = 60000): Promise<void> {
|
||||
const completion = this.#lastPromptCompletion;
|
||||
if (completion) {
|
||||
return this.#withIdleTimeout(
|
||||
completion.promise,
|
||||
timeout,
|
||||
`Timeout waiting for agent to become idle. Stderr: ${this.#process?.peekStderr() ?? ""}`,
|
||||
);
|
||||
}
|
||||
|
||||
const { promise, resolve, reject } = Promise.withResolvers<void>();
|
||||
let settled = false;
|
||||
const unsubscribe = this.#onPromptCompletion(() => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
unsubscribe();
|
||||
clearTimeout(timeoutId);
|
||||
resolve();
|
||||
const unsubscribe = this.onEvent(event => {
|
||||
if (event.type === "agent_end") {
|
||||
settled = true;
|
||||
unsubscribe();
|
||||
clearTimeout(timeoutId);
|
||||
resolve();
|
||||
}
|
||||
});
|
||||
|
||||
const timeoutId = this.#startTimeout(timeout, () => {
|
||||
@@ -795,34 +768,20 @@ export class RpcClient {
|
||||
const { promise, resolve, reject } = Promise.withResolvers<AgentEvent[]>();
|
||||
const events: AgentEvent[] = [];
|
||||
let settled = false;
|
||||
let unsubscribeCompletion = () => {};
|
||||
const finish = () => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
unsubscribe();
|
||||
unsubscribeCompletion();
|
||||
clearTimeout(timeoutId);
|
||||
resolve(events);
|
||||
};
|
||||
const unsubscribe = this.onEvent(event => {
|
||||
events.push(event);
|
||||
if (event.type === "agent_end") {
|
||||
finish();
|
||||
settled = true;
|
||||
unsubscribe();
|
||||
clearTimeout(timeoutId);
|
||||
resolve(events);
|
||||
}
|
||||
});
|
||||
|
||||
const completion = this.#lastPromptCompletion;
|
||||
if (completion) {
|
||||
void completion.promise.then(finish);
|
||||
} else {
|
||||
unsubscribeCompletion = this.#onPromptCompletion(finish);
|
||||
}
|
||||
|
||||
const timeoutId = this.#startTimeout(timeout, () => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
unsubscribe();
|
||||
unsubscribeCompletion();
|
||||
reject(new Error(`Timeout collecting events. Stderr: ${this.#process?.peekStderr() ?? ""}`));
|
||||
});
|
||||
return promise;
|
||||
@@ -832,122 +791,12 @@ export class RpcClient {
|
||||
* Send prompt and wait for completion, returning all events.
|
||||
*/
|
||||
async promptAndWait(message: string, images?: ImageContent[], timeout = 60000): Promise<AgentEvent[]> {
|
||||
const events: AgentEvent[] = [];
|
||||
const unsubscribe = this.onEvent(event => {
|
||||
events.push(event);
|
||||
});
|
||||
try {
|
||||
await this.prompt(message, images);
|
||||
await this.waitForIdle(timeout);
|
||||
return events;
|
||||
} finally {
|
||||
unsubscribe();
|
||||
}
|
||||
const eventsPromise = this.collectEvents(timeout);
|
||||
await this.prompt(message, images);
|
||||
return eventsPromise;
|
||||
}
|
||||
|
||||
// =========================================================================
|
||||
async #sendPrompt(message: string, images?: ImageContent[]): Promise<void> {
|
||||
const id = this.#nextRequestId();
|
||||
this.#createPromptCompletion(id);
|
||||
try {
|
||||
const response = await this.#send({ type: "prompt", message, images }, 30_000, id);
|
||||
if (!response.success) {
|
||||
this.#resolvePromptCompletion(id);
|
||||
return;
|
||||
}
|
||||
if (response.command === "prompt" && response.data?.agentInvoked === false) {
|
||||
this.#resolvePromptCompletion(id);
|
||||
}
|
||||
} catch (error) {
|
||||
this.#resolvePromptCompletion(id);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
#nextRequestId(): string {
|
||||
return `req_${++this.#requestId}`;
|
||||
}
|
||||
|
||||
#createPromptCompletion(id: string): PromptCompletion {
|
||||
const { promise, resolve } = Promise.withResolvers<void>();
|
||||
const completion: PromptCompletion = {
|
||||
promise,
|
||||
resolve: () => {
|
||||
if (completion.settled) return;
|
||||
completion.settled = true;
|
||||
resolve();
|
||||
this.#notifyPromptCompletion();
|
||||
},
|
||||
settled: false,
|
||||
};
|
||||
this.#promptCompletions.set(id, completion);
|
||||
this.#lastPromptCompletion = completion;
|
||||
return completion;
|
||||
}
|
||||
|
||||
#resolvePromptCompletion(id: string | undefined): void {
|
||||
if (!id) {
|
||||
this.#resolveAllPromptCompletions();
|
||||
return;
|
||||
}
|
||||
const completion = this.#promptCompletions.get(id);
|
||||
if (!completion) {
|
||||
this.#notifyPromptCompletion();
|
||||
return;
|
||||
}
|
||||
this.#promptCompletions.delete(id);
|
||||
completion.resolve();
|
||||
}
|
||||
|
||||
#resolveAllPromptCompletions(): void {
|
||||
if (this.#promptCompletions.size === 0) {
|
||||
this.#notifyPromptCompletion();
|
||||
return;
|
||||
}
|
||||
const completions = Array.from(this.#promptCompletions.values());
|
||||
this.#promptCompletions.clear();
|
||||
for (const completion of completions) {
|
||||
completion.resolve();
|
||||
}
|
||||
}
|
||||
|
||||
#onPromptCompletion(listener: () => void): () => void {
|
||||
this.#promptCompletionListeners.add(listener);
|
||||
return () => {
|
||||
this.#promptCompletionListeners.delete(listener);
|
||||
};
|
||||
}
|
||||
|
||||
#notifyPromptCompletion(): void {
|
||||
for (const listener of [...this.#promptCompletionListeners]) {
|
||||
listener();
|
||||
}
|
||||
}
|
||||
|
||||
#withIdleTimeout<T>(source: Promise<T>, timeoutMs: number, timeoutMessage: string): Promise<T> {
|
||||
const { promise, resolve, reject } = Promise.withResolvers<T>();
|
||||
let settled = false;
|
||||
const timeoutId = this.#startTimeout(timeoutMs, () => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
reject(new Error(timeoutMessage));
|
||||
});
|
||||
void source.then(
|
||||
value => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
clearTimeout(timeoutId);
|
||||
resolve(value);
|
||||
},
|
||||
error => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
clearTimeout(timeoutId);
|
||||
reject(error);
|
||||
},
|
||||
);
|
||||
return promise;
|
||||
}
|
||||
// Internal
|
||||
// =========================================================================
|
||||
|
||||
@@ -963,13 +812,6 @@ export class RpcClient {
|
||||
}
|
||||
}
|
||||
|
||||
if (isRpcPromptResultFrame(data)) {
|
||||
if (!data.agentInvoked) {
|
||||
this.#resolvePromptCompletion(data.id);
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
if (isRpcHostToolCallRequest(data)) {
|
||||
void this.#handleHostToolCall(data);
|
||||
return;
|
||||
@@ -1026,16 +868,14 @@ export class RpcClient {
|
||||
for (const listener of this.#eventListeners) {
|
||||
listener(data);
|
||||
}
|
||||
|
||||
if (data.type === "agent_end") {
|
||||
this.#resolveAllPromptCompletions();
|
||||
}
|
||||
}
|
||||
|
||||
#send(command: RpcCommandBody, timeoutMs = 30_000, id = this.#nextRequestId()): Promise<RpcResponse> {
|
||||
#send(command: RpcCommandBody, timeoutMs = 30_000): Promise<RpcResponse> {
|
||||
if (!this.#process?.stdin) {
|
||||
throw new Error("Client not started");
|
||||
}
|
||||
|
||||
const id = `req_${++this.#requestId}`;
|
||||
const fullCommand = { ...command, id } as RpcCommand;
|
||||
const { promise, resolve, reject } = Promise.withResolvers<RpcResponse>();
|
||||
let settled = false;
|
||||
|
||||
@@ -80,7 +80,7 @@ export type RpcSessionChangeResult =
|
||||
export type RpcSessionChangeSession = Pick<AgentSession, "newSession" | "switchSession" | "branch">;
|
||||
|
||||
export type RpcSkillCommandSession = Pick<AgentSession, "promptCustomMessage" | "skills" | "skillsSettings">;
|
||||
export type RpcSkillCommandResult = { agentInvoked: false };
|
||||
export type RpcSkillCommandResult = { agentInvoked: true };
|
||||
|
||||
export async function tryRunRpcSkillCommand(
|
||||
session: RpcSkillCommandSession,
|
||||
@@ -102,7 +102,7 @@ export async function tryRunRpcSkillCommand(
|
||||
details: built.details,
|
||||
attribution: "user",
|
||||
});
|
||||
return { agentInvoked: false };
|
||||
return { agentInvoked: true };
|
||||
}
|
||||
|
||||
export function reportLocalOnlyPromptResult(input: {
|
||||
@@ -111,13 +111,10 @@ export function reportLocalOnlyPromptResult(input: {
|
||||
output: (obj: object) => void;
|
||||
onError: (error: Error) => void;
|
||||
hasExtensionAgentMessageTask?: () => boolean;
|
||||
waitForExtensionAgentMessageTasks?: () => Promise<void>;
|
||||
}): void {
|
||||
void input.prompt
|
||||
.then(async agentInvoked => {
|
||||
if (agentInvoked) return;
|
||||
await input.waitForExtensionAgentMessageTasks?.();
|
||||
if (!input.hasExtensionAgentMessageTask?.()) {
|
||||
.then(agentInvoked => {
|
||||
if (!agentInvoked && !input.hasExtensionAgentMessageTask?.()) {
|
||||
input.output({ type: "prompt_result", id: input.id, agentInvoked: false });
|
||||
}
|
||||
})
|
||||
@@ -128,7 +125,6 @@ export function reportLocalOnlyPromptResult(input: {
|
||||
|
||||
type RpcExtensionUserMessageScope = {
|
||||
hasAgentMessageTask: boolean;
|
||||
pendingAgentMessageTasks: Set<Promise<void>>;
|
||||
};
|
||||
|
||||
/**
|
||||
@@ -146,42 +142,11 @@ export class RpcExtensionUserMessageTracker {
|
||||
}
|
||||
}
|
||||
|
||||
trackAgentMessageTask(task: Promise<boolean>): void {
|
||||
for (const scope of this.#activePromptScopes) {
|
||||
this.#trackAgentMessageTaskForScope(scope, task);
|
||||
}
|
||||
}
|
||||
|
||||
#trackAgentMessageTaskForScope(scope: RpcExtensionUserMessageScope, task: Promise<boolean>): void {
|
||||
const scopedTask = task.then(
|
||||
agentInvoked => {
|
||||
if (agentInvoked) {
|
||||
scope.hasAgentMessageTask = true;
|
||||
}
|
||||
},
|
||||
() => {},
|
||||
);
|
||||
scope.pendingAgentMessageTasks.add(scopedTask);
|
||||
void scopedTask.finally(() => {
|
||||
scope.pendingAgentMessageTasks.delete(scopedTask);
|
||||
});
|
||||
}
|
||||
|
||||
async #waitForAgentMessageTasks(scope: RpcExtensionUserMessageScope): Promise<void> {
|
||||
while (scope.pendingAgentMessageTasks.size > 0) {
|
||||
await Promise.allSettled(Array.from(scope.pendingAgentMessageTasks));
|
||||
}
|
||||
}
|
||||
|
||||
watchPrompt<T>(startPrompt: () => Promise<T>): {
|
||||
prompt: Promise<T>;
|
||||
hasAgentMessageTask: () => boolean;
|
||||
waitForAgentMessageTasks: () => Promise<void>;
|
||||
} {
|
||||
const scope: RpcExtensionUserMessageScope = {
|
||||
hasAgentMessageTask: false,
|
||||
pendingAgentMessageTasks: new Set(),
|
||||
};
|
||||
const scope: RpcExtensionUserMessageScope = { hasAgentMessageTask: false };
|
||||
this.#activePromptScopes.add(scope);
|
||||
let prompt: Promise<T>;
|
||||
try {
|
||||
@@ -195,7 +160,6 @@ export class RpcExtensionUserMessageTracker {
|
||||
this.#activePromptScopes.delete(scope);
|
||||
}),
|
||||
hasAgentMessageTask: () => scope.hasAgentMessageTask,
|
||||
waitForAgentMessageTasks: () => this.#waitForAgentMessageTasks(scope),
|
||||
};
|
||||
}
|
||||
}
|
||||
@@ -214,7 +178,6 @@ export function watchAndReportLocalOnlyPromptResult(input: {
|
||||
output: input.output,
|
||||
onError: input.onError,
|
||||
hasExtensionAgentMessageTask: trackedPrompt.hasAgentMessageTask,
|
||||
waitForExtensionAgentMessageTasks: trackedPrompt.waitForAgentMessageTasks,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -656,9 +619,6 @@ export async function runRpcMode(
|
||||
markAgentInvokingMessage: () => {
|
||||
extensionUserMessageTracker.markAgentMessageTask();
|
||||
},
|
||||
trackAgentInvokingUserMessage: task => {
|
||||
extensionUserMessageTracker.trackAgentMessageTask(task);
|
||||
},
|
||||
uiContext: rpcUiContext,
|
||||
});
|
||||
|
||||
|
||||
@@ -25,8 +25,6 @@ export interface InitializeExtensionsOptions {
|
||||
uiContext?: ExtensionUIContext;
|
||||
/** Optional lifecycle hook for extension-originated messages that can start an agent turn. */
|
||||
markAgentInvokingMessage?: () => void;
|
||||
/** Optional lifecycle hook for extension-originated user messages with async delivery outcomes. */
|
||||
trackAgentInvokingUserMessage?: (task: Promise<boolean>) => void;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -39,14 +37,7 @@ export async function initializeExtensions(session: AgentSession, options: Initi
|
||||
const runner = session.extensionRunner;
|
||||
if (!runner) return;
|
||||
|
||||
const {
|
||||
reportSendError,
|
||||
reportRuntimeError,
|
||||
onShutdown,
|
||||
uiContext,
|
||||
markAgentInvokingMessage,
|
||||
trackAgentInvokingUserMessage,
|
||||
} = options;
|
||||
const { reportSendError, reportRuntimeError, onShutdown, uiContext, markAgentInvokingMessage } = options;
|
||||
const shutdown = onShutdown ?? (() => {});
|
||||
|
||||
runner.initialize(
|
||||
@@ -61,21 +52,8 @@ export async function initializeExtensions(session: AgentSession, options: Initi
|
||||
});
|
||||
},
|
||||
sendUserMessage: (content, sendOptions) => {
|
||||
const sendTask = session.sendUserMessage(content, sendOptions);
|
||||
if (trackAgentInvokingUserMessage) {
|
||||
trackAgentInvokingUserMessage(sendTask);
|
||||
} else if (!sendOptions?.deliverAs) {
|
||||
markAgentInvokingMessage?.();
|
||||
} else {
|
||||
void sendTask
|
||||
.then(agentInvoked => {
|
||||
if (agentInvoked) {
|
||||
markAgentInvokingMessage?.();
|
||||
}
|
||||
})
|
||||
.catch(() => {});
|
||||
}
|
||||
void sendTask.catch(e => {
|
||||
markAgentInvokingMessage?.();
|
||||
session.sendUserMessage(content, sendOptions).catch(e => {
|
||||
reportSendError("extension_send_user", e instanceof Error ? e : new Error(String(e)));
|
||||
});
|
||||
},
|
||||
|
||||
@@ -2018,14 +2018,11 @@ export async function runSubprocess(options: ExecutorOptions): Promise<SingleRes
|
||||
pendingExtensionMessages.push(sendPromise);
|
||||
},
|
||||
sendUserMessage: (content, options) => {
|
||||
const sendPromise = session.sendUserMessage(content, options).then(
|
||||
() => undefined,
|
||||
e => {
|
||||
logger.error("Extension sendUserMessage failed", {
|
||||
error: e instanceof Error ? e.message : String(e),
|
||||
});
|
||||
},
|
||||
);
|
||||
const sendPromise = session.sendUserMessage(content, options).catch(e => {
|
||||
logger.error("Extension sendUserMessage failed", {
|
||||
error: e instanceof Error ? e.message : String(e),
|
||||
});
|
||||
});
|
||||
pendingExtensionMessages.push(sendPromise);
|
||||
},
|
||||
appendEntry: (customType, data) => {
|
||||
|
||||
@@ -330,9 +330,7 @@ class FakeAgentSession {
|
||||
|
||||
async sendCustomMessage(_message: string, _options?: unknown): Promise<void> {}
|
||||
|
||||
async sendUserMessage(_content: string, _options?: unknown): Promise<boolean> {
|
||||
return false;
|
||||
}
|
||||
async sendUserMessage(_content: string, _options?: unknown): Promise<void> {}
|
||||
|
||||
async compact(_instructions?: string, _options?: unknown): Promise<void> {}
|
||||
|
||||
|
||||
@@ -110,9 +110,7 @@ class FakeAgentSession {
|
||||
}
|
||||
setPlanModeState(): void {}
|
||||
async sendCustomMessage(): Promise<void> {}
|
||||
async sendUserMessage(): Promise<boolean> {
|
||||
return false;
|
||||
}
|
||||
async sendUserMessage(): Promise<void> {}
|
||||
async compact(): Promise<void> {}
|
||||
async fork(): Promise<boolean> {
|
||||
return false;
|
||||
|
||||
@@ -133,9 +133,7 @@ class LazyFakeSession {
|
||||
}
|
||||
setPlanModeState(): void {}
|
||||
async sendCustomMessage(): Promise<void> {}
|
||||
async sendUserMessage(): Promise<boolean> {
|
||||
return false;
|
||||
}
|
||||
async sendUserMessage(): Promise<void> {}
|
||||
async compact(): Promise<void> {}
|
||||
async fork(): Promise<boolean> {
|
||||
return false;
|
||||
|
||||
@@ -1,64 +0,0 @@
|
||||
import { describe, expect, test } from "bun:test";
|
||||
import * as fs from "node:fs/promises";
|
||||
import * as os from "node:os";
|
||||
import * as path from "node:path";
|
||||
import { RpcClient } from "@oh-my-pi/pi-coding-agent/modes/rpc/rpc-client";
|
||||
|
||||
const FAKE_RPC_SERVER = `
|
||||
function writeFrame(frame: unknown): void {
|
||||
process.stdout.write(JSON.stringify(frame) + "\\n");
|
||||
}
|
||||
|
||||
writeFrame({ type: "ready" });
|
||||
const decoder = new TextDecoder();
|
||||
let buffer = "";
|
||||
for await (const chunk of Bun.stdin.stream()) {
|
||||
buffer += decoder.decode(chunk, { stream: true });
|
||||
let newline = buffer.indexOf("\\n");
|
||||
while (newline !== -1) {
|
||||
const line = buffer.slice(0, newline).trim();
|
||||
buffer = buffer.slice(newline + 1);
|
||||
if (line.length > 0) {
|
||||
const command = JSON.parse(line) as { type?: string; id?: string; message?: string };
|
||||
if (command.type === "prompt") {
|
||||
if (command.message === "immediate") {
|
||||
writeFrame({
|
||||
type: "response",
|
||||
command: "prompt",
|
||||
id: command.id,
|
||||
success: true,
|
||||
data: { agentInvoked: false },
|
||||
});
|
||||
} else {
|
||||
writeFrame({ type: "response", command: "prompt", id: command.id, success: true });
|
||||
setTimeout(() => {
|
||||
writeFrame({ type: "prompt_result", id: command.id, agentInvoked: false });
|
||||
}, 10);
|
||||
}
|
||||
}
|
||||
}
|
||||
newline = buffer.indexOf("\\n");
|
||||
}
|
||||
}
|
||||
`;
|
||||
|
||||
describe("RpcClient prompt completion", () => {
|
||||
test("waits for local-only prompt completions without agent_end", async () => {
|
||||
const dir = await fs.mkdtemp(path.join(os.tmpdir(), "omp-rpc-client-prompt-result-"));
|
||||
const serverPath = path.join(dir, "fake-rpc-server.ts");
|
||||
await Bun.write(serverPath, FAKE_RPC_SERVER);
|
||||
const client = new RpcClient({ cliPath: serverPath, cwd: dir });
|
||||
try {
|
||||
await client.start();
|
||||
|
||||
await client.prompt("immediate");
|
||||
await client.waitForIdle(1000);
|
||||
|
||||
const events = await client.promptAndWait("deferred", undefined, 1000);
|
||||
expect(events).toEqual([]);
|
||||
} finally {
|
||||
client.stop();
|
||||
await fs.rm(dir, { recursive: true, force: true });
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -140,112 +140,6 @@ describe("reportLocalOnlyPromptResult", () => {
|
||||
expect(sentOptions).toEqual({ triggerTurn: true });
|
||||
});
|
||||
|
||||
test("emits prompt_result when extension sendUserMessage only queues follow-up", async () => {
|
||||
let extensionActions: ExtensionActions | undefined;
|
||||
let sentOptions: { deliverAs?: "steer" | "followUp" } | undefined;
|
||||
const delivery = Promise.withResolvers<boolean>();
|
||||
const extensionUserMessages = new RpcExtensionUserMessageTracker();
|
||||
const session = {
|
||||
extensionRunner: {
|
||||
initialize: (actions: ExtensionActions) => {
|
||||
extensionActions = actions;
|
||||
},
|
||||
onError: () => {},
|
||||
emit: async () => {},
|
||||
},
|
||||
sendUserMessage: async (_content: unknown, options?: { deliverAs?: "steer" | "followUp" }) => {
|
||||
sentOptions = options;
|
||||
return await delivery.promise;
|
||||
},
|
||||
} as unknown as AgentSession;
|
||||
|
||||
await initializeExtensions(session, {
|
||||
reportSendError: (_action, error) => {
|
||||
throw error;
|
||||
},
|
||||
reportRuntimeError: error => {
|
||||
throw error.error;
|
||||
},
|
||||
trackAgentInvokingUserMessage: task => {
|
||||
extensionUserMessages.trackAgentMessageTask(task);
|
||||
},
|
||||
});
|
||||
|
||||
const output: object[] = [];
|
||||
const trackedPrompt = extensionUserMessages.watchPrompt(() => {
|
||||
if (!extensionActions) throw new Error("extensions not initialized");
|
||||
extensionActions.sendUserMessage("queued locally", { deliverAs: "followUp" });
|
||||
return Promise.resolve(false);
|
||||
});
|
||||
reportLocalOnlyPromptResult({
|
||||
id: "req_queued",
|
||||
prompt: trackedPrompt.prompt,
|
||||
output: frame => output.push(frame),
|
||||
onError: error => {
|
||||
throw error;
|
||||
},
|
||||
hasExtensionAgentMessageTask: trackedPrompt.hasAgentMessageTask,
|
||||
waitForExtensionAgentMessageTasks: trackedPrompt.waitForAgentMessageTasks,
|
||||
});
|
||||
|
||||
await waitForPromptHandlers(trackedPrompt.prompt);
|
||||
expect(output).toEqual([]);
|
||||
|
||||
delivery.resolve(false);
|
||||
await waitForPromptHandlers(delivery.promise);
|
||||
await waitForPromptHandlers(trackedPrompt.prompt);
|
||||
|
||||
expect(sentOptions).toEqual({ deliverAs: "followUp" });
|
||||
expect(output).toEqual([{ type: "prompt_result", id: "req_queued", agentInvoked: false }]);
|
||||
});
|
||||
|
||||
test("suppresses prompt_result when extension sendUserMessage starts agent work", async () => {
|
||||
let extensionActions: ExtensionActions | undefined;
|
||||
const extensionUserMessages = new RpcExtensionUserMessageTracker();
|
||||
const session = {
|
||||
extensionRunner: {
|
||||
initialize: (actions: ExtensionActions) => {
|
||||
extensionActions = actions;
|
||||
},
|
||||
onError: () => {},
|
||||
emit: async () => {},
|
||||
},
|
||||
sendUserMessage: async () => true,
|
||||
} as unknown as AgentSession;
|
||||
|
||||
await initializeExtensions(session, {
|
||||
reportSendError: (_action, error) => {
|
||||
throw error;
|
||||
},
|
||||
reportRuntimeError: error => {
|
||||
throw error.error;
|
||||
},
|
||||
trackAgentInvokingUserMessage: task => {
|
||||
extensionUserMessages.trackAgentMessageTask(task);
|
||||
},
|
||||
});
|
||||
|
||||
const output: object[] = [];
|
||||
const trackedPrompt = extensionUserMessages.watchPrompt(() => {
|
||||
if (!extensionActions) throw new Error("extensions not initialized");
|
||||
extensionActions.sendUserMessage("start work");
|
||||
return Promise.resolve(false);
|
||||
});
|
||||
reportLocalOnlyPromptResult({
|
||||
id: "req_agent",
|
||||
prompt: trackedPrompt.prompt,
|
||||
output: frame => output.push(frame),
|
||||
onError: error => {
|
||||
throw error;
|
||||
},
|
||||
hasExtensionAgentMessageTask: trackedPrompt.hasAgentMessageTask,
|
||||
waitForExtensionAgentMessageTasks: trackedPrompt.waitForAgentMessageTasks,
|
||||
});
|
||||
await waitForPromptHandlers(trackedPrompt.prompt);
|
||||
|
||||
expect(output).toEqual([]);
|
||||
});
|
||||
|
||||
test("does not emit when prompt invokes the agent", async () => {
|
||||
const output: object[] = [];
|
||||
const prompt = Promise.resolve(true);
|
||||
|
||||
@@ -30,7 +30,7 @@ describe("tryRunRpcSkillCommand", () => {
|
||||
"/skill:reviewer focus on risks",
|
||||
);
|
||||
|
||||
expect(handled).toEqual({ agentInvoked: false });
|
||||
expect(handled).toEqual({ agentInvoked: true });
|
||||
expect(message?.customType).toBe(SKILL_PROMPT_MESSAGE_TYPE);
|
||||
expect(message?.content).toContain("Review the supplied code carefully.");
|
||||
expect(message?.content).toContain("User: focus on risks");
|
||||
|
||||
@@ -146,7 +146,6 @@ describe("runSubprocess yield reminders", () => {
|
||||
messageInFlight = true;
|
||||
await Bun.sleep(20);
|
||||
messageInFlight = false;
|
||||
return true;
|
||||
};
|
||||
mutableSession.extensionRunner = {
|
||||
initialize: (actions: ExtensionActions) => {
|
||||
|
||||
@@ -98,7 +98,6 @@ function createFakeSession(config: FakeSessionConfig = {}): FakeSessionHandle {
|
||||
},
|
||||
sendUserMessage: async (content, options) => {
|
||||
steerCalls.push({ content: String(content), options });
|
||||
return false;
|
||||
},
|
||||
getLastAssistantMessage: () => (config.lastAssistantMessage ?? undefined) as never,
|
||||
abort: async () => {
|
||||
|
||||
Reference in New Issue
Block a user