Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions docs/architecture/overview.md
Original file line number Diff line number Diff line change
Expand Up @@ -172,7 +172,7 @@ sequenceDiagram
participant M as Materializer
participant WAL as WriteAheadLog
participant PS as PageStore
participant GC as GroupCommitter (hstore-group-commit)
participant GC as GroupCommitter (leading committer)
participant F as ChangeFeed
W->>TX: createNode / insert / put ...
TX->>WS: Op.apply: path-copy trees (TreeWriter), emit MemberChange / SlotChange
Expand Down Expand Up @@ -201,14 +201,14 @@ sequenceDiagram
GC-->>TX: await() returns published generation
```

The ordering above is what makes recovery safe. Under `PAGE_REFERENCES` the WAL carries only `(pageId, length, crc32c)` for each new page. The flusher therefore makes data pages durable *before* the log records that reference them. Recovery replays committed transactions in order and stops at the first commit whose referenced pages fail verification (the *durable-prefix* rule in `Recovery.intact`). The full protocol and the crash-point matrix are in [wal-and-recovery.md](../transactions/wal-and-recovery.md).
The ordering above is what makes recovery safe. Under `PAGE_REFERENCES` the WAL carries only `(pageId, length, crc32c)` for each new page. The commit leader therefore makes data pages durable *before* the log records that reference them. Recovery replays committed transactions in order and stops at the first commit whose referenced pages fail verification (the *durable-prefix* rule in `Recovery.intact`). The full protocol and the crash-point matrix are in [wal-and-recovery.md](../transactions/wal-and-recovery.md).

## Threading model

| Thread | Kind | Created in | Role |
|---|---|---|---|
| Caller threads | any (the server uses virtual threads) | — | Run reads without locks. Run a transaction's operations without locks, then take `commitLock` only for append (encoding, WAL append, page writes). They then block in `GroupCommitter.await` outside the lock. |
| `hstore-group-commit` | platform, daemon | `GroupCommitter` constructor | Takes every queued `Pending` as one batch and calls `makeDurable(batchSize)` once: data sync in `PAGE_REFERENCES` mode, then WAL sync, both only under `Durability.SYNC`. It then runs each publication in generation order and wakes the waiters. One fsync pair covers the whole batch (`averageGroupCommit` in the statistics). A failure poisons the pipeline, and every later commit fails until restart. |
| (none; commit leader) | the committing thread | `GroupCommitter.await` | There is no group-commit thread. The first committer waiting on an unpublished commit leads: it takes every queued `Pending` as one batch and calls `makeDurable(batchSize)` once (data sync in `PAGE_REFERENCES` mode, then WAL sync, both only under `Durability.SYNC`). It then runs each publication in generation order and wakes the waiters. One fsync pair covers the whole batch (`averageGroupCommit` in the statistics). A failure poisons the pipeline, and every later commit fails until restart. |
| `hstore-maintenance` | virtual | `StorageEngine` constructor | Every 500 ms, under the `maintenance` lock (`tryLock`): checkpoint once the WAL written since the last checkpoint exceeds `checkpoint_wal_mb`, then reclaim retired segments that no pinned reader can still see. |
| `feed-subscriber` | virtual, one per subscription | `ChangeFeed.subscribe` | Tails the feed segments, waiting on a condition signalled by `append` (with a 250 ms timeout). Subscribers are `SemanticPlane`, `MaterializedViews` and `Statistics`. |
| Server connections | virtual, one per socket | `Server` (`Executors.newVirtualThreadPerTaskExecutor`) | Each connection runs a `Session`. |
Expand Down
3 changes: 1 addition & 2 deletions docs/development.md
Original file line number Diff line number Diff line change
Expand Up @@ -116,8 +116,7 @@ keep the matrix green. Add a `CrashPoint` when you introduce a new durability st
* Prefer records for data, sealed interfaces for closed hierarchies, and exhaustive `switch` with pattern
matching and record deconstruction. Use `_` for unused bindings.
* Streams and small pure functions where natural. Clear loops on hot paths.
* Virtual threads for request handling. Platform threads only where blocking I/O must not pin a carrier,
as in the group-commit flusher.
* Virtual threads for request handling and for the engine's background work.
* **Deterministic hashing.** Every fingerprint, summary or hash that is persisted or compared across
processes must be computed from explicit field values. Never use `Enum.hashCode()`, `Record.hashCode()`,
`Object.hashCode()` or `String.hashCode()` of a record's `toString` in an `EntryMeasure`, a summary or an
Expand Down
2 changes: 1 addition & 1 deletion docs/operations/logging.md
Original file line number Diff line number Diff line change
Expand Up @@ -96,7 +96,7 @@ This is every message the server emits. Placeholders are in `{braces}`.
| WARNING | hstore | `background maintenance failed` (with DETAIL) |
| LOG | txn | `created branch {name} (#{id}) from generation {g}` |
| LOG | txn | `branch #{id} is now {MERGED or DROPPED}` |
| ERROR | commit | `commit pipeline failed; the engine must be restarted` (with DETAIL). The flusher thread could not make a batch durable. All waiting and future committers fail. |
| ERROR | commit | `commit pipeline failed; the engine must be restarted` (with DETAIL). The committer leading a batch could not make it durable or publish it. All waiting and future committers fail. |
| LOG | feed | `change feed released {bytes} bytes; generations after {g} are retained`: retention after a checkpoint deleted whole feed segments (`feed_retention_generations`, lowered by holds). |
| WARNING | feed | `change feed subscriber stopped at generation {g}` (with DETAIL) |
| DEBUG | feed, wal | `directory sync is not supported for {dir}`: the file system rejected an `fsync` of a directory after a segment was created or deleted. |
Expand Down
4 changes: 2 additions & 2 deletions docs/storage/change-feed.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ sequenceDiagram
autonumber
participant TM as TransactionManager.append (commit lock)
participant WAL as WriteAheadLog
participant G as GroupCommitter flusher
participant G as GroupCommitter leader
participant CF as ChangeFeed
participant S as Subscribers
TM->>TM: event = CommitEvent(g, txn, wallTime, branch, members, slots)
Expand All @@ -53,7 +53,7 @@ sequenceDiagram
Two copies exist:

1. The **WAL `Feed` record** (type 5), written inside the commit before the `Commit` record. It is the durable copy and is covered by the commit's atomicity.
2. The **feed segments**, appended by the group-commit flusher during publication: after the commit is durable (under `SYNC`) and **immediately after** the generation becomes `current`. Appends happen on the single flusher thread in queue order, so the feed is in generation order.
2. The **feed segments**, appended during publication by the committer leading the batch: after the commit is durable (under `SYNC`) and **immediately after** the generation becomes `current`. Only one leader runs at a time and it publishes in queue order, so the feed is in generation order.

Appending after `current.set(next)` guarantees that a subscriber told about generation `g` can read it: any transaction it starts sees `g` or later. Subscribers such as `MaterializedViews` look up the atoms an event mentions, so an event that arrived before its commit was visible would make them skip changes they can't yet see. The opposite direction is not guaranteed: a reader whose snapshot is `g` may find `feed.lastGeneration()` still at `g - 1` while `g` is being published. Code that combines a snapshot with the feed replays only up to `min(snapshot, feed.lastGeneration())` (`SemanticPlane`, `MaterializedViews.refresh`) and picks up the rest from the next event. A committing client is acknowledged only after both steps, so its own commit is always in the feed by then. The commit is durable before either step, so an event is never rolled back by a process crash.

Expand Down
25 changes: 13 additions & 12 deletions docs/transactions/transactions.md
Original file line number Diff line number Diff line change
Expand Up @@ -223,7 +223,7 @@ Under `SNAPSHOT` both commits succeed (R replays `Insert(p)` into B, which L nev
|---|---|
| Read `head`, duplicate-request check | Executing the transaction's operations (concurrent, lock-free) |
| Fast/slow decision, read validation, replay | Spilling bulk-loaded subtrees |
| Materialization of new nodes into page images (`Materializer`) | `fsync` of data segments and WAL (`GroupCommitter` flusher) |
| Materialization of new nodes into page images (`Materializer`) | `fsync` of data segments and WAL (`GroupCommitter` leader) |
| WAL append of `Begin`, `Page`/`PageRef`, `Root`, `BranchMeta`, `Feed` records | Publication: feed append, `current.set`, `history.put` |
| Data page writes into the OS page cache | Waking the committing thread |
| WAL append of `Commit`, `head.set(next)`, `committer.submit(pending)` | |
Expand All @@ -234,7 +234,7 @@ Because the lock is released before waiting for durability, the next committer c

## 9. Group commit

`GroupCommitter` owns one platform daemon thread named `hstore-group-commit` and a FIFO `Deque<Pending>`.
`GroupCommitter` has no thread of its own. It keeps a FIFO `Deque<Pending>`, and the committers waiting on it take turns leading.

```text
Pending(generation, publication: Runnable, published: boolean)
Expand All @@ -243,22 +243,23 @@ Barrier.makeDurable(batchSize) -- TransactionManager.makeDurable

Protocol:

1. `submit(pending)` (called inside the commit lock, so queue order equals WAL commit order) appends to the queue and signals `submitted`.
2. The flusher waits until the queue is non-empty, copies **the whole queue** as the batch, releases the lock, and calls `barrier.makeDurable(batch.size())` once.
3. It then runs each batch member's `publication` in order: `feed.append(event)` for non-empty events, then `current.set(next)`, `history.put`, `trimHistory()`, and `CrashPoint.CATALOG_PUBLISH`. Appending to the feed **before** making the generation current gives readers an invariant: any thread that observes `current == g` also observes `feed.lastGeneration()` covering every non-empty commit up to `g`. The semantic plane relies on this for `FRESH` queries (see [change feed](../storage/change-feed.md#5-consumers)).
4. Under the lock it marks the batch `published`, pops it from the queue, updates `batches`/`flushedCommits`, and `signalAll`s `progressed`.
5. `await(pending)` blocks (uninterruptibly) until its own `published` flag is set.
1. `submit(pending)` (called inside the commit lock, so queue order equals WAL commit order) appends to the queue.
2. `await(pending)` loops until its own `published` flag is set. If another committer is leading, it waits (uninterruptibly) on `progressed`. Otherwise it becomes the **leader**: it copies **the whole queue** as the batch, releases the lock, and calls `barrier.makeDurable(batch.size())` once.
3. The leader then runs each batch member's `publication` in order: `current.set(next)`, `history.put`, then `feed.append(event)` for non-empty events, `trimHistory()`, and `CrashPoint.CATALOG_PUBLISH`. Making the generation current **before** appending to the feed means a subscriber can always read the commit it is told about (see [change feed](../storage/change-feed.md#2-when-events-are-written)).
4. It retakes the lock, marks the batch `published`, pops it from the queue, updates `batches`/`flushedCommits`, stops leading, and `signalAll`s `progressed`. A waiter whose commit wasn't in that batch becomes the next leader.

`makeDurable` for `Durability.SYNC`: if `walMode == PAGE_REFERENCES`, `pages.sync()` (force every dirty data segment); reach `DATA_SYNC`; `wal.sync()`; reach `WAL_SYNC`. For `Durability.ASYNC` it returns immediately: publication still happens in order on the flusher thread, but nothing is forced. A single `wal.sync()` covers every commit record appended so far, including records appended after the batch was copied; those commits are simply published in the next batch.
A lone committer therefore finishes on its own thread, with no hand-off to another thread. Committers that arrive while a leader is syncing queue up and are published together by the next leader, sharing one `fsync`. Only one leader runs at a time, so publication stays in generation order. The durability step runs on the committing thread, which in the server is usually a virtual thread; the JDK compensates for virtual threads blocked in file I/O by adding a carrier thread temporarily.

Failure handling: any `Throwable` from the barrier or a publication is stored in `failure`, logged at `ERROR`, and `progressed` is signalled; the flusher exits. Two checks consult it:
`makeDurable` for `Durability.SYNC`: if `walMode == PAGE_REFERENCES`, `pages.sync()` (force every dirty data segment); reach `DATA_SYNC`; `wal.sync()`; reach `WAL_SYNC`. For `Durability.ASYNC` it returns immediately: the leader still publishes in order, but nothing is forced. A single `wal.sync()` covers every commit record appended so far, including records appended after the batch was copied; those commits are simply published in the next batch.

Failure handling: any `Throwable` from the barrier or a publication is stored in `failure`, logged at `ERROR`, and `progressed` is signalled. The leader rethrows it, and every waiting and later committer gets it too. Two checks consult it:

| Check | Used by | Fails when |
|---|---|---|
| `failIfFailed()` | `await`, `drain` | a durability failure was recorded: rethrows an `Error` as-is (this is how `CrashPoint.SimulatedCrash` reaches the committing thread in tests), wraps anything else in `RETRYABLE_IO` |
| `failIfBroken()` | `submit` | `failIfFailed()` fails, or the committer is closed (`IllegalStateException: the commit pipeline is closed`) |

`close()` sets `closed`, wakes the flusher, and joins it. The flusher drains the remaining queue before exiting because its wait condition is `queue.isEmpty() && !closed`. Because `await` does not treat `closed` as an error, a commit that was submitted before shutdown is reported to its caller exactly as it ends up on disk: published and acknowledged, never "failed but durable". `submit` still rejects work after `close()`. It runs inside the commit lock after the `Commit` record has been appended and `head` advanced, so a commit rejected there is in the WAL and recovery would replay it (the same "outcome unknown" case as a crash at `COMMIT_APPEND`, see [WAL and recovery](wal-and-recovery.md#10-what-the-crash-tests-prove)). `TransactionManager.shutdown()` closes the committer while holding the commit lock, so no commit can reach `submit` after close through the normal paths; `StorageEngine.requireOpen()` rejects new transactions once the engine is closed.
`close()` sets `closed` and leads until the queue is empty, unless the pipeline has already failed. `drain()` leads the same way. Because `await` does not treat `closed` as an error, a commit that was submitted before shutdown is reported to its caller exactly as it ends up on disk: published and acknowledged, never "failed but durable". `submit` still rejects work after `close()`. It runs inside the commit lock after the `Commit` record has been appended and `head` advanced, so a commit rejected there is in the WAL and recovery would replay it (the same "outcome unknown" case as a crash at `COMMIT_APPEND`, see [WAL and recovery](wal-and-recovery.md#10-what-the-crash-tests-prove)). `TransactionManager.shutdown()` closes the committer while holding the commit lock, so no commit can reach `submit` after close through the normal paths; `StorageEngine.requireOpen()` rejects new transactions once the engine is closed.

`averageBatch()` (`flushedCommits / batches`) is exported as `Statistics.averageGroupCommit` and appears on the Studio dashboard.

Expand All @@ -271,7 +272,7 @@ sequenceDiagram
participant TM as TransactionManager (commitLock)
participant W as WriteAheadLog
participant D as PageStore
participant G as GroupCommitter flusher
participant G as GroupCommitter leader
participant F as ChangeFeed
C->>TM: commit(txn)
activate TM
Expand All @@ -298,7 +299,7 @@ sequenceDiagram
participant A as Txn A
participant B as Txn B
participant L as commitLock
participant G as flusher
participant G as commit leader
participant Disk as fsync
A->>L: lock, append records + Commit(g1), submit(g1), unlock
A->>G: await(g1)
Expand Down
10 changes: 5 additions & 5 deletions docs/transactions/wal-and-recovery.md
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,7 @@ flowchart TD
DW --> C[wal.append Commit t g]
C --> CA((COMMIT_APPEND))
CA --> H[head = g, committer.submit]
H --> F[flusher: makeDurable]
H --> F[commit leader: makeDurable]
F --> DS[[PAGE_REFERENCES and SYNC: pages.sync]]
DS --> DSP((DATA_SYNC))
DSP --> WS[[SYNC: wal.sync]]
Expand Down Expand Up @@ -285,7 +285,7 @@ After each completed checkpoint, `StorageEngine.timedCheckpoint` logs `checkpoin

1. Commits an ordered edge with 200 members (each a fresh node).
2. Arms a fault injector that throws `CrashPoint.SimulatedCrash` (an `Error`) the first time the chosen point is reached.
3. Begins a transaction that creates 300 more nodes and inserts each at position 0, then commits: the commit must throw `SimulatedCrash`. Crash points up to `COMMIT_APPEND` throw on the committing thread inside `append`; `DATA_SYNC`, `WAL_SYNC` and `CATALOG_PUBLISH` throw on the flusher thread, break the pipeline, and are rethrown to the waiting committer by `GroupCommitter.failIfFailed`.
3. Begins a transaction that creates 300 more nodes and inserts each at position 0, then commits: the commit must throw `SimulatedCrash`. Crash points up to `COMMIT_APPEND` throw on the committing thread inside `append`; `DATA_SYNC`, `WAL_SYNC` and `CATALOG_PUBLISH` throw on the committer leading the batch (here the committing thread itself), break the pipeline, and are rethrown to every waiting committer by `GroupCommitter.failIfFailed`.
4. `halt()`s the engine (no checkpoint, files closed as after a process crash; the OS page cache is preserved).
5. Reopens and asserts the edge has 200 or 500 members, never anything else; exactly `500` iff `point >= COMMIT_APPEND`; the node count equals the member count; every member has degree 1 (the reverse index agrees with the forward index). It then writes again, reopens, and checks that the post-crash write survived (recovery left the log appendable).

Expand All @@ -295,9 +295,9 @@ After each completed checkpoint, `StorageEngine.timedCheckpoint` logs `checkpoin
| `WAL_APPEND` | committer | `Begin`, page records, `Root`s, `Feed`; no `Commit` | 200, counted as unfinished | 200, counted as unfinished |
| `DATA_WRITE` | committer | as above plus data pages (orphaned) | 200 | 200 |
| `COMMIT_APPEND` | committer | `Commit` frame written, not synced | 500 (new) | 500, `PageRef`s verify |
| `DATA_SYNC` | flusher | data synced (refs mode only), WAL not synced | 500 | 500 |
| `WAL_SYNC` | flusher | WAL synced, not published | 500 | 500 |
| `CATALOG_PUBLISH` | flusher | feed appended (not forced), published to `current` | 500, feed frame kept or re-appended from the `Feed` record | 500, same |
| `DATA_SYNC` | commit leader | data synced (refs mode only), WAL not synced | 500 | 500 |
| `WAL_SYNC` | commit leader | WAL synced, not published | 500 | 500 |
| `CATALOG_PUBLISH` | commit leader | feed appended (not forced), published to `current` | 500, feed frame kept or re-appended from the `Feed` record | 500, same |

`COMMIT_APPEND` is the instructive row: the client received an error, yet the commit is recovered. This is the standard "commit outcome unknown" case of every database; clients that must know use `TxnOptions.withRequestId(...)` and retry, and the retry returns `DUPLICATE` with the original generation.

Expand Down
Loading
Loading