From e1389871125a20e84ab59123d0176e711e10a17d Mon Sep 17 00:00:00 2001 From: "jinli.yl" Date: Thu, 14 May 2026 18:01:43 +0800 Subject: [PATCH] refactor(file_graph): replace local file graph with pure Python implementation - Rename local_file_graph.py to nx_file_graph.py and update class name to NxFileGraph - Add proper import error handling for networkx dependency in NxFileGraph - Create new LocalFileGraph implementation using pure Python dict-based storage - Update base_file_graph to include abstract methods for node/link operations - Modify upsert_nodes to handle multiple nodes and edge management - Implement proper edge bookkeeping with inverse and pending edge tracking - Add rebuild_links functionality to reconstruct graph relationships - Update logging to show real vs virtual nodes and edge counts - Maintain backward compatibility with existing graph persistence mechanisms --- reme2/component/file_graph/__init__.py | 15 +- reme2/component/file_graph/base_file_graph.py | 41 ++-- .../component/file_graph/local_file_graph.py | 188 ++++++++---------- 3 files changed, 104 insertions(+), 140 deletions(-) diff --git a/reme2/component/file_graph/__init__.py b/reme2/component/file_graph/__init__.py index 63a4f59b..f803111f 100644 --- a/reme2/component/file_graph/__init__.py +++ b/reme2/component/file_graph/__init__.py @@ -1,22 +1,13 @@ -"""File graph module. - -Backend-agnostic graph engine that owns the wikilink graph (nodes + -typed edges + traversal queries). Sister to ``file_store``, which owns -chunks and projections (vector / FTS). - -Two backends ship today: - - ``LocalFileGraph`` (``@R.register("local")``) — networkx - ``MultiDiGraph`` + JSONL persistence. Default. - - ``Neo4jFileGraph`` (``@R.register("neo4j")``) — property-graph - backed by Neo4j. Requires the ``neo4j`` driver. -""" +"""File graph module.""" from .base_file_graph import BaseFileGraph from .local_file_graph import LocalFileGraph +from .nx_file_graph import NxFileGraph from .neo4j_file_graph import Neo4jFileGraph __all__ = [ "BaseFileGraph", "LocalFileGraph", + "NxFileGraph", "Neo4jFileGraph", ] diff --git a/reme2/component/file_graph/base_file_graph.py b/reme2/component/file_graph/base_file_graph.py index 0ea21cdb..a0132c50 100644 --- a/reme2/component/file_graph/base_file_graph.py +++ b/reme2/component/file_graph/base_file_graph.py @@ -1,20 +1,3 @@ -"""Abstract base for the file-graph engine. - -The file-graph owns ``FileNode`` records keyed by vault-relative path -and serves graph traversal over wikilink links. file_graph trusts -``FileLink.path`` directly — there is no internal wikilink resolution. -Pre-resolution forms (raw stem-form wikilinks like ``[[Foo]]``) are a -vault convention resolved by the external ``utils.wikilink_resolver`` -before nodes are upserted. - -Contract — six abstract methods, two blocks: - - Node CRUD upsert_node, delete_node, get_node, iter_nodes - Link access get_outlinks, get_inlinks -""" - -from __future__ import annotations - from abc import abstractmethod from pathlib import Path @@ -24,36 +7,38 @@ from ...schema import FileLink, FileNode class BaseFileGraph(BaseComponent): - """Pluggable file-graph backend — node CRUD + link adjacency.""" - component_type = ComponentEnum.FILE_GRAPH def __init__(self, graph_name: str = "default", **kwargs): super().__init__(**kwargs) self.graph_name: str = graph_name or self.name - self.graph_path: Path = self.working_path / self.component_type.value / graph_name + self.graph_path: Path = self.working_path / self.component_type.value self.graph_path.mkdir(parents=True, exist_ok=True) # -- Node CRUD --------------------------------------------------------- @abstractmethod - async def upsert_node(self, node: FileNode) -> None: - """Add or replace a node and its outgoing links.""" + async def upsert_nodes(self, nodes: list[FileNode]) -> None: + """Upsert nodes into the graph.""" @abstractmethod - async def delete_node(self, path: str) -> FileNode | None: - """Remove a node + its incident links. Returns the removed node, or None.""" + async def delete_nodes(self, paths: list[str]) -> None: + """Delete nodes from the graph.""" @abstractmethod - async def get_node(self, path: str) -> FileNode | None: - """Single-node lookup.""" + async def get_nodes(self, paths: list[str]) -> list[FileNode]: + """Get nodes from the graph.""" + + @abstractmethod + async def rebuild_links(self) -> None: + """Rebuild all links in the graph from each node's payload.""" # -- Link access ------------------------------------------------------- @abstractmethod async def get_outlinks(self, path: str) -> list[FileLink]: - """Resolved outgoing links from ``path``.""" + """Get outlinks from the graph.""" @abstractmethod async def get_inlinks(self, path: str) -> list[FileLink]: - """Resolved incoming links to ``path``.""" + """Get inlinks from the graph.""" diff --git a/reme2/component/file_graph/local_file_graph.py b/reme2/component/file_graph/local_file_graph.py index ddbf6341..17b06724 100644 --- a/reme2/component/file_graph/local_file_graph.py +++ b/reme2/component/file_graph/local_file_graph.py @@ -1,25 +1,5 @@ -"""Local file-graph backend — networkx ``MultiDiGraph``. +"""Pure-Python file-graph backend (no external deps).""" -Each ``FileNode`` lives as ``data['node']`` on the graph; each -``FileLink`` whose ``link.path`` exists in the graph lives as -``data['link']`` on a directed graph edge between the two file nodes. - -file_graph trusts ``link.path`` directly — there is no internal -wikilink resolution. The parser pipeline (with the external resolver) -is responsible for producing safe ``FileLink`` records where -``link.path`` is a real vault-relative target path. Stem ambiguity is -already handled there by emitting one link per candidate. - -Late-arriving target: when ``upsert_node(B)`` runs, other nodes whose -links have ``link.path == B.path`` get those in-edges restored -(one O(N×L) sweep per upsert; cheap for vault sizes). - -Persistence: pickle to ``graph.pkl`` on close, restored on _start. -""" - -from __future__ import annotations - -import pickle from pathlib import Path from .base_file_graph import BaseFileGraph @@ -29,112 +9,120 @@ from ...schema import FileLink, FileNode @R.register("local") class LocalFileGraph(BaseFileGraph): - """Networkx-backed file graph. Trusts ``FileLink.path`` for adjacency.""" + """Dict-backed file graph; trusts ``FileLink.path`` for adjacency.""" def __init__(self, **kwargs): super().__init__(**kwargs) - self._graph = None # networkx.MultiDiGraph; set in _start - self._graph_file: Path = self.graph_path / f"{self.graph_name}.pkl" + self._nodes: dict[str, FileNode] = {} + self._inverse: dict[str, set[str]] = {} + # Virtual edges: src→target where target isn't in ``_nodes`` yet. + self._pending: dict[str, set[str]] = {} + self._graph_file: Path = self.graph_path / f"{self.graph_name}.jsonl" # -- Lifecycle --------------------------------------------------------- async def _start(self) -> None: await super()._start() - import networkx as nx - self._graph = self._load_graph(nx) or nx.MultiDiGraph() + self._load() + await self.rebuild_links() + edges = sum(len(s) for s in self._inverse.values()) + pending = sum(len(s) for s in self._pending.values()) self.logger.info( f"LocalFileGraph '{self.graph_name}' ready: " - f"{self._graph.number_of_nodes()} nodes, " - f"{self._graph.number_of_edges()} edges", + f"{len(self._nodes)} nodes, {edges} edges, {pending} pending", ) async def _close(self) -> None: - self._save_graph() + self._dump() await super()._close() - def _load_graph(self, nx): + def _load(self) -> None: if not self._graph_file.exists(): - return None - try: - with open(self._graph_file, "rb") as f: - graph = pickle.load(f) - if not isinstance(graph, nx.MultiDiGraph): - self.logger.warning( - f"{self._graph_file} is not a MultiDiGraph; ignoring", - ) - return None - return graph - except Exception as e: - self.logger.exception(f"Failed to load {self._graph_file}: {e}") - return None + return + with open(self._graph_file, "r", encoding="utf-8") as f: + self._nodes.update( + (n.path, n) + for n in (FileNode.model_validate_json(line) for line in f if line.strip()) + ) - def _save_graph(self) -> None: - try: - tmp = self._graph_file.with_suffix(".tmp") - with open(tmp, "wb") as f: - pickle.dump(self._graph, f, protocol=pickle.HIGHEST_PROTOCOL) - tmp.replace(self._graph_file) - except Exception as e: - self.logger.exception(f"Failed to write {self._graph_file}: {e}") + def _dump(self) -> None: + tmp = self._graph_file.with_suffix(".tmp") + with open(tmp, "w", encoding="utf-8") as f: + f.writelines(f"{n.model_dump_json()}\n" for n in self._nodes.values()) + tmp.replace(self._graph_file) + + # -- Edge bookkeeping -------------------------------------------------- + + def _add_edge(self, src: str, target: str) -> None: + bucket = self._inverse if target in self._nodes else self._pending + bucket.setdefault(target, set()).add(src) + + def _remove_edge(self, src: str, target: str) -> None: + for bucket in (self._inverse, self._pending): + srcs = bucket.get(target) + if srcs is None or src not in srcs: + continue + srcs.discard(src) + if not srcs: + del bucket[target] # -- Node CRUD --------------------------------------------------------- - async def upsert_node(self, node: FileNode) -> None: - path = node.path - if self._graph.has_node(path): - self._graph.remove_node(path) - self._graph.add_node(path, node=node) + async def upsert_nodes(self, nodes: list[FileNode]) -> None: + for node in nodes: + path = node.path + old = self._nodes.get(path) + if old is not None: + for link in old.links: + if link.path: + self._remove_edge(path, link.path) + self._nodes[path] = node + for link in node.links: + if link.path: + self._add_edge(path, link.path) + # Promote virtual edges aimed at this newly-arrived target. + promoted = self._pending.pop(path, None) + if promoted: + self._inverse.setdefault(path, set()).update(promoted) - # Out-links from this node — trust link.path. - for link in node.links: - if link.path and self._graph.has_node(link.path): - self._graph.add_edge(path, link.path, link=link) - - # Late-arriving target: restore in-links from other nodes whose - # links resolve here. Scan all other nodes' links; cheap for - # vault sizes (O(N×L_avg) per upsert). - for src_path, src_data in self._graph.nodes(data=True): - if src_path == path: + async def delete_nodes(self, paths: list[str]) -> None: + for path in paths: + node = self._nodes.pop(path, None) + if node is None: continue - src_node = src_data.get("node") - if src_node is None: - continue - for link in src_node.links: - if link.path == path: - self._graph.add_edge(src_path, path, link=link) + for link in node.links: + if link.path: + self._remove_edge(path, link.path) + # Demote inbound edges to virtual; sources still link here. + demoted = self._inverse.pop(path, None) + if demoted: + self._pending.setdefault(path, set()).update(demoted) - async def delete_node(self, path: str) -> FileNode | None: - if not self._graph.has_node(path): - return None - node = self._graph.nodes[path].get("node") - self._graph.remove_node(path) - return node + async def get_nodes(self, paths: list[str]) -> list[FileNode]: + return [self._nodes[p] for p in paths if p in self._nodes] - async def get_node(self, path: str) -> FileNode | None: - if not self._graph.has_node(path): - return None - return self._graph.nodes[path].get("node") + async def rebuild_links(self) -> None: + self._inverse.clear() + self._pending.clear() + for src, node in self._nodes.items(): + for link in node.links: + if link.path: + self._add_edge(src, link.path) # -- Link access ------------------------------------------------------- - async def get_outlinks(self, path: str) -> list[tuple[FileNode, FileLink]]: - if not self._graph.has_node(path): + async def get_outlinks(self, path: str) -> list[FileLink]: + node = self._nodes.get(path) + if node is None: return [] - out: list[tuple[FileNode, FileLink]] = [] - for _src, dst, data in self._graph.out_edges(path, data=True): - target = self._graph.nodes[dst].get("node") - link = data.get("link") - if target is not None and isinstance(link, FileLink): - out.append((target, link)) - return out + return [link for link in node.links if link.path and link.path in self._nodes] - async def get_inlinks(self, path: str) -> list[tuple[FileNode, FileLink]]: - if not self._graph.has_node(path): + async def get_inlinks(self, path: str) -> list[FileLink]: + if path not in self._nodes: return [] - out: list[tuple[FileNode, FileLink]] = [] - for src, _dst, data in self._graph.in_edges(path, data=True): - source = self._graph.nodes[src].get("node") - link = data.get("link") - if source is not None and isinstance(link, FileLink): - out.append((source, link)) - return out + return [ + link + for src in self._inverse.get(path, ()) + for link in self._nodes[src].links + if link.path == path + ]