diff --git a/README.md b/README.md index 7411f899..70702f2e 100644 --- a/README.md +++ b/README.md @@ -45,7 +45,7 @@ The design starts with one question: **What is the smallest capability Open Max | Lifecycle policy and events | Process hooks and permission files | | Compaction integration | The built-in model summary plus the `compaction` hook event | | Model endpoints | `providers.json`, including local servers, gateways, and private proxies | -| Shortcuts or a completely different UI | Prompt templates, or a custom frontend speaking `openmax-stdio/5` | +| Shortcuts or a completely different UI | Prompt templates, or a custom frontend speaking `openmax-stdio/6` | These are deliberate boundaries, not placeholders for hidden orchestration products. Open Max does not carry an MCP host, nested-agent scheduler, plan mode, background-job product, built-in TODO database, user-keybinding engine, pluggable compactor, or TUI plugin ABI. The agent composes those richer workflows from the same host tools a developer can inspect, edit, test, and remove. @@ -198,7 +198,7 @@ Sessions, settings, tools, and skills stay under `~/.openmax/` and your project - [Configuration](docs/configuration.md): settings, approvals, providers, project trust - [Usage](docs/usage.md): CLI flags, keybindings, slash commands - [Extending](docs/extending.md): tools, skills, templates, hooks, permissions, validation, freezing -- [stdio protocol](docs/stdio-protocol.md): the `openmax-stdio/5` contract for custom frontends +- [stdio protocol](docs/stdio-protocol.md): the `openmax-stdio/6` contract for custom frontends ## Development diff --git a/crates/core/src/agent.rs b/crates/core/src/agent.rs index 826e8a56..053988f3 100644 --- a/crates/core/src/agent.rs +++ b/crates/core/src/agent.rs @@ -2293,6 +2293,13 @@ impl TokenBatcher { match delta { StreamDelta::Content(t) => self.content.push_str(&t), StreamDelta::Reasoning(t) => self.thinking.push_str(&t), + // Flush first so the wire keeps the order the client saw: the + // reasoning already streamed, then the notice that it is void. + StreamDelta::Retry { attempt, max_attempts, reason } => { + self.flush(); + self.core.send_agent(&self.session_id, AgentEvent::Retry { attempt, max_attempts, reason }); + return; + } } if self.last_flush.elapsed() >= FLUSH_INTERVAL { self.flush(); diff --git a/crates/core/src/client.rs b/crates/core/src/client.rs index abf2e850..ce9dc0f1 100644 --- a/crates/core/src/client.rs +++ b/crates/core/src/client.rs @@ -5,10 +5,18 @@ //! `stream_options`, and provider-specific headers. The request body is //! serialized once before the retry loop so a retry resends identical bytes. //! -//! Retries are deliberately narrow: connect failures, timeouts, and 429 only, -//! and only BEFORE the first token arrives. Once a stream has produced output, -//! retrying would duplicate text the caller already saw, so a mid-stream -//! failure is reported instead. +//! Retries cover what a network can do to one request before the reply +//! exists: a transport failure on send, a 429, and a stream that died before +//! any reply text arrived. Each attempt resends the same bytes after an +//! exponential backoff and tells the caller through [`StreamDelta::Retry`]. +//! Once reply text has streamed, a retry would duplicate what the caller +//! already showed, so a failure after that point is reported as a +//! truncation instead; reasoning deltas are display-only and do not count. +//! A stream this client ends on purpose (the size cap, an out-of-range tool +//! index, cancellation) is never retried: the next attempt would end the +//! same way. A reply carrying only tool calls has no reply text, so a server +//! that never sends a completion signal costs the whole budget on such a +//! reply before its truncation is reported. //! //! There is no overall request timeout, only a connect timeout. A local or //! slow endpoint can legitimately take minutes to generate, and a deadline @@ -113,6 +121,15 @@ fn serialize_chat_request_body( pub enum StreamDelta { Content(String), Reasoning(String), + /// The request is being resent: `attempt` is the one about to go out of + /// the `max_attempts` budget, after `reason` ended the previous one. It + /// is announced before the backoff wait, which is what the caller is + /// waiting through; a cancellation during the wait ends the request + /// instead, and the attempt never goes out. Any reasoning streamed for + /// the failed attempt is void; no content was. A request gives up before + /// the budget when [`UNREACHED_ATTEMPTS`] in a row never reached the + /// endpoint. + Retry { attempt: u32, max_attempts: u32, reason: String }, } /// `finish_reason` for a stream the server never terminated: no `[DONE]` line @@ -323,14 +340,11 @@ impl ChatClient { tools_wire, )?; - // Retry only pre-stream transport failures (connect/timeout) and 429 - // (request rejected before work starts). Do not retry 502/503/504: a - // proxy may already have forwarded the POST and started a completion, - // so a second attempt can duplicate backend work with no idempotency key. - // Once SSE bytes start, failures fail cleanly without a second prefill. - const MAX_ATTEMPTS: u32 = 3; let mut attempt = 0u32; - let resp = loop { + // Attempts in a row that never reached the endpoint; any other + // outcome, including a later transport fault, resets it. + let mut unreached = 0u32; + loop { attempt += 1; let mut req = self .http @@ -350,192 +364,212 @@ impl ChatClient { // it runs; keep cancellation responsive throughout. let send_result = tokio::select! { r = req.send() => r, - _ = cancelled.cancelled() => { - return Ok(CompletionResult { - content: String::new(), - tool_calls: Vec::new(), - finish_reason: "cancelled".into(), - usage: None, - }); - } + _ = cancelled.cancelled() => return Ok(cancelled_response()), }; let resp = match send_result { Ok(r) => r, Err(e) => { let msg = format!("request failed: {}", describe_transport(&e)); - if attempt < MAX_ATTEMPTS && is_transient_transport(&e) { - backoff_sleep(attempt).await; - if cancelled.is_cancelled() { - return Ok(CompletionResult { - content: String::new(), - tool_calls: Vec::new(), - finish_reason: "cancelled".into(), - usage: None, - }); + unreached = if e.is_connect() || e.is_timeout() { unreached + 1 } else { 0 }; + if attempt < MAX_ATTEMPTS && unreached < UNREACHED_ATTEMPTS && is_transient_transport(&e) { + if !retry_after(attempt, &msg, &cancelled, &mut on_delta).await { + return Ok(cancelled_response()); } continue; } return Err(msg); } }; + unreached = 0; let status = resp.status(); - if status.is_success() { - break resp; + if !status.is_success() { + let code = status.as_u16(); + let body = read_body(resp, &cancelled).await + .map_err(|e| format!("backend returned {status}: {e}"))?; + let Some(body) = body else { return Ok(cancelled_response()); }; + let text = String::from_utf8_lossy(&body); + let err = format!("backend returned {status}: {}", describe_backend(&text)); + if attempt < MAX_ATTEMPTS && is_retryable_status(code) { + if !retry_after(attempt, &err, &cancelled, &mut on_delta).await { + return Ok(cancelled_response()); + } + continue; + } + return Err(err); } - let code = status.as_u16(); - let body = read_body(resp, &cancelled).await - .map_err(|e| format!("backend returned {status}: {e}"))?; - let Some(body) = body else { return Ok(cancelled_response()); }; - let text = String::from_utf8_lossy(&body); - let err = format!("backend returned {status}: {}", describe_backend(&text)); - if attempt < MAX_ATTEMPTS && is_retryable_status(code) { - backoff_sleep(attempt).await; - if cancelled.is_cancelled() { - return Ok(CompletionResult { - content: String::new(), - tool_calls: Vec::new(), - finish_reason: "cancelled".into(), - usage: None, - }); + let is_json = resp + .headers() + .get("content-type") + .and_then(|v| v.to_str().ok()) + .is_some_and(|ct| ct.contains("application/json")); + + // Some servers ignore `stream` and return a complete JSON body. + if is_json { + let Some(body) = read_body(resp, &cancelled).await? else { return Ok(cancelled_response()); }; + let v: Value = serde_json::from_slice(&body).map_err(|e| format!("bad JSON response: {e}"))?; + return parse_complete_response(&v, &mut on_delta); + } + + let (reply, interrupted) = read_sse(resp, &cancelled, &mut on_delta).await; + match interrupted { + Some(reason) if reply.content.is_empty() && attempt < MAX_ATTEMPTS => { + if !retry_after(attempt, &reason, &cancelled, &mut on_delta).await { + return Ok(cancelled_response()); + } } - continue; + _ => return Ok(reply), } - return Err(err); - }; - let is_json = resp - .headers() - .get("content-type") - .and_then(|v| v.to_str().ok()) - .is_some_and(|ct| ct.contains("application/json")); - - // Some servers ignore `stream` and return a complete JSON body. - if is_json { - let Some(body) = read_body(resp, &cancelled).await? else { return Ok(cancelled_response()); }; - let v: Value = serde_json::from_slice(&body).map_err(|e| format!("bad JSON response: {e}"))?; - return parse_complete_response(&v, &mut on_delta); } + } +} - let mut content = String::new(); - let mut partials: Vec = Vec::new(); - let mut finish_reason = String::from("stop"); - // Did the server ever say it was done (a `[DONE]` line or a - // finish_reason chunk)? Without one, the stream ending means the - // connection dropped mid-answer, not that the model finished. - let mut saw_terminator = false; - let mut usage: Option = None; - // Byte buffer: chunks can split multi-byte UTF-8 sequences, so text - // conversion only happens on complete lines ('\n' is never part of a - // multi-byte sequence). - let mut buf: Vec = Vec::new(); - let mut stream = resp.bytes_stream(); - let mut received = 0usize; - let mut scanned = 0usize; - - 'outer: loop { - let next = tokio::select! { - c = stream.next() => c, - _ = cancelled.cancelled() => { - finish_reason = "cancelled".into(); - break; - } - }; - let Some(chunk) = next else { break }; - let chunk = match chunk { - Ok(chunk) if chunk.len() <= MAX_RESPONSE_BYTES - received => chunk, - _ => { - saw_terminator = false; - finish_reason = TRUNCATED.into(); - break; - } - }; - received += chunk.len(); - buf.extend_from_slice(&chunk); - - while let Some(rel) = buf[scanned..].iter().position(|&b| b == b'\n') { - let pos = scanned + rel; - scanned = 0; - let rest = buf.split_off(pos + 1); - let consumed = std::mem::replace(&mut buf, rest); - let line = trim_bytes(&consumed[..pos]); - if line.is_empty() || line.first() == Some(&b':') { - continue; +/// Parse one SSE stream to its end. The second value names the interruption +/// when the network ended the stream before the server did: EOF with no +/// terminator, or a read error. It is `None` for a finished reply and for a +/// stream this client cut itself, which a fresh attempt would cut the same way. +async fn read_sse( + resp: reqwest::Response, + cancelled: &crate::state::CancelToken, + on_delta: &mut impl FnMut(StreamDelta), +) -> (CompletionResult, Option) { + let mut content = String::new(); + let mut partials: Vec = Vec::new(); + let mut finish_reason = String::from("stop"); + // Did the server ever say it was done (a `[DONE]` line or a + // finish_reason chunk)? Without one, the stream ending means the + // connection dropped mid-answer, not that the model finished. + let mut saw_terminator = false; + let mut usage: Option = None; + // Byte buffer: chunks can split multi-byte UTF-8 sequences, so text + // conversion only happens on complete lines ('\n' is never part of a + // multi-byte sequence). + let mut buf: Vec = Vec::new(); + let mut stream = resp.bytes_stream(); + let mut received = 0usize; + let mut scanned = 0usize; + // Set when the network, not the server or this client, ended the + // stream: the one truncation a fresh attempt can undo. + let mut interrupted: Option = None; + + 'outer: loop { + let next = tokio::select! { + c = stream.next() => c, + _ = cancelled.cancelled() => { + finish_reason = "cancelled".into(); + break; + } + }; + let Some(chunk) = next else { + if !saw_terminator { + interrupted = Some("the stream ended before the reply finished".into()); + } + break; + }; + let chunk = match chunk { + Ok(chunk) if chunk.len() <= MAX_RESPONSE_BYTES - received => chunk, + Ok(_) => { + saw_terminator = false; + finish_reason = TRUNCATED.into(); + break; + } + Err(e) => { + saw_terminator = false; + finish_reason = TRUNCATED.into(); + // A body that fails to decode will fail the same way again; + // only the connection is worth a fresh attempt. + if !e.is_decode() { + interrupted = Some(describe_transport(&e)); } - let data = strip_data_prefix(line); - if data == b"[DONE]" { - saw_terminator = true; - break 'outer; + break; + } + }; + received += chunk.len(); + buf.extend_from_slice(&chunk); + + while let Some(rel) = buf[scanned..].iter().position(|&b| b == b'\n') { + let pos = scanned + rel; + scanned = 0; + let rest = buf.split_off(pos + 1); + let consumed = std::mem::replace(&mut buf, rest); + let line = trim_bytes(&consumed[..pos]); + if line.is_empty() || line.first() == Some(&b':') { + continue; + } + let data = strip_data_prefix(line); + if data == b"[DONE]" { + saw_terminator = true; + break 'outer; + } + let Ok(chunk) = serde_json::from_slice::(data) else { continue }; + if let Some(u) = chunk.usage { + usage = Some(u.into_usage()); + } + // The usage-bearing final chunk has an empty choices array. + let Some(choice) = chunk.choices.into_iter().next() else { continue }; + + if let Some(reason) = choice.finish_reason { + // Servers that end a stream here and never send `[DONE]` + // are still finished: this is a terminator too. + saw_terminator = true; + finish_reason = reason; + } + let delta = choice.delta; + if let Some(text) = delta.content { + if !text.is_empty() { + content.push_str(&text); + on_delta(StreamDelta::Content(text)); } - let Ok(chunk) = serde_json::from_slice::(data) else { continue }; - if let Some(u) = chunk.usage { - usage = Some(u.into_usage()); + } + // Reasoning models surface thinking under different keys. + if let Some(text) = delta.reasoning_content { + if !text.is_empty() { + on_delta(StreamDelta::Reasoning(text)); } - // The usage-bearing final chunk has an empty choices array. - let Some(choice) = chunk.choices.into_iter().next() else { continue }; - - if let Some(reason) = choice.finish_reason { - // Servers that end a stream here and never send `[DONE]` - // are still finished: this is a terminator too. - saw_terminator = true; - finish_reason = reason; + } else if let Some(text) = delta.reasoning { + if !text.is_empty() { + on_delta(StreamDelta::Reasoning(text)); } - let delta = choice.delta; - if let Some(text) = delta.content { - if !text.is_empty() { - content.push_str(&text); - on_delta(StreamDelta::Content(text)); + } + if let Some(calls) = delta.tool_calls { + for tc in calls { + let idx = tc.index.unwrap_or(0) as usize; + if idx >= MAX_TOOL_CALLS { + saw_terminator = false; + finish_reason = TRUNCATED.into(); + break 'outer; } - } - // Reasoning models surface thinking under different keys. - if let Some(text) = delta.reasoning_content { - if !text.is_empty() { - on_delta(StreamDelta::Reasoning(text)); + while partials.len() <= idx { + partials.push(PartialToolCall::default()); } - } else if let Some(text) = delta.reasoning { - if !text.is_empty() { - on_delta(StreamDelta::Reasoning(text)); + if let Some(id) = tc.id { + partials[idx].id.push_str(&id); } - } - if let Some(calls) = delta.tool_calls { - for tc in calls { - let idx = tc.index.unwrap_or(0) as usize; - if idx >= MAX_TOOL_CALLS { - saw_terminator = false; - finish_reason = TRUNCATED.into(); - break 'outer; - } - while partials.len() <= idx { - partials.push(PartialToolCall::default()); - } - if let Some(id) = tc.id { - partials[idx].id.push_str(&id); + if let Some(function) = tc.function { + if let Some(name) = function.name { + partials[idx].name.push_str(&name); } - if let Some(function) = tc.function { - if let Some(name) = function.name { - partials[idx].name.push_str(&name); - } - if let Some(args) = function.arguments { - partials[idx].arguments.push_str(&args); - } + if let Some(args) = function.arguments { + partials[idx].arguments.push_str(&args); } } } } - scanned = buf.len(); } + scanned = buf.len(); + } - // The stream ran out without the server ever finishing it: report the - // truncation rather than the default "stop", which would make a - // cut-off answer indistinguishable from a complete one. Cancellation - // ends the stream from this side, so it keeps its own reason. - if !saw_terminator && finish_reason != "cancelled" { - finish_reason = TRUNCATED.into(); - } - let tool_calls = finalize_tool_calls(partials); - if !tool_calls.is_empty() && finish_reason == "stop" { - finish_reason = "tool_calls".into(); - } - Ok(CompletionResult { content, tool_calls, finish_reason, usage }) + // The stream ran out without the server ever finishing it: report the + // truncation rather than the default "stop", which would make a + // cut-off answer indistinguishable from a complete one. Cancellation + // ends the stream from this side, so it keeps its own reason. + if !saw_terminator && finish_reason != "cancelled" { + finish_reason = TRUNCATED.into(); } + let tool_calls = finalize_tool_calls(partials); + if !tool_calls.is_empty() && finish_reason == "stop" { + finish_reason = "tool_calls".into(); + } + (CompletionResult { content, tool_calls, finish_reason, usage }, interrupted) } fn cancelled_response() -> CompletionResult { @@ -698,6 +732,31 @@ pub fn truncate(s: &str, max: usize) -> String { } } +/// Attempts per request. With [`BACKOFF_UNIT`] doubling up to +/// [`BACKOFF_CAP`], the last attempt goes out about a minute after the +/// first. A transient fault on the path to an endpoint (a reset or a TLS +/// alert on send, a stream cut mid-reply) can recur for minutes and clears +/// within seconds most times and within a minute at worst, so attempts +/// packed into one second only ever observe the fault and end a turn the +/// next connection would have completed. Every wait is announced and +/// cancellable. +const MAX_ATTEMPTS: u32 = 8; +/// Attempts in a row that may fail before any connection exists (refused, +/// unresolvable, a connect timeout, a failed TLS handshake) before the +/// request gives up early. That is an address with nothing listening far +/// more often than a blip, and a user who forgot to start a local server +/// should hear so after the two short waits, not a minute. A refused +/// address fails in about three seconds; one that drops packets pays the +/// connect timeout on each attempt as well. +const UNREACHED_ATTEMPTS: u32 = 3; + +#[cfg(not(test))] +const BACKOFF_UNIT: std::time::Duration = std::time::Duration::from_secs(1); +/// Tests exercise the attempt count, never the wall clock. +#[cfg(test)] +const BACKOFF_UNIT: std::time::Duration = std::time::Duration::from_millis(1); +const BACKOFF_CAP: std::time::Duration = std::time::Duration::from_secs(16); + /// Only statuses that mean the server rejected the request before doing work. /// 5xx from a proxy is not safe to retry without an idempotency key. fn is_retryable_status(code: u16) -> bool { @@ -708,10 +767,22 @@ fn is_transient_transport(err: &reqwest::Error) -> bool { err.is_connect() || err.is_timeout() || err.is_request() } -async fn backoff_sleep(attempt: u32) { - // ~100ms, ~200ms (capped); keep total retry budget small for local UX. - let ms = 100u64.saturating_mul(1u64 << (attempt.saturating_sub(1).min(2))); - tokio::time::sleep(std::time::Duration::from_millis(ms.min(400))).await; +/// Announce the retry, then wait out the backoff for `attempt` (1-based, the +/// one that just failed). Returns false when cancelled during the wait: a +/// backoff can reach [`BACKOFF_CAP`], and a user who cancels must not sit +/// through it. +async fn retry_after( + attempt: u32, + reason: &str, + cancelled: &crate::state::CancelToken, + on_delta: &mut impl FnMut(StreamDelta), +) -> bool { + on_delta(StreamDelta::Retry { attempt: attempt + 1, max_attempts: MAX_ATTEMPTS, reason: reason.to_string() }); + let wait = BACKOFF_UNIT.saturating_mul(1u32 << attempt.saturating_sub(1).min(8)).min(BACKOFF_CAP); + tokio::select! { + _ = tokio::time::sleep(wait) => true, + _ = cancelled.cancelled() => false, + } } #[cfg(test)] @@ -917,21 +988,180 @@ mod tests { assert_eq!(result.finish_reason, "stop"); } - /// The client reports what arrived without hiding it: the calls come back - /// alongside the truncation, and refusing to dispatch them is the agent - /// loop's decision (a call from a stream the model never finished may not - /// be the call it meant to make). + /// An endpoint that answers successive connections with successive + /// bodies, each close-delimited, then stops listening. An empty body + /// closes the connection after reading the request without answering: + /// the shape of a transport fault on send. Returns the URL and the + /// number of connections it served. + fn spawn_sse_sequence(bodies: Vec) -> (String, Arc>) { + use std::io::{Read as _, Write as _}; + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let addr = listener.local_addr().unwrap(); + let served = Arc::new(std::sync::Mutex::new(0usize)); + let count = served.clone(); + std::thread::spawn(move || { + for sse in bodies { + let Ok((mut stream, _)) = listener.accept() else { return }; + let mut buf = Vec::new(); + let mut byte = [0u8; 1]; + while !buf.ends_with(b"\r\n\r\n") { + match stream.read(&mut byte) { + Ok(1) => buf.push(byte[0]), + _ => return, + } + } + let headers = String::from_utf8_lossy(&buf).to_string(); + let content_length: usize = headers + .lines() + .find_map(|l| { + let (k, v) = l.split_once(':')?; + k.eq_ignore_ascii_case("content-length").then(|| v.trim().parse().ok())? + }) + .unwrap_or(0); + let mut body = vec![0u8; content_length]; + if content_length > 0 && stream.read_exact(&mut body).is_err() { + return; + } + *count.lock().unwrap() += 1; + if sse.is_empty() { + continue; + } + let _ = stream.write_all( + format!("HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nConnection: close\r\n\r\n{sse}").as_bytes(), + ); + } + }); + (format!("http://{addr}/v1"), served) + } + + async fn stream_sequence(bodies: Vec) -> (CompletionResult, Vec, usize) { + let (url, served) = spawn_sse_sequence(bodies); + let mut deltas: Vec = Vec::new(); + let result = ChatClient::new(url, None, "m".into(), 0.0, 64) + .stream_chat(&[ChatMessage::user("hi")], "[]", Arc::new(crate::state::CancelToken::default()), |d| { + deltas.push(match d { + StreamDelta::Content(t) => format!("content:{t}"), + StreamDelta::Reasoning(t) => format!("reasoning:{t}"), + StreamDelta::Retry { attempt, max_attempts, reason } => format!("retry:{attempt}/{max_attempts}:{reason}"), + }) + }) + .await + .unwrap(); + let served = *served.lock().unwrap(); + (result, deltas, served) + } + + /// The bug this guards: a connection that died while the model was still + /// reasoning ended the whole turn with the task half done, on a fault the + /// next connection did not see. Nothing the caller keeps had arrived, so + /// the client starts the reply over, announces the retry after the + /// reasoning it voids, and the result is the finished reply alone. #[tokio::test] - async fn a_truncated_stream_still_returns_the_tool_calls_it_carried() { - let result = stream_once( - "data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":0,\"id\":\"c1\",\"function\":{\"name\":\"read_file\",\"arguments\":\"{\\\"path\\\":\\\"a.rs\\\"}\"}}]}}]}\n\n", - ) + async fn a_stream_interrupted_before_any_reply_text_is_started_over() { + let (result, deltas, served) = stream_sequence(vec![ + "data: {\"choices\":[{\"delta\":{\"reasoning_content\":\"let me\"},\"finish_reason\":null}]}\n\n".into(), + concat!( + "data: {\"choices\":[{\"delta\":{\"content\":\"all of it\"},\"finish_reason\":null}]}\n\n", + "data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n", + "data: [DONE]\n\n", + ) + .into(), + ]) .await; + assert_eq!(result.finish_reason, "stop"); + assert_eq!(result.content, "all of it"); + assert_eq!(served, 2); + assert_eq!( + deltas, + vec![ + "reasoning:let me".to_string(), + format!("retry:2/{MAX_ATTEMPTS}:the stream ended before the reply finished"), + "content:all of it".to_string(), + ] + ); + } + + /// A connection the far side drops before answering is a transport fault + /// after the connection existed: each one is resent after a backoff, + /// announced, and the turn goes on with the reply that arrives. + #[tokio::test] + async fn a_request_dropped_on_send_is_resent_until_a_reply_arrives() { + let good = concat!( + "data: {\"choices\":[{\"delta\":{\"content\":\"all of it\"},\"finish_reason\":\"stop\"}]}\n\n", + "data: [DONE]\n\n", + ); + let (result, deltas, served) = + stream_sequence(vec![String::new(), String::new(), String::new(), good.into()]).await; + assert_eq!(served, 4); + assert_eq!(result.finish_reason, "stop"); + assert_eq!(result.content, "all of it"); + assert_eq!(deltas.len(), 4); + for (i, delta) in deltas.iter().take(3).enumerate() { + assert!( + delta.starts_with(&format!("retry:{}/{MAX_ATTEMPTS}:request failed: ", i + 2)), + "resend {} is announced with its cause: {delta}", + i + 2 + ); + } + assert_eq!(deltas[3], "content:all of it"); + } + + /// Nothing listening is not a fault worth a minute of backoff: after + /// [`UNREACHED_ATTEMPTS`] in a row the request gives up early. The port + /// was bound and released by this test, so nothing answers on it; the + /// deadline keeps a stray listener from turning that into a hang. + #[tokio::test] + async fn an_address_that_never_answers_gives_up_before_the_budget() { + let url = { + let released = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + format!("http://{}/v1", released.local_addr().unwrap()) + }; + let mut deltas: Vec = Vec::new(); + let client = ChatClient::new(url, None, "m".into(), 0.0, 64); + let messages = [ChatMessage::user("hi")]; + let request = client.stream_chat(&messages, "[]", Arc::new(crate::state::CancelToken::default()), |d| { + if let StreamDelta::Retry { attempt, max_attempts, .. } = d { + deltas.push(format!("{attempt}/{max_attempts}")); + } + }); + let result = tokio::time::timeout(std::time::Duration::from_secs(10), request) + .await + .expect("a released port answers with a refusal, not silence"); + let Err(err) = result else { panic!("no listener is a failed request") }; + assert!(err.starts_with("request failed: "), "{err}"); + assert_eq!(deltas, vec![format!("2/{MAX_ATTEMPTS}"), format!("3/{MAX_ATTEMPTS}")]); + } + + /// The client reports what arrived without hiding it: once the retry + /// budget is spent, the calls come back alongside the truncation, and + /// refusing to dispatch them is the agent loop's decision (a call from a + /// stream the model never finished may not be the call it meant to make). + #[tokio::test] + async fn a_truncated_stream_still_returns_the_tool_calls_it_carried() { + let body = "data: {\"choices\":[{\"delta\":{\"tool_calls\":[{\"index\":0,\"id\":\"c1\",\"function\":{\"name\":\"read_file\",\"arguments\":\"{\\\"path\\\":\\\"a.rs\\\"}\"}}]}}]}\n\n"; + let (result, deltas, served) = stream_sequence(vec![body.to_string(); MAX_ATTEMPTS as usize]).await; + assert_eq!(served, MAX_ATTEMPTS as usize, "every attempt in the budget was spent first"); + assert_eq!(deltas.len(), MAX_ATTEMPTS as usize - 1, "each resend was announced"); assert_eq!(result.tool_calls.len(), 1); assert_eq!(result.tool_calls[0].function.name, "read_file"); assert_eq!(result.finish_reason, TRUNCATED, "an unfinished stream is not a clean tool_calls stop"); } + /// Reply text the caller has already shown is never duplicated: a stream + /// that dies after content streamed is a truncation, not a retry. + #[tokio::test] + async fn a_stream_interrupted_after_reply_text_is_not_retried() { + let (result, deltas, served) = stream_sequence(vec![ + "data: {\"choices\":[{\"delta\":{\"content\":\"half an ans\"},\"finish_reason\":null}]}\n\n".into(), + "data: {\"choices\":[{\"delta\":{\"content\":\"unreached\"},\"finish_reason\":\"stop\"}]}\n\n".into(), + ]) + .await; + assert_eq!(served, 1); + assert_eq!(result.finish_reason, TRUNCATED); + assert_eq!(result.content, "half an ans"); + assert_eq!(deltas, vec!["content:half an ans".to_string()]); + } + #[test] fn retryable_status_codes() { assert!(is_retryable_status(429)); diff --git a/crates/core/src/spec.rs b/crates/core/src/spec.rs index 3ba0bd38..5e04c3c2 100644 --- a/crates/core/src/spec.rs +++ b/crates/core/src/spec.rs @@ -849,14 +849,14 @@ spending a read on the address - lexical ranking cannot separate those two, because they are about the same thing in the same words. "#; -const STDIO: &str = r#"# stdio protocol (openmax-stdio/5) +const STDIO: &str = r#"# stdio protocol (openmax-stdio/6) `openmax --stdio` speaks line-delimited JSON both ways: commands on stdin, `AgentEvent` envelopes on stdout. This is the stable contract for custom frontends, editor integrations, and one openmax driving another. Handshake: the first stdout line is -{"type":"hello","proto":"openmax-stdio/5","protocol_version":5,"session_id":"...","version":"...","project":"/abs/path","continued":false}. +{"type":"hello","proto":"openmax-stdio/6","protocol_version":6,"session_id":"...","version":"...","project":"/abs/path","continued":false}. `protocol_version` is compared as an integer; any wire change bumps it. Commands, one JSON object per line: @@ -898,7 +898,10 @@ the frozen tool schemas sent on every request, context_tokens), refreeze receipt, or a policy/providers/settings/approval notice - surfaced here so a frontend can render what the model sees; `call_id` links it to the tool result it rode, or is empty for a note inserted before the next prompt -like a turn-start receipt), `diff` (call_id, +like a turn-start receipt), `retry` (attempt, max_attempts, reason: the +model request is being resent after a transport failure, a 429, or a stream +that died before any reply text; thinking already streamed for that attempt +is void), `diff` (call_id, path, diff, added, removed), `approval_request` (approval_id, name, summary, detail, reason, source_path, source_sha, and an optional `env`), `approval_settled` (approval_id, outcome), `refrozen` (tools, skills, changes: the refreeze receipt naming @@ -946,9 +949,10 @@ only guaranteed terminator. A command that starts no turn (empty text, an untrusted project) still gets one, with stop_reason `refused`, after the `protocol_error` that says why. A turn that dies unexpectedly reports `error` and then `done` with stop_reason `error`; a provider stream that ends -mid-answer with no completion signal reports its partial `message_done`, then -`error`, then `done` with stop_reason `truncated`, and no tool call it carried -is run. The single exception is a +mid-answer with no completion signal is resent (each resend announced by a +`retry`) while no reply text has streamed, and otherwise reports its partial +`message_done`, then `error`, then `done` with stop_reason `truncated`, and +no tool call it carried is run. The single exception is a `user` sent while a turn is in flight: that is refused with a `protocol_error` and no `done`, because the running turn owns the next one. @@ -1338,6 +1342,7 @@ mod tests { }, AgentEvent::ToolEnd { call_id: String::new(), ok: true, output: String::new() }, AgentEvent::HarnessNote { call_id: String::new(), text: String::new() }, + AgentEvent::Retry { attempt: 0, max_attempts: 0, reason: String::new() }, AgentEvent::Diff { call_id: String::new(), path: String::new(), diff --git a/crates/core/src/types.rs b/crates/core/src/types.rs index 3416be3a..c9d99def 100644 --- a/crates/core/src/types.rs +++ b/crates/core/src/types.rs @@ -108,6 +108,13 @@ pub enum AgentEvent { /// links it to the tool result it rides when it rode one, else empty (a /// note inserted before the next prompt, e.g. a turn-start receipt). HarnessNote { call_id: String, text: String }, + /// The model request is being resent: `attempt` of `max_attempts` goes + /// out after a backoff, `reason` having ended the previous one (a + /// transport failure, a 429, or a stream that died before any reply + /// text). Emitted before the wait; a cancel during it ends the turn and + /// the attempt never goes out. Thinking streamed for the failed attempt + /// is void; no Token preceded it. + Retry { attempt: u32, max_attempts: u32, reason: String }, Diff { call_id: String, path: String, diff: String, added: usize, removed: usize }, /// Mutating tool waiting on the user. `detail` is a short args preview /// (paths, command head) for the TUI card; may be empty. @@ -138,7 +145,7 @@ pub enum AgentEvent { /// as structured data, not folded into `detail`, so the frontend /// controls its own un-clippable placement. Additive and defaulted: /// a stream without the key deserializes unchanged, and the key is - /// omitted from the wire whenever it is empty, so `openmax-stdio/5` + /// omitted from the wire whenever it is empty, so `openmax-stdio/6` /// bytes are byte-identical for every call that grants no env. #[serde(default, skip_serializing_if = "Vec::is_empty")] env: Vec, @@ -220,7 +227,7 @@ mod tests { /// Golden wire format for every `AgentEvent`, wrapped in its envelope /// exactly as `--stdio` and `--print --json` emit it. These strings are - /// the `openmax-stdio/5` contract: session_id first, then the `type` + /// the `openmax-stdio/6` contract: session_id first, then the `type` /// discriminator, then variant fields in declaration order. A change here /// is a protocol break and must bump `PROTO_VERSION`. /// `every_agent_event_variant_is_pinned_here` fails to compile if a variant @@ -345,6 +352,10 @@ mod tests { env(AgentEvent::SchemasOverBudget { schema_tokens: 6800, budget_tokens: 2150 }), r#"{"session_id":"s1","type":"schemas_over_budget","schema_tokens":6800,"budget_tokens":2150}"# ); + assert_eq!( + env(AgentEvent::Retry { attempt: 2, max_attempts: 8, reason: "request failed: connection reset".into() }), + r#"{"session_id":"s1","type":"retry","attempt":2,"max_attempts":8,"reason":"request failed: connection reset"}"# + ); assert_eq!( env(AgentEvent::HookFailed { @@ -428,6 +439,7 @@ mod tests { AgentEvent::ToolStart { .. } => {} AgentEvent::ToolEnd { .. } => {} AgentEvent::HarnessNote { .. } => {} + AgentEvent::Retry { .. } => {} AgentEvent::Diff { .. } => {} AgentEvent::ApprovalRequest { .. } => {} AgentEvent::ApprovalSettled { .. } => {} diff --git a/crates/tui/src/app.rs b/crates/tui/src/app.rs index 7a2b97b6..fe5730f6 100644 --- a/crates/tui/src/app.rs +++ b/crates/tui/src/app.rs @@ -2747,6 +2747,17 @@ impl App { )); } } + AgentEvent::Retry { attempt, max_attempts, reason } => { + // The reasoning shown so far belongs to the attempt that + // failed; the fresh attempt starts its own. The reason can + // carry a backend body; the note is one line. + self.thinking_tail.clear(); + self.thinking_source.clear(); + self.thinking_wrapped.clear(); + self.thinking_chars = 0; + self.note(&format!("{}; retrying ({attempt} of {max_attempts})", open_max_core::text::one_line(&reason))); + self.dirty.mark_tail(); + } AgentEvent::SchemasOverBudget { schema_tokens, budget_tokens } => { // Says what it costs and what to do, not how compaction reacts: // that depends on whether any room is left at all. @@ -5499,6 +5510,29 @@ mod tests { assert_eq!(home_shortened("/srv/data", Some("/")), "/srv/data"); } + /// The reasoning tail shown during a turn belongs to one attempt. A + /// retry starts the reply over, so what the failed attempt streamed is + /// dropped before the fresh attempt's reasoning arrives; otherwise the + /// two would read as one thought. + #[test] + fn a_retry_drops_the_failed_attempts_reasoning_tail() { + let (mut app, dir) = app_fixture(); + app.running = true; + app.turn_started = Some(std::time::Instant::now()); + app.on_agent_event(AgentEvent::Thinking { text: "abandoned line".into() }); + assert_eq!(app.thinking_tail, "abandoned line"); + app.on_agent_event(AgentEvent::Retry { + attempt: 2, + max_attempts: 8, + reason: "the stream ended before the reply finished".into(), + }); + assert!(app.thinking_tail.is_empty()); + assert_eq!(app.thinking_chars, 0); + app.on_agent_event(AgentEvent::Thinking { text: "fresh line".into() }); + assert_eq!(app.thinking_tail, "fresh line"); + fs::remove_dir_all(dir).unwrap(); + } + #[test] fn streaming_output_grows_above_the_fixed_prompt() { let (mut app, dir) = app_fixture(); diff --git a/crates/tui/src/headless.rs b/crates/tui/src/headless.rs index f9398f9e..06985997 100644 --- a/crates/tui/src/headless.rs +++ b/crates/tui/src/headless.rs @@ -264,6 +264,11 @@ async fn run_turn_events( ); } } + AgentEvent::Retry { attempt, max_attempts, reason } => { + if !json { + let _ = writeln!(stderr, "openmax: {}; retrying ({attempt} of {max_attempts})", one_line(reason)); + } + } AgentEvent::SchemasOverBudget { schema_tokens, budget_tokens } => { // Advisory: the turn still runs, so the exit code is untouched. if !json { diff --git a/crates/tui/src/stdio.rs b/crates/tui/src/stdio.rs index cd34b477..b1209568 100644 --- a/crates/tui/src/stdio.rs +++ b/crates/tui/src/stdio.rs @@ -4,7 +4,7 @@ //! speak line-delimited JSON (an editor plugin, an orchestrator, another //! openmax) can drive a complete interactive session, approvals included. //! -//! Protocol (`openmax-stdio/5`), one JSON object per line. The normative +//! Protocol (`openmax-stdio/6`), one JSON object per line. The normative //! reference (every field of every line) is `docs/stdio-protocol.md`, kept in //! step with `openmax --spec stdio`; `crates/core/src/types.rs` golden tests //! pin the event wire. @@ -18,7 +18,7 @@ //! {"cmd":"quit"} finish the turn, then exit //! //! stdout lines: -//! {"type":"hello","proto":"openmax-stdio/5","protocol_version":5,"session_id":"...","version":"...","project":"..."} +//! {"type":"hello","proto":"openmax-stdio/6","protocol_version":6,"session_id":"...","version":"...","project":"..."} //! AgentEvent envelopes exactly as `--print --json` emits them //! {"type":"protocol_error","message":"..."} bad input; session unharmed //! @@ -46,11 +46,11 @@ use open_max_core::types::{AgentEvent, AgentEventEnvelope}; use serde::Deserialize; use tokio::sync::mpsc; -pub const PROTO: &str = "openmax-stdio/5"; +pub const PROTO: &str = "openmax-stdio/6"; /// Machine-comparable protocol major. A client negotiates on this integer; /// `PROTO` embeds the same number as a human-readable id (checked in tests). /// Bump on any wire change (event field, command shape, framing line). -pub const PROTO_VERSION: u32 = 5; +pub const PROTO_VERSION: u32 = 6; // Unknown `cmd` values are protocol errors; extra fields on a known command // are ignored (lenient by design, so clients can annotate lines freely). @@ -465,7 +465,7 @@ fn transcript_value( }) } -/// Validate one JSONL line against the `openmax-stdio/5` contract using the +/// Validate one JSONL line against the `openmax-stdio/6` contract using the /// authoritative types (`Command` for stdin, `AgentEvent` for stdout events), /// so there is no second schema to drift. Returns a short label on success /// (`cmd user`, `event token`, `hello`) or a human reason on failure. diff --git a/crates/tui/tests/cli.rs b/crates/tui/tests/cli.rs index 9734ffd8..c1341a8d 100644 --- a/crates/tui/tests/cli.rs +++ b/crates/tui/tests/cli.rs @@ -602,8 +602,8 @@ fn stdio_handshake_speaks_the_contract() { reader.read_line(&mut hello).unwrap(); let hello: serde_json::Value = serde_json::from_str(&hello).unwrap(); assert_eq!(hello["type"], "hello"); - assert_eq!(hello["proto"], "openmax-stdio/5"); - assert_eq!(hello["protocol_version"], 5); + assert_eq!(hello["proto"], "openmax-stdio/6"); + assert_eq!(hello["protocol_version"], 6); assert!(hello["session_id"].is_string()); writeln!(stdin, r#"{{"cmd":"quit"}}"#).unwrap(); @@ -1084,11 +1084,14 @@ fn a_truncated_stream_reports_truncation_instead_of_a_clean_stop() { /// call, so the arguments parse and nothing looks broken. A stream with no /// completion signal is not a response the model asked to act on (more calls /// may have been coming, or this one may still have been under revision), so -/// the call must not run. +/// the call must not run. Reply text streams first, so this is the +/// interruption the client does not start over. #[test] fn a_truncated_stream_never_runs_the_tool_call_it_carried() { let (project, home) = fresh_dirs("truncated-native-call"); - let (base_url, _requests, _server) = spawn_scripted_server(vec![(WRITE_CALL_SSE.to_string(), false)]); + let prose = "data: {\"choices\":[{\"delta\":{\"content\":\"writing it\"},\"finish_reason\":null}]}\n\n"; + let (base_url, _requests, _server) = + spawn_scripted_server(vec![(format!("{prose}{WRITE_CALL_SSE}"), false)]); // auto, so a refusal here is the truncation and not the approval gate. write_settings_with_mode(&home, &base_url, "auto"); diff --git a/docs/stdio-protocol.md b/docs/stdio-protocol.md index 1e1daec4..90fb93b6 100644 --- a/docs/stdio-protocol.md +++ b/docs/stdio-protocol.md @@ -1,4 +1,4 @@ -# stdio protocol (`openmax-stdio/5`) +# stdio protocol (`openmax-stdio/6`) `openmax --stdio` speaks line-delimited JSON both ways, so any process that reads and writes JSONL (an editor plugin, an orchestrator, another openmax) can @@ -15,7 +15,7 @@ This file is the normative reference for every field of every line. The first stdout line is: ```json -{"type":"hello","proto":"openmax-stdio/5","protocol_version":5,"session_id":"...","version":"0.2.0","project":"/abs/path","continued":false} +{"type":"hello","proto":"openmax-stdio/6","protocol_version":6,"session_id":"...","version":"0.2.0","project":"/abs/path","continued":false} ``` `protocol_version` is an integer a client compares directly; `proto` carries @@ -72,6 +72,7 @@ one). | `tool_start` | `call_id`, `name`, `args` (object) | | `tool_end` | `call_id`, `ok` (bool), `output` | | `harness_note` | `call_id`, `text` (a note the harness wrote into the model's transcript: a refreeze receipt, or a policy, providers, settings, or approval notice. `call_id` links it to the tool result it rode; it is empty for a note inserted before the next prompt, such as a turn-start receipt) | +| `retry` | `attempt`, `max_attempts`, `reason` (the model request is being resent: `attempt` is the one that goes out of the `max_attempts` budget after a backoff wait, `reason` having ended the previous one with a transport failure, a 429, or a stream that died before any reply text arrived. Emitted before the wait; a `cancel` during it ends the turn and the attempt never goes out. `thinking` lines already emitted for that attempt are void and the reply starts over; no `token` line precedes a retried stream. A request gives up before the budget when three attempts in a row never reached the endpoint) | | `diff` | `call_id`, `path`, `diff`, `added`, `removed` | | `approval_request` | `approval_id`, `name`, `summary`, `detail`, `reason` (`gate`, or `unapproved_source` which unattended clients must never auto-approve), `source_path`, `source_sha`, and optional `env` (see below) | | `approval_settled` | `approval_id`, `outcome` (`approved`, `declined`, `timed_out`, or `cancelled`) | @@ -137,7 +138,7 @@ followed by `done` with `stop_reason` `refused`, so a client that blocks on | `stop_reason` | Meaning | | --- | --- | | provider `finish_reason` | Passed through verbatim on a normal turn, commonly `stop` or `length`. Treat any unlisted value as a normal end | -| `truncated` | The provider stream ended with no completion signal; the reply is incomplete, any tool calls it carried were refused, and an `error` line precedes it | +| `truncated` | The provider stream ended with no completion signal (after reply text had streamed, or on the last retry) or exceeded a client limit; the reply is incomplete, any tool calls it carried were refused, and an `error` line precedes it | | `max_iterations` | The turn hit the tool-call ceiling | | `budget_exhausted` | The per-turn `max_agent_tokens` cap refused the next request at admission; nothing was sent, and resubmitting continues the work | | `unverified` | A blocking `turn_end` hook refused the completion more times than the harness honors (8), or its refusal could not be persisted; the reply stands unverified |