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