diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index b3005e664..723248a82 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -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 diff --git a/packages/coding-agent/src/modes/rpc/rpc-client.ts b/packages/coding-agent/src/modes/rpc/rpc-client.ts index b94bd55d5..82f2e01d2 100644 --- a/packages/coding-agent/src/modes/rpc/rpc-client.ts +++ b/packages/coding-agent/src/modes/rpc/rpc-client.ts @@ -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) { diff --git a/packages/coding-agent/test/rpc-client.restart.test.ts b/packages/coding-agent/test/rpc-client.restart.test.ts index 587d077cf..d9cd8dab9 100644 --- a/packages/coding-agent/test/rpc-client.restart.test.ts +++ b/packages/coding-agent/test/rpc-client.restart.test.ts @@ -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 | undefined; + let resolveExit: ((exitCode: number) => void) | undefined; + let killCalls = 0; + const exited = new Promise(resolve => { + resolveExit = resolve; + }); + const stdout = new ReadableStream({ + 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, + ); + + 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, diff --git a/python/omp-rpc/src/omp_rpc/client.py b/python/omp-rpc/src/omp_rpc/client.py index c1b9cc3f8..1418f59f1 100644 --- a/python/omp-rpc/src/omp_rpc/client.py +++ b/python/omp-rpc/src/omp_rpc/client.py @@ -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, diff --git a/python/omp-rpc/src/omp_rpc/protocol.py b/python/omp-rpc/src/omp_rpc/protocol.py index 5cc4ad3c6..0034f3dd0 100644 --- a/python/omp-rpc/src/omp_rpc/protocol.py +++ b/python/omp-rpc/src/omp_rpc/protocol.py @@ -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() diff --git a/python/omp-rpc/tests/test_client.py b/python/omp-rpc/tests/test_client.py index f9da3d766..68114923d 100644 --- a/python/omp-rpc/tests/test_client.py +++ b/python/omp-rpc/tests/test_client.py @@ -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']}") diff --git a/python/omp-rpc/tests/test_protocol.py b/python/omp-rpc/tests/test_protocol.py index 5d7bb158f..58b2f91cd 100644 --- a/python/omp-rpc/tests/test_protocol.py +++ b/python/omp-rpc/tests/test_protocol.py @@ -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(