How to Use Neo4j for Custom Queries in Flowsint
To use Neo4j for custom queries in Flowsint, instantiate the Neo4jGraphRepository class and call its query() method with parameterized Cypher statements, which forwards the request to the underlying Neo4j driver.
Flowsint persists investigation data in a Neo4j graph database and exposes a high-level Python API for interacting with it. The Neo4jGraphRepository class, defined in flowsint-core/src/flowsint_core/core/graph/repository.py, serves as the primary interface for executing arbitrary Cypher queries that extend beyond the standard node and relationship CRUD operations.
Neo4jGraphRepository Architecture
The repository abstracts the Neo4j driver behind a singleton connection pattern, ensuring efficient connection reuse across the application while providing a clean Pythonic API for graph operations.
Core Implementation Details
The Neo4jGraphRepository wraps a singleton Neo4jConnection instance. When instantiated, it either receives an explicit connection object or obtains the shared singleton via Neo4jConnection.get_instance() from the same package. All methods build parameterized Cypher strings and forward them to the driver through self._connection.query or transaction-based methods like execute_write.
The query() Method
The generic query method (lines 51‑66 in repository.py) serves as the escape hatch for custom Cypher that is not covered by higher-level helpers. It accepts a raw Cypher string and a dictionary of parameters, returning a list of dictionaries representing Neo4j records.
def query(
self, cypher: str, parameters: Dict[str, Any] = {}
) -> List[Dict[str, Any]]:
"""
Execute a custom Cypher query.
"""
if not self._connection:
return []
return self._connection.query(cypher, parameters)
Source: [flowsint-core/src/flowsint_core/core/graph/repository.py](https://github.com/reconurge/flowsint/blob/main/flowsint-core/src/flowsint_core/core/graph/repository.py#L51-L66)
Because the repository is utilized throughout the core service layer in flowsint-core/src/flowsint_core/core/graph/service.py, you can invoke query directly on a repository instance or obtain it via the higher-level Neo4jGraphService, which mirrors the same method signature.
Executing Custom Cypher Queries
When using Neo4j for custom queries in Flowsint, follow this workflow: obtain a repository instance, write parameterized Cypher using $param syntax, pass a dictionary of parameters, and process the returned list of dictionaries. Each dictionary corresponds to a Neo4j record where keys match the RETURN clause aliases.
The repository automatically respects sketch isolation (sketch_id) and soft-delete semantics (deleted_at IS NULL). Custom queries should explicitly include these filters when operating on investigation data to maintain data integrity.
Practical Implementation Examples
Querying Nodes by Label
The following example demonstrates a simple node lookup using custom Cypher against the repository:
from flowsint_core.core.graph import Neo4jGraphRepository
repo = Neo4jGraphRepository() # uses the singleton connection
cypher = """
MATCH (n:Domain {nodeLabel: $label, sketch_id: $sketch_id})
RETURN elementId(n) AS id, properties(n) AS data
"""
params = {"label": "example.com", "sketch_id": "sketch-123"}
results = repo.query(cypher, params)
# results → [{'id': '12345', 'data': {'nodeLabel': 'example.com', ...}}]
Batch Updating Properties
For efficient batch operations, use the UNWIND clause to process multiple updates in a single transaction:
cypher = """
UNWIND $updates AS upd
MATCH (n) WHERE elementId(n) = upd.id AND n.sketch_id = $sketch_id
SET n += upd.props
RETURN count(n) AS updated
"""
params = {
"sketch_id": "sketch-123",
"updates": [
{"id": "12345", "props": {"confidence": 0.9}},
{"id": "67890", "props": {"confidence": 0.85}},
],
}
result = repo.query(cypher, params)
print(result[0]["updated"]) # → 2
Executing Soft Deletes
Flowsint uses soft deletes rather than hard removals. Custom queries should respect this pattern by setting the deleted_at timestamp:
cypher = """
OPTIONAL MATCH (n {sketch_id: $sketch_id})
WHERE n.deleted_at IS NULL
SET n.deleted_at = datetime()
RETURN count(n) AS soft_deleted
"""
result = repo.query(cypher, {"sketch_id": "sketch-123"})
print(f"Soft‑deleted {result[0]['soft_deleted']} nodes")
FastAPI Endpoint Integration
In a web service context, use FastAPI dependency injection to obtain the repository:
from fastapi import APIRouter, Depends
from flowsint_core.core.graph.service import Neo4jGraphRepository
router = APIRouter()
@router.get("/custom")
def run_custom_query(
repo: Neo4jGraphRepository = Depends(),
label: str = "example.com",
sketch_id: str = "sketch-123",
):
cypher = """
MATCH (n {nodeLabel: $label, sketch_id: $sketch_id})
RETURN elementId(n) AS id, properties(n) AS data
"""
return repo.query(cypher, {"label": label, "sketch_id": sketch_id})
Security and Data Integrity Considerations
Always use parameterized queries with the $parameter syntax rather than Python f-strings or concatenation to prevent Cypher injection attacks. The query method passes parameters directly to the Neo4j driver without modification, ensuring safe execution.
When writing custom queries, include sketch_id filters to maintain multi-tenancy boundaries between investigations. Additionally, check for deleted_at IS NULL unless your specific use case requires accessing archived data.
Summary
- Neo4jGraphRepository in
flowsint-core/src/flowsint_core/core/graph/repository.pyprovides thequery()method for executing arbitrary Cypher. - The method signature accepts a Cypher string and parameter dictionary, returning a list of record dictionaries.
- The repository manages a singleton connection via
Neo4jConnectionto optimize driver performance. - Custom queries must respect
sketch_idisolation anddeleted_atsoft-delete semantics for data consistency. - Integration with FastAPI uses standard dependency injection patterns against the repository or service layer.
Frequently Asked Questions
How do I access the raw Neo4j driver for advanced configurations?
You generally should not access the raw driver directly. Instead, instantiate Neo4jGraphRepository, which encapsulates the connection logic in flowsint-core/src/flowsint_core/core/graph/connection.py. If you need specific driver features, extend the repository class rather than bypassing it, as the singleton pattern ensures proper connection pooling and transaction management across the application.
What is the difference between Neo4jGraphRepository and Neo4jGraphService?
Neo4jGraphRepository provides low-level database access and the generic query() method for custom Cypher. Neo4jGraphService, defined in flowsint-core/src/flowsint_core/core/graph/service.py, composes the repository and adds domain-specific business logic, validation, and orchestration. For custom queries, you can use either class, as the service exposes the same query signature while adding higher-level abstractions for standard operations.
Do I need to manually close connections when using the query method?
No. The repository utilizes a singleton Neo4jConnection that manages the Neo4j driver lifecycle automatically. When you call repo.query(), the underlying connection remains open for subsequent operations, and the driver handles connection pooling internally. This design prevents connection leaks and reduces latency for repeated queries within the same process.
How should I index properties for performance in custom queries?
Flowsint creates indexes on commonly queried properties such as sketch_id, as shown in neo4j-migrations/001_indexes.cypher. When writing custom queries that filter on specific properties, ensure those properties are indexed in your Neo4j schema to avoid expensive full-graph scans. The repository does not automatically create indexes for custom query parameters, so database optimization remains the responsibility of the developer.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →