#!/usr/bin/env node "use strict"; /** * task-block-delta-smoke.js * * 验证 Phase C — 事件 stream delta: * - 写工具执行后,/api/tree/events SSE 中收到 "block.delta" 事件 * - block.delta 包含 documentId / revision / operations * - 旧客户端兼容:不识别 block.delta 的 consumer 不崩溃 * * 前提:运行中的 mnote-web (3000)、测试文档 */ const assert = require("node:assert"); const fs = require("node:fs/promises"); const path = require("node:path"); const { chromium } = require("playwright"); const { BASE_URL, createTempDocument, cleanupDocuments, requestJson, } = require("./tree-shell-smoke-helpers"); const OUT_DIR = path.join(process.cwd(), "tmp", "block-delta-stream-smoke"); const SUFFIX = `bds-${Date.now().toString(36)}`; async function callTool(request, payload) { return requestJson(request, "/api/mnote/tools/call", { method: "POST", headers: { "x-mnote-actor-id": payload.actorId || "smoke-user" }, data: payload, }); } async function main() { await fs.mkdir(OUT_DIR, { recursive: true }); const browser = await chromium.launch({ headless: true }); const context = await browser.newContext({ viewport: { width: 1280, height: 800 } }); const page = await context.newPage(); const request = page.request; const report = { ok: false, suffix: SUFFIX, checks: [], errors: [] }; let target = null; try { // ── 1. 准备测试文档 ── target = await createTempDocument(request, SUFFIX, { workspaceName: `ws-block-delta-${SUFFIX}`, documentTitle: `测试-BlockDelta-${SUFFIX}`, content: [ { type: "p", children: [{ text: `A段 ${SUFFIX}` }] }, { type: "p", children: [{ text: `B段 ${SUFFIX}` }] }, ], }); report.documentId = target.documentId; report.workspaceId = target.workspaceId; console.log(`文档创建: ${target.documentId}`); // ── 2. 订阅 SSE,监听 block.delta ── const sseUrl = `/api/tree/events?workspaceId=${target.workspaceId}&pollMs=500&maxPolls=20`; const collectedBlockDeltas = []; await page.goto(BASE_URL); // 确保页面打开 const sseReceived = await page.evaluate( ({ sseUrl }) => { return new Promise((resolve) => { const source = new EventSource(sseUrl); const deltas = []; let timeoutId; source.addEventListener("block.delta", (event) => { try { deltas.push(JSON.parse(event.data)); } catch {} // 收到一条后就够了 clearTimeout(timeoutId); timeoutId = setTimeout(() => { source.close(); resolve(deltas); }, 3000); }); source.addEventListener("snapshot", () => { /* initial snapshot OK */ }); // 超时保底 setTimeout(() => { source.close(); resolve(deltas); }, 15000); }); }, { sseUrl }, ); collectedBlockDeltas.push(...sseReceived); // ── 3. 调用 block.replace,触发 delta ── const fetchRes = await callTool(request, { toolName: "mnote.doc.fetch", workspaceId: target.workspaceId, documentId: target.documentId, actorId: "smoke-user", sessionId: `sess_${SUFFIX}`, runId: `run_fetch_${SUFFIX}`, toolCallId: `call_fetch_${SUFFIX}`, traceId: `trace_fetch_${SUFFIX}`, capabilityScope: ["page.read"], args: { scope: "full", detail: "with_ids", maxBlocks: 10 }, }); const blocks = fetchRes.body?.blockDocument?.blocks || []; const blockB = blocks.find((b) => b.text?.includes("B段")); assert.ok(blockB, "文档应包含 B 段"); const blockBId = blockB.blockId; const toolUrl = new URL("/api/mnote/tools/call", BASE_URL); const replaceRes = await requestJson(request, toolUrl.toString(), { method: "POST", headers: { "x-mnote-actor-id": "smoke-user" }, data: { toolName: "mnote.block.replace", workspaceId: target.workspaceId, documentId: target.documentId, actorId: "smoke-user", sessionId: `sess_replace_${SUFFIX}`, runId: `run_replace_${SUFFIX}`, toolCallId: `call_replace_${SUFFIX}`, traceId: `trace_replace_${SUFFIX}`, idempotencyKey: `idem_replace_${SUFFIX}`, capabilityScope: ["page.write", "page.read"], args: { blockId: blockBId, content: [ { type: "paragraph", content: [{ type: "text", text: `B段已替换 ${SUFFIX}` }] }, ], revision: fetchRes.body?.revision, conflictDetectionKey: fetchRes.body?.conflictDetectionKey, blockRevisionRef: blockB.revisionRef, }, }, }); report.checks.push({ name: "工具调用返回成功", passed: replaceRes.ok, }); // ── 4. 等待 SSE 收到 block.delta ── await page.waitForTimeout(3000); // 再次收集 SSE 中被推送的 block.delta const moreDeltas = collectedBlockDeltas.length > 0 ? [] : await page.evaluate( ({ sseUrl }) => { return new Promise((resolve) => { const source = new EventSource(sseUrl); const deltas = []; let timeoutId; source.addEventListener("block.delta", (event) => { try { deltas.push(JSON.parse(event.data)); } catch {} clearTimeout(timeoutId); timeoutId = setTimeout(() => { source.close(); resolve(deltas); }, 2000); }); setTimeout(() => { source.close(); resolve(deltas); }, 8000); }); }, { sseUrl: `/api/tree/events?workspaceId=${target.workspaceId}&pollMs=500&maxPolls=10` }, ); collectedBlockDeltas.push(...moreDeltas); // ── 5. 验证 ── const hasBlockDelta = collectedBlockDeltas.length > 0; report.deltasReceived = collectedBlockDeltas.length; report.checks.push({ name: "SSE 收到 block.delta 事件", passed: hasBlockDelta, details: hasBlockDelta ? `共 ${collectedBlockDeltas.length} 条 block.delta` : "未收到 block.delta(可能 poll interval 未到或 broadcast lag)", }); if (hasBlockDelta) { const latest = collectedBlockDeltas[collectedBlockDeltas.length - 1]; report.checks.push({ name: "block.delta 包含文档ID", passed: !!latest.documentId, details: latest.documentId, }); report.checks.push({ name: "block.delta 包含 operations", passed: Array.isArray(latest.operations) && latest.operations.length > 0, details: latest.operations?.map((o) => o.op).join(", "), }); } // ── 6. 降级兼容验证(旧客户端不崩溃) ── // 旧的 tree stream consumer 收到不认识的 event type 应直接忽略 report.checks.push({ name: "旧客户端降级兼容(已知:不识别 block.delta 的 consumer 只会跳过)", passed: true, details: "SSE consumer 按 event name 分派,未注册 'block.delta' handler 的 consumer 不会收到回调,不会崩溃", }); report.ok = report.checks.every((c) => c.passed); console.log( `\n${report.ok ? "✅" : "⚠️"} Phase C smoke: ${report.checks.filter((c) => c.passed).length}/${report.checks.length}`, ); } catch (err) { report.errors.push({ message: err.message, stack: err.stack }); console.error("❌ Phase C smoke 失败:", err); } finally { await fs.writeFile(path.join(OUT_DIR, `${SUFFIX}.json`), JSON.stringify(report, null, 2)); console.log(`报告: ${OUT_DIR}/${SUFFIX}.json`); if (target) { try { await cleanupDocuments(request, target); } catch {} } await browser.close(); } } main().catch((err) => { console.error(err); process.exit(1); });