Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Session revisions and sink retractions

Last updated: 6 Oct 20265 min read
tutorial
AdvancedBy AITrove Editorial

A session result is a changing claim over event membership, so replay and late merging need explicit supersession records at the sink.

Separate event identity from result identity

Deduplicate source events by immutable event ID before grouping. A result ID derived from sorted member IDs and a rule version identifies one exact session revision, but adding or bridging an event changes that ID. Store the prior published IDs for the same account and interval, then replace them as one logical commit. Event envelopes define source revisions and tombstones upstream.

Make the correction atomic to readers

When a late bridge merges sessions A and B into C, publish deletion markers for A and B and the upsert for C in one transaction or one candidate generation. If the storage system cannot make that group atomic, expose a generation pointer only after all three changes are durable. Readers should never see A, B and C as simultaneously current. An audit stream may retain every revision, but serving queries must filter to current rows.

Control duplicate trigger firings

An early trigger, an event-time trigger and a late firing may emit the same membership more than once. The sink should key by a deterministic result identity or by a stable logical key plus revision, and an attempt should carry a monotonic generation. Replaying the same checkpoint and source offsets must leave the same current table. A new event legitimately changes membership and must not be suppressed by a dedupe key based only on account ID.

Repair expired sessions elsewhere

Once a session has passed its accepted lateness and state has been cleaned, a historical correction cannot safely pass through the current live window. Route the affected account and interval to a bounded rebuild from raw events. Compute a candidate set, compare counts and amounts with the published generation, then switch readers. Atomic backfill publication avoids mixed old and new results.

Prove accounting invariants

For every input generation, each unique source event belongs to exactly one current session per key. Current session totals must sum to the accepted unique event amounts, even after a bridge and replay. Validate that no deleted revision remains current, and preserve the mapping from superseded IDs to the replacement ID. These checks catch a sink that accepts the new row but loses one of the two retractions.

Implementation

python
import hashlib

def session_revision(account_id, event_ids, rule_version):
    members = sorted(set(event_ids))
    payload = "|".join([account_id, rule_version] +
                       [f"{len(event_id)}:{event_id}" for event_id in members])
    return hashlib.sha256(payload.encode()).hexdigest()

first_id = session_revision("account-29", ["activity-47"], "gap-11-v2")
second_id = session_revision("account-29", ["activity-83"], "gap-11-v2")
merged_id = session_revision("account-29",
                             ["activity-47", "activity-61", "activity-83"],
                             "gap-11-v2")
current = {first_id: 74, second_id: 51}
for stale_id in (first_id, second_id):
    current.pop(stale_id)
current[merged_id] = 134
assert current == {merged_id: 134}
assert session_revision("account-29",
                        ["activity-83", "activity-61", "activity-47"],
                        "gap-11-v2") == merged_id

Performance and operating cost

Building an ID from M member IDs costs O(M log M) time for sorting and O(M) temporary space. A production operator may maintain incremental membership state, but the sink still needs merge or transaction support and an index on current IDs. Audit history grows with every revision; bound its retention separately from the live serving table.

Common Mistakes

  • Do not use account ID alone as a dedupe key for all sessions.
  • Do not expose the replacement before the obsolete results are retired.
  • Do not mistake repeat trigger firings for distinct activity.

Read next

ai-data
data-engineering
Storage details