Chapter 10 · Atomicity Across Business Workflows: Counters, Reservations, Idempotency, and Event-Driven Consistency
Design a Saga-Like Firestore Workflow with Recovery, Replay, Monitoring, and Manual Repair Paths
Assemble AtlasMart checkout as a saga-like Firestore workflow with recovery scanning, replay, monitoring, dead-letter evidence and manual repair.
Learning outcomes
Design a saga-like AtlasMart order workflow whose transitions, compensations and recovery paths are explicit and testable.
Preserve outbox/inbox, dead-letter and repair evidence so automated replay and human intervention are both safe.
Build a deterministic replay tool that is idempotent and cannot silently skip invariants.
Define workflow health metrics without fabricating production latency, retry or cost values from local emulator tests.
Use the Emulator Suite, a Firebase demo project, or an isolated test project for destructive, security-sensitive, billing-sensitive, migration, backup/restore, or write-heavy exercises unless the lesson explicitly marks managed verification as required. Treat shown output as expected evidence unless it is explicitly identified as captured output, and re-check current Firebase/Google Cloud edition, mode, quota, pricing, and security documentation before production execution.
AtlasMart continues the same environment used in Chapters
01–09: project ID demo-atlasmart-firestore,
Standard edition / Native mode /
(default) database for mandatory labs, Firestore
emulator 127.0.0.1:8080, Authentication emulator
127.0.0.1:9099, Emulator UI
127.0.0.1:4000, Firebase CLI
15.30.0, Firebase JavaScript SDK
12.19.0, Firebase Admin Node SDK
14.4.0 with
@google-cloud/firestore 9.1.0, and Node.js 22+.
Mandatory work remains local/no-cost. Cloud Functions,
Eventarc, managed TTL deletion, production IAM, billing,
regional delivery latency and external payment systems are
discussed accurately but are not falsely claimed to have run
in the local emulator.
The local lab simulates duplicate and reordered events deterministically with ordinary Node code so the learner can prove idempotency, compensation and repair behavior without deploying cloud infrastructure. Firestore-triggered Cloud Functions and Eventarc Standard can deliver events at least once; Firestore event ordering is not guaranteed. Firestore TTL deletion is asynchronous and documents are typically removed within about 24 hours after expiration, so TTL is a retention mechanism—not an exact reservation scheduler. Any production p95/p99, event-delivery delay, TTL cleanup delay, trigger retry count or cost must be measured in the actual edition/region/billing configuration rather than inferred from emulator timing.
1. The AtlasMart problem: failures happen between valid commits
At the end of Chapter 9, AtlasMart could protect one transactional invariant. By the end of Chapter 10 it needs to survive a crash after inventory reservation but before payment processing, a duplicate payment event, an expired reservation, a compensation failure, and an operator replay. These are not exceptional edge cases; they are the normal failure surface of a multi-step workflow.
A saga-like workflow decomposes one business process into committed steps with compensating actions for steps that must be reversed. Firestore does not provide a built-in saga coordinator; AtlasMart stores its own durable state machine and evidence. The design goal is convergence to a valid terminal state, not pretending all steps were one global transaction.
| State | Meaning | Allowed next states | Durable evidence |
|---|---|---|---|
| NEW | Idempotency command accepted; no inventory reserved yet | RESERVED, REJECTED | workflowCommands/{key} |
| RESERVED | Inventory is held by a reservation document | PAYMENT_PENDING, COMPENSATING | reservations/{orderId} + order state |
| PAYMENT_PENDING | External payment placeholder/event expected | COMPLETED, COMPENSATING, MANUAL_REVIEW | workflowOutbox + workflowEvents |
| COMPLETED | Order is terminally successful | none except explicit administrative correction | orders/{id} terminal state |
| COMPENSATING | Workflow is releasing inventory or reversing downstream effects | CANCELLED, MANUAL_REVIEW | compensationAttempt, reason, correlationId |
| CANCELLED | Compensation completed | none | terminal order + released reservation |
| MANUAL_REVIEW | Automatic progress is unsafe or repeatedly failed | operator-defined repair transition | dead-letter/repair document |
2. Persist the workflow timeline
process.env.FIRESTORE_EMULATOR_HOST = "127.0.0.1:8080";process.env.GCLOUD_PROJECT = "demo-atlasmart-firestore";import { initializeApp } from "firebase-admin/app";import { getFirestore, FieldValue, Timestamp } from "firebase-admin/firestore";initializeApp({ projectId: "demo-atlasmart-firestore" });const db = getFirestore();
async function appendWorkflowEvent({orderId,eventId,type,fromState,toState,correlationId,details={}}) { await db.doc(`workflowEvents/${eventId}`).create({ orderId,eventId,type,fromState,toState,correlationId,details, schemaVersion:3,recordedAt:FieldValue.serverTimestamp() });}
For replay and incident response, do not rely only on the current order document. The current state answers “where are we?”; the event/attempt evidence answers “how did we get here, what already ran, and what can be retried safely?”
3. One transition function owns allowed state movement
const allowed={ NEW:new Set(["RESERVED","REJECTED"]), RESERVED:new Set(["PAYMENT_PENDING","COMPENSATING"]), PAYMENT_PENDING:new Set(["COMPLETED","COMPENSATING","MANUAL_REVIEW"]), COMPENSATING:new Set(["CANCELLED","MANUAL_REVIEW"]), COMPLETED:new Set(), CANCELLED:new Set(), REJECTED:new Set(), MANUAL_REVIEW:new Set()};async function transitionOrder(orderId,toState,{eventId,correlationId,details={}}) { const orderRef=db.doc(`orders/${orderId}`); const eventRef=db.doc(`workflowEvents/${eventId}`); return db.runTransaction(async tx=>{ const [order,event]=await Promise.all([tx.get(orderRef),tx.get(eventRef)]); if (event.exists) return {action:"duplicate",state:order.get("state")}; if (!order.exists) throw new Error("ORDER_NOT_FOUND"); const from=order.get("state"); if (!allowed[from]?.has(toState)) throw new Error(`INVALID_TRANSITION_${from}_TO_${toState}`); tx.update(orderRef,{state:toState,lastEventId:eventId,updatedAt:FieldValue.serverTimestamp()}); tx.create(eventRef,{orderId,eventId,type:"STATE_TRANSITION",fromState:from,toState,correlationId,details,schemaVersion:3,recordedAt:FieldValue.serverTimestamp()}); return {action:"transitioned",from,to:toState}; });}
The same event cannot move state twice because its event record is created atomically with the transition. A different event that proposes an invalid movement fails and is parked for review.
4. Recovery scanner: find workflows that stopped progressing
async function findStalledOrders(cutoff) { // Production requires the appropriate index; keep the lab fixture bounded. const snap=await db.collection("orders") .where("state","in",["RESERVED","PAYMENT_PENDING","COMPENSATING"]) .where("updatedAt","<",cutoff) .limit(100) .get(); return snap.docs.map(d=>({id:d.id,...d.data()}));}
A scanner is not itself a repair. Each candidate goes through state-specific logic: requeue an outbox event, expire a reservation, retry compensation, or escalate to manual review. The scanner records a repair attempt ID so repeated scans are observable and idempotent.
5. Replay only from durable evidence
async function requestReplay({deadLetterId,operator,reason}) { const dlRef=db.doc(`workflowDeadLetters/${deadLetterId}`); const replayRef=db.doc(`workflowOutbox/replay-${deadLetterId}`); await db.runTransaction(async tx=>{ const dl=await tx.get(dlRef); if (!dl.exists) throw new Error("DEAD_LETTER_NOT_FOUND"); if (dl.get("state")!=="NEEDS_REVIEW") return; tx.set(replayRef,{ type:"REPLAY",sourceDeadLetter:deadLetterId, originalEventId:dl.get("id"),orderId:dl.get("orderId"), state:"READY",operator,reason,correlationId:dl.get("correlationId"), createdAt:FieldValue.serverTimestamp(),schemaVersion:3 },{merge:false}); tx.update(dlRef,{state:"REPLAY_REQUESTED",replayRequestedBy:operator,replayRequestedAt:FieldValue.serverTimestamp()}); });}
Manual repair should never mean “edit the order in the console until it looks right.” A replay request is itself durable, attributable and subject to the same transition rules. For truly exceptional corrections, preserve the before/after state, operator identity and reason.
6. Failure injection matrix
| Injected failure | Expected durable state | Safe recovery |
|---|---|---|
| Crash after reservation transaction, before outbox worker runs | RESERVED + READY outbox record | Worker/recovery scanner processes existing outbox |
| Payment-success event delivered twice | One COMPLETED transition + dedupe evidence | Duplicate returns no-op |
| Payment failure after reservation | COMPENSATING then CANCELLED; inventory released once | Retry compensation transaction |
| Compensation throws repeatedly | MANUAL_REVIEW + dead-letter/attempt evidence | Operator replay/repair |
| Reservation passes expiresAt before payment | HELD reservation becomes RELEASED; order CANCELLED if still eligible | Explicit expiry sweeper; TTL may clean later |
| Out-of-order stale event arrives after COMPLETED | No state regression; event parked/rejected | Review/ignore according to event version/state policy |
7. Observability: monitor the workflow, not only Firestore API errors
Useful workflow signals include counts by state, age of the
oldest non-terminal order, READY/PROCESSING outbox backlog,
duplicate-event count, compensation attempts, dead-letter
backlog, replay success/failure and reservation age past
expiresAt. Add correlationId,
orderId, eventId,
idempotencyKey, attempt number and state transition
to structured logs. These application metrics explain business
stuckness that raw Firestore request latency alone cannot.
The local emulator lab can prove transition correctness and produce your own sample durations, but it cannot establish production event latency, p95/p99 workflow completion, TTL delay or monthly cost. Define production objectives only after measuring the deployed architecture in its actual region/edition/event stack.
8. Reproducible end-to-end lab
{ "name": "atlasmart-firestore-ch10", "private": true, "type": "module", "engines": { "node": ">=22" }, "dependencies": { "firebase-admin": "14.4.0" }, "devDependencies": { "firebase-tools": "15.30.0" }}
{ "firestore": { "rules": "firestore.rules", "indexes": "firestore.indexes.json" }, "emulators": { "firestore": { "port": 8080 }, "auth": { "port": 9099 }, "ui": { "enabled": true, "port": 4000 } }}
rules_version = '2';service cloud.firestore { match /databases/{database}/documents { match /catalogItems/{productId} { allow read: if true; allow write: if false; } match /profiles/{uid} { allow read, write: if request.auth != null && request.auth.uid == uid; } match /orders/{orderId} { allow read: if request.auth != null && resource.data.customerId == request.auth.uid; allow write: if false; } // Workflow state, reservations, outbox/inbox, dedupe and repair evidence are server-owned. match /workflowCommands/{id} { allow read, write: if false; } match /reservations/{id} { allow read, write: if false; } match /workflowEvents/{id} { allow read, write: if false; } match /workflowOutbox/{id} { allow read, write: if false; } match /workflowDeadLetters/{id} { allow read, write: if false; } match /counters/{counterId}/{document=**} { allow read, write: if false; } match /{document=**} { allow read, write: if false; } }}
mkdir atlasmart-firestore-ch10 && cd atlasmart-firestore-ch10npm init -ynpm install firebase-admin@14.4.0npm install --save-dev firebase-tools@15.30.0# Save firebase.json, firestore.rules and firestore.indexes.json from this lesson.printf '{"indexes":[],"fieldOverrides":[]}' > firestore.indexes.jsonnpx firebase-tools@15.30.0 emulators:start --project demo-atlasmart-firestore --only firestore,auth
Create three orders with deterministic IDs:
-
ord-saga-success: reserve → payment pending → duplicate success event → completed. -
ord-saga-fail: reserve → payment pending → payment failure → compensation → cancelled; replay the failure and verify no second stock release. -
ord-saga-stuck: reserve → mark an outbox event failed repeatedly → park in manual review → create a replay request.
async function assertState(path,expected) { const s=await db.doc(path).get(); if (!s.exists || s.get("state")!==expected) throw new Error(`${path}: expected ${expected}, got ${s.get("state")}`);}await assertState("orders/ord-saga-success","COMPLETED");await assertState("orders/ord-saga-fail","CANCELLED");await assertState("orders/ord-saga-stuck","MANUAL_REVIEW");const successEvents=await db.collection("workflowEvents").where("orderId","==","ord-saga-success").get();const dead=await db.collection("workflowDeadLetters").where("orderId","==","ord-saga-stuck").get();console.log({successEventCount:successEvents.size,deadLetterCount:dead.size});
Cleanup should delete only the demo fixture after verification. In production, retention and audit requirements may forbid immediate deletion; Chapter 11 will separate hierarchy deletion, archival and lifecycle policy.
Production judgment
A reliable workflow is one that can explain and recover from partial progress. Keep the strict invariant in the smallest transaction, make asynchronous steps idempotent, record state changes and attempts, compensate intentionally, and preserve a human repair path. Do not use TTL as a scheduler, do not delete failure evidence prematurely, and do not claim exactly-once processing simply because duplicate tests happened to pass.
Chapter 11 builds on this state ownership by examining subcollections, collection groups, recursive deletion and lifecycle—important because workflow evidence itself eventually needs retention and deletion rules.
Knowledge check
- What makes a saga-like workflow different from one transaction?
- Why keep a transition history in addition to current state?
- What should a recovery scanner do with stalled workflows?
- Why is direct console editing a poor repair procedure?
- Which metrics reveal workflow health?
Review the answers
1. It consists of multiple committed steps over time, with explicit compensation/recovery rather than one global rollback.
2. It supports deduplication, incident diagnosis, replay decisions, audit and manual repair.
3. Classify them by state and invoke idempotent state-specific recovery; it should not blindly mutate them.
4. It bypasses state-transition rules and leaves weak or no durable evidence of who changed what and why.
5. State counts/age, outbox backlog, duplicates, compensation attempts, dead letters, replay outcomes and expired-reservation lag, with correlation IDs.
Summary
Chapter 10 extends atomic correctness into operational correctness. AtlasMart now has explicit atomic boundaries, sharded telemetry, strict reservations, duplicate-safe event processing and a saga-like recovery model with human repair. The next chapter applies the same ownership discipline to Firestore hierarchy and data lifecycle.
Authoritative references
- Transactions and batched writes — atomic transaction/batch boundaries and retry behavior.
- Distributed counters — shard-based write distribution and read aggregation tradeoffs.
- Manage data retention with TTL policies — asynchronous deletion behavior, limits, pricing and monitoring.
- Cloud Firestore triggers — at-least-once delivery, non-guaranteed ordering, trigger scope and idempotency requirement.
- Eventarc Standard retry events — at-least-once delivery, duplicate handling, idempotency and dead-letter guidance.
- Firestore best practices — hotspot and scaling guidance.
- Firestore Enterprise overview — edition/mode boundaries to re-check before porting workflow assumptions.
- Firebase release notes — current SDK/tool versions.
- Admin Node.js release notes — 14.4.0 and Firestore client dependency baseline.
- Firebase CLI release notes — 15.30.0 baseline.