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::{idle_after, max_load, should_idle_flush, system_is_idle};
12
13use crate::daemon::notify::{event_name, EventHub, Notice};
14use crate::daemon::paths::{daemon_dir, events_socket_path, pid_path, socket_path};
15use crate::daemon::protocol::{
16    decode_request, encode_response, ok_empty, MessageDto, Request, Response,
17};
18use crate::envelope::DEFAULT_SENDER;
19use crate::error::{Error, Result};
20use crate::home::UnifierHome;
21use crate::store::HotStore;
22use crate::tick::TickStartOutcome;
23
24pub fn run(home: UnifierHome) -> Result<()> {
25    let shutdown = Arc::new(AtomicBool::new(false));
26    home.ensure()?;
27    fs::create_dir_all(daemon_dir(&home))?;
28
29    let events_sock = events_socket_path(&home);
30    if events_sock.exists() {
31        fs::remove_file(&events_sock)?;
32    }
33    let events_listener = UnixListener::bind(&events_sock)?;
34    events_listener.set_nonblocking(true)?;
35    let hub = Arc::new(EventHub::default());
36
37    let sock = socket_path(&home);
38    if sock.exists() {
39        fs::remove_file(&sock)?;
40    }
41    let listener = UnixListener::bind(&sock)?;
42    listener.set_nonblocking(true)?;
43
44    let store = Arc::new(std::sync::Mutex::new(HotStore::load(&home)?));
45    write_pid(&home)?;
46
47    let idle_after = idle_after();
48    let max_load = max_load();
49    let mut last_activity = Instant::now();
50
51    while !shutdown.load(Ordering::Relaxed) {
52        match events_listener.accept() {
53            Ok((stream, _)) => {
54                if let Err(e) = hub.add(stream) {
55                    eprintln!("event subscriber error: {e}");
56                }
57            }
58            Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {}
59            Err(e) => return Err(e.into()),
60        }
61
62        match listener.accept() {
63            Ok((stream, _)) => {
64                let home = home.clone();
65                let store = Arc::clone(&store);
66                let shutdown = Arc::clone(&shutdown);
67                let hub = Arc::clone(&hub);
68                if let Err(e) = handle_client(stream, &home, &store, &shutdown, &hub) {
69                    eprintln!("daemon client error: {e}");
70                }
71                last_activity = Instant::now();
72            }
73            Err(e) if e.kind() == std::io::ErrorKind::WouldBlock => {
74                let _ = crate::namespace::reap_dead(&home);
75                maybe_idle_flush(&home, &store, last_activity, idle_after, max_load);
76                thread::sleep(Duration::from_millis(50));
77            }
78            Err(e) => return Err(e.into()),
79        }
80    }
81
82    cleanup(&home)?;
83    Ok(())
84}
85
86fn maybe_idle_flush(
87    home: &UnifierHome,
88    store: &Arc<std::sync::Mutex<HotStore>>,
89    last_activity: Instant,
90    idle_after: Duration,
91    max_load: f64,
92) {
93    if !should_idle_flush(
94        last_activity,
95        Instant::now(),
96        idle_after,
97        system_is_idle(max_load),
98        true,
99        false,
100    ) {
101        return;
102    }
103    let Ok(mut store) = store.lock() else {
104        return;
105    };
106    if store.active_tick.is_some() || !store.is_dirty() {
107        return;
108    }
109    if let Err(e) = store.flush(home) {
110        eprintln!("idle flush error: {e}");
111    }
112}
113
114fn handle_client(
115    mut stream: UnixStream,
116    home: &UnifierHome,
117    store: &Arc<std::sync::Mutex<HotStore>>,
118    shutdown: &Arc<AtomicBool>,
119    hub: &Arc<EventHub>,
120) -> Result<()> {
121    let mut reader = BufReader::new(stream.try_clone()?);
122    let mut line = String::new();
123    reader.read_line(&mut line)?;
124    if line.trim().is_empty() {
125        return Ok(());
126    }
127
128    let request = decode_request(&line).map_err(|e| Error::msg(format!("invalid request: {e}")))?;
129    let response = match dispatch(home, store, shutdown, hub, request) {
130        Ok(resp) => resp,
131        Err(e) => Response::Err {
132            error: e.to_string(),
133        },
134    };
135    stream.write_all(encode_response(&response)?.as_bytes())?;
136    stream.flush()?;
137    Ok(())
138}
139
140fn dispatch(
141    home: &UnifierHome,
142    store: &Arc<std::sync::Mutex<HotStore>>,
143    shutdown: &Arc<AtomicBool>,
144    hub: &Arc<EventHub>,
145    request: Request,
146) -> Result<Response> {
147    match request {
148        Request::Ping => Ok(ok_empty()),
149        Request::Shutdown => {
150            shutdown.store(true, Ordering::Relaxed);
151            let mut store = store.lock().map_err(lock_err)?;
152            if store.is_dirty() && store.active_tick.is_none() {
153                store.flush(home)?;
154            } else if store.active_tick.is_some() {
155                return Err(Error::msg(
156                    "cannot shutdown with active tick; run tick end first",
157                ));
158            }
159            Ok(Response::Ok {
160                value: None,
161                uuid: None,
162                found: None,
163                dirty: Some(false),
164                tick: None,
165                queued: None,
166                messages: vec![],
167            })
168        }
169        Request::Flush => {
170            let mut store = store.lock().map_err(lock_err)?;
171            let was_dirty = store.is_dirty();
172            if store.active_tick.is_some() {
173                return Err(Error::msg("cannot flush while a tick is active"));
174            }
175            if was_dirty {
176                store.flush(home)?;
177            }
178            Ok(Response::Ok {
179                value: None,
180                uuid: None,
181                found: None,
182                dirty: Some(was_dirty),
183                tick: None,
184                queued: None,
185                messages: vec![],
186            })
187        }
188        Request::Put { key, value } => {
189            let mut store = store.lock().map_err(lock_err)?;
190            store.put_key(&key, &value)?;
191            Ok(ok_empty())
192        }
193        Request::Get { key } => {
194            let store = store.lock().map_err(lock_err)?;
195            match store.get_key(&key)? {
196                Some(value) => Ok(Response::Ok {
197                    value: Some(value),
198                    uuid: None,
199                    found: None,
200                    dirty: None,
201                    tick: None,
202                    queued: None,
203                    messages: vec![],
204                }),
205                None => Ok(Response::Err {
206                    error: format!("key not found: {key}"),
207                }),
208            }
209        }
210        Request::Del { key } => {
211            let mut store = store.lock().map_err(lock_err)?;
212            if store.delete_key(&key)? {
213                Ok(ok_empty())
214            } else {
215                Ok(Response::Err {
216                    error: format!("key not found: {key}"),
217                })
218            }
219        }
220        Request::Send {
221            from,
222            recipient,
223            message,
224        } => {
225            let from = from.unwrap_or_else(|| DEFAULT_SENDER.to_string());
226            let mut store = store.lock().map_err(lock_err)?;
227            let id = store.send_from(&from, &recipient, &message)?;
228            hub.broadcast(&Notice::mailbox(id, &from, &recipient));
229            Ok(Response::Ok {
230                value: None,
231                uuid: Some(id),
232                found: None,
233                dirty: None,
234                tick: None,
235                queued: None,
236                messages: vec![],
237            })
238        }
239        Request::Cron { schedule, message } => {
240            let mut store = store.lock().map_err(lock_err)?;
241            let id = store.post_cron(&schedule, &message)?;
242            Ok(Response::Ok {
243                value: None,
244                uuid: Some(id),
245                found: None,
246                dirty: None,
247                tick: None,
248                queued: None,
249                messages: vec![],
250            })
251        }
252        Request::Poll { recipient } => {
253            let store = store.lock().map_err(lock_err)?;
254            let messages = store.poll_mailbox(&recipient)?;
255            Ok(messages_response(messages))
256        }
257        Request::PollCron => {
258            let store = store.lock().map_err(lock_err)?;
259            let messages = store.poll_cron()?;
260            Ok(messages_response(messages))
261        }
262        Request::List { path } => {
263            let store = store.lock().map_err(lock_err)?;
264            let messages = store.list_dir(home, &path)?;
265            Ok(messages_response(messages))
266        }
267        Request::Ack { id_or_path } => {
268            let mut store = store.lock().map_err(lock_err)?;
269            let found = store.ack(home, &id_or_path)?;
270            Ok(Response::Ok {
271                value: None,
272                uuid: None,
273                found: Some(found),
274                dirty: None,
275                tick: None,
276                queued: None,
277                messages: vec![],
278            })
279        }
280        Request::TickStart { label } => {
281            let mut store = store.lock().map_err(lock_err)?;
282            match store.tick_start(&label)? {
283                TickStartOutcome::Started { tick } => Ok(Response::Ok {
284                    value: None,
285                    uuid: None,
286                    found: None,
287                    dirty: None,
288                    tick: Some(tick),
289                    queued: None,
290                    messages: vec![],
291                }),
292                TickStartOutcome::Queued { position, .. } => {
293                    eprintln!("tick queue: queued start {label:?} at position {position}");
294                    Ok(Response::Ok {
295                        value: None,
296                        uuid: None,
297                        found: None,
298                        dirty: None,
299                        tick: None,
300                        queued: Some(position),
301                        messages: vec![],
302                    })
303                }
304            }
305        }
306        Request::TickEnd => {
307            let mut store = store.lock().map_err(lock_err)?;
308            let tick = store.tick_end(home)?;
309            Ok(Response::Ok {
310                value: None,
311                uuid: None,
312                found: None,
313                dirty: None,
314                tick: Some(tick),
315                queued: None,
316                messages: vec![],
317            })
318        }
319        Request::TickStatus => {
320            let store = store.lock().map_err(lock_err)?;
321            let status = store.tick_status();
322            Ok(Response::Ok {
323                value: Some(format!(
324                    "committed={} active={:?} queued={} locks={:?}",
325                    status.committed_tick, status.active_tick, status.queued, status.locked_keys
326                )),
327                uuid: None,
328                found: None,
329                dirty: None,
330                tick: status.active_tick,
331                queued: Some(status.queued),
332                messages: vec![],
333            })
334        }
335        Request::TickLock { key } => {
336            let mut store = store.lock().map_err(lock_err)?;
337            store.tick_lock(&key)?;
338            Ok(ok_empty())
339        }
340        Request::TickUnlock { key } => {
341            let mut store = store.lock().map_err(lock_err)?;
342            let found = store.tick_unlock(&key)?;
343            Ok(Response::Ok {
344                value: None,
345                uuid: None,
346                found: Some(found),
347                dirty: None,
348                tick: None,
349                queued: None,
350                messages: vec![],
351            })
352        }
353        Request::Event { payload } => {
354            let mut store = store.lock().map_err(lock_err)?;
355            let id = store.post_event(&payload)?;
356            hub.broadcast(&Notice::event(id, event_name(&payload)));
357            Ok(Response::Ok {
358                value: None,
359                uuid: Some(id),
360                found: None,
361                dirty: None,
362                tick: None,
363                queued: None,
364                messages: vec![],
365            })
366        }
367        Request::AgentMessage { from, to, payload } => {
368            let mut store = store.lock().map_err(lock_err)?;
369            let id = store.send_agent_message(&from, &to, &payload)?;
370            hub.broadcast(&Notice::mailbox(id, &from, &to));
371            Ok(Response::Ok {
372                value: None,
373                uuid: Some(id),
374                found: None,
375                dirty: None,
376                tick: None,
377                queued: None,
378                messages: vec![],
379            })
380        }
381    }
382}
383
384fn messages_response(messages: Vec<crate::postbox::Message>) -> Response {
385    Response::Ok {
386        value: None,
387        uuid: None,
388        found: None,
389        dirty: None,
390        tick: None,
391        queued: None,
392        messages: messages.into_iter().map(MessageDto::from).collect(),
393    }
394}
395
396fn write_pid(home: &UnifierHome) -> Result<()> {
397    let pid = std::process::id();
398    fs::write(pid_path(home), format!("{pid}\n"))?;
399    Ok(())
400}
401
402fn cleanup(home: &UnifierHome) -> Result<()> {
403    let _ = fs::remove_file(pid_path(home));
404    let _ = fs::remove_file(socket_path(home));
405    let _ = fs::remove_file(events_socket_path(home));
406    Ok(())
407}
408
409fn lock_err<E: std::fmt::Display>(e: E) -> Error {
410    Error::msg(format!("daemon store lock poisoned: {e}"))
411}