Score delayed receipt events against explicit state cutoffs, restore a checkpoint and reconcile the decision sink.
Project: recover a receipt stream scorer after offline uploads
Freeze the business contract
A handheld uploads receipts after a store reconnects. The live scorer uses a seven-minute merchant window and a fixed receipt-risk model. Define whether each receipt reads merchant history up to its occurrence time or to its scoring time; this project chooses occurrence time. Pin the model, preprocessing, threshold policy and state schema. An event older than the allowed delay goes to reconciliation, not an unmarked live route. The event-time contract is the test oracle.
Inject disorder and repetition
Produce 470 receipts with known event IDs, then reorder arrival, replay 47 uploads and delay one source partition. Record the watermark and latest state cutoff for each score. Compare outputs to a reference run in occurrence order. A duplicate upload must not add twice to merchant state or emit a second customer-facing decision. The delayed partition must either remain within the declared late tolerance or generate a reconciliation record. Feature parity is checked at the state cutoff, not merely on the final aggregate.
Crash at the difficult point
Stop the scorer after an output reaches the sink but before its acknowledgment returns. Restore the last checkpoint and replay source events. Compare expected event-and-revision keys with published keys, detect duplicate rows and verify any conflicting content is rejected. If the model pointer moved during downtime, the replay still uses its pinned revision. The recovery ledger makes the gap and duplicate visible.
Deliver the incident packet
Publish input ranges, state and model revisions, duplicate and late counts, output gaps, repaired decisions and explicit exceptions. Check that the reconciled sink has one active decision per eligible key; a deliberate correction carries a new scoring revision and a predecessor link. Test backlog recovery time and memory before declaring the system production ready. Incident evidence should show when the store disconnected, not just when uploads arrived.
Implementation
def receipt_stream_release(events, output_keys, late_ids):
eligible = {event["event_id"] for event in events
if event["event_id"] not in late_ids}
if len(output_keys) != len(set(output_keys)):
return "hold:duplicate-output"
if set(output_keys) != eligible:
return "hold:missing-or-extra-output"
return "accept:reconciled"
events = [{"event_id": "receipt-47"}, {"event_id": "receipt-82"},
{"event_id": "receipt-93"}]
assert receipt_stream_release(events, ["receipt-47", "receipt-82"],
{"receipt-93"}) == "accept:reconciled"
assert receipt_stream_release(events, ["receipt-47", "receipt-47"],
{"receipt-93"}) == "hold:duplicate-output"
Performance and operating cost
The set comparison is O(e + o) expected time and O(e + o) memory for events and outputs; bound the reconciliation window in production. Retaining a longer event-time window increases keyed state and checkpoint bytes. An offline surge raises backlog and catch-up compute even if steady traffic is small, so size recovery from the surge rather than the daily average.
Common Mistakes
- Scoring delayed receipts against whatever state exists at upload time.
- Counting repeated uploads as independent merchant activity.
- Restoring state without reconciling already published decisions.
- Allowing an old event to pick up a newly promoted model during replay.
Read next
- Stream inference: event time, feature state and decision identity
- Stream scoring replay: checkpoints, duplicates and output reconciliation
- Training-serving parity: compare feature values at one prediction clock
- Model incidents: build a release and evidence timeline before rollback
- Batch inference: make partitions idempotent and outputs identifiable
