Files
penguin-harness/packages/server/test/sse-stream.test.ts
T
Yaowei Zheng 45bfae6e94 Initialize repository with harness code and assets
Initial import of all source code, config, and README assets: the
packages workspace (cli, core, server, web, docs, landing, skills),
build scripts, tooling config, and CI workflows.

Includes the data-layout revision made on this branch: the local data
root defaults to ~/.penguin/data (PENGUIN_HOME still overrides; the
installer keeps its binaries in ~/.penguin), and every Agent lives
under <project>/agents/<agent>/ — path helpers, the three
agent-enumeration scans, the system prompt, built-in Skills, tests
and docs all follow the new layout.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_018ihk8iQuo3kv2aPjAYEPuR
2026-07-19 14:06:53 +08:00

151 lines
6.1 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/**
* Integration tests for the SSE endpoint (FD-1 / FD-2):
* - Subscribing immediately delivers the current task_state snapshot (first
* frame, ahead of any pending-approval replay);
* - A Last-Event-ID with a mismatched/unknown epoch → sends resync_required
* first; a buffer hit → replays events after it.
*/
import { afterEach, beforeEach, describe, expect, it } from "vitest";
import { approvalDecision, assistantText, toolCall, userText } from "@prismshadow/penguin-core";
import type { ApproveFn, OmniMessage } from "@prismshadow/penguin-core";
import type { SessionRow } from "../src/db/repos/sessions.js";
import type { RuntimeSession } from "../src/runtime/session-manager.js";
import { apiClient, createTestApp, provisionUser, waitFor } from "./helpers.js";
import type { TestApp } from "./helpers.js";
const SID = "session-2026-07-06-10-00-00-aabb0001";
interface SseFrame {
event?: string;
id?: string;
data: string;
}
/** Reads the first `count` frames of an SSE response (skipping heartbeat comment lines), then cancels the stream. */
async function readSseFrames(res: Response, count: number, timeoutMs = 3000): Promise<SseFrame[]> {
expect(res.status).toBe(200);
const reader = res.body!.getReader();
const decoder = new TextDecoder();
const frames: SseFrame[] = [];
let buf = "";
const deadline = Date.now() + timeoutMs;
try {
while (frames.length < count) {
if (Date.now() > deadline) throw new Error(`SSE 读取超时(已有 ${frames.length} 帧)`);
const { done, value } = await reader.read();
if (done) break;
buf += decoder.decode(value, { stream: true });
let idx: number;
while ((idx = buf.indexOf("\n\n")) !== -1) {
const raw = buf.slice(0, idx);
buf = buf.slice(idx + 2);
if (raw.startsWith(":") || raw.trim() === "") continue; // heartbeat/empty frame
const frame: SseFrame = { data: "" };
for (const line of raw.split("\n")) {
if (line.startsWith("event:")) frame.event = line.slice(6).trim();
else if (line.startsWith("id:")) frame.id = line.slice(3).trim();
else if (line.startsWith("data:")) frame.data += line.slice(5).trim();
}
frames.push(frame);
}
}
} finally {
await reader.cancel().catch(() => {});
}
return frames;
}
/** Fake Session that requests one approval (for the running state and approval-replay scenarios). */
function approvalFakeSession(sessionId: string): RuntimeSession {
return {
sessionId,
toolPermission: () => "rw",
generateTitle: async () => ({ title: null, usage: null }),
compactability: () => "ok" as const,
async *run(_input: OmniMessage[], opts: { approve: ApproveFn; signal: AbortSignal }) {
const tc = toolCall({ name: "exec_command", arguments: "{}", toolCallId: "tc-sse" });
yield tc;
const decision = await opts.approve(tc);
yield approvalDecision(decision, "tc-sse");
yield assistantText("done");
},
async *compact() {},
};
}
describe("sse-stream", () => {
let t: TestApp;
let cookie: string;
let row: SessionRow;
beforeEach(async () => {
t = await createTestApp();
({ cookie } = await provisionUser(t.app, "streamer"));
row = {
sessionId: SID,
// streamer's own initial Project (default_project belongs to admin; others
// get 404 via the index lookup).
projectId: "streamer-default_project",
agentId: "default_agent",
modelId: "m1",
provider: "custom",
workspace: "/tmp/w",
approvalMode: "always-ask",
title: null,
createdAt: new Date().toISOString(),
};
t.deps.sessionsRepo.insert(row);
});
afterEach(async () => {
await t.cleanup();
});
const getStream = (headers: Record<string, string> = {}) =>
t.app.request(`/api/sessions/${SID}/stream`, { headers: { cookie, ...headers } });
it("FD-1:新订阅第一条恒为 task_state 快照(idle)", async () => {
const frames = await readSseFrames(await getStream(), 1);
expect(frames[0]!.event).toBe("server_event");
expect(JSON.parse(frames[0]!.data)).toEqual({ type: "task_state", state: "idle" });
expect(frames[0]!.id).toMatch(/^[0-9a-f]{8}-\d+$/); // FD-2: opaque string id
});
it("FD-1:运行中订阅收到 task_state: running,再补发未决审批", async () => {
t.deps.manager.adopt(row, approvalFakeSession(SID));
await t.deps.manager.startTask(SID, [userText("go")]);
await waitFor(() => t.deps.manager.pendingApprovalCount(SID) === 1);
const frames = await readSseFrames(await getStream(), 2);
expect(JSON.parse(frames[0]!.data)).toEqual({ type: "task_state", state: "running" });
const approval = JSON.parse(frames[1]!.data) as {
type: string;
toolCall: { payload: { tool_call_id: string } };
};
expect(approval.type).toBe("approval_request");
expect(approval.toolCall.payload.tool_call_id).toBe("tc-sse");
t.deps.manager.abortTask(SID);
await waitFor(() => t.deps.manager.statusOf(SID) === "idle");
});
it("FD-2:Last-Event-ID 纪元不一致 → 先发 resync_required 再发 task_state 快照", async () => {
t.deps.channels.get(SID).publish(userText("旧事件"));
const frames = await readSseFrames(
await getStream({ "Last-Event-ID": "deadbeef-1" }), // guaranteed to differ from the current channel epoch
2,
);
expect(JSON.parse(frames[0]!.data)).toEqual({ type: "resync_required" });
expect(JSON.parse(frames[1]!.data)).toEqual({ type: "task_state", state: "idle" });
});
it("FD-2:同纪元 Last-Event-ID 命中缓冲 → 补发其后事件,再发 task_state 快照", async () => {
const channel = t.deps.channels.get(SID);
const first = channel.publish(userText("m1"));
channel.publish(userText("m2"));
const frames = await readSseFrames(await getStream({ "Last-Event-ID": first.id }), 2);
expect(frames[0]!.event).toBeUndefined(); // replayed OmniMessage
expect((JSON.parse(frames[0]!.data) as { payload: { text: string } }).payload.text).toBe("m2");
expect(JSON.parse(frames[1]!.data)).toEqual({ type: "task_state", state: "idle" });
});
});