1use 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}