diff --git a/api/controllers/console/knowledge_fs/error.py b/api/controllers/console/knowledge_fs/error.py index 16c456ef126..fd691311681 100644 --- a/api/controllers/console/knowledge_fs/error.py +++ b/api/controllers/console/knowledge_fs/error.py @@ -18,7 +18,7 @@ class KnowledgeFSOperationUnavailableHTTPError(BaseHTTPException): class KnowledgeFSUpstreamUnavailableHTTPError(BaseHTTPException): error_code = "knowledge_fs_upstream_unavailable" description = "KnowledgeFS is unavailable." - code = 502 + code = 503 class KnowledgeFSInvalidRequestHTTPError(BaseHTTPException): diff --git a/api/controllers/console/knowledge_fs/resources.py b/api/controllers/console/knowledge_fs/resources.py index 2bdfaf595c1..f6c931d6d9a 100644 --- a/api/controllers/console/knowledge_fs/resources.py +++ b/api/controllers/console/knowledge_fs/resources.py @@ -58,6 +58,7 @@ from services.knowledge_fs.download_service import ( KnowledgeFSDownloadObjectNotFoundError, KnowledgeFSDownloadService, KnowledgeFSDownloadTooLargeError, + KnowledgeFSDownloadUnavailableError, ) from services.knowledge_fs.initial_source_preview import KnowledgeFSInitialSourcePreviewService from services.knowledge_fs.initial_source_preview_job import ( @@ -1432,6 +1433,8 @@ class KnowledgeFSSpaceLogicalDocumentsDownloadApi(Resource): raise RequestEntityTooLarge(str(exc)) from exc except KnowledgeFSDownloadObjectNotFoundError as exc: raise NotFound("KnowledgeFS document object not found") from exc + except KnowledgeFSDownloadUnavailableError as exc: + raise ServiceUnavailable("KnowledgeFS object storage is unavailable") from exc @console_ns.route("/knowledge-fs/spaces//logical-documents//download") @@ -1459,6 +1462,8 @@ class KnowledgeFSSpaceLogicalDocumentDownloadApi(Resource): body = KnowledgeFSDownloadService().load_stream(descriptor) except KnowledgeFSDownloadObjectNotFoundError as exc: raise NotFound("KnowledgeFS document object not found") from exc + except KnowledgeFSDownloadUnavailableError as exc: + raise ServiceUnavailable("KnowledgeFS object storage is unavailable") from exc response = Response(body, content_type=descriptor.mime_type or "application/octet-stream") response.content_length = descriptor.size_bytes response.headers["Content-Disposition"] = f"attachment; filename*=UTF-8''{quote(descriptor.filename, safe='')}" diff --git a/api/services/knowledge_fs/download_service.py b/api/services/knowledge_fs/download_service.py index 29e5d3950c2..9b2d3ed2ad9 100644 --- a/api/services/knowledge_fs/download_service.py +++ b/api/services/knowledge_fs/download_service.py @@ -9,7 +9,11 @@ from tempfile import NamedTemporaryFile from zipfile import ZIP_DEFLATED, ZipFile from services.file_service import FileService -from services.knowledge_fs.object_storage import KnowledgeFSObjectStorageService +from services.knowledge_fs.object_storage import ( + KnowledgeFSObjectStorageCorruptError, + KnowledgeFSObjectStorageError, + KnowledgeFSObjectStorageService, +) from services.knowledge_fs.product_dto import KnowledgeFSDocumentDownloadDescriptor KNOWLEDGE_FS_BATCH_DOWNLOAD_MAX_BYTES = 500 * 1024 * 1024 @@ -23,18 +27,27 @@ class KnowledgeFSDownloadObjectNotFoundError(FileNotFoundError): pass +class KnowledgeFSDownloadUnavailableError(RuntimeError): + pass + + class KnowledgeFSDownloadService: def __init__(self, *, object_storage: KnowledgeFSObjectStorageService | None = None) -> None: self._object_storage = object_storage or KnowledgeFSObjectStorageService() def load_stream(self, descriptor: KnowledgeFSDocumentDownloadDescriptor) -> Generator[bytes, None, None]: - metadata = self._object_storage.head_object(key=descriptor.object_key) - if metadata is None or metadata.size_bytes != descriptor.size_bytes: - raise KnowledgeFSDownloadObjectNotFoundError(descriptor.document_id) - stream = self._object_storage.load_stream(key=descriptor.object_key) - if stream is None: - raise KnowledgeFSDownloadObjectNotFoundError(descriptor.document_id) - return stream + try: + metadata = self._object_storage.head_object(key=descriptor.object_key) + if metadata is None or metadata.size_bytes != descriptor.size_bytes: + raise KnowledgeFSDownloadObjectNotFoundError(descriptor.document_id) + stream = self._object_storage.load_stream(key=descriptor.object_key) + if stream is None: + raise KnowledgeFSDownloadObjectNotFoundError(descriptor.document_id) + return stream + except KnowledgeFSObjectStorageCorruptError as exc: + raise KnowledgeFSDownloadObjectNotFoundError(descriptor.document_id) from exc + except KnowledgeFSObjectStorageError as exc: + raise KnowledgeFSDownloadUnavailableError("KnowledgeFS object storage is unavailable") from exc @contextmanager def build_zip_tempfile( diff --git a/api/tests/unit_tests/controllers/test_knowledge_fs_product_controllers.py b/api/tests/unit_tests/controllers/test_knowledge_fs_product_controllers.py index 8e676b15366..cdbc52ba0bf 100644 --- a/api/tests/unit_tests/controllers/test_knowledge_fs_product_controllers.py +++ b/api/tests/unit_tests/controllers/test_knowledge_fs_product_controllers.py @@ -10,7 +10,7 @@ from unittest.mock import MagicMock import pytest from flask import Flask -from werkzeug.exceptions import Forbidden, NotFound, RequestEntityTooLarge +from werkzeug.exceptions import Forbidden, NotFound, RequestEntityTooLarge, ServiceUnavailable from controllers.common import wraps as common_wraps from controllers.console import console_ns @@ -19,6 +19,7 @@ from controllers.console.wraps import RBACPermission, RBACResourceScope from controllers.service_api import service_api_ns from controllers.service_api.knowledge_fs import resources as service_resources from services.knowledge_fs.credential_service import KnowledgeFSServiceCredentialProfile +from services.knowledge_fs.download_service import KnowledgeFSDownloadUnavailableError from services.knowledge_fs.product_dto import ( KnowledgeFSDocumentDownloadDescriptor, KnowledgeFSDocumentStagedUploadAcceptedResponse, @@ -278,6 +279,88 @@ def test_single_document_download_preserves_resource_not_found_after_permission_ download_service_factory.assert_not_called() +def test_single_document_download_maps_storage_unavailable_after_permission_allows( + monkeypatch: pytest.MonkeyPatch, +) -> None: + descriptor = KnowledgeFSDocumentDownloadDescriptor( + document_id="document-1", + filename="guide.md", + mime_type="text/markdown", + object_key="namespaces/tenant-1/spaces/space-1/documents/guide.md", + sha256="a" * 64, + size_bytes=4, + ) + facade = SimpleNamespace(prepare_logical_document_download=MagicMock(return_value=descriptor)) + download_service = SimpleNamespace( + load_stream=MagicMock(side_effect=KnowledgeFSDownloadUnavailableError("storage unavailable")) + ) + permission_gate = MagicMock() + monkeypatch.setattr(common_wraps.dify_config, "RBAC_ENABLED", True) + monkeypatch.setattr( + common_wraps, + "current_account_with_tenant", + lambda: (SimpleNamespace(id="account-1"), "tenant-1"), + ) + monkeypatch.setattr(common_wraps, "enforce_rbac_access", permission_gate) + monkeypatch.setattr(console_resources, "_actor", lambda: ("account-1", "tenant-1")) + monkeypatch.setattr(console_resources, "_console_services", lambda: SimpleNamespace(facade=facade)) + monkeypatch.setattr(console_resources, "KnowledgeFSDownloadService", lambda: download_service) + permission_wrapper = _rbac_wrapper(console_resources.KnowledgeFSSpaceLogicalDocumentDownloadApi.get) + app = Flask(__name__) + + with app.test_request_context(), pytest.raises(ServiceUnavailable): + permission_wrapper( + console_resources.KnowledgeFSSpaceLogicalDocumentDownloadApi(), + control_space_id="control-1", + document_id="document-1", + ) + + permission_gate.assert_called_once() + download_service.load_stream.assert_called_once_with(descriptor) + + +def test_batch_document_download_maps_storage_unavailable_after_permission_allows( + monkeypatch: pytest.MonkeyPatch, +) -> None: + descriptor = KnowledgeFSDocumentDownloadDescriptor( + document_id="document-1", + filename="guide.md", + mime_type="text/markdown", + object_key="namespaces/tenant-1/spaces/space-1/documents/guide.md", + sha256="a" * 64, + size_bytes=4, + ) + facade = SimpleNamespace(prepare_logical_document_download=MagicMock(return_value=descriptor)) + zip_context = MagicMock() + zip_context.__enter__.side_effect = KnowledgeFSDownloadUnavailableError("storage unavailable") + download_service = SimpleNamespace(build_zip_tempfile=MagicMock(return_value=zip_context)) + permission_gate = MagicMock() + monkeypatch.setattr(common_wraps.dify_config, "RBAC_ENABLED", True) + monkeypatch.setattr( + common_wraps, + "current_account_with_tenant", + lambda: (SimpleNamespace(id="account-1"), "tenant-1"), + ) + monkeypatch.setattr(common_wraps, "enforce_rbac_access", permission_gate) + monkeypatch.setattr(console_resources, "_actor", lambda: ("account-1", "tenant-1")) + monkeypatch.setattr(console_resources, "_console_services", lambda: SimpleNamespace(facade=facade)) + monkeypatch.setattr(console_resources, "KnowledgeFSDownloadService", lambda: download_service) + permission_wrapper = _rbac_wrapper(console_resources.KnowledgeFSSpaceLogicalDocumentsDownloadApi.post) + app = Flask(__name__) + + with ( + app.test_request_context(json={"document_ids": ["document-1"]}), + pytest.raises(ServiceUnavailable), + ): + permission_wrapper( + console_resources.KnowledgeFSSpaceLogicalDocumentsDownloadApi(), + control_space_id="control-1", + ) + + permission_gate.assert_called_once() + download_service.build_zip_tempfile.assert_called_once_with([descriptor]) + + def test_knowledge_fs_request_and_response_schemas_are_registered() -> None: assert { "KnowledgeFSSpaceCreatePayload", diff --git a/api/tests/unit_tests/controllers/test_knowledge_fs_resource_delegation.py b/api/tests/unit_tests/controllers/test_knowledge_fs_resource_delegation.py index 8ed9993cb80..a5ca4f60b89 100644 --- a/api/tests/unit_tests/controllers/test_knowledge_fs_resource_delegation.py +++ b/api/tests/unit_tests/controllers/test_knowledge_fs_resource_delegation.py @@ -1348,6 +1348,8 @@ def test_console_error_adapter_maps_every_domain_boundary_to_the_stable_http_con with pytest.raises(http_error): fail() + assert KnowledgeFSUpstreamUnavailableHTTPError.code == HTTPStatus.SERVICE_UNAVAILABLE + def test_service_error_adapter_maps_every_domain_boundary_to_the_stable_http_contract() -> None: from pydantic import ValidationError diff --git a/api/tests/unit_tests/services/test_knowledge_fs_download_service.py b/api/tests/unit_tests/services/test_knowledge_fs_download_service.py index bf47d8c0469..c4e81ef4999 100644 --- a/api/tests/unit_tests/services/test_knowledge_fs_download_service.py +++ b/api/tests/unit_tests/services/test_knowledge_fs_download_service.py @@ -11,6 +11,11 @@ from services.knowledge_fs.download_service import ( KnowledgeFSDownloadObjectNotFoundError, KnowledgeFSDownloadService, KnowledgeFSDownloadTooLargeError, + KnowledgeFSDownloadUnavailableError, +) +from services.knowledge_fs.object_storage import ( + KnowledgeFSObjectStorageCorruptError, + KnowledgeFSObjectStorageUnavailableError, ) from services.knowledge_fs.product_dto import KnowledgeFSDocumentDownloadDescriptor @@ -56,6 +61,32 @@ def test_load_stream_rejects_missing_or_changed_object() -> None: service.load_stream(descriptor(document_id="document-1", filename="a.txt", object_key="object-1", size_bytes=5)) +@pytest.mark.parametrize( + ("storage_error", "expected_error"), + [ + ( + KnowledgeFSObjectStorageCorruptError("object body is missing"), + KnowledgeFSDownloadObjectNotFoundError, + ), + ( + KnowledgeFSObjectStorageUnavailableError("storage is unavailable"), + KnowledgeFSDownloadUnavailableError, + ), + ], +) +def test_load_stream_translates_object_storage_errors( + storage_error: Exception, + expected_error: type[Exception], +) -> None: + object_storage = SimpleNamespace( + head_object=lambda **_: (_ for _ in ()).throw(storage_error), + ) + service = KnowledgeFSDownloadService(object_storage=object_storage) + + with pytest.raises(expected_error): + service.load_stream(descriptor(document_id="document-1", filename="a.txt", object_key="object-1", size_bytes=4)) + + def test_build_zip_streams_objects_and_deduplicates_names() -> None: service = KnowledgeFSDownloadService( object_storage=FakeObjectStorage({"object-1": b"first", "object-2": b"second"}) diff --git a/web/features/new-rag/__tests__/documents-page.spec.tsx b/web/features/new-rag/__tests__/documents-page.spec.tsx index 5298ded0dcb..3850bc8a4f4 100644 --- a/web/features/new-rag/__tests__/documents-page.spec.tsx +++ b/web/features/new-rag/__tests__/documents-page.spec.tsx @@ -3164,6 +3164,63 @@ describe('DocumentsPage', () => { }) }) + it('disables bulk download when any selected document has no active revision', async () => { + const user = userEvent.setup() + documentsQuery.data = { + pages: [ + { + items: [ + document({ id: 'ready', title: 'Ready.pdf' }), + document({ + active: null, + activeRevision: undefined, + id: 'pending', + title: 'Pending.pdf', + }), + ], + }, + ], + } + + render() + await user.click(screen.getByRole('checkbox', { name: 'Ready.pdf' })) + await user.click(screen.getByRole('checkbox', { name: 'Pending.pdf' })) + + const actions = screen.getByRole('group', { + name: 'dataset.newKnowledge.bulkDocumentActions', + }) + expect( + within(actions).getByRole('button', { name: 'dataset.newKnowledge.downloadDocuments' }), + ).toBeDisabled() + expect(downloadDocumentsMutation).not.toHaveBeenCalled() + }) + + it('disables bulk download when more than 100 documents are selected', async () => { + const user = userEvent.setup() + documentsQuery.data = { + pages: [ + { + items: Array.from({ length: 101 }, (_, index) => + document({ id: `document-${index}`, title: `Document ${index}.pdf` }), + ), + }, + ], + } + + render() + await user.click( + screen.getByRole('checkbox', { name: 'dataset.newKnowledge.selectAllDocuments' }), + ) + + const actions = screen.getByRole('group', { + name: 'dataset.newKnowledge.bulkDocumentActions', + }) + expect( + within(actions).getByRole('button', { name: 'dataset.newKnowledge.downloadDocuments' }), + ).toBeDisabled() + expect(downloadDocumentsMutation).not.toHaveBeenCalled() + }) + it('removes selected documents through one bulk deletion request', async () => { const user = userEvent.setup() documentsQuery.data = { diff --git a/web/features/new-rag/documents-page.tsx b/web/features/new-rag/documents-page.tsx index 6e588a93167..e0917c5d084 100644 --- a/web/features/new-rag/documents-page.tsx +++ b/web/features/new-rag/documents-page.tsx @@ -73,6 +73,7 @@ import { useKnowledgeModelSetupGuard } from './use-knowledge-model-setup-guard' import { useQueryDataUpdateCount } from './use-query-data-update-count' const TASK_PAGE_SIZE = 100 +const KNOWLEDGE_FS_BATCH_DOWNLOAD_MAX_DOCUMENTS = 100 const MAX_TASK_EVENT_STREAMS = 6 const MAX_AUTO_CURSOR_PAGES = 20 const FAILED_TASK_POLL_REQUEST_TIMEOUT = 3000 @@ -824,13 +825,18 @@ export function DocumentsPage({ knowledgeSpaceId }: { knowledgeSpaceId: string } ), [availableDocumentIds, selectedDocumentIds], ) - const downloadableSelectedDocumentIds = useMemo( - () => - documents.flatMap((document) => - validSelectedDocumentIds.has(document.id) && document.active ? [document.id] : [], - ), - [documents, validSelectedDocumentIds], - ) + const downloadableSelectedDocumentIds = useMemo(() => { + const selectedDocuments = documents.filter((document) => + validSelectedDocumentIds.has(document.id), + ) + if ( + selectedDocuments.length !== validSelectedDocumentIds.size || + selectedDocuments.length > KNOWLEDGE_FS_BATCH_DOWNLOAD_MAX_DOCUMENTS || + selectedDocuments.some((document) => !document.active) + ) + return [] + return selectedDocuments.map((document) => document.id) + }, [documents, validSelectedDocumentIds]) const filteredDocuments = useMemo(() => { const normalizedSearch = search.trim().toLocaleLowerCase() return documents.filter((document) => {