Skip to main content

unifier/daemon/
server.rs

1//! Hot daemon server: owns a [`HotStore`] and serves requests over a Unix socket.
2
3use std::fs;
4use std::io::{BufRead, BufReader, Write};
5use std::os::unix::net::{UnixListener, UnixStream};
6use std::sync::atomic::{AtomicBool, Ordering};
7use std::sync::Arc;
8use std::thread;
9use std::time::{Duration, Instant};
10
11use crate::daemon::idle::{
12    idle_after, idle_stop_after, max_load, should_idle_flush, should_idle_stop, system_is_idle,
13};
14
15use crate::daemon::notify::{event_name, EventHub, Notice};
16use crate::daemon::paths::{daemon_dir, events_socket_path, pid_path, socket_path};
17use crate::daemon::protocol::{
18    decode_request, encode_response, ok_empty, MessageDto, Request, Response,
19};
20use crate::daemon::www;
21use crate::envelope::DEFAULT_SENDER;
22use crate::error::{Error, Result};
23use crate::home::UnifierHome;
24use crate::store::HotStore;
25use crate::tick::TickStartOutcome;
26
27pub fn run(home: UnifierHome) -> Result<()> {
28    let shutdown = Arc::new(AtomicBool::new(false));
29    home.ensure()?;
30    fs::create_dir_all(daemon_dir(&home))?;
31
32    let events_sock = events_socket_path(&home);
33    if events_sock.exists() {
34        fs::remove_file(&events_sock)?;
35    }
36    let events_listener = UnixListener::bind(&events_sock)?;
37    events_listener.set_nonblocking(true)?;
38    let hub = Arc::new(EventHub::default());
39
40    let sock = socket_path(&home);
41    if sock.exists() {
42        fs::remove_file(&sock)?;
43    }
44    let listener = UnixListener::bind(&sock)?;
45    listener.set_nonblocking(true)?;
46
47    let store = Arc::new(std::sync::Mutex::new(HotStore::load(&home)?));
48    write_pid(&home)?;
49
50    let http_activity = Arc::new(AtomicBool::new(false));
51    let http_port = match www::spawn(home.clone(), Arc::clone(&shutdown), Arc::clone(&http_activity))
52    {
53        Ok(port) => {
54            eprintln!("www listening on http://127.0.0.1:{port}");
55            Some(port)
56        }
57        Err(e) => {
58            eprintln!("www server disabled: {e}");
59            None
60        }
61    };
62    let _ = http_port;
63
64    let idle_after = idle_after();
65    let idle_stop_after = idle_stop_after();
66    let max_load = max_load();
67    let mut last_activity = Instant::now();
68    let mut last_event_gc = Instant::now();
69
70    while !shutdown.load(Ordering::Relaxed) {
71        if !home.path().is_dir() {
72            // TempDir / store root removed — exit so orphan daemons do not pile up.
73            let _ = cleanup(&home);
74            return Ok(());
75        }
76
77        if http_activity.swap(false, Ordering::Relaxed) {
78            last_activity = Instant::now();
79        }
80
81        match events_listener.accept() {
82            Ok((stream, _)) => {
83                if let Err(e) = hub.add(stream) {
84                    eprintln!("event subscriber error: {e}");
85                }
86                last_activity = Instant::now();
87            }
88            Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {}
89            Err(e) => return Err(e.into()),
90        }
91
92        match listener.accept() {
93            Ok((stream, _)) => {
94                let home = home.clone();
95                let store = Arc::clone(&store);
96                let shutdown = Arc::clone(&shutdown);
97                let hub = Arc::clone(&hub);
98                if let Err(e) = handle_client(stream, &home, &store, &shutdown, &hub) {
99                    eprintln!("daemon client error: {e}");
100                }
101                last_activity = Instant::now();
102            }
103            Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {
104                let _ = crate::namespace::reap_dead(&home);
105                if last_event_gc.elapsed() >= Duration::from_secs(30) {
106                    if let Ok(mut store) = store.lock() {
107                        if let Err(e) = store.expire_events(&home) {
108                            eprintln!("event gc error: {e}");
109                        }
110                    }
111                    last_event_gc = Instant::now();
112                }
113                maybe_idle_flush(&home, &store, last_activity, idle_after, max_load);
114                if maybe_idle_stop(
115                    &home,
116                    &store,
117                    &hub,
118                    &shutdown,
119                    last_activity,
120                    idle_stop_after,
121                ) {
122                    break;
123                }
124                thread::sleep(Duration::from_millis(50));
125            }
126            Err(e) => return Err(e.into()),
127        }
128    }
129
130    cleanup(&home)?;
131    Ok(())
132}
133
134fn maybe_idle_flush(
135    home: &UnifierHome,
136    store: &Arc<std::sync::Mutex<HotStore>>,
137    last_activity: Instant,
138    idle_after: Duration,
139    max_load: f64,
140) {
141    if !should_idle_flush(
142        last_activity,
143        Instant::now(),
144        idle_after,
145        system_is_idle(max_load),
146        true,
147        false,
148    ) {
149        return;
150    }
151    let Ok(mut store) = store.lock() else {
152        return;
153    };
154    if store.active_tick.is_some() || !store.is_dirty() {
155        return;
156    }
157    if let Err(e) = store.flush(home) {
158        eprintln!("idle flush error: {e}");
159    }
160}
161
162/// Returns true when the daemon should exit after an idle stop.
163fn maybe_idle_stop(
164    home: &UnifierHome,
165    store: &Arc<std::sync::Mutex<HotStore>>,
166    hub: &Arc<EventHub>,
167    shutdown: &Arc<AtomicBool>,
168    last_activity: Instant,
169    idle_stop_after: Duration,
170) -> bool {
171    let tick_active = store
172        .lock()
173        .map(|s| s.active_tick.is_some())
174        .unwrap_or(true);
175    if !should_idle_stop(
176        last_activity,
177        Instant::now(),
178        idle_stop_after,
179        tick_active,
180        hub.subscriber_count(),
181    ) {
182        return false;
183    }
184    if let Ok(mut store) = store.lock() {
185        if store.is_dirty() {
186            if let Err(e) = store.flush(home) {
187                eprintln!("idle stop flush error: {e}");
188            }
189        }
190    }
191    shutdown.store(true, Ordering::Relaxed);
192    true
193}
194
195fn handle_client(
196    mut stream: UnixStream,
197    home: &UnifierHome,
198    store: &Arc<std::sync::Mutex<HotStore>>,
199    shutdown: &Arc<AtomicBool>,
200    hub: &Arc<EventHub>,
201) -> Result<()> {
202    let mut reader = BufReader::new(stream.try_clone()?);
203    let mut line = String::new();
204    reader.read_line(&mut line)?;
205    if line.trim().is_empty() {
206        return Ok(());
207    }
208
209    let request = decode_request(&line).map_err(|e| Error::msg(format!("invalid request: {e}")))?;
210    let response = match dispatch(home, store, shutdown, hub, request) {
211        Ok(resp) => resp,
212        Err(e) => Response::Err {
213            error: e.to_string(),
214        },
215    };
216    stream.write_all(encode_response(&response)?.as_bytes())?;
217    stream.flush()?;
218    Ok(())
219}
220
221fn dispatch(
222    home: &UnifierHome,
223    store: &Arc<std::sync::Mutex<HotStore>>,
224    shutdown: &Arc<AtomicBool>,
225    hub: &Arc<EventHub>,
226    request: Request,
227) -> Result<Response> {
228    match request {
229        Request::Ping => Ok(ok_empty()),
230        Request::Shutdown => {
231            shutdown.store(true, Ordering::Relaxed);
232            let mut store = store.lock().map_err(lock_err)?;
233            if store.is_dirty() && store.active_tick.is_none() {
234                store.flush(home)?;
235            } else if store.active_tick.is_some() {
236                return Err(Error::msg(
237                    "cannot shutdown with active tick; run tick end first",
238                ));
239            }
240            Ok(Response::Ok {
241                value: None,
242                uuid: None,
243                found: None,
244                dirty: Some(false),
245                tick: None,
246                queued: None,
247                messages: vec![],
248            })
249        }
250        Request::Flush => {
251            let mut store = store.lock().map_err(lock_err)?;
252            let was_dirty = store.is_dirty();
253            if store.active_tick.is_some() {
254                return Err(Error::msg("cannot flush while a tick is active"));
255            }
256            if was_dirty {
257                store.flush(home)?;
258            }
259            Ok(Response::Ok {
260                value: None,
261                uuid: None,
262                found: None,
263                dirty: Some(was_dirty),
264                tick: None,
265                queued: None,
266                messages: vec![],
267            })
268        }
269        Request::Put { key, value } => {
270            let mut store = store.lock().map_err(lock_err)?;
271            store.put_key(&key, &value)?;
272            Ok(ok_empty())
273        }
274        Request::Get { key } => {
275            let store = store.lock().map_err(lock_err)?;
276            match store.get_key(&key)? {
277                Some(value) => Ok(Response::Ok {
278                    value: Some(value),
279                    uuid: None,
280                    found: None,
281                    dirty: None,
282                    tick: None,
283                    queued: None,
284                    messages: vec![],
285                }),
286                None => Ok(Response::Err {
287                    error: format!("key not found: {key}"),
288                }),
289            }
290        }
291        Request::Del { key } => {
292            let mut store = store.lock().map_err(lock_err)?;
293            if store.delete_key(&key)? {
294                Ok(ok_empty())
295            } else {
296                Ok(Response::Err {
297                    error: format!("key not found: {key}"),
298                })
299            }
300        }
301        Request::Send {
302            from,
303            recipient,
304            message,
305        } => {
306            let from = from.unwrap_or_else(|| DEFAULT_SENDER.to_string());
307            let mut store = store.lock().map_err(lock_err)?;
308            let id = store.send_from(&from, &recipient, &message)?;
309            hub.broadcast(&Notice::mailbox(id, &from, &recipient));
310            Ok(Response::Ok {
311                value: None,
312                uuid: Some(id),
313                found: None,
314                dirty: None,
315                tick: None,
316                queued: None,
317                messages: vec![],
318            })
319        }
320        Request::Cron { schedule, message } => {
321            let mut store = store.lock().map_err(lock_err)?;
322            let id = store.post_cron(&schedule, &message)?;
323            Ok(Response::Ok {
324                value: None,
325                uuid: Some(id),
326                found: None,
327                dirty: None,
328                tick: None,
329                queued: None,
330                messages: vec![],
331            })
332        }
333        Request::Poll { recipient } => {
334            let store = store.lock().map_err(lock_err)?;
335            let messages = store.poll_mailbox(&recipient)?;
336            Ok(messages_response(messages))
337        }
338        Request::PollCron => {
339            let store = store.lock().map_err(lock_err)?;
340            let messages = store.poll_cron()?;
341            Ok(messages_response(messages))
342        }
343        Request::List { path } => {
344            let store = store.lock().map_err(lock_err)?;
345            let messages = store.list_dir(home, &path)?;
346            Ok(messages_response(messages))
347        }
348        Request::Ack { id_or_path } => {
349            let mut store = store.lock().map_err(lock_err)?;
350            let found = store.ack(home, &id_or_path)?;
351            Ok(Response::Ok {
352                value: None,
353                uuid: None,
354                found: Some(found),
355                dirty: None,
356                tick: None,
357                queued: None,
358                messages: vec![],
359            })
360        }
361        Request::TickStart { label } => {
362            let mut store = store.lock().map_err(lock_err)?;
363            match store.tick_start(&label)? {
364                TickStartOutcome::Started { tick } => Ok(Response::Ok {
365                    value: None,
366                    uuid: None,
367                    found: None,
368                    dirty: None,
369                    tick: Some(tick),
370                    queued: None,
371                    messages: vec![],
372                }),
373                TickStartOutcome::Queued { position, .. } => {
374                    eprintln!("tick queue: queued start {label:?} at position {position}");
375                    Ok(Response::Ok {
376                        value: None,
377                        uuid: None,
378                        found: None,
379                        dirty: None,
380                        tick: None,
381                        queued: Some(position),
382                        messages: vec![],
383                    })
384                }
385            }
386        }
387        Request::TickEnd => {
388            let mut store = store.lock().map_err(lock_err)?;
389            let tick = store.tick_end(home)?;
390            Ok(Response::Ok {
391                value: None,
392                uuid: None,
393                found: None,
394                dirty: None,
395                tick: Some(tick),
396                queued: None,
397                messages: vec![],
398            })
399        }
400        Request::TickStatus => {
401            let store = store.lock().map_err(lock_err)?;
402            let status = store.tick_status();
403            Ok(Response::Ok {
404                value: Some(format!(
405                    "committed={} active={:?} queued={} locks={:?}",
406                    status.committed_tick, status.active_tick, status.queued, status.locked_keys
407                )),
408                uuid: None,
409                found: None,
410                dirty: None,
411                tick: status.active_tick,
412                queued: Some(status.queued),
413                messages: vec![],
414            })
415        }
416        Request::TickLock { key } => {
417            let mut store = store.lock().map_err(lock_err)?;
418            store.tick_lock(&key)?;
419            Ok(ok_empty())
420        }
421        Request::TickUnlock { key } => {
422            let mut store = store.lock().map_err(lock_err)?;
423            let found = store.tick_unlock(&key)?;
424            Ok(Response::Ok {
425                value: None,
426                uuid: None,
427                found: Some(found),
428                dirty: None,
429                tick: None,
430                queued: None,
431                messages: vec![],
432            })
433        }
434        Request::Event { payload, ttl } => {
435            let mut store = store.lock().map_err(lock_err)?;
436            let id = store.post_event(&payload, ttl)?;
437            hub.broadcast(&Notice::event(id, event_name(&payload)));
438            Ok(Response::Ok {
439                value: None,
440                uuid: Some(id),
441                found: None,
442                dirty: None,
443                tick: None,
444                queued: None,
445                messages: vec![],
446            })
447        }
448        Request::AgentMessage { from, to, payload } => {
449            let mut store = store.lock().map_err(lock_err)?;
450            let id = store.send_agent_message(&from, &to, &payload)?;
451            hub.broadcast(&Notice::mailbox(id, &from, &to));
452            Ok(Response::Ok {
453                value: None,
454                uuid: Some(id),
455                found: None,
456                dirty: None,
457                tick: None,
458                queued: None,
459                messages: vec![],
460            })
461        }
462        Request::WebStatus => {
463            let value = match www::base_url(home) {
464                Some(url) => url,
465                None => return Err(Error::msg("web server is not listening")),
466            };
467            Ok(Response::Ok {
468                value: Some(value),
469                uuid: None,
470                found: None,
471                dirty: None,
472                tick: None,
473                queued: None,
474                messages: vec![],
475            })
476        }
477        Request::WebList => {
478            let entries = www::list(home)?;
479            let lines: Vec<String> = entries
480                .into_iter()
481                .map(|e| {
482                    let url = www::entry_url(home, &e.name).unwrap_or_default();
483                    format!("{}\t{}\t{}\t{}", e.name, e.content_type, e.bytes, url)
484                })
485                .collect();
486            Ok(Response::Ok {
487                value: Some(lines.join("\n")),
488                uuid: None,
489                found: None,
490                dirty: None,
491                tick: None,
492                queued: None,
493                messages: vec![],
494            })
495        }
496        Request::WebRm { name } => {
497            let found = www::remove(home, &name)?;
498            Ok(Response::Ok {
499                value: None,
500                uuid: None,
501                found: Some(found),
502                dirty: None,
503                tick: None,
504                queued: None,
505                messages: vec![],
506            })
507        }
508    }
509}
510
511fn messages_response(messages: Vec<crate::postbox::Message>) -> Response {
512    Response::Ok {
513        value: None,
514        uuid: None,
515        found: None,
516        dirty: None,
517        tick: None,
518        queued: None,
519        messages: messages.into_iter().map(MessageDto::from).collect(),
520    }
521}
522
523fn write_pid(home: &UnifierHome) -> Result<()> {
524    let pid = std::process::id();
525    fs::write(pid_path(home), format!("{pid}\n"))?;
526    Ok(())
527}
528
529fn cleanup(home: &UnifierHome) -> Result<()> {
530    let _ = fs::remove_file(pid_path(home));
531    let _ = fs::remove_file(socket_path(home));
532    let _ = fs::remove_file(events_socket_path(home));
533    let _ = fs::remove_file(crate::daemon::paths::http_port_path(home));
534    Ok(())
535}
536
537fn lock_err<E: std::fmt::Display>(e: E) -> Error {
538    Error::msg(format!("daemon store lock poisoned: {e}"))
539}