refactor(agent): split oversized service modules
This commit is contained in:
parent
2e9de6d2a2
commit
30b312f32f
|
|
@ -11,10 +11,8 @@
|
|||
//! Everything here is blocking and runs on the probe worker thread
|
||||
//! (`provider_probe_host.rs`); nothing touches the UI thread.
|
||||
|
||||
use std::io::{BufRead, BufReader};
|
||||
use std::path::Path;
|
||||
use std::process::{Command, Stdio};
|
||||
use std::sync::mpsc;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use base64::Engine;
|
||||
|
|
@ -23,14 +21,21 @@ use op_ai::chat_models::ModelEntry;
|
|||
use op_i18n::Locale;
|
||||
|
||||
use crate::model_discovery::{
|
||||
codex_models_from_app_server, codex_models_from_cache, discover_opencode, extract_json_object,
|
||||
parse_copilot_model_list, resolve_cli, write_lsp_frame,
|
||||
codex_models_from_app_server, codex_models_from_cache, discover_opencode, resolve_cli,
|
||||
};
|
||||
use crate::provider_probe_models::{
|
||||
claude_initialize_query, codex_home, codex_models_from_latest_md, ClaudeAccount,
|
||||
ClaudeInitResult,
|
||||
};
|
||||
|
||||
#[path = "provider_probe_copilot.rs"]
|
||||
mod copilot;
|
||||
use copilot::connect_copilot;
|
||||
#[cfg(test)]
|
||||
use copilot::{
|
||||
copilot_connection_info, friendly_copilot_error, parse_copilot_auth_status, CopilotAuth,
|
||||
};
|
||||
|
||||
/// Probe result — the TS `ConnectResult` shape plus the install
|
||||
/// guidance the not-installed path surfaces.
|
||||
#[derive(Debug, Clone, Default)]
|
||||
|
|
@ -656,179 +661,6 @@ fn opencode_provider_summary(locale: Locale, models: &[ModelEntry]) -> String {
|
|||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------
|
||||
// GitHub Copilot (connect-agent.ts:742-839)
|
||||
// ---------------------------------------------------------------
|
||||
|
||||
/// Auth status from `auth.getStatus` — the wire method the official
|
||||
/// copilot-sdk's `getAuthStatus()` sends (client.js:549-555).
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||
struct CopilotAuth {
|
||||
login: Option<String>,
|
||||
auth_type: Option<String>,
|
||||
status_message: Option<String>,
|
||||
}
|
||||
|
||||
const COPILOT_PROBE_TIMEOUT: Duration = Duration::from_secs(8);
|
||||
|
||||
fn connect_copilot(locale: Locale) -> ProbeOutcome {
|
||||
let Some(exe) = resolve_cli("copilot") else {
|
||||
return ProbeOutcome::not_installed(
|
||||
AgentProvider::GithubCopilot,
|
||||
tw(
|
||||
locale,
|
||||
"providerProbe.cliNotFound",
|
||||
&[("name", "GitHub Copilot")],
|
||||
),
|
||||
);
|
||||
};
|
||||
let Some((models, auth)) = copilot_probe_stdio(&exe) else {
|
||||
return ProbeOutcome::failed(friendly_copilot_error(locale, "Connection timed out"));
|
||||
};
|
||||
if models.is_empty() {
|
||||
return ProbeOutcome::failed(t(locale, "providerProbe.noModelsCopilot"));
|
||||
}
|
||||
let hint = config_path(
|
||||
"~/.config/github-copilot/config.json",
|
||||
"%USERPROFILE%\\.config\\github-copilot\\config.json",
|
||||
);
|
||||
let info = copilot_connection_info(locale, auth.as_ref());
|
||||
ProbeOutcome {
|
||||
connected: true,
|
||||
models,
|
||||
connection_info: Some(info),
|
||||
hint_path: Some(hint),
|
||||
..ProbeOutcome::default()
|
||||
}
|
||||
}
|
||||
|
||||
/// TS connectCopilot's status mapping (connect-agent.ts:799-819).
|
||||
fn copilot_connection_info(locale: Locale, auth: Option<&CopilotAuth>) -> String {
|
||||
if let Some(auth) = auth {
|
||||
if let Some(login) = auth.login.as_deref() {
|
||||
let method = auth
|
||||
.auth_type
|
||||
.as_deref()
|
||||
.map(|t| format!(" ({t})"))
|
||||
.unwrap_or_default();
|
||||
return tw(
|
||||
locale,
|
||||
"providerProbe.connectedAs",
|
||||
&[("login", login), ("method", &method)],
|
||||
);
|
||||
}
|
||||
if let Some(message) = auth.status_message.as_deref() {
|
||||
return message.to_string();
|
||||
}
|
||||
}
|
||||
t(locale, "providerProbe.connectedViaGithub")
|
||||
}
|
||||
|
||||
/// One `copilot --stdio` session: `connect` (id 1) → `models.list`
|
||||
/// (id 2) → `auth.getStatus` (id 3), all pipelined. Returns `None`
|
||||
/// when the model list never answered; auth is best-effort (TS
|
||||
/// logs and continues when `getAuthStatus` fails).
|
||||
fn copilot_probe_stdio(exe: &Path) -> Option<(Vec<ModelEntry>, Option<CopilotAuth>)> {
|
||||
let mut cmd = Command::new(exe);
|
||||
cmd.arg("--stdio")
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::null());
|
||||
crate::chat_spawn::hide_console_window(&mut cmd);
|
||||
let mut child = cmd.spawn().ok()?;
|
||||
let mut stdin = child.stdin.take()?;
|
||||
let stdout = child.stdout.take()?;
|
||||
|
||||
let (tx, rx) = mpsc::channel();
|
||||
std::thread::spawn(move || {
|
||||
for line in BufReader::new(stdout).lines().map_while(Result::ok) {
|
||||
if tx.send(line).is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
let connect = r#"{"jsonrpc":"2.0","id":1,"method":"connect","params":{}}"#;
|
||||
let list = r#"{"jsonrpc":"2.0","id":2,"method":"models.list","params":{}}"#;
|
||||
let auth = r#"{"jsonrpc":"2.0","id":3,"method":"auth.getStatus","params":{}}"#;
|
||||
write_lsp_frame(&mut stdin, connect).ok()?;
|
||||
write_lsp_frame(&mut stdin, list).ok()?;
|
||||
write_lsp_frame(&mut stdin, auth).ok()?;
|
||||
use std::io::Write as _;
|
||||
stdin.flush().ok();
|
||||
|
||||
let deadline = Instant::now() + COPILOT_PROBE_TIMEOUT;
|
||||
let mut models: Option<Vec<ModelEntry>> = None;
|
||||
let mut auth_status: Option<CopilotAuth> = None;
|
||||
while let Some(remaining) = deadline.checked_duration_since(Instant::now()) {
|
||||
if models.is_some() && auth_status.is_some() {
|
||||
break;
|
||||
}
|
||||
match rx.recv_timeout(remaining) {
|
||||
Ok(line) => {
|
||||
if models.is_none() {
|
||||
if let Some(parsed) = parse_copilot_model_list(&line) {
|
||||
models = Some(parsed);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
if auth_status.is_none() {
|
||||
if let Some(parsed) = parse_copilot_auth_status(&line) {
|
||||
auth_status = Some(parsed);
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(_) => break,
|
||||
}
|
||||
}
|
||||
let _ = child.kill();
|
||||
let _ = child.wait();
|
||||
models.map(|m| (m, auth_status))
|
||||
}
|
||||
|
||||
/// Parse the `id:3` (`auth.getStatus`) response.
|
||||
fn parse_copilot_auth_status(line: &str) -> Option<CopilotAuth> {
|
||||
let json: serde_json::Value = serde_json::from_str(extract_json_object(line)?).ok()?;
|
||||
if json.get("id")?.as_i64()? != 3 {
|
||||
return None;
|
||||
}
|
||||
let result = json.get("result")?;
|
||||
Some(CopilotAuth {
|
||||
login: result
|
||||
.get("login")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(str::to_string),
|
||||
auth_type: result
|
||||
.get("authType")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(str::to_string),
|
||||
status_message: result
|
||||
.get("statusMessage")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(str::to_string),
|
||||
})
|
||||
}
|
||||
|
||||
/// TS `friendlyCopilotError` (connect-agent.ts:842-853).
|
||||
fn friendly_copilot_error(locale: Locale, raw: &str) -> String {
|
||||
let lower = raw.to_ascii_lowercase();
|
||||
if lower.contains("not found") || lower.contains("enoent") {
|
||||
return t(locale, "providerProbe.copilotNotFoundInstall");
|
||||
}
|
||||
if lower.contains("not authenticated")
|
||||
|| lower.contains("authenticate first")
|
||||
|| lower.contains("auth")
|
||||
|| lower.contains("unauthenticated")
|
||||
|| lower.contains("login")
|
||||
{
|
||||
return t(locale, "providerProbe.notAuthenticatedCopilot");
|
||||
}
|
||||
if lower.contains("timed out") || lower.contains("timedout") {
|
||||
return t(locale, "providerProbe.timedOut");
|
||||
}
|
||||
raw.to_string()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "provider_probe_tests.rs"]
|
||||
mod tests;
|
||||
|
|
|
|||
184
crates/op-host-services/src/provider_probe_copilot.rs
Normal file
184
crates/op-host-services/src/provider_probe_copilot.rs
Normal file
|
|
@ -0,0 +1,184 @@
|
|||
//! GitHub Copilot connect-time probe.
|
||||
|
||||
use std::io::{BufRead, BufReader, Write as _};
|
||||
use std::path::Path;
|
||||
use std::process::{Command, Stdio};
|
||||
use std::sync::mpsc;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use op_ai::agent_settings_state::AgentProvider;
|
||||
use op_ai::chat_models::ModelEntry;
|
||||
use op_i18n::Locale;
|
||||
|
||||
use super::{config_path, t, tw, ProbeOutcome};
|
||||
use crate::model_discovery::{
|
||||
extract_json_object, parse_copilot_model_list, resolve_cli, write_lsp_frame,
|
||||
};
|
||||
|
||||
/// Auth status from `auth.getStatus` — the wire method the official
|
||||
/// copilot-sdk's `getAuthStatus()` sends (client.js:549-555).
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||
pub(super) struct CopilotAuth {
|
||||
pub(super) login: Option<String>,
|
||||
pub(super) auth_type: Option<String>,
|
||||
pub(super) status_message: Option<String>,
|
||||
}
|
||||
|
||||
const COPILOT_PROBE_TIMEOUT: Duration = Duration::from_secs(8);
|
||||
|
||||
pub(super) fn connect_copilot(locale: Locale) -> ProbeOutcome {
|
||||
let Some(exe) = resolve_cli("copilot") else {
|
||||
return ProbeOutcome::not_installed(
|
||||
AgentProvider::GithubCopilot,
|
||||
tw(
|
||||
locale,
|
||||
"providerProbe.cliNotFound",
|
||||
&[("name", "GitHub Copilot")],
|
||||
),
|
||||
);
|
||||
};
|
||||
let Some((models, auth)) = copilot_probe_stdio(&exe) else {
|
||||
return ProbeOutcome::failed(friendly_copilot_error(locale, "Connection timed out"));
|
||||
};
|
||||
if models.is_empty() {
|
||||
return ProbeOutcome::failed(t(locale, "providerProbe.noModelsCopilot"));
|
||||
}
|
||||
let hint = config_path(
|
||||
"~/.config/github-copilot/config.json",
|
||||
"%USERPROFILE%\\.config\\github-copilot\\config.json",
|
||||
);
|
||||
let info = copilot_connection_info(locale, auth.as_ref());
|
||||
ProbeOutcome {
|
||||
connected: true,
|
||||
models,
|
||||
connection_info: Some(info),
|
||||
hint_path: Some(hint),
|
||||
..ProbeOutcome::default()
|
||||
}
|
||||
}
|
||||
|
||||
/// TS connectCopilot's status mapping (connect-agent.ts:799-819).
|
||||
pub(super) fn copilot_connection_info(locale: Locale, auth: Option<&CopilotAuth>) -> String {
|
||||
if let Some(auth) = auth {
|
||||
if let Some(login) = auth.login.as_deref() {
|
||||
let method = auth
|
||||
.auth_type
|
||||
.as_deref()
|
||||
.map(|t| format!(" ({t})"))
|
||||
.unwrap_or_default();
|
||||
return tw(
|
||||
locale,
|
||||
"providerProbe.connectedAs",
|
||||
&[("login", login), ("method", &method)],
|
||||
);
|
||||
}
|
||||
if let Some(message) = auth.status_message.as_deref() {
|
||||
return message.to_string();
|
||||
}
|
||||
}
|
||||
t(locale, "providerProbe.connectedViaGithub")
|
||||
}
|
||||
|
||||
/// One `copilot --stdio` session: `connect` (id 1) → `models.list`
|
||||
/// (id 2) → `auth.getStatus` (id 3), all pipelined. Returns `None`
|
||||
/// when the model list never answered; auth is best-effort (TS
|
||||
/// logs and continues when `getAuthStatus` fails).
|
||||
fn copilot_probe_stdio(exe: &Path) -> Option<(Vec<ModelEntry>, Option<CopilotAuth>)> {
|
||||
let mut cmd = Command::new(exe);
|
||||
cmd.arg("--stdio")
|
||||
.stdin(Stdio::piped())
|
||||
.stdout(Stdio::piped())
|
||||
.stderr(Stdio::null());
|
||||
crate::chat_spawn::hide_console_window(&mut cmd);
|
||||
let mut child = cmd.spawn().ok()?;
|
||||
let mut stdin = child.stdin.take()?;
|
||||
let stdout = child.stdout.take()?;
|
||||
|
||||
let (tx, rx) = mpsc::channel();
|
||||
std::thread::spawn(move || {
|
||||
for line in BufReader::new(stdout).lines().map_while(Result::ok) {
|
||||
if tx.send(line).is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
let connect = r#"{"jsonrpc":"2.0","id":1,"method":"connect","params":{}}"#;
|
||||
let list = r#"{"jsonrpc":"2.0","id":2,"method":"models.list","params":{}}"#;
|
||||
let auth = r#"{"jsonrpc":"2.0","id":3,"method":"auth.getStatus","params":{}}"#;
|
||||
write_lsp_frame(&mut stdin, connect).ok()?;
|
||||
write_lsp_frame(&mut stdin, list).ok()?;
|
||||
write_lsp_frame(&mut stdin, auth).ok()?;
|
||||
stdin.flush().ok();
|
||||
|
||||
let deadline = Instant::now() + COPILOT_PROBE_TIMEOUT;
|
||||
let mut models: Option<Vec<ModelEntry>> = None;
|
||||
let mut auth_status: Option<CopilotAuth> = None;
|
||||
while let Some(remaining) = deadline.checked_duration_since(Instant::now()) {
|
||||
if models.is_some() && auth_status.is_some() {
|
||||
break;
|
||||
}
|
||||
match rx.recv_timeout(remaining) {
|
||||
Ok(line) => {
|
||||
if models.is_none() {
|
||||
if let Some(parsed) = parse_copilot_model_list(&line) {
|
||||
models = Some(parsed);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
if auth_status.is_none() {
|
||||
if let Some(parsed) = parse_copilot_auth_status(&line) {
|
||||
auth_status = Some(parsed);
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(_) => break,
|
||||
}
|
||||
}
|
||||
let _ = child.kill();
|
||||
let _ = child.wait();
|
||||
models.map(|m| (m, auth_status))
|
||||
}
|
||||
|
||||
/// Parse the `id:3` (`auth.getStatus`) response.
|
||||
pub(super) fn parse_copilot_auth_status(line: &str) -> Option<CopilotAuth> {
|
||||
let json: serde_json::Value = serde_json::from_str(extract_json_object(line)?).ok()?;
|
||||
if json.get("id")?.as_i64()? != 3 {
|
||||
return None;
|
||||
}
|
||||
let result = json.get("result")?;
|
||||
Some(CopilotAuth {
|
||||
login: result
|
||||
.get("login")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(str::to_string),
|
||||
auth_type: result
|
||||
.get("authType")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(str::to_string),
|
||||
status_message: result
|
||||
.get("statusMessage")
|
||||
.and_then(|v| v.as_str())
|
||||
.map(str::to_string),
|
||||
})
|
||||
}
|
||||
|
||||
/// TS `friendlyCopilotError` (connect-agent.ts:842-853).
|
||||
pub(super) fn friendly_copilot_error(locale: Locale, raw: &str) -> String {
|
||||
let lower = raw.to_ascii_lowercase();
|
||||
if lower.contains("not found") || lower.contains("enoent") {
|
||||
return t(locale, "providerProbe.copilotNotFoundInstall");
|
||||
}
|
||||
if lower.contains("not authenticated")
|
||||
|| lower.contains("authenticate first")
|
||||
|| lower.contains("auth")
|
||||
|| lower.contains("unauthenticated")
|
||||
|| lower.contains("login")
|
||||
{
|
||||
return t(locale, "providerProbe.notAuthenticatedCopilot");
|
||||
}
|
||||
if lower.contains("timed out") || lower.contains("timedout") {
|
||||
return t(locale, "providerProbe.timedOut");
|
||||
}
|
||||
raw.to_string()
|
||||
}
|
||||
|
|
@ -10,9 +10,9 @@ use std::io::Write;
|
|||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use base64::Engine as _;
|
||||
use op_ai::chat_provider::{
|
||||
ChatAttachment, ChatDelta, ChatHistoryRole, ChatProvider, ChatRequest, StopReason,
|
||||
};
|
||||
#[cfg(test)]
|
||||
use op_ai::chat_provider::StopReason;
|
||||
use op_ai::chat_provider::{ChatAttachment, ChatDelta, ChatHistoryRole, ChatProvider, ChatRequest};
|
||||
use op_editor_core::chat::MAX_ATTACHMENT_BYTES;
|
||||
use op_editor_core::{BuiltinAgentConfig, EditorCommand, EditorState, NodeId};
|
||||
use op_orchestrator::{
|
||||
|
|
@ -30,6 +30,13 @@ use crate::web_canvas_server::{SseHub, WebCanvasState};
|
|||
mod error;
|
||||
use error::WebChatStandardError;
|
||||
|
||||
#[path = "web_chat_standard_events.rs"]
|
||||
mod events;
|
||||
use events::{
|
||||
progress_label, web_identity_seed, write_agent_identity_event, write_delta_event,
|
||||
write_done_event, write_error_event, write_thinking_event,
|
||||
};
|
||||
|
||||
const STANDARD_MODIFY_STEP: &str =
|
||||
r#"<step title="Checking guidelines">Analyzing modification request...</step>"#;
|
||||
|
||||
|
|
@ -619,197 +626,6 @@ impl DocSink for WebDesignDocSink<'_> {
|
|||
fn end_undo_batch(&mut self) {}
|
||||
}
|
||||
|
||||
fn write_delta_event<W: Write>(out: &mut W, text: &str) -> std::io::Result<()> {
|
||||
out.write_all(
|
||||
crate::ai_proxy::delta_to_sse(&ChatDelta::TextDelta(text.to_string())).as_bytes(),
|
||||
)?;
|
||||
out.flush()
|
||||
}
|
||||
|
||||
fn write_thinking_event<W: Write>(out: &mut W, text: &str) -> std::io::Result<()> {
|
||||
out.write_all(
|
||||
crate::ai_proxy::delta_to_sse(&ChatDelta::Thinking(text.to_string())).as_bytes(),
|
||||
)?;
|
||||
out.flush()
|
||||
}
|
||||
|
||||
fn write_agent_identity_event<W: Write>(
|
||||
out: &mut W,
|
||||
identity: &op_orchestrator::agent_identity::AgentIdentity,
|
||||
) -> std::io::Result<()> {
|
||||
let payload = serde_json::json!({
|
||||
"agent": {
|
||||
"name": identity.name,
|
||||
"color": identity.color,
|
||||
}
|
||||
});
|
||||
out.write_all(format!("data: {payload}\n\n").as_bytes())?;
|
||||
out.flush()
|
||||
}
|
||||
|
||||
fn web_identity_seed() -> u64 {
|
||||
std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.map(|duration| duration.subsec_nanos() as u64)
|
||||
.unwrap_or(0)
|
||||
}
|
||||
|
||||
fn write_done_event<W: Write>(out: &mut W) -> std::io::Result<()> {
|
||||
out.write_all(
|
||||
crate::ai_proxy::delta_to_sse(&ChatDelta::Done {
|
||||
stop_reason: StopReason::EndTurn,
|
||||
})
|
||||
.as_bytes(),
|
||||
)?;
|
||||
out.flush()
|
||||
}
|
||||
|
||||
fn write_error_event<W: Write>(out: &mut W, message: &str) -> std::io::Result<()> {
|
||||
out.write_all(
|
||||
crate::ai_proxy::delta_to_sse(&ChatDelta::Error(message.to_string())).as_bytes(),
|
||||
)?;
|
||||
out.flush()
|
||||
}
|
||||
|
||||
fn progress_label(p: &Progress) -> String {
|
||||
match p {
|
||||
Progress::Planning => "• Planning…".into(),
|
||||
Progress::Planned { subtasks } => {
|
||||
format!("• Plan — {} section(s)", subtasks.len())
|
||||
}
|
||||
Progress::ScaffoldDone => "• Scaffold ready".into(),
|
||||
Progress::SubtaskStarted { id, label } => format!("• Subtask `{id}` — {label}"),
|
||||
Progress::SubtaskDone { id, node_count } => {
|
||||
format!("• Subtask `{id}` done ({node_count} nodes)")
|
||||
}
|
||||
Progress::SubtaskFailed { id, error } => format!("• Subtask `{id}` failed: {error}"),
|
||||
Progress::SubtaskSkills {
|
||||
id,
|
||||
included,
|
||||
dropped,
|
||||
budget_used,
|
||||
budget_max,
|
||||
} => format_subtask_skills(id, included, dropped, *budget_used, *budget_max),
|
||||
Progress::SubtaskRetry {
|
||||
attempt, reason, ..
|
||||
} => {
|
||||
format!(" ▸ retry #{attempt}: {reason}")
|
||||
}
|
||||
Progress::GeometryEcho { issue_count, .. } => {
|
||||
format!(" ▸ geometry echo: {issue_count} issue(s) → retry")
|
||||
}
|
||||
Progress::SubtaskNodes { id, nodes_so_far } => {
|
||||
format!("• Subtask `{id}` — {nodes_so_far} node(s) so far")
|
||||
}
|
||||
Progress::ConcurrentGroupsStarted {
|
||||
group_count,
|
||||
workers,
|
||||
} => format!("• {group_count} screen groups · {workers} workers"),
|
||||
Progress::ScreenGroupsSequential {
|
||||
group_count,
|
||||
requested_workers,
|
||||
} => format!(
|
||||
"• {group_count} screen groups · sequential (parallel setting: {requested_workers})"
|
||||
),
|
||||
Progress::WorkerScoped(worker) => {
|
||||
let detail = progress_label(worker.event.as_ref());
|
||||
format!(
|
||||
"• {} · {} — {}",
|
||||
worker.identity.name,
|
||||
worker.screen,
|
||||
detail.trim_start_matches("• ")
|
||||
)
|
||||
}
|
||||
Progress::CleanupDone => "• Cleanup done".into(),
|
||||
// The classic path's quality credential. `remaining` is `None` on
|
||||
// purpose: the promise-delivery check runs later in the pipeline, so
|
||||
// claiming anything about leftover work here would be a guess.
|
||||
Progress::QualityChecked { checks, repairs } => {
|
||||
crate::quality_credential::quality_credential_line(
|
||||
&op_ai::chat_provider::QualitySummary {
|
||||
checks: checks.clone(),
|
||||
repairs: repairs.clone(),
|
||||
},
|
||||
None,
|
||||
)
|
||||
.unwrap_or_default()
|
||||
.trim_start()
|
||||
.to_string()
|
||||
}
|
||||
Progress::ValidationStarted => "• Validation started".into(),
|
||||
Progress::ValidationPreCheckDone { applied, .. } => {
|
||||
format!("• Pre-validation applied {applied} fix(es)")
|
||||
}
|
||||
Progress::ValidationRoundStarted { round } => format!("• Vision round {round} started"),
|
||||
Progress::ValidationRoundDone {
|
||||
round,
|
||||
applied,
|
||||
quality_score,
|
||||
} => {
|
||||
format!("• Vision round {round} done — {applied} fix(es), quality {quality_score}/100")
|
||||
}
|
||||
Progress::ValidationDone { total_applied } => {
|
||||
format!("• Validation done — {total_applied} fix(es) total")
|
||||
}
|
||||
Progress::VisualRefStarted => "• Visual-ref pipeline started".into(),
|
||||
Progress::VisualRefDesignSystem { var_count } => {
|
||||
format!("• Design system ready — {var_count} variable(s) seeded")
|
||||
}
|
||||
Progress::VisualRefHtmlGenerated { byte_len } => {
|
||||
format!("• Visual-ref HTML generated ({byte_len} bytes)")
|
||||
}
|
||||
Progress::VisualRefScreenshotReady { skipped } => {
|
||||
if *skipped {
|
||||
"• Visual-ref screenshot skipped".into()
|
||||
} else {
|
||||
"• Visual-ref screenshot captured".into()
|
||||
}
|
||||
}
|
||||
Progress::VisualRefFallback { reason } => format!("• Visual-ref fallback: {reason}"),
|
||||
Progress::UnfilledScreens { names } => {
|
||||
format!(
|
||||
"• {} screen(s) left unfilled: {}",
|
||||
names.len(),
|
||||
names.join(", ")
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Format a `SubtaskSkills` payload into the concise summary line plus
|
||||
/// indented `▸ skills:` / `▸ dropped:` detail sub-lines (spec Component 5).
|
||||
fn format_subtask_skills(
|
||||
id: &str,
|
||||
included: &[op_orchestrator::SkillBrief],
|
||||
dropped: &[(String, String)],
|
||||
budget_used: u32,
|
||||
budget_max: u32,
|
||||
) -> String {
|
||||
let mut out = format!(
|
||||
"• Subtask `{id}` · {} skills · {budget_used}/{budget_max} tok · {} dropped",
|
||||
included.len(),
|
||||
dropped.len(),
|
||||
);
|
||||
if !included.is_empty() {
|
||||
let names: Vec<String> = included
|
||||
.iter()
|
||||
.map(|s| {
|
||||
if s.truncated {
|
||||
format!("{} (truncated)", s.name)
|
||||
} else {
|
||||
s.name.clone()
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
out.push_str(&format!("\n ▸ skills: {}", names.join(", ")));
|
||||
}
|
||||
if !dropped.is_empty() {
|
||||
let drops: Vec<String> = dropped.iter().map(|(n, r)| format!("{n} ({r})")).collect();
|
||||
out.push_str(&format!("\n ▸ dropped: {}", drops.join(", ")));
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[path = "web_chat_standard_tests.rs"]
|
||||
mod tests;
|
||||
|
|
|
|||
197
crates/op-host-services/src/web_chat_standard_events.rs
Normal file
197
crates/op-host-services/src/web_chat_standard_events.rs
Normal file
|
|
@ -0,0 +1,197 @@
|
|||
//! SSE event serialization and progress labels for the standard web turn.
|
||||
|
||||
use std::io::Write;
|
||||
|
||||
use op_ai::chat_provider::{ChatDelta, StopReason};
|
||||
use op_orchestrator::Progress;
|
||||
|
||||
pub(super) fn write_delta_event<W: Write>(out: &mut W, text: &str) -> std::io::Result<()> {
|
||||
out.write_all(
|
||||
crate::ai_proxy::delta_to_sse(&ChatDelta::TextDelta(text.to_string())).as_bytes(),
|
||||
)?;
|
||||
out.flush()
|
||||
}
|
||||
|
||||
pub(super) fn write_thinking_event<W: Write>(out: &mut W, text: &str) -> std::io::Result<()> {
|
||||
out.write_all(
|
||||
crate::ai_proxy::delta_to_sse(&ChatDelta::Thinking(text.to_string())).as_bytes(),
|
||||
)?;
|
||||
out.flush()
|
||||
}
|
||||
|
||||
pub(super) fn write_agent_identity_event<W: Write>(
|
||||
out: &mut W,
|
||||
identity: &op_orchestrator::agent_identity::AgentIdentity,
|
||||
) -> std::io::Result<()> {
|
||||
let payload = serde_json::json!({
|
||||
"agent": {
|
||||
"name": identity.name,
|
||||
"color": identity.color,
|
||||
}
|
||||
});
|
||||
out.write_all(format!("data: {payload}\n\n").as_bytes())?;
|
||||
out.flush()
|
||||
}
|
||||
|
||||
pub(super) fn web_identity_seed() -> u64 {
|
||||
std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)
|
||||
.map(|duration| duration.subsec_nanos() as u64)
|
||||
.unwrap_or(0)
|
||||
}
|
||||
|
||||
pub(super) fn write_done_event<W: Write>(out: &mut W) -> std::io::Result<()> {
|
||||
out.write_all(
|
||||
crate::ai_proxy::delta_to_sse(&ChatDelta::Done {
|
||||
stop_reason: StopReason::EndTurn,
|
||||
})
|
||||
.as_bytes(),
|
||||
)?;
|
||||
out.flush()
|
||||
}
|
||||
|
||||
pub(super) fn write_error_event<W: Write>(out: &mut W, message: &str) -> std::io::Result<()> {
|
||||
out.write_all(
|
||||
crate::ai_proxy::delta_to_sse(&ChatDelta::Error(message.to_string())).as_bytes(),
|
||||
)?;
|
||||
out.flush()
|
||||
}
|
||||
|
||||
pub(super) fn progress_label(p: &Progress) -> String {
|
||||
match p {
|
||||
Progress::Planning => "• Planning…".into(),
|
||||
Progress::Planned { subtasks } => {
|
||||
format!("• Plan — {} section(s)", subtasks.len())
|
||||
}
|
||||
Progress::ScaffoldDone => "• Scaffold ready".into(),
|
||||
Progress::SubtaskStarted { id, label } => format!("• Subtask `{id}` — {label}"),
|
||||
Progress::SubtaskDone { id, node_count } => {
|
||||
format!("• Subtask `{id}` done ({node_count} nodes)")
|
||||
}
|
||||
Progress::SubtaskFailed { id, error } => format!("• Subtask `{id}` failed: {error}"),
|
||||
Progress::SubtaskSkills {
|
||||
id,
|
||||
included,
|
||||
dropped,
|
||||
budget_used,
|
||||
budget_max,
|
||||
} => format_subtask_skills(id, included, dropped, *budget_used, *budget_max),
|
||||
Progress::SubtaskRetry {
|
||||
attempt, reason, ..
|
||||
} => {
|
||||
format!(" ▸ retry #{attempt}: {reason}")
|
||||
}
|
||||
Progress::GeometryEcho { issue_count, .. } => {
|
||||
format!(" ▸ geometry echo: {issue_count} issue(s) → retry")
|
||||
}
|
||||
Progress::SubtaskNodes { id, nodes_so_far } => {
|
||||
format!("• Subtask `{id}` — {nodes_so_far} node(s) so far")
|
||||
}
|
||||
Progress::ConcurrentGroupsStarted {
|
||||
group_count,
|
||||
workers,
|
||||
} => format!("• {group_count} screen groups · {workers} workers"),
|
||||
Progress::ScreenGroupsSequential {
|
||||
group_count,
|
||||
requested_workers,
|
||||
} => format!(
|
||||
"• {group_count} screen groups · sequential (parallel setting: {requested_workers})"
|
||||
),
|
||||
Progress::WorkerScoped(worker) => {
|
||||
let detail = progress_label(worker.event.as_ref());
|
||||
format!(
|
||||
"• {} · {} — {}",
|
||||
worker.identity.name,
|
||||
worker.screen,
|
||||
detail.trim_start_matches("• ")
|
||||
)
|
||||
}
|
||||
Progress::CleanupDone => "• Cleanup done".into(),
|
||||
// The classic path's quality credential. `remaining` is `None` on
|
||||
// purpose: the promise-delivery check runs later in the pipeline, so
|
||||
// claiming anything about leftover work here would be a guess.
|
||||
Progress::QualityChecked { checks, repairs } => {
|
||||
crate::quality_credential::quality_credential_line(
|
||||
&op_ai::chat_provider::QualitySummary {
|
||||
checks: checks.clone(),
|
||||
repairs: repairs.clone(),
|
||||
},
|
||||
None,
|
||||
)
|
||||
.unwrap_or_default()
|
||||
.trim_start()
|
||||
.to_string()
|
||||
}
|
||||
Progress::ValidationStarted => "• Validation started".into(),
|
||||
Progress::ValidationPreCheckDone { applied, .. } => {
|
||||
format!("• Pre-validation applied {applied} fix(es)")
|
||||
}
|
||||
Progress::ValidationRoundStarted { round } => format!("• Vision round {round} started"),
|
||||
Progress::ValidationRoundDone {
|
||||
round,
|
||||
applied,
|
||||
quality_score,
|
||||
} => {
|
||||
format!("• Vision round {round} done — {applied} fix(es), quality {quality_score}/100")
|
||||
}
|
||||
Progress::ValidationDone { total_applied } => {
|
||||
format!("• Validation done — {total_applied} fix(es) total")
|
||||
}
|
||||
Progress::VisualRefStarted => "• Visual-ref pipeline started".into(),
|
||||
Progress::VisualRefDesignSystem { var_count } => {
|
||||
format!("• Design system ready — {var_count} variable(s) seeded")
|
||||
}
|
||||
Progress::VisualRefHtmlGenerated { byte_len } => {
|
||||
format!("• Visual-ref HTML generated ({byte_len} bytes)")
|
||||
}
|
||||
Progress::VisualRefScreenshotReady { skipped } => {
|
||||
if *skipped {
|
||||
"• Visual-ref screenshot skipped".into()
|
||||
} else {
|
||||
"• Visual-ref screenshot captured".into()
|
||||
}
|
||||
}
|
||||
Progress::VisualRefFallback { reason } => format!("• Visual-ref fallback: {reason}"),
|
||||
Progress::UnfilledScreens { names } => {
|
||||
format!(
|
||||
"• {} screen(s) left unfilled: {}",
|
||||
names.len(),
|
||||
names.join(", ")
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Format a `SubtaskSkills` payload into the concise summary line plus
|
||||
/// indented `▸ skills:` / `▸ dropped:` detail sub-lines (spec Component 5).
|
||||
fn format_subtask_skills(
|
||||
id: &str,
|
||||
included: &[op_orchestrator::SkillBrief],
|
||||
dropped: &[(String, String)],
|
||||
budget_used: u32,
|
||||
budget_max: u32,
|
||||
) -> String {
|
||||
let mut out = format!(
|
||||
"• Subtask `{id}` · {} skills · {budget_used}/{budget_max} tok · {} dropped",
|
||||
included.len(),
|
||||
dropped.len(),
|
||||
);
|
||||
if !included.is_empty() {
|
||||
let names: Vec<String> = included
|
||||
.iter()
|
||||
.map(|s| {
|
||||
if s.truncated {
|
||||
format!("{} (truncated)", s.name)
|
||||
} else {
|
||||
s.name.clone()
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
out.push_str(&format!("\n ▸ skills: {}", names.join(", ")));
|
||||
}
|
||||
if !dropped.is_empty() {
|
||||
let drops: Vec<String> = dropped.iter().map(|(n, r)| format!("{n} ({r})")).collect();
|
||||
out.push_str(&format!("\n ▸ dropped: {}", drops.join(", ")));
|
||||
}
|
||||
out
|
||||
}
|
||||
Loading…
Reference in a new issue