fix(coding-agent): preserve small RPC terminal frames
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 `agent_end` 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 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)).
|
||||
|
||||
## [17.0.4] - 2026-07-18
|
||||
|
||||
|
||||
@@ -121,13 +121,16 @@ function overflowFrame(frame: object): object {
|
||||
|
||||
/** Serialize a complete JSONL frame while enforcing the transport byte ceiling. */
|
||||
export function encodeRpcFrame(frame: object, streamedMessageCount = 0, streamedMessages?: readonly unknown[]): string {
|
||||
const compacted = compactTerminalFrame(frame, streamedMessageCount, streamedMessages);
|
||||
let json = JSON.stringify(compacted);
|
||||
let json = JSON.stringify(frame);
|
||||
if (serializedFrameBytes(json) <= MAX_RPC_FRAME_BYTES) return `${json}\n`;
|
||||
if (isRecord(compacted) && compacted.type === "response") {
|
||||
return `${JSON.stringify(overflowFrame(compacted))}\n`;
|
||||
if (isRecord(frame) && frame.type === "response") {
|
||||
return `${JSON.stringify(overflowFrame(frame))}\n`;
|
||||
}
|
||||
|
||||
const compacted = compactTerminalFrame(frame, streamedMessageCount, streamedMessages);
|
||||
json = JSON.stringify(compacted);
|
||||
if (serializedFrameBytes(json) <= MAX_RPC_FRAME_BYTES) return `${json}\n`;
|
||||
|
||||
for (const pass of SHRINK_PASSES) {
|
||||
json = JSON.stringify(shrinkValue(compacted, pass));
|
||||
if (serializedFrameBytes(json) <= MAX_RPC_FRAME_BYTES) return `${json}\n`;
|
||||
|
||||
@@ -5,6 +5,14 @@ function decode(frame: string): Record<string, unknown> {
|
||||
return JSON.parse(frame) as Record<string, unknown>;
|
||||
}
|
||||
|
||||
function oversizedMessageHistory(prefix: string) {
|
||||
const payload = "x".repeat(1024);
|
||||
return Array.from({ length: 1024 }, (_, index) => ({
|
||||
role: "assistant",
|
||||
content: [{ type: "text", text: `${prefix}-${index}-${payload}` }],
|
||||
}));
|
||||
}
|
||||
|
||||
describe("RPC frame encoding", () => {
|
||||
it("preserves frames that already fit", () => {
|
||||
const frame = { id: "request-1", type: "response", command: "get_state", success: true, data: { ok: true } };
|
||||
@@ -39,93 +47,77 @@ describe("RPC frame encoding", () => {
|
||||
expect(decoded).toEqual({
|
||||
type: "agent_end",
|
||||
messages: [aborted],
|
||||
messageCount: 1,
|
||||
});
|
||||
});
|
||||
|
||||
it("compacts stateful agent_end messages that match earlier message events", () => {
|
||||
it("preserves terminal histories that fit for clients reading agent_end messages", () => {
|
||||
const streamed = { role: "assistant", content: [{ type: "text", text: "done" }] };
|
||||
const encoder = new RpcFrameEncoder();
|
||||
encoder.encode({ type: "agent_start" });
|
||||
encoder.encode({ type: "message_end", message: streamed });
|
||||
const decoded = decode(
|
||||
encoder.encode({
|
||||
type: "agent_end",
|
||||
messages: [{ role: "assistant", content: [{ type: "text", text: "done" }] }],
|
||||
}),
|
||||
);
|
||||
const frame = { type: "agent_end", messages: [streamed] };
|
||||
|
||||
expect(decoded).toEqual({
|
||||
type: "agent_end",
|
||||
messages: [],
|
||||
messageCount: 1,
|
||||
});
|
||||
expect(encoder.encode(frame)).toBe(`${JSON.stringify(frame)}\n`);
|
||||
});
|
||||
|
||||
it("matches terminal messages in the JSON shape sent by message_end", () => {
|
||||
it("matches oversized terminal messages in the JSON shape sent by message_end", () => {
|
||||
const messages = oversizedMessageHistory("wire-shape");
|
||||
const encoder = new RpcFrameEncoder();
|
||||
encoder.encode({ type: "agent_start" });
|
||||
encoder.encode({
|
||||
type: "message_end",
|
||||
message: {
|
||||
role: "assistant",
|
||||
content: [{ type: "text", text: "done" }],
|
||||
disabledFeatures: undefined,
|
||||
toolCallAbortMessages: undefined,
|
||||
},
|
||||
});
|
||||
const decoded = decode(
|
||||
for (const message of messages) {
|
||||
encoder.encode({
|
||||
type: "agent_end",
|
||||
messages: [{ role: "assistant", content: [{ type: "text", text: "done" }] }],
|
||||
}),
|
||||
);
|
||||
type: "message_end",
|
||||
message: {
|
||||
...message,
|
||||
disabledFeatures: undefined,
|
||||
toolCallAbortMessages: undefined,
|
||||
},
|
||||
});
|
||||
}
|
||||
const encoded = encoder.encode({ type: "agent_end", messages });
|
||||
|
||||
expect(decoded).toEqual({
|
||||
expect(Buffer.byteLength(encoded, "utf8")).toBeLessThanOrEqual(MAX_RPC_FRAME_BYTES);
|
||||
expect(decode(encoded)).toEqual({
|
||||
type: "agent_end",
|
||||
messages: [],
|
||||
messageCount: 1,
|
||||
messageCount: messages.length,
|
||||
});
|
||||
});
|
||||
|
||||
it("does not let later mutation rewrite the message_end snapshot", () => {
|
||||
const streamed = { role: "assistant", content: [{ type: "text", text: "before" }] };
|
||||
const messages = oversizedMessageHistory("before");
|
||||
const encoder = new RpcFrameEncoder();
|
||||
encoder.encode({ type: "agent_start" });
|
||||
encoder.encode({ type: "message_end", message: streamed });
|
||||
streamed.content[0].text = "after";
|
||||
const decoded = decode(encoder.encode({ type: "agent_end", messages: [streamed] }));
|
||||
for (const message of messages) encoder.encode({ type: "message_end", message });
|
||||
messages[0].content[0].text = "after";
|
||||
const decoded = decode(encoder.encode({ type: "agent_end", messages }));
|
||||
|
||||
expect(decoded).toEqual({
|
||||
type: "agent_end",
|
||||
messages: [streamed],
|
||||
messageCount: 1,
|
||||
});
|
||||
expect(decoded.messageCount).toBe(messages.length);
|
||||
expect(Array.isArray(decoded.messages)).toBe(true);
|
||||
expect((decoded.messages as unknown[]).length).toBeGreaterThan(0);
|
||||
});
|
||||
|
||||
it("keeps the active run snapshot when a continuing agent_end arrives late", () => {
|
||||
const active = { role: "assistant", content: [{ type: "text", text: "active" }] };
|
||||
const active = oversizedMessageHistory("active");
|
||||
const stale = { role: "assistant", content: [{ type: "text", text: "stale" }] };
|
||||
const encoder = new RpcFrameEncoder();
|
||||
encoder.encode({ type: "agent_start" });
|
||||
encoder.encode({ type: "message_end", message: active });
|
||||
for (const message of active) encoder.encode({ type: "message_end", message });
|
||||
|
||||
expect(decode(encoder.encode({ type: "agent_end", messages: [stale], willContinue: true }))).toEqual({
|
||||
type: "agent_end",
|
||||
messages: [stale],
|
||||
messageCount: 1,
|
||||
willContinue: true,
|
||||
});
|
||||
expect(decode(encoder.encode({ type: "agent_end", messages: [active] }))).toEqual({
|
||||
expect(decode(encoder.encode({ type: "agent_end", messages: active }))).toEqual({
|
||||
type: "agent_end",
|
||||
messages: [],
|
||||
messageCount: 1,
|
||||
});
|
||||
expect(decode(encoder.encode({ type: "agent_end", messages: [active] }))).toEqual({
|
||||
type: "agent_end",
|
||||
messages: [active],
|
||||
messageCount: 1,
|
||||
messageCount: active.length,
|
||||
});
|
||||
const replayed = decode(encoder.encode({ type: "agent_end", messages: active }));
|
||||
expect(replayed.messageCount).toBe(active.length);
|
||||
expect(Array.isArray(replayed.messages)).toBe(true);
|
||||
expect((replayed.messages as unknown[]).length).toBeGreaterThan(0);
|
||||
});
|
||||
|
||||
it("bounds a single multi-byte message without losing its event discriminator", () => {
|
||||
|
||||
Reference in New Issue
Block a user