Skip to content

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.

Invariant: a.balance + b.balance === 100. Two concurrent transfers:

transfer(100, A -> B) transfer(50, B -> A)
read A = 100 read B = 100
write 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:

  1. A global mutex around all balance mutations — correct, but every unrelated account operation serializes behind it.
  2. 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.
  3. 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.

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.

This section describes the real mechanism — worth understanding because its edge cases define your obligations as the author.

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.

When the outermost Effect.tx body completes:

  1. Validate: for every journaled ref, is ref.version still equal to the recorded snapshot version? Any mismatch = conflict.
  2. No conflict + success: commit — write each journaled value into its ref, bump its version, wake any transactions waiting on those refs.
  3. 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.

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.

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
})

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.

import { Effect, TxRef } from "effect"
const counter = yield* TxRef.make(0)
// read — journals the version
const current = yield* TxRef.get(counter)
// write — dual: data-last or data-first
yield* 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 step
const previous = yield* TxRef.modify(counter, (n) => [n, n + 1] as const)
// ^ [resultA, newValue]: yield the old, store n+1
// sync construction outside Effect contexts
const 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.

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.

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 blocks
const release = Effect.tx(TxSemaphore.release(slots))

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.

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 (or SynchronizedRef when 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.

Things worth knowing before you build a bank on this:

  • There is no standalone STM module — no STM<A, E> type, no atomically. Transactions live entirely in Effect (tx, txRetry) plus the Tx* modules. Code searching for “STM” finds nothing; search TxRef.
  • 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 Transaction service 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.