diff --git a/api/controllers/console/knowledge_fs/resources.py b/api/controllers/console/knowledge_fs/resources.py index 74c716ad54d..afb0c7ddd51 100644 --- a/api/controllers/console/knowledge_fs/resources.py +++ b/api/controllers/console/knowledge_fs/resources.py @@ -15,7 +15,14 @@ from urllib.parse import quote, urlencode from flask import Response, jsonify, request, send_file from flask_restx import Resource from pydantic import BaseModel, TypeAdapter, ValidationError -from werkzeug.exceptions import Conflict, NotFound, RequestEntityTooLarge, ServiceUnavailable, UnprocessableEntity +from werkzeug.exceptions import ( + BadRequest, + Conflict, + NotFound, + RequestEntityTooLarge, + ServiceUnavailable, + UnprocessableEntity, +) from configs import dify_config from controllers.common.fields import BinaryFileResponse @@ -68,8 +75,6 @@ from services.knowledge_fs.download_service import ( from services.knowledge_fs.initial_source_preview import KnowledgeFSInitialSourcePreviewService from services.knowledge_fs.initial_source_preview_job import ( KnowledgeFSInitialSourcePreviewJobAlreadyRunningError, - KnowledgeFSInitialSourcePreviewJobNotFoundError, - KnowledgeFSInitialSourcePreviewJobService, ) from services.knowledge_fs.object_storage import ( KnowledgeFSObjectStorageError, @@ -156,6 +161,7 @@ from services.knowledge_fs.product_dto import ( KnowledgeFSMetadataFieldListResponse, KnowledgeFSMetadataFieldResponse, KnowledgeFSMetadataFieldUpdatePayload, + KnowledgeFSNamespacePreviewCreatePayload, KnowledgeFSOverviewActivityListQuery, KnowledgeFSOverviewActivityListResponse, KnowledgeFSOverviewAttentionListQuery, @@ -240,6 +246,7 @@ from services.knowledge_fs.product_dto import ( KnowledgeFSUploadSessionCreatePayload, KnowledgeFSUploadSessionCreateResponse, KnowledgeFSUploadSessionMutationResponse, + knowledge_fs_initial_preview_configuration_fingerprint, normalize_knowledge_fs_source_url, ) from services.knowledge_fs.product_remote import ( @@ -864,12 +871,30 @@ class KnowledgeFSInitialSourcePreviewJobsApi(Resource): @_knowledge_fs_errors def post(self): account, tenant_id = current_account_with_tenant() - result = KnowledgeFSInitialSourcePreviewJobService(session_factory.get_session_maker()).start( - tenant_id=tenant_id, - account=account, - payload=_payload(KnowledgeFSInitialWebsiteSourcePreviewPayload), + payload = _payload(KnowledgeFSInitialWebsiteSourcePreviewPayload) + KnowledgeFSInitialSourcePreviewService(session_factory.get_session_maker()).require_visible_credential( + tenant_id=tenant_id, account=account, payload=payload ) - return dump_response(KnowledgeFSInitialSourcePreviewJobCreateResponse, result), HTTPStatus.ACCEPTED + root_url = payload.parameters.get("url") + if not isinstance(root_url, str) or not root_url.strip(): + raise BadRequest("Website preview URL is required") + result = _console_services().facade.create_namespace_source_preview( + tenant_id=tenant_id, + account_id=account.id, + payload=KnowledgeFSNamespacePreviewCreatePayload( + credentialId=payload.credential_id, + datasource=payload.datasource, + pluginId=payload.plugin_id, + provider=payload.provider, + parameters=payload.parameters, + rootUrl=root_url, + configurationFingerprint=knowledge_fs_initial_preview_configuration_fingerprint(payload), + ), + ) + return dump_response( + KnowledgeFSInitialSourcePreviewJobCreateResponse, + KnowledgeFSInitialSourcePreviewJobCreateResponse(jobId=result.job_id), + ), HTTPStatus.ACCEPTED @console_ns.route("/knowledge-fs/source-provider-preview/jobs/") @@ -886,12 +911,31 @@ class KnowledgeFSInitialSourcePreviewJobApi(Resource): def get(self, job_id: str): account, tenant_id = current_account_with_tenant() try: - result = KnowledgeFSInitialSourcePreviewJobService.get( - tenant_id=tenant_id, - account_id=account.id, - job_id=job_id, + remote = _console_services().facade.get_namespace_source_preview( + tenant_id=tenant_id, account_id=account.id, job_id=job_id ) - except KnowledgeFSInitialSourcePreviewJobNotFoundError as exc: + result = KnowledgeFSInitialSourcePreviewJobResponse( + jobId=remote.job_id, + status={"queued": "pending", "consumed": "completed"}.get(remote.status, remote.status), + result=( + KnowledgeFSInitialSourcePreviewResponse( + kind="website_crawl", + configurationFingerprint=remote.configuration_fingerprint, + pages=[ + { + "pageId": page.page_id, + "sourceUrl": page.source_url, + "title": page.title, + "description": page.description, + } + for page in remote.pages + ], + ) + if remote.status in {"completed", "consumed"} + else None + ), + ) + except KnowledgeFSProductResourceNotFoundError as exc: raise NotFound() from exc return dump_response(KnowledgeFSInitialSourcePreviewJobResponse, result) @@ -907,12 +951,14 @@ class KnowledgeFSInitialSourcePreviewJobApi(Resource): def delete(self, job_id: str): account, tenant_id = current_account_with_tenant() try: - result = KnowledgeFSInitialSourcePreviewJobService.cancel( - tenant_id=tenant_id, - account_id=account.id, - job_id=job_id, + remote = _console_services().facade.cancel_namespace_source_preview( + tenant_id=tenant_id, account_id=account.id, job_id=job_id ) - except KnowledgeFSInitialSourcePreviewJobNotFoundError as exc: + result = KnowledgeFSInitialSourcePreviewJobResponse( + jobId=remote.job_id, + status={"queued": "pending", "consumed": "completed"}.get(remote.status, remote.status), + ) + except KnowledgeFSProductResourceNotFoundError as exc: raise NotFound() from exc return dump_response(KnowledgeFSInitialSourcePreviewJobResponse, result) diff --git a/api/knowledge-fs-contract.lock.json b/api/knowledge-fs-contract.lock.json index 0d4f1dea0d9..37530312b77 100644 --- a/api/knowledge-fs-contract.lock.json +++ b/api/knowledge-fs-contract.lock.json @@ -1,9 +1,9 @@ { "schemaVersion": 5, - "subtreeTree": "1d3eabf30c7ef2aef0383595be622670c2ec56a8", - "openapiSha256": "5f9bee6a3593d8bd05308c7a57ab3d02ef451133cde9347a40543fae934a9a59", + "subtreeTree": "4b5ccdecb3012f6c9e628b03844d6c9a968a8ccb", + "openapiSha256": "8476e3c8a0143f250c9cc99c8be25bf2f7fa8f19214284643f17c2c6542b3f3e", "capabilityV2AuthManifestSha256": "e322a2fa779d1f40b95c54c1021cffecaec77abbd7b34899573dcdf4ff353109", "capabilityV2AuthTestVectorSha256": "ae0de37b1ff05c40f905cf17a7b410d8971acacf64db07d5ee3d6fecfa559ce3", - "productOperationManifestSha256": "5d1241a83bcca12ebbd848928dd3cdda0d2ecaeb5f5336e955eb24a8c8db175b", + "productOperationManifestSha256": "e4611c61c69478f20c1e1e482c44d8c48e48adfc67f1e8941ce36225e10b32c0", "productOperationGapManifestSha256": "332c80165bd5cf8e79bc511374dde771a8175510404db603a8067f1a22d36df8" } diff --git a/api/knowledge-fs-product-operations.json b/api/knowledge-fs-product-operations.json index 970d1476d66..7ccb46ac3a7 100644 --- a/api/knowledge-fs-product-operations.json +++ b/api/knowledge-fs-product-operations.json @@ -63,6 +63,10 @@ {"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":"createNamespaceSourcePreview","kfsOperationId":"createNamespaceSourcePreview","method":"POST","path":"/namespace/source-preview-jobs","action":"namespace_source_previews.create","resource":"namespace","transport":"json","stream":{"productKind":"json","kfsResponseKind":"buffered"},"limits":{"productMaxRequestBytes":65536,"productMaxResponseBytes":524288,"kfsMaxResponseBytes":1048576}}, + {"productOperationId":"getNamespaceSourcePreview","kfsOperationId":"getNamespaceSourcePreview","method":"GET","path":"/namespace/source-preview-jobs/{jobId}","action":"namespace_source_previews.read","resource":"namespace","transport":"json","stream":{"productKind":"json","kfsResponseKind":"buffered"},"limits":{"productMaxRequestBytes":16384,"productMaxResponseBytes":524288,"kfsMaxResponseBytes":1048576}}, + {"productOperationId":"cancelNamespaceSourcePreview","kfsOperationId":"cancelNamespaceSourcePreview","method":"DELETE","path":"/namespace/source-preview-jobs/{jobId}","action":"namespace_source_previews.cancel","resource":"namespace","transport":"json","stream":{"productKind":"json","kfsResponseKind":"buffered"},"limits":{"productMaxRequestBytes":16384,"productMaxResponseBytes":524288,"kfsMaxResponseBytes":1048576}}, + {"productOperationId":"consumeNamespaceSourcePreview","kfsOperationId":"consumeNamespaceSourcePreview","method":"POST","path":"/knowledge-spaces/{id}/sources/{sourceId}/namespace-preview-import","action":"namespace_source_previews.consume","resource":"source","transport":"json","stream":{"productKind":"json","kfsResponseKind":"buffered"},"limits":{"productMaxRequestBytes":65536,"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}}, diff --git a/api/services/knowledge_fs/capability_broker.py b/api/services/knowledge_fs/capability_broker.py index 48062cec2f7..0750afbed07 100644 --- a/api/services/knowledge_fs/capability_broker.py +++ b/api/services/knowledge_fs/capability_broker.py @@ -132,6 +132,49 @@ class KnowledgeFSCapabilityBroker: knowledge_space_revision=knowledge_space_revision, ) + def issue_namespace_interactive( + self, *, tenant_id: str, account_id: str, operation_id: str, trace_id: str | None = None + ) -> KnowledgeFSIssuedProductCapability: + """Issue a namespace-only capability before a control Space exists. + + Console authentication and datasource credential ownership are established by the caller; + the signed resource is restricted to this tenant and operation. Namespace previews never + authorize a KnowledgeSpace child resource. + """ + self._cutover_gate.require_capability_v2(tenant_id=tenant_id) + _, capability_operation_id = _operation_contract(operation_id) + issuer = self._require_issuer() + normalized_trace_id = _trace_id(trace_id) + operation = KNOWLEDGE_FS_CAPABILITY_OPERATIONS[capability_operation_id] + if operation.resource_type != "namespace": + raise KnowledgeFSOperationUnavailableError("KnowledgeFS operation is not namespace-scoped") + request = CapabilityIssueRequest( + actor=f"dify-account:{account_id}", + authz_revision=CapabilityAuthzRevision( + membership_epoch=0, space_acl_epoch=0, external_access_epoch=0, credential_revision=None + ), + caller_kind="interactive", + content_policy_revision=0, + control_space_id=f"namespace:{tenant_id}", + grant_id=str( + uuid.uuid5(uuid.NAMESPACE_URL, f"dify-kfs-namespace:{tenant_id}:{account_id}:{normalized_trace_id}") + ), + namespace_id=tenant_id, + operation_id=capability_operation_id, + principal_id=account_id, + resource=CapabilityResource(type="namespace", id=tenant_id), + trace_id=normalized_trace_id, + ) + issued = issuer.issue(request) + return KnowledgeFSIssuedProductCapability( + token=issued.token, + expires_at=datetime.fromtimestamp(issued.claims.exp, tz=UTC), + operation_id=operation_id, + knowledge_space_id=tenant_id, + knowledge_space_revision=0, + trace_id=normalized_trace_id, + ) + def issue_service( self, *, diff --git a/api/services/knowledge_fs/data_facade.py b/api/services/knowledge_fs/data_facade.py index 6206d5c23c8..bbcbba00f84 100644 --- a/api/services/knowledge_fs/data_facade.py +++ b/api/services/knowledge_fs/data_facade.py @@ -81,6 +81,10 @@ from services.knowledge_fs.product_dto import ( KnowledgeFSMetadataFieldListResponse, KnowledgeFSMetadataFieldResponse, KnowledgeFSMetadataFieldUpdatePayload, + KnowledgeFSNamespacePreviewConsumePayload, + KnowledgeFSNamespacePreviewCreatePayload, + KnowledgeFSNamespacePreviewImportResponse, + KnowledgeFSNamespacePreviewJobResponse, KnowledgeFSOverviewActivityListResponse, KnowledgeFSOverviewAttentionListResponse, KnowledgeFSOverviewBaseStatsResponse, @@ -1655,6 +1659,65 @@ class KnowledgeFSDataFacade: ) return KnowledgeFSSourceWorkflowResponse.model_validate(raw) + def create_namespace_source_preview( + self, *, tenant_id: str, account_id: str, payload: KnowledgeFSNamespacePreviewCreatePayload + ) -> KnowledgeFSNamespacePreviewJobResponse: + return KnowledgeFSNamespacePreviewJobResponse.model_validate( + self._namespace_interactive( + tenant_id=tenant_id, + account_id=account_id, + operation_id="createNamespaceSourcePreview", + payload=payload, + ) + ) + + def get_namespace_source_preview( + self, *, tenant_id: str, account_id: str, job_id: str + ) -> KnowledgeFSNamespacePreviewJobResponse: + return KnowledgeFSNamespacePreviewJobResponse.model_validate( + self._namespace_interactive( + tenant_id=tenant_id, + account_id=account_id, + operation_id="getNamespaceSourcePreview", + path_parameters=(("jobId", job_id),), + ) + ) + + def cancel_namespace_source_preview( + self, *, tenant_id: str, account_id: str, job_id: str + ) -> KnowledgeFSNamespacePreviewJobResponse: + return KnowledgeFSNamespacePreviewJobResponse.model_validate( + self._namespace_interactive( + tenant_id=tenant_id, + account_id=account_id, + operation_id="cancelNamespaceSourcePreview", + path_parameters=(("jobId", job_id),), + ) + ) + + def consume_namespace_source_preview( + self, + *, + tenant_id: str, + account_id: str, + control_space_id: str, + source_id: str, + payload: KnowledgeFSNamespacePreviewConsumePayload, + idempotency_key: str, + ) -> KnowledgeFSNamespacePreviewImportResponse: + return KnowledgeFSNamespacePreviewImportResponse.model_validate( + self._interactive_child( + tenant_id=tenant_id, + account_id=account_id, + control_space_id=control_space_id, + operation_id="consumeNamespaceSourcePreview", + resource_id=source_id, + path_parameters=(("sourceId", source_id),), + payload=payload, + headers=(("Idempotency-Key", idempotency_key),), + ) + ) + def import_source_workflow( self, *, @@ -2475,6 +2538,30 @@ class KnowledgeFSDataFacade: headers=headers, ) + def _namespace_interactive( + self, + *, + tenant_id: str, + account_id: str, + operation_id: str, + payload: BaseModel | None = None, + path_parameters: tuple[tuple[str, str], ...] = (), + ) -> JsonValue: + _assert_json_bff_ready(operation_id) + issued = self._broker.issue_namespace_interactive( + tenant_id=tenant_id, account_id=account_id, operation_id=operation_id + ) + return self._execute( + operation_id=operation_id, + namespace_id=tenant_id, + knowledge_space_id=tenant_id, + knowledge_space_revision=0, + capability_token=issued.token, + trace_id=issued.trace_id, + payload=payload, + path_parameters=path_parameters, + ) + def _interactive_child( self, *, diff --git a/api/services/knowledge_fs/product_dto.py b/api/services/knowledge_fs/product_dto.py index 4369c3e5051..2b83d53d618 100644 --- a/api/services/knowledge_fs/product_dto.py +++ b/api/services/knowledge_fs/product_dto.py @@ -147,11 +147,12 @@ def normalize_knowledge_fs_source_url(source_url: str) -> str: class KnowledgeFSInitialWebsiteSelectionPayload(BaseModel): + page_id: str | None = Field(default=None, min_length=1, max_length=128, alias="pageId") source_url: str = Field(min_length=1, max_length=4_096) canonical_url: str | None = Field(default=None, min_length=1, max_length=4_096) title: str | None = Field(default=None, max_length=500) - model_config = ConfigDict(extra="forbid") + model_config = ConfigDict(extra="forbid", validate_by_alias=True, validate_by_name=True) @field_validator("source_url") @classmethod @@ -287,6 +288,7 @@ def knowledge_fs_initial_preview_configuration_fingerprint( class KnowledgeFSInitialSourcePreviewPageResponse(ResponseModel): content: str | None = Field(default=None, exclude=True) description: str | None = None + page_id: str | None = Field(default=None, validation_alias=AliasChoices("page_id", "pageId")) source_url: str = Field(validation_alias=AliasChoices("source_url", "sourceUrl")) title: str | None = None @@ -339,6 +341,49 @@ class KnowledgeFSInitialSourcePreviewJobResponse(ResponseModel): status: Literal["pending", "running", "completed", "failed", "canceled"] +class KnowledgeFSNamespacePreviewCreatePayload(BaseModel): + credential_id: str = Field(min_length=1, max_length=255, alias="credentialId") + datasource: str = Field(min_length=1, max_length=255) + plugin_id: str = Field(min_length=1, max_length=1024, alias="pluginId") + provider: str = Field(min_length=1, max_length=255) + parameters: dict[str, JsonValue] = Field(default_factory=dict) + root_url: str = Field(min_length=1, max_length=4096, alias="rootUrl") + configuration_fingerprint: str = Field(min_length=64, max_length=64, alias="configurationFingerprint") + model_config = ConfigDict(extra="forbid", validate_by_alias=True, validate_by_name=True) + + +class KnowledgeFSNamespacePreviewConsumePayload(BaseModel): + preview_job_id: str = Field(min_length=1, max_length=255, alias="previewJobId") + page_ids: list[str] = Field(min_length=1, max_length=200, alias="pageIds") + configuration_fingerprint: str = Field(min_length=64, max_length=64, alias="configurationFingerprint") + model_config = ConfigDict(extra="forbid", validate_by_alias=True, validate_by_name=True) + + +class KnowledgeFSNamespacePreviewPageResponse(ResponseModel): + page_id: str = Field(validation_alias=AliasChoices("page_id", "pageId")) + source_url: str = Field(validation_alias=AliasChoices("source_url", "sourceUrl")) + title: str | None = None + description: str | None = None + + +class KnowledgeFSNamespacePreviewJobResponse(ResponseModel): + job_id: str = Field(validation_alias=AliasChoices("job_id", "jobId")) + status: Literal["queued", "running", "completed", "failed", "canceled", "consumed"] + configuration_fingerprint: str = Field( + validation_alias=AliasChoices("configuration_fingerprint", "configurationFingerprint") + ) + expires_at: str = Field(validation_alias=AliasChoices("expires_at", "expiresAt")) + error_code: str | None = Field(default=None, validation_alias=AliasChoices("error_code", "errorCode")) + import_workflow_id: str | None = Field( + default=None, validation_alias=AliasChoices("import_workflow_id", "importWorkflowId") + ) + pages: list[KnowledgeFSNamespacePreviewPageResponse] = Field(default_factory=list) + + +class KnowledgeFSNamespacePreviewImportResponse(ResponseModel): + workflow_id: str = Field(validation_alias=AliasChoices("workflow_id", "workflowId")) + + KnowledgeFSInitialSourcePayload = Annotated[ KnowledgeFSInitialWebsiteSourcePayload | KnowledgeFSInitialOnlineDocumentSourcePayload diff --git a/api/services/knowledge_fs/product_operations.py b/api/services/knowledge_fs/product_operations.py index 15bcfe1be2f..f468626bdbb 100644 --- a/api/services/knowledge_fs/product_operations.py +++ b/api/services/knowledge_fs/product_operations.py @@ -807,6 +807,50 @@ KNOWLEDGE_FS_PRODUCT_OPERATIONS: Final[MappingProxyType[str, KnowledgeFSProductO max_response_bytes=512 * 1024, stream_kind="json", ), + "createNamespaceSourcePreview": _operation( + "POST", + "createNamespaceSourcePreview", + KnowledgeFSProductPermission.DOCUMENT_WRITE, + "/namespace/source-preview-jobs", + "json", + resource_resolver="namespace", + max_request_bytes=64 * 1024, + max_response_bytes=512 * 1024, + stream_kind="json", + ), + "getNamespaceSourcePreview": _operation( + "GET", + "getNamespaceSourcePreview", + KnowledgeFSProductPermission.READ, + "/namespace/source-preview-jobs/{jobId}", + "json", + resource_resolver="namespace", + max_request_bytes=16 * 1024, + max_response_bytes=512 * 1024, + stream_kind="json", + ), + "cancelNamespaceSourcePreview": _operation( + "DELETE", + "cancelNamespaceSourcePreview", + KnowledgeFSProductPermission.DOCUMENT_WRITE, + "/namespace/source-preview-jobs/{jobId}", + "json", + resource_resolver="namespace", + max_request_bytes=16 * 1024, + max_response_bytes=512 * 1024, + stream_kind="json", + ), + "consumeNamespaceSourcePreview": _operation( + "POST", + "consumeNamespaceSourcePreview", + KnowledgeFSProductPermission.DOCUMENT_WRITE, + "/knowledge-spaces/{id}/sources/{sourceId}/namespace-preview-import", + "json", + resource_resolver="source", + max_request_bytes=64 * 1024, + max_response_bytes=512 * 1024, + stream_kind="json", + ), "importSelectedSourceCrawl": _operation( "POST", "createSourceCrawlImportWorkflow", diff --git a/api/services/knowledge_fs_capability.py b/api/services/knowledge_fs_capability.py index f0da46b7f84..39e71080d35 100644 --- a/api/services/knowledge_fs_capability.py +++ b/api/services/knowledge_fs_capability.py @@ -737,6 +737,34 @@ KNOWLEDGE_FS_CAPABILITY_OPERATIONS: Final[Mapping[str, KnowledgeFSCapabilityOper "/knowledge-spaces/{id}/sources/{sourceId}/crawl-preview", "source", ), + "createNamespaceSourcePreview": KnowledgeFSCapabilityOperation( + "namespace_source_previews.create", + _CONTROL_PLANE_CALLERS, + "POST", + "/namespace/source-preview-jobs", + "namespace", + ), + "getNamespaceSourcePreview": KnowledgeFSCapabilityOperation( + "namespace_source_previews.read", + _CONTROL_PLANE_CALLERS, + "GET", + "/namespace/source-preview-jobs/{jobId}", + "namespace", + ), + "cancelNamespaceSourcePreview": KnowledgeFSCapabilityOperation( + "namespace_source_previews.cancel", + _CONTROL_PLANE_CALLERS, + "DELETE", + "/namespace/source-preview-jobs/{jobId}", + "namespace", + ), + "consumeNamespaceSourcePreview": KnowledgeFSCapabilityOperation( + "namespace_source_previews.consume", + _CONTROL_PLANE_CALLERS, + "POST", + "/knowledge-spaces/{id}/sources/{sourceId}/namespace-preview-import", + "source", + ), "createSourceImportWorkflow": KnowledgeFSCapabilityOperation( "source_workflows.import.create", _STANDARD_CALLERS, diff --git a/api/tasks/knowledge_fs_initial_source_tasks.py b/api/tasks/knowledge_fs_initial_source_tasks.py index 74037e5197b..71873ed1e85 100644 --- a/api/tasks/knowledge_fs_initial_source_tasks.py +++ b/api/tasks/knowledge_fs_initial_source_tasks.py @@ -19,12 +19,12 @@ from repositories.sqlalchemy_knowledge_fs_control_space_repository import ( SQLAlchemyKnowledgeFSControlSpaceRepository, ) from services.credential_permission_service import CredentialPermissionService -from services.knowledge_fs.initial_source_preview_job import KnowledgeFSInitialSourcePreviewJobService from services.knowledge_fs.product_dto import ( KnowledgeFSCrawlImportPayload, KnowledgeFSInitialOnlineDocumentSourcePayload, KnowledgeFSInitialSourcePayload, KnowledgeFSInitialWebsiteSourcePayload, + KnowledgeFSNamespacePreviewConsumePayload, KnowledgeFSOnlineDocumentWorkflowImportPayload, KnowledgeFSOnlineDriveWorkflowImportPayload, KnowledgeFSSourceConnectionCreatePayload, @@ -281,14 +281,28 @@ def _start_workflow( payload: KnowledgeFSInitialSourcePayload, ): if isinstance(payload, KnowledgeFSInitialWebsiteSourcePayload): - pages = None if payload.preview_job_id: - pages = KnowledgeFSInitialSourcePreviewJobService.selected_content( + if not payload.preview_configuration_fingerprint or any( + not selection.page_id for selection in payload.selection + ): + raise ValueError("Website namespace preview selection requires page IDs and fingerprint") + imported = facade.consume_namespace_source_preview( tenant_id=tenant_id, account_id=account_id, - job_id=payload.preview_job_id, - source_urls=[selection.source_url for selection in payload.selection], - configuration_fingerprint=knowledge_fs_initial_preview_configuration_fingerprint(payload), + control_space_id=control_space_id, + source_id=source_id, + payload=KnowledgeFSNamespacePreviewConsumePayload( + previewJobId=payload.preview_job_id, + pageIds=[selection.page_id for selection in payload.selection], + configurationFingerprint=payload.preview_configuration_fingerprint, + ), + idempotency_key=f"{request_id}:crawl-import", + ) + return facade.get_source_workflow( + tenant_id=tenant_id, + account_id=account_id, + control_space_id=control_space_id, + run_id=imported.workflow_id, ) workflow = facade.import_selected_source_crawl( tenant_id=tenant_id, @@ -297,14 +311,9 @@ def _start_workflow( source_id=source_id, payload=KnowledgeFSCrawlImportPayload( sourceUrls=[selection.source_url for selection in payload.selection], - pages=pages, ), idempotency_key=f"{request_id}:crawl-import", ) - if payload.preview_job_id: - KnowledgeFSInitialSourcePreviewJobService.cleanup_content( - tenant_id=tenant_id, account_id=account_id, job_id=payload.preview_job_id - ) return workflow if isinstance(payload, KnowledgeFSInitialOnlineDocumentSourcePayload): import_payload = KnowledgeFSSourceWorkflowImportPayload( diff --git a/knowledge-fs/apps/api/src/index.ts b/knowledge-fs/apps/api/src/index.ts index 629c5045436..05e910953fd 100644 --- a/knowledge-fs/apps/api/src/index.ts +++ b/knowledge-fs/apps/api/src/index.ts @@ -701,6 +701,9 @@ const sourceProduct = logicalRevisions: sourceLogicalRevisions, providers: sourceProviderCatalog, repository: databaseRepositories.sourceProductWorkflows, + ...(databaseRepositories.namespaceSourcePreviews + ? { namespacePreviews: databaseRepositories.namespaceSourcePreviews } + : {}), workerId: `source-product-workflow:${randomUUID()}`, } : undefined; diff --git a/knowledge-fs/apps/api/src/repository-options.ts b/knowledge-fs/apps/api/src/repository-options.ts index 4d7c86e059a..7a040a4fae1 100644 --- a/knowledge-fs/apps/api/src/repository-options.ts +++ b/knowledge-fs/apps/api/src/repository-options.ts @@ -20,6 +20,7 @@ import { type KnowledgeSpaceProvisioningRepository, type KnowledgeSpaceUnpublishedProfileActivationRepository, type LegacySpacePublicationBootstrapRepository, + type NamespaceSourcePreviewRepository, type PageIndexFindabilityRepository, type PageIndexUpgradeBackfillRepository, type QualityControlRepository, @@ -73,6 +74,7 @@ import { createDatabaseKnowledgeSpaceUnpublishedProfileActivationRepository, createDatabaseLegacySpacePublicationBootstrapRepository, createDatabaseLogicalDocumentRepository, + createDatabaseNamespaceSourcePreviewRepository, createDatabasePageIndexFindabilityRepository, createDatabasePageIndexUpgradeBackfillRepository, createDatabaseParseArtifactRepository, @@ -152,6 +154,7 @@ export interface ApiDatabaseRepositoryBundle { readonly sourceCredentialBackfills?: SourceCredentialBackfillRepository | undefined; readonly sourceConnections?: SourceConnectionRepository | undefined; readonly sourceProductWorkflows?: SourceProductWorkflowRepository | undefined; + readonly namespaceSourcePreviews?: NamespaceSourcePreviewRepository | undefined; readonly sourceRetiredSecretCleanups?: SourceRetiredSecretCleanupRepository | undefined; readonly tidbFtsPostingBackfills?: TidbFtsPostingBackfillRepository | undefined; readonly uploadSessions?: UploadSessionRepository | undefined; @@ -277,6 +280,7 @@ export function createApiDatabaseRepositories({ maxClaimBatchSize: maxListLimit, maxListLimit: 200, }); + const namespaceSourcePreviews = createDatabaseNamespaceSourcePreviewRepository({ database }); const knowledgeSpaceAccess = createKnowledgeSpaceAccessService({ repository: createDatabaseKnowledgeSpaceAccessRepository({ database, @@ -467,6 +471,7 @@ export function createApiDatabaseRepositories({ sourceCredentialBackfills, sourceConnections, sourceProductWorkflows, + namespaceSourcePreviews, sourceRetiredSecretCleanups, ...(tidbFtsPostingBackfills ? { tidbFtsPostingBackfills } : {}), uploadSessions, diff --git a/knowledge-fs/packages/api/src/dify-capability-v2.ts b/knowledge-fs/packages/api/src/dify-capability-v2.ts index e47863f3d59..45f6a71d6df 100644 --- a/knowledge-fs/packages/api/src/dify-capability-v2.ts +++ b/knowledge-fs/packages/api/src/dify-capability-v2.ts @@ -1022,6 +1022,43 @@ export const DIFY_CAPABILITY_V2_OPERATIONS: readonly DifyCapabilityV2Operation[] resource: { pathParameter: "id" }, resourceType: "knowledge_space", }, + { + action: "namespace_source_previews.create", + allowedCallerKinds: CONTROL_PLANE_CALLERS, + method: "POST", + operationId: "createNamespaceSourcePreview", + pathTemplate: "/namespace/source-preview-jobs", + resource: { namespace: true }, + resourceType: "namespace", + }, + { + action: "namespace_source_previews.read", + allowedCallerKinds: CONTROL_PLANE_CALLERS, + method: "GET", + operationId: "getNamespaceSourcePreview", + pathTemplate: "/namespace/source-preview-jobs/{jobId}", + resource: { namespace: true }, + resourceType: "namespace", + }, + { + action: "namespace_source_previews.cancel", + allowedCallerKinds: CONTROL_PLANE_CALLERS, + method: "DELETE", + operationId: "cancelNamespaceSourcePreview", + pathTemplate: "/namespace/source-preview-jobs/{jobId}", + resource: { namespace: true }, + resourceType: "namespace", + }, + { + action: "namespace_source_previews.consume", + allowedCallerKinds: CONTROL_PLANE_CALLERS, + method: "POST", + operationId: "consumeNamespaceSourcePreview", + parentResource: { pathParameter: "id" }, + pathTemplate: "/knowledge-spaces/{id}/sources/{sourceId}/namespace-preview-import", + resource: { pathParameter: "sourceId" }, + resourceType: "source", + }, { action: "source_workflows.preview.create", allowedCallerKinds: STANDARD_CALLERS, diff --git a/knowledge-fs/packages/api/src/gateway-options.ts b/knowledge-fs/packages/api/src/gateway-options.ts index 4eb00ef6509..56612088cc6 100644 --- a/knowledge-fs/packages/api/src/gateway-options.ts +++ b/knowledge-fs/packages/api/src/gateway-options.ts @@ -82,6 +82,7 @@ import type { ModelCapabilityPreflight, } from "./model-capability-preflight"; import type { ModelInputModalityResolver } from "./model-input-modality-resolver"; +import type { NamespaceSourcePreviewRepository } from "./namespace-source-preview"; import type { OnlineDocumentConnector } from "./online-document-connector"; import type { OnlineDriveConnector } from "./online-drive-connector"; import type { PageIndexUpgradeBackfillRepository } from "./page-index-upgrade-backfill"; @@ -348,6 +349,7 @@ export interface KnowledgeGatewayOptions { readonly onWorkflowRuntime?: (runtime: SourceProductWorkflowRuntime) => void; readonly providers: SourceProviderCatalog; readonly repository: SourceProductWorkflowRepository; + readonly namespacePreviews?: NamespaceSourcePreviewRepository | undefined; readonly workerId: string; }; /** diff --git a/knowledge-fs/packages/api/src/index.ts b/knowledge-fs/packages/api/src/index.ts index 3ffb5bffa62..220b9ea3716 100644 --- a/knowledge-fs/packages/api/src/index.ts +++ b/knowledge-fs/packages/api/src/index.ts @@ -59,6 +59,9 @@ import { createInMemoryGoldenQuestionRepository, } from "./golden-question-repository"; export * from "./conflict-detection"; +export * from "./namespace-source-preview"; +export * from "./namespace-source-preview-handlers"; +export * from "./namespace-source-preview-routes"; export * from "./contextual-enrichment-flow"; export * from "./core-resource-response-schemas"; export * from "./cursor-utils"; @@ -581,6 +584,11 @@ import { registerLegacySpacePublicationBootstrapHandlers } from "./legacy-space- import { createLocalNodeQueryGenerator } from "./local-node-query-generator"; import { registerLogicalDocumentHandlers } from "./logical-document-handlers"; import { registerModelCapabilityHandlers } from "./model-capability-handlers"; +import { + createInMemoryNamespaceSourcePreviewRepository, + createNamespaceSourcePreviewService, +} from "./namespace-source-preview"; +import { registerNamespaceSourcePreviewHandlers } from "./namespace-source-preview-handlers"; import { registerOperationPolicyHandlers } from "./operation-policy-handlers"; import { registerPageIndexUpgradeBackfillHandlers } from "./page-index-upgrade-backfill-handlers"; import { @@ -2207,6 +2215,19 @@ export function createKnowledgeGateway({ repository: sourceProduct.repository, workflows: sourceProductWorkflows, }); + if (websiteCrawlConnector) { + const namespacePreviews = createNamespaceSourcePreviewService({ + repository: + sourceProduct.namespacePreviews ?? createInMemoryNamespaceSourcePreviewRepository(), + sources: sourceRepository, + storage: adapter.objectStorage, + websiteCrawl: websiteCrawlConnector, + workflows: sourceProductWorkflows, + }); + registerNamespaceSourcePreviewHandlers({ app, service: namespacePreviews }); + const timer = setInterval(() => void namespacePreviews.tick(), 1_000); + timer.unref?.(); + } } registerBackgroundTaskHandlers({ access: accessService, diff --git a/knowledge-fs/packages/api/src/namespace-source-preview-handlers.ts b/knowledge-fs/packages/api/src/namespace-source-preview-handlers.ts new file mode 100644 index 00000000000..1ec9ad53827 --- /dev/null +++ b/knowledge-fs/packages/api/src/namespace-source-preview-handlers.ts @@ -0,0 +1,125 @@ +import type { OpenAPIHono } from "@hono/zod-openapi"; +import type { KnowledgeGatewayEnv } from "./gateway-openapi-contracts"; +import type { + NamespaceSourcePreviewJob, + NamespaceSourcePreviewPage, + createNamespaceSourcePreviewService, +} from "./namespace-source-preview"; +import { NamespaceSourcePreviewError } from "./namespace-source-preview"; +import { + cancelNamespaceSourcePreviewRoute, + consumeNamespaceSourcePreviewRoute, + createNamespaceSourcePreviewRoute, + getNamespaceSourcePreviewRoute, +} from "./namespace-source-preview-routes"; +import type { LooseOpenApiContext } from "./openapi-handler-utils"; + +type Service = ReturnType; +export function registerNamespaceSourcePreviewHandlers(input: { + app: OpenAPIHono; + service: Service; +}): void { + const register = input.app.openapi.bind(input.app) as ( + // biome-ignore lint/suspicious/noExplicitAny: bounded OpenAPI registration adapter + route: any, + // biome-ignore lint/suspicious/noExplicitAny: bounded OpenAPI registration adapter + handler: (context: any) => unknown, + ) => void; + const present = ( + job: NamespaceSourcePreviewJob, + pages: readonly NamespaceSourcePreviewPage[] = [], + ) => ({ + jobId: job.id, + status: job.status, + configurationFingerprint: job.configurationFingerprint, + expiresAt: job.expiresAt, + ...(job.errorCode ? { errorCode: job.errorCode } : {}), + ...(job.importWorkflowId ? { importWorkflowId: job.importWorkflowId } : {}), + pages: pages.map((p) => ({ + pageId: p.pageId, + sourceUrl: p.sourceUrl, + ...(p.title ? { title: p.title } : {}), + ...(p.description ? { description: p.description } : {}), + })), + }); + register(createNamespaceSourcePreviewRoute, async (context) => { + try { + const body = context.req.valid("json"); + return context.json( + present( + await input.service.create( + context.get("subject"), + { + credentialId: body.credentialId, + pluginId: body.pluginId, + provider: body.provider, + datasource: body.datasource, + parameters: body.parameters, + rootUrl: body.rootUrl, + }, + body.configurationFingerprint, + ), + ), + 202, + ); + } catch (e) { + return failure(context, e); + } + }); + register(getNamespaceSourcePreviewRoute, async (context) => { + try { + const job = await input.service.get(context.get("subject"), context.req.valid("param").jobId); + return context.json( + present(job, await input.service.pages(context.get("subject"), job.id)), + 200, + ); + } catch (e) { + return failure(context, e); + } + }); + register(cancelNamespaceSourcePreviewRoute, async (context) => { + try { + return context.json( + present( + await input.service.cancel(context.get("subject"), context.req.valid("param").jobId), + ), + 200, + ); + } catch (e) { + return failure(context, e); + } + }); + register(consumeNamespaceSourcePreviewRoute, async (context) => { + try { + const params = context.req.valid("param"); + const body = context.req.valid("json"); + const headers = context.req.valid("header"); + return context.json( + { + workflowId: await input.service.consume(context.get("subject"), { + jobId: body.previewJobId, + pageIds: body.pageIds, + configurationFingerprint: body.configurationFingerprint, + knowledgeSpaceId: params.id, + sourceId: params.sourceId, + idempotencyKey: headers["Idempotency-Key"], + }), + }, + 202, + ); + } catch (e) { + return failure(context, e); + } + }); +} +function failure(context: LooseOpenApiContext, error: unknown) { + if (error instanceof NamespaceSourcePreviewError) { + const status = error.code.includes("NOT_FOUND") + ? 404 + : error.code.includes("EXPIRED") || error.code.includes("MISMATCH") + ? 409 + : 400; + return context.json({ code: error.code, error: error.message }, status as 400); + } + throw error; +} diff --git a/knowledge-fs/packages/api/src/namespace-source-preview-routes.ts b/knowledge-fs/packages/api/src/namespace-source-preview-routes.ts new file mode 100644 index 00000000000..747061941f6 --- /dev/null +++ b/knowledge-fs/packages/api/src/namespace-source-preview-routes.ts @@ -0,0 +1,122 @@ +import { createRoute, z } from "@hono/zod-openapi"; +import { ForbiddenResponse, UnauthorizedResponse } from "./gateway-openapi-contracts"; +import { ErrorResponseSchema } from "./gateway-route-schemas"; + +const ErrorResponse = { + content: { "application/json": { schema: ErrorResponseSchema } }, + description: "Request failed", +} as const; +const JobParams = z.object({ jobId: z.string().uuid() }); +const Job = z.object({ + jobId: z.string().uuid(), + status: z.enum(["queued", "running", "completed", "failed", "canceled", "consumed"]), + configurationFingerprint: z.string(), + expiresAt: z.string(), + errorCode: z.string().optional(), + importWorkflowId: z.string().uuid().optional(), +}); +const Page = z.object({ + pageId: z.string(), + sourceUrl: z.string().url(), + title: z.string().optional(), + description: z.string().optional(), +}); + +export const createNamespaceSourcePreviewRoute = createRoute({ + method: "post", + operationId: "createNamespaceSourcePreview", + path: "/namespace/source-preview-jobs", + request: { + body: { + required: true, + content: { + "application/json": { + schema: z + .object({ + credentialId: z.string().min(1).max(255), + pluginId: z.string().min(1).max(1024), + provider: z.string().min(1).max(255), + datasource: z.string().min(1).max(255), + parameters: z.record(z.unknown()), + rootUrl: z.string().url().max(4096), + configurationFingerprint: z.string().min(1).max(128), + }) + .strict(), + }, + }, + }, + }, + responses: { + 202: { + content: { "application/json": { schema: Job } }, + description: "Website preview queued", + }, + 400: ErrorResponse, + 401: UnauthorizedResponse, + 403: ForbiddenResponse, + }, +}); +export const getNamespaceSourcePreviewRoute = createRoute({ + method: "get", + operationId: "getNamespaceSourcePreview", + path: "/namespace/source-preview-jobs/{jobId}", + request: { params: JobParams }, + responses: { + 200: { + content: { "application/json": { schema: Job.extend({ pages: z.array(Page) }) } }, + description: "Website preview status", + }, + 404: ErrorResponse, + 401: UnauthorizedResponse, + 403: ForbiddenResponse, + }, +}); +export const cancelNamespaceSourcePreviewRoute = createRoute({ + method: "delete", + operationId: "cancelNamespaceSourcePreview", + path: "/namespace/source-preview-jobs/{jobId}", + request: { params: JobParams }, + responses: { + 200: { + content: { "application/json": { schema: Job } }, + description: "Website preview canceled", + }, + 404: ErrorResponse, + 401: UnauthorizedResponse, + 403: ForbiddenResponse, + }, +}); +export const consumeNamespaceSourcePreviewRoute = createRoute({ + method: "post", + operationId: "consumeNamespaceSourcePreview", + path: "/knowledge-spaces/{id}/sources/{sourceId}/namespace-preview-import", + request: { + params: z.object({ id: z.string().uuid(), sourceId: z.string().uuid() }), + headers: z.object({ "Idempotency-Key": z.string().min(8).max(255) }), + body: { + required: true, + content: { + "application/json": { + schema: z + .object({ + previewJobId: z.string().uuid(), + pageIds: z.array(z.string().min(1).max(128)).min(1).max(200), + configurationFingerprint: z.string().min(1).max(128), + }) + .strict(), + }, + }, + }, + }, + responses: { + 202: { + content: { "application/json": { schema: z.object({ workflowId: z.string().uuid() }) } }, + description: "Preview import accepted", + }, + 400: ErrorResponse, + 404: ErrorResponse, + 409: ErrorResponse, + 401: UnauthorizedResponse, + 403: ForbiddenResponse, + }, +}); diff --git a/knowledge-fs/packages/api/src/namespace-source-preview.test.ts b/knowledge-fs/packages/api/src/namespace-source-preview.test.ts new file mode 100644 index 00000000000..f7c65085352 --- /dev/null +++ b/knowledge-fs/packages/api/src/namespace-source-preview.test.ts @@ -0,0 +1,208 @@ +import { describe, expect, it, vi } from "vitest"; + +import { + createInMemoryNamespaceSourcePreviewRepository, + createNamespaceSourcePreviewService, +} from "./namespace-source-preview"; + +describe("namespace website source preview", () => { + it("stores content in KFS and consumes only selected page ids", async () => { + const objects = new Map(); + const createCrawlImport = vi.fn(async () => ({ id: "22222222-2222-4222-8222-222222222222" })); + const service = createNamespaceSourcePreviewService({ + repository: createInMemoryNamespaceSourcePreviewRepository(), + storage: { + putObject: vi.fn(async ({ key, body }) => { + objects.set(key, body); + }), + getObject: vi.fn(async (key) => objects.get(key) ?? null), + deleteObject: vi.fn(async (key) => { + objects.delete(key); + }), + listObjects: vi.fn(async ({ prefix }) => ({ + objects: [...objects.keys()] + .filter((key) => key.startsWith(prefix)) + .map((key) => ({ key })), + })), + } as never, + websiteCrawl: { + crawl: vi.fn(async () => ({ + pages: [{ sourceUrl: "https://example.com/a", title: "A", content: "large body" }], + })), + }, + workflows: { createCrawlImport } as never, + sources: { + get: vi.fn(async () => ({ + type: "web", + metadata: { + credentialId: "credential", + pluginId: "plugin", + provider: "firecrawl", + datasource: "crawl", + initialPreview: { configurationFingerprint: "f".repeat(64) }, + }, + })), + } as never, + now: () => new Date("2026-09-04T00:00:00.000Z"), + }); + const subject = { tenantId: "tenant", subjectId: "account", scopes: [] }; + const job = await service.create( + subject, + { + credentialId: "credential", + pluginId: "plugin", + provider: "firecrawl", + datasource: "crawl", + parameters: { url: "https://example.com" }, + rootUrl: "https://example.com", + }, + "f".repeat(64), + ); + expect(await service.tick()).toBe(true); + const [page] = await service.pages(subject, job.id); + expect(page).toMatchObject({ sourceUrl: "https://example.com/a" }); + if (!page) throw new Error("Expected a preview page"); + + const workflowId = await service.consume(subject, { + jobId: job.id, + pageIds: [page.pageId], + configurationFingerprint: "f".repeat(64), + knowledgeSpaceId: "11111111-1111-4111-8111-111111111111", + sourceId: "33333333-3333-4333-8333-333333333333", + idempotencyKey: "request:crawl-import", + }); + + expect(workflowId).toBe("22222222-2222-4222-8222-222222222222"); + expect(createCrawlImport).toHaveBeenCalledWith( + expect.objectContaining({ + sourceUrls: ["https://example.com/a"], + pages: [{ sourceUrl: "https://example.com/a", title: "A", content: "large body" }], + }), + ); + expect(JSON.stringify(createCrawlImport.mock.calls[0])).toContain("large body"); + expect(objects.size).toBe(0); + }); + + it("rejects a page larger than the configured per-page limit", async () => { + const objects = new Map(); + const repository = createInMemoryNamespaceSourcePreviewRepository(); + const service = createNamespaceSourcePreviewService({ + repository, + storage: memoryStorage(objects) as never, + websiteCrawl: crawlPages([{ sourceUrl: "https://example.com/large", content: "12345" }]), + workflows: {} as never, + sources: {} as never, + maxPageBytes: 4, + maxJobBytes: 10, + now: () => new Date("2026-09-04T00:00:00.000Z"), + }); + const { subject, job } = await createJob(service); + + await service.tick(); + + await expect(service.get(subject, job.id)).resolves.toMatchObject({ + status: "failed", + errorCode: "PREVIEW_PAGE_TOO_LARGE", + }); + expect(objects.size).toBe(0); + }); + + it("rejects a preview whose aggregate content exceeds the job limit", async () => { + const objects = new Map(); + const repository = createInMemoryNamespaceSourcePreviewRepository(); + const service = createNamespaceSourcePreviewService({ + repository, + storage: memoryStorage(objects) as never, + websiteCrawl: crawlPages([ + { sourceUrl: "https://example.com/a", content: "123" }, + { sourceUrl: "https://example.com/b", content: "456" }, + ]), + workflows: {} as never, + sources: {} as never, + maxPageBytes: 4, + maxJobBytes: 5, + now: () => new Date("2026-09-04T00:00:00.000Z"), + }); + const { subject, job } = await createJob(service); + + await service.tick(); + + await expect(service.get(subject, job.id)).resolves.toMatchObject({ + status: "failed", + errorCode: "PREVIEW_JOB_TOO_LARGE", + }); + expect(objects.size).toBe(0); + }); + + it("retries object cleanup after preview metadata persistence fails", async () => { + const objects = new Map(); + const base = createInMemoryNamespaceSourcePreviewRepository(); + const repository = { + ...base, + complete: vi.fn(async () => { + throw new Error("database unavailable"); + }), + }; + let deleteFailures = 2; + const storage = memoryStorage(objects); + storage.deleteObject = vi.fn(async (key: string) => { + if (deleteFailures-- > 0) throw new Error("storage unavailable"); + objects.delete(key); + }); + const service = createNamespaceSourcePreviewService({ + repository, + storage: storage as never, + websiteCrawl: crawlPages([{ sourceUrl: "https://example.com/a", content: "body" }]), + workflows: {} as never, + sources: {} as never, + now: () => new Date("2026-09-04T00:00:00.000Z"), + }); + const { subject, job } = await createJob(service); + + await service.tick(); + expect(objects.size).toBe(1); + await expect(service.get(subject, job.id)).resolves.toMatchObject({ status: "failed" }); + + await service.tick(); + expect(objects.size).toBe(0); + await expect(service.get(subject, job.id)).resolves.toMatchObject({ + contentCleanedAt: "2026-09-04T00:00:00.000Z", + }); + }); +}); + +function memoryStorage(objects: Map) { + return { + putObject: vi.fn(async ({ key, body }: { key: string; body: Uint8Array }) => { + objects.set(key, body); + }), + getObject: vi.fn(async (key: string) => objects.get(key) ?? null), + deleteObject: vi.fn(async (key: string) => { + objects.delete(key); + }), + listObjects: vi.fn(async ({ prefix }: { prefix: string }) => ({ + objects: [...objects.keys()].filter((key) => key.startsWith(prefix)).map((key) => ({ key })), + })), + }; +} + +function crawlPages(pages: Array<{ sourceUrl: string; content: string }>) { + return { crawl: vi.fn(async () => ({ pages })) } as never; +} + +async function createJob(service: ReturnType) { + const subject = { tenantId: "tenant", subjectId: "account", scopes: [] }; + const job = await service.create( + subject, + { + credentialId: "credential", + pluginId: "plugin", + provider: "firecrawl", + datasource: "crawl", + parameters: {}, + rootUrl: "https://example.com", + }, + "f".repeat(64), + ); + return { subject, job }; +} diff --git a/knowledge-fs/packages/api/src/namespace-source-preview.ts b/knowledge-fs/packages/api/src/namespace-source-preview.ts new file mode 100644 index 00000000000..419ad4c8573 --- /dev/null +++ b/knowledge-fs/packages/api/src/namespace-source-preview.ts @@ -0,0 +1,642 @@ +import { createHash, randomUUID } from "node:crypto"; + +import type { + AuthSubject, + DatabaseAdapter, + DatabaseQueryValue, + ObjectStorageAdapter, + Source, +} from "@knowledge/core"; + +import { optionalStringColumn, stringColumn } from "./database-row-utils"; +import { + databasePlaceholder, + jsonInsertPlaceholder, + quoteDatabaseIdentifier, +} from "./database-sql-utils"; +import type { SourceProductWorkflowService } from "./source-product-workflow"; +import type { SourceRepository } from "./source-repository"; +import type { WebsiteCrawlConnector } from "./website-crawl-connector"; + +export type NamespaceSourcePreviewStatus = + | "queued" + | "running" + | "completed" + | "failed" + | "canceled" + | "consumed"; + +export interface NamespaceSourcePreviewConfig { + readonly credentialId: string; + readonly datasource: string; + readonly parameters: Readonly>; + readonly pluginId: string; + readonly provider: string; + readonly rootUrl: string; +} + +export interface NamespaceSourcePreviewPage { + readonly contentHash: string; + readonly contentObjectKey: string; + readonly description?: string | undefined; + readonly pageId: string; + readonly sourceUrl: string; + readonly title?: string | undefined; +} + +export interface NamespaceSourcePreviewJob { + readonly accountId: string; + readonly configurationFingerprint: string; + readonly config: NamespaceSourcePreviewConfig; + readonly consumedAt?: string | undefined; + readonly contentCleanedAt?: string | undefined; + readonly createdAt: string; + readonly errorCode?: string | undefined; + readonly expiresAt: string; + readonly id: string; + readonly importWorkflowId?: string | undefined; + readonly status: NamespaceSourcePreviewStatus; + readonly tenantId: string; + readonly updatedAt: string; +} + +export interface NamespaceSourcePreviewRepository { + create(job: NamespaceSourcePreviewJob): Promise; + get(input: { + tenantId: string; + accountId: string; + jobId: string; + }): Promise; + listPages(jobId: string): Promise; + expire(now: string): Promise; + claim(now: string): Promise; + complete(input: { + jobId: string; + now: string; + pages: readonly NamespaceSourcePreviewPage[]; + }): Promise; + fail(input: { jobId: string; now: string; errorCode: string }): Promise; + cancel(input: { + tenantId: string; + accountId: string; + jobId: string; + now: string; + }): Promise; + consume(input: { jobId: string; now: string; workflowId: string }): Promise; + claimCleanup(): Promise; + completeCleanup(input: { jobId: string; now: string }): Promise; +} + +export class NamespaceSourcePreviewError extends Error { + constructor( + readonly code: string, + message: string, + ) { + super(message); + this.name = "NamespaceSourcePreviewError"; + } +} + +const jobsTable = "namespace_source_preview_jobs"; +const pagesTable = "namespace_source_preview_pages"; + +export function createInMemoryNamespaceSourcePreviewRepository(): NamespaceSourcePreviewRepository { + const jobs = new Map(); + const pages = new Map(); + return { + create: async (job) => { + jobs.set(job.id, job); + }, + get: async ({ tenantId, accountId, jobId }) => { + const job = jobs.get(jobId); + return job?.tenantId === tenantId && job.accountId === accountId ? job : null; + }, + listPages: async (jobId) => [...(pages.get(jobId) ?? [])], + expire: async (now) => { + const job = [...jobs.values()].find( + (candidate) => + ["queued", "running", "completed"].includes(candidate.status) && + candidate.expiresAt <= now, + ); + if (!job) return []; + jobs.set(job.id, { ...job, status: "failed", errorCode: "PREVIEW_EXPIRED", updatedAt: now }); + return [...(pages.get(job.id) ?? [])]; + }, + claim: async (now) => { + const job = [...jobs.values()] + .filter((j) => j.status === "queued" && j.expiresAt > now) + .sort((a, b) => a.createdAt.localeCompare(b.createdAt))[0]; + if (!job) return null; + const claimed = { ...job, status: "running" as const, updatedAt: now }; + jobs.set(job.id, claimed); + return claimed; + }, + complete: async ({ jobId, now, pages: items }) => { + const job = jobs.get(jobId); + if (job?.status !== "running") return false; + pages.set(jobId, [...items]); + jobs.set(jobId, { ...job, status: "completed", updatedAt: now }); + return true; + }, + fail: async ({ jobId, now, errorCode }) => { + const job = jobs.get(jobId); + if (job) jobs.set(jobId, { ...job, status: "failed", updatedAt: now, errorCode }); + }, + cancel: async ({ tenantId, accountId, jobId, now }) => { + const job = jobs.get(jobId); + if (!job || job.tenantId !== tenantId || job.accountId !== accountId) return false; + jobs.set(jobId, { ...job, status: "canceled", updatedAt: now }); + return true; + }, + consume: async ({ jobId, now, workflowId }) => { + const job = jobs.get(jobId); + if (job) + jobs.set(jobId, { + ...job, + status: "consumed", + updatedAt: now, + consumedAt: now, + importWorkflowId: workflowId, + }); + }, + claimCleanup: async () => + [...jobs.values()].find( + (job) => ["failed", "canceled", "consumed"].includes(job.status) && !job.contentCleanedAt, + ) ?? null, + completeCleanup: async ({ jobId, now }) => { + const job = jobs.get(jobId); + if (job) jobs.set(jobId, { ...job, contentCleanedAt: now }); + }, + }; +} + +export function createDatabaseNamespaceSourcePreviewRepository(input: { + database: DatabaseAdapter; +}): NamespaceSourcePreviewRepository { + const db = input.database; + const q = (name: string) => quoteDatabaseIdentifier(db, name); + const ph = (n: number) => databasePlaceholder(db, n); + return { + create: async (job) => { + await db.execute({ + operation: "insert", + tableName: jobsTable, + maxRows: 0, + params: [ + job.id, + job.tenantId, + job.accountId, + job.status, + JSON.stringify(job.config), + job.configurationFingerprint, + job.expiresAt, + job.createdAt, + job.updatedAt, + ], + sql: `INSERT INTO ${q(jobsTable)} (${q("id")},${q("tenant_id")},${q("account_id")},${q("status")},${q("provider_config")},${q("configuration_fingerprint")},${q("expires_at")},${q("created_at")},${q("updated_at")}) VALUES (${ph(1)},${ph(2)},${ph(3)},${ph(4)},${jsonInsertPlaceholder(db, 5, "provider_config")},${ph(6)},${ph(7)},${ph(8)},${ph(9)})`, + }); + }, + get: async ({ tenantId, accountId, jobId }) => { + const result = await db.execute({ + operation: "select", + tableName: jobsTable, + maxRows: 1, + params: [jobId, tenantId, accountId], + sql: `SELECT * FROM ${q(jobsTable)} WHERE ${q("id")}=${ph(1)} AND ${q("tenant_id")}=${ph(2)} AND ${q("account_id")}=${ph(3)}`, + }); + return result.rows[0] ? rowJob(result.rows[0]) : null; + }, + listPages: async (jobId) => + ( + await db.execute({ + operation: "select", + tableName: pagesTable, + maxRows: 200, + params: [jobId], + sql: `SELECT * FROM ${q(pagesTable)} WHERE ${q("job_id")}=${ph(1)} ORDER BY ${q("created_at")},${q("page_id")}`, + }) + ).rows.map(rowPage), + expire: async (now) => + db.transaction(async (tx) => { + const found = await tx.execute({ + operation: "select", + tableName: jobsTable, + maxRows: 1, + params: [now], + sql: `SELECT ${q("id")} FROM ${q(jobsTable)} WHERE ${q("status")} IN ('queued','running','completed') AND ${q("expires_at")}<=${ph(1)} ORDER BY ${q("expires_at")} LIMIT 1 FOR UPDATE`, + }); + const jobId = found.rows[0] ? stringColumn(found.rows[0], "id") : undefined; + if (!jobId) return []; + const expiredPages = await tx.execute({ + operation: "select", + tableName: pagesTable, + maxRows: 200, + params: [jobId], + sql: `SELECT * FROM ${q(pagesTable)} WHERE ${q("job_id")}=${ph(1)}`, + }); + await tx.execute({ + operation: "update", + tableName: jobsTable, + maxRows: 0, + params: [now, jobId], + sql: `UPDATE ${q(jobsTable)} SET ${q("status")}='failed',${q("error_code")}='PREVIEW_EXPIRED',${q("updated_at")}=${ph(1)} WHERE ${q("id")}=${ph(2)}`, + }); + return expiredPages.rows.map(rowPage); + }), + claim: async (now) => + db.transaction(async (tx) => { + const found = await tx.execute({ + operation: "select", + tableName: jobsTable, + maxRows: 1, + params: [now], + sql: `SELECT * FROM ${q(jobsTable)} WHERE ${q("status")}='queued' AND ${q("expires_at")}>${ph(1)} ORDER BY ${q("created_at")} LIMIT 1 FOR UPDATE`, + }); + if (!found.rows[0]) return null; + const job = rowJob(found.rows[0]); + const updated = await tx.execute({ + operation: "update", + tableName: jobsTable, + maxRows: 0, + params: [now, job.id], + sql: `UPDATE ${q(jobsTable)} SET ${q("status")}='running',${q("updated_at")}=${ph(1)} WHERE ${q("id")}=${ph(2)} AND ${q("status")}='queued'`, + }); + return updated.rowsAffected === 1 ? { ...job, status: "running", updatedAt: now } : null; + }), + complete: async ({ jobId, now, pages }) => + db.transaction(async (tx) => { + for (const page of pages) + await tx.execute({ + operation: "insert", + tableName: pagesTable, + maxRows: 0, + params: [ + jobId, + page.pageId, + page.sourceUrl, + page.title ?? null, + page.description ?? null, + page.contentHash, + page.contentObjectKey, + now, + ], + sql: `INSERT INTO ${q(pagesTable)} (${q("job_id")},${q("page_id")},${q("source_url")},${q("title")},${q("description")},${q("content_hash")},${q("content_object_key")},${q("created_at")}) VALUES (${[1, 2, 3, 4, 5, 6, 7, 8].map(ph).join(",")})`, + }); + const completed = await tx.execute({ + operation: "update", + tableName: jobsTable, + maxRows: 0, + params: [now, jobId], + sql: `UPDATE ${q(jobsTable)} SET ${q("status")}='completed',${q("updated_at")}=${ph(1)} WHERE ${q("id")}=${ph(2)} AND ${q("status")}='running'`, + }); + return completed.rowsAffected === 1; + }), + fail: async ({ jobId, now, errorCode }) => { + await db.execute({ + operation: "update", + tableName: jobsTable, + maxRows: 0, + params: [errorCode, now, jobId], + sql: `UPDATE ${q(jobsTable)} SET ${q("status")}='failed',${q("error_code")}=${ph(1)},${q("updated_at")}=${ph(2)} WHERE ${q("id")}=${ph(3)} AND ${q("status")} IN ('queued','running')`, + }); + }, + cancel: async ({ tenantId, accountId, jobId, now }) => + ( + await db.execute({ + operation: "update", + tableName: jobsTable, + maxRows: 0, + params: [now, jobId, tenantId, accountId], + sql: `UPDATE ${q(jobsTable)} SET ${q("status")}='canceled',${q("updated_at")}=${ph(1)} WHERE ${q("id")}=${ph(2)} AND ${q("tenant_id")}=${ph(3)} AND ${q("account_id")}=${ph(4)} AND ${q("status")} IN ('queued','running','completed')`, + }) + ).rowsAffected === 1, + consume: async ({ jobId, now, workflowId }) => { + await db.execute({ + operation: "update", + tableName: jobsTable, + maxRows: 0, + params: [workflowId, now, now, jobId], + sql: `UPDATE ${q(jobsTable)} SET ${q("status")}='consumed',${q("import_workflow_id")}=${ph(1)},${q("consumed_at")}=${ph(2)},${q("updated_at")}=${ph(3)} WHERE ${q("id")}=${ph(4)} AND ${q("status")}='completed'`, + }); + }, + claimCleanup: async () => { + const result = await db.execute({ + operation: "select", + tableName: jobsTable, + maxRows: 1, + params: [], + sql: `SELECT * FROM ${q(jobsTable)} WHERE ${q("status")} IN ('failed','canceled','consumed') AND ${q("content_cleaned_at")} IS NULL ORDER BY ${q("updated_at")} LIMIT 1`, + }); + return result.rows[0] ? rowJob(result.rows[0]) : null; + }, + completeCleanup: async ({ jobId, now }) => { + await db.execute({ + operation: "update", + tableName: jobsTable, + maxRows: 0, + params: [now, jobId], + sql: `UPDATE ${q(jobsTable)} SET ${q("content_cleaned_at")}=${ph(1)} WHERE ${q("id")}=${ph(2)} AND ${q("status")} IN ('failed','canceled','consumed')`, + }); + }, + }; +} + +function rowJob(row: Readonly>): NamespaceSourcePreviewJob { + const configRaw = row.provider_config; + const config = ( + typeof configRaw === "string" ? JSON.parse(configRaw) : configRaw + ) as NamespaceSourcePreviewConfig; + return { + id: stringColumn(row, "id"), + tenantId: stringColumn(row, "tenant_id"), + accountId: stringColumn(row, "account_id"), + status: stringColumn(row, "status") as NamespaceSourcePreviewStatus, + config, + configurationFingerprint: stringColumn(row, "configuration_fingerprint"), + expiresAt: stringColumn(row, "expires_at"), + createdAt: stringColumn(row, "created_at"), + updatedAt: stringColumn(row, "updated_at"), + ...(optionalStringColumn(row, "consumed_at") + ? { consumedAt: optionalStringColumn(row, "consumed_at") } + : {}), + ...(optionalStringColumn(row, "content_cleaned_at") + ? { contentCleanedAt: optionalStringColumn(row, "content_cleaned_at") } + : {}), + ...(optionalStringColumn(row, "error_code") + ? { errorCode: optionalStringColumn(row, "error_code") } + : {}), + ...(optionalStringColumn(row, "import_workflow_id") + ? { importWorkflowId: optionalStringColumn(row, "import_workflow_id") } + : {}), + }; +} +function rowPage(row: Readonly>): NamespaceSourcePreviewPage { + return { + pageId: stringColumn(row, "page_id"), + sourceUrl: stringColumn(row, "source_url"), + contentHash: stringColumn(row, "content_hash"), + contentObjectKey: stringColumn(row, "content_object_key"), + ...(optionalStringColumn(row, "title") ? { title: optionalStringColumn(row, "title") } : {}), + ...(optionalStringColumn(row, "description") + ? { description: optionalStringColumn(row, "description") } + : {}), + }; +} + +export function createNamespaceSourcePreviewService(input: { + repository: NamespaceSourcePreviewRepository; + storage: ObjectStorageAdapter; + websiteCrawl: WebsiteCrawlConnector; + workflows: SourceProductWorkflowService; + sources: SourceRepository; + now?: () => Date; + maxPageBytes?: number; + maxJobBytes?: number; +}) { + const now = input.now ?? (() => new Date()); + const maxPageBytes = input.maxPageBytes ?? 20 * 1024 * 1024; + const maxJobBytes = input.maxJobBytes ?? 200 * 1024 * 1024; + if (!Number.isSafeInteger(maxPageBytes) || maxPageBytes < 1) + throw new Error("Namespace preview maxPageBytes must be a positive safe integer"); + if (!Number.isSafeInteger(maxJobBytes) || maxJobBytes < maxPageBytes) + throw new Error("Namespace preview maxJobBytes must be a safe integer >= maxPageBytes"); + const objectPrefix = (job: NamespaceSourcePreviewJob) => + `__namespace-source-previews/${encodeURIComponent(job.tenantId)}/${job.id}/`; + const cleanup = async (job: NamespaceSourcePreviewJob) => { + const prefix = objectPrefix(job); + for (;;) { + const result = await input.storage.listObjects({ prefix, limit: 100 }); + for (const object of result.objects) await input.storage.deleteObject(object.key); + if (result.objects.length === 0) break; + } + await input.repository.completeCleanup({ jobId: job.id, now: now().toISOString() }); + }; + const requireJob = async (subject: AuthSubject, jobId: string) => { + const job = await input.repository.get({ + tenantId: subject.tenantId, + accountId: subject.subjectId, + jobId, + }); + if (!job) throw new NamespaceSourcePreviewError("PREVIEW_NOT_FOUND", "Preview job not found"); + return job; + }; + return { + create: async ( + subject: AuthSubject, + config: NamespaceSourcePreviewConfig, + fingerprint: string, + ) => { + const timestamp = now(); + const job: NamespaceSourcePreviewJob = { + id: randomUUID(), + tenantId: subject.tenantId, + accountId: subject.subjectId, + status: "queued", + config, + configurationFingerprint: fingerprint, + createdAt: timestamp.toISOString(), + updatedAt: timestamp.toISOString(), + expiresAt: new Date(timestamp.getTime() + 60 * 60_000).toISOString(), + }; + await input.repository.create(job); + return job; + }, + get: (subject: AuthSubject, jobId: string) => requireJob(subject, jobId), + pages: async (subject: AuthSubject, jobId: string) => { + const job = await requireJob(subject, jobId); + return job.status === "completed" || job.status === "consumed" + ? input.repository.listPages(job.id) + : []; + }, + cancel: async (subject: AuthSubject, jobId: string) => { + const job = await requireJob(subject, jobId); + await input.repository.cancel({ + tenantId: subject.tenantId, + accountId: subject.subjectId, + jobId, + now: now().toISOString(), + }); + const canceled = await requireJob(subject, jobId); + await cleanup(canceled); + return requireJob(subject, jobId); + }, + tick: async () => { + const timestamp = now().toISOString(); + await input.repository.expire(timestamp); + const cleanupJob = await input.repository.claimCleanup(); + if (cleanupJob) { + await cleanup(cleanupJob).catch(() => {}); + return true; + } + const job = await input.repository.claim(timestamp); + if (!job) return false; + const source: Source = { + id: randomUUID(), + knowledgeSpaceId: "00000000-0000-0000-0000-000000000000", + name: "Namespace website preview", + type: "web", + uri: job.config.rootUrl, + status: "disabled", + permissionScope: [], + version: 1, + createdAt: timestamp, + updatedAt: timestamp, + metadata: { + credentialId: job.config.credentialId, + datasource: job.config.datasource, + parameters: job.config.parameters, + pluginId: job.config.pluginId, + provider: job.config.provider, + datasourceParameterMode: "exact", + }, + }; + const saved: string[] = []; + try { + const result = await input.websiteCrawl.crawl({ + source, + tenantId: job.tenantId, + userId: job.accountId, + }); + const pages: NamespaceSourcePreviewPage[] = []; + let totalBytes = 0; + for (const page of result.pages) { + const body = new TextEncoder().encode(page.content); + if (body.byteLength > maxPageBytes) + throw new NamespaceSourcePreviewError( + "PREVIEW_PAGE_TOO_LARGE", + "Preview page exceeds the per-page content limit", + ); + if (totalBytes > maxJobBytes - body.byteLength) + throw new NamespaceSourcePreviewError( + "PREVIEW_JOB_TOO_LARGE", + "Preview job exceeds the total content limit", + ); + totalBytes += body.byteLength; + const hash = createHash("sha256").update(body).digest("hex"); + const pageId = createHash("sha256").update(page.sourceUrl).digest("hex").slice(0, 32); + const key = `${objectPrefix(job)}${pageId}-${hash}.bin`; + await input.storage.putObject({ + key, + body, + contentType: "text/markdown", + metadata: { contentHash: hash, lifecycle: "namespace-source-preview", jobId: job.id }, + }); + saved.push(key); + pages.push({ + pageId, + sourceUrl: page.sourceUrl, + contentHash: hash, + contentObjectKey: key, + ...(page.title ? { title: page.title } : {}), + ...(page.description ? { description: page.description } : {}), + }); + } + const completed = await input.repository.complete({ + jobId: job.id, + now: now().toISOString(), + pages, + }); + if (!completed) + throw new NamespaceSourcePreviewError( + "PREVIEW_STATE_CONFLICT", + "Preview job could not be completed from its current state", + ); + } catch (error) { + for (const key of saved) await input.storage.deleteObject(key).catch(() => {}); + await input.repository.fail({ + jobId: job.id, + now: now().toISOString(), + errorCode: + error instanceof NamespaceSourcePreviewError ? error.code : "PREVIEW_PROVIDER_FAILED", + }); + const failed = await input.repository.get({ + tenantId: job.tenantId, + accountId: job.accountId, + jobId: job.id, + }); + if (failed) await cleanup(failed).catch(() => {}); + } + return true; + }, + consume: async ( + subject: AuthSubject, + request: { + jobId: string; + pageIds: readonly string[]; + configurationFingerprint: string; + knowledgeSpaceId: string; + sourceId: string; + idempotencyKey: string; + }, + ) => { + const job = await requireJob(subject, request.jobId); + if (job.importWorkflowId) return job.importWorkflowId; + if (job.expiresAt <= now().toISOString()) + throw new NamespaceSourcePreviewError("PREVIEW_EXPIRED", "Preview job expired"); + if ( + job.status !== "completed" || + job.configurationFingerprint !== request.configurationFingerprint + ) + throw new NamespaceSourcePreviewError( + "PREVIEW_CONFIGURATION_MISMATCH", + "Preview configuration changed", + ); + const source = await input.sources.get({ + knowledgeSpaceId: request.knowledgeSpaceId, + id: request.sourceId, + }); + if ( + !source || + source.type !== "web" || + source.metadata.credentialId !== job.config.credentialId || + source.metadata.pluginId !== job.config.pluginId || + source.metadata.provider !== job.config.provider || + source.metadata.datasource !== job.config.datasource || + (source.metadata.initialPreview as { configurationFingerprint?: unknown } | undefined) + ?.configurationFingerprint !== job.configurationFingerprint + ) + throw new NamespaceSourcePreviewError( + "PREVIEW_TARGET_MISMATCH", + "Preview target source does not match", + ); + const byId = new Map((await input.repository.listPages(job.id)).map((p) => [p.pageId, p])); + const selected: NamespaceSourcePreviewPage[] = []; + for (const pageId of request.pageIds) { + const page = byId.get(pageId); + if (!page) + throw new NamespaceSourcePreviewError("PREVIEW_PAGE_NOT_FOUND", "Preview page not found"); + selected.push(page); + } + const pages = await Promise.all( + selected.map(async (p) => { + const body = await input.storage.getObject(p.contentObjectKey); + if (!body) + throw new NamespaceSourcePreviewError("PREVIEW_EXPIRED", "Preview content expired"); + return { + sourceUrl: p.sourceUrl, + content: new TextDecoder().decode(body), + ...(p.title ? { title: p.title } : {}), + ...(p.description ? { description: p.description } : {}), + }; + }), + ); + const workflow = await input.workflows.createCrawlImport({ + subject, + callerKind: "interactive", + knowledgeSpaceId: request.knowledgeSpaceId, + sourceId: request.sourceId, + idempotencyKey: request.idempotencyKey, + sourceUrls: pages.map((p) => p.sourceUrl), + pages, + }); + await input.repository.consume({ + jobId: job.id, + now: now().toISOString(), + workflowId: workflow.id, + }); + const consumed = await requireJob(subject, job.id); + await cleanup(consumed).catch(() => {}); + return workflow.id; + }, + }; +} diff --git a/knowledge-fs/packages/database/migrations/0052_namespace_source_previews.postgres.sql b/knowledge-fs/packages/database/migrations/0052_namespace_source_previews.postgres.sql new file mode 100644 index 00000000000..1406e3f78ac --- /dev/null +++ b/knowledge-fs/packages/database/migrations/0052_namespace_source_previews.postgres.sql @@ -0,0 +1,37 @@ +-- Knowledge Platform schema migration +-- Migration id: 0052_namespace_source_previews +-- Dialect: postgres + +CREATE TABLE IF NOT EXISTS "namespace_source_preview_jobs" ( + "id" uuid PRIMARY KEY, + "tenant_id" varchar(255) NOT NULL, + "account_id" varchar(255) NOT NULL, + "status" varchar(32) NOT NULL, + "provider_config" jsonb NOT NULL, + "configuration_fingerprint" varchar(128) NOT NULL, + "expires_at" timestamptz NOT NULL, + "consumed_at" timestamptz NULL, + "content_cleaned_at" timestamptz NULL, + "import_workflow_id" uuid NULL, + "error_code" varchar(128) NULL, + "created_at" timestamptz NOT NULL, + "updated_at" timestamptz NOT NULL +); +CREATE INDEX IF NOT EXISTS "namespace_source_preview_jobs_claim_idx" + ON "namespace_source_preview_jobs" ("status", "expires_at", "created_at"); +CREATE INDEX IF NOT EXISTS "namespace_source_preview_jobs_owner_idx" + ON "namespace_source_preview_jobs" ("tenant_id", "account_id", "created_at"); +CREATE INDEX IF NOT EXISTS "namespace_source_preview_jobs_cleanup_idx" + ON "namespace_source_preview_jobs" ("status", "content_cleaned_at", "updated_at"); + +CREATE TABLE IF NOT EXISTS "namespace_source_preview_pages" ( + "job_id" uuid NOT NULL REFERENCES "namespace_source_preview_jobs" ("id") ON DELETE CASCADE, + "page_id" varchar(128) NOT NULL, + "source_url" varchar(4096) NOT NULL, + "title" varchar(500) NULL, + "description" text NULL, + "content_hash" varchar(64) NOT NULL, + "content_object_key" varchar(2048) NOT NULL, + "created_at" timestamptz NOT NULL, + PRIMARY KEY ("job_id", "page_id") +); diff --git a/knowledge-fs/packages/database/migrations/0052_namespace_source_previews.tidb.sql b/knowledge-fs/packages/database/migrations/0052_namespace_source_previews.tidb.sql new file mode 100644 index 00000000000..8dba2877a94 --- /dev/null +++ b/knowledge-fs/packages/database/migrations/0052_namespace_source_previews.tidb.sql @@ -0,0 +1,35 @@ +-- Knowledge Platform schema migration +-- Migration id: 0052_namespace_source_previews +-- Dialect: tidb + +CREATE TABLE IF NOT EXISTS `namespace_source_preview_jobs` ( + `id` char(36) PRIMARY KEY, + `tenant_id` varchar(255) NOT NULL, + `account_id` varchar(255) NOT NULL, + `status` varchar(32) NOT NULL, + `provider_config` json NOT NULL, + `configuration_fingerprint` varchar(128) NOT NULL, + `expires_at` datetime(6) NOT NULL, + `consumed_at` datetime(6) NULL, + `content_cleaned_at` datetime(6) NULL, + `import_workflow_id` char(36) NULL, + `error_code` varchar(128) NULL, + `created_at` datetime(6) NOT NULL, + `updated_at` datetime(6) NOT NULL, + INDEX `namespace_source_preview_jobs_claim_idx` (`status`, `expires_at`, `created_at`), + INDEX `namespace_source_preview_jobs_owner_idx` (`tenant_id`, `account_id`, `created_at`), + INDEX `namespace_source_preview_jobs_cleanup_idx` (`status`, `content_cleaned_at`, `updated_at`) +); + +CREATE TABLE IF NOT EXISTS `namespace_source_preview_pages` ( + `job_id` char(36) NOT NULL, + `page_id` varchar(128) NOT NULL, + `source_url` varchar(4096) NOT NULL, + `title` varchar(500) NULL, + `description` text NULL, + `content_hash` varchar(64) NOT NULL, + `content_object_key` varchar(2048) NOT NULL, + `created_at` datetime(6) NOT NULL, + PRIMARY KEY (`job_id`, `page_id`), + CONSTRAINT `namespace_source_preview_pages_job_fk` FOREIGN KEY (`job_id`) REFERENCES `namespace_source_preview_jobs` (`id`) ON DELETE CASCADE +); diff --git a/knowledge-fs/scripts/export-openapi.mjs b/knowledge-fs/scripts/export-openapi.mjs index a9a4d1d4735..52ee702bb5d 100644 --- a/knowledge-fs/scripts/export-openapi.mjs +++ b/knowledge-fs/scripts/export-openapi.mjs @@ -16,7 +16,7 @@ function argumentValue(name) { process.env.NODE_ENV = "test"; const [ { createNodePlatformAdapter }, - { createKnowledgeGateway, registerSourceProductHandlers }, + { createKnowledgeGateway, registerNamespaceSourcePreviewHandlers, registerSourceProductHandlers }, { createInMemoryCapabilityGrantProvenanceRepository }, ] = await Promise.all([ import("../packages/adapters/src/node.ts"), @@ -67,6 +67,7 @@ registerSourceProductHandlers({ repository: unavailableService, workflows: unavailableService, }); +registerNamespaceSourcePreviewHandlers({ app, service: unavailableService }); const response = await app.request("/openapi.json"); if (!response.ok) { throw new Error(`OpenAPI export failed with HTTP ${response.status}`); diff --git a/packages/contracts/generated/api/console/knowledge-fs/types.gen.ts b/packages/contracts/generated/api/console/knowledge-fs/types.gen.ts index da963337a57..bf56808be57 100644 --- a/packages/contracts/generated/api/console/knowledge-fs/types.gen.ts +++ b/packages/contracts/generated/api/console/knowledge-fs/types.gen.ts @@ -1308,6 +1308,7 @@ export type KnowledgeFsInitialSourcePreviewFileResponse = { export type KnowledgeFsInitialSourcePreviewPageResponse = { description?: string | null + page_id?: string | null source_url: string title?: string | null } @@ -2222,6 +2223,7 @@ export type KnowledgeFsInitialWebsiteCrawlOptionsPayload = { export type KnowledgeFsInitialWebsiteSelectionPayload = { canonical_url?: string | null + pageId?: string | null source_url: string title?: string | null } diff --git a/packages/contracts/generated/api/console/knowledge-fs/zod.gen.ts b/packages/contracts/generated/api/console/knowledge-fs/zod.gen.ts index a8bb604b00e..4c233ce3215 100644 --- a/packages/contracts/generated/api/console/knowledge-fs/zod.gen.ts +++ b/packages/contracts/generated/api/console/knowledge-fs/zod.gen.ts @@ -708,6 +708,7 @@ export const zKnowledgeFsInitialSourcePreviewFileResponse = z.object({ */ export const zKnowledgeFsInitialSourcePreviewPageResponse = z.object({ description: z.string().nullish(), + page_id: z.string().nullish(), source_url: z.string(), title: z.string().nullish(), }) @@ -792,6 +793,14 @@ export const zKnowledgeFsSpaceCreateResponse = z.object({ /** * KnowledgeFSProductPermission + * + * Product-level capability of a caller on one KnowledgeFS space. + * + * These values are what the console returns as ``permission_keys`` so the web can gate + * its UI. Enterprise RBAC is not evaluated in this vocabulary: KnowledgeFS spaces reuse + * the legacy knowledge base (``dataset_*``) permission points, and each capability is + * granted when the caller holds the dataset permission it maps to (see + * :data:`RBAC_PERMISSION_BY_PRODUCT_PERMISSION`). */ export const zKnowledgeFsProductPermission = z.enum([ 'knowledge_space_access_config', @@ -2276,6 +2285,7 @@ export const zKnowledgeFsInitialWebsiteCrawlOptionsPayload = z.object({ */ export const zKnowledgeFsInitialWebsiteSelectionPayload = z.object({ canonical_url: z.string().min(1).max(4096).nullish(), + pageId: z.string().min(1).max(128).nullish(), source_url: z.string().min(1).max(4096), title: z.string().max(500).nullish(), }) diff --git a/web/features/new-rag/create/__tests__/page.spec.tsx b/web/features/new-rag/create/__tests__/page.spec.tsx index d3f5838782c..2a9bbb47afb 100644 --- a/web/features/new-rag/create/__tests__/page.spec.tsx +++ b/web/features/new-rag/create/__tests__/page.spec.tsx @@ -620,10 +620,12 @@ describe('CreateKnowledgePage', () => { job_id: 'website-preview-1', status: 'completed', result: { + configuration_fingerprint: 'a'.repeat(64), kind: 'website_crawl', pages: [ { description: 'Getting started', + page_id: 'page-getting-started', source_url: 'https://docs.dify.ai/getting-started', title: 'Getting started', }, @@ -1665,11 +1667,13 @@ describe('CreateKnowledgePage', () => { }, pluginId: 'langgenius/firecrawl_datasource', previewJobId: 'website-preview-1', + previewConfigurationFingerprint: 'a'.repeat(64), provider: 'firecrawl', providerDisplayName: 'Firecrawl', root_url: 'https://docs.dify.ai/', selection: [ { + pageId: 'page-getting-started', source_url: 'https://docs.dify.ai/getting-started', title: 'Getting started', }, diff --git a/web/features/new-rag/create/source-setup.tsx b/web/features/new-rag/create/source-setup.tsx index 12466377ac0..f7e1fd56a49 100644 --- a/web/features/new-rag/create/source-setup.tsx +++ b/web/features/new-rag/create/source-setup.tsx @@ -283,7 +283,7 @@ function CreateSourceSetupSession({ setPreviewPages( (response.result.pages ?? []).map((page) => ({ description: page.description ?? undefined, - pageId: page.source_url, + pageId: page.page_id ?? page.source_url, sourceUrl: page.source_url, title: page.title ?? page.source_url, })), @@ -392,7 +392,7 @@ function CreateSourceSetupSession({ setPreviewPages( (response.result.pages ?? []).map((page) => ({ description: page.description ?? undefined, - pageId: page.source_url, + pageId: page.page_id ?? page.source_url, sourceUrl: page.source_url, title: page.title ?? page.source_url, })), @@ -457,6 +457,7 @@ function CreateSourceSetupSession({ : {}), root_url: sourceUri, selection: selectedPages.map((page) => ({ + pageId: page.pageId, source_url: page.sourceUrl, ...(page.title ? { title: page.title } : {}), })),