use super::*;
use axum::body::Bytes;
use std::collections::HashMap;
static FOCUSED: std::sync::LazyLock<std::sync::Mutex<HashMap<String, (u64, bool)>>> =
std::sync::LazyLock::new(Default::default);
const FOCUS_PAGES: usize = 64;
#[derive(Deserialize, Default)]
pub(crate) struct FocusBody {
focused: Option<bool>,
page: Option<String>,
seq: Option<u64>,
}
pub(crate) async fn focus(State(app): S, headers: HeaderMap, body: Bytes) -> Response {
if let Some(no) = refuse_reader(&app, &headers) {
return no;
}
let b: FocusBody = serde_json::from_slice(&body).unwrap_or_default();
let focused = b.focused.unwrap_or(true);
if let Some(page) = b.page.filter(|p| !p.is_empty() && p.len() <= 64) {
let seq = b.seq.unwrap_or(0);
let mut m = FOCUSED.lock().unwrap();
if m.get(&page).is_none_or(|(last, _)| seq >= *last) {
if m.len() >= FOCUS_PAGES && !m.contains_key(&page) {
if let Some(k) = m.iter().find(|(_, v)| !v.1).map(|(k, _)| k.clone()) {
m.remove(&k);
} else {
m.clear();
}
}
m.insert(page, (seq, focused));
}
}
if focused {
*app.last_focus.lock().unwrap() = Instant::now();
}
StatusCode::NO_CONTENT.into_response()
}
pub(crate) fn focused(app: &App) -> bool {
if app.last_focus.lock().unwrap().elapsed() < FOCUS_FOR {
return true;
}
let mut m = FOCUSED.lock().unwrap();
if app.pages.load(Ordering::Relaxed) == 0 {
m.clear();
return false;
}
m.values().any(|(_, f)| *f)
}
pub(crate) fn focus_age(app: &App) -> std::time::Duration {
if focused(app) {
std::time::Duration::ZERO
} else {
app.last_focus.lock().unwrap().elapsed()
}
}
const FOCUS_FOR: std::time::Duration = std::time::Duration::from_secs(4);
pub(crate) fn notify_desktop(app: &App, doc: &Doc) {
if std::env::var("SNYVI_NOTIFY")
.map(|v| v == "0")
.unwrap_or(false)
{
return;
}
if focused(app) {
return;
}
let url = format!("{}/d/{}", config::base_url(), doc.id);
let has_window = app.has_window();
crate::platform::notify_open(
&doc.title,
&format!("{} · {}", doc.project, doc.workflow_title),
move || {
if has_window && crate::desktop::hand_to_window(&url) {
return;
}
platform::open_url(&url);
},
);
}
#[derive(Deserialize)]
pub(crate) struct EventsQ {
#[serde(default)]
pub(crate) window: Option<String>,
#[serde(default)]
pub(crate) agent: Option<String>,
}
impl EventsQ {
pub(crate) fn is_window(&self) -> bool {
self.window
.as_deref()
.is_some_and(|v| !matches!(v, "" | "0" | "false" | "False"))
}
pub(crate) fn window_version(&self) -> Option<String> {
self.window.clone().filter(|_| self.is_window())
}
pub(crate) fn agent(&self) -> Option<String> {
let name = self.agent.as_deref()?.trim();
if name.is_empty() {
return None;
}
Some(name.chars().take(64).collect())
}
}
pub(crate) struct StreamMark {
pub(crate) app: Arc<App>,
pub(crate) window: bool,
pub(crate) agent: Option<String>,
}
impl StreamMark {
pub(crate) fn new(
app: Arc<App>,
window: bool,
agent: Option<String>,
window_version: Option<String>,
) -> StreamMark {
app.streams.fetch_add(1, Ordering::Relaxed);
if agent.is_none() {
app.pages.fetch_add(1, Ordering::Relaxed);
}
if window {
app.windows.fetch_add(1, Ordering::Relaxed);
if let Some(u) = &app.update {
if u.window_is_too_old(window_version.as_deref()) {
if let Some(bin) = u.app_path() {
eprintln!("snyvi: the window is {}, older than this release wants; relaunching it", window_version.as_deref().unwrap_or("?"));
app.relaunch_window.store(true, Ordering::Relaxed);
if let Err(e) = platform::spawn_detached(&bin, &["--quit"]) {
eprintln!("snyvi: could not ask the window to quit: {e}");
app.relaunch_window.store(false, Ordering::Relaxed);
}
}
}
}
}
if let Some(name) = &agent {
let changed = {
let mut m = app.online.lock().unwrap_or_else(|e| e.into_inner());
*m.entry(name.clone()).or_insert(0) += 1;
json!(*m)
};
emit(&app, "agents", json!({ "online": changed }));
}
StreamMark { app, window, agent }
}
}
impl Drop for StreamMark {
fn drop(&mut self) {
self.app.streams.fetch_sub(1, Ordering::Relaxed);
if self.agent.is_none() {
self.app.pages.fetch_sub(1, Ordering::Relaxed);
}
if self.window {
let left = self
.app
.windows
.fetch_sub(1, Ordering::Relaxed)
.saturating_sub(1);
if left == 0 && self.app.relaunch_window.load(Ordering::Relaxed) {
relaunch_window_when_gone(self.app.clone());
}
}
self.app.restart_wake.notify_one();
if let Some(name) = &self.agent {
let changed = {
let mut m = self.app.online.lock().unwrap_or_else(|e| e.into_inner());
if let Some(n) = m.get_mut(name) {
*n -= 1;
if *n == 0 {
m.remove(name);
}
}
json!(*m)
};
emit(&self.app, "agents", json!({ "online": changed }));
}
}
}
pub(crate) fn relaunch_window_when_gone(app: Arc<App>) {
let Ok(rt) = tokio::runtime::Handle::try_current() else {
return;
};
rt.spawn(async move {
tokio::time::sleep(WINDOW_GONE_FOR).await;
if app.windows.load(Ordering::Relaxed) != 0
|| !app.relaunch_window.swap(false, Ordering::Relaxed)
{
return;
}
if let Some(exe) = app.exe.as_ref().map(|e| e.path.clone()) {
if let Err(e) = platform::spawn_detached(&exe, &["app"]) {
eprintln!("snyvi: could not start the window again: {e}");
}
}
});
}
pub(crate) const WINDOW_GONE_FOR: std::time::Duration = std::time::Duration::from_secs(2);
pub(crate) async fn events(
State(app): S,
Query(q): Query<EventsQ>,
) -> Sse<impl tokio_stream::Stream<Item = Result<Event, Infallible>>> {
let rx = app.events.subscribe();
let stop = BroadcastStream::new(app.shutdown.subscribe()).map(|_| None);
let mark = StreamMark::new(app.clone(), q.is_window(), q.agent(), q.window_version());
let first = tokio_stream::once(Some(Ok(Event::default()
.event("update")
.data(update_json(&app).to_string()))));
let stream = first
.chain(BroadcastStream::new(rx).filter_map(move |m| {
let _keep = &mark;
let ev = match m {
Ok(msg) => {
let (name, data) = msg.split_once('\n').unwrap_or(("doc", msg.as_str()));
Event::default().event(name).data(data)
}
Err(_lagged) => Event::default().event("resync").data("{}"),
};
Some(Some(Ok(ev)))
}))
.merge(stop)
.take_while(Option::is_some)
.map(Option::unwrap);
Sse::new(stream).keep_alive(KeepAlive::default())
}