9f386dd6e5
Adds an opt-in onSseEvent callback across HTTP-streaming providers (Anthropic, OpenAI Responses/Completions, Azure OpenAI Responses, OpenAI Codex SSE, Google Gemini CLI, GitLab Duo, Kimi, Synthetic) so callers can inspect raw SSE frames without altering parsed output. Provider fetch wrapping only tees response bodies when an observer is wired; standalone packages/ai consumers without onSseEvent are not penalized. Adds streamIdleTimeoutMs (env: PI_STREAM_IDLE_TIMEOUT_MS, with PI_OPENAI_STREAM_IDLE_TIMEOUT_MS as a backward-compatible alias). Anthropic now enforces a steady-state idle watchdog (default 120s) in addition to the first-event watchdog. OpenAI Responses, Azure Responses, and Codex (SSE + WebSocket) gain a semantic-progress predicate so response.in_progress-style keepalives no longer keep stalled tool calls alive forever. Adds a coding-agent debug-panel raw SSE viewer backed by a per-session bounded buffer (1000 records / 512KB) that AgentSession populates unconditionally so users can post-hoc inspect a stuck stream from the TUI.
72 lines
2.3 KiB
TypeScript
72 lines
2.3 KiB
TypeScript
import { describe, expect, it } from "bun:test";
|
|
import type { Model } from "@oh-my-pi/pi-ai";
|
|
import { RawSseDebugBuffer, rawSseRecordLines, resolveRawSseDebugBuffer } from "../../src/debug/raw-sse-buffer";
|
|
|
|
const model: Model<"anthropic-messages"> = {
|
|
id: "claude-test",
|
|
name: "Claude Test",
|
|
api: "anthropic-messages",
|
|
provider: "anthropic",
|
|
baseUrl: "https://api.anthropic.com",
|
|
reasoning: true,
|
|
input: ["text"],
|
|
cost: { input: 1, output: 1, cacheRead: 0, cacheWrite: 0 },
|
|
contextWindow: 200_000,
|
|
maxTokens: 8_192,
|
|
};
|
|
|
|
describe("RawSseDebugBuffer", () => {
|
|
it("records response metadata and raw SSE frame lines for diagnostics", () => {
|
|
const buffer = new RawSseDebugBuffer();
|
|
|
|
buffer.recordResponse(
|
|
{ status: 200, requestId: "req_123", headers: {}, metadata: { lastTransport: "sse" } },
|
|
model,
|
|
);
|
|
buffer.recordEvent(
|
|
{
|
|
event: "content_block_delta",
|
|
data: '{"type":"content_block_delta"}',
|
|
raw: ["event: content_block_delta", 'data: {"type":"content_block_delta"}'],
|
|
},
|
|
model,
|
|
);
|
|
|
|
const snapshot = buffer.snapshot();
|
|
expect(snapshot.totalEvents).toBe(1);
|
|
expect(snapshot.records).toHaveLength(2);
|
|
const [responseLine] = rawSseRecordLines(snapshot.records[0]);
|
|
expect(responseLine).toContain("provider=anthropic model=claude-test");
|
|
expect(rawSseRecordLines(snapshot.records[1])).toEqual([
|
|
"event: content_block_delta",
|
|
'data: {"type":"content_block_delta"}',
|
|
]);
|
|
expect(buffer.toRawText()).toContain("event: content_block_delta");
|
|
});
|
|
|
|
it("notifies subscribers when new frames arrive", () => {
|
|
const buffer = new RawSseDebugBuffer();
|
|
let updates = 0;
|
|
const unsubscribe = buffer.subscribe(() => {
|
|
updates += 1;
|
|
});
|
|
|
|
buffer.recordEvent({ event: null, data: "{}", raw: ["data: {}"] }, model);
|
|
unsubscribe();
|
|
buffer.recordEvent({ event: null, data: "{}", raw: ["data: {}"] }, model);
|
|
|
|
expect(updates).toBe(1);
|
|
expect(buffer.snapshot().totalEvents).toBe(2);
|
|
});
|
|
|
|
it("creates a fallback buffer for session objects without a preinstalled buffer", () => {
|
|
const owner = {};
|
|
const buffer = resolveRawSseDebugBuffer(owner);
|
|
|
|
buffer.recordEvent({ event: "message", data: "{}", raw: ["event: message", "data: {}"] }, model);
|
|
|
|
expect(resolveRawSseDebugBuffer(owner)).toBe(buffer);
|
|
expect(buffer.snapshot().totalEvents).toBe(1);
|
|
});
|
|
});
|