mirror of
https://github.com/windmill-labs/windmill.git
synced 2026-10-09 00:02:30 +00:00
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
This commit is contained in:
@@ -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() {
|
||||
|
||||
@@ -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<S3ProxyKey, BatchState> = 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
|
||||
}
|
||||
@@ -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<B> OnResponse<B> for MyOnResponse {
|
||||
_span: &tracing::Span,
|
||||
) {
|
||||
if *LOG_REQUESTS {
|
||||
if let Some(s3req) = response.extensions().get::<S3ProxyRequest>() {
|
||||
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() {
|
||||
|
||||
Reference in New Issue
Block a user