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::{
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 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
162fn 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}