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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions .repository-projection.json
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
87 changes: 87 additions & 0 deletions packages/local-host-rs/examples/semantic_stream_fixture.rs
Original file line number Diff line number Diff line change
@@ -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("**:"))
}
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
4 changes: 4 additions & 0 deletions packages/local-host-rs/src/headless/messages.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
1 change: 1 addition & 0 deletions packages/local-host-rs/src/headless/messages/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -826,6 +826,7 @@ impl AgentState {
None
}
FromAgentMessage::ProcessBudgetCheckpoint { .. } => None,
FromAgentMessage::AssistantTextObserved { .. } => None,
FromAgentMessage::ManagedGatewayReceipt {
request_id,
record_id,
Expand Down
6 changes: 6 additions & 0 deletions packages/local-host-rs/src/headless/proto.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down Expand Up @@ -610,6 +613,7 @@ mod tests {
| FromAgentMessage::WorkspaceCapabilitySetApplied { .. }
| FromAgentMessage::Ready { .. }
| FromAgentMessage::ResponseStart { .. }
| FromAgentMessage::AssistantTextObserved { .. }
| FromAgentMessage::ResponseChunk { .. }
| FromAgentMessage::ResponseEnd { .. }
| FromAgentMessage::TurnCompleted { .. }
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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()),
Expand Down
118 changes: 58 additions & 60 deletions packages/local-host-rs/src/headless_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
///
Expand Down Expand Up @@ -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<String>,
semantic_started_at: Option<Instant>,
semantic_raw_first_ms: Option<u64>,
semantic_first_released_ms: Option<u64>,
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<maestro_ai::ManagedGatewayReceipt>,
/// Controller assignment proven against the provider-bound system prompt.
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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())),
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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 })?;
Expand Down Expand Up @@ -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 })?;
Expand All @@ -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)
Expand All @@ -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)?;
}
Expand Down Expand Up @@ -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(
Expand Down
Loading
Loading