import { expect, test } from '@playwright/test'; import { Client } from 'pg'; test.skip(true, '2.0.0 S5 单人版已移除市场、handoff 或多用户前端入口,本 spec 隔离保留待恢复'); /** * E4 Market KB 召回触达 Knowledge 物化副本 e2e(真实后端 + 真 PG)。 * * 覆盖链路: * 1. 用户获取市场 KB 授权,并经 handoff 绑定到作品; * 2. Knowledge 把市场 KB 物化为安装者本地 installed_ref KB; * 3. 管理端召回该市场资产; * 4. Market 发布来源状态事件,Knowledge AFTER_COMMIT 消费后把本地 projection 置 recalled/blocked; * 5. 用户检索该作品知识时,后端来源状态门 fail-closed,把该 KB 作为 omittedSources 返回。 */ const API = 'http://localhost:48080/app-api/muse'; const ADMIN_API = 'http://localhost:48080/admin-api/muse'; const AUTH = { Authorization: 'Bearer test1', 'tenant-id': '1', 'X-API-Version': '1' }; const ADMIN_AUTH = { Authorization: 'Bearer test1', 'tenant-id': '1', 'X-API-Version': '1' }; const WORK_ID = 4; const KNOWLEDGE_HANDOFF_ASSET_NAME = '活体市场知识库·handoff e2e'; function createPgClient(): Client { return new Client({ host: process.env.MUSE_POSTGRES_HOST, port: Number(process.env.MUSE_POSTGRES_PORT ?? '5433'), database: process.env.MUSE_POSTGRES_DATABASE ?? 'muse_slice_live', user: process.env.MUSE_POSTGRES_USERNAME ?? 'root', password: process.env.MUSE_POSTGRES_PASSWORD, ssl: false, connectionTimeoutMillis: 10_000, }); } async function readKnowledgeHandoffAssetId(): Promise { const client = createPgClient(); await client.connect(); try { const result = await client.query( `SELECT id FROM muse_market_asset WHERE tenant_id=1 AND name=$1 AND asset_type='knowledge_base' AND deleted=false ORDER BY id DESC LIMIT 1`, [KNOWLEDGE_HANDOFF_ASSET_NAME] ); expect(result.rowCount, '应存在 knowledge_base handoff e2e 市场资产').toBe(1); return Number(result.rows[0].id); } finally { await client.end(); } } async function resetKnowledgeHandoffFixture(assetId: number): Promise { const client = createPgClient(); await client.connect(); try { const installed = await client.query( `SELECT id FROM muse_knowledge_base WHERE tenant_id=1 AND owner_user_id=1 AND source_market_asset_id=$1 AND kb_type='installed_ref' AND deleted=false`, [assetId] ); for (const row of installed.rows) { const kbId = Number(row.id); await client.query('DELETE FROM muse_knowledge_source_binding_projection WHERE tenant_id=1 AND kb_id=$1', [kbId]); await client.query('DELETE FROM muse_knowledge_binding WHERE tenant_id=1 AND kb_id=$1', [kbId]); await client.query('DELETE FROM muse_knowledge_ragflow_call WHERE tenant_id=1 AND kb_id=$1', [kbId]); await client.query('DELETE FROM muse_knowledge_processing_task WHERE tenant_id=1 AND kb_id=$1', [kbId]); await client.query('DELETE FROM muse_knowledge_ragflow_binding WHERE tenant_id=1 AND kb_id=$1', [kbId]); await client.query( `DELETE FROM muse_knowledge_document_version WHERE tenant_id=1 AND document_id IN (SELECT id FROM muse_knowledge_document WHERE tenant_id=1 AND kb_id=$1)`, [kbId] ); await client.query('DELETE FROM muse_knowledge_document WHERE tenant_id=1 AND kb_id=$1', [kbId]); await client.query('DELETE FROM muse_knowledge_base WHERE tenant_id=1 AND id=$1', [kbId]); } const sourceSnapshotId = `source-market_kb-${assetId}-v1`; await client.query( 'DELETE FROM muse_knowledge_bind_precheck WHERE tenant_id=1 AND owner_user_id=1 AND work_id=$1 AND source_snapshot_id=$2', [WORK_ID, sourceSnapshotId] ); await client.query('DELETE FROM muse_market_handoff_event WHERE tenant_id=1 AND owner_user_id=1 AND asset_id=$1', [ assetId, ]); await client.query('DELETE FROM muse_market_authorization_summary WHERE tenant_id=1 AND owner_user_id=1 AND asset_id=$1', [ assetId, ]); await client.query('DELETE FROM muse_market_installation WHERE tenant_id=1 AND user_id=1 AND asset_id=$1', [assetId]); await client.query('DELETE FROM muse_market_authorization_snapshot WHERE tenant_id=1 AND owner_user_id=1 AND asset_id=$1', [ assetId, ]); await client.query('DELETE FROM muse_market_source_status_event WHERE tenant_id=1 AND asset_id=$1', [assetId]); await client.query('DELETE FROM muse_market_command WHERE tenant_id=1 AND target_id=$1', [assetId]); await client.query('DELETE FROM muse_market_command WHERE tenant_id=1 AND target_id=$1', [String(assetId)]); await client.query( `UPDATE muse_market_asset SET listing_status='listed', status='active', tags=(coalesce(tags, '{}'::jsonb) - 'recheckReasons' - 'lastGovernanceAction' - 'lastGovernanceScope') || '{"forkStatus":"ready","actionPolicy":"allowed"}'::jsonb, update_time=CURRENT_TIMESTAMP WHERE tenant_id=1 AND id=$1 AND deleted=false`, [assetId] ); } finally { await client.end(); } } async function readWorkRevision(workId: number): Promise { const client = createPgClient(); await client.connect(); try { const result = await client.query('SELECT revision FROM muse_content_work WHERE tenant_id=1 AND id=$1', [workId]); expect(result.rowCount, `应存在 work ${workId}`).toBe(1); return Number(result.rows[0].revision ?? 1); } finally { await client.end(); } } async function readMaterializedKnowledge(assetId: number): Promise<{ kbId: number; bindingId: number; projectionStatus: string; actionPolicy: string; } | null> { const client = createPgClient(); await client.connect(); try { const result = await client.query( `SELECT kb.id AS kb_id, b.id AS binding_id, p.status AS projection_status, p.action_policy FROM muse_knowledge_base kb JOIN muse_knowledge_binding b ON b.tenant_id=kb.tenant_id AND b.kb_id=kb.id AND b.work_id=$2 AND b.binding_status='active' AND b.deleted=false JOIN muse_knowledge_source_binding_projection p ON p.tenant_id=kb.tenant_id AND p.kb_id=kb.id AND p.work_id=$2 AND p.binding_id=b.id AND p.deleted=false WHERE kb.tenant_id=1 AND kb.owner_user_id=1 AND kb.source_market_asset_id=$1 AND kb.kb_type='installed_ref' AND kb.deleted=false ORDER BY kb.id DESC LIMIT 1`, [assetId, WORK_ID] ); if (result.rowCount === 0) return null; return { kbId: Number(result.rows[0].kb_id), bindingId: Number(result.rows[0].binding_id), projectionStatus: String(result.rows[0].projection_status), actionPolicy: String(result.rows[0].action_policy), }; } finally { await client.end(); } } async function waitForProjectionBlocked(assetId: number): Promise<{ kbId: number; bindingId: number }> { const deadline = Date.now() + 20_000; while (Date.now() < deadline) { const row = await readMaterializedKnowledge(assetId); if (row?.projectionStatus === 'recalled' && row.actionPolicy === 'blocked') { return { kbId: row.kbId, bindingId: row.bindingId }; } await new Promise((resolve) => setTimeout(resolve, 500)); } const last = await readMaterializedKnowledge(assetId); throw new Error(`Knowledge projection 未被召回阻断: ${JSON.stringify(last)}`); } test.setTimeout(90_000); test('召回市场 KB 后,已物化 installed_ref 被阻断并在检索中省略', async ({ request }) => { const assetId = await readKnowledgeHandoffAssetId(); await resetKnowledgeHandoffFixture(assetId); try { // 1. 来源侧授权 + bind-precheck。 await request .post(`${API}/marketplace/assets/${assetId}/purchase`, { headers: AUTH, data: { commandId: `e2e-kb-recall-acq-${Date.now()}` }, }) .catch(() => undefined); const bp = await request.post(`${API}/marketplace/assets/${assetId}/bind-precheck`, { headers: AUTH, data: { commandId: `e2e-kb-recall-bp-${Date.now()}`, targetOwner: 'knowledge', targetAction: 'bind', targetWorkId: WORK_ID, }, }); const bpBody = await bp.json(); expect(bpBody.code, `bind-precheck: ${JSON.stringify(bpBody)}`).toBe(0); expect(bpBody.data?.handoffReady).toBe(true); // 2. 创建一次性 handoff token。 const handoff = await request.post(`${API}/marketplace/handoffs`, { headers: AUTH, data: { commandId: `e2e-kb-recall-ho-${Date.now()}`, assetId: String(assetId), targetOwner: 'knowledge', targetAction: 'bind', targetWorkId: WORK_ID, authorizationSummaryId: bpBody.data.authorizationSummaryId, returnUrl: 'http://localhost/knowledge/4', }, }); const handoffBody = await handoff.json(); expect(handoffBody.code, `createHandoff: ${JSON.stringify(handoffBody)}`).toBe(0); expect(handoffBody.data?.handoffToken).toBeTruthy(); // 3. Knowledge 兑现 precheck 并绑定,binding.kb_id 必须落本地 installed_ref KB,而不是市场 asset id。 const kbPrecheck = await request.post(`${API}/works/${WORK_ID}/knowledge-bindings/prechecks`, { headers: AUTH, data: { commandId: `e2e-kb-recall-kbpc-${Date.now()}`, sourceType: 'market_kb', sourceId: String(assetId), sourceVersion: Number(bpBody.data.sourceVersion ?? 1), sourceStatus: bpBody.data.sourceStatus ?? 'available', handoffToken: handoffBody.data.handoffToken, authorizationSummaryId: String(bpBody.data.authorizationSummaryId), authorizationSnapshotId: String(bpBody.data.authorizationSnapshot ?? `authz-${bpBody.data.authorizationSummaryId}`), purposes: ['search'], }, }); const kbPrecheckBody = await kbPrecheck.json(); expect(kbPrecheckBody.code, `knowledge precheck: ${JSON.stringify(kbPrecheckBody)}`).toBe(0); expect(kbPrecheckBody.data?.kbBindPrecheckId).toBeTruthy(); const expectedWorkRevision = await readWorkRevision(WORK_ID); const bind = await request.post(`${API}/works/${WORK_ID}/knowledge-bindings`, { headers: AUTH, data: { commandId: `e2e-kb-recall-bind-${Date.now()}`, kbBindPrecheckId: kbPrecheckBody.data.kbBindPrecheckId, expectedWorkRevision, }, }); const bindBody = await bind.json(); expect(bindBody.code, `knowledge bind: ${JSON.stringify(bindBody)}`).toBe(0); const materialized = await readMaterializedKnowledge(assetId); expect(materialized, '绑定后应存在 installed_ref KB + active binding/projection').not.toBeNull(); expect(materialized!.projectionStatus).toBe('active'); expect(materialized!.actionPolicy).toBe('allowed'); // 4. 管理端 preview + recall 真实执行;recall 事务提交后通过 Spring 事件通知 Knowledge。 const preview = await request.post(`${ADMIN_API}/market/assets/${assetId}/governance-impact`, { headers: ADMIN_AUTH, data: { commandId: `e2e-kb-recall-preview-${Date.now()}`, actionType: 'recall', scope: 'full_recall', }, }); const previewBody = await preview.json(); expect(previewBody.code, `admin preview: ${JSON.stringify(previewBody)}`).toBe(0); expect(previewBody.data?.previewId).toBeTruthy(); const recall = await request.post(`${ADMIN_API}/market/assets/${assetId}/recall`, { headers: ADMIN_AUTH, data: { commandId: `e2e-kb-recall-action-${Date.now()}`, expectedStatus: 'listed', impactPreviewId: previewBody.data.previewId, reason: 'e2e market KB recall propagation', basis: 'e2e verified recall event', scope: 'full_recall', exportPreservation: true, }, }); const recallBody = await recall.json(); expect(recallBody.code, `admin recall: ${JSON.stringify(recallBody)}`).toBe(0); const blocked = await waitForProjectionBlocked(assetId); expect(blocked.kbId).toBe(materialized!.kbId); // 5. 检索后端来源状态门必须省略该 KB,不能继续把已召回市场 KB 送入 RAGFlow。 const retrieval = await request.post(`${API}/works/${WORK_ID}/knowledge-retrievals`, { headers: AUTH, data: { commandId: `e2e-kb-recall-retrieval-${Date.now()}`, question: '召回后的市场知识库是否仍可检索?', topK: 3, }, }); const retrievalBody = await retrieval.json(); expect(retrievalBody.code, `retrieval: ${JSON.stringify(retrievalBody)}`).toBe(0); expect(retrievalBody.data?.status).toBe('empty'); expect(retrievalBody.data?.omittedReason).toBe('not_authorized'); expect( (retrievalBody.data?.omittedSources ?? []).some((source: { kbId?: string; reason?: string; sourceStatus?: string }) => source.kbId === String(blocked.kbId) && source.reason === 'not_authorized' && source.sourceStatus === 'recalled' ) ).toBe(true); } finally { await resetKnowledgeHandoffFixture(assetId); } });