From 59b94535922f8a5b93650347b00a6fe1d80ee75a Mon Sep 17 00:00:00 2001 From: FFXN Date: Fri, 7 Aug 2026 17:18:44 +0800 Subject: [PATCH] fix(knowledge_fs): standardize payload field names in initial source tasks --- api/knowledge-fs-contract.lock.json | 6 +- api/knowledge-fs-product-operations.json | 1 + .../knowledge_fs_initial_source_tasks.py | 21 +- .../extensions/test_knowledge_fs_celery.py | 54 ++++ .../test_knowledge_fs_initial_source_tasks.py | 261 ++++++++++++++++-- 5 files changed, 302 insertions(+), 41 deletions(-) create mode 100644 api/tests/unit_tests/extensions/test_knowledge_fs_celery.py diff --git a/api/knowledge-fs-contract.lock.json b/api/knowledge-fs-contract.lock.json index 5ab0b8ae0fd..83fb6b13aa3 100644 --- a/api/knowledge-fs-contract.lock.json +++ b/api/knowledge-fs-contract.lock.json @@ -1,9 +1,9 @@ { "schemaVersion": 5, - "subtreeTree": "dd7e958bff7830fbf794fcc02f6472ea13a3e3f8", - "openapiSha256": "86843fc82966be7860f3e2f61254e05a02f1a81a8204bb8dd8489483b460cb20", + "subtreeTree": "bfba34eb0dd57039e5eb81dd9446d4bd2a521fa9", + "openapiSha256": "02df34e84bb26a7d510d620f4e3c2bf7cd1d6c5bb3656675fd9c20ce6e4ed41e", "capabilityV2AuthManifestSha256": "fc0a47e23cce12544882f0298522b4933002e892b84ce1815df7e81d36a7a0c7", "capabilityV2AuthTestVectorSha256": "ae0de37b1ff05c40f905cf17a7b410d8971acacf64db07d5ee3d6fecfa559ce3", - "productOperationManifestSha256": "a2d5e2b72b87652205d505ff85ce2ca689dbd96a466495883bbf1c4b8a54ff32", + "productOperationManifestSha256": "179dcc31fb2d701d39c12adabc28a2d105f8370e4fdf7174dd57491971d05a70", "productOperationGapManifestSha256": "332c80165bd5cf8e79bc511374dde771a8175510404db603a8067f1a22d36df8" } diff --git a/api/knowledge-fs-product-operations.json b/api/knowledge-fs-product-operations.json index 8aa8ec950aa..26dbf318da8 100644 --- a/api/knowledge-fs-product-operations.json +++ b/api/knowledge-fs-product-operations.json @@ -55,6 +55,7 @@ {"productOperationId":"listSourceConnections","kfsOperationId":"listSourceConnections","method":"GET","path":"/knowledge-spaces/{id}/source-connections","action":"source_connections.list","resource":"knowledge_space","transport":"json","stream":{"productKind":"json","kfsResponseKind":"buffered"},"limits":{"productMaxRequestBytes":16384,"productMaxResponseBytes":1048576,"kfsMaxResponseBytes":1048576}}, {"productOperationId":"refreshSourceConnection","kfsOperationId":"refreshSourceConnection","method":"POST","path":"/knowledge-spaces/{id}/source-connections/{connectionId}/refresh","action":"source_connections.refresh","resource":"knowledge_space","transport":"json","stream":{"productKind":"json","kfsResponseKind":"buffered"},"limits":{"productMaxRequestBytes":32768,"productMaxResponseBytes":524288,"kfsMaxResponseBytes":1048576}}, {"productOperationId":"previewSourceCrawl","kfsOperationId":"createSourceCrawlPreviewWorkflow","method":"POST","path":"/knowledge-spaces/{id}/sources/{sourceId}/crawl-preview","action":"source_workflows.preview.create","resource":"source","transport":"json","stream":{"productKind":"json","kfsResponseKind":"buffered"},"limits":{"productMaxRequestBytes":16384,"productMaxResponseBytes":524288,"kfsMaxResponseBytes":1048576}}, + {"productOperationId":"importSelectedSourceCrawl","kfsOperationId":"createSourceCrawlImportWorkflow","method":"POST","path":"/knowledge-spaces/{id}/sources/{sourceId}/crawl-import","action":"source_workflows.crawl_import.create","resource":"source","transport":"json","stream":{"productKind":"json","kfsResponseKind":"buffered"},"limits":{"productMaxRequestBytes":1048576,"productMaxResponseBytes":524288,"kfsMaxResponseBytes":1048576}}, {"productOperationId":"importSourceWorkflow","kfsOperationId":"createSourceImportWorkflow","method":"POST","path":"/knowledge-spaces/{id}/sources/{sourceId}/workflow-imports","action":"source_workflows.import.create","resource":"source","transport":"json","stream":{"productKind":"json","kfsResponseKind":"buffered"},"limits":{"productMaxRequestBytes":4194304,"productMaxResponseBytes":524288,"kfsMaxResponseBytes":1048576}}, {"productOperationId":"getSourceSyncPolicy","kfsOperationId":"getSourceSyncPolicy","method":"GET","path":"/knowledge-spaces/{id}/sources/{sourceId}/sync-policy","action":"source_sync_policies.read","resource":"source","transport":"json","stream":{"productKind":"json","kfsResponseKind":"buffered"},"limits":{"productMaxRequestBytes":0,"productMaxResponseBytes":262144,"kfsMaxResponseBytes":1048576}}, {"productOperationId":"updateSourceSyncPolicy","kfsOperationId":"putSourceSyncPolicy","method":"PUT","path":"/knowledge-spaces/{id}/sources/{sourceId}/sync-policy","action":"source_sync_policies.update","resource":"source","transport":"json","stream":{"productKind":"json","kfsResponseKind":"buffered"},"limits":{"productMaxRequestBytes":32768,"productMaxResponseBytes":262144,"kfsMaxResponseBytes":1048576}}, diff --git a/api/tasks/knowledge_fs_initial_source_tasks.py b/api/tasks/knowledge_fs_initial_source_tasks.py index ee314dcaa2a..2aa926b78a9 100644 --- a/api/tasks/knowledge_fs_initial_source_tasks.py +++ b/api/tasks/knowledge_fs_initial_source_tasks.py @@ -88,10 +88,7 @@ def start_initial_website_source_import( ) if control_space is None: raise RuntimeError("KnowledgeFS control-space was not found") - if ( - control_space.state is not KnowledgeFSControlSpaceState.ACTIVE - or control_space.knowledge_space_id is None - ): + if control_space.state is not KnowledgeFSControlSpaceState.ACTIVE or control_space.knowledge_space_id is None: if control_space.state is KnowledgeFSControlSpaceState.PROVISIONING: raise KnowledgeFSInitialSourceNotReadyError("KnowledgeFS Space is still provisioning") raise RuntimeError( @@ -119,7 +116,7 @@ def start_initial_website_source_import( account_id=account_id, control_space_id=control_space_id, payload=KnowledgeFSSourceCreatePayload( - connection_id=connection.id, + connectionId=connection.id, metadata={ "clientRequestId": request_id, "crawlOptions": { @@ -142,7 +139,7 @@ def start_initial_website_source_import( control_space_id=control_space_id, source_id=source.id, payload=KnowledgeFSCrawlImportPayload( - source_urls=[selection.source_url for selection in payload.selection], + sourceUrls=[selection.source_url for selection in payload.selection], ), idempotency_key=f"{request_id}:crawl-import", ) @@ -171,22 +168,22 @@ def start_initial_website_source_import( sync_policy = KnowledgeFSSourceSyncPolicyPayload( enabled=False, mode="manual", - expected_revision=expected_revision, - expected_source_version=imported_source.version, + expectedRevision=expected_revision, + expectedSourceVersion=imported_source.version, ) elif payload.sync_policy == "daily": sync_policy = KnowledgeFSSourceSyncPolicyPayload( enabled=True, mode="interval", - expected_revision=expected_revision, - expected_source_version=imported_source.version, + expectedRevision=expected_revision, + expectedSourceVersion=imported_source.version, ) else: sync_policy = KnowledgeFSSourceSyncPolicyPayload( enabled=True, mode="provider", - expected_revision=expected_revision, - expected_source_version=imported_source.version, + expectedRevision=expected_revision, + expectedSourceVersion=imported_source.version, ) facade.update_source_sync_policy( tenant_id=tenant_id, diff --git a/api/tests/unit_tests/extensions/test_knowledge_fs_celery.py b/api/tests/unit_tests/extensions/test_knowledge_fs_celery.py new file mode 100644 index 00000000000..1c951385cdd --- /dev/null +++ b/api/tests/unit_tests/extensions/test_knowledge_fs_celery.py @@ -0,0 +1,54 @@ +from types import SimpleNamespace +from unittest.mock import MagicMock, patch + +from dify_app import DifyApp +from extensions.ext_celery import init_app + + +def test_celery_registers_initial_source_task_when_knowledge_fs_lifecycle_is_ready() -> None: + config = MagicMock() + config.BROKER_USE_SSL = False + config.REDIS_KEY_PREFIX = "test" + config.HUMAN_INPUT_TIMEOUT_TASK_INTERVAL = 1 + config.CELERY_BROKER_URL = "redis://localhost:6379/0" + config.CELERY_BACKEND = "redis" + config.CELERY_RESULT_BACKEND = "redis://localhost:6379/0" + config.CELERY_USE_SENTINEL = False + config.LOG_FORMAT = "%(message)s" + config.LOG_TZ = "UTC" + config.LOG_FILE = None + config.CELERY_TASK_ANNOTATIONS = {} + config.CELERY_BEAT_SCHEDULER_TIME = 1 + config.KNOWLEDGE_FS_LIFECYCLE_POLL_INTERVAL_SECONDS = 2 + config.ENABLE_CLEAN_EMBEDDING_CACHE_TASK = False + config.ENABLE_CLEAN_UNUSED_DATASETS_TASK = False + config.ENABLE_CREATE_TIDB_SERVERLESS_TASK = False + config.ENABLE_UPDATE_TIDB_SERVERLESS_STATUS_TASK = False + config.ENABLE_CLEAN_MESSAGES = False + config.ENABLE_MAIL_CLEAN_DOCUMENT_NOTIFY_TASK = False + config.ENABLE_DATASETS_QUEUE_MONITOR = False + config.ENABLE_HUMAN_INPUT_TIMEOUT_TASK = False + config.ENABLE_CHECK_UPGRADABLE_PLUGIN_TASK = False + config.MARKETPLACE_ENABLED = False + config.WORKFLOW_LOG_CLEANUP_ENABLED = False + config.ENABLE_WORKFLOW_RUN_CLEANUP_TASK = False + config.ENABLE_WORKFLOW_SCHEDULE_POLLER_TASK = False + config.WORKFLOW_SCHEDULE_POLLER_INTERVAL = 1 + config.ENABLE_TRIGGER_PROVIDER_REFRESH_TASK = False + config.TRIGGER_PROVIDER_REFRESH_INTERVAL = 15 + config.ENABLE_API_TOKEN_LAST_USED_UPDATE_TASK = False + config.API_TOKEN_LAST_USED_UPDATE_INTERVAL = 30 + config.ENTERPRISE_ENABLED = False + config.ENTERPRISE_TELEMETRY_ENABLED = False + + with ( + patch("extensions.ext_celery.dify_config", config), + patch( + "services.knowledge_fs.lifecycle_readiness.get_configured_knowledge_fs_lifecycle_worker_readiness", + return_value=SimpleNamespace(ready=True), + ), + ): + celery_app = init_app(DifyApp(__name__)) + + assert "tasks.knowledge_fs_initial_source_tasks" in celery_app.conf["imports"] + assert "tasks.knowledge_fs_lifecycle_tasks" in celery_app.conf["imports"] diff --git a/api/tests/unit_tests/tasks/test_knowledge_fs_initial_source_tasks.py b/api/tests/unit_tests/tasks/test_knowledge_fs_initial_source_tasks.py index e869883aa39..fc0d6da2c89 100644 --- a/api/tests/unit_tests/tasks/test_knowledge_fs_initial_source_tasks.py +++ b/api/tests/unit_tests/tasks/test_knowledge_fs_initial_source_tasks.py @@ -1,16 +1,37 @@ +from contextlib import contextmanager from types import SimpleNamespace from unittest.mock import MagicMock, patch +import pytest + from models.knowledge_fs import KnowledgeFSControlSpaceState from services.knowledge_fs.product_dto import KnowledgeFSInitialWebsiteSourcePayload from services.knowledge_fs.product_remote import KnowledgeFSProductResourceNotFoundError -from tasks.knowledge_fs_initial_source_tasks import start_initial_website_source_import +from tasks.knowledge_fs_initial_source_tasks import ( + KnowledgeFSInitialSourceNotReadyError, + import_initial_website_source, + start_initial_website_source_import, +) -def test_initial_website_source_import_recrawls_exact_selection_and_configures_sync() -> None: - session_context = MagicMock() - session_context.__enter__.return_value = object() - session_maker = MagicMock(return_value=session_context) +def _payload(sync_policy: str = "daily") -> KnowledgeFSInitialWebsiteSourcePayload: + return KnowledgeFSInitialWebsiteSourcePayload.model_validate( + { + "kind": "website_crawl", + "name": "Dify docs", + "provider": "firecrawl", + "root_url": "https://docs.dify.ai", + "crawl_options": {"include_subpages": True, "limit": 25}, + "selection": [ + {"source_url": "https://docs.dify.ai/a", "title": "A"}, + {"source_url": "https://docs.dify.ai/b", "title": "B"}, + ], + "sync_policy": sync_policy, + } + ) + + +def _facade() -> MagicMock: facade = MagicMock() facade.list_sources.return_value = SimpleNamespace(data=[], next_cursor=None) facade.list_source_providers.return_value = SimpleNamespace( @@ -33,39 +54,38 @@ def test_initial_website_source_import_recrawls_exact_selection_and_configures_s ) facade.get_source.return_value = SimpleNamespace(version=3) facade.get_source_sync_policy.side_effect = KnowledgeFSProductResourceNotFoundError("not found") + return facade - payload = KnowledgeFSInitialWebsiteSourcePayload.model_validate( - { - "kind": "website_crawl", - "name": "Dify docs", - "provider": "firecrawl", - "root_url": "https://docs.dify.ai", - "crawl_options": {"include_subpages": True, "limit": 25}, - "selection": [ - {"source_url": "https://docs.dify.ai/a", "title": "A"}, - {"source_url": "https://docs.dify.ai/b", "title": "B"}, - ], - "sync_policy": "daily", - } - ) +@contextmanager +def _runtime( + facade: MagicMock, + *, + state: KnowledgeFSControlSpaceState = KnowledgeFSControlSpaceState.ACTIVE, + knowledge_space_id: str | None = "space-1", +): + session_context = MagicMock() + session_context.__enter__.return_value = object() + session_maker = MagicMock(return_value=session_context) with ( patch( "tasks.knowledge_fs_initial_source_tasks.session_factory.get_session_maker", return_value=session_maker, ), - patch( - "tasks.knowledge_fs_initial_source_tasks.SQLAlchemyKnowledgeFSControlSpaceRepository" - ) as repository_type, + patch("tasks.knowledge_fs_initial_source_tasks.SQLAlchemyKnowledgeFSControlSpaceRepository") as repository_type, patch("tasks.knowledge_fs_initial_source_tasks.get_knowledge_fs_runtime") as get_runtime, ): repository_type.return_value.get.return_value = SimpleNamespace( - state=KnowledgeFSControlSpaceState.ACTIVE, - knowledge_space_id="space-1", + state=state, + knowledge_space_id=knowledge_space_id, ) get_runtime.return_value.facade = facade + yield repository_type - workflow_id = start_initial_website_source_import( + +def _start(facade: MagicMock, payload: KnowledgeFSInitialWebsiteSourcePayload) -> str: + with _runtime(facade): + return start_initial_website_source_import( tenant_id="tenant-1", account_id="account-1", control_space_id="control-1", @@ -73,9 +93,15 @@ def test_initial_website_source_import_recrawls_exact_selection_and_configures_s payload=payload, ) - assert workflow_id == "workflow-1" + +def test_initial_website_source_import_recrawls_exact_selection_and_configures_daily_sync() -> None: + facade = _facade() + + assert _start(facade, _payload()) == "workflow-1" + source_payload = facade.create_source.call_args.kwargs["payload"] assert source_payload.status == "disabled" + assert source_payload.connection_id == "connection-1" assert source_payload.metadata["preview"] is True import_payload = facade.import_selected_source_crawl.call_args.kwargs["payload"] assert import_payload.source_urls == [ @@ -85,4 +111,187 @@ def test_initial_website_source_import_recrawls_exact_selection_and_configures_s sync_payload = facade.update_source_sync_policy.call_args.kwargs["payload"] assert sync_payload.mode == "interval" assert sync_payload.enabled is True + assert sync_payload.expected_revision == 0 assert sync_payload.expected_source_version == 3 + + +def test_initial_website_source_import_reuses_source_across_pages_and_preserves_failure() -> None: + facade = _facade() + existing_source = SimpleNamespace( + id="existing-source", + metadata={"clientRequestId": "initial-website-source:operation-1"}, + ) + facade.list_sources.side_effect = [ + SimpleNamespace( + data=[SimpleNamespace(id="other", metadata={})], + next_cursor="next-source", + ), + SimpleNamespace(data=[existing_source], next_cursor=None), + ] + facade.import_selected_source_crawl.return_value = SimpleNamespace( + id="failed-workflow", + state="failed", + ) + + assert _start(facade, _payload()) == "failed-workflow" + facade.create_source.assert_not_called() + facade.list_source_connections.assert_not_called() + facade.update_source_sync_policy.assert_not_called() + + +@pytest.mark.parametrize( + ("sync_policy", "expected_mode", "expected_enabled"), + [ + ("manual", "manual", False), + ("provider", "provider", True), + ], +) +def test_initial_website_source_import_configures_remaining_sync_modes( + sync_policy: str, + expected_mode: str, + expected_enabled: bool, +) -> None: + facade = _facade() + facade.list_source_connections.side_effect = [ + SimpleNamespace( + data=[SimpleNamespace(id="inactive", provider_id="other", status="disabled")], + next_cursor="next-connection", + ), + SimpleNamespace( + data=[ + SimpleNamespace( + id="connection-2", + provider_id="plugin-daemon-website", + status="active", + ) + ], + next_cursor=None, + ), + ] + facade.get_source_sync_policy.side_effect = None + facade.get_source_sync_policy.return_value = SimpleNamespace(revision=7) + + assert _start(facade, _payload(sync_policy)) == "workflow-1" + sync_payload = facade.update_source_sync_policy.call_args.kwargs["payload"] + assert sync_payload.mode == expected_mode + assert sync_payload.enabled is expected_enabled + assert sync_payload.expected_revision == 7 + + +@pytest.mark.parametrize( + ("state", "knowledge_space_id", "error_type", "message"), + [ + ( + KnowledgeFSControlSpaceState.PROVISIONING, + None, + KnowledgeFSInitialSourceNotReadyError, + "still provisioning", + ), + ( + KnowledgeFSControlSpaceState.ERROR, + None, + RuntimeError, + "cannot accept", + ), + ( + KnowledgeFSControlSpaceState.ACTIVE, + None, + RuntimeError, + "cannot accept", + ), + ], +) +def test_initial_website_source_import_rejects_unavailable_spaces( + state: KnowledgeFSControlSpaceState, + knowledge_space_id: str | None, + error_type: type[Exception], + message: str, +) -> None: + facade = _facade() + with ( + _runtime(facade, state=state, knowledge_space_id=knowledge_space_id), + pytest.raises(error_type, match=message), + ): + start_initial_website_source_import( + tenant_id="tenant-1", + account_id="account-1", + control_space_id="control-1", + operation_id="operation-1", + payload=_payload(), + ) + + +def test_initial_website_source_import_rejects_missing_control_space() -> None: + facade = _facade() + with _runtime(facade) as repository_type: + repository_type.return_value.get.return_value = None + with pytest.raises(RuntimeError, match="control-space was not found"): + start_initial_website_source_import( + tenant_id="tenant-1", + account_id="account-1", + control_space_id="control-1", + operation_id="operation-1", + payload=_payload(), + ) + + +def test_initial_website_source_import_rejects_unavailable_provider_and_connection() -> None: + unavailable_provider = _facade() + unavailable_provider.list_source_providers.return_value = SimpleNamespace(data=[]) + with pytest.raises(RuntimeError, match="provider is unavailable"): + _start(unavailable_provider, _payload()) + + unavailable_connection = _facade() + unavailable_connection.list_source_connections.return_value = SimpleNamespace( + data=[SimpleNamespace(id="inactive", provider_id="other", status="disabled")], + next_cursor=None, + ) + with pytest.raises(RuntimeError, match="connection is unavailable"): + _start(unavailable_connection, _payload()) + + +def test_initial_website_source_import_retries_running_workflow() -> None: + facade = _facade() + facade.import_selected_source_crawl.return_value = SimpleNamespace( + id="running-workflow", + state="running", + ) + + with pytest.raises(KnowledgeFSInitialSourceNotReadyError, match="still running"): + _start(facade, _payload()) + + +def test_initial_website_source_task_returns_result_and_retries_not_ready_error() -> None: + serialized_payload = _payload().model_dump(mode="json") + with patch( + "tasks.knowledge_fs_initial_source_tasks.start_initial_website_source_import", + return_value="workflow-1", + ): + assert ( + import_initial_website_source.run( + tenant_id="tenant-1", + account_id="account-1", + control_space_id="control-1", + operation_id="operation-1", + payload=serialized_payload, + ) + == "workflow-1" + ) + + retry_error = RuntimeError("retry requested") + with ( + patch( + "tasks.knowledge_fs_initial_source_tasks.start_initial_website_source_import", + side_effect=KnowledgeFSInitialSourceNotReadyError("not ready"), + ), + patch.object(import_initial_website_source, "retry", side_effect=retry_error) as retry, + pytest.raises(RuntimeError, match="retry requested"), + ): + import_initial_website_source.run( + tenant_id="tenant-1", + account_id="account-1", + control_space_id="control-1", + operation_id="operation-1", + payload=serialized_payload, + ) + retry.assert_called_once()