diff --git a/server/src/compact.rs b/server/src/compact.rs index 9d9fc94e32..578fea3ecb 100644 --- a/server/src/compact.rs +++ b/server/src/compact.rs @@ -113,16 +113,25 @@ impl CompactWorker { } } - pub async fn send(&self, parts: &Vec>) { - if parts.len() > CONFIG.compact_parts_limit { + pub async fn try_send(&self, parts: &Vec>) -> bool { + if parts.len() >= CONFIG.compact_parts_limit { let task = CompactTask { workspace: parts[0].data.workspace, key: parts[0].data.key.clone(), }; + self.send(task).await + } else { + false + } + } - let res = self.ingest_tx.send(task.clone()).await; - if let Err(err) = res { + pub async fn send(&self, task: CompactTask) -> bool { + let res = self.ingest_tx.send(task).await; + match res { + Ok(_) => true, + Err(err) => { warn!(%err, "failed to schedule compact"); + false } } } diff --git a/server/src/handlers.rs b/server/src/handlers.rs index 43d86e4a90..207a58116e 100644 --- a/server/src/handlers.rs +++ b/server/src/handlers.rs @@ -434,7 +434,7 @@ pub async fn get(request: HttpRequest) -> HandlerResult { } None => { let compact = request.app_data::>().unwrap(); - compact.send(&parts).await; + compact.try_send(&parts).await; let stream = merge::stream(s3.clone(), parts).await?; response.body(SizedStream::new(stream.content_length, stream.stream))