Skip to content
Open
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
7 changes: 7 additions & 0 deletions crates/app-server/src/chat_completions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -390,6 +390,13 @@ pub(crate) async fn chat_completions_handler(
// Build upstream request. The shared platform builder sets no timeouts,
// so the proxy would hang forever on an accept-and-stall upstream;
// bound both the connect and the whole non-streaming round trip.
//
// `config` (the read guard) is dropped here on purpose: everything it
// feeds is resolved by now, and holding it across the upstream round
// trip (up to UPSTREAM_TOTAL_TIMEOUT) would block config writes and,
// through the write-preferring lock, every later resolve for that
// whole window.
drop(config);
let upstream_req = codewhale_release::platform_http_client_builder()
.connect_timeout(UPSTREAM_CONNECT_TIMEOUT)
.timeout(UPSTREAM_TOTAL_TIMEOUT)
Expand Down
573 changes: 536 additions & 37 deletions crates/app-server/src/lib.rs

Large diffs are not rendered by default.

6 changes: 3 additions & 3 deletions crates/cli/src/update.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,9 +32,9 @@ const GITHUB_RELEASE_DOWNLOAD_BASE_URL: &str =
const UPDATE_HTTP_ATTEMPTS: usize = 3;
const UPDATE_HTTP_RETRY_DELAY_MS: u64 = 100;
/// Ceiling for one asset download. Release binaries are tens of megabytes
/// and some of the networks this exists for are slow: 600s covers a full
/// 60 MiB at ~100 KiB/s (the same download budget the audit gives skill
/// tarballs; the old 300s needed an implausible >1.6 Mbps to finish).
/// and some of the networks this exists for are slow: 600s covers ~58 MiB
/// at ~100 KiB/s (the same download budget the audit gives skill tarballs;
/// the old 300s needed an implausible >1.6 Mbps to finish).
const UPDATE_DOWNLOAD_TIMEOUT: Duration = Duration::from_secs(600);
/// Ceiling for one checksum-manifest probe. The manifest is a few hundred
/// bytes, so this is only a backstop against a source that accepts the
Expand Down
14 changes: 13 additions & 1 deletion crates/core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,19 @@ use uuid::Uuid;
/// Per-tool dispatch budget for the headless runtime. 30 minutes: tools
/// legitimately run long (builds, test suites, MCP-backed calls), and this
/// wrapper is a runaway backstop, not an expected duration — the previous
/// 300s value cut off healthy in-flight tool work.
/// 300s value cut off healthy in-flight tool work. Part of the 1800s family
/// (TUI client envelope, vision request envelope, model stream cap
/// `STREAM_MAX_DURATION_SECS`, dynamic-tool result wait, sub-agent tool
/// timeout, MCP execute timeout, the mirrors in app-server and the MCP stdio
/// proxy, the background-task wall clock `TaskExecutionLimits::wall_time`,
/// the sub-agent `default_wall_time_secs`, and the fleet `builder` role
/// preset) that comments keep in sync; this comment anchors the family
/// roster — there is no shared constant across the crates yet.
///
/// Several family members are configurable defaults rather than constants
/// (`STREAM_MAX_DURATION_SECS`, the MCP `execute_timeout`, and the sub-agent
/// `default_wall_time_secs`): a family-wide bump changes their defaults, not
/// their ceilings, and user overrides survive it.
fn tool_dispatch_timeout() -> Duration {
if cfg!(test) {
Duration::from_millis(50)
Expand Down
20 changes: 20 additions & 0 deletions crates/mcp/src/stdio_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -899,6 +899,26 @@ impl Drop for Connection {
}
}

#[cfg(test)]
mod budget_pins {
// The engine-side stdio proxy has no behavioral timeout tests (the
// budgets are compile-time constants), so pin the wiring values: a
// drift here silently re-bounds every MCP `tools/call` the engine
// proxy carries. CALL_TOOL_TIMEOUT must mirror the TUI pool's default
// execute timeout (1800s), and the generic request budget must stay
// separate and much shorter.
use super::{CALL_TOOL_TIMEOUT, HANDSHAKE_TIMEOUT, REQUEST_TIMEOUT};
use std::time::Duration;

#[test]
fn call_tool_budget_mirrors_the_pool_default_and_stays_separate() {
assert_eq!(CALL_TOOL_TIMEOUT, Duration::from_secs(1800));
assert_eq!(REQUEST_TIMEOUT, Duration::from_secs(120));
assert!(CALL_TOOL_TIMEOUT > REQUEST_TIMEOUT);
assert_eq!(HANDSHAKE_TIMEOUT, Duration::from_secs(30));
}
}

#[cfg(test)]
mod tests {
use std::collections::HashMap;
Expand Down
1 change: 1 addition & 0 deletions crates/tui/src/automation_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2250,6 +2250,7 @@ mod tests {
allow_shell: true,
trust_mode: true,
execution_limits: crate::task_manager::TaskExecutionLimits::default(),
human_waits_answerable: false,
}
}

Expand Down
178 changes: 173 additions & 5 deletions crates/tui/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -354,7 +354,7 @@ pub(crate) const NON_STREAMING_REQUEST_ENVELOPE: Duration = Duration::from_secs(
static TEST_NON_STREAMING_ENVELOPE_MS: std::sync::atomic::AtomicU64 =
std::sync::atomic::AtomicU64::new(0);

fn non_streaming_request_envelope() -> Duration {
pub(crate) fn non_streaming_request_envelope() -> Duration {
#[cfg(test)]
{
let ms = TEST_NON_STREAMING_ENVELOPE_MS.load(std::sync::atomic::Ordering::SeqCst);
Expand Down Expand Up @@ -3022,6 +3022,11 @@ impl DeepSeekClient {
crate::retry_status::failed(last.to_string());
self.mark_request_failure("non-streaming request envelope exceeded")
.await;
// A provider that just wedged past the whole envelope is
// exactly the case where the /models health probe
// matters; without it, connection health stays degraded
// until the next real request succeeds.
self.maybe_probe_recovery().await;
return Err(anyhow::Error::new(last));
}
}
Expand Down Expand Up @@ -4557,6 +4562,76 @@ mod tests {
);
}

#[tokio::test]
async fn envelope_exceed_probes_recovery_once_degraded() {
// A provider that wedges past the whole envelope must not leave
// connection health stale: once the failure threshold has degraded
// the connection, the envelope exit probes /models so health
// recovers without waiting for the next real request. One exceed
// alone stays under the threshold and must not probe (the
// Healthy-state short-circuit), so this drives exactly two.
let _env_lock = crate::test_support::lock_test_env();
let _envelope = NonStreamingEnvelopeGuard::millis(2000);
let server = MockServer::start().await;
Mock::given(method("POST"))
.respond_with(
ResponseTemplate::new(200)
.set_body_json(json!({
"model": "deepseek-v4-pro",
"choices": [
{ "message": { "role": "assistant", "content": "late" } }
]
}))
.set_delay(Duration::from_secs(5)),
)
.expect(2)
.mount(&server)
.await;
Mock::given(method("GET"))
.and(path("/v1/models"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({ "data": [] })))
.expect(1)
.mount(&server)
.await;

let client = deepseek_request_boundary_client(&server.uri(), server.uri());
let make_request = || MessageRequest {
model: "deepseek-v4-pro".to_string(),
messages: vec![Message {
role: Role::User,
content: vec![ContentBlock::Text {
text: "envelope".to_string(),
cache_control: None,
}],
}],
max_tokens: 16,
system: None,
tools: None,
tool_choice: None,
metadata: None,
thinking: None,
reasoning_effort: Some("off".to_string()),
stream: Some(false),
temperature: None,
top_p: None,
};

for _ in 0..2 {
let err = client
.create_message(make_request())
.await
.expect_err("a provider that never answers must hit the envelope");
assert!(
err.to_string().to_lowercase().contains("timed out"),
"envelope timeout must be reported as such; got {err:#}"
);
}

// verify() enforces the expectations above: exactly two stalled
// requests and exactly one /models recovery probe.
server.verify().await;
}

#[tokio::test]
async fn stream_open_retry_path_sets_no_total_deadline() {
// The injected budget is process-global; serialize against other
Expand All @@ -4574,7 +4649,7 @@ mod tests {
.respond_with(
ResponseTemplate::new(200)
.set_body_string("data: [DONE]\n\n")
.set_delay(Duration::from_millis(2500)),
.set_delay(Duration::from_millis(4000)),
)
.mount(&server)
.await;
Expand Down Expand Up @@ -4606,19 +4681,19 @@ mod tests {
// A caller that pins its own, larger total (`list_models` pins 30s)
// must keep it: `.timeout()` on the builder is a pure overwrite, so
// an unconditional envelope would silently replace the pinned
// budget with the injected 2s and fail this 2.5s-late response.
// budget with the injected 2s and fail this 4s-late response.
Mock::given(method("GET"))
.respond_with(
ResponseTemplate::new(200)
.set_body_string("{}")
.set_delay(Duration::from_millis(2500)),
.set_delay(Duration::from_millis(4000)),
)
.mount(&server)
.await;

let client = deepseek_request_boundary_client(&server.uri(), server.uri());
let response = client
.send_with_retry_total(Duration::from_secs(5), || {
.send_with_retry_total(Duration::from_secs(12), || {
client.http_client.get(format!("{}/models", server.uri()))
})
.await
Expand Down Expand Up @@ -8325,6 +8400,11 @@ mod tests {

#[tokio::test]
async fn deepseek_anthropic_translate_uses_messages_endpoint() {
// The Anthropic dialect resolves its budget through the
// process-global `non_streaming_request_envelope()`, so this must
// hold the same lock the injecting tests take; otherwise a parallel
// `NonStreamingEnvelopeGuard` can cut this request off mid-flight.
let _env_lock = crate::test_support::lock_test_env();
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/v1/messages"))
Expand Down Expand Up @@ -8452,6 +8532,9 @@ mod tests {

#[tokio::test]
async fn minimax_anthropic_request_uses_messages_endpoint() {
// Same process-global envelope hazard as the DeepSeek sibling above,
// and its neighbour test injects a 1s budget.
let _env_lock = crate::test_support::lock_test_env();
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/anthropic/v1/messages"))
Expand Down Expand Up @@ -8509,6 +8592,91 @@ mod tests {
assert!(body.get("output_config").is_none(), "{body}");
}

#[tokio::test]
async fn anthropic_non_streaming_request_carries_the_envelope_and_probes() {
// The Anthropic dialect sends one direct request with no retry
// loop; its total budget must match every other non-streaming
// completion and must actually cut off a stalled provider. The
// transport-level cutoff never reaches the status checks, so this
// arm must mark the failure and probe /models itself — otherwise
// connection health stays stale until the next request. One
// envelope exceed alone stays under the degradation threshold, so
// the health-marking is asserted indirectly: exactly one POST and
// one probe GET reach the server (verified via server.verify()).
let _env_lock = crate::test_support::lock_test_env();
let _envelope = NonStreamingEnvelopeGuard::millis(1000);
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/anthropic/v1/messages"))
.respond_with(
ResponseTemplate::new(200)
.set_body_json(json!({
"id": "msg_slow",
"type": "message",
"role": "assistant",
"content": [{"type": "text", "text": "late"}],
"model": "MiniMax-M3",
"stop_reason": "end_turn",
"stop_sequence": null,
"usage": {"input_tokens": 1, "output_tokens": 1}
}))
.set_delay(Duration::from_millis(4000)),
)
.expect(2)
.mount(&server)
.await;
// The recovery probe: the transport arm runs it immediately after
// marking the failure. Routing the client's base_url at the mock
// (the way the health-check tests do) puts the probe on
// `<mock>/anthropic/v1/models`, so the pin stays hermetic.
Mock::given(method("GET"))
.and(path("/anthropic/v1/models"))
.respond_with(ResponseTemplate::new(200).set_body_json(json!({ "data": [] })))
.expect(1)
.mount(&server)
.await;

let mut client =
minimax_anthropic_client_with_base_url(format!("{}/anthropic", server.uri()));
client.test_messages_transport_base_url = Some(format!("{}/anthropic", server.uri()));
let make_request = || MessageRequest {
model: "MiniMax-M3".to_string(),
messages: vec![Message {
role: Role::User,
content: vec![ContentBlock::Text {
text: "hello".to_string(),
cache_control: None,
}],
}],
max_tokens: 32,
system: None,
tools: None,
tool_choice: None,
metadata: None,
thinking: None,
reasoning_effort: Some("off".to_string()),
stream: Some(false),
temperature: None,
top_p: None,
};
for _ in 0..2 {
let err = client
.create_message(make_request())
.await
.expect_err("a provider that answers past the envelope must be cut off");
let rendered = format!("{err:#}");
assert!(
rendered.contains("timed out"),
"the cutoff must surface as a timeout, not a generic failure: {rendered}"
);
}
// Enforces the expectations above: exactly two stalled POSTs (one
// alone stays under the degradation threshold) and — the point of
// this pin — exactly one /models recovery probe fired by the
// transport arm's own failure handling.
server.verify().await;
}

#[test]
fn custom_api_key_header_is_allowed_without_primary_provider_key() {
let mut extra = HashMap::new();
Expand Down
24 changes: 21 additions & 3 deletions crates/tui/src/client/anthropic.rs
Original file line number Diff line number Diff line change
Expand Up @@ -222,11 +222,29 @@ impl DeepSeekClient {
.http_client
.post(&url)
.header("Accept", "text/event-stream")
.timeout(crate::client::NON_STREAMING_REQUEST_ENVELOPE)
.timeout(crate::client::non_streaming_request_envelope())
.json(body)
.send()
.await
.context("Anthropic Messages API request failed")?;
.await;
let response = match response {
Ok(response) => response,
Err(error) => {
let rendered = if error.is_timeout() {
"anthropic request exceeded the envelope".to_string()
} else {
format!("anthropic transport failure: {error}")
};
// A transport-level failure never reaches the status checks
// below, so it must mark health and probe on its own —
// otherwise a provider that stalls past the envelope leaves
// connection health stale until the next request retries.
self.mark_request_failure(&rendered).await;
self.maybe_probe_recovery().await;
return Err(
anyhow::Error::new(error).context("Anthropic Messages API request failed")
);
}
};
self.check_anthropic_response(response).await
}

Expand Down
6 changes: 5 additions & 1 deletion crates/tui/src/config/subagent_limits.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,11 @@ pub const MAX_SUBAGENT_API_TIMEOUT_SECS: u64 = 3600;
/// legitimately outlasts 5 minutes, and the old default killed healthy
/// in-flight tools mid-run. The child's own wall-time budget remains the
/// spend backstop, and the heartbeat floor (tool_timeout + 30s) follows this
/// constant automatically.
/// constant automatically. Part of the 1800s family that comments keep in
/// sync (roster anchored at the core dispatch backstop's comment);
/// "single source of truth" applies within the sub-agent family.
/// In a background task the task `wall_time` backstop bounds the whole run
/// and preempts this timeout when they coincide — documented trade-off.
pub const DEFAULT_SUBAGENT_TOOL_TIMEOUT_SECS: u64 = 1800;
/// Default wall-clock interval without manager-visible sub-agent progress
/// before a running child can be auto-cancelled to release its slot (#2614).
Expand Down
Loading
Loading