mirror of
https://github.com/langgenius/dify.git
synced 2026-09-05 00:31:19 +08:00
refactor(api): derive cleanup tenant from dataset
This commit is contained in:
parent
a07e5c6eb4
commit
b124ec86fc
@ -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)
|
||||
|
||||
@ -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
|
||||
|
||||
@ -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(
|
||||
|
||||
@ -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(
|
||||
|
||||
@ -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(
|
||||
|
||||
@ -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):
|
||||
|
||||
@ -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")
|
||||
|
||||
@ -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()
|
||||
|
||||
@ -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()
|
||||
|
||||
@ -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
|
||||
|
||||
Loading…
Reference in New Issue
Block a user