diff --git a/src/fetch/tests.rs b/src/fetch/tests.rs index e31f3933..a079cf51 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,17 +9,21 @@ 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}; +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 and released claim, in order. + status_updates: Mutex>, /// Should operations fail? fail: bool, @@ -27,21 +32,28 @@ 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, } } fn error() -> Self { Self { - pending: Mutex::new(None), + pending: Mutex::new(Vec::new()), + status_updates: Mutex::new(Vec::new()), fail: true, } } @@ -89,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> { @@ -124,12 +142,35 @@ 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 release_claims(&self, claims: &[Claim]) -> Result { + if self.fail { + return Err(anyhow!("mock store error")); + } + + self.status_updates.lock().await.extend( + claims + .iter() + .map(|claim| (claim.id.clone(), ActivationStatus::Pending)), + ); + + Ok(claims.len() as u64) } async fn set_status_batch( @@ -267,3 +308,210 @@ 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()); +} + +/// 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() { + 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 + ); +} + +/// 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 2b5f9952..a25e6925 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}; @@ -13,8 +13,10 @@ 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::timed; +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); /// Abstraction for a single fetch thread. pub struct FetchThread { @@ -52,7 +54,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..."); @@ -84,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())); @@ -113,30 +118,52 @@ 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, + }; - Err(QueueError::Timeout) => { + // 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(); + + 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() ); + 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", ); + self.store.undo_claims(&unsent, "fetch.undo_claim").await; + // We cannot recover from a closed channel return false; } @@ -165,22 +192,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); + + let result = loop { + match self.sender.try_send(item) { + // Pushed to channel successfully + Ok(()) => break Ok(()), - match timed!(timeout, "push.queue.wait_duration") { - // The channel was full so the send timed out - Err(_) => Err(QueueError::Timeout), + // The channel may close early if the push pool encounters an error + Err(TrySendError::Disconnected(_)) => break Err(QueueError::Closed), - // The channel may close early if the push pool encounters an error - Ok(Err(_)) => Err(QueueError::Closed), + // The queue is full, so keep the activation here and try again + Err(TrySendError::Full(returned)) => { + if Instant::now() >= deadline { + break Err(QueueError::Timeout); + } - // Pushed to channel successfully - Ok(_) => Ok(()), - } + item = returned; + sleep(QUEUE_RETRY_INTERVAL).await; + } + } + }; + + metrics::histogram!("push.queue.wait_duration").record(start.elapsed()); + + result } } diff --git a/src/push/tests.rs b/src/push/tests.rs index 93e40d48..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,6 +80,10 @@ impl ActivationStore for MockStore { Ok(None) } + async fn release_claims(&self, claims: &[Claim]) -> Result { + Ok(claims.len() as u64) + } + async fn set_status_batch(&self, _ids: &[String], _status: ActivationStatus) -> Result { Ok(0) } diff --git a/src/push/thread.rs b/src/push/thread.rs index e75b8318..fd3b828a 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,9 @@ 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_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 2a0c96eb..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,6 +965,48 @@ impl ActivationStore for PostgresStore { .await } + /// Return claimed activations to pending and clear their claim leases. + /// + /// 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_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 AS t + SET claim_expires_at = null, + status = $1 + 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(&ids) + .bind(&expirations) + .bind(ActivationStatus::Claimed.to_string()) + .execute(&mut *conn) + .await?; + + Ok(released.rows_affected()) + }) + .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..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,6 +836,42 @@ impl ActivationStore for SqliteStore { Ok(Some(row.into())) } + /// Return claimed activations to pending and clear their claim leases. + /// + /// 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_claims(&self, claims: &[Claim]) -> Result { + if claims.is_empty() { + return Ok(0); + } + + 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(")"); + + let released = query_builder.build().execute(&mut *conn).await?; + Ok(released.rows_affected()) + } + #[instrument(skip_all)] async fn set_status_batch( &self, diff --git a/src/store/tests.rs b/src/store/tests.rs index a1abec18..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, @@ -1561,6 +1561,108 @@ 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_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()); + + 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); + 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_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()); + + 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 267dc419..f70bb47c 100644 --- a/src/store/traits.rs +++ b/src/store/traits.rs @@ -4,11 +4,11 @@ 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}; -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,6 +97,54 @@ pub trait ActivationStore: Send + Sync { delay_on_retry: Option, ) -> Result, Error>; + /// Return claimed activations to pending and clear their claim leases. + /// + /// 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 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 delivered. + async fn undo_claims(&self, claims: &[Claim], metric: &'static str) { + if claims.is_empty() { + return; + } + + match self.release_claims(claims).await { + Ok(released) => { + metrics::counter!(metric, "result" => "ok").increment(released); + + // 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(claims.len() as u64); + + error!( + task_ids = ?claims.iter().map(|c| &c.id).collect::>(), + error = ?e, + "Failed to undo claims on activations that were never pushed" + ); + } + } + } + /// Update the status of multiple activations in one batch. async fn set_status_batch( &self, 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)>,