fix(coding-agent): reconstruct compacted RPC prompt histories
This commit is contained in:
@@ -127,7 +127,7 @@
|
||||
- Fixed authoritative providers (e.g. `openai-codex`) keeping unsupported bundled models selectable when a fresh model cache and an expired OAuth token coincided: built-in discovery now forces the OAuth refresh so the provider's model manager is constructed and prunes stale bundled entries (e.g. `gpt-5.4-nano`) instead of waiting out the cache TTL. ([#5364](https://github.com/can1357/oh-my-pi/issues/5364))
|
||||
### Fixed
|
||||
|
||||
- Bounded RPC JSONL frames to 1 MiB, compacted oversized `agent_end` frames to retain only messages not already streamed, and guaranteed worker reaping plus pending-request rejection after output-reader failures or explicit stops ([#5405](https://github.com/can1357/oh-my-pi/issues/5405)).
|
||||
- Bounded RPC JSONL frames to 1 MiB, compacted oversized `agent_end` frames without losing complete Python prompt results, and guaranteed worker reaping plus pending-request rejection after output-reader failures or explicit stops ([#5405](https://github.com/can1357/oh-my-pi/issues/5405)).
|
||||
|
||||
## [17.0.4] - 2026-07-18
|
||||
|
||||
|
||||
@@ -300,27 +300,30 @@ export class RpcClient {
|
||||
}
|
||||
this.#handleLine(line);
|
||||
}
|
||||
// Stream ended without the ready signal — the child exited or is
|
||||
// exiting. Defer to the exit handler below: ptree resolves
|
||||
// `exited` only after stderr is fully drained (nonzero exits), so
|
||||
// rejecting here would snapshot a partial stderr tail and lose
|
||||
// the actual startup error.
|
||||
if (readySettled) {
|
||||
let error: Error;
|
||||
try {
|
||||
const exitCode = await child.exited;
|
||||
error = new Error(`Agent process exited with code ${exitCode}. Stderr: ${child.peekStderr()}`);
|
||||
} catch (cause) {
|
||||
error = new Error(`Agent output stream ended. Stderr: ${child.peekStderr()}`, { cause });
|
||||
}
|
||||
await reapAfterOutputFailure(error);
|
||||
return;
|
||||
}
|
||||
await child.exited.catch(() => {});
|
||||
// A closed stdout is terminal even if the child remains alive. Startup
|
||||
// failures are reaped by the readyPromise catch below; established
|
||||
// workers are reaped here so pending requests cannot hang indefinitely.
|
||||
if (!readySettled) {
|
||||
readySettled = true;
|
||||
readyReject(new Error(`Agent process exited before ready. Stderr: ${child.peekStderr()}`));
|
||||
readyReject(new Error(`Agent output stream ended before ready. Stderr: ${child.peekStderr()}`));
|
||||
return;
|
||||
}
|
||||
const exitResult = await Promise.race([
|
||||
child.exited.then(
|
||||
exitCode => ({ exitCode }),
|
||||
cause => ({ cause }),
|
||||
),
|
||||
Bun.sleep(100).then(() => null),
|
||||
]);
|
||||
const error =
|
||||
exitResult === null
|
||||
? new Error(`Agent output stream ended unexpectedly. Stderr: ${child.peekStderr()}`)
|
||||
: "exitCode" in exitResult
|
||||
? new Error(`Agent process exited with code ${exitResult.exitCode}. Stderr: ${child.peekStderr()}`)
|
||||
: new Error(`Agent output stream ended. Stderr: ${child.peekStderr()}`, {
|
||||
cause: exitResult.cause,
|
||||
});
|
||||
await reapAfterOutputFailure(error);
|
||||
})().catch(async (cause: unknown) => {
|
||||
const error = cause instanceof Error ? cause : new Error(String(cause));
|
||||
if (!readySettled) {
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
import { describe, expect, test } from "bun:test";
|
||||
import { describe, expect, spyOn, test } from "bun:test";
|
||||
import * as path from "node:path";
|
||||
import { RpcClient } from "@oh-my-pi/pi-coding-agent/modes/rpc/rpc-client";
|
||||
import { TempDir } from "@oh-my-pi/pi-utils";
|
||||
import { ptree, TempDir } from "@oh-my-pi/pi-utils";
|
||||
|
||||
const MOCK_AGENT = path.join(import.meta.dir, "fixtures", "mock-rpc-agent.ts");
|
||||
|
||||
@@ -114,6 +114,56 @@ describe("RpcClient lifecycle (issue #4079 B)", () => {
|
||||
}
|
||||
}, 10_000);
|
||||
|
||||
test("rejects pending requests and reaps a worker that closes stdout without exiting", async () => {
|
||||
let stdoutController: ReadableStreamDefaultController<Uint8Array> | undefined;
|
||||
let resolveExit: ((exitCode: number) => void) | undefined;
|
||||
let killCalls = 0;
|
||||
const exited = new Promise<number>(resolve => {
|
||||
resolveExit = resolve;
|
||||
});
|
||||
const stdout = new ReadableStream<Uint8Array>({
|
||||
start(controller) {
|
||||
stdoutController = controller;
|
||||
controller.enqueue(new TextEncoder().encode(`${JSON.stringify({ type: "ready" })}\n`));
|
||||
},
|
||||
});
|
||||
const fakeChild = {
|
||||
stdout,
|
||||
stdin: {
|
||||
write() {
|
||||
stdoutController?.close();
|
||||
stdoutController = undefined;
|
||||
return 0;
|
||||
},
|
||||
flush() {
|
||||
return 0;
|
||||
},
|
||||
},
|
||||
exited,
|
||||
peekStderr() {
|
||||
return "";
|
||||
},
|
||||
kill() {
|
||||
killCalls += 1;
|
||||
resolveExit?.(0);
|
||||
},
|
||||
};
|
||||
const spawn = spyOn(ptree, "spawn").mockImplementation(
|
||||
() => fakeChild as unknown as ReturnType<typeof ptree.spawn>,
|
||||
);
|
||||
|
||||
try {
|
||||
using client = new RpcClient({ cliPath: MOCK_AGENT });
|
||||
await client.start();
|
||||
|
||||
await expect(client.getState()).rejects.toThrow("Agent output stream ended unexpectedly");
|
||||
await expect(client.getState()).rejects.toThrow("Client not started");
|
||||
expect(killCalls).toBe(1);
|
||||
} finally {
|
||||
spawn.mockRestore();
|
||||
}
|
||||
}, 5_000);
|
||||
|
||||
test("reports exit code and stderr when a ready worker exits", async () => {
|
||||
using client = new RpcClient({
|
||||
cliPath: MOCK_AGENT,
|
||||
|
||||
@@ -1065,9 +1065,12 @@ class RpcClient:
|
||||
|
||||
def _build_prompt_turn(self, events: tuple[RpcAgentEvent, ...]) -> PromptTurn:
|
||||
final_messages: tuple[AgentMessage, ...] = ()
|
||||
for event in reversed(events):
|
||||
for event_index in range(len(events) - 1, -1, -1):
|
||||
event = events[event_index]
|
||||
if isinstance(event, AgentEndEvent):
|
||||
final_messages = event.messages
|
||||
final_messages = self._complete_agent_end_messages(
|
||||
events[:event_index], event
|
||||
)
|
||||
break
|
||||
|
||||
assistant_message: AssistantMessage | None = None
|
||||
@@ -1093,6 +1096,36 @@ class RpcClient:
|
||||
else None,
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def _complete_agent_end_messages(
|
||||
events: tuple[RpcAgentEvent, ...], terminal: AgentEndEvent
|
||||
) -> tuple[AgentMessage, ...]:
|
||||
if (
|
||||
terminal.message_count is None
|
||||
or terminal.message_count <= len(terminal.messages)
|
||||
):
|
||||
return terminal.messages
|
||||
|
||||
run_start = 0
|
||||
for event_index in range(len(events) - 1, -1, -1):
|
||||
if isinstance(events[event_index], AgentStartEvent):
|
||||
run_start = event_index + 1
|
||||
break
|
||||
|
||||
streamed_messages = tuple(
|
||||
event.message
|
||||
for event in events[run_start:]
|
||||
if isinstance(event, MessageEndEvent)
|
||||
)
|
||||
streamed_prefix_count = terminal.message_count - len(terminal.messages)
|
||||
if streamed_prefix_count > len(streamed_messages):
|
||||
raise RpcError(
|
||||
"Compacted agent_end references "
|
||||
f"{streamed_prefix_count} streamed messages, but only "
|
||||
f"{len(streamed_messages)} were retained"
|
||||
)
|
||||
return streamed_messages[:streamed_prefix_count] + terminal.messages
|
||||
|
||||
def _wait_for_agent_end(
|
||||
self,
|
||||
start_index: int,
|
||||
|
||||
@@ -2,7 +2,7 @@ from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import mimetypes
|
||||
from dataclasses import dataclass
|
||||
from dataclasses import dataclass, field
|
||||
from pathlib import Path
|
||||
from typing import Any, Final, Literal, NotRequired, TypedDict, TypeAlias, cast
|
||||
|
||||
@@ -905,6 +905,7 @@ class AgentStartEvent:
|
||||
class AgentEndEvent:
|
||||
messages: tuple[AgentMessage, ...]
|
||||
type: Literal["agent_end"] = "agent_end"
|
||||
message_count: int | None = field(default=None, kw_only=True)
|
||||
|
||||
|
||||
@dataclass(slots=True, frozen=True)
|
||||
@@ -1500,7 +1501,8 @@ def parse_notification(payload: JsonObject) -> RpcNotification:
|
||||
return AgentEndEvent(
|
||||
messages=parse_agent_messages(
|
||||
cast(JsonValue | None, payload.get("messages"))
|
||||
)
|
||||
),
|
||||
message_count=_optional_int(payload, "messageCount"),
|
||||
)
|
||||
if event_type == "turn_start":
|
||||
return TurnStartEvent()
|
||||
|
||||
@@ -86,7 +86,12 @@ FAKE_SERVER = textwrap.dedent(
|
||||
"dumpTools": [{"name": "read", "description": "Read files", "parameters": {"type": "object"}}] + registered_host_tools,
|
||||
}
|
||||
|
||||
def emit_prompt_turn(text: str, delay: float = 0.0, include_extra_events: bool = False):
|
||||
def emit_prompt_turn(
|
||||
text: str,
|
||||
delay: float = 0.0,
|
||||
include_extra_events: bool = False,
|
||||
compact_terminal: bool = False,
|
||||
):
|
||||
global last_assistant_text, messages
|
||||
print(json.dumps({"type": "agent_start"}), flush=True)
|
||||
print(json.dumps({"type": "turn_start"}), flush=True)
|
||||
@@ -198,9 +203,24 @@ FAKE_SERVER = textwrap.dedent(
|
||||
assistant = assistant_message(text)
|
||||
print(json.dumps({"type": "message_end", "message": assistant}), flush=True)
|
||||
print(json.dumps({"type": "turn_end", "message": assistant, "toolResults": []}), flush=True)
|
||||
print(json.dumps({"type": "agent_end", "messages": [assistant]}), flush=True)
|
||||
last_assistant_text = text
|
||||
messages = [assistant]
|
||||
if compact_terminal:
|
||||
terminal = assistant_message("terminal")
|
||||
print(
|
||||
json.dumps(
|
||||
{
|
||||
"type": "agent_end",
|
||||
"messages": [terminal],
|
||||
"messageCount": 2,
|
||||
}
|
||||
),
|
||||
flush=True,
|
||||
)
|
||||
last_assistant_text = "terminal"
|
||||
messages = [assistant, terminal]
|
||||
else:
|
||||
print(json.dumps({"type": "agent_end", "messages": [assistant]}), flush=True)
|
||||
last_assistant_text = text
|
||||
messages = [assistant]
|
||||
|
||||
def respond(request_id, command, data=None, success=True, error=None):
|
||||
payload = {"id": request_id, "type": "response", "command": command, "success": success}
|
||||
@@ -384,7 +404,12 @@ FAKE_SERVER = textwrap.dedent(
|
||||
if message == "notifications":
|
||||
print(json.dumps({"type": "extension_error", "extensionPath": "/tmp/ext.py", "event": "run", "error": "boom"}), flush=True)
|
||||
print(json.dumps({"type": "unknown_future_event", "value": 1}), flush=True)
|
||||
emit_prompt_turn("pong", delay=0.3 if message == "slow" else 0.0, include_extra_events=message == "all events")
|
||||
emit_prompt_turn(
|
||||
"pong",
|
||||
delay=0.3 if message == "slow" else 0.0,
|
||||
include_extra_events=message == "all events",
|
||||
compact_terminal=message == "compacted turn",
|
||||
)
|
||||
elif command_type == "host_tool_update":
|
||||
print(
|
||||
json.dumps(
|
||||
@@ -645,6 +670,16 @@ class RpcClientTests(unittest.TestCase):
|
||||
self.assertEqual(turn.require_assistant_text(), "pong")
|
||||
self.assertGreaterEqual(len(turn.events), 3)
|
||||
|
||||
def test_prompt_and_wait_reconstructs_compacted_terminal_messages(self) -> None:
|
||||
with self.make_client() as client:
|
||||
turn = client.prompt_and_wait("compacted turn", timeout=2.0)
|
||||
|
||||
self.assertEqual(
|
||||
[message["content"][0]["text"] for message in turn.messages],
|
||||
["pong", "terminal"],
|
||||
)
|
||||
self.assertEqual(turn.require_assistant_text(), "terminal")
|
||||
|
||||
def test_custom_tools_are_registered_and_executed_via_rpc(self) -> None:
|
||||
def echo_host(args: dict[str, str], context) -> str:
|
||||
context.send_update(f"working:{args['message']}")
|
||||
|
||||
@@ -125,11 +125,17 @@ class ProtocolParsingTests(unittest.TestCase):
|
||||
"timestamp": 1,
|
||||
}
|
||||
],
|
||||
"messageCount": 1,
|
||||
}
|
||||
)
|
||||
|
||||
self.assertIsInstance(notification, AgentEndEvent)
|
||||
self.assertEqual(assistant_text(notification.messages[0]), "hello")
|
||||
self.assertEqual(notification.message_count, 1)
|
||||
|
||||
legacy = AgentEndEvent(notification.messages, "agent_end")
|
||||
self.assertEqual(legacy.type, "agent_end")
|
||||
self.assertIsNone(legacy.message_count)
|
||||
|
||||
def test_parse_extension_ui_request(self) -> None:
|
||||
notification = parse_notification(
|
||||
|
||||
Reference in New Issue
Block a user