oh-my-muse/muse-studio/e2e/market-kb-recall.spec.ts
lili 3585219637
Some checks failed
Backend Maven CI / backend-local (push) Has been cancelled
feat(mvp): 收束1.0.0线A交付闭环
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-27 10:52:10 -07:00

326 lines
13 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.

import { expect, test } from '@playwright/test';
import { Client } from 'pg';
/**
* 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<number> {
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<void> {
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<number> {
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);
}
});