knowledge_graph
FalkorDB-backed knowledge graph with hybrid vector+graph retrieval.
- class knowledge_graph.KnowledgeGraphManager(redis_client, openrouter, embedding_model='google/gemini-embedding-001', admin_user_ids=None, dedup_threshold=0.9, read_redis_client=None)
Bases:
objectManages the FalkorDB knowledge graph for entity/relationship CRUD and hybrid retrieval.
- Parameters:
redis_client (aioredis.Redis)
openrouter (OpenRouterClient)
embedding_model (str)
dedup_threshold (float)
read_redis_client (aioredis.Redis | None)
- GRAPH_NAME = 'knowledge'
- __init__(redis_client, openrouter, embedding_model='google/gemini-embedding-001', admin_user_ids=None, dedup_threshold=0.9, read_redis_client=None)
Initialize the instance.
- Parameters:
redis_client (aioredis.Redis) – Redis connection client.
openrouter (OpenRouterClient) – The openrouter value.
embedding_model (str) – The embedding model value.
admin_user_ids (set[str] | None) – The admin user ids value.
read_redis_client (aioredis.Redis | None) – Optional replica-routed client (see
redis_replica_selector) used ONLY for read-only Cypher (GRAPH.RO_QUERY). Reads through it may be seconds stale, so any read that gates a write (entity dedup, MERGE-on-read, index-existence checks before DDL) must go through the master handle instead. Defaults to the master client when omitted.dedup_threshold (float)
- Return type:
None
- property indexes_ready: bool
True once
ensure_indexes()has completed all phases.Retrieval code checks this flag before issuing HNSW KNN queries; if False it means the index warm-up is still in progress and vector search will return no results from FalkorDB.
- async wait_for_foreground_idle()
Block until no foreground queries are pending or executing.
Called by background workers before each query attempt. Because background callers hold _query_semaphore for only one query at a time, a foreground caller that arrives mid-background-query only waits for the current background query to finish (not the entire batch).
- Return type:
- async wait_for_idle()
Alias for
wait_for_foreground_idle().- Return type:
- async query(q, params=None, is_background=False, **kwargs)
Execute a Cypher query against the knowledge graph.
Foreground callers (is_background=False) increment the waiter counter before acquiring the concurrency semaphore so that background workers will yield to them at their next boundary.
- async ro_query(q, params=None, is_background=False, **kwargs)
Execute a read-only Cypher query against the knowledge graph.
Routes to the replica-backed graph handle (with master fallback) and applies the same concurrency and priority logic as
query(). Callers must not depend on reading their own just-written data — see theread_redis_clientconstructor note.
- async embed(text)
Public proxy for the internal _embed method.
- async embed_batch(texts)
Public proxy for the internal _embed_batch method.
- async ensure_indexes()
Create vector + range indexes for every entity label.
Optimized with two pre-flight checks to eliminate unnecessary FalkorDB round-trips on warm restarts:
Strategy 1 — existence pre-check:
CALL db.indexes()is issued once before the Phase 1 loop. Labels whose HNSW index is already present are skipped entirely; noCREATE VECTOR INDEXcall is made for them.Strategy 2 — zero-node skip: A single
MATCHaggregation counts nodes per label. Labels with zero nodes are also skipped because building an empty index wastes a full server round-trip.Strategy 4 — readiness flag:
self._indexes_readyis set toTrueonly after all phases complete successfully. Retrieval code checks this flag before issuing KNN queries.
Both pre-checks fail safe: if the introspection query errors (e.g. older FalkorDB that lacks
db.indexes()), the corresponding skip-set is empty and all labels are attempted as before.Warming is serialised across services by a Redis boot lock. Acquisition is bounded (
_acquire_index_lock()): an endpoint that rejects writes — a demoted master, a read-only replica — used to spin here forever at one WARNING every two seconds, leaving_indexes_readyFalseand every KG vector search silently returning[]for the rest of the process’s life. It now gives up with a diagnosis and re-arms itself in the background (_schedule_index_warm_retry()).- Return type:
- async add_entity(name, entity_type, description, category='general', scope_id='_', created_by='unknown', pinned=False, metadata='{}', user_id='000000000000', embedding=None)
Create or update an entity.
Returns
{"name": ..., "uuid": ...}.
- async update_entity_description(name, entity_type, new_description, category=None, scope_id=None)
Update an entity’s description and re-embed.
- async edit_entity(uuid, description=None, append_text=None, pinned=None, category=None, metadata_updates=None)
Selectively update fields on an existing entity.
Looks up by uuid. Only the provided fields are changed; everything else is preserved.
description replaces the text entirely. append_text is concatenated to the existing description. (Mutually exclusive – caller must pick one.)
metadata_updates is shallow-merged into the existing metadata JSON (new keys added, existing overwritten, unmentioned preserved).
Returns the full updated entity dict via
get_entity(), orNoneif the UUID was not found.
- async delete_entity(name, entity_type, category, scope_id='_')
Delete the specified entity.
- async delete_entity_by_uuid(uuid)
Delete an entity by UUID (detach-deletes all relationships).
- async pin_entity(name, entity_type, pinned=True, category=None, scope_id=None)
Set or clear the pinned flag on an entity.
When category and/or scope_id are provided, only entities matching those filters are updated. This avoids pinning the wrong entity when the same name exists in multiple scopes.
- async get_entity(name='', entity_type=None, category=None, scope_id=None, uuid=None)
Fetch an entity with its immediate connections.
Can look up by name or by uuid.
- async inspect_entity(name='', uuid=None, max_depth=2, neighbor_limit=50)
Deep inspection of an entity and its full neighborhood.
Returns the entity’s properties plus all outgoing and incoming relationships (up to max_depth hops), with each neighbor’s core properties included.
- async list_entities(entity_type=None, category=None, scope_id=None, limit=50, offset=0, search=None)
List entities with optional filtering, pagination, and text search.
- async add_relationship(source_uuid, target_uuid, relation_type, weight=0.5, description='', evidence='')
Create or reinforce a relationship between two entities identified by UUID.
The edge inherits
priority = min(source, target)and thecategory/scope_idfrom the lower-priority endpoint. Cross-category edges are fully supported.
- async delete_relationship(source_uuid, target_uuid, relation_type)
Delete the specified relationship.
- async list_relationships(entity_uuid=None, relation_type=None, category=None, limit=50, order_by=True, timeout=None)
List relationships.
- async resolve_entity_cross_category(name, entity_type)
Find an entity by name across all categories.
Used for cross-category linking.
- async retrieve_context(query, query_embedding=None, user_ids=None, channel_id=None, guild_id=None, max_hops=2, max_per_user=60, max_channel=15, max_guild=15, max_general=30, max_per_lore=20, seed_top_k=64, seed_similarity_threshold=0.38, seed_limit=15, min_edge_weight=0.0, default_edge_weight=0.8, semantic_hop_decay=0.8, expansion_neighbor_limit=500, dynamic_threshold_enabled=True, dynamic_threshold_target_ratio=0.1, dynamic_threshold_min=0.2, dynamic_threshold_min_stored=5, full_user_memory_ids=None, user_seed_min=10, user_candidate_limit=100, lore_candidate_limit=40, lore_seed_min=5, lore_amplified=False, max_per_meta=20, meta_candidate_limit=40, meta_seed_min=5, meta_amplified=False, max_context_entities=None)
- Return type:
- Parameters:
query (str)
channel_id (str | None)
guild_id (str | None)
max_hops (int)
max_per_user (int)
max_channel (int)
max_guild (int)
max_general (int)
max_per_lore (int)
seed_top_k (int)
seed_similarity_threshold (float)
seed_limit (int)
min_edge_weight (float)
default_edge_weight (float)
semantic_hop_decay (float)
expansion_neighbor_limit (int)
dynamic_threshold_enabled (bool)
dynamic_threshold_target_ratio (float)
dynamic_threshold_min (float)
dynamic_threshold_min_stored (int)
user_seed_min (int)
user_candidate_limit (int)
lore_candidate_limit (int)
lore_seed_min (int)
lore_amplified (bool)
max_per_meta (int)
meta_candidate_limit (int)
meta_seed_min (int)
meta_amplified (bool)
max_context_entities (int | None)
- async search_entities(query, query_embedding=None, category=None, scope_id=None, entity_type=None, top_k=10)
- async reconsolidate_embeddings_on_recall(targets, query_embedding, *, learning_rate=0.03, max_step=0.25)
Nudge recalled entities’ embeddings toward the query with a tanh clamp.
Fire-and-forget from prompt build; failures are logged at debug.
Batched: instead of one read + one write round-trip per entity (up to 2 ×
mementropic_reconsolidation_max_entitiessequential round-trips, each embedding write taking the FalkorDB write lock and touching the HNSW index while the next message’s retrieval waits), the targets are grouped by graph label and each group is serviced by oneUNWINDread followed by oneUNWINDwrite.The mutation set is unchanged: the same entities are written, with the same vectors and the same single
nowtimestamp, in the same order. Grouping is by label rather than one global pair of queries because theuuidrange indexes are created per label (see_create_range_indexes_for_label()); a label-lessMATCH (e {uuid: ...})could not use them and would degrade to a full node scan per row. Each batched query therefore keeps the exact indexedMATCH (e:Label {uuid: ...})pattern the per-entity path used, only driven byUNWIND.Per-entity fault isolation is preserved by falling back: if the batched read or the batched write raises, the original one-query-per-entity path runs for the affected work, so a single poisoned row can still not cost the rest of the batch.