From 4d913326ea6ac95157e823836adf3d8489c32757 Mon Sep 17 00:00:00 2001 From: Enoch Tang Date: Thu, 17 Sep 2026 17:12:49 -0400 Subject: [PATCH 1/4] undo the claim when submitting to the push pool fails --- src/fetch/tests.rs | 110 ++++++++++++++++++++++++++++++++++++++++++-- src/fetch/thread.rs | 8 +++- src/push/thread.rs | 20 +------- src/store/traits.rs | 25 +++++++++- 4 files changed, 140 insertions(+), 23 deletions(-) diff --git a/src/fetch/tests.rs b/src/fetch/tests.rs index e31f3933..3e6b36e5 100644 --- a/src/fetch/tests.rs +++ b/src/fetch/tests.rs @@ -1,4 +1,5 @@ use std::sync::Arc; +use std::time::Instant; use anyhow::{Error, anyhow}; use chrono::{DateTime, Utc}; @@ -8,6 +9,7 @@ use tonic::async_trait; use crate::config::Config; use crate::config::fetch::FetchConfig; +use crate::config::push::{PushConfig, PushQueueConfig}; use crate::store::activation::{Activation, ActivationStatus}; use crate::store::traits::ActivationStore; use crate::store::types::{BucketRange, FailedTasksForwarder, TopicPartition}; @@ -20,6 +22,9 @@ struct MockStore { /// A single (optional) pending activation. pending: Mutex>, + /// Every `set_status` call, in order. + status_updates: Mutex>, + /// Should operations fail? fail: bool, } @@ -28,6 +33,7 @@ impl MockStore { fn empty() -> Self { Self { pending: Mutex::new(None), + status_updates: Mutex::new(Vec::new()), fail: false, } } @@ -35,6 +41,7 @@ impl MockStore { fn one(activation: Activation) -> Self { Self { pending: Mutex::new(Some(activation)), + status_updates: Mutex::new(Vec::new()), fail: false, } } @@ -42,6 +49,7 @@ impl MockStore { fn error() -> Self { Self { pending: Mutex::new(None), + status_updates: Mutex::new(Vec::new()), fail: true, } } @@ -124,12 +132,21 @@ impl ActivationStore for MockStore { async fn set_status( &self, - _id: &str, - _status: ActivationStatus, + id: &str, + status: ActivationStatus, _max_attempts: Option, _delay_on_retry: Option, ) -> Result, Error> { - unimplemented!() + if self.fail { + return Err(anyhow!("mock store error")); + } + + self.status_updates + .lock() + .await + .push((id.to_owned(), status)); + + Ok(None) } async fn set_status_batch( @@ -267,3 +284,90 @@ async fn fetch_pool_no_pending() { handle.abort(); } + +/// Build a single fetch thread covering every bucket. +/// +/// The claim tests drive `fetch_once` directly rather than `FetchPool::start`, +/// because `start` registers a process-wide `elegant_departure` shutdown guard. +/// Any test that lets a fetch loop exit or aborts it trips that guard for the +/// whole test binary, after which later fetch loops return without doing work. +fn fetch_thread( + store: Arc, + sender: Sender<(Activation, Instant)>, + config: Arc, +) -> FetchThread { + FetchThread { + sender, + store, + config, + bucket: bucket_range_for_fetch_thread(0, 1), + } +} + +/// A submit timeout leaves the activation claimed in the store with nothing +/// holding it, so the fetch thread has to release it itself. Otherwise the row +/// waits for `handle_claim_expiration`, whose lease is sized for the worst case +/// push and can strand a single activation for minutes. +#[tokio::test] +async fn fetch_once_undoes_claim_on_submit_timeout() { + let mut activations = make_activations(2); + let filler = activations.remove(0); + let claimed = activations.remove(0); + + let store = Arc::new(MockStore::one(claimed.clone())); + + let config = Arc::new(Config { + push: PushConfig { + queue: PushQueueConfig { + size: 1, + timeout: Duration::from_millis(20), + }, + ..Default::default() + }, + ..Default::default() + }); + + let (sender, receiver) = flume::bounded(config.push.queue.size); + + // Occupy the only queue slot so the submit cannot complete + sender + .send_async((filler, Instant::now())) + .await + .expect("an empty queue should accept one activation"); + + let mut thread = fetch_thread(store.clone(), sender, config); + + // A submit timeout is recoverable, so the loop keeps running + assert!(thread.fetch_once(Duration::from_millis(0)).await); + + assert_eq!( + vec![(claimed.id.clone(), ActivationStatus::Pending)], + *store.status_updates.lock().await + ); + + // Only the filler ever reached the push pool + assert_eq!(1, receiver.len()); +} + +/// Same contract on the path where the push pool has gone away. +#[tokio::test] +async fn fetch_once_undoes_claim_on_closed_queue() { + let claimed = make_activations(1).remove(0); + let store = Arc::new(MockStore::one(claimed.clone())); + + let config = test_config(); + let (sender, receiver) = flume::bounded(config.push.queue.size); + + // Dropping the receiving end makes the submit fail immediately + drop(receiver); + + let mut thread = fetch_thread(store.clone(), sender, config); + + // A closed queue is unrecoverable, so the loop stops + assert!(!thread.fetch_once(Duration::from_millis(0)).await); + + assert_eq!( + vec![(claimed.id.clone(), ActivationStatus::Pending)], + *store.status_updates.lock().await + ); +} diff --git a/src/fetch/thread.rs b/src/fetch/thread.rs index 2b5f9952..5630efe1 100644 --- a/src/fetch/thread.rs +++ b/src/fetch/thread.rs @@ -52,7 +52,7 @@ impl FetchThread { } } - async fn fetch_once(&mut self, fetch_backoff: Duration) -> bool { + pub(super) async fn fetch_once(&mut self, fetch_backoff: Duration) -> bool { let start = Instant::now(); debug!("Fetching next batch of pending activations..."); @@ -125,6 +125,9 @@ impl FetchThread { self.config.push.queue.timeout.as_millis() ); + // Nothing else will push this activation, so release the claim + self.store.undo_claim(&id, "fetch.undo_claim").await; + // Wait for push queue to empty backoff = true; } @@ -137,6 +140,9 @@ impl FetchThread { "Submit to push pool failed due to closed channel", ); + // Nothing else will push this activation, so release the claim + self.store.undo_claim(&id, "fetch.undo_claim").await; + // We cannot recover from a closed channel return false; } diff --git a/src/push/thread.rs b/src/push/thread.rs index e75b8318..052da107 100644 --- a/src/push/thread.rs +++ b/src/push/thread.rs @@ -9,7 +9,7 @@ use flume::Receiver; use tracing::{debug, error}; use crate::push::updater::Updater; -use crate::store::activation::{Activation, ActivationStatus}; +use crate::store::activation::Activation; use crate::store::traits::ActivationStore; use crate::timed; use crate::worker::WorkerMap; @@ -84,23 +84,7 @@ impl PushThread { ); // Revert claimed task back to pending - if let Err(e) = self - .store - .set_status(id, ActivationStatus::Pending, None, None) - .await - { - metrics::counter!("push.undo_claim", "result" => "error").increment(1); - - error!( - task_id = %id, - error = ?e, - "Failed to undo claim on send failure" - ); - - return; - } - - metrics::counter!("push.undo_claim", "result" => "ok").increment(1); + self.store.undo_claim(id, "push.undo_claim").await; return; } diff --git a/src/store/traits.rs b/src/store/traits.rs index 267dc419..5ca8fc83 100644 --- a/src/store/traits.rs +++ b/src/store/traits.rs @@ -4,7 +4,7 @@ use anyhow::{Error, anyhow}; use async_trait::async_trait; use chrono::{DateTime, Utc}; use tokio::join; -use tracing::warn; +use tracing::{error, warn}; use crate::killswitch::KillswitchSelector; use crate::store::activation::{Activation, ActivationStatus}; @@ -97,6 +97,29 @@ pub trait ActivationStore: Send + Sync { delay_on_retry: Option, ) -> Result, Error>; + /// Release a claim on an activation, returning it to pending. + /// + /// Both the fetch and push pools need this when an activation was claimed + /// but could not be handed to a worker. + async fn undo_claim(&self, id: &str, metric: &'static str) { + if let Err(e) = self + .set_status(id, ActivationStatus::Pending, None, None) + .await + { + metrics::counter!(metric, "result" => "error").increment(1); + + error!( + task_id = %id, + error = ?e, + "Failed to undo claim on an activation that was never pushed" + ); + + return; + } + + metrics::counter!(metric, "result" => "ok").increment(1); + } + /// Update the status of multiple activations in one batch. async fn set_status_batch( &self, From 94e931daec0f7adc751172a832ff283e6a0d106a Mon Sep 17 00:00:00 2001 From: Enoch Tang Date: Fri, 18 Sep 2026 11:48:41 -0400 Subject: [PATCH 2/4] implement release_claim to only match claimed tasks --- src/fetch/tests.rs | 13 ++++++++++ src/push/tests.rs | 4 +++ src/store/adapters/postgres.rs | 29 +++++++++++++++++++++ src/store/adapters/sqlite.rs | 25 ++++++++++++++++++ src/store/tests.rs | 42 ++++++++++++++++++++++++++++++ src/store/traits.rs | 47 +++++++++++++++++++++++----------- 6 files changed, 145 insertions(+), 15 deletions(-) diff --git a/src/fetch/tests.rs b/src/fetch/tests.rs index 3e6b36e5..e9bbcbc9 100644 --- a/src/fetch/tests.rs +++ b/src/fetch/tests.rs @@ -149,6 +149,19 @@ impl ActivationStore for MockStore { Ok(None) } + async fn release_claim(&self, id: &str) -> Result { + if self.fail { + return Err(anyhow!("mock store error")); + } + + self.status_updates + .lock() + .await + .push((id.to_owned(), ActivationStatus::Pending)); + + Ok(true) + } + async fn set_status_batch( &self, _ids: &[String], diff --git a/src/push/tests.rs b/src/push/tests.rs index 93e40d48..a81d5b1c 100644 --- a/src/push/tests.rs +++ b/src/push/tests.rs @@ -80,6 +80,10 @@ impl ActivationStore for MockStore { Ok(None) } + async fn release_claim(&self, _id: &str) -> Result { + Ok(true) + } + async fn set_status_batch(&self, _ids: &[String], _status: ActivationStatus) -> Result { Ok(0) } diff --git a/src/store/adapters/postgres.rs b/src/store/adapters/postgres.rs index 2a0c96eb..fdb12b03 100644 --- a/src/store/adapters/postgres.rs +++ b/src/store/adapters/postgres.rs @@ -965,6 +965,35 @@ impl ActivationStore for PostgresStore { .await } + /// Return a claimed activation to pending and clear its claim lease. + /// + /// The `status` guard makes this a no-op once the activation has moved on, + /// so an undo that races a successful push cannot pull the row back out from + /// under a worker that is already running it. + #[instrument(skip_all)] + #[framed] + async fn release_claim(&self, id: &str) -> Result { + retry_query(&self.config.retry, "release_claim", || async { + let mut conn = self.acquire_write_conn_metric("release_claim").await?; + + let released = sqlx::query( + "UPDATE inflight_taskactivations + SET claim_expires_at = null, + status = $1 + WHERE id = $2 + AND status = $3", + ) + .bind(ActivationStatus::Pending.to_string()) + .bind(id) + .bind(ActivationStatus::Claimed.to_string()) + .execute(&mut *conn) + .await?; + + Ok(released.rows_affected() > 0) + }) + .await + } + #[instrument(skip_all)] #[framed] async fn set_status_batch( diff --git a/src/store/adapters/sqlite.rs b/src/store/adapters/sqlite.rs index f14aee1e..d76dfebb 100644 --- a/src/store/adapters/sqlite.rs +++ b/src/store/adapters/sqlite.rs @@ -836,6 +836,31 @@ impl ActivationStore for SqliteStore { Ok(Some(row.into())) } + /// Return a claimed activation to pending and clear its claim lease. + /// + /// The `status` guard makes this a no-op once the activation has moved on, + /// so an undo that races a successful push cannot pull the row back out from + /// under a worker that is already running it. + #[instrument(skip_all)] + async fn release_claim(&self, id: &str) -> Result { + let mut conn = self.acquire_write_conn_metric("release_claim").await?; + + let released = sqlx::query( + "UPDATE inflight_taskactivations + SET claim_expires_at = null, + status = $1 + WHERE id = $2 + AND status = $3", + ) + .bind(ActivationStatus::Pending) + .bind(id) + .bind(ActivationStatus::Claimed) + .execute(&mut *conn) + .await?; + + Ok(released.rows_affected() > 0) + } + #[instrument(skip_all)] async fn set_status_batch( &self, diff --git a/src/store/tests.rs b/src/store/tests.rs index a1abec18..61ec88fc 100644 --- a/src/store/tests.rs +++ b/src/store/tests.rs @@ -1561,6 +1561,48 @@ async fn test_handle_processing_deadline_no_retries_remaining(#[case] adapter: & store.remove_db().await.unwrap(); } +/// The fetch and push pools release their own claims rather than waiting on +/// `handle_claim_expiration`, whose lease is sized for the worst case push. +#[tokio::test] +#[rstest] +#[case::sqlite("sqlite")] +#[case::postgres("postgres")] +async fn test_release_claim_reverts_claimed_to_pending(#[case] adapter: &str) { + let store = create_test_store(adapter).await; + let mut batch = make_activations(1); + batch[0].status = ActivationStatus::Claimed; + batch[0].claim_expires_at = Some(Utc.with_ymd_and_hms(2030, 1, 1, 1, 1, 1).unwrap()); + assert!(store.store(&batch).await.is_ok()); + + assert!(store.release_claim(&batch[0].id).await.unwrap()); + + let task = store.get_by_id(&batch[0].id).await.unwrap().unwrap(); + assert_eq!(task.status, ActivationStatus::Pending); + assert_eq!(task.claim_expires_at, None); + assert_eq!(task.processing_attempts, 0); + store.remove_db().await.unwrap(); +} + +/// Releasing a claim races a successful push. If the push won, the activation is +/// already `Processing` and a worker is running it, so the release has to be a +/// no-op instead of making the row claimable a second time. +#[tokio::test] +#[rstest] +#[case::sqlite("sqlite")] +#[case::postgres("postgres")] +async fn test_release_claim_ignores_activations_that_moved_on(#[case] adapter: &str) { + let store = create_test_store(adapter).await; + let mut batch = make_activations(1); + batch[0].status = ActivationStatus::Processing; + assert!(store.store(&batch).await.is_ok()); + + assert!(!store.release_claim(&batch[0].id).await.unwrap()); + + let task = store.get_by_id(&batch[0].id).await.unwrap().unwrap(); + assert_eq!(task.status, ActivationStatus::Processing); + store.remove_db().await.unwrap(); +} + #[tokio::test] #[rstest] #[case::sqlite("sqlite")] diff --git a/src/store/traits.rs b/src/store/traits.rs index 5ca8fc83..ea4c43d6 100644 --- a/src/store/traits.rs +++ b/src/store/traits.rs @@ -97,27 +97,44 @@ pub trait ActivationStore: Send + Sync { delay_on_retry: Option, ) -> Result, Error>; + /// Return a claimed activation to pending and clear its claim lease. + /// + /// Only rows still in `Claimed` are touched. The bool reports whether a row + /// was updated, so callers can distinguish a released claim from one that + /// already moved on, such as a push thread marking it `Processing`. + async fn release_claim(&self, id: &str) -> Result; + /// Release a claim on an activation, returning it to pending. /// /// Both the fetch and push pools need this when an activation was claimed /// but could not be handed to a worker. async fn undo_claim(&self, id: &str, metric: &'static str) { - if let Err(e) = self - .set_status(id, ActivationStatus::Pending, None, None) - .await - { - metrics::counter!(metric, "result" => "error").increment(1); - - error!( - task_id = %id, - error = ?e, - "Failed to undo claim on an activation that was never pushed" - ); - - return; + match self.release_claim(id).await { + Ok(true) => { + metrics::counter!(metric, "result" => "ok").increment(1); + } + + // Somebody else advanced the activation between the push attempt and + // this update, so it is no longer ours to return to pending. + Ok(false) => { + metrics::counter!(metric, "result" => "not_claimed").increment(1); + + warn!( + task_id = %id, + "Skipped undoing a claim that was already released or advanced" + ); + } + + Err(e) => { + metrics::counter!(metric, "result" => "error").increment(1); + + error!( + task_id = %id, + error = ?e, + "Failed to undo claim on an activation that was never pushed" + ); + } } - - metrics::counter!(metric, "result" => "ok").increment(1); } /// Update the status of multiple activations in one batch. From 012ad053c7b2205c3ab5acab99a519e48b9c4ba5 Mon Sep 17 00:00:00 2001 From: Enoch Tang Date: Fri, 18 Sep 2026 16:01:51 -0400 Subject: [PATCH 3/4] submit to the push pool without parking the activation `send_async` moves the activation onto a flume waiter that a push thread can take at any moment. Dropping that future on timeout neither returns the activation nor reports whether it was delivered, so a submit timeout could not tell a full queue from a late delivery. The fetch thread then undid a claim for an activation a worker had already started, and the row became claimable a second time. `try_send` keeps the activation on this thread and hands it back in `TrySendError::Full`, so every non-Ok exit proves the push pool never saw it. The cost is polling every millisecond while the queue is full instead of parking on the channel. Co-Authored-By: Claude Opus 5 (1M context) --- src/fetch/tests.rs | 49 +++++++++++++++++++++++++++++++++++++++++++++ src/fetch/thread.rs | 42 ++++++++++++++++++++++++++------------ 2 files changed, 78 insertions(+), 13 deletions(-) diff --git a/src/fetch/tests.rs b/src/fetch/tests.rs index e9bbcbc9..252a11b6 100644 --- a/src/fetch/tests.rs +++ b/src/fetch/tests.rs @@ -362,6 +362,55 @@ async fn fetch_once_undoes_claim_on_submit_timeout() { assert_eq!(1, receiver.len()); } +/// A slot opening mid-wait lets the submit through, so no claim is released. +#[tokio::test] +async fn fetch_once_submits_when_queue_drains() { + let mut activations = make_activations(2); + let filler = activations.remove(0); + let claimed = activations.remove(0); + + let store = Arc::new(MockStore::one(claimed.clone())); + + let config = Arc::new(Config { + push: PushConfig { + queue: PushQueueConfig { + size: 1, + timeout: Duration::from_millis(500), + }, + ..Default::default() + }, + ..Default::default() + }); + + let (sender, receiver) = flume::bounded(config.push.queue.size); + + // Occupy the only queue slot so the first submit attempt is refused + sender + .send_async((filler, Instant::now())) + .await + .expect("an empty queue should accept one activation"); + + // Free the slot partway through the submit window + let drain = tokio::spawn({ + let receiver = receiver.clone(); + + async move { + sleep(Duration::from_millis(50)).await; + receiver.recv_async().await.expect("the filler is queued"); + } + }); + + let mut thread = fetch_thread(store.clone(), sender, config); + + assert!(thread.fetch_once(Duration::from_millis(0)).await); + drain.await.expect("drain task should finish"); + + // The activation reached the push pool, so its claim still stands + assert!(store.status_updates.lock().await.is_empty()); + + assert_eq!(claimed.id, receiver.recv_async().await.unwrap().0.id); +} + /// Same contract on the path where the push pool has gone away. #[tokio::test] async fn fetch_once_undoes_claim_on_closed_queue() { diff --git a/src/fetch/thread.rs b/src/fetch/thread.rs index 5630efe1..2acfcf8b 100644 --- a/src/fetch/thread.rs +++ b/src/fetch/thread.rs @@ -5,7 +5,7 @@ use std::time::Instant; use anyhow::Result; use chrono::Utc; use elegant_departure::get_shutdown_guard; -use flume::Sender; +use flume::{Sender, TrySendError}; use tokio::time::{Duration, sleep}; use tracing::{debug, info, warn}; @@ -14,7 +14,9 @@ use crate::push::QueueError; use crate::store::activation::Activation; use crate::store::traits::ActivationStore; use crate::store::types::{BucketRange, TopicPartition}; -use crate::timed; + +/// How long to wait between attempts to submit into a full push queue. +const QUEUE_RETRY_INTERVAL: Duration = Duration::from_millis(1); /// Abstraction for a single fetch thread. pub struct FetchThread { @@ -171,22 +173,36 @@ impl FetchThread { true } + /// Submit one claimed activation to the push pool. An error proves it was never handed off. async fn push_task(&self, activation: Activation, time: Instant) -> Result<(), QueueError> { metrics::gauge!("push.queue.depth").set(self.sender.len() as f64); - let duration = self.config.push.queue.timeout; - let future = self.sender.send_async((activation, time)); - let timeout = tokio::time::timeout(duration, future); + let start = Instant::now(); + let deadline = start + self.config.push.queue.timeout; + let mut item = (activation, time); - match timed!(timeout, "push.queue.wait_duration") { - // The channel was full so the send timed out - Err(_) => Err(QueueError::Timeout), + let result = loop { + match self.sender.try_send(item) { + // Pushed to channel successfully + Ok(()) => break Ok(()), - // The channel may close early if the push pool encounters an error - Ok(Err(_)) => Err(QueueError::Closed), + // The channel may close early if the push pool encounters an error + Err(TrySendError::Disconnected(_)) => break Err(QueueError::Closed), - // Pushed to channel successfully - Ok(_) => Ok(()), - } + // The queue is full, so keep the activation here and try again + Err(TrySendError::Full(returned)) => { + if Instant::now() >= deadline { + break Err(QueueError::Timeout); + } + + item = returned; + sleep(QUEUE_RETRY_INTERVAL).await; + } + } + }; + + metrics::histogram!("push.queue.wait_duration").record(start.elapsed()); + + result } } From e0f91fa1e332173a29923535636a4ec206cde7f2 Mon Sep 17 00:00:00 2001 From: Enoch Tang Date: Thu, 24 Sep 2026 16:23:14 -0400 Subject: [PATCH 4/4] fence claim releases and release the rest of the batch on submit failure A release matched on id and status alone, so one that arrived after its lease expired could undo a later claim on the same row. The row went back to pending while the new claimant was delivering it, and the task could run twice. Releases now also match on the claim_expires_at the activation was fetched with. Every claim writes a fresh expiry, so a stale release no longer matches. On a submit failure the fetch thread released only the current activation. A closed queue returned early and stranded the rest of the batch until lease expiry, and a timeout made each remaining activation wait out its own timeout. The first failure now releases the current activation and the rest of the batch in one query. Co-Authored-By: Claude Opus 5.5 (1M context) --- src/fetch/tests.rs | 130 +++++++++++++++++++++++++++------ src/fetch/thread.rs | 39 +++++++--- src/push/tests.rs | 6 +- src/push/thread.rs | 4 +- src/store/activation.rs | 10 +++ src/store/adapters/postgres.rs | 39 ++++++---- src/store/adapters/sqlite.rs | 51 ++++++++----- src/store/tests.rs | 70 ++++++++++++++++-- src/store/traits.rs | 56 ++++++++------ src/store/types.rs | 13 ++++ 10 files changed, 318 insertions(+), 100 deletions(-) diff --git a/src/fetch/tests.rs b/src/fetch/tests.rs index 252a11b6..a079cf51 100644 --- a/src/fetch/tests.rs +++ b/src/fetch/tests.rs @@ -12,17 +12,17 @@ use crate::config::fetch::FetchConfig; use crate::config::push::{PushConfig, PushQueueConfig}; use crate::store::activation::{Activation, ActivationStatus}; use crate::store::traits::ActivationStore; -use crate::store::types::{BucketRange, FailedTasksForwarder, TopicPartition}; +use crate::store::types::{BucketRange, Claim, FailedTasksForwarder, TopicPartition}; use crate::test_utils::make_activations; use super::*; -/// Store stub that returns one activation once OR is always empty OR always fails. +/// Store stub that returns one batch once OR is always empty OR always fails. struct MockStore { - /// A single (optional) pending activation. - pending: Mutex>, + /// Pending activations, all claimed by the first fetch. + pending: Mutex>, - /// Every `set_status` call, in order. + /// Every `set_status` call and released claim, in order. status_updates: Mutex>, /// Should operations fail? @@ -32,15 +32,19 @@ struct MockStore { impl MockStore { fn empty() -> Self { Self { - pending: Mutex::new(None), + pending: Mutex::new(Vec::new()), status_updates: Mutex::new(Vec::new()), fail: false, } } fn one(activation: Activation) -> Self { + Self::many(vec![activation]) + } + + fn many(activations: Vec) -> Self { Self { - pending: Mutex::new(Some(activation)), + pending: Mutex::new(activations), status_updates: Mutex::new(Vec::new()), fail: false, } @@ -48,7 +52,7 @@ impl MockStore { fn error() -> Self { Self { - pending: Mutex::new(None), + pending: Mutex::new(Vec::new()), status_updates: Mutex::new(Vec::new()), fail: true, } @@ -97,17 +101,23 @@ impl ActivationStore for MockStore { return Err(anyhow!("mock store error")); } - Ok(match self.pending.lock().await.take() { - Some(mut a) => { - a.status = if mark_processing { - ActivationStatus::Processing + let claimed = self + .pending + .lock() + .await + .drain(..) + .map(|mut a| { + if mark_processing { + a.status = ActivationStatus::Processing; } else { - ActivationStatus::Claimed - }; - vec![a] - } - None => vec![], - }) + a.status = ActivationStatus::Claimed; + a.claim_expires_at = Some(Utc::now() + chrono::Duration::minutes(1)); + } + a + }) + .collect(); + + Ok(claimed) } async fn mark_activation_processing(&self, _id: &str) -> Result<(), Error> { @@ -149,17 +159,18 @@ impl ActivationStore for MockStore { Ok(None) } - async fn release_claim(&self, id: &str) -> Result { + async fn release_claims(&self, claims: &[Claim]) -> Result { if self.fail { return Err(anyhow!("mock store error")); } - self.status_updates - .lock() - .await - .push((id.to_owned(), ActivationStatus::Pending)); + self.status_updates.lock().await.extend( + claims + .iter() + .map(|claim| (claim.id.clone(), ActivationStatus::Pending)), + ); - Ok(true) + Ok(claims.len() as u64) } async fn set_status_batch( @@ -433,3 +444,74 @@ async fn fetch_once_undoes_claim_on_closed_queue() { *store.status_updates.lock().await ); } + +/// A timeout means the queue stayed full for the whole window, so the rest of +/// the batch would only wait out the same timeout one by one. Release the whole +/// batch in one go and back off instead. +#[tokio::test] +async fn fetch_once_releases_rest_of_batch_on_submit_timeout() { + let mut activations = make_activations(4); + let filler = activations.remove(0); + let claimed = activations; + + let store = Arc::new(MockStore::many(claimed.clone())); + + let config = Arc::new(Config { + push: PushConfig { + queue: PushQueueConfig { + size: 1, + timeout: Duration::from_millis(20), + }, + ..Default::default() + }, + ..Default::default() + }); + + let (sender, receiver) = flume::bounded(config.push.queue.size); + + // Occupy the only queue slot so the first submit cannot complete + sender + .send_async((filler, Instant::now())) + .await + .expect("an empty queue should accept one activation"); + + let mut thread = fetch_thread(store.clone(), sender, config); + + let start = Instant::now(); + assert!(thread.fetch_once(Duration::from_millis(0)).await); + + // Only the first activation waited out the timeout + assert!(start.elapsed() < Duration::from_millis(60)); + + let expected: Vec<_> = claimed + .iter() + .map(|a| (a.id.clone(), ActivationStatus::Pending)) + .collect(); + assert_eq!(expected, *store.status_updates.lock().await); + + assert_eq!(1, receiver.len()); +} + +/// A closed queue stops the fetch loop, so every activation left in the batch +/// has to be released on the way out rather than stranded until lease expiry. +#[tokio::test] +async fn fetch_once_releases_rest_of_batch_on_closed_queue() { + let claimed = make_activations(3); + let store = Arc::new(MockStore::many(claimed.clone())); + + let config = test_config(); + let (sender, receiver) = flume::bounded(config.push.queue.size); + + // Dropping the receiving end makes the submit fail immediately + drop(receiver); + + let mut thread = fetch_thread(store.clone(), sender, config); + + assert!(!thread.fetch_once(Duration::from_millis(0)).await); + + let expected: Vec<_> = claimed + .iter() + .map(|a| (a.id.clone(), ActivationStatus::Pending)) + .collect(); + assert_eq!(expected, *store.status_updates.lock().await); +} diff --git a/src/fetch/thread.rs b/src/fetch/thread.rs index 2acfcf8b..a25e6925 100644 --- a/src/fetch/thread.rs +++ b/src/fetch/thread.rs @@ -13,7 +13,7 @@ use crate::config::Config; use crate::push::QueueError; use crate::store::activation::Activation; use crate::store::traits::ActivationStore; -use crate::store::types::{BucketRange, TopicPartition}; +use crate::store::types::{BucketRange, Claim, TopicPartition}; /// How long to wait between attempts to submit into a full push queue. const QUEUE_RETRY_INTERVAL: Duration = Duration::from_millis(1); @@ -86,8 +86,11 @@ impl FetchThread { // Use this instant to track claimed → pushed latency let start = Instant::now(); - for activation in activations { + let mut activations = activations.into_iter(); + + while let Some(activation) = activations.next() { let id = activation.id.clone(); + let claim = activation.claim(); if activation.processing_attempts < 1 { let latency = cmp::max(0, activation.received_latency(Utc::now())); @@ -115,35 +118,51 @@ impl FetchThread { ); } - match self.push_task(activation, start).await { - Ok(()) => metrics::counter!("fetch.submit", "result" => "ok").increment(1), + let error = match self.push_task(activation, start).await { + Ok(()) => { + metrics::counter!("fetch.submit", "result" => "ok").increment(1); + continue; + } + + Err(e) => e, + }; + + // Neither this activation nor the rest of the batch reached the + // push pool, and nothing else will push them, so release every + // claim now instead of waiting for the leases to expire + let unsent: Vec = claim + .into_iter() + .chain(activations.by_ref().filter_map(|a| a.claim())) + .collect(); - Err(QueueError::Timeout) => { + match error { + QueueError::Timeout => { metrics::counter!("fetch.submit", "result" => "timeout").increment(1); warn!( task_id = %id, + unsent = unsent.len(), "Submit to push pool timed out after {} milliseconds", self.config.push.queue.timeout.as_millis() ); - // Nothing else will push this activation, so release the claim - self.store.undo_claim(&id, "fetch.undo_claim").await; + self.store.undo_claims(&unsent, "fetch.undo_claim").await; // Wait for push queue to empty backoff = true; + break; } - Err(QueueError::Closed) => { + QueueError::Closed => { metrics::counter!("fetch.submit", "result" => "closed").increment(1); warn!( task_id = %id, + unsent = unsent.len(), "Submit to push pool failed due to closed channel", ); - // Nothing else will push this activation, so release the claim - self.store.undo_claim(&id, "fetch.undo_claim").await; + self.store.undo_claims(&unsent, "fetch.undo_claim").await; // We cannot recover from a closed channel return false; diff --git a/src/push/tests.rs b/src/push/tests.rs index a81d5b1c..7bb3dea6 100644 --- a/src/push/tests.rs +++ b/src/push/tests.rs @@ -12,7 +12,7 @@ use crate::config::push::{PushConfig, PushQueueConfig}; use crate::push::updater::test_eager_updater; use crate::store::activation::{Activation, ActivationStatus}; use crate::store::traits::ActivationStore; -use crate::store::types::{FailedTasksForwarder, TopicPartition}; +use crate::store::types::{Claim, FailedTasksForwarder, TopicPartition}; use crate::test_utils::make_activations; use crate::worker::test_worker_map; @@ -80,8 +80,8 @@ impl ActivationStore for MockStore { Ok(None) } - async fn release_claim(&self, _id: &str) -> Result { - Ok(true) + async fn release_claims(&self, claims: &[Claim]) -> Result { + Ok(claims.len() as u64) } async fn set_status_batch(&self, _ids: &[String], _status: ActivationStatus) -> Result { diff --git a/src/push/thread.rs b/src/push/thread.rs index 052da107..fd3b828a 100644 --- a/src/push/thread.rs +++ b/src/push/thread.rs @@ -84,7 +84,9 @@ impl PushThread { ); // Revert claimed task back to pending - self.store.undo_claim(id, "push.undo_claim").await; + self.store + .undo_claims(activation.claim().as_slice(), "push.undo_claim") + .await; return; } diff --git a/src/store/activation.rs b/src/store/activation.rs index c15bdbef..2b27254c 100644 --- a/src/store/activation.rs +++ b/src/store/activation.rs @@ -6,6 +6,8 @@ use derive_builder::Builder; use sentry_protos::taskbroker::v1::{OnAttemptsExceeded, TaskActivationStatus}; use sqlx::Type; +use crate::store::types::Claim; + /// The members of this enum should be a superset of the members /// of `ActivationStatus` in `sentry_protos`. #[derive(Clone, Copy, Debug, PartialEq, Eq, Type, Hash)] @@ -170,6 +172,14 @@ pub struct Activation { } impl Activation { + /// The claim this activation was fetched under, if it holds one. + pub fn claim(&self) -> Option { + self.claim_expires_at.map(|expires_at| Claim { + id: self.id.clone(), + expires_at, + }) + } + /// The number of milliseconds between an activation's received timestamp and the provided datetime. pub fn received_latency(&self, now: DateTime) -> i64 { now.signed_duration_since(self.received_at) diff --git a/src/store/adapters/postgres.rs b/src/store/adapters/postgres.rs index fdb12b03..6ad81ee8 100644 --- a/src/store/adapters/postgres.rs +++ b/src/store/adapters/postgres.rs @@ -24,7 +24,7 @@ use crate::push::compute_claim_duration_ms; use crate::store::activation::{Activation, ActivationStatus}; use crate::store::retry::retry_query; use crate::store::traits::ActivationStore; -use crate::store::types::{BucketRange, DepthCounts, FailedTasksForwarder, TopicPartition}; +use crate::store::types::{BucketRange, Claim, DepthCounts, FailedTasksForwarder, TopicPartition}; /// Run migrations. pub async fn migrate(config: &StoreConfig) -> Result<()> { @@ -965,31 +965,44 @@ impl ActivationStore for PostgresStore { .await } - /// Return a claimed activation to pending and clear its claim lease. + /// Return claimed activations to pending and clear their claim leases. /// - /// The `status` guard makes this a no-op once the activation has moved on, - /// so an undo that races a successful push cannot pull the row back out from - /// under a worker that is already running it. + /// Each row must still be `Claimed` with the expiry it was claimed under. + /// That makes this a no-op once the activation was pushed, or once its lease + /// expired and a later claim took the row, so a late release can never make + /// a row claimable while another claimant is delivering it. #[instrument(skip_all)] #[framed] - async fn release_claim(&self, id: &str) -> Result { - retry_query(&self.config.retry, "release_claim", || async { - let mut conn = self.acquire_write_conn_metric("release_claim").await?; + async fn release_claims(&self, claims: &[Claim]) -> Result { + if claims.is_empty() { + return Ok(0); + } + + let (ids, expirations): (Vec, Vec>) = claims + .iter() + .map(|claim| (claim.id.clone(), claim.expires_at)) + .unzip(); + + retry_query(&self.config.retry, "release_claims", || async { + let mut conn = self.acquire_write_conn_metric("release_claims").await?; let released = sqlx::query( - "UPDATE inflight_taskactivations + "UPDATE inflight_taskactivations AS t SET claim_expires_at = null, status = $1 - WHERE id = $2 - AND status = $3", + FROM UNNEST($2::text[], $3::timestamptz[]) AS c(id, claim_expires_at) + WHERE t.id = c.id + AND t.status = $4 + AND t.claim_expires_at = c.claim_expires_at", ) .bind(ActivationStatus::Pending.to_string()) - .bind(id) + .bind(&ids) + .bind(&expirations) .bind(ActivationStatus::Claimed.to_string()) .execute(&mut *conn) .await?; - Ok(released.rows_affected() > 0) + Ok(released.rows_affected()) }) .await } diff --git a/src/store/adapters/sqlite.rs b/src/store/adapters/sqlite.rs index d76dfebb..9dd19656 100644 --- a/src/store/adapters/sqlite.rs +++ b/src/store/adapters/sqlite.rs @@ -31,7 +31,7 @@ use crate::killswitch::KillswitchSelector; use crate::push::compute_claim_duration_ms; use crate::store::activation::{Activation, ActivationStatus}; use crate::store::traits::ActivationStore; -use crate::store::types::{BucketRange, FailedTasksForwarder, TopicPartition}; +use crate::store::types::{BucketRange, Claim, FailedTasksForwarder, TopicPartition}; /// Database representation of an [`Activation`], used for both reads and /// writes. @@ -836,29 +836,40 @@ impl ActivationStore for SqliteStore { Ok(Some(row.into())) } - /// Return a claimed activation to pending and clear its claim lease. + /// Return claimed activations to pending and clear their claim leases. /// - /// The `status` guard makes this a no-op once the activation has moved on, - /// so an undo that races a successful push cannot pull the row back out from - /// under a worker that is already running it. + /// Each row must still be `Claimed` with the expiry it was claimed under. + /// That makes this a no-op once the activation was pushed, or once its lease + /// expired and a later claim took the row, so a late release can never make + /// a row claimable while another claimant is delivering it. #[instrument(skip_all)] - async fn release_claim(&self, id: &str) -> Result { - let mut conn = self.acquire_write_conn_metric("release_claim").await?; + async fn release_claims(&self, claims: &[Claim]) -> Result { + if claims.is_empty() { + return Ok(0); + } - let released = sqlx::query( - "UPDATE inflight_taskactivations - SET claim_expires_at = null, - status = $1 - WHERE id = $2 - AND status = $3", - ) - .bind(ActivationStatus::Pending) - .bind(id) - .bind(ActivationStatus::Claimed) - .execute(&mut *conn) - .await?; + let mut conn = self.acquire_write_conn_metric("release_claims").await?; + + let mut query_builder = QueryBuilder::new( + "UPDATE inflight_taskactivations SET claim_expires_at = null, status = ", + ); + query_builder.push_bind(ActivationStatus::Pending); + query_builder.push(" WHERE status = "); + query_builder.push_bind(ActivationStatus::Claimed); + query_builder.push(" AND ("); + + let mut separated = query_builder.separated(" OR "); + for claim in claims { + separated.push("(id = "); + separated.push_bind_unseparated(&claim.id); + separated.push_unseparated(" AND claim_expires_at = "); + separated.push_bind_unseparated(claim.expires_at.timestamp()); + separated.push_unseparated(")"); + } + separated.push_unseparated(")"); - Ok(released.rows_affected() > 0) + let released = query_builder.build().execute(&mut *conn).await?; + Ok(released.rows_affected()) } #[instrument(skip_all)] diff --git a/src/store/tests.rs b/src/store/tests.rs index 61ec88fc..8d2e1598 100644 --- a/src/store/tests.rs +++ b/src/store/tests.rs @@ -16,7 +16,7 @@ use crate::config::{Config, DEFAULT_TOPIC}; use crate::store::activation::{ActivationBuilder, ActivationStatus}; use crate::store::adapters::sqlite::{SqliteStore, create_sqlite_pool}; use crate::store::traits::ActivationStore; -use crate::store::types::TopicPartition; +use crate::store::types::{Claim, TopicPartition}; use crate::test_utils::{ StatusCount, TaskActivationBuilder, assert_counts, create_integration_config, create_test_store, generate_temp_filename, generate_unique_namespace, make_activations, @@ -1567,14 +1567,15 @@ async fn test_handle_processing_deadline_no_retries_remaining(#[case] adapter: & #[rstest] #[case::sqlite("sqlite")] #[case::postgres("postgres")] -async fn test_release_claim_reverts_claimed_to_pending(#[case] adapter: &str) { +async fn test_release_claims_reverts_claimed_to_pending(#[case] adapter: &str) { let store = create_test_store(adapter).await; let mut batch = make_activations(1); batch[0].status = ActivationStatus::Claimed; batch[0].claim_expires_at = Some(Utc.with_ymd_and_hms(2030, 1, 1, 1, 1, 1).unwrap()); assert!(store.store(&batch).await.is_ok()); - assert!(store.release_claim(&batch[0].id).await.unwrap()); + let claim = batch[0].claim().unwrap(); + assert_eq!(1, store.release_claims(&[claim]).await.unwrap()); let task = store.get_by_id(&batch[0].id).await.unwrap().unwrap(); assert_eq!(task.status, ActivationStatus::Pending); @@ -1590,19 +1591,78 @@ async fn test_release_claim_reverts_claimed_to_pending(#[case] adapter: &str) { #[rstest] #[case::sqlite("sqlite")] #[case::postgres("postgres")] -async fn test_release_claim_ignores_activations_that_moved_on(#[case] adapter: &str) { +async fn test_release_claims_ignores_activations_that_moved_on(#[case] adapter: &str) { let store = create_test_store(adapter).await; let mut batch = make_activations(1); batch[0].status = ActivationStatus::Processing; assert!(store.store(&batch).await.is_ok()); - assert!(!store.release_claim(&batch[0].id).await.unwrap()); + let claim = Claim { + id: batch[0].id.clone(), + expires_at: Utc.with_ymd_and_hms(2030, 1, 1, 1, 1, 1).unwrap(), + }; + assert_eq!(0, store.release_claims(&[claim]).await.unwrap()); let task = store.get_by_id(&batch[0].id).await.unwrap().unwrap(); assert_eq!(task.status, ActivationStatus::Processing); store.remove_db().await.unwrap(); } +/// A release can arrive after its lease expired and another fetch thread claimed +/// the row again. That claim carries a new expiry, so the stale release must +/// leave it alone, or the row becomes claimable while it is being delivered. +#[tokio::test] +#[rstest] +#[case::sqlite("sqlite")] +#[case::postgres("postgres")] +async fn test_release_claims_ignores_later_claims(#[case] adapter: &str) { + let store = create_test_store(adapter).await; + let mut batch = make_activations(1); + batch[0].status = ActivationStatus::Claimed; + batch[0].claim_expires_at = Some(Utc.with_ymd_and_hms(2030, 1, 1, 1, 1, 1).unwrap()); + assert!(store.store(&batch).await.is_ok()); + + let stale = Claim { + id: batch[0].id.clone(), + expires_at: Utc.with_ymd_and_hms(2030, 1, 1, 1, 0, 0).unwrap(), + }; + assert_eq!(0, store.release_claims(&[stale]).await.unwrap()); + + let task = store.get_by_id(&batch[0].id).await.unwrap().unwrap(); + assert_eq!(task.status, ActivationStatus::Claimed); + assert_eq!(task.claim_expires_at, batch[0].claim_expires_at); + store.remove_db().await.unwrap(); +} + +/// The claims come straight from the claim query, so the expiry it returns has +/// to match the stored value exactly, at whatever precision the adapter keeps. +#[tokio::test] +#[rstest] +#[case::sqlite("sqlite")] +#[case::postgres("postgres")] +async fn test_release_claims_matches_claims_from_the_claim_query(#[case] adapter: &str) { + let store = create_test_store(adapter).await; + let batch = make_activations(3); + assert!(store.store(&batch).await.is_ok()); + + let claimed = store + .claim_activations_for_push(Some(10), None) + .await + .unwrap(); + assert_eq!(3, claimed.len()); + + let claims: Vec = claimed.iter().filter_map(|a| a.claim()).collect(); + assert_eq!(3, claims.len()); + assert_eq!(3, store.release_claims(&claims).await.unwrap()); + + for activation in &batch { + let task = store.get_by_id(&activation.id).await.unwrap().unwrap(); + assert_eq!(task.status, ActivationStatus::Pending); + assert_eq!(task.claim_expires_at, None); + } + store.remove_db().await.unwrap(); +} + #[tokio::test] #[rstest] #[case::sqlite("sqlite")] diff --git a/src/store/traits.rs b/src/store/traits.rs index ea4c43d6..f70bb47c 100644 --- a/src/store/traits.rs +++ b/src/store/traits.rs @@ -8,7 +8,7 @@ use tracing::{error, warn}; use crate::killswitch::KillswitchSelector; use crate::store::activation::{Activation, ActivationStatus}; -use crate::store::types::{BucketRange, DepthCounts, FailedTasksForwarder, TopicPartition}; +use crate::store::types::{BucketRange, Claim, DepthCounts, FailedTasksForwarder, TopicPartition}; #[async_trait] pub trait ActivationStore: Send + Sync { @@ -97,41 +97,49 @@ pub trait ActivationStore: Send + Sync { delay_on_retry: Option, ) -> Result, Error>; - /// Return a claimed activation to pending and clear its claim lease. + /// Return claimed activations to pending and clear their claim leases. /// - /// Only rows still in `Claimed` are touched. The bool reports whether a row - /// was updated, so callers can distinguish a released claim from one that - /// already moved on, such as a push thread marking it `Processing`. - async fn release_claim(&self, id: &str) -> Result; + /// A row is only released while it still holds the given claim, meaning it + /// is `Claimed` with the same `claim_expires_at`. If the lease already + /// expired and another fetch thread claimed the row again, that claim has a + /// different expiry and is left alone. Returns the number of rows released. + async fn release_claims(&self, claims: &[Claim]) -> Result; - /// Release a claim on an activation, returning it to pending. + /// Release claims on activations that were never handed to a worker. /// /// Both the fetch and push pools need this when an activation was claimed - /// but could not be handed to a worker. - async fn undo_claim(&self, id: &str, metric: &'static str) { - match self.release_claim(id).await { - Ok(true) => { - metrics::counter!(metric, "result" => "ok").increment(1); - } + /// but could not be delivered. + async fn undo_claims(&self, claims: &[Claim], metric: &'static str) { + if claims.is_empty() { + return; + } - // Somebody else advanced the activation between the push attempt and - // this update, so it is no longer ours to return to pending. - Ok(false) => { - metrics::counter!(metric, "result" => "not_claimed").increment(1); + match self.release_claims(claims).await { + Ok(released) => { + metrics::counter!(metric, "result" => "ok").increment(released); - warn!( - task_id = %id, - "Skipped undoing a claim that was already released or advanced" - ); + // The rest were pushed, released, or reclaimed after their lease + // expired, so they are no longer ours to return to pending + let skipped = (claims.len() as u64).saturating_sub(released); + + if skipped > 0 { + metrics::counter!(metric, "result" => "not_claimed").increment(skipped); + + warn!( + task_ids = ?claims.iter().map(|c| &c.id).collect::>(), + skipped, + "Skipped undoing claims that were already released or reclaimed" + ); + } } Err(e) => { - metrics::counter!(metric, "result" => "error").increment(1); + metrics::counter!(metric, "result" => "error").increment(claims.len() as u64); error!( - task_id = %id, + task_ids = ?claims.iter().map(|c| &c.id).collect::>(), error = ?e, - "Failed to undo claim on an activation that was never pushed" + "Failed to undo claims on activations that were never pushed" ); } } diff --git a/src/store/types.rs b/src/store/types.rs index 7a490cdb..7da10a48 100644 --- a/src/store/types.rs +++ b/src/store/types.rs @@ -1,3 +1,5 @@ +use chrono::{DateTime, Utc}; + pub type BucketRange = (i16, i16); /// A Kafka topic paired with one of its partition indices. Partition indices @@ -30,6 +32,17 @@ impl From<&(String, i32)> for TopicPartition { } } +/// One claim on an activation, as handed out by the fetch query. +/// +/// Every claim writes a fresh `claim_expires_at`, so the expiry doubles as a +/// fencing token. A release that matches on it only undoes the claim it was +/// issued, never a later claim on the same row. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct Claim { + pub id: String, + pub expires_at: DateTime, +} + pub struct FailedTasksForwarder { pub to_discard: Vec<(String, Vec)>, pub to_deadletter: Vec<(String, Vec)>,