use std::{
env,
io::{Read, Write},
net::SocketAddr,
path::PathBuf,
sync::{Arc, Mutex},
};
use anyhow::{Context, Result};
use axum::{
extract::{ws::Message, Path, Query, RawQuery, State, WebSocketUpgrade},
http::{
header::{AUTHORIZATION, WWW_AUTHENTICATE},
HeaderMap, HeaderValue, Request, StatusCode,
},
middleware::{self, Next},
response::{IntoResponse, Response},
routing::{get, post},
Extension, Json, Router,
};
use base64::{
engine::general_purpose::{STANDARD as BASE64, URL_SAFE_NO_PAD as BASE64URL},
Engine,
};
use futures_util::{future, SinkExt, StreamExt};
use portable_pty::{native_pty_system, PtySize};
use rand::{distr::Alphanumeric, Rng};
use regex::Regex;
use serde::Deserialize;
use serde_json::json;
#[derive(rust_embed::RustEmbed)]
#[folder = "web/static"]
#[exclude = "vendor/*.map"]
#[exclude = "install/*"]
#[exclude = ".well-known/*"]
struct StaticAssets;
mod access;
mod cli;
mod config;
mod config_file;
mod configure;
mod db;
mod files;
mod host_suggestions;
mod local_stt;
mod local_tts;
mod mcp;
mod mcp_settings;
mod nodes;
mod pages_settings;
mod proxy;
mod push;
mod release_asset;
mod screen_model;
mod service;
mod session_history;
mod shell_integration;
mod speech_text;
mod ssl;
mod stt_debug;
mod terminal_cursor;
mod tmux;
mod transcribe;
mod twa;
mod update;
mod update_cli;
mod utf8_stream;
#[derive(Clone, Debug, PartialEq)]
enum InstallPhase {
Idle,
InstallingTools,
Running,
Success,
Failed(String),
}
impl InstallPhase {
fn is_active(&self) -> bool {
matches!(self, Self::InstallingTools | Self::Running)
}
}
#[derive(Clone)]
struct MissingHostPackages {
packages: Vec<String>,
install_command: Option<String>,
}
struct BackgroundJobState {
phase: InstallPhase,
output_tail: Vec<String>,
missing_host_packages: Option<MissingHostPackages>,
}
impl BackgroundJobState {
fn idle() -> Arc<tokio::sync::Mutex<Self>> {
Arc::new(tokio::sync::Mutex::new(Self {
phase: InstallPhase::Idle,
output_tail: vec![],
missing_host_packages: None,
}))
}
}
const JOB_OUTPUT_TAIL_LINES: usize = 200;
fn phase_parts(phase: &InstallPhase) -> (&'static str, Option<String>) {
match phase {
InstallPhase::Idle => ("idle", None),
InstallPhase::InstallingTools => ("installing_tools", None),
InstallPhase::Running => ("running", None),
InstallPhase::Success => ("success", None),
InstallPhase::Failed(e) => ("failed", Some(e.clone())),
}
}
fn strip_ansi(line: &str) -> String {
let mut out = String::with_capacity(line.len());
let mut chars = line.chars();
while let Some(c) = chars.next() {
if c != '\u{1b}' {
out.push(c);
continue;
}
match chars.next() {
Some('[') => {
for c in chars.by_ref() {
if ('\u{40}'..='\u{7e}').contains(&c) {
break;
}
}
}
Some(']') => {
while let Some(c) = chars.next() {
if c == '\u{7}' {
break;
}
if c == '\u{1b}' {
chars.next();
break;
}
}
}
_ => {}
}
}
out
}
async fn stream_into_job(
job: &Arc<tokio::sync::Mutex<BackgroundJobState>>,
mut command: tokio::process::Command,
) -> Result<(), String> {
use std::process::Stdio;
use tokio::io::{AsyncBufReadExt, BufReader};
let mut child = match command
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
{
Ok(c) => c,
Err(e) => return Err(format!("spawn error: {e}")),
};
let stdout = child.stdout.take().expect("stdout piped");
let stderr = child.stderr.take().expect("stderr piped");
let job_for_stdout = job.clone();
let job_for_stderr = job.clone();
let stdout_task = tokio::spawn(async move {
let mut reader = BufReader::new(stdout).lines();
while let Ok(Some(line)) = reader.next_line().await {
let mut guard = job_for_stdout.lock().await;
if guard.output_tail.len() >= JOB_OUTPUT_TAIL_LINES {
guard.output_tail.remove(0);
}
guard.output_tail.push(strip_ansi(&line));
}
});
let stderr_task = tokio::spawn(async move {
let mut last_line = String::new();
let mut reader = BufReader::new(stderr).lines();
while let Ok(Some(line)) = reader.next_line().await {
let line = strip_ansi(&line);
let mut guard = job_for_stderr.lock().await;
if guard.output_tail.len() >= JOB_OUTPUT_TAIL_LINES {
guard.output_tail.remove(0);
}
guard.output_tail.push(line.clone());
drop(guard);
last_line = line;
}
last_line
});
let _ = stdout_task.await;
let stderr_summary = stderr_task.await.unwrap_or_default();
match child.wait().await {
Ok(s) if s.success() => Ok(()),
Ok(s) => Err(format!(
"exit {}: {}",
s.code().unwrap_or(-1),
stderr_summary
)),
Err(e) => Err(format!("wait error: {e}")),
}
}
async fn run_streaming_job(
job: Arc<tokio::sync::Mutex<BackgroundJobState>>,
command: tokio::process::Command,
) {
let outcome = stream_into_job(&job, command).await;
let mut guard = job.lock().await;
guard.phase = match outcome {
Ok(()) => InstallPhase::Success,
Err(e) => InstallPhase::Failed(e),
};
}
async fn record_host_package_gap(
job: &Arc<tokio::sync::Mutex<BackgroundJobState>>,
gap: twa::HostPackageGap,
) {
let mut guard = job.lock().await;
guard.phase = InstallPhase::Failed(gap.message);
guard.missing_host_packages = Some(MissingHostPackages {
packages: gap.packages,
install_command: gap.install_command,
});
}
async fn push_job_line(job: &Arc<tokio::sync::Mutex<BackgroundJobState>>, line: String) {
let mut guard = job.lock().await;
if guard.output_tail.len() >= JOB_OUTPUT_TAIL_LINES {
guard.output_tail.remove(0);
}
guard.output_tail.push(line);
}
async fn run_twa_build_job(
job: Arc<tokio::sync::Mutex<BackgroundJobState>>,
setup: Option<tokio::process::Command>,
build: tokio::process::Command,
) {
if let Some(setup) = setup {
{
let mut guard = job.lock().await;
guard.phase = InstallPhase::InstallingTools;
}
push_job_line(
&job,
"Installing the Android build tools. This takes a few minutes.".to_string(),
)
.await;
if let Err(e) = stream_into_job(&job, setup).await {
let mut guard = job.lock().await;
guard.phase = InstallPhase::Failed(format!("installing the build tools failed: {e}"));
return;
}
let mut guard = job.lock().await;
guard.phase = InstallPhase::Running;
}
push_job_line(&job, "Building the Android package.".to_string()).await;
run_streaming_job(job, build).await;
}
#[derive(Clone)]
struct AppState {
session_name_re: Arc<Regex>,
auth: Option<AuthConfig>,
cache_bust: String,
db: Arc<db::Db>,
internal_token: Arc<String>,
config: Arc<config::Config>,
data_dir: PathBuf,
config_dir: PathBuf,
update: update::UpdateState,
build_hash: String,
twa_build: Arc<tokio::sync::Mutex<BackgroundJobState>>,
session_history: Arc<session_history::SessionHistoryStore>,
mcp: Arc<mcp_settings::McpServer>,
pages: Arc<pages_settings::PagesSettings>,
}
const SESSION_COOKIE_NAME: &str = "mobux_session";
#[derive(Clone)]
struct AuthConfig {
user: String,
pass: String,
session_cookie_name: String,
session_cookie_value: String,
}
#[tokio::main]
async fn main() -> Result<()> {
let options = match cli::parse(env::args().skip(1)) {
cli::Parsed::Run(options) => options,
cli::Parsed::Service(command) => std::process::exit(service::run(&command)),
cli::Parsed::Update(command) => std::process::exit(update_cli::run(command).await),
cli::Parsed::Configure(command) => std::process::exit(configure::run(&command)),
cli::Parsed::Help => {
print!("{}", cli::help_text(PKG_VERSION));
return Ok(());
}
cli::Parsed::Version => {
println!("mobux {PKG_VERSION}");
return Ok(());
}
cli::Parsed::Invalid(message) => {
eprintln!("mobux: {message}\n");
eprint!("{}", cli::help_text(PKG_VERSION));
std::process::exit(2);
}
};
rustls::crypto::aws_lc_rs::default_provider()
.install_default()
.map_err(|_| anyhow::anyhow!("failed to install rustls crypto provider"))?;
let settings = Arc::new(resolve_config(&options)?);
let config_dir = config::config_dir();
let auth = load_auth_config(&config_dir, &settings);
if auth.is_none() {
eprintln!("{}", config::NO_AUTH_WARNING);
}
let file_roots = Arc::new(files::FileRoots::from_config(&settings.files)?);
let proxy_targets = Arc::new(proxy::ProxyTargets::from_config(
&settings,
SESSION_COOKIE_NAME,
)?);
let data_dir = resolve_data_dir(&settings)?;
std::fs::create_dir_all(&data_dir)
.with_context(|| format!("creating data dir: {}", data_dir.display()))?;
let db_path = data_dir.join("mobux.db");
println!("data dir: {}", data_dir.display());
let db = Arc::new(db::Db::open(&db_path)?);
let _ = db.vapid_keys()?;
let internal_token: String = (&mut rand::rng())
.sample_iter(Alphanumeric)
.take(32)
.map(char::from)
.collect();
let port = settings.server.port;
let use_tls = settings.tls.enabled;
let build_hash = StaticAssets::get("build-info.json")
.and_then(|f| serde_json::from_slice::<serde_json::Value>(&f.data).ok())
.and_then(|v| v["hash"].as_str().map(str::to_owned))
.unwrap_or_else(|| "unknown".to_string());
let session_name_re = Arc::new(Regex::new(r"^[a-zA-Z0-9_-]+$")?);
let config_file = Arc::new(config_file::ConfigFile::new(
options
.config_path
.clone()
.unwrap_or_else(config::config_file_path),
));
let env = config::EnvSnapshot::from_env();
let mcp_server = Arc::new(mcp_settings::McpServer::new(
mcp::Context::new(session_name_re.clone(), db.clone(), &settings),
settings.clone(),
config_file.clone(),
mcp_settings::ManagedBy::detect(&env, &options.overrides),
));
let pages = Arc::new(pages_settings::PagesSettings::new(
file_roots.clone(),
proxy_targets.clone(),
config_file,
pages_settings::Managed::detect(&env),
));
let update_state = update::UpdateState::new(settings.update.check_url.clone());
update::spawn_checker(update_state.clone());
let state = AppState {
session_name_re,
auth,
cache_bust: format!(
"{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs()
),
db,
internal_token: Arc::new(internal_token),
config: settings.clone(),
data_dir: data_dir.clone(),
config_dir: config_dir.clone(),
update: update_state,
build_hash,
twa_build: BackgroundJobState::idle(),
session_history: Arc::new(session_history::SessionHistoryStore::new(&data_dir)),
mcp: mcp_server,
pages,
};
let internal_app = Router::new()
.route("/internal/trigger", post(api_internal_trigger))
.with_state(state.clone());
let internal_listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
let internal_port = internal_listener.local_addr()?.port();
tokio::spawn(async move {
if let Err(e) = axum::serve(internal_listener, internal_app).await {
eprintln!("internal listener error: {e:#}");
}
});
if let Err(e) = tmux::install_bell_hook(internal_port, &state.internal_token).await {
eprintln!("warning: failed to install tmux alert-bell hook: {e:#}");
} else {
println!("tmux alert-bell hook installed (internal port {internal_port})");
}
state.mcp.start_configured().await;
let state_for_mw = state.clone();
if !file_roots.is_empty() {
println!("files: serving {} root(s) under /files/", file_roots.len());
}
if !proxy_targets.is_empty() {
println!(
"proxy: forwarding {} target(s) under /proxy/",
proxy_targets.len()
);
}
let routes = app_routes(state.clone());
if settings.access.is_configured() {
serve_access_listener(&settings.access, routes.clone()).await?;
}
let app = routes.layer(middleware::from_fn_with_state(
state_for_mw,
auth_middleware,
));
let addr = SocketAddr::from(([0, 0, 0, 0], port));
if state.auth.is_some() {
println!("auth: enabled (HTTP Basic)");
} else {
println!(
"auth: disabled (pass --pin, or set MOBUX_PIN or MOBUX_AUTH_USER/MOBUX_AUTH_PASS)"
);
}
if let Some(warning) = clear_text_auth_warning(
state.auth.is_some(),
use_tls,
settings.server.behind_tls_proxy,
) {
eprintln!("{warning}");
}
println!("telemetry: /api/telemetry active, logs to stderr");
if settings.app.dev {
println!("dev mode: ON (MOBUX_DEV)");
}
if use_tls {
let cert_file = settings.tls.cert_file.trim();
let key_file = settings.tls.key_file.trim();
let (cert_path, key_path) = if !cert_file.is_empty() && !key_file.is_empty() {
eprintln!("[ssl] Using provided cert: {cert_file}, key: {key_file}");
(PathBuf::from(cert_file), PathBuf::from(key_file))
} else {
let challenges = if ssl::acme_mode_enabled(&settings.tls) {
let c = ssl::new_acme_challenges();
spawn_acme_http_server(c.clone(), settings.tls.acme_http_port).await?;
Some(c)
} else {
None
};
let paths = ssl::ensure_certs(&config_dir, &settings.tls, challenges).await?;
(paths.cert, paths.key)
};
let tls_config = ssl::load_rustls_config(&cert_path, &key_path)?;
let rustls_config =
axum_server::tls_rustls::RustlsConfig::from_config(std::sync::Arc::new(tls_config));
println!("mobux listening on https://{}", addr);
axum_server::bind_rustls(addr, rustls_config)
.serve(app.into_make_service_with_connect_info::<SocketAddr>())
.await?;
} else {
println!("mobux listening on http://{}", addr);
let listener = tokio::net::TcpListener::bind(addr).await?;
axum::serve(
listener,
app.into_make_service_with_connect_info::<SocketAddr>(),
)
.await?;
}
Ok(())
}
async fn serve_access_listener(access_config: &config::AccessConfig, routes: Router) -> Result<()> {
let addr = SocketAddr::from(([127, 0, 0, 1], access_config.port));
let listener = tokio::net::TcpListener::bind(addr).await.with_context(|| {
format!(
"access: cannot bind the Access listener on port {}",
access_config.port
)
})?;
let gate = Arc::new(access::AccessGuard::new(
access_config,
reqwest::Client::new(),
));
let app = routes.layer(middleware::from_fn_with_state(gate, access::guard));
println!(
"access: listening on http://{addr} for Cloudflare Access team {}",
access_config.team_origin()
);
tokio::spawn(async move {
if let Err(e) = axum::serve(
listener,
app.into_make_service_with_connect_info::<SocketAddr>(),
)
.await
{
eprintln!("access listener error: {e:#}");
}
});
Ok(())
}
fn resolve_config(options: &cli::RunOptions) -> Result<config::Config> {
let explicit = options.config_path.is_some();
let path = options
.config_path
.clone()
.unwrap_or_else(config::config_file_path);
let file = match config::load_partial_from(&path) {
Ok(Some(file)) => file,
Ok(None) if explicit => {
return Err(anyhow::anyhow!("{}: no such config file", path.display()))
}
Ok(None) => config::PartialConfig::default(),
Err(err) => return Err(err.into()),
};
let env = config::EnvSnapshot::from_env();
if let Some(message) = config::port_deprecation(&env, &options.overrides) {
eprintln!("[config] {message}");
}
let settings = config::resolve(
config::Config::default(),
file,
&env,
options.overrides.clone(),
);
config::check_access(&settings).map_err(|message| anyhow::anyhow!("config: {message}"))?;
Ok(settings)
}
fn resolve_data_dir(settings: &config::Config) -> Result<PathBuf> {
let data_dir = settings.paths.data_dir.trim();
if !data_dir.is_empty() {
return Ok(PathBuf::from(data_dir));
}
let dirs = directories::ProjectDirs::from("", "", "mobux")
.ok_or_else(|| anyhow::anyhow!("could not resolve user home directory for data dir"))?;
Ok(dirs.data_dir().to_path_buf())
}
async fn spawn_acme_http_server(challenges: ssl::AcmeChallenges, port: u16) -> Result<()> {
let addr = SocketAddr::from(([0, 0, 0, 0], port));
let router = Router::new()
.route(
"/.well-known/acme-challenge/{token}",
get(serve_acme_challenge),
)
.layer(Extension(challenges));
let listener = tokio::net::TcpListener::bind(addr).await.map_err(|e| {
anyhow::anyhow!(
"ACME mode: failed to bind HTTP listener on {addr} for HTTP-01 challenges \
(set MOBUX_ACME_HTTP_PORT to override): {e}"
)
})?;
eprintln!("[ssl] ACME: HTTP-01 challenge server listening on http://{addr}");
tokio::spawn(async move {
if let Err(e) = axum::serve(listener, router).await {
eprintln!("[ssl] ACME HTTP server exited with error: {e}");
}
});
Ok(())
}
async fn serve_acme_challenge(
Path(token): Path<String>,
Extension(challenges): Extension<ssl::AcmeChallenges>,
) -> Response {
match ssl::lookup_acme_challenge(&challenges, &token) {
Some(value) => (
StatusCode::OK,
[(axum::http::header::CONTENT_TYPE, "text/plain")],
value,
)
.into_response(),
None => (StatusCode::NOT_FOUND, "unknown acme challenge token").into_response(),
}
}
fn ensure_session_cookie_value(config_dir: &std::path::Path) -> String {
let path = config_dir.join("session-cookie");
if let Ok(existing) = std::fs::read_to_string(&path) {
let trimmed = existing.trim();
if trimmed.len() >= 32 {
return trimmed.to_string();
}
}
let value: String = rand::rng()
.sample_iter(&Alphanumeric)
.take(32)
.map(char::from)
.collect();
if let Some(parent) = path.parent() {
let _ = std::fs::create_dir_all(parent);
}
if let Err(e) = std::fs::write(&path, &value) {
eprintln!(
"[auth] WARN: could not persist session cookie to {}: {e}. \
Restarts will re-prompt clients for credentials.",
path.display()
);
} else {
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let _ = std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600));
}
}
value
}
fn load_auth_config(config_dir: &std::path::Path, settings: &config::Config) -> Option<AuthConfig> {
let credentials = settings.credentials()?;
Some(AuthConfig {
user: credentials.user,
pass: credentials.pass,
session_cookie_name: SESSION_COOKIE_NAME.to_string(),
session_cookie_value: ensure_session_cookie_value(config_dir),
})
}
fn is_public_path(path: &str) -> bool {
path == "/api/update/test-index" || PUBLIC_PATHS.iter().any(|public| public.matches(path))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PublicPath {
Exact(&'static str),
Prefix(&'static str),
}
impl PublicPath {
pub fn matches(self, path: &str) -> bool {
match self {
PublicPath::Exact(exact) => path == exact,
PublicPath::Prefix(prefix) => path.starts_with(prefix),
}
}
}
pub const PUBLIC_PATHS: &[PublicPath] = &[
PublicPath::Exact("/api/identify"),
PublicPath::Exact("/install"),
PublicPath::Prefix("/install/"),
PublicPath::Prefix("/.well-known/"),
PublicPath::Prefix("/static/icon-"),
PublicPath::Exact("/static/manifest.json"),
PublicPath::Exact("/sw.js"),
];
pub const ACCESS_PUBLIC_PATHS: &[PublicPath] = &[
PublicPath::Exact("/.well-known/assetlinks.json"),
PublicPath::Exact("/static/manifest.json"),
PublicPath::Prefix("/static/icon-"),
PublicPath::Exact("/sw.js"),
];
fn is_access_public_path(path: &str) -> bool {
ACCESS_PUBLIC_PATHS
.iter()
.any(|public| public.matches(path))
}
fn clear_text_auth_warning(
auth_enabled: bool,
use_tls: bool,
behind_tls_proxy: bool,
) -> Option<String> {
if !auth_enabled || use_tls || behind_tls_proxy {
return None;
}
let rule = "─".repeat(68);
Some(format!(
"{rule}\n\
WARNING auth is on and TLS is off. The password and the session\n\
cookie travel in clear text, and the cookie loses its Secure flag.\n\
Turn HTTPS on with --tls, MOBUX_TLS=1, or {{\"tls\": {{\"enabled\": true}}}}\n\
in config.json. Plain HTTP is only safe behind a TLS proxy — say so\n\
with --behind-tls-proxy to keep the cookie Secure and silence this.\n\
{rule}"
))
}
fn cookie_path(base_path: &str) -> String {
let trimmed = base_path.trim().trim_end_matches('/');
if trimmed.is_empty() {
return "/".to_string();
}
if trimmed.starts_with('/') {
return trimmed.to_string();
}
format!("/{trimmed}")
}
fn build_session_cookie(name: &str, value: &str, settings: &config::Config) -> String {
let path = cookie_path(&settings.server.base_path);
let secure = if settings.tls.enabled || settings.server.behind_tls_proxy {
"; Secure"
} else {
""
};
format!("{name}={value}; Path={path}; HttpOnly; SameSite=Lax{secure}; Max-Age=2592000")
}
async fn auth_middleware(
State(state): State<AppState>,
req: Request<axum::body::Body>,
next: Next,
) -> Response {
let Some(auth) = &state.auth else {
return next.run(req).await;
};
if is_public_path(req.uri().path()) {
return next.run(req).await;
}
let cookie_ok = req
.headers()
.get(axum::http::header::COOKIE)
.and_then(|v| v.to_str().ok())
.map(|cookie| {
cookie
.split(';')
.filter_map(|p| p.trim().split_once('='))
.any(|(k, v)| k == auth.session_cookie_name && v == auth.session_cookie_value)
})
.unwrap_or(false);
if cookie_ok {
return next.run(req).await;
}
let basic_ok = req
.headers()
.get(AUTHORIZATION)
.and_then(|v| v.to_str().ok())
.and_then(|v| v.strip_prefix("Basic "))
.and_then(|b64| BASE64.decode(b64).ok())
.and_then(|bytes| String::from_utf8(bytes).ok())
.and_then(|pair| {
let mut parts = pair.splitn(2, ':');
let user = parts.next()?.to_string();
let pass = parts.next()?.to_string();
Some((user, pass))
})
.map(|(user, pass)| user == auth.user && pass == auth.pass)
.unwrap_or(false);
if basic_ok {
let mut resp = next.run(req).await;
let set_cookie = build_session_cookie(
&auth.session_cookie_name,
&auth.session_cookie_value,
&state.config,
);
if let Ok(v) = HeaderValue::from_str(&set_cookie) {
resp.headers_mut().append(axum::http::header::SET_COOKIE, v);
}
return resp;
}
let mut resp = (StatusCode::UNAUTHORIZED, "Authentication required").into_response();
resp.headers_mut().insert(
WWW_AUTHENTICATE,
HeaderValue::from_static("Basic realm=\"mobux\""),
);
resp
}
async fn root_redirect() -> impl IntoResponse {
(
axum::http::StatusCode::TEMPORARY_REDIRECT,
[
(axum::http::header::LOCATION, "app"),
(axum::http::header::CACHE_CONTROL, "no-store"),
],
)
}
async fn api_sessions(
State(state): State<AppState>,
Query(q): Query<NodeQuery>,
) -> Result<Json<Vec<tmux::Session>>, AppError> {
let target = resolve_node_target(&state, q.node.as_deref()).await?;
let sessions = tmux::list_sessions(target.as_deref())
.await
.map_err(AppError::bad_request)?;
Ok(Json(sessions))
}
#[derive(serde::Serialize)]
struct Identify {
app: String,
version: String,
}
async fn api_identify() -> Json<Identify> {
Json(Identify {
app: "mobux".to_string(),
version: env!("CARGO_PKG_VERSION").to_string(),
})
}
async fn api_build_info(
State(state): State<AppState>,
via_access: Option<Extension<access::ViaAccess>>,
) -> Json<serde_json::Value> {
let via_access = via_access.is_some();
Json(json!({
"version": PKG_VERSION,
"build_hash": state.build_hash,
"dev_mode": state.config.app.dev,
"files": state.pages.files().names(),
"proxies": state.pages.proxies().names(),
"via_access": via_access,
"upload_limit_bytes": upload_limit_bytes(via_access),
}))
}
const UPLOAD_LIMIT_BYTES: u64 = 200 * 1024 * 1024;
const ACCESS_UPLOAD_LIMIT_BYTES: u64 = 100 * 1024 * 1024;
const MULTIPART_ENVELOPE_BYTES: u64 = 64 * 1024;
fn upload_limit_bytes(via_access: bool) -> u64 {
if via_access {
return ACCESS_UPLOAD_LIMIT_BYTES;
}
UPLOAD_LIMIT_BYTES
}
fn upload_too_large(limit: u64) -> AppError {
AppError {
status: StatusCode::PAYLOAD_TOO_LARGE,
message: format!(
"upload refused: the file is larger than the {} MB limit on this connection",
limit / (1024 * 1024)
),
}
}
async fn api_update_status(State(state): State<AppState>) -> Json<update::UpdateStatus> {
let mut status = state.update.status().await;
status.last_run_error = update::last_run_error(&state.data_dir);
Json(status)
}
async fn api_update_check(State(state): State<AppState>) -> Json<update::UpdateStatus> {
let mut status = state.update.refresh().await;
if status.error.is_none() {
update::clear_last_run_error(&state.data_dir);
}
status.last_run_error = update::last_run_error(&state.data_dir);
Json(status)
}
async fn api_update_run(State(state): State<AppState>) -> Response {
let status = state.update.status().await;
let Some(latest) = status.latest.clone() else {
let err = update::RunError::NoUpdateAvailable {
message: "no latest version known yet; run a check first".to_string(),
};
return (StatusCode::CONFLICT, Json(json!({ "error": err }))).into_response();
};
if !status.available {
let err = update::RunError::NoUpdateAvailable {
message: format!(
"already on the latest version ({})",
update::UpdateState::current_version()
),
};
return (StatusCode::CONFLICT, Json(json!({ "error": err }))).into_response();
}
if !state.update.try_begin_run() {
let err = update::RunError::AlreadyRunning {
message: "an update is already in progress".to_string(),
};
return (StatusCode::CONFLICT, Json(json!({ "error": err }))).into_response();
}
match update::spawn_updater(
&state.update,
&state.data_dir,
&state.config.app.service_name,
&latest,
state.config.server.port,
state.config.tls.enabled,
) {
Ok(log_path) => {
(
StatusCode::ACCEPTED,
Json(json!({
"started": true,
"version": latest,
"log": log_path.to_string_lossy(),
})),
)
.into_response()
}
Err(err) => {
state.update.end_run();
let status = match err {
update::RunError::NotSystemd { .. } => StatusCode::PRECONDITION_FAILED,
_ => StatusCode::INTERNAL_SERVER_ERROR,
};
(status, Json(json!({ "error": err }))).into_response()
}
}
}
async fn api_update_test_index() -> impl IntoResponse {
let body = std::env::var("MOBUX_UPDATE_TEST_INDEX").unwrap_or_default();
([(axum::http::header::CONTENT_TYPE, "text/plain")], body)
}
#[derive(Deserialize)]
struct CreateReq {
name: String,
}
async fn api_create_session(
State(state): State<AppState>,
Query(q): Query<NodeQuery>,
Json(payload): Json<CreateReq>,
) -> Result<Json<serde_json::Value>, AppError> {
let name = payload.name.trim();
validate_session_name(&state, name)?;
let target = resolve_node_target(&state, q.node.as_deref()).await?;
tmux::new_session(
name,
target.as_deref(),
&state.config.session.shell,
&state.data_dir,
)
.await
.map_err(AppError::bad_request)?;
Ok(Json(json!({"ok": true, "name": name})))
}
async fn api_kill_session(
State(state): State<AppState>,
Path(name): Path<String>,
Query(q): Query<NodeQuery>,
) -> Result<Json<serde_json::Value>, AppError> {
validate_session_name(&state, &name)?;
let target = resolve_node_target(&state, q.node.as_deref()).await?;
tmux::kill_session(&name, target.as_deref())
.await
.map_err(AppError::bad_request)?;
Ok(Json(json!({"ok": true})))
}
#[derive(Deserialize)]
struct RenameReq {
name: String,
}
async fn api_rename_session(
State(state): State<AppState>,
Path(old_name): Path<String>,
Query(q): Query<NodeQuery>,
Json(payload): Json<RenameReq>,
) -> Result<Json<serde_json::Value>, AppError> {
validate_session_name(&state, &old_name)?;
validate_session_name(&state, &payload.name)?;
let target = resolve_node_target(&state, q.node.as_deref()).await?;
tmux::rename_session(&old_name, &payload.name, target.as_deref())
.await
.map_err(AppError::bad_request)?;
Ok(Json(json!({"ok": true})))
}
#[derive(serde::Serialize)]
#[serde(rename_all = "camelCase")]
struct PaneJson {
#[serde(flatten)]
pane: tmux::Pane,
markers_seen: bool,
}
async fn api_list_panes(
State(state): State<AppState>,
Path(name): Path<String>,
Query(q): Query<NodeQuery>,
) -> Result<Json<Vec<PaneJson>>, AppError> {
validate_session_name(&state, &name)?;
let target = resolve_node_target(&state, q.node.as_deref()).await?;
let panes = tmux::list_panes(&name, target.as_deref())
.await
.map_err(AppError::bad_request)?;
let markers_seen = match (&target, panes.first()) {
(None, Some(pane)) => {
let history = state.session_history.clone();
let tmux_session = pane.tmux_session.clone();
tokio::task::spawn_blocking(move || history.markers_seen(&name, &tmux_session))
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?
}
_ => false,
};
Ok(Json(
panes
.into_iter()
.map(|pane| PaneJson { pane, markers_seen })
.collect(),
))
}
async fn api_select_pane(
State(state): State<AppState>,
Path((name, pane)): Path<(String, String)>,
Query(q): Query<NodeQuery>,
) -> Result<Json<serde_json::Value>, AppError> {
validate_session_name(&state, &name)?;
let target = resolve_node_target(&state, q.node.as_deref()).await?;
tmux::select_pane(&name, &pane, target.as_deref())
.await
.map_err(AppError::bad_request)?;
Ok(Json(json!({"ok": true})))
}
const HISTORY_MAX_LINES: u32 = 10000;
#[derive(Deserialize)]
struct SessionHistoryQuery {
node: Option<String>,
scope: Option<tmux::HistoryScope>,
lines: Option<u32>,
}
async fn api_session_history(
State(state): State<AppState>,
Path(name): Path<String>,
Query(q): Query<SessionHistoryQuery>,
) -> Result<(HeaderMap, String), AppError> {
validate_session_name(&state, &name)?;
let target = resolve_node_target(&state, q.node.as_deref()).await?;
let lines = q
.lines
.unwrap_or(HISTORY_MAX_LINES)
.clamp(1, HISTORY_MAX_LINES);
let history =
tmux::capture_history(&name, lines, q.scope.unwrap_or_default(), target.as_deref())
.await
.map_err(AppError::bad_request)?;
let mut headers = HeaderMap::new();
if history.continues {
headers.insert("x-history-continues", HeaderValue::from_static("1"));
}
Ok((headers, history.text))
}
#[derive(Deserialize)]
struct ConversationHistoryQuery {
cursor: Option<String>,
limit: Option<usize>,
tail: Option<usize>,
before: Option<String>,
}
async fn api_session_conversation(
State(state): State<AppState>,
Path(name): Path<String>,
Query(q): Query<ConversationHistoryQuery>,
) -> Result<Json<serde_json::Value>, AppError> {
validate_session_name(&state, &name)?;
if q.tail.is_some() && q.cursor.is_some() {
return Err(AppError::bad_request(anyhow::anyhow!(
"tail and cursor are mutually exclusive"
)));
}
if q.tail.is_some() && q.limit.is_some() {
return Err(AppError::bad_request(anyhow::anyhow!(
"tail carries its own count; limit is ambiguous alongside it"
)));
}
if q.tail.is_some() && q.before.is_some() {
return Err(AppError::bad_request(anyhow::anyhow!(
"tail and before are mutually exclusive"
)));
}
if q.cursor.is_some() && q.before.is_some() {
return Err(AppError::bad_request(anyhow::anyhow!(
"cursor and before name opposite directions"
)));
}
let cursor = match q.cursor {
Some(raw) => Some(
session_history::decode_cursor(&raw)
.ok_or_else(|| AppError::bad_request(anyhow::anyhow!("invalid cursor")))?,
),
None => None,
};
let before = match q.before {
Some(raw) => Some(
session_history::decode_cursor(&raw)
.ok_or_else(|| AppError::bad_request(anyhow::anyhow!("invalid cursor")))?,
),
None => None,
};
let limit = q
.limit
.unwrap_or(session_history::DEFAULT_LIMIT)
.clamp(1, session_history::MAX_LIMIT);
let history = state.session_history.clone();
let walking_back = before.is_some() && q.tail.is_none();
let page = match (q.tail, before) {
(Some(tail), _) => {
let count = tail.clamp(1, session_history::MAX_LIMIT);
tokio::task::spawn_blocking(move || history.read_tail(&name, count)).await
}
(None, Some(before)) => {
tokio::task::spawn_blocking(move || history.read_before(&name, before, limit)).await
}
(None, None) => {
tokio::task::spawn_blocking(move || history.read_page(&name, cursor, limit)).await
}
}
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?;
let mut body = json!({
"entries": page.entries,
"prevCursor": session_history::encode_cursor(page.prev_seq, page.prev_offset),
"hasOlder": page.has_older,
});
if !walking_back {
body["nextCursor"] = json!(session_history::encode_cursor(
page.next_seq,
page.next_offset
));
}
Ok(Json(body))
}
#[derive(Deserialize)]
struct CommandReq {
command: String,
}
async fn api_tmux_command(
State(state): State<AppState>,
Path(name): Path<String>,
Query(q): Query<NodeQuery>,
Json(payload): Json<CommandReq>,
) -> Result<Json<serde_json::Value>, AppError> {
validate_session_name(&state, &name)?;
let target = resolve_node_target(&state, q.node.as_deref()).await?;
let result = tmux::run_command(&name, &payload.command, target.as_deref())
.await
.map_err(AppError::bad_request)?;
Ok(Json(json!({"ok": true, "output": result})))
}
async fn api_telemetry(body: String) -> StatusCode {
let ts = chrono::Local::now().format("%H:%M:%S%.3f");
eprintln!("[telemetry {ts}] {body}");
StatusCode::NO_CONTENT
}
fn sanitize_upload_filename(filename: &str) -> String {
filename
.chars()
.map(|c| {
if c.is_alphanumeric() || c == '.' || c == '-' || c == '_' {
c
} else {
'_'
}
})
.collect()
}
async fn api_upload(
State(state): State<AppState>,
Query(q): Query<NodeQuery>,
via_access: Option<Extension<access::ViaAccess>>,
headers: HeaderMap,
mut multipart: axum::extract::Multipart,
) -> Result<Json<serde_json::Value>, AppError> {
use std::fs;
use std::path::PathBuf;
let limit = upload_limit_bytes(via_access.is_some());
let declared = headers
.get(axum::http::header::CONTENT_LENGTH)
.and_then(|v| v.to_str().ok())
.and_then(|v| v.parse::<u64>().ok());
if declared.is_some_and(|length| length > limit + MULTIPART_ENVELOPE_BYTES) {
return Err(upload_too_large(limit));
}
let target = resolve_node_target(&state, q.node.as_deref()).await?;
let upload_dir = PathBuf::from("/tmp/mobux-uploads");
if target.is_none() {
fs::create_dir_all(&upload_dir).map_err(|e| AppError::bad_request(e.into()))?;
}
if let Some(mut field) = multipart
.next_field()
.await
.map_err(|e| AppError::bad_request(e.into()))?
{
let filename = field.file_name().unwrap_or("upload").to_string();
let safe_name = sanitize_upload_filename(&filename);
let ts = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis();
let dest_filename = format!("{ts}-{safe_name}");
let mut data = Vec::new();
while let Some(chunk) = field
.chunk()
.await
.map_err(|e| AppError::bad_request(e.into()))?
{
if (data.len() + chunk.len()) as u64 > limit {
return Err(upload_too_large(limit));
}
data.extend_from_slice(&chunk);
}
let dest_display = match target.as_deref() {
None => {
let dest = upload_dir.join(&dest_filename);
fs::write(&dest, &data).map_err(|e| AppError::bad_request(e.into()))?;
dest.to_string_lossy().into_owned()
}
Some(ssh_target) => {
let upload_dir_str = upload_dir
.to_str()
.ok_or_else(|| AppError::internal(anyhow::anyhow!("upload dir not utf-8")))?;
tmux::write_remote_file(ssh_target, upload_dir_str, &dest_filename, &data)
.await
.map_err(AppError::internal)?
}
};
return Ok(Json(json!({
"path": dest_display,
"size": data.len(),
"name": safe_name,
})));
}
Err(AppError::bad_request(anyhow::anyhow!("no file in upload")))
}
async fn api_transcribe(
State(state): State<AppState>,
mut multipart: axum::extract::Multipart,
) -> Result<Json<serde_json::Value>, AppError> {
let mut audio_bytes: Option<Vec<u8>> = None;
let mut filename = "speech.wav".to_string();
while let Some(field) = multipart
.next_field()
.await
.map_err(|e| anyhow::anyhow!("multipart: {e}"))
.map_err(AppError::bad_request)?
{
if field.name() == Some("audio") {
if let Some(fname) = field.file_name() {
filename = fname.to_string();
}
audio_bytes = Some(
field
.bytes()
.await
.map_err(|e| anyhow::anyhow!("read field: {e}"))
.map_err(AppError::bad_request)?
.to_vec(),
);
} else {
let _ = field.bytes().await;
}
}
let audio = audio_bytes
.ok_or_else(|| AppError::bad_request(anyhow::anyhow!("missing 'audio' field")))?;
let (provider, debug_ctx) = tokio::task::spawn_blocking({
let db = state.db.clone();
move || active_provider(&db)
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?;
let debug_audio = audio.clone();
let debug_filename = filename.clone();
let data_dir = state.data_dir.clone();
let started = std::time::Instant::now();
let result = match &provider {
transcribe::Provider::InProcess { model } => {
local_stt::transcribe(state.data_dir.clone(), model.clone(), audio)
.await
.map_err(transcribe::TranscribeError::ProviderUnavailable)
}
transcribe::Provider::Remote(cfg) => {
transcribe::transcribe_with_provider(cfg, audio, &filename).await
}
};
let elapsed = started.elapsed();
let debug_outcome = match &result {
Ok(text) => Ok(text.clone()),
Err(e) => Err(e.to_string()),
};
tokio::task::spawn_blocking(move || {
stt_debug::store_clip(
&data_dir,
&debug_audio,
&debug_filename,
&debug_ctx,
elapsed,
&debug_outcome,
);
});
match result {
Ok(text) => Ok(Json(json!({ "text": text }))),
Err(transcribe::TranscribeError::ProviderUnavailable(msg)) => Err(AppError {
status: StatusCode::SERVICE_UNAVAILABLE,
message: msg,
}),
Err(e) => Err(AppError::internal(anyhow::anyhow!("{e}"))),
}
}
#[derive(Deserialize)]
struct PushSubscribeReq {
endpoint: String,
p256dh: String,
auth: String,
label: Option<String>,
}
#[derive(Deserialize)]
struct PushUnsubscribeReq {
endpoint: String,
}
fn decode_b64url(field: &str, value: &str) -> Result<Vec<u8>, AppError> {
BASE64URL
.decode(value)
.map_err(|e| AppError::bad_request(anyhow::anyhow!("invalid base64url in '{field}': {e}")))
}
async fn api_push_vapid_public_key(
State(state): State<AppState>,
) -> Result<Json<serde_json::Value>, AppError> {
let keys = state
.db
.vapid_keys()
.map_err(|e| AppError::internal(anyhow::anyhow!("loading vapid keys: {e}")))?;
Ok(Json(json!({ "key": BASE64URL.encode(&keys.public_key) })))
}
async fn api_push_subscribe(
State(state): State<AppState>,
Json(payload): Json<PushSubscribeReq>,
) -> Result<StatusCode, AppError> {
if payload.endpoint.trim().is_empty() {
return Err(AppError::bad_request(anyhow::anyhow!(
"endpoint must not be empty"
)));
}
let p256dh = decode_b64url("p256dh", &payload.p256dh)?;
let auth = decode_b64url("auth", &payload.auth)?;
let label = payload
.label
.map(|l| l.trim().to_string())
.filter(|l| !l.is_empty());
state
.db
.insert_subscription(db::NewSubscription {
endpoint: payload.endpoint,
p256dh,
auth,
label,
})
.map_err(|e| AppError::internal(anyhow::anyhow!("storing subscription: {e}")))?;
Ok(StatusCode::NO_CONTENT)
}
async fn api_push_unsubscribe(
State(state): State<AppState>,
Json(payload): Json<PushUnsubscribeReq>,
) -> Result<StatusCode, AppError> {
if payload.endpoint.trim().is_empty() {
return Err(AppError::bad_request(anyhow::anyhow!(
"endpoint must not be empty"
)));
}
state
.db
.remove_subscription(&payload.endpoint)
.map_err(|e| AppError::internal(anyhow::anyhow!("removing subscription: {e}")))?;
Ok(StatusCode::NO_CONTENT)
}
async fn api_push_devices(
State(state): State<AppState>,
) -> Result<Json<Vec<serde_json::Value>>, AppError> {
let subs = state
.db
.list_subscriptions()
.map_err(|e| AppError::internal(anyhow::anyhow!("listing subscriptions: {e}")))?;
let out: Vec<serde_json::Value> = subs
.into_iter()
.map(|s| {
json!({
"id": s.id,
"label": s.label,
"created_at": s.created_at,
"last_seen_at": s.last_seen_at,
})
})
.collect();
Ok(Json(out))
}
#[derive(Deserialize)]
struct PushNotifyRequest {
title: Option<String>,
body: String,
tag: Option<String>,
url: Option<String>,
}
async fn api_push_notify(
State(state): State<AppState>,
Json(req): Json<PushNotifyRequest>,
) -> Result<StatusCode, AppError> {
if req.body.trim().is_empty() {
return Err(AppError::bad_request(anyhow::anyhow!("body is required")));
}
let payload = push::Payload {
title: req.title.unwrap_or_else(|| "mobux".to_string()),
body: req.body,
tag: req.tag,
url: req.url,
};
tokio::spawn(push::notify(
state.db.clone(),
state.config.push.vapid_contact.clone(),
payload,
));
Ok(StatusCode::NO_CONTENT)
}
#[derive(Deserialize)]
struct InternalTriggerQuery {
kind: String,
session: String,
window: Option<String>,
}
#[derive(serde::Deserialize)]
struct SttModelsQuery {
kind: Option<String>,
host: Option<String>,
port: Option<String>,
}
async fn api_internal_trigger(
State(state): State<AppState>,
headers: HeaderMap,
axum::extract::Query(q): axum::extract::Query<InternalTriggerQuery>,
) -> StatusCode {
let token = headers
.get("X-Mobux-Token")
.and_then(|v| v.to_str().ok())
.unwrap_or("");
if token != state.internal_token.as_str() {
return StatusCode::UNAUTHORIZED;
}
if !state.session_name_re.is_match(&q.session) {
return StatusCode::BAD_REQUEST;
}
match q.kind.as_str() {
"bell" => {
let prefs = state.db.notification_prefs().unwrap_or_default();
if prefs.bell {
push::fire_bell(
state.db.clone(),
state.config.push.vapid_contact.clone(),
&q.session,
q.window.as_deref(),
);
}
}
_ => return StatusCode::BAD_REQUEST,
}
StatusCode::NO_CONTENT
}
#[derive(serde::Serialize, Deserialize)]
struct NotifPrefsJson {
bell: bool,
bell_emoji: bool,
program_exit: bool,
program_exit_nonzero: bool,
}
impl From<db::NotificationPrefs> for NotifPrefsJson {
fn from(p: db::NotificationPrefs) -> Self {
Self {
bell: p.bell,
bell_emoji: p.bell_emoji,
program_exit: p.program_exit,
program_exit_nonzero: p.program_exit_nonzero,
}
}
}
impl From<NotifPrefsJson> for db::NotificationPrefs {
fn from(j: NotifPrefsJson) -> Self {
Self {
bell: j.bell,
bell_emoji: j.bell_emoji,
program_exit: j.program_exit,
program_exit_nonzero: j.program_exit_nonzero,
}
}
}
async fn api_get_notification_prefs(
State(state): State<AppState>,
) -> Result<Json<NotifPrefsJson>, AppError> {
let prefs = state
.db
.notification_prefs()
.map_err(|e| AppError::internal(anyhow::anyhow!("reading prefs: {e}")))?;
Ok(Json(prefs.into()))
}
async fn api_set_notification_prefs(
State(state): State<AppState>,
Json(req): Json<NotifPrefsJson>,
) -> Result<StatusCode, AppError> {
state
.db
.set_notification_prefs(req.into())
.map_err(|e| AppError::internal(anyhow::anyhow!("writing prefs: {e}")))?;
Ok(StatusCode::NO_CONTENT)
}
async fn api_get_mcp_settings(
State(state): State<AppState>,
) -> Result<Json<mcp_settings::Status>, AppError> {
state
.mcp
.status()
.await
.map(Json)
.map_err(|message| AppError {
status: StatusCode::CONFLICT,
message,
})
}
#[derive(Deserialize)]
struct McpSettingsPut {
port: u16,
}
async fn api_set_mcp_settings(
State(state): State<AppState>,
Json(req): Json<McpSettingsPut>,
) -> Result<Json<mcp_settings::Status>, AppError> {
use mcp_settings::SetError;
state.mcp.set(req.port).await.map(Json).map_err(|err| {
let (status, message) = match err {
SetError::Invalid(message) => (StatusCode::BAD_REQUEST, message),
SetError::Managed(message) | SetError::Bind(message) => (StatusCode::CONFLICT, message),
SetError::Io(message) => (StatusCode::INTERNAL_SERVER_ERROR, message),
};
AppError { status, message }
})
}
async fn api_get_pages_settings(State(state): State<AppState>) -> Json<pages_settings::Status> {
Json(state.pages.status())
}
async fn api_set_pages_settings(
State(state): State<AppState>,
Json(req): Json<pages_settings::Change>,
) -> Result<Json<pages_settings::Status>, AppError> {
use pages_settings::SetError;
state.pages.set(req).await.map(Json).map_err(|err| {
let (status, message) = match err {
SetError::Invalid(message) => (StatusCode::BAD_REQUEST, message),
SetError::Managed(message) => (StatusCode::CONFLICT, message),
SetError::Io(message) => (StatusCode::INTERNAL_SERVER_ERROR, message),
};
AppError { status, message }
})
}
#[derive(serde::Serialize, Deserialize)]
struct UiPrefsJson {
renderer: String,
theme: String,
default_view: String,
osc133_hint_dismissed: bool,
listen_voice: String,
listen_rate: f64,
listen_pitch: f64,
#[serde(default)]
selected_node: String,
}
impl From<db::UiPreferences> for UiPrefsJson {
fn from(p: db::UiPreferences) -> Self {
let renderer = if p.renderer == "sterk" {
"sterk"
} else {
"xterm"
}
.to_string();
let default_view = match p.default_view.as_str() {
"reader" => "reader",
"read" => "read",
_ => "xterm",
}
.to_string();
Self {
renderer,
theme: p.theme,
default_view,
osc133_hint_dismissed: p.osc133_hint_dismissed,
listen_voice: p.listen_voice,
listen_rate: p.listen_rate.clamp(0.5, 2.0),
listen_pitch: p.listen_pitch.clamp(0.5, 2.0),
selected_node: p.selected_node,
}
}
}
impl UiPrefsJson {
fn validate(self) -> Result<db::UiPreferences, String> {
if self.renderer != "xterm" && self.renderer != "sterk" {
return Err(format!(
"invalid renderer {:?}: must be \"xterm\" or \"sterk\"",
self.renderer
));
}
if !matches!(self.default_view.as_str(), "xterm" | "reader" | "read") {
return Err(format!(
"invalid default_view {:?}: must be \"xterm\", \"reader\" or \"read\"",
self.default_view
));
}
Ok(db::UiPreferences {
renderer: self.renderer,
theme: self.theme,
default_view: self.default_view,
osc133_hint_dismissed: self.osc133_hint_dismissed,
listen_voice: self.listen_voice,
listen_rate: self.listen_rate.clamp(0.5, 2.0),
listen_pitch: self.listen_pitch.clamp(0.5, 2.0),
selected_node: self.selected_node,
})
}
}
async fn api_get_ui_preferences(
State(state): State<AppState>,
) -> Result<Json<UiPrefsJson>, AppError> {
let prefs = state
.db
.ui_preferences()
.map_err(|e| AppError::internal(anyhow::anyhow!("reading ui preferences: {e}")))?;
Ok(Json(prefs.into()))
}
async fn api_set_ui_preferences(
State(state): State<AppState>,
Json(req): Json<UiPrefsJson>,
) -> Result<StatusCode, AppError> {
let prefs = req
.validate()
.map_err(|e| AppError::bad_request(anyhow::anyhow!(e)))?;
state
.db
.set_ui_preferences(prefs)
.map_err(|e| AppError::internal(anyhow::anyhow!("writing ui preferences: {e}")))?;
Ok(StatusCode::NO_CONTENT)
}
#[derive(Deserialize)]
struct ShellIntegrationReq {
shell: shell_integration::Shell,
}
async fn api_shell_integration_status() -> Result<Json<shell_integration::Status>, AppError> {
let s = shell_integration::status().map_err(AppError::internal)?;
Ok(Json(s))
}
async fn api_shell_integration_install(
Json(req): Json<ShellIntegrationReq>,
) -> Result<Json<shell_integration::Status>, AppError> {
let s = shell_integration::install(req.shell).map_err(AppError::internal)?;
Ok(Json(s))
}
async fn api_shell_integration_uninstall(
Json(req): Json<ShellIntegrationReq>,
) -> Result<Json<shell_integration::Status>, AppError> {
let s = shell_integration::uninstall(req.shell).map_err(AppError::internal)?;
Ok(Json(s))
}
const PKG_VERSION: &str = env!("CARGO_PKG_VERSION");
async fn settings_page() -> impl IntoResponse {
(
axum::http::StatusCode::TEMPORARY_REDIRECT,
[
(axum::http::header::LOCATION, "app#/settings"),
(axum::http::header::CACHE_CONTROL, "no-store"),
],
)
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum SessionLocation {
Local,
Node(String),
}
fn resolve_session_location(
name: &str,
local: &[tmux::Session],
nodes: &[(&str, &[tmux::Session])],
) -> Option<SessionLocation> {
let mut found = Vec::new();
if local.iter().any(|s| s.name == name) {
found.push(SessionLocation::Local);
}
for (node_name, sessions) in nodes {
if sessions.iter().any(|s| s.name == name) {
found.push(SessionLocation::Node((*node_name).to_string()));
}
}
let mut it = found.into_iter();
let first = it.next()?;
if it.next().is_some() {
return None; }
Some(first)
}
const SESSION_PROBE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
async fn with_probe_timeout<T, Fut>(fut: Fut) -> Option<T>
where
Fut: std::future::Future<Output = Result<T>>,
{
tokio::time::timeout(SESSION_PROBE_TIMEOUT, fut)
.await
.ok()
.and_then(Result::ok)
}
async fn probe_sessions_or_absent(target: Option<&str>) -> Vec<tmux::Session> {
with_probe_timeout(tmux::list_sessions(target))
.await
.unwrap_or_default()
}
async fn locate_session(state: &AppState, name: &str) -> Result<Option<SessionLocation>, AppError> {
let db_nodes = tokio::task::spawn_blocking({
let db = state.db.clone();
move || db.list_nodes()
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?;
let local_probe = probe_sessions_or_absent(None);
let node_probes = db_nodes
.iter()
.map(|n| probe_sessions_or_absent(Some(n.target.as_str())));
let (local_sessions, node_session_lists) =
tokio::join!(local_probe, future::join_all(node_probes));
let node_sessions: Vec<(String, Vec<tmux::Session>)> = db_nodes
.into_iter()
.zip(node_session_lists)
.map(|(n, sessions)| (n.name, sessions))
.collect();
let node_refs: Vec<(&str, &[tmux::Session])> = node_sessions
.iter()
.map(|(name, sessions)| (name.as_str(), sessions.as_slice()))
.collect();
Ok(resolve_session_location(name, &local_sessions, &node_refs))
}
async fn terminal_page(
State(state): State<AppState>,
Path(name): Path<String>,
RawQuery(query): RawQuery,
) -> Result<impl IntoResponse, AppError> {
validate_session_name(&state, &name)?;
let query = query
.filter(|query| !query.is_empty())
.map(|query| format!("?{query}"))
.unwrap_or_default();
let location = match locate_session(&state, &name).await? {
Some(SessionLocation::Local) => format!("../app{query}#/s/{name}"),
Some(SessionLocation::Node(node)) => format!("../app{query}#/s/{node}/{name}"),
None => format!("../app{query}#/"),
};
Ok((
axum::http::StatusCode::TEMPORARY_REDIRECT,
[
(axum::http::header::LOCATION, location),
(axum::http::header::CACHE_CONTROL, "no-store".to_string()),
],
))
}
async fn serve_sw(State(state): State<AppState>) -> impl axum::response::IntoResponse {
use axum::http::header;
let body = format!(
"{}\n// sw-version: {}\n",
include_str!("../web/static/sw.js"),
state.cache_bust,
);
(
[
(header::CONTENT_TYPE, "text/javascript"),
(header::CACHE_CONTROL, "no-store, must-revalidate"),
],
body,
)
}
async fn serve_static(Path(path): Path<String>) -> Response {
use axum::http::header;
match StaticAssets::get(&path) {
Some(file) => {
let mime = file.metadata.mimetype();
let mut resp = (StatusCode::OK, file.data).into_response();
let h = resp.headers_mut();
if let Ok(v) = HeaderValue::from_str(mime) {
h.insert(header::CONTENT_TYPE, v);
}
h.insert(
header::CACHE_CONTROL,
HeaderValue::from_static("no-store, must-revalidate"),
);
resp
}
None => (StatusCode::NOT_FOUND, "not found").into_response(),
}
}
fn app_routes(state: AppState) -> Router {
let app = Router::new()
.route("/", get(root_redirect))
.route("/api/identify", get(api_identify))
.route("/api/build-info", get(api_build_info))
.route("/api/sessions", get(api_sessions).post(api_create_session))
.route("/api/sessions/{name}/kill", post(api_kill_session))
.route("/api/sessions/{name}/rename", post(api_rename_session))
.route("/api/sessions/{name}/panes", get(api_list_panes))
.route(
"/api/sessions/{name}/panes/{pane}/select",
post(api_select_pane),
)
.route("/api/sessions/{name}/history", get(api_session_history))
.route(
"/api/sessions/{name}/conversation",
get(api_session_conversation),
)
.route("/api/sessions/{name}/command", post(api_tmux_command))
.route(
"/api/settings/nodes",
get(api_get_settings_nodes).put(api_set_settings_nodes),
)
.route("/api/nodes", get(api_nodes_status))
.route("/api/host-suggestions", get(api_host_suggestions))
.route(
"/api/telemetry",
post(api_telemetry).layer(axum::extract::DefaultBodyLimit::max(64 * 1024)),
)
.route(
"/api/upload",
post(api_upload).layer(axum::extract::DefaultBodyLimit::max(
(UPLOAD_LIMIT_BYTES + MULTIPART_ENVELOPE_BYTES) as usize,
)),
)
.route(
"/transcribe",
post(api_transcribe).layer(axum::extract::DefaultBodyLimit::max(8 * 1024 * 1024)),
)
.route("/api/push/vapid-public-key", get(api_push_vapid_public_key))
.route(
"/api/push/subscribe",
post(api_push_subscribe).delete(api_push_unsubscribe),
)
.route("/api/push/devices", get(api_push_devices))
.route("/api/push/notify", post(api_push_notify))
.route(
"/api/settings/notifications",
get(api_get_notification_prefs).put(api_set_notification_prefs),
)
.route(
"/api/settings/mcp",
get(api_get_mcp_settings).put(api_set_mcp_settings),
)
.route(
"/api/settings/pages",
get(api_get_pages_settings).put(api_set_pages_settings),
)
.route(
"/api/settings/preferences",
get(api_get_ui_preferences).put(api_set_ui_preferences),
)
.route(
"/api/settings/stt",
get(api_get_stt_config).put(api_set_stt_config),
)
.route("/api/stt/status", get(api_stt_status))
.route("/api/stt/models", get(api_stt_models))
.route(
"/api/stt/install",
post(api_stt_install).layer(axum::extract::DefaultBodyLimit::max(1024)),
)
.route("/api/stt/install/status", get(api_stt_install_status))
.route("/api/install/apk/build", post(api_install_apk_build))
.route("/api/install/apk/status", get(api_install_apk_status))
.route("/api/tts/status", get(api_tts_status))
.route("/api/tts/prepare", post(api_tts_prepare))
.route(
"/api/tts/speak",
post(api_tts_speak).layer(axum::extract::DefaultBodyLimit::max(64 * 1024)),
)
.route(
"/api/shell-integration/status",
get(api_shell_integration_status),
)
.route(
"/api/shell-integration/install",
post(api_shell_integration_install),
)
.route(
"/api/shell-integration/uninstall",
post(api_shell_integration_uninstall),
)
.route("/api/update/status", get(api_update_status))
.route("/api/update/check", post(api_update_check))
.route("/api/update/run", post(api_update_run))
.route("/settings", get(settings_page))
.route("/s/{name}", get(terminal_page))
.route("/ws/{name}", get(terminal_ws))
.route("/sw.js", get(serve_sw))
.route("/install", get(install_page))
.route("/install/mobux.apk", get(serve_install_apk))
.route("/install/mobux-ca.crt", get(serve_install_ca))
.route("/.well-known/assetlinks.json", get(serve_assetlinks))
.route("/app", get(serve_spa_index))
.route("/app/{*rest}", get(spa_deep_link_redirect))
.route("/static/{*path}", get(serve_static))
.merge(mcp::absent());
let app = if std::env::var_os("MOBUX_UPDATE_TEST_INDEX").is_some() {
app.route("/api/update/test-index", get(api_update_test_index))
} else {
app
};
app.merge(files::router(state.pages.files().clone()))
.merge(proxy::router(state.pages.proxies().clone()))
.fallback(get(serve_spa_index))
.with_state(state)
}
async fn serve_spa_index() -> Response {
use axum::http::header;
match StaticAssets::get("spa/index.html") {
Some(file) => {
let mut resp = (StatusCode::OK, file.data).into_response();
let h = resp.headers_mut();
h.insert(
header::CONTENT_TYPE,
HeaderValue::from_static("text/html; charset=utf-8"),
);
h.insert(
header::CACHE_CONTROL,
HeaderValue::from_static("no-store, must-revalidate"),
);
resp
}
None => (
StatusCode::NOT_FOUND,
"SPA not built — run `node web/build.js` (or `make build`).",
)
.into_response(),
}
}
fn relative_app_target(rest: &str) -> String {
format!("{}app", "../".repeat(rest.split('/').count()))
}
async fn spa_deep_link_redirect(Path(rest): Path<String>) -> impl IntoResponse {
(
axum::http::StatusCode::TEMPORARY_REDIRECT,
[
(axum::http::header::LOCATION, relative_app_target(&rest)),
(axum::http::header::CACHE_CONTROL, "no-store".to_string()),
],
)
}
async fn install_page() -> impl IntoResponse {
(
axum::http::StatusCode::TEMPORARY_REDIRECT,
[
(axum::http::header::LOCATION, "app#/install"),
(axum::http::header::CACHE_CONTROL, "no-store"),
],
)
}
async fn serve_install_apk(State(state): State<AppState>) -> Response {
let path = twa::resolve_artifact(twa::apk_path(&state.data_dir), twa::CHECKOUT_APK_PATH);
serve_file_or_404(
path.to_string_lossy().as_ref(),
"application/vnd.android.package-archive",
Some("mobux.apk"),
)
.await
}
async fn serve_install_ca(State(state): State<AppState>) -> Response {
if ssl::acme_mode_enabled(&state.config.tls) {
return (StatusCode::NOT_FOUND, "ACME mode: no local CA to install").into_response();
}
let path = ssl::ca_cert_path(&state.config_dir);
serve_file_or_404(
path.to_string_lossy().as_ref(),
"application/x-x509-ca-cert",
Some("mobux-ca.crt"),
)
.await
}
async fn serve_assetlinks(State(state): State<AppState>) -> Response {
let path = twa::resolve_artifact(
twa::assetlinks_path(&state.data_dir),
twa::CHECKOUT_ASSETLINKS_PATH,
);
serve_file_or_404(path.to_string_lossy().as_ref(), "application/json", None).await
}
async fn serve_file_or_404(
path: &str,
content_type: &'static str,
download_name: Option<&'static str>,
) -> Response {
use axum::http::header;
let bytes = match tokio::fs::read(path).await {
Ok(b) => b,
Err(_) => return (StatusCode::NOT_FOUND, "not found").into_response(),
};
let mut resp = (
StatusCode::OK,
[(header::CONTENT_TYPE, content_type)],
bytes,
)
.into_response();
if let Some(name) = download_name {
let disp = format!("attachment; filename=\"{name}\"");
if let Ok(v) = HeaderValue::from_str(&disp) {
resp.headers_mut().insert(header::CONTENT_DISPOSITION, v);
}
}
resp
}
#[derive(Deserialize)]
struct NodeQuery {
node: Option<String>,
}
async fn resolve_node_target(
state: &AppState,
node: Option<&str>,
) -> Result<Option<String>, AppError> {
let Some(name) = node else { return Ok(None) };
let name = name.to_string();
let found = tokio::task::spawn_blocking({
let db = state.db.clone();
let name = name.clone();
move || db.get_node(&name)
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?;
let target = found.map(|n| n.target).ok_or_else(|| {
AppError::bad_request(anyhow::anyhow!(
"unknown node {name:?} — check Settings › Nodes"
))
})?;
if target.starts_with('-') {
return Err(AppError::bad_request(anyhow::anyhow!(
"node {name:?} has an invalid target {target:?} — check Settings › Nodes"
)));
}
Ok(Some(target))
}
#[derive(Deserialize)]
struct TerminalWsQuery {
node: Option<String>,
build: Option<String>,
}
const MAX_LOGGED_BUILD_LEN: usize = 64;
fn truncate_for_log(s: &str, max_chars: usize) -> &str {
match s.char_indices().nth(max_chars) {
Some((byte_idx, _)) => &s[..byte_idx],
None => s,
}
}
fn format_ws_attach_line(
outcome: &str,
session: &str,
node: Option<&str>,
target: &str,
build: Option<&str>,
server_build: &str,
user_agent: &str,
) -> String {
let build = build.map(|b| truncate_for_log(b, MAX_LOGGED_BUILD_LEN));
format!(
"[ws attach] {outcome} session={session:?} node={node:?} target={target} build={build:?} server_build={server_build} ua={user_agent:?}"
)
}
fn log_ws_attach(
state: &AppState,
outcome: &str,
session: &str,
node: Option<&str>,
target: &str,
build: Option<&str>,
user_agent: &str,
) {
eprintln!(
"{}",
format_ws_attach_line(
outcome,
session,
node,
target,
build,
&state.build_hash,
user_agent,
)
);
}
async fn terminal_ws(
State(state): State<AppState>,
Path(name): Path<String>,
Query(q): Query<TerminalWsQuery>,
headers: HeaderMap,
ws: WebSocketUpgrade,
) -> Result<Response, AppError> {
let user_agent = headers
.get(axum::http::header::USER_AGENT)
.and_then(|v| v.to_str().ok())
.unwrap_or("<none>")
.to_string();
let node = q.node.as_deref();
let build = q.build.as_deref();
if let Err(err) = validate_session_name(&state, &name) {
let outcome = format!("REJECTED[{}: {}]", err.status.as_u16(), err.message);
log_ws_attach(&state, &outcome, &name, node, "-", build, &user_agent);
return Err(err);
}
let ssh_target = match resolve_node_target(&state, node).await {
Ok(target) => target,
Err(err) => {
let outcome = format!("REJECTED[{}: {}]", err.status.as_u16(), err.message);
log_ws_attach(&state, &outcome, &name, node, "-", build, &user_agent);
return Err(err);
}
};
log_ws_attach(
&state,
"ok",
&name,
node,
ssh_target.as_deref().unwrap_or("local"),
build,
&user_agent,
);
let session_history = state.session_history.clone();
Ok(ws.on_upgrade(move |socket| async move {
if let Err(err) = handle_ws(socket, name, ssh_target, session_history).await {
eprintln!("ws error: {err:#}");
}
}))
}
const WS_PING_INTERVAL: std::time::Duration = std::time::Duration::from_secs(30);
#[derive(Deserialize)]
struct ResizeMsg {
#[serde(rename = "type")]
kind: String,
cols: u16,
rows: u16,
}
async fn handle_ws(
socket: axum::extract::ws::WebSocket,
session_name: String,
ssh_target: Option<String>,
session_history: Arc<session_history::SessionHistoryStore>,
) -> Result<()> {
let mut feeder_guard = session_history.try_acquire_feeder(&session_name);
let pty_system = native_pty_system();
let pair = pty_system.openpty(PtySize {
rows: 35,
cols: 120,
pixel_width: 0,
pixel_height: 0,
})?;
let tmux_bin = match (&ssh_target, std::env::var("MOBUX_TMUX_SOCKET")) {
(None, Ok(s)) if !s.is_empty() => format!("tmux -L {}", s),
_ => "tmux".to_string(),
};
let (mut history_tap, mut history_rx) = start_history_feed(
&ssh_target,
&tmux_bin,
&session_name,
feeder_guard.is_some(),
)
.await;
let mut tap_started_at = tokio::time::Instant::now();
let mut segmenter = match feeder_guard {
Some(_) => {
Some(new_segmenter(&tmux_bin, &ssh_target, &session_name, history_tap.as_ref()).await)
}
None => None,
};
let mut snapshot_at: Option<tokio::time::Instant> = None;
let mut pane_check_at = history_tap
.as_ref()
.map(|_| tokio::time::Instant::now() + PANE_RECHECK);
let cmd = nodes::build_attach_command(ssh_target.as_deref(), &tmux_bin, &session_name);
let mut child = pair.slave.spawn_command(cmd)?;
let mut reader = pair.master.try_clone_reader()?;
let writer = pair.master.take_writer()?;
let master = pair.master;
let writer = Arc::new(Mutex::new(writer));
let master = Arc::new(Mutex::new(master));
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<Vec<u8>>();
std::thread::spawn(move || {
let mut buf = vec![0u8; 8192];
loop {
match reader.read(&mut buf) {
Ok(0) => break,
Ok(n) => {
if tx.send(buf[..n].to_vec()).is_err() {
break;
}
}
Err(_) => break,
}
}
});
let (mut ws_sender, mut ws_receiver) = socket.split();
let mut ws_text = utf8_stream::Utf8Stream::new();
let mut keepalive = tokio::time::interval_at(
tokio::time::Instant::now() + WS_PING_INTERVAL,
WS_PING_INTERVAL,
);
keepalive.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
loop {
tokio::select! {
_ = keepalive.tick() => {
if ws_sender.send(Message::Ping(Default::default())).await.is_err() {
break;
}
}
maybe_hist = async { history_rx.as_mut().unwrap().recv().await }, if history_rx.is_some() => {
match maybe_hist {
Some(chunk) => {
if let Some(seg) = segmenter.as_mut() {
let produced = seg.feed(&chunk, session_history::now_ms());
snapshot_at = snapshot_deadline(seg);
record_history(&session_history, &session_name, seg.tmux_session(), produced).await;
}
}
None => {
history_rx = None;
history_tap = None;
snapshot_at = None;
pane_check_at = None;
if let Some(mut seg) = segmenter.take() {
let produced = seg.finish(session_history::now_ms());
record_history(&session_history, &session_name, seg.tmux_session(), produced).await;
if tap_started_at.elapsed() >= TAP_RESTART_MIN_AGE {
let (tap, rx) =
start_history_feed(&ssh_target, &tmux_bin, &session_name, true)
.await;
history_tap = tap;
history_rx = rx;
tap_started_at = tokio::time::Instant::now();
}
segmenter = Some(new_segmenter(&tmux_bin, &ssh_target, &session_name, history_tap.as_ref()).await);
pane_check_at = history_tap
.as_ref()
.map(|_| tokio::time::Instant::now() + PANE_RECHECK);
}
}
}
}
_ = tokio::time::sleep_until(snapshot_at.unwrap_or_else(tokio::time::Instant::now)), if snapshot_at.is_some() => {
snapshot_at = None;
if let Some(seg) = segmenter.as_mut() {
let produced = seg.tick(session_history::now_ms());
snapshot_at = snapshot_deadline(seg);
record_history(&session_history, &session_name, seg.tmux_session(), produced).await;
}
}
_ = tokio::time::sleep_until(pane_check_at.unwrap_or_else(tokio::time::Instant::now)), if pane_check_at.is_some() => {
pane_check_at = None;
if let (Some(tap), Some(seg)) = (history_tap.as_ref(), segmenter.as_mut()) {
match tmux::pane_screen(&tmux_bin, tap.pane_id()).await {
Ok(pane) => {
let produced = seg.resize(pane.rows, pane.cols, session_history::now_ms());
record_history(&session_history, &session_name, seg.tmux_session(), produced).await;
}
Err(err) => eprintln!("history: pane size for '{session_name}': {err:#}"),
}
pane_check_at = Some(tokio::time::Instant::now() + PANE_RECHECK);
}
}
maybe_out = rx.recv() => {
match maybe_out {
Some(chunk) => {
if feeder_guard.is_none() {
feeder_guard = session_history.try_acquire_feeder(&session_name);
if feeder_guard.is_some() {
let (tap, rx) =
start_history_feed(&ssh_target, &tmux_bin, &session_name, true)
.await;
history_tap = tap;
history_rx = rx;
tap_started_at = tokio::time::Instant::now();
segmenter = Some(new_segmenter(&tmux_bin, &ssh_target, &session_name, history_tap.as_ref()).await);
pane_check_at = history_tap
.as_ref()
.map(|_| tokio::time::Instant::now() + PANE_RECHECK);
}
}
if history_rx.is_none() {
if let Some(seg) = segmenter.as_mut() {
let produced = seg.feed(&chunk, session_history::now_ms());
snapshot_at = snapshot_deadline(seg);
record_history(&session_history, &session_name, seg.tmux_session(), produced).await;
}
}
let text = ws_text.decode(&chunk);
if text.is_empty() {
continue;
}
if ws_sender.send(Message::Text(text.into())).await.is_err() {
break;
}
}
None => {
let tail = ws_text.finish();
if !tail.is_empty() {
let _ = ws_sender.send(Message::Text(tail.into())).await;
}
break;
}
}
}
maybe_in = ws_receiver.next() => {
match maybe_in {
Some(Ok(msg)) => {
match msg {
Message::Text(t) => {
if let Ok(rz) = serde_json::from_str::<ResizeMsg>(&t) {
if rz.kind == "resize" && rz.cols > 0 && rz.rows > 0 {
if history_tap.is_some() {
pane_check_at = Some(tokio::time::Instant::now() + PANE_RESIZE_SETTLE);
}
if let Ok(m) = master.lock() {
let _ = m.resize(PtySize { rows: rz.rows, cols: rz.cols, pixel_width: 0, pixel_height: 0});
}
continue;
}
}
if let Ok(mut w) = writer.lock() {
let _ = w.write_all(t.as_bytes());
let _ = w.flush();
}
}
Message::Binary(b) => {
if let Ok(mut w) = writer.lock() {
let _ = w.write_all(&b);
let _ = w.flush();
}
}
Message::Close(_) => break,
Message::Ping(_) | Message::Pong(_) => {}
}
}
Some(Err(_)) | None => break,
}
}
}
}
if let Some(mut seg) = segmenter.take() {
let produced = seg.finish(session_history::now_ms());
record_history(
&session_history,
&session_name,
seg.tmux_session(),
produced,
)
.await;
}
drop(history_tap);
let _ = child.kill();
let _ = child.wait();
Ok(())
}
const PANE_RECHECK: std::time::Duration = std::time::Duration::from_secs(2);
const PANE_RESIZE_SETTLE: std::time::Duration = std::time::Duration::from_millis(200);
const TAP_RESTART_MIN_AGE: std::time::Duration = std::time::Duration::from_secs(2);
async fn new_segmenter(
tmux_bin: &str,
ssh_target: &Option<String>,
session_name: &str,
tap: Option<&tmux::PanePipeTap>,
) -> session_history::Segmenter {
let tmux_session = match ssh_target {
Some(_) => String::new(),
None => match tmux::session_stamp(tmux_bin, session_name).await {
Ok(stamp) => stamp,
Err(err) => {
eprintln!("history: recording without a session stamp: {err:#}");
String::new()
}
},
};
let Some(tap) = tap else {
return session_history::Segmenter::for_client_stream().recording(tmux_session);
};
let pane = match tmux::pane_screen(tmux_bin, tap.pane_id()).await {
Ok(pane) => pane,
Err(err) => {
eprintln!("history: recording without a screen model: {err:#}");
return session_history::Segmenter::for_client_stream().recording(tmux_session);
}
};
let mut segmenter =
session_history::Segmenter::for_pane(pane.rows, pane.cols).recording(tmux_session);
if pane.alternate {
segmenter.assume_alternate_screen();
}
segmenter
}
fn snapshot_deadline(segmenter: &session_history::Segmenter) -> Option<tokio::time::Instant> {
let wait_ms = segmenter.snapshot_due_at()? - session_history::now_ms();
Some(tokio::time::Instant::now() + std::time::Duration::from_millis(wait_ms.max(0) as u64))
}
async fn record_history(
history: &Arc<session_history::SessionHistoryStore>,
session_name: &str,
tmux_session: &str,
produced: Vec<session_history::PendingEntry>,
) {
if produced.is_empty() {
return;
}
let history = history.clone();
let name = session_name.to_string();
let tmux_session = tmux_session.to_string();
let _ = tokio::task::spawn_blocking(move || {
for entry in produced {
let _ = history.append(&name, &tmux_session, entry);
}
})
.await;
}
async fn start_history_feed(
ssh_target: &Option<String>,
tmux_bin: &str,
session_name: &str,
active: bool,
) -> (
Option<tmux::PanePipeTap>,
Option<tokio::sync::mpsc::UnboundedReceiver<Vec<u8>>>,
) {
if !active || ssh_target.is_some() {
return (None, None);
}
match tmux::PanePipeTap::start(tmux_bin, session_name).await {
Some((tap, rx)) => (Some(tap), Some(rx)),
None => (None, None),
}
}
fn is_valid_session_name(session_name_re: &Regex, name: &str) -> bool {
!name.is_empty() && session_name_re.is_match(name)
}
fn validate_session_name(state: &AppState, name: &str) -> Result<(), AppError> {
if !is_valid_session_name(&state.session_name_re, name) {
return Err(AppError::bad_request(anyhow::anyhow!(
"invalid session name"
)));
}
Ok(())
}
#[derive(serde::Serialize, Deserialize)]
struct NodeJson {
name: String,
target: String,
}
#[derive(serde::Serialize)]
struct NodesGetJson {
nodes: Vec<NodeJson>,
}
#[derive(Deserialize)]
struct NodesPutJson {
nodes: Vec<NodeJson>,
}
async fn api_get_settings_nodes(
State(state): State<AppState>,
) -> Result<Json<NodesGetJson>, AppError> {
let nodes = tokio::task::spawn_blocking({
let db = state.db.clone();
move || db.list_nodes()
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?;
Ok(Json(NodesGetJson {
nodes: nodes
.into_iter()
.map(|n| NodeJson {
name: n.name,
target: n.target,
})
.collect(),
}))
}
async fn api_set_settings_nodes(
State(state): State<AppState>,
Json(req): Json<NodesPutJson>,
) -> Result<StatusCode, AppError> {
let mut seen = std::collections::HashSet::new();
for n in &req.nodes {
if n.name.trim().is_empty() || n.target.trim().is_empty() {
return Err(AppError::bad_request(anyhow::anyhow!(
"node name and target must not be empty"
)));
}
if !seen.insert(n.name.clone()) {
return Err(AppError::bad_request(anyhow::anyhow!(
"duplicate node name: {}",
n.name
)));
}
}
let new_count = req.nodes.len();
tokio::task::spawn_blocking({
let db = state.db.clone();
let pairs = req.nodes.into_iter().map(|n| (n.name, n.target)).collect();
move || -> anyhow::Result<usize> {
let previous_count = db.list_nodes()?.len();
db.replace_nodes(pairs)?;
Ok(previous_count)
}
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)
.map(|previous_count| {
if previous_count > 0 && new_count == 0 {
eprintln!(
"warning: PUT /api/settings/nodes emptied the node list (was {previous_count}, now 0) — was this intentional?"
);
}
})?;
Ok(StatusCode::NO_CONTENT)
}
#[derive(serde::Serialize)]
struct NodeStatusJson {
name: String,
target: String,
reachable: bool,
}
#[derive(serde::Serialize)]
struct NodesStatusJson {
nodes: Vec<NodeStatusJson>,
}
async fn api_nodes_status(
State(state): State<AppState>,
) -> Result<Json<NodesStatusJson>, AppError> {
let db_nodes = tokio::task::spawn_blocking({
let db = state.db.clone();
move || db.list_nodes()
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?;
let probes = db_nodes
.iter()
.map(|n| nodes::probe_reachable(&n.target, std::time::Duration::from_secs(3)));
let reachable = future::join_all(probes).await;
let out = db_nodes
.into_iter()
.zip(reachable)
.map(|(n, reachable)| NodeStatusJson {
name: n.name,
target: n.target,
reachable,
})
.collect();
Ok(Json(NodesStatusJson { nodes: out }))
}
#[derive(serde::Serialize)]
struct HostSuggestionsJson {
hosts: Vec<host_suggestions::HostSuggestion>,
}
async fn api_host_suggestions() -> Json<HostSuggestionsJson> {
let ssh = host_suggestions::ssh_config_hosts();
let (tailscale, mdns) = tokio::join!(
host_suggestions::tailscale_hosts(std::time::Duration::from_millis(1500)),
host_suggestions::avahi_hosts(std::time::Duration::from_millis(2000)),
);
Json(HostSuggestionsJson {
hosts: host_suggestions::merge([ssh, tailscale, mdns]),
})
}
#[derive(serde::Serialize)]
struct SttProviderJson {
host: String,
port: String,
model: String,
has_key: bool,
}
#[derive(serde::Serialize)]
struct SttConfigGetJson {
#[serde(rename = "activeKind")]
active_kind: String,
providers: std::collections::HashMap<String, SttProviderJson>,
#[serde(rename = "localEngine")]
local_engine: bool,
}
#[derive(serde::Deserialize)]
struct SttConfigPutJson {
kind: String,
host: String,
port: String,
model: String,
#[serde(default)]
api_key: Option<String>,
}
async fn api_get_stt_config(
State(state): State<AppState>,
) -> Result<Json<SttConfigGetJson>, AppError> {
let (active_kind, providers) = tokio::task::spawn_blocking({
let db = state.db.clone();
move || -> anyhow::Result<_> {
let active_kind = db.stt_active_kind()?;
let rows = db.stt_all_providers()?;
Ok((active_kind, rows))
}
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?;
let mut map = std::collections::HashMap::new();
for row in &providers {
map.insert(
row.kind.clone(),
SttProviderJson {
host: row.host.clone(),
port: row.port.clone(),
model: row.model.clone(),
has_key: row.api_key.as_deref().is_some_and(|k| !k.is_empty()),
},
);
}
Ok(Json(SttConfigGetJson {
active_kind,
providers: map,
local_engine: local_stt::ENABLED,
}))
}
async fn api_set_stt_config(
State(state): State<AppState>,
Json(req): Json<SttConfigPutJson>,
) -> Result<StatusCode, AppError> {
let model = if req.kind == transcribe::LOCAL_KIND {
local_stt::resolve_model(&req.model).to_string()
} else {
req.model
};
let row = db::SttProviderRow {
kind: req.kind.clone(),
host: req.host,
port: req.port,
model,
api_key: req.api_key,
};
tokio::task::spawn_blocking({
let db = state.db.clone();
let kind = req.kind.clone();
move || -> anyhow::Result<()> {
db.set_stt_provider(row)?;
db.set_stt_active_kind(&kind)?;
Ok(())
}
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?;
Ok(StatusCode::NO_CONTENT)
}
struct SttStatus {
kind: String,
url: String,
model: String,
local: Option<local_stt::Phase>,
remote_reachable: bool,
}
impl SttStatus {
fn state(&self) -> &'static str {
match &self.local {
Some(phase) => phase.state(),
None if self.remote_reachable => "ready",
None => "unreachable",
}
}
fn ready(&self) -> bool {
self.state() == "ready"
}
fn into_json(self) -> serde_json::Value {
let state = self.state();
let ready = self.ready();
let mut body = json!({
"kind": self.kind,
"url": self.url,
"model": self.model,
"state": state,
"reachable": ready,
"installed": !matches!(self.local, Some(local_stt::Phase::NotDownloaded)),
"engine_available": !matches!(
self.local,
Some(local_stt::Phase::Disabled | local_stt::Phase::UnsupportedCpu(_))
),
});
match &self.local {
Some(phase) => {
body["message"] = serde_json::Value::String(phase.message());
if let local_stt::Phase::Downloading {
file,
downloaded,
total,
} = phase
{
body["progress"] = json!({
"file": file,
"downloaded": downloaded,
"total": total,
});
}
if let local_stt::Phase::Failed(err) = phase {
body["error"] = serde_json::Value::String(err.clone());
}
}
None if !ready => {
body["message"] = serde_json::Value::String(format!(
"The endpoint at {} did not answer.",
self.url
));
}
None => {}
}
body
}
}
fn active_provider(
db: &db::Db,
) -> anyhow::Result<(transcribe::Provider, stt_debug::ProviderContext)> {
let kind = db.stt_active_kind()?;
let row = db
.stt_provider(&kind)?
.unwrap_or_else(|| db::SttProviderRow::default_for(&kind));
let url = row.transcription_url();
let provider = transcribe::select_provider(&kind, &url, &row.model, row.api_key.as_deref());
let model = match &provider {
transcribe::Provider::InProcess { model } => model.clone(),
transcribe::Provider::Remote(cfg) => cfg.model.clone(),
};
let debug_ctx = stt_debug::ProviderContext {
kind,
model,
host: row.host,
port: row.port,
url,
};
Ok((provider, debug_ctx))
}
async fn api_stt_status(
State(state): State<AppState>,
) -> Result<Json<serde_json::Value>, AppError> {
let (provider, ctx) = tokio::task::spawn_blocking({
let db = state.db.clone();
move || active_provider(&db)
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?;
let status = match &provider {
transcribe::Provider::InProcess { model } => {
let phase = local_stt::phase(&state.data_dir, model);
if matches!(
phase,
local_stt::Phase::NotDownloaded
| local_stt::Phase::Downloading { .. }
| local_stt::Phase::Loading
) {
local_stt::prepare_in_background(state.data_dir.clone(), model.clone());
}
SttStatus {
kind: ctx.kind,
url: String::new(),
model: model.clone(),
local: Some(phase),
remote_reachable: false,
}
}
transcribe::Provider::Remote(cfg) => {
let reachable = transcribe::probe_transcribe(cfg).await;
SttStatus {
kind: ctx.kind,
url: cfg.url.clone(),
model: cfg.model.clone(),
local: None,
remote_reachable: reachable,
}
}
};
Ok(Json(status.into_json()))
}
async fn api_stt_install(State(state): State<AppState>) -> Result<impl IntoResponse, AppError> {
let (provider, _) = tokio::task::spawn_blocking({
let db = state.db.clone();
move || active_provider(&db)
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?;
let transcribe::Provider::InProcess { model } = provider else {
return Ok((
StatusCode::PRECONDITION_FAILED,
Json(json!({
"status": "not_local",
"error": "Only the local provider has a model to download.",
})),
));
};
if !local_stt::ENABLED {
return Ok((
StatusCode::PRECONDITION_FAILED,
Json(json!({
"status": "unsupported",
"error": local_stt::UNSUPPORTED_MESSAGE,
})),
));
}
local_stt::prepare_in_background(state.data_dir.clone(), model.clone());
Ok((
StatusCode::ACCEPTED,
Json(json!({"status": "started", "model": model})),
))
}
async fn api_stt_install_status(
State(state): State<AppState>,
) -> Result<Json<serde_json::Value>, AppError> {
let (provider, _) = tokio::task::spawn_blocking({
let db = state.db.clone();
move || active_provider(&db)
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?;
let phase = match &provider {
transcribe::Provider::InProcess { model } => local_stt::phase(&state.data_dir, model),
transcribe::Provider::Remote(_) => local_stt::Phase::Disabled,
};
Ok(Json(install_status_json(&phase)))
}
fn install_status_json(phase: &local_stt::Phase) -> serde_json::Value {
let (name, error) = match phase {
local_stt::Phase::Ready => ("success", None),
local_stt::Phase::Downloading { .. }
| local_stt::Phase::Verifying
| local_stt::Phase::Loading => ("running", None),
local_stt::Phase::NotDownloaded => ("idle", None),
local_stt::Phase::Disabled => ("failed", Some(local_stt::UNSUPPORTED_MESSAGE.to_string())),
local_stt::Phase::UnsupportedCpu(why) => ("failed", Some(why.clone())),
local_stt::Phase::Failed(err) => ("failed", Some(err.clone())),
};
json!({
"phase": name,
"output": [phase.message()],
"error": error,
})
}
fn apk_domain(config: &config::Config, headers: &HeaderMap) -> Result<String, String> {
let access_hostname = config
.access
.is_configured()
.then_some(config.access.hostname.as_str());
let host = headers
.get(axum::http::header::HOST)
.and_then(|h| h.to_str().ok());
twa::pinned_domain(access_hostname, Some(config.app.domain.as_str()), host)
}
async fn api_install_apk_build(
State(state): State<AppState>,
headers: HeaderMap,
) -> Result<Response, AppError> {
let domain = match apk_domain(&state.config, &headers) {
Ok(d) => d,
Err(e) => {
return Ok((
StatusCode::BAD_REQUEST,
Json(json!({"status": "invalid_domain", "error": e})),
)
.into_response())
}
};
{
let mut guard = state.twa_build.lock().await;
if guard.phase.is_active() {
return Ok((
StatusCode::CONFLICT,
Json(json!({"status": "already_running", "domain": domain})),
)
.into_response());
}
guard.phase = InstallPhase::Running;
guard.output_tail = vec![format!("Preparing the Android package for {domain}")];
guard.missing_host_packages = None;
}
let job = state.twa_build.clone();
let work_dir = twa::work_dir(&state.data_dir);
let data_dir = state.data_dir.clone();
let install_dir = state.data_dir.join("install");
let wellknown_dir = state.data_dir.join(".well-known");
let build_domain = domain.clone();
tokio::spawn(async move {
let scripts = match twa::materialize(&data_dir) {
Ok(s) => s,
Err(e) => {
let mut guard = job.lock().await;
guard.phase = InstallPhase::Failed(format!("{e:#}"));
return;
}
};
let toolchain_present = {
let build = scripts.build.clone();
tokio::task::spawn_blocking(move || twa::toolchain_present(&build))
.await
.unwrap_or(false)
};
if !toolchain_present {
let missing = twa::missing_host_packages_on_path();
if !missing.is_empty() {
record_host_package_gap(&job, twa::host_package_gap(&missing)).await;
return;
}
}
let setup = (!toolchain_present).then(|| {
let mut command = tokio::process::Command::new("bash");
command
.arg(&scripts.setup)
.current_dir(&work_dir)
.stdin(std::process::Stdio::null());
command
});
let mut build = tokio::process::Command::new("bash");
build
.arg(&scripts.build)
.current_dir(&work_dir)
.stdin(std::process::Stdio::null())
.env("MOBUX_DOMAIN", &build_domain)
.env("TWA_WORK_DIR", &work_dir)
.env("TWA_INSTALL_DIR", &install_dir)
.env("TWA_WELLKNOWN_DIR", &wellknown_dir);
run_twa_build_job(job.clone(), setup, build).await;
let mut guard = job.lock().await;
if guard.phase != InstallPhase::Success {
return;
}
if let Err(e) = twa::record_build(&data_dir, &build_domain) {
guard.phase = InstallPhase::Failed(format!("{e:#}"));
}
});
Ok((
StatusCode::ACCEPTED,
Json(json!({"status": "started", "domain": domain})),
)
.into_response())
}
async fn api_install_apk_status(
State(state): State<AppState>,
headers: HeaderMap,
) -> Result<Json<serde_json::Value>, AppError> {
let domain = apk_domain(&state.config, &headers);
let guard = state.twa_build.lock().await;
let (phase_str, error) = phase_parts(&guard.phase);
let gap = guard.missing_host_packages.as_ref();
Ok(Json(json!({
"phase": phase_str,
"output": guard.output_tail,
"error": error,
"missing_host_packages": gap.map(|g| g.packages.clone()),
"install_command": gap.map(|g| g.install_command.clone()),
"apk_available": twa::resolve_artifact(
twa::apk_path(&state.data_dir),
twa::CHECKOUT_APK_PATH,
)
.is_file(),
"domain": domain.as_deref().ok(),
"domain_error": domain.err(),
"apk_domain": twa::built_domain(&state.data_dir),
})))
}
#[derive(Debug, Deserialize)]
struct SpeakRequest {
text: String,
#[serde(default)]
kind: String,
#[serde(default)]
expand: bool,
#[serde(default)]
language: String,
}
async fn api_tts_status(State(state): State<AppState>) -> Json<serde_json::Value> {
let voice = local_tts::voice();
let phase = local_tts::phase(&state.data_dir);
Json(json!({
"enabled": local_tts::ENABLED,
"voice": voice,
"state": phase.state(),
"message": phase.message(),
}))
}
async fn api_tts_prepare(
State(state): State<AppState>,
) -> Result<Json<serde_json::Value>, AppError> {
let voice = local_tts::voice();
local_tts::ensure_ready(state.data_dir.clone())
.await
.map_err(|e| AppError::precondition(anyhow::anyhow!(e)))?;
let phase = local_tts::phase(&state.data_dir);
Ok(Json(json!({
"voice": voice,
"state": phase.state(),
"message": phase.message(),
})))
}
async fn api_tts_speak(
State(state): State<AppState>,
Json(req): Json<SpeakRequest>,
) -> Result<Response, AppError> {
let speech = speech_text::normalize(
speech_text::Kind::parse(&req.kind),
&req.text,
&speech_text::Options {
expand: req.expand,
language: req.language,
},
);
check_speakable(&speech)?;
if !local_tts::ENABLED {
return Ok(browser_speech(
&speech,
&local_tts::Phase::Disabled.message(),
));
}
match local_tts::synthesize(state.data_dir.clone(), speech.clone()).await {
Ok(clip) => Ok(([(axum::http::header::CONTENT_TYPE, "audio/wav")], clip).into_response()),
Err(err) => Ok(browser_speech(&speech, &err)),
}
}
const MAX_SPOKEN_CHARS: usize = 4000;
fn check_speakable(speech: &speech_text::Speech) -> Result<(), AppError> {
let chars = speech.text.chars().count();
if chars <= MAX_SPOKEN_CHARS {
return Ok(());
}
Err(AppError {
status: StatusCode::PAYLOAD_TOO_LARGE,
message: format!(
"This block is {chars} characters once cleaned up for speech; the voice reads at most {MAX_SPOKEN_CHARS} at a time. Read a shorter block."
),
})
}
fn browser_speech(speech: &speech_text::Speech, why: &str) -> Response {
Json(json!({
"engine": "browser",
"reason": why,
"text": speech.text,
"sentences": speech.sentences,
}))
.into_response()
}
async fn api_stt_models(
State(state): State<AppState>,
Query(q): Query<SttModelsQuery>,
) -> Result<Json<serde_json::Value>, AppError> {
use std::time::Duration;
let fallback_for_kind = |kind: &str| -> Vec<String> {
match kind {
transcribe::LOCAL_KIND => local_stt::model_ids(),
"openai" => vec![
"whisper-1".to_string(),
"gpt-4o-transcribe".to_string(),
"gpt-4o-mini-transcribe".to_string(),
],
_ => vec![
"Systran/faster-whisper-base.en".to_string(),
"Systran/faster-whisper-small.en".to_string(),
"Systran/faster-whisper-medium.en".to_string(),
],
}
};
let (base_url, api_key, kind) = if q.host.as_deref().map(|h| !h.is_empty()).unwrap_or(false) {
let raw_host = q.host.as_deref().unwrap_or("").trim_end_matches('/');
let host = if raw_host.contains("://") {
raw_host.to_string()
} else {
format!("http://{}", raw_host)
};
let port = q.port.as_deref().unwrap_or("");
let base = if port.is_empty() {
host
} else {
format!("{}:{}", host, port)
};
let k = q.kind.clone().unwrap_or_default();
let api_key = if k == "openai" {
let kc = k.clone();
tokio::task::spawn_blocking({
let db = state.db.clone();
move || db.stt_provider(&kc)
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?
.and_then(|r| r.api_key)
.filter(|k| !k.is_empty())
} else {
None
};
(base, api_key, k)
} else {
let requested = q.kind.clone().filter(|k| !k.is_empty());
tokio::task::spawn_blocking({
let db = state.db.clone();
move || -> anyhow::Result<_> {
let kind = match requested {
Some(kind) => kind,
None => db.stt_active_kind()?,
};
let row = db
.stt_provider(&kind)?
.unwrap_or_else(|| db::SttProviderRow::default_for(&kind));
let base = {
let raw = row.host.trim_end_matches('/');
let h = if raw.contains("://") {
raw.to_string()
} else {
format!("http://{}", raw)
};
if row.port.is_empty() {
h
} else {
format!("{}:{}", h, row.port)
}
};
let key = row.api_key.filter(|k| !k.is_empty());
Ok((base, key, kind))
}
})
.await
.map_err(|e| AppError::internal(anyhow::anyhow!("spawn_blocking: {e}")))?
.map_err(AppError::internal)?
};
if base_url.is_empty() || kind == transcribe::LOCAL_KIND {
return Ok(Json(serde_json::json!({
"models": fallback_for_kind(&kind)
})));
}
let models_url = format!("{}/v1/models", base_url.trim_end_matches('/'));
let client = match reqwest::Client::builder()
.timeout(Duration::from_secs(5))
.build()
{
Ok(c) => c,
Err(_) => {
return Ok(Json(
serde_json::json!({ "models": fallback_for_kind(&kind) }),
));
}
};
let mut req = client.get(&models_url);
if let Some(key) = &api_key {
req = req.bearer_auth(key);
}
let ids: Vec<String> = match req.send().await {
Ok(resp) if resp.status().is_success() => resp
.json::<serde_json::Value>()
.await
.ok()
.and_then(|v| v.get("data").cloned())
.and_then(|d| d.as_array().cloned())
.map(|arr| {
arr.iter()
.filter_map(|m| m.get("id").and_then(|id| id.as_str()).map(String::from))
.collect()
})
.filter(|v: &Vec<String>| !v.is_empty())
.unwrap_or_else(|| fallback_for_kind(&kind)),
_ => fallback_for_kind(&kind),
};
Ok(Json(serde_json::json!({ "models": ids })))
}
#[derive(Debug)]
struct AppError {
status: StatusCode,
message: String,
}
impl AppError {
fn bad_request(err: anyhow::Error) -> Self {
Self {
status: StatusCode::BAD_REQUEST,
message: err.to_string(),
}
}
fn internal(err: anyhow::Error) -> Self {
Self {
status: StatusCode::INTERNAL_SERVER_ERROR,
message: err.to_string(),
}
}
fn precondition(err: anyhow::Error) -> Self {
Self {
status: StatusCode::PRECONDITION_FAILED,
message: err.to_string(),
}
}
}
impl IntoResponse for AppError {
fn into_response(self) -> Response {
(self.status, self.message).into_response()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn an_access_port_in_use_stops_startup_naming_the_port() {
let held = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = held.local_addr().unwrap().port();
let access = config::AccessConfig {
port,
team_domain: "http://127.0.0.1:1".to_string(),
aud: "aud".to_string(),
allowed_emails: vec!["owner@example.com".to_string()],
..Default::default()
};
let error = serve_access_listener(&access, Router::new())
.await
.unwrap_err();
assert!(
format!("{error:#}").contains(&format!("port {port}")),
"{error:#}"
);
}
#[test]
fn a_block_too_long_to_speak_is_refused_with_a_413() {
let speech = |text: String| speech_text::Speech {
sentences: vec![text.clone()],
text,
};
assert!(check_speakable(&speech("a".repeat(MAX_SPOKEN_CHARS))).is_ok());
let err = check_speakable(&speech("a".repeat(MAX_SPOKEN_CHARS + 1)))
.expect_err("past the cap is refused");
assert_eq!(err.status, StatusCode::PAYLOAD_TOO_LARGE);
assert!(err.message.contains("at most 4000"), "{}", err.message);
}
#[test]
fn ws_attach_log_line_escapes_newlines_in_node_and_build() {
let line = format_ws_attach_line(
"ok",
"mysession",
Some("devbox\n[ws attach] ok session=\"forged\""),
"local",
Some("abcd\nEVIL"),
"serverhash",
"curl/8.0",
);
assert_eq!(line.lines().count(), 1);
assert!(line.contains(r#"node=Some("devbox\n[ws attach]"#));
assert!(line.contains(r#"build=Some("abcd\nEVIL")"#));
}
#[test]
fn ws_attach_log_line_absent_node_and_build_render_none() {
let line = format_ws_attach_line(
"ok",
"mysession",
None,
"local",
None,
"serverhash",
"curl/8.0",
);
assert!(line.contains("node=None"));
assert!(line.contains("build=None"));
}
#[test]
fn ws_attach_log_line_truncates_oversized_build() {
let huge = "x".repeat(10_000);
let line = format_ws_attach_line(
"ok",
"mysession",
None,
"local",
Some(&huge),
"serverhash",
"curl/8.0",
);
assert!(
line.len() < 500,
"an oversized build param must not produce an unbounded log line: {} bytes",
line.len()
);
}
#[test]
fn truncate_for_log_cuts_on_char_boundary() {
let s = "éééé";
assert_eq!(truncate_for_log(s, 3), "ééé");
}
#[test]
fn set_ui_preferences_rejects_invalid_renderer() {
let bad = UiPrefsJson {
renderer: "bogus".to_string(),
theme: "nord".to_string(),
default_view: "xterm".to_string(),
osc133_hint_dismissed: false,
listen_voice: String::new(),
listen_rate: 1.0,
listen_pitch: 1.0,
selected_node: String::new(),
};
let err = bad.validate().expect_err("bogus renderer must be rejected");
assert!(
err.contains("renderer"),
"error should name the field: {err}"
);
}
#[test]
fn set_ui_preferences_rejects_invalid_default_view() {
let bad = UiPrefsJson {
renderer: "xterm".to_string(),
theme: "nord".to_string(),
default_view: "bogus".to_string(),
osc133_hint_dismissed: false,
listen_voice: String::new(),
listen_rate: 1.0,
listen_pitch: 1.0,
selected_node: String::new(),
};
let err = bad
.validate()
.expect_err("bogus default_view must be rejected");
assert!(
err.contains("default_view"),
"error should name the field: {err}"
);
}
#[test]
fn set_ui_preferences_accepts_valid_enums_and_clamps_numerics() {
let ok = UiPrefsJson {
renderer: "sterk".to_string(),
theme: "nord".to_string(),
default_view: "reader".to_string(),
osc133_hint_dismissed: true,
listen_voice: "Daniel".to_string(),
listen_rate: 99.0, listen_pitch: -5.0, selected_node: "gpu-box".to_string(),
}
.validate()
.expect("valid enums must be accepted");
assert_eq!(ok.renderer, "sterk");
assert_eq!(ok.default_view, "reader");
assert_eq!(ok.listen_rate, 2.0);
assert_eq!(ok.listen_pitch, 0.5);
assert_eq!(ok.selected_node, "gpu-box");
}
#[test]
fn set_ui_preferences_accepts_read_as_a_default_view() {
let ok = UiPrefsJson {
renderer: "xterm".to_string(),
theme: "nord".to_string(),
default_view: "read".to_string(),
osc133_hint_dismissed: false,
listen_voice: String::new(),
listen_rate: 1.0,
listen_pitch: 1.0,
selected_node: String::new(),
}
.validate()
.expect("read must be accepted as a default_view");
assert_eq!(ok.default_view, "read");
let round_tripped: UiPrefsJson = ok.into();
assert_eq!(round_tripped.default_view, "read");
}
#[test]
fn get_ui_preferences_normalizes_a_corrupt_row_instead_of_erroring() {
let corrupt = db::UiPreferences {
renderer: "not-a-real-renderer".to_string(),
theme: "nord".to_string(),
default_view: "not-a-real-view".to_string(),
osc133_hint_dismissed: false,
listen_voice: String::new(),
listen_rate: 500.0,
listen_pitch: -500.0,
selected_node: String::new(),
};
let json: UiPrefsJson = corrupt.into();
assert_eq!(json.renderer, "xterm");
assert_eq!(json.default_view, "xterm");
assert_eq!(json.listen_rate, 2.0);
assert_eq!(json.listen_pitch, 0.5);
}
#[tokio::test]
async fn serve_static_is_no_store() {
use axum::http::header;
let resp = serve_static(Path("style.css".to_string())).await;
assert_eq!(resp.status(), StatusCode::OK);
let cc = resp
.headers()
.get(header::CACHE_CONTROL)
.and_then(|v| v.to_str().ok())
.unwrap_or("");
assert!(
cc.contains("no-store"),
"static assets must be no-store, got Cache-Control: {cc:?}"
);
assert!(
!cc.contains("immutable"),
"static assets must never be immutable, got Cache-Control: {cc:?}"
);
}
#[test]
fn twa_declares_record_audio_when_web_uses_getusermedia() {
use std::fs;
let static_dir = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("web/static");
let mut uses_mic = false;
let mut stack = vec![static_dir];
while let Some(dir) = stack.pop() {
let Ok(entries) = fs::read_dir(&dir) else {
continue;
};
for entry in entries.flatten() {
let path = entry.path();
if path.is_dir() {
stack.push(path);
} else if path.extension().and_then(|e| e.to_str()) == Some("js") {
if let Ok(src) = fs::read_to_string(&path) {
if src.contains("getUserMedia") {
uses_mic = true;
}
}
}
}
}
if uses_mic {
let init_js = fs::read_to_string(
std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("twa/init.js"),
)
.expect("twa/init.js must exist");
assert!(
init_js.contains("android.permission.RECORD_AUDIO"),
"web/static uses getUserMedia but twa/init.js does not inject \
android.permission.RECORD_AUDIO — the TWA mic prompt will be \
denied at the OS layer"
);
}
}
#[test]
fn no_client_storage_in_web_sources() {
use std::fs;
const FORBIDDEN: &[&str] = &["localStorage", "sessionStorage"];
let root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("web");
let mut offenders: Vec<String> = Vec::new();
let mut stack: Vec<std::path::PathBuf> = vec![root];
while let Some(dir) = stack.pop() {
let Ok(entries) = fs::read_dir(&dir) else {
continue;
};
for entry in entries.flatten() {
let path = entry.path();
if path.is_dir() {
let name = path.file_name().and_then(|n| n.to_str()).unwrap_or("");
if name == "node_modules"
|| path.ends_with("static/vendor")
|| path.ends_with("static/spa")
{
continue;
}
stack.push(path);
continue;
}
let ext = path.extension().and_then(|e| e.to_str());
if !matches!(ext, Some("js") | Some("jsx") | Some("mjs") | Some("cjs")) {
continue;
}
let Ok(src) = fs::read_to_string(&path) else {
continue;
};
for line in src.lines() {
for token in FORBIDDEN {
if line.contains(token) {
offenders.push(format!("{}: {}", path.display(), line.trim()));
}
}
}
}
}
assert!(
offenders.is_empty(),
"web sources must not use localStorage/sessionStorage — mobux keeps \
no client-side persistent storage (durable state is server-synced \
via /api/settings/preferences, the rest is in-memory):\n{}",
offenders.join("\n")
);
}
#[test]
fn base64url_round_trip_p256_point() {
let bytes: Vec<u8> = (0..65u8).collect();
let encoded = BASE64URL.encode(&bytes);
assert!(
!encoded.contains('='),
"URL_SAFE_NO_PAD must not emit padding"
);
assert!(
!encoded.contains('+') && !encoded.contains('/'),
"URL_SAFE_NO_PAD must use URL-safe alphabet"
);
let decoded = BASE64URL.decode(encoded).expect("round-trip decode");
assert_eq!(decoded, bytes);
}
#[test]
fn base64url_decode_rejects_bad_input() {
assert!(BASE64URL.decode("AAAA=").is_err());
assert!(BASE64URL.decode("AA+/").is_err());
}
#[test]
fn decode_b64url_helper_returns_400_on_garbage() {
let err = decode_b64url("p256dh", "!!not-valid!!").expect_err("must error");
assert_eq!(err.status, StatusCode::BAD_REQUEST);
assert!(
err.message.contains("p256dh"),
"error mentions field name: {}",
err.message
);
}
#[test]
fn session_name_regex_rejects_tmux_unsafe_chars() {
let re = Regex::new(r"^[a-zA-Z0-9_-]+$").unwrap();
for ok in ["foo", "my_app", "build-2", "ABC", "0"] {
assert!(re.is_match(ok), "should accept {ok:?}");
}
for bad in ["my.app", "a:b", "with space", ""] {
assert!(!re.is_match(bad), "should reject {bad:?}");
}
}
fn test_state(dev_mode: bool) -> (AppState, tempfile::TempDir) {
let dir = tempfile::tempdir().expect("tempdir");
let db = Arc::new(db::Db::open(&dir.path().join("mobux.db")).expect("open db"));
let mut settings = config::Config::default();
settings.tls.enabled = false;
settings.app.dev = dev_mode;
let session_name_re = Arc::new(Regex::new(r"^[a-zA-Z0-9_-]+$").unwrap());
let settings = Arc::new(settings);
let config_file = Arc::new(config_file::ConfigFile::new(
dir.path().join(config::CONFIG_FILE_NAME),
));
let mcp = Arc::new(mcp_settings::McpServer::new(
mcp::Context::new(session_name_re.clone(), db.clone(), &settings),
settings.clone(),
config_file.clone(),
None,
));
let pages = Arc::new(pages_settings::PagesSettings::new(
Arc::new(files::FileRoots::default()),
Arc::new(proxy::ProxyTargets::from_config(&settings, SESSION_COOKIE_NAME).unwrap()),
config_file,
pages_settings::Managed::default(),
));
let state = AppState {
session_name_re,
auth: None,
cache_bust: "test".to_string(),
db,
internal_token: Arc::new("test-token".to_string()),
config: settings,
data_dir: dir.path().to_path_buf(),
config_dir: dir.path().to_path_buf(),
update: update::UpdateState::new(String::new()),
build_hash: "test".to_string(),
twa_build: BackgroundJobState::idle(),
session_history: Arc::new(session_history::SessionHistoryStore::new(dir.path())),
mcp,
pages,
};
(state, dir)
}
fn listener_app(state: AppState, via_access: bool) -> Router {
let app = Router::new()
.route("/api/build-info", get(api_build_info))
.route("/api/upload", post(api_upload))
.with_state(state);
if !via_access {
return app;
}
app.layer(middleware::from_fn(
|mut request: Request<axum::body::Body>, next: middleware::Next| async move {
request.extensions_mut().insert(access::ViaAccess);
next.run(request).await
},
))
}
async fn build_info(via_access: bool) -> serde_json::Value {
use tower::ServiceExt;
let (state, _dir) = test_state(false);
let response = listener_app(state, via_access)
.oneshot(
Request::builder()
.uri("/api/build-info")
.body(axum::body::Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
let bytes = axum::body::to_bytes(response.into_body(), usize::MAX)
.await
.unwrap();
serde_json::from_slice(&bytes).unwrap()
}
#[tokio::test]
async fn build_info_reports_the_access_listener_and_its_upload_limit() {
let info = build_info(true).await;
assert_eq!(info["via_access"], json!(true));
assert_eq!(info["upload_limit_bytes"], json!(100 * 1024 * 1024));
}
#[tokio::test]
async fn build_info_on_the_main_listener_keeps_the_existing_upload_limit() {
let info = build_info(false).await;
assert_eq!(info["via_access"], json!(false));
assert_eq!(info["upload_limit_bytes"], json!(200 * 1024 * 1024));
}
async fn put_mcp(state: AppState, body: &str) -> (StatusCode, String) {
use tower::ServiceExt;
let response = Router::new()
.route(
"/api/settings/mcp",
get(api_get_mcp_settings).put(api_set_mcp_settings),
)
.with_state(state)
.oneshot(
Request::builder()
.method("PUT")
.uri("/api/settings/mcp")
.header("content-type", "application/json")
.body(axum::body::Body::from(body.to_string()))
.unwrap(),
)
.await
.unwrap();
let status = response.status();
let bytes = axum::body::to_bytes(response.into_body(), usize::MAX)
.await
.unwrap();
(status, String::from_utf8(bytes.to_vec()).unwrap())
}
async fn put_pages(state: AppState, body: &str) -> (StatusCode, String) {
use tower::ServiceExt;
let response = Router::new()
.route(
"/api/settings/pages",
get(api_get_pages_settings).put(api_set_pages_settings),
)
.with_state(state)
.oneshot(
Request::builder()
.method("PUT")
.uri("/api/settings/pages")
.header("content-type", "application/json")
.body(axum::body::Body::from(body.to_string()))
.unwrap(),
)
.await
.unwrap();
let status = response.status();
let bytes = axum::body::to_bytes(response.into_body(), usize::MAX)
.await
.unwrap();
(status, String::from_utf8(bytes.to_vec()).unwrap())
}
#[tokio::test]
async fn concurrent_pages_and_mcp_saves_both_land_in_the_file() {
let (state, dir) = test_state(false);
let path = dir.path().join(config::CONFIG_FILE_NAME);
for n in 0..8 {
let root = dir.path().join(format!("site{n}"));
std::fs::create_dir_all(&root).unwrap();
let port = std::net::TcpListener::bind("127.0.0.1:0")
.unwrap()
.local_addr()
.unwrap()
.port();
let pages = json!({"files": [{"name": format!("site{n}"), "path": root}]}).to_string();
let mcp = json!({"port": port}).to_string();
let ((mcp_status, mcp_body), (pages_status, pages_body)) = tokio::join!(
put_mcp(state.clone(), &mcp),
put_pages(state.clone(), &pages),
);
assert_eq!(mcp_status, StatusCode::OK, "{mcp_body}");
assert_eq!(pages_status, StatusCode::OK, "{pages_body}");
let written = config::load_from(&path).unwrap();
assert_eq!(written.mcp.port, port, "round {n}");
assert!(
written.files.roots.contains_key(&format!("site{n}")),
"round {n}"
);
}
put_mcp(state, r#"{"port": 0}"#).await;
}
#[tokio::test]
async fn an_unconfigured_page_is_404_and_the_spa_still_serves() {
use tower::ServiceExt;
let (mut state, _dir) = test_state(false);
state.auth = Some(AuthConfig {
user: "me".to_string(),
pass: "12345".to_string(),
session_cookie_name: SESSION_COOKIE_NAME.to_string(),
session_cookie_value: "cookie".to_string(),
});
let app =
app_routes(state.clone()).layer(middleware::from_fn_with_state(state, auth_middleware));
let get = |uri: &str| {
let request = Request::builder()
.uri(uri)
.header(axum::http::header::COOKIE, "mobux_session=cookie")
.body(axum::body::Body::empty())
.unwrap();
let app = app.clone();
async move {
let response = app.oneshot(request).await.unwrap();
let status = response.status();
let bytes = axum::body::to_bytes(response.into_body(), usize::MAX)
.await
.unwrap();
(status, bytes.to_vec())
}
};
for (uri, body) in [
("/files/x/", "no file root by that name; 0 configured\n"),
("/files/x", "no file root by that name; 0 configured\n"),
("/proxy/x/", "no proxy target by that name; 0 configured\n"),
] {
let (status, bytes) = get(uri).await;
assert_eq!(status, StatusCode::NOT_FOUND, "{uri}");
assert_eq!(String::from_utf8_lossy(&bytes), body, "{uri}");
}
let shell = serve_spa_index().await;
let shell_status = shell.status();
let shell_body = axum::body::to_bytes(shell.into_body(), usize::MAX)
.await
.unwrap()
.to_vec();
for uri in ["/app", "/somewhere/else"] {
let (status, bytes) = get(uri).await;
assert_eq!(status, shell_status, "{uri}");
assert_eq!(bytes, shell_body, "{uri}");
}
}
#[tokio::test]
async fn a_bad_page_answers_400_with_the_rule_and_writes_nothing() {
let (state, dir) = test_state(false);
for (body, rule) in [
(
r#"{"files": [{"name": "site", "path": "relative"}]}"#,
"must be an absolute path",
),
(
r#"{"proxies": [{"name": "up", "port": 70000}]}"#,
"port from 1 to 65535",
),
] {
let (status, text) = put_pages(state.clone(), body).await;
assert_eq!(status, StatusCode::BAD_REQUEST, "{body}: {text}");
assert!(text.contains(rule), "{body}: {text}");
}
assert!(!dir.path().join(config::CONFIG_FILE_NAME).exists());
}
#[tokio::test]
async fn a_page_section_the_environment_sets_answers_409() {
let (state, dir) = test_state(false);
let state = AppState {
pages: Arc::new(pages_settings::PagesSettings::new(
state.pages.files().clone(),
state.pages.proxies().clone(),
Arc::new(config_file::ConfigFile::new(
dir.path().join(config::CONFIG_FILE_NAME),
)),
pages_settings::Managed {
files: None,
proxies: Some(mcp_settings::ManagedBy::Env),
},
)),
..state
};
let (status, text) =
put_pages(state, r#"{"proxies": [{"name": "up", "port": 8000}]}"#).await;
assert_eq!(status, StatusCode::CONFLICT);
assert!(text.contains("MOBUX_PROXY"), "{text}");
}
#[tokio::test]
async fn an_mcp_port_in_use_answers_409_with_the_reason_and_writes_nothing() {
let (state, dir) = test_state(false);
let taken = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let port = taken.local_addr().unwrap().port();
let (status, body) = put_mcp(state, &format!("{{\"port\": {port}}}")).await;
assert_eq!(status, StatusCode::CONFLICT);
assert!(
body.contains(&format!("cannot listen on 127.0.0.1:{port}")),
"{body}"
);
assert!(!dir.path().join(config::CONFIG_FILE_NAME).exists());
}
#[tokio::test]
async fn an_mcp_port_below_1024_answers_400() {
let (state, _dir) = test_state(false);
let (status, body) = put_mcp(state, r#"{"port": 80}"#).await;
assert_eq!(status, StatusCode::BAD_REQUEST);
assert!(body.contains("from 1024 to 65535"), "{body}");
}
#[tokio::test]
async fn an_upload_over_the_access_limit_is_refused_with_a_one_line_413() {
use tower::ServiceExt;
let (state, _dir) = test_state(false);
let response = listener_app(state, true)
.oneshot(
Request::builder()
.method("POST")
.uri("/api/upload")
.header("content-type", "multipart/form-data; boundary=x")
.header("content-length", (150 * 1000 * 1000).to_string())
.body(axum::body::Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::PAYLOAD_TOO_LARGE);
let bytes = axum::body::to_bytes(response.into_body(), usize::MAX)
.await
.unwrap();
assert_eq!(
String::from_utf8(bytes.to_vec()).unwrap(),
"upload refused: the file is larger than the 100 MB limit on this connection"
);
}
#[tokio::test]
async fn the_same_body_on_the_main_listener_is_not_refused_for_size() {
use tower::ServiceExt;
let (state, _dir) = test_state(false);
let response = listener_app(state, false)
.oneshot(
Request::builder()
.method("POST")
.uri("/api/upload")
.header("content-type", "multipart/form-data; boundary=x")
.header("content-length", (150 * 1000 * 1000).to_string())
.body(axum::body::Body::from("--x--\r\n"))
.unwrap(),
)
.await
.unwrap();
assert_ne!(response.status(), StatusCode::PAYLOAD_TOO_LARGE);
}
#[tokio::test]
async fn served_files_sit_behind_the_auth_layer() {
use tower::ServiceExt;
let (mut state, dir) = test_state(false);
state.auth = Some(AuthConfig {
user: "me".to_string(),
pass: "12345".to_string(),
session_cookie_name: "mobux_session".to_string(),
session_cookie_value: "cookie".to_string(),
});
let root = dir.path().join("site");
std::fs::create_dir_all(&root).unwrap();
std::fs::write(root.join("index.html"), "home").unwrap();
let mut files = config::FilesConfig::default();
files
.roots
.insert("site".to_string(), root.display().to_string());
let roots = Arc::new(files::FileRoots::from_config(&files).unwrap());
let app: Router = Router::new()
.merge(files::router(roots))
.with_state(state.clone())
.layer(middleware::from_fn_with_state(state, auth_middleware));
let request = |cookie: Option<&str>| {
let mut builder = Request::builder().uri("/files/site/");
if let Some(cookie) = cookie {
builder = builder.header(axum::http::header::COOKIE, cookie);
}
builder.body(axum::body::Body::empty()).unwrap()
};
let response = app.clone().oneshot(request(None)).await.unwrap();
assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
let response = app
.oneshot(request(Some("mobux_session=cookie")))
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
}
#[tokio::test]
async fn proxied_ports_sit_behind_the_auth_layer() {
use tower::ServiceExt;
let upstream = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = upstream.local_addr().unwrap().port();
tokio::spawn(async move {
axum::serve(upstream, Router::new().route("/", get(|| async { "up" })))
.await
.unwrap()
});
let (mut state, _dir) = test_state(false);
state.auth = Some(AuthConfig {
user: "me".to_string(),
pass: "12345".to_string(),
session_cookie_name: SESSION_COOKIE_NAME.to_string(),
session_cookie_value: "cookie".to_string(),
});
let mut settings = config::Config::default();
settings.proxy.targets.insert("up".to_string(), port);
let targets =
Arc::new(proxy::ProxyTargets::from_config(&settings, SESSION_COOKIE_NAME).unwrap());
let app: Router = Router::new()
.merge(proxy::router(targets))
.with_state(state.clone())
.layer(middleware::from_fn_with_state(state, auth_middleware));
let request = |cookie: Option<&str>| {
let mut builder = Request::builder().uri("/proxy/up/");
if let Some(cookie) = cookie {
builder = builder.header(axum::http::header::COOKIE, cookie);
}
builder.body(axum::body::Body::empty()).unwrap()
};
let response = app.clone().oneshot(request(None)).await.unwrap();
assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
let response = app
.oneshot(request(Some("mobux_session=cookie")))
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
}
#[tokio::test]
async fn telemetry_endpoint_active_without_dev_mode() {
let status = api_telemetry("hello".to_string()).await;
assert_eq!(status, StatusCode::NO_CONTENT);
}
fn cookie_config(tls: bool, behind_tls_proxy: bool, base_path: &str) -> config::Config {
let mut settings = config::Config::default();
settings.tls.enabled = tls;
settings.server.behind_tls_proxy = behind_tls_proxy;
settings.server.base_path = base_path.to_string();
settings
}
fn cookie(tls: bool, behind_tls_proxy: bool, base_path: &str) -> String {
build_session_cookie(
"mobux_session",
"abc123",
&cookie_config(tls, behind_tls_proxy, base_path),
)
}
#[test]
fn session_cookie_defaults_to_no_secure() {
let settings = config::Config::default();
assert!(!settings.tls.enabled, "TLS is off by default");
let cookie = build_session_cookie("mobux_session", "abc123", &settings);
assert!(
!cookie.contains("Secure"),
"the default bind is plain HTTP, so no Secure: {cookie}"
);
assert!(
cookie.contains("Path=/;"),
"no base path means the site root: {cookie}"
);
}
#[test]
fn session_cookie_is_secure_whenever_the_browser_is_on_https() {
assert!(!cookie(false, false, "").contains("Secure"));
assert!(cookie(true, false, "").contains("Secure"));
assert!(cookie(false, true, "").contains("Secure"));
assert!(cookie(true, true, "").contains("Secure"));
}
#[test]
fn session_cookie_always_carries_httponly_and_samesite() {
for cookie in [cookie(false, false, ""), cookie(true, true, "/mobux")] {
assert!(cookie.contains("HttpOnly"), "{cookie}");
assert!(cookie.contains("SameSite=Lax"), "{cookie}");
assert!(cookie.contains("Max-Age=2592000"), "{cookie}");
}
}
#[test]
fn session_cookie_path_follows_the_base_path() {
assert!(cookie(false, false, "").contains("Path=/;"));
assert!(cookie(false, false, "/mobux").contains("Path=/mobux;"));
assert!(cookie(false, false, "/user/host/8080").contains("Path=/user/host/8080;"));
}
#[test]
fn cookie_path_normalises_the_slashes() {
assert_eq!(cookie_path(""), "/");
assert_eq!(cookie_path("/"), "/");
assert_eq!(cookie_path(" "), "/");
assert_eq!(cookie_path("/mobux/"), "/mobux");
assert_eq!(cookie_path("mobux"), "/mobux");
assert_eq!(cookie_path("/user/host/8080/"), "/user/host/8080");
}
fn location_of(response: Response) -> String {
response
.headers()
.get(axum::http::header::LOCATION)
.expect("a Location header")
.to_str()
.expect("an ASCII Location")
.to_string()
}
#[tokio::test]
async fn the_root_redirect_is_a_sibling_of_the_request() {
let response = root_redirect().await.into_response();
assert_eq!(response.status(), StatusCode::TEMPORARY_REDIRECT);
assert_eq!(location_of(response), "app");
}
#[tokio::test]
async fn the_settings_redirect_is_a_sibling_of_the_request() {
let response = settings_page().await.into_response();
assert_eq!(response.status(), StatusCode::TEMPORARY_REDIRECT);
assert_eq!(location_of(response), "app#/settings");
}
#[tokio::test]
async fn the_install_redirect_is_a_sibling_of_the_request() {
let response = install_page().await.into_response();
assert_eq!(response.status(), StatusCode::TEMPORARY_REDIRECT);
assert_eq!(location_of(response), "app#/install");
}
#[tokio::test]
async fn a_session_link_climbs_one_level_to_reach_the_app() {
let (state, _dir) = test_state(false);
let response = terminal_page(
State(state),
Path("no-such-session".to_string()),
RawQuery(None),
)
.await
.expect("a redirect")
.into_response();
assert_eq!(response.status(), StatusCode::TEMPORARY_REDIRECT);
assert_eq!(location_of(response), "../app#/");
}
#[tokio::test]
async fn a_session_link_keeps_its_window_query_ahead_of_the_fragment() {
let (state, _dir) = test_state(false);
let response = terminal_page(
State(state),
Path("no-such-session".to_string()),
RawQuery(Some("w=3".to_string())),
)
.await
.expect("a redirect")
.into_response();
assert_eq!(location_of(response), "../app?w=3#/");
}
#[test]
fn an_app_deep_link_climbs_back_to_the_single_document() {
assert_eq!(relative_app_target("foo"), "../app");
assert_eq!(relative_app_target("foo/bar"), "../../app");
assert_eq!(relative_app_target("foo/"), "../../app");
assert_eq!(relative_app_target("a/b/c"), "../../../app");
}
#[tokio::test]
async fn an_app_deep_link_redirects_rather_than_serving_the_shell() {
let response = spa_deep_link_redirect(Path("settings".to_string()))
.await
.into_response();
assert_eq!(response.status(), StatusCode::TEMPORARY_REDIRECT);
assert_eq!(location_of(response), "../app");
}
#[tokio::test]
async fn the_router_fallback_answers_with_a_document_never_a_redirect() {
let response = serve_spa_index().await;
assert!(
!response.status().is_redirection(),
"the fallback redirected: {}",
response.status()
);
assert!(response
.headers()
.get(axum::http::header::LOCATION)
.is_none());
}
#[test]
fn clear_text_warning_fires_only_with_auth_on_and_tls_off() {
assert!(clear_text_auth_warning(true, false, false).is_some());
assert!(clear_text_auth_warning(true, true, false).is_none());
assert!(clear_text_auth_warning(false, false, false).is_none());
assert!(clear_text_auth_warning(false, true, false).is_none());
assert!(clear_text_auth_warning(true, false, true).is_none());
}
#[test]
fn clear_text_warning_says_how_to_turn_https_on() {
let warning = clear_text_auth_warning(true, false, false).expect("a warning");
for hint in [
"--tls",
"MOBUX_TLS=1",
r#"{"tls": {"enabled": true}}"#,
"--behind-tls-proxy",
] {
assert!(
warning.contains(hint),
"warning is missing {hint}: {warning}"
);
}
}
#[tokio::test]
async fn stt_install_starts_the_local_model_fetch() {
let (state, _dir) = test_state(false);
let resp = api_stt_install(State(state)).await.unwrap().into_response();
let expected = if local_stt::ENABLED {
StatusCode::ACCEPTED
} else {
StatusCode::PRECONDITION_FAILED
};
assert_eq!(resp.status(), expected);
}
#[test]
fn a_cpu_that_cannot_run_the_engine_reports_as_unavailable() {
let body = local_status(local_stt::Phase::UnsupportedCpu(
"no FEAT_FP16 here".to_string(),
))
.into_json();
assert_eq!(body["state"], "unsupported");
assert_eq!(body["engine_available"], false);
assert_eq!(body["message"], "no FEAT_FP16 here");
let install = install_status_json(&local_stt::Phase::UnsupportedCpu(
"no FEAT_FP16 here".to_string(),
));
assert_eq!(install["phase"], "failed");
assert_eq!(install["error"], "no FEAT_FP16 here");
}
#[tokio::test]
async fn stt_install_refuses_a_remote_provider() {
let (state, _dir) = test_state(false);
state
.db
.set_stt_active_kind("openai")
.expect("switch to a configured endpoint");
let resp = api_stt_install(State(state)).await.unwrap().into_response();
assert_eq!(resp.status(), StatusCode::PRECONDITION_FAILED);
}
fn apk_build_request(host: &str) -> HeaderMap {
let mut headers = HeaderMap::new();
headers.insert(
axum::http::header::HOST,
HeaderValue::from_str(host).unwrap(),
);
headers
}
#[tokio::test]
async fn apk_build_returns_409_while_a_build_is_running() {
let (state, _dir) = test_state(false);
{
let mut guard = state.twa_build.lock().await;
guard.phase = InstallPhase::Running;
guard.output_tail = vec!["Building the Android package".to_string()];
}
let resp =
api_install_apk_build(State(state.clone()), apk_build_request("box.example.com"))
.await
.expect("handler must not error");
assert_eq!(resp.status(), StatusCode::CONFLICT);
let guard = state.twa_build.lock().await;
assert_eq!(guard.phase, InstallPhase::Running);
assert_eq!(guard.output_tail.len(), 1);
}
#[tokio::test]
async fn apk_build_rejects_a_host_it_cannot_turn_into_a_domain() {
let (state, _dir) = test_state(false);
let resp = api_install_apk_build(State(state.clone()), HeaderMap::new())
.await
.expect("handler must not error");
assert_eq!(resp.status(), StatusCode::BAD_REQUEST);
assert_eq!(
state.twa_build.lock().await.phase,
InstallPhase::Idle,
"a rejected request must not move the job out of idle"
);
}
#[tokio::test]
async fn apk_status_reports_phase_output_and_availability() {
let (state, dir) = test_state(false);
let resp = api_install_apk_status(State(state.clone()), apk_build_request("box:5151"))
.await
.unwrap();
assert_eq!(resp.0["phase"], "idle");
assert_eq!(resp.0["output"], json!([]));
assert_eq!(resp.0["error"], serde_json::Value::Null);
assert_eq!(resp.0["apk_available"], false);
assert_eq!(resp.0["domain"], "box:5151");
{
let mut guard = state.twa_build.lock().await;
guard.phase = InstallPhase::Failed("exit 1: bubblewrap not found".to_string());
guard.output_tail = vec!["Building signed APK".to_string()];
}
let apk = twa::apk_path(dir.path());
std::fs::create_dir_all(apk.parent().unwrap()).unwrap();
std::fs::write(&apk, b"apk").unwrap();
let resp = api_install_apk_status(State(state), apk_build_request("box:5151"))
.await
.unwrap();
assert_eq!(resp.0["phase"], "failed");
assert_eq!(resp.0["error"], "exit 1: bubblewrap not found");
assert_eq!(resp.0["output"], json!(["Building signed APK"]));
assert_eq!(resp.0["apk_available"], true);
}
fn local_status(phase: local_stt::Phase) -> SttStatus {
SttStatus {
kind: transcribe::LOCAL_KIND.to_string(),
url: String::new(),
model: local_stt::DEFAULT_MODEL.to_string(),
local: Some(phase),
remote_reachable: false,
}
}
#[test]
fn stt_status_reports_a_download_in_flight_as_progress() {
let status = local_status(local_stt::Phase::Downloading {
file: "model.safetensors".to_string(),
downloaded: 40_000_000,
total: 160_000_000,
});
assert_eq!(status.state(), "warming");
let body = status.into_json();
assert_eq!(body["state"], "warming");
assert_eq!(body["reachable"], false);
assert_eq!(body["progress"]["downloaded"], 40_000_000u64);
assert_eq!(body["progress"]["total"], 160_000_000u64);
assert_eq!(body["progress"]["file"], "model.safetensors");
assert!(body["message"].as_str().unwrap().contains("25%"));
}
#[test]
fn stt_status_separates_ready_loading_and_never_downloaded() {
assert_eq!(local_status(local_stt::Phase::Ready).state(), "ready");
assert_eq!(local_status(local_stt::Phase::Loading).state(), "warming");
assert_eq!(
local_status(local_stt::Phase::NotDownloaded).state(),
"not_installed"
);
assert_eq!(
local_status(local_stt::Phase::Failed("disk full".to_string())).state(),
"failed"
);
}
#[test]
fn stt_status_names_a_build_without_the_in_process_engine() {
let status = local_status(local_stt::Phase::Disabled);
assert_eq!(status.state(), "unsupported");
let body = status.into_json();
assert_eq!(body["engine_available"], false);
assert!(
body["message"]
.as_str()
.unwrap()
.contains("--features local-stt"),
"the message must carry the fix, not just the fault: {}",
body["message"]
);
}
#[test]
fn stt_status_surfaces_a_failed_download() {
let body = local_status(local_stt::Phase::Failed("404 Not Found".to_string())).into_json();
assert_eq!(body["state"], "failed");
assert_eq!(body["error"], "404 Not Found");
}
#[test]
fn stt_status_keeps_calling_a_remote_provider_unreachable() {
let status = SttStatus {
kind: "openai".to_string(),
url: "https://api.openai.com:443/v1/audio/transcriptions".to_string(),
model: "whisper-1".to_string(),
local: None,
remote_reachable: false,
};
assert_eq!(status.state(), "unreachable");
let body = status.into_json();
assert_eq!(body["state"], "unreachable");
assert!(body["progress"].is_null());
}
#[test]
fn stt_status_calls_a_responding_remote_provider_ready() {
let status = SttStatus {
kind: "network".to_string(),
url: "http://lab:8081/v1/audio/transcriptions".to_string(),
model: "Systran/faster-whisper-base.en".to_string(),
local: None,
remote_reachable: true,
};
assert_eq!(status.state(), "ready");
assert_eq!(status.into_json()["reachable"], true);
}
#[test]
fn install_status_maps_every_engine_phase_onto_the_poller_shape() {
assert_eq!(
install_status_json(&local_stt::Phase::Ready)["phase"],
"success"
);
assert_eq!(
install_status_json(&local_stt::Phase::Loading)["phase"],
"running"
);
assert_eq!(
install_status_json(&local_stt::Phase::NotDownloaded)["phase"],
"idle"
);
let failed = install_status_json(&local_stt::Phase::Failed("no space left".to_string()));
assert_eq!(failed["phase"], "failed");
assert_eq!(failed["error"], "no space left");
assert_eq!(
install_status_json(&local_stt::Phase::Loading)["output"]
.as_array()
.unwrap()
.len(),
1
);
}
#[tokio::test]
async fn stt_status_on_a_fresh_install_reports_the_local_engine() {
let (state, _dir) = test_state(false);
let resp = api_stt_status(State(state)).await.unwrap();
assert_eq!(resp.0["kind"], transcribe::LOCAL_KIND);
assert_eq!(resp.0["model"], local_stt::DEFAULT_MODEL);
assert_eq!(resp.0["url"], "");
assert_eq!(resp.0["engine_available"], local_stt::ENABLED);
assert_eq!(resp.0["reachable"], false);
}
#[tokio::test]
async fn resolve_node_target_rejects_a_target_starting_with_a_dash() {
let (state, _dir) = test_state(false);
state
.db
.replace_nodes(vec![("evil".to_string(), "-oProxyCommand=pwn".to_string())])
.expect("seed node");
let err = resolve_node_target(&state, Some("evil"))
.await
.expect_err("a leading-dash target must be rejected");
assert_eq!(err.status, StatusCode::BAD_REQUEST);
assert!(
err.message.contains("invalid target"),
"error names the problem: {}",
err.message
);
}
#[tokio::test]
async fn resolve_node_target_accepts_a_normal_target() {
let (state, _dir) = test_state(false);
state
.db
.replace_nodes(vec![("devbox".to_string(), "user@devbox".to_string())])
.expect("seed node");
let target = resolve_node_target(&state, Some("devbox"))
.await
.expect("a well-formed target resolves");
assert_eq!(target.as_deref(), Some("user@devbox"));
}
#[test]
fn sanitize_upload_filename_keeps_unicode_letters_strips_shell_metacharacters() {
assert_eq!(sanitize_upload_filename("naïve—file.txt"), "naïve_file.txt");
let adversarial = "$(touch PWNED);`x`\n'.txt";
let safe = sanitize_upload_filename(adversarial);
for dangerous in ['$', '(', ')', ';', '`', '\n', '\''] {
assert!(
!safe.contains(dangerous),
"{safe:?} must not contain {dangerous:?}"
);
}
assert!(safe.ends_with(".txt"));
assert!(safe.contains("touch"));
assert!(safe.contains("PWNED"));
}
#[tokio::test]
async fn stt_models_returns_fallback_when_no_config() {
let (state, _dir) = test_state(false);
let q = SttModelsQuery {
kind: None,
host: None,
port: None,
};
let result = api_stt_models(State(state), Query(q)).await;
let Json(val) = result.expect("handler should not error");
let models = val["models"].as_array().expect("models array");
assert!(!models.is_empty(), "fallback models must not be empty");
}
#[tokio::test]
async fn stt_models_answers_for_the_kind_that_was_asked_about() {
let (state, _dir) = test_state(false);
state
.db
.set_stt_active_kind("network")
.expect("a self-hosted endpoint is active");
let Json(val) = api_stt_models(
State(state.clone()),
Query(SttModelsQuery {
kind: Some(transcribe::LOCAL_KIND.to_string()),
host: None,
port: None,
}),
)
.await
.expect("handler should not error");
let models: Vec<String> = val["models"]
.as_array()
.expect("models array")
.iter()
.filter_map(|v| v.as_str().map(String::from))
.collect();
assert_eq!(models, local_stt::model_ids());
state
.db
.set_stt_active_kind(transcribe::LOCAL_KIND)
.expect("switch the active kind");
let Json(val) = api_stt_models(
State(state),
Query(SttModelsQuery {
kind: Some("network".to_string()),
host: None,
port: None,
}),
)
.await
.expect("handler should not error");
assert!(val["models"]
.as_array()
.expect("models array")
.iter()
.any(|m| m.as_str().is_some_and(|m| m.starts_with("Systran/"))));
}
#[tokio::test]
async fn saving_a_local_model_the_engine_cannot_run_stores_the_one_it_will_load() {
let (state, _dir) = test_state(false);
api_set_stt_config(
State(state.clone()),
Json(SttConfigPutJson {
kind: transcribe::LOCAL_KIND.to_string(),
host: String::new(),
port: String::new(),
model: "Systran/faster-whisper-small".to_string(),
api_key: None,
}),
)
.await
.expect("saving must not error");
assert_eq!(
state
.db
.stt_provider(transcribe::LOCAL_KIND)
.expect("read local")
.expect("local row")
.model,
local_stt::DEFAULT_MODEL
);
}
#[tokio::test]
async fn saving_a_local_model_from_the_catalog_keeps_it() {
let (state, _dir) = test_state(false);
api_set_stt_config(
State(state.clone()),
Json(SttConfigPutJson {
kind: transcribe::LOCAL_KIND.to_string(),
host: String::new(),
port: String::new(),
model: "small.en".to_string(),
api_key: None,
}),
)
.await
.expect("saving must not error");
assert_eq!(
state
.db
.stt_provider(transcribe::LOCAL_KIND)
.expect("read local")
.expect("local row")
.model,
"small.en"
);
}
#[tokio::test]
async fn the_local_model_list_never_offers_a_checkpoint_the_engine_cannot_run() {
let (state, _dir) = test_state(false);
state
.db
.set_stt_provider(db::SttProviderRow {
kind: transcribe::LOCAL_KIND.to_string(),
host: "http://127.0.0.1".to_string(),
port: "5200".to_string(),
model: "Systran/faster-whisper-small".to_string(),
api_key: None,
})
.expect("a row left over from the container-backed provider");
let Json(val) = api_stt_models(
State(state),
Query(SttModelsQuery {
kind: Some(transcribe::LOCAL_KIND.to_string()),
host: Some("http://127.0.0.1".to_string()),
port: Some("5200".to_string()),
}),
)
.await
.expect("handler should not error");
let models: Vec<String> = val["models"]
.as_array()
.expect("models array")
.iter()
.filter_map(|v| v.as_str().map(String::from))
.collect();
assert_eq!(models, local_stt::model_ids());
for model in &models {
assert!(
local_stt::is_known_model(model),
"the local picker offered {model}, which the engine cannot run"
);
}
}
#[tokio::test]
async fn stt_models_returns_openai_fallback_for_openai_kind() {
let (state, _dir) = test_state(false);
let q = SttModelsQuery {
kind: Some("openai".to_string()),
host: Some("https://api.openai.com".to_string()),
port: Some("443".to_string()),
};
let result = api_stt_models(State(state), Query(q)).await;
let Json(val) = result.expect("handler should not error");
let models: Vec<String> = val["models"]
.as_array()
.expect("models array")
.iter()
.filter_map(|v| v.as_str().map(String::from))
.collect();
assert!(!models.is_empty());
for m in &models {
assert!(!m.is_empty(), "model id must not be empty");
}
}
fn session(name: &str) -> tmux::Session {
tmux::Session {
name: name.to_string(),
windows: 1,
attached: 0,
created_unix: 0,
}
}
#[test]
fn resolve_session_location_unique_local_match() {
let local = [session("mobux")];
assert_eq!(
resolve_session_location("mobux", &local, &[]),
Some(SessionLocation::Local)
);
}
#[test]
fn resolve_session_location_unique_node_match() {
let local = [session("other")];
let devbox = [session("mobux")];
let nodes: [(&str, &[tmux::Session]); 1] = [("devbox", &devbox)];
assert_eq!(
resolve_session_location("mobux", &local, &nodes),
Some(SessionLocation::Node("devbox".to_string()))
);
}
#[test]
fn resolve_session_location_no_match_anywhere() {
let local = [session("other")];
let devbox = [session("also-other")];
let nodes: [(&str, &[tmux::Session]); 1] = [("devbox", &devbox)];
assert_eq!(resolve_session_location("mobux", &local, &nodes), None);
}
#[test]
fn resolve_session_location_ambiguous_local_and_node_is_not_a_tiebreak() {
let local = [session("mobux")];
let devbox = [session("mobux")];
let nodes: [(&str, &[tmux::Session]); 1] = [("devbox", &devbox)];
assert_eq!(resolve_session_location("mobux", &local, &nodes), None);
}
#[test]
fn resolve_session_location_ambiguous_across_two_nodes() {
let local: [tmux::Session; 0] = [];
let alpha = [session("mobux")];
let beta = [session("mobux")];
let nodes: [(&str, &[tmux::Session]); 2] = [("alpha", &alpha), ("beta", &beta)];
assert_eq!(resolve_session_location("mobux", &local, &nodes), None);
}
#[test]
fn resolve_session_location_ignores_other_names_on_other_locations() {
let local = [session("mobux")];
let devbox = [session("unrelated")];
let nodes: [(&str, &[tmux::Session]); 1] = [("devbox", &devbox)];
assert_eq!(
resolve_session_location("mobux", &local, &nodes),
Some(SessionLocation::Local)
);
}
#[tokio::test]
async fn probe_timeout_cuts_off_a_future_that_never_resolves() {
let start = std::time::Instant::now();
let result: Option<()> = with_probe_timeout(async {
tokio::time::sleep(SESSION_PROBE_TIMEOUT * 10).await;
Ok(())
})
.await;
let elapsed = start.elapsed();
assert_eq!(result, None, "a hung future must count as absent");
assert!(
elapsed < SESSION_PROBE_TIMEOUT + std::time::Duration::from_millis(500),
"must not wait past its own timeout: took {elapsed:?}",
);
}
#[tokio::test]
async fn terminal_page_errors_when_the_node_inventory_cannot_be_read() {
let (state, dir) = test_state(false);
{
let conn = rusqlite::Connection::open(dir.path().join("mobux.db"))
.expect("raw open of the same sqlite file");
conn.execute("DROP TABLE nodes", [])
.expect("drop nodes table to force list_nodes() to fail");
}
let result =
terminal_page(State(state), Path("whatever".to_string()), RawQuery(None)).await;
match result {
Err(e) => assert_eq!(e.status, StatusCode::INTERNAL_SERVER_ERROR),
Ok(_) => panic!("expected an error when the node inventory can't be read"),
}
}
#[tokio::test]
async fn build_info_reflects_dev_mode() {
let (state, _dir) = test_state(true);
let Json(val) = api_build_info(State(state), None).await;
assert_eq!(val["dev_mode"], true);
let (state, _dir) = test_state(false);
let Json(val) = api_build_info(State(state), None).await;
assert_eq!(val["dev_mode"], false);
}
#[tokio::test]
async fn build_info_lists_served_page_names() {
let (state, _dir) = test_state(false);
let Json(val) = api_build_info(State(state), None).await;
assert_eq!(val["files"], json!([]));
assert_eq!(val["proxies"], json!([]));
let (state, dir) = test_state(false);
let root = dir.path().join("secret-site");
std::fs::create_dir_all(&root).unwrap();
let (status, body) = put_pages(
state.clone(),
&json!({
"files": [{"name": "site", "path": root}],
"proxies": [{"name": "up", "port": 8291}],
})
.to_string(),
)
.await;
assert_eq!(status, StatusCode::OK, "{body}");
let Json(val) = api_build_info(State(state), None).await;
assert_eq!(val["files"], json!(["site"]));
assert_eq!(val["proxies"], json!(["up"]));
let body = val.to_string();
assert!(!body.contains("secret-site"), "leaked a path: {body}");
assert!(!body.contains("8291"), "leaked a port: {body}");
}
#[cfg(unix)]
fn stub_script(dir: &std::path::Path, name: &str, body: &str) -> PathBuf {
use std::os::unix::fs::PermissionsExt;
let path = dir.join(name);
std::fs::write(&path, format!("#!/bin/bash\nset -eu\n{body}\n")).unwrap();
std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o755)).unwrap();
path
}
#[cfg(unix)]
fn stub_command(script: &std::path::Path, log: &std::path::Path) -> tokio::process::Command {
let mut command = tokio::process::Command::new("bash");
command.arg(script).env("STUB_LOG", log);
command
}
#[cfg(unix)]
async fn await_phase(job: &Arc<tokio::sync::Mutex<BackgroundJobState>>, want: &InstallPhase) {
for _ in 0..500 {
if job.lock().await.phase == *want {
return;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
panic!("job never reached {want:?}");
}
#[cfg(unix)]
fn stub_log(path: &std::path::Path) -> String {
std::fs::read_to_string(path).unwrap_or_default()
}
#[cfg(unix)]
#[tokio::test]
async fn a_missing_toolchain_is_installed_before_the_build_runs() {
let dir = tempfile::tempdir().unwrap();
let log = dir.path().join("log");
let release = dir.path().join("release");
let setup = stub_script(
dir.path(),
"setup",
&format!(
"echo setup >> \"$STUB_LOG\"\nwhile [ ! -f '{}' ]; do sleep 0.01; done",
release.display()
),
);
let build = stub_script(dir.path(), "build", "echo build >> \"$STUB_LOG\"");
let job = BackgroundJobState::idle();
let handle = tokio::spawn(run_twa_build_job(
job.clone(),
Some(stub_command(&setup, &log)),
stub_command(&build, &log),
));
await_phase(&job, &InstallPhase::InstallingTools).await;
assert_eq!(stub_log(&log), "setup\n", "the build must wait for setup");
assert!(InstallPhase::InstallingTools.is_active());
assert_eq!(
phase_parts(&InstallPhase::InstallingTools).0,
"installing_tools"
);
std::fs::write(&release, b"go").unwrap();
handle.await.unwrap();
let guard = job.lock().await;
assert_eq!(guard.phase, InstallPhase::Success);
assert_eq!(stub_log(&log), "setup\nbuild\n");
assert!(
guard
.output_tail
.iter()
.any(|l| l.contains("Installing the Android build tools")),
"the tail must name the wait: {:?}",
guard.output_tail
);
}
#[cfg(unix)]
#[tokio::test]
async fn an_installed_toolchain_skips_the_install_phase() {
let dir = tempfile::tempdir().unwrap();
let log = dir.path().join("log");
let build = stub_script(dir.path(), "build", "echo build >> \"$STUB_LOG\"");
let job = BackgroundJobState::idle();
run_twa_build_job(job.clone(), None, stub_command(&build, &log)).await;
let guard = job.lock().await;
assert_eq!(guard.phase, InstallPhase::Success);
assert_eq!(stub_log(&log), "build\n");
assert!(
!guard
.output_tail
.iter()
.any(|l| l.contains("Installing the Android build tools")),
"nothing to install, so nothing to announce: {:?}",
guard.output_tail
);
}
#[test]
fn job_output_drops_the_colour_the_scripts_log_with() {
assert_eq!(
strip_ansi("\u{1b}[1;34m[setup-twa]\u{1b}[0m Installing SDKMAN"),
"[setup-twa] Installing SDKMAN"
);
assert_eq!(
strip_ansi("\u{1b}]8;;https://example.com\u{1b}\\link\u{1b}]8;;\u{1b}\\"),
"link"
);
assert_eq!(
strip_ansi("\u{1b}]0;a title\u{7}plain"),
"plain",
"window-title sequences must not leak into the log pane"
);
assert_eq!(strip_ansi("no escapes here"), "no escapes here");
}
#[tokio::test]
async fn a_host_package_gap_is_reported_through_the_status_endpoint() {
let (state, _dir) = test_state(false);
record_host_package_gap(
&state.twa_build,
twa::host_package_gap_with(Some(twa::PackageManager::Apt), &["unzip", "zip"]),
)
.await;
let mut headers = HeaderMap::new();
headers.insert(
axum::http::header::HOST,
"box.example.com".parse().expect("host header"),
);
let Json(body) = api_install_apk_status(State(state), headers)
.await
.expect("status");
assert_eq!(body["phase"], "failed");
assert_eq!(body["missing_host_packages"], json!(["unzip", "zip"]));
assert_eq!(body["install_command"], "sudo apt-get install -y unzip zip");
assert!(body["error"]
.as_str()
.expect("an error message")
.contains("unzip and zip"));
}
#[tokio::test]
async fn a_press_during_the_toolchain_install_gets_the_running_state() {
let (state, _dir) = test_state(false);
state.twa_build.lock().await.phase = InstallPhase::InstallingTools;
let mut headers = HeaderMap::new();
headers.insert(
axum::http::header::HOST,
"box.example.com".parse().expect("host header"),
);
let response = api_install_apk_build(State(state.clone()), headers)
.await
.expect("build response");
assert_eq!(response.status(), StatusCode::CONFLICT);
assert_eq!(
state.twa_build.lock().await.phase,
InstallPhase::InstallingTools,
"the refused press must not disturb the running job"
);
}
#[cfg(unix)]
#[tokio::test]
async fn a_failed_toolchain_install_stops_before_the_build() {
let dir = tempfile::tempdir().unwrap();
let log = dir.path().join("log");
let setup = stub_script(dir.path(), "setup", "echo 'no unzip' >&2\nexit 3");
let build = stub_script(dir.path(), "build", "echo build >> \"$STUB_LOG\"");
let job = BackgroundJobState::idle();
run_twa_build_job(
job.clone(),
Some(stub_command(&setup, &log)),
stub_command(&build, &log),
)
.await;
let guard = job.lock().await;
match &guard.phase {
InstallPhase::Failed(e) => {
assert!(e.contains("installing the build tools failed"), "{e}");
assert!(e.contains("no unzip"), "{e}");
}
other => panic!("expected a failure, got {other:?}"),
}
assert_eq!(stub_log(&log), "", "the build must not run");
}
}