Merge a reviewed entity pair
To combine two known duplicate entities, choose the source and target explicitly, merge that pair, and verify the source marker and target aliases. Review semantic candidates before applying the same operation to user data.
|
Available on NAMS: No. The procedure on this page needs the Bolt backend. With the hosted NAMS backend its operations are unavailable: depending on the call, the client raises |
|
Two write paths, one decision (0.7.0). |
Prerequisites
Use a dedicated AuraDB instance from the first memory tutorial’s Aura setup, with exported NEO4J_URI, NEO4J_USERNAME, and NEO4J_PASSWORD. Set OPENAI_API_KEY for the selected embedding provider. Prepare the example files below and install the provider extra:
Create a local folder and virtual environment for the examples on this page. The commands reuse an existing environment without changing its files:
mkdir -p ~/agent-memory-tutorials
cd ~/agent-memory-tutorials
if [ -e .venv ]; then
printf '%s\n' 'Using the existing virtual environment.'
else
python3 -m venv .venv
fi
source .venv/bin/activate
Expected: ~/agent-memory-tutorials is your working directory and its virtual environment is active. Install the published SDK with the command below.
Each complete code block labelled Save as names a file to create in this folder using your editor. Copy the entire block, including imports and the entry point. Expand each helper disclosure and use Copy code to copy its full source. Keep all files together so their imports resolve.
When continuing from another tutorial or guide, retain the existing environment, configuration, session files, and .tutorial-state/. Reuse unchanged helper files; compare an existing file before replacing it, and finish any pending cleanup or recovery before changing the code that owns its state.
python -m pip install 'neo4j-agent-memory[openai]==0.7.0'
The maintained program core_memory_recipes.py imports core_memory_settings.py from the same directory. That helper explicitly selects Bolt and disables automatic extraction/enrichment. Each command creates a new ID suffix so it can run independently. Embedding requests send example text to OpenAI.
These tasks use Python with Bolt, and not every call they make is portable to NAMS. On the hosted backend, preferences, facts, direct relationship writes, client-side deduplication configuration, add_messages_batch (use bulk_add_messages on NAMS), session listing and summaries, and reasoning search and statistics are unavailable: depending on the call, the client raises NotSupportedError, AttributeError or TypeError, or silently drops a Bolt-only option. Consult backend capabilities before adapting a task to NAMS.
Save the three complete files below in agent-memory-tutorials/. If a file already exists from another guide, keep its matching contents; do not replace configuration or state files. The task-specific excerpt in the next section explains the operation to run.
Open the helper below and save its complete source as aura_connection.py in ~/agent-memory-tutorials.
Show aura_connection.py
aura_connection.py"""Read the Aura connection exported by the documentation setup commands."""
import os
class AuraConfigurationError(ValueError):
"""Missing or invalid tutorial settings, with no credential values in errors."""
def aura_config():
required = ("NEO4J_URI", "NEO4J_USERNAME", "NEO4J_PASSWORD")
missing = [name for name in required if not os.environ.get(name, "").strip()]
if missing:
raise AuraConfigurationError("Export the Aura connection settings: " + ", ".join(missing))
uri = os.environ["NEO4J_URI"]
if not uri.startswith("neo4j+s://"):
raise AuraConfigurationError("NEO4J_URI must use the Aura neo4j+s:// connection scheme")
database = os.environ.get("NEO4J_DATABASE", "neo4j")
if not database.strip():
raise AuraConfigurationError("NEO4J_DATABASE must not be empty")
return {
"uri": uri,
"username": os.environ["NEO4J_USERNAME"],
"password": os.environ["NEO4J_PASSWORD"],
"database": database,
}
Complete core_memory_settings.py
core_memory_settings.py"""Shared settings for the three Aura-backed memory tutorials."""
from aura_connection import aura_config
from neo4j_agent_memory import MemorySettings
def settings():
return MemorySettings(
backend="bolt",
neo4j=aura_config(),
embedding="openai/text-embedding-3-small",
llm=None,
extraction={"extractor_type": "none"},
resolution={"strategy": "none"},
geocoding={"enabled": False},
enrichment={"enabled": False},
)
Complete core_memory_recipes.py
core_memory_recipes.py"""Independent Bolt how-to recipes; each command creates a distinct test dataset."""
import argparse
import asyncio
from datetime import datetime, timedelta, timezone
from uuid import uuid4
from core_memory_settings import settings
from neo4j_agent_memory import MemoryClient
from neo4j_agent_memory.core.memory import ToolCallStatus
from neo4j_agent_memory.memory.long_term import DeduplicationConfig, LongTermMemory
from neo4j_agent_memory.schema.models import EntityRef, TraceOutcome
# tag::messages[]
async def messages(client, run_id):
session = f"docs-messages-{run_id}"
# Explicit timestamps record when each turn happened; untimed rows are stamped in list order.
start = datetime.now(timezone.utc)
stored = await client.short_term.add_messages_batch(
session,
[
{
"role": "user",
"content": "Please find wide walking shoes.",
"metadata": {"topic": "shopping"},
"timestamp": start.isoformat(),
},
{
"role": "assistant",
"content": "I will look for wide fits.",
"timestamp": (start + timedelta(seconds=1)).isoformat(),
},
],
user_identifier=f"docs-user-{run_id}",
extract_entities=False,
)
conversation = await client.short_term.get_conversation(session)
assert [m.id for m in conversation.messages] == [m.id for m in stored]
summary = await client.short_term.get_conversation_summary(
session,
include_entities=False,
summarizer=lambda _transcript: "A shopper requested wide walking shoes.",
)
assert summary.message_count == 2
print(summary.summary)
# This semantic search is database-wide; it is not used for user-scoped readback.
matches = await client.short_term.search_messages("wide walking shoes", threshold=0.0, limit=5)
for message in matches:
print(message.content, message.metadata.get("similarity"))
print(f"Verified: ordered message IDs and summary count; session={session}")
# end::messages[]
# tag::entities[]
async def entities(client, run_id):
person, _ = await client.long_term.add_entity(
f"Maya {run_id}",
"PERSON",
aliases=[f"M. {run_id}"],
description="An engineer at the fictional Northstar laboratory.",
attributes={"team": "robotics"},
resolve=False,
deduplicate=False,
)
company, _ = await client.long_term.add_entity(
f"Northstar {run_id}", "ORGANIZATION", resolve=False, deduplicate=False
)
edge = await client.long_term.add_relationship(person, company, "WORKS_AT", confidence=1.0)
fact = await client.long_term.add_fact(person.name, "works_at", company.name)
assert fact.as_triple == (person.name, "works_at", company.name)
rows = await client.query.cypher(
"MATCH (a:Entity {id: $source})-[r:RELATED_TO {id: $edge}]->(b:Entity) "
"RETURN a.name AS source, r.type AS relation, b.name AS target",
{"source": str(person.id), "edge": str(edge.id)},
)
assert rows == [{"source": person.name, "relation": "WORKS_AT", "target": company.name}]
facts = await client.long_term.get_facts_about(person.name)
assert any(item.as_triple == fact.as_triple for item in facts)
print(f"Verified: entity relationship and fact readback; person={person.id}")
# end::entities[]
# tag::preferences[]
async def preferences(client, run_id):
user = f"docs-preference-user-{run_id}"
product, _ = await client.long_term.add_entity(
f"Walking shoes {run_id}",
"OBJECT",
subtype="PRODUCT",
generate_embedding=False,
resolve=False,
deduplicate=False,
)
scope = EntityRef(id=str(product.id))
# Disable embedding-based global preference dedupe for independent user revisions.
old = await client.long_term.add_preference(
"budget",
"Spend at most USD 60",
user_identifier=user,
applies_to=[scope],
generate_embedding=False,
context="Explicit user statement",
)
replacement = await client.long_term.add_preference(
"budget",
"Spend at most USD 80",
user_identifier=user,
applies_to=[scope],
generate_embedding=False,
context="Explicit user revision",
)
await client.long_term.supersede_preference(old.id, replacement.id)
active = await client.long_term.get_preferences_for(user, applies_to=scope)
historical = await client.long_term.get_preferences_for(user, active_only=False)
assert [p.id for p in active] == [replacement.id]
assert {p.id for p in historical} == {old.id, replacement.id}
context = "\n".join(f"{p.category}: {p.preference}" for p in active)
print(context)
print(f"Verified: current preference and retained history; user={user}")
# end::preferences[]
# tag::reasoning[]
async def reasoning(client, run_id):
session = f"docs-trace-{run_id}"
trace = await client.reasoning.start_trace(session, "Check a deliberately unavailable product")
step = await client.reasoning.add_step(trace.id, action="lookup_product")
try:
raise LookupError("The fictional product is unavailable")
except LookupError as error:
await client.reasoning.record_tool_call(
step.id,
"lookup_product",
{"sku": "missing-example-sku"},
status=ToolCallStatus.ERROR,
error=str(error),
duration_ms=0,
)
await client.reasoning.complete_trace(
trace.id,
outcome=TraceOutcome(success=False, summary=str(error), error_kind="no_results"),
)
saved = await client.reasoning.get_trace_with_steps(trace.id)
assert saved is not None and saved.success is False
assert len(saved.steps) == 1 and len(saved.steps[0].tool_calls) == 1
assert saved.steps[0].tool_calls[0].status == ToolCallStatus.ERROR
print(f"Verified: failed tool call and completed failure trace; trace={trace.id}")
# end::reasoning[]
# tag::deduplication[]
async def deduplication(client, run_id):
store = LongTermMemory(
client.graph,
embedder=client.long_term.embedder,
# The same bands MemoryClient applies from ResolutionConfig (0.90 merge, 0.85 review).
deduplication=DeduplicationConfig(
auto_merge_threshold=0.90, flag_threshold=0.85, use_fuzzy_matching=False
),
)
# Create an isolated pair without similarity-based merging so manual review is reproducible.
target, _ = await store.add_entity(
f"Northstar Laboratory {run_id}", "ORGANIZATION", resolve=False, deduplicate=False
)
source, _ = await store.add_entity(
f"Northstar Lab {run_id}", "ORGANIZATION", resolve=False, deduplicate=False
)
# The application has established that this specific pair denotes the same organization.
merged = await store.merge_duplicate_entities(source.id, target.id)
assert merged is not None
rows = await client.query.cypher(
"MATCH (s:Entity {id: $source}), (t:Entity {id: $target}) "
"RETURN s.merged_into AS merged_into, t.aliases AS aliases",
{"source": str(source.id), "target": str(target.id)},
)
assert rows[0]["merged_into"] == str(target.id)
assert source.name in rows[0]["aliases"]
stats = await store.get_deduplication_stats()
print(f"Verified: reviewed pair merged into {target.id}; merged nodes={stats.merged_entities}")
# end::deduplication[]
# tag::audit[]
async def audit(client, run_id):
session = f"docs-audit-{run_id}"
client_name = f"Anthem {run_id}"
consultant_name = f"Sara {run_id}"
trace = await client.reasoning.start_trace(session, "Recommend a consulting team")
step = await client.reasoning.add_step(trace.id, action="recommend_team")
await client.reasoning.record_tool_call(
step.id,
tool_name="recommend_team",
arguments={"client_name": client_name},
result=[{"consultant": consultant_name}],
touched_entities=[
EntityRef(name=client_name, type="CLIENT"),
EntityRef(name=consultant_name, type="PERSON"),
],
)
await client.reasoning.complete_trace(
trace.id,
outcome=TraceOutcome(success=True, summary="Matched one consultant"),
)
rows = await client.query.cypher(
"MATCH (:Entity {name: $client_name})<-[:TOUCHED]-(s:ReasoningStep)"
"<-[:HAS_STEP]-(rt:ReasoningTrace) "
"RETURN rt.task AS task, s.action AS action, rt.outcome AS outcome",
{"client_name": client_name},
)
assert rows == [
{"task": trace.task, "action": "recommend_team", "outcome": "Matched one consultant"}
]
print(f"Verified: touched-entity audit query found the trace; client={client_name}")
# end::audit[]
async def main():
commands = {
"messages": messages,
"entities": entities,
"preferences": preferences,
"reasoning": reasoning,
"deduplication": deduplication,
"audit": audit,
}
parser = argparse.ArgumentParser()
parser.add_argument("command", choices=commands)
args = parser.parse_args()
async with MemoryClient(settings()) as client:
await commands[args.command](client, uuid4().hex[:8])
if __name__ == "__main__":
asyncio.run(main())
1. Apply the task function
Configure custom deduplication thresholds and merge a reviewed duplicate pair. On a MemoryClient, add_entity takes its bands from MemorySettings.resolution: auto_merge_threshold (default 0.90) and review_threshold (default 0.85), the same values message ingestion uses. MemoryClient does not accept a deduplication= argument, and there is no MemorySettings.deduplication section; NAM_DEDUPLICATION__* environment variables are ignored. DeduplicationConfig, imported from neo4j_agent_memory.memory.long_term, applies when you construct LongTermMemory yourself, as this recipe does with the connected Bolt graph and embedder; DeduplicationConfig.from_resolution_config(settings.resolution) derives it from the client’s bands. The example uses a deliberately created pair and bypasses automatic deduplication so the manual review step is reproducible.
The function below is included from the complete maintained program, which already supplies imports, client setup, unique IDs, and asyncio.run. Run the file in step 2; this extract is not a separate top-level script.
async def deduplication(client, run_id):
store = LongTermMemory(
client.graph,
embedder=client.long_term.embedder,
# The same bands MemoryClient applies from ResolutionConfig (0.90 merge, 0.85 review).
deduplication=DeduplicationConfig(
auto_merge_threshold=0.90, flag_threshold=0.85, use_fuzzy_matching=False
),
)
# Create an isolated pair without similarity-based merging so manual review is reproducible.
target, _ = await store.add_entity(
f"Northstar Laboratory {run_id}", "ORGANIZATION", resolve=False, deduplicate=False
)
source, _ = await store.add_entity(
f"Northstar Lab {run_id}", "ORGANIZATION", resolve=False, deduplicate=False
)
# The application has established that this specific pair denotes the same organization.
merged = await store.merge_duplicate_entities(source.id, target.id)
assert merged is not None
rows = await client.query.cypher(
"MATCH (s:Entity {id: $source}), (t:Entity {id: $target}) "
"RETURN s.merged_into AS merged_into, t.aliases AS aliases",
{"source": str(source.id), "target": str(target.id)},
)
assert rows[0]["merged_into"] == str(target.id)
assert source.name in rows[0]["aliases"]
stats = await store.get_deduplication_stats()
print(f"Verified: reviewed pair merged into {target.id}; merged nodes={stats.merged_entities}")
2. Run the recipe
python core_memory_recipes.py deduplication
The command retains its example records for inspection. Remove the disposable tutorial database using that tutorial’s cleanup command when finished. Do not substitute an unrestricted delete query against an existing database.
Review semantic candidates
To let new writes use the configured policy, leave deduplicate=True on store.add_entity. Inspect the returned DeduplicationResult.action ("none", "merged" or "flagged"), matched_entity_id, and similarity_score. Default thresholds are 0.90 for automatic merging and 0.85 for review, in both ResolutionConfig and DeduplicationConfig; their suitability depends on the corpus and embedding model. When an ontology declares resolution_threshold or review_threshold on an entity type, OntologyResolver uses those bands for that type.
find_potential_duplicates(limit=…) returns (entity1, entity2, confidence) tuples for pending pairs, highest confidence first. In 0.7.0 each pair appears once, with the flagged (newer) entity first and its existing match second, and confidence is the SAME_AS edge’s stored score. Pairs flagged during message ingestion join the same queue. After an application review, call review_duplicate(source_id, target_id, confirm=True) to merge or confirm=False to reject the flag. An empty review queue is a valid result. get_deduplication_stats() returns a dataclass with total_entities, merged_entities, same_as_relationships, and pending_reviews.
Account for resolution during ingestion
Since 0.7.0, storing a message resolves the mentions extracted from it against the entities already in the graph, using the same resolver and bands as add_entity. Each mention has one of three outcomes:
-
Merged: no new node. The message’s
MENTIONSedge points at the matched entity, and the mention’s surface form is appended to that entity’saliases. -
Review: a new node plus a pending
SAME_ASedge to the matched entity, whichfind_potential_duplicates()returns. -
Created: a new node, as before 0.7.0.
The decision is recorded in each new node’s metadata as resolution.action, resolution.score and resolution.match_type. A merged mention is therefore not a pair for this page’s manual merge; inspect the survivor’s aliases instead. To restore the earlier behavior, one node per distinct surface form, pass resolution={"resolve_on_ingest": False} in MemorySettings. add_entity keeps deduplicating either way. See Tune entity resolution.
Check relationship coverage before merging
In SDK 0.7.0, the merge query adds the source name as a target alias and copies the source’s edges of these types onto the target — inbound MENTIONS, RELATED_TO in both directions, SAME_AS, outbound EXTRACTED_FROM and EXTRACTED_BY provenance, inbound APPLIES_TO, and inbound TOUCHED — then marks the source with merged_into/merged_at. It copies rather than deletes, so the source keeps its own edges. The copies are marked differently by type: RELATED_TO copies keep the source edge’s id and gain a migrated_from property set to the source entity ID; EXTRACTED_FROM, EXTRACTED_BY, APPLIES_TO, and TOUCHED copies gain migrated_from only; MENTIONS copies carry no marker, and SAME_AS copies are recorded with match_type: 'merged' but no source reference. To audit or undo a merge, rely on the source’s merged_into marker and the retained source edges rather than on a marker on every copy. Relationship types outside that list are not transferred automatically; inspect those edges and decide how your application will reconcile them before merging a populated graph.
Known aliases, type consistency, and a labeled sample of true/false duplicates help assess a policy. Do not choose automatic thresholds from an unmeasured "safe" universal value.
See also
-
Resolution and deduplication — the distinction between the two mechanisms.
-
The long-term API reference — configuration and return models.
-
Tune entity resolution — the bands, alias gazetteer and per-type overrides that both write paths use.