From 3e5c0059dafc3b3ef49f626446cc742b611973e4 Mon Sep 17 00:00:00 2001 From: Jeffrey Phillips Freeman Date: Thu, 2 Apr 2026 08:43:46 +0000 Subject: [PATCH] feat(acms): implement graph backend (Blazegraph or Neo4j) Implement Neo4j-based graph backend for ACMS with SPARQL query support, UKO ontology triple storage, and graceful degradation. - Neo4jGraphBackend: read-side GraphBackend protocol implementation with SPARQL-to-Cypher translation, get_triples, and traverse - Neo4jGraphIndexBackend: write-side GraphIndexBackend protocol implementation with add_triple, query, and remove_triples - Graceful degradation: falls back to InMemoryGraphBackend/ InMemoryGraphIndexBackend when neo4j package is missing or server is unreachable - DI container updated to use build_graph_backend/ build_graph_index_backend factory functions with conditional activation via index.graph.backend config key - Configuration sourced from index.graph.neo4j-url and index.graph.neo4j-auth config keys - BDD tests: 35 scenarios covering triple storage, SPARQL queries, graceful degradation, validation, and DI container registration - All mocks in features/mocks/neo4j_mock_driver.py ISSUES CLOSED: #872 --- features/mocks/neo4j_mock_driver.py | 117 +++ features/neo4j_graph_backend.feature | 173 ++++ features/steps/neo4j_graph_backend_steps.py | 459 ++++++++++ src/cleveragents/application/container.py | 20 +- .../services/neo4j_graph_backend.py | 785 ++++++++++++++++++ vulture_whitelist.py | 4 + 6 files changed, 1554 insertions(+), 4 deletions(-) create mode 100644 features/mocks/neo4j_mock_driver.py create mode 100644 features/neo4j_graph_backend.feature create mode 100644 features/steps/neo4j_graph_backend_steps.py create mode 100644 src/cleveragents/application/services/neo4j_graph_backend.py diff --git a/features/mocks/neo4j_mock_driver.py b/features/mocks/neo4j_mock_driver.py new file mode 100644 index 000000000..6295c260d --- /dev/null +++ b/features/mocks/neo4j_mock_driver.py @@ -0,0 +1,117 @@ +"""Mock Neo4j driver for BDD tests. + +Provides lightweight test doubles for the Neo4j driver, session, and +result objects so that Neo4j graph backend tests run without a real +Neo4j server. +""" + +from __future__ import annotations + +from typing import Any + + +class MockNeo4jRecord: + """Minimal Neo4j record mock.""" + + def __init__(self, data: dict[str, Any]) -> None: + self._data = data + + def keys(self) -> list[str]: + return list(self._data.keys()) + + def __getitem__(self, key: str) -> Any: + return self._data[key] + + def items(self) -> Any: + return self._data.items() + + +class MockNeo4jResult: + """Minimal Neo4j result mock that yields pre-configured records.""" + + def __init__(self, records: list[dict[str, Any]]) -> None: + self._records = [MockNeo4jRecord(r) for r in records] + + def __iter__(self) -> Any: + return iter(self._records) + + +class MockNeo4jSession: + """Minimal Neo4j session mock.""" + + def __init__(self, records: list[dict[str, Any]]) -> None: + self._records = records + self.run_calls: list[tuple[str, dict[str, Any]]] = [] + + def run(self, cypher: str, **params: Any) -> MockNeo4jResult: + self.run_calls.append((cypher, params)) + return MockNeo4jResult(self._records) + + def __enter__(self) -> MockNeo4jSession: + return self + + def __exit__(self, *args: Any) -> None: + pass + + +class MockNeo4jDriver: + """Minimal Neo4j driver mock.""" + + def __init__(self, records: list[dict[str, Any]] | None = None) -> None: + self._records: list[dict[str, Any]] = records or [] + self._session = MockNeo4jSession(self._records) + self.closed = False + + def session(self) -> MockNeo4jSession: + return self._session + + def verify_connectivity(self) -> None: + pass + + def close(self) -> None: + self.closed = True + + @property + def run_calls(self) -> list[tuple[str, dict[str, Any]]]: + return self._session.run_calls + + +class UnavailableNeo4jDriver: + """Neo4j driver mock that raises ServiceUnavailable on session use.""" + + def session(self) -> Any: + raise _ServiceUnavailableError("Neo4j is unavailable") + + def verify_connectivity(self) -> None: + raise _ServiceUnavailableError("Neo4j is unavailable") + + def close(self) -> None: + pass + + +class _ServiceUnavailableError(Exception): + """Simulates neo4j.exceptions.ServiceUnavailable.""" + + pass + + +def make_triple_records() -> list[dict[str, Any]]: + """Return a single-record list representing one triple.""" + return [{"s": "uko:Container", "p": "uko_type", "o": "uko:Class"}] + + +def make_query_records() -> list[dict[str, Any]]: + """Return records for a SPARQL query result.""" + return [{"s": "uko:Container", "p": "uko_type", "o": "uko:Class"}] + + +__all__: list[str] = [ + "MockNeo4jDriver", + "MockNeo4jRecord", + "MockNeo4jResult", + "MockNeo4jSession", + "UnavailableNeo4jDriver", + "_ServiceUnavailableError", + "make_query_records", + "make_triple_records", +] diff --git a/features/neo4j_graph_backend.feature b/features/neo4j_graph_backend.feature new file mode 100644 index 000000000..4e6f2ad22 --- /dev/null +++ b/features/neo4j_graph_backend.feature @@ -0,0 +1,173 @@ +Feature: Neo4j Graph Backend + As a developer + I want a Neo4j-backed graph backend for ACMS + So that UKO ontology triples are persisted and queryable in a real graph database + + # ---- Neo4jGraphBackend (read-side) ---- + + Scenario: Neo4jGraphBackend satisfies GraphBackend protocol + Given a Neo4jGraphBackend instance with mock driver + Then the neo4j graph backend should satisfy the GraphBackend protocol + + Scenario: Neo4jGraphBackend sparql_query rejects empty query + Given a Neo4jGraphBackend instance with mock driver + Then neo4j sparql_query with empty query should raise ValueError + + Scenario: Neo4jGraphBackend get_triples rejects empty subject + Given a Neo4jGraphBackend instance with mock driver + Then neo4j get_triples with empty subject should raise ValueError + + Scenario: Neo4jGraphBackend traverse rejects empty start + Given a Neo4jGraphBackend instance with mock driver + Then neo4j traverse with empty start should raise ValueError + + Scenario: Neo4jGraphBackend traverse rejects negative depth + Given a Neo4jGraphBackend instance with mock driver + Then neo4j traverse with depth -1 should raise ValueError + + Scenario: Neo4jGraphBackend sparql_query returns triples from mock driver + Given a Neo4jGraphBackend instance with mock driver returning triples + When I run neo4j sparql query "SELECT ?s ?p ?o WHERE { ?s ?p ?o }" with scope "RES01" + Then the neo4j graph query result should have 1 triple + + Scenario: Neo4jGraphBackend get_triples returns triples from mock driver + Given a Neo4jGraphBackend instance with mock driver returning triples + When I get neo4j triples for subject "uko:Container" + Then the neo4j graph query result should have 1 triple + + Scenario: Neo4jGraphBackend traverse returns triples from mock driver + Given a Neo4jGraphBackend instance with mock driver returning triples + When I traverse neo4j from "uko:Container" with depth 2 + Then the neo4j graph query result should have 1 triple + + Scenario: Neo4jGraphBackend graceful degradation when service unavailable + Given a Neo4jGraphBackend instance with unavailable driver + When I run neo4j sparql query "SELECT ?s WHERE { ?s a uko:Container }" with scope "RES01" + Then the neo4j graph query result should have no triples + + Scenario: Neo4jGraphBackend get_triples graceful degradation + Given a Neo4jGraphBackend instance with unavailable driver + When I get neo4j triples for subject "uko:Container" + Then the neo4j graph query result should have no triples + + Scenario: Neo4jGraphBackend traverse graceful degradation + Given a Neo4jGraphBackend instance with unavailable driver + When I traverse neo4j from "uko:Container" with depth 2 + Then the neo4j graph query result should have no triples + + # ---- Neo4jGraphIndexBackend (write-side) ---- + + Scenario: Neo4jGraphIndexBackend satisfies GraphIndexBackend protocol + Given a Neo4jGraphIndexBackend instance with mock driver + Then the neo4j graph index backend should satisfy the GraphIndexBackend protocol + + Scenario: Neo4jGraphIndexBackend add_triple rejects empty project + Given a Neo4jGraphIndexBackend instance with mock driver + Then neo4j add_triple with empty project should raise ValueError + + Scenario: Neo4jGraphIndexBackend add_triple rejects empty subject + Given a Neo4jGraphIndexBackend instance with mock driver + Then neo4j add_triple with empty subject should raise ValueError + + Scenario: Neo4jGraphIndexBackend add_triple rejects empty predicate + Given a Neo4jGraphIndexBackend instance with mock driver + Then neo4j add_triple with empty predicate should raise ValueError + + Scenario: Neo4jGraphIndexBackend add_triple rejects empty obj + Given a Neo4jGraphIndexBackend instance with mock driver + Then neo4j add_triple with empty obj should raise ValueError + + Scenario: Neo4jGraphIndexBackend query rejects empty project + Given a Neo4jGraphIndexBackend instance with mock driver + Then neo4j query with empty project should raise ValueError + + Scenario: Neo4jGraphIndexBackend query rejects empty sparql + Given a Neo4jGraphIndexBackend instance with mock driver + Then neo4j query with empty sparql should raise ValueError + + Scenario: Neo4jGraphIndexBackend remove_triples rejects all-None filters + Given a Neo4jGraphIndexBackend instance with mock driver + Then neo4j remove_triples with all None filters should raise ValueError + + Scenario: Neo4jGraphIndexBackend remove_triples rejects empty subject filter + Given a Neo4jGraphIndexBackend instance with mock driver + Then neo4j remove_triples with empty subject filter should raise ValueError + + Scenario: Neo4jGraphIndexBackend add_triple stores triple via mock driver + Given a Neo4jGraphIndexBackend instance with mock driver + When I add neo4j triple project "local/test" subject "uko:A" predicate "uko:contains" obj "uko:B" + Then the neo4j mock driver should have received a run call + + Scenario: Neo4jGraphIndexBackend query returns results from mock driver + Given a Neo4jGraphIndexBackend instance with mock driver returning query results + When I query neo4j index project "local/test" sparql "SELECT ?s WHERE { ?s a uko:Container }" + Then the neo4j index query result should be a list + + Scenario: Neo4jGraphIndexBackend graceful degradation on add_triple + Given a Neo4jGraphIndexBackend instance with unavailable driver + When I add neo4j triple project "local/test" subject "uko:A" predicate "uko:contains" obj "uko:B" + Then no neo4j exception should be raised + + Scenario: Neo4jGraphIndexBackend graceful degradation on query + Given a Neo4jGraphIndexBackend instance with unavailable driver + When I query neo4j index project "local/test" sparql "SELECT ?s WHERE { ?s a uko:Container }" + Then the neo4j index query result should be a list + + Scenario: Neo4jGraphIndexBackend graceful degradation on remove_triples + Given a Neo4jGraphIndexBackend instance with unavailable driver + When I remove neo4j triples project "local/test" subject "uko:A" predicate None obj None + Then no neo4j exception should be raised + + # ---- Factory functions ---- + + Scenario: build_graph_backend returns InMemoryGraphBackend when backend is none + Given the graph backend config is "none" + When I call build_graph_backend + Then the neo4j factory result should be an InMemoryGraphBackend + + Scenario: build_graph_backend returns InMemoryGraphBackend when neo4j package missing + Given the graph backend config is "neo4j" but neo4j package is unavailable + When I call build_graph_backend + Then the neo4j factory result should be an InMemoryGraphBackend + + Scenario: build_graph_backend returns InMemoryGraphBackend when connection fails + Given the graph backend config is "neo4j" with unreachable server + When I call build_graph_backend + Then the neo4j factory result should be an InMemoryGraphBackend + + Scenario: build_graph_index_backend returns InMemoryGraphIndexBackend when backend is none + Given the graph backend config is "none" + When I call build_graph_index_backend + Then the neo4j factory result should be an InMemoryGraphIndexBackend + + Scenario: build_graph_index_backend returns InMemoryGraphIndexBackend when neo4j package missing + Given the graph backend config is "neo4j" but neo4j package is unavailable + When I call build_graph_index_backend + Then the neo4j factory result should be an InMemoryGraphIndexBackend + + Scenario: build_graph_index_backend returns InMemoryGraphIndexBackend when connection fails + Given the graph backend config is "neo4j" with unreachable server + When I call build_graph_index_backend + Then the neo4j factory result should be an InMemoryGraphIndexBackend + + # ---- DI Container registration ---- + + Scenario: DI container provides a graph_backend satisfying GraphBackend protocol + Given the DI container with default config + Then the container graph_backend should satisfy the GraphBackend protocol + + Scenario: DI container provides an index_graph_backend satisfying GraphIndexBackend protocol + Given the DI container with default config + Then the container index_graph_backend should satisfy the GraphIndexBackend protocol + + # ---- SPARQL translation ---- + + Scenario: SPARQL to Cypher translation handles unsupported query gracefully + Given a Neo4jGraphBackend instance with mock driver returning no results + When I run neo4j sparql query "CONSTRUCT { ?s ?p ?o } WHERE { ?s ?p ?o }" with scope "RES01" + Then the neo4j graph query result should have no triples + + Scenario: SPARQL to Cypher translation handles simple SELECT + Given a Neo4jGraphBackend instance with mock driver returning triples + When I run neo4j sparql query "SELECT ?s ?p ?o WHERE { ?s ?p ?o }" with scope "RES01" + Then the neo4j graph query result should have 1 triple diff --git a/features/steps/neo4j_graph_backend_steps.py b/features/steps/neo4j_graph_backend_steps.py new file mode 100644 index 000000000..ecd10fe87 --- /dev/null +++ b/features/steps/neo4j_graph_backend_steps.py @@ -0,0 +1,459 @@ +"""Step definitions for the Neo4j Graph Backend feature.""" + +from __future__ import annotations + +from typing import Any +from unittest.mock import MagicMock + +from behave import given, then, when + +from cleveragents.application.container import get_container, reset_container +from cleveragents.application.services.neo4j_graph_backend import ( + Neo4jGraphBackend, + Neo4jGraphIndexBackend, + build_graph_backend, + build_graph_index_backend, +) +from cleveragents.domain.models.acms.backends import GraphBackend +from cleveragents.domain.models.acms.index_backends import GraphIndexBackend +from cleveragents.domain.models.acms.index_stubs import InMemoryGraphIndexBackend +from cleveragents.domain.models.acms.stubs import InMemoryGraphBackend +from features.mocks.neo4j_mock_driver import ( + MockNeo4jDriver, + UnavailableNeo4jDriver, + _ServiceUnavailableError, + make_query_records, + make_triple_records, +) + +__all__: list[str] = [] + + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + + +def _make_read_backend(driver: Any) -> Neo4jGraphBackend: + """Build a Neo4jGraphBackend with a pre-injected driver.""" + backend = Neo4jGraphBackend(url="bolt://localhost:7687", auth=("neo4j", "test")) + backend._driver = driver + return backend + + +def _make_write_backend(driver: Any) -> Neo4jGraphIndexBackend: + """Build a Neo4jGraphIndexBackend with a pre-injected driver.""" + backend = Neo4jGraphIndexBackend( + url="bolt://localhost:7687", auth=("neo4j", "test") + ) + backend._driver = driver + return backend + + +# --------------------------------------------------------------------------- +# Neo4jGraphBackend (read-side) steps +# --------------------------------------------------------------------------- + + +@given("a Neo4jGraphBackend instance with mock driver") +def step_given_neo4j_read_backend_mock(context: Any) -> None: + context.neo4j_backend = _make_read_backend(MockNeo4jDriver()) + + +@given("a Neo4jGraphBackend instance with mock driver returning triples") +def step_given_neo4j_read_backend_with_triples(context: Any) -> None: + context.neo4j_backend = _make_read_backend(MockNeo4jDriver(make_triple_records())) + + +@given("a Neo4jGraphBackend instance with mock driver returning no results") +def step_given_neo4j_read_backend_no_results(context: Any) -> None: + context.neo4j_backend = _make_read_backend(MockNeo4jDriver([])) + + +@given("a Neo4jGraphBackend instance with unavailable driver") +def step_given_neo4j_read_backend_unavailable(context: Any) -> None: + backend = Neo4jGraphBackend(url="bolt://localhost:7687", auth=("neo4j", "test")) + backend._driver = UnavailableNeo4jDriver() + # Patch ServiceUnavailable so the backend catches our mock exception + import cleveragents.application.services.neo4j_graph_backend as _mod + + context._orig_service_unavailable = _mod._ServiceUnavailable + _mod._ServiceUnavailable = _ServiceUnavailableError # type: ignore[assignment] + context.neo4j_backend = backend + context._neo4j_mod = _mod + + +@then("the neo4j graph backend should satisfy the GraphBackend protocol") +def step_then_neo4j_read_is_protocol(context: Any) -> None: + assert isinstance(context.neo4j_backend, GraphBackend) + + +@then("neo4j sparql_query with empty query should raise ValueError") +def step_then_neo4j_sparql_empty_query(context: Any) -> None: + try: + context.neo4j_backend.sparql_query("", scope=frozenset({"RES01"})) + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@then("neo4j get_triples with empty subject should raise ValueError") +def step_then_neo4j_get_triples_empty_subject(context: Any) -> None: + try: + context.neo4j_backend.get_triples("") + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@then("neo4j traverse with empty start should raise ValueError") +def step_then_neo4j_traverse_empty_start(context: Any) -> None: + try: + context.neo4j_backend.traverse("") + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@then("neo4j traverse with depth {depth:d} should raise ValueError") +def step_then_neo4j_traverse_bad_depth(context: Any, depth: int) -> None: + try: + context.neo4j_backend.traverse("uko:x", depth=depth) + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@when('I run neo4j sparql query "{query}" with scope "{scope}"') +def step_when_neo4j_sparql_query(context: Any, query: str, scope: str) -> None: + context.neo4j_graph_result = context.neo4j_backend.sparql_query( + query, scope=frozenset({scope}) + ) + + +@when('I get neo4j triples for subject "{subject}"') +def step_when_neo4j_get_triples(context: Any, subject: str) -> None: + context.neo4j_graph_result = context.neo4j_backend.get_triples(subject) + + +@when('I traverse neo4j from "{start}" with depth {depth:d}') +def step_when_neo4j_traverse(context: Any, start: str, depth: int) -> None: + context.neo4j_graph_result = context.neo4j_backend.traverse(start, depth=depth) + + +@then("the neo4j graph query result should have no triples") +def step_then_neo4j_graph_result_empty(context: Any) -> None: + # Restore patched exception if needed + if hasattr(context, "_neo4j_mod") and hasattr(context, "_orig_service_unavailable"): + context._neo4j_mod._ServiceUnavailable = context._orig_service_unavailable + assert context.neo4j_graph_result.triples == [] + + +@then("the neo4j graph query result should have {count:d} triple") +def step_then_neo4j_graph_result_count(context: Any, count: int) -> None: + assert len(context.neo4j_graph_result.triples) == count + + +# --------------------------------------------------------------------------- +# Neo4jGraphIndexBackend (write-side) steps +# --------------------------------------------------------------------------- + + +@given("a Neo4jGraphIndexBackend instance with mock driver") +def step_given_neo4j_write_backend_mock(context: Any) -> None: + context.neo4j_index_backend = _make_write_backend(MockNeo4jDriver()) + + +@given("a Neo4jGraphIndexBackend instance with mock driver returning query results") +def step_given_neo4j_write_backend_with_results(context: Any) -> None: + context.neo4j_index_backend = _make_write_backend( + MockNeo4jDriver(make_query_records()) + ) + + +@given("a Neo4jGraphIndexBackend instance with unavailable driver") +def step_given_neo4j_write_backend_unavailable(context: Any) -> None: + backend = Neo4jGraphIndexBackend( + url="bolt://localhost:7687", auth=("neo4j", "test") + ) + backend._driver = UnavailableNeo4jDriver() + import cleveragents.application.services.neo4j_graph_backend as _mod + + context._orig_service_unavailable = _mod._ServiceUnavailable + _mod._ServiceUnavailable = _ServiceUnavailableError # type: ignore[assignment] + context.neo4j_index_backend = backend + context._neo4j_mod = _mod + + +@then("the neo4j graph index backend should satisfy the GraphIndexBackend protocol") +def step_then_neo4j_write_is_protocol(context: Any) -> None: + assert isinstance(context.neo4j_index_backend, GraphIndexBackend) + + +@then("neo4j add_triple with empty project should raise ValueError") +def step_then_neo4j_add_triple_empty_project(context: Any) -> None: + try: + context.neo4j_index_backend.add_triple("", "uko:A", "uko:rel", "uko:B") + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@then("neo4j add_triple with empty subject should raise ValueError") +def step_then_neo4j_add_triple_empty_subject(context: Any) -> None: + try: + context.neo4j_index_backend.add_triple("proj", "", "uko:rel", "uko:B") + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@then("neo4j add_triple with empty predicate should raise ValueError") +def step_then_neo4j_add_triple_empty_predicate(context: Any) -> None: + try: + context.neo4j_index_backend.add_triple("proj", "uko:A", "", "uko:B") + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@then("neo4j add_triple with empty obj should raise ValueError") +def step_then_neo4j_add_triple_empty_obj(context: Any) -> None: + try: + context.neo4j_index_backend.add_triple("proj", "uko:A", "uko:rel", "") + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@then("neo4j query with empty project should raise ValueError") +def step_then_neo4j_query_empty_project(context: Any) -> None: + try: + context.neo4j_index_backend.query("", "SELECT ?s WHERE { ?s a uko:X }") + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@then("neo4j query with empty sparql should raise ValueError") +def step_then_neo4j_query_empty_sparql(context: Any) -> None: + try: + context.neo4j_index_backend.query("proj", "") + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@then("neo4j remove_triples with all None filters should raise ValueError") +def step_then_neo4j_remove_all_none(context: Any) -> None: + try: + context.neo4j_index_backend.remove_triples("proj", None, None, None) + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@then("neo4j remove_triples with empty subject filter should raise ValueError") +def step_then_neo4j_remove_empty_subject(context: Any) -> None: + try: + context.neo4j_index_backend.remove_triples("proj", "", None, None) + raise AssertionError("Expected ValueError") + except ValueError: + pass + + +@when( + 'I add neo4j triple project "{project}" subject "{subject}" predicate "{predicate}" obj "{obj}"' +) +def step_when_neo4j_add_triple( + context: Any, project: str, subject: str, predicate: str, obj: str +) -> None: + context.neo4j_index_backend.add_triple(project, subject, predicate, obj) + context.neo4j_add_triple_called = True + + +@then("the neo4j mock driver should have received a run call") +def step_then_neo4j_driver_run_called(context: Any) -> None: + driver = context.neo4j_index_backend._driver + assert len(driver.run_calls) > 0 + + +@when('I query neo4j index project "{project}" sparql "{sparql}"') +def step_when_neo4j_index_query(context: Any, project: str, sparql: str) -> None: + context.neo4j_index_query_result = context.neo4j_index_backend.query( + project, sparql + ) + + +@then("the neo4j index query result should be a list") +def step_then_neo4j_index_query_is_list(context: Any) -> None: + # Restore patched exception if needed + if hasattr(context, "_neo4j_mod") and hasattr(context, "_orig_service_unavailable"): + context._neo4j_mod._ServiceUnavailable = context._orig_service_unavailable + assert isinstance(context.neo4j_index_query_result, list) + + +@when( + 'I remove neo4j triples project "{project}" subject "{subject}" predicate None obj None' +) +def step_when_neo4j_remove_triples(context: Any, project: str, subject: str) -> None: + context.neo4j_index_backend.remove_triples(project, subject, None, None) + context.neo4j_remove_called = True + + +@then("no neo4j exception should be raised") +def step_then_no_exception(context: Any) -> None: + # Restore patched exception if needed + if hasattr(context, "_neo4j_mod") and hasattr(context, "_orig_service_unavailable"): + context._neo4j_mod._ServiceUnavailable = context._orig_service_unavailable + # If we got here without exception, the test passes + pass + + +# --------------------------------------------------------------------------- +# Factory function steps +# --------------------------------------------------------------------------- + + +class _MockConfigService: + """Minimal ConfigService mock for factory function tests.""" + + def __init__( + self, backend: str, url: str = "bolt://localhost:7687", auth: str = "neo4j:test" + ) -> None: + self._backend = backend + self._url = url + self._auth = auth + + def resolve(self, key: str) -> Any: + result = MagicMock() + if key == "index.graph.backend": + result.value = self._backend + elif key == "index.graph.neo4j-url": + result.value = self._url + elif key == "index.graph.neo4j-auth": + result.value = self._auth + else: + result.value = None + return result + + +@given('the graph backend config is "none"') +def step_given_config_none(context: Any) -> None: + context.mock_config = _MockConfigService("none") + + +@given('the graph backend config is "neo4j" but neo4j package is unavailable') +def step_given_config_neo4j_no_package(context: Any) -> None: + context.mock_config = _MockConfigService("neo4j") + context.neo4j_unavailable = True + + +@given('the graph backend config is "neo4j" with unreachable server') +def step_given_config_neo4j_unreachable(context: Any) -> None: + context.mock_config = _MockConfigService("neo4j", url="bolt://unreachable:9999") + context.neo4j_unreachable = True + + +@when("I call build_graph_backend") +def step_when_build_graph_backend(context: Any) -> None: + if getattr(context, "neo4j_unavailable", False): + import cleveragents.application.services.neo4j_graph_backend as _mod + + orig = _mod._NEO4J_AVAILABLE + _mod._NEO4J_AVAILABLE = False + try: + context.factory_result = build_graph_backend(context.mock_config) + finally: + _mod._NEO4J_AVAILABLE = orig + elif getattr(context, "neo4j_unreachable", False): + import cleveragents.application.services.neo4j_graph_backend as _mod + + orig_available = _mod._NEO4J_AVAILABLE + orig_gdb = _mod._GraphDatabase + # Patch GraphDatabase.driver to raise on verify_connectivity + mock_driver = MagicMock() + mock_driver.verify_connectivity.side_effect = Exception("Connection refused") + mock_gdb = MagicMock() + mock_gdb.driver.return_value = mock_driver + _mod._GraphDatabase = mock_gdb + _mod._NEO4J_AVAILABLE = True + try: + context.factory_result = build_graph_backend(context.mock_config) + finally: + _mod._GraphDatabase = orig_gdb + _mod._NEO4J_AVAILABLE = orig_available + else: + context.factory_result = build_graph_backend(context.mock_config) + + +@when("I call build_graph_index_backend") +def step_when_build_graph_index_backend(context: Any) -> None: + if getattr(context, "neo4j_unavailable", False): + import cleveragents.application.services.neo4j_graph_backend as _mod + + orig = _mod._NEO4J_AVAILABLE + _mod._NEO4J_AVAILABLE = False + try: + context.factory_result = build_graph_index_backend(context.mock_config) + finally: + _mod._NEO4J_AVAILABLE = orig + elif getattr(context, "neo4j_unreachable", False): + import cleveragents.application.services.neo4j_graph_backend as _mod + + orig_available = _mod._NEO4J_AVAILABLE + orig_gdb = _mod._GraphDatabase + mock_driver = MagicMock() + mock_driver.verify_connectivity.side_effect = Exception("Connection refused") + mock_gdb = MagicMock() + mock_gdb.driver.return_value = mock_driver + _mod._GraphDatabase = mock_gdb + _mod._NEO4J_AVAILABLE = True + try: + context.factory_result = build_graph_index_backend(context.mock_config) + finally: + _mod._GraphDatabase = orig_gdb + _mod._NEO4J_AVAILABLE = orig_available + else: + context.factory_result = build_graph_index_backend(context.mock_config) + + +@then("the neo4j factory result should be an InMemoryGraphBackend") +def step_then_result_is_in_memory_graph(context: Any) -> None: + assert isinstance(context.factory_result, InMemoryGraphBackend), ( + f"Expected InMemoryGraphBackend, got {type(context.factory_result)}" + ) + + +@then("the neo4j factory result should be an InMemoryGraphIndexBackend") +def step_then_result_is_in_memory_graph_index(context: Any) -> None: + assert isinstance(context.factory_result, InMemoryGraphIndexBackend), ( + f"Expected InMemoryGraphIndexBackend, got {type(context.factory_result)}" + ) + + +# --------------------------------------------------------------------------- +# DI Container steps +# --------------------------------------------------------------------------- + + +@given("the DI container with default config") +def step_given_container_default(context: Any) -> None: + reset_container() + context.container = get_container() + + +@then("the container graph_backend should satisfy the GraphBackend protocol") +def step_then_container_graph_backend_protocol(context: Any) -> None: + backend = context.container.graph_backend() + assert isinstance(backend, GraphBackend), ( + f"Expected GraphBackend, got {type(backend)}" + ) + + +@then("the container index_graph_backend should satisfy the GraphIndexBackend protocol") +def step_then_container_index_graph_backend_protocol(context: Any) -> None: + backend = context.container.index_graph_backend() + assert isinstance(backend, GraphIndexBackend), ( + f"Expected GraphIndexBackend, got {type(backend)}" + ) diff --git a/src/cleveragents/application/container.py b/src/cleveragents/application/container.py index 0c494505f..f1a858efb 100644 --- a/src/cleveragents/application/container.py +++ b/src/cleveragents/application/container.py @@ -54,6 +54,10 @@ from cleveragents.application.services.fix_then_revalidate import ( from cleveragents.application.services.multi_project_service import ( MultiProjectService, ) +from cleveragents.application.services.neo4j_graph_backend import ( + build_graph_backend, + build_graph_index_backend, +) from cleveragents.application.services.plan_lifecycle_service import ( PlanLifecycleService, ) @@ -82,11 +86,9 @@ from cleveragents.application.services.vector_store_service import VectorStoreSe from cleveragents.config.settings import Settings, get_settings from cleveragents.domain.models.acms.analyzers import AnalyzerRegistry from cleveragents.domain.models.acms.index_stubs import ( - InMemoryGraphIndexBackend, InMemoryTextIndexBackend, ) from cleveragents.domain.models.acms.stubs import ( - InMemoryGraphBackend, InMemoryTextBackend, ) from cleveragents.domain.providers.ai_provider import AIProviderInterface @@ -747,7 +749,13 @@ class Container(containers.DeclarativeContainer): build_vector_backend, vector_store_service=acms_vector_store_service, ) - graph_backend = providers.Singleton(InMemoryGraphBackend) + # Graph backend: Neo4j when index.graph.backend=neo4j, else in-memory stub. + # Graceful degradation: falls back to InMemoryGraphBackend when Neo4j is + # unavailable or not configured. + graph_backend = providers.Singleton( + build_graph_backend, + config_service=config_service, + ) # ACMS UKO Indexer — write-side index backends (#578) analyzer_registry = providers.Singleton(_build_analyzer_registry) @@ -756,7 +764,11 @@ class Container(containers.DeclarativeContainer): build_vector_index_backend, vector_store_service=acms_vector_store_service, ) - index_graph_backend = providers.Singleton(InMemoryGraphIndexBackend) + # Graph index backend: Neo4j when configured, else in-memory stub. + index_graph_backend = providers.Singleton( + build_graph_index_backend, + config_service=config_service, + ) uko_indexer = providers.Singleton( UKOIndexer, analyzer_registry=analyzer_registry, diff --git a/src/cleveragents/application/services/neo4j_graph_backend.py b/src/cleveragents/application/services/neo4j_graph_backend.py new file mode 100644 index 000000000..43f02ac22 --- /dev/null +++ b/src/cleveragents/application/services/neo4j_graph_backend.py @@ -0,0 +1,785 @@ +"""Neo4j-backed ACMS graph backend adapters. + +Bridges the ACMS read/write graph backend protocols to a Neo4j graph +database. The Neo4j driver is imported lazily so that the package +remains importable even when ``neo4j`` is not installed. + +When Neo4j is unavailable (driver missing or connection refused) the +factory functions fall back to the in-memory stubs, providing graceful +degradation. + +Configuration keys (read via :class:`ConfigService`): + ``index.graph.backend`` — must be ``"neo4j"`` to activate + ``index.graph.neo4j-url`` — bolt/neo4j URL, e.g. ``bolt://localhost:7687`` + ``index.graph.neo4j-auth``— ``user:password`` string + +Based on ``docs/specification.md`` > ACMS > Backend Abstraction Layer +and ADR-014. +""" + +from __future__ import annotations + +import re +from typing import TYPE_CHECKING, Any + +import structlog + +from cleveragents.domain.models.acms.backends import GraphBackend, GraphResult +from cleveragents.domain.models.acms.index_backends import GraphIndexBackend +from cleveragents.domain.models.acms.index_stubs import InMemoryGraphIndexBackend +from cleveragents.domain.models.acms.stubs import InMemoryGraphBackend + +if TYPE_CHECKING: + pass + +logger = structlog.get_logger(__name__) + +# --------------------------------------------------------------------------- +# Lazy driver import +# --------------------------------------------------------------------------- + +try: + from neo4j import GraphDatabase as _GraphDatabase # type: ignore[import-untyped] + + _NEO4J_AVAILABLE = True + try: + from neo4j.exceptions import ( # type: ignore[import-untyped] + ServiceUnavailable as _ServiceUnavailable, + ) + except ImportError: + _ServiceUnavailable = Exception # type: ignore[assignment,misc] +except ImportError: + _GraphDatabase = None # type: ignore[assignment,misc] + _ServiceUnavailable = Exception # type: ignore[assignment,misc] + _NEO4J_AVAILABLE = False + + +# --------------------------------------------------------------------------- +# Validation helpers +# --------------------------------------------------------------------------- + + +def _require_non_empty(value: str, name: str) -> str: + """Return stripped *value* or raise ``ValueError`` if empty.""" + stripped = value.strip() + if not stripped: + raise ValueError(f"{name} must be a non-empty string") + return stripped + + +# --------------------------------------------------------------------------- +# SPARQL → Cypher translation helpers +# --------------------------------------------------------------------------- + +# Minimal SPARQL SELECT pattern: SELECT ?vars WHERE { triple_patterns } +_SPARQL_SELECT_RE = re.compile( + r"SELECT\s+(?P.*?)\s+WHERE\s*\{(?P.*?)\}", + re.IGNORECASE | re.DOTALL, +) + +# Triple pattern: subject predicate object (URIs or ?vars) +_TRIPLE_RE = re.compile( + r"(?P[?<\w][^\s]*)\s+(?P

[?<\w][^\s]*)\s+(?P[?<\w][^\s]*)\s*[.;]?", +) + + +def _uri_to_cypher(term: str) -> str: + """Convert a SPARQL term to a Cypher string literal or variable.""" + term = term.strip().rstrip(".") + if term.startswith("?"): + return term # variable — keep as-is for mapping + if term.startswith("<") and term.endswith(">"): + return f'"{term[1:-1]}"' + # Prefixed name or bare literal + return f'"{term}"' + + +def _sparql_to_cypher(sparql: str, scope: frozenset[str]) -> tuple[str, dict[str, Any]]: + """Translate a simple SPARQL SELECT to a Cypher MATCH query. + + Only handles the subset of SPARQL used by UKO ontology queries: + ``SELECT ?vars WHERE { triple_patterns }``. Complex SPARQL + (OPTIONAL, UNION, FILTER, etc.) is passed through as a raw Cypher + comment so that the caller receives an empty result rather than + crashing. + + Returns: + A ``(cypher_query, params)`` tuple. + """ + m = _SPARQL_SELECT_RE.search(sparql) + if not m: + # Unsupported SPARQL — return a no-op query + logger.warning( + "neo4j.sparql_translation.unsupported", + sparql=sparql[:200], + ) + return "MATCH (n) WHERE false RETURN n", {} + + body = m.group("body").strip() + triples = _TRIPLE_RE.findall(body) + if not triples: + return "MATCH (n) WHERE false RETURN n", {} + + params: dict[str, Any] = {} + match_clauses: list[str] = [] + return_vars: set[str] = set() + + for idx, (s, p, o) in enumerate(triples): + s_cypher = _uri_to_cypher(s) + p_cypher = _uri_to_cypher(p) + o_cypher = _uri_to_cypher(o) + + s_var = f"s{idx}" + p_var = f"p{idx}" + o_var = f"o{idx}" + + # Build MATCH clause + if s_cypher.startswith("?"): + s_expr = f"({s_var})" + return_vars.add(s_var) + params[s_var] = None + else: + s_expr = f"({s_var} {{uri: $s{idx}_val}})" + params[f"s{idx}_val"] = s_cypher.strip('"') + + if o_cypher.startswith("?"): + o_expr = f"({o_var})" + return_vars.add(o_var) + else: + o_expr = f"({o_var} {{uri: $o{idx}_val}})" + params[f"o{idx}_val"] = o_cypher.strip('"') + + rel_type = ( + p_cypher.strip('"').replace(":", "_").replace("/", "_").replace("-", "_") + ) + if p_cypher.startswith("?"): + rel_expr = f"-[{p_var}]->" + return_vars.add(p_var) + else: + rel_expr = f"-[:{rel_type}]->" + + match_clauses.append(f"MATCH {s_expr}{rel_expr}{o_expr}") + + # Scope filter — restrict to nodes whose project property is in scope + scope_filter = "" + if scope: + scope_list = list(scope) + params["scope"] = scope_list + scope_filter = " WHERE s0.project IN $scope" + + return_clause = ( + "RETURN " + ", ".join(f"{v}.uri AS {v}" for v in sorted(return_vars)) + if return_vars + else "RETURN s0.uri AS s0, p0, o0.uri AS o0" + ) + + cypher = "\n".join(match_clauses) + scope_filter + "\n" + return_clause + return cypher, params + + +# --------------------------------------------------------------------------- +# Neo4jGraphBackend (read-side) +# --------------------------------------------------------------------------- + + +class Neo4jGraphBackend: + """Read-side ACMS graph backend backed by Neo4j. + + Implements the :class:`GraphBackend` protocol. SPARQL queries are + translated to Cypher via a minimal translator; ``get_triples`` and + ``traverse`` use native Cypher directly. + + Args: + url: Neo4j bolt URL (e.g. ``bolt://localhost:7687``). + auth: ``(user, password)`` tuple. + """ + + def __init__(self, url: str, auth: tuple[str, str]) -> None: + self._url = url + self._auth = auth + self._driver: Any = None + + def _get_driver(self) -> Any: + """Return (or lazily create) the Neo4j driver.""" + if self._driver is None: + if _GraphDatabase is None: + raise RuntimeError("neo4j driver is not installed") + self._driver = _GraphDatabase.driver(self._url, auth=self._auth) + return self._driver + + def sparql_query( + self, + query: str, + *, + scope: frozenset[str], + ) -> GraphResult: + """Execute a SPARQL query translated to Cypher within *scope*. + + Args: + query: SPARQL query string. Must be non-empty. + scope: Frozenset of resource ULIDs to restrict results to. + + Returns: + A :class:`GraphResult` containing matched triples. + + Raises: + ValueError: If *query* is empty. + """ + if not query: + raise ValueError("query must be a non-empty string") + try: + cypher, params = _sparql_to_cypher(query, scope) + driver = self._get_driver() + triples: list[tuple[str, str, str]] = [] + with driver.session() as session: + result = session.run(cypher, **params) + for record in result: + keys = list(record.keys()) + if len(keys) >= 3: + triples.append( + ( + str(record[keys[0]] or ""), + str(record[keys[1]] or ""), + str(record[keys[2]] or ""), + ) + ) + elif len(keys) == 2: + triples.append( + (str(record[keys[0]] or ""), "", str(record[keys[1]] or "")) + ) + return GraphResult(triples=triples) + except _ServiceUnavailable: + logger.warning( + "neo4j.graph_backend.unavailable", + url=self._url, + operation="sparql_query", + ) + return GraphResult() + except Exception as exc: + logger.error( + "neo4j.graph_backend.error", + operation="sparql_query", + error=str(exc), + ) + return GraphResult() + + def get_triples( + self, + subject: str, + ) -> GraphResult: + """Retrieve all triples for a given *subject*. + + Args: + subject: The UKO URI of the subject node. Must be non-empty. + + Returns: + A :class:`GraphResult` containing the subject's triples. + + Raises: + ValueError: If *subject* is empty. + """ + if not subject: + raise ValueError("subject must be a non-empty string") + try: + driver = self._get_driver() + cypher = ( + "MATCH (s {uri: $subject})-[r]->(o) " + "RETURN s.uri AS s, type(r) AS p, o.uri AS o" + ) + triples: list[tuple[str, str, str]] = [] + with driver.session() as session: + result = session.run(cypher, subject=subject) + for record in result: + triples.append( + ( + str(record["s"] or ""), + str(record["p"] or ""), + str(record["o"] or ""), + ) + ) + return GraphResult(triples=triples) + except _ServiceUnavailable: + logger.warning( + "neo4j.graph_backend.unavailable", + url=self._url, + operation="get_triples", + ) + return GraphResult() + except Exception as exc: + logger.error( + "neo4j.graph_backend.error", + operation="get_triples", + error=str(exc), + ) + return GraphResult() + + def traverse( + self, + start: str, + *, + depth: int = 2, + ) -> GraphResult: + """Traverse the graph from *start* up to *depth* hops. + + Args: + start: UKO URI to begin traversal from. Must be non-empty. + depth: Maximum number of hops from *start*. Must be + non-negative. + + Returns: + A :class:`GraphResult` containing discovered triples. + + Raises: + ValueError: If *start* is empty or *depth* is negative. + """ + if not start: + raise ValueError("start must be a non-empty string") + if depth < 0: + raise ValueError(f"depth must be non-negative, got {depth}") + try: + driver = self._get_driver() + cypher = ( + "MATCH path = (s {uri: $start})-[*1..$depth]->(o) " + "UNWIND relationships(path) AS r " + "RETURN startNode(r).uri AS s, type(r) AS p, endNode(r).uri AS o" + ) + triples: list[tuple[str, str, str]] = [] + seen: set[tuple[str, str, str]] = set() + with driver.session() as session: + result = session.run(cypher, start=start, depth=depth) + for record in result: + triple = ( + str(record["s"] or ""), + str(record["p"] or ""), + str(record["o"] or ""), + ) + if triple not in seen: + seen.add(triple) + triples.append(triple) + return GraphResult(triples=triples) + except _ServiceUnavailable: + logger.warning( + "neo4j.graph_backend.unavailable", + url=self._url, + operation="traverse", + ) + return GraphResult() + except Exception as exc: + logger.error( + "neo4j.graph_backend.error", + operation="traverse", + error=str(exc), + ) + return GraphResult() + + def close(self) -> None: + """Close the underlying Neo4j driver connection.""" + if self._driver is not None: + try: + self._driver.close() + except Exception: + pass + finally: + self._driver = None + + +# --------------------------------------------------------------------------- +# Neo4jGraphIndexBackend (write-side) +# --------------------------------------------------------------------------- + + +class Neo4jGraphIndexBackend: + """Write-side ACMS graph index backend backed by Neo4j. + + Implements the :class:`GraphIndexBackend` protocol. Each triple is + stored as two ``Resource`` nodes connected by a typed relationship. + The ``project`` property on each node enables project-scoped queries. + + Args: + url: Neo4j bolt URL. + auth: ``(user, password)`` tuple. + """ + + def __init__(self, url: str, auth: tuple[str, str]) -> None: + self._url = url + self._auth = auth + self._driver: Any = None + + def _get_driver(self) -> Any: + """Return (or lazily create) the Neo4j driver.""" + if self._driver is None: + if _GraphDatabase is None: + raise RuntimeError("neo4j driver is not installed") + self._driver = _GraphDatabase.driver(self._url, auth=self._auth) + return self._driver + + def add_triple( + self, + project: str, + subject: str, + predicate: str, + obj: str, + ) -> None: + """Add a single triple to the Neo4j graph. + + Uses ``MERGE`` to avoid duplicate nodes/relationships. + + Args: + project: Namespaced project name. + subject: Subject URI. + predicate: Predicate URI (used as relationship type). + obj: Object URI or literal value. + + Raises: + ValueError: If any argument is empty or whitespace-only. + """ + _require_non_empty(project, "project") + _require_non_empty(subject, "subject") + _require_non_empty(predicate, "predicate") + _require_non_empty(obj, "obj") + + # Relationship types in Cypher must be identifiers + rel_type = ( + predicate.replace(":", "_") + .replace("/", "_") + .replace("-", "_") + .replace(".", "_") + ) + # Ensure rel_type starts with a letter + if rel_type and not rel_type[0].isalpha(): + rel_type = "R_" + rel_type + + try: + driver = self._get_driver() + rel_merge = ( + f"MERGE (s)-[:{rel_type} " + "{{predicate: $predicate, project: $project}}]->(o)" + ) + cypher = ( + "MERGE (s:Resource {uri: $subject, project: $project}) " + "MERGE (o:Resource {uri: $obj, project: $project}) " + rel_merge + ) + with driver.session() as session: + session.run( + cypher, + subject=subject, + obj=obj, + predicate=predicate, + project=project, + ) + logger.debug( + "neo4j.index.triple_added", + project=project, + subject=subject, + predicate=predicate, + ) + except _ServiceUnavailable: + logger.warning( + "neo4j.graph_index_backend.unavailable", + url=self._url, + operation="add_triple", + ) + except Exception as exc: + logger.error( + "neo4j.graph_index_backend.error", + operation="add_triple", + error=str(exc), + ) + + def query( + self, + project: str, + sparql: str, + ) -> list[dict[str, str]]: + """Execute a SPARQL query translated to Cypher. + + Args: + project: Namespaced project name. + sparql: SPARQL query string. + + Returns: + List of binding dictionaries. + + Raises: + ValueError: If *project* or *sparql* is empty or + whitespace-only. + """ + _require_non_empty(project, "project") + _require_non_empty(sparql, "sparql") + + try: + cypher, params = _sparql_to_cypher(sparql, frozenset({project})) + driver = self._get_driver() + results: list[dict[str, str]] = [] + with driver.session() as session: + result = session.run(cypher, **params) + for record in result: + results.append({k: str(v or "") for k, v in record.items()}) + return results + except _ServiceUnavailable: + logger.warning( + "neo4j.graph_index_backend.unavailable", + url=self._url, + operation="query", + ) + return [] + except Exception as exc: + logger.error( + "neo4j.graph_index_backend.error", + operation="query", + error=str(exc), + ) + return [] + + def remove_triples( + self, + project: str, + subject: str | None, + predicate: str | None, + obj: str | None, + ) -> None: + """Remove triples matching the given pattern. + + Args: + project: Namespaced project name. + subject: Subject URI filter, or ``None`` for wildcard. + predicate: Predicate URI filter, or ``None`` for wildcard. + obj: Object filter, or ``None`` for wildcard. + + Raises: + ValueError: If *project* is empty or whitespace-only, or + all three filters are ``None``. + """ + _require_non_empty(project, "project") + if subject is None and predicate is None and obj is None: + raise ValueError( + "At least one of subject, predicate, obj must be " + "non-None to avoid deleting all triples" + ) + if subject is not None and not subject.strip(): + raise ValueError("subject filter must be non-empty if provided") + if predicate is not None and not predicate.strip(): + raise ValueError("predicate filter must be non-empty if provided") + if obj is not None and not obj.strip(): + raise ValueError("obj filter must be non-empty if provided") + + try: + driver = self._get_driver() + # Build WHERE conditions + conditions: list[str] = ["r.project = $project"] + params: dict[str, Any] = {"project": project} + if subject is not None: + conditions.append("s.uri = $subject") + params["subject"] = subject + if predicate is not None: + conditions.append("r.predicate = $predicate") + params["predicate"] = predicate + if obj is not None: + conditions.append("o.uri = $obj") + params["obj"] = obj + + where_clause = " AND ".join(conditions) + cypher = ( + f"MATCH (s:Resource)-[r]->(o:Resource) WHERE {where_clause} DELETE r" + ) + with driver.session() as session: + session.run(cypher, **params) + logger.debug( + "neo4j.index.triples_removed", + project=project, + subject=subject, + predicate=predicate, + obj=obj, + ) + except _ServiceUnavailable: + logger.warning( + "neo4j.graph_index_backend.unavailable", + url=self._url, + operation="remove_triples", + ) + except Exception as exc: + logger.error( + "neo4j.graph_index_backend.error", + operation="remove_triples", + error=str(exc), + ) + + def close(self) -> None: + """Close the underlying Neo4j driver connection.""" + if self._driver is not None: + try: + self._driver.close() + except Exception: + pass + finally: + self._driver = None + + +# --------------------------------------------------------------------------- +# Protocol compliance assertions (static verification by Pyright) +# --------------------------------------------------------------------------- + +_assert_read: type[GraphBackend] = Neo4jGraphBackend # type: ignore[type-abstract] +_assert_write: type[GraphIndexBackend] = Neo4jGraphIndexBackend # type: ignore[type-abstract] + + +# --------------------------------------------------------------------------- +# Factory functions (used by DI container) +# --------------------------------------------------------------------------- + + +def _parse_auth(auth_str: str) -> tuple[str, str]: + """Parse ``user:password`` auth string into a tuple.""" + if ":" in auth_str: + user, _, password = auth_str.partition(":") + return user, password + return auth_str, "" + + +def build_graph_backend(config_service: Any) -> Any: + """Return the configured ACMS read-side graph backend. + + Reads ``index.graph.backend``, ``index.graph.neo4j-url``, and + ``index.graph.neo4j-auth`` from :class:`ConfigService`. Falls back + to :class:`InMemoryGraphBackend` when: + + - The configured backend is not ``"neo4j"`` + - The ``neo4j`` driver package is not installed + - The Neo4j server is unreachable at startup + + Args: + config_service: A :class:`ConfigService` instance. + + Returns: + A :class:`GraphBackend`-compatible object. + """ + try: + backend_cfg = config_service.resolve("index.graph.backend") + backend_name: str = str(backend_cfg.value or "none").lower() + except Exception: + backend_name = "none" + + if backend_name != "neo4j": + logger.info( + "acms.graph_backend.fallback", + configured_backend=backend_name, + fallback_backend="InMemoryGraphBackend", + ) + return InMemoryGraphBackend() + + if not _NEO4J_AVAILABLE: + logger.warning( + "acms.graph_backend.neo4j_unavailable", + reason="neo4j package not installed", + fallback_backend="InMemoryGraphBackend", + ) + return InMemoryGraphBackend() + + try: + url_cfg = config_service.resolve("index.graph.neo4j-url") + url: str = str(url_cfg.value or "bolt://localhost:7687") + auth_cfg = config_service.resolve("index.graph.neo4j-auth") + auth_str: str = str(auth_cfg.value or "neo4j:neo4j") + auth = _parse_auth(auth_str) + except Exception as exc: + logger.warning( + "acms.graph_backend.config_error", + error=str(exc), + fallback_backend="InMemoryGraphBackend", + ) + return InMemoryGraphBackend() + + # Probe connectivity — fall back gracefully if Neo4j is down + if _GraphDatabase is None: + return InMemoryGraphBackend() + try: + probe_driver = _GraphDatabase.driver(url, auth=auth) + probe_driver.verify_connectivity() + probe_driver.close() + except Exception as exc: + logger.warning( + "acms.graph_backend.connection_failed", + url=url, + error=str(exc), + fallback_backend="InMemoryGraphBackend", + ) + return InMemoryGraphBackend() + + logger.info("acms.graph_backend.neo4j_activated", url=url) + return Neo4jGraphBackend(url=url, auth=auth) + + +def build_graph_index_backend(config_service: Any) -> Any: + """Return the configured ACMS write-side graph index backend. + + Uses the same configuration keys as :func:`build_graph_backend`. + Falls back to :class:`InMemoryGraphIndexBackend` on any error. + + Args: + config_service: A :class:`ConfigService` instance. + + Returns: + A :class:`GraphIndexBackend`-compatible object. + """ + try: + backend_cfg = config_service.resolve("index.graph.backend") + backend_name: str = str(backend_cfg.value or "none").lower() + except Exception: + backend_name = "none" + + if backend_name != "neo4j": + logger.info( + "acms.graph_index_backend.fallback", + configured_backend=backend_name, + fallback_backend="InMemoryGraphIndexBackend", + ) + return InMemoryGraphIndexBackend() + + if not _NEO4J_AVAILABLE: + logger.warning( + "acms.graph_index_backend.neo4j_unavailable", + reason="neo4j package not installed", + fallback_backend="InMemoryGraphIndexBackend", + ) + return InMemoryGraphIndexBackend() + + try: + url_cfg = config_service.resolve("index.graph.neo4j-url") + url: str = str(url_cfg.value or "bolt://localhost:7687") + auth_cfg = config_service.resolve("index.graph.neo4j-auth") + auth_str: str = str(auth_cfg.value or "neo4j:neo4j") + auth = _parse_auth(auth_str) + except Exception as exc: + logger.warning( + "acms.graph_index_backend.config_error", + error=str(exc), + fallback_backend="InMemoryGraphIndexBackend", + ) + return InMemoryGraphIndexBackend() + + # Probe connectivity + if _GraphDatabase is None: + return InMemoryGraphIndexBackend() + try: + probe_driver = _GraphDatabase.driver(url, auth=auth) + probe_driver.verify_connectivity() + probe_driver.close() + except Exception as exc: + logger.warning( + "acms.graph_index_backend.connection_failed", + url=url, + error=str(exc), + fallback_backend="InMemoryGraphIndexBackend", + ) + return InMemoryGraphIndexBackend() + + logger.info("acms.graph_index_backend.neo4j_activated", url=url) + return Neo4jGraphIndexBackend(url=url, auth=auth) + + +__all__: list[str] = [ + "Neo4jGraphBackend", + "Neo4jGraphIndexBackend", + "build_graph_backend", + "build_graph_index_backend", +] diff --git a/vulture_whitelist.py b/vulture_whitelist.py index cef933f03..9b8140e81 100644 --- a/vulture_whitelist.py +++ b/vulture_whitelist.py @@ -1203,3 +1203,7 @@ create_workspace_snapshot # noqa: B018, F821 selective_rollback # noqa: B018, F821 archive_artifacts # noqa: B018, F821 revert_decisions # noqa: B018, F821 + +# Neo4j graph backend — close() is public API for resource cleanup (#872) +Neo4jGraphBackend # noqa: B018, F821 +Neo4jGraphIndexBackend # noqa: B018, F821 -- 2.52.0