Skip to content
All posts

Turning Two Agents' Event Streams Into One Canonical Graph

Diagram of two durable event streams, sales-bot and support-bot, each ticking through offsets (sales-bot crashes and resumes without duplication), converging into one canonical graph entry for Jane Doe with four flagged property conflicts

Imagine a sales bot and a support bot that both talk to the same person without knowing it. The sales bot knows her as “Jane Doe”, VP Engineering, at jane@acme.com, while the support bot has “J. Doe”, VP Eng, at the same email, along with two companies the sales bot has never heard of, one of which the support bot retracts a day later. Both bots emit durable, at-least-once event streams of what they’ve seen.

Event logs are great at delivery, ordering, and replay, but they don’t resolve entities, and you can’t ask a Kafka topic what an agent believed last Tuesday. That work happens in a layer between the log and whatever reads it, which is what Materializing Event Logs describes. @nicia-ai/agent-stream-graph is the reference implementation, built on 0.35’s store.transactionWithReceipt() and a couple of 0.36 additions. It pulls together three things I’d built separately (bitemporal history, idempotent writes, and graph merge), and seeing them click into one pipeline was a lot of fun.

Streams redeliver changes after a crash, a reconnect, or a replay, so the second delivery of a change has to land on the same row as the first. In practice that means upsertById for nodes and getOrCreateByEndpoints for edges, and a bare create only when the source event carries its own unique id that you pass through as the TypeGraph id.

const project: Projector<typeof intelGraph> = async (belief, change) => {
switch (change.shape) {
case "person": {
if (change.operation === "delete") {
await belief.nodes.Person.delete(asNodeId(change.key));
return;
}
await belief.nodes.Person.upsertById(change.key, {
name: change.value.name,
email: change.value.email,
title: change.value.title,
});
return;
}
case "employment": {
await belief.edges.worksAt.getOrCreateByEndpoints(
{ kind: "Person", id: change.value.person },
{ kind: "Company", id: change.value.company },
{},
);
return;
}
}
};

consume() checkpoints the last processed offset and a recorded-time anchor after every change, and uses store.transactionWithReceipt() to know whether the projector actually wrote anything. The demo runs the sales bot’s stream, stops it after two of its four changes, then simulates the nastiest crash window there is: a change gets projected, and the process dies before the checkpoint lands.

(b) Durable consumer — resume from checkpoint, replay safely after crash
consumer ran, then crashed after 2 messages
durable cursor: last offset = 002
sales-bot belief so far: 1 person, 1 company
crash window: projected 003, then died before checkpoint
uncheckpointed anchor existed: 2026-07-14T16:32:30.397Z
durable cursor is still: 002
belief already has worksAt edges: 1
restarted — replayed 003, then processed 004 (2 messages)
sales-bot belief now: 1 person, 1 company — Jane Doe (VP Eng & Product)
worksAt edges after replay: 1 (no duplicate edge)
re-run (at-least-once): 0 messages processed; belief unchanged: 1 person, 1 company

The cursor still says 002 even though 003 was already projected, which is the crash window the demo sets up deliberately. On restart 003 replays, getOrCreateByEndpoints finds the edge it already wrote, and processing carries on. Running the whole stream a third time processes nothing because the cursor is already at the end.

A full rebuild is a different case: replay every event from the beginning into a fresh belief graph, for recovery or a schema migration. Before 0.36, that rewrote every row even when the replayed value matched what was already there, so each rebuild cost a wasted write per event and grew recorded history by the length of the log.

createStore(graph, backend, { coalesceUnchangedUpserts: true }) makes a value-identical redelivery a true no-op that writes nothing, adds no history row, and doesn’t advance the revision. A stream that actually changes a value and later changes it back still writes both times, since those changes really happened. 0.36 also adds tx.measure(), which scopes a receipt to the writes made by one callback, so a materializer’s own cursor bookkeeping in the same transaction can’t be mistaken for projector output.

Each bot’s belief graph has history enabled, so book.anchorFor(stream, offset) gives you a recorded-time coordinate for any processed offset. That means you can read exactly what this bot believed at that point in its own stream:

(c) What did each agent believe, at which offset?
sales-bot's own belief graph:
@offset 001: people: Jane Doe (VP Engineering) | companies: —
@offset 004: people: Jane Doe (VP Eng & Product) | companies: Acme Corp (Series A) (title corrected)
support-bot's own belief graph (same person, different surface form):
@offset 004: people: J. Doe (VP Eng) | companies: Umbrella (unverified), Acme (Series B), Globex (Series C)
@offset 005: people: J. Doe (VP Eng) | companies: Acme (Series B), Globex (Series C) (Umbrella retracted)
→ Same email, but neither agent alone knows 'Jane Doe' and 'J. Doe'
are one person. The retracted company also remains visible in past belief.

The support bot’s belief at offset 4 still shows “Umbrella (unverified)”, even though it was deleted at offset 5, because a read at a past offset reconstructs the belief as it was then.

Each bot’s belief exports through streaming interchange into a graph merge branch, and mergeIncremental() folds it into a canonical graph. It’s the same entity resolution as in the graph merge post, run once per stream as it arrives:

await importGraphStream(
agentBranch.store,
exportGraphStream(belief, { includeTemporal: true }),
{ onConflict: "update" },
);
const result = await mergeIncremental({
forkPoint,
target: canonical,
branches: [agentBranch],
options: {
resolve: {
Person: {
similarity: { kind: "fulltext", fields: ["name"] },
threshold: 0.9,
},
Company: {
similarity: { kind: "fulltext", fields: ["name"] },
threshold: 0.9,
},
},
onPropertyConflict: "flag",
onBasePropertyConflict: "flag",
branchOrder: [branchId],
persistProvenance: true,
},
});

The sales bot merges first, as a clean append. The support bot merges second, and that’s where the same-person, different-spelling problem gets resolved:

[wave 1] merged sales-bot — conflicts: 0
[wave 2] merged support-bot — conflicts: 4
conflict: Company.name on c1: support-bot="Acme"; kept "Acme Corp"
conflict: Company.stage on c1: support-bot="Series B"; kept "Series A"
conflict: Person.name on p1: support-bot="J. Doe"; kept "Jane Doe"
conflict: Person.title on p1: support-bot="VP Eng"; kept "VP Eng & Product"
canonical now: 1 person, 2 companies — Jane Doe (VP Eng & Product)
provenance — sales-bot contributed to 3 canonical entities
provenance — support-bot contributed to 3 canonical entities

“Jane Doe” and “J. Doe” collapse into one canonical person, and every disagreement is flagged instead of silently overwritten. Globex, which only the support bot ever saw, joins as a new company. Umbrella, which the support bot retracted before the merge, never shows up at all. And the canonical graph has history too, so it time-travels across merge waves:

canonical, time-travelled:
asOfRecorded(after wave 1): 1 person, 1 company
asOfRecorded(after wave 2): 1 person, 2 companies

So from two unreliable streams you end up with one canonical graph, and you can still ask any of the three graphs what it believed at any point along the way.

The demo above is small on purpose: one person, two bots, five offsets. The repo’s main demo (pnpm demo) runs the same mechanics against a set of wiki entries that mention overlapping concepts under different names, and converges them into one canonical concept graph with source attribution for every mention.

Stay in the loop

Occasional updates on new features, guides, and releases. No spam.