diff --git a/.repository-projection.json b/.repository-projection.json index 05ee0bf73..c45b486cc 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "c0fac5f2a1a1bddd988acc6c0c564d984c6a6f59", + "sourceSha": "16426c3c3b3dc59f6817c9a2836fd714d17a5533", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "61666aceff72f31fbe3dfd44477540e703aa4b85", + "priorProjectedBase": "5bfbbeb575801b605b42755a5a1154849c56a75a", "definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f", - "toolDigest": "ab19140af8e449c04288cb7c52e992982833e3882ddce59d22314d971feaf946", - "contentDigest": "cb9bcb8eeb4b7107c0dbf55cedd0d8ad81db6ad3edc8ff3c2c53587d62015449", + "toolDigest": "c244d99199a7ae3eb8ff644a99462163c23b0bb6a83ef50af01efbdca0b81d04", + "contentDigest": "42be9fd04b572f4a3cf4bb0854bcac5ae498a5d48dc99cf8d129ade8a4ab0536", "publicationEligible": true } diff --git a/CHANGELOG.md b/CHANGELOG.md index e57d445d5..be9196caa 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -51,6 +51,132 @@ versioning when releases are cut. and keep scheduled public runs inert so public publishing stays downstream of the internal source-of-truth release. +## [0.10.108] - 2026-09-29 + +### Fixed + +- Snapshot the committed Rust workspace and `dex-loop` inputs when syncing the public source tree (#11355). +- Request public release credentials for `dx-corp/code`, the actual publishing repository (#11362). + +## [0.10.107] - 2026-09-29 + +### Fixed + +- Include the exact `dex-loop` source in public release mirrors so the published source builds on its own (#11320). +- Check macOS release signing material before native compilation and resolve it from the runner home (#11311). + +## [0.10.106] - 2026-09-29 + +### Added + +- Persist provider reasoning with the model step (#11274). + +### Changed + +- Main-heavy capacity is 14 max / 2 min (#11278). + +### Fixed + +- Route cloud turns through the public protocol so signed releases pass the artifact guard (#11289). +- Use stable key sorting for find results (#11277). +- Isolate crash replay database ports (#11276). +- Send the guardian review to Messages on Anthropic (#11275). + +## [0.10.105] - 2026-09-29 + +### Added + +- Run Claude with adaptive thinking over /v1/messages (#11262). +- Workload-token-only mode when no login provider is enabled (#11252). +- Dex.find estate and connector scopes over tool-execution (#11250). +- Dex.make capability over the platform-api artifact owner (#11216). +- Platform-api dex.find/dex.read backend for history, skills and forms (#11246). +- Render one-loop tool progress without tool names (#11239). +- Add dex.describe over the capability catalogs (#11229). +- Dex.find dispatcher and owner refs for dex.read (#11236). +- Compute over tool-execution, resuming its approval parks under Dex approval (#11222). +- Tell Dex to look up, cite, and make one artifact (#11219). +- Govern change and act by their arguments (#11213). +- Platform-api endpoint for one-loop platform tools, starting with conversation history (#11210). + +### Changed + +- Treat killed Linux zombies as stopped in the process cleanup test (#11268). +- Name the stub's recorded call type for clippy (#11263). +- One-loop e2e fake gateway speaks Responses (#11259). +- Dex e2e fake gateway speaks Responses (#11258). +- Run the client-tool one-loop test on the no-remote-tools wiring (#11254). +- Record failed optimization executions during qualification (#11255). +- Fix manual replays and record measured gate timings (#11253). +- One-loop fake gateway speaks Responses (#11251). +- Stage decision-packet evaluation and bootstrap by world (#11248). +- Run dex-harness-parity in the platform-api shard on the production wiring (#11209). +- Score evals and canaries on capabilities, not old tool names (#11247). +- Fix optimization Clippy and downstream Cargo lock freshness (#11244). + +### Fixed + +- Verify Mise with the runner platform digest so macOS release builds can proceed (#11266). +- Require an accepted turn on the dex artifact endpoint (#11265). +- Class compute calls per action (shell guardian, write approval) (#11261). +- Make the computer attach replay instead of 409ing on the next call (#11264). +- Honor dex.read range.offset for platform-api refs (#11260). +- Register dex.find in production; stop advertising dex.change/dex.act (#11257). +- Render unknown Dex event kinds safely (#11256). +- Serve Dex model steps through dex-model's Responses client (#11249). +- Preserve guarded background replies and prove transport parity (#11245). +- Cover every fuzz prompt boundary (#11243). +- Preserve retained computer classifier names (#11242). +- Make the model stream Send again after stored-output reads (#11237). + +## [0.10.104] - 2026-09-29 + +### Added + +- Compute over tool-execution, resuming its approval parks under Dex approval (#11222). +- Tell Dex to look up, cite, and make one artifact (#11219). +- Govern change and act by their arguments (#11213). +- Platform-api endpoint for one-loop platform tools, starting with conversation history (#11210). +- Add make dex-send ENV=staging and the Dex staging runbook (#11212). +- Web client tools through the shared Dex loop SDK (#11208). +- Gate and wire candidate character assets (#11204). +- Add make dex-trace for one-command Dex thread timelines (#11202). +- Carry server-derived call context to platform tools (#11197). +- Platform-api artifact endpoint for the one loop (#11200). +- Add draft character asset pipeline and pilot (#11187). +- Own workload identity with per-tenant least-scope tokens (#11148). + +### Changed + +- Add seeded high-volume one-loop fuzz harness (#11225). +- Expect clientTools in the pinned submit argument list (#11211). +- Add pre-cutover tool usage counts to the parity inventory (#11205). +- Fail when a managed migration file is unregistered (#11201). +- One-loop tool parity inventory (#11195). +- Settle preflight on verified descendants and log every poll (#11194). +- Refresh service-scoped Rust publisher targets (#11191). +- Remove unused ToolCatalog renderer (#11192). +- Remove retired MCP catalog attachment helpers (#11190). +- Remove source-coupled Make harness assertions (#11189). +- Bump pillow (#11188). +- Remove unreachable swarm checkpoint persistence (#11186). + +### Fixed + +- Recover managed gateway stream timeouts after partial thinking, with bounded turn retry and request diagnostics (#11178). +- Namespace the change and act capability names (#11224). +- The model reads stored tool output (#11226). +- Compile from_send and give checkpoint tests a token (#11217). +- Preserve chat SSE errors across byte boundaries (#11223). +- Validate secrets and URLs at startup; probe dex-runtime credential for /readyz (#11220). +- Preserve streaming without content guardrail rules (#11215). +- Bind PNG cache to Pillow runtime (#11218). +- Classify the reqwest error before without_url consumes it (#11214). +- Mark retired coding_* capabilities non-executable in the registry (#11207). +- Register dex-runtime as its own tool-execute caller (#11203). +- Stop 503ing runtime dispatch for guardrailed workspaces (#11196). +- Align synthetic contracts with current Dex loop (#11198). + ## [0.10.103] - 2026-09-28 ### Changed diff --git a/Cargo.lock b/Cargo.lock index 723cf91ea..40b989ec7 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2004,6 +2004,20 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "dex-loop" +version = "0.1.0" +dependencies = [ + "futures-util", + "hex", + "serde", + "serde_json", + "sha2 0.10.9", + "thiserror 2.0.20", + "tokio", + "tokio-util", +] + [[package]] name = "diff" version = "0.1.13" @@ -3738,7 +3752,7 @@ checksum = "dae608c151f68243f2b000364e1f7b186d9c29845f7d2d85bd31b9ad77ad552b" [[package]] name = "maestro" -version = "0.10.103" +version = "0.10.108" dependencies = [ "anyhow", "ctor", @@ -3820,6 +3834,23 @@ dependencies = [ "tiktoken-rs", ] +[[package]] +name = "maestro-dex-host" +version = "0.1.0" +dependencies = [ + "anyhow", + "dex-loop", + "fd-lock", + "futures-util", + "maestro-ai", + "serde", + "serde_json", + "tempfile", + "thiserror 2.0.20", + "tokio", + "tokio-stream", +] + [[package]] name = "maestro-execpolicy" version = "0.1.0" @@ -4115,6 +4146,7 @@ dependencies = [ "maestro-ai", "maestro-codex", "maestro-context", + "maestro-dex-host", "maestro-execpolicy", "maestro-interaction", "maestro-local-host", diff --git a/Cargo.toml b/Cargo.toml index a2d01e214..0f221ff36 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [workspace] resolver = "2" -exclude = ["examples/hooks/wasm-plugin", "vendor/*"] +exclude = ["examples/hooks/wasm-plugin", "vendor/*", "vendor/dex-loop"] members = [ "packages/maestro-rs", "packages/runtime-rs", @@ -25,6 +25,7 @@ members = [ "packages/ui-preview-rs", "packages/interaction-rs", "packages/ambient-agent-rs", + "packages/dex-host-rs", ] [workspace.dependencies] @@ -83,6 +84,7 @@ maestro-ui = { path = "packages/ui-rs" } maestro-presentation = { path = "packages/presentation-rs" } unicode-width = "0.2" maestro-interaction = { path = "packages/interaction-rs" } +maestro-dex-host = { path = "packages/dex-host-rs" } ratatui = { version = "0.30.2", default-features = false } # Shared by maestro-tui and maestro-runtime-gateway; keep both crates on the # same version here so a security bump (e.g. the lopdf RUSTSEC-2026-0187 diff --git a/Dockerfile b/Dockerfile index 968305244..f98ccbd40 100644 --- a/Dockerfile +++ b/Dockerfile @@ -3,12 +3,21 @@ WORKDIR /app FROM chef AS planner COPY Cargo.toml Cargo.lock ./ +# Keep the local Dex host available in the published Maestro binary. Its +# dex-loop dependency belongs to mono's separate Rust workspace; narrow that +# workspace to this crate inside the image while retaining its shared versions +# and lint configuration. +COPY --from=dex-loop-workspace /Cargo.toml /rust/Cargo.toml +COPY --from=dex-loop-workspace /crates/dex-loop /rust/crates/dex-loop +RUN sed -i '/^members = \[/,/^\]/c\members = ["crates/dex-loop"]' /rust/Cargo.toml \ + && sed -i '/^\[patch.crates-io\]/,/^\[workspace.package\]/c\[workspace.package]' /rust/Cargo.toml COPY vendor/zstd-0.13.3 ./vendor/zstd-0.13.3 COPY vendor/zstd-safe-7.2.4 ./vendor/zstd-safe-7.2.4 COPY vendor/zstd-sys-2.0.16+zstd.1.5.7 ./vendor/zstd-sys-2.0.16+zstd.1.5.7 COPY packages/execpolicy-rs ./packages/execpolicy-rs COPY packages/context-rs ./packages/context-rs COPY packages/tui-rs ./packages/tui-rs +COPY packages/dex-host-rs ./packages/dex-host-rs COPY packages/local-host-rs ./packages/local-host-rs COPY packages/sandbox-rs ./packages/sandbox-rs COPY packages/workspace-rs ./packages/workspace-rs @@ -33,6 +42,7 @@ RUN cargo chef prepare --recipe-path recipe.json FROM chef AS native COPY --from=planner /app/recipe.json recipe.json +COPY --from=planner /rust /rust COPY vendor/zstd-0.13.3 ./vendor/zstd-0.13.3 COPY vendor/zstd-safe-7.2.4 ./vendor/zstd-safe-7.2.4 COPY vendor/zstd-sys-2.0.16+zstd.1.5.7 ./vendor/zstd-sys-2.0.16+zstd.1.5.7 @@ -41,6 +51,7 @@ COPY Cargo.toml Cargo.lock ./ COPY packages/execpolicy-rs ./packages/execpolicy-rs COPY packages/context-rs ./packages/context-rs COPY packages/tui-rs ./packages/tui-rs +COPY packages/dex-host-rs ./packages/dex-host-rs COPY packages/local-host-rs ./packages/local-host-rs COPY packages/sandbox-rs ./packages/sandbox-rs COPY packages/workspace-rs ./packages/workspace-rs diff --git a/package-lock.json b/package-lock.json index 063d77093..65bbc1f2d 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,12 +1,12 @@ { "name": "@evalops/maestro", - "version": "0.10.103", + "version": "0.10.108", "lockfileVersion": 3, "requires": true, "packages": { "": { "name": "@evalops/maestro", - "version": "0.10.103", + "version": "0.10.108", "license": "BUSL-1.1", "bin": { "deixic-code": "bin/deixic-code", diff --git a/package.json b/package.json index a7e886c58..9f8bf5f55 100644 --- a/package.json +++ b/package.json @@ -1,7 +1,7 @@ { "name": "@evalops/deixic-code", "description": "Deixic Code — native Rust coding agent, CLI, TUI, and runtime gateway", - "version": "0.10.103", + "version": "0.10.108", "private": false, "type": "module", "bin": { diff --git a/packages/ai-rs/src/anthropic.rs b/packages/ai-rs/src/anthropic.rs index 23076e558..25a93eb46 100644 --- a/packages/ai-rs/src/anthropic.rs +++ b/packages/ai-rs/src/anthropic.rs @@ -304,7 +304,7 @@ impl AnthropicClient { // This task runs independently and sends events through the channel. tokio::spawn(async move { // Buffer for accumulating incomplete SSE events - let mut buffer = String::new(); + let mut buffer = Vec::new(); // Parser state - local to this stream, avoiding thread-local race conditions let mut parser_state = SseParserState::default(); @@ -312,14 +312,23 @@ impl AnthropicClient { while let Some(chunk) = stream.next().await { match chunk { Ok(bytes) => { - // Convert bytes to UTF-8 and append to buffer - buffer.push_str(&String::from_utf8_lossy(&bytes)); + buffer.extend_from_slice(&bytes); // Process all complete SSE events in the buffer // SSE events are delimited by double newlines: "\n\n" - while let Some(pos) = buffer.find("\n\n") { - let event_data = buffer[..pos].to_string(); - buffer = buffer[pos + 2..].to_string(); + loop { + let event_data = match super::sse::take_frame(&mut buffer) { + Ok(Some(event_data)) => event_data, + Ok(None) => break, + Err(error) => { + let _ = tx.send(StreamEvent::Error { + message: format!( + "Invalid UTF-8 in Anthropic SSE event: {error}" + ), + }); + return; + } + }; // Parse SSE event and send to receiver if let Some(event) = parse_sse_event(&event_data, &mut parser_state) { @@ -612,9 +621,9 @@ fn parse_sse_event(data: &str, state: &mut SseParserState) -> Option( ); break; } + let mut event = match event { + StreamEvent::ProviderError { + kind: ProviderStreamErrorKind::TransientProtocol, + mut message, + } if committed_content + && message + .starts_with("managed_gateway_stream_error: provider_stream_timeout:") => + { + // The native turn discards an incomplete assistant reply + // before retrying. Preserve this exact gateway timeout so + // the outer turn can retry even after thinking was streamed. + message.push_str("; "); + message.push_str(PARTIAL_CONTENT_STREAM_FAILURE_MARKER); + StreamEvent::ProviderError { + kind: ProviderStreamErrorKind::TransientProtocol, + message, + } + } + other => other, + }; let terminal_error = matches!( &event, StreamEvent::Error { .. } | StreamEvent::ProviderError { .. } ); if terminal_error { + let gateway_timeout = matches!( + &event, + StreamEvent::ProviderError { + kind: ProviderStreamErrorKind::TransientProtocol, + message, + } if message.starts_with("managed_gateway_stream_error: provider_stream_timeout:") + ); + if !gateway_request_id.is_empty() { + let message = match &mut event { + StreamEvent::Error { message } + | StreamEvent::ProviderError { message, .. } => message, + _ => unreachable!("terminal error has an error message"), + }; + message.push_str("; gateway request ID: "); + message.push_str(&gateway_request_id); + } let error_message_len = match &event { StreamEvent::Error { message } | StreamEvent::ProviderError { message, .. } => { message.len() @@ -1301,6 +1337,12 @@ async fn forward_stream_with_idle_policy_with_span( duration_ms = stream_started.elapsed().as_millis() as u64, events_forwarded, gateway_request_id = %gateway_request_id, + gateway_error_code = if gateway_timeout { "provider_stream_timeout" } else { "" }, + turn_retryable = matches!( + &event, + StreamEvent::ProviderError { kind, message } + if is_retryable_partial_content_stream_failure(*kind, message) + ), ); } let terminal = terminal_error || matches!(&event, StreamEvent::MessageStop { .. }); @@ -2000,7 +2042,7 @@ mod tests { #[cfg(test)] mod stream_idle_policy_tests { use super::*; - use crate::types::ProviderStreamErrorKind; + use crate::types::{ManagedGatewayReceipt, ProviderStreamErrorKind}; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; use std::time::Duration; @@ -2499,6 +2541,60 @@ mod stream_idle_policy_tests { )); } + #[tokio::test] + async fn gateway_timeout_after_partial_thinking_can_retry_the_model_turn() { + let (attempt_tx, attempt_rx) = mpsc::unbounded_channel(); + attempt_tx + .send(StreamEvent::ManagedGatewayReceipt(ManagedGatewayReceipt { + request_id: "gateway-request-123".to_owned(), + record_id: "record-123".to_owned(), + lineage_id: "lineage-123".to_owned(), + record_status: "accepted".to_owned(), + provider_prompt_sha256: None, + provider_tools_sha256: None, + provider_tool_count: None, + })) + .unwrap(); + attempt_tx + .send(StreamEvent::ThinkingDelta { + index: 0, + thinking: "working".to_owned(), + }) + .unwrap(); + attempt_tx + .send(StreamEvent::ProviderError { + kind: ProviderStreamErrorKind::TransientProtocol, + message: "managed_gateway_stream_error: provider_stream_timeout: Provider stream exceeded its attempt deadline".to_owned(), + }) + .unwrap(); + drop(attempt_tx); + let (tx, mut rx) = mpsc::unbounded_channel(); + forward_stream_with_idle_policy( + Some(attempt_rx), + || async { unreachable!("the stream owner must not replay partial output") }, + IDLE, + RETRIES, + tx, + ) + .await; + + let events = drain(&mut rx); + assert!( + matches!( + events.last(), + Some(StreamEvent::ProviderError { + kind: ProviderStreamErrorKind::TransientProtocol, + message, + }) if is_retryable_partial_content_stream_failure( + ProviderStreamErrorKind::TransientProtocol, + message, + ) && message.contains("provider_stream_timeout") + && message.contains("gateway request ID: gateway-request-123") + ), + "{events:?}" + ); + } + #[tokio::test(start_paused = true)] async fn transient_502_recovers_after_backoff_without_forwarding_failed_attempts() { let started = tokio::time::Instant::now(); diff --git a/packages/ai-rs/src/google.rs b/packages/ai-rs/src/google.rs index 7448f2379..2a28c8bb8 100644 --- a/packages/ai-rs/src/google.rs +++ b/packages/ai-rs/src/google.rs @@ -232,22 +232,21 @@ async fn stream_google_response( // Parse SSE stream let mut stream = response.bytes_stream(); - let mut buffer = String::new(); + let mut buffer = Vec::new(); let mut input_tokens = 0u64; let mut output_tokens = 0u64; let mut cache_read_tokens = None; while let Some(chunk) = stream.next().await { let chunk = chunk.context("Failed to read chunk")?; - buffer.push_str(&String::from_utf8_lossy(&chunk)); + buffer.extend_from_slice(&chunk); // Process complete SSE events - while let Some(pos) = buffer.find("\n\n") { - let event = buffer[..pos].to_string(); - buffer = buffer[pos + 2..].to_string(); - + while let Some(event) = + super::sse::take_frame(&mut buffer).context("Invalid UTF-8 in Google SSE event")? + { // Parse SSE data - if let Some(data) = event.strip_prefix("data: ") { + if let Some(data) = google_sse_data(&event) { if data.trim() == "[DONE]" { break; } @@ -315,6 +314,12 @@ async fn stream_google_response( Ok(()) } +fn google_sse_data(event: &str) -> Option<&str> { + event + .lines() + .find_map(|line| line.strip_prefix("data:").map(str::trim)) +} + /// Google prompt counts include cached tokens. Keep cache absence distinct from zero. pub(super) fn cache_usage_event(input: u64, output: u64, cached: Option) -> StreamEvent { StreamEvent::Usage { @@ -428,6 +433,14 @@ struct UsageMetadata { mod tests { use super::*; + #[test] + fn sse_data_accepts_event_header_and_no_space_after_colon() { + assert_eq!( + google_sse_data("event: update\ndata:{\"value\":\"東京\"}"), + Some("{\"value\":\"東京\"}") + ); + } + #[test] fn provider_cache_usage_distinguishes_hit_miss_and_unknown() { for (field, expected) in [(None, None), (Some(0), Some(0)), (Some(80), Some(80))] { diff --git a/packages/ai-rs/src/lib.rs b/packages/ai-rs/src/lib.rs index 970588495..5e0e4bf54 100644 --- a/packages/ai-rs/src/lib.rs +++ b/packages/ai-rs/src/lib.rs @@ -106,6 +106,7 @@ mod provider_matrix; mod providers; pub mod sanitize; mod scripted; +mod sse; mod transform; mod types; mod vertex; diff --git a/packages/ai-rs/src/openai.rs b/packages/ai-rs/src/openai.rs index af289f4aa..048c295ca 100644 --- a/packages/ai-rs/src/openai.rs +++ b/packages/ai-rs/src/openai.rs @@ -797,6 +797,19 @@ fn classify_error(error: &serde_json::Value) -> ApiError { } } +// Keep raw bytes until a complete line arrives. Decoding each HTTP chunk would +// replace a UTF-8 character split between chunks and corrupt a JSON error. +fn take_chat_sse_line(buffer: &mut Vec) -> Option> { + let end = buffer.iter().position(|byte| *byte == b'\n')?; + let line = buffer.drain(..=end).collect::>(); + Some(String::from_utf8(line)) +} + +fn chat_sse_data(line: &str) -> Option<&str> { + line.strip_prefix("data:") + .map(|data| data.strip_prefix(' ').unwrap_or(data)) +} + fn response_failed_error( response: Option<&serde_json::Value>, ) -> (ProviderStreamErrorKind, String) { @@ -2397,6 +2410,7 @@ impl OpenAiClient { // Spawn task to process SSE stream let model = config.model.clone(); + let managed_gateway = self.managed_gateway; let producer = if is_responses_api { // Use eventsource-stream for proper SSE parsing (Responses API) @@ -2752,6 +2766,21 @@ impl OpenAiClient { "response.failed" => { let (kind, message) = response_failed_error(event.response.as_ref()); + let message = if managed_gateway + && event + .response + .as_ref() + .and_then(|response| response.get("error")) + .and_then(|error| error.get("code")) + .and_then(serde_json::Value::as_str) + == Some("provider_stream_timeout") + { + format!( + "managed_gateway_stream_error: provider_stream_timeout: {message}" + ) + } else { + message + }; let _ = tx.send(StreamEvent::ProviderError { kind, message }); return; } @@ -2779,7 +2808,7 @@ impl OpenAiClient { // Chat Completions API - uses simpler line-based SSE let mut stream = response.bytes_stream(); tokio::spawn(async move { - let mut buffer = String::new(); + let mut buffer = Vec::new(); let mut message_id = String::new(); let mut current_tool_calls: Vec = Vec::new(); let mut content_started = false; @@ -2797,18 +2826,24 @@ impl OpenAiClient { }; match chunk { Ok(bytes) => { - buffer.push_str(&String::from_utf8_lossy(&bytes)); + buffer.extend_from_slice(&bytes); // Process complete SSE lines - while let Some(pos) = buffer.find('\n') { - let line = buffer[..pos].trim().to_string(); - buffer = buffer[pos + 1..].to_string(); + while let Some(line) = take_chat_sse_line(&mut buffer) { + let Ok(line) = line else { + let _ = tx.send(StreamEvent::ProviderError { + kind: ProviderStreamErrorKind::TransientProtocol, + message: "Chat stream contained invalid UTF-8".to_string(), + }); + return; + }; + let line = line.trim(); if line.is_empty() { continue; } - if line == "data: [DONE]" { + if chat_sse_data(line) == Some("[DONE]") { let Some(stop_reason) = terminal_reason else { let _ = tx.send(StreamEvent::ProviderError { kind: ProviderStreamErrorKind::TransientProtocol, @@ -2858,7 +2893,29 @@ impl OpenAiClient { return; } - if let Some(data) = line.strip_prefix("data: ") { + if let Some(data) = chat_sse_data(line) { + if let Ok(payload) = + serde_json::from_str::(data) + { + if let Some(error) = payload.get("error") { + let (kind, classified) = + response_failed_error(Some(&payload)); + let code = error + .get("code") + .and_then(serde_json::Value::as_str) + .unwrap_or("unknown"); + let source = if managed_gateway { + "managed_gateway_stream_error" + } else { + "chat_stream_error" + }; + let _ = tx.send(StreamEvent::ProviderError { + kind, + message: format!("{source}: {code}: {classified}"), + }); + return; + } + } if let Ok(chunk) = serde_json::from_str::(data) { if message_id.is_empty() { message_id = chunk.id.clone(); @@ -4451,6 +4508,142 @@ mod tests { .expect("managed stream should terminate") } + #[test] + fn chat_sse_lines_survive_seeded_network_chunk_boundaries() { + let frame = "event: error\r\ndata:{\"error\":{\"code\":\"provider_stream_timeout\",\"message\":\"réessayer 🌊\"}}\r\n\r\n"; + for seed in 0..512_u64 { + let mut state = seed + 1; + let mut buffer = Vec::new(); + let mut lines = Vec::new(); + let mut cursor = 0; + while cursor < frame.len() { + state ^= state << 13; + state ^= state >> 7; + state ^= state << 17; + let end = (cursor + (state as usize % 17) + 1).min(frame.len()); + buffer.extend_from_slice(&frame.as_bytes()[cursor..end]); + while let Some(line) = take_chat_sse_line(&mut buffer) { + lines.push(line.expect("valid UTF-8 after joining chunks")); + } + cursor = end; + } + assert!(buffer.is_empty(), "seed {seed}"); + assert_eq!(lines.concat(), frame, "seed {seed}"); + let payload = chat_sse_data(lines[1].trim()).expect("valid no-space data field"); + let error: serde_json::Value = serde_json::from_str(payload).expect("intact JSON"); + assert_eq!(error["error"]["code"], "provider_stream_timeout"); + assert_eq!(error["error"]["message"], "réessayer 🌊"); + } + } + + #[test] + fn chat_sse_invalid_utf8_fails_explicitly() { + let mut buffer = b"data: \xff\n".to_vec(); + assert!(take_chat_sse_line(&mut buffer).unwrap().is_err()); + } + + #[tokio::test] + async fn chat_gateway_timeout_accepts_no_space_after_data_colon() { + let sse = "event: error\ndata:{\"error\":{\"type\":\"server_error\",\"code\":\"provider_stream_timeout\",\"message\":\"attempt deadline\"}}\n\n"; + let (client, _request) = managed_gateway_test_client(sse, &managed_receipt_headers()); + let mut receiver = client + .stream_authorized_invocation( + &[], + &RequestConfig { + model: "evalops/openai/gpt-4o-mini".to_owned(), + ..Default::default() + }, + ) + .await + .expect("open chat stream"); + let events = tokio::time::timeout(std::time::Duration::from_secs(5), async { + let mut events = Vec::new(); + while let Some(event) = receiver.recv().await { + events.push(event); + } + events + }) + .await + .expect("gateway error stream should terminate"); + assert!( + matches!(events.last(), Some(StreamEvent::ProviderError { + kind: ProviderStreamErrorKind::TransientProtocol, + message, + }) if message.starts_with("managed_gateway_stream_error: provider_stream_timeout:")), + "{events:?}" + ); + } + + #[tokio::test] + async fn chat_gateway_timeout_frame_preserves_its_typed_error() { + let sse = "event: error\ndata: {\"error\":{\"type\":\"server_error\",\"code\":\"provider_stream_timeout\",\"message\":\"Provider stream exceeded its attempt deadline\"}}\n\n"; + let (client, _request) = managed_gateway_test_client(sse, &managed_receipt_headers()); + let mut receiver = client + .stream_authorized_invocation( + &[], + &RequestConfig { + model: "evalops/openai/gpt-4o-mini".to_owned(), + ..Default::default() + }, + ) + .await + .expect("open chat stream"); + let events = tokio::time::timeout(std::time::Duration::from_secs(5), async { + let mut events = Vec::new(); + while let Some(event) = receiver.recv().await { + events.push(event); + } + events + }) + .await + .expect("gateway error stream should terminate"); + assert!( + matches!( + events.last(), + Some(StreamEvent::ProviderError { + kind: ProviderStreamErrorKind::TransientProtocol, + message, + }) if message.starts_with("managed_gateway_stream_error: provider_stream_timeout:") + ), + "{events:?}" + ); + } + + #[tokio::test] + async fn responses_gateway_timeout_frame_uses_the_same_retryable_error() { + let sse = "event: response.failed\ndata: {\"type\":\"response.failed\",\"response\":{\"id\":\"response-1\",\"error\":{\"type\":\"server_error\",\"code\":\"provider_stream_timeout\",\"message\":\"Provider stream exceeded its attempt deadline\"}}}\n\n"; + let (client, _request) = managed_gateway_test_client(sse, &managed_receipt_headers()); + let mut receiver = client + .stream_authorized_invocation( + &[], + &RequestConfig { + model: "evalops/openai/gpt-5.6-terra".to_owned(), + ..Default::default() + }, + ) + .await + .expect("open responses stream"); + let events = tokio::time::timeout(std::time::Duration::from_secs(5), async { + let mut events = Vec::new(); + while let Some(event) = receiver.recv().await { + events.push(event); + } + events + }) + .await + .expect("gateway error stream should terminate"); + assert!( + matches!( + events.last(), + Some(StreamEvent::ProviderError { + kind: ProviderStreamErrorKind::TransientProtocol, + message, + }) if message.starts_with("managed_gateway_stream_error: provider_stream_timeout:") + ), + "{events:?}" + ); + } + fn managed_receipt_headers() -> [(&'static str, &'static str); 4] { [ ("X-Request-ID", "request-1"), diff --git a/packages/ai-rs/src/sse.rs b/packages/ai-rs/src/sse.rs new file mode 100644 index 000000000..c9049f15e --- /dev/null +++ b/packages/ai-rs/src/sse.rs @@ -0,0 +1,71 @@ +//! Byte framing shared by providers whose SSE transport arrives in arbitrary HTTP chunks. + +/// Decode only after the entire event arrives so split UTF-8 never becomes replacement text. +pub(crate) fn take_frame(buffer: &mut Vec) -> Result, std::str::Utf8Error> { + let lf = buffer.windows(2).position(|window| window == b"\n\n"); + let crlf = buffer.windows(4).position(|window| window == b"\r\n\r\n"); + let (end, delimiter_len) = match (lf, crlf) { + (Some(lf), Some(crlf)) if lf < crlf => (lf, 2), + (Some(_), Some(crlf)) => (crlf, 4), + (Some(lf), None) => (lf, 2), + (None, Some(crlf)) => (crlf, 4), + (None, None) => return Ok(None), + }; + let event = std::str::from_utf8(&buffer[..end])?.to_owned(); + buffer.drain(..end + delimiter_len); + Ok(Some(event)) +} + +#[cfg(test)] +mod tests { + use super::take_frame; + + #[test] + fn unicode_frames_survive_every_two_chunk_boundary() { + for delimiter in ["\n\n", "\r\n\r\n"] { + let frame = format!("data: {{\"text\":\"café 🌊 東京\"}}{delimiter}"); + for split in 0..=frame.len() { + let mut buffer = Vec::new(); + buffer.extend_from_slice(&frame.as_bytes()[..split]); + let first = take_frame(&mut buffer).expect("valid partial frame"); + buffer.extend_from_slice(&frame.as_bytes()[split..]); + let second = take_frame(&mut buffer).expect("valid complete frame"); + assert_eq!( + first.into_iter().chain(second).collect::>(), + [frame.trim_end_matches(['\r', '\n'])] + ); + assert!(buffer.is_empty()); + } + } + } + + #[test] + fn seeded_chunk_partitions_preserve_adjacent_events() { + let first = "event:update\r\ndata:{\"text\":\"é 🌊\"}\r\n\r\n"; + let second = "data:{\"text\":\"東京\"}\n\n"; + let wire = format!("{first}{second}"); + for seed in 0..512_u64 { + let mut random = seed + 1; + let mut cursor = 0; + let mut buffer = Vec::new(); + let mut frames = Vec::new(); + while cursor < wire.len() { + random = random.wrapping_mul(6364136223846793005).wrapping_add(1); + let end = (cursor + 1 + (random as usize % 9)).min(wire.len()); + buffer.extend_from_slice(&wire.as_bytes()[cursor..end]); + while let Some(frame) = take_frame(&mut buffer).expect("valid UTF-8 frame") { + frames.push(frame); + } + cursor = end; + } + assert_eq!(frames, [first.trim_end(), second.trim_end()]); + assert!(buffer.is_empty()); + } + } + + #[test] + fn invalid_utf8_is_an_error_instead_of_replacement_text() { + let mut buffer = b"data: {\"text\":\"\xff\"}\n\n".to_vec(); + assert!(take_frame(&mut buffer).is_err()); + } +} diff --git a/packages/dex-host-rs/Cargo.toml b/packages/dex-host-rs/Cargo.toml new file mode 100644 index 000000000..2c850aea5 --- /dev/null +++ b/packages/dex-host-rs/Cargo.toml @@ -0,0 +1,33 @@ +[package] +name = "maestro-dex-host" +version = "0.1.0" +edition = "2021" +license = "MIT" +description = "Local dex-loop host for Maestro: Log, Tools and Effects ports backed by the workspace filesystem" + +[lib] +name = "maestro_dex_host" +path = "src/lib.rs" + +[dependencies] +# dex-loop lives in the mono root's `rust/` Cargo workspace, not this one; +# this is a path dependency across the workspace boundary, the same shape +# `maestro-swarm` already uses to reach crates outside `products/maestro`. +dex-loop = { path = "../../vendor/dex-loop" } +# `AiRsModel` (`src/model.rs`) wraps this crate's `UnifiedClient::stream` as a +# `dex_loop::Model` -- default features off, matching every other consumer's +# workspace alias, so this stays Bedrock-free. +maestro-ai.workspace = true + +anyhow.workspace = true +fd-lock.workspace = true +futures-util = "0.3" +serde.workspace = true +serde_json.workspace = true +thiserror.workspace = true +tokio = { workspace = true, features = ["fs", "rt", "sync"] } +tokio-stream.workspace = true + +[dev-dependencies] +tempfile.workspace = true +tokio = { workspace = true, features = ["rt-multi-thread", "macros", "time"] } diff --git a/packages/dex-host-rs/src/effects.rs b/packages/dex-host-rs/src/effects.rs new file mode 100644 index 000000000..58a11e2bf --- /dev/null +++ b/packages/dex-host-rs/src/effects.rs @@ -0,0 +1,180 @@ +//! A durable claim/record ledger for local mutations, keyed by `CallId`. +//! +//! Maestro's tool registry has nothing like this today: a local mutation +//! either ran or it didn't, and a process that crashes mid-write has no +//! record of whether the write landed. This adds the `Effects` contract's +//! guarantee — a mutation is dispatched at most once per call id, even +//! across a restart — on top of one JSON file. +//! +//! Access is serialized by an in-process `tokio::sync::Mutex`, not a file +//! lock: unlike [`crate::log::LocalLog`], this ledger does not need to +//! defend against a second, concurrent process racing the same thread — a +//! local Maestro session already assumes one running process per workspace. + +use std::collections::HashMap; +use std::path::{Path, PathBuf}; +use std::sync::Arc; + +use dex_loop::{CallId, Claim, Effects, Fenced, Outcome, Output, ProposedCall, ToolResult}; +use serde::{Deserialize, Serialize}; +use tokio::sync::Mutex; + +/// `None` means claimed but not yet recorded — a mutation dispatched by a +/// process that has not yet appended its outcome (including one that +/// crashed between claim and record). +type Ledger = HashMap>; + +#[derive(Default, Serialize, Deserialize)] +struct Persisted(Ledger); + +#[derive(Clone)] +pub struct LocalEffects { + path: Arc, + ledger: Arc>, +} + +impl LocalEffects { + /// Loads any ledger already at `path` (a fresh thread has none). + pub async fn open(path: impl AsRef) -> Result { + let path = path.as_ref().to_path_buf(); + let ledger = match tokio::fs::read(&path).await { + Ok(bytes) => { + let Persisted(ledger) = + serde_json::from_slice(&bytes).map_err(std::io::Error::other)?; + ledger + } + Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ledger::default(), + Err(error) => return Err(error), + }; + Ok(Self { + path: Arc::new(path), + ledger: Arc::new(Mutex::new(ledger)), + }) + } + + async fn persist(&self, ledger: &Ledger) -> Result<(), Fenced> { + if let Some(parent) = self.path.parent() { + tokio::fs::create_dir_all(parent).await.map_err(|error| { + Fenced::new(format!("could not create ledger directory: {error}")) + })?; + } + let bytes = serde_json::to_vec(&Persisted(ledger.clone())) + .map_err(|error| Fenced::new(format!("ledger does not encode: {error}")))?; + // Write-then-rename: a crash mid-write leaves the previous ledger in + // place instead of a truncated one. + let temp = self.path.with_extension("json.tmp"); + tokio::fs::write(&temp, &bytes) + .await + .map_err(|error| Fenced::new(format!("could not write ledger: {error}")))?; + tokio::fs::rename(&temp, self.path.as_path()) + .await + .map_err(|error| Fenced::new(format!("could not commit ledger: {error}"))) + } +} + +fn running() -> ToolResult { + ToolResult { + outcome: Outcome::Running, + output: Output::Text("dispatched; no outcome recorded yet".into()), + receipt: None, + } +} + +impl Effects for LocalEffects { + async fn claim(&self, call: &ProposedCall) -> Result { + let mut ledger = self.ledger.lock().await; + match ledger.get(call.id.as_str()) { + Some(Some(result)) => Ok(Claim::Existing(result.clone())), + Some(None) => Ok(Claim::Existing(running())), + None => { + ledger.insert(call.id.as_str().to_owned(), None); + self.persist(&ledger).await?; + Ok(Claim::Granted) + } + } + } + + async fn record(&self, call: &CallId, result: &ToolResult) -> Result<(), Fenced> { + let mut ledger = self.ledger.lock().await; + ledger.insert(call.as_str().to_owned(), Some(result.clone())); + self.persist(&ledger).await + } +} + +#[cfg(test)] +mod tests { + use super::*; + use dex_loop::{PrincipalId, ProposedCall, ToolName}; + use tempfile::TempDir; + + fn call(id: &str) -> ProposedCall { + ProposedCall::new( + CallId::new(id), + ToolName::new("fs.write_file"), + serde_json::json!({"path": "a.txt", "content": "x"}), + PrincipalId::new("alice"), + ) + } + + #[tokio::test] + async fn first_claim_grants_and_second_claim_sees_running() { + let dir = TempDir::new().expect("tempdir"); + let effects = LocalEffects::open(dir.path().join("ledger.json")) + .await + .expect("open"); + let call = call("c1"); + assert_eq!(effects.claim(&call).await.expect("claim"), Claim::Granted); + let second = effects.claim(&call).await.expect("claim again"); + assert!(matches!(second, Claim::Existing(result) if result.outcome == Outcome::Running)); + } + + #[tokio::test] + async fn recorded_outcome_is_returned_on_a_later_claim() { + let dir = TempDir::new().expect("tempdir"); + let effects = LocalEffects::open(dir.path().join("ledger.json")) + .await + .expect("open"); + let call = call("c1"); + effects.claim(&call).await.expect("claim"); + effects + .record(&call.id, &ToolResult::text("done")) + .await + .expect("record"); + let claim = effects.claim(&call).await.expect("claim again"); + assert_eq!(claim, Claim::Existing(ToolResult::text("done"))); + } + + #[tokio::test] + async fn ledger_survives_a_reopen_at_the_same_path() { + let dir = TempDir::new().expect("tempdir"); + let path = dir.path().join("ledger.json"); + let first = LocalEffects::open(&path).await.expect("open"); + let call = call("c1"); + first.claim(&call).await.expect("claim"); + first + .record(&call.id, &ToolResult::text("done")) + .await + .expect("record"); + + // A fresh handle, as a restarted process would create. + let reopened = LocalEffects::open(&path).await.expect("reopen"); + let claim = reopened.claim(&call).await.expect("claim after reopen"); + assert_eq!(claim, Claim::Existing(ToolResult::text("done"))); + } + + #[tokio::test] + async fn distinct_call_ids_are_independent() { + let dir = TempDir::new().expect("tempdir"); + let effects = LocalEffects::open(dir.path().join("ledger.json")) + .await + .expect("open"); + assert_eq!( + effects.claim(&call("a")).await.expect("claim a"), + Claim::Granted + ); + assert_eq!( + effects.claim(&call("b")).await.expect("claim b"), + Claim::Granted + ); + } +} diff --git a/packages/dex-host-rs/src/lease.rs b/packages/dex-host-rs/src/lease.rs new file mode 100644 index 000000000..3d890c3f9 --- /dev/null +++ b/packages/dex-host-rs/src/lease.rs @@ -0,0 +1,210 @@ +//! Generation fencing for [`crate::log::LocalLog`], mirroring `dex-runtime`'s +//! `PgLog`: a lease holds one generation number, persisted next to the log +//! instead of in a Postgres row. `acquire` invalidates any earlier lease for +//! the same thread by bumping the persisted generation under an exclusive +//! file lock; every later write re-checks the persisted value against the +//! generation this lease captured, under the same lock, so a superseded +//! writer is refused instead of silently appending after a new owner took +//! over. + +use std::fs::{self, OpenOptions}; +use std::io; +use std::path::{Path, PathBuf}; + +use serde::{Deserialize, Serialize}; + +/// The durable state fenced together, the same way `PgLog`'s thread row +/// holds `lease_generation` and (via `MAX(cursor)`) the last cursor in one +/// row locked by one transaction. +#[derive(Clone, Copy, Debug, Default, Serialize, Deserialize)] +pub(crate) struct Meta { + pub(crate) generation: u64, + pub(crate) last_cursor: i64, +} + +/// Runs `body` with an exclusive lock on `lock_path`'s file, creating the +/// lock file (and its parent directory) if needed. Blocking: callers on an +/// async runtime must run this inside `spawn_blocking`. +pub(crate) fn with_locked_meta( + dir: &Path, + body: impl FnOnce(&mut Meta) -> io::Result, +) -> io::Result { + fs::create_dir_all(dir)?; + let lock_file = OpenOptions::new() + .create(true) + .truncate(false) + .read(true) + .write(true) + .open(lock_path(dir))?; + let mut lock = fd_lock::RwLock::new(lock_file); + let _guard = lock.write()?; + let mut meta = read_meta(dir)?; + let result = body(&mut meta)?; + write_meta(dir, &meta)?; + Ok(result) +} + +/// Runs `body` with a shared lock on the same lock file `with_locked_meta` +/// takes exclusively: any number of readers may hold this at once, but none +/// of them overlaps a writer's `write_all` to the log file, so a reader +/// never observes a torn trailing line from a write still in progress. +/// Blocking: callers on an async runtime must run this inside +/// `spawn_blocking`. +pub(crate) fn with_shared_lock( + dir: &Path, + body: impl FnOnce() -> io::Result, +) -> io::Result { + fs::create_dir_all(dir)?; + let lock_file = OpenOptions::new() + .create(true) + .truncate(false) + .read(true) + .write(true) + .open(lock_path(dir))?; + let lock = fd_lock::RwLock::new(lock_file); + let _guard = lock.read()?; + body() +} + +pub(crate) fn lock_path(dir: &Path) -> PathBuf { + dir.join("lease.lock") +} + +pub(crate) fn meta_path(dir: &Path) -> PathBuf { + dir.join("meta.json") +} + +pub(crate) fn read_meta(dir: &Path) -> io::Result { + match fs::read_to_string(meta_path(dir)) { + Ok(text) => serde_json::from_str(&text) + .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error)), + Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(Meta::default()), + Err(error) => Err(error), + } +} + +fn write_meta(dir: &Path, meta: &Meta) -> io::Result<()> { + let path = meta_path(dir); + let temp = path.with_extension("json.tmp"); + let contents = serde_json::to_vec(meta).map_err(io::Error::other)?; + fs::write(&temp, contents)?; + fs::rename(&temp, &path) +} + +/// Also used directly by tests that need a file handle open past this +/// module, e.g. to assert the lock file exists. +#[cfg(test)] +pub(crate) fn open_lock_file(dir: &Path) -> io::Result { + OpenOptions::new().read(true).open(lock_path(dir)) +} + +#[cfg(test)] +mod tests { + use super::*; + use tempfile::TempDir; + + #[test] + fn first_acquire_starts_at_generation_one() { + let dir = TempDir::new().expect("tempdir"); + let generation = with_locked_meta(dir.path(), |meta| { + meta.generation += 1; + Ok(meta.generation) + }) + .expect("acquire"); + assert_eq!(generation, 1); + } + + #[test] + fn concurrent_acquires_are_serialized_and_strictly_increasing() { + let dir = TempDir::new().expect("tempdir"); + let dir_path = dir.path().to_path_buf(); + let handles: Vec<_> = (0..8) + .map(|_| { + let dir_path = dir_path.clone(); + std::thread::spawn(move || { + with_locked_meta(&dir_path, |meta| { + meta.generation += 1; + Ok(meta.generation) + }) + .expect("acquire") + }) + }) + .collect(); + let mut generations: Vec = handles + .into_iter() + .map(|handle| handle.join().unwrap()) + .collect(); + generations.sort_unstable(); + assert_eq!(generations, (1..=8).collect::>()); + } + + #[test] + fn meta_persists_across_separate_locked_sections() { + let dir = TempDir::new().expect("tempdir"); + with_locked_meta(dir.path(), |meta| { + meta.generation = 5; + meta.last_cursor = 12; + Ok(()) + }) + .expect("write"); + let meta = read_meta(dir.path()).expect("read"); + assert_eq!(meta.generation, 5); + assert_eq!(meta.last_cursor, 12); + } + + #[test] + fn shared_lock_waits_for_an_in_progress_exclusive_lock() { + use std::sync::mpsc; + use std::time::Duration; + + let dir = TempDir::new().expect("tempdir"); + let dir_path = dir.path().to_path_buf(); + let (writer_started, wait_for_writer_started) = mpsc::channel(); + let (release_writer, wait_to_release_writer) = mpsc::channel::<()>(); + + let writer = std::thread::spawn({ + let dir_path = dir_path.clone(); + move || { + with_locked_meta(&dir_path, |meta| { + meta.generation += 1; + writer_started.send(()).unwrap(); + // Hold the exclusive lock until the reader has had a + // chance to observe it is blocked. + wait_to_release_writer + .recv_timeout(Duration::from_secs(5)) + .unwrap(); + Ok(()) + }) + .expect("writer") + } + }); + wait_for_writer_started + .recv_timeout(Duration::from_secs(5)) + .unwrap(); + + let reader = + std::thread::spawn(move || with_shared_lock(&dir_path, || Ok(())).expect("reader")); + // The reader cannot finish before the writer releases; if it did, + // the shared lock was not actually exclusive against the writer. + std::thread::sleep(Duration::from_millis(50)); + assert!( + !reader.is_finished(), + "reader must wait for the writer's exclusive lock" + ); + + release_writer.send(()).unwrap(); + writer.join().unwrap(); + reader.join().unwrap(); + } + + #[test] + fn lock_file_is_created_alongside_meta() { + let dir = TempDir::new().expect("tempdir"); + with_locked_meta(dir.path(), |meta| { + meta.generation += 1; + Ok(()) + }) + .expect("write"); + open_lock_file(dir.path()).expect("lock file exists"); + } +} diff --git a/packages/dex-host-rs/src/lib.rs b/packages/dex-host-rs/src/lib.rs new file mode 100644 index 000000000..c0c38f5a7 --- /dev/null +++ b/packages/dex-host-rs/src/lib.rs @@ -0,0 +1,37 @@ +//! A local `dex_loop` host: `Log`, `Tools` and `Effects` ports backed by the +//! workspace filesystem, with no database and no HTTP server, plus a `Model` +//! adapter over `maestro-ai` (`model.rs`) and a consumer that drives one +//! turn to completion unattended (`turn.rs`, `run_local_turn`). +//! +//! This is the first slice of "Maestro becomes another host with a local +//! Log" (see `docs/design/maestro-on-dex-loop.md` at the repository root). +//! It proves the kernel's park/approve/resume semantics run correctly +//! against a purely local host — see `tests/turn.rs` and `turn.rs`'s own +//! tests — without touching Maestro's existing turn loop (`maestro-runtime`, +//! `maestro-local-host`), which still runs every real turn today. +//! `run_local_turn` is reached from `tui-rs` only behind +//! `MAESTRO_DEX_LOOP=1` (`dex_loop_local.rs`); the flag is unset by default. +//! +//! Placed in Maestro's own Rust workspace (`products/maestro/packages/`, +//! not `rust/crates/`) because that is where a real cutover would keep it; +//! `dex-loop` is reached by a path dependency across the workspace +//! boundary, the same shape `maestro-swarm` already uses. + +mod effects; +mod lease; +mod log; +mod model; +mod tools; +mod turn; + +pub use effects::LocalEffects; +pub use log::{LocalLog, LogError}; +pub use model::AiRsModel; +pub use tools::{LocalTools, READ_FILE, WRITE_FILE}; +pub use turn::{LocalTurnOutcome, LocalTurnRequest, UNATTENDED_ANSWER, run_local_turn}; + +/// Re-exported so a consumer that only depends on this crate (e.g. +/// `tui-rs`'s `dex_loop_local.rs`) can build a `ThreadId`/`PrincipalId`/ +/// `TurnId`/`Exit` without also taking a direct path dependency on +/// `dex-loop` across the same workspace boundary. +pub use dex_loop; diff --git a/packages/dex-host-rs/src/log.rs b/packages/dex-host-rs/src/log.rs new file mode 100644 index 000000000..b274eef1f --- /dev/null +++ b/packages/dex-host-rs/src/log.rs @@ -0,0 +1,398 @@ +//! The local, file-backed `dex_loop::Log`: one append-only JSONL file per +//! thread, fenced the way `dex-runtime`'s `PgLog` fences a Postgres row. +//! +//! This writes a fresh, dex-loop-native event stream — one JSON object per +//! line, `{"cursor": N, "event": {...}}` — rather than reusing Maestro's +//! existing `maestro-session` JSONL format, so this first slice can be +//! staged without touching Maestro's current session persistence. See +//! `docs/design/maestro-on-dex-loop.md` for the cutover that unifies them. +//! +//! Unlike `PgLog`, consecutive `TextDelta` calls are not coalesced into one +//! row: `append_text` writes one row per call. `dex_loop::Log`'s contract +//! permits this ("the engine never assumes one row per call"); coalescing +//! is a production-log optimization, not a correctness requirement, and is +//! left for the real cutover. + +use std::fs::{self, OpenOptions}; +use std::io::{self, BufRead, Write}; +use std::path::{Path, PathBuf}; +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, Ordering}; + +use dex_loop::{Cursor, Event, Fenced, Log, ThreadId}; +use serde::{Deserialize, Serialize}; + +use crate::lease::{self, Meta}; + +/// Reading or decoding the durable log failed. Distinct from `Fenced`, +/// which is the `Log` trait's own error for a refused *write*; this is +/// returned by the host-side helpers (`acquire`, `read_all`) a caller uses +/// before the engine ever sees the log. +#[derive(Debug, thiserror::Error)] +pub enum LogError { + #[error("io error: {0}")] + Io(#[from] io::Error), + #[error("event at line {line} does not decode: {source}")] + Decode { + line: usize, + source: serde_json::Error, + }, + #[error("background task panicked: {0}")] + Join(#[from] tokio::task::JoinError), +} + +impl From for Fenced { + fn from(error: LogError) -> Self { + Fenced::new(error.to_string()) + } +} + +#[derive(Serialize, Deserialize)] +struct Row { + cursor: Cursor, + event: Event, +} + +fn thread_dir(root: &Path, thread: &ThreadId) -> PathBuf { + root.join(sanitize(&thread.org)) + .join(sanitize(&thread.workspace)) + .join(sanitize(&thread.thread)) +} + +/// Confines a `ThreadId` field to one path component: a stray `/` or `..` +/// in a tenant-controlled id must not let one thread's log escape into +/// another thread's directory, or above `root` entirely. +fn sanitize(part: &str) -> String { + let cleaned: String = part + .chars() + .map(|c| { + if c.is_ascii_alphanumeric() || c == '-' || c == '_' || c == '.' { + c + } else { + '_' + } + }) + .collect(); + if cleaned.is_empty() || cleaned == "." || cleaned == ".." { + "_".to_owned() + } else { + cleaned + } +} + +fn log_path(dir: &Path) -> PathBuf { + dir.join("log.jsonl") +} + +enum AppendOutcome { + Fenced { held: u64, current: u64 }, + Wrote(Vec), +} + +/// Blocking: appends `events` to `dir`'s log under the same lock that +/// guards the generation check, so the check and the write are atomic +/// together — the same guarantee `PgLog::write_rows` gets from doing both +/// inside one Postgres transaction. +fn append_locked( + dir: &Path, + expected_generation: u64, + events: &[Event], +) -> io::Result { + lease::with_locked_meta(dir, |meta: &mut Meta| { + if meta.generation != expected_generation { + return Ok(AppendOutcome::Fenced { + held: expected_generation, + current: meta.generation, + }); + } + if events.is_empty() { + return Ok(AppendOutcome::Wrote(Vec::new())); + } + let mut file = OpenOptions::new() + .create(true) + .append(true) + .open(log_path(dir))?; + let mut cursors = Vec::with_capacity(events.len()); + let mut buffer = Vec::new(); + for event in events { + meta.last_cursor += 1; + let cursor = Cursor(meta.last_cursor); + let row = Row { + cursor, + event: event.clone(), + }; + serde_json::to_writer(&mut buffer, &row).map_err(io::Error::other)?; + buffer.push(b'\n'); + cursors.push(cursor); + } + file.write_all(&buffer)?; + file.sync_data()?; + Ok(AppendOutcome::Wrote(cursors)) + }) +} + +/// Blocking: bumps `dir`'s generation, invalidating any earlier lease for +/// the same thread. +fn acquire_locked(dir: &Path) -> io::Result { + lease::with_locked_meta(dir, |meta: &mut Meta| { + meta.generation += 1; + Ok(meta.generation) + }) +} + +/// Blocking: every row in `dir`'s log, oldest first. Takes the same lock +/// `append_locked` takes exclusively, shared, so this never overlaps an +/// in-progress `write_all` and cannot observe a torn trailing line. +fn read_all_locked(dir: &Path) -> Result, LogError> { + lease::with_shared_lock(dir, || read_rows(dir)) + .map_err(LogError::from)? + .into_iter() + .enumerate() + .map(|(index, line)| { + serde_json::from_str::(&line) + .map(|row| (row.cursor, row.event)) + .map_err(|source| LogError::Decode { + line: index + 1, + source, + }) + }) + .collect() +} + +/// Every non-empty line in `dir`'s log file, unparsed. Split out of +/// `read_all_locked` so decode errors (which need the line index) happen +/// outside the lock. +fn read_rows(dir: &Path) -> io::Result> { + let path = log_path(dir); + let file = match fs::File::open(&path) { + Ok(file) => file, + Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(Vec::new()), + Err(error) => return Err(error), + }; + io::BufReader::new(file) + .lines() + .filter(|line| !matches!(line, Ok(line) if line.trim().is_empty())) + .collect() +} + +struct Inner { + dir: PathBuf, + generation: u64, + lost: AtomicBool, +} + +/// The `dex_loop::Log` for one thread under one lease. Cloning shares the +/// same lease: every clone observes the same "lost" state once any of them +/// sees a fenced write, matching `PgLog`'s `Clone` semantics. +#[derive(Clone)] +pub struct LocalLog { + inner: Arc, +} + +impl LocalLog { + /// Acquires a fresh lease for `thread` under `root`, bumping the + /// persisted generation so any earlier lease for the same thread is + /// fenced out on its next write. + pub async fn acquire(root: impl AsRef, thread: &ThreadId) -> Result { + let dir = thread_dir(root.as_ref(), thread); + let generation = { + let dir = dir.clone(); + tokio::task::spawn_blocking(move || acquire_locked(&dir)).await?? + }; + Ok(Self { + inner: Arc::new(Inner { + dir, + generation, + lost: AtomicBool::new(false), + }), + }) + } + + /// Every event in the thread's log, oldest first — for + /// `dex_loop::rehydrate`. Call this once, right after `acquire`, before + /// handing the log to `Engine::run`. + pub async fn read_all(&self) -> Result, LogError> { + let dir = self.inner.dir.clone(); + tokio::task::spawn_blocking(move || read_all_locked(&dir)).await? + } + + fn ensure_live(&self) -> Result<(), Fenced> { + if self.inner.lost.load(Ordering::Acquire) { + return Err(Fenced::new("lease lost")); + } + Ok(()) + } + + async fn append_events(&self, events: Vec) -> Result, Fenced> { + self.ensure_live()?; + if events.is_empty() { + return Ok(Vec::new()); + } + let dir = self.inner.dir.clone(); + let generation = self.inner.generation; + let outcome = tokio::task::spawn_blocking(move || append_locked(&dir, generation, &events)) + .await + .map_err(|error| Fenced::new(format!("background append task panicked: {error}")))? + .map_err(|error| Fenced::new(format!("log write failed: {error}")))?; + match outcome { + AppendOutcome::Fenced { held, current } => { + self.inner.lost.store(true, Ordering::Release); + Err(Fenced::new(format!( + "lease generation moved: held {held}, now {current}" + ))) + } + AppendOutcome::Wrote(cursors) => Ok(cursors), + } + } +} + +impl Log for LocalLog { + async fn append(&self, events: &[Event]) -> Result, Fenced> { + self.append_events(events.to_vec()).await + } + + /// No coalescing (see the module doc): each call is its own row. + async fn append_text(&self, text: String) -> Result<(), Fenced> { + self.append_events(vec![Event::TextDelta { text }]) + .await + .map(|_| ()) + } + + async fn control_since(&self, after: Cursor) -> Result, Fenced> { + self.ensure_live()?; + let all = self.read_all().await.map_err(Fenced::from)?; + Ok(all + .into_iter() + .filter(|(cursor, event)| *cursor > after && event.is_control()) + .collect()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use dex_loop::{PrincipalId, TurnId}; + use tempfile::TempDir; + + fn thread() -> ThreadId { + ThreadId { + org: "org-1".into(), + workspace: "ws-1".into(), + thread: "thread-1".into(), + } + } + + fn user_message(text: &str) -> Event { + Event::UserMessage { + turn: TurnId::new("t1"), + message_id: None, + principal: PrincipalId::new("alice"), + text: text.into(), + attachments: Vec::new(), + client_tools: Vec::new(), + authorized_tools: Vec::new(), + approval_mode: dex_loop::ApprovalMode::Interactive, + } + } + + #[tokio::test] + async fn appended_events_round_trip_with_increasing_cursors() { + let dir = TempDir::new().expect("tempdir"); + let log = LocalLog::acquire(dir.path(), &thread()) + .await + .expect("acquire"); + let cursors = log + .append(&[user_message("hi"), Event::Interrupted]) + .await + .expect("append"); + assert_eq!(cursors, vec![Cursor(1), Cursor(2)]); + let all = log.read_all().await.expect("read"); + assert_eq!(all.len(), 2); + assert_eq!(all[0].0, Cursor(1)); + assert_eq!(all[1].0, Cursor(2)); + } + + #[tokio::test] + async fn control_since_only_returns_control_events_after_the_cursor() { + let dir = TempDir::new().expect("tempdir"); + let log = LocalLog::acquire(dir.path(), &thread()) + .await + .expect("acquire"); + log.append(&[user_message("hi")]).await.expect("append"); + let interrupt = Event::Interrupt { + principal: PrincipalId::new("alice"), + }; + log.append(std::slice::from_ref(&interrupt)) + .await + .expect("append"); + let control = log.control_since(Cursor::START).await.expect("control"); + assert_eq!(control, vec![(Cursor(2), interrupt)]); + } + + #[tokio::test] + async fn a_second_acquire_fences_the_first_leases_next_write() { + let dir = TempDir::new().expect("tempdir"); + let first = LocalLog::acquire(dir.path(), &thread()) + .await + .expect("first acquire"); + let _second = LocalLog::acquire(dir.path(), &thread()) + .await + .expect("second acquire"); + let result = first.append(&[user_message("late")]).await; + assert!(matches!(result, Err(error) if error.reason.contains("lease generation moved"))); + } + + #[tokio::test] + async fn a_fenced_lease_stays_fenced_without_rechecking_disk() { + let dir = TempDir::new().expect("tempdir"); + let first = LocalLog::acquire(dir.path(), &thread()) + .await + .expect("first acquire"); + let _second = LocalLog::acquire(dir.path(), &thread()) + .await + .expect("second acquire"); + assert!(first.append(&[user_message("one")]).await.is_err()); + // A third acquire would let a *fresh* lease through, but `first` is + // latched lost in-process and must not silently recover. + assert!(first.append(&[user_message("two")]).await.is_err()); + } + + #[tokio::test] + async fn empty_thread_reads_back_as_no_events() { + let dir = TempDir::new().expect("tempdir"); + let log = LocalLog::acquire(dir.path(), &thread()) + .await + .expect("acquire"); + assert!(log.read_all().await.expect("read").is_empty()); + } + + #[tokio::test] + async fn thread_id_components_cannot_escape_the_root() { + let dir = TempDir::new().expect("tempdir"); + let escaping = ThreadId { + org: "../../etc".into(), + workspace: "ws".into(), + thread: "t".into(), + }; + // `sanitize` keeps literal `.` characters (legitimate in an org or + // workspace name), so a sanitized component can still *contain* two + // dots; the property that matters is that no component *is* `..`, + // which is what would let `Path::join` walk back up to a parent. + let path = thread_dir(dir.path(), &escaping); + assert!(path.starts_with(dir.path())); + assert!( + path.strip_prefix(dir.path()) + .expect("under root") + .components() + .all(|component| component != std::path::Component::ParentDir) + ); + + let log = LocalLog::acquire(dir.path(), &escaping) + .await + .expect("acquire"); + log.append(&[user_message("hi")]).await.expect("append"); + let mut entries = fs::read_dir(dir.path()).expect("read root"); + let child = entries.next().expect("one child").expect("entry"); + assert!(child.path().starts_with(dir.path())); + } +} diff --git a/packages/dex-host-rs/src/model.rs b/packages/dex-host-rs/src/model.rs new file mode 100644 index 000000000..b724dcfec --- /dev/null +++ b/packages/dex-host-rs/src/model.rs @@ -0,0 +1,338 @@ +//! A `dex_loop::Model` adapter over `maestro-ai`'s `UnifiedClient::stream`. +//! +//! First real port of `docs/design/maestro-on-dex-loop.md`'s cutover step 2 +//! ("Build the real `Model` adapter over `ai_rs::AiClient::stream`"). Scope +//! matches the rest of this crate's "first slice" framing: one call at a +//! time, no retries beyond what `UnifiedClient::stream` already does, no +//! prompt caching, no thinking-block translation, and cost is not yet +//! attributed (`Usage::cost_micros` is always `0` -- `maestro-ai` reports +//! cost as a separate `StreamEvent::ProviderCost`, not on the `Usage` event +//! this adapter reads). None of that is a correctness bug for the local +//! `run_local_turn` consumer this backs; it is exactly what "first slice" +//! means, called out here the way the crate doc comment calls out the rest. + +use std::collections::BTreeMap; +use std::sync::Arc; + +use dex_loop::{ + Context, Entry, Message as LoopMessage, ModelChunk, ModelError, Outcome, Output, ToolName, + ToolSpec, Usage, +}; +use futures_util::Stream; +use maestro_ai::{ + ContentBlock as AiContentBlock, Message as AiMessage, MessageContent, RequestConfig, Role, + StreamEvent, Tool as AiTool, UnifiedClient, +}; +use tokio_stream::wrappers::UnboundedReceiverStream; + +/// Wraps one `maestro_ai::UnifiedClient` as a `dex_loop::Model`. +#[derive(Clone)] +pub struct AiRsModel { + client: UnifiedClient, + model: String, + max_tokens: u32, + system: Option, +} + +impl AiRsModel { + #[must_use] + pub fn new(client: UnifiedClient, model: impl Into, max_tokens: u32) -> Self { + Self { + client, + model: model.into(), + max_tokens, + system: None, + } + } + + #[must_use] + pub fn with_system(mut self, system: impl Into) -> Self { + self.system = Some(system.into()); + self + } + + fn request_config(&self, tools: Vec) -> RequestConfig { + RequestConfig { + model: self.model.clone(), + max_tokens: self.max_tokens, + system: self.system.clone(), + tools: Arc::new(tools), + ..RequestConfig::default() + } + } +} + +fn to_ai_messages(history: &[Entry]) -> Vec { + history + .iter() + .filter_map(|entry| to_ai_message(&entry.message)) + .collect() +} + +fn to_ai_message(message: &LoopMessage) -> Option { + match message { + LoopMessage::User { text, .. } | LoopMessage::Summary { text } => Some(AiMessage { + role: Role::User, + content: MessageContent::text(text.clone()), + }), + LoopMessage::Assistant { text, calls, .. } => { + if calls.is_empty() { + return Some(AiMessage { + role: Role::Assistant, + content: MessageContent::text(text.clone()), + }); + } + let mut blocks = Vec::with_capacity(calls.len() + 1); + if !text.is_empty() { + blocks.push(AiContentBlock::Text { text: text.clone() }); + } + for call in calls { + blocks.push(AiContentBlock::ToolUse { + id: call.id.as_str().to_owned(), + name: call.tool.as_str().to_owned(), + input: call.args.clone(), + gemini_context: None, + }); + } + Some(AiMessage { + role: Role::Assistant, + content: MessageContent::Blocks(blocks), + }) + } + LoopMessage::Tool { + call, + outcome, + output, + .. + } => { + let content = match output { + Output::Text(text) => text.clone(), + // `LocalTools` never produces this variant (see `tools.rs`); + // a real out-of-line store needs its own resolver here. + Output::Ref(_) => { + "[tool output stored out of line; not available to this host]".to_owned() + } + }; + Some(AiMessage { + role: Role::User, + content: MessageContent::Blocks(vec![AiContentBlock::ToolResult { + tool_use_id: call.as_str().to_owned(), + content, + is_error: Some(matches!(outcome, Outcome::Failed)), + }]), + }) + } + } +} + +fn to_ai_tools(tools: &[&ToolSpec]) -> Vec { + tools + .iter() + .map(|spec| AiTool { + name: spec.name.as_str().to_owned(), + description: spec.label.clone(), + input_schema: spec.schema.clone(), + schema_enforcement: Default::default(), + }) + .collect() +} + +/// Accumulates one streamed response's tool-call JSON per content-block +/// index (providers stream `InputJsonDelta` chunks between a block's +/// `ContentBlockStart` and `ContentBlockStop`), and turns each provider +/// `StreamEvent` into zero or more `ModelChunk`s. +#[derive(Default)] +struct ChunkTranslator { + pending_tools: BTreeMap, // index -> (tool name, json buffer) +} + +impl ChunkTranslator { + fn translate(&mut self, event: StreamEvent) -> Vec> { + match event { + StreamEvent::TextDelta { text, .. } => vec![Ok(ModelChunk::Text(text))], + StreamEvent::ContentBlockStart { + index, + block: AiContentBlock::ToolUse { name, .. }, + } => { + self.pending_tools.insert(index, (name, String::new())); + Vec::new() + } + StreamEvent::InputJsonDelta { + index, + partial_json, + } => { + if let Some((_, buffer)) = self.pending_tools.get_mut(&index) { + buffer.push_str(&partial_json); + } + Vec::new() + } + StreamEvent::ContentBlockStop { index, .. } => { + let Some((name, buffer)) = self.pending_tools.remove(&index) else { + return Vec::new(); + }; + let args = if buffer.trim().is_empty() { + serde_json::Value::Object(serde_json::Map::new()) + } else { + match serde_json::from_str(&buffer) { + Ok(value) => value, + Err(error) => { + return vec![Err(ModelError { + message: format!( + "provider returned invalid tool-call JSON for {name}: {error}" + ), + })]; + } + } + }; + vec![Ok(ModelChunk::ToolCall { + name: ToolName::new(name), + args, + })] + } + StreamEvent::Usage { + input_tokens, + output_tokens, + .. + } => vec![Ok(ModelChunk::Usage(Usage { + input_tokens, + output_tokens, + // Not attributed here; see the module doc comment. + cost_micros: 0, + }))], + StreamEvent::ProviderError { message, .. } | StreamEvent::Error { message } => { + vec![Err(ModelError { message })] + } + _ => Vec::new(), + } + } +} + +impl dex_loop::Model for AiRsModel { + fn stream<'a>( + &'a self, + ctx: &'a Context, + tools: &'a [&'a ToolSpec], + ) -> impl Stream> + Send + 'a { + let messages = to_ai_messages(ctx.history()); + let config = self.request_config(to_ai_tools(tools)); + let client = self.client.clone(); + let (tx, rx) = tokio::sync::mpsc::unbounded_channel(); + tokio::spawn(async move { + match client.stream(&messages, &config).await { + Ok(mut source) => { + let mut translator = ChunkTranslator::default(); + while let Some(event) = source.recv().await { + for chunk in translator.translate(event) { + if tx.send(chunk).is_err() { + return; + } + } + } + } + Err(error) => { + let _ = tx.send(Err(ModelError { + message: format!("{error:#}"), + })); + } + } + }); + UnboundedReceiverStream::new(rx) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use dex_loop::{Model as _, PrincipalId, ThreadId}; + use futures_util::StreamExt; + use maestro_ai::{ScriptedBlock, ScriptedClient, ScriptedResponse}; + + fn thread() -> ThreadId { + ThreadId { + org: "org-1".into(), + workspace: "ws-1".into(), + thread: "thread-1".into(), + } + } + + #[tokio::test] + async fn translates_text_then_tool_call_then_usage() { + let scripted = ScriptedClient::new( + "scripted", + vec![ScriptedResponse { + blocks: vec![ + ScriptedBlock::Text("Reading the file.".to_owned()), + ScriptedBlock::ToolUse { + id: "call-1".to_owned(), + name: "fs.read_file".to_owned(), + input: serde_json::json!({"path": "a.txt"}), + }, + ], + stop_reason: maestro_ai::StopReason::ToolUse, + error: None, + }], + ); + let model = AiRsModel::new(UnifiedClient::Scripted(scripted), "scripted-model", 1024); + + let ctx = dex_loop::rehydrate( + thread(), + &[( + dex_loop::Cursor(1), + dex_loop::Event::UserMessage { + turn: dex_loop::TurnId::new("t1"), + message_id: None, + principal: PrincipalId::new("alice"), + text: "read a.txt".into(), + attachments: Vec::new(), + client_tools: Vec::new(), + authorized_tools: Vec::new(), + approval_mode: dex_loop::ApprovalMode::Interactive, + }, + )], + ); + let spec = ToolSpec { + name: ToolName::new("fs.read_file"), + label: "Read a file".into(), + schema: serde_json::json!({"type": "object"}), + read_only: true, + core: true, + governance: dex_loop::GovernanceClass::Plain, + executor: dex_loop::ExecutorKind::InProcess, + }; + let tools = [&spec]; + let chunks: Vec<_> = model.stream(&ctx, &tools).collect().await; + let chunks: Vec = chunks + .into_iter() + .collect::>() + .expect("no model errors"); + + // Text, then the assembled tool call, then the trailing usage chunk + // every scripted response emits. + assert_eq!(chunks.len(), 3, "{chunks:?}"); + assert!(matches!(&chunks[0], ModelChunk::Text(text) if text == "Reading the file.")); + assert!(matches!( + &chunks[1], + ModelChunk::ToolCall { name, args } + if name.as_str() == "fs.read_file" && args["path"] == "a.txt" + )); + assert!(matches!(&chunks[2], ModelChunk::Usage(_))); + } + + #[tokio::test] + async fn surfaces_provider_errors_as_model_errors() { + let scripted = ScriptedClient::new( + "scripted", + vec![ScriptedResponse::stream_error("upstream exploded")], + ); + let model = AiRsModel::new(UnifiedClient::Scripted(scripted), "scripted-model", 1024); + let ctx = dex_loop::rehydrate(thread(), &[]); + let tools: [&ToolSpec; 0] = []; + let chunks: Vec<_> = model.stream(&ctx, &tools).collect().await; + assert!( + chunks.iter().any( + |chunk| matches!(chunk, Err(error) if error.message.contains("upstream exploded")) + ), + "expected an error chunk, got {chunks:?}" + ); + } +} diff --git a/packages/dex-host-rs/src/tools.rs b/packages/dex-host-rs/src/tools.rs new file mode 100644 index 000000000..395c61a3c --- /dev/null +++ b/packages/dex-host-rs/src/tools.rs @@ -0,0 +1,301 @@ +//! The `Tools` port over the workspace filesystem. +//! +//! Two tools, standing in for `local-host-rs`'s much larger registry +//! (`packages/local-host-rs/src/tools/`): one read (`fs.read_file`) and one +//! mutation that needs approval (`fs.write_file`). A real cutover ports the +//! rest of that registry the same way — each tool becomes a `ToolSpec` plus +//! a `run` arm, and `NativeHost::requires_approval`'s per-call decision +//! becomes this module's `policy` — not a rewrite of the tools themselves. + +use std::path::{Path, PathBuf}; +use std::sync::Arc; + +use dex_loop::{ + ApprovalId, CancellationToken, Context, ExecutorKind, GovernanceClass, PrincipalId, + ProposedCall, ThreadId, ToolName, ToolResult, ToolSpec, Tools, Verdict, +}; +use serde_json::Value; + +pub const READ_FILE: &str = "fs.read_file"; +pub const WRITE_FILE: &str = "fs.write_file"; + +fn catalog() -> Vec { + vec![ + ToolSpec { + name: ToolName::new(READ_FILE), + label: "Read a workspace file".into(), + schema: serde_json::json!({ + "type": "object", + "additionalProperties": false, + "required": ["path"], + "properties": { + "path": { + "type": "string", + "description": "Path relative to the workspace root", + }, + }, + }), + read_only: true, + core: true, + governance: GovernanceClass::Plain, + executor: ExecutorKind::InProcess, + }, + ToolSpec { + name: ToolName::new(WRITE_FILE), + label: "Write a workspace file".into(), + schema: serde_json::json!({ + "type": "object", + "additionalProperties": false, + "required": ["path", "content"], + "properties": { + "path": { + "type": "string", + "description": "Path relative to the workspace root", + }, + "content": {"type": "string"}, + }, + }), + read_only: false, + core: true, + governance: GovernanceClass::Approval, + executor: ExecutorKind::InProcess, + }, + ] +} + +/// Confines `relative` to `root`: rejects an absolute path and any `..` +/// component instead of relying on `canonicalize`, which requires the +/// target to already exist and so cannot guard a write to a new file. +fn resolve_within(root: &Path, relative: &str) -> Result { + if relative.trim().is_empty() { + return Err("invalid call: args.path must not be empty".to_owned()); + } + let candidate = Path::new(relative); + let mut resolved = root.to_path_buf(); + for component in candidate.components() { + match component { + std::path::Component::Normal(part) => resolved.push(part), + std::path::Component::CurDir => {} + std::path::Component::ParentDir => { + return Err(format!("invalid call: {relative:?} must not contain ..")); + } + std::path::Component::RootDir | std::path::Component::Prefix(_) => { + return Err(format!( + "invalid call: {relative:?} must be relative to the workspace root" + )); + } + } + } + Ok(resolved) +} + +fn string_arg<'a>(args: &'a Value, key: &str) -> Result<&'a str, String> { + args.get(key) + .and_then(Value::as_str) + .ok_or_else(|| format!("invalid call: args.{key} must be a string")) +} + +async fn read_file(root: &Path, args: &Value) -> ToolResult { + let path = match string_arg(args, "path") { + Ok(path) => path, + Err(reason) => return ToolResult::error(reason), + }; + let resolved = match resolve_within(root, path) { + Ok(resolved) => resolved, + Err(reason) => return ToolResult::error(reason), + }; + match tokio::fs::read_to_string(&resolved).await { + Ok(contents) => ToolResult::text(contents), + Err(error) => ToolResult::error(format!("could not read {path}: {error}")), + } +} + +async fn write_file(root: &Path, args: &Value) -> ToolResult { + let path = match string_arg(args, "path") { + Ok(path) => path, + Err(reason) => return ToolResult::error(reason), + }; + let content = match string_arg(args, "content") { + Ok(content) => content, + Err(reason) => return ToolResult::error(reason), + }; + let resolved = match resolve_within(root, path) { + Ok(resolved) => resolved, + Err(reason) => return ToolResult::error(reason), + }; + if let Some(parent) = resolved.parent() { + if let Err(error) = tokio::fs::create_dir_all(parent).await { + return ToolResult::error(format!("could not create {}: {error}", parent.display())); + } + } + match tokio::fs::write(&resolved, content).await { + Ok(()) => ToolResult::text(format!("wrote {} bytes to {path}", content.len())), + Err(error) => ToolResult::error(format!("could not write {path}: {error}")), + } +} + +fn approval_summary(call: &ProposedCall) -> String { + match call.args.get("path").and_then(Value::as_str) { + Some(path) => format!("Write to {path}"), + None => "Write to a workspace file".to_owned(), + } +} + +/// The `Tools` port over the workspace filesystem at `root`. +#[derive(Clone)] +pub struct LocalTools { + root: Arc, + catalog: Arc<[ToolSpec]>, +} + +impl LocalTools { + pub fn new(root: impl AsRef) -> Self { + Self { + root: Arc::new(root.as_ref().to_path_buf()), + catalog: Arc::from(catalog()), + } + } +} + +impl Tools for LocalTools { + fn catalog(&self) -> &[ToolSpec] { + &self.catalog + } + + async fn search(&self, _principal: &PrincipalId, query: &str) -> Vec { + let query = query.to_ascii_lowercase(); + if query.trim().is_empty() { + return Vec::new(); + } + self.catalog + .iter() + .filter(|spec| spec.label.to_ascii_lowercase().contains(&query)) + .map(|spec| spec.name.clone()) + .collect() + } + + async fn policy(&self, _ctx: &Context, call: &ProposedCall) -> Verdict { + match self.spec(&call.tool).map(|spec| spec.governance) { + Some(GovernanceClass::Approval) => Verdict::NeedsApproval { + approval: ApprovalId::new(format!("approve-{}", call.id)), + summary: approval_summary(call), + }, + _ => Verdict::Allow, + } + } + + async fn run( + &self, + _thread: &ThreadId, + call: &ProposedCall, + _cancel: &CancellationToken, + ) -> ToolResult { + match call.tool.as_str() { + READ_FILE => read_file(&self.root, &call.args).await, + WRITE_FILE => write_file(&self.root, &call.args).await, + other => ToolResult::error(format!("unknown local tool: {other}")), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use dex_loop::Outcome; + use tempfile::TempDir; + + fn thread() -> ThreadId { + ThreadId { + org: "org-1".into(), + workspace: "ws-1".into(), + thread: "thread-1".into(), + } + } + + fn call(tool: &str, args: Value) -> ProposedCall { + ProposedCall::new( + dex_loop::CallId::new("c1"), + ToolName::new(tool), + args, + PrincipalId::new("alice"), + ) + } + + #[tokio::test] + async fn read_file_returns_contents() { + let dir = TempDir::new().expect("tempdir"); + std::fs::write(dir.path().join("a.txt"), "hello").expect("seed file"); + let tools = LocalTools::new(dir.path()); + let cancel = CancellationToken::new(); + let result = tools + .run( + &thread(), + &call(READ_FILE, serde_json::json!({"path": "a.txt"})), + &cancel, + ) + .await; + assert_eq!(result, ToolResult::text("hello")); + } + + #[tokio::test] + async fn read_file_rejects_path_traversal() { + let dir = TempDir::new().expect("tempdir"); + let tools = LocalTools::new(dir.path()); + let cancel = CancellationToken::new(); + let result = tools + .run( + &thread(), + &call(READ_FILE, serde_json::json!({"path": "../escape.txt"})), + &cancel, + ) + .await; + assert_eq!(result.outcome, Outcome::Failed); + } + + #[tokio::test] + async fn write_file_creates_parent_directories_and_writes_content() { + let dir = TempDir::new().expect("tempdir"); + let tools = LocalTools::new(dir.path()); + let cancel = CancellationToken::new(); + let result = tools + .run( + &thread(), + &call( + WRITE_FILE, + serde_json::json!({"path": "nested/b.txt", "content": "world"}), + ), + &cancel, + ) + .await; + assert_eq!(result.outcome, Outcome::Succeeded); + assert_eq!( + std::fs::read_to_string(dir.path().join("nested/b.txt")).expect("read back"), + "world" + ); + } + + #[tokio::test] + async fn write_file_requires_approval_and_read_file_does_not() { + let tools = LocalTools::new("."); + let read_verdict = tools + .policy( + &fake_ctx(), + &call(READ_FILE, serde_json::json!({"path": "a"})), + ) + .await; + assert_eq!(read_verdict, Verdict::Allow); + let write_verdict = tools + .policy( + &fake_ctx(), + &call(WRITE_FILE, serde_json::json!({"path": "a", "content": "x"})), + ) + .await; + assert!(matches!(write_verdict, Verdict::NeedsApproval { .. })); + } + + /// `Tools::policy` here does not read `ctx`, so an empty rehydrated + /// context is enough to exercise it without a real thread history. + fn fake_ctx() -> Context { + dex_loop::rehydrate(thread(), &[]) + } +} diff --git a/packages/dex-host-rs/src/turn.rs b/packages/dex-host-rs/src/turn.rs new file mode 100644 index 000000000..231ac8296 --- /dev/null +++ b/packages/dex-host-rs/src/turn.rs @@ -0,0 +1,355 @@ +//! The local-turn consumer: wires `LocalLog`/`LocalTools`/`LocalEffects` +//! plus a `dex_loop::Model` into one full `dex_loop::Engine` turn, driven to +//! completion (`Exit::Done`, `Exit::Interrupted` or `Exit::Failed`) without a +//! human in the loop. +//! +//! This is Maestro's `MAESTRO_DEX_LOOP=1` cutover step: a local turn now has +//! a consumer of this crate's ports, not just `tests/turn.rs`'s own +//! hand-driven approve/resume. Parked approvals are auto-approved and parked +//! questions are auto-answered with [`UNATTENDED_ANSWER`] -- the same +//! unattended-turn policy `print_mode.rs` (Maestro's existing non-interactive +//! "auto-approves tools" entry point) and `cloud_cli.rs`'s attached REPL +//! already use. This does not change trust posture: it matches the mode +//! this consumer's caller opts into. +//! +//! `Model` stays generic so tests exercise this driver against a scripted +//! fake instead of a live provider; production callers pass [`crate::AiRsModel`]. + +use std::path::Path; + +use dex_loop::{ + ApprovalMode, Budget, CancellationToken, Event, Exit, Lexicon, Log as _, Model, PrincipalId, + ThreadId, TurnId, rehydrate, +}; + +use crate::{LocalEffects, LocalLog, LocalTools}; + +/// Answer given to a parked `Question` when no human is attached to the +/// turn. Mirrors `cloud_cli::UNATTENDED_USER_INPUT_ANSWER`. +pub const UNATTENDED_ANSWER: &str = + "No user is watching this turn. Proceed with your best judgment and state the assumption."; + +/// One local turn to run. +pub struct LocalTurnRequest { + pub thread: ThreadId, + pub principal: PrincipalId, + pub turn: TurnId, + pub text: String, +} + +/// What running one local turn produced. +#[derive(Debug, Clone, PartialEq)] +pub struct LocalTurnOutcome { + pub exit: Exit, + /// The model's final answer, when the turn reached `Exit::Done`. + pub final_text: Option, +} + +/// Runs `request` to completion against a fresh `LocalLog`/`LocalTools`/ +/// `LocalEffects` rooted at `state_root`/`workspace_root`, auto-approving +/// every parked approval and auto-answering every parked question. +/// +/// # Errors +/// Returns an error if the log cannot be acquired (e.g. its lease is held by +/// another process), if a write loses that lease mid-turn (`Fenced`), or if +/// the engine reports `Exit::Failed`. +pub async fn run_local_turn( + state_root: &Path, + workspace_root: &Path, + model: M, + request: LocalTurnRequest, +) -> anyhow::Result { + let log = LocalLog::acquire(state_root.join("log"), &request.thread) + .await + .map_err(|error| anyhow::anyhow!("acquire the local dex-loop log: {error}"))?; + log.append(&[Event::UserMessage { + turn: request.turn, + message_id: None, + principal: request.principal.clone(), + text: request.text, + attachments: Vec::new(), + client_tools: Vec::new(), + authorized_tools: Vec::new(), + approval_mode: ApprovalMode::Interactive, + }]) + .await + .map_err(|_fenced| { + anyhow::anyhow!("the local dex-loop log's lease was superseded before the turn started") + })?; + + let tools = LocalTools::new(workspace_root); + let effects = LocalEffects::open(state_root.join("effects.json")) + .await + .map_err(|error| anyhow::anyhow!("open the local dex-loop effects ledger: {error}"))?; + let engine = dex_loop::Engine::new( + log.clone(), + model, + tools, + effects, + Lexicon::default(), + Budget::default(), + ); + let cancel = CancellationToken::new(); + + let exit = + drive_to_completion(&engine, &log, &request.thread, &request.principal, &cancel).await?; + let final_text = if exit == Exit::Done { + final_answer_text(&log).await? + } else { + None + }; + Ok(LocalTurnOutcome { exit, final_text }) +} + +type LocalEngine = dex_loop::Engine; + +/// Runs `engine` until it reaches a terminal `Exit`, auto-resolving every +/// `Parked` (approval) and `Asked` (question) it hits along the way. +/// `AwaitingClientTool` has no local consumer yet -- client-side tools are +/// not part of this crate's two-tool slice (`tools.rs`) -- so it is reported +/// as an error instead of hanging forever. +async fn drive_to_completion( + engine: &LocalEngine, + log: &LocalLog, + thread: &ThreadId, + principal: &PrincipalId, + cancel: &CancellationToken, +) -> anyhow::Result { + let mut ctx = rehydrate(thread.clone(), &read_log(log).await?); + loop { + let exit = engine + .run(&mut ctx, cancel) + .await + .map_err(|_fenced| anyhow::anyhow!("the local dex-loop log's lease was superseded"))?; + match exit { + Exit::Done | Exit::Interrupted | Exit::Failed => return Ok(exit), + Exit::AwaitingClientTool(call) => { + anyhow::bail!( + "call {call} is waiting on a client-side tool result; the local dex-loop \ + consumer has no client session to answer it" + ); + } + Exit::Parked(approval) => { + let entries = read_log(log).await?; + let pending = entries.iter().find_map(|(_, event)| match event { + Event::ApprovalRequested { + call, + approval: pending_approval, + args_digest, + .. + } if *pending_approval == approval => Some((call.clone(), args_digest.clone())), + _ => None, + }); + let Some((call, args_digest)) = pending else { + anyhow::bail!("no ApprovalRequested event found for {approval}"); + }; + log.append(&[Event::ApprovalDecided { + call, + approval, + args_digest, + approved: true, + principal: principal.clone(), + }]) + .await + .map_err(|_fenced| { + anyhow::anyhow!( + "the local dex-loop log's lease was superseded while auto-approving" + ) + })?; + ctx = rehydrate(thread.clone(), &read_log(log).await?); + } + // `LocalTools`' two-tool catalog (`tools.rs`) has no + // `ExecutorKind::User` tool today, so the engine cannot actually + // reach this arm through `run_local_turn` yet; it is here so a + // future ask-the-user tool does not silently hang instead of + // getting the same unattended answer `Parked` gets. + Exit::Asked(call) => { + log.append(&[Event::Answer { + call, + principal: principal.clone(), + text: UNATTENDED_ANSWER.to_owned(), + }]) + .await + .map_err(|_fenced| { + anyhow::anyhow!( + "the local dex-loop log's lease was superseded while auto-answering" + ) + })?; + ctx = rehydrate(thread.clone(), &read_log(log).await?); + } + } + } +} + +async fn read_log(log: &LocalLog) -> anyhow::Result> { + log.read_all() + .await + .map_err(|error| anyhow::anyhow!("read the local dex-loop log: {error}")) +} + +async fn final_answer_text(log: &LocalLog) -> anyhow::Result> { + let entries = read_log(log).await?; + Ok(entries + .into_iter() + .rev() + .find_map(|(_, event)| match event { + Event::Final { text } => Some(text), + _ => None, + })) +} + +#[cfg(test)] +mod tests { + use std::collections::VecDeque; + use std::sync::{Arc, Mutex, PoisonError}; + + use dex_loop::{ModelChunk, ModelError, ToolName}; + use futures_util::{Stream, stream}; + use tempfile::TempDir; + + use super::*; + use crate::{READ_FILE, WRITE_FILE}; + + /// One scripted model call: the chunks it streams, in order. Matches + /// `tests/turn.rs`'s own `Script` alias. + type Script = Vec>; + + /// Replays one scripted model call's chunks per `stream()` invocation, + /// the same double `dex-loop`'s own tests and `tests/turn.rs` use. + #[derive(Clone, Default)] + struct ScriptedModel { + scripts: Arc>>, + } + + impl ScriptedModel { + fn new(scripts: Vec