mod activity;
pub mod app;
pub(crate) mod bootstrap_setup;
mod cache;
pub mod session;
pub mod settings;
mod terminal_palette;
pub mod theme;
mod ui;
pub mod wizard;
use std::io::{Write, stdout};
use std::panic;
use super::client_sessions::{self, Connection as ClientConnection, Thread as ClientThread};
use super::remote_threads::{self, RemoteThread};
use anyhow::{Context, Result};
use crossterm::cursor::{Hide, Show};
use crossterm::event::{
DisableBracketedPaste, DisableMouseCapture, EnableBracketedPaste, EnableMouseCapture, Event,
EventStream, KeyEventKind, KeyModifiers, KeyboardEnhancementFlags, MouseButton, MouseEventKind,
PopKeyboardEnhancementFlags, PushKeyboardEnhancementFlags,
};
use crossterm::execute;
use crossterm::terminal::{
EnterAlternateScreen, LeaveAlternateScreen, disable_raw_mode, enable_raw_mode,
};
use futures_util::StreamExt;
use ratatui::Terminal;
use ratatui::backend::CrosstermBackend;
use tokio::sync::mpsc;
use app::{
Agent, AgentOp, ConsoleSession, EnvNode, HeldConnect, Load, LoadSessions, ProjectNode,
SshKeyOffer, SshKeyState, WorkspaceNode,
};
pub use app::{App, Effect, LaunchRequest, Screen, Target};
use crate::client::post_graphql;
use crate::commands::code::{self, LaunchArgs, Prepared, Progress};
use crate::config::Configs;
use crate::gql::{mutations, queries};
async fn create_default_project(
client: &reqwest::Client,
backboard: &str,
workspace_id: String,
) -> Result<wizard::ProjectOption> {
use crate::gql::mutations;
let created = post_graphql::<mutations::ProjectCreate, _>(
client,
backboard.to_string(),
mutations::project_create::Variables {
name: Some("Cloud Agents".to_string()),
description: Some("Home for Railway cloud agents".to_string()),
workspace_id: Some(workspace_id),
},
)
.await?
.project_create;
let environment = created
.environments
.edges
.first()
.ok_or_else(|| anyhow::anyhow!("the new project has no environment"))?;
Ok(wizard::ProjectOption {
project_id: created.id,
project_name: created.name,
environment_id: environment.node.id.clone(),
environment_name: environment.node.name.clone(),
})
}
fn save_setup(
app: &App,
outcome: &wizard::Outcome,
) -> Result<crate::commands::cloud_agent::prefs::AgentPrefs> {
use crate::commands::cloud_agent::prefs::{AgentPrefs, DefaultProject, SkillsPrefs};
let home = dirs::home_dir().ok_or_else(|| anyhow::anyhow!("no home directory"))?;
let prefs = AgentPrefs {
version: crate::commands::cloud_agent::prefs::CURRENT_VERSION,
agent: Some(outcome.agent.clone()),
skills: SkillsPrefs {
enabled: outcome.skills,
source: outcome.skills_source.clone(),
exclude: Vec::new(),
},
mcp: AgentPrefs::load_in(&home)
.map(|p| p.mcp)
.unwrap_or_default(),
default_project: outcome.project.as_ref().map(|p| DefaultProject {
project_id: p.project_id.clone(),
project_name: p.project_name.clone(),
environment_id: p.environment_id.clone(),
environment_name: p.environment_name.clone(),
}),
theme: Some(outcome.theme.clone()),
hide_tabs: outcome.hide_tabs,
sidebar_width: app.sidebar_width,
};
prefs.save_in(&home)?;
Ok(prefs)
}
fn save_settings(
outcome: &wizard::Outcome,
) -> Result<crate::commands::cloud_agent::prefs::AgentPrefs> {
use crate::commands::cloud_agent::prefs::{AgentPrefs, DefaultProject};
let home = dirs::home_dir().ok_or_else(|| anyhow::anyhow!("no home directory"))?;
let mut prefs = AgentPrefs::load_in(&home).unwrap_or_default();
prefs.version = crate::commands::cloud_agent::prefs::CURRENT_VERSION;
prefs.agent = Some(outcome.agent.clone());
prefs.skills.enabled = outcome.skills;
prefs.skills.source = outcome.skills_source.clone();
prefs.default_project = outcome.project.as_ref().map(|p| DefaultProject {
project_id: p.project_id.clone(),
project_name: p.project_name.clone(),
environment_id: p.environment_id.clone(),
environment_name: p.environment_name.clone(),
});
prefs.theme = Some(outcome.theme.clone());
prefs.hide_tabs = outcome.hide_tabs;
prefs.save_in(&home)?;
Ok(prefs)
}
fn apply_settings(app: &mut App, outcome: &wizard::Outcome) {
match save_settings(outcome) {
Ok(_) => app.configured = true,
Err(err) => app.toast_error(format!("Couldn't save your settings: {err:#}")),
}
app.set_harness(Some(&outcome.agent));
app.set_theme(Some(&outcome.theme));
app.skills_enabled = outcome.skills;
app.hide_tabs = outcome.hide_tabs;
match &outcome.project {
Some(project) => {
app.default_project = Some(project.project_id.clone());
app.target = Some(Target {
project_id: project.project_id.clone(),
project_name: project.project_name.clone(),
environment_id: project.environment_id.clone(),
environment_name: project.environment_name.clone(),
});
}
None => app.default_project = None,
}
}
fn elide(text: &str, width: usize) -> String {
let chars: Vec<char> = text.chars().collect();
if chars.len() <= width {
return text.to_string();
}
chars[..width.saturating_sub(1)].iter().collect::<String>() + "…"
}
fn save_default_project(target: &Target) -> Result<()> {
use crate::commands::cloud_agent::prefs::{AgentPrefs, DefaultProject};
let home = dirs::home_dir().ok_or_else(|| anyhow::anyhow!("no home directory"))?;
let mut prefs = AgentPrefs::load_in(&home).unwrap_or_default();
prefs.default_project = Some(DefaultProject {
project_id: target.project_id.clone(),
project_name: target.project_name.clone(),
environment_id: target.environment_id.clone(),
environment_name: target.environment_name.clone(),
});
prefs.save_in(&home)
}
fn ssh_command_for(environment_id: &str, agent_id: &str) -> String {
use crate::commands::ssh::native;
let mut args = vec!["ssh".to_string(), "-t".to_string()];
args.extend(native::relay_port_args());
args.push(native::relay_destination(&format!(
"agent:{environment_id}:{agent_id}"
)));
args.push(code::LOGIN_SHELL_COMMAND.to_string());
crate::util::shell::shell_join(&args)
}
const SPINNER_TICK: std::time::Duration = std::time::Duration::from_millis(110);
const FLOOD_FRAME: std::time::Duration = std::time::Duration::from_millis(33);
pub enum Outcome {
NeedsCredential(LaunchRequest),
FullScreen(FullScreenRequest),
OpenShell {
agent_id: String,
agent_name: String,
},
Quit,
}
pub struct FullScreenRequest {
pub ssh_target: String,
pub identity: Option<std::path::PathBuf>,
pub relay_opts: Vec<String>,
pub session_name: String,
pub agent_name: String,
}
enum Message {
RemoteThreadReady {
connect: app::AutoConnect,
info: Box<code::ConnectInfo>,
thread: Box<RemoteThread>,
},
ClientThreadSelected {
client_id: String,
thread: ClientThread,
},
ClientReady {
pane: Box<ClientPane>,
background: bool,
},
AgentsLoaded {
path: (usize, usize, usize),
environment_id: String,
result: Result<Vec<Agent>, String>,
asked_at: std::time::Instant,
},
MyAgentsLoaded {
result: Result<Vec<(String, Agent)>, String>,
asked_at: std::time::Instant,
},
RateLimited {
retry_after_secs: Option<u64>,
},
SessionsLoaded {
path: (usize, usize, usize, usize),
agent_id: String,
result: Result<SessionInventory, String>,
},
BootstrapDefaultLoaded(String, bootstrap_setup::DefaultState),
BootstrapStep(String),
BootstrapDone(
String,
Result<crate::controllers::agent_bootstrap::Bootstrap, String>,
),
BootstrapsLoaded(
String,
Result<Vec<crate::controllers::agent_bootstrap::Bootstrap>, String>,
),
BootstrapSelected(String, Result<(), String>),
LaunchStep(String),
LaunchReady(Box<Prepared>, Box<LaunchRequest>),
LaunchFailed(String),
AgentOpDone {
agent_id: String,
environment_id: String,
op: AgentOp,
error: Option<String>,
},
ReattachReady {
agent_id: String,
agent_name: String,
session_name: String,
info: Box<code::ConnectInfo>,
},
ReattachFailed {
session_name: String,
error: String,
},
ReattachTargetGone {
agent_id: String,
agent_name: String,
session_name: String,
},
AutoReattachReady {
agent_id: String,
agent_name: String,
session_name: String,
info: Box<code::ConnectInfo>,
},
AutoConnectFailed {
session_name: String,
error: String,
},
SessionKilled {
agent_id: String,
session_name: String,
error: Option<String>,
},
ProjectCreated(Result<wizard::ProjectOption, String>),
SshKeyRegistered {
result: Result<(), String>,
then: Option<HeldConnect>,
},
ClaudeMintDone {
ok: bool,
req: Box<LaunchRequest>,
},
SessionOutput(String),
ReportsLoaded {
agent_id: String,
result: Result<Vec<activity::Report>, String>,
},
}
struct ChannelProgress(mpsc::UnboundedSender<Message>);
pub(crate) struct ClientPane {
pub agent_id: String,
pub agent_name: String,
pub environment_id: String,
pub binary: std::path::PathBuf,
pub connection: ClientConnection,
pub thread: Option<ClientThread>,
pub prompt: Option<String>,
}
impl ClientPane {
fn name(&self) -> String {
client_sessions::name(
self.connection.harness(),
&self.agent_id,
self.thread.as_ref().map(|t| t.id.as_str()),
)
}
}
fn open_client(
app: &mut App,
pane: ClientPane,
background: bool,
tx: &mpsc::UnboundedSender<Message>,
) -> Result<()> {
let name = pane.name();
if app.deleted_threads.contains(&name) {
app.connecting.remove(&name);
return Ok(());
}
if pane.thread.is_some()
&& let Some(index) = app
.sessions
.iter()
.position(|s| s.durable_name == name && !s.ended())
{
app.connecting.remove(&name);
if !background {
app.active = Some(index);
app.focus = app::ManageFocus::Session;
}
return Ok(());
}
let notify = tx.clone();
let client_id = super::opencode::generate_password();
let bridge = if let ClientConnection::Codex(connection) = &pane.connection {
let id = client_id.clone();
let updates = tx.clone();
let bridge = super::codex::bridge::Bridge::start(connection.clone(), move |thread| {
let _ = updates.send(Message::ClientThreadSelected {
client_id: id.clone(),
thread,
});
})?;
Some(bridge)
} else {
None
};
let opencode_bridge = if let ClientConnection::OpenCode(connection, beta) = &pane.connection {
let id = client_id.clone();
let updates = tx.clone();
Some(super::opencode::bridge::Bridge::start(
connection.clone(),
*beta,
move |thread| {
let _ = updates.send(Message::ClientThreadSelected {
client_id: id.clone(),
thread,
});
},
)?)
} else {
None
};
let mut session = session::Session::spawn_client(
pane.agent_id.clone(),
pane.agent_name,
&pane.binary,
&pane.connection,
bridge
.as_ref()
.map(|bridge| bridge.url.as_str())
.or_else(|| opencode_bridge.as_ref().map(|bridge| bridge.url.as_str())),
pane.thread.as_ref().map(|t| t.id.as_str()),
pane.prompt.as_deref(),
24,
80,
move || {
let _ = notify.send(Message::SessionOutput(String::new()));
},
)?;
session.ssh_target = format!("agent:{}:{}", pane.environment_id, pane.agent_id);
if pane.thread.is_none() {
session.durable_name =
client_sessions::draft_name(pane.connection.harness(), &pane.agent_id, &client_id);
}
session.client_id = Some(client_id);
session.client_thread = pane.thread;
session.client_bridge = bridge;
session.opencode_bridge = opencode_bridge;
if background {
app.attach_session_background(session, pane.agent_id.clone());
} else {
app.attach_session(session, pane.agent_id.clone());
app.expand_agent_after_load(pane.agent_id.clone());
}
Ok(())
}
fn reconnect_client(
connect: app::AutoConnect,
background: bool,
tx: &mpsc::UnboundedSender<Message>,
) {
let tx = tx.clone();
tokio::spawn(async move {
let result = async {
let (harness, _, thread_id) = client_sessions::parse_name(&connect.session_name)
.ok_or_else(|| anyhow::anyhow!("Invalid client conversation identity"))?;
let info = code::connect_info(&connect.environment_id, &connect.agent_id).await?;
let mut connection = if harness == "codex" {
ClientConnection::Codex(super::codex::reconnect(&info).await?)
} else {
ClientConnection::OpenCode(
super::opencode::reconnect(&info, harness == "opencode2").await?,
harness == "opencode2",
)
};
let binary = match &connection {
ClientConnection::Codex(c) => {
super::codex::local::ensure_client(&c.version).await?
}
ClientConnection::OpenCode(_, beta) => {
super::opencode::local::ensure_client_quiet(*beta).await?
}
};
let thread = if let Some(id) = thread_id {
Some(connection.thread(id).await?)
} else {
None
};
if let Some(thread) = &thread {
if let ClientConnection::Codex(c) = &connection {
super::codex::trust_directory(c, &thread.directory).await?;
}
connection.set_directory(&thread.directory);
}
Ok::<_, anyhow::Error>(ClientPane {
agent_id: connect.agent_id,
agent_name: connect.agent_name,
environment_id: connect.environment_id,
binary,
connection,
thread,
prompt: None,
})
}
.await;
let message = match result {
Ok(pane) => Message::ClientReady {
pane: Box::new(pane),
background,
},
Err(error) if background => Message::AutoConnectFailed {
session_name: connect.session_name,
error: format!("{error:#}"),
},
Err(error) => Message::ReattachFailed {
session_name: connect.session_name,
error: format!("{error:#}"),
},
};
let _ = tx.send(message);
});
}
fn reconnect_remote_thread(
connect: app::AutoConnect,
tx: &mpsc::UnboundedSender<Message>,
client: &reqwest::Client,
backboard: &str,
) {
let tx = tx.clone();
let client = client.clone();
let backboard = backboard.to_owned();
tokio::spawn(async move {
let result = async {
let inventory = fetch_sessions(
&client,
&backboard,
&connect.agent_id,
&connect.environment_id,
)
.await?;
let thread = inventory
.remote
.into_iter()
.find(|row| row.name(&connect.agent_id) == connect.session_name)
.ok_or_else(|| {
anyhow::anyhow!(
"Conversation is no longer available on this VM{}",
if inventory.warnings.is_empty() {
String::new()
} else {
format!(": {}", inventory.warnings.join("; "))
}
)
})?;
let info = code::connect_info(&connect.environment_id, &connect.agent_id).await?;
Ok::<_, anyhow::Error>((info, thread))
}
.await;
let message = match result {
Ok((info, thread)) => Message::RemoteThreadReady {
connect,
info: Box::new(info),
thread: Box::new(thread),
},
Err(error) => Message::ReattachFailed {
session_name: connect.session_name,
error: format!("{error:#}"),
},
};
let _ = tx.send(message);
});
}
#[derive(Default)]
struct SessionInventory {
primary_harness: Option<String>,
rows: Vec<ConsoleSession>,
remote: Vec<RemoteThread>,
warnings: Vec<String>,
failed: Vec<String>,
}
impl Progress for ChannelProgress {
fn step(&self, text: &str) {
let _ = self.0.send(Message::LaunchStep(text.to_string()));
}
fn note(&self, _text: &str) {}
fn finish(&self) {}
}
pub(crate) struct InflightLaunch {
pub(crate) req: LaunchRequest,
rx: mpsc::UnboundedReceiver<Message>,
}
impl InflightLaunch {
pub(crate) async fn settle_for_abort(mut self, limit: std::time::Duration) -> Option<String> {
let deadline = tokio::time::Instant::now() + limit;
loop {
match tokio::time::timeout_at(deadline, self.rx.recv()).await {
Ok(Some(Message::LaunchReady(prepared, _))) => {
return Some(format!(
"A launch was already in flight and created agent {} — `railway ca ssh {}` reattaches it, `railway ca delete {}` removes it.",
prepared.agent_name, prepared.agent_name, prepared.agent_name
));
}
Ok(Some(Message::ClientReady { pane, .. })) => {
return Some(format!(
"Agent {} is ready — `railway code --{} connect {}` reconnects it.",
pane.agent_name,
pane.connection.harness(),
pane.agent_name
));
}
Ok(Some(Message::LaunchFailed(_))) => return None,
Ok(Some(_)) => continue,
Ok(None) => return None,
Err(_) => {
return Some(
"A launch was still in flight — check `railway ca list` for an agent it may have created.".to_string(),
);
}
}
}
}
}
struct BootstrapProgress(mpsc::UnboundedSender<Message>);
impl Progress for BootstrapProgress {
fn finish(&self) {}
fn step(&self, text: &str) {
let _ = self.0.send(Message::BootstrapStep(text.into()));
}
fn note(&self, _text: &str) {}
}
fn load_bootstrap_default(
app: &mut App,
tx: &mpsc::UnboundedSender<Message>,
client: &reqwest::Client,
backboard: &str,
) {
let targets = [
app.target.clone(),
app.bootstrap_target(),
app.harness_pick_target.clone(),
];
for target in targets.into_iter().flatten() {
let env = target.environment_id;
if app.bootstrap_defaults.contains_key(&env) {
continue;
}
app.bootstrap_defaults
.insert(env.clone(), bootstrap_setup::DefaultState::Loading);
let (tx, client, url) = (tx.clone(), client.clone(), backboard.to_owned());
tokio::spawn(async move {
use crate::controllers::agent_bootstrap as bootstrap;
use bootstrap_setup::DefaultState;
let result = async {
let configs = Configs::new()?;
let rows = bootstrap::list(&configs, &client, &url, &env).await?;
if rows.is_empty() {
return Ok(DefaultState::Missing);
}
Ok(
match rows
.into_iter()
.find(|b| b.is_default && b.status == "READY")
{
Some(b) => DefaultState::Ready(b.name),
None => DefaultState::Available,
},
)
}
.await;
let state =
result.unwrap_or_else(|e: anyhow::Error| DefaultState::Failed(format!("{e:#}")));
let _ = tx.send(Message::BootstrapDefaultLoaded(env, state));
});
}
}
fn spawn_prepare(req: LaunchRequest, sink: mpsc::UnboundedSender<Message>) {
tokio::spawn(async move {
let req = req;
let args = launch_args_for(&req);
let progress = ChannelProgress(sink.clone());
if !args.client_on_agent
&& matches!(req.harness.as_str(), "codex" | "opencode" | "opencode2")
{
let message = match code::client::prepare_pane(args, &req.harness, &progress).await {
Ok(pane) => Message::ClientReady {
pane: Box::new(pane),
background: false,
},
Err(err) => Message::LaunchFailed(format!("{err:#}")),
};
let _ = sink.send(message);
return;
}
let message = match code::prepare(&args, &progress, code::SessionStyle::Pane).await {
Ok(prepared) => Message::LaunchReady(Box::new(prepared), Box::new(req)),
Err(err) => Message::LaunchFailed(format!("{err:#}")),
};
let _ = sink.send(message);
});
}
pub(crate) fn begin_launch_early(req: LaunchRequest) -> InflightLaunch {
let (tx, rx) = mpsc::unbounded_channel();
spawn_prepare(req.clone(), tx);
InflightLaunch { req, rx }
}
pub async fn load_tree(client: &reqwest::Client, configs: &Configs) -> Result<Vec<WorkspaceNode>> {
let workspaces = crate::workspace::workspaces_with_client(client, configs).await?;
Ok(workspaces
.into_iter()
.map(|ws| WorkspaceNode {
id: ws.id().to_string(),
name: ws.name().to_string(),
expanded: false,
projects: ws
.projects()
.into_iter()
.filter(|p| p.deleted_at().is_none())
.map(|p| ProjectNode {
id: p.id().to_string(),
name: p.name().to_string(),
expanded: false,
envs: p
.environments()
.into_iter()
.filter(|e| e.can_access)
.map(|e| EnvNode {
id: e.id,
name: e.name,
expanded: false,
agents: Load::NotLoaded,
})
.collect(),
})
.filter(|p| !p.envs.is_empty())
.collect(),
})
.filter(|ws| !ws.projects.is_empty())
.collect())
}
async fn fetch_agents(
client: &reqwest::Client,
backboard: &str,
environment_id: &str,
) -> Result<Vec<Agent>> {
let res = post_graphql::<queries::CloudAgents, _>(
client,
backboard,
queries::cloud_agents::Variables {
environment_id: environment_id.to_owned(),
mine: Some(true),
},
)
.await?;
Ok(res
.cloud_agents
.into_iter()
.map(|a| Agent {
id: a.id,
name: a.name,
status: crate::controllers::cloud_agent::Status::from(a.status).label(),
sessions: LoadSessions::NotLoaded,
expanded: false,
})
.collect())
}
async fn fetch_my_agents(
client: &reqwest::Client,
backboard: &str,
) -> Result<Vec<(String, Agent)>> {
let res = post_graphql::<queries::MyCloudAgents, _>(
client,
backboard,
queries::my_cloud_agents::Variables {},
)
.await?;
Ok(res
.my_cloud_agents
.into_iter()
.map(|a| {
(
a.environment_id,
Agent {
id: a.id,
name: a.name,
status: crate::controllers::cloud_agent::Status::from(a.status).label(),
sessions: LoadSessions::NotLoaded,
expanded: false,
},
)
})
.collect())
}
fn spawn_my_agents_fetch(
tx: &mpsc::UnboundedSender<Message>,
client: &reqwest::Client,
backboard: &str,
) {
let tx = tx.clone();
let client = client.clone();
let backboard = backboard.to_string();
let asked_at = std::time::Instant::now();
tokio::spawn(async move {
match fetch_my_agents(&client, &backboard).await {
Ok(agents) => {
let _ = tx.send(Message::MyAgentsLoaded {
result: Ok(agents),
asked_at,
});
}
Err(err) => match rate_limit_from(&err) {
Some(retry_after_secs) => {
let _ = tx.send(Message::RateLimited { retry_after_secs });
}
None => {
let _ = tx.send(Message::MyAgentsLoaded {
result: Err(err.to_string()),
asked_at,
});
}
},
}
});
}
fn start_refresh(
app: &mut App,
tx: &mpsc::UnboundedSender<Message>,
client: &reqwest::Client,
backboard: &str,
) {
if app.refreshing {
return;
}
app.refresh_started();
if app.account_query_unavailable {
let effects = app.environments_to_refresh();
app.refresh_finished();
if !effects.is_empty() {
spawn_sweep(effects, tx, client, backboard, Default::default());
}
return;
}
spawn_my_agents_fetch(tx, client, backboard);
}
async fn fetch_sessions(
client: &reqwest::Client,
backboard: &str,
cloud_agent_id: &str,
environment_id: &str,
) -> Result<SessionInventory> {
let machine = post_graphql::<queries::CloudAgent, _>(
client,
backboard,
queries::cloud_agent::Variables {
id: cloud_agent_id.into(),
environment_id: environment_id.into(),
},
)
.await?;
anyhow::ensure!(
machine
.cloud_agent
.is_some_and(
|agent| crate::controllers::cloud_agent::Status::from(agent.status).label()
== "running"
),
"Machine is not running; showing cached conversations"
);
let discovery = async {
let info = code::connect_info(environment_id, cloud_agent_id).await?;
remote_threads::discover(&info).await
};
let native = async {
let Some(connection) =
code::saved_config::client_connection(cloud_agent_id, environment_id)
else {
return (None, Ok(Vec::new()));
};
(Some(connection.harness()), connection.list().await)
};
let platform = post_graphql::<queries::CloudAgentSessionThreads, _>(
client,
backboard,
queries::cloud_agent_session_threads::Variables {
cloud_agent_id: cloud_agent_id.to_owned(),
environment_id: environment_id.to_owned(),
},
);
let (res, discovery, (native_harness, native)) = tokio::join!(platform, discovery, native);
let res = res?;
let ws_url = res
.cloud_agent
.as_ref()
.and_then(|a| a.agent_ws_url.clone());
let mut snapshots: std::collections::HashMap<String, app::ThreadSnapshot> =
std::collections::HashMap::new();
for snapshot in res.cloud_agent.map(|a| a.sessions).unwrap_or_default() {
let Some(name) = snapshot.session_name else {
continue;
};
let candidate = app::ThreadSnapshot {
harness: snapshot.harness,
session_id: snapshot.session_id,
state: snapshot.state,
prompt: snapshot.prompt,
latest_prompt: snapshot.latest_prompt,
last_reply: None,
updated_at: snapshot.updated_at,
};
match snapshots.entry(name) {
std::collections::hash_map::Entry::Occupied(mut slot) => {
if candidate.updated_at > slot.get().updated_at {
slot.insert(candidate);
}
}
std::collections::hash_map::Entry::Vacant(slot) => {
slot.insert(candidate);
}
}
}
fill_daemon_replies(
client,
backboard,
cloud_agent_id,
environment_id,
ws_url,
&mut snapshots,
)
.await;
let sessions: Vec<ConsoleSession> = res
.cloud_agent_console_sessions
.map(|conn| {
conn.edges
.into_iter()
.map(|edge| ConsoleSession {
snapshot: snapshots.remove(&edge.node.name),
name: edge.node.name,
kind: format!("{:?}", edge.node.kind),
command: Some(edge.node.command),
running: edge.node.run_state.running,
attached: edge.node.attached,
created_at: edge.node.created_at,
})
.collect()
})
.unwrap_or_default();
let discovery = discovery.unwrap_or_else(|error| remote_threads::Discovery {
warnings: vec![format!("Couldn't read VM conversation history: {error:#}")],
failed: ["claude", "grok", "codex", "opencode", "opencode2"]
.map(str::to_owned)
.to_vec(),
..Default::default()
});
let mut inventory = merge_remote_threads(cloud_agent_id, sessions, discovery);
if let Some(harness) = native_harness {
match native {
Ok(threads) => {
for thread in &threads {
let row = ConsoleSession::client_thread(cloud_agent_id, harness, Some(thread));
if let Some(previous) = inventory.rows.iter_mut().find(|r| r.name == row.name) {
*previous = row;
} else {
inventory.rows.push(row);
}
}
}
Err(error) => {
inventory.failed.push(harness.into());
inventory
.warnings
.push(format!("Couldn't read {harness} history: {error:#}"));
}
}
if inventory.primary_harness.is_none()
&& inventory
.rows
.iter()
.all(|row| row.harness_slug().is_none_or(|h| h == harness))
{
inventory.primary_harness = Some(harness.into());
}
}
inventory.rows.sort_by(|a, b| {
b.snapshot
.as_ref()
.map(|s| &s.updated_at)
.cmp(&a.snapshot.as_ref().map(|s| &s.updated_at))
.then_with(|| a.name.cmp(&b.name))
});
Ok(inventory)
}
fn merge_remote_threads(
agent_id: &str,
mut consoles: Vec<ConsoleSession>,
discovery: remote_threads::Discovery,
) -> SessionInventory {
let mut remote = discovery.threads;
for console in &consoles {
let Some(snapshot) = &console.snapshot else {
continue;
};
if snapshot.harness != "railway-agent" || snapshot.session_id.is_empty() {
continue;
}
remote.push(RemoteThread {
harness: "railway".into(),
config_dir: "/app".into(),
active: console.running,
pane_id: None,
console_name: console.running.then(|| console.name.clone()),
background_id: None,
database: None,
thread: ClientThread {
id: snapshot.session_id.clone(),
title: snapshot
.prompt
.clone()
.filter(|s| !s.trim().is_empty())
.unwrap_or_else(|| client_sessions::NEW_THREAD.into()),
directory: "/app".into(),
created_at: console.created_at,
updated_at: snapshot.updated_at.clone(),
state: snapshot.state.clone(),
},
});
}
let mut claimed = std::collections::HashSet::new();
for console in &consoles {
if !console.running {
continue;
}
if let Some(snapshot) = &console.snapshot
&& let Some(marker) = console
.command
.as_deref()
.and_then(|c| c.split_once("RAILWAY_THREAD_PANE_ID="))
.map(|(_, tail)| tail.split(';').next().unwrap_or("").trim())
&& client_sessions::validate_id(marker).is_ok()
&& !remote.iter().any(|row| {
row.active
&& (row.pane_id.as_deref() == Some(marker)
|| row.console_name.as_deref() == Some(console.name.as_str()))
})
&& let Some(row) = remote
.iter_mut()
.find(|row| row.harness == snapshot.harness && row.thread.id == snapshot.session_id)
&& row.pane_id.is_none()
{
row.pane_id = Some(marker.into());
}
let direct: Vec<_> = remote
.iter()
.enumerate()
.filter(|(_, row)| {
row.active
&& (row.console_name.as_deref() == Some(console.name.as_str())
|| row.pane_id.as_ref().is_some_and(|id| {
console.command.as_ref().is_some_and(|cmd| {
cmd.contains(&format!("RAILWAY_THREAD_PANE_ID={id};"))
})
}))
})
.map(|(i, _)| i)
.collect();
let selected = if direct.len() == 1 {
direct.first().copied()
} else if direct.is_empty() {
console.snapshot.as_ref().and_then(|snapshot| {
remote.iter().position(|row| {
row.active
&& row.harness == snapshot.harness
&& row.thread.id == snapshot.session_id
})
})
} else {
None
};
if let Some(index) = selected {
let row = &mut remote[index];
row.console_name = Some(console.name.clone());
if let Some(snapshot) = &console.snapshot
&& snapshot.session_id == row.thread.id
&& snapshot.harness == row.harness
&& (!row.active || row.harness == "grok")
{
row.thread.state = snapshot.state.clone();
}
claimed.insert(console.name.clone());
}
}
for row in &mut remote {
if row
.console_name
.as_ref()
.is_some_and(|name| !claimed.contains(name))
{
row.console_name = None;
}
}
consoles.retain(|row| !claimed.contains(&row.name) && row.is_shell());
consoles.extend(
remote
.iter()
.map(|row| ConsoleSession::client_thread(agent_id, &row.harness, Some(&row.thread))),
);
SessionInventory {
primary_harness: discovery.primary_harness,
rows: consoles,
remote,
warnings: discovery.warnings,
failed: discovery.failed,
}
}
pub async fn run(
app: &mut App,
client: reqwest::Client,
backboard: String,
pending: Option<LaunchRequest>,
) -> Result<Outcome> {
let original_hook = panic::take_hook();
panic::set_hook(Box::new(move |info| {
restore_terminal();
original_hook(info);
}));
let mut terminal = setup_terminal()?;
let _cleanup = scopeguard::guard((), |_| restore_terminal());
let mut events = EventStream::new();
let (tx, mut rx) = mpsc::unbounded_channel::<Message>();
app.bootstrap_defaults.clear();
app.refreshing = false;
app.thread_polls.clear();
app.activity = Default::default();
app.thread_cache = cache::Cache::open(&backboard);
app.restore_cached_threads();
if let Some(pane) = app.autostart_client.take() {
open_client(app, pane, false, &tx)?;
}
if let Some(req) = pending {
start_launch(app, req, &tx);
}
let stop_fetching: StopFlag = Default::default();
start_refresh(app, &tx, &client, &backboard);
let mut last_frame = std::time::Instant::now() - FLOOD_FRAME;
loop {
load_bootstrap_default(app, &tx, &client, &backboard);
if rx.is_empty() || last_frame.elapsed() >= FLOOD_FRAME {
let mut rects = app.panes;
let mut copied: Option<String> = None;
terminal.draw(|f| {
let (r, text) = ui::render_with_layout(app, f);
rects = r;
copied = text;
})?;
last_frame = std::time::Instant::now();
app.panes = rects;
if app.pending_copy.take().is_some() {
finish_copy(app, copied);
}
sync_session_size(app, &terminal);
}
if let Some(inflight) = app.autostart_inflight.take() {
let InflightLaunch { req, mut rx } = inflight;
app.screen = Screen::Manage;
if let Some(Effect::LoadAgents {
environment_id,
path,
}) = app.reveal_environment(&req.environment_id)
{
spawn_env_agents_fetch(environment_id, path, &tx, &client, &backboard);
}
app.start_loading(&req);
let tx = tx.clone();
tokio::spawn(async move {
while let Some(message) = rx.recv().await {
if tx.send(message).is_err() {
break;
}
}
});
continue;
}
if let Some(req) = app.autostart.take() {
dispatch_launch(app, req, &tx, &client, &backboard);
continue;
}
let activity_in = app.activity.next(std::time::Instant::now()).map(|delay| {
delay.max(
app.refresh_paused_until
.map(|at| at.saturating_duration_since(std::time::Instant::now()))
.unwrap_or_default(),
)
});
let toast_remaining = app.toast_remaining();
let stall_remaining = app
.stall_check_remaining()
.unwrap_or(std::time::Duration::MAX);
let effect = tokio::select! {
Some((name, error)) = app.deletion_rx.recv() => {
app.thread_deleted(&name, error);
None
}
_ = tokio::time::sleep(activity_in.unwrap_or(std::time::Duration::MAX)), if activity_in.is_some() => {
for agent_id in app.activity.take_due(std::time::Instant::now()) {
let environment = app.sessions.iter().find(|pane| pane.agent_id == agent_id && !pane.ended())
.and_then(|pane| pane.ssh_target.strip_prefix("agent:")?.split_once(':').map(|(env, _)| env.to_owned()));
if let Some(environment) = environment {
let tx = tx.clone(); let client = client.clone(); let backboard = backboard.clone();
tokio::spawn(async move {
let result = activity::fetch(&client, &backboard, &agent_id, &environment).await;
if let Err(error) = &result && let Some(retry_after_secs) = rate_limit_from(error) {
let _ = tx.send(Message::RateLimited { retry_after_secs });
}
let _ = tx.send(Message::ReportsLoaded { agent_id, result: result.map_err(|e| e.to_string()) });
});
} else { app.activity.finished(&agent_id); }
}
None
}
Some(message) = rx.recv() => handle_message(app, message, &tx, &client, &backboard, &stop_fetching),
_ = tokio::time::sleep(SPINNER_TICK), if app.loading.active
|| !app.connecting.is_empty()
|| app.bootstrap_form.as_ref().is_some_and(|f| f.running)
|| app.bootstrap_picker.as_ref().is_some_and(|p| p.loading || p.saving)
|| app.wizard.as_ref().is_some_and(|w| w.busy.is_some())
|| app.settings.as_ref().is_some_and(|s| s.busy.is_some()) => {
app.tick();
None
}
_ = tokio::time::sleep(toast_remaining), if app.toast.is_some() => {
app.expire_toast();
None
}
_ = tokio::time::sleep(app::WATCH_TICK), if app.watching_agents() => {
app.watch_tick()
}
_ = tokio::time::sleep(stall_remaining),
if app.stall_check_remaining().is_some() => None,
_ = tokio::time::sleep(std::time::Duration::from_millis(100)),
if app.awaiting_exit_status() => app.reap_ended_sessions(),
event = events.next() => match event {
Some(Ok(Event::Key(key))) if key.kind == KeyEventKind::Press => app.on_key(key),
Some(Ok(Event::Paste(text))) => app.on_paste(text),
Some(Ok(Event::Mouse(mouse))) => {
let action = match mouse.kind {
MouseEventKind::Down(MouseButton::Left) => Some(app::MouseAction::Down),
MouseEventKind::Drag(MouseButton::Left) => Some(app::MouseAction::Drag),
MouseEventKind::Up(MouseButton::Left) => Some(app::MouseAction::Up),
MouseEventKind::ScrollUp => Some(app::MouseAction::ScrollUp),
MouseEventKind::ScrollDown => Some(app::MouseAction::ScrollDown),
_ => None,
};
let shift = mouse.modifiers.contains(KeyModifiers::SHIFT);
action.and_then(|action| {
app.on_mouse_shifted(action, mouse.column, mouse.row, shift)
})
}
Some(Ok(Event::Resize(..))) => { terminal.clear()?; None }
None => Some(Effect::Quit),
_ => None,
},
};
for connect in app.take_auto_connects() {
spawn_auto_connect(connect, &tx);
}
match effect {
None => {}
Some(Effect::Quit) => {
while let Some(mut session) = app.take_session(0) {
session.detach();
}
return Ok(Outcome::Quit);
}
Some(Effect::FullScreen {
agent_id,
session_name,
agent_name,
}) => {
if client_sessions::is_client(&session_name) {
app.activate_session(&agent_id);
app.maximized = true;
continue;
}
let Some(index) = app.sessions.iter().position(|s| s.agent_id == agent_id) else {
continue;
};
let Some(session) = app.detach_session(index) else {
continue;
};
return Ok(Outcome::FullScreen(FullScreenRequest {
ssh_target: session.ssh_target.clone(),
identity: session.identity.clone(),
relay_opts: session.relay_opts.clone(),
session_name,
agent_name,
}));
}
Some(Effect::OpenShell {
agent_id,
agent_name,
}) => {
if app.hold_for_ssh_key(HeldConnect::OpenShell {
agent_id: agent_id.clone(),
agent_name: agent_name.clone(),
}) {
continue;
}
return Ok(Outcome::OpenShell {
agent_id,
agent_name,
});
}
Some(Effect::Reattach {
agent_id,
agent_name,
environment_id,
session_name,
}) => {
if app.hold_for_ssh_key(HeldConnect::Reattach {
agent_id: agent_id.clone(),
agent_name: agent_name.clone(),
environment_id: environment_id.clone(),
session_name: session_name.clone(),
}) {
continue;
}
app.connecting.insert(session_name.clone());
if client_sessions::parse_name(&session_name).is_some_and(|(h, _, _)| {
matches!(h, "claude" | "grok" | "railway")
|| code::saved_config::client_connection(&agent_id, &environment_id)
.is_none_or(|c| c.harness() != h)
}) {
reconnect_remote_thread(
app::AutoConnect {
agent_id,
agent_name,
environment_id,
session_name,
},
&tx,
&client,
&backboard,
);
continue;
}
if client_sessions::is_client(&session_name) {
reconnect_client(
app::AutoConnect {
agent_id,
agent_name,
environment_id,
session_name,
},
false,
&tx,
);
continue;
}
let tx = tx.clone();
let client = client.clone();
let backboard = backboard.clone();
tokio::spawn(async move {
let (info, listed) = tokio::join!(
code::connect_info(&environment_id, &agent_id),
post_graphql::<queries::CloudAgentSessionThreads, _>(
&client,
&backboard,
queries::cloud_agent_session_threads::Variables {
cloud_agent_id: agent_id.clone(),
environment_id: environment_id.clone(),
},
),
);
let gone = matches!(
&listed,
Ok(sessions) if !sessions.cloud_agent_console_sessions.as_ref()
.is_some_and(|sessions| sessions.edges.iter()
.any(|s| s.node.name == session_name && s.node.run_state.running))
);
let message = match info {
Ok(_) if gone => {
super::telemetry::track_session_event("reattach_target_gone", None)
.await;
Message::ReattachTargetGone {
agent_id,
agent_name,
session_name,
}
}
Ok(info) => Message::ReattachReady {
agent_id,
agent_name,
session_name,
info: Box::new(info),
},
Err(err) => {
let error = format!("{err:#}");
super::telemetry::track_session_event(
"reattach_connect_failed",
Some(error.as_str()),
)
.await;
Message::ReattachFailed {
session_name,
error,
}
}
};
let _ = tx.send(message);
});
}
Some(Effect::StepOutForMint(req)) => {
return Ok(Outcome::NeedsCredential(req));
}
Some(Effect::RegisterSshKey { offer, then }) => {
let tx = tx.clone();
let client = client.clone();
tokio::spawn(async move {
let result = register_gate_key(&client, &offer).await;
if let Err(message) = &result {
crate::commands::ssh::tel::report_failure_for(
"cloud_agent_launch",
"ssh_key_register",
message,
)
.await;
}
let _ = tx.send(Message::SshKeyRegistered { result, then });
});
}
Some(Effect::CreateDefaultProject(workspace_id)) => {
let tx = tx.clone();
let client = client.clone();
let backboard = backboard.clone();
tokio::spawn(async move {
let result = create_default_project(&client, &backboard, workspace_id)
.await
.map_err(|e| format!("{e:#}"));
let _ = tx.send(Message::ProjectCreated(result));
});
}
Some(Effect::SaveSetup(outcome)) => {
match save_setup(app, &outcome) {
Ok(prefs) => {
tokio::spawn(async move {
super::telemetry::track_setup_saved("wizard", &prefs).await;
});
app.offer_ssh_key_setup();
app.status = "Saved — Setup again to change it".into();
}
Err(err) => {
let message = format!("{err:#}");
let telemetry_message = message.clone();
tokio::spawn(async move {
super::telemetry::track_setup_failed("wizard", &telemetry_message)
.await;
});
app.toast_error(format!("Couldn't save your setup: {message}"));
}
};
app.set_harness(Some(&outcome.agent));
app.set_theme(Some(&outcome.theme));
if let Some(project) = outcome.project {
app.default_project = Some(project.project_id.clone());
app.target = Some(Target {
project_id: project.project_id,
project_name: project.project_name,
environment_id: project.environment_id,
environment_name: project.environment_name,
});
}
}
Some(Effect::SaveSidebarWidth(width)) => {
let result = dirs::home_dir()
.context("No home directory")
.and_then(|home| super::prefs::AgentPrefs::save_sidebar_width_in(&home, width));
if let Err(error) = result {
app.toast_error(format!("Couldn't save sidebar width: {error:#}"));
}
}
Some(Effect::SaveSettings(outcome)) => {
apply_settings(app, &outcome);
}
Some(Effect::ScanEverywhere) => {
stop_fetching.store(false, std::sync::atomic::Ordering::Relaxed);
let effects = app.scan_environments();
match effects.len() {
0 => {
app.status = "Refreshing…".into();
app.refresh_announce = true;
start_refresh(app, &tx, &client, &backboard);
}
n => {
app.status =
format!("Looking for agents in {n} more environment{}…", plural(n));
spawn_sweep(effects, &tx, &client, &backboard, stop_fetching.clone());
}
}
}
Some(Effect::RefreshAll) => {
app.bootstrap_defaults.clear();
start_refresh(app, &tx, &client, &backboard);
}
Some(Effect::OpenUrl(url)) => {
match ::open::that_detached(&url) {
Ok(()) => app.toast(format!("Opened {}", elide(&url, 48))),
Err(err) => app.toast_error(format!("Couldn't open it: {err}")),
}
}
Some(Effect::SaveDefaultProject(target)) => {
app.status = match save_default_project(&target) {
Ok(()) => format!("Default project is now {}", target.label()),
Err(err) => format!("Couldn't save your default project: {err:#}"),
};
}
Some(Effect::CopySsh {
agent_id,
environment_id,
}) => {
let command = ssh_command_for(&environment_id, &agent_id);
match crate::util::clipboard::copy(&command) {
Ok(()) => app.toast("Copied the SSH shell command"),
Err(err) => app.toast_error(format!("Couldn't copy: {err}")),
}
}
Some(Effect::DeleteThread {
agent_id,
environment_id,
session_name,
}) => {
let mut consoles = std::collections::HashSet::new();
let mut panes = Vec::new();
for index in (0..app.sessions.len()).rev() {
if app.sessions[index].agent_id != agent_id
|| app.sessions[index].durable_name != session_name
{
continue;
}
if let Some(mut pane) = app.take_session(index) {
pane.sync_console_name();
if pane.client_bridge.is_none()
&& pane.opencode_bridge.is_none()
&& let Some(console) = &pane.console_name
{
consoles.insert(console.clone());
}
panes.push(pane);
}
}
let tx = app.deletion_tx.clone();
tokio::spawn(async move {
let result: Result<()> =
tokio::time::timeout(std::time::Duration::from_secs(150), async {
let (harness, scoped_agent, id) =
client_sessions::parse_name(&session_name)
.context("Invalid conversation identity")?;
anyhow::ensure!(
scoped_agent == agent_id,
"Conversation belongs to another VM"
);
let id = id.context("Conversation has no native thread yet")?;
client_sessions::validate_id(id)?;
tokio::task::spawn_blocking(move || {
for mut pane in panes {
pane.detach();
}
})
.await?;
if let Some(connection) =
code::saved_config::client_connection(&agent_id, &environment_id)
.filter(|c| c.harness() == harness && consoles.is_empty())
{
connection.delete_thread(id).await
} else {
let info = code::connect_info(&environment_id, &agent_id).await?;
remote_threads::delete(
&info,
harness,
id,
&consoles.into_iter().collect::<Vec<_>>(),
)
.await
}
})
.await
.unwrap_or_else(|_| {
Err(anyhow::anyhow!("Conversation deletion timed out"))
});
let _ = tx.send((session_name, result.err().map(|e| format!("{e:#}"))));
});
}
Some(Effect::KillSession {
agent_id,
environment_id,
session_name,
}) => {
if let Some(index) = app
.sessions
.iter()
.position(|s| s.durable_name == session_name)
{
close_session(app, index, &client, &backboard).await;
}
let tx = tx.clone();
tokio::spawn(async move {
let error = code::kill_session(&environment_id, &agent_id, &session_name)
.await
.err()
.map(|e| format!("{e:#}"));
let _ = tx.send(Message::SessionKilled {
agent_id,
session_name,
error,
});
});
}
Some(Effect::CloseSession { index }) => {
close_session(app, index, &client, &backboard).await
}
Some(Effect::Agent {
op,
agent_id,
environment_id,
}) => {
if op == AgentOp::Delete
&& let Some(index) = app.sessions.iter().position(|s| s.agent_id == agent_id)
{
close_session(app, index, &client, &backboard).await;
}
let tx = tx.clone();
let client = client.clone();
let backboard = backboard.clone();
tokio::spawn(async move {
let error = run_agent_op(&client, &backboard, op, &agent_id, &environment_id)
.await
.err()
.map(|e| format!("{e:#}"));
super::telemetry::track_agent_op(op, error.as_deref()).await;
let _ = tx.send(Message::AgentOpDone {
agent_id,
environment_id,
op,
error,
});
});
}
Some(Effect::LoadBootstraps { environment_id }) => {
let (tx, client, url) = (tx.clone(), client.clone(), backboard.clone());
tokio::spawn(async move {
let result = async {
let configs = Configs::new()?;
crate::controllers::agent_bootstrap::list(
&configs,
&client,
&url,
&environment_id,
)
.await
}
.await
.map_err(|e: anyhow::Error| format!("{e:#}"));
let _ = tx.send(Message::BootstrapsLoaded(environment_id, result));
});
}
Some(Effect::SelectBootstrap { environment_id, id }) => {
let (tx, client, url) = (tx.clone(), client.clone(), backboard.clone());
tokio::spawn(async move {
let result = async {
let mut configs = Configs::new()?;
if let Some(id) = id {
let rows = crate::controllers::agent_bootstrap::list(
&configs,
&client,
&url,
&environment_id,
)
.await?;
let b = rows.iter().find(|b| b.id == id).context(
"This bootstrap is no longer available. Refresh the list.",
)?;
b.require_ready()?;
configs
.set_agent_bootstrap_default(&environment_id, &id, false)
.await?;
} else {
configs
.clear_agent_bootstrap_default(&environment_id)
.await?;
}
Ok(())
}
.await
.map_err(|e: anyhow::Error| format!("{e:#}"));
let _ = tx.send(Message::BootstrapSelected(environment_id, result));
});
}
Some(Effect::CreateBootstrap(req)) => {
if req.snapshot.is_none()
&& app.hold_for_ssh_key(HeldConnect::Bootstrap(req.clone()))
{
if app.ssh_gate.is_none()
&& let Some(form) = app.bootstrap_form.as_mut()
{
form.running = false;
form.error = Some(
"No SSH key found. Run `ssh-keygen -t ed25519`, then retry setup."
.into(),
);
}
continue;
}
let tx = tx.clone();
tokio::spawn(async move {
let env = req.target.environment_id.clone();
let result = code::bootstrap_setup::create(req, &BootstrapProgress(tx.clone()))
.await
.map_err(|e| format!("{e:#}"));
let _ = tx.send(Message::BootstrapDone(env, result));
});
}
Some(Effect::Launch(req)) => dispatch_launch(app, req, &tx, &client, &backboard),
Some(Effect::LoadSessions {
agent_id,
environment_id,
path,
}) => {
app.mark_thread_refresh(&agent_id);
spawn_session_fetch(agent_id, environment_id, path, &tx, &client, &backboard);
}
Some(Effect::LoadAgents {
environment_id,
path,
}) => {
spawn_env_agents_fetch(environment_id, path, &tx, &client, &backboard);
}
}
}
}
async fn fill_daemon_replies(
client: &reqwest::Client,
backboard: &str,
cloud_agent_id: &str,
environment_id: &str,
ws_url: Option<String>,
snapshots: &mut std::collections::HashMap<String, app::ThreadSnapshot>,
) {
const MAX_DIALS: usize = 4;
let Some(ws_url) = ws_url else { return };
let mut targets: Vec<(String, String)> = snapshots
.values()
.filter(|snapshot| snapshot.harness == "railway-agent")
.map(|snapshot| (snapshot.updated_at.clone(), snapshot.session_id.clone()))
.collect();
targets.sort_by(|a, b| b.0.cmp(&a.0));
targets.dedup_by(|a, b| a.1 == b.1);
targets.truncate(MAX_DIALS);
if targets.is_empty() {
return;
}
let (cached, to_dial): (Vec<_>, Vec<_>) = {
let cache = reply_cache().lock().unwrap_or_else(|e| e.into_inner());
targets.into_iter().partition(|(updated_at, session_id)| {
cache
.get(session_id)
.is_some_and(|(at, _)| at == updated_at)
})
};
for (_, session_id) in &cached {
let reply = {
let cache = reply_cache().lock().unwrap_or_else(|e| e.into_inner());
cache.get(session_id).and_then(|(_, reply)| reply.clone())
};
if let Some(reply) = reply {
apply_reply(snapshots, session_id, reply);
}
}
if to_dial.is_empty() {
return;
}
let Ok(res) = post_graphql::<mutations::CloudAgentHarnessToken, _>(
client,
backboard,
mutations::cloud_agent_harness_token::Variables {
id: cloud_agent_id.to_owned(),
environment_id: environment_id.to_owned(),
},
)
.await
else {
return;
};
let token = res.cloud_agent_harness_token;
for (updated_at, session_id) in to_dial {
let reply = gate_last_reply(&ws_url, &token, &session_id).await;
{
let mut cache = reply_cache().lock().unwrap_or_else(|e| e.into_inner());
if cache.len() > 256 {
cache.clear();
}
cache.insert(session_id.clone(), (updated_at, reply.clone()));
}
if let Some(reply) = reply {
apply_reply(snapshots, &session_id, reply);
}
}
}
fn reply_cache()
-> &'static std::sync::Mutex<std::collections::HashMap<String, (String, Option<String>)>> {
static CACHE: std::sync::OnceLock<
std::sync::Mutex<std::collections::HashMap<String, (String, Option<String>)>>,
> = std::sync::OnceLock::new();
CACHE.get_or_init(Default::default)
}
fn apply_reply(
snapshots: &mut std::collections::HashMap<String, app::ThreadSnapshot>,
session_id: &str,
reply: String,
) {
for snapshot in snapshots.values_mut() {
if snapshot.session_id == session_id {
snapshot.last_reply = Some(reply.clone());
}
}
}
async fn gate_last_reply(ws_url: &str, token: &str, session_id: &str) -> Option<String> {
use futures_util::StreamExt;
use reqwest_websocket::{Message as WsMessage, RequestBuilderExt};
let url = format!("{ws_url}?session_id={session_id}&token={token}");
let response = reqwest::Client::default()
.get(&url)
.timeout(std::time::Duration::from_secs(8))
.upgrade()
.send()
.await
.ok()?;
let mut ws = response.into_websocket().await.ok()?;
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(6);
loop {
let frame = tokio::time::timeout_at(deadline, ws.next()).await.ok()??;
let Ok(WsMessage::Text(text)) = frame else {
continue;
};
let Ok(value) = serde_json::from_str::<serde_json::Value>(&text) else {
continue;
};
match value.get("type").and_then(|t| t.as_str()) {
Some("catchup") => return last_assistant_text(&value),
Some("ready" | "session_start") => return None,
_ => {}
}
}
}
fn last_assistant_text(catchup: &serde_json::Value) -> Option<String> {
let messages = catchup.get("data")?.get("messages")?.as_array()?;
for message in messages.iter().rev() {
if message.get("role").and_then(|r| r.as_str()) != Some("assistant") {
continue;
}
let Some(blocks) = message.get("content").and_then(|c| c.as_array()) else {
continue;
};
let text = blocks
.iter()
.filter(|block| block.get("type").and_then(|t| t.as_str()) == Some("text"))
.filter_map(|block| block.get("text").and_then(|t| t.as_str()))
.collect::<Vec<_>>()
.join(" ");
let text = text.split_whitespace().collect::<Vec<_>>().join(" ");
if !text.is_empty() {
return Some(text.chars().take(400).collect());
}
}
None
}
fn spawn_auto_connect(connect: app::AutoConnect, tx: &mpsc::UnboundedSender<Message>) {
if client_sessions::is_client(&connect.session_name) {
reconnect_client(connect, true, tx);
return;
}
let app::AutoConnect {
agent_id,
agent_name,
environment_id,
session_name,
} = connect;
let tx = tx.clone();
tokio::spawn(async move {
let message = match code::connect_info(&environment_id, &agent_id).await {
Ok(info) => Message::AutoReattachReady {
agent_id,
agent_name,
session_name,
info: Box::new(info),
},
Err(err) => {
let error = format!("{err:#}");
super::telemetry::track_session_event(
"auto_reattach_connect_failed",
Some(error.as_str()),
)
.await;
Message::AutoConnectFailed {
session_name,
error,
}
}
};
let _ = tx.send(message);
});
}
fn spawn_env_agents_fetch(
environment_id: String,
path: (usize, usize, usize),
tx: &mpsc::UnboundedSender<Message>,
client: &reqwest::Client,
backboard: &str,
) {
let tx = tx.clone();
let client = client.clone();
let backboard = backboard.to_string();
tokio::spawn(async move {
let asked_at = std::time::Instant::now();
match fetch_agents(&client, &backboard, &environment_id).await {
Ok(agents) => {
let _ = tx.send(Message::AgentsLoaded {
path,
environment_id,
result: Ok(agents),
asked_at,
});
}
Err(err) => match rate_limit_from(&err) {
Some(retry_after_secs) => {
let _ = tx.send(Message::RateLimited { retry_after_secs });
}
None => {
let _ = tx.send(Message::AgentsLoaded {
path,
environment_id,
result: Err(err.to_string()),
asked_at,
});
}
},
}
});
}
fn handle_message(
app: &mut App,
message: Message,
tx: &mpsc::UnboundedSender<Message>,
client: &reqwest::Client,
backboard: &str,
stop_fetching: &StopFlag,
) -> Option<Effect> {
match message {
Message::RemoteThreadReady {
connect,
info,
thread,
} => {
app.connecting.remove(&connect.session_name);
if app.deleted_threads.contains(&connect.session_name) {
return None;
}
let pane_id = super::opencode::generate_password();
let result = thread.resume_command(&pane_id).and_then(|command| {
let console_name = thread
.console_name
.clone()
.unwrap_or_else(|| session::durable_name(&thread.harness));
let notify = tx.clone();
let output_agent = connect.agent_id.clone();
let mut pane = session::Session::spawn(
connect.agent_id.clone(),
connect.agent_name,
thread.harness.clone(),
&info.ssh_target,
info.identity.as_deref(),
&info.relay_opts,
&command,
thread.console_name.is_some(),
&console_name,
24,
80,
move || {
let _ = notify.send(Message::SessionOutput(output_agent.clone()));
},
)?;
pane.console_name = Some(console_name);
pane.client_id = Some(if thread.console_name.is_some() {
thread.pane_id.clone().unwrap_or(pane_id)
} else {
pane_id
});
pane.durable_name = connect.session_name.clone();
pane.client_thread = Some(thread.thread.clone());
Ok::<_, anyhow::Error>(pane)
});
match result {
Ok(pane) => {
app.attach_session(pane, connect.agent_id.clone());
}
Err(error) => app.toast_error(format!("Couldn't open conversation: {error:#}")),
}
None
}
Message::ClientThreadSelected { client_id, thread } => {
if let Some(agent_id) = app.client_thread_selected(&client_id, thread) {
app.persist_threads(&agent_id);
}
None
}
Message::ClientReady { pane, background } => {
let name = pane.name();
let environment_id = pane.environment_id.clone();
if let Err(error) = open_client(app, *pane, background, tx) {
app.connecting.remove(&name);
if background {
app.auto_connect_failed(&name, &format!("{error:#}"));
} else {
app.launch_failed(format!("Could not open local client: {error:#}"));
}
}
app.reveal_environment(&environment_id)
}
Message::RateLimited { retry_after_secs } => {
app.rate_limited(retry_after_secs);
None
}
Message::AgentsLoaded {
path,
environment_id,
result,
asked_at,
} => {
app.agents_loaded_at(path, &environment_id, result, asked_at);
app.restore_cached_threads();
let prefetch = app.sessions_to_prefetch();
if !prefetch.is_empty() {
spawn_session_prefetch(prefetch, tx, client, backboard, stop_fetching.clone());
}
app.expand_pending()
}
Message::MyAgentsLoaded { result, asked_at } => match result {
Ok(agents) => {
let count = agents.len();
app.refresh_finished();
app.my_agents_loaded(agents, asked_at);
app.restore_cached_threads();
let prefetch = app.sessions_to_prefetch();
if !prefetch.is_empty() {
spawn_session_prefetch(prefetch, tx, client, backboard, stop_fetching.clone());
}
let watched = app.threads_to_refresh();
if !watched.is_empty() {
spawn_session_prefetch(watched, tx, client, backboard, stop_fetching.clone());
}
app.refreshed(count);
app.expand_pending()
}
Err(err) => {
app.refresh_finished();
app.refresh_announce = false;
let first_failure = !app.account_query_unavailable;
app.account_query_unavailable = true;
stop_fetching.store(false, std::sync::atomic::Ordering::Relaxed);
let mut sweep = app.initial_environments();
if sweep.is_empty() {
sweep = app.environments_to_refresh();
}
if !sweep.is_empty() {
spawn_sweep(sweep, tx, client, backboard, stop_fetching.clone());
} else if first_failure {
app.toast_error(format!("Couldn't load agents: {err}"));
}
None
}
},
Message::SessionsLoaded {
path,
agent_id,
result,
} => {
match result {
Ok(mut inventory) => {
if let Some(harness) = inventory.primary_harness.take() {
app.primary_harnesses.insert(agent_id.clone(), harness);
} else {
app.primary_harnesses.remove(&agent_id);
}
app.remote_threads_loaded(&agent_id, &inventory.remote);
app.preserve_failed_threads(&agent_id, &inventory.failed, &mut inventory.rows);
app.sessions_loaded(path, &agent_id, Ok(inventory.rows));
if !inventory.warnings.is_empty() {
app.status = inventory.warnings.join("; ");
}
}
Err(error) => app.sessions_loaded(path, &agent_id, Err(error)),
}
app.finish_agent_connect(&agent_id)
}
Message::BootstrapDefaultLoaded(env, state) => {
app.bootstrap_defaults.insert(env, state);
None
}
Message::BootstrapsLoaded(env, result) => {
if let Some(picker) = app
.bootstrap_picker
.as_mut()
.filter(|p| p.target.environment_id == env && p.loading)
{
picker.loaded(result.clone());
if picker.for_launch {
use bootstrap_setup::LaunchChoice;
picker.cursor = match &app.harness_bootstrap {
LaunchChoice::Named(name) => picker
.entries
.iter()
.position(|b| b.name == *name)
.map_or(0, |i| i + 1),
LaunchChoice::None => picker.entries.len() + 1,
LaunchChoice::Default => 0,
};
}
}
if let Some(form) = app
.bootstrap_form
.as_mut()
.filter(|f| f.target.environment_id == env && f.defaults_loading)
{
match result {
Ok(entries) => {
form.make_default = !entries.iter().any(|b| b.is_default);
form.defaults_loading = false;
}
Err(error) => {
form.error = Some(format!(
"Could not load the current default: {error}. Press Esc and try again."
));
}
}
}
None
}
Message::BootstrapSelected(env, result) => {
app.bootstrap_defaults.remove(&env);
if let Some(picker) = app
.bootstrap_picker
.as_mut()
.filter(|p| p.target.environment_id == env)
{
picker.saving = false;
match result {
Ok(()) => {
app.bootstrap_picker = None;
app.screen = Screen::Manage;
app.toast("Bootstrap default updated");
}
Err(error) => picker.error = Some(error),
}
}
None
}
Message::BootstrapStep(text) => {
if let Some(form) = app.bootstrap_form.as_mut() {
if form.steps.last() != Some(&text) {
form.steps.push(text);
}
}
None
}
Message::BootstrapDone(env, result) => {
app.bootstrap_defaults.remove(&env);
if let Some(form) = app
.bootstrap_form
.as_mut()
.filter(|f| f.target.environment_id == env)
{
form.running = false;
match result {
Ok(b) => {
form.finished = true;
form.steps = vec![if b.is_default {
format!("'{}' is your default for new Cloud Agents.", b.name)
} else {
format!("'{}' is ready to use. Your default is unchanged.", b.name)
}];
}
Err(error) => form.error = Some(error),
}
}
start_refresh(app, tx, client, backboard);
None
}
Message::LaunchStep(text) => {
app.loading_step(text);
None
}
Message::LaunchFailed(err) => {
app.launch_failed(err);
None
}
Message::AgentOpDone {
agent_id,
environment_id,
op,
error,
} => {
app.agent_op_finished(&agent_id, &environment_id, op, error);
app.reveal_environment(&environment_id)
}
Message::LaunchReady(prepared, req) => open_session(app, *prepared, *req, tx),
Message::ReattachReady {
agent_id,
agent_name,
session_name,
info,
} => {
let notify_tx = tx.clone();
let output_agent = agent_id.clone();
match session::Session::spawn(
agent_id.clone(),
agent_name,
"session".to_string(),
&info.ssh_target,
info.identity.as_deref(),
&info.relay_opts,
"",
true,
&session_name,
24,
80,
move || {
let _ = notify_tx.send(Message::SessionOutput(output_agent.clone()));
},
) {
Ok(session) => {
app.attach_session(session, agent_id);
tokio::spawn(super::telemetry::track_session_event("reattach", None));
None
}
Err(err) => {
app.connecting.remove(&session_name);
let message = format!("couldn't reattach: {err}");
let telemetry_message = message.clone();
tokio::spawn(async move {
super::telemetry::track_session_event(
"reattach_open_failed",
Some(telemetry_message.as_str()),
)
.await;
});
app.launch_failed(message);
None
}
}
}
Message::ReattachFailed {
session_name,
error,
} => {
app.connecting.remove(&session_name);
app.launch_failed(error);
None
}
Message::ReattachTargetGone {
agent_id,
agent_name,
session_name,
} => app.reattach_target_gone(&agent_id, &agent_name, &session_name),
Message::AutoReattachReady {
agent_id,
agent_name,
session_name,
info,
} => {
let notify_tx = tx.clone();
let output_agent = agent_id.clone();
match session::Session::spawn(
agent_id.clone(),
agent_name,
"session".to_string(),
&info.ssh_target,
info.identity.as_deref(),
&info.relay_opts,
"",
true,
&session_name,
24,
80,
move || {
let _ = notify_tx.send(Message::SessionOutput(output_agent.clone()));
},
) {
Ok(session) => {
app.attach_session_background(session, agent_id);
tokio::spawn(super::telemetry::track_session_event("auto_reattach", None));
}
Err(err) => {
let error = format!("{err:#}");
let telemetry_error = error.clone();
tokio::spawn(async move {
super::telemetry::track_session_event(
"auto_reattach_open_failed",
Some(telemetry_error.as_str()),
)
.await;
});
app.auto_connect_failed(&session_name, &error);
}
}
None
}
Message::AutoConnectFailed {
session_name,
error,
} => {
app.auto_connect_failed(&session_name, &error);
None
}
Message::ProjectCreated(result) => {
if let Some(w) = app.wizard.as_mut() {
w.project_created(result);
} else if let Some(outcome) = app
.settings
.as_mut()
.and_then(|s| s.project_created(result))
{
apply_settings(app, &outcome);
}
None
}
Message::ClaudeMintDone { ok, req } => {
if ok {
start_launch(app, *req, tx);
} else {
app.loading.active = false;
return Some(Effect::StepOutForMint(*req));
}
None
}
Message::SshKeyRegistered { result, then } => match result {
Ok(()) => {
app.ssh_key = SshKeyState::Ready;
app.toast("SSH key registered");
then.map(HeldConnect::into_effect)
}
Err(message) => {
if matches!(then, Some(HeldConnect::Bootstrap(_)))
&& let Some(form) = app.bootstrap_form.as_mut()
{
form.running = false;
form.error = Some(format!("SSH key registration failed: {message}"));
}
let first = message.lines().next().unwrap_or("registration failed");
app.toast_error(format!("Couldn't register the key: {first}"));
None
}
},
Message::SessionKilled {
agent_id,
session_name,
error,
} => {
let telemetry_error = error.clone();
tokio::spawn(async move {
super::telemetry::track_session_event("kill_session", telemetry_error.as_deref())
.await;
});
app.session_killed(&session_name, error);
app.refresh_agent_sessions(&agent_id)
}
Message::ReportsLoaded { agent_id, result } => {
app.activity.finished(&agent_id);
if let Ok(reports) = result {
activity::apply(app, &agent_id, &reports);
}
None
}
Message::SessionOutput(agent_id) => {
for pane in &mut app.sessions {
pane.sync_console_name();
}
if app.sessions.iter().any(|pane| {
pane.agent_id == agent_id
&& !pane.ended()
&& pane.client_id.is_some()
&& pane.client_bridge.is_none()
&& pane.opencode_bridge.is_none()
}) {
app.activity.changed(&agent_id, std::time::Instant::now());
}
app.reap_ended_sessions()
}
}
}
fn open_session(
app: &mut App,
prepared: Prepared,
req: LaunchRequest,
tx: &mpsc::UnboundedSender<Message>,
) -> Option<Effect> {
if !req.wants_new_session()
&& req.session_name.is_none()
&& app.activate_session(&prepared.agent_id)
{
app.screen = app::Screen::Manage;
app.status = "Switched to the open session".into();
return app.reveal_environment(&prepared.environment_id);
}
let durable_session = req
.session_name
.clone()
.unwrap_or_else(|| session::durable_name(prepared.harness));
let notify_tx = tx.clone();
let output_agent = prepared.agent_id.clone();
let (rows, cols) = (24u16, 80u16);
let pane_id = (prepared.harness != "shell").then(super::opencode::generate_password);
let remote_cmd = match &pane_id {
Some(id) => format!(
"export RAILWAY_THREAD_PANE_ID={id}; {}",
prepared.remote_cmd
),
None => prepared.remote_cmd.clone(),
};
match session::Session::spawn(
prepared.agent_id.clone(),
prepared.agent_name.clone(),
prepared.harness.to_string(),
&prepared.ssh_target,
prepared.identity.as_deref(),
&prepared.relay_opts,
&remote_cmd,
req.session_name.is_some(),
&durable_session,
rows,
cols,
move || {
let _ = notify_tx.send(Message::SessionOutput(output_agent.clone()));
},
) {
Ok(mut session) => {
if let Some(id) = pane_id {
session.durable_name =
client_sessions::draft_name(prepared.harness, &prepared.agent_id, &id);
session.client_id = Some(id);
session.console_name = Some(durable_session);
}
app.attach_session(session, prepared.agent_id.clone());
app.expand_agent_after_load(prepared.agent_id.clone());
app.reveal_environment(&prepared.environment_id)
}
Err(err) => {
let message = format!("couldn't open the session: {err}");
let telemetry_message = message.clone();
tokio::spawn(async move {
super::telemetry::track_session_event(
"launch_open_failed",
Some(telemetry_message.as_str()),
)
.await;
});
app.launch_failed(message);
None
}
}
}
const SWEEP_CONCURRENCY: usize = 5;
fn plural(n: usize) -> &'static str {
if n == 1 { "" } else { "s" }
}
type StopFlag = std::sync::Arc<std::sync::atomic::AtomicBool>;
fn rate_limit_from(err: &anyhow::Error) -> Option<Option<u64>> {
match err.downcast_ref::<crate::errors::RailwayError>() {
Some(crate::errors::RailwayError::Ratelimited { retry_after_secs }) => {
Some(*retry_after_secs)
}
_ => None,
}
}
fn spawn_sweep(
effects: Vec<Effect>,
tx: &mpsc::UnboundedSender<Message>,
client: &reqwest::Client,
backboard: &str,
stop: StopFlag,
) {
use std::sync::atomic::Ordering;
let tx = tx.clone();
let client = client.clone();
let backboard = backboard.to_string();
tokio::spawn(async move {
let permits = std::sync::Arc::new(tokio::sync::Semaphore::new(SWEEP_CONCURRENCY));
for effect in effects {
let Effect::LoadAgents {
environment_id,
path,
} = effect
else {
continue;
};
if stop.load(Ordering::Relaxed) {
return;
}
let Ok(permit) = permits.clone().acquire_owned().await else {
return;
};
let tx = tx.clone();
let client = client.clone();
let backboard = backboard.clone();
let stop = stop.clone();
tokio::spawn(async move {
let asked_at = std::time::Instant::now();
match fetch_agents(&client, &backboard, &environment_id).await {
Ok(agents) => {
let _ = tx.send(Message::AgentsLoaded {
path,
environment_id,
result: Ok(agents),
asked_at,
});
}
Err(err) => match rate_limit_from(&err) {
Some(retry_after_secs) => {
stop.store(true, Ordering::Relaxed);
let _ = tx.send(Message::RateLimited { retry_after_secs });
}
None => {
let _ = tx.send(Message::AgentsLoaded {
path,
environment_id,
result: Err(err.to_string()),
asked_at,
});
}
},
}
drop(permit);
});
}
});
}
fn spawn_session_fetch(
agent_id: String,
environment_id: String,
path: (usize, usize, usize, usize),
tx: &mpsc::UnboundedSender<Message>,
client: &reqwest::Client,
backboard: &str,
) {
let tx = tx.clone();
let client = client.clone();
let backboard = backboard.to_string();
tokio::spawn(async move {
let result = match fetch_sessions(&client, &backboard, &agent_id, &environment_id).await {
Ok(sessions) => Ok(sessions),
Err(err) => {
if let Some(retry_after_secs) = rate_limit_from(&err) {
let _ = tx.send(Message::RateLimited { retry_after_secs });
return;
}
Err(err.to_string())
}
};
let _ = tx.send(Message::SessionsLoaded {
path,
agent_id,
result,
});
});
}
fn spawn_session_prefetch(
effects: Vec<Effect>,
tx: &mpsc::UnboundedSender<Message>,
client: &reqwest::Client,
backboard: &str,
stop: StopFlag,
) {
use std::sync::atomic::Ordering;
let tx = tx.clone();
let client = client.clone();
let backboard = backboard.to_string();
tokio::spawn(async move {
let permits = std::sync::Arc::new(tokio::sync::Semaphore::new(SWEEP_CONCURRENCY));
for effect in effects {
let Effect::LoadSessions {
agent_id,
environment_id,
path,
} = effect
else {
continue;
};
if stop.load(Ordering::Relaxed) {
return;
}
let Ok(permit) = permits.clone().acquire_owned().await else {
return;
};
let tx = tx.clone();
let client = client.clone();
let backboard = backboard.clone();
let stop = stop.clone();
tokio::spawn(async move {
match fetch_sessions(&client, &backboard, &agent_id, &environment_id).await {
Ok(sessions) => {
let _ = tx.send(Message::SessionsLoaded {
path,
agent_id,
result: Ok(sessions),
});
}
Err(err) => match rate_limit_from(&err) {
Some(retry_after_secs) => {
stop.store(true, Ordering::Relaxed);
let _ = tx.send(Message::RateLimited { retry_after_secs });
}
None => {
let _ = tx.send(Message::SessionsLoaded {
path,
agent_id,
result: Err(err.to_string()),
});
}
},
}
drop(permit);
});
}
});
}
async fn register_gate_key(client: &reqwest::Client, offer: &SshKeyOffer) -> Result<(), String> {
let configs = Configs::new().map_err(|e| format!("{e:#}"))?;
crate::controllers::ssh::keys::register_ssh_key(
client,
&configs,
&offer.name,
&offer.public_key,
None,
)
.await
.map(|_| ())
.map_err(|e| format!("{e:#}"))
}
fn launch_args_for(req: &LaunchRequest) -> LaunchArgs {
(*req.base).clone().retargeted(
req.project_id.clone(),
req.environment_id.clone(),
&req.harness,
req.force_new,
req.prompt.clone(),
req.agent_id.clone(),
)
}
fn dispatch_launch(
app: &mut App,
req: LaunchRequest,
tx: &mpsc::UnboundedSender<Message>,
client: &reqwest::Client,
backboard: &str,
) {
app.screen = Screen::Manage;
if let Some(Effect::LoadAgents {
environment_id,
path,
}) = app.reveal_environment(&req.environment_id)
{
spawn_env_agents_fetch(environment_id, path, tx, client, backboard);
}
if app.hold_for_ssh_key(HeldConnect::Launch(req.clone())) {
return;
}
if req.harness == "claude" && code::claude_needs_local_mint() {
app.start_loading(&req);
let _ = tx.send(Message::LaunchStep(
"Minting a Claude token — approve the browser prompt if one appears".to_string(),
));
let tx = tx.clone();
tokio::task::spawn_blocking(move || {
let ok = code::mint_claude_credential_headless().is_ok();
let _ = tx.send(Message::ClaudeMintDone {
ok,
req: Box::new(req),
});
});
return;
}
start_launch(app, req, tx);
}
fn start_launch(app: &mut App, req: LaunchRequest, tx: &mpsc::UnboundedSender<Message>) {
app.start_loading(&req);
spawn_prepare(req, tx.clone());
}
async fn run_agent_op(
client: &reqwest::Client,
backboard: &str,
op: AgentOp,
agent_id: &str,
environment_id: &str,
) -> Result<()> {
let backboard = backboard.to_string();
match op {
AgentOp::Sleep => {
crate::controllers::cloud_agent::sleep(client, &backboard, environment_id, agent_id)
.await?;
}
AgentOp::Wake => {
post_graphql::<mutations::CloudAgentWake, _>(
client,
backboard,
mutations::cloud_agent_wake::Variables {
id: agent_id.to_string(),
},
)
.await?;
}
AgentOp::Delete => {
post_graphql::<mutations::CloudAgentDelete, _>(
client,
backboard,
mutations::cloud_agent_delete::Variables {
id: agent_id.to_string(),
},
)
.await?;
}
}
Ok(())
}
async fn close_session(app: &mut App, index: usize, _client: &reqwest::Client, _backboard: &str) {
if let Some(mut session) = app.take_session(index) {
session.detach();
}
}
fn sync_session_size(app: &mut App, terminal: &Terminal<CrosstermBackend<std::io::Stdout>>) {
let Some((rows, cols)) =
ui::session_pane_size(terminal.size().ok(), app.pane_is_full(), app.sidebar_width)
else {
return;
};
for session in app.sessions.iter_mut() {
session.resize(rows, cols);
}
}
fn finish_copy(app: &mut App, text: Option<String>) {
app.selection = None;
let Some(text) = text else {
return;
};
let lines = text.lines().count();
match crate::util::clipboard::copy(&text) {
Ok(()) => app.toast(format!(
"Copied {lines} line{}",
if lines == 1 { "" } else { "s" }
)),
Err(err) => app.toast_error(format!("Couldn't copy: {err}")),
}
}
fn setup_terminal() -> Result<Terminal<CrosstermBackend<std::io::Stdout>>> {
terminal_palette::capture();
enable_raw_mode()?;
crate::util::prompt::set_terminal_owned(true);
execute!(
stdout(),
EnterAlternateScreen,
EnableMouseCapture,
EnableBracketedPaste,
Hide
)?;
if matches!(
crossterm::terminal::supports_keyboard_enhancement(),
Ok(true)
) {
let _ = execute!(
stdout(),
PushKeyboardEnhancementFlags(KeyboardEnhancementFlags::DISAMBIGUATE_ESCAPE_CODES)
);
}
let backend = CrosstermBackend::new(stdout());
let mut terminal = Terminal::new(backend)?;
terminal.clear()?;
Ok(terminal)
}
fn restore_terminal() {
let _ = execute!(stdout(), PopKeyboardEnhancementFlags);
let _ = execute!(
stdout(),
DisableBracketedPaste,
DisableMouseCapture,
LeaveAlternateScreen,
Show
);
let _ = execute!(stdout(), DisableBracketedPaste, DisableMouseCapture);
let _ = disable_raw_mode();
crate::util::prompt::set_terminal_owned(false);
let _ = stdout().flush();
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn bootstrap_default_loading_and_selection_stay_in_the_tui() {
let target = Target {
project_id: "project".into(),
project_name: "Demo".into(),
environment_id: "env".into(),
environment_name: "production".into(),
};
let mut app = App::new(
vec![],
Some(target.clone()),
Some("claude"),
None,
None,
true,
);
let (tx, _rx) = mpsc::unbounded_channel();
let client = reqwest::Client::new();
let stop = StopFlag::default();
for has_default in [false, true] {
let mut form = bootstrap_setup::Form::new(target.clone(), 0);
form.snapshot = Some(bootstrap_setup::Snapshot {
agent_id: "vm".into(),
agent_name: "selected".into(),
});
form.defaults_loading = true;
app.bootstrap_form = Some(form);
app.screen = Screen::BootstrapSetup;
let entries = if has_default {
vec![crate::controllers::agent_bootstrap::Bootstrap {
id: "saved".into(),
name: "dev".into(),
environment_id: "env".into(),
status: "READY".into(),
failure_reason: None,
updated_at: chrono::Utc::now(),
is_default: true,
source_agent_id: None,
checkpoint_id: None,
}]
} else {
vec![]
};
handle_message(
&mut app,
Message::BootstrapsLoaded("env".into(), Ok(entries)),
&tx,
&client,
"unused",
&stop,
);
let form = app.bootstrap_form.as_ref().unwrap();
assert_eq!(form.make_default, !has_default);
assert!(!form.defaults_loading);
assert_eq!(app.screen, Screen::BootstrapSetup);
}
app.bootstrap_form = None;
app.bootstrap_picker = Some(bootstrap_setup::Picker::new(target, false));
app.screen = Screen::BootstrapPick;
handle_message(
&mut app,
Message::BootstrapSelected("env".into(), Err("cannot save".into())),
&tx,
&client,
"unused",
&stop,
);
assert_eq!(app.screen, Screen::BootstrapPick);
assert_eq!(
app.bootstrap_picker.as_ref().unwrap().error.as_deref(),
Some("cannot save")
);
handle_message(
&mut app,
Message::BootstrapSelected("env".into(), Ok(())),
&tx,
&client,
"unused",
&stop,
);
assert_eq!(app.screen, Screen::Manage);
assert!(app.bootstrap_picker.is_none());
}
#[test]
fn bootstrap_and_launch_progress_send_stages_without_detail_notes() {
let (tx, mut rx) = mpsc::unbounded_channel();
let progress = ChannelProgress(tx.clone());
progress.note("Using bootstrap 'dev'");
progress.step("Creating a cloud agent");
assert!(
matches!(rx.try_recv(), Ok(Message::LaunchStep(s)) if s == "Creating a cloud agent")
);
assert!(rx.try_recv().is_err());
let progress = BootstrapProgress(tx);
progress.note("Copied the selected coding agent settings");
progress.step("Saving checkpoint");
assert!(matches!(rx.try_recv(), Ok(Message::BootstrapStep(s)) if s == "Saving checkpoint"));
assert!(rx.try_recv().is_err());
}
#[test]
fn history_survives_console_exit_and_native_switches_replace_the_console_row() {
let old = remote_threads::tests::thread("claude", "old");
let mut current = remote_threads::tests::thread("claude", "current");
current.active = true;
current.pane_id = Some("pane-id".into());
let console = ConsoleSession {
name: "relay-name".into(),
kind: "SHELL".into(),
running: true,
attached: true,
created_at: None,
command: Some("export RAILWAY_THREAD_PANE_ID=pane-id; claude".into()),
snapshot: Some(app::ThreadSnapshot {
harness: "claude".into(),
session_id: "old".into(),
state: "working".into(),
prompt: Some("stale prompt".into()),
latest_prompt: None,
last_reply: None,
updated_at: "2026-09-10T09:00:00Z".into(),
}),
};
let stale = merge_remote_threads(
"vm",
vec![console.clone()],
remote_threads::Discovery {
threads: vec![old.clone()],
..Default::default()
},
);
assert!(
stale.remote[0].console_name.is_none(),
"a stale hook on a running shell must resume by ID"
);
let inventory = merge_remote_threads(
"vm",
vec![console.clone()],
remote_threads::Discovery {
threads: vec![old.clone(), current.clone()],
..Default::default()
},
);
assert_eq!(inventory.rows.len(), 2);
assert!(inventory.rows.iter().all(|row| row.kind == "THREAD"));
assert!(
inventory.remote[0].console_name.is_none(),
"old hook must not attach the wrong conversation"
);
assert_eq!(
inventory.remote[1].console_name.as_deref(),
Some("relay-name")
);
let mut exited = console;
exited.running = false;
let inventory = merge_remote_threads(
"vm",
vec![exited],
remote_threads::Discovery {
threads: vec![old, current],
..Default::default()
},
);
assert!(
inventory
.remote
.iter()
.all(|row| row.console_name.is_none())
);
assert_eq!(
inventory
.rows
.iter()
.filter(|row| row.is_interesting())
.count(),
2
);
}
#[test]
fn multiple_grok_tabs_in_one_process_do_not_guess_the_focused_conversation() {
let mut first = remote_threads::tests::thread("grok", "first");
first.active = true;
first.console_name = Some("relay".into());
let mut second = first.clone();
second.thread.id = "second".into();
let inventory = merge_remote_threads(
"vm",
vec![ConsoleSession {
name: "relay".into(),
kind: "SHELL".into(),
running: true,
attached: true,
created_at: None,
command: Some("grok".into()),
snapshot: None,
}],
remote_threads::Discovery {
threads: vec![first, second],
..Default::default()
},
);
assert!(
inventory
.remote
.iter()
.all(|row| row.console_name.is_none())
);
}
fn request() -> LaunchRequest {
LaunchRequest {
project_id: "proj_1".into(),
environment_id: "env_prod".into(),
agent_id: None,
session_name: None,
force_new: false,
new_session: false,
harness: "claude".into(),
prompt: None,
label: "devtools/production".into(),
base: Default::default(),
}
}
#[test]
fn a_new_session_request_pins_its_agent_and_creates_nothing() {
let args = launch_args_for(&LaunchRequest {
agent_id: Some("ca_1".into()),
new_session: true,
..request()
});
assert_eq!(args.agent_id.as_deref(), Some("ca_1"));
assert!(!args.new, "must not ask the pipeline to create an agent");
assert_eq!(args.environment.as_deref(), Some("env_prod"));
assert_eq!(args.project.as_deref(), Some("proj_1"));
}
#[test]
fn a_new_agent_request_creates_and_pins_nothing() {
let args = launch_args_for(&LaunchRequest {
force_new: true,
..request()
});
assert!(args.new);
assert_eq!(args.agent_id, None);
}
#[test]
fn the_copied_ssh_command_opens_a_shell_on_the_vm() {
let command = ssh_command_for("env_1", "ca_1");
let args = shlex::split(&command).unwrap();
assert_eq!(&args[..2], &["ssh", "-t"]);
assert_eq!(
args.last().map(String::as_str),
Some(code::LOGIN_SHELL_COMMAND)
);
assert!(!command.contains("RAILWAY_DURABLE_SESSION_NAME"));
assert!(command.contains("agent:env_1:ca_1@"), "{command}");
assert!(!command.contains(" agent:env_1:ca_1 "), "{command}");
}
#[test]
fn a_command_line_launch_keeps_the_flags_it_arrived_with() {
use clap::Parser;
let base = LaunchArgs::parse_from(["code", "--new", "--name", "api", "--variable", "K=V"]);
let args = launch_args_for(&LaunchRequest {
base: Box::new(base.clone()),
force_new: true,
..request()
});
assert_eq!(
args,
base.retargeted(
"proj_1".into(),
"env_prod".into(),
"claude",
true,
None,
None
)
);
assert_eq!(args.environment.as_deref(), Some("env_prod"));
}
#[test]
fn a_prompt_request_carries_its_task_and_agent() {
let args = launch_args_for(&LaunchRequest {
agent_id: Some("ca_9".into()),
prompt: Some("fix the tests".into()),
..request()
});
assert_eq!(args.initial_prompt.as_deref(), Some("fix the tests"));
assert_eq!(args.agent_id.as_deref(), Some("ca_9"));
assert!(!args.new);
}
}