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