From 2cc3903cdab35e457a26258b0d1901932dc5c023 Mon Sep 17 00:00:00 2001 From: Maurus Decimus <11444311+mdecimus@users.noreply.github.com> Date: Thu, 11 Jun 2026 15:54:00 +0200 Subject: [PATCH] FoundationDB: Fix read version cache expiration logic --- CHANGELOG.md | 1 + crates/store/src/backend/foundationdb/mod.rs | 85 +++++++-- crates/store/src/backend/foundationdb/read.rs | 173 +++++++++++++----- .../store/src/backend/foundationdb/write.rs | 8 +- tests/src/store/ops.rs | 124 +++++++++++++ 5 files changed, 325 insertions(+), 66 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8f582d61..12e0ae1d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,7 @@ If you are upgrading from v0.16.x, replace the binary (or run `docker pull`). If ## Fixed - JMAP: `*/changes` methods leak ids of non-shared objects (reported by @5ud0er). - Sieve: Do not allow invalid certs in `http_header` function. +- FoundationDB: Fix read version cache expiration logic. ## [0.16.8] - 2026-06-06 diff --git a/crates/store/src/backend/foundationdb/mod.rs b/crates/store/src/backend/foundationdb/mod.rs index 27e11335..09035df1 100644 --- a/crates/store/src/backend/foundationdb/mod.rs +++ b/crates/store/src/backend/foundationdb/mod.rs @@ -5,7 +5,10 @@ */ use foundationdb::{Database, FdbError}; -use std::time::{Duration, Instant}; +use std::{ + sync::atomic::{AtomicBool, AtomicI64, AtomicU64, Ordering}, + time::{Duration, Instant}, +}; pub mod blob; pub mod main; @@ -13,40 +16,96 @@ pub mod read; pub mod write; const MAX_VALUE_SIZE: usize = 100000; -pub const TRANSACTION_EXPIRY: Duration = Duration::from_secs(1); + +const REFRESH_READ_VERSION_AFTER: Duration = Duration::from_secs(1); +const MAX_READ_VERSION_AGE: Duration = Duration::from_secs(4); pub struct FdbStore { db: Database, - version: parking_lot::Mutex, + version: ReadVersion, } pub(crate) struct ReadVersion { - version: i64, - expires: Instant, + base: Instant, + version: AtomicI64, + obtained: AtomicU64, + refreshing: AtomicBool, } impl ReadVersion { - pub fn new(version: i64) -> Self { - Self { - version, - expires: Instant::now() + TRANSACTION_EXPIRY, + fn now(&self) -> u64 { + self.base.elapsed().as_nanos() as u64 + } + + fn current(&self) -> i64 { + self.version.load(Ordering::Acquire) + } + + fn age(&self) -> u64 { + self.now() + .saturating_sub(self.obtained.load(Ordering::Acquire)) + } + + fn store_max(&self, version: i64) { + let mut current = self.version.load(Ordering::Relaxed); + while version > current { + match self.version.compare_exchange_weak( + current, + version, + Ordering::Release, + Ordering::Relaxed, + ) { + Ok(_) => break, + Err(actual) => current = actual, + } } } - pub fn is_expired(&self) -> bool { - self.expires < Instant::now() + fn refreshed(&self, version: i64) { + self.store_max(version); + self.obtained.store(self.now(), Ordering::Release); + } + + fn raise_floor(&self, version: i64) { + self.store_max(version); + } + + fn expire(&self) { + self.obtained.store(0, Ordering::Release); + } + + fn try_begin_refresh(&self) -> Option> { + if self + .refreshing + .compare_exchange(false, true, Ordering::AcqRel, Ordering::Relaxed) + .is_ok() + { + Some(RefreshGuard(&self.refreshing)) + } else { + None + } } } impl Default for ReadVersion { fn default() -> Self { Self { - version: 0, - expires: Instant::now(), + base: Instant::now(), + version: AtomicI64::new(0), + obtained: AtomicU64::new(0), + refreshing: AtomicBool::new(false), } } } +pub(crate) struct RefreshGuard<'a>(&'a AtomicBool); + +impl Drop for RefreshGuard<'_> { + fn drop(&mut self) { + self.0.store(false, Ordering::Release); + } +} + #[inline(always)] fn into_error(error: FdbError) -> trc::Error { trc::StoreEvent::FoundationdbError diff --git a/crates/store/src/backend/foundationdb/read.rs b/crates/store/src/backend/foundationdb/read.rs index 28b18225..8b08d624 100644 --- a/crates/store/src/backend/foundationdb/read.rs +++ b/crates/store/src/backend/foundationdb/read.rs @@ -4,18 +4,21 @@ * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ -use super::{FdbStore, MAX_VALUE_SIZE, ReadVersion, into_error}; +use super::{ + FdbStore, MAX_READ_VERSION_AGE, MAX_VALUE_SIZE, REFRESH_READ_VERSION_AFTER, into_error, +}; use crate::{ Deserialize, IterateParams, Key, ValueKey, WITH_SUBSPACE, backend::deserialize_i64_le, - write::{ValueClass, key::KeySerializer}, + write::{MAX_COMMIT_ATTEMPTS, MAX_COMMIT_TIME, ValueClass, key::KeySerializer}, }; use foundationdb::{ - KeySelector, RangeOption, Transaction, + FdbError, KeySelector, RangeOption, Transaction, future::FdbSlice, options::{self}, }; use futures::TryStreamExt; +use std::time::Instant; #[allow(dead_code)] pub(crate) enum ChunkedValue { @@ -35,26 +38,44 @@ impl FdbStore { U: Deserialize, { let key = key.serialize(WITH_SUBSPACE); - let trx = self.read_trx().await?; + let mut retry_count = 0; + let start = Instant::now(); - match read_chunked_value(&key, &trx, true).await? { - ChunkedValue::Single(bytes) => { - U::deserialize_with_key(key.get(1..).unwrap_or_default(), &bytes).map(Some) + loop { + let trx = self.read_trx().await?; + + match read_chunked_value(&key, &trx, true).await { + Ok(ChunkedValue::Single(bytes)) => { + return U::deserialize_with_key(key.get(1..).unwrap_or_default(), &bytes) + .map(Some); + } + Ok(ChunkedValue::Chunked { bytes, .. }) => { + return U::deserialize_owned_with_key(key.get(1..).unwrap_or_default(), bytes) + .map(Some); + } + Ok(ChunkedValue::None) => return Ok(None), + Err(err) => { + self.on_read_error(trx, err, &mut retry_count, start).await?; + } } - ChunkedValue::Chunked { bytes, .. } => { - U::deserialize_owned_with_key(key.get(1..).unwrap_or_default(), bytes).map(Some) - } - ChunkedValue::None => Ok(None), } } pub(crate) async fn key_exists(&self, key: impl Key) -> trc::Result { let key = key.serialize(WITH_SUBSPACE); - let trx = self.read_trx().await?; + let mut retry_count = 0; + let start = Instant::now(); - match read_chunked_value(&key, &trx, true).await? { - ChunkedValue::Single(_) | ChunkedValue::Chunked { .. } => Ok(true), - ChunkedValue::None => Ok(false), + loop { + let trx = self.read_trx().await?; + + match read_chunked_value(&key, &trx, true).await { + Ok(ChunkedValue::Single(_) | ChunkedValue::Chunked { .. }) => return Ok(true), + Ok(ChunkedValue::None) => return Ok(false), + Err(err) => { + self.on_read_error(trx, err, &mut retry_count, start).await?; + } + } } } @@ -65,6 +86,8 @@ impl FdbStore { ) -> trc::Result<()> { let begin = params.begin.serialize(WITH_SUBSPACE); let end = params.end.serialize(WITH_SUBSPACE); + let mut retry_count = 0; + let start = Instant::now(); if !params.first { let mut last_key = vec![]; @@ -144,11 +167,25 @@ impl FdbStore { break 'outer; } Err(e) => { + drop(values); if e.code() == 1007 && !last_key_.is_empty() { // Transaction is too old to perform reads or be committed - drop(values); last_key = last_key_; continue 'outer; + } else if e.is_retryable() + && retry_count < MAX_COMMIT_ATTEMPTS + && start.elapsed() < MAX_COMMIT_TIME + { + // Transient error such as a cached read version ahead of lagging + // storage servers (code 1009); resume from the last key read, + // refresh the read version and back off before retrying. + if !last_key_.is_empty() { + last_key = last_key_; + } + self.version.expire(); + trx.on_error(e).await.map_err(into_error)?; + retry_count += 1; + continue 'outer; } else { return Err(into_error(e)); } @@ -157,20 +194,30 @@ impl FdbStore { } } } else { - let trx = self.read_trx().await?; - let mut values = trx.get_ranges_keyvalues( - RangeOption { - begin: KeySelector::first_greater_or_equal(&begin), - end: KeySelector::first_greater_than(&end), - mode: options::StreamingMode::Small, - reverse: !params.ascending, - ..Default::default() - }, - true, - ); + loop { + let trx = self.read_trx().await?; + let mut values = trx.get_ranges_keyvalues( + RangeOption { + begin: KeySelector::first_greater_or_equal(&begin), + end: KeySelector::first_greater_than(&end), + mode: options::StreamingMode::Small, + reverse: !params.ascending, + ..Default::default() + }, + true, + ); - if let Some(value) = values.try_next().await.map_err(into_error)? { - cb(value.key().get(1..).unwrap_or_default(), value.value())?; + match values.try_next().await { + Ok(Some(value)) => { + cb(value.key().get(1..).unwrap_or_default(), value.value())?; + break; + } + Ok(None) => break, + Err(e) => { + drop(values); + self.on_read_error(trx, e, &mut retry_count, start).await?; + } + } } } @@ -182,31 +229,61 @@ impl FdbStore { key: impl Into> + Sync + Send, ) -> trc::Result { let key = key.into().serialize(WITH_SUBSPACE); - if let Some(bytes) = self - .read_trx() - .await? - .get(&key, true) - .await - .map_err(into_error)? + let mut retry_count = 0; + let start = Instant::now(); + + loop { + let trx = self.read_trx().await?; + match trx.get(&key, true).await { + Ok(Some(bytes)) => return deserialize_i64_le(&key, &bytes), + Ok(None) => return Ok(0), + Err(e) => { + self.on_read_error(trx, e, &mut retry_count, start).await?; + } + } + } + } + + async fn on_read_error( + &self, + trx: Transaction, + err: FdbError, + retry_count: &mut u32, + start: Instant, + ) -> trc::Result<()> { + if err.is_retryable() + && *retry_count < MAX_COMMIT_ATTEMPTS + && start.elapsed() < MAX_COMMIT_TIME { - deserialize_i64_le(&key, &bytes) + // The cached read version may be ahead of lagging storage servers under heavy write + // load (code 1009); expire it so the retry obtains a fresh read version, then let + // FoundationDB back off before retrying. + self.version.expire(); + trx.on_error(err).await.map_err(into_error)?; + *retry_count += 1; + Ok(()) } else { - Ok(0) + Err(into_error(err)) } } pub(crate) async fn read_trx(&self) -> trc::Result { - let (is_expired, mut read_version) = { - let version = self.version.lock(); - (version.is_expired(), version.version) - }; let trx = self.db.create_trx().map_err(into_error)?; + let version = self.version.current(); + let age = self.version.age(); - if is_expired { - read_version = trx.get_read_version().await.map_err(into_error)?; - *self.version.lock() = ReadVersion::new(read_version); + if version != 0 && age < MAX_READ_VERSION_AGE.as_nanos() as u64 { + if age >= REFRESH_READ_VERSION_AFTER.as_nanos() as u64 + && let Some(_guard) = self.version.try_begin_refresh() + { + let read_version = trx.get_read_version().await.map_err(into_error)?; + self.version.refreshed(read_version); + } else { + trx.set_read_version(version); + } } else { - trx.set_read_version(read_version); + let read_version = trx.get_read_version().await.map_err(into_error)?; + self.version.refreshed(read_version); } Ok(trx) @@ -217,8 +294,8 @@ pub(crate) async fn read_chunked_value( key: &[u8], trx: &Transaction, snapshot: bool, -) -> trc::Result { - if let Some(bytes) = trx.get(key, snapshot).await.map_err(into_error)? { +) -> Result { + if let Some(bytes) = trx.get(key, snapshot).await? { if bytes.len() < MAX_VALUE_SIZE { Ok(ChunkedValue::Single(bytes)) } else { @@ -229,7 +306,7 @@ pub(crate) async fn read_chunked_value( .write(0u8) .finalize(); - while let Some(bytes) = trx.get(&key, snapshot).await.map_err(into_error)? { + while let Some(bytes) = trx.get(&key, snapshot).await? { value.extend_from_slice(&bytes); *key.last_mut().unwrap() += 1; } diff --git a/crates/store/src/backend/foundationdb/write.rs b/crates/store/src/backend/foundationdb/write.rs index 23a7e689..85516c4f 100644 --- a/crates/store/src/backend/foundationdb/write.rs +++ b/crates/store/src/backend/foundationdb/write.rs @@ -5,7 +5,7 @@ */ use super::{ - FdbStore, MAX_VALUE_SIZE, ReadVersion, into_error, + FdbStore, MAX_VALUE_SIZE, into_error, read::{ChunkedValue, read_chunked_value}, }; use crate::{ @@ -102,6 +102,7 @@ impl FdbStore { let (merge_result, is_chunked) = match read_chunked_value(&key, &trx, false) .await + .map_err(into_error) .caused_by(trc::location!())? { ChunkedValue::Single(slice) => ( @@ -270,10 +271,7 @@ impl FdbStore { match trx.commit().await { Ok(result) => { let commit_version = result.committed_version().map_err(into_error)?; - let mut version = self.version.lock(); - if commit_version > version.version { - *version = ReadVersion::new(commit_version); - } + self.version.raise_floor(commit_version); Ok(true) } Err(err) => { diff --git a/tests/src/store/ops.rs b/tests/src/store/ops.rs index 6e94d8fa..5f6d79ad 100644 --- a/tests/src/store/ops.rs +++ b/tests/src/store/ops.rs @@ -124,6 +124,130 @@ pub async fn test(test: &TestServer) { .await .unwrap(); + // Read-your-writes through the cached read version: overwrite a key in a tight loop + println!("Running FoundationDB read-your-writes test..."); + for n in 0u64..200 { + db.write( + BatchBuilder::new() + .with_account_id(0) + .with_collection(Collection::Email) + .with_document(0) + .set( + ValueClass::Registry(RegistryClass::Item { + object_id: 100, + item_id: 0, + }), + n.to_be_bytes().to_vec(), + ) + .build_all(), + ) + .await + .unwrap(); + + let got = db + .get_value::(ValueKey { + account_id: 0, + collection: 0, + document_id: 0, + class: ValueClass::Registry(RegistryClass::Item { + object_id: 100, + item_id: 0, + }), + }) + .await + .unwrap() + .unwrap(); + assert_eq!(got, n, "stale read: wrote {n} but read back {got}"); + } + db.write( + BatchBuilder::new() + .with_account_id(0) + .with_collection(Collection::Email) + .with_document(0) + .clear(ValueClass::Registry(RegistryClass::Item { + object_id: 100, + item_id: 0, + })) + .build_all(), + ) + .await + .unwrap(); + + // Read-version cache monotonicity under concurrency: while a writer increments a counter + println!("Running FoundationDB read-version monotonicity test..."); + let n_increments = 500u64; + + let writer = { + let db = db.clone(); + tokio::spawn(async move { + for _ in 0..n_increments { + db.write( + BatchBuilder::new() + .with_account_id(0) + .with_collection(Collection::Email) + .with_document(5000) + .add_and_get(ValueClass::Quota, 1) + .build_all(), + ) + .await + .unwrap(); + } + }) + }; + + let mut readers = Vec::new(); + for _ in 0..16 { + let db = db.clone(); + readers.push(tokio::spawn(async move { + let deadline = std::time::Instant::now() + std::time::Duration::from_millis(1500); + let mut last = 0i64; + while std::time::Instant::now() < deadline { + let current = db + .get_counter(ValueKey { + account_id: 0, + collection: 0, + document_id: 5000, + class: ValueClass::Quota, + }) + .await + .unwrap(); + assert!( + current >= last, + "read version regressed: counter went from {last} to {current}" + ); + last = current; + } + })); + } + + writer.await.unwrap(); + for reader in readers { + reader.await.unwrap(); + } + + assert_eq!( + db.get_counter(ValueKey { + account_id: 0, + collection: 0, + document_id: 5000, + class: ValueClass::Quota, + }) + .await + .unwrap(), + n_increments as i64, + "counter did not reach the expected total" + ); + db.write( + BatchBuilder::new() + .with_account_id(0) + .with_collection(Collection::Email) + .with_document(5000) + .clear(ValueClass::Quota) + .build_all(), + ) + .await + .unwrap(); + if std::env::var("SLOW_FDB_TRX").is_ok() { println!("Running FoundationDB slow transaction tests..."); // Create 900000 keys