- 上下文清单、任务列表与详情、恢复控制、运行观察、创作对话与消息流、正式变更决策面板、正文写作与候选审阅、资料研究等页面同步本轮后端能力。 - 事件标签表补齐后端事件类型;查询键与生成客户端随 HTTP 契约重新生成。 - 依赖调整后 lockfile 同步;交互用例与浏览器旅程 spec 一并更新。 - 证据:make 前端检查(lint + typecheck + Pi 工具桥 tsc)exit 0;make 前端测试 37 文件 / 120 用例通过。
273 lines
8.8 KiB
TypeScript
273 lines
8.8 KiB
TypeScript
import { afterEach, expect, it, vi } from "vitest";
|
||
import { act, cleanup, render, screen } from "@testing-library/react";
|
||
import { QueryClient, QueryClientProvider } from "@tanstack/react-query";
|
||
import { useMuse任务订阅 } from "../../src/功能/任务运行/任务订阅";
|
||
import type { 任务 } from "../../src/功能/任务运行/数据";
|
||
|
||
class Source {
|
||
static instances: Source[] = [];
|
||
handlers = new Map<string, ((event: MessageEvent) => void)[]>();
|
||
onerror: (() => void) | null = null;
|
||
close = vi.fn();
|
||
constructor(public url: string) {
|
||
Source.instances.push(this);
|
||
}
|
||
addEventListener(type: string, callback: (event: MessageEvent) => void) {
|
||
this.handlers.set(type, [...(this.handlers.get(type) ?? []), callback]);
|
||
}
|
||
emit(type: string, data: unknown) {
|
||
this.handlers
|
||
.get(type)
|
||
?.forEach((callback) =>
|
||
callback(new MessageEvent(type, { data: JSON.stringify(data) })),
|
||
);
|
||
}
|
||
}
|
||
const running: 任务 = {
|
||
task_id: "task",
|
||
state: "running",
|
||
work_id: "work",
|
||
task_type: "生成正文",
|
||
supports_step_advance: false,
|
||
last_sequence: 0,
|
||
steps: [],
|
||
execution_final_sequence: null,
|
||
accounting_revision: 0,
|
||
result_refs: [],
|
||
};
|
||
function Probe({ task = running }: { task?: 任务 }) {
|
||
const { events, notice } = useMuse任务订阅(task, "work");
|
||
return (
|
||
<>
|
||
<output>{events.map((e) => e.sequence).join(",")}</output>
|
||
<p>{notice}</p>
|
||
</>
|
||
);
|
||
}
|
||
function event(sequence: number, type = "text.delta") {
|
||
return {
|
||
event_id: `event-${sequence}`,
|
||
sequence,
|
||
event_type: type,
|
||
occurred_at: "2026-09-17T00:00:00Z",
|
||
payload: { 文本: "合成增量" },
|
||
};
|
||
}
|
||
function setup(snapshot: () => 任务) {
|
||
vi.useFakeTimers();
|
||
Source.instances = [];
|
||
vi.stubGlobal("EventSource", Source);
|
||
const counts = { snapshot: 0, history: 0 };
|
||
vi.stubGlobal(
|
||
"fetch",
|
||
vi.fn(async (request: Request) => {
|
||
const path = new URL(request.url).pathname;
|
||
let data: unknown;
|
||
if (path.endsWith("/events")) {
|
||
counts.history++;
|
||
data = {
|
||
items: [],
|
||
reset_required: false,
|
||
first_sequence: 1,
|
||
last_sequence: snapshot().last_sequence,
|
||
execution_final_sequence: snapshot().execution_final_sequence,
|
||
};
|
||
} else {
|
||
counts.snapshot++;
|
||
data = snapshot();
|
||
}
|
||
return new Response(JSON.stringify(data), {
|
||
headers: { "Content-Type": "application/json" },
|
||
});
|
||
}),
|
||
);
|
||
const cache = new QueryClient({
|
||
defaultOptions: { queries: { retry: false } },
|
||
});
|
||
return {
|
||
counts,
|
||
cache,
|
||
view: (task = running) =>
|
||
render(
|
||
<QueryClientProvider client={cache}>
|
||
<Probe task={task} />
|
||
</QueryClientProvider>,
|
||
),
|
||
};
|
||
}
|
||
afterEach(() => {
|
||
cleanup();
|
||
vi.useRealTimers();
|
||
vi.unstubAllGlobals();
|
||
});
|
||
const flush = () =>
|
||
act(async () => {
|
||
await vi.advanceTimersByTimeAsync(0);
|
||
});
|
||
|
||
// @muse-case {"case_id":"O08-stream-terminal-idle","then":"已终态任务只读取一页有限历史,不建立持续SSE;推进60秒不重读"}
|
||
it("O08-stream-terminal-idle:已终态快照没有持续事件连接", async () => {
|
||
const terminal: 任务 = {
|
||
...running,
|
||
state: "cancelled",
|
||
last_sequence: 50,
|
||
execution_final_sequence: 50,
|
||
};
|
||
const { counts, view } = setup(() => terminal);
|
||
view(terminal);
|
||
await flush();
|
||
expect(Source.instances).toHaveLength(0);
|
||
expect(counts).toEqual({ snapshot: 0, history: 1 });
|
||
await act(() => vi.advanceTimersByTimeAsync(60_000));
|
||
expect(counts).toEqual({ snapshot: 0, history: 1 });
|
||
});
|
||
|
||
// @muse-case {"case_id":"O08-stream-delta-final-close","then":"100条text.delta不刷新任务;状态合并读一次;最终游标只close一次且60秒无续连"}
|
||
it("O08-stream-delta-final-close:增量不刷快照,执行最终游标双端收尾", async () => {
|
||
let snapshot = running;
|
||
const { counts, view } = setup(() => snapshot);
|
||
view();
|
||
await flush();
|
||
const source = Source.instances[0];
|
||
act(() => {
|
||
for (let i = 1; i <= 100; i++) source.emit("task-event", event(i));
|
||
});
|
||
await act(() => vi.advanceTimersByTimeAsync(200));
|
||
expect(counts.snapshot).toBe(0);
|
||
snapshot = { ...running, last_sequence: 102 };
|
||
act(() => {
|
||
source.emit("task-event", event(101, "step.completed"));
|
||
source.emit("task-event", event(102, "candidate.saved"));
|
||
});
|
||
await act(() => vi.advanceTimersByTimeAsync(100));
|
||
expect(counts.snapshot).toBe(1);
|
||
snapshot = {
|
||
...running,
|
||
state: "completed",
|
||
last_sequence: 103,
|
||
execution_final_sequence: 103,
|
||
};
|
||
act(() => {
|
||
source.emit("task-event", event(103, "task.completed"));
|
||
source.emit("stream-complete", { execution_final_sequence: 103 });
|
||
});
|
||
await flush();
|
||
expect(source.close).toHaveBeenCalledTimes(1);
|
||
expect(counts.snapshot).toBe(2);
|
||
await act(() => vi.advanceTimersByTimeAsync(60_000));
|
||
expect(counts.snapshot).toBe(2);
|
||
expect(Source.instances).toHaveLength(1);
|
||
});
|
||
|
||
// @muse-case {"case_id":"O08-stream-reset-resume","then":"reset关闭旧源,只读一次权威快照后从新last_sequence续连,旧源及重复序号不追加"}
|
||
it("O08-stream-reset-resume:缺口按新快照恢复且重复事件只消费一次", async () => {
|
||
const snapshot = { ...running, last_sequence: 10 };
|
||
const { counts, view } = setup(() => snapshot);
|
||
view();
|
||
await flush();
|
||
const old = Source.instances[0];
|
||
act(() => old.emit("reset", {}));
|
||
await flush();
|
||
expect(old.close).toHaveBeenCalledTimes(1);
|
||
expect(counts.snapshot).toBe(1);
|
||
expect(Source.instances).toHaveLength(2);
|
||
const next = Source.instances[1];
|
||
expect(next.url).toContain("cursor=10");
|
||
act(() => {
|
||
old.emit("task-event", event(1));
|
||
next.emit("task-event", event(11));
|
||
next.emit("task-event", event(11));
|
||
});
|
||
expect(screen.getByRole("status").textContent).toBe("11");
|
||
});
|
||
|
||
// @muse-case {"case_id":"O08-stream-reset-after-inflight","then":"reset等旧快照请求结束后只发一次新权威读取,后续游标采用新快照而非旧在途结果"}
|
||
it("O08-stream-reset-after-inflight:缺口不把较早在途快照当作恢复结果", async () => {
|
||
const { counts, view } = setup(() => running);
|
||
let finish!: (response: Response) => void;
|
||
let snapshots = 0;
|
||
vi.stubGlobal(
|
||
"fetch",
|
||
vi.fn(async (request: Request) => {
|
||
const path = new URL(request.url).pathname;
|
||
if (path.endsWith("/events"))
|
||
return new Response(
|
||
JSON.stringify({
|
||
items: [],
|
||
reset_required: false,
|
||
first_sequence: 1,
|
||
last_sequence: 0,
|
||
}),
|
||
{ headers: { "Content-Type": "application/json" } },
|
||
);
|
||
snapshots++;
|
||
if (snapshots === 1)
|
||
return new Promise<Response>((resolve) => {
|
||
finish = resolve;
|
||
});
|
||
return new Response(JSON.stringify({ ...running, last_sequence: 10 }), {
|
||
headers: { "Content-Type": "application/json" },
|
||
});
|
||
}),
|
||
);
|
||
view();
|
||
await flush();
|
||
const old = Source.instances[0];
|
||
act(() => old.emit("task-event", event(1, "step.started")));
|
||
await act(() => vi.advanceTimersByTimeAsync(100));
|
||
expect(snapshots).toBe(1);
|
||
act(() => old.emit("reset", {}));
|
||
await flush();
|
||
expect(snapshots).toBe(1);
|
||
await act(async () =>
|
||
finish(
|
||
new Response(JSON.stringify({ ...running, last_sequence: 1 }), {
|
||
headers: { "Content-Type": "application/json" },
|
||
}),
|
||
),
|
||
);
|
||
await flush();
|
||
expect(snapshots).toBe(2);
|
||
expect(Source.instances).toHaveLength(2);
|
||
expect(Source.instances[1].url).toContain("cursor=10");
|
||
expect(counts.snapshot).toBe(0);
|
||
});
|
||
|
||
// @muse-case {"case_id":"O08-stream-unmount-aborts-read","then":"离开任务时关闭执行流并取消在途只读快照,迟到响应不产生新订阅"}
|
||
it("O08-stream-unmount-aborts-read:离开作用域结束订阅和在途读取", async () => {
|
||
const { view } = setup(() => running);
|
||
let signal: AbortSignal | undefined;
|
||
vi.stubGlobal(
|
||
"fetch",
|
||
vi.fn(async (request: Request) => {
|
||
if (new URL(request.url).pathname.endsWith("/events"))
|
||
return new Response(
|
||
JSON.stringify({
|
||
items: [],
|
||
reset_required: false,
|
||
first_sequence: 1,
|
||
last_sequence: 0,
|
||
}),
|
||
{ headers: { "Content-Type": "application/json" } },
|
||
);
|
||
signal = request.signal;
|
||
return new Promise<Response>((_, reject) =>
|
||
request.signal.addEventListener("abort", () =>
|
||
reject(new DOMException("aborted", "AbortError")),
|
||
),
|
||
);
|
||
}),
|
||
);
|
||
const mounted = view();
|
||
await flush();
|
||
const source = Source.instances[0];
|
||
act(() => source.emit("task-event", event(1, "step.started")));
|
||
await act(() => vi.advanceTimersByTimeAsync(100));
|
||
expect(signal?.aborted).toBe(false);
|
||
mounted.unmount();
|
||
await flush();
|
||
expect(signal?.aborted).toBe(true);
|
||
expect(source.close).toHaveBeenCalledTimes(1);
|
||
expect(Source.instances).toHaveLength(1);
|
||
});
|