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.
182 lines
5.6 KiB
TypeScript
182 lines
5.6 KiB
TypeScript
import { afterEach, describe, expect, it, vi } from "bun:test";
|
|
import { Agent, type AgentMessage } from "@oh-my-pi/pi-agent-core";
|
|
import type { Message, SimpleStreamOptions } from "@oh-my-pi/pi-ai";
|
|
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
|
|
import { AgentSession, type AgentSessionEvent } from "@oh-my-pi/pi-coding-agent/session/agent-session";
|
|
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
|
|
|
|
function createAgent(): Agent {
|
|
return new Agent({
|
|
initialState: {
|
|
systemPrompt: ["system prompt"],
|
|
messages: [],
|
|
tools: [],
|
|
},
|
|
});
|
|
}
|
|
|
|
describe("AgentSession message pipeline", () => {
|
|
const sessions: AgentSession[] = [];
|
|
|
|
afterEach(async () => {
|
|
vi.restoreAllMocks();
|
|
for (const session of sessions.splice(0)) {
|
|
await session.dispose();
|
|
}
|
|
});
|
|
|
|
it("applies transformContext before convertToLlm", async () => {
|
|
const inputMessages: AgentMessage[] = [{ role: "user", content: "hello", timestamp: Date.now() }];
|
|
const transformedMessages: AgentMessage[] = [
|
|
...inputMessages,
|
|
{ role: "user", content: "injected context", timestamp: Date.now() },
|
|
];
|
|
const convertedMessages: Message[] = [
|
|
{
|
|
role: "user",
|
|
content: [{ type: "text", text: "converted" }],
|
|
attribution: "user",
|
|
timestamp: Date.now(),
|
|
},
|
|
];
|
|
const transformContext = vi.fn(async (messages: AgentMessage[], signal?: AbortSignal) => {
|
|
expect(signal).toBe(abortController.signal);
|
|
return [...messages, ...transformedMessages.slice(messages.length)];
|
|
});
|
|
const convertToLlm = vi.fn(async (_messages: AgentMessage[]) => {
|
|
return convertedMessages;
|
|
});
|
|
const abortController = new AbortController();
|
|
const session = new AgentSession({
|
|
agent: createAgent(),
|
|
sessionManager: SessionManager.inMemory(),
|
|
settings: Settings.isolated({ "compaction.enabled": false }),
|
|
modelRegistry: {} as never,
|
|
transformContext,
|
|
convertToLlm,
|
|
});
|
|
sessions.push(session);
|
|
|
|
const result = await session.convertMessagesToLlm(inputMessages, abortController.signal);
|
|
|
|
expect(transformContext).toHaveBeenCalledWith(inputMessages, abortController.signal);
|
|
expect(convertToLlm).toHaveBeenCalledWith(transformedMessages);
|
|
expect(result).toEqual(convertedMessages);
|
|
});
|
|
|
|
it("composes session payload hooks into direct side-request options", async () => {
|
|
const sessionOnPayload = vi.fn(async (payload: unknown) => ({
|
|
...(payload as Record<string, unknown>),
|
|
session: true,
|
|
}));
|
|
const requestOnPayload = vi.fn(async () => undefined);
|
|
const session = new AgentSession({
|
|
agent: createAgent(),
|
|
sessionManager: SessionManager.inMemory(),
|
|
settings: Settings.isolated({ "compaction.enabled": false }),
|
|
modelRegistry: {} as never,
|
|
onPayload: sessionOnPayload,
|
|
});
|
|
sessions.push(session);
|
|
const options: SimpleStreamOptions = {
|
|
apiKey: "key",
|
|
onPayload: requestOnPayload,
|
|
};
|
|
|
|
const prepared = session.prepareSimpleStreamOptions(options);
|
|
const result = await prepared.onPayload?.({ original: true });
|
|
|
|
expect(sessionOnPayload).toHaveBeenCalledWith({ original: true }, undefined);
|
|
expect(requestOnPayload).toHaveBeenCalledWith({ original: true, session: true }, undefined);
|
|
expect(result).toEqual({ original: true, session: true });
|
|
});
|
|
|
|
it("records raw SSE diagnostics into the session buffer before request hooks", async () => {
|
|
const requestOnSseEvent = vi.fn();
|
|
const session = new AgentSession({
|
|
agent: createAgent(),
|
|
sessionManager: SessionManager.inMemory(),
|
|
settings: Settings.isolated({ "compaction.enabled": false }),
|
|
modelRegistry: {} as never,
|
|
onSseEvent: requestOnSseEvent,
|
|
});
|
|
sessions.push(session);
|
|
|
|
const prepared = session.prepareSimpleStreamOptions({});
|
|
prepared.onSseEvent?.({ event: "message", data: "{}", raw: ["event: message", "data: {}"] });
|
|
|
|
expect(session.rawSseDebugBuffer.snapshot().totalEvents).toBe(1);
|
|
expect(requestOnSseEvent).toHaveBeenCalledWith(
|
|
{ event: "message", data: "{}", raw: ["event: message", "data: {}"] },
|
|
undefined,
|
|
);
|
|
});
|
|
|
|
it("emits message_update to session listeners before slow extension handlers finish", async () => {
|
|
const { promise, resolve } = Promise.withResolvers<void>();
|
|
const extensionEmit = vi.fn(async (event: { type: string }) => {
|
|
if (event.type === "message_update") {
|
|
await promise;
|
|
}
|
|
});
|
|
const session = new AgentSession({
|
|
agent: createAgent(),
|
|
sessionManager: SessionManager.inMemory(),
|
|
settings: Settings.isolated({ "compaction.enabled": false }),
|
|
modelRegistry: {} as never,
|
|
extensionRunner: {
|
|
emit: extensionEmit,
|
|
} as never,
|
|
});
|
|
sessions.push(session);
|
|
|
|
const events: AgentSessionEvent[] = [];
|
|
session.subscribe(event => {
|
|
events.push(event);
|
|
});
|
|
|
|
const assistantMessage = {
|
|
role: "assistant",
|
|
content: [
|
|
{
|
|
type: "toolCall",
|
|
id: "call_1",
|
|
name: "edit",
|
|
arguments: {},
|
|
partialJson: '{"file":"preview.txt","steps":[{"kbd":["ggdGi"],"insert":"rep',
|
|
},
|
|
],
|
|
api: "test",
|
|
provider: "test",
|
|
model: "test",
|
|
usage: {
|
|
input: 0,
|
|
output: 0,
|
|
cacheRead: 0,
|
|
cacheWrite: 0,
|
|
totalTokens: 0,
|
|
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
|
},
|
|
timestamp: Date.now(),
|
|
} as const;
|
|
|
|
session.agent.emitExternalEvent({
|
|
type: "message_update",
|
|
message: assistantMessage as never,
|
|
assistantMessageEvent: {
|
|
type: "toolcall_delta",
|
|
contentIndex: 0,
|
|
delta: "rep",
|
|
},
|
|
} as never);
|
|
|
|
await Bun.sleep(0);
|
|
|
|
expect(events.some(event => event.type === "message_update")).toBe(true);
|
|
expect(extensionEmit).toHaveBeenCalledTimes(1);
|
|
|
|
resolve();
|
|
await Bun.sleep(0);
|
|
});
|
|
});
|