"""Projects, Graph Releases, and Ontology Schema endpoints.""" import asyncio import json from pathlib import Path from typing import Any from fastapi import APIRouter, Depends, HTTPException from psycopg.types.json import Jsonb from app.auth import CurrentUser from app.config import settings from app.db import get_conn from app.project_context import ProjectContext, get_project_context from app.project_lifecycle import ( GraphImportError, ProjectValidationError, SAFE_RESOURCE_ID, SAFE_SCHEMA_VERSION, delete_falkor_graph, import_falkor_graph, list_falkor_graphs, normalize_provision_payload, ) router = APIRouter() ROOT_DIR = Path(__file__).resolve().parents[2] async def _run_blocking_to_completion(function: Any, *args: Any) -> Any: """Run a blocking FalkorDB call without abandoning its worker on cancellation. ``asyncio.to_thread`` keeps running after its awaiting coroutine is cancelled. For graph mutations that would otherwise release an advisory lock too early and potentially affect a newly-created graph with the same name. Shielding and draining the worker keeps the lock owner alive until the exact external operation has actually stopped. """ worker = asyncio.create_task(asyncio.to_thread(function, *args)) try: return await asyncio.shield(worker) except asyncio.CancelledError as cancelled: while not worker.done(): try: await asyncio.shield(worker) except asyncio.CancelledError: continue except Exception: break # Retrieve a possible worker exception so it is not reported as an # unhandled task. The durable cleanup intent remains pending and the # caller's cancellation remains the externally visible outcome. if worker.done() and not worker.cancelled(): try: worker.result() except Exception: pass raise cancelled async def _run_async_cleanup_to_completion(coroutine: Any) -> Any: """Drain an async cleanup task before propagating caller cancellation.""" worker = asyncio.create_task(coroutine) try: return await asyncio.shield(worker) except asyncio.CancelledError as cancelled: while not worker.done(): try: await asyncio.shield(worker) except asyncio.CancelledError: continue except BaseException: break if worker.done() and not worker.cancelled(): try: worker.result() except BaseException: pass raise cancelled async def _release_session_locks(conn: Any, lock_keys: list[str]) -> None: """Release session locks or close the physical connection unconditionally.""" try: if conn.info.transaction_status != 0: await conn.rollback() async with conn.cursor() as cur: for lock_key in reversed(lock_keys): await cur.execute("SELECT pg_advisory_unlock(hashtext(%s))", (lock_key,)) await conn.commit() except BaseException: # A failed or cancelled unlock must never return this physical session # to the pool with a lock still attached. Drain close even if a second # cancellation arrives while it is in progress. close_worker = asyncio.create_task(conn.close()) while not close_worker.done(): try: await asyncio.shield(close_worker) except asyncio.CancelledError: continue except BaseException: break if close_worker.done() and not close_worker.cancelled(): try: close_worker.result() except BaseException: pass def _as_dict(value: Any) -> dict[str, Any]: if isinstance(value, dict): return value if isinstance(value, str): try: parsed = json.loads(value) return parsed if isinstance(parsed, dict) else {} except json.JSONDecodeError: return {} return {} def _named_items(value: Any) -> list[tuple[str, dict[str, Any]]]: if isinstance(value, dict): return [(str(name), _as_dict(meta)) for name, meta in value.items()] if isinstance(value, list): items: list[tuple[str, dict[str, Any]]] = [] for index, item in enumerate(value): meta = _as_dict(item) name = str(meta.get("name") or meta.get("entity_type") or meta.get("relation_type") or f"Item{index + 1}") items.append((name, meta)) return items return [] def _field_names(value: Any) -> list[str]: if isinstance(value, dict): return [str(key) for key in value.keys()] if isinstance(value, list): return [str(item) for item in value] return [] def _schema_source_from_json(schema: dict[str, Any]) -> str: namespace = str(schema.get("namespace") or "schema") version = str(schema.get("version") or "") lines = [f"namespace {namespace}"] if version: lines.append(f"version {version}") purpose = str(schema.get("purpose") or "").strip() if purpose: lines.extend(["", f"// {purpose}"]) for name, meta in _named_items(schema.get("entity_types")): definition = str(meta.get("definition") or meta.get("description") or "").strip() lines.extend(["", f"entity {name}" + (f" // {definition}" if definition else "")]) primary_key = str(meta.get("primary_key") or meta.get("primaryKey") or "").strip() if primary_key: lines.append(f" primary_key {primary_key}") for field in _field_names(meta.get("fields")): lines.append(f" property {field}: Text") for name, meta in _named_items(schema.get("relation_types")): source = str(meta.get("from") or meta.get("source") or "").strip() target = str(meta.get("to") or meta.get("target") or "").strip() definition = str(meta.get("definition") or meta.get("description") or "").strip() arrow = f": {source} -> {target}" if source or target else "" lines.extend(["", f"relation {name}{arrow}" + (f" // {definition}" if definition else "")]) for prop in _field_names(meta.get("properties")): lines.append(f" property {prop}: Text") enums = _as_dict(schema.get("enums")) if enums: lines.append("") for name, values in enums.items(): lines.append(f"enum {name} = {' | '.join(_field_names(values))}") recipes = _as_dict(schema.get("query_recipes")) if recipes: lines.append("") lines.append("// query_recipes") for name, query in recipes.items(): lines.append(f"// {name}: {query}") return "\n".join(lines).strip() + "\n" def _schema_source_from_file(row: dict[str, Any], schema: dict[str, Any]) -> str | None: project_id = str(row.get("project_id") or "").strip() namespace = str(row.get("namespace") or schema.get("namespace") or "").strip() version = str(schema.get("version") or row.get("version") or "").strip() if ( not SAFE_RESOURCE_ID.fullmatch(project_id) or not SAFE_RESOURCE_ID.fullmatch(namespace) or not SAFE_SCHEMA_VERSION.fullmatch(version) ): return None schema_root = (ROOT_DIR / "schema搭建").resolve() schema_dir = (schema_root / project_id).resolve() if schema_root not in schema_dir.parents: return None version_slug = version.replace(".", "_") candidates = [ schema_dir / f"{namespace}_schema.current.dsl.md", schema_dir / f"{namespace}_schema.v{version_slug}.dsl.md", schema_dir / f"{project_id}_schema.current.dsl.md", schema_dir / f"{project_id}_schema.v{version_slug}.dsl.md", ] candidates.extend(sorted(schema_dir.glob("*.current.dsl.md"))) candidates.extend(sorted(schema_dir.glob(f"*v{version_slug}*.dsl.md"))) candidates.extend(sorted(schema_dir.glob("*.dsl.md"))) seen: set[Path] = set() for path in candidates: if path in seen: continue seen.add(path) resolved_path = path.resolve() if schema_root not in resolved_path.parents: continue if resolved_path.exists() and resolved_path.is_file(): return resolved_path.read_text(encoding="utf-8") return None def _attach_schema_source(row: Any) -> dict[str, Any]: data = dict(row) schema = _as_dict(data.get("schema_jsonb")) semantic_version = str(schema.get("version") or "").strip() if semantic_version: # Restored legacy databases keep an integer revision in the physical # ``version`` column. The JSON document remains the source of truth for # the user-facing semantic version. data["version"] = semantic_version data["schema_source"] = _schema_source_from_file(data, schema) or _schema_source_from_json(schema) data["schema_source_format"] = "dsl" return data # Explicitly reviewed graph-project tables. Relational data-platform registry # tables (project_databases/project_table_definitions/project_table_overrides) # and every biz_* schema are intentionally absent from this allow-list. _PROJECT_SCOPED_TABLES = ( "mapping_profiles", "import_templates", "import_batches", "candidate_entities", "candidate_relations", "review_actions", "publish_jobs", "ontology_schema_publish_jobs", "field_proposals", "source_profiles", "question_traces", "simulation_questions", "acquisition_tasks", "inventory_issues", "vocabulary_terms", "agent_call_logs", "kg_evidence_links", "kg_statements", "kg_relations", "kg_events", "kg_concepts", "kg_entities", "kg_schema_proposals", "kg_schema_versions", "kg_place_spatial", "kg_route_metrics", "graph_releases", "ontology_schemas", ) _PROJECT_DELETE_ORDER = ( "review_actions", "publish_jobs", "candidate_relations", "candidate_entities", "mapping_profiles", "import_templates", "field_proposals", "source_profiles", "inventory_issues", "question_traces", "simulation_questions", "acquisition_tasks", "vocabulary_terms", "agent_call_logs", "kg_evidence_links", "kg_statements", "kg_relations", "kg_events", "kg_concepts", "kg_entities", "kg_schema_proposals", "kg_schema_versions", "kg_place_spatial", "kg_route_metrics", "import_batches", "ontology_schema_publish_jobs", "graph_releases", "ontology_schemas", ) _INDIRECT_TABLES = ( "raw_records", "validation_issues", "ontology_entity_types", "ontology_fields", "ontology_relation_types", "task_events", "notifications", "super_agent_tasks", ) _GRAPH_BOUND_TABLES = ( "amap_spatial_collect_tasks", "amap_spatial_pois", "amap_spatial_collect_runs", ) def _project_where(columns: set[str], tenant_id: str, project_id: str) -> tuple[str, tuple[str, ...]]: if "tenant_id" in columns: return "tenant_id=%s AND project_id=%s", (tenant_id, project_id) return "project_id=%s", (project_id,) def _require_project_admin(user: dict[str, Any]) -> None: if "admin" not in set(user.get("roles") or []): raise HTTPException(403, "只有系统管理员可以创建、重命名或删除知识图谱项目") def _resolve_graph_name(project_id: str, explicit_graph_name: Any = None) -> str: """Use one stable implicit graph namespace across lifecycle endpoints.""" return str(explicit_graph_name or "").strip() or project_id _INTEGER_VERSION_TYPES = frozenset({"smallint", "integer", "bigint", "int2", "int4", "int8"}) _TEXT_VERSION_TYPES = frozenset({"text", "character varying", "character", "varchar", "bpchar"}) async def _ontology_schema_storage_version( cur: Any, *, tenant_id: str, project_id: str, namespace: str, semantic_version: str, ) -> str | int: """Adapt a semantic schema version to either legacy integer or text storage. Legacy snapshots use an integer revision column while the normalized JSON already carries the actual semantic version. A project-scoped advisory lock is held by every caller, so allocating the next integer revision here is deterministic without changing existing rows or their schema ids. """ await cur.execute( """SELECT data_type, udt_name FROM information_schema.columns WHERE table_schema=%s AND table_name='ontology_schemas' AND column_name='version'""", (settings.db_schema,), ) column = await cur.fetchone() if not column: raise RuntimeError("ontology_schemas.version 列不存在") column_types = { str(column.get("data_type") or "").lower(), str(column.get("udt_name") or "").lower(), } if column_types & _INTEGER_VERSION_TYPES: await cur.execute( f"""SELECT COALESCE(MAX(version), 0) + 1 AS next_version FROM {settings.db_schema}.ontology_schemas WHERE tenant_id=%s AND project_id=%s AND namespace=%s""", (tenant_id, project_id, namespace), ) next_row = await cur.fetchone() return int((next_row or {}).get("next_version") or 1) if column_types & _TEXT_VERSION_TYPES: return semantic_version rendered_type = str(column.get("data_type") or column.get("udt_name") or "unknown") raise RuntimeError(f"ontology_schemas.version 列类型 {rendered_type!r} 不受支持") def _graph_resource_conflict_detail(payload: dict[str, Any], state: str) -> str: """Describe a graph namespace conflict in terms used by the project form.""" project_id = str(payload["project_id"]) graph_name = str(payload["graph_name"]) if graph_name == project_id: return f"项目英文标识“{project_id}”对应的图谱资源{state},请更换项目英文标识" return ( f"项目英文标识“{project_id}”显式指定的图谱资源“{graph_name}”{state}," "请改用未占用的图谱资源标识" ) async def _load_table_shapes(cur: Any) -> dict[str, set[str]]: names = list(_PROJECT_SCOPED_TABLES + _INDIRECT_TABLES + _GRAPH_BOUND_TABLES + ("projects",)) await cur.execute( """SELECT table_name, array_agg(column_name::text ORDER BY ordinal_position) AS columns FROM information_schema.columns WHERE table_schema=%s AND table_name=ANY(%s) GROUP BY table_name""", (settings.db_schema, names), ) return {str(row["table_name"]): set(row["columns"] or []) for row in await cur.fetchall()} async def _load_project(cur: Any, project_id: str, *, for_update: bool = False) -> dict[str, Any]: suffix = " FOR UPDATE" if for_update else "" await cur.execute( f"SELECT * FROM {settings.db_schema}.projects WHERE project_id=%s ORDER BY created_at LIMIT 2{suffix}", (project_id,), ) rows = await cur.fetchall() if not rows: raise HTTPException(404, "项目不存在") if len(rows) > 1: raise HTTPException(409, "存在多个同 project_id 的租户项目,无法安全执行该操作") return dict(rows[0]) async def _fetch_ids( cur: Any, table: str, id_column: str, columns: set[str], tenant_id: str, project_id: str, ) -> list[Any]: if not columns or id_column not in columns or "project_id" not in columns: return [] where, params = _project_where(columns, tenant_id, project_id) await cur.execute( f"SELECT {id_column} AS value FROM {settings.db_schema}.{table} WHERE {where}", params, ) return [row["value"] for row in await cur.fetchall()] async def _dependency_ids( cur: Any, shapes: dict[str, set[str]], tenant_id: str, project_id: str, ) -> dict[str, list[Any]]: batch_columns = shapes.get("import_batches", set()) batch_key = "id" if "id" in batch_columns else "batch_id" return { "schema_ids": await _fetch_ids( cur, "ontology_schemas", "id", shapes.get("ontology_schemas", set()), tenant_id, project_id ), "batch_ids": await _fetch_ids( cur, "import_batches", batch_key, batch_columns, tenant_id, project_id ), "template_ids": await _fetch_ids( cur, "import_templates", "template_id", shapes.get("import_templates", set()), tenant_id, project_id, ), "candidate_entity_ids": await _fetch_ids( cur, "candidate_entities", "id", shapes.get("candidate_entities", set()), tenant_id, project_id ), "candidate_relation_ids": await _fetch_ids( cur, "candidate_relations", "id", shapes.get("candidate_relations", set()), tenant_id, project_id ), "task_ids": await _fetch_ids( cur, "acquisition_tasks", "id", shapes.get("acquisition_tasks", set()), tenant_id, project_id ), } async def _entity_type_ids( cur: Any, shapes: dict[str, set[str]], schema_ids: list[Any], ) -> list[Any]: columns = shapes.get("ontology_entity_types", set()) if not schema_ids or not {"id", "schema_id"}.issubset(columns): return [] await cur.execute( f"SELECT id AS value FROM {settings.db_schema}.ontology_entity_types WHERE schema_id=ANY(%s)", (schema_ids,), ) return [row["value"] for row in await cur.fetchall()] async def _graph_references( cur: Any, tenant_id: str, project_id: str, ) -> list[dict[str, Any]]: await cur.execute( f"SELECT DISTINCT graph_name FROM {settings.db_schema}.graph_releases " "WHERE tenant_id=%s AND project_id=%s AND graph_name IS NOT NULL ORDER BY graph_name", (tenant_id, project_id), ) graph_names = [str(row["graph_name"]) for row in await cur.fetchall()] result: list[dict[str, Any]] = [] for graph_name in graph_names: await cur.execute( f"""SELECT DISTINCT project_id FROM {settings.db_schema}.graph_releases WHERE graph_name=%s AND NOT (tenant_id=%s AND project_id=%s) ORDER BY project_id""", (graph_name, tenant_id, project_id), ) other_projects = [str(row["project_id"]) for row in await cur.fetchall()] result.append( { "graph_name": graph_name, "shared": bool(other_projects), "other_project_ids": other_projects, "will_delete": not other_projects, } ) return result async def _count_direct_table( cur: Any, table: str, columns: set[str], tenant_id: str, project_id: str, ) -> int: if "project_id" not in columns: return 0 where, params = _project_where(columns, tenant_id, project_id) await cur.execute(f"SELECT count(*) AS count FROM {settings.db_schema}.{table} WHERE {where}", params) return int((await cur.fetchone())["count"]) async def _count_any(cur: Any, table: str, column: str, values: list[Any]) -> int: if not values: return 0 await cur.execute( f"SELECT count(*) AS count FROM {settings.db_schema}.{table} WHERE {column}=ANY(%s)", (values,), ) return int((await cur.fetchone())["count"]) async def _count_validation_issues( cur: Any, columns: set[str], ids: dict[str, list[Any]], ) -> int: clauses: list[str] = [] params: list[list[Any]] = [] candidates = ( ("batch_id", ids["batch_ids"]), ("candidate_entity_id", ids["candidate_entity_ids"]), ("candidate_relation_id", ids["candidate_relation_ids"]), ("candidate_id", ids["candidate_entity_ids"]), ) for column, values in candidates: if column in columns and values: clauses.append(f"{column}=ANY(%s)") params.append(values) if not clauses: return 0 await cur.execute( f"SELECT count(*) AS count FROM {settings.db_schema}.validation_issues WHERE " + " OR ".join(clauses), tuple(params), ) return int((await cur.fetchone())["count"]) async def _count_legacy_review_actions( cur: Any, columns: set[str], ids: dict[str, list[Any]], ) -> int: if "project_id" in columns or not {"candidate_id", "candidate_type"}.issubset(columns): return 0 clauses: list[str] = [] params: list[Any] = [] if ids["candidate_entity_ids"]: clauses.append("candidate_type IN ('entity', 'candidate_entity') AND candidate_id=ANY(%s)") params.append(ids["candidate_entity_ids"]) if ids["candidate_relation_ids"]: clauses.append("candidate_type IN ('relation', 'candidate_relation') AND candidate_id=ANY(%s)") params.append(ids["candidate_relation_ids"]) if not clauses: return 0 await cur.execute( f"SELECT count(*) AS count FROM {settings.db_schema}.review_actions WHERE " + " OR ".join(f"({clause})" for clause in clauses), tuple(params), ) return int((await cur.fetchone())["count"]) async def _build_deletion_impact( cur: Any, project: dict[str, Any], ) -> tuple[dict[str, Any], dict[str, set[str]], dict[str, list[Any]]]: tenant_id = str(project["tenant_id"]) project_id = str(project["project_id"]) shapes = await _load_table_shapes(cur) ids = await _dependency_ids(cur, shapes, tenant_id, project_id) entity_type_ids = await _entity_type_ids(cur, shapes, ids["schema_ids"]) ids["entity_type_ids"] = entity_type_ids counts: dict[str, int] = {"projects": 1} for table in _PROJECT_SCOPED_TABLES: if table in shapes: count = await _count_direct_table(cur, table, shapes[table], tenant_id, project_id) if count: counts[table] = count indirect_specs = ( ("mapping_profiles", "template_id", ids["template_ids"]), ("raw_records", "batch_id", ids["batch_ids"]), ("ontology_schema_publish_jobs", "schema_id", ids["schema_ids"]), ("ontology_entity_types", "schema_id", ids["schema_ids"]), ("ontology_fields", "entity_type_id", entity_type_ids), ("ontology_relation_types", "schema_id", ids["schema_ids"]), ("task_events", "task_id", ids["task_ids"]), ("notifications", "related_task_id", ids["task_ids"]), ("super_agent_task_links", "related_task_id", ids["task_ids"]), ) for result_key, column, values in indirect_specs: table = "super_agent_tasks" if result_key == "super_agent_task_links" else result_key if ( table in shapes and column in shapes[table] and not ( table in {"mapping_profiles", "ontology_schema_publish_jobs"} and "project_id" in shapes[table] ) ): count = await _count_any(cur, table, column, values) if count: counts[result_key] = count if "validation_issues" in shapes: count = await _count_validation_issues(cur, shapes["validation_issues"], ids) if count: counts["validation_issues"] = count if "review_actions" in shapes and "project_id" not in shapes["review_actions"]: count = await _count_legacy_review_actions(cur, shapes["review_actions"], ids) if count: counts["review_actions"] = count graphs = await _graph_references(cur, tenant_id, project_id) exclusive_graph_names = [item["graph_name"] for item in graphs if item["will_delete"]] for table in _GRAPH_BOUND_TABLES: if table not in shapes or "graph_name" not in shapes[table] or not exclusive_graph_names: continue count = await _count_any(cur, table, "graph_name", exclusive_graph_names) if count: counts[table] = count impact = { "project_id": project_id, "tenant_id": tenant_id, "display_name": project.get("display_name") or project_id, "counts": counts, "graphs": graphs, "protected_resources": [ "PostgreSQL 数据中心注册信息(project_databases)", "数据中心表定义与字段配置", "全部 biz_* Schema、业务表与业务记录", "被其他知识图谱项目共享的 FalkorDB 图", ], } return impact, shapes, ids async def _ensure_cleanup_table(cur: Any) -> None: await cur.execute( f"""CREATE TABLE IF NOT EXISTS {settings.db_schema}.graph_cleanup_tasks ( id BIGSERIAL PRIMARY KEY, project_id TEXT NOT NULL, graph_name TEXT NOT NULL, operation TEXT NOT NULL DEFAULT 'delete', status TEXT NOT NULL DEFAULT 'pending', attempts INTEGER NOT NULL DEFAULT 0, last_error TEXT, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now() )""" ) async def _record_cleanup_task( cur: Any, project_id: str, graph_name: str, *, operation: str = "delete", ) -> int: await _ensure_cleanup_table(cur) await cur.execute( f"""INSERT INTO {settings.db_schema}.graph_cleanup_tasks (project_id, graph_name, operation, status) VALUES (%s, %s, %s, 'pending') RETURNING id""", (project_id, graph_name, operation), ) return int((await cur.fetchone())["id"]) async def _set_cleanup_task(cur: Any, task_id: int, *, error: str | None = None) -> None: await cur.execute( f"""UPDATE {settings.db_schema}.graph_cleanup_tasks SET status=%s, attempts=attempts+1, last_error=%s, updated_at=now() WHERE id=%s""", ("pending" if error else "completed", error, task_id), ) async def _finish_cleanup_task( task_id: int, *, error: str | None = None, conn: Any | None = None, ) -> None: """Persist cleanup state, optionally reusing a lock-owning connection.""" if conn is not None: async with conn.cursor() as cur: await _set_cleanup_task(cur, task_id, error=error) await conn.commit() return async with get_conn() as task_conn: async with task_conn.cursor() as cur: await _set_cleanup_task(cur, task_id, error=error) await task_conn.commit() async def _delete_any(cur: Any, table: str, column: str, values: list[Any]) -> None: if values: await cur.execute( f"DELETE FROM {settings.db_schema}.{table} WHERE {column}=ANY(%s)", (values,), ) async def _delete_validation_issues( cur: Any, columns: set[str], ids: dict[str, list[Any]], ) -> None: clauses: list[str] = [] params: list[list[Any]] = [] for column, values in ( ("batch_id", ids["batch_ids"]), ("candidate_entity_id", ids["candidate_entity_ids"]), ("candidate_relation_id", ids["candidate_relation_ids"]), ("candidate_id", ids["candidate_entity_ids"]), ): if column in columns and values: clauses.append(f"{column}=ANY(%s)") params.append(values) if clauses: await cur.execute( f"DELETE FROM {settings.db_schema}.validation_issues WHERE " + " OR ".join(clauses), tuple(params), ) async def _delete_project_resources( cur: Any, project: dict[str, Any], impact: dict[str, Any], shapes: dict[str, set[str]], ids: dict[str, list[Any]], ) -> dict[str, int]: """Delete only reviewed graph-side resources inside the current transaction.""" tenant_id = str(project["tenant_id"]) project_id = str(project["project_id"]) before = dict(impact["counts"]) question_columns = shapes.get("question_traces", set()) if "acquisition_task_id" in question_columns and "project_id" in question_columns: where, params = _project_where(question_columns, tenant_id, project_id) await cur.execute( f"UPDATE {settings.db_schema}.question_traces SET acquisition_task_id=NULL WHERE {where}", params, ) task_columns = shapes.get("acquisition_tasks", set()) reset_columns = [ column for column in ("triggered_by_trace_id", "result_batch_id") if column in task_columns ] if reset_columns and "project_id" in task_columns: where, params = _project_where(task_columns, tenant_id, project_id) assignments = ", ".join(f"{column}=NULL" for column in reset_columns) await cur.execute( f"UPDATE {settings.db_schema}.acquisition_tasks SET {assignments} WHERE {where}", params, ) task_ids = ids["task_ids"] if "notifications" in shapes and "related_task_id" in shapes["notifications"]: await _delete_any(cur, "notifications", "related_task_id", task_ids) if "task_events" in shapes and "task_id" in shapes["task_events"]: await _delete_any(cur, "task_events", "task_id", task_ids) if "super_agent_tasks" in shapes and "related_task_id" in shapes["super_agent_tasks"] and task_ids: # Preserve the autonomous run ledger; only remove its dangling task link. await cur.execute( f"UPDATE {settings.db_schema}.super_agent_tasks SET related_task_id=NULL " "WHERE related_task_id=ANY(%s)", (task_ids,), ) if "validation_issues" in shapes: await _delete_validation_issues(cur, shapes["validation_issues"], ids) if "review_actions" in shapes and "project_id" not in shapes["review_actions"]: review_columns = shapes["review_actions"] if {"candidate_id", "candidate_type"}.issubset(review_columns): clauses: list[str] = [] params: list[Any] = [] if ids["candidate_entity_ids"]: clauses.append("candidate_type IN ('entity', 'candidate_entity') AND candidate_id=ANY(%s)") params.append(ids["candidate_entity_ids"]) if ids["candidate_relation_ids"]: clauses.append("candidate_type IN ('relation', 'candidate_relation') AND candidate_id=ANY(%s)") params.append(ids["candidate_relation_ids"]) if clauses: await cur.execute( f"DELETE FROM {settings.db_schema}.review_actions WHERE " + " OR ".join(f"({clause})" for clause in clauses), tuple(params), ) if "raw_records" in shapes and "batch_id" in shapes["raw_records"]: await _delete_any(cur, "raw_records", "batch_id", ids["batch_ids"]) if ( "mapping_profiles" in shapes and "project_id" not in shapes["mapping_profiles"] and "template_id" in shapes["mapping_profiles"] ): await _delete_any(cur, "mapping_profiles", "template_id", ids["template_ids"]) if "ontology_fields" in shapes and "entity_type_id" in shapes["ontology_fields"]: await _delete_any(cur, "ontology_fields", "entity_type_id", ids["entity_type_ids"]) if "ontology_relation_types" in shapes and "schema_id" in shapes["ontology_relation_types"]: await _delete_any(cur, "ontology_relation_types", "schema_id", ids["schema_ids"]) if "ontology_entity_types" in shapes and "schema_id" in shapes["ontology_entity_types"]: await _delete_any(cur, "ontology_entity_types", "schema_id", ids["schema_ids"]) schema_job_columns = shapes.get("ontology_schema_publish_jobs", set()) if schema_job_columns and "project_id" not in schema_job_columns and "schema_id" in schema_job_columns: await _delete_any(cur, "ontology_schema_publish_jobs", "schema_id", ids["schema_ids"]) for table in _PROJECT_DELETE_ORDER: columns = shapes.get(table, set()) if "project_id" not in columns: continue where, params = _project_where(columns, tenant_id, project_id) await cur.execute(f"DELETE FROM {settings.db_schema}.{table} WHERE {where}", params) exclusive_graph_names = [item["graph_name"] for item in impact["graphs"] if item["will_delete"]] for table in _GRAPH_BOUND_TABLES: columns = shapes.get(table, set()) if "graph_name" in columns and exclusive_graph_names: await _delete_any(cur, table, "graph_name", exclusive_graph_names) await cur.execute( f"DELETE FROM {settings.db_schema}.projects WHERE tenant_id=%s AND project_id=%s", (tenant_id, project_id), ) if cur.rowcount != 1: raise RuntimeError("项目元数据删除数量异常,事务已中止") return before @router.get("/projects") async def list_projects(_user: CurrentUser): async with get_conn() as conn: async with conn.cursor() as cur: await cur.execute( f"""SELECT * FROM {settings.db_schema}.projects WHERE status <> 'archived' ORDER BY CASE project_id WHEN 'CityGraph-new2' THEN 1 WHEN 'travel_agency' THEN 2 ELSE 99 END, created_at DESC""" ) return await cur.fetchall() @router.post("/projects/provision") async def provision_project(body: dict, _user: CurrentUser): """Atomically provision project metadata, ontology, release, and a new graph.""" _require_project_admin(_user) try: payload = normalize_provision_payload(body) except ProjectValidationError as exc: raise HTTPException(422, {"message": "项目数据校验失败", "errors": exc.errors}) from exc try: existing_graphs = await asyncio.to_thread(list_falkor_graphs) except Exception as exc: raise HTTPException(503, f"FalkorDB 不可用,未写入任何项目数据:{exc}") from exc if payload["graph_name"] in existing_graphs: raise HTTPException(409, _graph_resource_conflict_detail(payload, "已存在,不能覆盖")) s = settings.db_schema graph_created = False cleanup_pending = False cleanup_error: Exception | None = None provision_task_id: int | None = None project_row: dict[str, Any] | None = None schema_row: dict[str, Any] | None = None release_row: dict[str, Any] | None = None graph_counts: dict[str, int] = {"nodes": 0, "relations": 0} # Session-level locks let us durably commit a cleanup intent before the # external graph write, while retaining the same project/graph exclusion # across the following PostgreSQL transaction. The fixed project→graph # order is shared with delete and legacy release creation. lock_keys = [ f"kg-project:{payload['project_id']}", f"kg-graph:{payload['graph_name']}", ] acquired_lock_keys: list[str] = [] async with get_conn() as conn: try: async with conn.cursor() as cur: for lock_key in lock_keys: acquired_lock_keys.append(lock_key) await cur.execute("SELECT pg_advisory_lock(hashtext(%s))", (lock_key,)) await conn.commit() # Recheck every namespace under the session locks and persist the # compensation outbox before FalkorDB can be changed. async with conn.cursor() as cur: await cur.execute( f"SELECT 1 FROM {s}.projects WHERE project_id=%s LIMIT 1", (payload["project_id"],), ) if await cur.fetchone(): raise HTTPException(409, "项目英文标识已存在,不能覆盖") await cur.execute( f"SELECT 1 FROM {s}.graph_releases WHERE graph_name=%s LIMIT 1", (payload["graph_name"],), ) if await cur.fetchone(): raise HTTPException( 409, _graph_resource_conflict_detail(payload, "已被现有项目使用"), ) await cur.execute( f"SELECT 1 FROM {s}.graph_project_tombstones WHERE project_id=%s", (payload["project_id"],), ) if await cur.fetchone(): raise HTTPException( 409, "项目英文标识曾被删除,为隔离历史写请求不能复用,请使用新的英文标识", ) await cur.execute( f"SELECT 1 FROM {s}.graph_name_tombstones WHERE graph_name=%s", (payload["graph_name"],), ) if await cur.fetchone(): raise HTTPException( 409, _graph_resource_conflict_detail(payload, "曾被删除,为隔离历史写请求不能复用"), ) try: locked_graphs = await _run_blocking_to_completion(list_falkor_graphs) except Exception as exc: raise HTTPException(503, f"FalkorDB 不可用,未写入项目数据:{exc}") from exc if payload["graph_name"] in locked_graphs: raise HTTPException( 409, _graph_resource_conflict_detail(payload, "已存在,不能覆盖"), ) provision_task_id = await _record_cleanup_task( cur, payload["project_id"], payload["graph_name"], operation="rollback_provision", ) # This commit is the durable intent boundary. Session locks remain # owned by this physical connection after the transaction ends. await conn.commit() async with conn.cursor() as cur: schema = payload["schema"] schema_storage_version = await _ontology_schema_storage_version( cur, tenant_id=payload["tenant_id"], project_id=payload["project_id"], namespace=schema["namespace"], semantic_version=schema["version"], ) await cur.execute( f"""INSERT INTO {s}.projects ( tenant_id, project_id, display_name, description, status, default_namespace, metadata_jsonb, created_by, updated_at ) VALUES (%s, %s, %s, %s, 'provisioning', %s, %s, %s, now()) RETURNING *""", ( payload["tenant_id"], payload["project_id"], payload["display_name"], payload["description"], payload["schema"]["namespace"], Jsonb({"provision_source": "admin-project-wizard"}), _user["username"], ), ) project_row = dict(await cur.fetchone()) await cur.execute( f"""INSERT INTO {s}.ontology_schemas ( tenant_id, project_id, namespace, version, display_name, description, status, schema_jsonb, created_by, updated_at ) VALUES (%s, %s, %s, %s, %s, %s, 'draft', %s, %s, now()) RETURNING *""", ( payload["tenant_id"], payload["project_id"], schema["namespace"], schema_storage_version, schema.get("display_name") or f"{payload['display_name']} Schema", schema.get("description") or schema.get("purpose"), Jsonb(schema), _user["username"], ), ) schema_row = dict(await cur.fetchone()) schema_id = schema_row["id"] for entity_type, entity_meta in schema["entity_types"].items(): natural_key_rule = entity_meta.get("natural_key_rule") or {} if not natural_key_rule and entity_meta.get("primary_key"): natural_key_rule = {"fields": [entity_meta["primary_key"]]} await cur.execute( f"""INSERT INTO {s}.ontology_entity_types ( schema_id, entity_type, display_name, description, is_core, extends_type, natural_key_rule_jsonb, recommendable, evidence_required, metadata_jsonb, updated_at ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, now()) RETURNING id""", ( schema_id, entity_type, entity_meta.get("label") or entity_meta.get("display_name") or entity_type, entity_meta.get("description") or entity_meta.get("definition"), bool(entity_meta.get("is_core", False)), entity_meta.get("extends") or entity_meta.get("extends_type"), Jsonb(natural_key_rule), bool(entity_meta.get("recommendable", False)), bool(entity_meta.get("evidence_required", False)), Jsonb(entity_meta.get("metadata") or {}), ), ) entity_type_id = (await cur.fetchone())["id"] for field_name, field_meta in entity_meta.get("fields", {}).items(): enum_values = field_meta.get("enum_values") or field_meta.get("enum") or [] await cur.execute( f"""INSERT INTO {s}.ontology_fields ( entity_type_id, field_name, display_name, value_type, required, multi_value, default_value_jsonb, enum_values_jsonb, validation_rules_jsonb, recommendable, evidence_required, updated_at ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, now())""", ( entity_type_id, field_name, field_meta.get("label") or field_meta.get("display_name") or field_name, field_meta.get("type") or "string", bool(field_meta.get("required", False)), bool(field_meta.get("multi_value", False)), Jsonb(field_meta["default"]) if "default" in field_meta else None, Jsonb(enum_values), Jsonb(field_meta.get("validation") or field_meta.get("validation_rules") or {}), bool(field_meta.get("recommendable", False)), bool(field_meta.get("evidence_required", False)), ), ) for relation_type, relation_meta in schema["relation_types"].items(): await cur.execute( f"""INSERT INTO {s}.ontology_relation_types ( schema_id, relation_type, source_entity_type, target_entity_type, display_name, cardinality, required, evidence_required, publish_as, metadata_jsonb, updated_at ) VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, now())""", ( schema_id, relation_type, relation_meta["from"], relation_meta["to"], relation_meta.get("label") or relation_meta.get("display_name") or relation_type, relation_meta.get("cardinality") or "many_to_many", bool(relation_meta.get("required", False)), bool(relation_meta.get("evidence_required", False)), relation_meta.get("publish_as"), Jsonb({"properties": relation_meta.get("properties") or {}}), ), ) graph_release_id = f"{payload['project_id']}_v{schema['version']}" await cur.execute( f"""INSERT INTO {s}.graph_releases ( tenant_id, project_id, graph_release_id, graph_name, alias, status, schema_id, metadata_jsonb, created_by, updated_at ) VALUES (%s, %s, %s, %s, 'active', 'provisioning', %s, %s, %s, now()) RETURNING *""", ( payload["tenant_id"], payload["project_id"], graph_release_id, payload["graph_name"], schema_id, Jsonb(payload["counts"]), _user["username"], ), ) release_row = dict(await cur.fetchone()) try: graph_counts = await _run_blocking_to_completion( import_falkor_graph, payload["graph_name"], payload["graph_data"], ) except FileExistsError as exc: raise HTTPException( 409, _graph_resource_conflict_detail(payload, "已存在,不能覆盖"), ) from exc graph_created = graph_counts["nodes"] > 0 if graph_counts["nodes"] != payload["counts"]["nodes"]: raise RuntimeError("FalkorDB 节点数量校验失败") if graph_counts["relations"] != payload["counts"]["relations"]: raise RuntimeError("FalkorDB 关系数量校验失败") await cur.execute( f"""UPDATE {s}.ontology_schemas SET status='active', published_by=%s, published_at=now(), updated_at=now() WHERE id=%s RETURNING *""", (_user["username"], schema_id), ) schema_row = dict(await cur.fetchone()) await cur.execute( f"""UPDATE {s}.graph_releases SET status='active', metadata_jsonb=%s, published_at=now(), activated_at=now(), updated_at=now() WHERE id=%s RETURNING *""", (Jsonb({**payload["counts"], "graph_counts": graph_counts}), release_row["id"]), ) release_row = dict(await cur.fetchone()) await cur.execute( f"""UPDATE {s}.projects SET status='active', metadata_jsonb=%s, updated_at=now() WHERE tenant_id=%s AND project_id=%s RETURNING *""", ( Jsonb( { "provision_source": "admin-project-wizard", "schema_id": schema_id, "graph_name": payload["graph_name"], **payload["counts"], } ), payload["tenant_id"], payload["project_id"], ), ) project_row = dict(await cur.fetchone()) await _set_cleanup_task(cur, provision_task_id) # Metadata activation and outbox completion become visible together. await conn.commit() except HTTPException: if conn.info.transaction_status != 0: await conn.rollback() if provision_task_id is not None: try: await _finish_cleanup_task(provision_task_id, conn=conn) except Exception: if conn.info.transaction_status != 0: await conn.rollback() raise except asyncio.CancelledError: if conn.info.transaction_status != 0: await conn.rollback() # The durable pending intent is deliberately retained. A retry can # safely delete an unreferenced graph or complete an already-absent # task after this cancelled request releases its session locks. raise except Exception as exc: if conn.info.transaction_status != 0: await conn.rollback() cleanup_error = exc.cleanup_error if isinstance(exc, GraphImportError) else None if graph_created: try: await _run_blocking_to_completion(delete_falkor_graph, payload["graph_name"]) except Exception as cleanup_exc: cleanup_pending = True cleanup_error = cleanup_exc elif isinstance(exc, GraphImportError) and exc.cleanup_required: cleanup_pending = True if provision_task_id is not None: task_error = cleanup_error or RuntimeError("FalkorDB 图需要补偿清理") try: await _finish_cleanup_task( provision_task_id, error=str(task_error) if cleanup_pending else None, conn=conn, ) except Exception as task_exc: if conn.info.transaction_status != 0: await conn.rollback() cleanup_pending = True cleanup_error = task_exc detail: dict[str, Any] = { "message": f"项目创建失败,PostgreSQL 元数据已回滚:{exc}", "graph_name": payload["graph_name"], } if cleanup_pending: detail["message"] += ";FalkorDB 补偿清理已进入待重试状态" detail["cleanup_task_id"] = provision_task_id detail["cleanup_error"] = str(cleanup_error or "待补偿清理") raise HTTPException(500, detail) from exc finally: if acquired_lock_keys: await _run_async_cleanup_to_completion( _release_session_locks(conn, acquired_lock_keys) ) if schema_row is not None: schema_row["version"] = payload["schema"]["version"] actual_graph_name = str((release_row or {}).get("graph_name") or payload["graph_name"]) return { "project": project_row, "schema": schema_row, "graph_release": release_row, "graph": {"graph_name": actual_graph_name, **graph_counts}, "counts": payload["counts"], } @router.post("/projects") async def create_project(body: dict, _user: CurrentUser): _require_project_admin(_user) s = settings.db_schema project_id = str(body.get("project_id") or "").strip() if not SAFE_RESOURCE_ID.fullmatch(project_id): raise HTTPException(422, "project_id 只能包含安全字符") tenant_id = str(body.get("tenant_id") or project_id).strip() if not SAFE_RESOURCE_ID.fullmatch(tenant_id): raise HTTPException(422, "tenant_id 只能包含安全字符") display_name = str(body.get("display_name") or body.get("name") or project_id).strip() if not display_name: raise HTTPException(422, "display_name 不能为空") async with get_conn() as conn: async with conn.cursor() as cur: await cur.execute( "SELECT pg_advisory_xact_lock(hashtext(%s))", (f"kg-project:{project_id}",), ) await cur.execute( f"SELECT 1 FROM {s}.projects WHERE project_id=%s LIMIT 1", (project_id,), ) if await cur.fetchone(): raise HTTPException(409, "project_id 已存在,不允许覆盖") await cur.execute( f"SELECT 1 FROM {s}.graph_project_tombstones WHERE project_id=%s", (project_id,), ) if await cur.fetchone(): raise HTTPException(409, "project_id 曾被删除,为隔离历史写请求不能复用,请使用新 ID") await cur.execute( f"""INSERT INTO {s}.projects ( tenant_id, project_id, display_name, description, status, default_namespace, metadata_jsonb, created_by, updated_at ) VALUES ( %(tenant_id)s, %(project_id)s, %(display_name)s, %(description)s, %(status)s, %(default_namespace)s, %(metadata_jsonb)s, %(created_by)s, now() ) ON CONFLICT (tenant_id, project_id) DO NOTHING RETURNING *""", { "project_id": project_id, "tenant_id": tenant_id, "display_name": display_name, "description": body.get("description"), "status": body.get("status") or "active", "default_namespace": body.get("default_namespace") or project_id, "metadata_jsonb": Jsonb(body.get("metadata_jsonb") or {}), "created_by": _user["username"], }, ) row = await cur.fetchone() if not row: raise HTTPException(409, "项目已存在,不允许覆盖") await conn.commit() return row @router.patch("/projects/{project_id}") async def rename_project(project_id: str, body: dict, _user: CurrentUser): _require_project_admin(_user) extra_fields = set(body) - {"display_name"} if extra_fields: raise HTTPException(422, f"只允许修改 display_name;禁止字段:{', '.join(sorted(extra_fields))}") display_name = str(body.get("display_name") or "").strip() if not display_name: raise HTTPException(422, "display_name 不能为空") if len(display_name) > 200: raise HTTPException(422, "display_name 不能超过 200 个字符") async with get_conn() as conn: try: async with conn.cursor() as cur: project = await _load_project(cur, project_id, for_update=True) await cur.execute( f"""UPDATE {settings.db_schema}.projects SET display_name=%s, updated_at=now() WHERE tenant_id=%s AND project_id=%s RETURNING *""", (display_name, project["tenant_id"], project_id), ) row = await cur.fetchone() await conn.commit() except Exception: await conn.rollback() raise return row @router.get("/projects/{project_id}/deletion-impact") async def project_deletion_impact(project_id: str, _user: CurrentUser): _require_project_admin(_user) async with get_conn() as conn: async with conn.cursor() as cur: project = await _load_project(cur, project_id) impact, _shapes, _ids = await _build_deletion_impact(cur, project) try: existing_graphs = await asyncio.to_thread(list_falkor_graphs) for item in impact["graphs"]: item["exists_in_falkordb"] = item["graph_name"] in existing_graphs except Exception as exc: impact["falkordb_error"] = str(exc) for item in impact["graphs"]: item["exists_in_falkordb"] = None return impact @router.delete("/projects/{project_id}") async def delete_project( project_id: str, confirm_project_id: str, _user: CurrentUser, ): _require_project_admin(_user) if confirm_project_id != project_id: raise HTTPException(400, "确认项目 ID 不匹配,未执行任何删除") impact: dict[str, Any] cleanup_task_ids: dict[str, int] = {} locked_resource_keys: list[str] = [] deleted_graphs: list[str] = [] absent_graphs: list[str] = [] cleanup_pending: list[dict[str, Any]] = [] async with get_conn() as conn: try: async with conn.cursor() as cur: project_lock_key = f"kg-project:{project_id}" locked_resource_keys.append(project_lock_key) await cur.execute( "SELECT pg_advisory_lock(hashtext(%s))", (project_lock_key,), ) project = await _load_project(cur, project_id, for_update=True) preliminary_graphs = await _graph_references( cur, str(project["tenant_id"]), project_id, ) for item in preliminary_graphs: graph_lock_key = f"kg-graph:{item['graph_name']}" locked_resource_keys.append(graph_lock_key) await cur.execute( # A session lock intentionally survives the PostgreSQL # commit below and remains held until the corresponding # Falkor delete finishes. This prevents a new project # from reusing an absent graph name in that small gap. "SELECT pg_advisory_lock(hashtext(%s))", (graph_lock_key,), ) # Recompute every count and shared-graph decision while all # relevant transaction locks are held. The earlier GET result # is deliberately never trusted for deletion. impact, shapes, ids = await _build_deletion_impact(cur, project) for item in impact["graphs"]: if item["will_delete"]: cleanup_task_ids[item["graph_name"]] = await _record_cleanup_task( cur, project_id, item["graph_name"], ) await _delete_project_resources(cur, project, impact, shapes, ids) await cur.execute( f"""INSERT INTO {settings.db_schema}.graph_project_tombstones (project_id, tenant_id, deleted_by, deleted_at) VALUES (%s, %s, %s, now()) ON CONFLICT (project_id) DO UPDATE SET tenant_id=EXCLUDED.tenant_id, deleted_by=EXCLUDED.deleted_by, deleted_at=now()""", (project_id, project["tenant_id"], _user["username"]), ) for item in impact["graphs"]: if not item["will_delete"]: continue await cur.execute( f"""INSERT INTO {settings.db_schema}.graph_name_tombstones (graph_name, project_id, deleted_at) VALUES (%s, %s, now()) ON CONFLICT (graph_name) DO UPDATE SET project_id=EXCLUDED.project_id, deleted_at=now()""", (item["graph_name"], project_id), ) await conn.commit() for graph_name, task_id in cleanup_task_ids.items(): try: deleted = await _run_blocking_to_completion(delete_falkor_graph, graph_name) if deleted: deleted_graphs.append(graph_name) else: absent_graphs.append(graph_name) await _finish_cleanup_task(task_id, conn=conn) except Exception as exc: if conn.info.transaction_status != 0: await conn.rollback() error = str(exc) cleanup_pending.append( {"task_id": task_id, "graph_name": graph_name, "error": error} ) try: await _finish_cleanup_task(task_id, error=error, conn=conn) except Exception: if conn.info.transaction_status != 0: await conn.rollback() except Exception: if conn.info.transaction_status != 0: await conn.rollback() raise finally: if locked_resource_keys: await _run_async_cleanup_to_completion( _release_session_locks(conn, locked_resource_keys) ) shared_graphs = [item["graph_name"] for item in impact["graphs"] if item["shared"]] return { "status": "deleted_with_cleanup_pending" if cleanup_pending else "deleted", "project_id": project_id, "deleted_counts": impact["counts"], "deleted_graphs": deleted_graphs, "already_absent_graphs": absent_graphs, "shared_graphs_preserved": shared_graphs, "cleanup_pending": cleanup_pending, "protected_resources": impact["protected_resources"], } @router.get("/projects/cleanup-tasks") async def list_project_graph_cleanup_tasks( _user: CurrentUser, status: str = "pending", project_id: str | None = None, ): """List lifecycle compensation tasks so failed cleanup is operable.""" _require_project_admin(_user) if status not in {"pending", "completed", "all"}: raise HTTPException(422, "status 只能是 pending、completed 或 all") if project_id is not None and not SAFE_RESOURCE_ID.fullmatch(project_id): raise HTTPException(422, "project_id 包含不安全字符") clauses: list[str] = [] params: list[Any] = [] if status != "all": clauses.append("status=%s") params.append(status) if project_id is not None: clauses.append("project_id=%s") params.append(project_id) where = " WHERE " + " AND ".join(clauses) if clauses else "" async with get_conn() as conn: async with conn.cursor() as cur: await _ensure_cleanup_table(cur) await cur.execute( f"""SELECT id, project_id, graph_name, operation, status, attempts, last_error, created_at, updated_at FROM {settings.db_schema}.graph_cleanup_tasks{where} ORDER BY updated_at DESC, id DESC LIMIT 200""", tuple(params), ) return await cur.fetchall() @router.post("/projects/cleanup-tasks/{task_id}/retry") async def retry_project_graph_cleanup(task_id: int, _user: CurrentUser): """Retry one recorded FalkorDB compensation without trusting stale ownership.""" _require_project_admin(_user) async with get_conn() as conn: try: async with conn.cursor() as cur: await _ensure_cleanup_table(cur) await cur.execute( f"SELECT * FROM {settings.db_schema}.graph_cleanup_tasks WHERE id=%s FOR UPDATE", (task_id,), ) task = await cur.fetchone() if not task: raise HTTPException(404, "补偿清理任务不存在") if task["status"] == "completed": await conn.commit() return { "status": "completed", "task_id": task_id, "graph_name": task["graph_name"], "already_completed": True, } await cur.execute( "SELECT pg_advisory_xact_lock(hashtext(%s))", (f"kg-graph:{task['graph_name']}",), ) await cur.execute( f"SELECT count(*) AS count FROM {settings.db_schema}.graph_releases WHERE graph_name=%s", (task["graph_name"],), ) if int((await cur.fetchone())["count"]): # A crash can occur after metadata commit but before a # rollback-provision intent is observed as completed. Any # current release reference makes deletion unsafe, so # preserve the graph and resolve the stale cleanup task. await cur.execute( f"""UPDATE {settings.db_schema}.graph_cleanup_tasks SET status='completed', attempts=attempts+1, last_error='图已被 GraphRelease 引用,已安全保留', updated_at=now() WHERE id=%s""", (task_id,), ) await conn.commit() return { "status": "completed", "task_id": task_id, "graph_name": task["graph_name"], "deleted": False, "preserved_because_referenced": True, } try: deleted = await _run_blocking_to_completion( delete_falkor_graph, str(task["graph_name"]), ) except Exception as exc: await cur.execute( f"""UPDATE {settings.db_schema}.graph_cleanup_tasks SET status='pending', attempts=attempts+1, last_error=%s, updated_at=now() WHERE id=%s""", (str(exc), task_id), ) await conn.commit() raise HTTPException( 503, { "message": f"FalkorDB 清理仍未完成,可稍后重试:{exc}", "cleanup_task_id": task_id, "graph_name": task["graph_name"], }, ) from exc await cur.execute( f"""UPDATE {settings.db_schema}.graph_cleanup_tasks SET status='completed', attempts=attempts+1, last_error=NULL, updated_at=now() WHERE id=%s""", (task_id,), ) await conn.commit() return { "status": "completed", "task_id": task_id, "graph_name": task["graph_name"], "deleted": deleted, } except HTTPException: if conn.info.transaction_status != 0: await conn.rollback() raise except Exception: await conn.rollback() raise @router.get("/projects/{project_id}/graph-releases") async def list_graph_releases(project_id: str, _user: CurrentUser): async with get_conn() as conn: async with conn.cursor() as cur: await cur.execute( f"SELECT * FROM {settings.db_schema}.graph_releases " "WHERE project_id=%s AND status <> 'archived' " "ORDER BY " "CASE WHEN status='active' THEN 1 WHEN status='published' THEN 2 ELSE 99 END, " "created_at DESC", (project_id,), ) return await cur.fetchall() @router.post("/projects/{project_id}/graph-releases") async def create_graph_release(project_id: str, body: dict, _user: CurrentUser): _require_project_admin(_user) s = settings.db_schema graph_name = _resolve_graph_name(project_id, body.get("graph_name")) alias = str(body.get("alias") or "active").strip() tenant_id = str(body.get("tenant_id") or "").strip() status = str(body.get("status") or "active").strip() if not SAFE_RESOURCE_ID.fullmatch(project_id) or not SAFE_RESOURCE_ID.fullmatch(graph_name): raise HTTPException(422, "project_id 或 graph_name 包含不安全字符") if tenant_id and not SAFE_RESOURCE_ID.fullmatch(tenant_id): raise HTTPException(422, "tenant_id 包含不安全字符") if not SAFE_RESOURCE_ID.fullmatch(alias): raise HTTPException(422, "GraphRelease alias 包含不安全字符") async with get_conn() as conn: async with conn.cursor() as cur: await cur.execute( "SELECT pg_advisory_xact_lock(hashtext(%s))", (f"kg-project:{project_id}",), ) await cur.execute( "SELECT pg_advisory_xact_lock(hashtext(%s))", (f"kg-graph:{graph_name}",), ) if not tenant_id: await cur.execute( f"SELECT tenant_id FROM {s}.projects WHERE project_id=%s LIMIT 1", (project_id,), ) project = await cur.fetchone() if not project: raise HTTPException(404, "项目不存在,不能创建孤立 GraphRelease") tenant_id = project["tenant_id"] else: await cur.execute( f"SELECT 1 FROM {s}.projects WHERE tenant_id=%s AND project_id=%s LIMIT 1", (tenant_id, project_id), ) if not await cur.fetchone(): raise HTTPException(404, "项目不存在,不能创建孤立 GraphRelease") try: existing_graphs = await asyncio.to_thread(list_falkor_graphs) except Exception as exc: raise HTTPException(503, f"FalkorDB 不可用,未创建 GraphRelease:{exc}") from exc if graph_name not in existing_graphs: raise HTTPException(409, "FalkorDB 图不存在,请使用项目创建向导完成原子创建") await cur.execute( f"SELECT 1 FROM {s}.graph_name_tombstones WHERE graph_name=%s", (graph_name,), ) if await cur.fetchone(): raise HTTPException(409, "graph_name 曾被删除,为隔离历史写请求不能复用,请使用新名称") await cur.execute( f"""INSERT INTO {s}.graph_releases ( tenant_id, project_id, graph_release_id, graph_name, alias, status, metadata_jsonb, created_by, activated_at, updated_at ) VALUES ( %(tenant_id)s, %(project_id)s, %(graph_release_id)s, %(graph_name)s, %(alias)s, %(status)s, %(metadata_jsonb)s, %(created_by)s, now(), now() ) ON CONFLICT (tenant_id, project_id, alias) DO NOTHING RETURNING *""", { "tenant_id": tenant_id, "project_id": project_id, "graph_release_id": body.get("graph_release_id") or f"{project_id}_{alias}", "graph_name": graph_name, "alias": alias, "status": status, "metadata_jsonb": Jsonb(body.get("metadata_jsonb") or {}), "created_by": _user["username"], }, ) row = await cur.fetchone() if not row: raise HTTPException(409, "该项目的 GraphRelease alias 已存在,不允许覆盖") await conn.commit() return row @router.get("/ontology-schemas") async def list_ontology_schemas( _user: CurrentUser, ctx: ProjectContext = Depends(get_project_context), ): """List schema versions for the current project context.""" async with get_conn() as conn: async with conn.cursor() as cur: await cur.execute( f"""SELECT os.id, os.tenant_id, os.project_id, os.namespace, os.version, os.display_name, os.description, os.status, os.schema_jsonb, os.published_at, os.created_at, os.updated_at, gr.graph_release_id AS active_graph_release_id, gr.graph_name AS active_graph_name FROM {settings.db_schema}.ontology_schemas os LEFT JOIN {settings.db_schema}.graph_releases gr ON gr.schema_id=os.id AND gr.tenant_id=os.tenant_id AND gr.project_id=os.project_id AND gr.alias='active' AND gr.status <> 'archived' WHERE os.tenant_id=%s AND os.project_id=%s AND os.status <> 'archived' ORDER BY CASE WHEN os.status='active' THEN 1 ELSE 99 END, os.version DESC, os.updated_at DESC""", (ctx.tenant_id, ctx.project_id), ) rows = await cur.fetchall() return [_attach_schema_source(row) for row in rows] @router.get("/ontology-schemas/current") async def get_current_ontology_schema( _user: CurrentUser, ctx: ProjectContext = Depends(get_project_context), ): """Return the active graph release schema for the current project.""" async with get_conn() as conn: async with conn.cursor() as cur: await cur.execute( f"""SELECT os.*, gr.graph_release_id AS active_graph_release_id, gr.graph_name AS active_graph_name, gr.alias AS active_graph_alias FROM {settings.db_schema}.graph_releases gr JOIN {settings.db_schema}.ontology_schemas os ON os.id=gr.schema_id WHERE gr.tenant_id=%s AND gr.project_id=%s AND gr.alias='active' AND gr.status <> 'archived' AND os.status <> 'archived' ORDER BY gr.activated_at DESC NULLS LAST, gr.updated_at DESC LIMIT 1""", (ctx.tenant_id, ctx.project_id), ) row = await cur.fetchone() if row: return _attach_schema_source(row) await cur.execute( f"""SELECT os.*, NULL::text AS active_graph_release_id, NULL::text AS active_graph_name, NULL::text AS active_graph_alias FROM {settings.db_schema}.ontology_schemas os WHERE os.tenant_id=%s AND os.project_id=%s AND os.status <> 'archived' ORDER BY CASE WHEN os.status='active' THEN 1 ELSE 99 END, os.version DESC, os.updated_at DESC LIMIT 1""", (ctx.tenant_id, ctx.project_id), ) row = await cur.fetchone() if not row: raise HTTPException(404, "No ontology schema for current project") return _attach_schema_source(row) @router.get("/ontology-schemas/{schema_id}") async def get_ontology_schema( schema_id: int, _user: CurrentUser, ctx: ProjectContext = Depends(get_project_context), ): """Return one schema version for the current project context.""" async with get_conn() as conn: async with conn.cursor() as cur: await cur.execute( f"""SELECT os.*, gr.graph_release_id AS active_graph_release_id, gr.graph_name AS active_graph_name, gr.alias AS active_graph_alias FROM {settings.db_schema}.ontology_schemas os LEFT JOIN {settings.db_schema}.graph_releases gr ON gr.schema_id=os.id AND gr.tenant_id=os.tenant_id AND gr.project_id=os.project_id AND gr.alias='active' AND gr.status <> 'archived' WHERE os.id=%s AND os.tenant_id=%s AND os.project_id=%s AND os.status <> 'archived' LIMIT 1""", (schema_id, ctx.tenant_id, ctx.project_id), ) row = await cur.fetchone() if not row: raise HTTPException(404, "Schema not found for current project") return _attach_schema_source(row) @router.get("/health") async def health(): return {"status": "ok"}