From a072beb563ab8835ebb7b5b6ce92b53ffa89cb8e Mon Sep 17 00:00:00 2001 From: cj2026-bit <647646783@qq.com> Date: Wed, 26 Aug 2026 17:16:51 +0800 Subject: [PATCH 1/5] fix: isolate KB ingestion tasks and prevent delete resurrection --- backend/consts/const.py | 11 +- backend/data_process/app.py | 9 +- backend/data_process/tasks.py | 163 ++++++++++++- backend/data_process_service.py | 85 ++++--- backend/services/data_process_service.py | 34 +++ backend/services/file_management_service.py | 19 ++ backend/services/redis_service.py | 230 ++++++++++++++++++ backend/services/vectordatabase_service.py | 47 +++- test/backend/data_process/test_tasks.py | 134 +++++++++- .../services/test_data_process_service.py | 37 +++ test/backend/services/test_redis_service.py | 56 +++++ .../test_data_process_service_entrypoint.py | 15 ++ 12 files changed, 793 insertions(+), 47 deletions(-) diff --git a/backend/consts/const.py b/backend/consts/const.py index 7500a83150..715f824be3 100644 --- a/backend/consts/const.py +++ b/backend/consts/const.py @@ -317,10 +317,17 @@ class VectorDatabaseType(str, Enum): # Worker Configuration RAY_ADDRESS = os.getenv("RAY_ADDRESS", "auto") -QUEUES = os.getenv("QUEUES", "process_q,process_part_q,forward_q") +QUEUES = os.getenv( + "QUEUES", + "process_q,process_part_q,forward_q,forward_part_q,forward_aggregate_q", +) # Will be dynamically set based on PID if not provided WORKER_NAME = os.getenv("WORKER_NAME") -WORKER_CONCURRENCY = DP_PART_PROCESSOR_COUNT + 1 +# The data-process service sets a queue-specific value for each child worker. +# Keep the historical default when the variable is not provided. +WORKER_CONCURRENCY = int( + os.getenv("WORKER_CONCURRENCY", str(DP_PART_PROCESSOR_COUNT + 1)) +) RAY_WARM_ACTOR_POOL_SIZE_PART = int( os.getenv("RAY_WARM_ACTOR_POOL_SIZE_PART", "2")) RAY_WARM_ACTOR_POOL_SIZE_PROCESS = int( diff --git a/backend/data_process/app.py b/backend/data_process/app.py index c3346c0882..1d4e025d41 100644 --- a/backend/data_process/app.py +++ b/backend/data_process/app.py @@ -44,11 +44,18 @@ # Explicitly set result backend broker_url=REDIS_URL, result_backend=REDIS_BACKEND_URL, - # Two task queues for processing and forward steps + # Dedicated queues keep parent, part, and aggregate tasks from starving + # one another when a parent task waits for a Celery chord. task_routes={ f'{import_path}.process': {'queue': 'process_q'}, f'{import_path}.forward': {'queue': 'forward_q'}, f'{import_path}.process_and_forward': {'queue': 'process_q'}, + f'{import_path}.process_part': {'queue': 'process_part_q'}, + f'{import_path}.aggregate_parts': {'queue': 'process_part_q'}, + f'{import_path}.aggregate_store_chunks': {'queue': 'process_part_q'}, + f'{import_path}.forward_part': {'queue': 'forward_part_q'}, + f'{import_path}.aggregate_forward_parts': {'queue': 'forward_aggregate_q'}, + f'{import_path}.cleanup_source': {'queue': 'forward_q'}, }, task_serializer='json', accept_content=['json'], diff --git a/backend/data_process/tasks.py b/backend/data_process/tasks.py index 9135b1e263..4a9f539cf1 100644 --- a/backend/data_process/tasks.py +++ b/backend/data_process/tasks.py @@ -519,6 +519,35 @@ def _is_forward_task_cancelled(ctx: _ForwardContext) -> bool: return False +def _is_document_delete_requested( + index_name: Optional[str], + source: Optional[str], + file_id: Optional[str], +) -> bool: + """Check the durable file-level deletion fence used by delete + sync.""" + if not index_name: + return False + try: + return bool( + get_redis_service().is_document_delete_requested( + index_name=index_name, + path_or_url=source, + file_id=file_id, + ) + ) + except Exception as exc: + # Deletion fencing is best effort when Redis is unavailable. The task + # still follows the existing task-cancellation/lifecycle safeguards. + logger.warning( + "Failed to check document deletion fence for index=%s source=%s file_id=%s: %s", + index_name, + source, + file_id, + exc, + ) + return False + + def _build_forward_cancelled_result(ctx: _ForwardContext) -> Dict[str, Any]: return { 'task_id': ctx.task_id, @@ -1088,7 +1117,7 @@ def aggregate_store_chunks( } -@app.task(bind=True, base=LoggingTask, name='data_process.tasks.forward_part', queue='forward_q') +@app.task(bind=True, base=LoggingTask, name='data_process.tasks.forward_part', queue='forward_part_q') @trace_knowledge_operation("knowledge.forward.batch", "forward.batch") def forward_part( self, @@ -1102,11 +1131,30 @@ def forward_part( batch_index: Optional[int] = None, total_batches: Optional[int] = None, large_mode: Optional[bool] = False, + file_id: Optional[str] = None, ) -> Dict[str, Any]: """ Forward sub-task that indexes a chunk batch. """ try: + if _is_document_delete_requested(index_name, source, file_id): + logger.info( + "Skipping forward batch %s/%s because document was deleted: index=%s source=%s file_id=%s", + batch_index, + total_batches, + index_name, + source, + file_id, + ) + return { + "success": True, + "total_indexed": 0, + "total_submitted": 0, + "batch_index": batch_index, + "total_batches": total_batches, + "cancelled": True, + } + # Respect cancellation from parent task if available if parent_task_id: try: @@ -1210,18 +1258,42 @@ def forward_part( ) -@app.task(bind=True, base=LoggingTask, name='data_process.tasks.aggregate_forward_parts', queue='forward_q') +@app.task( + bind=True, + base=LoggingTask, + name='data_process.tasks.aggregate_forward_parts', + queue='forward_aggregate_q', +) @trace_knowledge_operation("knowledge.forward.aggregate", "forward.aggregate") def aggregate_forward_parts( self, parts_results: List[Dict[str, Any]], source: Optional[str] = None, index_name: Optional[str] = None, - original_filename: Optional[str] = None + original_filename: Optional[str] = None, + file_id: Optional[str] = None, ) -> Dict[str, Any]: """ Aggregate forward_part results. """ + if _is_document_delete_requested(index_name, source, file_id): + logger.info( + "Skipping forward aggregate because document was deleted: index=%s source=%s file_id=%s", + index_name, + source, + file_id, + ) + return { + "success": True, + "total_indexed": 0, + "total_submitted": 0, + "source": source, + "index_name": index_name, + "original_filename": original_filename, + "file_id": file_id, + "cancelled": True, + } + total_indexed = 0 total_submitted = 0 for result in parts_results or []: @@ -1558,6 +1630,21 @@ def process( start_time = time.time() task_id = self.request.id file_id = params.get("file_id") + if _is_document_delete_requested(index_name, source, file_id): + logger.info( + "Skipping process task because document was deleted: index=%s source=%s file_id=%s", + index_name, + source, + file_id, + ) + return { + "cancelled": True, + "source": source, + "index_name": index_name, + "original_filename": original_filename, + "file_id": file_id, + "task_id": task_id, + } _update_file_lifecycle( file_id=file_id, tenant_id=tenant_id, @@ -1617,6 +1704,22 @@ def process( tenant_id=tenant_id, params=params, ) + if _is_document_delete_requested(index_name, source, file_id): + logger.info( + "Stopping process task after source extraction because document was deleted: " + "index=%s source=%s file_id=%s", + index_name, + source, + file_id, + ) + return { + "cancelled": True, + "source": source, + "index_name": index_name, + "original_filename": original_filename, + "file_id": file_id, + "task_id": task_id, + } elapsed_time = time.time() - start_time processing_speed = file_size_mb / \ elapsed_time if file_size_mb > 0 and elapsed_time > 0 else 0 @@ -1657,6 +1760,23 @@ def process( raise NotImplementedError( f"Source type '{source_type}' not yet supported") + if _is_document_delete_requested(index_name, source, file_id): + logger.info( + "Stopping process task after extraction because document was deleted: " + "index=%s source=%s file_id=%s", + index_name, + source, + file_id, + ) + return { + "cancelled": True, + "source": source, + "index_name": index_name, + "original_filename": original_filename, + "file_id": file_id, + "task_id": task_id, + } + if split_async: chunk_count = split_chunk_count or 0 if chunk_count == 0: @@ -1933,6 +2053,13 @@ def forward( ) return _build_forward_cancelled_result(ctx) + if _is_document_delete_requested(index_name, source, file_id): + logger.info( + f"[{self.request.id}] FORWARD TASK: Document deletion fence is set; " + f"skipping chunk forwarding for source '{source}' in index '{index_name}'." + ) + return _build_forward_cancelled_result(ctx) + chunks, split_async, original_source, original_index_name, filename = _load_forward_chunks( self, processed_data=processed_data, @@ -1941,6 +2068,13 @@ def forward( filename=filename, ) + if _is_document_delete_requested(original_index_name, original_source, file_id): + logger.info( + f"[{self.request.id}] FORWARD TASK: Document was deleted while loading chunks; " + f"skipping ES indexing for source '{original_source}'." + ) + return _build_forward_cancelled_result(ctx) + # Calculate total chunks for progress tracking total_chunks = len(chunks) if chunks else 0 set_span_attributes(chunk_count=total_chunks, stage="forward.format") @@ -2005,6 +2139,15 @@ def forward( } ) + # Re-check immediately before the first ES write. A delete may race + # with formatting or progress initialization after the earlier check. + if _is_document_delete_requested(original_index_name, original_source, file_id): + logger.info( + f"[{self.request.id}] FORWARD TASK: Document was deleted before ES indexing; " + f"skipping source '{original_source}'." + ) + return _build_forward_cancelled_result(ctx) + try: redis_service = get_redis_service() redis_service.save_progress_info(task_id, 0, total_chunks) @@ -2054,17 +2197,19 @@ def forward( parent_total_chunks=total_chunks, source=original_source, original_filename=original_filename, + file_id=file_id, batch_index=idx + 1, total_batches=total_batches, # If request was split into multiple groups, force all groups to use large path. large_mode=True, - ).set(queue='forward_q') for idx, batch in enumerate(batches) + ).set(queue='forward_part_q') for idx, batch in enumerate(batches) ) callback = aggregate_forward_parts.s( source=original_source, index_name=original_index_name, - original_filename=original_filename - ).set(queue='forward_q') + original_filename=original_filename, + file_id=file_id, + ).set(queue='forward_aggregate_q') result = chord(group_tasks)(callback) with allow_join_result(): es_result = result.get() @@ -2276,6 +2421,12 @@ def cleanup_source( "error": None, } + if (forward_result or {}).get("cancelled"): + cleanup_info["skipped_reason"] = "forward_cancelled" + forward_result = dict(forward_result or {}) + forward_result["source_cleanup"] = cleanup_info + return forward_result + if not index_name or not source: cleanup_info["skipped_reason"] = "missing_index_name_or_source" forward_result = dict(forward_result or {}) diff --git a/backend/data_process_service.py b/backend/data_process_service.py index 1d955ebf80..961b10d8e2 100644 --- a/backend/data_process_service.py +++ b/backend/data_process_service.py @@ -179,6 +179,45 @@ def start_ray_cluster(self): logger.error(traceback.format_exc()) return False + @staticmethod + def _build_worker_configs(total_cpus: int) -> list[dict[str, Any]]: + """Build isolated Celery worker pools for each processing stage.""" + total_cpus = max(1, int(total_cpus)) + ray_actor_num_cpus = max(1, int(RAY_ACTOR_NUM_CPUS)) + process_worker_concurrency = min( + DP_PART_PROCESSOR_COUNT, + max(1, total_cpus // ray_actor_num_cpus), + ) + forward_worker_concurrency = min(8, total_cpus * 2) + forward_aggregate_worker_concurrency = min(2, total_cpus) + return [ + { + 'name': 'process-worker', + 'queue': 'process_q', + 'concurrency': process_worker_concurrency, + }, + { + 'name': 'process-part-worker', + 'queue': 'process_part_q', + 'concurrency': process_worker_concurrency, + }, + { + 'name': 'forward-worker', + 'queue': 'forward_q', + 'concurrency': forward_worker_concurrency, + }, + { + 'name': 'forward-part-worker', + 'queue': 'forward_part_q', + 'concurrency': forward_worker_concurrency, + }, + { + 'name': 'forward-aggregate-worker', + 'queue': 'forward_aggregate_q', + 'concurrency': forward_aggregate_worker_concurrency, + }, + ] + def start_workers(self): """Start Celery workers for process and forward queues""" if not self.config.get('start_workers', True): @@ -194,47 +233,23 @@ def start_workers(self): # Fallback to 1 if os.cpu_count() is None. total_cpus = int(RAY_NUM_CPUS) if RAY_NUM_CPUS else (os.cpu_count() or 1) - # Get the number of CPUs requested by each actor. - ray_actor_num_cpus = RAY_ACTOR_NUM_CPUS - - # Calculate concurrency for the process-worker. Each worker will spawn an actor, - # so we limit concurrency to avoid oversubscribing Ray's CPU resources. - process_worker_concurrency = min( - DP_PART_PROCESSOR_COUNT, - max(1, total_cpus // ray_actor_num_cpus), - ) - - # For forward-worker, it's I/O bound. A higher concurrency is fine, but we can cap it - # relative to CPU count to avoid creating excessive threads on small machines. - forward_worker_concurrency = min(8, total_cpus * 2) + workers_config = self._build_worker_configs(total_cpus) + concurrency_by_name = { + config['name']: config['concurrency'] for config in workers_config + } + process_worker_concurrency = concurrency_by_name['process-worker'] + forward_worker_concurrency = concurrency_by_name['forward-worker'] + forward_aggregate_worker_concurrency = concurrency_by_name['forward-aggregate-worker'] + ray_actor_num_cpus = max(1, int(RAY_ACTOR_NUM_CPUS)) logger.debug(f"Total available CPUs: {total_cpus}") logger.debug(f"CPUs per processing actor (RAY_ACTOR_NUM_CPUS): {ray_actor_num_cpus}") logger.debug(f"Process-worker concurrency set to: {process_worker_concurrency}") logger.debug(f"Forward-worker concurrency set to: {forward_worker_concurrency}") + logger.debug( + f"Forward-aggregate-worker concurrency set to: {forward_aggregate_worker_concurrency}" + ) - # Define worker configurations based on split architecture: - # - process-worker handles orchestration (process_q) - # - process-part-worker handles split sub-tasks (process_part_q) - # - forward-worker handles vectorization/storage (forward_q) - workers_config = [ - { - 'name': 'process-worker', - 'queue': 'process_q', - 'concurrency': process_worker_concurrency - }, - { - 'name': 'process-part-worker', - 'queue': 'process_part_q', - 'concurrency': process_worker_concurrency - }, - { - 'name': 'forward-worker', - 'queue': 'forward_q', - 'concurrency': forward_worker_concurrency - } - ] - # Start each worker in a separate process for config in workers_config: # Use full Python path and correct module path diff --git a/backend/services/data_process_service.py b/backend/services/data_process_service.py index 8051802ba2..34148f0f21 100644 --- a/backend/services/data_process_service.py +++ b/backend/services/data_process_service.py @@ -27,6 +27,7 @@ from data_process.tasks import submit_process_forward_chain from data_process.utils import get_all_task_ids_from_redis, get_task_info from database.attachment_db import delete_file, file_exists, get_file_size_from_minio, get_file_stream, upload_file +from services.redis_service import get_redis_service from utils.file_management_utils import convert_office_to_pdf from utils.knowledge_ingestion_errors import classify_ingestion_exception @@ -146,6 +147,14 @@ async def get_all_tasks(self, filter: bool = True) -> List[Dict[str, Any]]: all_tasks = [] try: self._get_celery_inspector() + delete_fence_service = None + try: + delete_fence_service = get_redis_service() + except Exception as redis_init_exc: + logger.debug( + "Deletion-fence filtering unavailable while listing tasks: %s", + redis_init_exc, + ) # Collect task IDs from different sources and keep runtime metadata task_ids = set() @@ -264,6 +273,31 @@ def get_reserved(): if not task_info.get('file_id') and runtime_meta.get('file_id'): task_info['file_id'] = runtime_meta.get('file_id') + # A task can remain in Celery's result backend or broker after + # the file row and external document were deleted. Hide it + # from the user-visible task/file list; the worker-side fence + # also prevents a late task from recreating the document. + if delete_fence_service: + try: + document_deleted = delete_fence_service.is_document_delete_requested( + index_name=task_info.get('index_name'), + path_or_url=task_info.get('path_or_url'), + file_id=task_info.get('file_id'), + ) + except Exception as fence_exc: + logger.debug( + "Unable to check deletion fence for task %s: %s", + task_id, + fence_exc, + ) + document_deleted = False + if document_deleted: + logger.info( + "Skipping task %s because its document deletion fence is set", + task_id, + ) + continue + if filter and not (task_info.get('index_name') and task_info.get('task_name')): # Keep user-visible queued tasks even before worker updates task meta. if task_info.get('task_name') not in {'process', 'forward', 'process_and_forward'}: diff --git a/backend/services/file_management_service.py b/backend/services/file_management_service.py index 76a37be9bb..38223c921d 100644 --- a/backend/services/file_management_service.py +++ b/backend/services/file_management_service.py @@ -450,6 +450,25 @@ async def upload_files_impl( }) if lifecycle_record_specs: + # A new upload may reuse an object path after a previous + # delete. Clear the old deletion fence before workers see the + # new lifecycle row; the file-id fence remains unique to the + # old record. + try: + from services.redis_service import get_redis_service + + redis_service = get_redis_service() + for spec in lifecycle_record_specs: + redis_service.clear_document_delete_marker( + index_name=spec["index_name"], + path_or_url=spec["object_name"], + ) + except Exception as deletion_fence_exc: + logger.warning( + "Failed to clear an old document deletion fence before upload: %s", + deletion_fence_exc, + ) + # Lifecycle persistence is a required upload precondition. The # repository creates the whole batch in one transaction; any # database error must stop before MinIO is touched. diff --git a/backend/services/redis_service.py b/backend/services/redis_service.py index 15a41b6772..b452a0be20 100644 --- a/backend/services/redis_service.py +++ b/backend/services/redis_service.py @@ -1,6 +1,8 @@ +import hashlib import json import logging import re +import time from typing import Any, Dict, List, Optional, Set, Tuple import redis @@ -19,9 +21,12 @@ class RedisService: """Redis service for managing cache and task data""" + DOCUMENT_DELETE_MARKER_TTL_SECONDS = 24 * 60 * 60 + def __init__(self): self._client = None self._backend_client = None + self._celery_control_app = None @property def client(self) -> redis.Redis: @@ -51,6 +56,231 @@ def backend_client(self) -> redis.Redis: # Cancellation helpers # ------------------------------------------------------------------ + @staticmethod + def _document_delete_marker_keys( + index_name: Optional[str] = None, + path_or_url: Optional[str] = None, + file_id: Optional[str] = None, + ) -> List[str]: + """Return stable Redis keys for a file deletion fence.""" + keys = [] + if file_id: + keys.append(f"kb-delete:file:{file_id}") + if index_name and path_or_url: + identity = f"{index_name}\0{path_or_url}".encode("utf-8") + digest = hashlib.sha256(identity).hexdigest() + keys.append(f"kb-delete:path:{digest}") + return keys + + def mark_document_delete_requested( + self, + index_name: str, + path_or_url: Optional[str] = None, + file_id: Optional[str] = None, + ttl_seconds: Optional[int] = None, + ) -> bool: + """Write a short-lived file-level deletion fence. + + The fence is deliberately separate from per-task cancellation keys. It + survives result/cache cleanup so queued broker messages cannot recreate + a document after its lifecycle row has been hard deleted. + """ + keys = self._document_delete_marker_keys(index_name, path_or_url, file_id) + if not keys: + return False + ttl = int(ttl_seconds or self.DOCUMENT_DELETE_MARKER_TTL_SECONDS) + if ttl <= 0: + return False + payload = json.dumps({ + "index_name": index_name, + "path_or_url": path_or_url, + "file_id": file_id, + "requested_at": time.time(), + }, ensure_ascii=False) + try: + for key in keys: + self.client.setex(key, ttl, payload) + logger.info( + "Marked document deletion fence: index=%s path=%s file_id=%s ttl=%ss", + index_name, + path_or_url, + file_id, + ttl, + ) + return True + except Exception as exc: + logger.warning( + "Failed to mark document deletion fence for index=%s path=%s file_id=%s: %s", + index_name, + path_or_url, + file_id, + exc, + ) + return False + + def clear_document_delete_marker( + self, + index_name: str, + path_or_url: Optional[str] = None, + file_id: Optional[str] = None, + ) -> int: + """Clear a path/file deletion fence before a new upload reuses it.""" + keys = self._document_delete_marker_keys(index_name, path_or_url, file_id) + if not keys: + return 0 + try: + return int(self.client.delete(*keys)) + except Exception as exc: + logger.warning( + "Failed to clear document deletion fence for index=%s path=%s file_id=%s: %s", + index_name, + path_or_url, + file_id, + exc, + ) + return 0 + + def is_document_delete_requested( + self, + index_name: Optional[str], + path_or_url: Optional[str] = None, + file_id: Optional[str] = None, + ) -> bool: + """Return whether a file is fenced from further processing.""" + keys = self._document_delete_marker_keys(index_name, path_or_url, file_id) + if not keys: + return False + try: + return any(bool(self.client.get(key)) for key in keys) + except Exception as exc: + # A Redis outage must preserve the legacy best-effort behavior. + logger.debug( + "Unable to check document deletion fence for index=%s path=%s file_id=%s: %s", + index_name, + path_or_url, + file_id, + exc, + ) + return False + + def _get_celery_control_app(self): + """Create a lightweight Celery control app without importing task modules.""" + if self._celery_control_app is None: + from celery import Celery + + broker_url = REDIS_URL + backend_url = REDIS_BACKEND_URL or REDIS_URL + if not broker_url: + raise ValueError("REDIS_URL is not configured") + self._celery_control_app = Celery( + "nexent-delete-control", + broker=broker_url, + backend=backend_url, + ) + return self._celery_control_app + + @staticmethod + def _task_kwargs(task: Dict[str, Any]) -> Dict[str, Any]: + kwargs = task.get("kwargs") or {} + if isinstance(kwargs, str): + try: + kwargs = json.loads(kwargs) + except (TypeError, json.JSONDecodeError): + kwargs = {} + return kwargs if isinstance(kwargs, dict) else {} + + @classmethod + def _runtime_task_matches_document( + cls, + task: Dict[str, Any], + index_name: str, + path_or_url: Optional[str], + file_id: Optional[str], + ) -> bool: + kwargs = cls._task_kwargs(task) + task_index = kwargs.get("index_name") + task_source = kwargs.get("source") or kwargs.get("path_or_url") + task_file_id = kwargs.get("file_id") + if file_id and task_file_id and task_file_id != file_id: + return False + if task_index != index_name: + return False + if file_id and task_file_id == file_id: + return True + return bool(path_or_url and task_source == path_or_url) + + def _collect_runtime_task_ids( + self, + index_name: str, + path_or_url: Optional[str], + file_id: Optional[str] = None, + ) -> Set[str]: + """Collect matching active/reserved task IDs for targeted revocation.""" + task_ids: Set[str] = set() + try: + inspector = self._get_celery_control_app().control.inspect(timeout=0.5) + for state_name in ("active", "reserved"): + snapshot = getattr(inspector, state_name)() or {} + for tasks in snapshot.values(): + for task in tasks or []: + if not isinstance(task, dict): + continue + task_id = task.get("id") + if task_id and self._runtime_task_matches_document( + task, index_name, path_or_url, file_id + ): + task_ids.add(str(task_id)) + except Exception as exc: + logger.warning( + "Failed to inspect runtime tasks for index=%s path=%s file_id=%s: %s", + index_name, + path_or_url, + file_id, + exc, + ) + return task_ids + + def _revoke_task(self, task_id: str) -> bool: + if not task_id: + return False + try: + self._get_celery_control_app().control.revoke( + task_id, + terminate=False, + ) + logger.info("Revoked Celery task %s", task_id) + return True + except Exception as exc: + logger.warning("Failed to revoke Celery task %s: %s", task_id, exc) + return False + + def prepare_document_deletion( + self, + index_name: str, + path_or_url: Optional[str] = None, + file_id: Optional[str] = None, + ) -> Dict[str, Any]: + """Fence a file and revoke discoverable runtime tasks before storage deletion.""" + result = { + "delete_fence_set": self.mark_document_delete_requested( + index_name=index_name, + path_or_url=path_or_url, + file_id=file_id, + ), + "runtime_tasks_found": 0, + "tasks_cancelled": 0, + "tasks_revoked": 0, + "warnings": [], + } + task_ids = self._collect_runtime_task_ids(index_name, path_or_url, file_id) + result["runtime_tasks_found"] = len(task_ids) + for task_id in task_ids: + if self.mark_task_cancelled(task_id): + result["tasks_cancelled"] += 1 + if self._revoke_task(task_id): + result["tasks_revoked"] += 1 + return result + def mark_task_cancelled(self, task_id: str, ttl_hours: int = 24) -> bool: """ Mark a Celery task as cancelled in Redis so that long-running diff --git a/backend/services/vectordatabase_service.py b/backend/services/vectordatabase_service.py index 4df5ec5795..7e68bdfc7c 100644 --- a/backend/services/vectordatabase_service.py +++ b/backend/services/vectordatabase_service.py @@ -2209,6 +2209,19 @@ def delete_lifecycle_record_without_object( tenant_id = lifecycle_record.get("tenant_id") index_name = lifecycle_record.get("index_name") + if index_name: + try: + get_redis_service().prepare_document_deletion( + index_name=index_name, + path_or_url=None, + file_id=file_id, + ) + except Exception as deletion_exc: + logger.warning( + "Failed to prepare deletion for lifecycle file %s: %s", + file_id, + deletion_exc, + ) current_status = str(lifecycle_record.get("status") or "").upper() deleteable_statuses = ( "UPLOADING", @@ -2294,7 +2307,22 @@ async def delete_document_by_scope( await ElasticSearchService._assert_source_only_deletable( index_name, path_or_url ) - ElasticSearchService._mark_file_delete_requested(index_name, path_or_url) + lifecycle_record = ElasticSearchService._mark_file_delete_requested( + index_name, path_or_url + ) + try: + get_redis_service().prepare_document_deletion( + index_name=index_name, + path_or_url=path_or_url, + file_id=(lifecycle_record or {}).get("file_id"), + ) + except Exception as deletion_exc: + logger.warning( + "Failed to prepare source-only deletion for index=%s path=%s: %s", + index_name, + path_or_url, + deletion_exc, + ) try: knowledge = get_knowledge_record({"index_name": index_name}) or {} except Exception: @@ -2322,7 +2350,22 @@ async def delete_document_by_scope( ), } - ElasticSearchService._mark_file_delete_requested(index_name, path_or_url) + lifecycle_record = ElasticSearchService._mark_file_delete_requested( + index_name, path_or_url + ) + try: + get_redis_service().prepare_document_deletion( + index_name=index_name, + path_or_url=path_or_url, + file_id=(lifecycle_record or {}).get("file_id"), + ) + except Exception as deletion_exc: + logger.warning( + "Failed to prepare full deletion for index=%s path=%s: %s", + index_name, + path_or_url, + deletion_exc, + ) result = ElasticSearchService.delete_documents( index_name, path_or_url, vdb_core ) diff --git a/test/backend/data_process/test_tasks.py b/test/backend/data_process/test_tasks.py index 89aeb8c818..b31323f613 100644 --- a/test/backend/data_process/test_tasks.py +++ b/test/backend/data_process/test_tasks.py @@ -2167,8 +2167,10 @@ def is_task_cancelled(self, *args, **kwargs): class _Sig: def __init__(self, kwargs): self.kwargs = kwargs + self.queue = None - def set(self, **_kw): + def set(self, **kw): + self.queue = kw.get("queue") return self captured = {"group_sigs": None} @@ -2217,6 +2219,81 @@ def _fake_allow_join_result(): assert len(captured["group_sigs"]) == 2 assert all(sig.kwargs.get("large_mode") is True for sig in captured["group_sigs"]) + assert all(sig.queue == "forward_part_q" for sig in captured["group_sigs"]) + + +def test_forward_large_chunks_routes_aggregate_to_dedicated_queue(monkeypatch): + tasks, _ = import_tasks_with_fake_ray(monkeypatch) + monkeypatch.setattr(tasks, "get_file_size", lambda *args, **kwargs: 0) + + class _RedisSvc: + def save_progress_info(self, *args, **kwargs): + return True + + def is_task_cancelled(self, *args, **kwargs): + return False + + def is_document_delete_requested(self, *args, **kwargs): + return False + + monkeypatch.setattr(tasks, "get_redis_service", lambda: _RedisSvc()) + + captured = {} + + class _Sig: + def __init__(self, kwargs): + self.kwargs = kwargs + self.queue = None + + def set(self, **kwargs): + self.queue = kwargs.get("queue") + return self + + monkeypatch.setattr(tasks, "forward_part", types.SimpleNamespace( + s=lambda **kwargs: _Sig(kwargs))) + monkeypatch.setattr(tasks, "aggregate_forward_parts", types.SimpleNamespace( + s=lambda **kwargs: _Sig(kwargs))) + + def _fake_group(signatures): + captured["parts"] = list(signatures) + return captured["parts"] + + def _fake_chord(group_tasks): + def _runner(callback): + captured["callback"] = callback + total = sum(len(sig.kwargs["chunks"]) for sig in group_tasks) + return types.SimpleNamespace(get=lambda: { + "success": True, + "total_indexed": total, + "total_submitted": total, + }) + return _runner + + @contextmanager + def _fake_allow_join_result(): + yield + + monkeypatch.setattr(tasks, "group", _fake_group) + monkeypatch.setattr(tasks, "chord", _fake_chord) + monkeypatch.setattr(tasks, "allow_join_result", _fake_allow_join_result) + monkeypatch.setattr(tasks, "_send_chunks_to_es", lambda **kwargs: { + "success": True, + "total_indexed": len(kwargs["chunks"]), + "total_submitted": len(kwargs["chunks"]), + }) + + out = tasks.forward( + FakeSelf("forward-aggregate-queue"), + processed_data={"chunks": [{"content": f"c-{i}", "metadata": {}} for i in range(70)]}, + index_name="idx", + source="/big.txt", + source_type="local", + file_id="file-1", + ) + + assert out["chunks_stored"] == 70 + assert captured["callback"].queue == "forward_aggregate_q" + assert captured["callback"].kwargs["file_id"] == "file-1" def test_process_sync_unsupported_raises_and_updates_state(monkeypatch): @@ -2893,6 +2970,61 @@ def is_task_cancelled(self, _task_id): assert out["total_submitted"] == 0 +def test_forward_part_returns_cancelled_when_document_is_deleted(monkeypatch): + tasks, _ = import_tasks_with_fake_ray(monkeypatch) + + class _Svc: + def is_document_delete_requested(self, **kwargs): + return kwargs["file_id"] == "file-deleted" + + def is_task_cancelled(self, _task_id): + return False + + monkeypatch.setattr(tasks, "get_redis_service", lambda: _Svc()) + monkeypatch.setattr( + tasks, + "_send_chunks_to_es", + lambda **kwargs: (_ for _ in ()).throw(AssertionError("deleted batch must not be sent")), + ) + self = types.SimpleNamespace( + request=types.SimpleNamespace(id="fp-deleted", retries=0), + retry=lambda **kwargs: (_ for _ in ()).throw(AssertionError("must not retry")), + ) + + out = tasks.forward_part( + self, + chunks=[{"content": "x"}], + index_name="idx", + source="knowledge_base/deleted.txt", + file_id="file-deleted", + batch_index=1, + total_batches=1, + ) + + assert out["cancelled"] is True + assert out["total_indexed"] == 0 + + +def test_aggregate_forward_parts_returns_cancelled_when_document_is_deleted(monkeypatch): + tasks, _ = import_tasks_with_fake_ray(monkeypatch) + + class _Svc: + def is_document_delete_requested(self, **kwargs): + return kwargs["file_id"] == "file-deleted" + + monkeypatch.setattr(tasks, "get_redis_service", lambda: _Svc()) + out = tasks.aggregate_forward_parts( + types.SimpleNamespace(request=types.SimpleNamespace(id="aggregate-deleted")), + parts_results=[{"total_indexed": 10, "total_submitted": 10}], + source="knowledge_base/deleted.txt", + index_name="idx", + file_id="file-deleted", + ) + + assert out["cancelled"] is True + assert out["total_indexed"] == 0 + + def test_forward_part_continues_when_parent_cancellation_lookup_fails(monkeypatch): tasks, _ = import_tasks_with_fake_ray(monkeypatch) monkeypatch.setattr(tasks, "get_redis_service", lambda: (_ for _ in ()).throw(RuntimeError("redis down"))) diff --git a/test/backend/services/test_data_process_service.py b/test/backend/services/test_data_process_service.py index 49d820e7df..6a90eb7ed6 100644 --- a/test/backend/services/test_data_process_service.py +++ b/test/backend/services/test_data_process_service.py @@ -2629,6 +2629,43 @@ async def _task_info(task_id): asyncio.run(_run()) + @patch('backend.services.data_process_service.get_redis_service') + @patch('backend.services.data_process_service.celery_app') + @patch('backend.services.data_process_service.get_all_task_ids_from_redis', return_value=['deleted-task']) + @patch('backend.services.data_process_service.get_task_info') + def test_get_all_tasks_hides_tasks_for_deleted_documents( + self, mock_get_task_info, _mock_ids, mock_celery_app, mock_get_redis_service + ): + mock_inspector = MagicMock() + mock_inspector.active.return_value = {} + mock_inspector.reserved.return_value = {} + mock_celery_app.control.inspect.return_value = mock_inspector + mock_celery_app.conf.broker_url = "redis://mock:6379/0" + mock_celery_app.conf.result_backend = "redis://mock:6379/0" + + async def _task_info(task_id): + return { + "id": task_id, + "task_name": "forward", + "index_name": "idx", + "path_or_url": "knowledge_base/deleted.txt", + "file_id": "file-deleted", + } + + mock_get_task_info.side_effect = _task_info + fence_service = MagicMock() + fence_service.is_document_delete_requested.return_value = True + mock_get_redis_service.return_value = fence_service + + rows = asyncio.run(self.service.get_all_tasks(filter=False)) + + self.assertEqual(rows, []) + fence_service.is_document_delete_requested.assert_called_once_with( + index_name="idx", + path_or_url="knowledge_base/deleted.txt", + file_id="file-deleted", + ) + # ---- Additional coverage tests ---- diff --git a/test/backend/services/test_redis_service.py b/test/backend/services/test_redis_service.py index 4c3434249c..2e8ab05278 100644 --- a/test/backend/services/test_redis_service.py +++ b/test/backend/services/test_redis_service.py @@ -143,6 +143,62 @@ def test_mark_and_check_task_cancelled(self, mock_from_url): mock_client.setex.assert_called_once() mock_client.get.assert_called_once() + @patch('backend.services.redis_service.REDIS_URL', 'redis://localhost:6379/0') + def test_document_delete_fence_can_be_set_checked_and_cleared(self): + mock_client = MagicMock() + mock_client.get.side_effect = [b'{"deleted": true}', None] + mock_client.delete.return_value = 1 + service = RedisService() + service._client = mock_client + + self.assertTrue(service.mark_document_delete_requested( + "idx", "knowledge_base/a.txt", file_id="file-1", ttl_seconds=60 + )) + self.assertTrue(service.is_document_delete_requested( + "idx", "knowledge_base/a.txt", file_id="file-1" + )) + self.assertEqual(service.clear_document_delete_marker( + "idx", "knowledge_base/a.txt", file_id="file-1" + ), 1) + self.assertEqual(mock_client.setex.call_count, 2) + mock_client.delete.assert_called_once() + + def test_prepare_document_deletion_marks_and_revokes_runtime_tasks(self): + service = RedisService() + service.mark_document_delete_requested = MagicMock(return_value=True) + service._collect_runtime_task_ids = MagicMock(return_value={"task-1", "task-2"}) + service.mark_task_cancelled = MagicMock(return_value=True) + service._revoke_task = MagicMock(return_value=True) + + result = service.prepare_document_deletion( + "idx", "knowledge_base/a.txt", file_id="file-1" + ) + + self.assertTrue(result["delete_fence_set"]) + self.assertEqual(result["runtime_tasks_found"], 2) + self.assertEqual(result["tasks_cancelled"], 2) + self.assertEqual(result["tasks_revoked"], 2) + self.assertEqual(service.mark_task_cancelled.call_count, 2) + self.assertEqual(service._revoke_task.call_count, 2) + + def test_runtime_task_matches_document_by_file_id_or_path(self): + task = { + "kwargs": json.dumps({ + "index_name": "idx", + "source": "knowledge_base/a.txt", + "file_id": "file-1", + }) + } + self.assertTrue(RedisService._runtime_task_matches_document( + task, "idx", "different-path", "file-1" + )) + self.assertTrue(RedisService._runtime_task_matches_document( + task, "idx", "knowledge_base/a.txt", None + )) + self.assertFalse(RedisService._runtime_task_matches_document( + task, "idx", "knowledge_base/a.txt", "file-2" + )) + def test_delete_knowledgebase_records(self): """Test delete_knowledgebase_records method""" # Setup diff --git a/test/backend/test_data_process_service_entrypoint.py b/test/backend/test_data_process_service_entrypoint.py index f2eb74fbdc..b3effda86d 100644 --- a/test/backend/test_data_process_service_entrypoint.py +++ b/test/backend/test_data_process_service_entrypoint.py @@ -93,6 +93,21 @@ def test_start_ray_cluster_returns_when_disabled(service_module): service_module.RayConfig.init_ray_for_service.assert_not_called() +def test_worker_configs_isolate_forward_parent_parts_and_aggregate(service_module): + configs = service_module.ServiceManager._build_worker_configs(4) + + assert [config["queue"] for config in configs] == [ + "process_q", + "process_part_q", + "forward_q", + "forward_part_q", + "forward_aggregate_q", + ] + assert configs[0]["concurrency"] == configs[1]["concurrency"] == 2 + assert configs[2]["concurrency"] == configs[3]["concurrency"] == 8 + assert configs[4]["concurrency"] == 2 + + def test_start_all_services_starts_enabled_services_in_order(service_module, monkeypatch): scheduler = types.SimpleNamespace(start=MagicMock()) scheduler_module = types.ModuleType("services.auto_summary_scheduler") From ad16436e36c87fef04ab1536931bf2eb7fe52014 Mon Sep 17 00:00:00 2001 From: cj2026-bit <647646783@qq.com> Date: Wed, 26 Aug 2026 17:53:40 +0800 Subject: [PATCH 2/5] test: cover deletion fences and cancellation branches --- backend/data_process/tasks.py | 1 + test/backend/data_process/test_tasks.py | 157 ++++++++++++++++++ .../services/test_data_process_service.py | 49 ++++++ .../services/test_file_management_service.py | 9 +- test/backend/services/test_redis_service.py | 88 ++++++++++ .../services/test_vectordatabase_service.py | 11 +- .../test_data_process_service_entrypoint.py | 29 ++++ 7 files changed, 340 insertions(+), 4 deletions(-) diff --git a/backend/data_process/tasks.py b/backend/data_process/tasks.py index 4a9f539cf1..e2a30f21f6 100644 --- a/backend/data_process/tasks.py +++ b/backend/data_process/tasks.py @@ -550,6 +550,7 @@ def _is_document_delete_requested( def _build_forward_cancelled_result(ctx: _ForwardContext) -> Dict[str, Any]: return { + 'cancelled': True, 'task_id': ctx.task_id, 'source': ctx.source, 'index_name': ctx.index_name, diff --git a/test/backend/data_process/test_tasks.py b/test/backend/data_process/test_tasks.py index b31323f613..7a27758598 100644 --- a/test/backend/data_process/test_tasks.py +++ b/test/backend/data_process/test_tasks.py @@ -3025,6 +3025,146 @@ def is_document_delete_requested(self, **kwargs): assert out["total_indexed"] == 0 +def test_process_returns_cancelled_when_document_is_deleted_before_start(monkeypatch): + tasks, _ = import_tasks_with_fake_ray(monkeypatch) + monkeypatch.setattr(tasks, "_is_document_delete_requested", lambda *args, **kwargs: True) + monkeypatch.setattr( + tasks, + "_update_file_lifecycle", + lambda **kwargs: (_ for _ in ()).throw(AssertionError("deleted process must not update lifecycle")), + ) + + out = tasks.process( + FakeSelf("process-deleted-before-start"), + source="knowledge_base/deleted.txt", + source_type="local", + index_name="idx", + original_filename="deleted.txt", + file_id="file-deleted", + ) + + assert out["cancelled"] is True + assert out["file_id"] == "file-deleted" + + +def test_process_returns_cancelled_when_deleted_after_local_extraction(monkeypatch, tmp_path): + tasks, _ = import_tasks_with_fake_ray(monkeypatch) + source = tmp_path / "deleted-after-extraction.txt" + source.write_text("content") + checks = iter([False, True]) + monkeypatch.setattr(tasks, "_is_document_delete_requested", lambda *args, **kwargs: next(checks)) + monkeypatch.setattr(tasks, "_update_file_lifecycle", lambda **kwargs: None) + monkeypatch.setattr( + tasks, + "_process_source_with_split", + lambda **kwargs: (False, [{"content": "chunk", "metadata": {}}], None), + ) + + out = tasks.process( + FakeSelf("process-deleted-after-local"), + source=str(source), + source_type="local", + chunking_strategy="basic", + index_name="idx", + file_id="file-deleted", + ) + + assert out["cancelled"] is True + + +def test_process_returns_cancelled_when_deleted_after_minio_extraction(monkeypatch): + tasks, _ = import_tasks_with_fake_ray(monkeypatch) + checks = iter([False, False, True]) + monkeypatch.setattr(tasks, "_is_document_delete_requested", lambda *args, **kwargs: next(checks)) + monkeypatch.setattr(tasks, "_update_file_lifecycle", lambda **kwargs: None) + monkeypatch.setattr(tasks, "_fetch_minio_source", lambda source: b"content") + monkeypatch.setattr( + tasks, + "_process_source_with_split", + lambda **kwargs: (False, [{"content": "chunk", "metadata": {}}], None), + ) + + out = tasks.process( + FakeSelf("process-deleted-after-minio"), + source="knowledge_base/deleted.txt", + source_type="minio", + chunking_strategy="basic", + index_name="idx", + file_id="file-deleted", + ) + + assert out["cancelled"] is True + + +def test_forward_returns_cancelled_when_document_is_deleted_before_load(monkeypatch): + tasks, _ = import_tasks_with_fake_ray(monkeypatch) + + class _Service: + def is_task_cancelled(self, _task_id): + return False + + def is_document_delete_requested(self, **kwargs): + return True + + monkeypatch.setattr(tasks, "get_redis_service", lambda: _Service()) + monkeypatch.setattr(tasks, "_update_file_lifecycle", lambda **kwargs: None) + + out = tasks.forward( + FakeSelf("forward-deleted-before-load"), + processed_data={"chunks": [{"content": "x", "metadata": {}}]}, + index_name="idx", + source="knowledge_base/deleted.txt", + file_id="file-deleted", + ) + + assert out["cancelled"] is True + + +def test_forward_returns_cancelled_when_deleted_after_loading_chunks(monkeypatch): + tasks, _ = import_tasks_with_fake_ray(monkeypatch) + checks = iter([False, True]) + monkeypatch.setattr(tasks, "_is_document_delete_requested", lambda *args, **kwargs: next(checks)) + monkeypatch.setattr(tasks, "_update_file_lifecycle", lambda **kwargs: None) + monkeypatch.setattr( + tasks, + "_load_forward_chunks", + lambda _self, **kwargs: ([{"content": "x", "metadata": {}}], False, "knowledge_base/deleted.txt", "idx", "deleted.txt"), + ) + + out = tasks.forward( + FakeSelf("forward-deleted-after-load"), + processed_data={"chunks": [{"content": "x", "metadata": {}}]}, + index_name="idx", + source="knowledge_base/deleted.txt", + file_id="file-deleted", + ) + + assert out["cancelled"] is True + + +def test_forward_returns_cancelled_before_first_es_write(monkeypatch): + tasks, _ = import_tasks_with_fake_ray(monkeypatch) + checks = iter([False, False, True]) + monkeypatch.setattr(tasks, "_is_document_delete_requested", lambda *args, **kwargs: next(checks)) + monkeypatch.setattr(tasks, "_update_file_lifecycle", lambda **kwargs: None) + monkeypatch.setattr(tasks, "get_file_size", lambda *args, **kwargs: 0) + monkeypatch.setattr( + tasks, + "_load_forward_chunks", + lambda _self, **kwargs: ([{"content": "x", "metadata": {}}], False, "knowledge_base/deleted.txt", "idx", "deleted.txt"), + ) + + out = tasks.forward( + FakeSelf("forward-deleted-before-es"), + processed_data={"chunks": [{"content": "x", "metadata": {}}]}, + index_name="idx", + source="knowledge_base/deleted.txt", + file_id="file-deleted", + ) + + assert out["cancelled"] is True + + def test_forward_part_continues_when_parent_cancellation_lookup_fails(monkeypatch): tasks, _ = import_tasks_with_fake_ray(monkeypatch) monkeypatch.setattr(tasks, "get_redis_service", lambda: (_ for _ in ()).throw(RuntimeError("redis down"))) @@ -3374,6 +3514,23 @@ def _delete(*_a, **_k): assert called["delete"] == 0 +def test_cleanup_source_skips_when_forward_was_cancelled(monkeypatch): + tasks, _ = import_tasks_with_fake_ray(monkeypatch) + monkeypatch.setattr( + tasks, + "_delete_source_file_via_http_sync", + lambda **kwargs: (_ for _ in ()).throw(AssertionError("cancelled forward must not delete source")), + ) + + out = tasks.cleanup_source( + FakeSelf("cleanup-cancelled"), + {"cancelled": True, "index_name": "idx", "source": "/a.txt"}, + ) + + assert out["source_cleanup"]["skipped_reason"] == "forward_cancelled" + assert out["source_cleanup"]["attempted"] is False + + def test_cleanup_source_calls_delete_with_scope_source_only(monkeypatch): tasks, _ = import_tasks_with_fake_ray(monkeypatch) monkeypatch.setattr(tasks, "ELASTICSEARCH_SERVICE", "http://api") diff --git a/test/backend/services/test_data_process_service.py b/test/backend/services/test_data_process_service.py index 6a90eb7ed6..1e08c074f2 100644 --- a/test/backend/services/test_data_process_service.py +++ b/test/backend/services/test_data_process_service.py @@ -2666,6 +2666,55 @@ async def _task_info(task_id): file_id="file-deleted", ) + @patch('backend.services.data_process_service.get_redis_service', side_effect=RuntimeError("redis unavailable")) + @patch('backend.services.data_process_service.celery_app') + @patch('backend.services.data_process_service.get_all_task_ids_from_redis', return_value=[]) + @patch('backend.services.data_process_service.get_task_info') + def test_get_all_tasks_continues_when_delete_fence_service_unavailable( + self, mock_get_task_info, _mock_ids, mock_celery_app, _mock_get_redis_service + ): + mock_inspector = MagicMock() + mock_inspector.active.return_value = {} + mock_inspector.reserved.return_value = {} + mock_celery_app.control.inspect.return_value = mock_inspector + mock_celery_app.conf.broker_url = "redis://mock:6379/0" + mock_celery_app.conf.result_backend = "redis://mock:6379/0" + + rows = asyncio.run(self.service.get_all_tasks(filter=False)) + + assert rows == [] + mock_get_task_info.assert_not_called() + + @patch('backend.services.data_process_service.get_redis_service') + @patch('backend.services.data_process_service.celery_app') + @patch('backend.services.data_process_service.get_all_task_ids_from_redis', return_value=['fenced-task']) + @patch('backend.services.data_process_service.get_task_info') + def test_get_all_tasks_keeps_task_when_delete_fence_lookup_fails( + self, mock_get_task_info, _mock_ids, mock_celery_app, mock_get_redis_service + ): + mock_inspector = MagicMock() + mock_inspector.active.return_value = {} + mock_inspector.reserved.return_value = {} + mock_celery_app.control.inspect.return_value = mock_inspector + mock_celery_app.conf.broker_url = "redis://mock:6379/0" + mock_celery_app.conf.result_backend = "redis://mock:6379/0" + async def _task_info(task_id): + return { + "id": task_id, + "task_name": "forward", + "index_name": "idx", + "path_or_url": "knowledge_base/a.txt", + } + + mock_get_task_info.side_effect = _task_info + fence_service = MagicMock() + fence_service.is_document_delete_requested.side_effect = RuntimeError("lookup failed") + mock_get_redis_service.return_value = fence_service + + rows = asyncio.run(self.service.get_all_tasks(filter=False)) + + assert [row["id"] for row in rows] == ["fenced-task"] + # ---- Additional coverage tests ---- diff --git a/test/backend/services/test_file_management_service.py b/test/backend/services/test_file_management_service.py index 63ba921a3f..0660d12e14 100644 --- a/test/backend/services/test_file_management_service.py +++ b/test/backend/services/test_file_management_service.py @@ -568,13 +568,20 @@ def check_hard_limit_post_write(self, *_args, **_kwargs): quota_module = types.ModuleType("services.quota_service") quota_module.QuotaService = FakeQuota + redis_module = types.ModuleType("services.redis_service") + redis_module.get_redis_service = MagicMock( + side_effect=RuntimeError("redis unavailable") + ) with patch.object(knowledge_storage_stub, "resolve_storage_context", return_value=context), \ patch("backend.services.file_management_service.create_file_records", return_value=[record]) as create_records, \ patch("backend.services.file_management_service.transition_file_record", side_effect=transition) as transition_record, \ patch("backend.services.file_management_service.upload_to_minio", AsyncMock(return_value=[ {"success": True, "file_id": "fid-1", "file_name": "a.txt", "object_name": "folder/a.txt", "file_size": 3} ])) as upload_mock, \ - patch.dict(sys.modules, {"services.quota_service": quota_module}): + patch.dict(sys.modules, { + "services.quota_service": quota_module, + "services.redis_service": redis_module, + }): result = await upload_files_impl( destination="minio", file=[mock_file], folder="folder", index_name="kb-1", user_id="user-1" ) diff --git a/test/backend/services/test_redis_service.py b/test/backend/services/test_redis_service.py index 2e8ab05278..d050b2b6d2 100644 --- a/test/backend/services/test_redis_service.py +++ b/test/backend/services/test_redis_service.py @@ -163,6 +163,80 @@ def test_document_delete_fence_can_be_set_checked_and_cleared(self): self.assertEqual(mock_client.setex.call_count, 2) mock_client.delete.assert_called_once() + def test_document_delete_fence_validation_and_redis_failures(self): + service = RedisService() + assert service.mark_document_delete_requested("idx") is False + assert service.mark_document_delete_requested( + "idx", "knowledge_base/a.txt", ttl_seconds=-1 + ) is False + + mock_client = MagicMock() + mock_client.setex.side_effect = RuntimeError("set failed") + service._client = mock_client + assert service.mark_document_delete_requested("idx", "knowledge_base/a.txt") is False + + mock_client.delete.side_effect = RuntimeError("delete failed") + assert service.clear_document_delete_marker("idx", "knowledge_base/a.txt") == 0 + mock_client.get.side_effect = RuntimeError("get failed") + assert service.is_document_delete_requested("idx", "knowledge_base/a.txt") is False + + def test_document_delete_fence_clear_without_keys(self): + service = RedisService() + assert service.clear_document_delete_marker("idx") == 0 + assert service.is_document_delete_requested(None) is False + + @patch('backend.services.redis_service.REDIS_URL', 'redis://broker/0') + @patch('backend.services.redis_service.REDIS_BACKEND_URL', 'redis://backend/1') + def test_celery_control_app_is_lazy_and_rejects_missing_broker(self): + service = RedisService() + app = service._get_celery_control_app() + assert app is service._get_celery_control_app() + assert app.conf.broker_url == 'redis://broker/0' + + with patch('backend.services.redis_service.REDIS_URL', ''): + missing = RedisService() + with self.assertRaises(ValueError): + missing._get_celery_control_app() + + def test_task_kwargs_normalization(self): + assert RedisService._task_kwargs({"kwargs": '{"file_id": "f1"}'}) == {"file_id": "f1"} + assert RedisService._task_kwargs({"kwargs": "bad-json"}) == {} + assert RedisService._task_kwargs({"kwargs": ["not-a-dict"]}) == {} + assert RedisService._task_kwargs({}) == {} + + def test_collect_runtime_task_ids_filters_matches_and_handles_inspector_failure(self): + service = RedisService() + inspector = MagicMock() + inspector.active.return_value = { + "worker": [ + {"id": "match-file", "kwargs": {"index_name": "idx", "file_id": "f1"}}, + {"id": "wrong-index", "kwargs": {"index_name": "other", "file_id": "f1"}}, + {"kwargs": {"index_name": "idx", "file_id": "f1"}}, + "not-a-task", + ] + } + inspector.reserved.return_value = { + "worker": [{"id": "match-path", "kwargs": {"index_name": "idx", "source": "a.txt"}}] + } + control_app = MagicMock() + control_app.control.inspect.return_value = inspector + service._get_celery_control_app = MagicMock(return_value=control_app) + + assert service._collect_runtime_task_ids("idx", "a.txt", "f1") == {"match-file", "match-path"} + + service._get_celery_control_app.side_effect = RuntimeError("inspect failed") + assert service._collect_runtime_task_ids("idx", "a.txt", "f1") == set() + + def test_revoke_task_success_failure_and_empty_id(self): + service = RedisService() + assert service._revoke_task("") is False + control_app = MagicMock() + service._get_celery_control_app = MagicMock(return_value=control_app) + assert service._revoke_task("task-1") is True + control_app.control.revoke.assert_called_once_with("task-1", terminate=False) + control_app.control.revoke.side_effect = RuntimeError("revoke failed") + assert service._revoke_task("task-2") is False + def test_prepare_document_deletion_marks_and_revokes_runtime_tasks(self): service = RedisService() service.mark_document_delete_requested = MagicMock(return_value=True) @@ -181,6 +255,20 @@ def test_prepare_document_deletion_marks_and_revokes_runtime_tasks(self): self.assertEqual(service.mark_task_cancelled.call_count, 2) self.assertEqual(service._revoke_task.call_count, 2) + def test_prepare_document_deletion_keeps_zero_counts_when_cancel_or_revoke_fails(self): + service = RedisService() + service.mark_document_delete_requested = MagicMock(return_value=False) + service._collect_runtime_task_ids = MagicMock(return_value={"task-1"}) + service.mark_task_cancelled = MagicMock(return_value=False) + service._revoke_task = MagicMock(return_value=False) + + result = service.prepare_document_deletion("idx", "a.txt", file_id="f1") + + assert result["delete_fence_set"] is False + assert result["runtime_tasks_found"] == 1 + assert result["tasks_cancelled"] == 0 + assert result["tasks_revoked"] == 0 + def test_runtime_task_matches_document_by_file_id_or_path(self): task = { "kwargs": json.dumps({ diff --git a/test/backend/services/test_vectordatabase_service.py b/test/backend/services/test_vectordatabase_service.py index fc87d000ca..010367bfc3 100644 --- a/test/backend/services/test_vectordatabase_service.py +++ b/test/backend/services/test_vectordatabase_service.py @@ -2456,9 +2456,12 @@ def test_mark_file_deleted_ignores_missing_record( ElasticSearchService._mark_file_deleted("test_index", "knowledge_base/missing.txt") mock_delete.assert_not_called() + @patch('backend.services.vectordatabase_service.get_redis_service', side_effect=RuntimeError("redis unavailable")) @patch('backend.services.vectordatabase_service.delete_file_record', return_value=True) @patch('backend.services.vectordatabase_service.transition_file_record') - def test_delete_lifecycle_record_without_object_hard_deletes(self, mock_transition, mock_delete): + def test_delete_lifecycle_record_without_object_hard_deletes( + self, mock_transition, mock_delete, _mock_redis + ): mock_transition.return_value = {"file_id": "fid-no-object", "status": "DELETE_REQUESTED"} result = ElasticSearchService.delete_lifecycle_record_without_object( @@ -2782,13 +2785,14 @@ def test_delete_source_file_releases_charge_only_after_minio_success( updated_by=None, ) + @patch('backend.services.vectordatabase_service.get_redis_service', side_effect=RuntimeError("redis unavailable")) @patch( 'backend.services.vectordatabase_service.get_all_files_status', new_callable=AsyncMock, ) @patch('backend.services.vectordatabase_service.delete_file') def test_delete_document_by_scope_source_only( - self, mock_delete_file, mock_get_status + self, mock_delete_file, mock_get_status, _mock_redis ): mock_get_status.return_value = { "knowledge_base/doc.pdf": {"state": "COMPLETED"} @@ -2813,8 +2817,9 @@ def test_delete_document_by_scope_source_only( 'delete_documents', return_value={"status": "success", "deleted_minio": True}, ) + @patch('backend.services.vectordatabase_service.get_redis_service', side_effect=RuntimeError("redis unavailable")) def test_delete_document_by_scope_does_not_require_storage_ledger( - self, mock_delete_documents + self, _mock_redis, mock_delete_documents ): result = asyncio.run( ElasticSearchService.delete_document_by_scope( diff --git a/test/backend/test_data_process_service_entrypoint.py b/test/backend/test_data_process_service_entrypoint.py index b3effda86d..3401902895 100644 --- a/test/backend/test_data_process_service_entrypoint.py +++ b/test/backend/test_data_process_service_entrypoint.py @@ -108,6 +108,35 @@ def test_worker_configs_isolate_forward_parent_parts_and_aggregate(service_modul assert configs[4]["concurrency"] == 2 +def test_start_workers_launches_each_isolated_queue(service_module, monkeypatch): + launched = [] + + class _Process: + def __init__(self, command, **kwargs): + self.pid = len(launched) + 100 + self.stdout = types.SimpleNamespace(readline=lambda: "") + launched.append((command, kwargs)) + + monkeypatch.setattr(service_module, "RAY_NUM_CPUS", "4") + monkeypatch.setattr(service_module, "RAY_ACTOR_NUM_CPUS", 2) + monkeypatch.setattr(service_module.subprocess, "Popen", _Process) + monkeypatch.setattr(service_module.threading, "Thread", lambda **kwargs: types.SimpleNamespace(start=lambda: None)) + + service_module.service_processes["workers"] = [] + manager = service_module.ServiceManager({"start_workers": True}) + + assert manager.start_workers() is True + assert [row["queue"] for row in service_module.service_processes["workers"]] == [ + "process_q", + "process_part_q", + "forward_q", + "forward_part_q", + "forward_aggregate_q", + ] + assert len(launched) == 5 + service_module.service_processes["workers"] = [] + + def test_start_all_services_starts_enabled_services_in_order(service_module, monkeypatch): scheduler = types.SimpleNamespace(start=MagicMock()) scheduler_module = types.ModuleType("services.auto_summary_scheduler") From 8125ca0dc458c0c703a2e7a2c11f9b5bba4b2407 Mon Sep 17 00:00:00 2001 From: cj2026-bit <647646783@qq.com> Date: Wed, 26 Aug 2026 18:17:41 +0800 Subject: [PATCH 3/5] test: fix CI coverage test assumptions --- test/backend/data_process/test_tasks.py | 3 ++- test/backend/services/test_redis_service.py | 1 - 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/test/backend/data_process/test_tasks.py b/test/backend/data_process/test_tasks.py index 7a27758598..16e9f945ec 100644 --- a/test/backend/data_process/test_tasks.py +++ b/test/backend/data_process/test_tasks.py @@ -3074,7 +3074,8 @@ def test_process_returns_cancelled_when_deleted_after_local_extraction(monkeypat def test_process_returns_cancelled_when_deleted_after_minio_extraction(monkeypatch): tasks, _ = import_tasks_with_fake_ray(monkeypatch) - checks = iter([False, False, True]) + # The MinIO path performs the initial and post-extraction checks only. + checks = iter([False, True]) monkeypatch.setattr(tasks, "_is_document_delete_requested", lambda *args, **kwargs: next(checks)) monkeypatch.setattr(tasks, "_update_file_lifecycle", lambda **kwargs: None) monkeypatch.setattr(tasks, "_fetch_minio_source", lambda source: b"content") diff --git a/test/backend/services/test_redis_service.py b/test/backend/services/test_redis_service.py index d050b2b6d2..a4a2392a36 100644 --- a/test/backend/services/test_redis_service.py +++ b/test/backend/services/test_redis_service.py @@ -191,7 +191,6 @@ def test_celery_control_app_is_lazy_and_rejects_missing_broker(self): service = RedisService() app = service._get_celery_control_app() assert app is service._get_celery_control_app() - assert app.conf.broker_url == 'redis://broker/0' with patch('backend.services.redis_service.REDIS_URL', ''): missing = RedisService() From d7ff75b8ee762c58de582a822344ac269c9490d4 Mon Sep 17 00:00:00 2001 From: cj2026-bit <647646783@qq.com> Date: Wed, 26 Aug 2026 18:46:30 +0800 Subject: [PATCH 4/5] refactor: remove redundant celery route entries --- backend/data_process/app.py | 8 ++------ 1 file changed, 2 insertions(+), 6 deletions(-) diff --git a/backend/data_process/app.py b/backend/data_process/app.py index 1d4e025d41..3ec2ad8d16 100644 --- a/backend/data_process/app.py +++ b/backend/data_process/app.py @@ -44,18 +44,14 @@ # Explicitly set result backend broker_url=REDIS_URL, result_backend=REDIS_BACKEND_URL, - # Dedicated queues keep parent, part, and aggregate tasks from starving - # one another when a parent task waits for a Celery chord. + # Explicitly route the newly isolated forward child and aggregate tasks. + # Other tasks keep their queue from the @app.task declaration. task_routes={ f'{import_path}.process': {'queue': 'process_q'}, f'{import_path}.forward': {'queue': 'forward_q'}, f'{import_path}.process_and_forward': {'queue': 'process_q'}, - f'{import_path}.process_part': {'queue': 'process_part_q'}, - f'{import_path}.aggregate_parts': {'queue': 'process_part_q'}, - f'{import_path}.aggregate_store_chunks': {'queue': 'process_part_q'}, f'{import_path}.forward_part': {'queue': 'forward_part_q'}, f'{import_path}.aggregate_forward_parts': {'queue': 'forward_aggregate_q'}, - f'{import_path}.cleanup_source': {'queue': 'forward_q'}, }, task_serializer='json', accept_content=['json'], From ee43fcbb674174196f42c5b5f52e06ec79cccf62 Mon Sep 17 00:00:00 2001 From: cj2026-bit <647646783@qq.com> Date: Wed, 26 Aug 2026 19:33:37 +0800 Subject: [PATCH 5/5] revert: defer in-progress file deletion handling --- backend/data_process/tasks.py | 147 ----------- backend/services/data_process_service.py | 34 --- backend/services/file_management_service.py | 19 -- backend/services/redis_service.py | 230 ------------------ backend/services/vectordatabase_service.py | 47 +--- test/backend/data_process/test_tasks.py | 217 ----------------- .../services/test_data_process_service.py | 86 ------- .../services/test_file_management_service.py | 9 +- test/backend/services/test_redis_service.py | 143 ----------- .../services/test_vectordatabase_service.py | 11 +- 10 files changed, 6 insertions(+), 937 deletions(-) diff --git a/backend/data_process/tasks.py b/backend/data_process/tasks.py index e2a30f21f6..ed6ba8a539 100644 --- a/backend/data_process/tasks.py +++ b/backend/data_process/tasks.py @@ -519,38 +519,8 @@ def _is_forward_task_cancelled(ctx: _ForwardContext) -> bool: return False -def _is_document_delete_requested( - index_name: Optional[str], - source: Optional[str], - file_id: Optional[str], -) -> bool: - """Check the durable file-level deletion fence used by delete + sync.""" - if not index_name: - return False - try: - return bool( - get_redis_service().is_document_delete_requested( - index_name=index_name, - path_or_url=source, - file_id=file_id, - ) - ) - except Exception as exc: - # Deletion fencing is best effort when Redis is unavailable. The task - # still follows the existing task-cancellation/lifecycle safeguards. - logger.warning( - "Failed to check document deletion fence for index=%s source=%s file_id=%s: %s", - index_name, - source, - file_id, - exc, - ) - return False - - def _build_forward_cancelled_result(ctx: _ForwardContext) -> Dict[str, Any]: return { - 'cancelled': True, 'task_id': ctx.task_id, 'source': ctx.source, 'index_name': ctx.index_name, @@ -1132,30 +1102,11 @@ def forward_part( batch_index: Optional[int] = None, total_batches: Optional[int] = None, large_mode: Optional[bool] = False, - file_id: Optional[str] = None, ) -> Dict[str, Any]: """ Forward sub-task that indexes a chunk batch. """ try: - if _is_document_delete_requested(index_name, source, file_id): - logger.info( - "Skipping forward batch %s/%s because document was deleted: index=%s source=%s file_id=%s", - batch_index, - total_batches, - index_name, - source, - file_id, - ) - return { - "success": True, - "total_indexed": 0, - "total_submitted": 0, - "batch_index": batch_index, - "total_batches": total_batches, - "cancelled": True, - } - # Respect cancellation from parent task if available if parent_task_id: try: @@ -1272,29 +1223,10 @@ def aggregate_forward_parts( source: Optional[str] = None, index_name: Optional[str] = None, original_filename: Optional[str] = None, - file_id: Optional[str] = None, ) -> Dict[str, Any]: """ Aggregate forward_part results. """ - if _is_document_delete_requested(index_name, source, file_id): - logger.info( - "Skipping forward aggregate because document was deleted: index=%s source=%s file_id=%s", - index_name, - source, - file_id, - ) - return { - "success": True, - "total_indexed": 0, - "total_submitted": 0, - "source": source, - "index_name": index_name, - "original_filename": original_filename, - "file_id": file_id, - "cancelled": True, - } - total_indexed = 0 total_submitted = 0 for result in parts_results or []: @@ -1631,21 +1563,6 @@ def process( start_time = time.time() task_id = self.request.id file_id = params.get("file_id") - if _is_document_delete_requested(index_name, source, file_id): - logger.info( - "Skipping process task because document was deleted: index=%s source=%s file_id=%s", - index_name, - source, - file_id, - ) - return { - "cancelled": True, - "source": source, - "index_name": index_name, - "original_filename": original_filename, - "file_id": file_id, - "task_id": task_id, - } _update_file_lifecycle( file_id=file_id, tenant_id=tenant_id, @@ -1705,22 +1622,6 @@ def process( tenant_id=tenant_id, params=params, ) - if _is_document_delete_requested(index_name, source, file_id): - logger.info( - "Stopping process task after source extraction because document was deleted: " - "index=%s source=%s file_id=%s", - index_name, - source, - file_id, - ) - return { - "cancelled": True, - "source": source, - "index_name": index_name, - "original_filename": original_filename, - "file_id": file_id, - "task_id": task_id, - } elapsed_time = time.time() - start_time processing_speed = file_size_mb / \ elapsed_time if file_size_mb > 0 and elapsed_time > 0 else 0 @@ -1761,23 +1662,6 @@ def process( raise NotImplementedError( f"Source type '{source_type}' not yet supported") - if _is_document_delete_requested(index_name, source, file_id): - logger.info( - "Stopping process task after extraction because document was deleted: " - "index=%s source=%s file_id=%s", - index_name, - source, - file_id, - ) - return { - "cancelled": True, - "source": source, - "index_name": index_name, - "original_filename": original_filename, - "file_id": file_id, - "task_id": task_id, - } - if split_async: chunk_count = split_chunk_count or 0 if chunk_count == 0: @@ -2054,13 +1938,6 @@ def forward( ) return _build_forward_cancelled_result(ctx) - if _is_document_delete_requested(index_name, source, file_id): - logger.info( - f"[{self.request.id}] FORWARD TASK: Document deletion fence is set; " - f"skipping chunk forwarding for source '{source}' in index '{index_name}'." - ) - return _build_forward_cancelled_result(ctx) - chunks, split_async, original_source, original_index_name, filename = _load_forward_chunks( self, processed_data=processed_data, @@ -2069,13 +1946,6 @@ def forward( filename=filename, ) - if _is_document_delete_requested(original_index_name, original_source, file_id): - logger.info( - f"[{self.request.id}] FORWARD TASK: Document was deleted while loading chunks; " - f"skipping ES indexing for source '{original_source}'." - ) - return _build_forward_cancelled_result(ctx) - # Calculate total chunks for progress tracking total_chunks = len(chunks) if chunks else 0 set_span_attributes(chunk_count=total_chunks, stage="forward.format") @@ -2140,15 +2010,6 @@ def forward( } ) - # Re-check immediately before the first ES write. A delete may race - # with formatting or progress initialization after the earlier check. - if _is_document_delete_requested(original_index_name, original_source, file_id): - logger.info( - f"[{self.request.id}] FORWARD TASK: Document was deleted before ES indexing; " - f"skipping source '{original_source}'." - ) - return _build_forward_cancelled_result(ctx) - try: redis_service = get_redis_service() redis_service.save_progress_info(task_id, 0, total_chunks) @@ -2198,7 +2059,6 @@ def forward( parent_total_chunks=total_chunks, source=original_source, original_filename=original_filename, - file_id=file_id, batch_index=idx + 1, total_batches=total_batches, # If request was split into multiple groups, force all groups to use large path. @@ -2209,7 +2069,6 @@ def forward( source=original_source, index_name=original_index_name, original_filename=original_filename, - file_id=file_id, ).set(queue='forward_aggregate_q') result = chord(group_tasks)(callback) with allow_join_result(): @@ -2422,12 +2281,6 @@ def cleanup_source( "error": None, } - if (forward_result or {}).get("cancelled"): - cleanup_info["skipped_reason"] = "forward_cancelled" - forward_result = dict(forward_result or {}) - forward_result["source_cleanup"] = cleanup_info - return forward_result - if not index_name or not source: cleanup_info["skipped_reason"] = "missing_index_name_or_source" forward_result = dict(forward_result or {}) diff --git a/backend/services/data_process_service.py b/backend/services/data_process_service.py index 34148f0f21..8051802ba2 100644 --- a/backend/services/data_process_service.py +++ b/backend/services/data_process_service.py @@ -27,7 +27,6 @@ from data_process.tasks import submit_process_forward_chain from data_process.utils import get_all_task_ids_from_redis, get_task_info from database.attachment_db import delete_file, file_exists, get_file_size_from_minio, get_file_stream, upload_file -from services.redis_service import get_redis_service from utils.file_management_utils import convert_office_to_pdf from utils.knowledge_ingestion_errors import classify_ingestion_exception @@ -147,14 +146,6 @@ async def get_all_tasks(self, filter: bool = True) -> List[Dict[str, Any]]: all_tasks = [] try: self._get_celery_inspector() - delete_fence_service = None - try: - delete_fence_service = get_redis_service() - except Exception as redis_init_exc: - logger.debug( - "Deletion-fence filtering unavailable while listing tasks: %s", - redis_init_exc, - ) # Collect task IDs from different sources and keep runtime metadata task_ids = set() @@ -273,31 +264,6 @@ def get_reserved(): if not task_info.get('file_id') and runtime_meta.get('file_id'): task_info['file_id'] = runtime_meta.get('file_id') - # A task can remain in Celery's result backend or broker after - # the file row and external document were deleted. Hide it - # from the user-visible task/file list; the worker-side fence - # also prevents a late task from recreating the document. - if delete_fence_service: - try: - document_deleted = delete_fence_service.is_document_delete_requested( - index_name=task_info.get('index_name'), - path_or_url=task_info.get('path_or_url'), - file_id=task_info.get('file_id'), - ) - except Exception as fence_exc: - logger.debug( - "Unable to check deletion fence for task %s: %s", - task_id, - fence_exc, - ) - document_deleted = False - if document_deleted: - logger.info( - "Skipping task %s because its document deletion fence is set", - task_id, - ) - continue - if filter and not (task_info.get('index_name') and task_info.get('task_name')): # Keep user-visible queued tasks even before worker updates task meta. if task_info.get('task_name') not in {'process', 'forward', 'process_and_forward'}: diff --git a/backend/services/file_management_service.py b/backend/services/file_management_service.py index 38223c921d..76a37be9bb 100644 --- a/backend/services/file_management_service.py +++ b/backend/services/file_management_service.py @@ -450,25 +450,6 @@ async def upload_files_impl( }) if lifecycle_record_specs: - # A new upload may reuse an object path after a previous - # delete. Clear the old deletion fence before workers see the - # new lifecycle row; the file-id fence remains unique to the - # old record. - try: - from services.redis_service import get_redis_service - - redis_service = get_redis_service() - for spec in lifecycle_record_specs: - redis_service.clear_document_delete_marker( - index_name=spec["index_name"], - path_or_url=spec["object_name"], - ) - except Exception as deletion_fence_exc: - logger.warning( - "Failed to clear an old document deletion fence before upload: %s", - deletion_fence_exc, - ) - # Lifecycle persistence is a required upload precondition. The # repository creates the whole batch in one transaction; any # database error must stop before MinIO is touched. diff --git a/backend/services/redis_service.py b/backend/services/redis_service.py index b452a0be20..15a41b6772 100644 --- a/backend/services/redis_service.py +++ b/backend/services/redis_service.py @@ -1,8 +1,6 @@ -import hashlib import json import logging import re -import time from typing import Any, Dict, List, Optional, Set, Tuple import redis @@ -21,12 +19,9 @@ class RedisService: """Redis service for managing cache and task data""" - DOCUMENT_DELETE_MARKER_TTL_SECONDS = 24 * 60 * 60 - def __init__(self): self._client = None self._backend_client = None - self._celery_control_app = None @property def client(self) -> redis.Redis: @@ -56,231 +51,6 @@ def backend_client(self) -> redis.Redis: # Cancellation helpers # ------------------------------------------------------------------ - @staticmethod - def _document_delete_marker_keys( - index_name: Optional[str] = None, - path_or_url: Optional[str] = None, - file_id: Optional[str] = None, - ) -> List[str]: - """Return stable Redis keys for a file deletion fence.""" - keys = [] - if file_id: - keys.append(f"kb-delete:file:{file_id}") - if index_name and path_or_url: - identity = f"{index_name}\0{path_or_url}".encode("utf-8") - digest = hashlib.sha256(identity).hexdigest() - keys.append(f"kb-delete:path:{digest}") - return keys - - def mark_document_delete_requested( - self, - index_name: str, - path_or_url: Optional[str] = None, - file_id: Optional[str] = None, - ttl_seconds: Optional[int] = None, - ) -> bool: - """Write a short-lived file-level deletion fence. - - The fence is deliberately separate from per-task cancellation keys. It - survives result/cache cleanup so queued broker messages cannot recreate - a document after its lifecycle row has been hard deleted. - """ - keys = self._document_delete_marker_keys(index_name, path_or_url, file_id) - if not keys: - return False - ttl = int(ttl_seconds or self.DOCUMENT_DELETE_MARKER_TTL_SECONDS) - if ttl <= 0: - return False - payload = json.dumps({ - "index_name": index_name, - "path_or_url": path_or_url, - "file_id": file_id, - "requested_at": time.time(), - }, ensure_ascii=False) - try: - for key in keys: - self.client.setex(key, ttl, payload) - logger.info( - "Marked document deletion fence: index=%s path=%s file_id=%s ttl=%ss", - index_name, - path_or_url, - file_id, - ttl, - ) - return True - except Exception as exc: - logger.warning( - "Failed to mark document deletion fence for index=%s path=%s file_id=%s: %s", - index_name, - path_or_url, - file_id, - exc, - ) - return False - - def clear_document_delete_marker( - self, - index_name: str, - path_or_url: Optional[str] = None, - file_id: Optional[str] = None, - ) -> int: - """Clear a path/file deletion fence before a new upload reuses it.""" - keys = self._document_delete_marker_keys(index_name, path_or_url, file_id) - if not keys: - return 0 - try: - return int(self.client.delete(*keys)) - except Exception as exc: - logger.warning( - "Failed to clear document deletion fence for index=%s path=%s file_id=%s: %s", - index_name, - path_or_url, - file_id, - exc, - ) - return 0 - - def is_document_delete_requested( - self, - index_name: Optional[str], - path_or_url: Optional[str] = None, - file_id: Optional[str] = None, - ) -> bool: - """Return whether a file is fenced from further processing.""" - keys = self._document_delete_marker_keys(index_name, path_or_url, file_id) - if not keys: - return False - try: - return any(bool(self.client.get(key)) for key in keys) - except Exception as exc: - # A Redis outage must preserve the legacy best-effort behavior. - logger.debug( - "Unable to check document deletion fence for index=%s path=%s file_id=%s: %s", - index_name, - path_or_url, - file_id, - exc, - ) - return False - - def _get_celery_control_app(self): - """Create a lightweight Celery control app without importing task modules.""" - if self._celery_control_app is None: - from celery import Celery - - broker_url = REDIS_URL - backend_url = REDIS_BACKEND_URL or REDIS_URL - if not broker_url: - raise ValueError("REDIS_URL is not configured") - self._celery_control_app = Celery( - "nexent-delete-control", - broker=broker_url, - backend=backend_url, - ) - return self._celery_control_app - - @staticmethod - def _task_kwargs(task: Dict[str, Any]) -> Dict[str, Any]: - kwargs = task.get("kwargs") or {} - if isinstance(kwargs, str): - try: - kwargs = json.loads(kwargs) - except (TypeError, json.JSONDecodeError): - kwargs = {} - return kwargs if isinstance(kwargs, dict) else {} - - @classmethod - def _runtime_task_matches_document( - cls, - task: Dict[str, Any], - index_name: str, - path_or_url: Optional[str], - file_id: Optional[str], - ) -> bool: - kwargs = cls._task_kwargs(task) - task_index = kwargs.get("index_name") - task_source = kwargs.get("source") or kwargs.get("path_or_url") - task_file_id = kwargs.get("file_id") - if file_id and task_file_id and task_file_id != file_id: - return False - if task_index != index_name: - return False - if file_id and task_file_id == file_id: - return True - return bool(path_or_url and task_source == path_or_url) - - def _collect_runtime_task_ids( - self, - index_name: str, - path_or_url: Optional[str], - file_id: Optional[str] = None, - ) -> Set[str]: - """Collect matching active/reserved task IDs for targeted revocation.""" - task_ids: Set[str] = set() - try: - inspector = self._get_celery_control_app().control.inspect(timeout=0.5) - for state_name in ("active", "reserved"): - snapshot = getattr(inspector, state_name)() or {} - for tasks in snapshot.values(): - for task in tasks or []: - if not isinstance(task, dict): - continue - task_id = task.get("id") - if task_id and self._runtime_task_matches_document( - task, index_name, path_or_url, file_id - ): - task_ids.add(str(task_id)) - except Exception as exc: - logger.warning( - "Failed to inspect runtime tasks for index=%s path=%s file_id=%s: %s", - index_name, - path_or_url, - file_id, - exc, - ) - return task_ids - - def _revoke_task(self, task_id: str) -> bool: - if not task_id: - return False - try: - self._get_celery_control_app().control.revoke( - task_id, - terminate=False, - ) - logger.info("Revoked Celery task %s", task_id) - return True - except Exception as exc: - logger.warning("Failed to revoke Celery task %s: %s", task_id, exc) - return False - - def prepare_document_deletion( - self, - index_name: str, - path_or_url: Optional[str] = None, - file_id: Optional[str] = None, - ) -> Dict[str, Any]: - """Fence a file and revoke discoverable runtime tasks before storage deletion.""" - result = { - "delete_fence_set": self.mark_document_delete_requested( - index_name=index_name, - path_or_url=path_or_url, - file_id=file_id, - ), - "runtime_tasks_found": 0, - "tasks_cancelled": 0, - "tasks_revoked": 0, - "warnings": [], - } - task_ids = self._collect_runtime_task_ids(index_name, path_or_url, file_id) - result["runtime_tasks_found"] = len(task_ids) - for task_id in task_ids: - if self.mark_task_cancelled(task_id): - result["tasks_cancelled"] += 1 - if self._revoke_task(task_id): - result["tasks_revoked"] += 1 - return result - def mark_task_cancelled(self, task_id: str, ttl_hours: int = 24) -> bool: """ Mark a Celery task as cancelled in Redis so that long-running diff --git a/backend/services/vectordatabase_service.py b/backend/services/vectordatabase_service.py index 7e68bdfc7c..4df5ec5795 100644 --- a/backend/services/vectordatabase_service.py +++ b/backend/services/vectordatabase_service.py @@ -2209,19 +2209,6 @@ def delete_lifecycle_record_without_object( tenant_id = lifecycle_record.get("tenant_id") index_name = lifecycle_record.get("index_name") - if index_name: - try: - get_redis_service().prepare_document_deletion( - index_name=index_name, - path_or_url=None, - file_id=file_id, - ) - except Exception as deletion_exc: - logger.warning( - "Failed to prepare deletion for lifecycle file %s: %s", - file_id, - deletion_exc, - ) current_status = str(lifecycle_record.get("status") or "").upper() deleteable_statuses = ( "UPLOADING", @@ -2307,22 +2294,7 @@ async def delete_document_by_scope( await ElasticSearchService._assert_source_only_deletable( index_name, path_or_url ) - lifecycle_record = ElasticSearchService._mark_file_delete_requested( - index_name, path_or_url - ) - try: - get_redis_service().prepare_document_deletion( - index_name=index_name, - path_or_url=path_or_url, - file_id=(lifecycle_record or {}).get("file_id"), - ) - except Exception as deletion_exc: - logger.warning( - "Failed to prepare source-only deletion for index=%s path=%s: %s", - index_name, - path_or_url, - deletion_exc, - ) + ElasticSearchService._mark_file_delete_requested(index_name, path_or_url) try: knowledge = get_knowledge_record({"index_name": index_name}) or {} except Exception: @@ -2350,22 +2322,7 @@ async def delete_document_by_scope( ), } - lifecycle_record = ElasticSearchService._mark_file_delete_requested( - index_name, path_or_url - ) - try: - get_redis_service().prepare_document_deletion( - index_name=index_name, - path_or_url=path_or_url, - file_id=(lifecycle_record or {}).get("file_id"), - ) - except Exception as deletion_exc: - logger.warning( - "Failed to prepare full deletion for index=%s path=%s: %s", - index_name, - path_or_url, - deletion_exc, - ) + ElasticSearchService._mark_file_delete_requested(index_name, path_or_url) result = ElasticSearchService.delete_documents( index_name, path_or_url, vdb_core ) diff --git a/test/backend/data_process/test_tasks.py b/test/backend/data_process/test_tasks.py index 16e9f945ec..fe82418a4a 100644 --- a/test/backend/data_process/test_tasks.py +++ b/test/backend/data_process/test_tasks.py @@ -2233,9 +2233,6 @@ def save_progress_info(self, *args, **kwargs): def is_task_cancelled(self, *args, **kwargs): return False - def is_document_delete_requested(self, *args, **kwargs): - return False - monkeypatch.setattr(tasks, "get_redis_service", lambda: _RedisSvc()) captured = {} @@ -2293,7 +2290,6 @@ def _fake_allow_join_result(): assert out["chunks_stored"] == 70 assert captured["callback"].queue == "forward_aggregate_q" - assert captured["callback"].kwargs["file_id"] == "file-1" def test_process_sync_unsupported_raises_and_updates_state(monkeypatch): @@ -2970,202 +2966,6 @@ def is_task_cancelled(self, _task_id): assert out["total_submitted"] == 0 -def test_forward_part_returns_cancelled_when_document_is_deleted(monkeypatch): - tasks, _ = import_tasks_with_fake_ray(monkeypatch) - - class _Svc: - def is_document_delete_requested(self, **kwargs): - return kwargs["file_id"] == "file-deleted" - - def is_task_cancelled(self, _task_id): - return False - - monkeypatch.setattr(tasks, "get_redis_service", lambda: _Svc()) - monkeypatch.setattr( - tasks, - "_send_chunks_to_es", - lambda **kwargs: (_ for _ in ()).throw(AssertionError("deleted batch must not be sent")), - ) - self = types.SimpleNamespace( - request=types.SimpleNamespace(id="fp-deleted", retries=0), - retry=lambda **kwargs: (_ for _ in ()).throw(AssertionError("must not retry")), - ) - - out = tasks.forward_part( - self, - chunks=[{"content": "x"}], - index_name="idx", - source="knowledge_base/deleted.txt", - file_id="file-deleted", - batch_index=1, - total_batches=1, - ) - - assert out["cancelled"] is True - assert out["total_indexed"] == 0 - - -def test_aggregate_forward_parts_returns_cancelled_when_document_is_deleted(monkeypatch): - tasks, _ = import_tasks_with_fake_ray(monkeypatch) - - class _Svc: - def is_document_delete_requested(self, **kwargs): - return kwargs["file_id"] == "file-deleted" - - monkeypatch.setattr(tasks, "get_redis_service", lambda: _Svc()) - out = tasks.aggregate_forward_parts( - types.SimpleNamespace(request=types.SimpleNamespace(id="aggregate-deleted")), - parts_results=[{"total_indexed": 10, "total_submitted": 10}], - source="knowledge_base/deleted.txt", - index_name="idx", - file_id="file-deleted", - ) - - assert out["cancelled"] is True - assert out["total_indexed"] == 0 - - -def test_process_returns_cancelled_when_document_is_deleted_before_start(monkeypatch): - tasks, _ = import_tasks_with_fake_ray(monkeypatch) - monkeypatch.setattr(tasks, "_is_document_delete_requested", lambda *args, **kwargs: True) - monkeypatch.setattr( - tasks, - "_update_file_lifecycle", - lambda **kwargs: (_ for _ in ()).throw(AssertionError("deleted process must not update lifecycle")), - ) - - out = tasks.process( - FakeSelf("process-deleted-before-start"), - source="knowledge_base/deleted.txt", - source_type="local", - index_name="idx", - original_filename="deleted.txt", - file_id="file-deleted", - ) - - assert out["cancelled"] is True - assert out["file_id"] == "file-deleted" - - -def test_process_returns_cancelled_when_deleted_after_local_extraction(monkeypatch, tmp_path): - tasks, _ = import_tasks_with_fake_ray(monkeypatch) - source = tmp_path / "deleted-after-extraction.txt" - source.write_text("content") - checks = iter([False, True]) - monkeypatch.setattr(tasks, "_is_document_delete_requested", lambda *args, **kwargs: next(checks)) - monkeypatch.setattr(tasks, "_update_file_lifecycle", lambda **kwargs: None) - monkeypatch.setattr( - tasks, - "_process_source_with_split", - lambda **kwargs: (False, [{"content": "chunk", "metadata": {}}], None), - ) - - out = tasks.process( - FakeSelf("process-deleted-after-local"), - source=str(source), - source_type="local", - chunking_strategy="basic", - index_name="idx", - file_id="file-deleted", - ) - - assert out["cancelled"] is True - - -def test_process_returns_cancelled_when_deleted_after_minio_extraction(monkeypatch): - tasks, _ = import_tasks_with_fake_ray(monkeypatch) - # The MinIO path performs the initial and post-extraction checks only. - checks = iter([False, True]) - monkeypatch.setattr(tasks, "_is_document_delete_requested", lambda *args, **kwargs: next(checks)) - monkeypatch.setattr(tasks, "_update_file_lifecycle", lambda **kwargs: None) - monkeypatch.setattr(tasks, "_fetch_minio_source", lambda source: b"content") - monkeypatch.setattr( - tasks, - "_process_source_with_split", - lambda **kwargs: (False, [{"content": "chunk", "metadata": {}}], None), - ) - - out = tasks.process( - FakeSelf("process-deleted-after-minio"), - source="knowledge_base/deleted.txt", - source_type="minio", - chunking_strategy="basic", - index_name="idx", - file_id="file-deleted", - ) - - assert out["cancelled"] is True - - -def test_forward_returns_cancelled_when_document_is_deleted_before_load(monkeypatch): - tasks, _ = import_tasks_with_fake_ray(monkeypatch) - - class _Service: - def is_task_cancelled(self, _task_id): - return False - - def is_document_delete_requested(self, **kwargs): - return True - - monkeypatch.setattr(tasks, "get_redis_service", lambda: _Service()) - monkeypatch.setattr(tasks, "_update_file_lifecycle", lambda **kwargs: None) - - out = tasks.forward( - FakeSelf("forward-deleted-before-load"), - processed_data={"chunks": [{"content": "x", "metadata": {}}]}, - index_name="idx", - source="knowledge_base/deleted.txt", - file_id="file-deleted", - ) - - assert out["cancelled"] is True - - -def test_forward_returns_cancelled_when_deleted_after_loading_chunks(monkeypatch): - tasks, _ = import_tasks_with_fake_ray(monkeypatch) - checks = iter([False, True]) - monkeypatch.setattr(tasks, "_is_document_delete_requested", lambda *args, **kwargs: next(checks)) - monkeypatch.setattr(tasks, "_update_file_lifecycle", lambda **kwargs: None) - monkeypatch.setattr( - tasks, - "_load_forward_chunks", - lambda _self, **kwargs: ([{"content": "x", "metadata": {}}], False, "knowledge_base/deleted.txt", "idx", "deleted.txt"), - ) - - out = tasks.forward( - FakeSelf("forward-deleted-after-load"), - processed_data={"chunks": [{"content": "x", "metadata": {}}]}, - index_name="idx", - source="knowledge_base/deleted.txt", - file_id="file-deleted", - ) - - assert out["cancelled"] is True - - -def test_forward_returns_cancelled_before_first_es_write(monkeypatch): - tasks, _ = import_tasks_with_fake_ray(monkeypatch) - checks = iter([False, False, True]) - monkeypatch.setattr(tasks, "_is_document_delete_requested", lambda *args, **kwargs: next(checks)) - monkeypatch.setattr(tasks, "_update_file_lifecycle", lambda **kwargs: None) - monkeypatch.setattr(tasks, "get_file_size", lambda *args, **kwargs: 0) - monkeypatch.setattr( - tasks, - "_load_forward_chunks", - lambda _self, **kwargs: ([{"content": "x", "metadata": {}}], False, "knowledge_base/deleted.txt", "idx", "deleted.txt"), - ) - - out = tasks.forward( - FakeSelf("forward-deleted-before-es"), - processed_data={"chunks": [{"content": "x", "metadata": {}}]}, - index_name="idx", - source="knowledge_base/deleted.txt", - file_id="file-deleted", - ) - - assert out["cancelled"] is True - - def test_forward_part_continues_when_parent_cancellation_lookup_fails(monkeypatch): tasks, _ = import_tasks_with_fake_ray(monkeypatch) monkeypatch.setattr(tasks, "get_redis_service", lambda: (_ for _ in ()).throw(RuntimeError("redis down"))) @@ -3515,23 +3315,6 @@ def _delete(*_a, **_k): assert called["delete"] == 0 -def test_cleanup_source_skips_when_forward_was_cancelled(monkeypatch): - tasks, _ = import_tasks_with_fake_ray(monkeypatch) - monkeypatch.setattr( - tasks, - "_delete_source_file_via_http_sync", - lambda **kwargs: (_ for _ in ()).throw(AssertionError("cancelled forward must not delete source")), - ) - - out = tasks.cleanup_source( - FakeSelf("cleanup-cancelled"), - {"cancelled": True, "index_name": "idx", "source": "/a.txt"}, - ) - - assert out["source_cleanup"]["skipped_reason"] == "forward_cancelled" - assert out["source_cleanup"]["attempted"] is False - - def test_cleanup_source_calls_delete_with_scope_source_only(monkeypatch): tasks, _ = import_tasks_with_fake_ray(monkeypatch) monkeypatch.setattr(tasks, "ELASTICSEARCH_SERVICE", "http://api") diff --git a/test/backend/services/test_data_process_service.py b/test/backend/services/test_data_process_service.py index 1e08c074f2..49d820e7df 100644 --- a/test/backend/services/test_data_process_service.py +++ b/test/backend/services/test_data_process_service.py @@ -2629,92 +2629,6 @@ async def _task_info(task_id): asyncio.run(_run()) - @patch('backend.services.data_process_service.get_redis_service') - @patch('backend.services.data_process_service.celery_app') - @patch('backend.services.data_process_service.get_all_task_ids_from_redis', return_value=['deleted-task']) - @patch('backend.services.data_process_service.get_task_info') - def test_get_all_tasks_hides_tasks_for_deleted_documents( - self, mock_get_task_info, _mock_ids, mock_celery_app, mock_get_redis_service - ): - mock_inspector = MagicMock() - mock_inspector.active.return_value = {} - mock_inspector.reserved.return_value = {} - mock_celery_app.control.inspect.return_value = mock_inspector - mock_celery_app.conf.broker_url = "redis://mock:6379/0" - mock_celery_app.conf.result_backend = "redis://mock:6379/0" - - async def _task_info(task_id): - return { - "id": task_id, - "task_name": "forward", - "index_name": "idx", - "path_or_url": "knowledge_base/deleted.txt", - "file_id": "file-deleted", - } - - mock_get_task_info.side_effect = _task_info - fence_service = MagicMock() - fence_service.is_document_delete_requested.return_value = True - mock_get_redis_service.return_value = fence_service - - rows = asyncio.run(self.service.get_all_tasks(filter=False)) - - self.assertEqual(rows, []) - fence_service.is_document_delete_requested.assert_called_once_with( - index_name="idx", - path_or_url="knowledge_base/deleted.txt", - file_id="file-deleted", - ) - - @patch('backend.services.data_process_service.get_redis_service', side_effect=RuntimeError("redis unavailable")) - @patch('backend.services.data_process_service.celery_app') - @patch('backend.services.data_process_service.get_all_task_ids_from_redis', return_value=[]) - @patch('backend.services.data_process_service.get_task_info') - def test_get_all_tasks_continues_when_delete_fence_service_unavailable( - self, mock_get_task_info, _mock_ids, mock_celery_app, _mock_get_redis_service - ): - mock_inspector = MagicMock() - mock_inspector.active.return_value = {} - mock_inspector.reserved.return_value = {} - mock_celery_app.control.inspect.return_value = mock_inspector - mock_celery_app.conf.broker_url = "redis://mock:6379/0" - mock_celery_app.conf.result_backend = "redis://mock:6379/0" - - rows = asyncio.run(self.service.get_all_tasks(filter=False)) - - assert rows == [] - mock_get_task_info.assert_not_called() - - @patch('backend.services.data_process_service.get_redis_service') - @patch('backend.services.data_process_service.celery_app') - @patch('backend.services.data_process_service.get_all_task_ids_from_redis', return_value=['fenced-task']) - @patch('backend.services.data_process_service.get_task_info') - def test_get_all_tasks_keeps_task_when_delete_fence_lookup_fails( - self, mock_get_task_info, _mock_ids, mock_celery_app, mock_get_redis_service - ): - mock_inspector = MagicMock() - mock_inspector.active.return_value = {} - mock_inspector.reserved.return_value = {} - mock_celery_app.control.inspect.return_value = mock_inspector - mock_celery_app.conf.broker_url = "redis://mock:6379/0" - mock_celery_app.conf.result_backend = "redis://mock:6379/0" - async def _task_info(task_id): - return { - "id": task_id, - "task_name": "forward", - "index_name": "idx", - "path_or_url": "knowledge_base/a.txt", - } - - mock_get_task_info.side_effect = _task_info - fence_service = MagicMock() - fence_service.is_document_delete_requested.side_effect = RuntimeError("lookup failed") - mock_get_redis_service.return_value = fence_service - - rows = asyncio.run(self.service.get_all_tasks(filter=False)) - - assert [row["id"] for row in rows] == ["fenced-task"] - # ---- Additional coverage tests ---- diff --git a/test/backend/services/test_file_management_service.py b/test/backend/services/test_file_management_service.py index 0660d12e14..63ba921a3f 100644 --- a/test/backend/services/test_file_management_service.py +++ b/test/backend/services/test_file_management_service.py @@ -568,20 +568,13 @@ def check_hard_limit_post_write(self, *_args, **_kwargs): quota_module = types.ModuleType("services.quota_service") quota_module.QuotaService = FakeQuota - redis_module = types.ModuleType("services.redis_service") - redis_module.get_redis_service = MagicMock( - side_effect=RuntimeError("redis unavailable") - ) with patch.object(knowledge_storage_stub, "resolve_storage_context", return_value=context), \ patch("backend.services.file_management_service.create_file_records", return_value=[record]) as create_records, \ patch("backend.services.file_management_service.transition_file_record", side_effect=transition) as transition_record, \ patch("backend.services.file_management_service.upload_to_minio", AsyncMock(return_value=[ {"success": True, "file_id": "fid-1", "file_name": "a.txt", "object_name": "folder/a.txt", "file_size": 3} ])) as upload_mock, \ - patch.dict(sys.modules, { - "services.quota_service": quota_module, - "services.redis_service": redis_module, - }): + patch.dict(sys.modules, {"services.quota_service": quota_module}): result = await upload_files_impl( destination="minio", file=[mock_file], folder="folder", index_name="kb-1", user_id="user-1" ) diff --git a/test/backend/services/test_redis_service.py b/test/backend/services/test_redis_service.py index a4a2392a36..4c3434249c 100644 --- a/test/backend/services/test_redis_service.py +++ b/test/backend/services/test_redis_service.py @@ -143,149 +143,6 @@ def test_mark_and_check_task_cancelled(self, mock_from_url): mock_client.setex.assert_called_once() mock_client.get.assert_called_once() - @patch('backend.services.redis_service.REDIS_URL', 'redis://localhost:6379/0') - def test_document_delete_fence_can_be_set_checked_and_cleared(self): - mock_client = MagicMock() - mock_client.get.side_effect = [b'{"deleted": true}', None] - mock_client.delete.return_value = 1 - service = RedisService() - service._client = mock_client - - self.assertTrue(service.mark_document_delete_requested( - "idx", "knowledge_base/a.txt", file_id="file-1", ttl_seconds=60 - )) - self.assertTrue(service.is_document_delete_requested( - "idx", "knowledge_base/a.txt", file_id="file-1" - )) - self.assertEqual(service.clear_document_delete_marker( - "idx", "knowledge_base/a.txt", file_id="file-1" - ), 1) - self.assertEqual(mock_client.setex.call_count, 2) - mock_client.delete.assert_called_once() - - def test_document_delete_fence_validation_and_redis_failures(self): - service = RedisService() - assert service.mark_document_delete_requested("idx") is False - assert service.mark_document_delete_requested( - "idx", "knowledge_base/a.txt", ttl_seconds=-1 - ) is False - - mock_client = MagicMock() - mock_client.setex.side_effect = RuntimeError("set failed") - service._client = mock_client - assert service.mark_document_delete_requested("idx", "knowledge_base/a.txt") is False - - mock_client.delete.side_effect = RuntimeError("delete failed") - assert service.clear_document_delete_marker("idx", "knowledge_base/a.txt") == 0 - mock_client.get.side_effect = RuntimeError("get failed") - assert service.is_document_delete_requested("idx", "knowledge_base/a.txt") is False - - def test_document_delete_fence_clear_without_keys(self): - service = RedisService() - assert service.clear_document_delete_marker("idx") == 0 - assert service.is_document_delete_requested(None) is False - - @patch('backend.services.redis_service.REDIS_URL', 'redis://broker/0') - @patch('backend.services.redis_service.REDIS_BACKEND_URL', 'redis://backend/1') - def test_celery_control_app_is_lazy_and_rejects_missing_broker(self): - service = RedisService() - app = service._get_celery_control_app() - assert app is service._get_celery_control_app() - - with patch('backend.services.redis_service.REDIS_URL', ''): - missing = RedisService() - with self.assertRaises(ValueError): - missing._get_celery_control_app() - - def test_task_kwargs_normalization(self): - assert RedisService._task_kwargs({"kwargs": '{"file_id": "f1"}'}) == {"file_id": "f1"} - assert RedisService._task_kwargs({"kwargs": "bad-json"}) == {} - assert RedisService._task_kwargs({"kwargs": ["not-a-dict"]}) == {} - assert RedisService._task_kwargs({}) == {} - - def test_collect_runtime_task_ids_filters_matches_and_handles_inspector_failure(self): - service = RedisService() - inspector = MagicMock() - inspector.active.return_value = { - "worker": [ - {"id": "match-file", "kwargs": {"index_name": "idx", "file_id": "f1"}}, - {"id": "wrong-index", "kwargs": {"index_name": "other", "file_id": "f1"}}, - {"kwargs": {"index_name": "idx", "file_id": "f1"}}, - "not-a-task", - ] - } - inspector.reserved.return_value = { - "worker": [{"id": "match-path", "kwargs": {"index_name": "idx", "source": "a.txt"}}] - } - control_app = MagicMock() - control_app.control.inspect.return_value = inspector - service._get_celery_control_app = MagicMock(return_value=control_app) - - assert service._collect_runtime_task_ids("idx", "a.txt", "f1") == {"match-file", "match-path"} - - service._get_celery_control_app.side_effect = RuntimeError("inspect failed") - assert service._collect_runtime_task_ids("idx", "a.txt", "f1") == set() - - def test_revoke_task_success_failure_and_empty_id(self): - service = RedisService() - assert service._revoke_task("") is False - control_app = MagicMock() - service._get_celery_control_app = MagicMock(return_value=control_app) - assert service._revoke_task("task-1") is True - control_app.control.revoke.assert_called_once_with("task-1", terminate=False) - control_app.control.revoke.side_effect = RuntimeError("revoke failed") - assert service._revoke_task("task-2") is False - - def test_prepare_document_deletion_marks_and_revokes_runtime_tasks(self): - service = RedisService() - service.mark_document_delete_requested = MagicMock(return_value=True) - service._collect_runtime_task_ids = MagicMock(return_value={"task-1", "task-2"}) - service.mark_task_cancelled = MagicMock(return_value=True) - service._revoke_task = MagicMock(return_value=True) - - result = service.prepare_document_deletion( - "idx", "knowledge_base/a.txt", file_id="file-1" - ) - - self.assertTrue(result["delete_fence_set"]) - self.assertEqual(result["runtime_tasks_found"], 2) - self.assertEqual(result["tasks_cancelled"], 2) - self.assertEqual(result["tasks_revoked"], 2) - self.assertEqual(service.mark_task_cancelled.call_count, 2) - self.assertEqual(service._revoke_task.call_count, 2) - - def test_prepare_document_deletion_keeps_zero_counts_when_cancel_or_revoke_fails(self): - service = RedisService() - service.mark_document_delete_requested = MagicMock(return_value=False) - service._collect_runtime_task_ids = MagicMock(return_value={"task-1"}) - service.mark_task_cancelled = MagicMock(return_value=False) - service._revoke_task = MagicMock(return_value=False) - - result = service.prepare_document_deletion("idx", "a.txt", file_id="f1") - - assert result["delete_fence_set"] is False - assert result["runtime_tasks_found"] == 1 - assert result["tasks_cancelled"] == 0 - assert result["tasks_revoked"] == 0 - - def test_runtime_task_matches_document_by_file_id_or_path(self): - task = { - "kwargs": json.dumps({ - "index_name": "idx", - "source": "knowledge_base/a.txt", - "file_id": "file-1", - }) - } - self.assertTrue(RedisService._runtime_task_matches_document( - task, "idx", "different-path", "file-1" - )) - self.assertTrue(RedisService._runtime_task_matches_document( - task, "idx", "knowledge_base/a.txt", None - )) - self.assertFalse(RedisService._runtime_task_matches_document( - task, "idx", "knowledge_base/a.txt", "file-2" - )) - def test_delete_knowledgebase_records(self): """Test delete_knowledgebase_records method""" # Setup diff --git a/test/backend/services/test_vectordatabase_service.py b/test/backend/services/test_vectordatabase_service.py index 010367bfc3..fc87d000ca 100644 --- a/test/backend/services/test_vectordatabase_service.py +++ b/test/backend/services/test_vectordatabase_service.py @@ -2456,12 +2456,9 @@ def test_mark_file_deleted_ignores_missing_record( ElasticSearchService._mark_file_deleted("test_index", "knowledge_base/missing.txt") mock_delete.assert_not_called() - @patch('backend.services.vectordatabase_service.get_redis_service', side_effect=RuntimeError("redis unavailable")) @patch('backend.services.vectordatabase_service.delete_file_record', return_value=True) @patch('backend.services.vectordatabase_service.transition_file_record') - def test_delete_lifecycle_record_without_object_hard_deletes( - self, mock_transition, mock_delete, _mock_redis - ): + def test_delete_lifecycle_record_without_object_hard_deletes(self, mock_transition, mock_delete): mock_transition.return_value = {"file_id": "fid-no-object", "status": "DELETE_REQUESTED"} result = ElasticSearchService.delete_lifecycle_record_without_object( @@ -2785,14 +2782,13 @@ def test_delete_source_file_releases_charge_only_after_minio_success( updated_by=None, ) - @patch('backend.services.vectordatabase_service.get_redis_service', side_effect=RuntimeError("redis unavailable")) @patch( 'backend.services.vectordatabase_service.get_all_files_status', new_callable=AsyncMock, ) @patch('backend.services.vectordatabase_service.delete_file') def test_delete_document_by_scope_source_only( - self, mock_delete_file, mock_get_status, _mock_redis + self, mock_delete_file, mock_get_status ): mock_get_status.return_value = { "knowledge_base/doc.pdf": {"state": "COMPLETED"} @@ -2817,9 +2813,8 @@ def test_delete_document_by_scope_source_only( 'delete_documents', return_value={"status": "success", "deleted_minio": True}, ) - @patch('backend.services.vectordatabase_service.get_redis_service', side_effect=RuntimeError("redis unavailable")) def test_delete_document_by_scope_does_not_require_storage_ledger( - self, _mock_redis, mock_delete_documents + self, mock_delete_documents ): result = asyncio.run( ElasticSearchService.delete_document_by_scope(