From 03823a9a01bfafb3c2eaa0b5e6db3a795e37a14e Mon Sep 17 00:00:00 2001 From: jeremyhi Date: Mon, 28 Sep 2026 09:17:45 +0000 Subject: [PATCH] feat(log-store): add the enqueued acknowledgement mode to the object store WAL (#9358) * feat(log-store): add the enqueued acknowledgement mode to the object store WAL Add `ack_mode` (`durable` by default, or `enqueued`) and the backlog thresholds `max_unpersisted_bytes` and `max_unpersisted_age` to the object store WAL config, validated by the datanode and the store. In the `enqueued` mode `append_batch` returns on admission with the entry ids assigned and the object is created in the background. At a backlog threshold the next append is held back until an upload completes. A transient create failure is repeated under the same sequence with the same bytes; any conflicting object poisons the store. `stop` uploads the backlog, or returns the error that dropped it once stop began. `obsolete` clamps the watermark to the durable entry id, and an id the store handed out needs no sequence floor. Add `LogStore::wait_durable` with a default that returns at once. The object store WAL answers it once the region is durable and indexed through the entry id, and fails it after a backlog was lost. Signed-off-by: jeremyhi * fix(log-store): poison an enqueued store on permanent create failures Repeat a failed create in the enqueued mode only when the storage error is retryable; any other storage error poisons the store. A conflicting object poisons an enqueued store without reading its epoch, so a failed header read cannot turn the conflict into a retry. A durability wait for an id above the highest id the store handed out now waits for the handed-out ids of the region below it instead of returning at once. Signed-off-by: jeremyhi * fix(log-store): poison on permanent create failures after stop begins A create that fails with a storage error that is not retryable poisons an enqueued store even after stop began; only a transient failure drops the backlog without poisoning. Durability waiters whose callers stopped waiting are pruned before a new waiter is queued. The backlog age test no longer depends on a follow-up append finishing within the threshold. Signed-off-by: jeremyhi * fix(log-store): answer a durability wait once no earlier entry is pending A durability wait now returns once the region holds no entry at or below the target that is handed out but not durable, instead of waiting for the largest id handed out to the region. A later object of the region that is still being created no longer holds back a wait whose target it does not cover. Signed-off-by: jeremyhi * test(log-store): order the pending durability wait check after the actor Signed-off-by: jeremyhi * docs(store-api): state that the default wait_durable keeps each store's guarantee The default `LogStore::wait_durable` returns at once, which keeps each log store's own acknowledgement guarantee; Raft Engine with `sync_write = false` acknowledges before its periodic sync, so the documentation no longer claims that every entry id a caller holds is durable. The object store WAL configuration test now also serializes the new options and reads them back. Signed-off-by: jeremyhi * docs(log-store): limit the acknowledgement guarantees to the durable mode Signed-off-by: jeremyhi * fix(log-store): repeat enqueued creates that a retry layer marks persistent An object store wrapped in the OpenDAL retry layer reports a temporary error that outlasted its retries as persistent rather than temporary. The enqueued mode now repeats a create after any storage error that is not permanent, so a transient outage behind the retry layer no longer poisons the store and drops the acknowledged backlog; after stop began such a failure still drops the backlog without poisoning. Signed-off-by: jeremyhi * test(log-store): cover a persistent create failure after stop begins Signed-off-by: jeremyhi --------- Signed-off-by: jeremyhi --- config/config.md | 10 +- config/datanode.example.toml | 17 +- config/standalone.example.toml | 17 +- src/common/wal/src/config.rs | 20 +- src/common/wal/src/config/object_store.rs | 24 + src/datanode/src/datanode.rs | 32 + src/log-store/src/object_store_wal/batch.rs | 4 + src/log-store/src/object_store_wal/store.rs | 1069 +++++++++++++++++-- src/store-api/src/logstore.rs | 18 + 9 files changed, 1140 insertions(+), 71 deletions(-) diff --git a/config/config.md b/config/config.md index cebc5091b4c..89b3e13dcdd 100644 --- a/config/config.md +++ b/config/config.md @@ -125,7 +125,10 @@ | `wal.overwrite_entry_start_id` | Bool | `false` | Ignore missing entries during read WAL.
**It's only used when the provider is `kafka`**.

This option ensures that when Kafka messages are deleted, the system
can still successfully replay memtable data without throwing an
out-of-range error.
However, enabling this option might lead to unexpected data loss,
as the system will skip over missing entries instead of treating
them as critical errors. | | `wal.storage_provider` | String | `""` | The name of the storage provider that holds the WAL objects, an empty name selects the default object store.
**It's only used when the provider is `experimental_object_store`**. | | `wal.prefix` | String | `wal` | The path prefix of the WAL objects inside the storage provider.
The objects are written under `/datanodes//epochs/`, which is derived from this prefix.
**It's only used when the provider is `experimental_object_store`**. | -| `wal.flush_interval` | String | `100ms` | The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`.
Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters.
When writes are sparse, a shorter interval lowers the acknowledgement latency of appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node.
**It's only used when the provider is `experimental_object_store`**. | +| `wal.flush_interval` | String | `100ms` | The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`.
Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters.
When writes are sparse, a shorter interval lowers the acknowledgement latency of `durable` appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node, and in `enqueued` mode the backlog thresholds can seal earlier. An `enqueued` append does not wait for its batch to seal once it is admitted, but the backlog thresholds can delay admission.
**It's only used when the provider is `experimental_object_store`**. | +| `wal.ack_mode` | String | `durable` | When an append to the object store WAL returns.
- `durable`: an append returns after the object holding its entries is durable (the default).
- `enqueued`: an append returns once its entries are admitted and their ids are assigned, and the object is created in the background. A crash loses the unpersisted backlog.
**It's only used when the provider is `experimental_object_store`**. | +| `wal.max_unpersisted_bytes` | String | `64MB` | The size of the unpersisted backlog at which new appends stall until an upload completes, in `enqueued` mode.
Nothing is dropped and nothing is rejected; the threshold does not bound what a crash during an outage can lose.
**It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. | +| `wal.max_unpersisted_age` | String | `8s` | The age of the oldest unpersisted entry at which new appends stall until an upload completes, in `enqueued` mode.
**It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. | | `wal.on_corrupted_segment` | String | `skip` | What a read does with a segment that still does not decode after a second fetch, because its checksum does not match or its content disagrees with its footer entry.
- `skip`: the segment is skipped and recorded as a WAL hole of its region, a metric is incremented and a warning is logged; the other regions of the object are unaffected (the default).
- `fail`: the read fails, so the region does not open.
**It's only used when the provider is `experimental_object_store`**. | | `metadata_store` | -- | -- | Metadata storage options. | | `metadata_store.file_size` | String | `64MB` | The size of the metadata store log file. | @@ -599,7 +602,10 @@ | `wal.overwrite_entry_start_id` | Bool | `false` | Ignore missing entries during read WAL.
**It's only used when the provider is `kafka`**.

This option ensures that when Kafka messages are deleted, the system
can still successfully replay memtable data without throwing an
out-of-range error.
However, enabling this option might lead to unexpected data loss,
as the system will skip over missing entries instead of treating
them as critical errors. | | `wal.storage_provider` | String | `""` | The name of the storage provider that holds the WAL objects, an empty name selects the default object store.
**It's only used when the provider is `experimental_object_store`**. | | `wal.prefix` | String | `wal` | The path prefix of the WAL objects inside the storage provider.
The objects are written under `/datanodes//epochs/`, which is derived from this prefix.
**It's only used when the provider is `experimental_object_store`**. | -| `wal.flush_interval` | String | `100ms` | The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`.
Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters.
When writes are sparse, a shorter interval lowers the acknowledgement latency of appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node.
**It's only used when the provider is `experimental_object_store`**. | +| `wal.flush_interval` | String | `100ms` | The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`.
Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters.
When writes are sparse, a shorter interval lowers the acknowledgement latency of `durable` appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node, and in `enqueued` mode the backlog thresholds can seal earlier. An `enqueued` append does not wait for its batch to seal once it is admitted, but the backlog thresholds can delay admission.
**It's only used when the provider is `experimental_object_store`**. | +| `wal.ack_mode` | String | `durable` | When an append to the object store WAL returns.
- `durable`: an append returns after the object holding its entries is durable (the default).
- `enqueued`: an append returns once its entries are admitted and their ids are assigned, and the object is created in the background. A crash loses the unpersisted backlog.
**It's only used when the provider is `experimental_object_store`**. | +| `wal.max_unpersisted_bytes` | String | `64MB` | The size of the unpersisted backlog at which new appends stall until an upload completes, in `enqueued` mode.
Nothing is dropped and nothing is rejected; the threshold does not bound what a crash during an outage can lose.
**It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. | +| `wal.max_unpersisted_age` | String | `8s` | The age of the oldest unpersisted entry at which new appends stall until an upload completes, in `enqueued` mode.
**It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. | | `wal.on_corrupted_segment` | String | `skip` | What a read does with a segment that still does not decode after a second fetch, because its checksum does not match or its content disagrees with its footer entry.
- `skip`: the segment is skipped and recorded as a WAL hole of its region, a metric is incremented and a warning is logged; the other regions of the object are unaffected (the default).
- `fail`: the read fails, so the region does not open.
**It's only used when the provider is `experimental_object_store`**. | | `query` | -- | -- | The query engine options. | | `query.parallelism` | Integer | `0` | Parallelism of the query engine.
Default to 0, which means the number of CPU cores. | diff --git a/config/datanode.example.toml b/config/datanode.example.toml index 253cc3f2a57..82b93fd6ccc 100644 --- a/config/datanode.example.toml +++ b/config/datanode.example.toml @@ -231,10 +231,25 @@ overwrite_entry_start_id = false ## The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`. ## Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters. -## When writes are sparse, a shorter interval lowers the acknowledgement latency of appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node. +## When writes are sparse, a shorter interval lowers the acknowledgement latency of `durable` appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node, and in `enqueued` mode the backlog thresholds can seal earlier. An `enqueued` append does not wait for its batch to seal once it is admitted, but the backlog thresholds can delay admission. ## **It's only used when the provider is `experimental_object_store`**. #+ flush_interval = "100ms" +## When an append to the object store WAL returns. +## - `durable`: an append returns after the object holding its entries is durable (the default). +## - `enqueued`: an append returns once its entries are admitted and their ids are assigned, and the object is created in the background. A crash loses the unpersisted backlog. +## **It's only used when the provider is `experimental_object_store`**. +#+ ack_mode = "durable" + +## The size of the unpersisted backlog at which new appends stall until an upload completes, in `enqueued` mode. +## Nothing is dropped and nothing is rejected; the threshold does not bound what a crash during an outage can lose. +## **It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. +#+ max_unpersisted_bytes = "64MB" + +## The age of the oldest unpersisted entry at which new appends stall until an upload completes, in `enqueued` mode. +## **It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. +#+ max_unpersisted_age = "8s" + ## What a read does with a segment that still does not decode after a second fetch, because its checksum does not match or its content disagrees with its footer entry. ## - `skip`: the segment is skipped and recorded as a WAL hole of its region, a metric is incremented and a warning is logged; the other regions of the object are unaffected (the default). ## - `fail`: the read fails, so the region does not open. diff --git a/config/standalone.example.toml b/config/standalone.example.toml index fad4f64f08a..e5eaa86c4c5 100644 --- a/config/standalone.example.toml +++ b/config/standalone.example.toml @@ -414,10 +414,25 @@ overwrite_entry_start_id = false ## The interval of flushing buffered entries to the object store, at least `10ms`, defaults to `100ms`. ## Each non-empty timer-triggered flush creates one object, and a batch that reaches `max_batch_bytes` is flushed immediately, so under sustained load the batch seals on size and the interval no longer matters. -## When writes are sparse, a shorter interval lowers the acknowledgement latency of appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node. +## When writes are sparse, a shorter interval lowers the acknowledgement latency of `durable` appends and raises the number of object requests: timer-triggered sealing creates at most one object per interval per node, and in `enqueued` mode the backlog thresholds can seal earlier. An `enqueued` append does not wait for its batch to seal once it is admitted, but the backlog thresholds can delay admission. ## **It's only used when the provider is `experimental_object_store`**. #+ flush_interval = "100ms" +## When an append to the object store WAL returns. +## - `durable`: an append returns after the object holding its entries is durable (the default). +## - `enqueued`: an append returns once its entries are admitted and their ids are assigned, and the object is created in the background. A crash loses the unpersisted backlog. +## **It's only used when the provider is `experimental_object_store`**. +#+ ack_mode = "durable" + +## The size of the unpersisted backlog at which new appends stall until an upload completes, in `enqueued` mode. +## Nothing is dropped and nothing is rejected; the threshold does not bound what a crash during an outage can lose. +## **It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. +#+ max_unpersisted_bytes = "64MB" + +## The age of the oldest unpersisted entry at which new appends stall until an upload completes, in `enqueued` mode. +## **It's only used when the provider is `experimental_object_store` and `ack_mode` is `enqueued`**. +#+ max_unpersisted_age = "8s" + ## What a read does with a segment that still does not decode after a second fetch, because its checksum does not match or its content disagrees with its footer entry. ## - `skip`: the segment is skipped and recorded as a WAL hole of its region, a metric is incremented and a warning is logged; the other regions of the object are unaffected (the default). ## - `fail`: the read fails, so the region does not open. diff --git a/src/common/wal/src/config.rs b/src/common/wal/src/config.rs index f79a28607cc..0cbebfb20ab 100644 --- a/src/common/wal/src/config.rs +++ b/src/common/wal/src/config.rs @@ -138,7 +138,7 @@ mod tests { use super::*; use crate::TopicSelectorType; - use crate::config::object_store::CorruptedSegmentAction; + use crate::config::object_store::{AckMode, CorruptedSegmentAction}; use crate::config::{DatanodeKafkaConfig, MetasrvKafkaConfig}; #[test] @@ -294,6 +294,9 @@ mod tests { assert_eq!(config.prefix, "wal"); assert_eq!(config.flush_interval, Duration::from_millis(100)); assert_eq!(config.max_batch_bytes, ReadableSize::mb(8)); + assert_eq!(config.ack_mode, AckMode::Durable); + assert_eq!(config.max_unpersisted_bytes, ReadableSize::mb(64)); + assert_eq!(config.max_unpersisted_age, Duration::from_secs(8)); assert_eq!(config.on_corrupted_segment, CorruptedSegmentAction::Skip); assert!(matches!( MetasrvWalConfig::try_from(datanode_wal_config).unwrap_err(), @@ -306,6 +309,9 @@ mod tests { prefix = "cluster-a/wal" flush_interval = "500ms" max_batch_bytes = "4MB" + ack_mode = "enqueued" + max_unpersisted_bytes = "32MB" + max_unpersisted_age = "4s" on_corrupted_segment = "fail" "#; let datanode_wal_config: DatanodeWalConfig = toml::from_str(toml_str).unwrap(); @@ -314,12 +320,24 @@ mod tests { prefix: "cluster-a/wal".to_string(), flush_interval: Duration::from_millis(500), max_batch_bytes: ReadableSize::mb(4), + ack_mode: AckMode::Enqueued, + max_unpersisted_bytes: ReadableSize::mb(32), + max_unpersisted_age: Duration::from_secs(4), on_corrupted_segment: CorruptedSegmentAction::Fail, }; assert_eq!( datanode_wal_config, DatanodeWalConfig::ObjectStore(expected) ); + let serialized = toml::to_string(&datanode_wal_config).unwrap(); + let table: toml::Table = toml::from_str(&serialized).unwrap(); + assert_eq!(table["ack_mode"].as_str(), Some("enqueued")); + assert_eq!(table["max_unpersisted_bytes"].as_str(), Some("32MiB")); + assert_eq!(table["max_unpersisted_age"].as_str(), Some("4s")); + assert_eq!( + datanode_wal_config, + toml::from_str::(&serialized).unwrap() + ); // The persisted tag is not accepted as a config provider. let toml_str = r#" diff --git a/src/common/wal/src/config/object_store.rs b/src/common/wal/src/config/object_store.rs index 04a212fb53f..dd2c4544129 100644 --- a/src/common/wal/src/config/object_store.rs +++ b/src/common/wal/src/config/object_store.rs @@ -24,6 +24,18 @@ use serde::{Deserialize, Serialize}; /// node id and the generation from the metasrv. pub const STANDALONE_GENERATION: u64 = 0; +/// When an append to the object store WAL returns. +#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum AckMode { + /// An append returns after the object holding its entries is durable. + #[default] + Durable, + /// An append returns once its entries are admitted and their ids are + /// assigned; the object is created in the background. + Enqueued, +} + /// What a read does with a segment that still does not decode after a /// second fetch, because its checksum does not match or its content /// disagrees with its footer entry. @@ -56,6 +68,15 @@ pub struct ObjectStoreWalConfig { pub flush_interval: Duration, /// The max size of a single batch object, defaults to 8MiB. pub max_batch_bytes: ReadableSize, + /// When an append returns, defaults to `durable`. + pub ack_mode: AckMode, + /// Size of the unpersisted backlog at which appends stall until an + /// upload completes in `enqueued` mode, defaults to 64MiB. + pub max_unpersisted_bytes: ReadableSize, + /// Age of the oldest unpersisted entry at which appends stall until an + /// upload completes in `enqueued` mode, defaults to 8s. + #[serde(with = "humantime_serde")] + pub max_unpersisted_age: Duration, /// What a read does with a segment that still does not decode after a /// second fetch, because its checksum does not match or its content /// disagrees with its footer entry, defaults to `skip`. @@ -69,6 +90,9 @@ impl Default for ObjectStoreWalConfig { prefix: "wal".to_string(), flush_interval: Duration::from_millis(100), max_batch_bytes: ReadableSize::mb(8), + ack_mode: AckMode::Durable, + max_unpersisted_bytes: ReadableSize::mb(64), + max_unpersisted_age: Duration::from_secs(8), on_corrupted_segment: CorruptedSegmentAction::Skip, } } diff --git a/src/datanode/src/datanode.rs b/src/datanode/src/datanode.rs index ae4b0723e1a..b3d2c981f1e 100644 --- a/src/datanode/src/datanode.rs +++ b/src/datanode/src/datanode.rs @@ -780,6 +780,22 @@ fn validate_object_store_wal_config(config: &ObjectStoreWalConfig) -> Result<()> reason: "must be greater than 0", } ); + ensure!( + config.max_unpersisted_bytes.as_bytes() > 0, + InvalidObjectStoreWalConfigSnafu { + field: "max_unpersisted_bytes", + value: config.max_unpersisted_bytes.to_string(), + reason: "must be greater than 0", + } + ); + ensure!( + config.max_unpersisted_age > Duration::ZERO, + InvalidObjectStoreWalConfigSnafu { + field: "max_unpersisted_age", + value: format!("{:?}", config.max_unpersisted_age), + reason: "must be greater than 0", + } + ); Ok(()) } @@ -1163,12 +1179,28 @@ mod tests { }, "max_batch_bytes", ), + ( + ObjectStoreWalConfig { + max_unpersisted_bytes: ReadableSize(0), + ..Default::default() + }, + "max_unpersisted_bytes", + ), + ( + ObjectStoreWalConfig { + max_unpersisted_age: Duration::ZERO, + ..Default::default() + }, + "max_unpersisted_age", + ), ]; for (config, expected_field) in cases { let expected_value = match expected_field { "prefix" => config.prefix.clone(), "flush_interval" => format!("{:?}", config.flush_interval), + "max_unpersisted_bytes" => config.max_unpersisted_bytes.to_string(), + "max_unpersisted_age" => format!("{:?}", config.max_unpersisted_age), _ => config.max_batch_bytes.to_string(), }; let data_home = create_temp_dir("object-store-wal-invalid-config"); diff --git a/src/log-store/src/object_store_wal/batch.rs b/src/log-store/src/object_store_wal/batch.rs index b07ac8e8cb0..078af921def 100644 --- a/src/log-store/src/object_store_wal/batch.rs +++ b/src/log-store/src/object_store_wal/batch.rs @@ -145,6 +145,10 @@ impl OpenBatch { self.entries.is_empty() } + pub(crate) fn holds_region(&self, region_id: RegionId) -> bool { + self.positions.contains_key(®ion_id) + } + /// Returns the estimated size of the admitted entries. pub(crate) fn estimated_bytes(&self) -> usize { self.estimated_bytes diff --git a/src/log-store/src/object_store_wal/store.rs b/src/log-store/src/object_store_wal/store.rs index ee925ae0a38..13d4220c6ab 100644 --- a/src/log-store/src/object_store_wal/store.rs +++ b/src/log-store/src/object_store_wal/store.rs @@ -19,12 +19,12 @@ use std::fmt; use std::ops::Range; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex, PoisonError, RwLock}; -use std::time::Duration; +use std::time::{Duration, Instant}; use async_stream::try_stream; use bytes::Bytes; use common_telemetry::info; -use common_wal::config::object_store::ObjectStoreWalConfig; +use common_wal::config::object_store::{AckMode, ObjectStoreWalConfig}; use futures::future::BoxFuture; use futures::stream::FuturesUnordered; use futures::{StreamExt, TryStreamExt}; @@ -46,7 +46,7 @@ use crate::error::{ StaleWalObjectSnafu, UnconfirmedWalEpochStartSnafu, WalObjectSequenceExhaustedSnafu, WalObjectSequenceUnsettledSnafu, }; -use crate::object_store_wal::batch::{OBJECT_SEQ_LIMIT, OpenBatch, sequence_floor}; +use crate::object_store_wal::batch::{OBJECT_SEQ_LIMIT, OpenBatch, entry_id, sequence_floor}; use crate::object_store_wal::catalog::ObjectCatalog; use crate::object_store_wal::format::{ ChainLink, EncodedObject, FixedTrailer, FooterEntry, HEADER_LEN, Header, MIN_OBJECT_LEN, @@ -60,6 +60,9 @@ const COMMAND_BUFFER: usize = 1024; const APPEND_BUFFER: usize = 16; const MIN_FLUSH_INTERVAL: Duration = Duration::from_millis(10); const MAX_IN_FLIGHT_CREATES: usize = 4; +/// Delay before a create that failed transiently is attempted again in the +/// `enqueued` acknowledgement mode, where no caller is left to retry it. +const CREATE_RETRY_DELAY: Duration = Duration::from_millis(100); /// Number of sealed batches that may wait to be created or indexed. One /// admission can seal two batches, the open batch before an append that would /// exhaust its positions and the append's own batch, so the actor takes an @@ -79,13 +82,15 @@ const RECOVERY_TAIL_WINDOW: usize = 64 * 1024; /// when it reaches the size limit or the flush interval elapses and creates /// the object under the next sequence while it keeps admitting entries into /// the next batch; up to [`MAX_IN_FLIGHT_CREATES`] creates run at a time. -/// Objects are indexed in the catalog and their batches acknowledged in -/// sequence order, so an acknowledged entry never has a missing predecessor. -/// An append returns once its object is durable and indexed. While -/// [`MAX_SEALED_BATCHES`] batches wait to become durable no append is -/// admitted. +/// Objects are indexed in the catalog in sequence order. In the `durable` +/// acknowledgement mode an append returns once its object is durable and +/// indexed, so an acknowledged entry never has a missing predecessor; in the +/// `enqueued` mode it returns on admission and the object is created in the +/// background. While [`MAX_SEALED_BATCHES`] batches +/// wait to become durable no append is admitted. pub(crate) struct ObjectStoreLogStore { prefix: String, + ack_mode: AckMode, io: Arc, catalog: Arc>, obsolete_entry_ids: ObsoleteEntryIds, @@ -109,6 +114,7 @@ impl fmt::Debug for ObjectStoreLogStore { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { f.debug_struct("ObjectStoreLogStore") .field("prefix", &self.prefix) + .field("ack_mode", &self.ack_mode) .finish_non_exhaustive() } } @@ -145,6 +151,16 @@ impl ObjectStoreLogStore { ); let max_batch_bytes = positive_bytes(config.max_batch_bytes.as_bytes(), "max batch bytes")?; + let max_unpersisted_bytes = positive_bytes( + config.max_unpersisted_bytes.as_bytes(), + "max unpersisted bytes", + )?; + ensure!( + config.max_unpersisted_age > Duration::ZERO, + InvalidWalObjectStoreSnafu { + reason: "max unpersisted age is zero", + } + ); let Recovered { mut catalog, next_object_seq, @@ -196,12 +212,18 @@ impl ObjectStoreLogStore { stopped: stopped.clone(), command_rx, append_rx, + ack_mode: config.ack_mode, + max_unpersisted_bytes, + max_unpersisted_age: config.max_unpersisted_age, open_batch: OpenBatch::new(max_batch_bytes), issued_entry_ids: durable_entry_ids, pending: Vec::new(), sealed: VecDeque::new(), creates: FuturesUnordered::new(), + stalled: None, + durable_waiters: Vec::new(), stop: Vec::new(), + stop_error: None, next_object_seq: start .object_seq .checked_add(1) @@ -219,6 +241,7 @@ impl ObjectStoreLogStore { common_runtime::spawn_global(actor.run()); Ok(Arc::new(Self { prefix, + ack_mode: config.ack_mode, io, catalog, obsolete_entry_ids, @@ -235,6 +258,16 @@ impl ObjectStoreLogStore { })) } + /// Returns the largest entry id of the provider's region whose object is + /// durable and indexed, or zero for a region without such entries. In the + /// `enqueued` acknowledgement mode an entry id that an append returned + /// stays above this value until the object holding it is created. + pub(crate) fn durable_entry_id(&self, provider: &Provider) -> Result { + self.check_terminal()?; + let region_id = self.region_of(provider)?; + Ok(self.durable_entry_id_of(region_id)) + } + fn durable_entry_id_of(&self, region_id: RegionId) -> EntryId { let catalog = self.catalog.read().unwrap_or_else(PoisonError::into_inner); catalog.region_max_entry_id(region_id).unwrap_or(0) @@ -362,7 +395,8 @@ impl LogStore for ObjectStoreLogStore { type Error = Error; /// Stops the store. Creates in flight run to completion and acknowledge - /// their entries if they succeed. + /// their entries if they succeed; in the `enqueued` mode the remaining + /// backlog is uploaded first, and a failure of that upload is returned. async fn stop(&self) -> Result<()> { self.stopped.store(true, Ordering::Release); let (response_tx, response_rx) = oneshot::channel(); @@ -496,12 +530,17 @@ impl LogStore for ObjectStoreLogStore { .collect()) } - /// Moves the obsolete watermark of the region up to `entry_id` and makes - /// every id the region is assigned from now on greater than `entry_id`: - /// the sequence of the next object is raised above the object `entry_id` - /// names unless it is there already, so a watermark assigned under - /// another prefix, or one whose object this prefix no longer holds, is - /// never passed by a new id. Zero names no object and needs no floor. + /// Moves the obsolete watermark of the region up to `entry_id`. In the + /// `enqueued` acknowledgement mode the watermark never passes the durable + /// entry id: an entry that is not durable yet is replayed after a crash, + /// and hiding it would skip it. + /// + /// Every id the region is assigned from now on is greater than + /// `entry_id`, whatever watermark is recorded: the sequence of the next + /// object is raised above the object `entry_id` names unless it is there + /// already, so a watermark assigned under another prefix, or one whose + /// object this prefix no longer holds, is never passed by a new id. Zero + /// names no object and needs no floor. /// /// The watermark and the floor are applied together by the actor: a call /// cancelled before its command is queued publishes neither, while a @@ -509,7 +548,8 @@ impl LogStore for ObjectStoreLogStore { /// ids under the next sequence the floor cannot move, and the call fails /// with [`Error::WalObjectSequenceUnsettled`] and records neither. The /// watermark is kept in memory; the caller re-establishes it after a - /// restart. + /// restart. In the `enqueued` mode an `entry_id` this store handed out + /// needs no floor: it is never handed out again while the store runs. async fn obsolete( &self, provider: &Provider, @@ -517,10 +557,15 @@ impl LogStore for ObjectStoreLogStore { entry_id: EntryId, ) -> Result<()> { self.check_region(provider, region_id)?; + let watermark = match self.ack_mode { + AckMode::Durable => entry_id, + AckMode::Enqueued => entry_id.min(self.durable_entry_id_of(region_id)), + }; let (response_tx, response_rx) = oneshot::channel(); let command = Command::Obsolete { region_id, entry_id, + watermark, response: response_tx, }; let answered = match self.command_tx.send(command).await { @@ -533,7 +578,7 @@ impl LogStore for ObjectStoreLogStore { // still queued behind the stop: nothing is assigned an id any // more, so the watermark alone is consistent. None => { - record_obsolete(&self.obsolete_entry_ids, region_id, entry_id); + record_obsolete(&self.obsolete_entry_ids, region_id, watermark); Ok(()) } } @@ -564,9 +609,29 @@ impl LogStore for ObjectStoreLogStore { } fn latest_entry_id(&self, provider: &Provider) -> Result { + self.durable_entry_id(provider) + } + + /// Waits until every entry of the provider's region with an id at or + /// below `entry_id` is durable and indexed. Returns at once in the + /// `durable` acknowledgement mode, where a caller only holds durable ids. + async fn wait_durable(&self, provider: &Provider, entry_id: EntryId) -> Result<()> { self.check_terminal()?; let region_id = self.region_of(provider)?; - Ok(self.durable_entry_id_of(region_id)) + if entry_id <= self.durable_entry_id_of(region_id) { + return Ok(()); + } + let (response_tx, response_rx) = oneshot::channel(); + self.command_tx + .send(Command::WaitDurable { + region_id, + entry_id, + response: response_tx, + }) + .await + .ok() + .context(ObjectStoreWalStoppedSnafu)?; + response_rx.await.ok().context(ObjectStoreWalStoppedSnafu)? } } @@ -575,12 +640,19 @@ type AppendResponse = oneshot::Sender>; type QueuedAppend = (Vec, AppendResponse); enum Command { + /// Answered once the region is durable through `entry_id`. + WaitDurable { + region_id: RegionId, + entry_id: EntryId, + response: oneshot::Sender>, + }, /// Answered once the watermark is recorded and no id of the region at /// or below `entry_id` can be assigned, or with the reason neither was /// done. Obsolete { region_id: RegionId, entry_id: EntryId, + watermark: EntryId, response: oneshot::Sender>, }, Stop { @@ -598,9 +670,17 @@ struct PendingAppend { response: AppendResponse, } +/// A caller waiting until a region is durable through an entry id. +struct DurableWaiter { + region_id: RegionId, + entry_id: EntryId, + response: oneshot::Sender>, +} + /// Where the conditional create of a sealed batch stands. enum CreateState { - /// The create has not started because no slot was free. + /// The create has not started: no slot was free, or in the `enqueued` + /// mode an earlier attempt failed transiently and is repeated. Pending, InFlight, Created, @@ -612,6 +692,9 @@ struct SealedBatch { object_seq: u64, bytes: Bytes, footer: Vec, + first_admitted_at: Instant, + /// Number of creates that were attempted for the batch. + attempts: u32, waiters: Vec, #[cfg(any(test, feature = "testing"))] seal_waiters: Vec>>, @@ -619,6 +702,10 @@ struct SealedBatch { } impl SealedBatch { + fn is_in_flight(&self) -> bool { + matches!(self.state, CreateState::InFlight) + } + fn fail(self, error: impl Fn() -> Error) { for waiter in self.waiters { let _ = waiter.response.send(Err(error())); @@ -646,18 +733,22 @@ type CreateOutcome = (u64, Result); /// | Situation | Sequence | Waiters | Store | /// | --- | --- | --- | --- | /// | created, or identical retry | advances | acknowledged in order | healthy | -/// | transient error | never reused: a create that reported an error may have written its object | every batch that is not indexed fails, created ones included; retries are assigned new ids | healthy: the next batch links to the last indexed object, so no chain reaches a failed batch | -/// | an object of an earlier epoch holds the sequence | never reused | as for a transient error | healthy: that object can never be on the chain | -/// | an object of the same or a later epoch holds the sequence, encoding or catalog error | unchanged | every batch that is not indexed fails | poisoned | +/// | transient error, `durable` mode | never reused: a create that reported an error may have written its object | every batch that is not indexed fails, created ones included; retries are assigned new ids | healthy: the next batch links to the last indexed object, so no chain reaches a failed batch | +/// | transient error, `enqueued` mode | unchanged | already acknowledged | healthy: the create is repeated with the same bytes after [`CREATE_RETRY_DELAY`], which an identical retry accepts | +/// | transient error, `enqueued` mode after `stop` began | never reused | already acknowledged | the backlog is dropped and `stop` reports the error | +/// | an object of an earlier epoch holds the sequence, `durable` mode | never reused | as for a transient error | healthy: that object can never be on the chain | +/// | an object of the same or a later epoch holds the sequence, any conflicting object in the `enqueued` mode, encoding or catalog error | unchanged | every batch that is not indexed fails | poisoned | /// | created at the last representable sequence | cannot advance | acknowledged | poisoned: no later batch can be allocated a sequence | /// /// A create still in flight when its batch failed runs to completion and its /// outcome is ignored: whatever it stored is off the chain. After `stop` began -/// nothing is admitted and no create starts; creates in flight run to -/// completion and acknowledge if they succeed. +/// nothing is admitted and no create starts, except that the `enqueued` mode +/// uploads its backlog; creates in flight run to completion and acknowledge +/// if they succeed. /// /// Appends arrive on their own channel, which the actor does not read while -/// [`MAX_SEALED_BATCHES`] batches wait, so commands such as `stop` are handled +/// [`MAX_SEALED_BATCHES`] batches wait or an append is held back at a backlog +/// threshold of the `enqueued` mode, so commands such as `stop` are handled /// while admission is held back. struct Actor { io: Arc, @@ -667,16 +758,29 @@ struct Actor { stopped: Arc, command_rx: mpsc::Receiver, append_rx: mpsc::Receiver, + ack_mode: AckMode, + max_unpersisted_bytes: usize, + max_unpersisted_age: Duration, open_batch: OpenBatch, + /// Largest entry id ever handed out per region, whether it became + /// durable or failed. Unlike the accepted ids of the open batch it never + /// moves down. issued_entry_ids: HashMap, - /// Waiters of the open batch. + /// Waiters of the open batch in the `durable` mode. pending: Vec, /// Batches that are not durable yet, in sequence order. sealed: VecDeque, /// Creates in flight, including those of batches that already failed. creates: FuturesUnordered>, + /// The append held back in the `enqueued` mode while the unpersisted + /// backlog is at a threshold. Later appends wait in their channel. + stalled: Option, + durable_waiters: Vec, /// Callers of `stop`, answered once nothing is in flight. stop: Vec>>, + /// The failure `stop` reports in the `enqueued` mode once an + /// acknowledged backlog was dropped, recorded when it happens. + stop_error: Option>, /// Sequence of the next sealed batch, `None` once the sequence is /// exhausted. The open batch assigns its entry ids from it. It is above /// every sealed batch and never moves back while the store runs. @@ -705,7 +809,8 @@ impl Actor { loop { tokio::select! { _ = interval.tick() => { - // Nothing starts after stop began. + // Nothing starts after stop began; the `enqueued` mode + // sealed its backlog when stop was requested. if !self.is_stopped() { self.flush_open_batch(); } @@ -713,12 +818,15 @@ impl Actor { Some((object_seq, result)) = self.creates.next(), if !self.creates.is_empty() => { self.on_create_completed(object_seq, result); } - Some((entries, response)) = self.append_rx.recv(), if self.sealed.len() + 2 <= MAX_SEALED_BATCHES => { + Some((entries, response)) = self.append_rx.recv(), if self.stalled.is_none() && self.sealed.len() + 2 <= MAX_SEALED_BATCHES => { self.handle_append(entries, response); } command = self.command_rx.recv() => match command { - Some(Command::Obsolete { region_id, entry_id, response }) => { - self.handle_obsolete(region_id, entry_id, response); + Some(Command::WaitDurable { region_id, entry_id, response }) => { + self.handle_wait_durable(region_id, entry_id, response); + } + Some(Command::Obsolete { region_id, entry_id, watermark, response }) => { + self.handle_obsolete(region_id, entry_id, watermark, response); } Some(Command::Stop { response }) => { self.handle_stop(response); @@ -749,11 +857,17 @@ impl Actor { let _ = response.send(Err(shared(&error))); return; } + if self.ack_mode == AckMode::Enqueued && self.backlog_at_threshold() { + self.stalled = Some((entries, response)); + self.ensure_create_in_flight(); + return; + } self.admit(entries, response); } /// Admits `entries` into the open batch, which assigns their ids under - /// the next object sequence. The caller waits for the object. + /// the next object sequence. In the `durable` mode the caller waits for + /// the object, in the `enqueued` mode it is answered now. fn admit(&mut self, entries: Vec, response: AppendResponse) { // A region would run past the position range of the open batch: the // batch is sealed and the entries open the next one. The size limit @@ -785,10 +899,21 @@ impl Actor { return; } }; - self.pending.push(PendingAppend { - last_entry_ids, - response, - }); + for (region_id, entry_id) in &last_entry_ids { + self.issued_entry_ids + .entry(*region_id) + .and_modify(|issued| *issued = (*issued).max(*entry_id)) + .or_insert(*entry_id); + } + match self.ack_mode { + AckMode::Durable => self.pending.push(PendingAppend { + last_entry_ids, + response, + }), + AckMode::Enqueued => { + let _ = response.send(Ok(AppendBatchResponse { last_entry_ids })); + } + } #[cfg(any(test, feature = "testing"))] self.admitted_appends.send_modify(|count| *count += 1); if self.open_batch.should_seal() { @@ -796,6 +921,59 @@ impl Actor { } } + /// Returns true once the unpersisted backlog, the open batch and every + /// sealed batch that is not durable, reaches the size or the age threshold. + fn backlog_at_threshold(&self) -> bool { + let bytes = self.open_batch.estimated_bytes() + + self + .sealed + .iter() + .map(|batch| batch.bytes.len()) + .sum::(); + if bytes >= self.max_unpersisted_bytes { + return true; + } + let oldest = self + .sealed + .front() + .map(|batch| batch.first_admitted_at) + .or_else(|| self.open_batch.first_admitted_at()); + oldest.is_some_and(|admitted_at| admitted_at.elapsed() >= self.max_unpersisted_age) + } + + /// Seals the open batch when no create is in flight, so that a stalled + /// append has an upload to wait for. + fn ensure_create_in_flight(&mut self) { + if !self.sealed.iter().any(SealedBatch::is_in_flight) { + self.flush_open_batch(); + } + } + + /// Admits the stalled append once the backlog is below the thresholds and + /// the sealed batches leave room for the two an admission can seal. + fn release_stalled(&mut self) { + if self.stalled.is_none() { + return; + } + if self.is_stopped() { + if let Some((_, response)) = self.stalled.take() { + let _ = response.send(Err(ObjectStoreWalStoppedSnafu.build())); + } + return; + } + if self.backlog_at_threshold() { + // The append needs an upload to complete. + self.ensure_create_in_flight(); + return; + } + if self.sealed.len() + 2 > MAX_SEALED_BATCHES { + return; + } + if let Some((entries, response)) = self.stalled.take() { + self.admit(entries, response); + } + } + /// Seals the open batch as the object `next_object_seq` and starts its /// create when a slot is free. Returns whether a batch was sealed. fn flush_open_batch(&mut self) -> bool { @@ -816,7 +994,7 @@ impl Actor { return false; }; - let (entries, _) = self.open_batch.seal(); + let (entries, first_admitted_at) = self.open_batch.seal(); let header = Header { object_seq, epoch: self.epoch, @@ -854,6 +1032,8 @@ impl Actor { object_seq, bytes: encoded.bytes, footer: encoded.footer, + first_admitted_at, + attempts: 0, waiters: std::mem::take(&mut self.pending), #[cfg(any(test, feature = "testing"))] seal_waiters: Vec::new(), @@ -865,9 +1045,10 @@ impl Actor { /// Starts the creates of pending batches in sequence order while fewer /// than [`MAX_IN_FLIGHT_CREATES`] are in flight, counting the creates of - /// batches that already failed. Nothing starts once stop began. + /// batches that already failed. Nothing starts once stop began, except + /// the backlog of the `enqueued` mode. fn start_creates(&mut self) { - if self.is_stopped() { + if self.is_stopped() && self.ack_mode == AckMode::Durable { return; } let mut in_flight = self.creates.len(); @@ -880,15 +1061,25 @@ impl Actor { } batch.state = CreateState::InFlight; in_flight += 1; + let delay = if batch.attempts == 0 { + Duration::ZERO + } else { + CREATE_RETRY_DELAY + }; + batch.attempts += 1; let io = self.io.clone(); let object_seq = batch.object_seq; let bytes = batch.bytes.clone(); let epoch = self.epoch; + let ack_mode = self.ack_mode; #[cfg(any(test, feature = "testing"))] let mut creates_held = self.creates_held.clone(); #[cfg(any(test, feature = "testing"))] let creates_fail = self.creates_fail.clone(); self.creates.push(Box::pin(async move { + if !delay.is_zero() { + tokio::time::sleep(delay).await; + } // The store was dropped while the create was parked: it never // runs. #[cfg(any(test, feature = "testing"))] @@ -908,8 +1099,12 @@ impl Actor { }); return (object_seq, result); } + // Any conflict poisons an `enqueued` store, so only the + // `durable` mode reads the epoch of the existing object. let result = match io.put_if_absent(object_seq, bytes).await { - Err(error @ Error::WalObjectConflict { .. }) => { + Err(error @ Error::WalObjectConflict { .. }) + if ack_mode == AckMode::Durable => + { stale_conflict(io.as_ref(), object_seq, epoch, error).await } result => result, @@ -932,10 +1127,28 @@ impl Actor { }; match result { Ok(_) => self.sealed[index].state = CreateState::Created, + // Nobody is left to retry in the `enqueued` mode, so the store + // repeats a create that failed transiently under the same sequence + // with the same bytes, which an identical retry accepts. + Err(ref error) + if self.ack_mode == AckMode::Enqueued + && !self.is_stopped() + && is_transient(error) => + { + self.sealed[index].state = CreateState::Pending; + } // The object store did not confirm the object, which may still // exist or land later, or an earlier epoch holds the sequence and - // can never be on the chain. The caller retries the append itself. - Err(error @ (Error::WalObjectStore { .. } | Error::StaleWalObject { .. })) => { + // can never be on the chain. A caller of the `durable` mode + // retries the append itself; after stop began a transient failure + // drops the `enqueued` backlog. + Err(error @ Error::WalObjectStore { .. }) + if self.ack_mode == AckMode::Durable + || (self.is_stopped() && is_transient(&error)) => + { + self.roll_back(Arc::new(error)) + } + Err(error @ Error::StaleWalObject { .. }) if self.ack_mode == AckMode::Durable => { self.roll_back(Arc::new(error)) } Err(error) => { @@ -955,6 +1168,7 @@ impl Actor { } } self.start_creates(); + self.release_stalled(); } /// Indexes the created object at the front and acknowledges its waiters. @@ -988,9 +1202,45 @@ impl Actor { for waiter in batch.seal_waiters { let _ = waiter.send(Ok(())); } + self.resolve_durable_waiters(); true } + fn resolve_durable_waiters(&mut self) { + let waiters = std::mem::take(&mut self.durable_waiters); + for waiter in waiters { + if self.is_durable_through(waiter.region_id, waiter.entry_id) { + let _ = waiter.response.send(Ok(())); + } else { + self.durable_waiters.push(waiter); + } + } + } + + /// Returns true when no id of the region at or below `entry_id` waits to + /// become durable. An id of a batch that failed will never be durable and + /// no longer waits. + fn is_durable_through(&self, region_id: RegionId, entry_id: EntryId) -> bool { + self.lowest_pending_entry_id(region_id) + .is_none_or(|pending| pending > entry_id) + } + + /// Returns the lowest id of the region that was handed out and is not + /// durable yet. The sealed batches are in sequence order and the open + /// batch is above all of them. + fn lowest_pending_entry_id(&self, region_id: RegionId) -> Option { + self.sealed + .iter() + .flat_map(|batch| batch.footer.iter()) + .find(|entry| entry.region_id == region_id) + .map(|entry| entry.min_entry_id) + .or_else(|| { + self.next_object_seq + .filter(|_| self.open_batch.holds_region(region_id)) + .map(|object_seq| entry_id(object_seq, 1)) + }) + } + /// Fails every batch that is not indexed after a create failed /// transiently, including batches already created and batches whose /// create is in flight. Their sequences are not reused: a create that @@ -1013,14 +1263,20 @@ impl Actor { } self.reset_open_batch(); self.fail_unacknowledged(failure); + // An acknowledged backlog was dropped: `stop` reports it, whether + // its caller has arrived yet or not. + if self.ack_mode == AckMode::Enqueued { + self.stop_error.get_or_insert(error); + } } /// Records `error` as terminal and fails every waiter that is not /// acknowledged with it, or with the stopped error if the store was /// stopped meanwhile. Creates in flight run to completion, but their - /// outcome is ignored: the entries of an object they create were never - /// acknowledged, like those of a crash between creation and - /// acknowledgement. Returns the recorded error. + /// outcome is ignored: in the `durable` mode the entries of an object they + /// create were never acknowledged, like those of a crash between creation + /// and acknowledgement; in the `enqueued` mode the acknowledged backlog is + /// discarded and `stop` reports it. Returns the recorded error. fn poison(&mut self, error: Error) -> Arc { let error = set_terminal(&self.terminal_error, error); let stopped = self.is_stopped(); @@ -1036,6 +1292,9 @@ impl Actor { } self.reset_open_batch(); self.fail_unacknowledged(failure); + if self.ack_mode == AckMode::Enqueued { + self.stop_error.get_or_insert(error.clone()); + } error } @@ -1045,11 +1304,56 @@ impl Actor { self.open_batch.reset(); } - /// Fails the waiters of the open batch. + /// Fails the waiters of the open batch, the stalled append and the + /// durability waiters. fn fail_unacknowledged(&mut self, error: impl Fn() -> Error) { for pending in self.pending.drain(..) { let _ = pending.response.send(Err(error())); } + if let Some((_, response)) = self.stalled.take() { + let _ = response.send(Err(error())); + } + for waiter in self.durable_waiters.drain(..) { + let _ = waiter.response.send(Err(error())); + } + } + + fn handle_wait_durable( + &mut self, + region_id: RegionId, + entry_id: EntryId, + response: oneshot::Sender>, + ) { + if let Some(error) = terminal(&self.terminal_error) { + let _ = response.send(Err(shared(&error))); + return; + } + let durable = { + let catalog = self.catalog.read().unwrap_or_else(PoisonError::into_inner); + catalog.region_max_entry_id(region_id).unwrap_or(0) + }; + if entry_id <= durable { + let _ = response.send(Ok(())); + return; + } + // An acknowledged backlog was dropped: no entry that is not durable + // can be certified any more, whether it was in that backlog or not. + if let Some(error) = &self.stop_error { + let _ = response.send(Err(shared(error))); + return; + } + if self.is_durable_through(region_id, entry_id) { + let _ = response.send(Ok(())); + return; + } + // A caller that stopped waiting leaves a closed response behind. + self.durable_waiters + .retain(|waiter| !waiter.response.is_closed()); + self.durable_waiters.push(DurableWaiter { + region_id, + entry_id, + response, + }); } /// Makes sure no id of the region at or below `entry_id` is assigned @@ -1059,10 +1363,11 @@ impl Actor { &mut self, region_id: RegionId, entry_id: EntryId, + watermark: EntryId, response: oneshot::Sender>, ) { - let result = self.raise_sequence_floor(entry_id).map(|_| { - record_obsolete(&self.obsolete_entry_ids, region_id, entry_id); + let result = self.raise_sequence_floor(region_id, entry_id).map(|_| { + record_obsolete(&self.obsolete_entry_ids, region_id, watermark); }); let _ = response.send(result); } @@ -1071,7 +1376,7 @@ impl Actor { /// `entry_id` unless it is there already. The move fails while the open /// batch has handed out ids under the next sequence. A floor that does /// not fit an entry id poisons the store like an exhausted sequence. - fn raise_sequence_floor(&mut self, entry_id: EntryId) -> Result<()> { + fn raise_sequence_floor(&mut self, region_id: RegionId, entry_id: EntryId) -> Result<()> { // Nothing is assigned an id after stop began. if self.is_stopped() { return Ok(()); @@ -1079,6 +1384,16 @@ impl Actor { if let Some(error) = terminal(&self.terminal_error) { return Err(shared(&error)); } + // In the `enqueued` mode an id the store handed out is never handed + // out again while it runs: a create that fails transiently is + // repeated under its sequence and a permanent failure poisons the + // store. Every later id of the region is greater, so the id needs no + // floor even before it is durable. + if self.ack_mode == AckMode::Enqueued + && entry_id <= self.issued_entry_ids.get(®ion_id).copied().unwrap_or(0) + { + return Ok(()); + } let sequence_floor = sequence_floor(entry_id); let Some(next_object_seq) = self.next_object_seq else { return Ok(()); @@ -1107,24 +1422,35 @@ impl Actor { } } - /// Begins stopping. Nothing is admitted from now on; the open batch and - /// the batches whose create has not started are dropped. Stop is answered by [`finish_stop`](Self::finish_stop) - /// once nothing is in flight. + /// Begins stopping. Nothing is admitted from now on; the `durable` mode + /// drops the open batch and the batches whose create has not started, + /// the `enqueued` mode seals its backlog so it is uploaded. Stop is + /// answered by [`finish_stop`](Self::finish_stop) once nothing is in flight. fn handle_stop(&mut self, response: oneshot::Sender>) { self.stop.push(response); - self.reset_open_batch(); - for pending in self.pending.drain(..) { - let _ = pending - .response - .send(Err(ObjectStoreWalStoppedSnafu.build())); - } - if let Some(index) = self - .sealed - .iter() - .position(|batch| matches!(batch.state, CreateState::Pending)) - { - for batch in self.sealed.drain(index..) { - batch.fail(|| ObjectStoreWalStoppedSnafu.build()); + match self.ack_mode { + AckMode::Durable => { + self.reset_open_batch(); + for pending in self.pending.drain(..) { + let _ = pending + .response + .send(Err(ObjectStoreWalStoppedSnafu.build())); + } + if let Some(index) = self + .sealed + .iter() + .position(|batch| matches!(batch.state, CreateState::Pending)) + { + for batch in self.sealed.drain(index..) { + batch.fail(|| ObjectStoreWalStoppedSnafu.build()); + } + } + } + AckMode::Enqueued => { + if let Some((_, response)) = self.stalled.take() { + let _ = response.send(Err(ObjectStoreWalStoppedSnafu.build())); + } + self.flush_open_batch(); } } } @@ -1135,8 +1461,14 @@ impl Actor { if self.stop.is_empty() || !self.sealed.is_empty() || !self.creates.is_empty() { return false; } + for waiter in self.durable_waiters.drain(..) { + let _ = waiter + .response + .send(Err(ObjectStoreWalStoppedSnafu.build())); + } + let error = self.stop_error.take(); for response in self.stop.drain(..) { - let _ = response.send(Ok(())); + let _ = response.send(error.as_ref().map_or(Ok(()), |error| Err(shared(error)))); } true } @@ -1179,6 +1511,13 @@ fn set_terminal(terminal_error: &TerminalError, error: Error) -> Arc { .clone() } +/// Returns true for a storage error that a later attempt may not meet: one +/// the object store reports as temporary, or as persistent, which is how its +/// retry layer reports a temporary error that outlasted its retries. +fn is_transient(error: &Error) -> bool { + matches!(error, Error::WalObjectStore { error, .. } if !error.is_permanent()) +} + /// Wraps an error that several callers receive. fn shared(error: &Arc) -> Error { ObjectStoreWalSnafu.into_error(error.clone()) @@ -1623,6 +1962,14 @@ mod tests { } } + /// The same batching as `config`, acknowledging appends on admission. + fn enqueued(config: ObjectStoreWalConfig) -> ObjectStoreWalConfig { + ObjectStoreWalConfig { + ack_mode: AckMode::Enqueued, + ..config + } + } + /// Every append reaches the size limit, so it is persisted on its own. fn eager() -> ObjectStoreWalConfig { config(Duration::from_secs(3600), 1) @@ -1787,6 +2134,14 @@ mod tests { prefix: "/absolute".to_string(), ..config(Duration::from_secs(1), 1) }, + ObjectStoreWalConfig { + max_unpersisted_bytes: ReadableSize(0), + ..config(Duration::from_secs(1), 1) + }, + ObjectStoreWalConfig { + max_unpersisted_age: Duration::ZERO, + ..config(Duration::from_secs(1), 1) + }, ] { let error = ObjectStoreLogStore::try_new(memory_store(), &config, 1, 2) .await @@ -2393,6 +2748,7 @@ mod tests { let (append_tx, _) = mpsc::channel(APPEND_BUFFER); let store = ObjectStoreLogStore { prefix: PREFIX.to_string(), + ack_mode: AckMode::Durable, io: Arc::new(ObjectStoreIo::new(memory_store(), PREFIX).unwrap()), catalog: Arc::default(), obsolete_entry_ids: ObsoleteEntryIds::default(), @@ -4268,6 +4624,561 @@ mod tests { ); } + #[tokio::test] + async fn test_store_enqueued_append_returns_before_the_object_exists() { + let store = open(memory_store(), &enqueued(manual())).await; + let region_id = region(1); + + let response = timeout(WAIT, append(&store, region_id, "a1")) + .await + .unwrap() + .unwrap(); + assert_eq!( + HashMap::from([(region_id, id(1, 1))]), + response.last_entry_ids + ); + assert_eq!(vec![0], object_seqs(store.io.as_ref()).await); + assert_eq!(0, latest(&store, region_id)); + assert_eq!(0, store.durable_entry_id(&provider(region_id)).unwrap()); + assert!(read_entries(&store, region_id, 1).await.is_empty()); + let response = append(&store, region_id, "a2").await.unwrap(); + assert_eq!( + HashMap::from([(region_id, id(1, 2))]), + response.last_entry_ids + ); + + store.seal_open_batch().await.unwrap(); + assert_eq!(vec![0, 1], object_seqs(store.io.as_ref()).await); + assert_eq!(id(1, 2), latest(&store, region_id)); + assert_eq!( + id(1, 2), + store.durable_entry_id(&provider(region_id)).unwrap() + ); + assert_eq!( + expected_entries(region_id, &[(id(1, 1), "a1"), (id(1, 2), "a2")]), + read_entries(&store, region_id, 1).await + ); + } + + #[tokio::test] + async fn test_store_enqueued_durable_id_advances_after_create_and_indexing() { + let (store, io, mut parked) = + open_parking_creates(memory_store(), &enqueued(eager())).await; + let region_one = region(1); + let region_two = region(2); + + let response = timeout(WAIT, append(&store, region_one, "a1")) + .await + .unwrap() + .unwrap(); + assert_eq!( + HashMap::from([(region_one, id(1, 1))]), + response.last_entry_ids + ); + let (object_seq, release) = next_create(&mut parked).await; + assert_eq!(1, object_seq); + assert_eq!(0, store.durable_entry_id(&provider(region_one)).unwrap()); + let wait = |entry_id| { + let store = store.clone(); + tokio::spawn(async move { store.wait_durable(&provider(region_one), entry_id).await }) + }; + // An id the store never handed out waits for the ids of the region + // that were handed out below it. + let waits = [wait(id(1, 1)), wait(id(1, 7))]; + // The other region has nothing to wait for. + timeout(WAIT, store.wait_durable(&provider(region_two), id(1, 1))) + .await + .unwrap() + .unwrap(); + for _ in 0..16 { + tokio::task::yield_now().await; + } + assert!(waits.iter().all(|wait| !wait.is_finished())); + + release.send(true).unwrap(); + for wait in waits { + timeout(WAIT, wait).await.unwrap().unwrap().unwrap(); + } + assert_eq!( + id(1, 1), + store.durable_entry_id(&provider(region_one)).unwrap() + ); + assert_eq!(0, store.durable_entry_id(&provider(region_two)).unwrap()); + assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await); + } + + #[tokio::test] + async fn test_store_enqueued_wait_ignores_later_pending_entries_of_the_region() { + let (store, _io, mut parked) = + open_parking_creates(memory_store(), &enqueued(eager())).await; + let region_a = region(1); + let region_b = region(2); + // Object 1 holds region A's id(1, 1) and object 2 region B's id(2, 1). + for (region_id, data) in [(region_a, "a1"), (region_b, "b1")] { + let response = append(&store, region_id, data).await.unwrap(); + let (_, release) = next_create(&mut parked).await; + release.send(true).unwrap(); + let entry_id = response.last_entry_ids[®ion_id]; + timeout(WAIT, store.wait_durable(&provider(region_id), entry_id)) + .await + .unwrap() + .unwrap(); + } + let response = append(&store, region_a, "a2").await.unwrap(); + assert_eq!( + HashMap::from([(region_a, id(3, 1))]), + response.last_entry_ids + ); + let (object_seq, release) = next_create(&mut parked).await; + assert_eq!(3, object_seq); + + // Every entry of region A up to id(2, 1) is durable, so the wait does + // not depend on the parked create of object 3. + timeout(WAIT, store.wait_durable(&provider(region_a), id(2, 1))) + .await + .unwrap() + .unwrap(); + // A wait for id(3, 1) itself is answered only once object 3 is + // indexed; the round trip orders the check after the actor handled it. + let (response_tx, mut response_rx) = oneshot::channel(); + store + .command_tx + .send(Command::WaitDurable { + region_id: region_a, + entry_id: id(3, 1), + response: response_tx, + }) + .await + .unwrap(); + round_trip_actor(&store).await; + assert!(matches!( + response_rx.try_recv(), + Err(oneshot::error::TryRecvError::Empty) + )); + release.send(true).unwrap(); + timeout(WAIT, response_rx).await.unwrap().unwrap().unwrap(); + } + + #[tokio::test] + async fn test_store_enqueued_obsolete_never_passes_the_durable_id() { + let store = open(memory_store(), &enqueued(manual())).await; + let region_id = region(1); + append(&store, region_id, "a1").await.unwrap(); + append(&store, region_id, "a2").await.unwrap(); + + store + .obsolete(&provider(region_id), region_id, id(1, 2)) + .await + .unwrap(); + assert_eq!( + Some(&0), + store.obsolete_entry_ids.lock().unwrap().get(®ion_id) + ); + store.seal_open_batch().await.unwrap(); + assert_eq!( + expected_entries(region_id, &[(id(1, 1), "a1"), (id(1, 2), "a2")]), + read_entries(&store, region_id, 1).await + ); + store + .obsolete(&provider(region_id), region_id, id(1, 2)) + .await + .unwrap(); + assert_eq!( + Some(&id(1, 2)), + store.obsolete_entry_ids.lock().unwrap().get(®ion_id) + ); + assert!(read_entries(&store, region_id, 1).await.is_empty()); + } + + /// Appends one entry in the enqueued mode with `config`, whose create + /// parks, then appends a second one that the backlog threshold must + /// stall. Returns the store, the parked creates, the stalled append and + /// the release of the first create. + async fn stall_second_append( + config: ObjectStoreWalConfig, + ) -> ( + Arc, + Arc, + mpsc::UnboundedReceiver<(u64, oneshot::Sender)>, + tokio::task::JoinHandle>, + oneshot::Sender, + ) { + let (store, io, mut parked) = open_parking_creates(memory_store(), &config).await; + let region_id = region(1); + let response = timeout(WAIT, append(&store, region_id, "a1")) + .await + .unwrap() + .unwrap(); + assert_eq!( + HashMap::from([(region_id, id(1, 1))]), + response.last_entry_ids + ); + if config.max_unpersisted_age < Duration::from_secs(1) { + tokio::time::sleep(config.max_unpersisted_age * 2).await; + } + + let stalled = spawn_append_batch(&store, vec![entry(&store, region_id, "a2")]); + // The stall seals the open batch so that an upload is in flight. + let (object_seq, release) = next_create(&mut parked).await; + assert_eq!(1, object_seq); + for _ in 0..16 { + tokio::task::yield_now().await; + } + assert!(!stalled.is_finished()); + assert!(parked.try_recv().is_err()); + (store, io, parked, stalled, release) + } + + #[tokio::test] + async fn test_store_enqueued_backlog_bytes_stall_admission_until_an_upload_completes() { + let config = ObjectStoreWalConfig { + max_unpersisted_bytes: ReadableSize(1), + ..enqueued(manual()) + }; + let (store, io, mut parked, stalled, release) = stall_second_append(config).await; + let region_id = region(1); + // Two more appends queue behind the stalled one. + let third = spawn_append_batch(&store, vec![entry(&store, region_id, "a3")]); + let fourth = spawn_append_batch(&store, vec![entry(&store, region_id, "a4")]); + for _ in 0..16 { + tokio::task::yield_now().await; + } + assert!(!third.is_finished() && !fourth.is_finished()); + assert!(parked.try_recv().is_err()); + + // The upload releases the second append, whose entry reaches the + // threshold again; the stall seals it so that the next upload can + // release the third, and so on. + release.send(true).unwrap(); + let response = timeout(WAIT, stalled).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_id, id(2, 1))]), + response.last_entry_ids + ); + assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await); + assert_eq!(id(1, 1), latest(&store, region_id)); + for (append, object_seq) in [(third, 2), (fourth, 3)] { + let (parked_seq, release) = next_create(&mut parked).await; + assert_eq!(object_seq, parked_seq); + release.send(true).unwrap(); + let response = timeout(WAIT, append).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_id, id(object_seq + 1, 1))]), + response.last_entry_ids + ); + } + assert_eq!(vec![0, 1, 2, 3], object_seqs(io.as_ref()).await); + assert_eq!( + expected_entries( + region_id, + &[(id(1, 1), "a1"), (id(2, 1), "a2"), (id(3, 1), "a3")] + ), + read_entries(&store, region_id, 1).await + ); + } + + #[tokio::test] + async fn test_store_enqueued_backlog_age_stalls_admission_until_an_upload_completes() { + let config = ObjectStoreWalConfig { + max_unpersisted_age: Duration::from_millis(50), + ..enqueued(manual()) + }; + let (store, io, _parked, stalled, release) = stall_second_append(config).await; + let region_id = region(1); + + release.send(true).unwrap(); + let response = timeout(WAIT, stalled).await.unwrap().unwrap().unwrap(); + assert_eq!( + HashMap::from([(region_id, id(2, 1))]), + response.last_entry_ids + ); + assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await); + assert_eq!( + expected_entries(region_id, &[(id(1, 1), "a1")]), + read_entries(&store, region_id, 1).await + ); + } + + #[tokio::test] + async fn test_store_enqueued_stop_fails_a_stalled_append_and_uploads_the_backlog() { + let config = ObjectStoreWalConfig { + max_unpersisted_bytes: ReadableSize(1), + ..enqueued(manual()) + }; + let (store, io, mut parked, stalled, release) = stall_second_append(config).await; + + let stop = { + let store = store.clone(); + tokio::spawn(async move { store.stop().await }) + }; + assert_stopped(&timeout(WAIT, stalled).await.unwrap().unwrap().unwrap_err()); + release.send(true).unwrap(); + timeout(WAIT, stop).await.unwrap().unwrap().unwrap(); + assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await); + assert!(parked.try_recv().is_err()); + } + + #[tokio::test] + async fn test_store_enqueued_repeats_a_create_that_failed_transiently() { + let (io, _) = RecordingIo::over(memory_store()); + let store = open_over(io.clone(), &enqueued(eager())).await; + let region_id = region(1); + + // The create stores the object but reports a failure; the repeat + // writes the same bytes under the same sequence. + io.fail_after_next_put.store(true, Ordering::Relaxed); + let response = append(&store, region_id, "a1").await.unwrap(); + assert_eq!( + HashMap::from([(region_id, id(1, 1))]), + response.last_entry_ids + ); + timeout(WAIT, store.wait_durable(&provider(region_id), id(1, 1))) + .await + .unwrap() + .unwrap(); + assert_eq!(vec![0, 1], object_seqs(io.as_ref()).await); + assert_eq!( + expected_entries(region_id, &[(id(1, 1), "a1")]), + read_entries(&store, region_id, 1).await + ); + + // An object store with a retry layer reports a temporary error that + // outlasted its retries as persistent; the create is repeated too. + io.fail_next_put_persistently(); + append(&store, region_id, "a2").await.unwrap(); + timeout(WAIT, store.wait_durable(&provider(region_id), id(2, 1))) + .await + .unwrap() + .unwrap(); + assert_eq!(vec![0, 1, 2], object_seqs(io.as_ref()).await); + } + + #[tokio::test] + async fn test_store_enqueued_permanent_failure_poisons() { + // A conflict of any epoch poisons the store without reading the + // existing object, since the acknowledged entries cannot move to + // another sequence. So does a permanent storage error. + for foreign_epoch in [Some(0), Some(2), None] { + let object_store = memory_store(); + let (io, reads) = RecordingIo::over(object_store.clone()); + let store = open_over(io.clone(), &enqueued(eager())).await; + let region_id = region(1); + match foreign_epoch { + Some(epoch) => put_foreign(&object_store, 1, epoch).await, + None => io.fail_next_put_permanently(), + } + let assert_poisoned = |error: Error| { + let poisoned = match foreign_epoch { + Some(_) => matches!(unwrap_shared(&error), Error::WalObjectConflict { .. }), + None => matches!(unwrap_shared(&error), Error::WalObjectStore { .. }), + }; + assert!(poisoned, "unexpected error: {error:?}"); + }; + + // The append was acknowledged; the failure surfaces afterwards. + let second = entry(&store, region_id, "a2"); + append(&store, region_id, "a1").await.unwrap(); + let error = timeout(WAIT, store.wait_durable(&provider(region_id), id(1, 1))) + .await + .unwrap() + .unwrap_err(); + assert_poisoned(error); + assert_poisoned(store.append_batch(vec![second]).await.unwrap_err()); + assert!(store.latest_entry_id(&provider(region_id)).is_err()); + // The acknowledged entry was dropped, which stop reports. + assert_poisoned(store.stop().await.unwrap_err()); + assert!(reads.lock().unwrap().is_empty()); + } + } + + #[tokio::test] + async fn test_store_enqueued_stop_reports_a_backlog_lost_before_the_stop_command() { + let (store, io, mut parked) = + open_parking_creates(memory_store(), &enqueued(manual())).await; + let region_id = region(1); + append(&store, region_id, "a1").await.unwrap(); + let seal = { + let store = store.clone(); + tokio::spawn(async move { store.seal_open_batch().await }) + }; + let (object_seq, release) = next_create(&mut parked).await; + assert_eq!(1, object_seq); + + // Stop began, but the actor has not received the stop command when + // the create fails: the backlog is dropped and the failure is kept + // for the stop that follows. + store.begin_stop(); + release.send(false).unwrap(); + assert_stopped(&timeout(WAIT, seal).await.unwrap().unwrap().unwrap_err()); + // The lost entry cannot be certified as durable before the stop + // command is handled; a durable id still is. + let error = timeout(WAIT, store.wait_durable(&provider(region_id), id(1, 1))) + .await + .unwrap() + .unwrap_err(); + assert!( + matches!(unwrap_shared(&error), Error::WalObjectStore { .. }), + "unexpected error: {error:?}" + ); + timeout(WAIT, store.wait_durable(&provider(region_id), 0)) + .await + .unwrap() + .unwrap(); + let error = store.stop().await.unwrap_err(); + assert!( + matches!(unwrap_shared(&error), Error::WalObjectStore { .. }), + "unexpected error: {error:?}" + ); + assert_eq!(vec![0], object_seqs(io.as_ref()).await); + assert!(parked.try_recv().is_err()); + } + + #[tokio::test] + async fn test_store_enqueued_stop_uploads_the_backlog() { + let store = open(memory_store(), &enqueued(manual())).await; + let region_id = region(1); + append(&store, region_id, "a1").await.unwrap(); + append(&store, region_id, "a2").await.unwrap(); + assert_eq!(vec![0], object_seqs(store.io.as_ref()).await); + + store.stop().await.unwrap(); + assert_eq!(vec![0, 1], object_seqs(store.io.as_ref()).await); + assert_eq!(id(1, 2), latest(&store, region_id)); + assert_eq!( + expected_entries(region_id, &[(id(1, 1), "a1"), (id(1, 2), "a2")]), + read_entries(&store, region_id, 1).await + ); + + // A backlog that cannot be uploaded is reported by stop. A transient + // failure, temporary or persistent, drops it; a permanent one also + // poisons the store. + for failure in ["temporary", "persistent", "permanent"] { + let (io, _) = RecordingIo::over(memory_store()); + let store = open_over(io.clone(), &enqueued(manual())).await; + append(&store, region_id, "a1").await.unwrap(); + match failure { + "temporary" => store.fail_creates(), + "persistent" => io.fail_next_put_persistently(), + _ => io.fail_next_put_permanently(), + } + let error = store.stop().await.unwrap_err(); + assert!( + matches!(unwrap_shared(&error), Error::WalObjectStore { .. }), + "unexpected error: {error:?}" + ); + assert_eq!(vec![0], object_seqs(io.as_ref()).await); + assert_eq!( + failure == "permanent", + store.latest_entry_id(&provider(region_id)).is_err() + ); + store.stop().await.unwrap(); + } + } + + #[tokio::test] + async fn test_store_enqueued_issued_id_needs_no_floor_before_it_is_durable() { + let store = open(memory_store(), &enqueued(manual())).await; + let region_id = region(1); + append(&store, region_id, "a1").await.unwrap(); + + // The open batch holds id(1, 1) under the next sequence. The id was + // handed out, so it is accepted without a floor, with the watermark + // capped to what is durable; an id the store never handed out under + // that sequence is refused. + store + .obsolete(&provider(region_id), region_id, id(1, 1)) + .await + .unwrap(); + assert_eq!( + Some(&0), + store.obsolete_entry_ids.lock().unwrap().get(®ion_id) + ); + let error = store + .obsolete(&provider(region_id), region_id, id(1, 2)) + .await + .unwrap_err(); + assert!( + matches!( + error, + Error::WalObjectSequenceUnsettled { object_seq: 1, .. } + ), + "unexpected error: {error:?}" + ); + + store.seal_open_batch().await.unwrap(); + assert_eq!( + expected_entries(region_id, &[(id(1, 1), "a1")]), + read_entries(&store, region_id, 0).await + ); + store + .obsolete(&provider(region_id), region_id, id(1, 1)) + .await + .unwrap(); + assert!(read_entries(&store, region_id, 0).await.is_empty()); + let response = append(&store, region_id, "a2").await.unwrap(); + assert_eq!( + HashMap::from([(region_id, id(2, 1))]), + response.last_entry_ids + ); + } + + #[tokio::test] + async fn test_store_enqueued_inherited_watermark_raises_the_floor() { + let object_store = memory_store(); + let store = open(object_store.clone(), &enqueued(manual())).await; + let region_id = region(1); + + // The region has nothing durable here, so the recorded watermark + // stays at zero, but its ids must still start above the watermark. + store + .obsolete(&provider(region_id), region_id, id(5, 1)) + .await + .unwrap(); + assert_eq!( + Some(&0), + store.obsolete_entry_ids.lock().unwrap().get(®ion_id) + ); + let response = append(&store, region_id, "r1").await.unwrap(); + assert_eq!( + HashMap::from([(region_id, id(6, 1))]), + response.last_entry_ids + ); + store.seal_open_batch().await.unwrap(); + timeout(WAIT, store.wait_durable(&provider(region_id), id(6, 1))) + .await + .unwrap() + .unwrap(); + assert_eq!(id(6, 1), latest(&store, region_id)); + store.stop().await.unwrap(); + + // Replay from the watermark sees the entry. + let store = open(object_store, &enqueued(manual())).await; + assert_eq!( + expected_entries(region_id, &[(id(6, 1), "r1")]), + read_entries(&store, region_id, id(5, 1) + 1).await + ); + } + + #[tokio::test] + async fn test_store_durability_wait_pending_at_stop_fails_with_stopped() { + let store = open(memory_store(), &manual()).await; + let region_id = region(1); + let append = spawn_append_batch(&store, vec![entry(&store, region_id, "a1")]); + store.wait_for_admitted_appends(1).await.unwrap(); + let wait = { + let store = store.clone(); + tokio::spawn(async move { store.wait_durable(&provider(region_id), id(1, 1)).await }) + }; + for _ in 0..16 { + tokio::task::yield_now().await; + } + assert!(!wait.is_finished()); + + store.stop().await.unwrap(); + assert_stopped(&timeout(WAIT, append).await.unwrap().unwrap().unwrap_err()); + assert_stopped(&timeout(WAIT, wait).await.unwrap().unwrap().unwrap_err()); + } + /// Object access whose next create stores the object but reports a /// transient failure, as a create whose response is lost. struct LostResponseIo { @@ -4446,11 +5357,13 @@ mod tests { /// Object access that records every range read as (sequence, offset, /// length) and, on request, reports the next conditional create as failed - /// after it wrote the object. + /// after it wrote the object, or fails it with a given error before it + /// writes. struct RecordingIo { inner: ObjectStoreIo, reads: RangeReads, fail_after_next_put: AtomicBool, + fail_next_put: Mutex>, } impl RecordingIo { @@ -4460,14 +5373,38 @@ mod tests { inner: ObjectStoreIo::new(object_store, PREFIX).unwrap(), reads: reads.clone(), fail_after_next_put: AtomicBool::new(false), + fail_next_put: Mutex::new(None), }; (Arc::new(io), reads) } + + fn fail_next_put_permanently(&self) { + let error = object_store::Error::new( + object_store::ErrorKind::PermissionDenied, + "injected failure", + ); + *self.fail_next_put.lock().unwrap() = Some(error); + } + + /// Fails the next create as a retry layer reports a temporary error + /// that outlasted its retries. + fn fail_next_put_persistently(&self) { + let error = object_store::Error::new(object_store::ErrorKind::Unexpected, "injected") + .set_temporary() + .set_persistent(); + *self.fail_next_put.lock().unwrap() = Some(error); + } } #[async_trait::async_trait] impl WalObjectIo for RecordingIo { async fn put_if_absent(&self, object_seq: u64, content: Bytes) -> Result { + if let Some(error) = self.fail_next_put.lock().unwrap().take() { + return Err(error).context(WalObjectStoreSnafu { + operation: "write", + path: self.object_path(object_seq), + }); + } let result = self.inner.put_if_absent(object_seq, content).await?; if self.fail_after_next_put.swap(false, Ordering::Relaxed) { return injected_failure("write", self.object_path(object_seq)); diff --git a/src/store-api/src/logstore.rs b/src/store-api/src/logstore.rs index a2c1e3ee803..6084f207641 100644 --- a/src/store-api/src/logstore.rs +++ b/src/store-api/src/logstore.rs @@ -105,6 +105,24 @@ pub trait LogStore: Send + Sync + 'static + std::fmt::Debug { /// Returns the latest entry id in the log store. fn latest_entry_id(&self, provider: &Provider) -> Result; + + /// Waits until every entry of the provider's region with an id at or below + /// `entry_id` has the persistence this log store guarantees for an + /// acknowledged append, so that a caller can persist `entry_id` as a + /// replay watermark. + /// + /// The default returns at once and keeps each log store's own + /// acknowledgement guarantee, whatever it is: Raft Engine with + /// `sync_write = false`, for example, acknowledges an append before its + /// periodic sync. A log store that acknowledges an append before its + /// entries have that persistence overrides this method. + async fn wait_durable( + &self, + _provider: &Provider, + _entry_id: EntryId, + ) -> Result<(), Self::Error> { + Ok(()) + } } /// The response of an `append` operation.