From d8d60800861b8f4c34b8bf10e25989fbbce2405c Mon Sep 17 00:00:00 2001 From: Cl0ver Date: Mon, 28 Sep 2026 11:31:45 -0400 Subject: [PATCH 1/9] build: add the FalkorDB client as an optional dependency The FalkorDB store needs the falkordb client, which most people do not, so it is an extra: pip install 'chatlore[falkordb]'. mypy skips its missing type information, as for the other untyped libraries. --- pyproject.toml | 7 ++++++- uv.lock | 50 ++++++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 56 insertions(+), 1 deletion(-) diff --git a/pyproject.toml b/pyproject.toml index 1c2f53c..709f5b1 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -53,6 +53,11 @@ Changelog = "https://github.com/cl0ver012/chatlore/blob/main/CHANGELOG.md" [project.scripts] chatlore = "chatlore.cli:run" +[project.optional-dependencies] +falkordb = [ + "falkordb>=1.7", +] + [build-system] requires = ["hatchling>=1.25"] build-backend = "hatchling.build" @@ -103,7 +108,7 @@ pretty = true plugins = ["pydantic.mypy"] [[tool.mypy.overrides]] -module = ["sqlite_vec", "fastembed", "networkx"] +module = ["sqlite_vec", "fastembed", "networkx", "falkordb", "falkordb.*"] ignore_missing_imports = true [tool.pytest.ini_options] diff --git a/uv.lock b/uv.lock index c304ae1..0d1624d 100644 --- a/uv.lock +++ b/uv.lock @@ -365,6 +365,11 @@ dependencies = [ { name = "uvicorn" }, ] +[package.optional-dependencies] +falkordb = [ + { name = "falkordb" }, +] + [package.dev-dependencies] dev = [ { name = "httpx" }, @@ -378,6 +383,7 @@ dev = [ [package.metadata] requires-dist = [ + { name = "falkordb", marker = "extra == 'falkordb'", specifier = ">=1.7" }, { name = "fastapi", specifier = ">=0.141" }, { name = "fastembed", specifier = ">=0.4" }, { name = "mcp", specifier = ">=2.2" }, @@ -391,6 +397,7 @@ requires-dist = [ { name = "typer", specifier = ">=0.15" }, { name = "uvicorn", specifier = ">=0.53" }, ] +provides-extras = ["falkordb"] [package.metadata.requires-dev] dev = [ @@ -579,6 +586,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/02/08/9c41fb51ab5b43eb21674aff13df270e8ba6c4b29c8624e328dc7a9482af/distlib-0.4.3-py2.py3-none-any.whl", hash = "sha256:4b0ce306c966eb73bc3a7b6abad017c556dadd92c44701562cd528ac7fde4d5b", size = 470628, upload-time = "2026-06-12T08:04:50.506Z" }, ] +[[package]] +name = "falkordb" +version = "1.7.1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "python-dateutil" }, + { name = "redis" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/98/d3/799a947869d9b8dd4cfcc6dc9654869874b1130375a78033b5a92362dc67/falkordb-1.7.1.tar.gz", hash = "sha256:09dd89dfb668c6fe7741c0ec67fcdb5c7a4b009e87f065a644199170f4fa5766", size = 121375, upload-time = "2026-08-13T13:17:39.652Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/bf/43/41967fad8b625e2b319e03a1c8897535dd1e05723a051342422b48b870da/falkordb-1.7.1-py3-none-any.whl", hash = "sha256:0e62d535edd5abf7b6d35e2493fba71118734313436ff803fcd5d1b20e9e2194", size = 38283, upload-time = "2026-08-13T13:17:38.257Z" }, +] + [[package]] name = "fastapi" version = "0.141.1" @@ -1746,6 +1766,18 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/9d/7a/d968e294073affff457b041c2be9868a40c1c71f4a35fcc1e45e5493067b/pytest_cov-7.1.0-py3-none-any.whl", hash = "sha256:a0461110b7865f9a271aa1b51e516c9a95de9d696734a2f71e3e78f46e1d4678", size = 22876, upload-time = "2026-03-21T20:11:14.438Z" }, ] +[[package]] +name = "python-dateutil" +version = "2.9.0.post0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "six" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/66/c0/0c8b6ad9f17a802ee498c46e004a0eb49bc148f2fd230864601a86dcf6db/python-dateutil-2.9.0.post0.tar.gz", hash = "sha256:37dd54208da7e1cd875388217d5e00ebd4179249f90fb72437e91a35459a0ad3", size = 342432, upload-time = "2024-03-01T18:36:20.211Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/ec/57/56b9bcc3c9c6a792fcbaf139543cee77261f3651ca9da0c93f5c1221264b/python_dateutil-2.9.0.post0-py2.py3-none-any.whl", hash = "sha256:a8b2bc7bffae282281c8140a97d3aa9c14da0b136dfe83f850eea9a5f7470427", size = 229892, upload-time = "2024-03-01T18:36:18.57Z" }, +] + [[package]] name = "python-discovery" version = "1.6.0" @@ -1841,6 +1873,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/f1/12/de94a39c2ef588c7e6455cfbe7343d3b2dc9d6b6b2f40c4c6565744c873d/pyyaml-6.0.3-cp314-cp314t-win_arm64.whl", hash = "sha256:ebc55a14a21cb14062aa4162f906cd962b28e2e9ea38f9b4391244cd8de4ae0b", size = 149341, upload-time = "2025-09-25T21:32:56.828Z" }, ] +[[package]] +name = "redis" +version = "8.1.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/a8/99/604f0b666d4c616d891cf77ebb9db6bb21601344c051aebf1b72b9ff915f/redis-8.1.0.tar.gz", hash = "sha256:6e1a19beef9225c83efd689c7e6b7da2d5215b1f42cd13b7fc3714d0a09c7b25", size = 5254356, upload-time = "2026-07-30T08:51:00.269Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/66/9d/c5731f6e3608663d4d3656fd8d3aecee8b509c3082818f5a13eae925baea/redis-8.1.0-py3-none-any.whl", hash = "sha256:a4fe1aac3d3b3cc791d4b3d5931c5a956045dc951ee74d1c913ee3ac4d2ee9fb", size = 560618, upload-time = "2026-07-30T08:50:58.497Z" }, +] + [[package]] name = "referencing" version = "0.37.0" @@ -2013,6 +2054,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/e0/f9/0595336914c5619e5f28a1fb793285925a8cd4b432c9da0a987836c7f822/shellingham-1.5.4-py2.py3-none-any.whl", hash = "sha256:7ecfff8f2fd72616f7481040475a65b2bf8af90a56c89140852d1120324e8686", size = 9755, upload-time = "2023-10-24T04:13:38.866Z" }, ] +[[package]] +name = "six" +version = "1.17.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/94/e7/b2c673351809dca68a0e064b6af791aa332cf192da575fd474ed7d6f16a2/six-1.17.0.tar.gz", hash = "sha256:ff70335d468e7eb6ec65b95b99d3a2836546063f63acc5171de367e834932a81", size = 34031, upload-time = "2024-12-04T17:35:28.174Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/b7/ce/149a00dd41f10bc29e5921b496af8b574d8413afcd5e30dfa0ed46c2cc5e/six-1.17.0-py2.py3-none-any.whl", hash = "sha256:4721f391ed90541fddacab5acf947aa0d3dc7d27b2e1e8eda2be8970586c3274", size = 11050, upload-time = "2024-12-04T17:35:26.475Z" }, +] + [[package]] name = "sniffio" version = "1.3.1" From 057e099e4cb766247c5d15be0502710f1eb20049 Mon Sep 17 00:00:00 2001 From: Cl0ver Date: Mon, 28 Sep 2026 11:31:45 -0400 Subject: [PATCH 2/9] feat(store): keep the graph in FalkorDB as well as SQLite CHATLORE_STORE=falkordb keeps the graph in a FalkorDB server instead of the SQLite file in the library folder, with CHATLORE_FALKORDB_URL and CHATLORE_FALKORDB_GRAPH saying where. Everything above the GraphStore interface is unchanged, and the store passes the same contract tests. Nodes carry their id and label as indexed properties and their props as JSON, since FalkorDB properties cannot hold maps or nulls, with simple values copied out for find_nodes and search filters. Where FalkorDB behaves differently from SQLite, the store makes up for it, each found by running against a real server: - full-text search does not fold accents, so indexed text and queries are folded here, and snippets are built here; stemming and stop words are turned off, so a word matches only itself, as in SQLite; - a vector of the wrong length and an edge to a missing node are accepted silently, so both are checked first; - a query returns at most 10,000 rows, so long reads come in pages; - within one query the first of two writes to the same node or edge wins, so the last one given is kept, as in SQLite; - a client made from a URL does not close its sockets, so they are closed here. chatlore index --rebuild drops either store, and chatlore doctor says which is in use. With CHATLORE_STORE=falkordb every test gets a graph of its own, so the whole suite runs against FalkorDB; the few tests that assumed the SQLite file now drop the store instead. --- src/chatlore/cli.py | 14 +- src/chatlore/store/__init__.py | 73 ++++- src/chatlore/store/falkordb.py | 516 +++++++++++++++++++++++++++++++++ tests/conftest.py | 25 ++ tests/store/conftest.py | 31 +- tests/store/test_falkordb.py | 89 ++++++ tests/store/test_sqlite.py | 7 +- tests/test_api.py | 6 +- tests/test_cli_search.py | 4 +- 9 files changed, 749 insertions(+), 16 deletions(-) create mode 100644 src/chatlore/store/falkordb.py create mode 100644 tests/store/test_falkordb.py diff --git a/src/chatlore/cli.py b/src/chatlore/cli.py index 44f6d21..a35d0cb 100644 --- a/src/chatlore/cli.py +++ b/src/chatlore/cli.py @@ -58,12 +58,13 @@ ) from chatlore.search import hybrid_search from chatlore.store import ( - DATABASE_NAME, EdgeType, GraphStore, Label, Node, TextHit, + describe_store, + drop_store, open_store, ) @@ -151,10 +152,18 @@ def doctor() -> None: table.add_row("python", f"{platform.python_version()} ({sys.executable})") table.add_row("platform", platform.platform()) table.add_row("data dir", f"{display_path(home)} ({state})") + table.add_row("store", _describe_store(home)) table.add_row("model", _describe_llm()) console.print(table) +def _describe_store(home: Path) -> str: + try: + return describe_store(home) + except ValueError as error: + return str(error) + + def _describe_llm() -> str: """Name the configured model and whether it has a key, without showing the key.""" try: @@ -349,8 +358,7 @@ def index( """Bring the search database in line with the library.""" home = default_home() if rebuild: - for name in (DATABASE_NAME, f"{DATABASE_NAME}-wal", f"{DATABASE_NAME}-shm"): - (home / name).unlink(missing_ok=True) + drop_store(home) count = 0 library = Library(home) if rebuild: diff --git a/src/chatlore/store/__init__.py b/src/chatlore/store/__init__.py index bb4af69..898f44d 100644 --- a/src/chatlore/store/__init__.py +++ b/src/chatlore/store/__init__.py @@ -1,8 +1,17 @@ -"""Graph storage: one interface, pluggable backends.""" +"""Graph storage: one interface, pluggable backends. + +SQLite, in the library folder, is the default. ``CHATLORE_STORE=falkordb`` keeps +the graph in a FalkorDB server instead: ``CHATLORE_FALKORDB_URL`` says where +(``redis://localhost:6379`` by default) and ``CHATLORE_FALKORDB_GRAPH`` which +graph (``chatlore``). The FalkorDB client is an optional dependency, installed +with ``chatlore[falkordb]``. +""" from __future__ import annotations +import os from pathlib import Path +from typing import TYPE_CHECKING from chatlore.store.base import ( ConversationSummary, @@ -17,10 +26,21 @@ ) from chatlore.store.sqlite import SQLiteStore +if TYPE_CHECKING: + from chatlore.store.falkordb import FalkorDBStore + DATABASE_NAME = "chatlore.db" +STORE_ENV = "CHATLORE_STORE" +FALKORDB_URL_ENV = "CHATLORE_FALKORDB_URL" +FALKORDB_GRAPH_ENV = "CHATLORE_FALKORDB_GRAPH" +FALKORDB_DEFAULT_URL = "redis://localhost:6379" +FALKORDB_DEFAULT_GRAPH = "chatlore" __all__ = [ "DATABASE_NAME", + "FALKORDB_GRAPH_ENV", + "FALKORDB_URL_ENV", + "STORE_ENV", "ConversationSummary", "Direction", "Edge", @@ -31,11 +51,60 @@ "SQLiteStore", "TextHit", "VectorHit", + "describe_store", + "drop_store", "open_store", ] def open_store(home: Path) -> GraphStore: - """Open the store that lives in a ChatLore home directory, creating it if needed.""" + """Open the library's store, creating it if needed.""" home.mkdir(parents=True, exist_ok=True) + if _backend() == "falkordb": + return _falkordb() return SQLiteStore(home / DATABASE_NAME) + + +def drop_store(home: Path) -> None: + """Delete the library's store, so the next ``open_store`` starts empty.""" + if _backend() == "falkordb": + store = _falkordb() + try: + store.drop() + finally: + store.close() + return + for name in (DATABASE_NAME, f"{DATABASE_NAME}-wal", f"{DATABASE_NAME}-shm"): + (home / name).unlink(missing_ok=True) + + +def describe_store(home: Path) -> str: + """Where the library's graph is kept, for people to read.""" + if _backend() == "falkordb": + url, graph = _falkordb_settings() + return f"FalkorDB graph '{graph}' at {url}" + return f"SQLite at {home / DATABASE_NAME}" + + +def _backend() -> str: + backend = (os.environ.get(STORE_ENV) or "sqlite").strip().lower() + if backend not in ("sqlite", "falkordb"): + raise ValueError(f"{STORE_ENV} must be sqlite or falkordb, not {backend!r}") + return backend + + +def _falkordb_settings() -> tuple[str, str]: + return ( + os.environ.get(FALKORDB_URL_ENV) or FALKORDB_DEFAULT_URL, + os.environ.get(FALKORDB_GRAPH_ENV) or FALKORDB_DEFAULT_GRAPH, + ) + + +def _falkordb() -> FalkorDBStore: + try: + from chatlore.store.falkordb import FalkorDBStore + except ImportError as error: + raise RuntimeError( + f"{STORE_ENV}=falkordb needs the FalkorDB client: pip install 'chatlore[falkordb]'" + ) from error + return FalkorDBStore(*_falkordb_settings()) diff --git a/src/chatlore/store/falkordb.py b/src/chatlore/store/falkordb.py new file mode 100644 index 0000000..8db374c --- /dev/null +++ b/src/chatlore/store/falkordb.py @@ -0,0 +1,516 @@ +"""A ``GraphStore`` on FalkorDB, the graph database that runs as a Redis module. + +Every node carries the label ``Node`` and these properties: + +- ``_id`` and ``_label``, the node's id and ChatLore label, both indexed; +- ``_props``, all of its props as JSON, since FalkorDB properties cannot hold + maps or nulls; +- a ``p_`` copy of each prop that is a string, number, or boolean, so + ``find_nodes`` and search filters can match on it (``text`` is left out: it + is only ever read back from ``_props``); +- ``_ft_title`` and ``_ft_text``, folded copies of the title and text of a node + that has text, in a full-text index; +- ``_embedding``, the node's vector, in a vector index. + +Edges are relationships of their own type with their props as JSON in +``_props``, and store-wide state lives on ``Meta`` nodes. + +FalkorDB differs from the SQLite store in ways this module makes up for. Its +full-text search does not fold accents, so the indexed copies and the queries +are folded here, and matching words are marked in snippets here too. It accepts +a vector of the wrong length and an edge to a missing node without complaint, so +both are checked first. It returns at most 10,000 rows per query, so long reads +are fetched in pages. Every query is atomic, but FalkorDB has no transactions +that span queries, so ``transaction`` only groups calls for the reader. +""" + +from __future__ import annotations + +import json +import re +import threading +import unicodedata +from collections.abc import Iterable, Iterator, Mapping, Sequence +from contextlib import contextmanager +from typing import Any + +from falkordb import FalkorDB +from redis.exceptions import ResponseError + +from chatlore.models import Conversation +from chatlore.store.base import ( + ConversationSummary, + Direction, + Edge, + EdgeType, + GraphStore, + Label, + Node, + TextHit, + VectorHit, +) +from chatlore.store.mapping import conversation_to_graph, graph_to_conversation + +_PAGE = 5_000 +"""Rows fetched or written per query, well under FalkorDB's 10,000-row limit.""" +_TOKEN = re.compile(r"[^\W_]+") +_NAME = re.compile(r"[A-Za-z_][A-Za-z0-9_]*") +_SNIPPET_TOKENS = 16 +_UNSEARCHABLE = frozenset({"text"}) +_DIMENSION = "embedding_dimension" + +# Indexes are created once per server and graph in each process. +_prepared: set[tuple[str, str]] = set() +_preparing = threading.Lock() + + +class FalkorDBStore(GraphStore): + """A ``GraphStore`` on one graph of a FalkorDB server.""" + + def __init__(self, url: str, graph: str) -> None: + self.url = url + self.name = graph + self._db = FalkorDB.from_url(url) + self._graph = self._db.select_graph(graph) + with _preparing: + if (url, graph) not in _prepared: + self._create_indexes() + _prepared.add((url, graph)) + dimension = self.get_meta(_DIMENSION) + self._dimension = int(dimension) if dimension is not None else None + + def close(self) -> None: + self._db.close() + # A client made from a URL does not own its connection pool, so closing it + # leaves the sockets open; they are closed here. + self._db.connection.connection_pool.disconnect() + + @contextmanager + def transaction(self) -> Iterator[None]: + """Group writes for the reader. Each query is atomic on its own in FalkorDB.""" + yield + + def drop(self) -> None: + """Delete the whole graph.""" + if self.name in self._db.list_graphs(): + self._graph.delete() + _prepared.discard((self.url, self.name)) + + # -- nodes and edges ----------------------------------------------------- + + def upsert_nodes(self, nodes: Iterable[Node]) -> None: + # One row per node, the last one given winning, as with SQLite: within one + # query FalkorDB would apply the first. + rows = list({node.id: _node_row(node) for node in nodes}.values()) + for page in _pages(rows): + # Replacing every property keeps the indexes in step; the embedding + # is carried over, as the SQLite store keeps it apart from the node. + self._query( + "UNWIND $rows AS row " + "MERGE (n:Node {_id: row._id}) " + "WITH n, row, n._embedding AS embedding " + "SET n = row " + "SET n._embedding = embedding", + {"rows": page}, + ) + + def upsert_edges(self, edges: Iterable[Edge]) -> None: + # One row per edge, the last one given winning, as for nodes. + by_type: dict[str, dict[tuple[str, str], dict[str, Any]]] = {} + for edge in edges: + by_type.setdefault(_name(edge.type), {})[(edge.src, edge.dst)] = { + "src": edge.src, + "dst": edge.dst, + "props": _dumps(edge.props), + } + ends = {end for rows in by_type.values() for pair in rows for end in pair} + missing = ends - self._existing(ends) + if missing: + raise ValueError( + f"FOREIGN KEY constraint failed: no node {sorted(missing)[0]!r} for an edge" + ) + for edge_type, rows in by_type.items(): + for page in _pages(list(rows.values())): + self._query( + "UNWIND $rows AS row " + "MATCH (a:Node {_id: row.src}), (b:Node {_id: row.dst}) " + f"MERGE (a)-[r:{edge_type}]->(b) " + "SET r._props = row.props", + {"rows": page}, + ) + + def get_node(self, node_id: str) -> Node | None: + rows = self._query( + "MATCH (n:Node {_id: $id}) RETURN n._id, n._label, n._props", {"id": node_id} + ) + return _node(rows[0]) if rows else None + + def delete_nodes(self, node_ids: Iterable[str]) -> None: + for page in _pages(list(node_ids)): + # Edges and index entries go with the node. + self._query("UNWIND $ids AS id MATCH (n:Node {_id: id}) DETACH DELETE n", {"ids": page}) + + def neighbors( + self, + node_id: str, + edge_types: Sequence[str] | None = None, + direction: Direction = "out", + limit: int = 100, + ) -> list[tuple[Edge, Node]]: + types = ":" + "|".join(_name(edge_type) for edge_type in edge_types) if edge_types else "" + pattern = { + "out": f"(a)-[r{types}]->(b:Node)", + "in": f"(a)<-[r{types}]-(b:Node)", + "both": f"(a)-[r{types}]-(b:Node)", + }[direction] + rows = self._paged( + f"MATCH (a:Node {{_id: $id}}) MATCH {pattern} " + "RETURN type(r), startNode(r)._id, endNode(r)._id, r._props, " + "b._id, b._label, b._props " + "ORDER BY type(r), endNode(r)._id, startNode(r)._id", + {"id": node_id}, + limit, + ) + return [(Edge(row[1], row[0], row[2], _loads(row[3])), _node(row[4:7])) for row in rows] + + def find_nodes( + self, label: str, where: Mapping[str, Any] | None = None, limit: int = 1_000_000 + ) -> list[Node]: + clauses = ["n._label = $label"] + params: dict[str, Any] = {"label": label} + for number, (key, value) in enumerate((where or {}).items()): + if not _NAME.fullmatch(key): + raise ValueError(f"invalid property name: {key!r}") + clauses.append(f"n.p_{key} = $value{number}") + params[f"value{number}"] = value + rows = self._paged( + f"MATCH (n:Node) WHERE {' AND '.join(clauses)} " + "RETURN n._id, n._label, n._props ORDER BY n._id", + params, + limit, + ) + return [_node(row) for row in rows] + + def count_nodes(self, label: str | None = None) -> int: + if label is None: + rows = self._query("MATCH (n:Node) RETURN count(n)") + else: + rows = self._query( + "MATCH (n:Node) WHERE n._label = $label RETURN count(n)", {"label": label} + ) + return int(rows[0][0]) + + # -- conversations ------------------------------------------------------- + + def upsert_conversation(self, conversation: Conversation) -> None: + nodes, edges = conversation_to_graph(conversation) + keep = {message.id for message in conversation.messages} + # As in SQLite: messages still present are updated in place, so their + # chunks and embeddings survive; only messages that disappeared go. + self._delete_messages(conversation.id, keep=keep) + self._query( + f"MATCH (:Node {{_id: $id}})-[:{EdgeType.HAS_MESSAGE}]->(:Node)" + f"-[r:{EdgeType.REPLIES_TO}]->() DELETE r", + {"id": conversation.id}, + ) + self.upsert_nodes(nodes) + self.upsert_edges(edges) + + def get_conversation(self, conversation_id: str) -> Conversation | None: + node = self.get_node(conversation_id) + if node is None or node.label != Label.CONVERSATION: + return None + messages = [ + neighbor + for _, neighbor in self.neighbors( + conversation_id, [EdgeType.HAS_MESSAGE], "out", limit=1_000_000 + ) + ] + return graph_to_conversation(node, messages) + + def list_conversations( + self, source: str | None = None, limit: int = 50, offset: int = 0 + ) -> list[ConversationSummary]: + where = "n._label = $label" + (" AND n.p_source = $source" if source is not None else "") + rows = self._query( + f"MATCH (n:Node) WHERE {where} RETURN n._id, n._props " + "ORDER BY coalesce(n.p_updated_at, n.p_created_at, '') DESC, n._id " + "SKIP $offset LIMIT $limit", + {"label": Label.CONVERSATION.value, "source": source, "offset": offset, "limit": limit}, + ) + summaries: list[ConversationSummary] = [] + for node_id, raw in rows: + props = _loads(raw) + summaries.append( + ConversationSummary( + id=node_id, + source=str(props.get("source")), + title=props.get("title"), + created_at=props.get("created_at"), + updated_at=props.get("updated_at"), + message_count=int(props.get("message_count", 0)), + ) + ) + return summaries + + def delete_conversation(self, conversation_id: str) -> None: + self._delete_messages(conversation_id) + self.delete_nodes([conversation_id]) + + # -- search -------------------------------------------------------------- + + def search_text( + self, + query: str, + limit: int = 20, + sources: Sequence[str] | None = None, + labels: Sequence[str] | None = None, + ) -> list[TextHit]: + words = fold_words(query) + if not words: + return [] + clauses = [] + if sources: + clauses.append("node.p_source IN $sources") + if labels: + clauses.append("node._label IN $labels") + where = f"WHERE {' AND '.join(clauses)} " if clauses else "" + rows = self._query( + "CALL db.idx.fulltext.queryNodes('Node', $query) YIELD node, score " + f"{where}" + "RETURN node._id, node._label, node._props, score " + "ORDER BY score DESC, node._id LIMIT $limit", + { + "query": " ".join(words), + "sources": list(sources or []), + "labels": list(labels or []), + "limit": limit, + }, + ) + hits: list[TextHit] = [] + for node_id, label, raw, score in rows: + props = _loads(raw) + hits.append( + TextHit( + node_id=node_id, + label=label, + conversation_id=props.get("conversation_id"), + source=props.get("source"), + title=props.get("title") or None, + snippet=snippet(str(props.get("text", "")), set(words)), + score=float(score), + ) + ) + return hits + + def set_embedding(self, node_id: str, embedding: Sequence[float]) -> None: + vector = [float(value) for value in embedding] + if not vector: + raise ValueError("embedding is empty") + if self.get_node(node_id) is None: + raise KeyError(node_id) + if self._dimension is None: + self._query( + "CREATE VECTOR INDEX FOR (n:Node) ON (n._embedding) " + f"OPTIONS {{dimension: {len(vector)}, similarityFunction: 'euclidean'}}" + ) + self.set_meta(_DIMENSION, str(len(vector))) + self._dimension = len(vector) + elif len(vector) != self._dimension: + raise ValueError(f"embedding has {len(vector)} dimensions, store has {self._dimension}") + self._query( + "MATCH (n:Node {_id: $id}) SET n._embedding = vecf32($vector)", + {"id": node_id, "vector": vector}, + ) + + def nodes_without_embedding(self, label: str, limit: int = 100) -> list[Node]: + rows = self._paged( + "MATCH (n:Node) WHERE n._label = $label AND n._embedding IS NULL " + "RETURN n._id, n._label, n._props ORDER BY n._id", + {"label": label}, + limit, + ) + return [_node(row) for row in rows] + + def count_embeddings(self) -> int: + rows = self._query("MATCH (n:Node) WHERE n._embedding IS NOT NULL RETURN count(n)") + return int(rows[0][0]) + + def clear_embeddings(self) -> None: + if self._dimension is not None: + self._query("DROP VECTOR INDEX FOR (n:Node) ON (n._embedding)") + self._query("MATCH (n:Node) WHERE n._embedding IS NOT NULL SET n._embedding = NULL") + self._query( + "MATCH (m:Meta) WHERE m.key IN $keys DELETE m", + {"keys": [_DIMENSION, "embedding_model"]}, + ) + self._dimension = None + + def get_meta(self, key: str) -> str | None: + rows = self._query("MATCH (m:Meta {key: $key}) RETURN m.value", {"key": key}) + return str(rows[0][0]) if rows else None + + def set_meta(self, key: str, value: str) -> None: + self._query("MERGE (m:Meta {key: $key}) SET m.value = $value", {"key": key, "value": value}) + + def search_vector( + self, embedding: Sequence[float], limit: int = 20, labels: Sequence[str] | None = None + ) -> list[VectorHit]: + if self._dimension is None: + return [] + vector = [float(value) for value in embedding] + if len(vector) != self._dimension: + raise ValueError(f"embedding has {len(vector)} dimensions, store has {self._dimension}") + wanted = set(labels or []) + # Label filtering happens after the k-nearest query, so ask for extra candidates. + k = limit * 4 if wanted else limit + if k <= 0: + return [] + rows = self._query( + "CALL db.idx.vector.queryNodes('Node', '_embedding', $k, vecf32($vector)) " + "YIELD node, score RETURN node._id, node._label, score ORDER BY score, node._id", + {"k": k, "vector": vector}, + ) + hits = [ + VectorHit(node_id, label, float(score)) + for node_id, label, score in rows + if not wanted or label in wanted + ] + return hits[:limit] + + # -- internals ----------------------------------------------------------- + + def _query(self, query: str, params: dict[str, Any] | None = None) -> list[list[Any]]: + rows: list[list[Any]] = self._graph.query(query, params).result_set + return rows + + def _paged(self, query: str, params: dict[str, Any], limit: int) -> list[list[Any]]: + """Up to ``limit`` rows of an ordered query, fetched a page at a time.""" + rows: list[list[Any]] = [] + while len(rows) < limit: + size = min(_PAGE, limit - len(rows)) + page = self._query( + f"{query} SKIP $skip LIMIT $size", {**params, "skip": len(rows), "size": size} + ) + rows.extend(page) + if len(page) < size: + break + return rows + + def _existing(self, ids: set[str]) -> set[str]: + found: set[str] = set() + for page in _pages(sorted(ids)): + rows = self._query("MATCH (n:Node) WHERE n._id IN $ids RETURN n._id", {"ids": page}) + found.update(row[0] for row in rows) + return found + + def _delete_messages(self, conversation_id: str, keep: set[str] | None = None) -> None: + """Delete a conversation's messages, except ``keep``, along with their chunks.""" + rows = self._query( + f"MATCH (:Node {{_id: $id}})-[:{EdgeType.HAS_MESSAGE}]->(m:Node) RETURN m._id", + {"id": conversation_id}, + ) + doomed = [row[0] for row in rows if not keep or row[0] not in keep] + chunks: list[str] = [] + for page in _pages(doomed): + chunk_rows = self._query( + f"MATCH (m:Node)-[:{EdgeType.HAS_CHUNK}]->(c:Node) WHERE m._id IN $ids " + "RETURN c._id", + {"ids": page}, + ) + chunks.extend(row[0] for row in chunk_rows) + self.delete_nodes([*chunks, *doomed]) + + def _create_indexes(self) -> None: + statements = [ + "CREATE INDEX FOR (n:Node) ON (n._id)", + "CREATE INDEX FOR (n:Node) ON (n._label)", + "CREATE INDEX FOR (n:Node) ON (n.p_conversation_id)", + "CREATE INDEX FOR (m:Meta) ON (m.key)", + # No stop words and no stemming, so a word matches that word only, + # as in the SQLite store. The title counts twice. + "CALL db.idx.fulltext.createNodeIndex({label: 'Node', stopwords: []}, " + "{field: '_ft_title', nostem: true, weight: 2}, " + "{field: '_ft_text', nostem: true})", + ] + for statement in statements: + try: + self._query(statement) + except ResponseError as error: + if "already indexed" not in str(error): + raise + + +def fold_words(text: str) -> list[str]: + """The words of ``text``, lowercased and without accents, as they are indexed. + + Letters and digits make words and everything else separates them, so + ``customer_id`` is two words and ``café's`` is ``cafe`` and ``s``. + """ + decomposed = unicodedata.normalize("NFKD", text) + plain = "".join(char for char in decomposed if not unicodedata.combining(char)) + return _TOKEN.findall(plain.casefold()) + + +def snippet(text: str, words: set[str]) -> str: + """About sixteen words of ``text`` around the most matches, matching words in [ ].""" + spans = [(match.start(), match.end()) for match in _TOKEN.finditer(text)] + if not spans: + return text + matched = [bool(set(fold_words(text[start:end])) & words) for start, end in spans] + window = min(_SNIPPET_TOKENS, len(spans)) + best = max( + range(len(spans) - window + 1), + key=lambda first: (sum(matched[first : first + window]), -first), + ) + last = best + window - 1 + begin = 0 if best == 0 else spans[best][0] + end = len(text) if last == len(spans) - 1 else spans[last][1] + pieces: list[str] = [] + cursor = begin + for (start, stop), hit in zip(spans[best : last + 1], matched[best : last + 1], strict=True): + pieces.append(text[cursor:start]) + pieces.append(f"[{text[start:stop]}]" if hit else text[start:stop]) + cursor = stop + pieces.append(text[cursor:end]) + body = "".join(pieces).strip() + return ("..." if best > 0 else "") + body + ("..." if end < len(text) else "") + + +def _node_row(node: Node) -> dict[str, Any]: + row: dict[str, Any] = {"_id": node.id, "_label": node.label, "_props": _dumps(node.props)} + for key, value in node.props.items(): + if key in _UNSEARCHABLE or not _NAME.fullmatch(key): + continue + if isinstance(value, str | int | float | bool): + row[f"p_{key}"] = value + text = node.props.get("text") + if isinstance(text, str) and text.strip(): + row["_ft_text"] = " ".join(fold_words(text)) + row["_ft_title"] = " ".join(fold_words(str(node.props.get("title") or ""))) + return row + + +def _node(row: Sequence[Any]) -> Node: + return Node(row[0], row[1], _loads(row[2])) + + +def _name(value: str) -> str: + """An edge type, which FalkorDB cannot take as a parameter, checked before use.""" + if not _NAME.fullmatch(value): + raise ValueError(f"invalid edge type: {value!r}") + return value + + +def _pages[T](items: list[T]) -> Iterator[list[T]]: + for start in range(0, len(items), _PAGE): + yield items[start : start + _PAGE] + + +def _dumps(value: dict[str, Any]) -> str: + return json.dumps(value, ensure_ascii=False, separators=(",", ":"), default=str) + + +def _loads(value: str | None) -> dict[str, Any]: + data = json.loads(value) if value else {} + return data if isinstance(data, dict) else {} diff --git a/tests/conftest.py b/tests/conftest.py index d937ca1..136d2b2 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -2,7 +2,10 @@ from __future__ import annotations +import os +import uuid import zipfile +from collections.abc import Iterator from pathlib import Path import pytest @@ -10,6 +13,28 @@ FIXTURES = Path(__file__).parent / "fixtures" +@pytest.fixture(autouse=True) +def falkordb_graph(monkeypatch: pytest.MonkeyPatch) -> Iterator[None]: + """With ``CHATLORE_STORE=falkordb``, give each test a graph of its own. + + That runs the whole suite against a FalkorDB server, as CI does. + """ + if os.environ.get("CHATLORE_STORE") != "falkordb": + yield + return + from chatlore.store import FALKORDB_DEFAULT_URL, FALKORDB_URL_ENV + from chatlore.store.falkordb import FalkorDBStore + + graph = f"test_{uuid.uuid4().hex}" + monkeypatch.setenv("CHATLORE_FALKORDB_GRAPH", graph) + yield + store = FalkorDBStore(os.environ.get(FALKORDB_URL_ENV) or FALKORDB_DEFAULT_URL, graph) + try: + store.drop() + finally: + store.close() + + @pytest.fixture def fixtures() -> Path: """Directory holding the synthetic export fixtures.""" diff --git a/tests/store/conftest.py b/tests/store/conftest.py index 29e7a66..b956f8f 100644 --- a/tests/store/conftest.py +++ b/tests/store/conftest.py @@ -1,18 +1,39 @@ -"""Backends under test. Every contract test runs against each entry.""" +"""Backends under test. Every contract test runs against each entry. + +FalkorDB needs a server: set ``CHATLORE_TEST_FALKORDB_URL`` to run against one, +for example ``redis://localhost:6379`` with ``docker run -p 6379:6379 +falkordb/falkordb``. Each test gets a graph of its own, deleted afterwards. +""" from __future__ import annotations +import os +import uuid from collections.abc import Iterator import pytest from chatlore.store import GraphStore, SQLiteStore +FALKORDB_URL = os.environ.get("CHATLORE_TEST_FALKORDB_URL") + -@pytest.fixture(params=["sqlite"]) +@pytest.fixture(params=["sqlite", "falkordb"]) def store(request: pytest.FixtureRequest) -> Iterator[GraphStore]: - backend = SQLiteStore(":memory:") + if request.param == "sqlite": + backend: GraphStore = SQLiteStore(":memory:") + try: + yield backend + finally: + backend.close() + return + if not FALKORDB_URL: + pytest.skip("set CHATLORE_TEST_FALKORDB_URL to test against FalkorDB") + from chatlore.store.falkordb import FalkorDBStore + + falkordb = FalkorDBStore(FALKORDB_URL, f"test_{uuid.uuid4().hex}") try: - yield backend + yield falkordb finally: - backend.close() + falkordb.drop() + falkordb.close() diff --git a/tests/store/test_falkordb.py b/tests/store/test_falkordb.py new file mode 100644 index 0000000..2e89572 --- /dev/null +++ b/tests/store/test_falkordb.py @@ -0,0 +1,89 @@ +"""Tests for choosing a store and for the FalkorDB backend's text helpers. + +The FalkorDB store itself runs the contract tests when a server is available; +these need none. +""" + +from __future__ import annotations + +from pathlib import Path +from typing import Any + +import pytest + +from chatlore.store import ( + DATABASE_NAME, + SQLiteStore, + describe_store, + drop_store, + open_store, +) + +falkordb = pytest.importorskip("chatlore.store.falkordb") + + +def test_words_are_folded_as_they_are_indexed() -> None: + assert falkordb.fold_words("The Café's RÉSUMÉ, customer_id 42!") == [ + "the", + "cafe", + "s", + "resume", + "customer", + "id", + "42", + ] + assert falkordb.fold_words("(draft). --") == ["draft"] + assert falkordb.fold_words("") == [] + + +def test_snippets_mark_matching_words_around_the_most_matches() -> None: + text = " ".join(f"w{number}" for number in range(40)) + " Postgres index here." + + short = falkordb.snippet("Add an index on customer_id.", {"index"}) + long = falkordb.snippet(text, {"postgres", "index"}) + + assert short == "Add an [index] on customer_id." + assert long.startswith("...w26 ") + assert long.endswith("[Postgres] [index]...") + assert falkordb.snippet("Café au lait", {"cafe"}) == "[Café] au lait" + assert falkordb.snippet("", {"x"}) == "" + + +def test_sqlite_is_the_default_store(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.delenv("CHATLORE_STORE", raising=False) + + with open_store(tmp_path) as store: + assert isinstance(store, SQLiteStore) + assert (tmp_path / DATABASE_NAME).exists() + assert describe_store(tmp_path) == f"SQLite at {tmp_path / DATABASE_NAME}" + + drop_store(tmp_path) + + assert not (tmp_path / DATABASE_NAME).exists() + + +def test_falkordb_is_chosen_by_setting(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + opened: list[tuple[Any, ...]] = [] + + class Fake: + def __init__(self, *settings: Any) -> None: + opened.append(settings) + + monkeypatch.setattr(falkordb, "FalkorDBStore", Fake) + monkeypatch.setenv("CHATLORE_STORE", "FalkorDB") + monkeypatch.delenv("CHATLORE_FALKORDB_URL", raising=False) + monkeypatch.setenv("CHATLORE_FALKORDB_GRAPH", "mine") + + store = open_store(tmp_path) + + assert isinstance(store, Fake) + assert opened == [("redis://localhost:6379", "mine")] + assert describe_store(tmp_path) == "FalkorDB graph 'mine' at redis://localhost:6379" + assert not (tmp_path / DATABASE_NAME).exists() + + +def test_an_unknown_store_is_refused(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("CHATLORE_STORE", "neo4j") + + with pytest.raises(ValueError, match="sqlite or falkordb"): + open_store(tmp_path) diff --git a/tests/store/test_sqlite.py b/tests/store/test_sqlite.py index 3718d0c..9e05abc 100644 --- a/tests/store/test_sqlite.py +++ b/tests/store/test_sqlite.py @@ -5,6 +5,8 @@ import contextlib from pathlib import Path +import pytest + from chatlore.models import ContentPart, Conversation, Message, Role, SourceKind from chatlore.store import Edge, EdgeType, Label, Node, SQLiteStore, open_store from chatlore.store.sqlite import fts_query @@ -22,7 +24,10 @@ def test_data_persists_across_connections(tmp_path: Path) -> None: assert [h.node_id for h in store.search_vector([0.5, 0.5])] == ["n"] -def test_open_store_creates_the_home_directory(tmp_path: Path) -> None: +def test_open_store_creates_the_home_directory( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + monkeypatch.delenv("CHATLORE_STORE", raising=False) home = tmp_path / "fresh" / "home" with open_store(home) as store: diff --git a/tests/test_api.py b/tests/test_api.py index 18e8177..336bd15 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -13,6 +13,7 @@ from chatlore.api import create_app from chatlore.cli import app as cli +from chatlore.store import drop_store from tests.fakes import FakeEmbedder, FakeLLM runner = CliRunner() @@ -134,10 +135,9 @@ def test_chat_with_nothing_to_go_on_does_not_call_the_model( client: TestClient, llm: FakeLLM, home: Path ) -> None: asked = len(llm.chat_requests) - empty = home.parent / "empty" + drop_store(home) - with TestClient(create_app(empty)) as fresh: - events = _events(fresh.post("/chat", json={"question": "Anything?"}).text) + events = _events(client.post("/chat", json={"question": "Anything?"}).text) assert events == [("sources", []), ("done", {"cited": [], "found": False})] assert len(llm.chat_requests) == asked diff --git a/tests/test_cli_search.py b/tests/test_cli_search.py index 2af0dc6..a338294 100644 --- a/tests/test_cli_search.py +++ b/tests/test_cli_search.py @@ -9,7 +9,7 @@ from typer.testing import CliRunner from chatlore.cli import app -from chatlore.store import open_store +from chatlore.store import drop_store, open_store from tests.fakes import FakeEmbedder runner = CliRunner() @@ -54,7 +54,7 @@ def test_dry_run_does_not_create_the_database(home: Path, fixtures: Path) -> Non def test_index_rebuilds_from_the_library(home: Path, fixtures: Path) -> None: runner.invoke(app, ["import", str(fixtures / "claude")]) runner.invoke(app, ["import", str(fixtures / "gemini")]) - (home / "chatlore.db").unlink() + drop_store(home) result = runner.invoke(app, ["index", "--rebuild"]) From 2743bbdf6cb4d69abf52127968e56787f2a1885d Mon Sep 17 00:00:00 2001 From: Cl0ver Date: Mon, 28 Sep 2026 11:31:45 -0400 Subject: [PATCH 3/9] ci: test the FalkorDB store against a FalkorDB server The new job starts FalkorDB 4.20.7 as a service, runs the contract tests on both stores, then runs every test with FalkorDB as the store. --- .github/workflows/ci.yml | 31 +++++++++++++++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 224940e..4a13f87 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -68,3 +68,34 @@ jobs: uvx --from "$wheel" chatlore --version uvx --from "$wheel" chatlore demo --no-serve uvx --from "$wheel" chatlore --home ~/.chatlore-demo search "litestream" + + falkordb: + name: FalkorDB backend + runs-on: ubuntu-latest + services: + falkordb: + image: falkordb/falkordb:v4.20.7 + ports: ["6379:6379"] + env: + CHATLORE_FALKORDB_URL: redis://localhost:6379 + CHATLORE_TEST_FALKORDB_URL: redis://localhost:6379 + steps: + - name: Check out + uses: actions/checkout@v4 + + - name: Set up uv + uses: astral-sh/setup-uv@v6 + with: + python-version: "3.12" + enable-cache: true + + - name: Install dependencies + run: uv sync --locked --extra falkordb + + - name: Contract tests on both stores + run: uv run pytest tests/store + + - name: Every test with FalkorDB as the store + run: uv run pytest + env: + CHATLORE_STORE: falkordb From 238d9ad1603724e011c9963f8e752b471f1cd65c Mon Sep 17 00:00:00 2001 From: Cl0ver Date: Mon, 28 Sep 2026 11:31:45 -0400 Subject: [PATCH 4/9] docs: describe keeping the graph in FalkorDB docs/falkordb.md covers setting it up, moving a library into it from the caches without computing anything again, how the graph looks for Cypher queries, and how it differs from SQLite, with timings from a library of 12,000 notes. Roadmap: M10 is done. --- CHANGELOG.md | 9 ++++++ README.md | 5 ++-- docs/falkordb.md | 78 ++++++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 90 insertions(+), 2 deletions(-) create mode 100644 docs/falkordb.md diff --git a/CHANGELOG.md b/CHANGELOG.md index 986267f..2ff322a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,15 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added + +- A FalkorDB store, chosen with `CHATLORE_STORE=falkordb`, with `CHATLORE_FALKORDB_URL` and + `CHATLORE_FALKORDB_GRAPH` saying where. Every command, the web interface, and the MCP server + work on it as on SQLite, and it passes the same contract tests. The client is an optional + dependency: `pip install 'chatlore[falkordb]'`. Guide in `docs/falkordb.md`. +- `chatlore doctor` shows which store is in use, and `chatlore index --rebuild` rebuilds either. +- CI runs the contract tests on both stores, and every test with FalkorDB as the store. + ## [0.1.0] - 2026-09-24 ### Added diff --git a/README.md b/README.md index e5486ed..8f368c5 100644 --- a/README.md +++ b/README.md @@ -59,7 +59,8 @@ see [docs/models.md](docs/models.md). The knowledge graph is described in [docs/extraction.md](docs/extraction.md), chat and the API in [docs/chat.md](docs/chat.md), the web interface in [docs/web.md](docs/web.md), setting up assistants over MCP in [docs/mcp.md](docs/mcp.md), and archives and -Markdown export in [docs/export.md](docs/export.md). +Markdown export in [docs/export.md](docs/export.md), and keeping the graph in +FalkorDB instead of SQLite in [docs/falkordb.md](docs/falkordb.md). ## What ChatLore will do @@ -95,7 +96,7 @@ extraction is an optional enrichment you can re-run with a better model later. | M7 | MCP server for Claude Desktop, Claude Code, Cursor, ChatGPT | done for local assistants; ChatGPT with the hosted demo | | M8 | Easy to try: PyPI package, demo library, export and import | done, v0.1.0 | | M9 | Hosted demo | planned | -| M10 | FalkorDB backend | planned | +| M10 | FalkorDB backend | done | ## Development setup diff --git a/docs/falkordb.md b/docs/falkordb.md new file mode 100644 index 0000000..0ee0c23 --- /dev/null +++ b/docs/falkordb.md @@ -0,0 +1,78 @@ +# Keeping the graph in FalkorDB + +```bash +pip install 'chatlore[falkordb]' +docker run -d -p 6379:6379 --name falkordb falkordb/falkordb +export CHATLORE_STORE=falkordb +chatlore index # copy the library into FalkorDB +chatlore process # chunks and embeddings +chatlore extract # entities and topics, from the cached answers where possible +``` + +By default ChatLore keeps its graph in a SQLite file in the library folder, +which needs nothing else. [FalkorDB](https://www.falkordb.com) is a graph +database that runs as a server. Use it when you want to query the graph with +Cypher, look at it with FalkorDB's browser, or share one graph between several +machines. Every command, the web interface, and the MCP server work the same +on both. + +## Settings + +| Variable | Meaning | +|---|---| +| `CHATLORE_STORE` | `sqlite` (default) or `falkordb`. | +| `CHATLORE_FALKORDB_URL` | The server, `redis://localhost:6379` by default. Add a password as `redis://:password@host:6379`. | +| `CHATLORE_FALKORDB_GRAPH` | The graph's name, `chatlore` by default. Give each library its own. | + +`chatlore doctor` shows which store is in use. + +## Moving a library into FalkorDB + +The conversations stay in the library folder either way, and the caches of +embeddings and model answers too, so the graph can be rebuilt in FalkorDB +without computing anything again: + +1. Set the variables above. +2. `chatlore index` copies every conversation into the graph. +3. `chatlore process` makes the chunks and takes their embeddings from the cache. +4. `chatlore extract` rebuilds entities, relationships, and topics from the + cached model answers. It needs a model key, but only asks the model about + text it has not read yet. + +Going back to SQLite is the same with `CHATLORE_STORE` unset. An archive from +`chatlore export` imports into either store. `chatlore index --rebuild` deletes +the graph and builds it again from the library. + +## What is in the graph + +Every node has the label `Node`, with its id in `_id`, its ChatLore label +(`Conversation`, `Message`, `Chunk`, `Entity`, `Topic`) in `_label`, and all +of its properties as JSON in `_props`. Simple values are copied into `p_` +properties as well, so they can be queried: an entity's name is `p_name`, a +topic's title is `p_title`. Edges have their ChatLore type, such as `MENTIONS`, +`RELATED_TO`, or `IN_TOPIC`, with their properties as JSON in `_props`. + +```cypher +MATCH (e:Node {_label: 'Entity'})-[r:RELATED_TO]-(o:Node) +WHERE e.p_name = 'Tidewater' +RETURN o.p_name, r._props +``` + +Everything ChatLore adds for search, the `_ft_` text copies and the +`_embedding` vectors, is indexed in FalkorDB. Treat the graph as ChatLore's: +changes made to it directly are lost on the next `chatlore index --rebuild`. + +## Differences from SQLite + +- FalkorDB has no transactions spanning several queries. Each write is atomic, + but an import interrupted halfway leaves the conversations written so far; + running it again finishes the job, as with SQLite. +- Word search matches the same words as SQLite, without stemming or stop words, + and ignores accents. Its scores are FalkorDB's own, so results with equal + matches can come back in a different order. +- The server has to be running. Commands fail with a connection error when it + is not. +- Each conversation takes a few round trips to the server, so importing is + slower. On a made-up library of 12,000 notes, importing took 31 seconds + against 8 with SQLite, chunking 26 seconds against 145, and a search the same + time on both. From e4396bc6bbfcf755124dc4f13a2fe47b9c697a23 Mon Sep 17 00:00:00 2001 From: Cl0ver Date: Mon, 28 Sep 2026 14:50:21 -0400 Subject: [PATCH 5/9] fix(store): reuse FalkorDB connections and hide the password The web interface and the MCP server open the store for every request, and each FalkorDB store connected anew: several round trips, 3.5 seconds to a FalkorDB Cloud server in another region, before the first query. Stores now share one client per server for the whole process, closed when it exits, which brought a stats request there from 5 seconds to 1. Networks drop idle connections, which then timed out, so a connection unused for half a minute is checked before it is used. chatlore doctor printed the server URL with its password; it now shows *** instead. The docs say to keep the server close, with the timings seen against a distant one. --- docs/falkordb.md | 6 +++++- src/chatlore/store/__init__.py | 13 ++++++++++++- src/chatlore/store/falkordb.py | 35 ++++++++++++++++++++++++++-------- tests/store/test_falkordb.py | 12 ++++++++++++ 4 files changed, 56 insertions(+), 10 deletions(-) diff --git a/docs/falkordb.md b/docs/falkordb.md index 0ee0c23..927fea7 100644 --- a/docs/falkordb.md +++ b/docs/falkordb.md @@ -13,7 +13,7 @@ By default ChatLore keeps its graph in a SQLite file in the library folder, which needs nothing else. [FalkorDB](https://www.falkordb.com) is a graph database that runs as a server. Use it when you want to query the graph with Cypher, look at it with FalkorDB's browser, or share one graph between several -machines. Every command, the web interface, and the MCP server work the same +machines on a network. Every command, the web interface, and the MCP server work the same on both. ## Settings @@ -72,6 +72,10 @@ changes made to it directly are lost on the next `chatlore index --rebuild`. matches can come back in a different order. - The server has to be running. Commands fail with a connection error when it is not. +- Every page and command asks the server many small questions, so keep it + close: on the same machine or network. Against a FalkorDB Cloud server a + continent away, 0.16 seconds per query, loading the demo took a minute and a + half, a search 8 to 10 seconds, and the graph view almost two minutes. - Each conversation takes a few round trips to the server, so importing is slower. On a made-up library of 12,000 notes, importing took 31 seconds against 8 with SQLite, chunking 26 seconds against 145, and a search the same diff --git a/src/chatlore/store/__init__.py b/src/chatlore/store/__init__.py index 898f44d..66a9aef 100644 --- a/src/chatlore/store/__init__.py +++ b/src/chatlore/store/__init__.py @@ -12,6 +12,7 @@ import os from pathlib import Path from typing import TYPE_CHECKING +from urllib.parse import urlsplit, urlunsplit from chatlore.store.base import ( ConversationSummary, @@ -82,10 +83,20 @@ def describe_store(home: Path) -> str: """Where the library's graph is kept, for people to read.""" if _backend() == "falkordb": url, graph = _falkordb_settings() - return f"FalkorDB graph '{graph}' at {url}" + return f"FalkorDB graph '{graph}' at {_without_password(url)}" return f"SQLite at {home / DATABASE_NAME}" +def _without_password(url: str) -> str: + """The URL with any password replaced, so it can be shown.""" + parts = urlsplit(url) + if parts.password is None: + return url + user = f"{parts.username}:" if parts.username else ":" + host = parts.netloc.rsplit("@", 1)[1] + return urlunsplit(parts._replace(netloc=f"{user}***@{host}")) + + def _backend() -> str: backend = (os.environ.get(STORE_ENV) or "sqlite").strip().lower() if backend not in ("sqlite", "falkordb"): diff --git a/src/chatlore/store/falkordb.py b/src/chatlore/store/falkordb.py index 8db374c..e2c3b36 100644 --- a/src/chatlore/store/falkordb.py +++ b/src/chatlore/store/falkordb.py @@ -26,6 +26,7 @@ from __future__ import annotations +import atexit import json import re import threading @@ -59,9 +60,30 @@ _UNSEARCHABLE = frozenset({"text"}) _DIMENSION = "embedding_dimension" -# Indexes are created once per server and graph in each process. +# One client per server for the whole process. Connecting takes several round +# trips, seconds to a server far away, and the web interface and MCP server open +# the store for every request. Indexes are created once per server and graph. +_clients: dict[str, FalkorDB] = {} _prepared: set[tuple[str, str]] = set() -_preparing = threading.Lock() +_lock = threading.Lock() + + +def _client(url: str) -> FalkorDB: + with _lock: + if url not in _clients: + # Networks drop connections that sit idle, so one unused for half a + # minute is checked before it is used again. + _clients[url] = FalkorDB.from_url(url, health_check_interval=30, socket_keepalive=True) + return _clients[url] + + +@atexit.register +def _disconnect() -> None: + """Close every client's sockets; a client made from a URL does not close them itself.""" + with _lock: + for client in _clients.values(): + client.connection.connection_pool.disconnect() + _clients.clear() class FalkorDBStore(GraphStore): @@ -70,9 +92,9 @@ class FalkorDBStore(GraphStore): def __init__(self, url: str, graph: str) -> None: self.url = url self.name = graph - self._db = FalkorDB.from_url(url) + self._db = _client(url) self._graph = self._db.select_graph(graph) - with _preparing: + with _lock: if (url, graph) not in _prepared: self._create_indexes() _prepared.add((url, graph)) @@ -80,10 +102,7 @@ def __init__(self, url: str, graph: str) -> None: self._dimension = int(dimension) if dimension is not None else None def close(self) -> None: - self._db.close() - # A client made from a URL does not own its connection pool, so closing it - # leaves the sockets open; they are closed here. - self._db.connection.connection_pool.disconnect() + """Nothing to release: the connection stays open for the next store on this server.""" @contextmanager def transaction(self) -> Iterator[None]: diff --git a/tests/store/test_falkordb.py b/tests/store/test_falkordb.py index 2e89572..85df6ae 100644 --- a/tests/store/test_falkordb.py +++ b/tests/store/test_falkordb.py @@ -82,6 +82,18 @@ def __init__(self, *settings: Any) -> None: assert not (tmp_path / DATABASE_NAME).exists() +def test_the_password_is_not_shown(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("CHATLORE_STORE", "falkordb") + monkeypatch.setenv("CHATLORE_FALKORDB_URL", "redis://me:p%40ss@db.example:6380") + monkeypatch.setenv("CHATLORE_FALKORDB_GRAPH", "mine") + + shown = describe_store(tmp_path) + + assert shown == "FalkorDB graph 'mine' at redis://me:***@db.example:6380" + monkeypatch.setenv("CHATLORE_FALKORDB_URL", "redis://:secret@db.example") + assert describe_store(tmp_path).endswith("at redis://:***@db.example") + + def test_an_unknown_store_is_refused(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: monkeypatch.setenv("CHATLORE_STORE", "neo4j") From 6f6a072a751cfed356cc8817f3621dfae3457f2c Mon Sep 17 00:00:00 2001 From: Cl0ver Date: Mon, 28 Sep 2026 14:50:46 -0400 Subject: [PATCH 6/9] docs: rewrap a paragraph in the FalkorDB guide --- docs/falkordb.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/falkordb.md b/docs/falkordb.md index 927fea7..76e8c24 100644 --- a/docs/falkordb.md +++ b/docs/falkordb.md @@ -13,8 +13,8 @@ By default ChatLore keeps its graph in a SQLite file in the library folder, which needs nothing else. [FalkorDB](https://www.falkordb.com) is a graph database that runs as a server. Use it when you want to query the graph with Cypher, look at it with FalkorDB's browser, or share one graph between several -machines on a network. Every command, the web interface, and the MCP server work the same -on both. +machines on a network. Every command, the web interface, and the MCP server +work the same on both. ## Settings From 2ad2066fe8b78888b2fe9ceeedc95d6c5d32d481 Mon Sep 17 00:00:00 2001 From: Cl0ver Date: Mon, 28 Sep 2026 15:05:57 -0400 Subject: [PATCH 7/9] perf(store): fetch many nodes and neighbours in one query A search against a FalkorDB server a continent away took 10 seconds: it sent 46 queries, 40 of them one lookup per passage found by meaning, and each waited 0.16 seconds on the network. On SQLite those lookups cost nothing, so the pattern went unnoticed; the graph view did the same with about 400 queries. GraphStore gains get_nodes and neighbors_many, with a default that asks one node at a time, so any backend keeps working. FalkorDB answers each in one query, in pages. Search, the passages gathered for a question, the graph view, and the entity and topic lists use them. The same search now takes 7 queries and 1 second; the graph view takes 5 queries, and the rest of its time is the download of what it shows. --- CHANGELOG.md | 4 +++ docs/falkordb.md | 10 +++--- src/chatlore/api.py | 24 +++++++------ src/chatlore/chat.py | 12 +++---- src/chatlore/mcp_server.py | 6 ++-- src/chatlore/search.py | 14 +++++--- src/chatlore/store/base.py | 31 +++++++++++++++++ src/chatlore/store/falkordb.py | 63 ++++++++++++++++++++++++++++------ tests/store/test_contract.py | 35 +++++++++++++++++++ 9 files changed, 160 insertions(+), 39 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 2ff322a..530f416 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,6 +14,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 work on it as on SQLite, and it passes the same contract tests. The client is an optional dependency: `pip install 'chatlore[falkordb]'`. Guide in `docs/falkordb.md`. - `chatlore doctor` shows which store is in use, and `chatlore index --rebuild` rebuilds either. +- `GraphStore.get_nodes` and `GraphStore.neighbors_many` fetch many nodes, or the neighbours of + many nodes, at once. Search, the passages gathered for a question, the graph view, and entity + and topic lists use them, so a store on a server answers in a few queries instead of one per + item: a search went from 46 queries to 7. - CI runs the contract tests on both stores, and every test with FalkorDB as the store. ## [0.1.0] - 2026-09-24 diff --git a/docs/falkordb.md b/docs/falkordb.md index 76e8c24..1cd1fab 100644 --- a/docs/falkordb.md +++ b/docs/falkordb.md @@ -72,10 +72,12 @@ changes made to it directly are lost on the next `chatlore index --rebuild`. matches can come back in a different order. - The server has to be running. Commands fail with a connection error when it is not. -- Every page and command asks the server many small questions, so keep it - close: on the same machine or network. Against a FalkorDB Cloud server a - continent away, 0.16 seconds per query, loading the demo took a minute and a - half, a search 8 to 10 seconds, and the graph view almost two minutes. +- Keep the server close: on the same machine or network. Reading pages fetch + what they need in a few queries, but each one waits for the network. Against + a free FalkorDB Cloud server a continent away, with 0.16 seconds per round + trip and about 20 KB per second, a search took about a second, loading the + demo a minute and a half, and the graph view, which downloads about 500 KB, + 20 seconds. - Each conversation takes a few round trips to the server, so importing is slower. On a made-up library of 12,000 notes, importing took 31 seconds against 8 with SQLite, chunking 26 seconds against 145, and a search the same diff --git a/src/chatlore/api.py b/src/chatlore/api.py index d1a9d50..0fd3d0e 100644 --- a/src/chatlore/api.py +++ b/src/chatlore/api.py @@ -224,7 +224,8 @@ def entities( with store() as graph: if q: hits = graph.search_text(q, limit=limit, labels=[Label.ENTITY]) - nodes = [node for hit in hits if (node := graph.get_node(hit.node_id))] + found = graph.get_nodes(hit.node_id for hit in hits) + nodes = [found[hit.node_id] for hit in hits if hit.node_id in found] else: nodes = sorted( graph.find_nodes(Label.ENTITY), key=lambda node: -int(node.props["mentions"]) @@ -277,7 +278,8 @@ def topics( with store() as graph: if q: hits = graph.search_text(q, limit=limit, labels=[Label.TOPIC]) - nodes = [node for hit in hits if (node := graph.get_node(hit.node_id))] + found = graph.get_nodes(hit.node_id for hit in hits) + nodes = [found[hit.node_id] for hit in hits if hit.node_id in found] else: nodes = sorted( graph.find_nodes(Label.TOPIC), key=lambda node: -int(node.props["size"]) @@ -336,24 +338,24 @@ def graph( } chosen = dict(list(chosen.items())[:limit]) else: - linked = [ - node - for node in graph.find_nodes(Label.ENTITY) - if graph.neighbors(node.id, [EdgeType.RELATED_TO], "both", limit=1) - ] + every = graph.find_nodes(Label.ENTITY) + related = graph.neighbors_many( + (node.id for node in every), [EdgeType.RELATED_TO], "both" + ) + linked = [node for node in every if related[node.id]] linked.sort(key=lambda node: -int(node.props["mentions"])) chosen = {node.id: node for node in linked[:limit]} nodes, edges, topics = [], [], {} + in_topics = graph.neighbors_many(chosen, [EdgeType.IN_TOPIC]) + links = graph.neighbors_many(chosen, [EdgeType.RELATED_TO], "out") for node in chosen.values(): - in_topic = graph.neighbors(node.id, [EdgeType.IN_TOPIC]) + in_topic = in_topics[node.id] topic_node = in_topic[0][1] if in_topic else None if topic_node is not None: topics[topic_node.id] = topic_node.props.get("title") nodes.append({**_entity(node), "topic": topic_node.id if topic_node else None}) - for edge, other in graph.neighbors( - node.id, [EdgeType.RELATED_TO], "out", limit=1_000_000 - ): + for edge, other in links[node.id]: if other.id in chosen: edges.append( { diff --git a/src/chatlore/chat.py b/src/chatlore/chat.py index a4ab3cb..dd791a3 100644 --- a/src/chatlore/chat.py +++ b/src/chatlore/chat.py @@ -170,11 +170,9 @@ def _chunks_of( the ones about the question help. Without it, its most recent chunks. """ lists: list[list[str]] = [] + mentions = store.neighbors_many((entity.id for entity in entities), [EdgeType.MENTIONS], "in") for entity in entities: - chunks = [ - chunk - for _, chunk in store.neighbors(entity.id, [EdgeType.MENTIONS], "in", limit=1_000_000) - ] + chunks = [chunk for _, chunk in mentions[entity.id]] if nearness is not None: ordered_chunks = sorted( (chunk.id for chunk in chunks if chunk.id in nearness), key=nearness.__getitem__ @@ -203,8 +201,9 @@ def _notes(store: GraphStore, entities: Sequence[Node]) -> list[Note]: ] counts: dict[str, int] = {} topics: dict[str, Node] = {} + in_topics = store.neighbors_many((node.id for node in entities), [EdgeType.IN_TOPIC]) for node in entities: - for _, topic in store.neighbors(node.id, [EdgeType.IN_TOPIC]): + for _, topic in in_topics[node.id]: counts[topic.id] = counts.get(topic.id, 0) + 1 topics[topic.id] = topic for topic_id in sorted(counts, key=lambda key: -counts[key])[:MAX_TOPICS]: @@ -236,8 +235,9 @@ def retrieve( best = sorted(scores, key=lambda chunk_id: -scores[chunk_id])[:limit] sources: list[Source] = [] + chunks = store.get_nodes(best) for chunk_id in best: - node = store.get_node(chunk_id) + node = chunks.get(chunk_id) if node is None: continue props = node.props diff --git a/src/chatlore/mcp_server.py b/src/chatlore/mcp_server.py index ac48dbd..0f8a11f 100644 --- a/src/chatlore/mcp_server.py +++ b/src/chatlore/mcp_server.py @@ -224,7 +224,8 @@ def topics( with open_store(library) as graph: if words: hits = graph.search_text(words, limit=limit, labels=[Label.TOPIC]) - nodes = [node for hit in hits if (node := graph.get_node(hit.node_id))] + found = graph.get_nodes(hit.node_id for hit in hits) + nodes = [found[hit.node_id] for hit in hits if hit.node_id in found] else: nodes = sorted( graph.find_nodes(Label.TOPIC), key=lambda node: -int(node.props["size"]) @@ -318,7 +319,8 @@ def _find_entity(graph: GraphStore, name: str) -> Node | None: return node exact = entity_key(name) hits = graph.search_text(name, limit=20, labels=[Label.ENTITY]) - nodes = [node for hit in hits if (node := graph.get_node(hit.node_id)) is not None] + found = graph.get_nodes(hit.node_id for hit in hits) + nodes = [found[hit.node_id] for hit in hits if hit.node_id in found] for node in nodes: if entity_key(str(node.props["name"])) == exact: return node diff --git a/src/chatlore/search.py b/src/chatlore/search.py index 6643901..b70c001 100644 --- a/src/chatlore/search.py +++ b/src/chatlore/search.py @@ -73,10 +73,12 @@ def hybrid_search( wanted = set(sources or []) seen: set[str] = set() # Source filtering happens after the nearest-neighbour query, so ask for extra. - for vector_hit in store.search_vector( + vector_hits = store.search_vector( embedding, limit=candidates * (4 if wanted else 1), labels=[Label.CHUNK] - ): - node = store.get_node(vector_hit.node_id) + ) + chunk_nodes = store.get_nodes(hit.node_id for hit in vector_hits) + for vector_hit in vector_hits: + node = chunk_nodes.get(vector_hit.node_id) if node is None or (wanted and node.props.get("source") not in wanted): continue message_id = str(node.props.get("message_id") or node.id) @@ -112,8 +114,10 @@ def hybrid_search( break mentioned: set[str] = set() - for entity_hit in store.search_text(query, limit=limit, labels=[Label.ENTITY]): - chunks = store.neighbors(entity_hit.node_id, [EdgeType.MENTIONS], "in", limit=candidates) + entity_hits = store.search_text(query, limit=limit, labels=[Label.ENTITY]) + mentions = store.neighbors_many((hit.node_id for hit in entity_hits), [EdgeType.MENTIONS], "in") + for entity_hit in entity_hits: + chunks = mentions[entity_hit.node_id][:candidates] for _, chunk in sorted(chunks, key=lambda pair: pair[1].id): if wanted and chunk.props.get("source") not in wanted: continue diff --git a/src/chatlore/store/base.py b/src/chatlore/store/base.py index 9384613..1fc5ccd 100644 --- a/src/chatlore/store/base.py +++ b/src/chatlore/store/base.py @@ -48,6 +48,8 @@ class EdgeType(StrEnum): Direction = Literal["out", "in", "both"] +_ALL = 1_000_000_000 + @dataclass(frozen=True, slots=True) class Node: @@ -154,6 +156,35 @@ def neighbors( ) -> list[tuple[Edge, Node]]: """Return edges touching ``node_id`` with the node at the other end.""" + def get_nodes(self, node_ids: Iterable[str]) -> dict[str, Node]: + """Return the nodes among ``node_ids`` that exist, by id. + + One lookup per node here; a backend that pays for every round trip, such + as a database server, fetches them together. + """ + found: dict[str, Node] = {} + for node_id in node_ids: + node = self.get_node(node_id) + if node is not None: + found[node_id] = node + return found + + def neighbors_many( + self, + node_ids: Iterable[str], + edge_types: Sequence[str] | None = None, + direction: Direction = "out", + ) -> dict[str, list[tuple[Edge, Node]]]: + """``neighbors`` of each node, all of them and in the same order, by id. + + One query per node here; a backend that pays for every round trip + fetches them together. + """ + return { + node_id: self.neighbors(node_id, edge_types, direction, limit=_ALL) + for node_id in dict.fromkeys(node_ids) + } + @abstractmethod def find_nodes( self, label: str, where: Mapping[str, Any] | None = None, limit: int = 1_000_000 diff --git a/src/chatlore/store/falkordb.py b/src/chatlore/store/falkordb.py index e2c3b36..a001121 100644 --- a/src/chatlore/store/falkordb.py +++ b/src/chatlore/store/falkordb.py @@ -59,6 +59,7 @@ _SNIPPET_TOKENS = 16 _UNSEARCHABLE = frozenset({"text"}) _DIMENSION = "embedding_dimension" +_ALL = 1_000_000_000 # One client per server for the whole process. Connecting takes several round # trips, seconds to a server far away, and the web interface and MCP server open @@ -164,6 +165,16 @@ def get_node(self, node_id: str) -> Node | None: ) return _node(rows[0]) if rows else None + def get_nodes(self, node_ids: Iterable[str]) -> dict[str, Node]: + found: dict[str, Node] = {} + for page in _pages(list(dict.fromkeys(node_ids))): + rows = self._query( + "MATCH (n:Node) WHERE n._id IN $ids RETURN n._id, n._label, n._props", + {"ids": page}, + ) + found.update((row[0], _node(row)) for row in rows) + return found + def delete_nodes(self, node_ids: Iterable[str]) -> None: for page in _pages(list(node_ids)): # Edges and index entries go with the node. @@ -176,21 +187,33 @@ def neighbors( direction: Direction = "out", limit: int = 100, ) -> list[tuple[Edge, Node]]: - types = ":" + "|".join(_name(edge_type) for edge_type in edge_types) if edge_types else "" - pattern = { - "out": f"(a)-[r{types}]->(b:Node)", - "in": f"(a)<-[r{types}]-(b:Node)", - "both": f"(a)-[r{types}]-(b:Node)", - }[direction] rows = self._paged( - f"MATCH (a:Node {{_id: $id}}) MATCH {pattern} " - "RETURN type(r), startNode(r)._id, endNode(r)._id, r._props, " - "b._id, b._label, b._props " - "ORDER BY type(r), endNode(r)._id, startNode(r)._id", + f"MATCH (a:Node {{_id: $id}}) MATCH {_pattern(edge_types, direction)} " + f"RETURN {_NEIGHBOR} ORDER BY type(r), endNode(r)._id, startNode(r)._id", {"id": node_id}, limit, ) - return [(Edge(row[1], row[0], row[2], _loads(row[3])), _node(row[4:7])) for row in rows] + return [_neighbor(row) for row in rows] + + def neighbors_many( + self, + node_ids: Iterable[str], + edge_types: Sequence[str] | None = None, + direction: Direction = "out", + ) -> dict[str, list[tuple[Edge, Node]]]: + wanted = list(dict.fromkeys(node_ids)) + found: dict[str, list[tuple[Edge, Node]]] = {node_id: [] for node_id in wanted} + for page in _pages(wanted): + rows = self._paged( + f"MATCH (a:Node) WHERE a._id IN $ids MATCH {_pattern(edge_types, direction)} " + f"RETURN a._id, {_NEIGHBOR} " + "ORDER BY a._id, type(r), endNode(r)._id, startNode(r)._id", + {"ids": page}, + _ALL, + ) + for row in rows: + found[row[0]].append(_neighbor(row[1:])) + return found def find_nodes( self, label: str, where: Mapping[str, Any] | None = None, limit: int = 1_000_000 @@ -514,6 +537,24 @@ def _node(row: Sequence[Any]) -> Node: return Node(row[0], row[1], _loads(row[2])) +_NEIGHBOR = "type(r), startNode(r)._id, endNode(r)._id, r._props, b._id, b._label, b._props" +"""What a neighbour query returns: the edge, then the node at its other end.""" + + +def _neighbor(row: Sequence[Any]) -> tuple[Edge, Node]: + return Edge(row[1], row[0], row[2], _loads(row[3])), _node(row[4:7]) + + +def _pattern(edge_types: Sequence[str] | None, direction: Direction) -> str: + """The pattern from node ``a`` over ``r`` to its neighbour ``b``.""" + types = ":" + "|".join(_name(edge_type) for edge_type in edge_types) if edge_types else "" + return { + "out": f"(a)-[r{types}]->(b:Node)", + "in": f"(a)<-[r{types}]-(b:Node)", + "both": f"(a)-[r{types}]-(b:Node)", + }[direction] + + def _name(value: str) -> str: """An edge type, which FalkorDB cannot take as a parameter, checked before use.""" if not _NAME.fullmatch(value): diff --git a/tests/store/test_contract.py b/tests/store/test_contract.py index 884c812..c800e09 100644 --- a/tests/store/test_contract.py +++ b/tests/store/test_contract.py @@ -269,3 +269,38 @@ def test_meta_round_trips(store: GraphStore) -> None: store.set_meta("embedding_model", "model-a") store.set_meta("embedding_model", "model-b") assert store.get_meta("embedding_model") == "model-b" + + +def test_nodes_and_neighbors_can_be_fetched_many_at_once(store: GraphStore) -> None: + store.upsert_nodes( + [ + Node("a", Label.ENTITY), + Node("b", Label.ENTITY), + Node("c", Label.TOPIC), + Node("d", Label.ENTITY), + ] + ) + store.upsert_edges( + [ + Edge("a", EdgeType.RELATED_TO, "b", {"weight": 2}), + Edge("b", EdgeType.RELATED_TO, "d"), + Edge("a", EdgeType.IN_TOPIC, "c"), + Edge("d", EdgeType.IN_TOPIC, "c"), + ] + ) + + nodes = store.get_nodes(["d", "missing", "a", "a"]) + many = store.neighbors_many(["b", "a", "ghost"], [EdgeType.RELATED_TO], "both") + + assert nodes == {"d": Node("d", Label.ENTITY), "a": Node("a", Label.ENTITY)} + assert store.get_nodes([]) == {} + assert many == { + "b": store.neighbors("b", [EdgeType.RELATED_TO], "both"), + "a": store.neighbors("a", [EdgeType.RELATED_TO], "both"), + "ghost": [], + } + assert [n.id for _, n in many["b"]] == ["a", "d"] + assert store.neighbors_many(["c"], direction="in") == { + "c": store.neighbors("c", direction="in") + } + assert store.neighbors_many([]) == {} From 71e050dee6005080b8ff8033dcf53f2508af7f78 Mon Sep 17 00:00:00 2001 From: Cl0ver Date: Mon, 28 Sep 2026 16:42:29 -0400 Subject: [PATCH 8/9] perf(store): batch the remaining reads and writes Against a FalkorDB server a continent away, every query waits for the network, so pages still paid for many round trips and for data they did not use. Each change below returns exactly what it returned before, and a contract test holds each batch method to its one-at-a-time version on both stores. - The graph view fetched every neighbour's full record, and each entity's topic with the topic's whole report, to draw links it only needed the ends of. edges_many returns just the links; the outgoing links now come from the ones already fetched, and each topic is read once. About 500 KB became about 120 KB: 28 seconds to 2.4. - Finding the entities a question names read every entity. An entity's id is made from the same normal form of its name the question is matched in, so only the ids its words could name are fetched, and each is checked against its name again. 5.7 seconds to 0.9. - Stats counted each label in its own query; count_by_label counts them in one. - The FalkorDB store read the embedding dimension on every open; it now reads it the first time a request needs it. - Importing, indexing, and archive imports wrote each conversation in about five queries; a ConversationWriter hands them to upsert_conversations 200 at a time, and writes what is waiting even when the import stops with an error. Embeddings are written a batch at a time with set_embeddings. Loading the demo went from 96 seconds to 41. - Deleting a conversation's messages read the chunks without pages, so more than 10,000 rows would have been cut off; it pages now. --- CHANGELOG.md | 15 +++- docs/falkordb.md | 12 +-- src/chatlore/api.py | 51 +++++++----- src/chatlore/archive.py | 15 ++-- src/chatlore/chat.py | 26 ++++-- src/chatlore/cli.py | 26 ++++-- src/chatlore/pipeline.py | 3 +- src/chatlore/store/__init__.py | 39 ++++++++- src/chatlore/store/base.py | 30 +++++++ src/chatlore/store/falkordb.py | 146 +++++++++++++++++++++++---------- tests/store/test_contract.py | 68 +++++++++++++++ 11 files changed, 338 insertions(+), 93 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 530f416..7f5b3e2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,10 +14,17 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 work on it as on SQLite, and it passes the same contract tests. The client is an optional dependency: `pip install 'chatlore[falkordb]'`. Guide in `docs/falkordb.md`. - `chatlore doctor` shows which store is in use, and `chatlore index --rebuild` rebuilds either. -- `GraphStore.get_nodes` and `GraphStore.neighbors_many` fetch many nodes, or the neighbours of - many nodes, at once. Search, the passages gathered for a question, the graph view, and entity - and topic lists use them, so a store on a server answers in a few queries instead of one per - item: a search went from 46 queries to 7. +- Batch methods on `GraphStore`, each with a default that goes one item at a time: + `get_nodes`, `neighbors_many`, and `edges_many` read many nodes, neighbours, or links at once; + `count_by_label` counts every label in one call; `upsert_conversations` and `set_embeddings` + write many at once. Search, the passages gathered for a question, the graph view, stats, + entity and topic lists, importing, indexing, and embedding use them, so a store on a server + answers in a few queries instead of one per item. Against a distant FalkorDB server a search + went from 46 queries and 10.6 seconds to 7 queries and 1 second, the graph view from about 400 + queries and 105 seconds to 4 queries and 2.4 seconds, and loading the demo from 96 to 41 + seconds, with the same results. +- Finding the entities a question names looks up only the ones its words could name, by id, + instead of reading every entity. - CI runs the contract tests on both stores, and every test with FalkorDB as the store. ## [0.1.0] - 2026-09-24 diff --git a/docs/falkordb.md b/docs/falkordb.md index 1cd1fab..fc5975a 100644 --- a/docs/falkordb.md +++ b/docs/falkordb.md @@ -72,12 +72,12 @@ changes made to it directly are lost on the next `chatlore index --rebuild`. matches can come back in a different order. - The server has to be running. Commands fail with a connection error when it is not. -- Keep the server close: on the same machine or network. Reading pages fetch - what they need in a few queries, but each one waits for the network. Against - a free FalkorDB Cloud server a continent away, with 0.16 seconds per round - trip and about 20 KB per second, a search took about a second, loading the - demo a minute and a half, and the graph view, which downloads about 500 KB, - 20 seconds. +- Keep the server close: on the same machine or network. ChatLore asks for + what a page needs in a few queries and writes in batches, but each query + still waits for the network. Against a free FalkorDB Cloud server a + continent away, with 0.16 seconds per round trip and about 20 KB per second, + a search took about a second, the graph view 2.4 seconds, gathering the + passages for a question under a second, and loading the demo 41 seconds. - Each conversation takes a few round trips to the server, so importing is slower. On a made-up library of 12,000 notes, importing took 31 seconds against 8 with SQLite, chunking 26 seconds against 145, and a search the same diff --git a/src/chatlore/api.py b/src/chatlore/api.py index 0fd3d0e..69bddaf 100644 --- a/src/chatlore/api.py +++ b/src/chatlore/api.py @@ -30,7 +30,7 @@ from chatlore.llm import LLMError, make_llm from chatlore.paths import default_home from chatlore.search import hybrid_search -from chatlore.store import EdgeType, GraphStore, Label, Node, open_store +from chatlore.store import Edge, EdgeType, GraphStore, Label, Node, open_store WEB = Path(__file__).parent / "web" """The web interface's files, served at /.""" @@ -121,13 +121,14 @@ def health() -> dict[str, str]: @app.get("/stats") def stats() -> dict[str, int]: with store() as graph: + counts = graph.count_by_label() return { - "conversations": graph.count_nodes(Label.CONVERSATION), - "messages": graph.count_nodes(Label.MESSAGE), - "chunks": graph.count_nodes(Label.CHUNK), + "conversations": counts[Label.CONVERSATION], + "messages": counts[Label.MESSAGE], + "chunks": counts[Label.CHUNK], "embeddings": graph.count_embeddings(), - "entities": graph.count_nodes(Label.ENTITY), - "topics": graph.count_nodes(Label.TOPIC), + "entities": counts[Label.ENTITY], + "topics": counts[Label.TOPIC], } @app.get("/search") @@ -337,30 +338,42 @@ def graph( for _, node in sorted(members, key=lambda pair: -int(pair[1].props["mentions"])) } chosen = dict(list(chosen.items())[:limit]) - else: + links: dict[str, list[Edge]] | None = None + if not entity and not topic: every = graph.find_nodes(Label.ENTITY) - related = graph.neighbors_many( + related = graph.edges_many( (node.id for node in every), [EdgeType.RELATED_TO], "both" ) linked = [node for node in every if related[node.id]] linked.sort(key=lambda node: -int(node.props["mentions"])) chosen = {node.id: node for node in linked[:limit]} + # The outgoing links are among the ones just fetched. + links = { + node_id: [edge for edge in related[node_id] if edge.src == node_id] + for node_id in chosen + } - nodes, edges, topics = [], [], {} - in_topics = graph.neighbors_many(chosen, [EdgeType.IN_TOPIC]) - links = graph.neighbors_many(chosen, [EdgeType.RELATED_TO], "out") + nodes, edges, topic_of = [], [], {} + if links is None: + links = graph.edges_many(chosen, [EdgeType.RELATED_TO], "out") + for node_id, in_topic in graph.edges_many(chosen, [EdgeType.IN_TOPIC]).items(): + if in_topic: + topic_of[node_id] = in_topic[0].dst + topic_nodes = graph.get_nodes(dict.fromkeys(topic_of.values())) + topics = { + topic_id: topic_nodes[topic_id].props.get("title") + for topic_id in dict.fromkeys(topic_of.values()) + if topic_id in topic_nodes + } for node in chosen.values(): - in_topic = in_topics[node.id] - topic_node = in_topic[0][1] if in_topic else None - if topic_node is not None: - topics[topic_node.id] = topic_node.props.get("title") - nodes.append({**_entity(node), "topic": topic_node.id if topic_node else None}) - for edge, other in links[node.id]: - if other.id in chosen: + topic_id = topic_of.get(node.id) + nodes.append({**_entity(node), "topic": topic_id if topic_id in topics else None}) + for edge in links[node.id]: + if edge.dst in chosen: edges.append( { "source": node.id, - "target": other.id, + "target": edge.dst, "weight": edge.props.get("weight", 1), "description": (edge.props.get("descriptions") or [None])[0], } diff --git a/src/chatlore/archive.py b/src/chatlore/archive.py index 6705eed..ebdfa44 100644 --- a/src/chatlore/archive.py +++ b/src/chatlore/archive.py @@ -37,7 +37,7 @@ from chatlore.extraction import ExtractionCache from chatlore.library import AddOutcome, Library from chatlore.models import Conversation, Role -from chatlore.store import Edge, GraphStore, Label, Node, open_store +from chatlore.store import ConversationWriter, Edge, GraphStore, Label, Node, open_store FORMAT = "chatlore-archive" FORMAT_VERSION = 1 @@ -216,12 +216,13 @@ def import_archive(path: Path, home: Path, dry_run: bool = False) -> ImportRepor return report with Library(home) as library, open_store(home) as store, store.transaction(): - for conversation in read_conversations(path): - outcome = library.add(conversation) - report.outcomes[outcome.value] += 1 - report.messages += len(conversation.messages) - if outcome is not AddOutcome.UNCHANGED: - store.upsert_conversation(conversation) + with ConversationWriter(store) as writer: + for conversation in read_conversations(path): + outcome = library.add(conversation) + report.outcomes[outcome.value] += 1 + report.messages += len(conversation.messages) + if outcome is not AddOutcome.UNCHANGED: + writer.add(conversation) report.nodes, report.edges = restore_graph(store, *read_graph(path)) restore_caches(path, home) return report diff --git a/src/chatlore/chat.py b/src/chatlore/chat.py index dd791a3..e86ddca 100644 --- a/src/chatlore/chat.py +++ b/src/chatlore/chat.py @@ -18,6 +18,7 @@ from dataclasses import dataclass from chatlore.extraction import entity_key +from chatlore.ids import entity_id from chatlore.llm import ChatMessage, StreamingLLM from chatlore.store import EdgeType, GraphStore, Label, Node @@ -145,18 +146,29 @@ def mentioned_entities(store: GraphStore, question: str) -> list[Node]: Every run of up to four words is compared with the entities' names the way entities are matched to each other, so "claude code", "Claude-Code", and "claude codes" all find the entity Claude Code. + + An entity's id is made from that same form of its name, so only the entities + the question could name are fetched, not all of them, and each is checked + against its name again. """ - by_key = {entity_key(str(node.props["name"])): node for node in store.find_nodes(Label.ENTITY)} words = _WORD.findall(question) - found: dict[str, Node] = {} + keys: list[str] = [] for size in range(min(MAX_NAME_WORDS, len(words)), 0, -1): for start in range(len(words) - size + 1): key = entity_key(" ".join(words[start : start + size])) - if len(key) < 2 or key in _NOT_NAMES: - continue - node = by_key.get(key) - if node is not None and node.id not in found: - found[node.id] = node + if len(key) >= 2 and key not in _NOT_NAMES: + keys.append(key) + nodes = store.get_nodes(entity_id(key) for key in keys) + found: dict[str, Node] = {} + for key in keys: + node = nodes.get(entity_id(key)) + if ( + node is not None + and node.label == Label.ENTITY + and entity_key(str(node.props["name"])) == key + and node.id not in found + ): + found[node.id] = node return list(found.values()) diff --git a/src/chatlore/cli.py b/src/chatlore/cli.py index a35d0cb..872bfe0 100644 --- a/src/chatlore/cli.py +++ b/src/chatlore/cli.py @@ -58,6 +58,7 @@ ) from chatlore.search import hybrid_search from chatlore.store import ( + ConversationWriter, EdgeType, GraphStore, Label, @@ -212,7 +213,12 @@ def import_( try: kind = detect_source(path) if source == "auto" else get_importer(source).kind importer = get_importer(kind.value) - with Library(default_home()) as library, _store(dry_run) as store, _writes(store): + with ( + Library(default_home()) as library, + _store(dry_run) as store, + _writes(store), + _writer(store) as writer, + ): for conversation in importer.parse(path, issues.append): messages += len(conversation.messages) if dry_run: @@ -220,8 +226,8 @@ def import_( continue outcome = library.add(conversation) outcomes[outcome.value] += 1 - if outcome is not AddOutcome.UNCHANGED and store is not None: - store.upsert_conversation(conversation) + if outcome is not AddOutcome.UNCHANGED and writer is not None: + writer.add(conversation) except ImporterError as error: console.print(f"[red]Import failed:[/red] {error}") raise typer.Exit(code=1) from error @@ -364,9 +370,9 @@ def index( if rebuild: library.rebuild_index() with open_store(home) as store: - with store.transaction(): + with store.transaction(), ConversationWriter(store) as writer: for conversation in library: - store.upsert_conversation(conversation) + writer.add(conversation) count += 1 total = store.count_nodes("Conversation") console.print(f"Indexed {count} conversations. Database holds {total}.") @@ -998,6 +1004,16 @@ def _writes(store: GraphStore | None) -> Iterator[None]: yield +@contextmanager +def _writer(store: GraphStore | None) -> Iterator[ConversationWriter | None]: + """Batched conversation writes; nothing without a store.""" + if store is None: + yield None + return + with ConversationWriter(store) as writer: + yield writer + + @contextmanager def _store(dry_run: bool) -> Iterator[GraphStore | None]: if dry_run: diff --git a/src/chatlore/pipeline.py b/src/chatlore/pipeline.py index 44c520a..f22f993 100644 --- a/src/chatlore/pipeline.py +++ b/src/chatlore/pipeline.py @@ -163,8 +163,7 @@ def sync_embeddings( known.update(fresh) with store.transaction(): - for node in batch: - store.set_embedding(node.id, known[hashes[node.id]]) + store.set_embeddings({node.id: known[hashes[node.id]] for node in batch}) store.set_meta(EMBEDDING_MODEL_KEY, embedder.name) report.embedded += len(missing) diff --git a/src/chatlore/store/__init__.py b/src/chatlore/store/__init__.py index 66a9aef..d3f4bd7 100644 --- a/src/chatlore/store/__init__.py +++ b/src/chatlore/store/__init__.py @@ -11,9 +11,11 @@ import os from pathlib import Path -from typing import TYPE_CHECKING +from types import TracebackType +from typing import TYPE_CHECKING, Self from urllib.parse import urlsplit, urlunsplit +from chatlore.models import Conversation from chatlore.store.base import ( ConversationSummary, Direction, @@ -43,6 +45,7 @@ "FALKORDB_URL_ENV", "STORE_ENV", "ConversationSummary", + "ConversationWriter", "Direction", "Edge", "EdgeType", @@ -58,6 +61,40 @@ ] +class ConversationWriter: + """Upserts conversations a batch at a time, so a store on a server gets few requests. + + Whatever is still waiting is written when the ``with`` block ends, even when + it ends with an error, so the store keeps up with the library. + """ + + def __init__(self, store: GraphStore, size: int = 200) -> None: + self._store = store + self._size = size + self._waiting: list[Conversation] = [] + + def __enter__(self) -> Self: + return self + + def __exit__( + self, + kind: type[BaseException] | None, + error: BaseException | None, + traceback: TracebackType | None, + ) -> None: + self.flush() + + def add(self, conversation: Conversation) -> None: + self._waiting.append(conversation) + if len(self._waiting) >= self._size: + self.flush() + + def flush(self) -> None: + if self._waiting: + waiting, self._waiting = self._waiting, [] + self._store.upsert_conversations(waiting) + + def open_store(home: Path) -> GraphStore: """Open the library's store, creating it if needed.""" home.mkdir(parents=True, exist_ok=True) diff --git a/src/chatlore/store/base.py b/src/chatlore/store/base.py index 1fc5ccd..bf51c30 100644 --- a/src/chatlore/store/base.py +++ b/src/chatlore/store/base.py @@ -185,6 +185,22 @@ def neighbors_many( for node_id in dict.fromkeys(node_ids) } + def edges_many( + self, + node_ids: Iterable[str], + edge_types: Sequence[str] | None = None, + direction: Direction = "out", + ) -> dict[str, list[Edge]]: + """The edges ``neighbors_many`` finds, without the nodes at their other ends. + + For a caller that needs only the links, so a backend on a server need not + send every neighbour's props along. + """ + return { + node_id: [edge for edge, _ in pairs] + for node_id, pairs in self.neighbors_many(node_ids, edge_types, direction).items() + } + @abstractmethod def find_nodes( self, label: str, where: Mapping[str, Any] | None = None, limit: int = 1_000_000 @@ -195,6 +211,10 @@ def find_nodes( def count_nodes(self, label: str | None = None) -> int: """Return how many nodes exist, optionally of one label.""" + def count_by_label(self) -> dict[str, int]: + """How many nodes of each ``Label`` exist, in one call.""" + return {label.value: self.count_nodes(label) for label in Label} + # -- conversations ------------------------------------------------------- @abstractmethod @@ -205,6 +225,11 @@ def upsert_conversation(self, conversation: Conversation) -> None: that disappeared are deleted together with their chunks. """ + def upsert_conversations(self, conversations: Iterable[Conversation]) -> None: + """``upsert_conversation`` for each, which a backend may do together.""" + for conversation in conversations: + self.upsert_conversation(conversation) + @abstractmethod def get_conversation(self, conversation_id: str) -> Conversation | None: """Rebuild a conversation from the graph.""" @@ -235,6 +260,11 @@ def search_text( def set_embedding(self, node_id: str, embedding: Sequence[float]) -> None: """Attach an embedding to a node. All embeddings share one dimension.""" + def set_embeddings(self, embeddings: Mapping[str, Sequence[float]]) -> None: + """``set_embedding`` for each node, which a backend may do together.""" + for node_id, embedding in embeddings.items(): + self.set_embedding(node_id, embedding) + @abstractmethod def nodes_without_embedding(self, label: str, limit: int = 100) -> list[Node]: """Return up to ``limit`` nodes of ``label`` that have no embedding yet.""" diff --git a/src/chatlore/store/falkordb.py b/src/chatlore/store/falkordb.py index a001121..3433961 100644 --- a/src/chatlore/store/falkordb.py +++ b/src/chatlore/store/falkordb.py @@ -60,6 +60,8 @@ _UNSEARCHABLE = frozenset({"text"}) _DIMENSION = "embedding_dimension" _ALL = 1_000_000_000 +_CONVERSATIONS_PER_WRITE = 200 +_VECTORS_PER_WRITE = 500 # One client per server for the whole process. Connecting takes several round # trips, seconds to a server far away, and the web interface and MCP server open @@ -99,8 +101,10 @@ def __init__(self, url: str, graph: str) -> None: if (url, graph) not in _prepared: self._create_indexes() _prepared.add((url, graph)) - dimension = self.get_meta(_DIMENSION) - self._dimension = int(dimension) if dimension is not None else None + # The embeddings' dimension is read when first needed, so requests that do + # not touch vectors cost one round trip less. + self._dimension_read = False + self._known_dimension: int | None = None def close(self) -> None: """Nothing to release: the connection stays open for the next store on this server.""" @@ -242,21 +246,40 @@ def count_nodes(self, label: str | None = None) -> int: ) return int(rows[0][0]) + def count_by_label(self) -> dict[str, int]: + rows = self._query("MATCH (n:Node) RETURN n._label, count(n)") + counts = {label: int(count) for label, count in rows} + return {label.value: counts.get(label.value, 0) for label in Label} + # -- conversations ------------------------------------------------------- def upsert_conversation(self, conversation: Conversation) -> None: - nodes, edges = conversation_to_graph(conversation) - keep = {message.id for message in conversation.messages} - # As in SQLite: messages still present are updated in place, so their - # chunks and embeddings survive; only messages that disappeared go. - self._delete_messages(conversation.id, keep=keep) - self._query( - f"MATCH (:Node {{_id: $id}})-[:{EdgeType.HAS_MESSAGE}]->(:Node)" - f"-[r:{EdgeType.REPLIES_TO}]->() DELETE r", - {"id": conversation.id}, - ) - self.upsert_nodes(nodes) - self.upsert_edges(edges) + self.upsert_conversations([conversation]) + + def upsert_conversations(self, conversations: Iterable[Conversation]) -> None: + # The last copy of a conversation given wins, as when upserting one by one. + latest = list({conversation.id: conversation for conversation in conversations}.values()) + for page in _pages(latest, _CONVERSATIONS_PER_WRITE): + # As in SQLite: messages still present are updated in place, so their + # chunks and embeddings survive; only messages that disappeared go. + keep = { + conversation.id: {message.id for message in conversation.messages} + for conversation in page + } + self._delete_messages(keep) + self._query( + f"MATCH (c:Node)-[:{EdgeType.HAS_MESSAGE}]->(:Node)" + f"-[r:{EdgeType.REPLIES_TO}]->() WHERE c._id IN $ids DELETE r", + {"ids": list(keep)}, + ) + nodes: list[Node] = [] + edges: list[Edge] = [] + for conversation in page: + conversation_nodes, conversation_edges = conversation_to_graph(conversation) + nodes.extend(conversation_nodes) + edges.extend(conversation_edges) + self.upsert_nodes(nodes) + self.upsert_edges(edges) def get_conversation(self, conversation_id: str) -> Conversation | None: node = self.get_node(conversation_id) @@ -296,7 +319,7 @@ def list_conversations( return summaries def delete_conversation(self, conversation_id: str) -> None: - self._delete_messages(conversation_id) + self._delete_messages({conversation_id: None}) self.delete_nodes([conversation_id]) # -- search -------------------------------------------------------------- @@ -346,24 +369,40 @@ def search_text( return hits def set_embedding(self, node_id: str, embedding: Sequence[float]) -> None: - vector = [float(value) for value in embedding] - if not vector: + self.set_embeddings({node_id: embedding}) + + def set_embeddings(self, embeddings: Mapping[str, Sequence[float]]) -> None: + vectors = { + node_id: [float(value) for value in vector] for node_id, vector in embeddings.items() + } + if not vectors: + return + # Everything is checked before anything is written. + if not all(vectors.values()): raise ValueError("embedding is empty") - if self.get_node(node_id) is None: - raise KeyError(node_id) - if self._dimension is None: + missing = set(vectors) - self._existing(set(vectors)) + if missing: + raise KeyError(sorted(missing)[0]) + sizes = sorted({len(vector) for vector in vectors.values()}) + dimension = self._dimension + if dimension is None and len(sizes) == 1: self._query( "CREATE VECTOR INDEX FOR (n:Node) ON (n._embedding) " - f"OPTIONS {{dimension: {len(vector)}, similarityFunction: 'euclidean'}}" + f"OPTIONS {{dimension: {sizes[0]}, similarityFunction: 'euclidean'}}" + ) + self.set_meta(_DIMENSION, str(sizes[0])) + self._set_dimension(sizes[0]) + elif sizes != [dimension]: + expected = dimension if dimension is not None else sizes[0] + wrong = next(size for size in sizes if size != expected) + raise ValueError(f"embedding has {wrong} dimensions, store has {expected}") + rows = [{"id": node_id, "vector": vector} for node_id, vector in vectors.items()] + for page in _pages(rows, _VECTORS_PER_WRITE): + self._query( + "UNWIND $rows AS row MATCH (n:Node {_id: row.id}) " + "SET n._embedding = vecf32(row.vector)", + {"rows": page}, ) - self.set_meta(_DIMENSION, str(len(vector))) - self._dimension = len(vector) - elif len(vector) != self._dimension: - raise ValueError(f"embedding has {len(vector)} dimensions, store has {self._dimension}") - self._query( - "MATCH (n:Node {_id: $id}) SET n._embedding = vecf32($vector)", - {"id": node_id, "vector": vector}, - ) def nodes_without_embedding(self, label: str, limit: int = 100) -> list[Node]: rows = self._paged( @@ -386,7 +425,7 @@ def clear_embeddings(self) -> None: "MATCH (m:Meta) WHERE m.key IN $keys DELETE m", {"keys": [_DIMENSION, "embedding_model"]}, ) - self._dimension = None + self._set_dimension(None) def get_meta(self, key: str) -> str | None: rows = self._query("MATCH (m:Meta {key: $key}) RETURN m.value", {"key": key}) @@ -446,19 +485,42 @@ def _existing(self, ids: set[str]) -> set[str]: found.update(row[0] for row in rows) return found - def _delete_messages(self, conversation_id: str, keep: set[str] | None = None) -> None: - """Delete a conversation's messages, except ``keep``, along with their chunks.""" - rows = self._query( - f"MATCH (:Node {{_id: $id}})-[:{EdgeType.HAS_MESSAGE}]->(m:Node) RETURN m._id", - {"id": conversation_id}, - ) - doomed = [row[0] for row in rows if not keep or row[0] not in keep] + @property + def _dimension(self) -> int | None: + """The dimension every embedding has, or None before the first.""" + if not self._dimension_read: + value = self.get_meta(_DIMENSION) + self._set_dimension(int(value) if value is not None else None) + return self._known_dimension + + def _set_dimension(self, dimension: int | None) -> None: + self._known_dimension = dimension + self._dimension_read = True + + def _delete_messages(self, keep: Mapping[str, set[str] | None]) -> None: + """Delete messages of the conversations in ``keep``, except the ids it keeps. + + The chunks of the messages deleted go with them. + """ + doomed: list[str] = [] + for page in _pages(list(keep)): + rows = self._paged( + f"MATCH (c:Node)-[:{EdgeType.HAS_MESSAGE}]->(m:Node) WHERE c._id IN $ids " + "RETURN c._id, m._id ORDER BY m._id", + {"ids": page}, + _ALL, + ) + for conversation_id, message_id in rows: + kept = keep[conversation_id] + if not kept or message_id not in kept: + doomed.append(message_id) chunks: list[str] = [] for page in _pages(doomed): - chunk_rows = self._query( + chunk_rows = self._paged( f"MATCH (m:Node)-[:{EdgeType.HAS_CHUNK}]->(c:Node) WHERE m._id IN $ids " - "RETURN c._id", + "RETURN c._id ORDER BY c._id", {"ids": page}, + _ALL, ) chunks.extend(row[0] for row in chunk_rows) self.delete_nodes([*chunks, *doomed]) @@ -562,9 +624,9 @@ def _name(value: str) -> str: return value -def _pages[T](items: list[T]) -> Iterator[list[T]]: - for start in range(0, len(items), _PAGE): - yield items[start : start + _PAGE] +def _pages[T](items: list[T], size: int = _PAGE) -> Iterator[list[T]]: + for start in range(0, len(items), size): + yield items[start : start + size] def _dumps(value: dict[str, Any]) -> str: diff --git a/tests/store/test_contract.py b/tests/store/test_contract.py index c800e09..ddada8d 100644 --- a/tests/store/test_contract.py +++ b/tests/store/test_contract.py @@ -304,3 +304,71 @@ def test_nodes_and_neighbors_can_be_fetched_many_at_once(store: GraphStore) -> N "c": store.neighbors("c", direction="in") } assert store.neighbors_many([]) == {} + + +def test_edges_many_are_the_edges_of_neighbors_many(store: GraphStore) -> None: + store.upsert_nodes([Node("a", Label.ENTITY), Node("b", Label.ENTITY), Node("c", Label.TOPIC)]) + store.upsert_edges( + [ + Edge("a", EdgeType.RELATED_TO, "b", {"weight": 3}), + Edge("b", EdgeType.IN_TOPIC, "c"), + Edge("a", EdgeType.IN_TOPIC, "c"), + ] + ) + + for direction in ("out", "in", "both"): + many = store.neighbors_many(["a", "b", "c"], direction=direction) + assert store.edges_many(["a", "b", "c"], direction=direction) == { + node_id: [edge for edge, _ in pairs] for node_id, pairs in many.items() + } + assert store.edges_many(["a"], [EdgeType.RELATED_TO]) == { + "a": [Edge("a", EdgeType.RELATED_TO, "b", {"weight": 3})] + } + + +def test_nodes_are_counted_by_label_in_one_call(store: GraphStore) -> None: + store.upsert_conversation(_conversation()) + store.upsert_nodes([Node("e", Label.ENTITY), Node("t", Label.TOPIC)]) + + counts = store.count_by_label() + + assert counts == {label.value: store.count_nodes(label) for label in Label} + assert counts[Label.MESSAGE] == 2 + assert counts[Label.FACT] == 0 + + +def test_conversations_can_be_upserted_together(store: GraphStore) -> None: + store.upsert_conversation(_conversation("conv_a", texts=("one", "two", "three"))) + store.upsert_nodes([Node("chunk_a", Label.CHUNK, {"conversation_id": "conv_a"})]) + store.upsert_edges([Edge("conv_a_m0", EdgeType.HAS_CHUNK, "chunk_a")]) + first_b = _conversation("conv_b", SourceKind.CLAUDE, ("cat", "names"), "Cats") + latest_b = _conversation("conv_b", SourceKind.CLAUDE, ("cat", "biscuit"), "Cats") + + store.upsert_conversations( + [_conversation("conv_a", texts=("one", "two changed")), first_b, latest_b] + ) + + assert store.get_conversation("conv_b") == latest_b + rebuilt = store.get_conversation("conv_a") + assert rebuilt is not None + assert [m.text for m in rebuilt.messages] == ["one", "two changed"] + assert store.count_nodes(Label.MESSAGE) == 4 + assert store.get_node("chunk_a") is not None # an unchanged message keeps its chunks + assert store.search_text("three") == [] + store.upsert_conversations([]) + + +def test_embeddings_can_be_set_together(store: GraphStore) -> None: + store.upsert_nodes([Node("x", Label.CHUNK), Node("y", Label.CHUNK), Node("z", Label.ENTITY)]) + + store.set_embeddings({"x": [1.0, 0.0], "y": [0.0, 1.0], "z": [0.9, 0.1]}) + store.set_embeddings({}) + + assert store.count_embeddings() == 3 + assert [h.node_id for h in store.search_vector([1.0, 0.0], limit=2)] == ["x", "z"] + with pytest.raises(ValueError, match="dimensions"): + store.set_embeddings({"x": [1.0, 0.0, 0.0]}) + with pytest.raises(KeyError): + store.set_embeddings({"x": [1.0, 0.0], "ghost": [0.0, 1.0]}) + with pytest.raises(ValueError, match="empty"): + store.set_embeddings({"x": []}) From 7d138b036634716ed32c190c7487c4a42bc8ed5d Mon Sep 17 00:00:00 2001 From: Cl0ver Date: Mon, 28 Sep 2026 16:45:30 -0400 Subject: [PATCH 9/9] fix(tests): keep CI green without the FalkorDB extra The plain CI jobs do not install the falkordb extra, so mypy could not find redis, which the FalkorDB store imports; it is skipped like the other untyped libraries. With FalkorDB as the store, the per-test graph was deleted using the URL a test had just pointed elsewhere; the URL is now read before the test runs. --- pyproject.toml | 2 +- tests/conftest.py | 4 +++- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 709f5b1..d6da464 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -108,7 +108,7 @@ pretty = true plugins = ["pydantic.mypy"] [[tool.mypy.overrides]] -module = ["sqlite_vec", "fastembed", "networkx", "falkordb", "falkordb.*"] +module = ["sqlite_vec", "fastembed", "networkx", "falkordb", "falkordb.*", "redis", "redis.*"] ignore_missing_imports = true [tool.pytest.ini_options] diff --git a/tests/conftest.py b/tests/conftest.py index 136d2b2..57e7eb1 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -26,9 +26,11 @@ def falkordb_graph(monkeypatch: pytest.MonkeyPatch) -> Iterator[None]: from chatlore.store.falkordb import FalkorDBStore graph = f"test_{uuid.uuid4().hex}" + # Read now: a test may point the setting elsewhere, and it is still set here after. + url = os.environ.get(FALKORDB_URL_ENV) or FALKORDB_DEFAULT_URL monkeypatch.setenv("CHATLORE_FALKORDB_GRAPH", graph) yield - store = FalkorDBStore(os.environ.get(FALKORDB_URL_ENV) or FALKORDB_DEFAULT_URL, graph) + store = FalkorDBStore(url, graph) try: store.drop() finally: