#[cfg(all(feature = "tui", feature = "a2a"))]
use zeph_tui::{App, EventReader};
#[cfg(all(feature = "tui", feature = "a2a"))]
#[allow(clippy::too_many_lines)]
pub(crate) async fn run_tui_remote(
url: String,
config_path: Option<&std::path::Path>,
) -> anyhow::Result<()> {
use futures::StreamExt;
use std::time::Duration;
let config_file = crate::bootstrap::resolve_config_path(config_path);
let config = zeph_core::config::Config::load(&config_file)
.unwrap_or_else(|_| zeph_core::config::Config::default());
config.validate()?;
let auth_token = config.a2a.auth_token.clone();
let client = zeph_a2a::A2aClient::new(zeph_core::http::default_client()).with_security(
zeph_a2a::SecurityPolicy {
require_tls: config.a2a.require_tls,
ssrf_protection: config.a2a.ssrf_protection,
},
);
let remote_daemon_url = url.clone();
let (user_tx, mut user_rx) = tokio::sync::mpsc::channel::<String>(32);
let (agent_tx, agent_rx) = tokio::sync::mpsc::channel::<zeph_tui::AgentEvent>(256);
let tui_cancel = tokio_util::sync::CancellationToken::new();
let tui_supervisor = zeph_common::TaskSupervisor::new(tui_cancel.clone());
let agent_tx_pump = agent_tx.clone();
let sse_fut = async move {
while let Some(text) = user_rx.recv().await {
let message = zeph_a2a::Message::user_text(&text);
let params = zeph_a2a::SendMessageParams {
message,
configuration: None,
};
let _ = agent_tx_pump.send(zeph_tui::AgentEvent::Typing).await;
let stream_result = client
.stream_message(&url, params, auth_token.as_deref())
.await;
match stream_result {
Ok(mut stream) => {
while let Some(event) = stream.next().await {
match event {
Ok(zeph_a2a::TaskEvent::ArtifactUpdate(artifact_evt)) => {
let text: String = artifact_evt
.artifact
.parts
.iter()
.filter_map(|p| {
if let zeph_a2a::Part::Text { text, .. } = p {
Some(text.as_str())
} else {
None
}
})
.collect();
let is_final = artifact_evt.is_final;
if !text.is_empty() {
let _ =
agent_tx_pump.send(zeph_tui::AgentEvent::Chunk(text)).await;
}
if is_final {
let _ = agent_tx_pump.send(zeph_tui::AgentEvent::Flush).await;
}
}
Ok(zeph_a2a::TaskEvent::StatusUpdate(status_evt)) => {
match status_evt.status.state {
zeph_a2a::TaskState::Completed => {
let _ =
agent_tx_pump.send(zeph_tui::AgentEvent::Flush).await;
}
zeph_a2a::TaskState::Failed => {
let _ = agent_tx_pump
.send(zeph_tui::AgentEvent::FullMessage(
"Error: task failed".into(),
))
.await;
}
_ => {}
}
}
Ok(_) => {}
Err(e) => {
let _ = agent_tx_pump
.send(zeph_tui::AgentEvent::FullMessage(format!(
"Connection error: {e}"
)))
.await;
break;
}
}
}
}
Err(e) => {
let _ = agent_tx_pump
.send(zeph_tui::AgentEvent::FullMessage(format!(
"Connection error: {e}"
)))
.await;
}
}
}
};
let sse_cell = std::sync::Arc::new(parking_lot::Mutex::new(Some(sse_fut)));
tui_supervisor.spawn(zeph_common::TaskDescriptor {
name: "a2a_sse_pump",
restart: zeph_common::RestartPolicy::RunOnce,
factory: move || {
let f = sse_cell.lock().take();
async move {
if let Some(f) = f {
f.await;
}
}
},
});
let (event_tx, event_rx) = tokio::sync::mpsc::channel(256);
let reader = EventReader::new(event_tx, Duration::from_millis(100));
std::thread::spawn(move || reader.run());
let (tui_theme, tui_theme_name, tui_color_mode) = {
use zeph_tui::theme::{EffectiveColorMode, Theme, resolve_color_mode, resolve_palette};
let theme_cfg = &config.tui.theme;
let mode = resolve_color_mode(theme_cfg.color_mode);
let theme_name = theme_cfg.name.clone();
let palette_result =
tokio::task::spawn_blocking(move || resolve_palette(&theme_name)).await;
match palette_result {
Ok(Ok(p)) => (
Theme::from_palette_with_mode(&p, mode),
theme_cfg.name.clone(),
mode,
),
Ok(Err(e)) => {
tracing::warn!("TUI theme '{}' could not be loaded: {e}", theme_cfg.name);
(
Theme::default(),
"zephyr".to_owned(),
EffectiveColorMode::Truecolor,
)
}
Err(e) => {
tracing::warn!("TUI theme '{}' resolution panicked: {e}", theme_cfg.name);
(
Theme::default(),
"zephyr".to_owned(),
EffectiveColorMode::Truecolor,
)
}
}
};
let mut tui_app = App::new(user_tx, agent_rx)
.with_tool_density(config.tui.tool_density)
.with_theme(tui_theme)
.with_theme_name(tui_theme_name)
.with_effective_color_mode(tui_color_mode)
.with_remote_daemon_url(remote_daemon_url);
tui_app.set_show_source_labels(config.tui.show_source_labels);
zeph_tui::run_tui(tui_app, event_rx).await?;
tui_cancel.cancel();
tui_supervisor
.shutdown_all(std::time::Duration::from_secs(5))
.await;
Ok(())
}