Transactional State — the Tx* Modules
Lock-free composable transactions in Effect v4 — how Effect.tx journals TxRef reads and writes, conflict detection via version checks, Effect.txRetry blocking semantics, the full Tx* module inventory, and when STM beats Ref or queues.
Two fibers need to move money between accounts: debit one, credit the other. With locks, you serialize everything and invent deadlock-avoidance rules. With raw refs, you get interleavings that lose money. STM (software transactional memory) offers a third way: write code as if it touches all variables atomically; the runtime detects conflicts and retries transparently.
v4 ships this natively in Effect. There is no separate STM monad anymore and
no STM<E, A> type to lift from — transactions are plain Effects, delimited
by Effect.tx, operating on TxRefs. This chapter documents exactly what’s
there, because the surface is leaner than v3 and honest about it.
The problem, precisely
Section titled “The problem, precisely”Invariant: a.balance + b.balance === 100. Two concurrent transfers:
transfer(100, A -> B) transfer(50, B -> A)read A = 100 read B = 100write A = 0 write B = 50 read A = ??? // 0? 100? torn?write B = 150 ...Any interleaving that lets one transfer read state another has half-written breaks conservation. Classic fixes:
- A global mutex around all balance mutations — correct, but every unrelated account operation serializes behind it.
- Per-account locks with an ordering discipline — finer-grained, but now lock ordering is a convention enforced by review, and one violation hangs the system.
- Optimistic concurrency: each transaction snapshots what it reads, computes privately, then commits only if nothing it read changed since the snapshot. Conflicts are rare → retries are cheap; no locks exist at all.
Effect v4 implements option 3.
The inventory
Section titled “The inventory”All transactional modules live at the top level of effect (no unstable path,
one version for the ecosystem). The family:
| Module | What it transactionalizes |
|---|---|
TxRef |
single value — the primitive everything else builds on |
TxHashMap / TxHashSet |
hash map / set (backed by a TxRef of immutable structures) |
TxChunk |
transactional sequence |
TxQueue |
transactional queue |
TxPubSub |
broadcast inside transactions |
TxPriorityQueue |
priority-ordered take |
TxDeferred |
one-shot cell whose completion participates in the transaction |
TxSemaphore |
counting semaphore with transactional acquire/release |
TxReentrantLock |
reentrant mutual exclusion |
TxSubscriptionRef |
ref + change stream, transactionally |
Note the naming shift from v3: STM/TRef/TQueue/… → TxRef,
TxQueue, …. The collections mirror their plain counterparts’ APIs (get,
set, modify, size, …), so learning one is learning them all.
How transactions actually work
Section titled “How transactions actually work”This section describes the real mechanism — worth understanding because its edge cases define your obligations as the author.
The journal
Section titled “The journal”Every fiber running inside Effect.tx carries a Transaction service
containing two fields: a journal (Map<TxRef, {version, value}>) and a
retry flag. Each entry records, per accessed TxRef: the version observed
at first access, and the pending new value if written.
- Reads are recorded too.
TxRef.get(ref)journals(ref.version)without a write. If another transaction later changes a value you merely looked at, you’re in conflict — this is what makes invariants spanning multiple refs safe (the balance-conservation check reads both accounts, so any concurrent change to either triggers retry). - Writes go to the journal, not the ref. Inside the transaction, subsequent reads see your own writes; outside observers see nothing yet.
Commit and conflict detection
Section titled “Commit and conflict detection”When the outermost Effect.tx body completes:
- Validate: for every journaled ref, is
ref.versionstill equal to the recorded snapshot version? Any mismatch = conflict. - No conflict + success: commit — write each journaled value into its ref, bump its version, wake any transactions waiting on those refs.
- Conflict, or failure, or retry-flag set: discard the journal and run the body again (failure exits instead — see below).
The whole validation-commit step is safe because it runs synchronously on one fiber: JS’s single thread is the commit lock. There are no locks anywhere in userland; consistency comes from check-then-write being atomic on the event loop.
Retry semantics
Section titled “Retry semantics”Two distinct triggers cause a re-run:
- Conflict-driven: another transaction committed a change to any ref you journaled before you finished. Your result is discarded and the body reruns against fresh state.
- Explicit blocking: the body calls
Effect.txRetry(usually after finding data not-yet-satisfactory — an empty queue, insufficient balance). This sets the flag and aborts the current attempt; crucially, the transaction then waits until some accessed ref actually changes (it registers itself on each journaled ref’s waiter list and sleeps) rather than busy-looping. When a committing transaction touches one of those refs, waiters are rescheduled.
That second behavior is the crown jewel of STM: blocking until state makes progress possible, with no condition variables and no missed-wakeup bugs.
The lifecycle at a glance
Section titled “The lifecycle at a glance”| Phase | What happens |
|---|---|
Enter outermost Effect.tx |
create journal + retry flag; provide as Transaction service |
TxRef.get |
record (ref → {version: ref.version}) if absent; return current (journal-aware) value |
TxRef.set/update/modify |
journal the new value under the recorded version; observers outside see nothing yet |
nested Effect.tx |
joins — reuses the same journal and retry flag |
body calls Effect.txRetry |
set flag, abort attempt; park on all journaled refs until one changes |
| body completes | validate every journaled version against the live refs |
| consistent + success | commit: write values, bump versions, wake waiters |
| inconsistent / flagged | discard journal, rerun the body from scratch |
| body fails (typed error/die) | discard journal, propagate — nothing commits |
import { Effect, Fiber, TxRef } from "effect"
const awaitBalance = Effect.gen(function* () { const balance = yield* TxRef.make(0)
const program = yield* Effect.forkChild( Effect.tx(Effect.gen(function* () { const b = yield* TxRef.get(balance) if (b < 100) return yield* Effect.txRetry // park until balance changes return b })) )
yield* Effect.sleep("10 millis") yield* Effect.tx(TxRef.update(balance, (b) => b + 100)) // wakes the parked tx return yield* Fiber.join(program) // => 100})Nesting and composition
Section titled “Nesting and composition”Effect.tx inside an active transaction does not open a nested boundary —
it reuses the current journal and retry state. Composition falls out naturally:
call helper functions that each wrap themselves in Effect.tx; compose ten of
them; the outermost call commits everything atomically or nothing. This is the
property impossible with per-variable locks and the reason STM exists.
TxRef API tour
Section titled “TxRef API tour”import { Effect, TxRef } from "effect"
const counter = yield* TxRef.make(0)
// read — journals the versionconst current = yield* TxRef.get(counter)
// write — dual: data-last or data-firstyield* TxRef.set(counter, 42)yield* TxRef.set(42)(counter)yield* TxRef.update(counter, (n) => n + 1)
// modify — compute a result AND a new value in one journaled stepconst previous = yield* TxRef.modify(counter, (n) => [n, n + 1] as const)// ^ [resultA, newValue]: yield the old, store n+1
// sync construction outside Effect contextsconst cfg = TxRef.makeUnsafe({ attempts: 0 })get and update are literally defined in terms of modify — one primitive,
four ergonomic shapes. All require an ambient transaction when executed
directly (wrap them in Effect.tx); used inside a bigger transaction they
join it.
Worked example: atomic transfers
Section titled “Worked example: atomic transfers”Three accounts, invariant over total supply, concurrent transfers, plus a blocking withdrawal:
import { Context, Effect, Layer, Schema, TxRef } from "effect"
class Accounts extends Context.Service< Accounts, { readonly transfer: ( from: string, to: string, amountCents: number ) => Effect.Effect<void, InsufficientFunds> readonly deposit: (to: string, amountCents: number) => Effect.Effect<void> readonly withdrawBlocking: ( from: string, amountCents: number ) => Effect.Effect<number> // parks until funds exist }>()("Accounts") {}
class InsufficientFunds extends Schema.TaggedError<InsufficientFunds>()( "InsufficientFunds", { account: Schema.String, needed: Schema.Number }) {}
export const AccountsLive = Layer.effect( Accounts, Effect.gen(function* () { const ledger = yield* TxRef.make(new Map<string, number>([ ["alice", 100_00], ["bob", 50_00], ["treasury", 500_00] ]))
const adjust = (account: string, delta: number) => TxRef.modify(ledger, (map) => { const balance = (map.get(account) ?? 0) + delta const next = new Map(map) next.set(account, balance) return [balance, next] })
const lookup = (account: string) => Effect.map(TxRef.get(ledger), (m) => m.get(account) ?? 0)
return { deposit: (to, cents) => Effect.tx(adjust(to, cents).pipe(Effect.asVoid)),
transfer: (from, to, cents) => Effect.tx(Effect.gen(function* () { const from_ = yield* lookup(from) if (from_ < cents) { return yield* new InsufficientFunds({ account: from, needed: cents }) } yield* adjust(from, -cents) yield* adjust(to, cents) // both adjustments share one journal — they commit together or not at all })),
withdrawBlocking: (from, cents) => Effect.tx(Effect.gen(function* () { const balance = yield* lookup(from) if (balance < cents) { return yield* Effect.txRetry // park until any ledger commit } yield* adjust(from, -cents) return balance })) } }))
const day = Effect.gen(function* () { const accounts = yield* Accounts
// race these — no lost updates regardless of interleaving yield* Effect.all([ accounts.transfer("alice", "bob", 30_00), accounts.transfer("bob", "treasury", 20_00), accounts.deposit("alice", 10_00), accounts.withdrawBlocking("treasury", 5_00) ], { concurrency: "unbounded", discard: true })}).pipe(Effect.provide(AccountsLive))Walk through the interesting failure mode: transfer-1 journals reads of
ledger (via both adjust calls), computes debited+credited states, and just
before commit, deposit’s transaction commits its own version. Transfer-1’s
validation sees the version bump, discards, and reruns against post-deposit
state — recomputing balances from truth rather than patching stale numbers.
Total supply stays conserved without either party holding a lock across the
critical region.
Beyond refs: the collections in practice
Section titled “Beyond refs: the collections in practice”The wrappers matter because they give structure-aware conflict semantics and
ergonomics. TxHashMap is representative — full Map-like API, all
transaction-aware:
import { Effect, TxHashMap } from "effect"
const cache = yield* TxHashMap.empty<string, number>()
const memoized = Effect.fn("memoized")(function* (key: string) { return yield* Effect.tx(Effect.gen(function* () { const hit = yield* TxHashMap.get(cache, key) if (hit._tag === "Some") return hit.value const computed = key.length * 7 // pretend expense yield* TxHashMap.set(cache, key, computed) return computed }))})
const size = yield* TxHashMap.size(cache)const keys = yield* TxHashMap.keys(cache)Also available: setMany, removeMany, union, modifyAt, snapshot
(read the whole structure into an immutable HashMap — useful for reporting),
plus iteration over entries/values.
TxSemaphore gives transactional permits — acquire can block via the
retry mechanism, so “wait until resources free up” composes with other
transactional conditions instead of needing a separate coordination channel:
import { Effect, TxSemaphore } from "effect"
const slots = yield* TxSemaphore.make(3)
const claimSlot = Effect.tx(TxSemaphore.acquire(slots))const tryClaim = Effect.tx(TxSemaphore.tryAcquire(slots)) // boolean, never blocksconst release = Effect.tx(TxSemaphore.release(slots))A transactional work queue
Section titled “A transactional work queue”TxQueue mirrors the plain queue’s shape (bounded/unbounded/dropping/
sliding, offer/take/poll/takeAll/takeN) but its operations are
journal-aware. The difference shows when you combine queue operations with
other state inside one transaction:
import { Effect, TxHashMap, TxQueue } from "effect"
interface Task { readonly id: string; readonly run: string }
const scheduler = Effect.gen(function* () { const work = yield* TxQueue.unbounded<Task>() const inFlight = yield* TxHashMap.empty<string, Task>()
// atomically: pull a task AND record it as in-flight. // with plain structures these are two steps another fiber could interleave; // here a conflicting commit simply retries both halves together const claimNext = Effect.tx(Effect.gen(function* () { const task = yield* TxQueue.take(work) // blocks via retry when empty yield* TxHashMap.set(inFlight, task.id, task) return task }))
const complete = (task: Task) => Effect.tx(Effect.gen(function* () { yield* TxHashMap.remove(inFlight, task.id) yield* Effect.log(`done: ${task.id}`) }))
return { claimNext, complete }})Note what replaced hand-rolled synchronization: TxQueue.take on an empty
queue parks the whole transaction (txRetry semantics) until some other
transaction offers work — no condition variable, no polling.
The remaining modules follow the same recipe. TxDeferred is a one-shot cell
whose completion is journaled (succeed/fail/done write through the
transaction; poll reads an Option<Result<A, E>>) — so “complete X only if
Y also commits” is expressible. TxReentrantLock, TxPriorityQueue,
TxPubSub, TxChunk, and TxSubscriptionRef round out the family with the
transactional takes on their namesakes; their APIs mirror the plain modules,
so the inventory table above plus one source file each is all you need.
STM vs the alternatives
Section titled “STM vs the alternatives”| Concern | TxRef + Effect.tx |
Ref / SynchronizedRef |
Queue-based serialization |
|---|---|---|---|
| Multi-variable atomicity | native — one journal | manual: one Ref holding composite state, or careful locking | single consumer owns all state |
| Blocking until state allows progress | Effect.txRetry — no waker plumbing |
poll/re-check loops | natural (consumers wait on take) |
| Contention behavior | optimistic: conflicts retry | last-write-wins per cell (Ref); serialized effectful updates (SynchronizedRef) |
fully serialized |
| Composability of helpers | excellent — nested tx joins |
poor — helpers must agree on locking | requires routing all mutations through messages |
| Overhead on uncontended paths | lowest — a few map ops | low | highest — message indirection |
| Failure modes | side effects re-run (keep bodies pure) | races if composite state updated non-atomically | backpressure, queue sizing |
Decision heuristics:
- Single mutable value, simple updates →
Ref(orSynchronizedRefwhen updates need effects). - Invariants spanning multiple values, or wait-until-possible logic →
TxRef/STM. This is its irreplaceable niche: compositional atomicity with built-in blocking. - High-throughput command processing, rate limiting, work distribution → a queue with a single owning fiber (chapter 15); deterministic ordering, natural backpressure, trivially auditable.
Honest status notes
Section titled “Honest status notes”Things worth knowing before you build a bank on this:
- There is no standalone
STMmodule — noSTM<A, E>type, noatomically. Transactions live entirely inEffect(tx,txRetry) plus theTx*modules. Code searching for “STM” finds nothing; searchTxRef. - Errors don’t roll back committed siblings. A failing transaction simply doesn’t commit its own journal — there is no partial rollback of other transactions, and no isolation levels to configure. One journal, one snapshot, one retry loop.
- No automatic retry budgeting. A pathologically contended transaction retries indefinitely by design (that’s liveness-by-retry, standard STM). Bound hot spots structurally — finer-grained refs — rather than expecting a max-attempts knob.
- The
Transactionservice is inspectable (yield* Effect.Transaction) should you need custom instrumentation of journal size or retry counts.
For most applications, the practical summary is: reach for STM when you have
an invariant across more than one variable and want blocking semantics;
otherwise Ref covers single cells and queues cover pipelines.