diff --git a/quickwit/Cargo.lock b/quickwit/Cargo.lock index 1d4e449dd3b..33575ed483b 100644 --- a/quickwit/Cargo.lock +++ b/quickwit/Cargo.lock @@ -8521,6 +8521,7 @@ dependencies = [ "tower 0.5.3", "tracing", "tracing-opentelemetry", + "tracing-subscriber", ] [[package]] diff --git a/quickwit/quickwit-aws/src/retry.rs b/quickwit/quickwit-aws/src/retry.rs index 6885d7896be..d5670ac450b 100644 --- a/quickwit/quickwit-aws/src/retry.rs +++ b/quickwit/quickwit-aws/src/retry.rs @@ -62,13 +62,18 @@ where E: AwsRetryable } } -pub async fn aws_retry(retry_params: &RetryParams, f: impl Fn() -> Fut) -> Result +pub async fn aws_retry( + request_name: &'static str, + retry_params: &RetryParams, + f: impl Fn() -> Fut, +) -> Result where Fut: Future>, E: AwsRetryable + Debug + 'static, { retry_with_mockable_sleep( retry_params, + request_name, || f().map_err(AwsRetryableWrapper), TokioSleep, ) diff --git a/quickwit/quickwit-common/Cargo.toml b/quickwit/quickwit-common/Cargo.toml index 2060cd81ac5..16ded88a036 100644 --- a/quickwit/quickwit-common/Cargo.toml +++ b/quickwit/quickwit-common/Cargo.toml @@ -73,3 +73,4 @@ serde_json = { workspace = true } serial_test = { workspace = true } tempfile = { workspace = true } tokio = { workspace = true, features = ["test-util"] } +tracing-subscriber = { workspace = true } diff --git a/quickwit/quickwit-common/src/retry.rs b/quickwit/quickwit-common/src/retry.rs index 5129e56912a..e40b7e127a5 100644 --- a/quickwit/quickwit-common/src/retry.rs +++ b/quickwit/quickwit-common/src/retry.rs @@ -142,6 +142,7 @@ impl MockableSleep for TokioSleep { pub async fn retry_with_mockable_sleep( retry_params: &RetryParams, + request_name: &'static str, f: impl Fn() -> Fut, mockable_sleep: impl MockableSleep, ) -> Result @@ -167,8 +168,10 @@ where if num_attempts >= retry_params.max_attempts { warn!( + request=%request_name, num_attempts=%num_attempts, - "request failed" + error=?error, + "request failed after exhausting retries" ); return Err(error); } @@ -177,6 +180,7 @@ where None => retry_params.compute_delay(num_attempts), }; debug!( + request=%request_name, num_attempts=%num_attempts, delay_ms=%delay.as_millis(), error=?error, @@ -186,17 +190,22 @@ where } } -pub async fn retry(retry_params: &RetryParams, f: impl Fn() -> Fut) -> Result +pub async fn retry( + request_name: &'static str, + retry_params: &RetryParams, + f: impl Fn() -> Fut, +) -> Result where Fut: Future>, E: Retryable + Debug + 'static, { - retry_with_mockable_sleep(retry_params, f, TokioSleep).await + retry_with_mockable_sleep(retry_params, request_name, f, TokioSleep).await } #[cfg(test)] mod tests { - use std::sync::RwLock; + use std::io::Write; + use std::sync::{Arc, Mutex, RwLock}; use std::time::Duration; use futures::future::ready; @@ -220,6 +229,19 @@ mod tests { struct NoopSleep; + struct TestWriter(Arc>>); + + impl Write for TestWriter { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0.lock().unwrap().extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + #[async_trait::async_trait] impl MockableSleep for NoopSleep { async fn sleep(&self, _duration: Duration) { @@ -236,6 +258,7 @@ mod tests { max_delay: Duration::from_millis(2), max_attempts: 30, }, + "test.retry", || ready(values_it.write().unwrap().next().unwrap()), noop_mock, ) @@ -284,6 +307,33 @@ mod tests { assert_eq!(simulate_retries(retry_sequence).await, Ok(())); } + #[tokio::test(flavor = "current_thread")] + async fn test_retry_failure_log_includes_request_name_and_error() { + let output = Arc::new(Mutex::new(Vec::new())); + let writer_output = output.clone(); + let subscriber = tracing_subscriber::fmt() + .with_ansi(false) + .without_time() + .with_writer(move || TestWriter(writer_output.clone())) + .finish(); + let _subscriber_guard = tracing::subscriber::set_default(subscriber); + + let result: Result<(), Retry> = retry_with_mockable_sleep( + &RetryParams::no_retries(), + "s3.get_object", + || ready(Err(Retry::Transient(42))), + NoopSleep, + ) + .await; + + assert_eq!(result, Err(Retry::Transient(42))); + let logs = String::from_utf8(output.lock().unwrap().clone()).unwrap(); + assert!(logs.contains("request failed after exhausting retries")); + assert!(logs.contains("request=s3.get_object")); + assert!(logs.contains("error=Transient(42)")); + assert!(logs.contains("num_attempts=1")); + } + fn test_retry_delay_does_not_overflow_aux(retry_params: RetryParams) { for i in 1..100 { let delay = retry_params.compute_delay(i); diff --git a/quickwit/quickwit-indexing/src/source/kinesis/api.rs b/quickwit/quickwit-indexing/src/source/kinesis/api.rs index 54e303c7c9a..6ac6017e30b 100644 --- a/quickwit/quickwit-indexing/src/source/kinesis/api.rs +++ b/quickwit/quickwit-indexing/src/source/kinesis/api.rs @@ -27,7 +27,7 @@ pub(crate) async fn get_records( ) -> anyhow::Result { // TODO: Return an error other than `anyhow::Error` so that expired shard iterators can be // handled properly. - let response = aws_retry(retry_params, || async { + let response = aws_retry("kinesis.get_records", retry_params, || async { kinesis_client .get_records() .shard_iterator(shard_iterator.clone()) @@ -59,7 +59,7 @@ pub(crate) async fn get_shard_iterator( ShardIteratorType::TrimHorizon }; - let response = aws_retry(retry_params, || async { + let response = aws_retry("kinesis.get_shard_iterator", retry_params, || async { kinesis_client .get_shard_iterator() .stream_name(stream_name) @@ -93,7 +93,7 @@ pub(crate) async fn list_shards( None }; let limit_per_request = limit_per_request.map(|limit| limit as i32); - let response = aws_retry(retry_params, || async { + let response = aws_retry("kinesis.list_shards", retry_params, || async { kinesis_client .list_shards() .set_stream_name(stream_name.clone()) @@ -132,7 +132,7 @@ pub(crate) mod tests { stream_name: &str, num_shards: usize, ) -> anyhow::Result<()> { - aws_retry(&DEFAULT_RETRY_PARAMS, || async { + aws_retry("kinesis.create_stream", &DEFAULT_RETRY_PARAMS, || async { kinesis_client .create_stream() .stream_name(stream_name) @@ -151,7 +151,7 @@ pub(crate) mod tests { kinesis_client: &KinesisClient, stream_name: &str, ) -> anyhow::Result<()> { - aws_retry(&DEFAULT_RETRY_PARAMS, || async { + aws_retry("kinesis.delete_stream", &DEFAULT_RETRY_PARAMS, || async { kinesis_client .delete_stream() .stream_name(stream_name.to_string()) @@ -169,7 +169,7 @@ pub(crate) mod tests { kinesis_client: &KinesisClient, stream_name: &str, ) -> anyhow::Result { - let response = aws_retry(&DEFAULT_RETRY_PARAMS, || async { + let response = aws_retry("kinesis.describe_stream", &DEFAULT_RETRY_PARAMS, || async { kinesis_client .describe_stream() .stream_name(stream_name.to_string()) @@ -193,7 +193,7 @@ pub(crate) mod tests { let mut has_more_streams = true; let limit_per_request = limit_per_request.map(|limit| limit as i32); while has_more_streams { - let response = aws_retry(&DEFAULT_RETRY_PARAMS, || async { + let response = aws_retry("kinesis.list_streams", &DEFAULT_RETRY_PARAMS, || async { kinesis_client .list_streams() .set_exclusive_start_stream_name(exclusive_start_stream_name.clone()) @@ -218,7 +218,7 @@ pub(crate) mod tests { shard_id: &str, adjacent_shard_id: &str, ) -> anyhow::Result<()> { - aws_retry(&DEFAULT_RETRY_PARAMS, || async { + aws_retry("kinesis.merge_shards", &DEFAULT_RETRY_PARAMS, || async { kinesis_client .merge_shards() .stream_name(stream_name) @@ -240,7 +240,7 @@ pub(crate) mod tests { shard_id: &str, starting_hash_key: &str, ) -> anyhow::Result<()> { - aws_retry(&DEFAULT_RETRY_PARAMS, || async { + aws_retry("kinesis.split_shard", &DEFAULT_RETRY_PARAMS, || async { kinesis_client .split_shard() .stream_name(stream_name) diff --git a/quickwit/quickwit-indexing/src/source/queue_sources/sqs_queue.rs b/quickwit/quickwit-indexing/src/source/queue_sources/sqs_queue.rs index 7b5dbb5643f..5ea3db76bf2 100644 --- a/quickwit/quickwit-indexing/src/source/queue_sources/sqs_queue.rs +++ b/quickwit/quickwit-indexing/src/source/queue_sources/sqs_queue.rs @@ -67,7 +67,7 @@ impl Queue for SqsQueue { // state that it starts when the message is returned. let initial_deadline = Instant::now() + suggested_deadline; let clamped_max_messages = std::cmp::min(max_messages, 10) as i32; - let receive_output = aws_retry(&self.receive_retries, || async { + let receive_output = aws_retry("sqs.receive_message", &self.receive_retries, || async { self.sqs_client .receive_message() .queue_url(&self.queue_url) @@ -133,13 +133,17 @@ impl Queue for SqsQueue { let mut batch_errors = Vec::new(); let mut message_errors = Vec::new(); for batch in entry_batches { - let res = aws_retry(&self.acknowledge_retries, || { - self.sqs_client - .delete_message_batch() - .queue_url(&self.queue_url) - .set_entries(Some(batch.clone())) - .send() - }) + let res = aws_retry( + "sqs.delete_message_batch", + &self.acknowledge_retries, + || { + self.sqs_client + .delete_message_batch() + .queue_url(&self.queue_url) + .set_entries(Some(batch.clone())) + .send() + }, + ) .await; match res { Ok(res) => { @@ -187,14 +191,18 @@ impl Queue for SqsQueue { ) -> anyhow::Result { let visibility_timeout = std::cmp::min(suggested_deadline.as_secs() as i32, 43200); let new_deadline = Instant::now() + suggested_deadline; - aws_retry(&self.modify_deadline_retries, || { - self.sqs_client - .change_message_visibility() - .queue_url(&self.queue_url) - .visibility_timeout(visibility_timeout) - .receipt_handle(ack_id) - .send() - }) + aws_retry( + "sqs.change_message_visibility", + &self.modify_deadline_retries, + || { + self.sqs_client + .change_message_visibility() + .queue_url(&self.queue_url) + .visibility_timeout(visibility_timeout) + .receipt_handle(ack_id) + .send() + }, + ) .await?; Ok(new_deadline) } diff --git a/quickwit/quickwit-storage/src/object_storage/azure_blob_storage.rs b/quickwit/quickwit-storage/src/object_storage/azure_blob_storage.rs index bc195da47f6..ef2944603d0 100644 --- a/quickwit/quickwit-storage/src/object_storage/azure_blob_storage.rs +++ b/quickwit/quickwit-storage/src/object_storage/azure_blob_storage.rs @@ -215,7 +215,7 @@ impl AzureBlobStorage { ) -> StorageResult { let name = self.blob_name(path); let capacity = range_opt.as_ref().map(Range::len).unwrap_or(0); - retry(&self.retry_params, || async { + retry("azure.get", &self.retry_params, || async { let _timer = HistogramTimer::new(&crate::metrics::OBJECT_STORAGE_GET_OBJECT_DURATION); let (mut response_stream, _in_flight_guards) = if let Some(range) = range_opt.as_ref() { let stream = self @@ -247,7 +247,7 @@ impl AzureBlobStorage { crate::metrics::OBJECT_STORAGE_PUT_PARTS.inc(); crate::metrics::OBJECT_STORAGE_UPLOAD_NUM_BYTES.inc_by(payload.len()); let _timer = HistogramTimer::new(&crate::metrics::OBJECT_STORAGE_PUT_OBJECT_DURATION); - retry(&self.retry_params, || async { + retry("azure.put", &self.retry_params, || async { let data = Bytes::from(payload.read_all().await?.to_vec()); let hash = azure_storage_blobs::prelude::Hash::from(md5::compute(&data[..]).0); self.container_client @@ -284,7 +284,7 @@ impl AzureBlobStorage { async move { let _timer = HistogramTimer::new(&crate::metrics::OBJECT_STORAGE_UPLOAD_PART_DURATION); - retry(&self.retry_params, || async { + retry("azure.put", &self.retry_params, || async { // zero pad block ids to make them sortable as strings let block_id = format!("block:{:05}", num); let (data, hash_digest) = @@ -469,7 +469,7 @@ impl Storage for AzureBlobStorage { path: &Path, range: Range, ) -> StorageResult> { - retry(&self.retry_params, || async { + retry("azure.get", &self.retry_params, || async { let range = range.clone(); let name = self.blob_name(path); let page_stream = self diff --git a/quickwit/quickwit-storage/src/object_storage/s3_compatible_storage.rs b/quickwit/quickwit-storage/src/object_storage/s3_compatible_storage.rs index d015bf799d6..eed4bfc6dab 100644 --- a/quickwit/quickwit-storage/src/object_storage/s3_compatible_storage.rs +++ b/quickwit/quickwit-storage/src/object_storage/s3_compatible_storage.rs @@ -417,7 +417,7 @@ impl S3CompatibleObjectStorage { len: u64, ) -> StorageResult<()> { let bucket = &self.bucket; - aws_retry(&self.retry_params, || async { + aws_retry("s3.put_object", &self.retry_params, || async { self.put_single_part_single_try(bucket, key, payload.clone(), len) .await }) @@ -427,7 +427,7 @@ impl S3CompatibleObjectStorage { } async fn create_multipart_upload(&self, key: &str) -> StorageResult { - let upload_id = aws_retry(&self.retry_params, || async { + let upload_id = aws_retry("s3.create_multipart_upload", &self.retry_params, || async { self.s3_client .create_multipart_upload() .bucket(self.bucket.clone()) @@ -584,7 +584,7 @@ impl S3CompatibleObjectStorage { stream::iter(parts.into_iter().map(|part| { let payload = payload.clone(); let upload_id = upload_id.clone(); - aws_retry(&self.retry_params, move || { + aws_retry("s3.upload_part", &self.retry_params, move || { self.upload_part(upload_id.clone(), key, part.clone(), payload.clone()) }) })) @@ -623,22 +623,26 @@ impl S3CompatibleObjectStorage { let completed_upload = CompletedMultipartUpload::builder() .set_parts(Some(completed_parts)) .build(); - aws_retry(&self.retry_params, || async { - self.s3_client - .complete_multipart_upload() - .bucket(self.bucket.clone()) - .key(key) - .multipart_upload(completed_upload.clone()) - .upload_id(upload_id) - .send() - .await - }) + aws_retry( + "s3.complete_multipart_upload", + &self.retry_params, + || async { + self.s3_client + .complete_multipart_upload() + .bucket(self.bucket.clone()) + .key(key) + .multipart_upload(completed_upload.clone()) + .upload_id(upload_id) + .send() + .await + }, + ) .await?; Ok(()) } async fn abort_multipart_upload(&self, key: &str, upload_id: &str) -> StorageResult<()> { - aws_retry(&self.retry_params, || async { + aws_retry("s3.abort_multipart_upload", &self.retry_params, || async { self.s3_client .abort_multipart_upload() .bucket(self.bucket.clone()) @@ -687,7 +691,7 @@ impl S3CompatibleObjectStorage { path: &Path, range_opt: Option>, ) -> StorageResult { - let get_object_output = aws_retry(&self.retry_params, || { + let get_object_output = aws_retry("s3.get_object", &self.retry_params, || { self.get_object(path, range_opt.clone()) }) .await?; @@ -761,7 +765,7 @@ impl S3CompatibleObjectStorage { for (path_chunk, delete) in &mut delete_requests_it { let delete_objects_res: StorageResult = - aws_retry(&self.retry_params, || async { + aws_retry("s3.delete_objects", &self.retry_params, || async { crate::metrics::OBJECT_STORAGE_BULK_DELETE_REQUESTS_TOTAL.inc(); let _timer = HistogramTimer::new( &crate::metrics::OBJECT_STORAGE_BULK_DELETE_REQUEST_DURATION, @@ -899,8 +903,10 @@ impl Storage for S3CompatibleObjectStorage { #[instrument(name = "storage.s3.copy_to", level = "debug", skip(self, output))] async fn copy_to(&self, path: &Path, output: &mut dyn SendableAsync) -> StorageResult<()> { let _permit = REQUEST_SEMAPHORE.acquire().await; - let get_object_output = - aws_retry(&self.retry_params, || self.get_object(path, None)).await?; + let get_object_output = aws_retry("s3.get_object", &self.retry_params, || { + self.get_object(path, None) + }) + .await?; let mut body_read = BufReader::new(get_object_output.body.into_async_read()); let num_bytes_copied = tokio::io::copy_buf(&mut body_read, output).await?; crate::metrics::OBJECT_STORAGE_DOWNLOAD_NUM_BYTES.inc_by(num_bytes_copied); @@ -913,7 +919,7 @@ impl Storage for S3CompatibleObjectStorage { let _permit = REQUEST_SEMAPHORE.acquire().await; let bucket = self.bucket.clone(); let key = self.key(path); - let delete_res = aws_retry(&self.retry_params, || async { + let delete_res = aws_retry("s3.delete_object", &self.retry_params, || async { crate::metrics::OBJECT_STORAGE_DELETE_REQUESTS_TOTAL.inc(); let _timer = HistogramTimer::new(&crate::metrics::OBJECT_STORAGE_DELETE_REQUEST_DURATION); @@ -965,7 +971,7 @@ impl Storage for S3CompatibleObjectStorage { range: Range, ) -> crate::StorageResult> { let permit = REQUEST_SEMAPHORE.acquire().await; - let get_object_output = aws_retry(&self.retry_params, || { + let get_object_output = aws_retry("s3.get_object", &self.retry_params, || { self.get_object(path, Some(range.clone())) }) .await?; @@ -1003,7 +1009,7 @@ impl Storage for S3CompatibleObjectStorage { let _permit = REQUEST_SEMAPHORE.acquire().await; let bucket = self.bucket.clone(); let key = self.key(path); - let head_object_output = aws_retry(&self.retry_params, || async { + let head_object_output = aws_retry("s3.head_object", &self.retry_params, || async { self.s3_client .head_object() .bucket(&bucket)