Turning Two Agents' Event Streams Into One Canonical Graph

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.
Projecting events idempotently
Section titled “Projecting events idempotently”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; } }};Resuming after a crash
Section titled “Resuming after a crash”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 companyThe 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.
Rebuilding from offset zero
Section titled “Rebuilding from offset zero”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.
What did each bot believe, and when?
Section titled “What did each bot believe, and when?”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.
One canonical graph
Section titled “One canonical graph”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 companiesSo 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.
A bigger example
Section titled “A bigger example”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.
Try it
Section titled “Try it”- Materializing Event Logs: idempotent projectors, cursor bookkeeping, transaction receipts, and mapping event time onto the two temporal axes
- Graph Merge: the entity resolution this builds on
agent-stream-graph:pnpm demoandpnpm demo:mechanics- GitHub
Stay in the loop
Occasional updates on new features, guides, and releases. No spam.