From 73615b9219cadbf1e355a7b128ffb8f9823d3f11 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" Date: Fri, 25 Sep 2026 04:20:13 +0000 Subject: [PATCH] chore: sync public mirror from internal --- .repository-projection.json | 6 +- README.md | 2 +- .../examples/semantic_stream_fixture.rs | 87 +++ .../src/headless/generated_protocol.rs | 1 + .../local-host-rs/src/headless/messages.rs | 4 + .../src/headless/messages/state.rs | 1 + packages/local-host-rs/src/headless/proto.rs | 6 + packages/local-host-rs/src/headless_server.rs | 118 ++-- .../src/headless_server/semantic_stream.rs | 231 ++++++++ packages/local-host-rs/src/lib.rs | 1 + packages/local-host-rs/src/semantic_text.rs | 517 ++++++++++++++++++ packages/local-host-rs/src/ui_wake.rs | 4 +- .../fixtures/headless-protocol-v1.json | 1 + packages/runtime-contracts-rs/src/protocol.rs | 5 +- packages/tui-rs/src/app.rs | 16 + packages/tui-rs/src/app/startup_update.rs | 258 +++++++++ packages/tui-rs/src/entrypoint.rs | 4 + packages/tui-rs/src/markdown.rs | 11 +- packages/tui-rs/src/themes/mod.rs | 19 +- packages/tui-rs/src/update_cli.rs | 328 +++++++++-- packages/tui-rs/tests/pty_e2e.rs | 327 ++++++++++- proto/maestro/v1/headless.proto | 5 + .../v1/protocol-compatibility-manifest.json | 9 +- 23 files changed, 1833 insertions(+), 128 deletions(-) create mode 100644 packages/local-host-rs/examples/semantic_stream_fixture.rs create mode 100644 packages/local-host-rs/src/headless_server/semantic_stream.rs create mode 100644 packages/local-host-rs/src/semantic_text.rs create mode 100644 packages/tui-rs/src/app/startup_update.rs diff --git a/.repository-projection.json b/.repository-projection.json index 0fc76727d..5d38bc6ad 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "58427e18485e66532d24ee186dc4c1aab3169a0b", + "sourceSha": "29f82f72c4a2fb5964d303282bb74638f134ec0c", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "6936cf00084799f473919193aa03333d33bddc30", + "priorProjectedBase": "2d1fe9cde03a021bb4d22539ab9440f2ab5ef30b", "definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f", "toolDigest": "89cbdbe1d79917bad52655817aab0eb545b89183915380eaedc3217f014e6715", - "contentDigest": "9d42748e8e0da39268d1af21ef1b2d7d4f71ba4814e1541f9cf2246dee33a956", + "contentDigest": "cd4484854c325b7a6c6766c9755e49588a1810aa5fa395ed0c05323deb2c7314", "publicationEligible": true } diff --git a/README.md b/README.md index 945d7d24f..aae613b7e 100644 --- a/README.md +++ b/README.md @@ -49,7 +49,7 @@ canonical `deixic-code` command and the retained `maestro` alias. See the The installer verifies the release checksum manifest and Cosign signatures when the release provides them, stages binaries and web assets under a versioned data directory, and swaps only the launcher. Set `MAESTRO_REQUIRE_SIGNED_INSTALL=1` to refuse legacy releases without signed metadata. -Installed interactive sessions check for updates on startup and apply newer releases before opening the TUI. The check is bounded and failures never block startup. Set `MAESTRO_AUTO_UPDATE=0` to opt out, `MAESTRO_AUTO_UPDATE=check` to show availability without installing, or use `deixic-code update --check` for an explicit check. Use `deixic-code update --channel beta` or `deixic-code update --channel alpha` for a one-time channel update; channel installers persist that choice for startup checks. Signed-release installs require Cosign verification during automatic updates; global npm and Bun installs update through their original package manager. +Installed interactive sessions check for a newer release after the TUI paints its first frame. The check is bounded (`MAESTRO_STARTUP_UPDATE_TIMEOUT_MS`, default 350 ms), runs in the background, never delays terminal setup, and is cancelled on quit. When a newer release exists, the transcript shows a notice with the version and `deixic-code update`; nothing is installed automatically. Set `MAESTRO_AUTO_UPDATE=0` to opt out, `MAESTRO_AUTO_UPDATE=apply` to restore the earlier behavior (check, install, and restart before the TUI opens, which delays startup by up to the same bound), or use `deixic-code update --check` for an explicit check. Use `deixic-code update --channel beta` or `deixic-code update --channel alpha` for a one-time channel update; channel installers persist that choice for startup checks. Signed-release installs require Cosign verification during automatic updates; global npm and Bun installs update through their original package manager. Release assets retain the compatibility names `maestro-darwin-arm64`, `maestro-darwin-x64`, `maestro-linux-arm64`, and `maestro-linux-x64`. The npm diff --git a/packages/local-host-rs/examples/semantic_stream_fixture.rs b/packages/local-host-rs/examples/semantic_stream_fixture.rs new file mode 100644 index 000000000..85089d3f9 --- /dev/null +++ b/packages/local-host-rs/examples/semantic_stream_fixture.rs @@ -0,0 +1,87 @@ +//! Deterministic browser fixture: one upstream token stream, two release policies. + +use maestro_local_host::semantic_text::{FlushReason, SemanticTextRelease}; +use serde_json::json; + +fn main() { + let input = concat!( + "## Risk by area\n\nThe largest risk is delayed evidence.\n\n", + "**Next steps**\n\n- Check the run receipt.\n- Verify the rollout.\n\n", + "### Comparison\n\n| Area | Risk |\n| --- | --- |\n| API | Medium |\n| UI | Low |\n\n", + "### Example\n\n```rust\nfn main() { println!(\"ready\"); }\n```\n", + ); + let mut policy = SemanticTextRelease::default(); + let mut legacy = String::new(); + let mut semantic = String::new(); + let mut frames = Vec::new(); + let mut holds = Vec::new(); + let mut title_holds = Vec::new(); + let mut released_chunks = 0; + let mut first_visible_ms = None; + let mut forced_flushes = 0; + let mut legacy_title_frames = 0; + let mut semantic_title_frames = 0; + let rate_ms = 64; + for (index, chunk) in input.as_bytes().chunks(16).enumerate() { + let chunk = std::str::from_utf8(chunk).expect("ASCII fixture"); + let at_ms = (index as u64 + 1) * rate_ms; + legacy.push_str(chunk); + for release in policy.push(chunk, at_ms) { + holds.push(release.held_ms); + if release.semantic_unit { + title_holds.push(release.held_ms); + } + forced_flushes += u64::from(release.forced.is_some()); + first_visible_ms.get_or_insert(at_ms); + released_chunks += 1; + semantic.push_str(&release.text); + } + legacy_title_frames += u64::from(ends_in_standalone_title(&legacy)); + semantic_title_frames += u64::from(ends_in_standalone_title(&semantic)); + frames.push(json!({"atMs": at_ms, "legacy": legacy, "semantic": semantic})); + } + let end_ms = (frames.len() as u64 + 1) * rate_ms; + for release in policy.flush(FlushReason::Finalization, end_ms) { + holds.push(release.held_ms); + if release.semantic_unit { + title_holds.push(release.held_ms); + } + forced_flushes += u64::from(release.forced.is_some()); + first_visible_ms.get_or_insert(end_ms); + released_chunks += 1; + semantic.push_str(&release.text); + } + assert_eq!(legacy, input); + assert_eq!(semantic, input); + frames.push(json!({"atMs": end_ms, "legacy": legacy, "semantic": semantic})); + holds.sort_unstable(); + title_holds.sort_unstable(); + let measured = json!({ + "inputRateCharsPerSecond": 250, + "upstreamChunkBytes": 16, + "upstreamIntervalMs": rate_ms, + "upstreamChunks": input.len().div_ceil(16), + "semanticPublishedChunks": released_chunks, + "semanticMedianHoldMs": holds.get(holds.len() / 2).copied().unwrap_or(0), + "semanticMaxHoldMs": holds.last().copied().unwrap_or(0), + "titleMedianHoldMs": title_holds.get(title_holds.len() / 2).copied().unwrap_or(0), + "titleMaxHoldMs": title_holds.last().copied().unwrap_or(0), + "transportFirstTextMs": rate_ms, + "firstVisibleSemanticContentMs": first_visible_ms, + "peakBytesHeld": policy.peak_held_bytes(), + "forcedFlushes": forced_flushes, + "legacyStandaloneTitleFrames": legacy_title_frames, + "semanticStandaloneTitleFrames": semantic_title_frames, + }); + let fixture = json!({"measured": measured, "frames": frames}); + println!("{}", serde_json::to_string_pretty(&fixture).unwrap()); +} + +fn ends_in_standalone_title(body: &str) -> bool { + let last = body.trim_end().lines().last().unwrap_or_default(); + let value = last.trim(); + value.starts_with("# ") + || value.starts_with("## ") + || value.starts_with("### ") + || value.starts_with("**") && (value.ends_with("**") || value.ends_with("**:")) +} diff --git a/packages/local-host-rs/src/headless/generated_protocol.rs b/packages/local-host-rs/src/headless/generated_protocol.rs index bd9a64c77..1a4e9898e 100644 --- a/packages/local-host-rs/src/headless/generated_protocol.rs +++ b/packages/local-host-rs/src/headless/generated_protocol.rs @@ -68,6 +68,7 @@ pub const HEADLESS_FROM_AGENT_MESSAGE_TYPES: &[&str] = &[ "hello_ok", "ready", "response_start", + "assistant_text_observed", "response_chunk", "response_end", "tool_call", diff --git a/packages/local-host-rs/src/headless/messages.rs b/packages/local-host-rs/src/headless/messages.rs index 97d021d77..0f2732771 100644 --- a/packages/local-host-rs/src/headless/messages.rs +++ b/packages/local-host-rs/src/headless/messages.rs @@ -1007,6 +1007,10 @@ pub enum FromAgentMessage { ResponseStart { response_id: String, }, + /// First non-empty raw assistant delta, with no model text on this event. + AssistantTextObserved { + response_id: String, + }, /// Response chunk (text or thinking) ResponseChunk { response_id: String, diff --git a/packages/local-host-rs/src/headless/messages/state.rs b/packages/local-host-rs/src/headless/messages/state.rs index 564f28a25..bf75f7505 100644 --- a/packages/local-host-rs/src/headless/messages/state.rs +++ b/packages/local-host-rs/src/headless/messages/state.rs @@ -826,6 +826,7 @@ impl AgentState { None } FromAgentMessage::ProcessBudgetCheckpoint { .. } => None, + FromAgentMessage::AssistantTextObserved { .. } => None, FromAgentMessage::ManagedGatewayReceipt { request_id, record_id, diff --git a/packages/local-host-rs/src/headless/proto.rs b/packages/local-host-rs/src/headless/proto.rs index 0237024c4..97f239b13 100644 --- a/packages/local-host-rs/src/headless/proto.rs +++ b/packages/local-host-rs/src/headless/proto.rs @@ -341,6 +341,9 @@ mod tests { FromAgentMessage::ResponseStart { response_id: "response-1".into(), }, + FromAgentMessage::AssistantTextObserved { + response_id: "response-1".into(), + }, FromAgentMessage::ResponseChunk { response_id: "response-1".into(), content: "chunk".into(), @@ -610,6 +613,7 @@ mod tests { | FromAgentMessage::WorkspaceCapabilitySetApplied { .. } | FromAgentMessage::Ready { .. } | FromAgentMessage::ResponseStart { .. } + | FromAgentMessage::AssistantTextObserved { .. } | FromAgentMessage::ResponseChunk { .. } | FromAgentMessage::ResponseEnd { .. } | FromAgentMessage::TurnCompleted { .. } @@ -1161,6 +1165,7 @@ mod tests { FromPayload::HelloOk(_) => "hello_ok", FromPayload::Ready(_) => "ready", FromPayload::ResponseStart(_) => "response_start", + FromPayload::AssistantTextObserved(_) => "assistant_text_observed", FromPayload::ResponseChunk(_) => "response_chunk", FromPayload::ResponseEnd(_) => "response_end", FromPayload::ToolCall(_) => "tool_call", @@ -1247,6 +1252,7 @@ mod tests { FromPayload::HelloOk(Default::default()), FromPayload::Ready(Default::default()), FromPayload::ResponseStart(Default::default()), + FromPayload::AssistantTextObserved(Default::default()), FromPayload::ResponseChunk(Default::default()), FromPayload::ResponseEnd(Default::default()), FromPayload::ToolCall(Default::default()), diff --git a/packages/local-host-rs/src/headless_server.rs b/packages/local-host-rs/src/headless_server.rs index f98d3008e..47f997327 100644 --- a/packages/local-host-rs/src/headless_server.rs +++ b/packages/local-host-rs/src/headless_server.rs @@ -59,6 +59,11 @@ use crate::headless::messages::{ UtilityFileSearchMatch, }; use crate::headless::{HEADLESS_PROTOCOL_VERSION, native_server_capabilities}; +use crate::semantic_text::{ + FlushReason as SemanticFlushReason, Release as SemanticRelease, SemanticTextRelease, +}; + +mod semantic_stream; /// Test-only provider override for the local SDK conformance fixture. /// @@ -124,6 +129,16 @@ struct RuntimeMeta { turn_active: bool, transcript_grade: crate::transcript::TranscriptGrade, response_chunks: Vec<(String, bool)>, + semantic_text: SemanticTextRelease, + semantic_response_id: Option, + semantic_started_at: Option, + semantic_raw_first_ms: Option, + semantic_first_released_ms: Option, + semantic_published_chunks: u64, + semantic_standalone_releases: u64, + semantic_forced_flushes: u64, + semantic_peak_held_bytes: usize, + semantic_total_hold_ms: u64, /// Last safe managed-Gateway evidence for the active turn. managed_gateway_receipt: Option, /// Controller assignment proven against the provider-bound system prompt. @@ -168,19 +183,6 @@ impl RuntimeMeta { self.decided_tool_execution_ids .insert(tool_execution_id.to_string()) } - - fn record_response_chunk( - &mut self, - content: &str, - is_thinking: bool, - ) -> crate::transcript::TranscriptGrade { - let grade = self.transcript_grade; - if grade != crate::transcript::TranscriptGrade::Delta { - self.response_chunks - .push((content.to_string(), is_thinking)); - } - grade - } } struct HeadlessState { @@ -269,6 +271,16 @@ impl HeadlessState { turn_active: false, transcript_grade: crate::transcript::TranscriptGrade::Delta, response_chunks: Vec::new(), + semantic_text: SemanticTextRelease::default(), + semantic_response_id: None, + semantic_started_at: None, + semantic_raw_first_ms: None, + semantic_first_released_ms: None, + semantic_published_chunks: 0, + semantic_standalone_releases: 0, + semantic_forced_flushes: 0, + semantic_peak_held_bytes: 0, + semantic_total_hold_ms: 0, managed_gateway_receipt: None, prompt_experiment: None, receipt_tasks: Arc::new(Mutex::new(Vec::new())), @@ -531,25 +543,34 @@ impl HeadlessState { self.ready_emitted = true; } let event_task = tokio::spawn(async move { - while let Some(msg) = event_rx.recv().await { - if let Err(err) = handle_agent_event( - msg, - &meta_bg, - &tool_tx_bg, - &reported_model, - routed_provider.as_deref(), - ) - .await - { - let _ = emit(&FromAgentMessage::Error { - request_id: None, - message: format!("headless event bridge failed: {err:#}"), - fatal: false, - terminal: true, - error_type: Some(HeadlessErrorType::Protocol), - }); + let mut semantic_tick = + tokio::time::interval(std::time::Duration::from_millis(100)); + semantic_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + loop { + tokio::select! { + msg = event_rx.recv() => { + let Some(msg) = msg else { break }; + if let Err(err) = handle_agent_event( + msg, &meta_bg, &tool_tx_bg, &reported_model, + routed_provider.as_deref(), + ).await { + let _ = semantic_stream::flush_and_emit(&meta_bg, SemanticFlushReason::Error); + let _ = emit(&FromAgentMessage::Error { + request_id: None, + message: format!("headless event bridge failed: {err:#}"), + fatal: false, + terminal: true, + error_type: Some(HeadlessErrorType::Protocol), + }); + } + } + _ = semantic_tick.tick() => { + let _ = semantic_stream::tick_and_emit(&meta_bg); + } } } + let _ = + semantic_stream::flush_and_emit(&meta_bg, SemanticFlushReason::Cancellation); }); self.tool_tx = Some(tool_tx); self.event_task = Some(event_task); @@ -2751,6 +2772,9 @@ async fn handle_agent_event( emit(&FromAgentMessage::ProcessBudgetCheckpoint { budget })?; } } + if let Some(reason) = semantic_stream::boundary_reason(&msg) { + semantic_stream::flush_and_emit(meta, reason)?; + } match msg { FromAgent::ManagedAuthorizationRequest { request_id } => { emit(&FromAgentMessage::ManagedAuthorizationRequest { request_id })?; @@ -2852,6 +2876,7 @@ async fn handle_agent_event( .lock() .unwrap_or_else(std::sync::PoisonError::into_inner); meta.response_chunks.clear(); + semantic_stream::reset_response(&mut meta, &response_id); meta.managed_gateway_receipt = None; drop(meta); emit(&FromAgentMessage::ResponseStart { response_id })?; @@ -2861,21 +2886,10 @@ async fn handle_agent_event( content, is_thinking, } => { - let grade = { - let mut meta = meta - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner); - meta.record_response_chunk(&content, is_thinking) - }; - if grade == crate::transcript::TranscriptGrade::Delta { - emit(&FromAgentMessage::ResponseChunk { - response_id, - content, - is_thinking, - })?; - } + semantic_stream::publish_response_chunk(meta, response_id, content, is_thinking)?; } FromAgent::ResponseEnd { response_id, usage } => { + semantic_stream::flush_and_emit(meta, SemanticFlushReason::Finalization)?; let session_id = meta .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) @@ -2890,6 +2904,7 @@ async fn handle_agent_event( configured_model = %configured_model, routed_provider = routed_provider.unwrap_or(""), ); + semantic_stream::log_response(meta, &response_id); for message in take_interrupted_tool_terminal_messages(meta) { emit(&message)?; } @@ -6217,23 +6232,6 @@ else if(x.method==="turn/start"){const turnId="turn-"+x.id;send({id:x.id,result: assert_eq!(snapshot["processed_queue_ids"], serde_json::json!([7, 9])); } - #[test] - fn delta_transcript_does_not_buffer_emitted_response_chunks() { - let mut meta = RuntimeMeta { - transcript_grade: crate::transcript::TranscriptGrade::Delta, - ..RuntimeMeta::default() - }; - - for _ in 0..10_000 { - assert_eq!( - meta.record_response_chunk("already emitted", false), - crate::transcript::TranscriptGrade::Delta, - ); - } - - assert!(meta.response_chunks.is_empty()); - } - #[test] fn client_tool_content_preserves_text_and_images() { let result = client_content_to_agent_result( diff --git a/packages/local-host-rs/src/headless_server/semantic_stream.rs b/packages/local-host-rs/src/headless_server/semantic_stream.rs new file mode 100644 index 000000000..eed03e629 --- /dev/null +++ b/packages/local-host-rs/src/headless_server/semantic_stream.rs @@ -0,0 +1,231 @@ +//! Headless publication of the shared incremental assistant text policy. + +use super::*; + +impl RuntimeMeta { + fn semantic_now_ms(&self) -> u64 { + self.semantic_started_at + .map_or(0, |start| start.elapsed().as_millis() as u64) + } + + pub(super) fn record_response_chunk( + &mut self, + content: &str, + is_thinking: bool, + ) -> (crate::transcript::TranscriptGrade, Vec) { + let grade = self.transcript_grade; + if grade != crate::transcript::TranscriptGrade::Delta { + self.response_chunks + .push((content.to_string(), is_thinking)); + } + let now_ms = self.semantic_now_ms(); + if !is_thinking && !content.trim().is_empty() { + self.semantic_raw_first_ms.get_or_insert(now_ms); + } + if is_thinking || grade != crate::transcript::TranscriptGrade::Delta { + return (grade, Vec::new()); + } + let releases = self.semantic_text.push(content, now_ms); + self.semantic_peak_held_bytes = self.semantic_text.peak_held_bytes(); + self.record_semantic_releases(&releases, now_ms); + (grade, releases) + } + + fn flush_semantic_text(&mut self, reason: SemanticFlushReason) -> Vec { + let now_ms = self.semantic_now_ms(); + let releases = self.semantic_text.flush(reason, now_ms); + self.record_semantic_releases(&releases, now_ms); + releases + } + + fn record_semantic_releases(&mut self, releases: &[SemanticRelease], now_ms: u64) { + for release in releases { + if !release.text.trim().is_empty() { + self.semantic_first_released_ms.get_or_insert(now_ms); + } + self.semantic_published_chunks += 1; + self.semantic_standalone_releases += u64::from(release.standalone_title); + self.semantic_forced_flushes += u64::from(release.forced.is_some()); + self.semantic_total_hold_ms += release.held_ms; + tracing::debug!( + target: "maestro.semantic_text", + event = "semantic_assistant_text_released", + response_id = self.semantic_response_id.as_deref().unwrap_or(""), + held_ms = release.held_ms, + bytes = release.text.len(), + forced_reason = ?release.forced, + standalone_title = release.standalone_title, + ); + } + } +} + +fn emit_releases(response_id: Option, releases: Vec) -> Result<()> { + if let Some(response_id) = response_id { + for release in releases { + emit(&FromAgentMessage::ResponseChunk { + response_id: response_id.clone(), + content: release.text, + is_thinking: false, + })?; + } + } + Ok(()) +} + +pub(super) fn flush_and_emit( + meta: &Arc>, + reason: SemanticFlushReason, +) -> Result<()> { + let (response_id, releases) = { + let mut runtime = meta + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + ( + runtime.semantic_response_id.clone(), + runtime.flush_semantic_text(reason), + ) + }; + emit_releases(response_id, releases) +} + +pub(super) fn tick_and_emit(meta: &Arc>) -> Result<()> { + let (response_id, releases) = { + let mut runtime = meta + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let now_ms = runtime.semantic_now_ms(); + let releases = runtime.semantic_text.tick(now_ms); + runtime.record_semantic_releases(&releases, now_ms); + (runtime.semantic_response_id.clone(), releases) + }; + emit_releases(response_id, releases) +} + +pub(super) fn boundary_reason(msg: &FromAgent) -> Option { + match msg { + FromAgent::ToolCall { .. } | FromAgent::ToolStart { .. } | FromAgent::ToolEnd { .. } => { + Some(SemanticFlushReason::Tool) + } + FromAgent::ResponseStart { .. } => Some(SemanticFlushReason::Finalization), + FromAgent::TurnCompleted { .. } => Some(SemanticFlushReason::TurnEnd), + FromAgent::TurnInterrupted { .. } => Some(SemanticFlushReason::Cancellation), + FromAgent::Error { .. } | FromAgent::ProviderError { .. } => { + Some(SemanticFlushReason::Error) + } + _ => None, + } +} + +pub(super) fn reset_response(meta: &mut RuntimeMeta, response_id: &str) { + meta.semantic_text = SemanticTextRelease::default(); + meta.semantic_response_id = Some(response_id.to_owned()); + meta.semantic_started_at = Some(Instant::now()); + meta.semantic_raw_first_ms = None; + meta.semantic_first_released_ms = None; + meta.semantic_published_chunks = 0; + meta.semantic_standalone_releases = 0; + meta.semantic_forced_flushes = 0; + meta.semantic_peak_held_bytes = 0; + meta.semantic_total_hold_ms = 0; +} + +pub(super) fn publish_response_chunk( + meta: &Arc>, + response_id: String, + content: String, + is_thinking: bool, +) -> Result<()> { + let (grade, first_raw_text, releases) = { + let mut runtime = meta + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let first_raw_text = + !is_thinking && !content.trim().is_empty() && runtime.semantic_raw_first_ms.is_none(); + let (grade, releases) = runtime.record_response_chunk(&content, is_thinking); + (grade, first_raw_text, releases) + }; + if first_raw_text { + emit(&FromAgentMessage::AssistantTextObserved { + response_id: response_id.clone(), + })?; + } + if grade == crate::transcript::TranscriptGrade::Delta { + if is_thinking { + emit(&FromAgentMessage::ResponseChunk { + response_id, + content, + is_thinking, + })?; + } else { + emit_releases(Some(response_id), releases)?; + } + } + Ok(()) +} + +pub(super) fn log_response(meta: &Arc>, response_id: &str) { + let runtime = meta + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + tracing::info!( + target: "maestro.semantic_text", + event = "semantic_assistant_text_response", + response_id, + transport_first_text_ms = ?runtime.semantic_raw_first_ms, + first_released_content_ms = ?runtime.semantic_first_released_ms, + standalone_title_releases = runtime.semantic_standalone_releases, + titles_seen = runtime.semantic_text.titles_seen(), + standalone_title_release_percent = if runtime.semantic_text.titles_seen() == 0 { + 0.0 + } else { + 100.0 * runtime.semantic_standalone_releases as f64 + / runtime.semantic_text.titles_seen() as f64 + }, + published_text_chunks = runtime.semantic_published_chunks, + forced_flushes = runtime.semantic_forced_flushes, + peak_bytes_held = runtime.semantic_peak_held_bytes, + total_hold_ms = runtime.semantic_total_hold_ms, + ); +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn headless_delta_bridge_publishes_title_with_following_text() { + let mut meta = RuntimeMeta { + transcript_grade: crate::transcript::TranscriptGrade::Delta, + semantic_started_at: Some(Instant::now()), + ..RuntimeMeta::default() + }; + let (_, title) = meta.record_response_chunk("## Risk\n", false); + assert!(title.is_empty()); + assert_eq!(meta.semantic_text.titles_seen(), 1); + assert!(meta.semantic_raw_first_ms.is_some()); + assert!(meta.semantic_first_released_ms.is_none()); + let (_, paragraph) = meta.record_response_chunk("The risk is delay.\n", false); + assert_eq!(paragraph.len(), 1); + assert_eq!(paragraph[0].text, "## Risk\nThe risk is delay.\n"); + assert!(meta.semantic_first_released_ms.is_some()); + assert_eq!(meta.semantic_published_chunks, 1); + assert_eq!(meta.semantic_standalone_releases, 0); + } + #[test] + fn delta_transcript_does_not_buffer_emitted_response_chunks() { + let mut meta = RuntimeMeta { + transcript_grade: crate::transcript::TranscriptGrade::Delta, + ..RuntimeMeta::default() + }; + + for _ in 0..10_000 { + assert_eq!( + meta.record_response_chunk("already emitted", false).0, + crate::transcript::TranscriptGrade::Delta, + ); + } + + assert!(meta.response_chunks.is_empty()); + } +} diff --git a/packages/local-host-rs/src/lib.rs b/packages/local-host-rs/src/lib.rs index 861a8d823..7f0f69e5c 100644 --- a/packages/local-host-rs/src/lib.rs +++ b/packages/local-host-rs/src/lib.rs @@ -57,6 +57,7 @@ mod private_code_authority; pub mod rlm; pub mod safety; pub mod sandbox_policy; +pub mod semantic_text; pub mod service_connections; pub mod session; pub mod skill_cli; diff --git a/packages/local-host-rs/src/semantic_text.rs b/packages/local-host-rs/src/semantic_text.rs new file mode 100644 index 000000000..204661279 --- /dev/null +++ b/packages/local-host-rs/src/semantic_text.rs @@ -0,0 +1,517 @@ +//! Incremental release policy for assistant text on the headless protocol. +//! The bytes in `Release::text` are always a prefix of the bytes pushed so far. + +const DEFAULT_MAX_BYTES: usize = 4_096; +const DEFAULT_MAX_HOLD_MS: u64 = 1_500; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum FlushReason { + Tool, + TurnEnd, + Error, + Cancellation, + Finalization, + SizeLimit, + LatencyLimit, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Release { + pub text: String, + pub held_ms: u64, + pub forced: Option, + pub standalone_title: bool, + pub semantic_unit: bool, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum Following { + None, + Awaiting, + TableHeader, + TableDelimiter, + Fence { marker: char, width: usize }, +} + +#[derive(Debug, Clone)] +pub struct SemanticTextRelease { + pending: String, + line: String, + line_non_title: bool, + following: Following, + held_since_ms: Option, + max_bytes: usize, + max_hold_ms: u64, + titles_seen: u64, + peak_held_bytes: usize, +} + +impl Default for SemanticTextRelease { + fn default() -> Self { + Self::new(DEFAULT_MAX_BYTES, DEFAULT_MAX_HOLD_MS) + } +} + +impl SemanticTextRelease { + pub fn new(max_bytes: usize, max_hold_ms: u64) -> Self { + Self { + pending: String::new(), + line: String::new(), + line_non_title: false, + following: Following::None, + held_since_ms: None, + max_bytes: max_bytes.max(4), + max_hold_ms, + titles_seen: 0, + peak_held_bytes: 0, + } + } + + pub fn held_bytes(&self) -> usize { + self.pending.len() + } + + pub fn titles_seen(&self) -> u64 { + self.titles_seen + } + + pub fn peak_held_bytes(&self) -> usize { + self.peak_held_bytes + } + + pub fn push(&mut self, delta: &str, now_ms: u64) -> Vec { + let mut releases = Vec::new(); + for ch in delta.chars() { + if self.pending.len() + ch.len_utf8() > self.max_bytes { + self.release(now_ms, Some(FlushReason::SizeLimit), &mut releases); + self.line.clear(); + self.line_non_title = true; + } + if self.pending.is_empty() { + self.held_since_ms = Some(now_ms); + } + self.pending.push(ch); + self.peak_held_bytes = self.peak_held_bytes.max(self.pending.len()); + if !self.line_non_title { + self.line.push(ch); + } + if ch == '\n' { + if self.line_non_title { + self.release(now_ms, None, &mut releases); + } else { + self.finish_line(now_ms, &mut releases); + } + self.line.clear(); + self.line_non_title = false; + } else if self.following == Following::None + && !self.line_non_title + && !could_be_title(&self.line) + { + // Ordinary prose need not wait for a line or a transport timer. + self.release(now_ms, None, &mut releases); + self.line.clear(); + self.line_non_title = true; + } else if self.line_non_title && self.following == Following::None { + self.release(now_ms, None, &mut releases); + } + if self.pending.len() >= self.max_bytes { + self.release(now_ms, Some(FlushReason::SizeLimit), &mut releases); + self.line.clear(); + self.line_non_title = true; + } + } + if self + .held_since_ms + .is_some_and(|start| now_ms.saturating_sub(start) >= self.max_hold_ms) + { + self.release(now_ms, Some(FlushReason::LatencyLimit), &mut releases); + if !self.line.is_empty() { + self.line.clear(); + self.line_non_title = true; + } + } + // Ordinary prose can pass through immediately, but a provider token + // is never a reason to publish one event per character. + let mut compact: Vec = Vec::new(); + for release in releases { + if !release.semantic_unit && release.forced.is_none() { + if let Some(last) = compact.last_mut() { + if !last.semantic_unit + && last.forced.is_none() + && last.text.len() + release.text.len() <= self.max_bytes + { + last.text.push_str(&release.text); + continue; + } + } + } + compact.push(release); + } + compact + } + + pub fn tick(&mut self, now_ms: u64) -> Vec { + let mut releases = Vec::new(); + if self + .held_since_ms + .is_some_and(|start| now_ms.saturating_sub(start) >= self.max_hold_ms) + { + self.release(now_ms, Some(FlushReason::LatencyLimit), &mut releases); + if !self.line.is_empty() { + self.line.clear(); + self.line_non_title = true; + } + } + releases + } + + pub fn flush(&mut self, reason: FlushReason, now_ms: u64) -> Vec { + let mut releases = Vec::new(); + self.release(now_ms, Some(reason), &mut releases); + self.following = Following::None; + self.line.clear(); + self.line_non_title = false; + releases + } + + fn finish_line(&mut self, now_ms: u64, releases: &mut Vec) { + let line = self.line.trim_end_matches(['\r', '\n']); + match self.following { + Following::None if standalone_title(line) => { + self.following = Following::Awaiting; + self.titles_seen += 1; + } + Following::None => self.release(now_ms, None, releases), + Following::Awaiting if line.trim().is_empty() => {} + Following::Awaiting if standalone_title(line) => {} + Following::Awaiting => { + if let Some((marker, width)) = opening_fence(line) { + self.following = Following::Fence { marker, width }; + } else if table_row(line) { + self.following = Following::TableHeader; + } else { + // A complete first paragraph line or list item is useful + // with the title. Later lines flow without title buffering. + self.release(now_ms, None, releases); + } + } + Following::TableHeader if table_delimiter(line) => { + self.following = Following::TableDelimiter; + } + Following::TableDelimiter if line.trim().is_empty() => {} + Following::TableDelimiter if table_row(line) => { + self.release(now_ms, None, releases); + } + Following::TableHeader | Following::TableDelimiter => { + self.release(now_ms, None, releases); + } + Following::Fence { marker, width } if closing_fence(line, marker, width) => { + self.release(now_ms, None, releases); + } + Following::Fence { .. } => {} + } + } + + fn release(&mut self, now_ms: u64, forced: Option, releases: &mut Vec) { + if self.pending.is_empty() { + return; + } + let incomplete_title = self.following == Following::None && standalone_title(&self.line); + if incomplete_title { + self.titles_seen += 1; + } + let standalone_title = forced.is_some() + && (self.following == Following::Awaiting && self.line.trim().is_empty() + || incomplete_title); + let semantic_unit = self.following != Following::None; + releases.push(Release { + text: std::mem::take(&mut self.pending), + held_ms: now_ms.saturating_sub(self.held_since_ms.unwrap_or(now_ms)), + forced, + standalone_title, + semantic_unit, + }); + self.held_since_ms = None; + self.following = Following::None; + } +} + +fn could_be_title(line: &str) -> bool { + let value = line.trim_start_matches(' '); + if value.starts_with('#') { + let width = value.bytes().take_while(|byte| *byte == b'#').count(); + return width <= 6 && (value.len() == width || value.as_bytes().get(width) == Some(&b' ')); + } + if value == "*" { + return true; + } + if let Some(body) = value.strip_prefix("**") { + if let Some(end) = body.find("**") { + return body[end + 2..].trim().trim_end_matches(':').is_empty(); + } + return true; + } + false +} + +fn standalone_title(line: &str) -> bool { + let value = line.trim(); + if value.starts_with('#') { + let width = value.bytes().take_while(|byte| *byte == b'#').count(); + return (1..=6).contains(&width) + && value.as_bytes().get(width) == Some(&b' ') + && !value[width + 1..].trim().is_empty(); + } + if let Some(body) = value.strip_prefix("**") { + let body = body.strip_suffix(':').unwrap_or(body); + if let Some(body) = body.strip_suffix("**") { + return !body.trim().is_empty(); + } + } + false +} + +fn opening_fence(line: &str) -> Option<(char, usize)> { + let value = line.trim_start_matches(' '); + let marker = value.chars().next()?; + if !matches!(marker, '`' | '~') { + return None; + } + let width = value.chars().take_while(|ch| *ch == marker).count(); + (width >= 3).then_some((marker, width)) +} + +fn closing_fence(line: &str, marker: char, width: usize) -> bool { + let value = line.trim(); + let count = value.chars().take_while(|ch| *ch == marker).count(); + count >= width && value[count..].trim().is_empty() +} + +fn table_row(line: &str) -> bool { + line.trim().contains('|') +} + +fn table_delimiter(line: &str) -> bool { + let value = line.trim().trim_matches('|'); + let mut cells = 0; + for cell in value.split('|') { + let cell = cell.trim().trim_matches(':'); + if cell.len() < 3 || !cell.bytes().all(|byte| byte == b'-') { + return false; + } + cells += 1; + } + cells >= 2 +} + +#[cfg(test)] +mod tests { + use super::*; + + fn run(input: &str, chunk: usize) -> Vec { + let mut policy = SemanticTextRelease::default(); + let mut out = Vec::new(); + for slice in input.as_bytes().chunks(chunk) { + // ASCII fixtures make arbitrary byte slicing valid UTF-8 here. + out.extend(policy.push(std::str::from_utf8(slice).unwrap(), 0)); + } + out.extend(policy.flush(FlushReason::Finalization, 0)); + assert_eq!( + out.iter() + .map(|release| release.text.as_str()) + .collect::(), + input + ); + out + } + + #[test] + fn titles_wait_for_the_next_useful_block() { + let mut policy = SemanticTextRelease::default(); + assert!(policy.push("## Risk by area", 0).is_empty()); + assert!(policy.push("\n\n", 0).is_empty()); + assert_eq!( + policy.push("The largest risk is...\n", 100)[0].text, + "## Risk by area\n\nThe largest risk is...\n" + ); + assert!(policy.push("**Next steps**\n", 200).is_empty()); + assert_eq!( + policy.push("- First...\n", 300)[0].text, + "**Next steps**\n- First...\n" + ); + assert!(policy.push("**Decision**:\n", 400).is_empty()); + assert_eq!( + policy.push("Proceed with care.\n", 500)[0].text, + "**Decision**:\nProceed with care.\n" + ); + } + + #[test] + fn table_and_fence_wait_until_coherent() { + let mut policy = SemanticTextRelease::default(); + assert!( + policy + .push("# Table\n| A | B |\n| --- | --- |\n", 0) + .is_empty() + ); + assert_eq!( + policy.push("| 1 | 2 |\n", 0)[0].text, + "# Table\n| A | B |\n| --- | --- |\n| 1 | 2 |\n" + ); + assert!(policy.push("# Code\n```rs\nfn main() {}\n", 0).is_empty()); + assert_eq!( + policy.push("```\n", 0)[0].text, + "# Code\n```rs\nfn main() {}\n```\n" + ); + } + + #[test] + fn inline_bold_and_prior_prose_flow() { + let out = run( + "Before heading\n# Heading\nFollowing.\n**Note:** read this\n", + 1, + ); + let text = out + .iter() + .map(|release| release.text.as_str()) + .collect::(); + assert_eq!( + text, + "Before heading\n# Heading\nFollowing.\n**Note:** read this\n" + ); + let heading_index = out + .iter() + .position(|release| release.text.contains("# Heading")) + .unwrap(); + assert_eq!( + out[..heading_index] + .iter() + .map(|release| release.text.as_str()) + .collect::(), + "Before heading\n" + ); + assert!(out[heading_index].text.contains("Following.")); + assert!( + out.iter() + .any(|release| release.text.contains("**Note:**") && !release.semantic_unit) + ); + } + + #[test] + fn limits_and_boundaries_never_lose_text() { + for reason in [ + FlushReason::Tool, + FlushReason::TurnEnd, + FlushReason::Error, + FlushReason::Cancellation, + FlushReason::Finalization, + ] { + let mut policy = SemanticTextRelease::default(); + assert!(policy.push("# Held\n", 0).is_empty()); + assert_eq!(policy.flush(reason, 5)[0].text, "# Held\n"); + } + let mut policy = SemanticTextRelease::new(12, 10); + assert!(policy.push("# Held\n", 0).is_empty()); + assert_eq!(policy.tick(10)[0].forced, Some(FlushReason::LatencyLimit)); + let mut incomplete = SemanticTextRelease::new(128, 10); + assert!(incomplete.push("# Held", 0).is_empty()); + assert!(incomplete.tick(10)[0].standalone_title); + assert_eq!(incomplete.titles_seen(), 1); + assert!(!policy.push("# Long title\n", 20).is_empty()); + let mut unclosed = SemanticTextRelease::new(128, 10); + assert!(unclosed.push("# Code\n```rust\nlet x = 1;\n", 0).is_empty()); + let timed = unclosed.tick(10); + assert_eq!(timed.len(), 1); + assert_eq!(timed[0].forced, Some(FlushReason::LatencyLimit)); + assert_eq!(timed[0].text, "# Code\n```rust\nlet x = 1;\n"); + assert_eq!( + run(&"normal prose ".repeat(20_000), 1) + .iter() + .map(|r| r.text.len()) + .sum::(), + 260_000 + ); + } + + #[test] + fn arbitrary_chunk_boundaries_preserve_semantics() { + let input = "# A\nFirst.\n**B**\n- item\n# C\n| A | B |\n| --- | --- |\n| 1 | 2 |\n# D\n```\ncode\n```\n"; + let canonical = run(input, 1) + .into_iter() + .filter(|release| release.semantic_unit) + .map(|release| release.text) + .collect::>(); + for chunk in 2..=input.len() { + let actual = run(input, chunk) + .into_iter() + .filter(|release| release.semantic_unit) + .map(|release| release.text) + .collect::>(); + assert_eq!(actual, canonical, "chunk size {chunk}"); + } + } + + #[test] + fn generated_unicode_streams_preserve_every_byte_at_every_character_split() { + let atoms = ["Ă©", "đź§Ş", "\n", "#", "*", "|", "`", " text", "- item", ":"]; + let mut seed = 0x5eed_cafe_u64; + for _case in 0..80 { + let mut input = String::new(); + for _ in 0..160 { + seed ^= seed << 13; + seed ^= seed >> 7; + seed ^= seed << 17; + input.push_str(atoms[(seed as usize) % atoms.len()]); + } + let mut policy = SemanticTextRelease::new(128, 30); + let mut output = String::new(); + for (index, ch) in input.chars().enumerate() { + let now_ms = index as u64 * 100; + for release in policy.push(ch.encode_utf8(&mut [0; 4]), now_ms) { + output.push_str(&release.text); + } + if index % 37 == 0 { + for release in policy.tick(now_ms + 30) { + output.push_str(&release.text); + } + } + if index % 53 == 0 { + for release in policy.flush(FlushReason::Tool, now_ms + 31) { + output.push_str(&release.text); + } + } + assert!(policy.held_bytes() < 128); + } + for release in policy.flush(FlushReason::Finalization, 500) { + output.push_str(&release.text); + } + assert_eq!(output.as_bytes(), input.as_bytes()); + } + } + + #[test] + fn one_megabyte_stream_stays_bounded_and_reports_throughput() { + let input = "ordinary assistant prose. ".repeat(43_692); + let mut policy = SemanticTextRelease::default(); + let mut output = String::with_capacity(input.len()); + let started = std::time::Instant::now(); + for chunk in input.as_bytes().chunks(32) { + for release in policy.push(std::str::from_utf8(chunk).unwrap(), 0) { + assert!(release.text.len() <= 4_096); + output.push_str(&release.text); + } + assert!(policy.held_bytes() <= 4_096); + } + for release in policy.flush(FlushReason::Finalization, 0) { + output.push_str(&release.text); + } + assert_eq!(output.as_bytes(), input.as_bytes()); + assert!(policy.peak_held_bytes() <= 4_096); + eprintln!( + "semantic_text_perf bytes={} elapsed_ms={} peak_held_bytes={}", + input.len(), + started.elapsed().as_millis(), + policy.peak_held_bytes() + ); + } +} diff --git a/packages/local-host-rs/src/ui_wake.rs b/packages/local-host-rs/src/ui_wake.rs index 9a0c7a891..60e3c81cc 100644 --- a/packages/local-host-rs/src/ui_wake.rs +++ b/packages/local-host-rs/src/ui_wake.rs @@ -11,7 +11,9 @@ static HOOK: Mutex>> = Mutex::new(None); /// Install or clear the UI wake hook. The TUI sets this while its loop runs. pub fn set_hook(hook: Option>) { - *HOOK.lock().unwrap_or_else(std::sync::PoisonError::into_inner) = hook; + *HOOK + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) = hook; } /// Wake the installed UI loop, if any. Safe to call when no hook is set. diff --git a/packages/runtime-contracts-rs/fixtures/headless-protocol-v1.json b/packages/runtime-contracts-rs/fixtures/headless-protocol-v1.json index 79010a9c0..93c892b9d 100644 --- a/packages/runtime-contracts-rs/fixtures/headless-protocol-v1.json +++ b/packages/runtime-contracts-rs/fixtures/headless-protocol-v1.json @@ -47,6 +47,7 @@ "hello_ok", "ready", "response_start", + "assistant_text_observed", "response_chunk", "response_end", "tool_call", diff --git a/packages/runtime-contracts-rs/src/protocol.rs b/packages/runtime-contracts-rs/src/protocol.rs index d5698ef58..76602901d 100644 --- a/packages/runtime-contracts-rs/src/protocol.rs +++ b/packages/runtime-contracts-rs/src/protocol.rs @@ -211,6 +211,7 @@ pub enum FromRuntimeMessageType { HelloOk, Ready, ResponseStart, + AssistantTextObserved, ResponseChunk, ResponseEnd, ToolCall, @@ -357,6 +358,7 @@ pub const HEADLESS_FROM_RUNTIME_MESSAGE_NAMES: &[&str] = &[ "hello_ok", "ready", "response_start", + "assistant_text_observed", "response_chunk", "response_end", "tool_call", @@ -524,6 +526,7 @@ const FROM_RUNTIME_MESSAGES: &[FromRuntimeMessageType] = &[ FromRuntimeMessageType::HelloOk, FromRuntimeMessageType::Ready, FromRuntimeMessageType::ResponseStart, + FromRuntimeMessageType::AssistantTextObserved, FromRuntimeMessageType::ResponseChunk, FromRuntimeMessageType::ResponseEnd, FromRuntimeMessageType::ToolCall, @@ -1267,7 +1270,7 @@ mod tests { assert_eq!(headless_protocol_capability_digest(), expected); assert_eq!( headless_protocol_capability_digest(), - "sha256:43703fc50891d32f11e169729feb0a9eed7b92546a53234c758ba451dc1f4e53" + "sha256:8353f3612f7e7f4f7f5c55020b2c255f7ad52b57417a027ae66ceb657c402ba8" ); } } diff --git a/packages/tui-rs/src/app.rs b/packages/tui-rs/src/app.rs index c542b3e3f..9a5984ca5 100644 --- a/packages/tui-rs/src/app.rs +++ b/packages/tui-rs/src/app.rs @@ -808,6 +808,9 @@ pub struct App { setup_login_rx: Option>>, setup_login_url_rx: Option>, setup_login_task: Option>, + /// Bounded release check spawned after the first frame; its result is a + /// transcript notice. Never awaited by the startup path. + startup_update_check: startup_update::StartupUpdateCheck, onboarding: onboarding::OnboardingSession, pending_agent_spawn: bool, /// True when `current_model` was chosen by `/setup` or `/model`, so a @@ -1792,6 +1795,7 @@ impl App { setup_login_rx: None, setup_login_url_rx: None, setup_login_task: None, + startup_update_check: startup_update::StartupUpdateCheck::idle(), onboarding: onboarding::OnboardingSession::default(), pending_agent_spawn: false, current_model_user_set: false, @@ -2274,6 +2278,8 @@ Always use tools when they would be helpful. Be concise and direct in your respo /// Exit code for the process (0 = success, non-zero = error). pub async fn run(&mut self) -> Result { let result = self.run_inner().await; + // Quit never waits for an in-flight release check. + self.startup_update_check.cancel(); let disable_theme_reporting = self.prepare_terminal_restore(); if disable_theme_reporting { let _ = terminal::disable_theme_reporting(); @@ -2321,6 +2327,12 @@ Always use tools when they would be helpful. Be concise and direct in your respo let startup_seqs = self.terminal_notifier.session_started(); Self::write_terminal_sequences(&startup_seqs); + // The bounded release check runs only now, after the first frame, and + // reports through the loop wake. The opt-in + // `MAESTRO_AUTO_UPDATE=apply` path ran before terminal setup instead + // (`update_cli::run_startup_update`). + self.spawn_startup_update_check(); + // Index @-mention files with a bounded, killable scan (see workspace.rs). // Kick it off on a background thread so agent spawn is not gated on it. // The scan was started before first-run presentation and continues after skip. @@ -2372,6 +2384,9 @@ Always use tools when they would be helpful. Be concise and direct in your respo if self.poll_setup_login() { needs_redraw = true; } + if self.poll_startup_update_notice() { + needs_redraw = true; + } if self.poll_onboarding() { needs_redraw = true; } @@ -5879,6 +5894,7 @@ mod session_commands; mod session_recording; mod session_transition; mod startup; +mod startup_update; #[cfg(test)] mod tests; diff --git a/packages/tui-rs/src/app/startup_update.rs b/packages/tui-rs/src/app/startup_update.rs new file mode 100644 index 000000000..9fbcf0825 --- /dev/null +++ b/packages/tui-rs/src/app/startup_update.rs @@ -0,0 +1,258 @@ +//! Release notice for installed interactive launches, off the first-frame path. +//! +//! `App::run_inner` paints the first frame, then spawns the bounded check +//! from `update_cli::startup_update_notice`. The result travels through a +//! oneshot the main loop drains with `try_recv`, and the task signals the +//! shared `LoopWake` so an idle loop wakes for it. Nothing on the startup path +//! awaits the check; quit aborts it. +//! +//! The opt-in `MAESTRO_AUTO_UPDATE=apply` path (check, install, restart before +//! terminal setup) stays in `update_cli::run_startup_update` and is not +//! touched here. + +use std::future::Future; + +use tokio::sync::oneshot; + +use super::App; +use crate::loop_wake::LoopWake; +use crate::update_cli::StartupUpdateNotice; + +pub(super) struct StartupUpdateCheck { + notice_rx: Option>, + task: Option>, +} + +impl StartupUpdateCheck { + /// No check in flight. `App::new` starts here; `run_inner` spawns later. + pub(super) fn idle() -> Self { + Self { + notice_rx: None, + task: None, + } + } + + /// Run `check` on its own task. The task signals `wake` when it ends, + /// with or without a notice, so the loop's idle wait returns for it. + pub(super) fn spawn( + check: impl Future> + Send + 'static, + wake: LoopWake, + ) -> Self { + let (notice_tx, notice_rx) = oneshot::channel(); + let task = tokio::spawn(async move { + if let Some(notice) = check.await { + let _ = notice_tx.send(notice); + } + wake.signal(); + }); + Self { + notice_rx: Some(notice_rx), + task: Some(task), + } + } + + /// Non-blocking. Returns the notice once, then never again. + pub(super) fn try_take(&mut self) -> Option { + let notice_rx = self.notice_rx.as_mut()?; + match notice_rx.try_recv() { + Ok(notice) => { + self.notice_rx = None; + Some(notice) + } + Err(oneshot::error::TryRecvError::Empty) => None, + Err(oneshot::error::TryRecvError::Closed) => { + self.notice_rx = None; + None + } + } + } + + /// Abort an in-flight check and drop any undelivered notice. + pub(super) fn cancel(&mut self) { + if let Some(task) = self.task.take() { + task.abort(); + } + self.notice_rx = None; + } + + #[cfg(test)] + pub(super) fn is_pending(&self) -> bool { + self.task.as_ref().is_some_and(|task| !task.is_finished()) + } +} + +impl App { + /// Spawn the production check. Called once, after the first frame. + pub(super) fn spawn_startup_update_check(&mut self) { + self.spawn_startup_update_check_with(crate::update_cli::startup_update_notice()); + } + + /// Spawn `check` in place of the production check. Returns immediately. + pub(super) fn spawn_startup_update_check_with( + &mut self, + check: impl Future> + Send + 'static, + ) { + self.startup_update_check.cancel(); + self.startup_update_check = StartupUpdateCheck::spawn(check, self.loop_wake.clone()); + } + + /// Drain a finished check into the transcript. Non-blocking. + pub(super) fn poll_startup_update_notice(&mut self) -> bool { + let Some(notice) = self.startup_update_check.try_take() else { + return false; + }; + let message = self.state.locale.format( + "Deixic Code {0} is available (current {1}); run `deixic-code update`.", + &[notice.latest, notice.current], + ); + self.state.add_system_message(message); + true + } +} + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use super::*; + use crate::app::tests::new_test_app; + use crate::loop_wake::{LoopWakeCause, await_loop_wake}; + + fn notice(latest: &str, current: &str) -> StartupUpdateNotice { + StartupUpdateNotice { + latest: latest.to_owned(), + current: current.to_owned(), + } + } + + /// A check that never resolves must not hold the first frame, the poll, + /// or the app's ability to continue. This is the startup contract. + #[tokio::test] + async fn first_frame_renders_while_the_startup_update_check_is_pending() { + let mut app = new_test_app(); + app.spawn_startup_update_check_with(std::future::pending()); + + app.render() + .expect("first frame renders with the release check still pending"); + assert!( + !app.poll_startup_update_notice(), + "a pending check must not produce a notice" + ); + assert!( + app.startup_update_check.is_pending(), + "the check is still running after the frame" + ); + assert!( + app.state + .messages + .iter() + .all(|message| !message.content.contains("is available")), + "no notice reaches the transcript before the check completes" + ); + + app.startup_update_check.cancel(); + assert!( + !app.startup_update_check.is_pending(), + "cancel drops the task handle" + ); + } + + /// Wait through the loop's own idle wait until the check task has ended. + /// Returns whether a producer signal (the task's `wake.signal()`) was + /// observed on the way. + async fn wait_for_check_to_end(app: &App) -> bool { + let mut producer_wake = false; + for _ in 0..400 { + let cause = await_loop_wake( + Duration::from_millis(25), + &app.loop_wake, + std::future::pending::<()>(), + ) + .await; + producer_wake |= cause == LoopWakeCause::Producer; + if !app.startup_update_check.is_pending() { + return producer_wake; + } + } + panic!("startup update check did not end"); + } + + #[tokio::test] + async fn a_finished_check_wakes_the_loop_and_posts_one_transcript_notice() { + let mut app = new_test_app(); + app.spawn_startup_update_check_with(async { Some(notice("9.9.9", "1.0.0")) }); + + assert!( + wait_for_check_to_end(&app).await, + "the check signals the loop wake" + ); + + assert!(app.poll_startup_update_notice()); + let last = app.state.messages.last().expect("notice message"); + assert!( + last.content + .contains("Deixic Code 9.9.9 is available (current 1.0.0)"), + "unexpected notice: {}", + last.content + ); + assert!( + !app.poll_startup_update_notice(), + "the notice is delivered once" + ); + } + + #[tokio::test] + async fn a_check_with_no_newer_release_posts_nothing() { + let mut app = new_test_app(); + let before = app.state.messages.len(); + app.spawn_startup_update_check_with(async { None }); + + assert!(wait_for_check_to_end(&app).await); + + assert!(!app.poll_startup_update_notice()); + assert_eq!(app.state.messages.len(), before); + assert!(!app.poll_startup_update_notice()); + } + + /// Set when the check future is dropped, which is what abort does to a + /// task that never resolves. + struct DropFlag(std::sync::Arc); + + impl Drop for DropFlag { + fn drop(&mut self) { + self.0.store(true, std::sync::atomic::Ordering::SeqCst); + } + } + + #[tokio::test] + async fn cancel_aborts_a_check_that_never_resolves() { + use std::sync::atomic::Ordering; + + let dropped = std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)); + let flag = DropFlag(std::sync::Arc::clone(&dropped)); + let mut check = StartupUpdateCheck::spawn( + async move { + let _flag = flag; + std::future::pending::>().await + }, + LoopWake::new(), + ); + tokio::task::yield_now().await; + assert!(check.is_pending()); + assert!(!dropped.load(Ordering::SeqCst)); + + check.cancel(); + for _ in 0..100 { + if dropped.load(Ordering::SeqCst) { + break; + } + tokio::time::sleep(Duration::from_millis(1)).await; + } + assert!( + dropped.load(Ordering::SeqCst), + "cancel must drop the pending check future" + ); + assert!(check.try_take().is_none()); + assert!(!check.is_pending()); + } +} diff --git a/packages/tui-rs/src/entrypoint.rs b/packages/tui-rs/src/entrypoint.rs index 3eeb06374..cfd6e8984 100644 --- a/packages/tui-rs/src/entrypoint.rs +++ b/packages/tui-rs/src/entrypoint.rs @@ -949,6 +949,10 @@ pub async fn run_cli(raw_args: Vec) -> Result<()> { if classify_agent_entry(&raw_args) == AgentEntry::ClapParsed && classify_clap_dispatch(&raw_args) == ClapDispatch::Interactive { + // Opt-in `MAESTRO_AUTO_UPDATE=apply` only: check, install, and restart + // before terminal setup. Every other mode returns `None` here without + // touching the network; the default notice check runs inside the TUI + // after its first frame (`App::spawn_startup_update_check`). if let Some(exit_code) = crate::update_cli::run_startup_update(&raw_args).await { std::process::exit(exit_code); } diff --git a/packages/tui-rs/src/markdown.rs b/packages/tui-rs/src/markdown.rs index fa5a4dcc9..9abe99a15 100644 --- a/packages/tui-rs/src/markdown.rs +++ b/packages/tui-rs/src/markdown.rs @@ -660,7 +660,16 @@ mod tests { let mut custom = crate::themes::light_theme(); custom.colors.md_heading = "#ff0000".into(); custom.colors.md_link = "#00ff00".into(); - assert_ne!(custom.get_color("md_link"), custom.get_color("md_heading")); + assert_ne!( + crate::palette::color_for_level(0, 255, 0, crate::palette::ColorLevel::Indexed), + crate::palette::color_for_level(255, 0, 0, crate::palette::ColorLevel::Indexed) + ); + if crate::palette::color_level() == crate::palette::ColorLevel::None { + assert_eq!(custom.get_color("md_link"), Some(Color::Reset)); + assert_eq!(custom.get_color("md_heading"), Some(Color::Reset)); + } else { + assert_ne!(custom.get_color("md_link"), custom.get_color("md_heading")); + } for theme in [crate::themes::light_theme(), custom] { let line = parse_markdown_line_with_theme("See [guide](https://example.com).", &theme); let link = line diff --git a/packages/tui-rs/src/themes/mod.rs b/packages/tui-rs/src/themes/mod.rs index e6463d47a..debe6dc38 100644 --- a/packages/tui-rs/src/themes/mod.rs +++ b/packages/tui-rs/src/themes/mod.rs @@ -865,13 +865,18 @@ mod ui_theme_tests { fn built_in_control_surfaces_are_opaque_including_light_on_dark_terminals() { for theme in [dark_theme(), light_theme(), high_contrast_theme()] { let ui = theme.ui_theme(); - assert_ne!( - ui.surface, - Color::Reset, - "{} must not inherit an incompatible terminal background", - theme.name - ); - assert_ne!(ui.surface, ui.text); + if palette::color_level() == palette::ColorLevel::None { + assert_eq!(ui.surface, Color::Reset); + assert_eq!(ui.text, Color::Reset); + } else { + assert_ne!( + ui.surface, + Color::Reset, + "{} must not inherit an incompatible terminal background", + theme.name + ); + assert_ne!(ui.surface, ui.text); + } assert_eq!( ui.surface, parse_color(&theme.colors.assistant_message_bg) diff --git a/packages/tui-rs/src/update_cli.rs b/packages/tui-rs/src/update_cli.rs index 7b3c4101d..cb83744b3 100644 --- a/packages/tui-rs/src/update_cli.rs +++ b/packages/tui-rs/src/update_cli.rs @@ -2147,36 +2147,114 @@ fn persist_rollback_suppression(path: &Path, version: &str) -> Result<()> { write_startup_state(path, &state) } -fn startup_update_mode() -> &'static str { - let mode = env::var("MAESTRO_AUTO_UPDATE") - .or_else(|_| env::var("MAESTRO_STARTUP_UPDATE")) - .unwrap_or_default() - .trim() - .to_ascii_lowercase(); +/// What an installed interactive launch does about a newer release. +/// +/// Read from `MAESTRO_AUTO_UPDATE`, then `MAESTRO_STARTUP_UPDATE`. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum StartupUpdateMode { + /// `0`, `false`, `off`, `skip`, or `disabled`: no startup check at all. + Off, + /// The default, and every value other than the off and apply spellings. + /// The TUI runs a bounded check after its first frame and shows a notice + /// in the transcript when a newer release exists. Nothing is installed. + Notice, + /// `apply` or `install`: check, install, and restart before the terminal + /// is set up. This is the only path that still delays the first frame, + /// and it is opt-in. + ApplyBeforeTui, +} + +pub(crate) fn startup_update_mode_from(value: Option<&str>) -> StartupUpdateMode { + let mode = value.unwrap_or_default().trim().to_ascii_lowercase(); if matches!(mode.as_str(), "0" | "false" | "off" | "skip" | "disabled") { - "off" - } else if matches!(mode.as_str(), "check" | "notice" | "notify") { - "check" + StartupUpdateMode::Off + } else if matches!(mode.as_str(), "apply" | "install") { + StartupUpdateMode::ApplyBeforeTui } else { - "apply" + StartupUpdateMode::Notice } } +fn startup_update_mode() -> StartupUpdateMode { + let value = env::var("MAESTRO_AUTO_UPDATE") + .or_else(|_| env::var("MAESTRO_STARTUP_UPDATE")) + .ok(); + startup_update_mode_from(value.as_deref()) +} + fn startup_update_enabled() -> bool { env::var_os("MAESTRO_SKIP_STARTUP_UPDATE").is_none() && env::var_os("CI").is_none() && env::var("NODE_ENV").ok().as_deref() != Some("test") - && startup_update_mode() != "off" + && startup_update_mode() != StartupUpdateMode::Off && std::io::stdin().is_terminal() && std::io::stdout().is_terminal() } -/// Best-effort update of an installed interactive Maestro before the TUI starts. +/// One bounded startup check against the trusted release sources. /// -/// Returns the restarted process exit code after a successful update. All check, -/// state, and install failures fail open so an unavailable update service can -/// never prevent Maestro from starting. -pub async fn run_startup_update(raw_args: &[std::ffi::OsString]) -> Option { +/// `total_timeout` is the existing `MAESTRO_STARTUP_UPDATE_TIMEOUT_MS` bound +/// (default 350 ms); it applies to the whole check, and each source gets an +/// equal share of it. +#[derive(Debug, Clone)] +pub(crate) struct StartupUpdateCheckRequest { + current: String, + urls: Vec, + channel: UpdateChannel, + total_timeout: Duration, + source_timeout: Duration, +} + +impl StartupUpdateCheckRequest { + fn new(current: String, urls: Vec, channel: UpdateChannel) -> Self { + let total_timeout = env_duration( + "MAESTRO_STARTUP_UPDATE_TIMEOUT_MS", + DEFAULT_STARTUP_CHECK_TIMEOUT, + false, + ); + Self::with_timeout(current, urls, channel, total_timeout) + } + + fn with_timeout( + current: String, + urls: Vec, + channel: UpdateChannel, + total_timeout: Duration, + ) -> Self { + let source_count = u128::try_from(urls.len().max(1)).unwrap_or(1); + let source_timeout = Duration::from_millis( + u64::try_from((total_timeout.as_millis() / source_count).max(1)).unwrap_or(1), + ); + Self { + current, + urls, + channel, + total_timeout, + source_timeout, + } + } + + /// The check, bounded by `total_timeout`. `None` on timeout, failure, or + /// when the installed version is already current. + async fn run(self) -> Option { + let check = tokio::time::timeout( + self.total_timeout, + check_for_update_urls_with_timeout( + &self.current, + self.urls, + self.source_timeout, + self.channel, + ), + ) + .await + .ok()?; + (check.status == "available" && check.latest_version.is_some()).then_some(check) + } +} + +/// The gate every startup check passes: machine policy, environment opt-outs, +/// a real terminal, and an installed (package or release) binary. +fn startup_update_check_request() -> Option<(InstallContext, StartupUpdateCheckRequest)> { if maestro_local_host::safety::vendor_network_disabled() { return None; } @@ -2187,38 +2265,70 @@ pub async fn run_startup_update(raw_args: &[std::ffi::OsString]) -> Option let channel = UpdateChannel::from_environment().ok()?; let current = current_version(); let urls = trusted_startup_update_urls(&context, channel); - let total_timeout = env_duration( - "MAESTRO_STARTUP_UPDATE_TIMEOUT_MS", - DEFAULT_STARTUP_CHECK_TIMEOUT, - false, - ); - let source_count = u128::try_from(urls.len().max(1)).unwrap_or(1); - let source_timeout = Duration::from_millis( - u64::try_from((total_timeout.as_millis() / source_count).max(1)).unwrap_or(1), - ); - let check = match tokio::time::timeout( - total_timeout, - check_for_update_urls_with_timeout(¤t, urls, source_timeout, channel), - ) + let request = StartupUpdateCheckRequest::new(current, urls, channel); + Some((context, request)) +} + +/// A newer release the TUI should mention. Carries raw values; the UI +/// formats them with its own locale. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct StartupUpdateNotice { + pub(crate) latest: String, + pub(crate) current: String, +} + +/// Background startup check for the interactive TUI. +/// +/// Runs after the first frame; the caller never awaits it on the startup +/// path. Returns `None` in every case except a completed check that found a +/// newer release while the mode is [`StartupUpdateMode::Notice`]. The +/// [`StartupUpdateMode::ApplyBeforeTui`] path already ran (or restarted the +/// process) before the terminal was set up, so it is not repeated here. +pub(crate) async fn startup_update_notice() -> Option { + // The gate reads the environment, checks the tty, and resolves the + // install layout from disk. Keep those syscalls off the async worker. + let request = tokio::task::spawn_blocking(|| { + if startup_update_mode() != StartupUpdateMode::Notice { + return None; + } + startup_update_check_request().map(|(_, request)| request) + }) .await - { - Ok(check) if check.status != "failed" => check, - _ => return None, - }; - if check.status != "available" { + .ok()??; + startup_update_notice_for(request).await +} + +async fn startup_update_notice_for( + request: StartupUpdateCheckRequest, +) -> Option { + let current = request.current.clone(); + let check = request.run().await?; + Some(StartupUpdateNotice { + latest: check.latest_version?, + current, + }) +} + +/// Explicit `MAESTRO_AUTO_UPDATE=apply`: update an installed interactive +/// Maestro before the TUI starts. +/// +/// This is the only startup path that waits on the network before terminal +/// setup, and it is bounded by `MAESTRO_STARTUP_UPDATE_TIMEOUT_MS`. Every other +/// mode returns `None` immediately; the default notice check runs inside the +/// TUI after the first frame (see [`startup_update_notice`]). +/// +/// Returns the restarted process exit code after a successful update. All check, +/// state, and install failures fail open so an unavailable update service can +/// never prevent Maestro from starting. +pub async fn run_startup_update(raw_args: &[std::ffi::OsString]) -> Option { + if startup_update_mode() != StartupUpdateMode::ApplyBeforeTui { return None; } + let (context, request) = startup_update_check_request()?; + let current = request.current.clone(); + let channel = request.channel; + let check = request.run().await?; let latest = check.latest_version.as_deref()?; - if startup_update_mode() == "check" { - eprintln!( - "{}", - crate::localization::cli_locale().format( - "Deixic Code {0} is available (current {1}); run `deixic-code update`.", - &[(latest).to_string(), (current).clone()] - ) - ); - return None; - } let state_path = startup_state_path_for(Some(&context))?; let update_lock = try_acquire_startup_update_lock(&state_path).ok()??; @@ -3672,6 +3782,136 @@ mod tests { assert!(legacy_package_install_context(&local).is_none()); } + #[test] + fn startup_update_mode_defaults_to_notice_and_only_apply_installs_before_the_tui() { + assert_eq!(startup_update_mode_from(None), StartupUpdateMode::Notice); + for value in ["", "check", "notice", "notify", "1", "true", "on", "yes"] { + assert_eq!( + startup_update_mode_from(Some(value)), + StartupUpdateMode::Notice, + "{value:?}" + ); + } + for value in ["0", "false", "off", "skip", "disabled", " OFF "] { + assert_eq!( + startup_update_mode_from(Some(value)), + StartupUpdateMode::Off, + "{value:?}" + ); + } + for value in ["apply", "install", "Apply", " install "] { + assert_eq!( + startup_update_mode_from(Some(value)), + StartupUpdateMode::ApplyBeforeTui, + "{value:?}" + ); + } + } + + fn startup_request(url: String, timeout: Duration) -> StartupUpdateCheckRequest { + StartupUpdateCheckRequest::with_timeout( + "0.10.52".to_owned(), + vec![url], + UpdateChannel::Stable, + timeout, + ) + } + + #[tokio::test] + async fn startup_update_notice_reports_a_newer_release() { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind test server"); + let address = listener.local_addr().expect("server address"); + let server = thread::spawn(move || { + let (mut stream, _) = listener.accept().expect("accept request"); + let mut request = [0_u8; 1024]; + let _ = stream.read(&mut request).expect("read request"); + let body = r#"{"version":"0.11.0"}"#; + write!( + stream, + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), + body + ) + .expect("write response"); + }); + + let request = startup_request( + format!("http://{address}/version.json"), + Duration::from_secs(5), + ); + let notice = startup_update_notice_for(request).await; + server.join().expect("join server"); + assert_eq!( + notice, + Some(StartupUpdateNotice { + latest: "0.11.0".to_owned(), + current: "0.10.52".to_owned(), + }) + ); + } + + #[tokio::test] + async fn startup_update_notice_is_silent_when_the_installed_version_is_current() { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind test server"); + let address = listener.local_addr().expect("server address"); + let server = thread::spawn(move || { + let (mut stream, _) = listener.accept().expect("accept request"); + let mut request = [0_u8; 1024]; + let _ = stream.read(&mut request).expect("read request"); + let body = r#"{"version":"0.10.52"}"#; + write!( + stream, + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), + body + ) + .expect("write response"); + }); + + let request = startup_request( + format!("http://{address}/version.json"), + Duration::from_secs(5), + ); + let notice = startup_update_notice_for(request).await; + server.join().expect("join server"); + assert_eq!(notice, None); + } + + /// The startup bound (`MAESTRO_STARTUP_UPDATE_TIMEOUT_MS`, default 350 ms) + /// still caps the whole check when a source accepts and never answers. + #[tokio::test] + async fn startup_update_notice_is_bounded_when_the_source_never_answers() { + use std::time::Instant; + + let listener = TcpListener::bind("127.0.0.1:0").expect("bind holding server"); + let address = listener.local_addr().expect("holding server address"); + let (done_tx, done_rx) = std::sync::mpsc::channel::<()>(); + let server = thread::spawn(move || { + let (stream, _) = listener.accept().expect("accept request"); + // Hold the connection open, unanswered, until the test finishes. + let _ = done_rx.recv(); + drop(stream); + }); + + let bound = Duration::from_millis(200); + let request = startup_request(format!("http://{address}/version.json"), bound); + let started = Instant::now(); + let notice = startup_update_notice_for(request).await; + let elapsed = started.elapsed(); + let _ = done_tx.send(()); + server.join().expect("join holding server"); + + assert_eq!(notice, None); + assert!( + elapsed >= bound, + "the check returned before its {bound:?} bound: {elapsed:?}" + ); + assert!( + elapsed < Duration::from_secs(5), + "the check ran past its {bound:?} bound: {elapsed:?}" + ); + } + #[test] fn startup_update_lock_is_nonblocking_and_released_on_drop() { let temporary = tempfile::tempdir().expect("temp directory"); diff --git a/packages/tui-rs/tests/pty_e2e.rs b/packages/tui-rs/tests/pty_e2e.rs index 8eb424833..2b79c164e 100644 --- a/packages/tui-rs/tests/pty_e2e.rs +++ b/packages/tui-rs/tests/pty_e2e.rs @@ -367,6 +367,28 @@ impl PtySession { extra_env: &[(&str, &str)], columns: u16, ) -> Self { + Self::launch( + mock, + workdir, + PtyLaunch { + executable: std::env::var_os("CARGO_BIN_EXE_maestro-tui") + .expect("Cargo must provide the maestro-tui integration-test binary"), + args, + extra_env, + removed_env: &[], + columns, + }, + ) + } + + fn launch(mock: &MockOpenAiServer, workdir: &std::path::Path, launch: PtyLaunch<'_>) -> Self { + let PtyLaunch { + executable, + args, + extra_env, + removed_env, + columns, + } = launch; let pty_system = native_pty_system(); let pair = pty_system .openpty(PtySize { @@ -384,14 +406,12 @@ impl PtySession { std::fs::write(&preferences, r#"{"onboardingSeen":true}"#).unwrap(); } - let mut command = CommandBuilder::new( - std::env::var_os("CARGO_BIN_EXE_maestro-tui") - .expect("Cargo must provide the maestro-tui integration-test binary"), - ); + let mut command = CommandBuilder::new(executable); command.args(args); command.cwd(workdir); - // CommandBuilder starts from an empty environment; pass through only - // what the child needs and pin everything else explicitly. + // CommandBuilder copies the parent environment (portable-pty's + // `get_base_env`). Pin what the child needs explicitly; a scenario + // that must see a variable unset lists it in `removed_env`. for key in ["PATH", "LANG", "USER", "LOGNAME", "TMPDIR"] { if let Ok(value) = std::env::var(key) { command.env(key, value); @@ -425,6 +445,9 @@ impl PtySession { for (name, value) in extra_env { command.env(name, value); } + for name in removed_env { + command.env_remove(name); + } let child = pair .slave @@ -650,6 +673,22 @@ impl PtySession { .any(|(pid, _, args)| args.contains(needle) && is_descendant(&table, *pid, root)) } + /// Ask the TUI to quit (Ctrl+D) and report whether it exited on its own + /// within `timeout`. `Drop` kills a child that did not. + fn quit_and_wait(&mut self, timeout: Duration) -> Option { + self.send_bytes(b"\x04"); + let deadline = Instant::now() + timeout; + loop { + match self.child.try_wait() { + Ok(Some(status)) => return Some(status), + Ok(None) if Instant::now() < deadline => { + std::thread::sleep(Duration::from_millis(50)); + } + _ => return None, + } + } + } + /// Ask the TUI to quit (Ctrl+D), then fall back to killing the child. fn shutdown(mut self) { self.send_bytes(b"\x04"); @@ -1096,6 +1135,282 @@ fn pty_confirmed_rewind_preserves_earlier_turn_and_persists_child_lineage() { /// the app's position reads. static PTY_TEST_SERIAL: Mutex<()> = Mutex::new(()); +/// Everything `PtySession::launch` needs beyond the harness defaults. +struct PtyLaunch<'a> { + executable: std::ffi::OsString, + args: &'a [&'a str], + extra_env: &'a [(&'a str, &'a str)], + /// Removed after the harness defaults and `extra_env`, so the child sees + /// these unset even when the test runner's environment defines them. + removed_env: &'a [&'a str], + columns: u16, +} + +/// A TCP listener that accepts every connection and never writes a byte. +/// Stands in for a release source that is up but does not answer. +struct HoldingServer { + url: String, + accepted: Arc, + _held: Arc>>, +} + +impl HoldingServer { + fn start() -> Self { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind holding server"); + let address = listener.local_addr().expect("holding server address"); + let accepted = Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let held: Arc>> = Arc::new(Mutex::new(Vec::new())); + let thread_accepted = Arc::clone(&accepted); + let thread_held = Arc::clone(&held); + std::thread::Builder::new() + .name("pty-e2e-holding-server".to_owned()) + .spawn(move || { + for stream in listener.incoming() { + let Ok(stream) = stream else { + break; + }; + thread_held + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(stream); + thread_accepted.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + } + }) + .expect("spawn holding server thread"); + Self { + url: format!("http://{address}/version.json"), + accepted, + _held: held, + } + } + + fn accepted(&self) -> usize { + self.accepted.load(std::sync::atomic::Ordering::SeqCst) + } + + fn wait_for_connection(&self, timeout: Duration) { + let deadline = Instant::now() + timeout; + while self.accepted() == 0 { + assert!( + Instant::now() < deadline, + "the release check never connected to the holding server within {timeout:?}" + ); + std::thread::sleep(Duration::from_millis(25)); + } + } +} + +/// Serve `{"version": }` to every request, the same shape the +/// update client's unit tests use. Returns the URL and a request counter. +fn start_version_server(version: &str) -> (String, Arc) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind version server"); + let address = listener.local_addr().expect("version server address"); + let body = format!(r#"{{"version":"{version}"}}"#); + let requests = Arc::new(std::sync::atomic::AtomicUsize::new(0)); + let thread_requests = Arc::clone(&requests); + std::thread::Builder::new() + .name("pty-e2e-version-server".to_owned()) + .spawn(move || { + for stream in listener.incoming() { + let Ok(mut stream) = stream else { + break; + }; + let mut request = [0_u8; 4096]; + let _ = stream.read(&mut request); + let _ = write!( + stream, + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), + body + ); + let _ = stream.flush(); + thread_requests.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + } + }) + .expect("spawn version server thread"); + (format!("http://{address}/version.json"), requests) +} + +/// Lay the test binary out as a global npm package install so +/// `update_cli::install_context` recognizes it. Returns the package root and +/// the executable to launch. +fn stage_fake_package_install( + workdir: &std::path::Path, +) -> (std::path::PathBuf, std::path::PathBuf) { + let source = std::env::var_os("CARGO_BIN_EXE_maestro-tui") + .expect("Cargo must provide the maestro-tui integration-test binary"); + let package_root = workdir + .join("lib") + .join("node_modules") + .join("@evalops") + .join("deixic-code"); + let vendor_dir = package_root.join("vendor").join("maestro").join("pty-e2e"); + std::fs::create_dir_all(&vendor_dir).expect("create vendor dir"); + let executable = vendor_dir.join("maestro"); + if std::fs::hard_link(&source, &executable).is_err() { + std::fs::copy(&source, &executable).expect("copy test binary into the fake package"); + } + let bin_dir = package_root.join("bin"); + std::fs::create_dir_all(&bin_dir).expect("create bin dir"); + std::fs::write(bin_dir.join("maestro"), "#!/bin/sh\nexit 0\n").expect("write launcher"); + (package_root, executable) +} + +/// Variables that would switch the startup check off or reroute it. Removed +/// so the child runs the default (notice) mode regardless of the runner's +/// environment. +const STARTUP_UPDATE_ENV_TO_REMOVE: &[&str] = &[ + "CI", + "NODE_ENV", + "MAESTRO_AUTO_UPDATE", + "MAESTRO_STARTUP_UPDATE", + "MAESTRO_SKIP_STARTUP_UPDATE", + "MAESTRO_UPDATE_URLS", + "MAESTRO_UPDATE_CHANNEL", + "MAESTRO_MANAGED_POLICY_PATH", +]; + +fn startup_update_launch<'a>( + executable: &std::path::Path, + package_root: &'a str, + update_url: &'a str, + timeout_ms: &'a str, +) -> (std::ffi::OsString, Vec<(&'a str, &'a str)>) { + ( + executable.as_os_str().to_owned(), + vec![ + ("MAESTRO_INSTALL_METHOD", "package"), + ("MAESTRO_PACKAGE_NAME", "@evalops/deixic-code"), + ("MAESTRO_PACKAGE_ROOT", package_root), + ("MAESTRO_UPDATE_URL", update_url), + ("MAESTRO_VERSION", "1.0.0"), + ("MAESTRO_STARTUP_UPDATE_TIMEOUT_MS", timeout_ms), + ], + ) +} + +/// Terminal setup and the first frame must not wait for the release check. +/// +/// The check is pointed at a socket that accepts and never answers, with the +/// bound raised to 60 s. Before the fix the check ran before `terminal::init`, +/// so the "Starting…" composer could not appear inside the 5 s wait. After +/// the fix the composer appears, the check connects afterwards and stays +/// pending, and Ctrl+D exits without waiting for the bound. +#[test] +fn pty_first_frame_renders_while_the_startup_update_check_is_pending() { + let _serial = PTY_TEST_SERIAL.lock().unwrap_or_else(|e| e.into_inner()); + let (release, gate) = std::sync::mpsc::channel(); + let mut mock = MockOpenAiServer::start(Vec::new()); + mock.managed_setup_base_url = start_mock_managed_setup_server_with_gate(Some(gate)); + let workdir = tempfile::tempdir().expect("temp workdir"); + let (package_root, executable) = stage_fake_package_install(workdir.path()); + let package_root = package_root.to_string_lossy().into_owned(); + let holding = HoldingServer::start(); + let (executable, extra_env) = + startup_update_launch(&executable, &package_root, &holding.url, "60000"); + + let started = Instant::now(); + let mut session = PtySession::launch( + &mock, + workdir.path(), + PtyLaunch { + executable, + args: &["--model", "openai/gpt-4o"], + extra_env: &extra_env, + removed_env: STARTUP_UPDATE_ENV_TO_REMOVE, + columns: 120, + }, + ); + session.wait_for_text("Starting…", Duration::from_secs(5)); + let first_frame = started.elapsed(); + assert_eq!( + holding.accepted(), + 0, + "the release check must not run before the first frame" + ); + + // Policy release lets the loop start; the check starts after its first + // frame and then stays pending against the silent socket. + release.send(()).unwrap(); + // Agent spawn runs between the first frame and the loop; under load it + // takes several seconds, so the connection wait uses the startup ceiling. + holding.wait_for_connection(READY_TIMEOUT); + let check_connected = started.elapsed(); + + let quit_started = Instant::now(); + let status = session.quit_and_wait(Duration::from_secs(10)); + let quit_elapsed = quit_started.elapsed(); + assert!( + status.as_ref().is_some_and(|status| status.success()), + "quit must not wait for the pending 60 s release check (exit status: {status:?})" + ); + eprintln!( + "startup composer at {first_frame:?}; release check connected at {check_connected:?} and was still pending; quit took {quit_elapsed:?}" + ); +} + +/// A newer release reaches the transcript as a notice after the first frame, +/// through the loop wake, without any install. +#[test] +fn pty_startup_update_notice_appears_after_the_first_frame() { + let _serial = PTY_TEST_SERIAL.lock().unwrap_or_else(|e| e.into_inner()); + let (release, gate) = std::sync::mpsc::channel(); + let mut mock = MockOpenAiServer::start(Vec::new()); + mock.managed_setup_base_url = start_mock_managed_setup_server_with_gate(Some(gate)); + let workdir = tempfile::tempdir().expect("temp workdir"); + let (package_root, executable) = stage_fake_package_install(workdir.path()); + let package_root = package_root.to_string_lossy().into_owned(); + let (update_url, version_requests) = start_version_server("99.0.0"); + let (executable, extra_env) = + startup_update_launch(&executable, &package_root, &update_url, "5000"); + + let started = Instant::now(); + let mut session = PtySession::launch( + &mock, + workdir.path(), + PtyLaunch { + executable, + args: &["--model", "openai/gpt-4o"], + extra_env: &extra_env, + removed_env: STARTUP_UPDATE_ENV_TO_REMOVE, + columns: 120, + }, + ); + session.wait_for_text("Starting…", Duration::from_secs(5)); + assert_eq!( + version_requests.load(std::sync::atomic::Ordering::SeqCst), + 0, + "the release check must not run before the first frame" + ); + release.send(()).unwrap(); + let request_deadline = Instant::now() + READY_TIMEOUT; + while version_requests.load(std::sync::atomic::Ordering::SeqCst) == 0 { + assert!( + Instant::now() < request_deadline, + "the release check never requested version.json" + ); + std::thread::sleep(Duration::from_millis(25)); + } + let requested_at = started.elapsed(); + // The transcript renders the code span without its backticks. The notice + // lands once the loop starts, after agent spawn, so use the startup ceiling. + session.wait_for_wrapped_text( + "Deixic Code 99.0.0 is available (current 1.0.0); run deixic-code update.", + READY_TIMEOUT, + ); + let notice_at = started.elapsed(); + assert_eq!( + mock.request_count(), + 0, + "the notice must not trigger a model request" + ); + let status = session.quit_and_wait(Duration::from_secs(10)); + assert!(status.is_some_and(|status| status.success())); + eprintln!( + "release check requested version.json at {requested_at:?}; notice visible at {notice_at:?}" + ); +} + fn bad_gateway_turn() -> ScriptedTurn { ScriptedTurn { status: "502 Bad Gateway", diff --git a/proto/maestro/v1/headless.proto b/proto/maestro/v1/headless.proto index 28d8f9f80..2406722b3 100644 --- a/proto/maestro/v1/headless.proto +++ b/proto/maestro/v1/headless.proto @@ -557,6 +557,10 @@ message ResponseChunkMessage { bool is_thinking = 3; } +message AssistantTextObservedMessage { + string response_id = 1; +} + message ResponseEndMessage { string response_id = 1; ResponseUsage usage = 2; @@ -917,5 +921,6 @@ message FromAgentEnvelope { DelegationEventMessage delegation_event = 33; WorkspaceCapabilitySetAppliedMessage workspace_capability_set_applied = 34; ManagedAuthorizationRequestMessage managed_authorization_request = 35; + AssistantTextObservedMessage assistant_text_observed = 36; } } diff --git a/proto/maestro/v1/protocol-compatibility-manifest.json b/proto/maestro/v1/protocol-compatibility-manifest.json index f89bc529d..f9e9814d2 100644 --- a/proto/maestro/v1/protocol-compatibility-manifest.json +++ b/proto/maestro/v1/protocol-compatibility-manifest.json @@ -5,7 +5,7 @@ "canonicalization": "rfc8785", "encoding": "utf-8" }, - "compatibilityDigest": "sha256:f5e8e7ac8d9e165b854eeae3baa063030fa269c8f847c86490c63d263b5426ce", + "compatibilityDigest": "sha256:1f7aa547ac86f4008bab638ff01321efb99a2274bcf8d823eea3a2cb2763db56", "buildIdentity": { "sourceSha": null, "buildDigest": null, @@ -37,8 +37,8 @@ "2026-08-08" ], "schemaDeclaredVersion": "2026-08-08", - "schemaContractDigest": "sha256:6f861ccdc76c8f85150a6a2aeceb99a0da5847f04ffd3d0b39011d94c81dbc53", - "runtimeContractDigest": "sha256:cea8cf73427930169aa188ee8eca2ccee0bb020f7d379150b44dbd0c6de0952a", + "schemaContractDigest": "sha256:88a8ce00ee666bde7d519eacf552f2c22608b553446f3b5a718ca66ce696894f", + "runtimeContractDigest": "sha256:2883475a1a847bdfc2cd9585a88a657ce788380951f2de1593ae6eaa5d97716e", "messages": { "toRuntime": [ "managed_authorization_result", @@ -78,6 +78,7 @@ "workspace_capability_set_applied", "ready", "response_start", + "assistant_text_observed", "response_chunk", "response_end", "turn_completed", @@ -201,7 +202,7 @@ }, "runtime": { "schemaVersion": "evalops.maestro.headless-protocol.v1", - "contractDigest": "sha256:43703fc50891d32f11e169729feb0a9eed7b92546a53234c758ba451dc1f4e53", + "contractDigest": "sha256:8353f3612f7e7f4f7f5c55020b2c255f7ad52b57417a027ae66ceb657c402ba8", "receipt": { "schemaVersion": "evalops.maestro.runtime-receipt.v1", "sourceDigest": "sha256:8fbcb58770dccabaf0d87b2532121addafd35e04f5a09f9dffdd2dea8897b921",