feat(log-store): add the object store WAL durable write path (#9320)

* feat(log-store): add the object store WAL durable write path

ObjectStoreLogStore can now write. append_batch admits entries into one
open batch and assigns object-sequence-major entry ids at admission. The
batch is sealed by size, by the flush interval, or before a region would
run past the position range, and sealed batches are uploaded with at most
four conditional creates in flight, started in sequence order. Created
objects are indexed in sequence order, and an append is acknowledged only
once its object is durable and indexed.

A transient create failure rolls back the failed batch and every later
batch unless a later object is already durable, in which case the store
poisons itself with a history-gap error. A conflicting object, an
encoding or catalog error, or taking the last representable sequence
poisons the store.

obsolete now goes through the actor and raises the sequence floor
together with the watermark, refusing with a retryable error while the
next sequence is not settled. stop drops the open batch and the batches
whose create has not started, and lets creates in flight finish.

Only the durable acknowledgement mode exists. A testing feature exposes
hooks to wait for admissions, seal the open batch, and hold or fail
creates.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* test(log-store): synchronize object store WAL write tests with the actor

Tests that assert nothing happened round-trip a command through the actor
instead of yielding the test task, the conflict test waits for the object
it expects, and the obsolete-behind-stop test holds the command channel
itself so the unanswered command is deterministic.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix(log-store): bound admitted WAL appends and never reuse a failed sequence

A create that reports an error may still have written its object, so the
sequences of failed batches are no longer handed out again: the next
batch keeps the sequence after the last sealed one and retries are
assigned new ids. A retry batched differently can no longer conflict
with that object. As the next sequence never moves back, the sequence
floor of obsolete only waits for an open batch that has handed out ids.

Appends now arrive on their own bounded channel, which the actor stops
reading while MAX_SEALED_BATCHES batches wait to become durable, so a
stalled object store holds callers back instead of growing the backlog,
while stop and obsolete are still handled.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix(log-store): reserve room for two seals before admitting a WAL append

One admission can seal the open batch before an append that would exhaust
its positions and then the append's own batch, so the actor takes an
append only while two more sealed batches fit under MAX_SEALED_BATCHES.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* feat(log-store): chain object store WAL objects and recover the latest chain

A conditional create can fail with an unknown outcome while its object is
stored, or still lands later, and Mito reuses the row sequences of a
failed append. Replaying such an object next to a later acknowledged one
can let the unacknowledged rows win after a restart.

Format version 2 gives every object header its writer's epoch, a link
to the object it extends (sequence and writer instance) and a header
CRC32, and allows objects without segments. Recovery replays only the
chain ending at the complete object with the largest epoch and sequence:
a link holds when its predecessor is present with the recorded writer
instance, or is missing below every present object. Objects off the
chain are orphans that are never replayed but keep their sequences.

Each open writes an empty object that starts an epoch above every
present object, linked to the recovered tip, before it accepts writes,
so a late object of an earlier instance never ends the chain. A start
object that meets an object of an earlier epoch moves to the next
sequence; one of an equal or later epoch fails the open.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix(log-store): break the WAL chain at every missing predecessor

Accepting a missing predecessor below every present object lets a late
object that lands below the chain change which links hold. Nothing
collects objects yet, so a missing predecessor now always breaks the
link, and recovery fails when objects are present but none completes a
chain.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix(log-store): keep the object store WAL format at version 1

The object store WAL has not been enabled anywhere, so no object in the
previous layout exists and the chained header can stay version 1.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* refactor(log-store): link object store WAL objects by epoch

Only one store instance writes under an epoch, since every open starts an
epoch above every present object and a start object that meets the same
or a later epoch fails the open. The epoch therefore identifies the
instance, and the random writer instance id is dropped from the header.
A link now records the sequence and the epoch of the object it extends,
and holds when the predecessor carries that epoch. The header shrinks to
46 bytes, and the store logs its epoch when it opens.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix(log-store): fail an open whose start object is already present

Without a random writer instance, two opens that recover the same objects
encode byte-identical start objects, and a conditional create treats the
same bytes as its own retry. A start object that is already present
therefore fails the open with a retryable error instead of letting both
opens claim the epoch; the next open counts the object and starts a
later epoch.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix(log-store): derive the WAL epoch from the claimed start sequence

Two opens can recover different listings when a late object lands across
a sequence gap between them, pick the same largest epoch plus one, and
both create their start objects under different sequences. The epoch of
an instance is now one above the sequence its start object claims, so a
successful create decides the epoch and no two instances share one. It
stays above every epoch recovery listed, and an object that carries an
epoch above the next sequence fails the open.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* fix(log-store): never run a held WAL create after the store is dropped

The create test hook ignored the closed hold channel, so a create parked
when the store was dropped could still run. It now returns without
creating. Drop the per-admission bookkeeping of issued entry ids, which
nothing reads, and move the parked I/O documentation to its helper.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* test(log-store): cover a WAL start object stored with an unknown outcome

Add a fault-injection test in which the create of the start object
stores the object but reports an error: the open fails without moving
to another sequence, and the next open counts the stored object and
claims a later epoch. Rename the test helper that writes a whole object
from a given header to put_object_with_header, and drop a needless
clone.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

* test(log-store): make the held-create drop test deterministic

Keep the actor running while the store drops the hold sender, so the
parked create always completes on the closed channel instead of racing
the actor's exit. Drop a comment that restates epoch_of.

Signed-off-by: jeremyhi <fengjiachun@gmail.com>

---------

Signed-off-by: jeremyhi <fengjiachun@gmail.com>
This commit is contained in:
jeremyhi
2026-09-24 11:07:29 +00:00
committed by GitHub
parent 3c2aac0a55
commit 194bc2fb3c
3 changed files with 2149 additions and 96 deletions
+3
View File
@@ -11,6 +11,9 @@ protoc-rust.workspace = true
[lints]
workspace = true
[features]
testing = []
[dependencies]
async-stream.workspace = true
async-trait.workspace = true
+43 -8
View File
@@ -350,6 +350,13 @@ pub enum Error {
location: Location,
},
#[snafu(display("Incomplete multipart WAL entry of region {}", region_id))]
IncompleteWalEntry {
region_id: RegionId,
#[snafu(implicit)]
location: Location,
},
#[snafu(display("WAL object sequence is exhausted, last sequence: {}", last_object_seq))]
WalObjectSequenceExhausted {
last_object_seq: u64,
@@ -357,6 +364,16 @@ pub enum Error {
location: Location,
},
#[snafu(display(
"WAL object sequence {} is not settled: the open batch has assigned entry ids under it",
object_seq
))]
WalObjectSequenceUnsettled {
object_seq: u64,
#[snafu(implicit)]
location: Location,
},
#[snafu(display(
"WAL entry positions of region {} in one object are exhausted",
region_id
@@ -374,6 +391,20 @@ pub enum Error {
location: Location,
},
#[snafu(display(
"WAL object {} was written by the earlier epoch {}, this store writes epoch {}",
path,
existing_epoch,
epoch
))]
StaleWalObject {
path: String,
existing_epoch: u64,
epoch: u64,
#[snafu(implicit)]
location: Location,
},
#[snafu(display(
"WAL object {} that starts epoch {} was already present, so the open cannot tell whether it wrote it",
path,
@@ -417,16 +448,15 @@ pub enum Error {
location: Location,
},
#[snafu(display("Object store WAL operation failed"))]
ObjectStoreWal {
source: Arc<Error>,
#[snafu(display("Object store WAL log store is stopped"))]
ObjectStoreWalStopped {
#[snafu(implicit)]
location: Location,
},
/// Appending to the object store WAL is unsupported.
#[snafu(display("Object store WAL operation is not supported"))]
UnsupportedObjectStoreWalOperation {
#[snafu(display("Object store WAL operation failed"))]
ObjectStoreWal {
source: Arc<Error>,
#[snafu(implicit)]
location: Location,
},
@@ -464,6 +494,7 @@ impl ErrorExt for Error {
| InvalidWalObjectStore { .. }
| MismatchedWalPrefix { .. }
| MismatchedWalRegion { .. }
| IncompleteWalEntry { .. }
| InvalidWalEntryRange { .. } => StatusCode::InvalidArguments,
StartWalTask { .. }
| StopWalTask { .. }
@@ -477,14 +508,15 @@ impl ErrorExt for Error {
| OrderedBatchProducerStopped { .. }
| WaitProduceResultReceiver { .. }
| WaitDumpIndex { .. }
| MetaLengthExceededLimit { .. } => StatusCode::Internal,
| MetaLengthExceededLimit { .. }
| ObjectStoreWalStopped { .. } => StatusCode::Internal,
CorruptedWalObject { .. }
| WalObjectConflict { .. }
| WalObjectSequenceExhausted { .. }
| WalEntryPositionExhausted { .. } => StatusCode::Unexpected,
WalObjectSequenceUnsettled { .. } => StatusCode::IllegalState,
UnsupportedObjectStoreWalOperation { .. } => StatusCode::Unsupported,
InvalidWalObject { source, .. } => source.status_code(),
ObjectStoreWal { source, .. } => source.status_code(),
@@ -493,6 +525,7 @@ impl ErrorExt for Error {
| WriteIndex { .. }
| ReadIndex { .. }
| WalObjectStore { .. }
| StaleWalObject { .. }
| UnconfirmedWalEpochStart { .. }
| Io { .. } => StatusCode::StorageUnavailable,
// Raft engine
@@ -529,6 +562,8 @@ impl ErrorExt for Error {
FetchEntry { .. }
| RaftEngine { .. }
| AddEntryLogBatch { .. }
| WalObjectSequenceUnsettled { .. }
| StaleWalObject { .. }
| UnconfirmedWalEpochStart { .. } => RetryHint::Retryable,
ProduceRecord { error, .. } => match error {
rskafka::client::producer::Error::Client(error) => {
File diff suppressed because it is too large Load Diff