Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 7 additions & 7 deletions atom/compass/audit/sync_sites.json
Original file line number Diff line number Diff line change
Expand Up @@ -2483,27 +2483,27 @@
"anchor": "pub const DEFAULT_WORKER_HTTP_TIMEOUT_SECS",
"category": "C1",

@jgong5 jgong5 Sep 29, 2026 •

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reservation, not blocking, and a landing-order note. This row keeps category C1 while its new, correct why says it bounds no request; changing the category changes the audit README counts, outside this PR, so it is recorded on #477. This file conflicts with #476: whichever lands second must keep this PR's file/line/anchor/why on all four rows and add #476's mechanism fields.

Detail

Category: the audit README defines C1 as "a bound that exists to declare something broken: raise it, or switch it off", and counts the row in "five router and server bounds". The why is correct: the client has one use (http_health_check), and reqwest 0.12.28 RequestConfig::fetch returns the request's timeout and falls back to the client default only when the request has none, so it replaces the default rather than taking the minimum.

Conflict: git merge-tree --write-tree 67117f738 b28823d9d reports CONFLICT (content) here. #476 (delivers #454) still has worker_manager.rs:24 const REQUEST_TIMEOUT and cliargs.rs:423/:398 on these rows, and gives this row mechanism: K8. #476 carries need human, so this PR will most likely land first, and #476's base-update merge must do the keeping. If it keeps #476's line numbers, test_anchor_lines_are_still_where_they_say fails with exactly the three messages recorded in this PR's body. K8 on this row needs a second look for the same reason as the category.

"peer": "deployment",
"why": "the router's thirty-second bound on one worker request, compiled in; it fires whenever simulated time runs slower than real"
"why": "the default timeout of the client the router uses only for HTTP health checks; each check sets its own timeout from --health-check-timeout-secs, which replaces this one, so it bounds no request"
},
{
"file": "atom/mesh/src/core/worker_manager.rs",
"line": 24,
"anchor": "const REQUEST_TIMEOUT",
"file": "atom/mesh/src/cliargs.rs",
"line": 334,
"anchor": "pub worker_request_timeout_secs",
"category": "C1",
"peer": "deployment",
"why": "the router's five-second bound on a fan-out request, compiled in; it fires whenever simulated time runs slower than real"
"why": "the router's five-second default bound on its /get_load and /flush_cache requests to workers; raise it from the command line when simulated time runs slower than real"
},
{
"file": "atom/mesh/src/cliargs.rs",
"line": 423,
"line": 430,
"anchor": "pub disable_health_check",
"category": "C1",
"peer": "deployment",
"why": "the router's worker health check, already switchable from the command line"
},
{
"file": "atom/mesh/src/cliargs.rs",
"line": 398,
"line": 405,
"anchor": "pub disable_circuit_breaker",
"category": "C1",
"peer": "deployment",
Expand Down
13 changes: 11 additions & 2 deletions atom/mesh/src/cliargs.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
use std::sync::Arc;
use std::sync::{atomic::Ordering, Arc};

use clap::{ArgAction, Parser, Subcommand, ValueEnum};

Expand All @@ -8,7 +8,10 @@ use crate::{
HealthCheckConfig, MetricsConfig, PolicyConfig, RetryConfig, RouterConfig, RoutingMode,
TokenizerCacheConfig,
},
core::ConnectionMode,
core::{
worker_manager::{DEFAULT_WORKER_REQUEST_TIMEOUT_SECS, WORKER_REQUEST_TIMEOUT_SECS},
ConnectionMode,
},
observability::metrics::PrometheusConfig,
routers::atom_standalone::AtomStandaloneRuntime,
server::{ServerConfig, ServerTlsConfig},
Expand Down Expand Up @@ -326,6 +329,10 @@ pub struct CliArgs {
#[arg(long, default_value_t = 1800, help_heading = "Request Handling")]
pub request_timeout_secs: u64,

/// Timeout in seconds for the router's own /get_load and /flush_cache requests to workers
#[arg(long, default_value_t = DEFAULT_WORKER_REQUEST_TIMEOUT_SECS, value_parser = clap::value_parser!(u64).range(1..), help_heading = "Request Handling")]
pub worker_request_timeout_secs: u64,

/// Grace period in seconds to wait for in-flight requests during shutdown
#[arg(long, default_value_t = 180, help_heading = "Request Handling")]
pub shutdown_grace_period_secs: u64,
Expand Down Expand Up @@ -551,6 +558,7 @@ impl CliArgs {
prefill_urls: Vec<(String, Option<u16>)>,
) -> ConfigResult<RouterConfig> {
self.validate_tls_args()?;
WORKER_REQUEST_TIMEOUT_SECS.store(self.worker_request_timeout_secs, Ordering::Relaxed);

// Determine routing mode based on PD disaggregation flag
let mode = if self.pd_disaggregation {
Expand Down Expand Up @@ -799,6 +807,7 @@ impl Default for CliArgs {
prometheus_duration_buckets: Vec::new(),
request_id_headers: Vec::new(),
request_timeout_secs: 1800,
worker_request_timeout_secs: DEFAULT_WORKER_REQUEST_TIMEOUT_SECS,
shutdown_grace_period_secs: 180,
max_payload_size: 536_870_912,
max_concurrent_requests: -1,
Expand Down
95 changes: 91 additions & 4 deletions atom/mesh/src/core/worker_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,14 @@
//!
//! Provides worker lifecycle operations and fan-out request utilities.

use std::{collections::HashMap, sync::Arc, time::Duration};
use std::{
collections::HashMap,
sync::{
atomic::{AtomicU64, Ordering},
Arc,
},
time::Duration,
};

use futures::{
future,
Expand All @@ -21,7 +28,18 @@ use crate::{
protocols::worker_spec::{FlushCacheResult, WorkerLoadInfo, WorkerLoadsResult},
};

const REQUEST_TIMEOUT: Duration = Duration::from_secs(5);
pub const DEFAULT_WORKER_REQUEST_TIMEOUT_SECS: u64 = 5;

/// Timeout for the `/flush_cache` and `/get_load` requests below; set from
/// `--worker-request-timeout-secs` when the router config is built.
/// Process-wide: the last router config built in a process wins.
pub static WORKER_REQUEST_TIMEOUT_SECS: AtomicU64 =

@jgong5 jgong5 Sep 29, 2026 •

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reservation, not blocking: the process-wide static is acceptable and the leaner choice. Suggestion: one comment line on the static naming the ceiling, e.g. "process-wide: the last router config built wins; move into RouterConfig if a process ever builds two."

What the alternative costs, and the static's costs, at this head

Moving the value into RouterConfig would touch config/types.rs, server.rs (two call sites) and app_context.rs (LoadMonitor::new), and change the signatures of flush_cache_all, get_all_worker_loads and LoadMonitor: about 20 lines in 3 more files, no difference in behaviour. Simplicity outweighs a cleaner abstraction at that price.

  • It skips RouterConfig::validate(); finding 1 fixes that at the clap layer.
  • The last to_router_config in a process wins. Today there is one caller per process: main in main.rs, and build_server_config in python.rs once per launch path in atom/entrypoints/atomesh/server.py.
  • The test leaves the static at 60 for the rest of the lib test process. No other lib test calls to_router_config or the fan-out functions, so nothing reads it.

AtomicU64::new(DEFAULT_WORKER_REQUEST_TIMEOUT_SECS);

fn request_timeout() -> Duration {
Duration::from_secs(WORKER_REQUEST_TIMEOUT_SECS.load(Ordering::Relaxed))
}

const MAX_CONCURRENT: usize = 32;

/// Result of a fan-out request to a single worker
Expand All @@ -47,7 +65,7 @@ async fn fan_out(
let method = method.clone();

async move {
let mut req = client.request(method, &full_url).timeout(REQUEST_TIMEOUT);
let mut req = client.request(method, &full_url).timeout(request_timeout());
if let Some(key) = api_key {
req = req.bearer_auth(key);
}
Expand Down Expand Up @@ -194,7 +212,7 @@ impl WorkerManager {
api_key: Option<&str>,
) -> isize {
let load_url = format!("{}/get_load", url);
let mut req = client.get(&load_url).timeout(REQUEST_TIMEOUT);
let mut req = client.get(&load_url).timeout(request_timeout());
if let Some(key) = api_key {
req = req.bearer_auth(key);
}
Expand Down Expand Up @@ -340,3 +358,72 @@ impl Drop for LoadMonitor {
}
}
}

#[cfg(test)]
mod tests {
use std::time::Instant;

use clap::Parser;

use super::*;
use crate::{cliargs::CliArgs, core::BasicWorkerBuilder};

#[tokio::test]
async fn worker_request_timeout_option_outlasts_a_40s_worker() {
let slow = Duration::from_secs(40);
let app = axum::Router::new()
.route(
"/get_load",
axum::routing::get(move || async move {
tokio::time::sleep(slow).await;
axum::Json(serde_json::json!([{ "num_tokens": 7 }]))
}),
)
.route(
"/flush_cache",
axum::routing::post(move || tokio::time::sleep(slow)),
);
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let url = format!("http://{}", listener.local_addr().unwrap());
tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });

let registry = WorkerRegistry::new();
registry.register(Arc::new(BasicWorkerBuilder::new(&url).build()));
let client = reqwest::Client::new();

for (flags, answered) in [
(&[][..], false),
(&["--worker-request-timeout-secs", "60"][..], true),
] {
let args = CliArgs::parse_from(["atomesh"].iter().chain(flags));
args.to_router_config(vec![]).unwrap();
let start = Instant::now();
let (loads, flush) = tokio::join!(
WorkerManager::get_all_worker_loads(&registry, &client),
WorkerManager::flush_cache_all(&registry, &client),
);
let elapsed = start.elapsed();
eprintln!(
"flags={flags:?} load={} flushed={} elapsed={elapsed:?}",
loads.loads[0].load,
flush.successful.len()
);
assert_eq!(loads.loads[0].load, if answered { 7 } else { -1 });
assert_eq!(flush.successful.len(), answered as usize);
assert_eq!(elapsed >= slow, answered);

@jgong5 jgong5 Sep 29, 2026 •

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Required finding 2: "defaults unchanged" is pinned only for a default of 40 s or more, because the no-flag leg's only time bound is elapsed < 40 s; the body's M4 row ("a changed default") overstates it. Fix, one line: assert!(answered || elapsed < Duration::from_secs(10)); (the leg measured 5.002 s), or assert CliArgs::parse_from(["atomesh"]).worker_request_timeout_secs == 5. Either reddens on the 5 -> 30 mutant; please record that red.

The mutant

DEFAULT_WORKER_REQUEST_TIMEOUT_SECS: u64 = 5 -> 30, one line, nothing else changed, line counts 849/415:

  • published: 1 passed, 45.05 s (no-flag leg 5.002 s)
  • mutant: 1 passed, 70.05 s (no-flag leg 30.002 s, load -1, 0 flushed)

}
}

#[test]
fn worker_request_timeout_defaults_to_5s_and_refuses_zero() {
let parse = |flags: &[&str]| CliArgs::try_parse_from(["atomesh"].iter().chain(flags));
assert_eq!(parse(&[]).unwrap().worker_request_timeout_secs, 5);
assert_eq!(
parse(&["--worker-request-timeout-secs", "1"])
.unwrap()
.worker_request_timeout_secs,
1
);
assert!(parse(&["--worker-request-timeout-secs", "0"]).is_err());
}
}