diff --git a/api/apps/services/dataset_api_service.py b/api/apps/services/dataset_api_service.py index 002de87a98..caf002f94d 100644 --- a/api/apps/services/dataset_api_service.py +++ b/api/apps/services/dataset_api_service.py @@ -2687,50 +2687,98 @@ async def delete_nav_node(dataset_id: str, tenant_id: str, name: str): index_nm, _ = pack from common.doc_store.doc_store_base import OrderByExpr + from rag.advanced_rag.knowlege_compile.dataset_nav import ( + _LOCK_BLOCKING_TIMEOUT_S, + _LOCK_TIMEOUT_S, + _nav_lock_key, + _remove_dataset_nav_doc_locked, + ) + from rag.utils.redis_conn import RedisDistributedLock - # Collect the node plus every descendant, level by level via parent_kwd. - names: set[str] = {name} - frontier: list[str] = [name] - for _ in range(64): # depth guard against a malformed (cyclic) tree - if not frontier: - break - try: - res = await thread_pool_exec( - settings.docStoreConn.search, - ["name"], - [], - {"compile_kwd": [_NAV_COMPILE_KWD], "parent_kwd": frontier}, - [], - OrderByExpr(), - 0, - 10000, - index_nm, - [dataset_id], - ) - rows = settings.docStoreConn.get_fields(res, ["name"]) or {} - except Exception: - logging.exception("delete_nav_node: subtree scan failed for kb=%s name=%s", dataset_id, name) - break - nxt: list[str] = [] - for row in rows.values(): - child = row.get("name") - if isinstance(child, str) and child and child not in names: - names.add(child) - nxt.append(child) - frontier = nxt + async def search_rows(condition): + res = await thread_pool_exec( + settings.docStoreConn.search, + ["id", "name", "type_kwd", "parent_kwd", "doc_id"], + [], + condition, + [], + OrderByExpr(), + 0, + 10000, + index_nm, + [dataset_id], + ) + return list((settings.docStoreConn.get_fields(res, ["id", "name", "type_kwd", "parent_kwd", "doc_id"]) or {}).values()) + + lock = RedisDistributedLock( + _nav_lock_key(dataset_id), + timeout=_LOCK_TIMEOUT_S, + blocking_timeout=_LOCK_BLOCKING_TIMEOUT_S, + ) + try: + await lock.spin_acquire() + except Exception: + logging.exception("delete_nav_node: lock acquire failed for kb=%s", dataset_id) + return False, "Failed to acquire the navigation tree lock." try: - deleted = await thread_pool_exec( - settings.docStoreConn.delete, - {"compile_kwd": [_NAV_COMPILE_KWD], "name": list(names)}, - index_nm, - dataset_id, - ) - except Exception: - logging.exception("delete_nav_node: docStore delete failed for kb=%s name=%s", dataset_id, name) - return False, "Failed to delete the navigation node." + # ES requires the keyword subfield for exact node-name matching. Infinity + # has a scalar varchar name field, while parent_kwd is exact-matchable. + name_field = "name.keyword" if settings.DOC_ENGINE.lower() in {"elasticsearch", "opensearch"} else "name" + target_condition = {"compile_kwd": [_NAV_COMPILE_KWD], name_field: [name]} + rows = await search_rows(target_condition) + if not rows and name_field != "name": + # Keep this fallback for installations with an older/missing mapping. + rows = [row for row in await search_rows({"compile_kwd": [_NAV_COMPILE_KWD]}) if row.get("name") == name] + if not rows: + return True, {"deleted": 0} - return True, {"deleted": int(deleted or 0)} + # Collect the target and descendants by stable row IDs, not by names. + # Names are display values and may collide across malformed/old data. + rows_by_id = {row.get("id"): row for row in rows if row.get("id")} + frontier = [row.get("name") for row in rows if row.get("name")] + for _ in range(64): + if not frontier: + break + children = await search_rows( + {"compile_kwd": [_NAV_COMPILE_KWD], "parent_kwd": frontier}, + ) + frontier = [] + for child in children: + row_id = child.get("id") + child_name = child.get("name") + if row_id and row_id not in rows_by_id: + rows_by_id[row_id] = child + if child_name: + frontier.append(child_name) + + deleted = 0 + # Reuse the existing document-removal path so direct parents are + # updated and empty ancestor clusters are cleaned consistently. + for row in rows_by_id.values(): + if row.get("type_kwd") == "nav_doc" and row.get("doc_id"): + await _remove_dataset_nav_doc_locked(tenant_id, dataset_id, row["doc_id"]) + deleted += 1 + + cluster_ids = [row_id for row_id, row in rows_by_id.items() if row.get("type_kwd") == "nav_cluster" and row_id] + if cluster_ids: + deleted_rows = await thread_pool_exec( + settings.docStoreConn.delete, + {"compile_kwd": [_NAV_COMPILE_KWD], "id": cluster_ids}, + index_nm, + dataset_id, + ) + deleted += int(deleted_rows or 0) + + return True, {"deleted": deleted} + except Exception: + logging.exception("delete_nav_node: deletion failed for kb=%s name=%s", dataset_id, name) + return False, "Failed to delete the navigation node." + finally: + try: + lock.release() + except Exception: + logging.exception("delete_nav_node: lock release failed for kb=%s", dataset_id) async def update_wiki_page(