use super::*;
#[derive(Clone, Copy, Debug)]
pub struct Pending {
pub since: Instant,
pub apply: bool,
pub back: bool,
pub now: bool,
}
#[derive(Clone, Debug)]
pub enum Leaving {
Stopped,
Restart { exe: Option<PathBuf>, apply: bool },
}
pub(crate) fn under_systemd() -> bool {
if std::env::var_os("INVOCATION_ID").is_none() {
return false;
}
std::process::Command::new("systemctl")
.args(["--user", "is-active", "--quiet", "snyvi"])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.status()
.map(|s| s.success())
.unwrap_or(false)
}
pub const PLANNED_RESTART_EXIT: i32 = 75;
pub(crate) const RESTART_MARKER: &str = "restart.json";
pub(crate) fn leave_for_restart(app: &App, apply: bool, back: bool) -> bool {
app.restarting.store(true, Ordering::Relaxed);
let left = leave(app, apply, back);
if !left {
app.restarting.store(false, Ordering::Relaxed);
emit_update(app);
}
left
}
pub(crate) fn leave(app: &App, apply: bool, back: bool) -> bool {
if apply || back {
let Some(u) = &app.update else {
eprintln!("snyvi: no updater on this daemon; not restarting");
return false;
};
if back {
match u.rollback(true) {
Ok(true) => {
eprintln!("snyvi: the previous version is back in place; restarting onto it")
}
Ok(false) => {
eprintln!("snyvi: there is no previous version to go back to");
return false;
}
Err(e) => {
eprintln!("snyvi: could not put the previous version back: {e:#}");
u.fail(format!("{e:#}"));
emit_update(app);
return false;
}
}
} else {
match u.apply() {
Ok(v) => eprintln!("snyvi: {v} is in place; restarting onto it"),
Err(e) => {
eprintln!("snyvi: could not apply the update: {e:#}");
u.fail(format!("{e:#}"));
u.drop_staged();
emit_update(app);
return false;
}
}
}
}
app.restarting.store(true, Ordering::Relaxed);
emit_update(app);
let (mut resume, offer) = app.panes.unspent();
resume.extend(app.panes.with_agent());
match app.store.mark_panes_resume(&resume) {
Ok(n) if n > 0 => {
eprintln!("snyvi: restarting; {n} panel(s) will resume their conversation")
}
Ok(_) => eprintln!("snyvi: restarting"),
Err(e) => eprintln!("snyvi: restarting; could not mark panels to resume: {e}"),
}
let _ = app.store.offer_panes_resume(&offer);
let marker = json!({ "apply": apply, "from": VERSION, "at": crate::store::now() });
let _ = std::fs::write(app.paths.data_dir.join(RESTART_MARKER), marker.to_string());
*app.leaving.lock().unwrap() = Leaving::Restart {
exe: app.exe.as_ref().map(|e| e.path.clone()),
apply,
};
let _ = app.shutdown.send(());
true
}
pub(crate) const RESTART_MARKER_FOR: i64 = 10 * 60;
pub(crate) fn read_restart_marker(paths: &Paths) -> Option<bool> {
let text = std::fs::read_to_string(paths.data_dir.join(RESTART_MARKER)).ok()?;
let v: serde_json::Value = serde_json::from_str(&text).ok()?;
let at = v["at"].as_i64()?;
if (crate::store::now() - at).abs() > RESTART_MARKER_FOR {
return None;
}
Some(v["apply"].as_bool().unwrap_or(false))
}
pub(crate) fn drop_restart_marker(paths: &Paths) {
let _ = std::fs::remove_file(paths.data_dir.join(RESTART_MARKER));
}
pub(crate) fn pending_json(app: &App) -> serde_json::Value {
let pending = *app.restart.lock().unwrap();
match pending {
Some(p) => {
let busy = if p.now { vec![] } else { app.panes.busy() };
json!({
"apply": p.apply,
"back": p.back,
"now": p.now,
"waiting": waiting_named(app, &busy),
"waiting_on": busy,
})
}
None => serde_json::Value::Null,
}
}
pub(crate) fn waiting_named(app: &App, ids: &[String]) -> Vec<serde_json::Value> {
ids.iter()
.map(|id| {
let st = app.panes.status(id);
let placed = app.store.pane(id).ok().flatten();
json!({
"pane": id,
"desk": placed.as_ref().map(|p| p.desk_name.clone()),
"desk_id": placed.as_ref().map(|p| p.desk_id),
"slot": placed.as_ref().map(|p| p.pane.slot),
"agent": st.agent,
"since": st.agent_since.or(st.since),
})
})
.collect()
}
const RESTART_TICK: std::time::Duration = std::time::Duration::from_secs(5);
const RESTART_IDLE_TICK: std::time::Duration = std::time::Duration::from_secs(30);
pub(crate) fn restart_json(app: &App) -> serde_json::Value {
match *app.restart.lock().unwrap() {
Some(p) => json!({
"pending": true,
"apply": p.apply,
"back": p.back,
"since_s": p.since.elapsed().as_secs(),
"waiting_on": if p.now { vec![] } else { app.panes.busy() },
}),
None => serde_json::Value::Null,
}
}
pub(crate) fn spawn_restart_watcher(app: Arc<App>) {
tokio::spawn(async move {
let mut changes = app.events.subscribe();
loop {
let every = if app.restart.lock().unwrap().is_some() {
RESTART_TICK
} else {
RESTART_IDLE_TICK
};
tokio::select! {
_ = app.restart_wake.notified() => {}
_ = tokio::time::sleep(every) => {}
m = changes.recv() => {
if !matches!(&m, Ok(s) if s.starts_with("panes\n")) { continue }
}
}
emit_update_if_changed(&app);
let pending = *app.restart.lock().unwrap();
match pending {
Some(p) => {
if !p.now && !app.panes.busy().is_empty() {
continue;
}
if (p.apply || p.back) && app.update.as_ref().is_some_and(|u| u.busy()) {
continue;
}
app.restart.lock().unwrap().take();
if leave_for_restart(&app, p.apply, p.back) {
return;
}
}
None => {
let Some(u) = &app.update else { continue };
let now = crate::store::now();
if let Some(why) = u.should_apply(now, &doors(&app)) {
let v = u.state().ready.unwrap_or_default();
eprintln!("snyvi: applying {v} now, since {}", why.say());
if leave_for_restart(&app, true, false) {
return;
}
}
}
}
}
});
}
pub(crate) fn doors(app: &App) -> crate::update::Doors {
let agents: usize = app
.online
.lock()
.unwrap_or_else(|e| e.into_inner())
.values()
.sum();
let streams = app.streams.load(Ordering::Relaxed);
crate::update::Doors {
windows: app.windows.load(Ordering::Relaxed),
pages: streams.saturating_sub(agents),
focus_age: focus_age(app),
busy: !app.panes.busy().is_empty(),
}
}
pub(crate) fn spawn_update_checker(app: Arc<App>) {
use crate::update::{jitter_d, CHECK_EVERY, CHECK_JITTER, FIRST_CHECK, FIRST_CHECK_JITTER};
let Some(u) = app.update.clone() else { return };
if u.channel == crate::update::Channel::Dev {
return;
}
let every = std::env::var("SNYVI_UPDATE_EVERY_S")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.map(std::time::Duration::from_secs);
tokio::spawn(async move {
let mut wait = match every {
Some(e) => e,
None => FIRST_CHECK + jitter_d(FIRST_CHECK_JITTER),
};
loop {
tokio::time::sleep(wait).await;
wait = match every {
Some(e) => e,
None => CHECK_EVERY - CHECK_JITTER + jitter_d(CHECK_JITTER * 2),
};
if !u.auto() {
continue;
}
let checker = u.clone();
let before = checker.state().ready;
let r =
tokio::task::spawn_blocking(move || checker.check(crate::update::Ask::TIMER, None))
.await;
match r {
Ok(Ok(c)) if c.ready.is_some() && c.ready != before => {
eprintln!(
"snyvi: {} is staged, verified, and applies at a quiet moment",
c.ready.as_deref().unwrap_or("?")
)
}
Ok(Ok(c)) if c.newer() && c.told => eprintln!(
"snyvi: {} is out; this install is updated by hand",
c.latest
),
Ok(Ok(_)) => {}
Ok(Err(e)) => eprintln!("snyvi: update check: {e:#}"),
Err(_) => {}
}
emit_update(&app);
app.restart_wake.notify_one();
}
});
}
pub(crate) fn update_json(app: &App) -> serde_json::Value {
let stale = app.exe.as_ref().is_some_and(Exe::stale);
let mut j = match &app.update {
Some(u) => u.json(crate::store::now(), stale),
None => json!({ "channel": "unknown", "auto": false, "show": stale, "stale": stale }),
};
j["restart"] = pending_json(app);
j["restarting"] = json!(app.restarting.load(Ordering::Relaxed));
j
}
pub(crate) fn emit_update(app: &App) {
let j = update_json(app);
*app.update_sent.lock().unwrap_or_else(|e| e.into_inner()) = j.to_string();
emit(app, "update", j);
}
pub(crate) fn emit_update_if_changed(app: &App) {
let j = update_json(app);
let text = j.to_string();
{
let mut sent = app.update_sent.lock().unwrap_or_else(|e| e.into_inner());
if *sent == text {
return;
}
*sent = text;
}
emit(app, "update", j);
}
pub fn relaunch(exe: Option<PathBuf>, apply: bool, paths: &Paths) -> ! {
if under_systemd() {
eprintln!("snyvi: planned restart under systemd; exiting {PLANNED_RESTART_EXIT} for the unit to start the new file");
std::process::exit(PLANNED_RESTART_EXIT);
}
let Some(exe) = exe.or_else(|| std::env::current_exe().ok()) else {
eprintln!("snyvi: cannot restart: the path of this executable is unknown");
std::process::exit(1);
};
let me = std::process::id();
if let Err(e) = platform::spawn_daemon(&exe) {
eprintln!(
"snyvi: cannot restart: starting {} failed: {e}",
exe.display()
);
std::process::exit(1);
}
match came_up(me, &exe) {
CameUp::Ours(v) => {
eprintln!(
"snyvi: {v} is up on {}; this one is done",
config::base_url()
);
std::process::exit(0);
}
CameUp::Other(what) => {
eprintln!(
"snyvi: something else answered on {} after the restart ({what}), not {}; stop it and run `snyvi restart`",
config::base_url(),
exe.display()
);
std::process::exit(1);
}
CameUp::Nothing => {}
}
eprintln!(
"snyvi: {} did not answer within {} s of a planned restart; see {}",
exe.display(),
CAME_UP_WITHIN.as_secs(),
platform::daemon_log(&paths.data_dir).display()
);
if apply {
let u = crate::update::Updater::new(paths, &exe, Box::new(crate::update::Http));
match u.rollback(false) {
Ok(true) => eprintln!("snyvi: put the previous version back"),
Ok(false) => eprintln!("snyvi: nothing to put back"),
Err(e) => eprintln!("snyvi: could not put the previous version back: {e:#}"),
}
if let Err(e) = platform::spawn_daemon(&exe) {
eprintln!("snyvi: starting {} again failed: {e}", exe.display());
std::process::exit(1);
}
if let CameUp::Ours(v) = came_up(me, &exe) {
eprintln!("snyvi: {v} is back up on {}", config::base_url());
std::process::exit(0);
}
eprintln!(
"snyvi: {} did not answer after the rollback either",
exe.display()
);
}
std::process::exit(1);
}
pub(crate) const CAME_UP_WITHIN: std::time::Duration = std::time::Duration::from_secs(30);
pub(crate) enum CameUp {
Ours(String),
Other(String),
Nothing,
}
pub(crate) fn came_up(me: u32, exe: &std::path::Path) -> CameUp {
let deadline = Instant::now() + CAME_UP_WITHIN;
let mut wait = std::time::Duration::from_millis(20);
let mut other = None;
while Instant::now() < deadline {
if let Some(h) = crate::client::health() {
let version = h["version"].as_str().unwrap_or("?").to_string();
if h["pid"].as_u64() != Some(u64::from(me)) {
match h["exe"].as_str() {
Some(e) if std::path::Path::new(e) != exe => {
other = Some(format!("snyvi {version} from {e}"));
}
_ => return CameUp::Ours(version),
}
}
}
std::thread::sleep(wait);
wait = (wait * 2).min(std::time::Duration::from_millis(250));
}
match other {
Some(what) => CameUp::Other(what),
None => CameUp::Nothing,
}
}
pub(crate) async fn health(State(app): S) -> Json<serde_json::Value> {
Json(json!({
"ok": true,
"version": VERSION,
"commit": BUILD_SHA,
"pid": std::process::id(),
"exe": app.exe.as_ref().map(|e| e.path.display().to_string()),
"docs": app.store.count().unwrap_or(0),
"window": app.has_window(),
"streams": app.streams.load(Ordering::Relaxed),
"panes": app.panes.running(),
"git_runs": crate::pane::git_runs(),
"agents": app.online(),
"v": app.asset_v(),
"languages": app.renderer.languages().len(),
"uptime_s": app.started.elapsed().as_secs(),
"started": app.started_at,
"stale": app.exe.as_ref().is_some_and(Exe::stale),
"restart": restart_json(&app),
"update": update_json(&app),
}))
}
pub(crate) async fn about(State(app): S) -> Json<serde_json::Value> {
let exe = app.exe.as_ref().map(|e| e.path.clone());
Json(json!({
"name": "snyvi",
"description": env!("CARGO_PKG_DESCRIPTION"),
"version": VERSION,
"commit": BUILD_SHA,
"target": BUILD_TARGET,
"binary": exe.as_deref().map(|p| p.display().to_string()),
"stale": app.exe.as_ref().is_some_and(Exe::stale),
"data_dir": app.paths.data_dir.display().to_string(),
"config_dir": app.paths.config_dir.display().to_string(),
"agents": std::iter::once(crate::setup::claude_code_status())
.chain(crate::agents::status_lines())
.collect::<Vec<_>>()
.join("\n"),
"license": env!("CARGO_PKG_LICENSE"),
"repository": env!("CARGO_PKG_REPOSITORY"),
"docs": app.store.count().unwrap_or(0),
"uptime_s": app.started.elapsed().as_secs(),
"update": update_json(&app),
}))
}
pub(crate) async fn shutdown(State(app): S, headers: HeaderMap) -> Response {
if !windowed(&app, &headers) {
return not_windowed();
}
let _ = app.shutdown.send(());
Json(json!({ "ok": true, "version": VERSION })).into_response()
}
#[derive(Deserialize)]
pub(crate) struct RestartBody {
#[serde(default)]
pub(crate) when: Option<String>,
#[serde(default)]
pub(crate) apply: bool,
#[serde(default)]
pub(crate) back: bool,
}
pub(crate) async fn restart(
State(app): S,
headers: HeaderMap,
Json(b): Json<RestartBody>,
) -> Response {
if !windowed(&app, &headers) {
return not_windowed();
}
let now = match b.when.as_deref().unwrap_or("idle") {
"idle" => false,
"now" => true,
other => {
return (
StatusCode::BAD_REQUEST,
Json(json!({ "error": format!("when must be idle or now, not {other:?}") })),
)
.into_response()
}
};
if b.apply && app.update.as_ref().is_none_or(|u| u.staged().is_none()) {
return (
StatusCode::CONFLICT,
Json(
json!({ "error": "nothing is staged to apply; `snyvi update` checks and stages" }),
),
)
.into_response();
}
if b.back && !app.update.as_ref().is_some_and(|u| u.can_go_back()) {
return (
StatusCode::CONFLICT,
Json(json!({ "error": "there is no previous version beside this one to go back to" })),
)
.into_response();
}
let now = {
let mut slot = app.restart.lock().unwrap();
let was = *slot;
let p = Pending {
since: was.map(|p| p.since).unwrap_or_else(Instant::now),
apply: b.apply || was.is_some_and(|p| p.apply),
back: b.back || was.is_some_and(|p| p.back),
now: now || was.is_some_and(|p| p.now),
};
*slot = Some(p);
p.now
};
let waiting_on = if now { vec![] } else { app.panes.busy() };
let waiting = waiting_named(&app, &waiting_on);
app.restart_wake.notify_one();
emit_update(&app);
Json(json!({ "ok": true, "waiting_on": waiting_on, "waiting": waiting, "version": VERSION }))
.into_response()
}
pub(crate) async fn cancel_restart(State(app): S, headers: HeaderMap) -> Response {
if !windowed(&app, &headers) {
return not_windowed();
}
let cancelled = app.restart.lock().unwrap().take().is_some();
if cancelled {
eprintln!("snyvi: the restart that was waiting is called off");
}
emit_update(&app);
Json(json!({ "ok": true, "cancelled": cancelled })).into_response()
}
#[derive(Deserialize)]
pub(crate) struct UpdateCheckBody {
#[serde(default)]
pub(crate) to: Option<String>,
#[serde(default = "yes")]
pub(crate) lift: bool,
}
pub(crate) fn yes() -> bool {
true
}
pub(crate) async fn update_check(
State(app): S,
headers: HeaderMap,
Json(b): Json<UpdateCheckBody>,
) -> Response {
if !windowed(&app, &headers) {
return not_windowed();
}
let Some(u) = app.update.clone() else {
return (StatusCode::CONFLICT, Json(json!({ "error": "this daemon cannot say what file it runs from, so it does not update itself" }))).into_response();
};
if u.channel == crate::update::Channel::Dev {
return (
StatusCode::CONFLICT,
Json(json!({ "error": "a development build does not update itself" })),
)
.into_response();
}
if let Some(to) = &b.to {
if semver::Version::parse(to).is_err() {
return (
StatusCode::BAD_REQUEST,
Json(json!({ "error": format!("{to:?} is not a version") })),
)
.into_response();
}
}
let to = b.to.clone();
let ask = crate::update::Ask {
fresh: true,
lift: b.lift,
};
let r = tokio::task::spawn_blocking(move || u.check(ask, to.as_deref())).await;
emit_update(&app);
app.restart_wake.notify_one();
match r {
Ok(Ok(c)) => {
if b.lift {
if let Some(u) = &app.update {
let _ = u.snooze(None, crate::store::now());
}
}
Json(json!({
"ok": true,
"running": c.running.to_string(),
"latest": c.latest.to_string(),
"newer": c.newer(),
"ready": c.ready,
"told": c.told,
"update": update_json(&app),
}))
.into_response()
}
Ok(Err(e)) => (
StatusCode::BAD_GATEWAY,
Json(json!({ "error": format!("{e:#}"), "update": update_json(&app) })),
)
.into_response(),
Err(e) => err(anyhow::anyhow!("the check did not finish: {e}")),
}
}
#[derive(Deserialize)]
pub(crate) struct UpdateAutoBody {
pub(crate) on: bool,
}
pub(crate) async fn update_auto(
State(app): S,
headers: HeaderMap,
Json(b): Json<UpdateAutoBody>,
) -> Response {
if !windowed(&app, &headers) {
return not_windowed();
}
let Some(u) = &app.update else {
return (
StatusCode::CONFLICT,
Json(json!({ "error": "this daemon does not update itself" })),
)
.into_response();
};
match u.set_auto(b.on) {
Ok(auto) => {
emit_update(&app);
Json(json!({ "ok": true, "auto": auto, "update": update_json(&app) })).into_response()
}
Err(e) => err(e),
}
}
#[derive(Deserialize)]
pub(crate) struct LaterBody {
#[serde(default)]
pub(crate) until: Option<i64>,
}
pub(crate) async fn update_later(
State(app): S,
headers: HeaderMap,
Json(b): Json<LaterBody>,
) -> Response {
if !windowed(&app, &headers) {
return not_windowed();
}
let Some(u) = &app.update else {
return (
StatusCode::CONFLICT,
Json(json!({ "error": "this daemon does not update itself" })),
)
.into_response();
};
match u.snooze(b.until, crate::store::now()) {
Ok(()) => {
emit_update(&app);
Json(json!({ "ok": true, "update": update_json(&app) })).into_response()
}
Err(e) => err(e),
}
}
pub(crate) async fn reset_census(State(app): S) -> Response {
match app.store.census() {
Ok(c) => Json(c).into_response(),
Err(e) => err(e),
}
}
#[derive(Deserialize)]
pub(crate) struct ResetBody {
pub(crate) documents: i64,
#[serde(default)]
pub(crate) pinned: bool,
#[serde(default)]
pub(crate) desks: i64,
}
pub(crate) fn docs(n: i64) -> String {
format!("{n} document{}", if n == 1 { "" } else { "s" })
}
pub(crate) fn pinned_docs(n: i64) -> String {
format!("{n} pinned document{}", if n == 1 { "" } else { "s" })
}
pub(crate) async fn reset(State(app): S, headers: HeaderMap, Json(b): Json<ResetBody>) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
let census = match app.store.census() {
Ok(c) => c,
Err(e) => return err(e),
};
if b.documents != census.documents {
return (
StatusCode::CONFLICT,
Json(json!({
"error": format!("the library has changed: {} now, not {}", docs(census.documents), b.documents),
"census": census,
})),
)
.into_response();
}
if b.desks != census.desks {
return (
StatusCode::CONFLICT,
Json(json!({
"error": format!("the desks have changed: {} now, not {}", census.desks, b.desks),
"census": census,
})),
)
.into_response();
}
if census.pinned > 0 && !b.pinned {
return (
StatusCode::CONFLICT,
Json(json!({
"error": format!("{} would go with it", pinned_docs(census.pinned)),
"census": census,
})),
)
.into_response();
}
let app2 = app.clone();
match tokio::task::spawn_blocking(move || app2.store.reset()).await {
Ok(Ok(())) => {}
Ok(Err(e)) => return err(e),
Err(e) => return err(anyhow::anyhow!("reset task: {e}")),
}
app.panes.clear();
for root in app.browse.list() {
app.browse.close(&root.id);
}
app.browse.forget_closed();
let _ = std::fs::remove_file(app.paths.config_dir.join("sessions.json"));
let _ = std::fs::remove_file(app.paths.config_dir.join("folders.json"));
match config::rotate_token(&app.paths) {
Ok(t) => *app.token.write().unwrap() = t,
Err(e) => return err(e),
}
emit(&app, "reset", json!({}));
Json(json!({ "ok": true, "removed": census })).into_response()
}