From a3073ad8244efd9043e27f6731f7b53dbda662c1 Mon Sep 17 00:00:00 2001 From: Diego Imbert <70353967+diegoimbert@users.noreply.github.com> Date: Fri, 3 Apr 2026 15:46:09 +0200 Subject: [PATCH] fix: debounce S3 proxy logs (#8694) * Debounce S3 proxy logs * missing workspace id * nit perf * nit * prevent DOS * handle 4xx/5xx statuses * fix magic numbers --- backend/windmill-api/src/lib.rs | 27 ++-- backend/windmill-api/src/s3_log_batching.rs | 130 ++++++++++++++++++++ backend/windmill-api/src/tracing_init.rs | 10 ++ 3 files changed, 160 insertions(+), 7 deletions(-) create mode 100644 backend/windmill-api/src/s3_log_batching.rs diff --git a/backend/windmill-api/src/lib.rs b/backend/windmill-api/src/lib.rs index fd356de67f..4f66945396 100644 --- a/backend/windmill-api/src/lib.rs +++ b/backend/windmill-api/src/lib.rs @@ -24,6 +24,7 @@ use tower_http::catch_panic::CatchPanicLayer; use crate::tracing_init::MyOnFailure; use crate::{ + s3_log_batching::{s3_proxy_log_middleware, FLUSH_INTERVAL_MS}, tracing_init::{MyMakeSpan, MyOnResponse}, users::OptAuthed, webhook_util::WebhookShared, @@ -154,6 +155,7 @@ mod teams_approvals_oss; pub mod native_triggers; mod public_app_layer; mod public_app_rate_limit; +mod s3_log_batching; mod static_assets; #[cfg(all(feature = "stripe", feature = "enterprise", feature = "private"))] pub mod stripe_ee; @@ -952,13 +954,24 @@ pub async fn run_server( let app = if disable_response_logs { app } else { - app.layer( - TraceLayer::new_for_http() - .on_response(MyOnResponse {}) - .make_span_with(MyMakeSpan {}) - .on_request(()) - .on_failure(MyOnFailure {}), - ) + tokio::spawn(async { + let mut interval = + tokio::time::interval(std::time::Duration::from_millis(FLUSH_INTERVAL_MS)); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + loop { + interval.tick().await; + crate::s3_log_batching::flush_s3_batches(); + } + }); + + app.layer(axum::middleware::from_fn(s3_proxy_log_middleware)) + .layer( + TraceLayer::new_for_http() + .on_response(MyOnResponse {}) + .make_span_with(MyMakeSpan {}) + .on_request(()) + .on_failure(MyOnFailure {}), + ) }; let app = if let Some(domain) = public_app_layer::PUBLIC_APP_DOMAIN.as_ref() { diff --git a/backend/windmill-api/src/s3_log_batching.rs b/backend/windmill-api/src/s3_log_batching.rs new file mode 100644 index 0000000000..4c1d0e8e44 --- /dev/null +++ b/backend/windmill-api/src/s3_log_batching.rs @@ -0,0 +1,130 @@ +// Ducklake issues many rapid, identical S3 proxy requests (same object, method, status) +// within milliseconds, flooding the logs. This module debounces those logs: +// - The first request in a batch logs immediately. +// - Subsequent identical requests within 500ms are silently counted. +// - When 500ms pass with no new request (or 5s max), a single summary line is emitted +// with the total count and average latency, then the batch is cleared. +// A background task (spawned in lib.rs) calls flush_s3_batches() every 250ms to emit +// and clean up expired batches. + +use axum::{extract::Request, middleware::Next, response::Response}; +use dashmap::DashMap; +use std::time::Instant; + +lazy_static::lazy_static! { + static ref S3_LOG_BATCHES: DashMap = DashMap::new(); +} + +#[derive(Clone)] +pub(crate) struct S3ProxyRequest { + pub uri: String, + pub method: String, +} + +#[derive(Hash, Eq, PartialEq)] +struct S3ProxyKey { + workspace: String, + object_key: String, + method: String, + status: u16, +} + +struct BatchState { + total_count: u64, + total_latency_ms: u128, + first_seen: Instant, + last_seen: Instant, + uri: String, +} + +// URI format: /api/w/{workspace}/s3_proxy/{object_key} +fn extract_workspace_and_object_key(uri: &str) -> (&str, &str) { + // Extract workspace from /api/w/{workspace}/... + let workspace = uri + .strip_prefix("/api/w/") + .and_then(|s| s.split('/').next()) + .unwrap_or(""); + let object_key = uri + .find("/s3_proxy/") + .map(|i| &uri[i + "/s3_proxy/".len()..]) + .unwrap_or(uri); + (workspace, object_key) +} + +const MAX_BATCH_ENTRIES: usize = 10_000; +pub const FLUSH_INTERVAL_MS: u64 = 250; +const IDLE_TIMEOUT_MS: u128 = 500; +const MAX_BATCH_AGE_SECS: u64 = 5; + +pub(crate) fn record_s3_log(uri: &str, method: &str, status: u16, latency_ms: u128) { + // Prevent DOS by limiting the number of entries in the batch + if S3_LOG_BATCHES.len() >= MAX_BATCH_ENTRIES { + tracing::info!(latency = latency_ms, status = status, "response"); + return; + } + + let (workspace, object_key) = extract_workspace_and_object_key(uri); + let key = S3ProxyKey { + workspace: workspace.to_string(), + object_key: object_key.to_string(), + method: method.to_string(), + status, + }; + + let mut entry = S3_LOG_BATCHES.entry(key).or_insert_with(|| BatchState { + total_count: 0, + total_latency_ms: 0, + first_seen: Instant::now(), + last_seen: Instant::now(), + uri: uri.to_string(), + }); + + entry.total_count += 1; + entry.total_latency_ms += latency_ms; + entry.last_seen = Instant::now(); + + if entry.total_count == 1 { + tracing::info!(latency = latency_ms, status = status, "response") + } +} + +pub fn flush_s3_batches() { + let now = Instant::now(); + S3_LOG_BATCHES.retain(|key, batch| { + let idle = now.duration_since(batch.last_seen).as_millis() > IDLE_TIMEOUT_MS; + let max_age = now.duration_since(batch.first_seen).as_secs() > MAX_BATCH_AGE_SECS; + + if idle || max_age { + if batch.total_count > 1 { + let avg_latency = batch.total_latency_ms / batch.total_count as u128; + tracing::info!( + count = batch.total_count, + avg_latency = avg_latency, + status = key.status, + method = %key.method, + uri = %batch.uri, + "s3_proxy batched response" + ); + } + false + } else { + true + } + }); +} + +pub async fn s3_proxy_log_middleware(request: Request, next: Next) -> Response { + let s3_info = if request.uri().path().contains("/s3_proxy/") { + Some(S3ProxyRequest { + uri: request.uri().to_string(), + method: request.method().to_string(), + }) + } else { + None + }; + let mut response = next.run(request).await; + if let Some(info) = s3_info { + response.extensions_mut().insert(info); + } + response +} diff --git a/backend/windmill-api/src/tracing_init.rs b/backend/windmill-api/src/tracing_init.rs index 8ae3eabf47..de6fba55d2 100644 --- a/backend/windmill-api/src/tracing_init.rs +++ b/backend/windmill-api/src/tracing_init.rs @@ -6,6 +6,7 @@ * LICENSE-AGPL for a copy of the license. */ +use crate::s3_log_batching::{record_s3_log, S3ProxyRequest}; use ::tracing::{field, Span}; use hyper::Response; use tower_http::trace::{MakeSpan, OnFailure, OnResponse}; @@ -29,6 +30,15 @@ impl OnResponse for MyOnResponse { _span: &tracing::Span, ) { if *LOG_REQUESTS { + if let Some(s3req) = response.extensions().get::() { + let status = response.status().as_u16(); + if response.status().is_success() || response.status().is_redirection() { + record_s3_log(&s3req.uri, &s3req.method, status, latency.as_millis()); + return; + } + // 4xx/5xx errors are rare and shouldn't be batched + } + let latency = latency.as_millis(); let status = response.status().as_u16(); if response.status().is_success() || response.status().is_redirection() {