From b124ec86fc2f254de9bbfc90754cfdf925873f87 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=9E=97=E7=8E=AE=20=28Jade=20Lin=29?= Date: Wed, 22 Jul 2026 14:02:21 +0800 Subject: [PATCH] refactor(api): derive cleanup tenant from dataset --- .../event_handlers/clean_when_document_deleted.py | 3 +-- api/services/dataset_service.py | 2 -- api/tasks/batch_clean_document_task.py | 10 +++++----- api/tasks/clean_dataset_task.py | 2 +- api/tasks/clean_document_task.py | 9 +++------ .../services/test_dataset_service_document.py | 2 -- .../event_handlers/test_clean_when_document_deleted.py | 5 ++--- .../services/test_dataset_service_document.py | 2 +- .../unit_tests/tasks/test_batch_clean_document_task.py | 2 -- api/tests/unit_tests/tasks/test_clean_document_task.py | 3 --- 10 files changed, 13 insertions(+), 27 deletions(-) diff --git a/api/events/event_handlers/clean_when_document_deleted.py b/api/events/event_handlers/clean_when_document_deleted.py index f8bd24061a1..0add109b06d 100644 --- a/api/events/event_handlers/clean_when_document_deleted.py +++ b/api/events/event_handlers/clean_when_document_deleted.py @@ -8,7 +8,6 @@ def handle(sender, **kwargs): dataset_id = kwargs.get("dataset_id") doc_form = kwargs.get("doc_form") file_id = kwargs.get("file_id") - tenant_id = kwargs.get("tenant_id") if not dataset_id or not doc_form: return - clean_document_task.delay(document_id, dataset_id, doc_form, file_id, tenant_id) + clean_document_task.delay(document_id, dataset_id, doc_form, file_id) diff --git a/api/services/dataset_service.py b/api/services/dataset_service.py index d652c908886..27650e68101 100644 --- a/api/services/dataset_service.py +++ b/api/services/dataset_service.py @@ -1978,7 +1978,6 @@ class DocumentService: dataset_id=document.dataset_id, doc_form=document.doc_form, file_id=file_id, - tenant_id=document.tenant_id, ) session.delete(document) @@ -2022,7 +2021,6 @@ class DocumentService: dataset_ref.dataset_id, doc_form, file_ids, - dataset_ref.tenant_id, ) @staticmethod diff --git a/api/tasks/batch_clean_document_task.py b/api/tasks/batch_clean_document_task.py index dddd9119fac..ff3cf803492 100644 --- a/api/tasks/batch_clean_document_task.py +++ b/api/tasks/batch_clean_document_task.py @@ -27,7 +27,6 @@ def batch_clean_document_task( dataset_id: str, doc_form: str | None, file_ids: list[str], - tenant_id: str | None = None, ) -> None: """ Clean document when document deleted. @@ -35,7 +34,6 @@ def batch_clean_document_task( :param dataset_id: dataset id :param doc_form: doc_form :param file_ids: file ids - :param tenant_id: tenant id Usage: batch_clean_document_task.delay(document_ids, dataset_id) """ @@ -49,6 +47,7 @@ def batch_clean_document_task( segment_ids: list[str] = [] total_image_upload_file_ids: list[str] = [] vector_cleanup_succeeded = False + dataset_tenant_id: str | None = None try: # ============ Step 1: Query segment and file data (short read-only transaction) ============ @@ -88,7 +87,6 @@ def batch_clean_document_task( if not dataset: logger.warning("Dataset not found for vector index cleanup, dataset_id: %s", dataset_id) else: - tenant_id = tenant_id or dataset.tenant_id index_processor = IndexProcessorFactory(doc_form).init_index_processor() index_processor.clean( dataset, @@ -98,6 +96,7 @@ def batch_clean_document_task( delete_summaries=True, session=session, ) + dataset_tenant_id = dataset.tenant_id vector_cleanup_succeeded = True except Exception: logger.exception( @@ -214,8 +213,9 @@ def batch_clean_document_task( dataset_id, ) - if vector_cleanup_succeeded and tenant_id: - schedule_billing_vector_space_refresh(tenant_id) + if vector_cleanup_succeeded: + assert dataset_tenant_id is not None + schedule_billing_vector_space_refresh(dataset_tenant_id) end_at = time.perf_counter() logger.info( diff --git a/api/tasks/clean_dataset_task.py b/api/tasks/clean_dataset_task.py index 839633459d6..5bf8784e3c2 100644 --- a/api/tasks/clean_dataset_task.py +++ b/api/tasks/clean_dataset_task.py @@ -190,7 +190,7 @@ def clean_dataset_task( session.commit() if vector_cleanup_succeeded: - schedule_billing_vector_space_refresh(tenant_id) + schedule_billing_vector_space_refresh(dataset.tenant_id) end_at = time.perf_counter() logger.info( click.style( diff --git a/api/tasks/clean_document_task.py b/api/tasks/clean_document_task.py index 2d08658e64c..e09743a0018 100644 --- a/api/tasks/clean_document_task.py +++ b/api/tasks/clean_document_task.py @@ -22,7 +22,6 @@ def clean_document_task( dataset_id: str, doc_form: str, file_id: str | None, - tenant_id: str | None = None, ) -> None: """ Clean document when document deleted. @@ -30,7 +29,6 @@ def clean_document_task( :param dataset_id: dataset id :param doc_form: doc_form :param file_id: file id - :param tenant_id: tenant id Usage: clean_document_task.delay(document_id, dataset_id) """ @@ -46,8 +44,7 @@ def clean_document_task( if not dataset: raise Exception("Document has no dataset") - tenant_id = tenant_id or dataset.tenant_id - + dataset_tenant_id = dataset.tenant_id segments = session.scalars(select(DocumentSegment).where(DocumentSegment.document_id == document_id)).all() # Use JOIN to fetch attachments with bindings in a single query attachments_with_bindings = session.execute( @@ -166,8 +163,8 @@ def clean_document_task( ) ) - if vector_cleanup_succeeded and tenant_id: - schedule_billing_vector_space_refresh(tenant_id) + if vector_cleanup_succeeded: + schedule_billing_vector_space_refresh(dataset_tenant_id) end_at = time.perf_counter() logger.info( diff --git a/api/tests/test_containers_integration_tests/services/test_dataset_service_document.py b/api/tests/test_containers_integration_tests/services/test_dataset_service_document.py index df7242f2a79..e722f943820 100644 --- a/api/tests/test_containers_integration_tests/services/test_dataset_service_document.py +++ b/api/tests/test_containers_integration_tests/services/test_dataset_service_document.py @@ -611,7 +611,6 @@ def test_delete_document_emits_signal_and_commits(db_session_with_containers: Se dataset_id=document.dataset_id, doc_form=document.doc_form, file_id=upload_file.id, - tenant_id=document.tenant_id, ) @@ -669,7 +668,6 @@ def test_delete_documents_deletes_rows_and_dispatches_cleanup_task(db_session_wi assert args[0] == [document_a.id, document_b.id] assert args[1] == dataset.id assert set(args[3]) == {upload_file_a.id, upload_file_b.id} - assert args[4] == dataset.tenant_id def test_get_documents_position_returns_next_position_when_documents_exist(db_session_with_containers: Session): diff --git a/api/tests/unit_tests/events/event_handlers/test_clean_when_document_deleted.py b/api/tests/unit_tests/events/event_handlers/test_clean_when_document_deleted.py index 3b252c7479a..098a7b67880 100644 --- a/api/tests/unit_tests/events/event_handlers/test_clean_when_document_deleted.py +++ b/api/tests/unit_tests/events/event_handlers/test_clean_when_document_deleted.py @@ -3,14 +3,13 @@ from unittest.mock import patch from events.event_handlers.clean_when_document_deleted import handle -def test_handler_passes_tenant_id_to_cleanup_task(): +def test_handler_dispatches_cleanup_task(): with patch("events.event_handlers.clean_when_document_deleted.clean_document_task.delay") as delay: handle( "document-1", dataset_id="dataset-1", doc_form="paragraph", file_id="file-1", - tenant_id="tenant-1", ) - delay.assert_called_once_with("document-1", "dataset-1", "paragraph", "file-1", "tenant-1") + delay.assert_called_once_with("document-1", "dataset-1", "paragraph", "file-1") diff --git a/api/tests/unit_tests/services/test_dataset_service_document.py b/api/tests/unit_tests/services/test_dataset_service_document.py index e3cce45779a..26a3ac08d5f 100644 --- a/api/tests/unit_tests/services/test_dataset_service_document.py +++ b/api/tests/unit_tests/services/test_dataset_service_document.py @@ -149,7 +149,7 @@ class TestDocumentServiceMutations: assert dataset.id in compiled.params.values() session.delete.assert_called_once_with(document) session.commit.assert_called_once() - clean_task.delay.assert_called_once_with(["doc-1"], dataset.id, dataset.doc_form, [], dataset.tenant_id) + clean_task.delay.assert_called_once_with(["doc-1"], dataset.id, dataset.doc_form, []) def test_rename_document_raises_when_dataset_is_missing(self, rename_account_context): session = MagicMock() diff --git a/api/tests/unit_tests/tasks/test_batch_clean_document_task.py b/api/tests/unit_tests/tasks/test_batch_clean_document_task.py index cd17ad0f911..6386c72188e 100644 --- a/api/tests/unit_tests/tasks/test_batch_clean_document_task.py +++ b/api/tests/unit_tests/tasks/test_batch_clean_document_task.py @@ -30,7 +30,6 @@ def test_successful_vector_cleanup_schedules_billing_refresh(): dataset_id="dataset-1", doc_form="paragraph", file_ids=[], - tenant_id="tenant-1", ) processor_factory.return_value.init_index_processor.return_value.clean.assert_called_once() @@ -54,7 +53,6 @@ def test_failed_vector_cleanup_does_not_schedule_billing_refresh(): dataset_id="dataset-1", doc_form="paragraph", file_ids=[], - tenant_id="tenant-1", ) schedule_refresh.assert_not_called() diff --git a/api/tests/unit_tests/tasks/test_clean_document_task.py b/api/tests/unit_tests/tasks/test_clean_document_task.py index 3e4919635b5..2f517ce1ba4 100644 --- a/api/tests/unit_tests/tasks/test_clean_document_task.py +++ b/api/tests/unit_tests/tasks/test_clean_document_task.py @@ -175,7 +175,6 @@ class TestVectorCleanupResilience: dataset_id=dataset_id, doc_form="paragraph", file_id=None, - tenant_id=tenant_id, ) # Assert @@ -238,7 +237,6 @@ class TestVectorCleanupResilience: dataset_id=dataset_id, doc_form="paragraph", file_id=None, - tenant_id=tenant_id, ) assert mock_index_processor_factory["processor"].clean.call_count == 1 @@ -291,7 +289,6 @@ class TestVectorCleanupResilience: dataset_id=dataset_id, doc_form="paragraph", file_id=None, - tenant_id=tenant_id, ) # Vector cleanup is gated on ``index_node_ids``; when there are no