From 9006e110eddcf5203f99ff07a177e3ea8e4efdb2 Mon Sep 17 00:00:00 2001 From: FFXN Date: Mon, 7 Sep 2026 18:07:34 +0800 Subject: [PATCH] fix(knowledge-import): handle temporary policy-write fence for website source imports --- api/tasks/knowledge_fs_source_import_tasks.py | 33 ++++---- .../test_knowledge_fs_source_import_tasks.py | 77 +++++++++++++++++++ ...rable-deletion-target-capabilities.test.ts | 28 +++++++ ...se-durable-deletion-target-capabilities.ts | 20 ++++- 4 files changed, 143 insertions(+), 15 deletions(-) diff --git a/api/tasks/knowledge_fs_source_import_tasks.py b/api/tasks/knowledge_fs_source_import_tasks.py index 8198a05e5d3..0b8d292b141 100644 --- a/api/tasks/knowledge_fs_source_import_tasks.py +++ b/api/tasks/knowledge_fs_source_import_tasks.py @@ -122,19 +122,26 @@ def finalize_source_import_once( expected_revision = current_policy.revision except KnowledgeFSProductResourceNotFoundError: expected_revision = 0 - facade.update_source_sync_policy( - tenant_id=tenant_id, - account_id=account_id, - control_space_id=control_space_id, - source_id=source_id, - payload=KnowledgeFSSourceSyncPolicyPayload( - enabled=desired_policy.enabled, - mode=desired_policy.mode, - customIntervalSeconds=desired_policy.custom_interval_seconds, - expectedRevision=expected_revision, - expectedSourceVersion=source.version, - ), - ) + try: + facade.update_source_sync_policy( + tenant_id=tenant_id, + account_id=account_id, + control_space_id=control_space_id, + source_id=source_id, + payload=KnowledgeFSSourceSyncPolicyPayload( + enabled=desired_policy.enabled, + mode=desired_policy.mode, + customIntervalSeconds=desired_policy.custom_interval_seconds, + expectedRevision=expected_revision, + expectedSourceVersion=source.version, + ), + ) + except KnowledgeFSProductRequestRejectedError as exc: + # A website selection replacement can still be draining its child document deletions after + # the parent workflow completes. KnowledgeFS reports that temporary policy-write fence as 400. + if exc.status_code != 400 or import_metadata.get("kind") != "website-crawl-import": + raise + raise KnowledgeFSSourceImportNotReadyError("Website source import cleanup is still running") from exc return workflow_id diff --git a/api/tests/unit_tests/tasks/test_knowledge_fs_source_import_tasks.py b/api/tests/unit_tests/tasks/test_knowledge_fs_source_import_tasks.py index 14785acedc8..88bd865f6bf 100644 --- a/api/tests/unit_tests/tasks/test_knowledge_fs_source_import_tasks.py +++ b/api/tests/unit_tests/tasks/test_knowledge_fs_source_import_tasks.py @@ -1,4 +1,5 @@ from types import SimpleNamespace +from typing import Literal from unittest.mock import MagicMock, patch import pytest @@ -194,6 +195,60 @@ def test_finalize_source_import_persists_failure_on_visible_source() -> None: facade.update_source_sync_policy.assert_not_called() +def test_finalize_website_source_import_waits_for_policy_write_fence() -> None: + facade = MagicMock() + facade.get_source_workflow.return_value = SimpleNamespace(id="import-1", state="completed") + facade.get_source.return_value = SimpleNamespace( + id="source-1", + metadata={ + "lastImport": { + "kind": "website-crawl-import", + "state": "completed", + "syncPolicy": {"enabled": True, "mode": "interval"}, + "workflowId": "import-1", + }, + "preview": False, + }, + status="active", + version=5, + ) + facade.get_source_sync_policy.side_effect = KnowledgeFSProductResourceNotFoundError("missing") + facade.update_source_sync_policy.side_effect = KnowledgeFSProductRequestRejectedError(status_code=400) + + with pytest.raises(KnowledgeFSSourceImportNotReadyError): + _run(facade) + + +@pytest.mark.parametrize("import_kind", ["online-document-import", "online-drive-import"]) +def test_finalize_non_website_source_import_does_not_reclassify_a_bad_request( + import_kind: str, +) -> None: + facade = MagicMock() + facade.get_source_workflow.return_value = SimpleNamespace(id="import-1", state="completed") + facade.get_source.return_value = SimpleNamespace( + id="source-1", + metadata={ + "lastImport": { + "kind": import_kind, + "state": "completed", + "syncPolicy": {"enabled": True, "mode": "interval"}, + "workflowId": "import-1", + }, + "preview": False, + }, + status="active", + version=5, + ) + facade.get_source_sync_policy.side_effect = KnowledgeFSProductResourceNotFoundError("missing") + rejection = KnowledgeFSProductRequestRejectedError(status_code=400) + facade.update_source_sync_policy.side_effect = rejection + + with pytest.raises(KnowledgeFSProductRequestRejectedError) as raised: + _run(facade) + + assert raised.value is rejection + + def test_finalize_source_import_retries_a_source_revision_conflict(monkeypatch: pytest.MonkeyPatch) -> None: error = KnowledgeFSProductRequestRejectedError(status_code=409) monkeypatch.setattr(task_module, "finalize_source_import_once", MagicMock(side_effect=error)) @@ -210,3 +265,25 @@ def test_finalize_source_import_retries_a_source_revision_conflict(monkeypatch: ) retry.assert_called_once_with(exc=error) + + +@pytest.mark.parametrize("status_code", [400, 422]) +def test_finalize_source_import_does_not_retry_a_permanent_rejection( + monkeypatch: pytest.MonkeyPatch, status_code: Literal[400, 422] +) -> None: + error = KnowledgeFSProductRequestRejectedError(status_code=status_code) + monkeypatch.setattr(task_module, "finalize_source_import_once", MagicMock(side_effect=error)) + retry = MagicMock(side_effect=Retry()) + monkeypatch.setattr(task_module.finalize_source_import, "retry", retry) + + with pytest.raises(KnowledgeFSProductRequestRejectedError) as raised: + task_module.finalize_source_import.run( + tenant_id="tenant-1", + account_id="account-1", + control_space_id="control-1", + source_id="source-1", + workflow_id="import-1", + ) + + assert raised.value is error + retry.assert_not_called() diff --git a/knowledge-fs/packages/api/src/database-durable-deletion-target-capabilities.test.ts b/knowledge-fs/packages/api/src/database-durable-deletion-target-capabilities.test.ts index c0b4cd6f566..b27d72449e5 100644 --- a/knowledge-fs/packages/api/src/database-durable-deletion-target-capabilities.test.ts +++ b/knowledge-fs/packages/api/src/database-durable-deletion-target-capabilities.test.ts @@ -3067,6 +3067,34 @@ describe("database durable deletion target capability edge branches", () => { }, ); + it.each(["postgres", "tidb"] as const)( + "keeps the parent Source syncing for a website replacement child deletion (%s)", + async (dialect) => { + const calls: DatabaseExecuteInput[] = []; + const workflowId = "018f0d60-7a49-7cc2-9c1b-5b36f18f2d20"; + const childJob = job({ + idempotencyKey: `source-remote-missing:${workflowId}:${targetDocumentId}`, + targetId: targetDocumentId, + targetType: "logical_document", + }); + + await capabilitiesFor(dialect, async (input) => { + calls.push(input); + return result([]); + }).quiesce({ job: childJob, signal: new AbortController().signal }); + + const sourceCancellation = calls.find( + (call) => call.operation === "update" && call.tableName === "sources", + ); + expect(sourceCancellation?.params).toContain(workflowId); + expect(sourceCancellation?.sql).toContain("source_workflow_runs"); + expect(sourceCancellation?.sql).toContain("crawl-import"); + expect(sourceCancellation?.sql).toContain("sync"); + expect(sourceCancellation?.sql).toContain("web"); + expect(sourceCancellation?.sql).toContain("NOT EXISTS"); + }, + ); + it.each(["postgres", "tidb"] as const)( "cancels only staged commits owned by the child deletion target (%s)", async (dialect) => { diff --git a/knowledge-fs/packages/api/src/database-durable-deletion-target-capabilities.ts b/knowledge-fs/packages/api/src/database-durable-deletion-target-capabilities.ts index 6799665a1ee..ea1e27834b9 100644 --- a/knowledge-fs/packages/api/src/database-durable-deletion-target-capabilities.ts +++ b/knowledge-fs/packages/api/src/database-durable-deletion-target-capabilities.ts @@ -1739,6 +1739,7 @@ async function cancelScopedWork( const p = (position: number) => databasePlaceholder(database, position); const nowIso = new Date().toISOString(); const nowMs = Date.now(); + const replacementWorkflowId = websiteReplacementWorkflowId(job); await database.transaction(async (transaction) => { await assertJobFence(database, transaction, job); await transaction.execute({ @@ -1798,14 +1799,16 @@ async function cancelScopedWork( params: job.targetType === "knowledge_space" ? [job.knowledgeSpaceId, nowIso] - : [job.knowledgeSpaceId, job.targetId, nowIso], + : job.targetType === "logical_document" && replacementWorkflowId + ? [job.knowledgeSpaceId, job.targetId, nowIso, job.tenantId, replacementWorkflowId] + : [job.knowledgeSpaceId, job.targetId, nowIso], sql: job.targetType === "knowledge_space" ? `UPDATE ${q("sources")} SET ${q("status")} = 'disabled', ${q("updated_at")} = ${p(2)}, ${q("version")} = ${q("version")} + 1 WHERE ${q("knowledge_space_id")} = ${p(1)} AND ${q("status")} = 'syncing';` : job.targetType === "source" ? `UPDATE ${q("sources")} SET ${q("status")} = 'disabled', ${q("updated_at")} = ${p(3)}, ${q("version")} = ${q("version")} + 1 WHERE ${q("knowledge_space_id")} = ${p(1)} AND ${q("id")} = ${p(2)} AND ${q("status")} = 'syncing';` : job.targetType === "logical_document" - ? `UPDATE ${q("sources")} SET ${q("status")} = 'disabled', ${q("updated_at")} = ${p(3)}, ${q("version")} = ${q("version")} + 1 WHERE ${q("knowledge_space_id")} = ${p(1)} AND ${q("id")} IN (SELECT ${q("source_id")} FROM ${q("logical_documents")} WHERE ${q("knowledge_space_id")} = ${p(1)} AND ${q("id")} = ${p(2)} AND ${q("source_id")} IS NOT NULL) AND ${q("status")} = 'syncing';` + ? `UPDATE ${q("sources")} SET ${q("status")} = 'disabled', ${q("updated_at")} = ${p(3)}, ${q("version")} = ${q("version")} + 1 WHERE ${q("knowledge_space_id")} = ${p(1)} AND ${q("id")} IN (SELECT ${q("source_id")} FROM ${q("logical_documents")} WHERE ${q("knowledge_space_id")} = ${p(1)} AND ${q("id")} = ${p(2)} AND ${q("source_id")} IS NOT NULL) AND ${q("status")} = 'syncing'${replacementWorkflowId ? ` AND NOT EXISTS (SELECT 1 FROM ${q("source_workflow_runs")} replacement_workflow WHERE replacement_workflow.${q("tenant_id")} = ${p(4)} AND replacement_workflow.${q("knowledge_space_id")} = ${p(1)} AND replacement_workflow.${q("id")} = ${p(5)} AND replacement_workflow.${q("source_id")} = ${q("sources")}.${q("id")} AND (replacement_workflow.${q("kind")} = 'crawl-import' OR (replacement_workflow.${q("kind")} = 'sync' AND ${q("sources")}.${q("type")} = 'web')))` : ""};` : `UPDATE ${q("sources")} SET ${q("status")} = 'disabled', ${q("updated_at")} = ${p(3)}, ${q("version")} = ${q("version")} + 1 WHERE ${q("knowledge_space_id")} = ${p(1)} AND ${q("id")} IN (SELECT ${q("source_id")} FROM ${q("document_assets")} WHERE ${q("knowledge_space_id")} = ${p(1)} AND ${q("id")} = ${p(2)} AND ${q("source_id")} IS NOT NULL) AND ${q("status")} = 'syncing';`, tableName: "sources", }); @@ -1813,6 +1816,19 @@ async function cancelScopedWork( }); } +function websiteReplacementWorkflowId( + job: DurableDeletionTargetOperationInput["job"], +): string | undefined { + if (job.targetType !== "logical_document") return undefined; + const prefix = "source-remote-missing:"; + const suffix = `:${job.targetId}`; + if (!job.idempotencyKey.startsWith(prefix) || !job.idempotencyKey.endsWith(suffix)) { + return undefined; + } + const workflowId = job.idempotencyKey.slice(prefix.length, -suffix.length); + return workflowId && !workflowId.includes(":") ? workflowId : undefined; +} + async function cancelWholeSpaceOpaqueWriters( transaction: DatabaseExecutor, database: DatabaseAdapter,