mod api_agent;
mod api_browse;
mod api_desk;
mod api_docs;
mod api_peer;
mod api_thread;
mod assets;
mod auth;
mod events;
mod lifecycle;
mod peer_link;
#[cfg(test)]
mod tests;
mod ws;
use api_agent::*;
use api_browse::*;
use api_desk::*;
use api_docs::*;
use api_peer::*;
use api_thread::*;
use assets::*;
use auth::*;
use events::*;
use lifecycle::*;
pub use lifecycle::{relaunch, Leaving};
use peer_link::*;
use ws::*;
use crate::browse::Browser;
use crate::capability::constant_eq;
use crate::config::{self, Paths};
use crate::platform;
use crate::receive::{self, Payload};
use crate::render::{self, Renderer};
use crate::store::{Doc, Store};
pub use crate::version::{BUILD_SHA, BUILD_TARGET, VERSION};
use axum::{
body::Body,
extract::{
ws::{Message, WebSocket, WebSocketUpgrade},
Path, Query, State,
},
http::{header, HeaderMap, HeaderValue, StatusCode},
response::{
sse::{Event, KeepAlive, Sse},
Html, IntoResponse, Response,
},
routing::{get, post},
Json, Router,
};
use serde::Deserialize;
use serde_json::json;
use std::borrow::Cow;
use std::convert::Infallible;
use std::path::PathBuf;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::Instant;
use tokio::sync::broadcast;
use tokio_stream::{wrappers::BroadcastStream, StreamExt};
pub struct App {
pub store: Store,
pub renderer: Renderer,
pub browse: Browser,
pub paths: Paths,
pub secrets: crate::secrets::Secrets,
pub token: std::sync::RwLock<String>,
pub window: String,
pub events: broadcast::Sender<String>,
pub shutdown: broadcast::Sender<()>,
pub started: Instant,
built_v: String,
pub ui: Ui,
pub last_focus: std::sync::Mutex<Instant>,
pub windows: AtomicUsize,
pub streams: AtomicUsize,
pub pages: AtomicUsize,
pub online: std::sync::Mutex<std::collections::BTreeMap<String, usize>>,
pub capabilities: crate::capability::Capabilities,
pub panes: Arc<crate::pane::Panes>,
pub asides: crate::aside::Asides,
pub git: crate::git::Cache,
pub exe: Option<Exe>,
pub started_at: i64,
pub restart: std::sync::Mutex<Option<Pending>>,
pub restart_wake: tokio::sync::Notify,
leaving: std::sync::Mutex<Leaving>,
pub update: Option<Arc<crate::update::Updater>>,
pub relaunch_window: std::sync::atomic::AtomicBool,
pub restarting: std::sync::atomic::AtomicBool,
update_sent: std::sync::Mutex<String>,
pub quota: std::sync::Mutex<Option<serde_json::Value>>,
pub peers: api_peer::Peers,
}
#[derive(Clone, Debug)]
pub struct Exe {
pub path: PathBuf,
pub(crate) stamp: (u64, u64, u64, u64),
}
impl Exe {
pub(crate) fn here() -> Option<Exe> {
let path = std::env::current_exe().ok()?;
let path = path.canonicalize().unwrap_or(path);
let stamp = stamp(&path)?;
Some(Exe { path, stamp })
}
pub fn stale(&self) -> bool {
stamp(&self.path) != Some(self.stamp)
}
}
fn stamp(path: &std::path::Path) -> Option<(u64, u64, u64, u64)> {
let m = std::fs::metadata(path).ok()?;
let mtime = m
.modified()
.ok()
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
.map(|d| d.as_nanos() as u64)
.unwrap_or(0);
#[cfg(unix)]
let (dev, ino) = {
use std::os::unix::fs::MetadataExt;
(m.dev(), m.ino())
};
#[cfg(not(unix))]
let (dev, ino) = (0, 0);
Some((dev, ino, m.len(), mtime))
}
impl App {
pub fn online(&self) -> serde_json::Value {
let m = self.online.lock().unwrap_or_else(|e| e.into_inner());
json!(*m)
}
pub fn has_window(&self) -> bool {
self.windows.load(Ordering::Relaxed) > 0
}
pub fn asset_v(&self) -> String {
self.ui.version(&self.built_v)
}
}
type S = State<Arc<App>>;
const QUEUE_MAX: usize = 500;
const QUEUE_BOOT: usize = 24;
fn waiting(app: &App) -> i64 {
app.store.waiting().unwrap_or(0)
}
fn new_app(
paths: &Paths,
store: Store,
token: String,
window: String,
update: Option<Arc<crate::update::Updater>>,
exe: Option<Exe>,
) -> Arc<App> {
let renderer = Renderer::new();
let (tx, _) = broadcast::channel(256);
let (stop_tx, _) = broadcast::channel::<()>(1);
let asset_v = {
let mut h = blake3::Hasher::new();
h.update(INDEX_HTML.as_bytes());
for (_, built_in, _) in ASSETS {
h.update(built_in.as_bytes());
}
h.update(VERSION.as_bytes());
h.update(MERMAID_JS_GZ);
h.finalize().to_hex()[..8].to_string()
};
let panes = crate::pane::Panes::new(&paths.data_dir, tx.clone());
Arc::new(App {
store,
renderer,
browse: Browser::load(paths.config_dir.join("folders.json")),
paths: paths.clone(),
secrets: crate::secrets::Secrets::new(paths.config_dir.join("keys.json")),
token: std::sync::RwLock::new(token),
window,
events: tx,
shutdown: stop_tx,
started: Instant::now(),
built_v: asset_v,
ui: Ui::from_env(),
last_focus: std::sync::Mutex::new(Instant::now() - std::time::Duration::from_secs(60)),
windows: AtomicUsize::new(0),
streams: AtomicUsize::new(0),
pages: AtomicUsize::new(0),
online: std::sync::Mutex::new(Default::default()),
capabilities: crate::capability::Capabilities::load(paths.config_dir.join("capabilities")),
panes,
asides: Default::default(),
git: Default::default(),
exe,
started_at: crate::store::now(),
restart: std::sync::Mutex::new(None),
restart_wake: tokio::sync::Notify::new(),
leaving: std::sync::Mutex::new(Leaving::Stopped),
update,
relaunch_window: std::sync::atomic::AtomicBool::new(false),
restarting: std::sync::atomic::AtomicBool::new(false),
update_sent: Default::default(),
quota: Default::default(),
peers: Default::default(),
})
}
fn peer_routes() -> Router<Arc<App>> {
Router::new()
.route("/api/peers", get(peers_list))
.route("/api/peers/pair", post(pair_start))
.route("/api/peers/join", post(pair_join))
.route("/api/peers/pair/{code}", get(pair_state))
.route("/api/peers/{id}/rename", post(peer_rename))
.route("/api/peers/{id}/mute", post(peer_mute))
.route("/api/peers/{id}/remove", post(peer_remove))
.route("/api/peers/{id}/restore", post(peer_restore))
.route("/api/peers/{id}/note", post(peer_note))
.route("/api/peers/notes/{id}", post(peer_note_settle))
.route("/api/peers/offers/{id}", post(offer_answer))
.route("/api/docs/{id}/send", post(doc_send))
.route("/api/peers/{id}/desk", post(peer_desk))
.route("/api/docs/{id}/keep", post(doc_keep))
.route("/api/docs/{id}/save", post(doc_save))
.route("/api/docs/{id}/unfile", post(doc_unfile))
.route("/api/peers/outbox/{id}/retry", post(outbox_retry))
}
fn pane_routes() -> Router<Arc<App>> {
Router::new()
.route("/api/panes/{id}/delete", post(close_pane))
.route("/api/panes/{id}/restore", post(restore_pane))
.route("/api/panes/{id}/rename", post(rename_pane))
.route("/api/panes/{id}/start", post(start_pane))
.route("/api/panes/{id}/stop", post(stop_pane))
.route("/api/panes/{id}/agent", post(pane_agent))
.route("/api/panes/{id}/notes", get(pane_notes))
.route("/api/panes/{id}/notes/{note}/tick", post(pane_tick_note))
.route("/api/panes/{id}/notes/{note}/mark", post(pane_mark_note))
.route("/api/panes/{id}/name", post(pane_name))
.route("/api/panes/{id}/brief", get(pane_brief))
.route("/api/panes/{id}/changes", get(pane_changes))
.route("/api/panes/{id}/keys/{name}", get(pane_key))
.route("/api/panes/{id}/leftoff", post(pane_left_off))
.route("/api/panes/{id}/suggest", post(pane_suggest_note))
.route("/api/panes/{id}/offer", post(pane_offer))
.route(
"/api/panes/{id}/paste",
post(paste_image).layer(axum::extract::DefaultBodyLimit::max(receive::MAX_BYTES)),
)
}
fn thread_routes() -> Router<Arc<App>> {
Router::new()
.route("/api/panes/{id}/thread", post(pane_start_thread))
.route("/api/panes/{id}/thread/move", post(pane_move_thread))
.route("/api/panes/{id}/ask", post(pane_ask))
.route("/api/panes/{id}/handover", post(pane_hand_over))
.route("/api/panes/{id}/suggest-panel", post(pane_suggest_panel))
.route("/api/panes/{id}/suggest-desk", post(pane_suggest_desk))
.route("/api/panes/{id}/seen", post(pane_seen))
.route("/api/panes/{id}/band", get(pane_band))
.route(
"/api/panes/{id}/turns/{turn}",
get(pane_wait_turn).post(pane_answer_turn),
)
.route("/api/panes/{id}/note", post(pane_note))
.route("/api/claude-mod", get(mod_setting).post(set_mod_setting))
.route("/api/desks/{id}/threads", get(desk_threads))
.route("/api/desks/{id}/threads/{row}/move", post(desk_move_thread))
.route("/api/desks/{id}/threads/{row}/{act}", post(desk_thread_act))
.route("/api/desks/{id}/turns/{row}/answer", post(desk_answer_turn))
.route("/api/desks/{id}/turns/{row}/{act}", post(desk_turn_act))
.route(
"/api/desks/{id}/suggestions/{row}/{act}",
post(desk_suggestion_act),
)
}
fn receive_route() -> axum::routing::MethodRouter<Arc<App>> {
post(receive_doc).layer(axum::extract::DefaultBodyLimit::max(
receive::MAX_BYTES + 64 * 1024,
))
}
fn router(app: Arc<App>) -> Router {
Router::new()
.route("/", get(shell_home))
.route("/inbox", get(shell_inbox))
.route("/api/home", get(home))
.route("/connect", get(shell_connect))
.route("/start", get(shell_start))
.route("/welcome", get(shell_welcome))
.route("/d/{id}", get(shell_doc))
.route("/b/{id}", get(shell_browse))
.route("/b/{id}/{*path}", get(shell_browse_file))
.route("/files/{id}/{*path}", get(doc_file))
.route("/assets/mermaid.js", get(asset_mermaid))
.route("/assets/{name}", get(asset_named))
.route("/assets/fonts/{name}", get(asset_font))
.route("/api/health", get(health))
.route("/api/about", get(about))
.route("/api/agents", get(agents))
.route("/api/agents/claude/connect", post(connect_claude))
.route("/api/tree", get(tree))
.route("/api/projects/{id}/tree", get(project_tree))
.route("/api/workflows/{id}/tree", get(workflow_tree))
.route("/api/inbox", get(inbox))
.route("/api/search", get(search))
.route("/api/docs", receive_route())
.route("/api/docs/{id}", get(doc_json))
.route("/api/docs/{id}/pin", post(pin))
.route("/api/docs/{id}/read", post(mark_read))
.route("/api/queue", get(queue))
.route("/api/queue/clear", post(clear_queue))
.route("/api/queue/unread", post(unread))
.route("/api/docs/{id}/delete", post(delete_doc))
.route("/api/docs/{id}/undelete", post(undelete_doc))
.route("/api/removed", get(removed_list))
.route("/api/docs/{id}/history", get(history))
.route("/api/projects/{id}/rename", post(rename_project))
.route("/api/workflows/{id}/rename", post(rename_workflow))
.route("/api/docs/{id}/split", get(doc_split))
.route("/api/docs/{id}/outline", get(doc_outline))
.route("/api/notes", get(asides).post(receive_aside))
.route("/api/notes/seen", post(see_asides))
.route("/api/notes/dismiss", post(dismiss_asides))
.route("/api/notes/restore", post(restore_asides))
.route("/api/focus", post(focus))
.route("/api/shutdown", post(shutdown))
.route("/api/restart", post(restart).delete(cancel_restart))
.route("/api/update/check", post(update_check))
.route("/api/update/auto", post(update_auto))
.route("/api/update/later", post(update_later))
.route("/api/reset", get(reset_census).post(reset))
.route("/api/reveal", post(reveal))
.route("/api/resolve", post(resolve_path))
.route("/api/browse", get(browse_list).post(browse_open))
.route("/api/browse/pick", post(browse_pick))
.route("/api/browse/pick/cancel", post(browse_pick_cancel))
.route("/api/browse/{id}/close", post(browse_close))
.route("/api/browse/{id}/reopen", post(browse_reopen))
.route("/api/browse/{id}/tree", get(browse_tree))
.route("/api/browse/{id}/file", get(browse_file))
.route("/api/browse/{id}/raw", get(browse_raw))
.route("/api/browse/{id}/raw/{*path}", get(browse_raw_path))
.route("/api/browse/{id}/find", get(browse_find))
.route("/api/browse/{id}/outline", get(browse_outline))
.route("/api/docs/{id}/raw", get(doc_raw))
.route("/api/docs/{id}/blob", get(doc_blob))
.route("/api/compare/{a}/{b}", get(compare))
.route("/api/events", get(events))
.route("/api/capability", post(mint_capability))
.route("/api/desk", get(desk_socket))
.route("/api/desks", get(desks).post(create_desk))
.route("/api/desks/order", post(order_desks))
.route("/api/desks/{id}/rename", post(rename_desk))
.route("/api/desks/{id}/layout", post(desk_layout))
.route("/api/desks/{id}/move", post(move_pane))
.route("/api/desks/{id}/delete", post(delete_desk))
.route("/api/desks/{id}/reopen", post(reopen_desk))
.route("/api/desks/{id}/panes", post(open_pane))
.route("/api/desks/{id}/docs", get(desk_docs))
.route("/api/desks/{id}/docs/{doc}/remove", post(remove_desk_doc))
.route("/api/desks/{id}/docs/{doc}/restore", post(restore_desk_doc))
.route("/api/desks/{id}/notes", get(desk_notes).post(add_desk_note))
.route("/api/desks/{id}/notes/{note}", post(set_desk_note))
.route(
"/api/desks/{id}/notes/{note}/remove",
post(remove_desk_note),
)
.route(
"/api/desks/{id}/notes/{note}/restore",
post(restore_desk_note),
)
.route("/api/desks/{id}/notes/{note}/keep", post(keep_desk_note))
.route("/api/desks/{id}/leftoff", post(desk_left_off))
.route("/api/desks/{id}/keys", get(desk_keys).post(add_desk_key))
.route("/api/desks/{id}/keys/{name}/remove", post(remove_desk_key))
.route("/api/desks/{id}/visit", post(visit_desk))
.route("/api/desks/{id}/git", get(desk_git))
.route("/api/desks/{id}/park", post(park_desk))
.route("/api/desks/{id}/week", post(desk_week))
.route(
"/api/desks/{id}/notes/{note}/image",
post(add_note_image).layer(axum::extract::DefaultBodyLimit::max(NOTE_IMAGE_BYTES)),
)
.route("/api/desks/{id}/notes/{note}/images", post(set_note_images))
.route("/api/desks/{id}/note-images/{name}", get(note_image))
.route("/api/brief", get(brief_setting).post(set_brief_setting))
.route("/api/asides", get(asides_setting).post(set_asides_setting))
.merge(pane_routes())
.merge(thread_routes())
.merge(peer_routes())
.route("/desks", get(shell_desk_list))
.route("/desk/{id}", get(shell_desk))
.fallback(not_found)
.layer(axum::middleware::from_fn(host_gate))
.with_state(app)
}
fn nodelay(
listener: tokio::net::TcpListener,
) -> axum::serve::TapIo<tokio::net::TcpListener, fn(&mut tokio::net::TcpStream)> {
use axum::serve::ListenerExt;
listener.tap_io(|tcp| {
let _ = tcp.set_nodelay(true);
})
}
async fn asked_to_stop(mut stop_rx: broadcast::Receiver<()>, told: broadcast::Sender<()>) {
let term = async {
#[cfg(unix)]
{
let mut sig = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
.expect("SIGTERM handler");
sig.recv().await;
}
#[cfg(not(unix))]
std::future::pending::<()>().await;
};
tokio::select! {
_ = tokio::signal::ctrl_c() => {},
_ = term => {},
_ = stop_rx.recv() => {},
}
let _ = told.send(());
}
pub async fn run(paths: Paths) -> anyhow::Result<Leaving> {
let token = config::load_or_create_token(&paths)?;
let window = config::load_or_create_window_secret(&paths)?;
let store = Store::open(&paths)?;
let planned = read_restart_marker(&paths);
let (mut resume, mut offer) = store.panes_resume().unwrap_or_default();
if planned.is_none() {
offer.append(&mut resume);
}
if let Some(apply) = planned {
eprintln!(
"snyvi: back from a planned restart{}; {} panel(s) to resume",
if apply {
" (an update was applied)"
} else {
""
},
resume.len()
);
}
let exe = Exe::here();
let update = exe.as_ref().map(|e| {
Arc::new(crate::update::Updater::new(
&paths,
&e.path,
Box::new(crate::update::Http),
))
});
let app = new_app(&paths, store, token, window, update, exe);
let stop_rx = app.shutdown.subscribe();
if let Some(u) = &app.update {
if u.channel == crate::update::Channel::Dev {
eprintln!("snyvi: a development build; it will not check for updates");
} else if !u.auto() {
eprintln!("snyvi: automatic updates are off; `snyvi update` still works");
}
}
app.panes.mark_resume(resume);
app.panes.mark_offer(offer);
let weak = Arc::downgrade(&app);
app.panes.on_cwd(Box::new(move |id, cwd| {
if let Some(app) = weak.upgrade() {
let _ = app.store.set_pane_cwd(id, cwd);
}
}));
let leaving = app.clone();
let told = app.shutdown.clone();
crate::watch::spawn_browse_watcher(app.clone());
crate::watch::spawn_ui_watcher(app.clone());
crate::claude_mod::start(&paths);
let router = router(app);
let addr = format!("127.0.0.1:{}", config::port());
let listener = tokio::net::TcpListener::bind(&addr).await?;
if let Some(u) = &leaving.update {
match u.note_started(planned.unwrap_or(false), crate::store::now()) {
Some(crate::update::Started::Applied(v)) => {
eprintln!("snyvi: now {v}, updated; the next automatic update is a day or so away")
}
Some(crate::update::Started::Failed(v)) => eprintln!(
"snyvi: {v} was put in place but this is {VERSION} running from the same path; {v} is marked failed and not tried again on its own"
),
None => {}
}
}
drop_restart_marker(&paths);
let _ = leaving.store.clear_panes_resume();
spawn_restart_watcher(leaving.clone());
spawn_update_checker(leaving.clone());
spawn_peer_link(leaving.clone());
tokio::task::spawn_blocking(|| {
if let Ok(true) = crate::hook::top_up() {
eprintln!("snyvi: added the panel status hooks to ~/.claude/settings.json");
}
});
eprintln!("snyvi {VERSION} listening on http://{addr}");
axum::serve(nodelay(listener), router)
.with_graceful_shutdown(asked_to_stop(stop_rx, told))
.await?;
let why = leaving.leaving.lock().unwrap().clone();
if matches!(why, Leaving::Stopped) {
let (resume, offer) = leaving.panes.unspent();
let mut with_agent = leaving.panes.with_agent();
with_agent.extend(resume);
with_agent.extend(offer);
if let Ok(n) = leaving.store.offer_panes_resume(&with_agent) {
if n > 0 {
eprintln!("snyvi: {n} panel(s) had Claude open; the window will offer each conversation back");
}
}
}
leaving.panes.shutdown();
Ok(why)
}
pub(crate) fn emit(app: &App, name: &str, data: serde_json::Value) {
let _ = app.events.send(format!("{name}\n{data}"));
}
pub(crate) fn doc_event(app: &App, received: &receive::Received) -> serde_json::Value {
let doc = &received.doc;
let mut ev = json!({ "doc": doc, "url": format!("{}/d/{}", config::base_url(), doc.id), "existing": received.existing, "supersedes": received.supersedes, "waiting": waiting(app) });
if app.pages.load(Ordering::Relaxed) > 0 {
ev["project"] = json!(app.store.project_row(doc.project_id).ok().flatten());
ev["rows"] = json!(project_rows(
app,
doc.project_id,
TREE_WORKFLOWS,
TREE_DOCS,
None
));
}
ev
}
async fn receive_and_emit(
app: &Arc<App>,
payload: Payload,
) -> Result<receive::Received, Box<Response>> {
let app2 = app.clone();
match tokio::task::spawn_blocking(move || {
receive::receive(&app2.store, &app2.renderer, payload).map(|r| (doc_event(&app2, &r), r))
})
.await
{
Ok(Ok((event, received))) => {
emit(app, "doc", event);
Ok(received)
}
Ok(Err(e)) => Err(Box::new(
(
StatusCode::BAD_REQUEST,
Json(json!({ "error": e.to_string() })),
)
.into_response(),
)),
Err(e) => Err(Box::new(err(anyhow::anyhow!(e)))),
}
}
fn err(e: anyhow::Error) -> Response {
(
StatusCode::INTERNAL_SERVER_ERROR,
Json(json!({ "error": e.to_string() })),
)
.into_response()
}