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;
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}