Distributed Transactions: Sagas vs Two-Phase Commit
The Flight Booking That Charged Without a Seat
Test your architecture intuition: Pitch a 7-axis solution, survive two aggressive reviewer objections, and inspect the staff-level Teacher Gold Answer.
Part 0. Start here
The problem: a sold-out hotel that cost $36,260
Wayfare, a travel site, sells trip packages: a flight seat and a few hotel nights, paid by card. Its first version books a trip the obvious way, inside the web request: charge the card, then take the seat, then take the rooms, and if a later step fails, refund the card.
Wayfare runs a one-hour Lisbon sale. 50 trips a minute are booked (our example's number), each $1,240. At minute 40 the hotel sells out. Nobody stops the sale, so for the last 20 minutes every attempt charges the card, takes a seat, fails at the hotel, gives the seat back and refunds.
| What happened | Arithmetic | Result |
|---|---|---|
| Refunds | 20 minutes × 50 trips a minute | 1,000 |
| Fee on each charge (one large PSP's US price, "2.9% + 30¢") | 1,240 × 0.029 = 35.96, + 0.30 | $36.26 |
| Fees kept by the PSP | 1,000 × 36.26 | $36,260 |
| Customers who see $1,240 missing | every one of the 1,000 | for 5 to 10 business days |
| If 1% of refund calls time out and the code treats a timeout as "refund failed, stop" | 1% × 1,000 | up to 10 customers charged for a trip they don't have |
That PSP, Stripe, states both costs plainly in its refund documentation: "Stripe's processing fees from the original transaction aren't returned", and "Your customer sees the refund as a credit approximately 5-10 business days later". The same page offers the way out: "You can cancel a payment before it's completed at no cost." A card payment can be authorized first (the money is held on the card, not taken) and captured later. An uncaptured authorization of an online, customer-initiated card payment stays valid for 7 days. Cancelling it (a void) is free.
"Up to 10" because a refund that timed out may still have happened. Nobody knows until someone asks the PSP.
Put the four steps (the seat, the rooms, the card, the airline ticket) in an order, and say what happens when each one fails, so that a sold-out hotel costs nothing, a crash anywhere never loses track of the trip, and no customer is ever charged for a trip they don't have. Would two-phase commit across the seat, room and trip databases do the job instead?
The big picture
Synthesizing vector architecture diagram...
What to notice: the arrows give the order the orchestrator runs the steps: S1, S2, S3, SF, S4, S6, S7, S5, S8. Every step in "Before the pivot: undoable" has an undo, and SF is the last of them. "The pivot" has none. The steps in "After the pivot: retriable" have no undo either: they are built to succeed when retried, so the trip can only move forward once the ticket exists.
What you'll be able to do after this page
- Say why one database transaction can't cover several databases and partner APIs, and name the two families of answers (Part 1).
- Walk through two-phase commit's prepare, vote, decision and commit, and compute what its lock window costs a hot row (Part 2).
- Say exactly what a coordinator crash blocks, for how long, and why an operator must never guess (Part 3).
- Build a saga from local steps and compensations, and undo it in reverse order (Part 4).
- Order the steps around a pivot, fence the holds before it, and tell the customer at the pivot (Part 5).
- Use holds that expire (try-confirm-cancel), and handle a cancel that arrives before its hold (Part 6).
- Treat a timeout as an unknown, resolve it, and handle a compensation that fails (Part 7).
- Recover a saga after its orchestrator crashes, with a log, keys, a lease and an epoch (Part 8).
- Name the three anomalies other bookings can see, and the countermeasure for each (Part 9).
- Compare orchestration with choreography, including ordering on the bus (Part 10).
- Compare 2PC, sagas, one-store transactions and choreography on equal terms (Part 11).
- Trace one booking through every layer, with the durable record and the guard at each (Part 12).
- Map all of it to AWS, and name the look-alikes (Part 13).
You may have arrived from a step that relies on this: step 2.1 of the payment processing loop (a payment is a saga, and "PSPs don't take part in anyone's 2PC"), step 2.2 of the digital wallet loop (a cross-shard transfer as try-confirm-cancel), steps 2.3 and 3.5 of the hotel reservation loop (fence the hold first, charge second; the refund as a durable step), step 3.1 of the job scheduler loop (a durable workflow engine with compensations) or step 3.4 of the Shopify case study (a void as the compensation, driven from an outbox). This page is the "why" behind all five, and behind the drill The Flight Booking That Charged Without a Seat, whose two questions it answers in full (Parts 5, 7 and 11). In all, 22 interview loops depend on this page's mechanisms in at least one step, 13 of them on sagas, 2PC, holds or compensations directly.
Part 1. One trip, three databases, two partners
The hook's code charged, then took a seat, then took the rooms. Each call committed on its own, so when the third failed, the first two had already happened. This Part says why no single transaction can fix that, and names the two ways out.
Charge, seat, rooms, refund
Synthesizing vector architecture diagram...
What to notice: the first message already moves money, so every failure after it is repaired with a refund. The two steps most likely to fail (the rooms, then the seat) come after the only step that is expensive to undo.
What one transaction covers
A database transaction is atomic because one database writes one log: the commit record is either on disk or not. A second database has its own log, and an HTTP partner has none that you can see. So BEGIN; charge; take seat; take rooms; COMMIT can include none of them. Writing to each in turn is a dual write: a crash between two writes leaves one done and the other not. The Change Streams & the Transactional Outbox loop primitive shows this failure for a database and a broker (Part 1) and fixes it with an outbox; here the second party is another service or a partner, and an outbox alone can't undo a charge.
Two words get mixed up here:
- Atomicity: all of the changes happen or none do.
- Isolation: nobody else sees the changes half-done.
A trip needs atomicity across systems. Whether it also needs isolation is a separate question, answered in Part 9.
Two families
| Family | The idea in one line | What it needs from every participant |
|---|---|---|
| Commit together (two-phase commit, or one store's own multi-item transaction) | Every participant first promises it can commit, holding locks; then all commit or all abort, on one decision | The ability to prepare and wait: our databases can, a PSP or an airline can't |
| Commit each step, and plan the undo (a saga) | Each step is a local transaction that commits at once; if a later step fails, a compensation undoes each earlier one | A step and an undo for each participant, both safe to repeat |
The example we follow
One story runs through the page: one trip booking: flight, hotel, card. Wayfare sells packages from seat and room allotments it owns. At 10:00:00.000 customer c-88 books trip trip-4417: seat 14C on flight FL-212 to Lisbon (240 a night, 1,240** on a card. Three services each own a PostgreSQL database; two partners sit outside:
seatsowns the rowflight_inventory(FL-212, available), one row per flight (every booking onFL-212updates it: the hot row), and a row per seat.roomsowns a row per hotel night with itsavailablecount, and a row per room hold.tripsowns the trip record, the saga log and the outbox; it runs the orchestrator.- The PSP (a partner) authorizes, captures, voids and refunds.
- The airline (a partner) issues tickets against a fare quote, here quote
q-77.
The saga has nine steps, run in the order S1, S2, S3, SF, S4, S6, S7, S5, S8. The numbers are names, not the order; Parts 5 and 6 explain the order. Every step carries a key trip-4417:<step>; every undo carries the step's key plus :undo.
| # | Step | Latency (ours) | Zone | Undo (compensation) |
|---|---|---|---|---|
| S1 | hold_seat: 14C HELD, available − 1, expires in 15 minutes | 15 ms | compensatable | C1 release_seat (available + 1), 15 ms |
| S2 | hold_rooms: 3 nights, available − 1 each, expires in 15 minutes | 20 ms | compensatable | C2 cancel_rooms, 20 ms |
| S3 | authorize_card $1,240 | 400 ms | compensatable | C3 void_authorization, 300 ms, free |
| SF | fence_holds: seat and rooms HELD → CONFIRMING, only if still HELD and unexpired | 20 ms (both in parallel) | compensatable, the last one | C1 and C2 release a CONFIRMING hold too |
| S4 | issue_ticket with quote q-77 | 1,200 ms | pivot | none: an issued ticket is Wayfare's cost (our assumption) |
| S6 | confirm_rooms: CONFIRMING → BOOKED | 15 ms | retriable | — |
| S7 | confirm_seat: CONFIRMING → SOLD | 15 ms | retriable | — |
| S5 | capture_card $1,240 | 300 ms | retriable | — |
| S8 | COMPLETED and an outbox row for the email, in the same record as S5's result | — | retriable | — |
The orchestrator writes a record to the saga log (5 ms) after every reply and before the next call. The network adds nothing in our example, so every time on this page is exact arithmetic from these numbers. Timeouts: 2 s for seats and rooms, 5 s for the PSP, 10 s for the airline. The full settings table is in Part 14.
The story has twelve beats: one trip, three databases, two partners (Part 1); all or nothing, with locks (2); the coordinator that died (3); undo instead of rollback (4); the point of no return (5); holds that expire (6); the reply that never came (7); the orchestrator that crashed (8); what other bookings see (9); nobody in charge (10); the same trip four ways (11); and one trip, end to end (12). The main line (the booking that succeeds) runs through Parts 4, 5, 6 and 12. Every failure, and every two-phase commit, runs on a copy rewound to 10:00:00.000, so one never swamps another. Every trace comes from running a private reference implementation of this setup, not from working by hand. The full event table is in Part 14.
At 10:00:00.000 the request arrives with the client's key k-4417 (the Idempotency & Effectively-Once Processing loop primitive), trips creates trip-4417, and the orchestrator writes STARTED to the saga log, done at .005 (times are seconds after 10:00:00.000).
A few words used on this page:
| Word | Meaning |
|---|---|
| Participant | A system that does part of the work: here seats, rooms, trips, the PSP and the airline |
| Coordinator | In two-phase commit, the one process that collects votes and decides |
| Prepare | A participant's promise: "I can commit this, I have forced it to disk, and I'll wait for your decision, holding my locks" |
| In-doubt | A prepared participant that hasn't heard the decision; it may not decide alone |
| Saga | A sequence of local transactions, each with a compensation |
| Compensation | A new local transaction that undoes a step's business effect |
| Pivot | The go/no-go step; after it the saga only moves forward |
| Hold | A reservation that expires on its own unless confirmed |
| Fence | Moving a hold to a state its expiry can't touch, before an irreversible step |
| Unknown | A step whose reply never came: it may or may not have happened |
| Orchestrator | The one component that sends each step and records the saga log |
| Epoch | A number that goes up each time a new instance takes over a saga; stale owners' writes are refused |
Background reading, not relied on for any fact here: Primitive #10: Two-phase commit and saga orchestration.
Your service writes the seat, then the rooms, then calls the PSP, all inside one database transaction. What does the transaction actually protect?
What to remember from Part 1
- A database transaction covers one database; a partner's API is outside every transaction.
- Across systems you choose: commit together (2PC) or commit each step with a planned undo (saga).
- Atomicity is not isolation: say which one you need.
Part 2. All or nothing, with locks
On a copy: suppose we want seats, rooms and trips to commit trip-4417 together, as one atomic unit, whatever happens. The protocol for that is two-phase commit (2PC). This Part runs it on our three databases, shows what it can't include, and measures what it costs when nothing fails.
Three databases, one decision
The trips service acts as the coordinator. Replay P0, from the reference implementation (1 ms per round trip; each PREPARE TRANSACTION, the decision record and each COMMIT PREPARED is one forced write of 5 ms):
Synthesizing vector architecture diagram...
What to notice: three forced writes sit on the path (the participants' prepare, the coordinator's decision, the participants' commit), and the FL-212 row stays locked from the first write at .005 until the commit is acknowledged at .062. That lock window, not the forced writes, is what other bookings pay for.
| Time | Event (P0) |
|---|---|
| .000 → .005 | The coordinator records global transaction g-4417 and its participants, so a recovering coordinator knows where to look |
| .005 → .020 | seats updates FL-212 and 14C: both rows locked from .005 |
| .020 → .040 | rooms updates the 3 nights (locked from .020) |
| .040 → .045 | trips inserts the trip row |
| .045 → .051 | PREPARE TRANSACTION on all three, in parallel (1 ms + 5 ms); all vote yes |
| .051 → .056 | The coordinator forces its decision: COMMIT |
| .056 → .062 | COMMIT PREPARED on all three, acknowledged. FL-212 was locked 57 ms (.005 → .062) |
Prepare, vote, decide, commit
textCOORDINATOR (the trips service), global transaction g-4417 1. record g-4417 and its participants (so recovery knows whom to ask) 2. run the writes on each participant, each in its own open transaction 3. send PREPARE to every participant; wait for every vote 4. all YES -> force the decision record COMMIT to disk any NO, or a vote missing at the timeout -> ABORT (nothing to force: no record means abort, "presumed abort", Part 3) 5. send the decision to every participant; resend until each acknowledges 6. once every participant acknowledged, forget the decision record PARTICIPANT (seats), on PREPARE g-4417 1. can I still commit? (constraints checked, rows still locked) 2. yes -> force the changes and a "prepared" record to disk, vote YES from now on: never decide alone, keep the locks, wait for the decision no -> roll back, vote NO
Each forced write has a job. The participant's prepared record is its promise: after a crash it must still be able to commit. The coordinator's decision record is the only place the outcome exists between the votes and the last acknowledgement: if it crashes after deciding COMMIT and before telling everyone, the record is what it reads on restart.
What it looks like in PostgreSQL and MySQL
The standard for this is XA: a transaction manager drives the protocol, and each database acts as a resource manager.
| PostgreSQL | MySQL (InnoDB) | |
|---|---|---|
| Do the writes | BEGIN; … | XA START 'g-4417'; … ; XA END 'g-4417' |
| Prepare | PREPARE TRANSACTION 'g-4417' | XA PREPARE 'g-4417' |
| Commit | COMMIT PREPARED 'g-4417' | XA COMMIT 'g-4417' |
| Abort | ROLLBACK PREPARED 'g-4417' | XA ROLLBACK 'g-4417' |
| List the in-doubt ones | SELECT gid, prepared FROM pg_prepared_xacts | XA RECOVER |
| One participant only | — | XA COMMIT 'g-4417' ONE PHASE |
| Default | max_prepared_transactions = 0: "Setting this parameter to zero (which is the default) disables the prepared-transaction feature"; if you use it, set it "at least as large as max_connections" | XA is supported by InnoDB; xa_detach_on_prepare is ON, so a prepared transaction is detached from its session |
| The documentation's warning | PREPARE TRANSACTION "is not intended for use in applications or interactive sessions"; "It is unwise to leave transactions in the prepared state for a long time" (it holds back VACUUM and, in the extreme, transaction ID wraparound) | "REPEATABLE READ may not be sufficient for distributed transactions" |
After PREPARE TRANSACTION, PostgreSQL says, "the transaction is no longer associated with the current session; instead, its state is fully stored on disk", and "the transaction continues to hold whatever locks it held". Part 3 needs both sentences.
What 2PC can't include
The trip also needs the card and the ticket. Put them inside the lock window (replay P0x), so that "all or nothing" covers them:
| Time | Event (P0x) |
|---|---|
| .005 → .020 | seats updates FL-212 (locked from .005) |
| .020 → .040 | rooms updates the nights |
| .040 → .440 | The PSP authorizes $1,240 (400 ms) |
| .440 → 1.640 | The airline issues the ticket (1,200 ms): it exists now, whatever the vote |
| 1.640 → 1.662 | Trip row, prepares, decision, commits, as in P0 |
FL-212 locked .005 → 1.662: 1.657 s |
Two things are wrong. The PSP and the airline never prepared: they did their work at once, so if a prepare fails at 1.651 the ticket is already issued and the authorization already holds money. Undoing them needs compensations: a saga is needed anyway. And the hot row is now locked for 1.657 s instead of 57 ms. 2PC can cover only our own databases; the partners force a saga regardless.
Isolation is still the participants'
2PC gives atomicity across participants: all commit or all abort. It adds no isolation. Each database isolates its own rows at its own level, and a reader can see seats committed and rooms not yet, between two COMMIT PREPAREDs. MySQL's own XA page says it: "REPEATABLE READ may not be sufficient for distributed transactions". "2PC gives serializability" is a common wrong answer.
The lock window
A row that every booking updates can commit one change at a time. Formula 1:
Replay H: bookings arrive at FL-212 at random, 2 a second during the sale (our number), for 60 seconds, and a waiting booking gives up after lock_timeout = 5 s.
| Design | FL-212 lock held per booking | Most commits a second | At 2 bookings a second, seeds 1 to 5 |
|---|---|---|---|
| Saga (S1 is one local transaction) | 15 ms | 1 ÷ 0.015 = 66.7 | 0 failed; almost nobody waits |
| 2PC over our databases (P0) | 57 ms | 1 ÷ 0.057 = 17.5 | 0 failed; about 0.01 waiting on average |
| 2PC with the partners inside (P0x) | 1.657 s | 1 ÷ 1.657 = 0.60 | 60 to 84 failed in the minute; 7.4 to 9.4 waiting on average, up to 9.8 in the last 30 s |
Under P0x the queue doesn't grow for ever: it levels off. The row serves 0.60 a second, so about 2 − 0.60 = 1.4 bookings a second wait out the 5 s and fail, and about 2 × 5 = 10 are waiting at any moment, each holding a database connection. P0 and the saga both have plenty of headroom at 2 a second; the gap shows when a flash sale pushes the rate past 17.5.
One store, one call
On another copy: if all three services' tables lived in one DynamoDB account and Region, the store's own transaction could replace the protocol (replay D). One TransactWriteItems writes 6 items (the flight counter with the condition available >= :one, seat 14C, the 3 room nights with conditions, and the trip), each under 1 KB, with the client token trip-4417:
| Time | Event (D), from the reference implementation |
|---|---|
| .005 → .025 | All 6 items written together (20 ms, our number): 12 write units (6 items × 2) |
| .010 | Another trip's transaction on the same flight row: cancelled with TransactionCanceledException (a transaction conflict), not queued; it consumed its write units anyway |
| 300 s | trip-4417 retried with the same token (within 10 minutes): success, no changes |
| 700 s | Retried again, token expired: a new transaction, cancelled by seat 14C's condition (status = FREE), 12 write units consumed, nothing written |
The facts it relies on, from AWS's documentation: TransactWriteItems "groups up to 100 write actions in a single all-or-nothing operation … within the same AWS account and in the same Region", the items total at most 4 MB, and no two actions may target the same item. "DynamoDB performs two underlying reads or writes of every item in the transaction: one to prepare the transaction and one to commit", so transactional writes cost 2 write units per KB, consumed even when a condition cancels the transaction. A client token "is valid for 10 minutes after the request that uses it finishes". Each item's partition keeps its own limit too (1,000 write units a second per partition, the figure the loops use), so a hot flight counter is still one hot item. Not BatchWriteItem: its up to 25 puts and deletes are each atomic, but "BatchWriteItem as a whole is not", and failed ones come back in UnprocessedItems. The last row shows why the conditions matter more than the token: after 10 minutes the token no longer protects the retry, and the seat's condition does.
The saga shrinks to four steps: the D transaction (compensatable: its undo is a second transaction that gives everything back), authorize, the ticket (the pivot), and capture. Its "booked" would come at 1.640 and its end at 1.945. It still needs a saga, because the partners are still outside. Several loops use exactly this one-store answer: step 2.3 of the ride-sharing dispatch loop (three items instead of "an ordered two-step with compensation"), and the effect-plus-idempotency-record writes in the chat, email, Drive, nearby-friends and bot-defense loops.
Why can't the PSP just "prepare" a charge?
What to remember from Part 2
- 2PC makes several databases commit together by making each promise first.
- Every participant must be able to prepare and wait: partners' APIs can't.
- Locks last from the first write to the commit, so the slowest step sets every locked row's throughput; one store's own transaction avoids the protocol, but not the partners.
Part 3. The coordinator that died
On a copy of P0: the coordinator crashes at .051, right after all three participants voted yes and before it wrote a decision. Another instance of the trips service takes over when the dead one's lease runs out, 10 s later, the same takeover the saga's orchestrator gets in Part 8. Even that short gap costs about ten bookings on FL-212. With nobody to take over, it costs about 170.
Ten seconds, or ninety, of a frozen flight
Synthesizing vector architecture diagram...
What to notice: between the crash and the scan, the three participants hold their locks and do nothing, because each has promised to wait. The new coordinator has to go and look: nothing in PostgreSQL asks it.
From the reference implementation:
| Replay | What happens | Rows locked |
|---|---|---|
| P1s: a takeover by lease | Crash at .051. At 10.051 a new instance takes over, scans pg_prepared_xacts on the three participants (one round trip, 10.052), finds no decision record for g-4417, and sends ROLLBACK PREPARED to all three: done at 10.058 | FL-212 and 14C from .005, the nights from .020, the trip row from .040, all to 10.058 |
| P1: a single process, nobody takes over | The same crash; the coordinator restarts at 90.051, scans at 90.052, rolls back by 90.058 | The same rows, for 90 s |
| P2: a crash after the decision | The coordinator forced COMMIT at .056 and crashed before sending it. On restart at 90.056 it scans (90.057), finds its COMMIT record, and sends COMMIT PREPARED: the trip commits at 90.063. With a lease takeover: 10.063 | The same rows, 90 s (or 10 s) |
Why a prepared participant waits
Blocking needs one specific moment: every participant has voted yes, and none has heard the decision. Before that moment, a participant can decide alone:
- it hasn't voted yet: it may abort (the coordinator will then decide ABORT);
- it learns that another participant voted no, or never prepared: the decision must be ABORT, so it may abort.
After that moment it may not. The coordinator might have decided COMMIT and told another participant already; aborting alone would split the trip. So it waits, holding its locks. Participants asking each other doesn't help either: all of them voted yes, so none of them knows more.
Presumed abort and the scan
The R* system's paper (Mohan, Lindsay and Obermarck, 1986) named the rule most systems use, presumed abort: it is "safe for a coordinator to 'forget' a transaction immediately after it makes the decision to abort it". No record means abort, so the coordinator never forces an ABORT record, and a coordinator that finds nothing on restart aborts.
The R* participants periodically ask the coordinator what happened. PostgreSQL participants never ask. They sit in pg_prepared_xacts until someone runs COMMIT PREPARED or ROLLBACK PREPARED. So the recovering coordinator must know every participant it used (step 1 of the protocol), scan each one's pg_prepared_xacts, and for every gid it finds: resend COMMIT PREPARED if its log says COMMIT, otherwise ROLLBACK PREPARED.
What exactly is blocked
"The whole cluster freezes" is a common wrong answer. What is locked is the rows the prepared transactions wrote: FL-212, seat 14C, three nights at H-31 and one trip row. Bookings on other flights and other hotels share no lock with them and are unaffected.
And prepared transactions hold no connections: PostgreSQL detached them from their sessions at PREPARE. What fills a connection pool is the new transactions that wait for those rows, each holding its own connection. Formula 2 (Little's law):
Each booking waits at most lock_timeout, then fails. The bookings that arrive in the last 5 s before the rollback are still waiting when the rows unlock, and succeed. So the failures are about the arrival rate × (the blocked time − lock_timeout).
Synthesizing vector architecture diagram...
Snapshot T4, P1 at 30 s. What to notice: the three transactions in "Prepared" hold the locks and no connection. The bookings in "Waiting" hold the connections. Nothing in "Unaffected" waits for anything.
Replay P1h, the 90 s of P1 under the sale's load, from the reference implementation (seeds 1 to 5):
| Arithmetic | Seeds 1 to 5 | |
|---|---|---|
FL-212 bookings that fail (2 a second) | 2 × (90.058 − .005 − 5) ≈ 170 | 148 to 183 (mean 163) |
FL-212 bookings that waited and then succeeded | about 2 × 5 = 10 arriving in the last 5 s | 6 to 12 |
FL-212 connections waiting at any moment | 2 × 5 = 10 | 8.5 to 10.3 on average; 5 to 16 at the 30 s mark |
H-31 night bookings that fail (0.2 a second, our number) | 0.2 × (90.058 − .020 − 5) ≈ 17 | 10 to 20 (mean 16) |
| Connections held by the prepared transactions | 0 | 0 |
| Bookings on other flights affected | 0 | 0 |
With the 10 s takeover of P1s, the same arithmetic gives 2 × (10.058 − .005 − 5) ≈ 10 failed bookings (seeds 1 to 5: 5 to 15, mean 9.4), and 8 to 11 more that waited up to 5 s and then succeeded.
The operator's temptation
On a copy of P2 (the coordinator decided COMMIT, then died): at 60 s, with the flight frozen, an operator runs ROLLBACK PREPARED on seats to unblock FL-212 (replay P3):
| Time | Event (P3) |
|---|---|
| .056 | The coordinator forced COMMIT, then crashed |
| 60.000 → 60.005 | The operator rolls back g-4417 on seats: FL-212 unlocks, available goes back to 5 |
| 90.056 → 90.057 | The coordinator restarts and scans: g-4417 is prepared on rooms and trips; its log says COMMIT |
| 90.063 | COMMIT PREPARED on rooms and trips; seats answers that the prepared transaction does not exist |
| The trip is booked, with its rooms, and without its seat |
On a copy of P1 (no decision) the same action happens to be right: the restart rolls back the other two, and the trip is cleanly aborted. The operator can't tell P1 from P2 without the coordinator's log. The R* paper describes exactly this operator interface, to "forcibly commit or abort" a prepared process, warns that its "misuse … could lead to inconsistencies", and suggests the operator "could use the telephone to find out the coordinator site's decision". The rule: read the coordinator's decision log before touching an in-doubt transaction. COMMIT recorded: COMMIT PREPARED. Nothing recorded, or ABORT: ROLLBACK PREPARED. Never guess. XA calls an outcome forced this way a heuristic decision.
When a participant crashes
If seats itself crashes after preparing, its prepared transaction survives: its state "is fully stored on disk", and PostgreSQL restores it, with its locks, when it restarts. FL-212 stays locked through the participant's recovery until the decision arrives. What a managed failover does with prepared transactions depends on the service; this page doesn't state it for Aurora, because AWS doesn't document it.
Coordinators that don't die: replication
Distributed databases such as Google Spanner and CockroachDB run 2PC inside, between their own shards, without this blocking. They don't use a better protocol; they replicate the coordinator. It isn't the clocks: Spanner's TrueTime gives transactions their timestamps (external consistency); what removes the blocking is the replication. In Spanner, "If a transaction involves more than one Paxos group, those groups' leaders coordinate to perform two-phase commit", and "The state of each transaction manager is stored in the underlying Paxos group". A dead leader is replaced by another replica with the same state, so "Running two-phase commit over Paxos mitigates the availability problems". CockroachDB replicates each transaction's record with Raft; other transactions check the coordinator's heartbeat, and if it stops, the record is moved to ABORTED. The Replication, Quorums & Read-Your-Writes loop primitive explains the elections (Part 6); background in Primitive #09.
Replacement isn't instant. Spanner's leader leases were "10 seconds by default" in its 2012 paper, which reports that after killing a leader, "Approximately 10 seconds after the kill time, all of the groups have leaders". Replay P4 puts our coordinator's state in such a group: its leader dies at .051, a new leader takes over at 10.051, finds no decision in the replicated log (presumed abort), and rolls back by 10.057: about 10 failed bookings on FL-212 (seeds 1 to 5: 5 to 15), the same as P1s, with no takeover machinery of your own and no human. And all of it is still one database: replication doesn't let a PSP or an airline prepare.
The coordinator is back in 90 s. An operator offers to roll back the prepared seats transaction at 60 s to unblock the flight. What do you need to know first?
What to remember from Part 3
- A crash after the votes blocks exactly the rows the prepared transactions locked, until the coordinator decides.
- Read the decision log before touching an in-doubt transaction; never guess.
- Replicated coordinators (Spanner, CockroachDB) and one-store transactions (DynamoDB) avoid the wait, inside one system.
Part 4. Undo instead of rollback
Back to the main line. trip-4417 is booked as a saga: each step commits in its own database at once, and nothing holds a lock across services. But on a copy where night 2 at H-31 has no rooms left (replay F1), the seat is already held when the rooms fail. There is no rollback across databases. Something must undo the seat.
Seat held, rooms gone
Synthesizing vector architecture diagram...
What to notice: the orchestrator writes ABORTING before the first undo, so a crash in the middle of undoing continues backwards, never forwards. The card was never touched, and the whole failure costs $0.
A saga
The name comes from Hector Garcia-Molina and Kenneth Salem's 1987 paper Sagas, written for long-lived transactions inside one database: "A LLT is a saga if it can be written as a sequence of transactions that can be interleaved with other transactions." Each step Táµ¢ has a compensating transaction Cáµ¢ that "undoes, from a semantic point of view, any of the actions performed by Táµ¢". Services and partners fit the same idea: each step is a local transaction in its own system, and a failure is repaired by running the compensations of the steps that happened, newest first.
The main line so far, from the reference implementation:
| # | Time | Step | Result |
|---|---|---|---|
| 1 | .000 → .005 | Log STARTED (epoch 1) | |
| 2 | .005 → .020, log → .025 | S1 hold_seat | 14C HELD, available 5 → 4, expires 10:15:00.020 |
| 3 | .025 → .045, log → .050 | S2 hold_rooms | 3 nights HELD, 4 → 3 each, expires 10:15:00.045 |
| 4 | .050 → .450, log → .455 | S3 authorize_card | $1,240 authorized, valid 7 days; no money moved |
Snapshot T1, at .480 (after the fence of Part 5):
| Lane | State |
|---|---|
| Trip | Log: STARTED, S1 DONE, S2 DONE, S3 DONE, SF DONE (epoch 1) |
| Inventory | 14C CONFIRMING, FL-212 available 4; 3 nights CONFIRMING, 3 left each |
| Money | Authorization 0 |
| Customer | "processing"; nothing irreversible has happened yet |
Semantic undo
A compensation is a new local transaction that restores the business meaning, not a physical rollback of the old bytes:
- Undoing a seat hold writes
available = available + 1and marks the holdRELEASED. It doesn't delete the hold row (Part 6 needs it) and doesn't write back the number it read (Part 9 shows why). - Undoing an authorization is a void, a call to the PSP.
- Undoing a transfer between accounts writes a balanced reversal, a pair of entries. The digital wallet loop's history has the counter-example: a compensation that inserted a one-sided credit line created money, and was replaced by a try-confirm-cancel flow with an in-transit account and balanced pairs (step 2.2 of the wallet loop).
A compensation is never one-sided and never a delete of the step's record. It is itself a complete local transaction, with its own key.
The idea reaches past servers. A phone app that shows a like at once and removes it when the server rejects the action runs a client-side compensation (step 2.5 of the mobile news feed loop).
Backwards, in order
On a copy where the card is declined (replay F2), two steps must be undone:
| Replay | Rejection | Undo, from the reference implementation | Customer told | Cost |
|---|---|---|---|---|
| F1: night 2 sold out | S2 at .045 | Log ABORTING → .050; C1 .050 → .065; C1 DONE + ABORTED → .070 | "sold out" at .070 | $0 |
| F2: card declined | S3 at .450 | Log ABORTING → .455; C2 .455 → .475, log → .480; C1 .480 → .495; C1 DONE + ABORTED → .500 | "declined" at .500 | $0 |
textON a step Sk rejected for a business reason (sold out, declined, fare changed): write "ABORTING (reason)" to the saga log FOR each step with a DONE record, newest first: IF it has a compensation Ck: run Ck with key trip-4417:<step>:undo (safe to repeat) write "Ck DONE" (the last one: "Ck DONE + ABORTED", one record) tell the customer
Why newest first: a later step may depend on an earlier one (a capture needs its authorization), so undoing in reverse never leaves a step whose basis is gone. Only an explicit rejection starts this loop. A timeout doesn't (Part 7).
The seat is held and the rooms are sold out. What does "undo the seat" write, and why not just delete the hold row?
What to remember from Part 4
- A saga commits each step locally and plans an undo for each.
- An undo is a new transaction that restores the meaning, not the old bytes.
- Undo in reverse order, each with its own key.
Part 5. The point of no return
The hook charged first and refunded later: the expensive undo was the first step. A saga lets us choose the order. This Part puts the cheap, likely-to-fail steps first, puts the one step that can't be undone at a chosen point, and answers the customer there.
Three zones
Azure's saga pattern guide names the three kinds of step:
- Compensatable steps "can be undone or compensated for by other transactions with the opposite effect": our holds and the authorization.
- The pivot: "Pivot transactions serve as the point of no return in the saga. After a pivot transaction succeeds, compensable transactions are no longer relevant." Ours is issuing the ticket.
- Retriable steps "follow the pivot transaction. Retryable transactions are idempotent and help ensure that the saga can reach its final state". Ours: confirm the rooms, confirm the seat, capture, send the email.
Synthesizing vector architecture diagram...
Snapshot T2, at 1.685. What to notice: everything in "Compensatable" is done and undoable, the one step in "Pivot" is done and isn't, and "Retriable" holds the four steps still to run, in the order shown. The customer node hangs off the pivot, not off the last step.
Ordering rules
-
The steps most likely to fail go first, if they're cheap to undo. The seat and the rooms are the ones that sell out (the hotel did, in the hook), and each hold is released in 15 to 20 ms, so both holds come before the card.
-
Cheap undos before expensive ones. A hold costs 15 to 20 ms to release. A void costs a call. A refund keeps the fee and takes 5 to 10 business days.
-
Authorize; don't charge. The authorization reserves the money; the capture takes it after the pivot. A void is free:
Undo of a card step Cost When the customer sees it Possible when Void (cancel the authorization) "at no cost" The hold disappears from the statement when the bank releases it (bank-dependent) Before capture, within the authorization's life (7 days here) Refund The fee is kept: 1,240 "approximately 5-10 business days later" After capture -
The irreversible step last among the risky ones: that is the pivot.
-
After the pivot, only steps that can't fail for a business reason.
Fence, then the pivot
Rule 5 has a trap. The holds expire in 15 minutes by their own database's clock (Part 6). If the saga stalls between the ticket and the confirmations for longer than that, a sweeper releases the holds, and "confirm the rooms" fails for a business reason after the ticket is issued. So just before the pivot, step SF fences both holds: one conditional update each, HELD → CONFIRMING only if still HELD and unexpired. A sweeper never touches CONFIRMING. SF is the last compensatable step: if a hold already expired, SF fails before anything irreversible, and the saga undoes itself. The hotel loop learned this the hard way: its confirm flow once charged the card before fencing the hold, and now moves HOLD → CONFIRMING first and charges second (step 2.3 of the hotel reservation loop). Part 6 shows the fence failing and saving the trip.
The main line, from the reference implementation:
| # | Time | Step | Result |
|---|---|---|---|
| 5f | .455 → .475, log → .480 | SF fence_holds | Seat and rooms HELD → CONFIRMING, 900.045 − .475 = 899.6 s before the rooms hold would expire |
| 5 | .480 → 1.680, log → 1.685 | S4 issue_ticket, quote q-77 | Ticket issued; the airline checks the quote when it issues, at 1.680. The customer is told "booked" at 1.685 |
| 6 | 1.685 → 1.700, log → 1.705 | S6 confirm_rooms | CONFIRMING → BOOKED |
| 1.705 → 1.720, log → 1.725 | S7 confirm_seat | CONFIRMING → SOLD | |
| 1.725 → 2.025, log → 2.030 | S5 capture_card | 36.26; S5 DONE + COMPLETED and the outbox row for the email in one record |
Telling the customer at the pivot
Once the ticket exists, the trip can't fail: every step left is built to succeed when retried. So the customer is told "booked" at 1.685, not at 2.030, and the rest runs without anyone waiting. Snapshot T3, at 2.030:
| Lane | State |
|---|---|
| Trip | Log ends with S5 DONE + COMPLETED + outbox row (9 records, epoch 1) |
| Inventory | 14C SOLD, FL-212 available 4; 3 nights BOOKED |
| Money | 36.26, the only fee in the whole story |
| Customer | "booked" since 1.685; the email leaves through the outbox |
Confirm, then capture
The confirmations run before the capture: secure the inventory, then take the money. If the order were reversed and a confirmation could still fail, the customer would have paid for a room Wayfare doesn't hold. With the fence, neither can fail; the order is the second line of defence.
Retriable, not "retried forever"
"Retriable" means designed to succeed when retried: a capture of a valid authorization, a confirmation of a fenced hold, an email from the outbox. It doesn't mean "retry anything forever". On a copy (replay F7), the PSP answers 503 to every capture from 2.025 to 122.025 (two minutes, our number). The capture is retried with backoff and full jitter (base 1 s, cap 30 s: the Retries, Timeouts, Backpressure & Load Shedding loop primitive, Part 4), always with the same key:
| Seed | Retries | Captured at | After the outage ended |
|---|---|---|---|
| 1 | 13 | 142.062 | 20.0 s |
| 2 | 11 | 127.683 | 5.7 s |
| 3 | 13 | 134.941 | 12.9 s |
| 4 | 12 | 127.569 | 5.5 s |
| 5 | 11 | 147.296 | 25.3 s |
The ticket is never undone, the seat and rooms are already confirmed, and the customer already saw "booked". Our model's 503 is refused before the PSP starts executing, so nothing is saved under the key and the retry runs fresh. That matters: Stripe saves "the resulting status code and body of the first request made for any given idempotency key, regardless of whether it succeeds or fails … including 500 errors", but only "after the execution of an endpoint begins". A capture whose key saved a 500 needs a status lookup, not the same key again.
A retriable step still has limits. If captures kept failing for 7 days, the authorization would expire; the saga then needs a human queue and the nightly reconciliation against the PSP's report (step 2.1 of the payment loop).
When the pivot fails
The pivot is the go/no-go step: it can still say no. On a copy (replay F3), the airline rejects quote q-77 at 1.680 because the fare changed:
| Time | Event (F3), from the reference implementation |
|---|---|
| 1.680 → 1.685 | S4 rejected; log ABORTING |
| 1.685 → 1.985, log → 1.990 | C3: the authorization is voided, fee $0 |
| 1.990 → 2.010, log → 2.015 | C2 releases the CONFIRMING rooms hold: 3 → 4 each |
| 2.015 → 2.030, log → 2.035 | C1 releases the CONFIRMING seat: 4 → 5; C1 DONE + ABORTED |
| Customer told "fare changed" at 2.035. **Fees: 36.26 and made the customer wait 5 to 10 business days |
The drill's first answer follows from this Part and Part 7: The Flight Booking That Charged Without a Seat charged before the flight step; with holds and an authorization first, a flight failure costs a void, and the money never leaves the card.
Where does "send the confirmation email" go, and what happens if you put it second?
What to remember from Part 5
- Put cheap, likely-to-fail, easy-to-undo steps first; the one irreversible step is the pivot.
- Fence the holds just before the pivot; after it, only retriable steps, so the customer can be told "booked".
- Hold and authorize; never charge what you may have to refund.
Part 6. Holds that expire
On a copy (replay F5): the hold request for the rooms gets stuck in a retrying proxy at .025. At 1.000, with S2 still in flight, the customer cancels. The orchestrator sends cancel_rooms, and rooms answers "no such hold, OK". At 3.025 the stuck request finally arrives. What stops it from holding three nights for nobody?
Try, confirm, cancel
A saga whose forward steps are reservations that expire has a name: try-confirm-cancel (TCC). The main line is one:
| Phase | Steps | What it does |
|---|---|---|
| Try | S1, S2, S3 | Reserve: a seat hold, a room hold, an authorization. Each expires on its own (15 minutes, 15 minutes, 7 days) |
| Fence | SF | Move the holds to CONFIRMING before the pivot |
| (the pivot decides) | S4 | Ticket issued: confirm. Rejected: cancel |
| Confirm | S6, S7, S5 | Turn each reservation into the real thing, only if it's still reserved |
| Cancel | C1, C2, C3 | Release each reservation, only if it's still reserved |
The digital wallet's cross-shard transfer has the same shape: the money moves into an in_transit account (try), then to the receiver (confirm) or back (cancel), and a timeout is never a cancel (step 2.2 of the wallet loop).
Why holds expire
A hold is a lease on inventory: its owner (the saga) must confirm it before it runs out. The Leases, Fencing Tokens & Distributed Locks loop primitive's rules apply (Parts 3 and 6):
- It expires by the participant database's own clock (
expires_atchecked bynow()in the same database), never by an app server's clock, whose time can differ. - It expires on its own. The expiry is what bounds a lost cancel, a dead orchestrator, and a hold that arrives after its cancel.
The cost is inventory off sale. Rooms off sale on average = holds a second × how long each hold stays open (Little's law again). On our main line a hold is confirmed about 1.7 s after it is taken, so at the hotel loop's 20 holds a second only about 20 × 1.7 = 34 are open at once. The 15-minute expiry is the worst case, reached only by holds nobody confirms or cancels: if every hold ran to expiry, 20 × 900 = 18,000 would be open. Shorter expiries put abandoned stock back sooner and give slow customers less time.
Synthesizing vector architecture diagram...
What to notice: CONFIRMING has no arrow to EXPIRED: the sweeper can't touch a fenced hold, and only a cancel releases it. And a cancel for a hold that doesn't exist yet leaves a TOMBSTONE, which is where a late try ends.
Empty rollback and hanging
Synthesizing vector architecture diagram...
What to notice: the cancel reaches rooms two seconds before the hold it cancels. Because it leaves a tombstone, the hold that arrives at 3.025 is refused instead of taking three nights.
Seata, an open-source transaction framework, names the two problems: an empty rollback is a cancel for a try that never arrived, and hanging (Seata's English pages also say "dangling") is a try that arrives after its cancel. From the reference implementation:
| Replay | 1.005 → 1.025: C2 on rooms | 3.025 → 3.045: the late hold | Rooms off sale |
|---|---|---|---|
| F5, tombstone on | No hold: empty rollback, tombstone trip-4417 written | Refused by the tombstone; its late reply is ignored | None |
| F5x, tombstone off | No hold: "nothing to do" | Commits: 3 nights HELD, expires 3.045 + 900 = 903.045 (10:15:03.045) | Until the next sweep, 910.000 (10:15:10.000): 15 minutes for nobody. Without an expiry: for ever |
In F5 the orchestrator logged ABORTING at 1.005, sent C2 for the in-flight step at once (it works whether or not the step arrived), then C1 for the done step, and ABORTED at 1.050. Snapshot T5, at 3.045: the trip is ABORTED, the seat RELEASED (available back to 5), no room hold exists, and the tombstone for trip-4417 has just refused the late hold.
texthold(hold_id): try (S2) row = the hold row for hold_id if row exists: answer from the row, never hold again (TOMBSTONE, RELEASED, EXPIRED -> "refused"; REJECTED -> "sold out"; otherwise "HELD") if any of the nights has available = 0: insert (hold_id, REJECTED); return "sold out" each night: available = available - 1 insert (hold_id, HELD, expires_at = now() + 15 minutes) cancel(hold_id): C2, key trip-4417:hold_rooms:undo row = the hold row for hold_id none -> insert (hold_id, TOMBSTONE) empty rollback HELD, CONFIRMING -> each night: available = available + 1; status = RELEASED RELEASED, EXPIRED -> nothing: already given back BOOKED -> refuse: the pivot has passed
The tombstone must live at least as long as a delayed try can: longer than the retry window of every proxy and client that might resend it. The expiry bounds the damage if the tombstone is missing. The wallet loop's "reject if not yet processed" is the same rule.
The fence, and twenty minutes without an orchestrator
The fence and the sweeper, as SQL in rooms:
sql-- SF: fence the hold, only if it is still held and unexpired (or already fenced) UPDATE room_hold SET status = 'CONFIRMING' WHERE hold_id = 'trip-4417' AND (status = 'CONFIRMING' OR (status = 'HELD' AND expires_at > now())); -- 1 row: fenced (or already fenced by an earlier run of SF). 0 rows: expired or cancelled, so SF fails before the pivot. -- the sweeper, every 10 s: only HELD holds, never CONFIRMING UPDATE room_hold SET status = 'EXPIRED' WHERE status = 'HELD' AND expires_at <= now() RETURNING hold_id; -- and add 1 back to each night of each returned hold, in the same transaction
The status = 'CONFIRMING' branch makes SF safe to re-run. On a copy, SF commits at .475, its reply is lost, and the orchestrator dies at 1.000 before logging it. Epoch 2 takes over at 11.000, finds S3 DONE as the last record, and re-runs SF: it finds both holds already CONFIRMING, counts that as done, and the trip completes at 12.575. With the stricter status = 'HELD' condition, the same re-run finds 0 rows and voids a good trip (ABORTED at 11.375).
On copies (replay F8), the whole orchestrator fleet is down for 20 minutes from just before the pivot:
| Replay | Crash after | What the holds do | After the takeover |
|---|---|---|---|
| Fence off (no SF: S1, S2, S3, S4, S6, S7, S5) | S3 DONE at .455 | Expire at 900.020 and 900.045; swept at 910.000 | 1,200.455: S4 re-run, ticket issued at 1,201.655, "booked" at 1,201.660; S6 at 1,201.675: rejected, the rooms hold expired. A ticket and an authorization, but no room, and seat 14C is back on sale: NEEDS_HUMAN at 1,201.680 |
| Fence on | SF DONE at .480 | CONFIRMING: not swept | 1,200.480: S4 → 1,201.680, "booked" at 1,201.685; S6, S7, S5; completed at 1,202.030 |
| Fence on, crash between S3 and SF | S3 DONE at .455 | Expire; swept at 910.000 | 1,200.455: SF finds both holds EXPIRED and fails before the pivot: void at 1,200.780, C2 and C1 find nothing to give back, ABORTED at 1,200.830. Nothing irreversible happened |
The fence's cost: a CONFIRMING hold has no expiry of its own. If the orchestrator never comes back, the stock stays off sale. So alarm on CONFIRMING holds older than the longest normal saga (a few seconds here; minutes if a partner is slow): that alarm is the backstop the expiry used to be. Shopify's flash-sale design has the same shape: a PAYING state that no timer touches (step 2.4 of the Shopify case study).
Your cancel reached rooms before the hold did. rooms said "no such hold, OK". What happens three seconds later?
What to remember from Part 6
- A hold is a lease on inventory: it must expire on its own, by the participant's clock.
- A cancel for something you haven't seen yet must leave a tombstone.
- Fence the holds before the pivot so their expiry can't release them after the point of no return; alarm on old fenced holds.
Part 7. The reply that never came
On a copy (replay F4): at .450 the PSP authorizes $1,240, and the reply is lost on the way back. At 5.050 the orchestrator's 5 s timeout fires. Did the card get authorized? The orchestrator can't know, and the drill's failure lives right here.
Timed out at 5 s
Synthesizing vector architecture diagram...
What to notice: nothing is undone at 5.050. The orchestrator writes UNKNOWN, asks the PSP again with the same key, and gets the saved answer: the card was authorized, so the saga continues forwards.
From the reference implementation:
| Time | Event (F4) |
|---|---|
| .450 | The PSP authorizes; the reply is lost |
| 5.050 → 5.055 | Timeout; log S3 UNKNOWN |
| 5.555 → 5.605 | Replay with key trip-4417:authorize after 0.5 s; the PSP returns the saved result in 50 ms: authorized |
| 5.605 → 5.610 | Log S3 DONE |
| 5.610 → 5.635 | SF fences both holds (still 894 s before they expire) |
| 5.635 → 6.840 | S4; "booked" at 6.840, 5.155 s later than the main line; completed at 7.185 |
Unknown, not failed
A timeout says only that the answer didn't arrive in time. The request may never have arrived, may have been refused, or may have succeeded with its reply lost. The Idempotency loop primitive owns this rule, the key record behind it, and "When a compensation times out" (Parts 2 to 4); step 1.2 of the payment loop and step 1.2 of the mobile stock trading loop apply it. For a saga it means one thing: only an explicit rejection starts backward recovery. A timeout is logged as UNKNOWN and resolved first.
What guessing costs (replay F4x, the naive rule "a timeout is a decline"):
| Time | Event (F4x) |
|---|---|
| 5.055 | Log S3 UNKNOWN |
| 5.060 | Log ABORTING: the timeout is treated as a decline |
| 5.060 → 5.105 | C2 and C1 release the holds; ABORTED |
| 5.105 | The customer is told "declined" |
| The $1,240 authorization stays on the customer's card for up to 7 days, with no trip and nothing in the saga log to void it |
Resolve, then decide
| Participant | How to resolve an unknown step |
|---|---|
| The PSP | Replay the same request with the same key: it returns the saved result (Stripe keeps keys at least 24 hours). A request still running with that key gets a 409; wait and replay. Or look the payment up by your reference |
Our own participants (seats, rooms) | Replay the step with its key: the hold row keyed by hold_id returns the saved answer |
| The airline (our assumption: it honours keys) | Replay with the key. A partner without keys: look the ticket up by your booking reference before re-issuing |
| A participant that offers "cancel if not yet processed" | Send that; its answer tells you which case you're in (the wallet loop's rule) |
Once resolved: if the step happened and the saga continues, write DONE and go on. If the saga is aborting, the unknown step is compensated too, once you know it happened; AWS's own saga sample does the same when a step fails, reverting that step as well as the ones before it. The one exception: when the saga is already aborting, the orchestrator may send an in-flight step's compensation without resolving first, because the compensation works either way (it releases the hold, or leaves a tombstone for it: Part 6).
When the undo fails
On a copy of F3 (replay F4r), the void's reply is lost:
| Time | Event (F4r) |
|---|---|
| 1.685 | Log ABORTING (fare changed); C3 sent with key trip-4417:authorize:undo |
| 1.985 | The PSP voids; the reply is lost |
| 6.685 → 6.690 | Timeout; log C3 UNKNOWN |
| 7.190 → 7.240 | Replay with the same undo key: the saved result, "voided" |
| 7.245 → 7.290 | Log C3 DONE; C2, C1; ABORTED |
This is the drill's first question: a compensating refund (or void) that fails halfway. What keeps the money safe:
- The intent is recorded first.
ABORTINGis in the log before C3 is sent, so a crash continues backwards. - The undo has its own key, so every retry is the same void, never a second one.
- Retry with backoff and jitter while the PSP is failing, the same key each time. After the retry budget, the saga stays
ABORTING, raises an alarm and goes to a human queue: never "failed, stop". While the void is stuck, the fenced holds stay off sale too, and theCONFIRMINGage alarm of Part 6 fires as well. - The nightly reconciliation against the PSP's report has the last word: every authorization the PSP holds must match a trip, or be voided. The hotel loop's refund is the same durable step (step 3.5 of the hotel loop), and the scheduler loop marks such a saga
COMPENSATION_STUCK(step 3.1 of the scheduler loop).
Two limits are shared. FL-212's row is taken by every hold, release and confirmation on the flight, so a burst of compensations after an outage competes with new bookings for the same lock. And the PSP's rate limit is shared: Stripe's live-mode limit is 100 requests a second per account, and each API endpoint gets 25 a second unless noted. So 1,000 voids after an outage use the cancel endpoint's own 25 a second for 1,000 ÷ 25 = 40 s, and share the account's 100 a second with every authorization and capture. The crowdfunding loop's mass refund draws on the same budget (step 3.4 of the crowdfunding loop).
Your authorize call timed out. Do you run the compensations? Which ones?
What to remember from Part 7
- A timeout is an unknown: find out before you undo.
- Undo the unknown step too, once you know it happened.
- A failed undo is retried with its own key, then alarmed and reconciled, never dropped.
Part 8. The orchestrator that crashed
On a copy (replay F6): the orchestrator (epoch 1) sends "issue ticket" at .480 and dies at .800, right after renewing its lease. The airline still issues the ticket at 1.680, and its reply goes nowhere. The last record in the saga log is SF DONE. What does the replacement do?
The saga log
The saga log is the saga. The orchestrator writes one record (5 ms) after every reply and before the next call:
| Record | Written when |
|---|---|
STARTED | Before the first step |
Sk DONE | After step k's success |
Sk UNKNOWN | After step k timed out |
ABORTING (reason) | Before the first compensation |
Ck DONE | After compensation k's success |
COMPLETED or ABORTED | In the same record as the last step's or last compensation's result |
Writing after the reply means a crash loses at most "which step was in flight". Writing before the next call means the log always shows every step that could have started.
Taking over
sql-- every saga-log write is conditional on the writer's epoch (a fencing token) BEGIN; SELECT owner_epoch FROM saga_owner WHERE trip_id = 'trip-4417' FOR SHARE; -- not 1: ROLLBACK and stop; another instance owns this saga now INSERT INTO saga_log (trip_id, epoch, record) VALUES ('trip-4417', 1, 'S4 DONE'); COMMIT; -- the takeover's UPDATE of owner_epoch waits for this lock, so it reads the log after this write, never before
textOWN each saga with a lease of 10 s, renewed on each log write and every 3 s ON taking over trip-4417 (the lease ran out): owner_epoch = owner_epoch + 1 (conditional write; we are epoch 2) read the saga log IF it contains ABORTING: continue backwards (compensations without a DONE record) ELSE: re-run the first step without a DONE record, with its SAME key, and continue
From the reference implementation:
Synthesizing vector architecture diagram...
What to notice: the new owner doesn't know whether the ticket was issued, and doesn't need to: the same key makes the airline return the ticket it already issued. The last two messages show the fence: a paused old owner that wakes up can't write to the log any more.
| Time | Event (F6) |
|---|---|
| .480 | S4 sent (log shows SF DONE) |
| .800 | The orchestrator renews its lease (it now runs to 10.800), then dies |
| 1.680 | The airline issues ticket TKT-1; the reply goes nowhere |
| 10.800 | The lease runs out; a new instance takes over with epoch 2, reads SF DONE, re-sends S4 with trip-4417:ticket |
| 12.000 → 12.005 | The airline returns TKT-1 (a replay takes its full 1.2 s in our model); log S4 DONE; "booked" |
| 12.005 → 12.350 | S6, S7, S5; COMPLETED |
| "Booked" 10.320 s later than the main line. Had the airline returned its saved result in 50 ms, "booked" would have come at 10.855 |
Snapshot T6, at 10.800: the log's last record is SF DONE; the owner epoch is now 2; 14C and the three nights are CONFIRMING (fenced, so no sweeper touches them during the outage); the authorization is uncaptured; the airline holds TKT-1, which the log doesn't know about yet; the customer still sees "processing".
In the paused variant, epoch 1 wasn't dead, only stalled, and wakes at 20.000 with the airline's reply waiting. Its write S4 DONE at 20.005 is refused: the owner epoch is 2. A lease alone would not stop it; the epoch on every log write does (the Leases, Fencing Tokens & Distributed Locks loop primitive, Parts 4 and 5).
Why the key matters
On a copy where the airline call carries no key (replay F6x), the re-run at 10.800 issues TKT-2 at 12.000. Wayfare now holds two tickets for one seat and pays $520 more, unless the airline's rules let it void one. Every step needs a key for exactly this moment. Where a partner has no key, look the result up by your own reference before re-sending (the payment loop's lookup).
Determinism and versions
Recovery works only if the orchestrator's decisions depend on nothing but the log: "if S3's logged result said declined, abort" is safe to recompute; "if a random number is below 0.5" or "if the clock says after 10:00" is not, unless the value itself was logged. And a running saga must keep the step list it started with: a deploy that inserts a step between S3 and S4 must not change what a half-finished saga does next. Step Functions does this for you: after UpdateStateMachine, "Running executions will continue to use the previous definition". Temporal replays each workflow's history through the code, so code changes that alter the order of decisions need its versioning APIs.
Commands through the outbox
The orchestrator's trips database holds the saga log and an outbox. Writing "S4 DONE" and the next command in one local transaction, then relaying the command, means a crash can't lose the command or send it without its record: the Change Streams loop primitive, Parts 2 and 3. The relay sends at least once, so every participant must accept a repeated command, which the keys already guarantee.
Durable workflow engines are this Part, packaged: Step Functions Standard keeps each execution's history as its saga log, Temporal keeps an event history per workflow and retries activities, and the job scheduler loop builds its own (step 3.1). Part 13 has their limits. After a Region failover, sagas restart from the replicated log in the new Region (step 3.5 of the payment loop; the Multi-Region Failover loop primitive).
The orchestrator died after sending "issue ticket" and before writing anything. What does its replacement do?
What to remember from Part 8
- The saga log is the saga: write it after every reply.
- Recovery re-runs the next unfinished step with the same key.
- Own each saga with a lease and a fencing epoch on every log write, so only one instance drives it.
Part 9. What other bookings see
On a copy (replay I1), FL-212 has 5 seats. trip-4417 holds one at .020 (5 → 4). trip-4418 starts at .300 and holds another at .320 (4 → 3). Then trip-4417 takes F3's path: the fare changed, and C1 releases its seat at 2.030. If C1 writes back the number it read, available becomes 5 while one seat is still held by trip-4418.
The undo that oversold
Synthesizing vector architecture diagram...
What to notice: both panels run the same three writes; only the last one differs. "Absolute undo" writes back the value trip-4417 read before trip-4418 held its seat, which erases that hold. "Delta undo" adds 1 to whatever is there.
| Replay | Undo form | available after C1 at 2.030 | Truth |
|---|---|---|---|
I1, undo = absolute | Restore the value read in S1 (5) | 5 | 4: one seat is sold twice later |
I1, undo = delta | available = available + 1 | 4 | 4 |
The 1987 paper made this exact point about seats: a compensation "cannot simply store in the database the number of seats that existed when Táµ¢ ran".
Three anomalies
A saga commits each step at once, so every intermediate state is visible to others. Azure's guide lists what goes wrong:
- Lost updates: one saga overwrites another's change (I1).
- Dirty reads: a saga reads data another saga changed but hasn't finished, and may still undo (I2).
- Fuzzy (non-repeatable) reads: two steps of one saga read the same data and see different values, because another change came in between (I3).
I2, on a copy with 1 seat left: trip-4417 holds it at .020 (available 0). trip-4419 arrives at .480, its S1 finds 0 at .500, and it is told "sold out" at .505. Then trip-4417 aborts (F3), and the seat comes back at 2.030: 1.530 s after trip-4419's hold was refused (.500). The seat was never sold; trip-4419 read a hold as if it were a sale.
I3, on a copy: the fare for FL-212 changes from 560** at 1.000, while S4 (sent at .480) is in flight. S4 carries quote q-77, and the airline checks it when it issues, at 1.680, so it rereads the fare, finds it changed, and rejects: exactly F3. The trip is voided and released by 2.035 at a cost of 560: the trip now costs 720 = **1,240 authorization. By default a capture takes at most the authorized amount (some online card payments can be overcaptured, within card-brand limits, on some pricing plans or on request), so either Wayfare absorbs $40, or it goes back to the customer after telling them "booked".
Countermeasures
| Anomaly | Countermeasure | On our trip | Its cost |
|---|---|---|---|
| Lost update | Commutative updates: undo with a delta, so order doesn't matter | C1 adds 1 (I1) | Undo code must never read-then-write |
| Dirty read | Semantic lock: a state that says "in progress", which readers treat as not final | The hold is HELD or CONFIRMING, never SOLD; the page can offer "1 seat on hold may free up: notify me" (I2) | Every reader must understand the pending state |
| Fuzzy read | Reread value: carry a version or quote ID into the step that acts, and fail that step if it changed | S4 carries q-77 (I3) | Some pivots are rejected, and the saga undoes itself |
| Dirty read, by order | Pessimistic view: reorder so risky updates happen in retriable steps, after the pivot | The seat becomes SOLD only in S7, after the ticket | Fewer steps can be moved |
Two more, one sentence each. A version file logs every operation on a record so they can be applied in the right order even if they arrive out of order. Risk-based concurrency by value picks the mechanism per operation: a saga for a $1,240 trip, and 2PC or a single store's transaction where the money at risk justifies the locks. For the isolation levels inside one database, see Primitive #21. The wallet loop's rule "never check the total across shards mid-transfer" is the same lesson (step 2.3 of the wallet loop).
Two trips hold FL-212 seats and one aborts. Its undo sets available back to the number it read. What goes wrong?
What to remember from Part 9
- Every step of a saga is visible to others at once.
- Undo with deltas; show holds as holds; reread what the pivot depends on.
- Choose the countermeasure per anomaly, not one for all.
Part 10. Nobody in charge
On a copy (replay C), the same trip runs with no orchestrator: each service reacts to events on a bus and publishes its own. Night 2 is sold out, as in F1. The saga aborts in 205 ms instead of 70, and "where is trip-4417?" has no single answer.
Two styles
Synthesizing vector architecture diagram...
What to notice: in "Orchestration" one component sends every command and knows the order. In "Choreography" the order exists only as the chain of who reacts to which event; no box knows the whole trip.
From the reference implementation (50 ms per bus hop, our number):
| Time | Event (C) |
|---|---|
| .000 → .005 | trips writes the trip and an outbox row; publishes TripRequested |
| .055 → .070 | seats receives it, holds 14C, publishes SeatHeld |
| .120 → .140 | rooms receives it; night 2 sold out; publishes RoomsUnavailable |
| .190 → .195 | trips receives it and marks the trip FAILED |
| .190 → .205 | seats receives it and releases 14C: everything is undone |
Same trip, same failure as F1: .205 against .070. The difference is three bus hops (150 ms), less the orchestrator's three log writes (15 ms). The larger difference is where the logic lives: seats must know that RoomsUnavailable means "release", and the answer to "where is trip-4417?" is spread over trips, seats and rooms.
Ordering on the bus
On a copy (replay Cx), the bus doesn't order events. SeatHeld (published at .070) is delayed on its way to rooms until 1.300 (our number). The customer cancels at 1.000, and TripCancelled reaches rooms at 1.055, before the event it cancels:
| Setting | 1.055 → 1.075: TripCancelled at rooms | 1.300 → 1.320: SeatHeld at rooms | Result |
|---|---|---|---|
| No ordering, no tombstone | "No such hold": nothing | Holds 3 nights | Rooms off sale until the sweep at 910.000: F5's race again |
| No ordering, tombstone | Tombstone written | Refused | Clean |
| FIFO per event group (group = trip ID) | Waits behind SeatHeld | Holds 3 nights at 1.320, then TripCancelled releases them at 1.340 | Clean, in order |
Amazon EventBridge now offers both: the Custom Event Bus has FIFO subscribers that deliver "in the order that they were published within an event group" (set by EventGroupId), and the Custom Event Bus - Classic doesn't order. Part 13 has the details. A cycle is the other trap: AWS's choreography guide warns of "cyclic dependencies" when services react to each other's events.
| On equal terms (same trip, same F1 failure, same 50 ms hop) | Orchestration | Choreography |
|---|---|---|
| Time to undo everything | .070 | .205 |
| "Where is my trip?" | One saga log | Three services' records |
| Where the order and the undo live | One place | In every service's event handlers |
| Adding a car-rental step between rooms and card | One service changes (the orchestrator), plus the new one | Every service that reacts to the neighbouring events changes |
| Ordering needs | None: the orchestrator sends one command at a time | FIFO per trip, or tombstones everywhere |
| Extra to run | The orchestrator and its log | Nothing central; more events, and a view to answer "where is it?" |
| Who notices a step that never answers | The orchestrator's timeout for that step | Nobody, unless each service or a separate monitor times out waiting for the next event |
The loops chose: "for money we want one place that knows the state" (step R2.7 of the payment loop). Uber's multi-entity changes are sagas too (step 2.7 of the Uber case study). Choreography fits short, loosely coupled reactions that need no multi-step undo.
A new "car rental" step must sit between rooms and card. How many services change under each style?
What to remember from Part 10
- Orchestration keeps the order and the undo in one place; choreography spreads them.
- Events reorder unless you order them per trip.
- For money and multi-step undo, prefer an orchestrator.
Part 11. The same trip four ways
So which one? Every row below runs the same trip, the same participants, the same partners and the same failures. 2PC covers only our three databases wherever it appears, because they are the only participants that can prepare; the PSP and airline steps are identical in every row. A crashed coordinator and a crashed orchestrator are both recovered by the same 10 s lease takeover.
On equal terms
| 2PC over our databases (P0) + saga steps for the partners | 2PC with the partners inside (P0x) | Saga, plain steps (take, then give back) | Saga with TCC holds (the main line) | One TransactWriteItems (D) + a 4-step saga | Choreography (C) | |
|---|---|---|---|---|---|---|
| Who can take part | Databases that prepare | The same; partners act at once, outside the vote | Anyone with a step and an undo | Anyone with a try, a confirm and a cancel | One DynamoDB account and Region, then saga steps | Anyone on the bus |
| Atomicity | Databases: yes. Trip: by the saga steps | No: a failed prepare leaves a ticket issued | Eventually, by compensation | Eventually, by compensation; holds expire | Items: yes. Trip: by the saga steps | Eventually, by compensation events |
| Isolation | The databases' rows: locked until commit | Same, for 1.657 s | None (Part 9) | Semantic locks (HELD) | Serializable for the items; none across steps | None |
FL-212 locked or held per booking | Locked 57 ms | Locked 1.657 s | Locked 15 ms; the seat is taken at once | Locked 15 ms; held 15 minutes, then fenced | One transaction; conflicts cancelled, not queued | Locked 15 ms; held |
| Hot row at 2 bookings a second (H) | 17.5 a second possible; 0 failed | 0.60 a second; 60 to 84 failed a minute | 66.7 a second; 0 failed | 66.7 a second; 0 failed | Conflicting transactions cancelled and retried | 66.7 a second; 0 failed |
| Coordinator or orchestrator crash, 10 s takeover | The locked rows block for 10.05 s (P1s): about 10 failed FL-212 bookings (seeds 1 to 5: 5 to 15), 8 to 11 delayed | The same, and the ticket may already exist | No other booking blocked; the trip waits | No other booking blocked (0 failed, seeds 1 to 5); trip-4417 "booked" 10.320 s later (F6) | The transaction itself doesn't block; the saga resumes | No coordinator; a crashed service's events wait for it |
| Nobody takes over for 90 s | About 170 failed (P1h) | The same | No other booking blocked; the takes wait for recovery at 90 s | No other booking blocked (0 failed); trip-4417 resumes at 90 s: its holds (15 minutes) outlive the outage, fenced ones wait for the orchestrator | Nothing blocks | Events wait |
| A partner times out | An unknown to resolve with a key: 2PC doesn't remove it | The same, inside the lock window | Unknown: resolve | Unknown: resolve | Unknown: resolve | Unknown: resolve |
| A failure after the card step | Void, $0 | Void, $0, plus the issued ticket | 36.26** if charged first | Void, $0 | Void, $0 | Void, $0 |
| Undo code | For the partner steps | For the partner steps | Every step | Every step (cancel) | The D undo and the partner steps | Every step, in every service |
| Customer told "booked" | After the 2PC and the partner steps | 1.662, only if nothing fails | At the pivot | 1.685 (at the pivot) | 1.640 | When trips sees the last event |
| Operations | A transaction manager, max_prepared_transactions, a runbook for in-doubt | The same, worse | An orchestrator and its log | The same, plus a sweeper and a CONFIRMING alarm | The same as the saga, one store | A bus, and a view to answer "where is my trip?" |
Read each row from both sides. 2PC does win something: the rows it covers are isolated until commit, and it needs no undo code for them. The saga loses that isolation (Part 9) and pays with undo code. The one-store transaction is the cheapest atomicity of all, but only if every item lives in one account and Region, and it still needs a saga for the partners. P4's replicated coordinator (Part 3) blocks for about the same 10 s as P1s, without machinery of your own, but inside one database.
Choose this when
- One store's own transaction (
TransactWriteItems, one PostgreSQL database'sBEGIN … COMMIT): every item lives in that store. Always the first choice when it applies. - A distributed SQL database (its internal 2PC over replicated groups): the data is in one such database and you want serializable transactions across its shards.
- 2PC (XA) across your own databases: rarely. Only when the databases can't be merged, the transactions are short, the rows aren't hot, and you run a transaction manager whose log survives it. AWS advises against XA on Aurora MySQL (Part 13).
- A saga with holds and an orchestrator: anything that involves a partner API, a long wait, or several services. The default for payments, bookings and transfers.
- Choreography: short reactions between loosely coupled services with no multi-step undo; never for money.
This is the drill's second question: why a saga over 2PC for a cross-service booking. The partners can't prepare, so 2PC can't include them; 2PC would hold the hot row locked across every round trip and for the whole outage if the coordinator dies after the votes; and a saga's crash blocks no other booking. The price is isolation and undo code.
What to remember from Part 11
- 2PC and one-store transactions give atomicity inside systems that can prepare.
- Across partners you need a saga; holds make it cheap to undo.
- Decide by who can prepare, how long locks would last, and what the business can undo.
Part 12. One trip, end to end
Each guarantee on this page comes from one layer's durable record and one layer's key. Follow trip-4417 from the click to the email, on the main line.
Synthesizing vector architecture diagram...
What to notice: the boxes run top to bottom in the saga's order, and each names what makes a repeat at that layer harmless: a key, a conditional state change or an event ID. Only "Airline: ticket" can't be undone; the boxes after it need no undo, because they are retriable.
| Layer | Durable state | What can repeat | The guard | The undo |
|---|---|---|---|---|
| Browser and API | The client's key k-4417 mapped to trip-4417 | A double click, a client retry | The key record (the Idempotency loop primitive, Part 3) | — |
| Orchestrator | The saga log in trips, the lease, the owner epoch | A step re-sent after a takeover | Every step's key; every log write conditional on the epoch | Backward recovery from ABORTING |
seats, rooms | The hold rows, keyed by hold_id, with expires_at by their own clock; the inventory counters | A command relayed twice; a late hold | The hold row (saved answer), the tombstone, the conditional fence | C1, C2: a delta and RELEASED |
| PSP | The authorization, keyed; the issuer's hold on the card | A replay after a lost reply | Its idempotency keys (kept at least 24 hours) | C3: a free void |
| Fence | CONFIRMING on both holds | A re-run after a takeover | "Only if HELD and unexpired"; a re-run that finds CONFIRMING counts as done | C1, C2 still release it |
| Airline | The ticket (our assumption: keyed, quote-checked) | A re-issue after a crash | The key, or a lookup by reference | None: the pivot |
| Confirmations and capture | SOLD, BOOKED, the capture | Retries after a 503 | Keys; a fenced hold can't have expired | None needed: retriable |
| Outbox and email | The outbox row, written with COMPLETED | The relay re-sends | The event ID | None needed |
Two rules come out of the table. The saga log decides, and the partners' records confirm: the nightly reconciliation compares the two. And only the pivot is irreversible; everything else repeats safely.
The scheduler loop's summary fits: the engine's state changes exactly once, and every task runs at least once (step 3.4 of the scheduler loop). "The saga is exactly-once" is only true of the log.
What to remember from Part 12
- Each layer has its own durable record and its own key.
- The saga log decides; the partners' records confirm (reconciliation).
- Only the pivot is irreversible; everything else repeats safely.
Part 13. On AWS
Where each piece runs in AWS production, stated only from AWS's public documentation (and, for the self-run engines, their own).
Managed services that use it
| Service | What it provides here | What AWS states |
|---|---|---|
| AWS Step Functions (Standard) | A durable saga orchestrator: each execution's history is its saga log | "Standard Workflows follow an exactly-once model, where your tasks and states are never run more than once, unless you have specified Retry behavior"; executions run up to one year; history is kept "for up to 90 days after your execution completes" (30 on request); 25,000 events per history; 256 KiB per state's input or output; .waitForTaskToken pauses a state until a callback; starting an execution with the same name and input as a running one returns an idempotent response; failed executions of the last 14 days can be redriven from the unsuccessful step; running executions keep the definition they started with. AWS's Prescriptive Guidance builds saga orchestration on a Standard workflow |
| Amazon DynamoDB transactions | 2PC inside one store (Part 2's D) | TransactWriteItems: up to 100 actions, 4 MB, one account and Region, never the same item twice; a prepare and a commit write per item, so 2 write units per KB, consumed even when cancelled; a conflict fails with TransactionCanceledException; a client token is valid for 10 minutes after the request finishes; with global tables, ACID "only within the AWS Region where the write API was invoked"; a failed condition cancels the transaction with TransactionCanceledException (reason ConditionalCheckFailed); on a single-item write it is ConditionalCheckFailedException, HTTP 400 |
| Amazon Aurora PostgreSQL | The saga log and outbox in one local transaction; or a 2PC participant | max_prepared_transactions defaults to 0 and is set in the DB cluster parameter group; with 0, PostgreSQL disables prepared transactions. You run the coordinator and its recovery yourself |
| Amazon Aurora MySQL | An XA participant, which AWS advises against | "We recommend that you don't use eXtended Architecture (XA) transactions with Aurora MySQL, because they can cause long recovery times if the XA was in the PREPARED state." If you must: don't leave one prepared, and keep them small |
| Amazon EventBridge | The bus for a choreographed saga (Part 10) | Custom Event Bus: FIFO subscribers deliver in publish order within an event group (EventGroupId); an undeliverable event "holds up the later events in its own group"; publish deduplication within 5 minutes by DeduplicationId or content; retention 1 to 365 days; subscribers retry for 300 s and 5 attempts by default; ingestion per event group per second is 1,500, "counted as the number of events plus their total size in KB". AWS says that, within the 5-minute window, this "gives exactly-once delivery", and adds "Keep consumers idempotent". Custom Event Bus - Classic: rules and targets, no ordering; retries "for 24 hours and up to 185 times", then "the event is dropped" (or sent to a dead-letter queue if you set one) |
| Amazon SQS | Commands and replies between the orchestrator and participants; per-trip order with FIFO message groups | At-least-once delivery from standard queues; order per MessageGroupId in FIFO queues; dead-letter queues and redrive: the Queues & Delivery Semantics loop primitive |
| AWS Lambda | The worker behind each Step Functions task | Timeout 1 to 900 s, default 3 s. A partner call that can outlast it uses a callback (.waitForTaskToken) instead of waiting inside the function |
Step Functions Standard vs Express, for a saga:
| Standard | Express | |
|---|---|---|
| How often a state runs | Exactly once, unless you configure Retry | Asynchronous: "At-least-once"; synchronous: "At-most-once" |
| How long it can run | Up to one year | Up to five minutes |
| State between steps | "Execution state internally persists between state transitions" | "Execution state doesn't persist between state transitions" |
Callbacks (.waitForTaskToken, .sync) | Yes | No |
Default StartExecution rate | Bucket 1,300, refill 300 a second in N. Virginia, Oregon and Ireland (800 and 150 elsewhere) | Refill 6,000 a second |
| State transitions | 5,000 a second in those Regions (800 elsewhere); new accounts get less; all soft quotas | No state-transition quota |
| Price | $0.025 per 1,000 state transitions (US East, N. Virginia) | Per request and duration |
| For our saga | Yes | No: a saga that waits for a partner, or must not repeat a step, doesn't fit |
At the hook's 50 trips a minute, 0.83 sagas a second × about 12 state transitions each (nine tasks and a few choices, our design) ≈ 10 transitions a second, against 5,000; and 0.83 StartExecution calls a second against a refill of 300. Both are shared by every workflow in the account and Region, and compensations are transitions too. The hotel loop rejected one workflow per hold for this reason: 500 holds a second is over the 300 a second default (step R2.7 of the hotel loop).
Running it yourself
| Option | What it is | Facts and sizing |
|---|---|---|
| Temporal on Amazon EKS | A durable execution engine: a workflow's event history is its saga log; activities are its steps | An official Helm chart for Kubernetes; persistence on Cassandra, PostgreSQL or MySQL (Aurora fits the SQL options). History limit 51,200 events or 50 MB (warnings at 10,240 or 10 MB). Activity retry defaults: first retry after 1 s, backoff coefficient 2.0, maximum interval 100 × the first, unlimited attempts; workflows don't retry by default. "Temporal guarantees that an Activity Task either runs or timeouts", so every activity that calls a partner still needs a key |
| Camunda 8 on Amazon EKS | A BPMN workflow engine with compensation events | Camunda documents a Self-Managed install on EKS with Helm. When a compensation event fires, it "invokes all compensation handlers at once without any specific order"; if order matters, trigger compensation per activity. Handlers of activities still running aren't invoked, so an in-flight step's undo is yours to model |
| PostgreSQL or MySQL on Amazon EC2, with XA | Classic 2PC with an external transaction manager, where you control recovery | Set max_prepared_transactions at least as large as max_connections if you use it; run a transaction manager whose log is durable and backed up (Narayana and Atomikos are maintained examples); alarm on the age of every row in pg_prepared_xacts; disks as in the Write-Ahead Log loop primitive |
| Sizing in words | — | State transitions = transitions per saga × sagas a second, against the account's quota. A hot row's commits a second ≤ 1 ÷ its lock hold. Stock off sale = holds a second × how long each hold stays open (the expiry is only the worst case). Waiting connections during a block = arrival rate × lock_timeout |
Look-alikes that are not this mechanism
| Look-alike | Why it looks like this | Why it isn't |
|---|---|---|
| Step Functions Express as a saga engine | "Step Functions" | At-least-once or at-most-once, five minutes, no callbacks (above) |
| SQS FIFO's "exactly-once processing" | "Exactly once" | Deduplication of sends within 5 minutes; no atomicity across services (the Idempotency loop primitive) |
| EventBridge Custom Event Bus's "exactly-once delivery" | "Exactly once" | Publish deduplication for 5 minutes plus ordered delivery per group; AWS adds "Keep consumers idempotent". It orders events; it doesn't make a saga's steps atomic |
| Kafka transactions on Amazon MSK | "Transactions" | Atomic across Kafka writes only; Kafka's participation in an external 2PC is opt-in and off by default (step 1.1 of the outbox loop) |
| DynamoDB global tables | "Multi-Region transactions" | Transactions are ACID only in the Region that receives the call; multi-Region strongly consistent tables don't support transactions (step 3.3 of the Uber case study) |
| Aurora Global Database | "Commits in two Regions" | Asynchronous replication of one database, not an atomic commit across services (the Multi-Region Failover loop primitive) |
| Amazon Aurora DSQL | "Distributed transactions" | Transactions across its own storage, within limits (the Uber loop lists them): one database's transaction, like DynamoDB's, never a transaction across your services and partners |
| EventBridge Scheduler, or an SQS delay, as the hold's timer | "A timer that releases the hold" | Timers. The hold's expiry must live in the participant's own data, checked by its conditional writes; a lost timer must not keep stock off sale |
What to remember from Part 13
- Step Functions Standard is a saga log with exactly-once states; Express is not for sagas that wait.
- DynamoDB transactions are 2PC inside one table set, one account, one Region.
- Aurora PostgreSQL can prepare transactions only if you turn it on and run the coordinator; AWS advises against XA on Aurora MySQL.
Part 14. What you've learned
Back to the sale
The Lisbon sale charged first and refunded later: $36,260 in kept fees and up to 10 customers charged for nothing. With the design on this page:
- The sold-out hotel costs $0. The seat and rooms are held, not taken; the card is authorized, not charged. F1 aborts at .070 by releasing the seat, and the card is never touched (Part 4). Had the fare changed at the pivot, a free void (F3, Part 5).
- The refund that timed out can't strand a customer. A timeout is logged
UNKNOWNand resolved with the same key; a failed void is retried with its own key, alarmed, and reconciled, never dropped: 0 unrefunded (Part 7). - Nothing irreversible happens early. The ticket is the pivot; the holds are fenced just before it; "booked" is said at 1.685, and the rest is retriable (Part 5). A 20-minute outage with the fence completes the trip; without it, a ticket with no room (Part 6).
- A crash anywhere resumes, never repeats (Part 8), and other bookings see holds as holds (Part 9).
- 2PC wasn't the answer: the partners can't prepare, and a coordinator crash after the votes freezes the flight's row for as long as it's down (Parts 2, 3 and 11).
What it costs
- No isolation between steps: semantic locks, deltas and rereads in every design that shares data (Part 9).
- Undo code for every step, each with its own key (Parts 4 and 7).
- Inventory off sale while held: holds a second × how long each hold stays open (Part 6).
- An orchestrator, its log, a lease and an epoch to run (Part 8).
- A sweeper, tombstones, and an alarm on old fenced holds (Part 6).
- Unknowns to resolve, a review queue, and a nightly reconciliation (Part 7).
The settings
| Setting | Our example | At real scale |
|---|---|---|
| Trip | trip-4417, customer c-88, seat 14C on FL-212 (240 each), total $1,240; client key k-4417 | Hotel loop: about 20 holds a second; payment loop: 50,000 payments a day |
| Traffic | Main line: one trip. Replay H: FL-212 bookings at random, 2 a second; H-31's nights 0.2 a second. Hook: 50 trips a minute | Shopify loop: 250 buyers a second on one product variant |
| Step latencies (fixed) | Log write 5 ms; S1 15 ms; S2 20 ms; S3 400 ms; SF 20 ms; S4 1,200 ms; S5 300 ms; S6 15 ms; S7 15 ms; C1 15 ms; C2 20 ms; C3 300 ms; the network adds nothing | Payment loop: PSP P99 300 to 600 ms |
| Timeouts | PSP 5 s; airline 10 s; seats, rooms 2 s | From percentiles (the Retries loop primitive, Part 2) |
| Holds | 15 minutes by the participant's clock; sweeper every 10 s; authorization valid 7 days | Hotel loop: 15 minutes, extendable to 25; Shopify: 10 minutes |
| PSP | Keys kept at least 24 hours, results saved; fee 2.9% + 30¢ kept on refund; void free; unknowns replayed after 0.5 s, saved result in 50 ms | Stripe's documentation |
| Airline (our assumptions) | Honours keys; a replay takes its full 1.2 s; checks the quote when it issues; an issued ticket is Wayfare's cost | Real partner APIs vary |
| Orchestrator | Saga log in trips; a lease of 10 s, renewed on each log write and every 3 s; an epoch on every log write | Step Functions Standard or Temporal |
| 2PC | Participants seats, rooms, trips; 1 ms a round trip; each forced write 5 ms; recovery by a 10 s lease takeover, or a 90 s restart (P1); lock_timeout 5 s | PostgreSQL's max_prepared_transactions defaults to 0 |
| Choreography | At least once, 50 ms a hop; no order, or FIFO per trip | EventBridge Custom Event Bus FIFO subscribers |
The whole story, event by event
From the reference implementation. Times are seconds after 10:00:00.000; replays run on copies.
| # | Time | Event | Result |
|---|---|---|---|
| 0 | hook | Charge first; the hotel sells out for 20 minutes | 1,000 refunds; $36,260 in kept fees; up to 10 unrefunded if 1% time out |
| 1 | .000 → .005 | Request with k-4417; log STARTED (epoch 1) | |
| 1b | text | Why not one transaction | Three databases and two HTTP partners; no BEGIN … COMMIT covers them |
| P0 | copy | 2PC over the three databases | Prepares .045 → .051, decision → .056, commits → .062; FL-212 locked 57 ms |
| P0x | copy | The same with the PSP and airline inside | FL-212 locked 1.657 s; the ticket exists before the vote |
| H | copy, 60 s | FL-212 at 2 a second | Caps 66.7 / 17.5 / 0.60 a second; P0x: 60 to 84 failed, about 10 waiting |
| D | copy | One TransactWriteItems | 12 write units; a conflict cancelled; a retry after 10 minutes stopped by the seat's condition; "booked" at 1.640 |
| P1s | copy | Coordinator crash at .051, 10 s takeover | Rollback done 10.058; about 10 failed bookings (5 to 15), 8 to 11 delayed |
| P1 | copy | The same crash, 90 s restart | Rollback done 90.058; rows locked 90 s |
| P1h | copy | P1 under the sale's load | About 170 failed on FL-212 (148 to 183), 17 on H-31 (10 to 20); about 10 connections waiting; prepared transactions hold 0 |
| T4 | P1, 30 s | Snapshot | 3 prepared transactions with locks, 0 connections; about 10 bookings waiting (5 to 16 at that moment) |
| P2 | copy | Crash after logging COMMIT | COMMIT PREPARED at 90.063 (10.063 with a takeover) |
| P3 | copy | P2 plus an operator rollback on seats at 60 s | Trip and rooms committed at 90.063, seat rolled back: a trip without its seat |
| P4 | copy | Replicated coordinator | New leader at 10.051; rolled back by 10.057; about 10 failed |
| 3p | text | A participant crashes after preparing | Its prepared transaction and locks survive its restart |
| 2 | .005 → .025 | S1 hold_seat | 14C HELD, available 5 → 4, expires 10:15:00.020 |
| 3 | .025 → .050 | S2 hold_rooms | 3 nights HELD, expires 10:15:00.045 |
| 4 | .050 → .455 | S3 authorize_card | $1,240 authorized, 7 days |
| F1 | copy | Night 2 sold out | S2 rejected .045; C1; "sold out" at .070; $0 |
| F2 | copy | Card declined | S3 rejected .450; C2, C1; "declined" at .500; $0 |
| 5f | .455 → .480 | SF fence_holds | HELD → CONFIRMING, 899.6 s before the rooms hold's expiry |
| T1 | .480 | Snapshot | Two fenced holds and an authorization; nothing irreversible |
| 5 | .480 → 1.685 | S4 issue_ticket (pivot) | Ticket at 1.680; "booked" at 1.685 |
| T2 | 1.685 | Snapshot | The pivot passed; three retriable steps and the email left |
| 6 | 1.685 → 2.030 | S6, S7, S5, S8 | Rooms BOOKED 1.700, seat SOLD 1.720, captured 2.025, COMPLETED 2.030 |
| T3 | 2.030 | Snapshot | Sold, booked, captured ($36.26 fee), email in the outbox |
| F3 | copy | Fare changed at the pivot | C3 void, C2, C1; ABORTED 2.035; $0 |
| F7 | copy | Capture 503 from 2.025 to 122.025 | 11 to 13 retries; captured between 127.569 and 147.296 (seeds 1 to 5); ticket never undone |
| 5z | text | Captures fail for 7 days | The authorization expires: human queue and reconciliation |
| 7 | main | The main line as TCC | Try S1 to S3, fence SF, confirm S6, S7, S5, cancel C1 to C3 |
| 6f | main | The fence | Conditional HELD → CONFIRMING; never swept; alarm on old CONFIRMING holds |
| F5 | copy | Cancel at 1.000 overtakes a hold stuck until 3.025 | Empty rollback, tombstone at 1.025; ABORTED 1.050; late hold refused at 3.045 |
| T5 | F5, 3.045 | Snapshot | Trip aborted; seat released; no room hold; the tombstone refused the late hold |
| F5x | copy | The same without a tombstone | Late hold commits 3.045, expires 10:15:03.045, swept at 10:15:10.000 |
| F8 | copy | Fleet down 20 minutes, fence off | Holds swept at 910.000; ticket 1,201.655; S6 rejected: a ticket, no room, NEEDS_HUMAN |
| F8 | copy | The same, fence on | Holds kept; "booked" 1,201.685; completed 1,202.030 |
| F8 | copy | Fence on, crash between S3 and SF | SF fails at 1,200.475; void; ABORTED 1,200.830; nothing irreversible |
| F4 | copy | Authorize reply lost | UNKNOWN 5.055; replay → authorized 5.605; "booked" 6.840 (+5.155 s) |
| F4x | copy | The timeout treated as a decline | "Declined" at 5.105; $1,240 held on the card for up to 7 days, nothing to void it |
| F4r | copy | F3 and the void's reply lost | C3 UNKNOWN 6.690; replay → voided 7.240; ABORTED 7.290 |
| F6 | copy | Orchestrator dies at .800, S4 in flight | Epoch 2 at 10.800 re-sends S4 with its key; same ticket; "booked" 12.005; COMPLETED 12.350; a woken epoch-1 write refused |
| T6 | F6, 10.800 | Snapshot | Epoch 2 reads SF DONE; holds fenced; the airline has a ticket the log doesn't know about |
| F6x | copy | The same without a key | A second ticket: $520 more |
| 8d | text | Determinism and versions | Decisions from logged values only; running executions keep their definition |
| I1 | copy | Two trips on FL-212, one aborts | Absolute undo: available 5 (one too many); delta: 4 |
| I2 | copy | One seat; trip-4419 turned away at .500 | The seat returns at 2.030, 1.530 s later |
| I3 | copy | Fare $560 at 1.000 | With q-77: rejected, as F3. Without: a 1,240 authorization |
| 9v | text | Version file, risk by value | One sentence each |
| C | copy | Choreography, night 2 sold out | Undone at .205, against .070 |
| Cx | copy | TripCancelled before a delayed SeatHeld | No order: a hanging hold; tombstone or FIFO per trip: clean |
| 11 | copies | The equal-terms table | 2PC blocks the locked rows for the takeover (about 10 failed); the saga blocks no other booking |
| 12 | main | End to end | Each layer's record, key and undo |
The cheat card
| Topic | Remember |
|---|---|
| The two families | Commit together (2PC, one store's transaction) or commit each step and plan the undo (saga) |
| 2PC | Prepare (force, keep locks, wait), decide (force the record), commit; atomicity, not isolation |
| Who can prepare | Databases, yes; PSPs, airlines, email providers, no. Partners force a saga |
| The lock window | A row commits at most 1 ÷ (lock hold) a second; under 2PC the hold runs from the first write to the commit's acknowledgement |
| Coordinator crash | Blocks only the rows the prepared transactions locked; waiters hold the connections; failures ≈ rate × (block − lock_timeout) |
| In-doubt | Read the coordinator's log; COMMIT → COMMIT PREPARED, else ROLLBACK PREPARED; never guess |
| Presumed abort | No decision record means abort; PostgreSQL participants never ask: scan pg_prepared_xacts |
| Non-blocking 2PC | Replicate the coordinator (Spanner, CockroachDB): about the lease time (10 s), inside one database |
| One store | TransactWriteItems: 100 actions, 4 MB, one account and Region, 2 units per KB, conflicts cancelled, token 10 minutes |
| Saga order | Cheap, likely-to-fail, easy-to-undo first; authorize, don't charge; fence; pivot; confirm; capture; email |
| Compensation | A new local transaction with its own key; a delta, a void, a balanced pair; never a delete, never one-sided |
| Holds | Expire by the participant's clock; tombstone on an early cancel; fence before the pivot; alarm on old fenced holds |
| Unknowns | A timeout is an unknown: log it, replay the key, then decide; compensate the unknown step if it happened |
| Failed undo | Retry with its key and backoff; then ABORTING + alarm + human queue; reconcile nightly |
| Recovery | Log after every reply; lease + epoch on every log write; re-run the first step without DONE, same key |
| Isolation | Lost update → delta; dirty read → semantic lock; fuzzy read → reread (quote ID); pessimistic view → reorder |
| Choreography | No central log; order per trip (FIFO group) or tombstones; not for multi-step money |
| Step Functions | Standard for sagas (exactly-once states, 1 year, callbacks, redrive 14 days); Express isn't (5 minutes, at-least-once or at-most-once) |
| Void vs refund | Void: free, before capture. Refund: the fee is kept, 5 to 10 business days |
Failure checklist
- Does every partner call run as its own saga step, outside any database transaction?
- Is every step and every compensation keyed (
<saga>:<step>and:undo), and does every participant return the saved result for a repeat? - Are compensatable steps first, the one irreversible step the pivot, and only retriable steps after it?
- Does the card get authorized before the pivot and captured after it, never charged early?
- Is every hold expiring by the participant's own clock, with a sweeper that touches only unfenced holds?
- Are the holds fenced (
CONFIRMING) just before the pivot, with an alarm on fenced holds older than the longest normal saga? - Does a cancel for a missing hold leave a tombstone that refuses a late try?
- Is a timeout logged as
UNKNOWNand resolved before any compensation runs? - Is a failed compensation retried with its own key, then alarmed and queued, never dropped, and is there a nightly reconciliation against the PSP?
- Is the saga log written after every reply, with every write conditional on the owner's epoch?
- Does recovery re-run the first unfinished step with the same key, and does a running saga keep its definition?
- Do undos use deltas, do readers treat holds as not final, and does the pivot carry the quote or version it depends on?
- If 2PC is used anywhere: does the coordinator's decision log survive it, does recovery scan every participant, and does the runbook forbid guessing?
- Is
lock_timeoutset on the rows a prepared transaction may lock, so waiters fail instead of filling the pool?
Think-first drills
Drill 1. A 2PC coordinator restarts 45 s after crashing with three prepared participants. The locked flight row receives 3 bookings a second, each waiting up to a 2 s lock_timeout. How many bookings fail, and how many connections are busy waiting?
Drill 2. A saga: hold 20 ms, authorize 500 ms, issue 1,500 ms (the pivot), then confirm 20 ms and capture 300 ms, with a 5 ms log write after each step. When can the customer be told "booked", and when does the saga finish?
Drill 3. 600 trips fail at the hotel after a $900 charge, at 2.9% + 30¢. How much fee is kept, and what does authorize-then-void save?
Interview questions
| Question | Model answer |
|---|---|
| 2PC or a saga for a booking across your services and a payment provider: why? | A saga. 2PC needs every participant to prepare and wait for a decision; a PSP (and an airline) can't, so 2PC could cover only our databases and we'd need a saga for the partners anyway. 2PC also holds the hot inventory row locked from the first write to the commit, including every round trip, and for the whole outage if the coordinator dies after the votes. A saga commits each step locally, with holds and an authorization as cheap undos, a pivot, and retriable steps after it; a crash blocks no other booking. The price: no isolation between steps (semantic locks, deltas, rereads) and undo code for every step. |
| The 2PC coordinator crashes after the votes: what exactly is blocked, for how long, and what must the on-call engineer not do? | Only the rows the prepared transactions wrote, until a coordinator reads the decision log and finishes: the takeover time with a lease (about 10 s), or the restart time without. Prepared transactions hold locks but no connections; new transactions waiting on those rows hold the connections, about arrival rate × lock_timeout of them, and those arriving earlier than lock_timeout before the end fail. Other rows are unaffected. The engineer must not roll back (or commit) an in-doubt transaction by guessing: look up the gid in the coordinator's decision log; COMMIT → COMMIT PREPARED, nothing or ABORT → ROLLBACK PREPARED. A wrong guess splits the transaction. |
| Order the steps of a trip booking and name the pivot. What changes if the airline lets you void a ticket for 24 hours? | Hold the seat and rooms (expiring holds), authorize the card, fence the holds, issue the ticket (the pivot), then confirm the holds, capture, and send the email from the outbox. Tell the customer at the pivot. If tickets can be voided for 24 hours, the ticket becomes compensatable within that window, so the pivot moves: the capture (or the end of the void window) becomes the point of no return, and a failure after the ticket costs a ticket void instead of a human queue. |
| A step times out: what do you do, and when do you compensate it? | Log it as UNKNOWN; never treat it as a failure. Resolve it: replay with the same key (the participant returns its saved result), or look it up by reference, or send "cancel if not yet processed" where offered. If it happened and the saga continues, mark it done. If the saga aborts, compensate it too, once you know it happened; if the saga is already aborting, you may send its compensation straight away, since it works either way. Treating the timeout as a decline leaves an authorization on the card with nothing to void it. |
| A compensation arrives before its step: what happens, and how do you prevent it? | That's an empty rollback: the participant has nothing to cancel. If it just says "OK", the late step arrives later and commits, and nobody will ever cancel it (hanging): stock off sale until its expiry, or for ever without one. Prevent it with a tombstone: the cancel records "cancelled" for that ID, and a later try with that ID is refused; keep the tombstone longer than any retry window; and give every hold an expiry as the bound when all else fails. Ordered delivery per saga (FIFO groups) also stops the reordering. |
| Two sagas touch the same inventory row: which anomalies can happen, and how do you prevent each? | Lost update: an undo that restores a remembered value erases the other saga's change; use deltas (commutative updates). Dirty read: one saga sees another's hold as a sale and turns a customer away; show holds as a pending state (semantic lock) and never as final. Fuzzy read: a value read early (a fare) changed before the step that acts on it; carry a version or quote ID and reread at that step. Or reorder so risky updates happen after the pivot (pessimistic view). |
Where to go next
- Idempotency & Effectively-Once Processing: the keys, the durable guard, "a timeout is an unknown" and "When a compensation times out" (Parts 2 to 4), and the payment six ways, including the 2PC and saga rows (Part 10).
- Change Streams & the Transactional Outbox: the outbox each step's command and reply leave through (Parts 2 and 3), and the outbox vs 2PC with a broker (Part 10).
- Leases, Fencing Tokens & Distributed Locks: holds as leases, whose clock decides (Part 3), the epoch that fences a paused owner (Parts 4 and 5), and orphaned locks (Part 6).
- Replication, Quorums & Read-Your-Writes: the elections behind a replicated coordinator (Part 6).
- Sharding, Hot Keys & Rebalancing: transactions across shards, and choosing keys to avoid them (Part 9); hot rows.
- Retries, Timeouts, Backpressure & Load Shedding: timeouts from percentiles (Part 2), one retry layer (Part 3), full jitter (Part 4).
- Queues & Delivery Semantics: the ack point, dead-letter queues and redrive for the commands between the orchestrator and its participants.
- Multi-Region Failover: sagas restarted from the replicated log in another Region.
- Background: Primitive #10: Two-phase commit and saga orchestration, Primitive #21: Database isolation levels, Primitive #09: Distributed consensus.
- Drill: The Flight Booking That Charged Without a Seat (question 1 in Parts 5 and 7, question 2 in Parts 2, 3 and 11).
- Loops: payment processing (steps 1.2, 2.1, R2.7, R2.11, 3.5, R3.9), the digital wallet (steps 2.2, 2.3, R2.7, R2.8, R2.11), hotel reservations (steps 1.6, 2.3, 3.1, 3.5, R2.7), the transactional outbox and ledger (steps 1.1, R1.8, 2.2, R3.11), crowdfunding (steps 2.3, 3.4), the job scheduler (steps 3.1, 3.4, R3.11), the key-value store (step R2.3), chat (step 1.3), the news feed (step 3.3), the gaming leaderboard (step 3.5), nearby friends (step 1.2), ride-sharing dispatch (step 2.3), YouTube (step 2.1), Google Drive (step 2.5), ad-click aggregation (step 2.4), the email service (step R1.6), the mobile news feed (step 2.5), mobile stock trading (step 1.2), and the Uber (steps 2.6, 2.7, 3.2, 3.3, R2.9), Shopify (steps 2.4, 3.4, R1.8), Figma (step 3.4) and bot defense (step R1.6) case studies.