mirror of
https://github.com/langgenius/dify.git
synced 2026-09-08 11:04:27 +08:00
fix(knowledge-import): handle temporary policy-write fence for website source imports
This commit is contained in:
parent
d0364e6963
commit
9006e110ed
@ -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
|
||||
|
||||
|
||||
|
||||
@ -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()
|
||||
|
||||
@ -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) => {
|
||||
|
||||
@ -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,
|
||||
|
||||
Loading…
Reference in New Issue
Block a user