mirror of
https://github.com/langgenius/dify.git
synced 2026-08-24 12:32:54 -04:00
fix(knowledge-fs): align document download contracts
This commit is contained in:
@@ -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):
|
||||
|
||||
@@ -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/<string:control_space_id>/logical-documents/<string:document_id>/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='')}"
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"})
|
||||
|
||||
@@ -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(<DocumentsPage knowledgeSpaceId="space-1" />)
|
||||
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(<DocumentsPage knowledgeSpaceId="space-1" />)
|
||||
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 = {
|
||||
|
||||
@@ -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) => {
|
||||
|
||||
Reference in New Issue
Block a user