1use std::cell::Cell;
2use std::collections::{HashMap, HashSet};
3use std::fs::{self, File, OpenOptions};
4use std::io::{self, BufRead, BufReader, Read, Write};
5use std::os::unix::fs::PermissionsExt;
6use std::os::unix::net::UnixStream as StdUnixStream;
7use std::os::unix::process::CommandExt;
8use std::path::{Path, PathBuf};
9use std::process::{Command, Stdio};
10use std::sync::Arc;
11use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
12use std::time::{Duration, Instant};
13
14use fs2::FileExt;
15use serde_json::json;
16use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader as AsyncBufReader};
17use tokio::net::{UnixListener, UnixStream};
18use tokio::sync::{Notify, mpsc, oneshot};
19
20use crate::clock::{Clock, FileClock, SystemClock};
21use crate::config::{Config, ConfigError, Repository, TicketSourceConfig};
22use crate::db::{Db, StoreError};
23use crate::frontmatter::FrontmatterError;
24use crate::ids::IdError;
25use crate::protocol::{
26 Capability, ErrorBody, ErrorCode, Request, RequestEnvelope, RequestId, ResponseEnvelope,
27};
28use crate::run_ref::{RandomRunIds, RunIdSource};
29use crate::run_store::RunStore;
30use crate::runner::local::{process_identity_matches, process_start_time};
31use crate::vendor::{CatalogError, VendorErrorClassifier};
32use crate::work_state::TicketFeeder;
33use crate::work_state::local::LocalSqlite;
34
35use super::dispatcher::{
36 DaemonControl, DispatcherMessage, DispatcherState, RequestOrigin, internal, protocol_error,
37 run_dispatcher, unauthorized,
38};
39use super::logging::{LogLevel, OperationalLog};
40use super::recovery::recover_inflight_runs;
41use super::scheduler::{
42 index_projects, reconcile_merged_ticket_triggers, reconcile_tickets,
43 restore_reported_output_stalls,
44};
45
46const MAX_ENVELOPE_BYTES: u64 = 1024 * 1024;
47const STARTUP_TIMEOUT: Duration = Duration::from_secs(5);
48const LOCK_GRACE: Duration = Duration::from_secs(2);
53const LOCK_POLL: Duration = Duration::from_millis(20);
54const CLIENT_TIMEOUT: Duration = Duration::from_secs(5);
55const DISPATCH_CHANNEL_CAPACITY: usize = 64;
56const EVENT_RETENTION: i64 = 10_000;
59
60static NEXT_REQUEST_ID: AtomicU64 = AtomicU64::new(1);
61
62pub struct ClientResponse {
63 pub response: ResponseEnvelope,
64 pub started: bool,
65}
66
67pub fn request(request: Request) -> Result<ClientResponse, DaemonError> {
68 let cwd = std::env::current_dir().map_err(DaemonError::CurrentDirectory)?;
69 let repository = Repository::discover(&cwd)?;
70 Config::validate_client_essentials(&repository)?;
71
72 if matches!(&request, Request::Post(_)) {
73 Config::load(&repository)?;
74 }
75
76 if let Ok(response) = send_existing(&repository, request.clone()) {
77 return Ok(ClientResponse {
78 response,
79 started: false,
80 });
81 }
82
83 Config::load(&repository)?;
84 spawn_daemon(&repository)?;
85 let deadline = Instant::now() + STARTUP_TIMEOUT;
86 loop {
87 match send_existing(&repository, request.clone()) {
88 Ok(response) => {
89 return Ok(ClientResponse {
90 response,
91 started: true,
92 });
93 }
94 Err(error) if Instant::now() >= deadline => return Err(error),
95 Err(_) => std::thread::sleep(Duration::from_millis(20)),
96 }
97 }
98}
99
100pub fn request_running(request: Request) -> Result<Option<ResponseEnvelope>, DaemonError> {
103 let cwd = std::env::current_dir().map_err(DaemonError::CurrentDirectory)?;
104 let repository = Repository::discover(&cwd)?;
105 Config::validate_client_essentials(&repository)?;
106 match send_existing(&repository, request) {
107 Ok(response) => Ok(Some(response)),
108 Err(DaemonError::Connect(_)) => Ok(None),
109 Err(error) => Err(error),
110 }
111}
112
113pub fn serve_current_repository() -> Result<(), DaemonError> {
114 let executable = std::env::current_exe().map_err(DaemonError::CurrentExecutable)?;
115 loop {
116 match serve_current_repository_once()? {
117 ServeExit::Stopped => return Ok(()),
118 ServeExit::Restart { root, daemon_log } => {
119 if let Ok(log) = OperationalLog::open(&daemon_log) {
120 log.emit(LogLevel::Info, "sloop::daemon", "restart_exec");
121 }
122 let error = Command::new(&executable)
123 .args(["daemon", "--foreground"])
124 .current_dir(root)
125 .exec();
126 if let Ok(log) = OperationalLog::open(&daemon_log) {
127 log.emit_with_fields(
128 LogLevel::Error,
129 "sloop::daemon",
130 "restart_exec_failed",
131 json!({"path": executable, "error": error.to_string()}),
132 );
133 }
134 }
135 }
136 }
137}
138
139enum ServeExit {
140 Stopped,
141 Restart { root: PathBuf, daemon_log: PathBuf },
142}
143
144fn serve_current_repository_once() -> Result<ServeExit, DaemonError> {
145 let cwd = std::env::current_dir().map_err(DaemonError::CurrentDirectory)?;
146 let repository = Repository::discover(&cwd)?;
147 let config = Config::load(&repository)?;
148 let classifier = Arc::new(VendorErrorClassifier::built_in().map_err(DaemonError::Catalog)?);
149 fs::create_dir_all(&repository.state_dir).map_err(|source| DaemonError::Io {
150 path: repository.state_dir.clone(),
151 source,
152 })?;
153 fs::set_permissions(&repository.state_dir, fs::Permissions::from_mode(0o700)).map_err(
154 |source| DaemonError::Io {
155 path: repository.state_dir.clone(),
156 source,
157 },
158 )?;
159 let runtime_root = repository
160 .runtime_dir
161 .parent()
162 .expect("repository runtime directories have a parent");
163 fs::create_dir_all(runtime_root).map_err(|source| DaemonError::Io {
164 path: runtime_root.to_path_buf(),
165 source,
166 })?;
167 fs::set_permissions(runtime_root, fs::Permissions::from_mode(0o700)).map_err(|source| {
168 DaemonError::Io {
169 path: runtime_root.to_path_buf(),
170 source,
171 }
172 })?;
173 fs::create_dir(&repository.runtime_dir)
174 .or_else(|source| {
175 if source.kind() == io::ErrorKind::AlreadyExists {
176 Ok(())
177 } else {
178 Err(source)
179 }
180 })
181 .map_err(|source| DaemonError::Io {
182 path: repository.runtime_dir.clone(),
183 source,
184 })?;
185 fs::set_permissions(&repository.runtime_dir, fs::Permissions::from_mode(0o700)).map_err(
186 |source| DaemonError::Io {
187 path: repository.runtime_dir.clone(),
188 source,
189 },
190 )?;
191
192 let lock = acquire_daemon_lock(&repository.lock_path)?;
193 let legacy_lock_path = repository.runtime_dir.join("daemon.lock");
194 let legacy_lock = acquire_daemon_lock(&legacy_lock_path)?;
195 let identity = json!({
196 "pid": std::process::id(),
197 "started_at_ms": process_start_time(std::process::id()),
198 "socket": repository.operator_socket,
199 });
200 let _ = lock.set_len(0);
201 let _ = {
202 use std::io::Write as _;
203 writeln!(&lock, "{identity}")
204 };
205
206 let clock: Arc<dyn Clock> = match std::env::var_os("SLOOP_TEST_CLOCK_PATH") {
207 Some(path) => Arc::new(FileClock::new(path.into())),
208 None => Arc::new(SystemClock),
209 };
210 let db = Db::open(&repository.db_path, clock.now_ms())
211 .map_err(StoreError::from)
212 .map_err(DaemonError::Store)?;
213 let local_work_state = LocalSqlite::from_db(db.clone());
214 let run_store = RunStore::from_db(db.clone());
215 run_store
216 .clear_restart_draining(clock.now_ms())
217 .map_err(DaemonError::Store)?;
218 run_store
219 .trim_events(EVENT_RETENTION)
220 .map_err(DaemonError::Store)?;
221 if let Some(agent) = &config.agent {
222 local_work_state
223 .backfill_ticket_targets(&agent.default_target, clock.now_ms())
224 .map_err(DaemonError::Store)?;
225 }
226 let _ = index_projects(
227 &repository.root,
228 &config.project_dir,
229 &local_work_state,
230 clock.now_ms(),
231 &config.project_prefix,
232 )?;
233 reconcile_tickets(
234 &repository.root,
235 &local_work_state,
236 &run_store,
237 clock.now_ms(),
238 config.delete_missing_after_ms,
239 )?;
240
241 let runtime = tokio::runtime::Builder::new_multi_thread()
242 .enable_all()
243 .build()
244 .map_err(DaemonError::Runtime)?;
245 let root = repository.root.clone();
246 let daemon_log = repository.daemon_log.clone();
247 let control = runtime.block_on(serve(
248 repository,
249 config,
250 db,
251 lock,
252 legacy_lock,
253 clock,
254 classifier,
255 ))?;
256 drop(runtime);
257 Ok(match control {
258 DaemonControl::Stop => ServeExit::Stopped,
259 DaemonControl::Restart => ServeExit::Restart { root, daemon_log },
260 })
261}
262
263async fn serve(
264 repository: Repository,
265 config: Config,
266 db: Db,
267 _lock: fs::File,
268 _legacy_lock: fs::File,
269 clock: Arc<dyn Clock>,
270 classifier: Arc<VendorErrorClassifier>,
271) -> Result<DaemonControl, DaemonError> {
272 if repository.operator_socket.exists() {
273 fs::remove_file(&repository.operator_socket).map_err(|source| DaemonError::Io {
274 path: repository.operator_socket.clone(),
275 source,
276 })?;
277 }
278
279 let listener =
280 UnixListener::bind(&repository.operator_socket).map_err(|source| DaemonError::Io {
281 path: repository.operator_socket.clone(),
282 source,
283 })?;
284 fs::set_permissions(
285 &repository.operator_socket,
286 fs::Permissions::from_mode(0o600),
287 )
288 .map_err(|source| DaemonError::Io {
289 path: repository.operator_socket.clone(),
290 source,
291 })?;
292
293 let log = OperationalLog::open(&repository.daemon_log).map_err(|source| DaemonError::Io {
294 path: repository.daemon_log.clone(),
295 source,
296 })?;
297 log.emit(LogLevel::Info, "sloop::daemon", "daemon_started");
298
299 let run_store = RunStore::from_db(db.clone());
300 let paused = run_store.paused().map_err(DaemonError::Store)?;
301 let (dispatcher_tx, dispatcher_rx) = mpsc::channel(DISPATCH_CHANNEL_CAPACITY);
302 let (events_tx, events_rx) = mpsc::channel(DISPATCH_CHANNEL_CAPACITY);
303 let (shutdown_tx, mut shutdown_rx) = mpsc::channel::<DaemonControl>(1);
304 let shutdown_flag = Arc::new(AtomicBool::new(false));
305 let ticket_source = match &config.ticket_source {
306 TicketSourceConfig::Markdown => {
307 TicketFeeder::markdown(&repository.root, &config.ticket_dir)
308 }
309 TicketSourceConfig::Exec(argv) => TicketFeeder::exec(&repository.root, argv.clone()),
310 };
311 let work_state_author_enabled = ticket_source.supports_authoring();
312 let local_work_state = LocalSqlite::from_db(db.clone());
313 let work_state = Arc::new(LocalSqlite::from_db_with_clock_and_reporter(
314 db,
315 clock.clone(),
316 ticket_source.exec_reporter(),
317 )) as Arc<dyn crate::work_state::WorkState>;
318 let mut state = DispatcherState {
319 pid: std::process::id(),
320 paused,
321 draining: false,
322 restart_acknowledged: false,
323 restart_signalled: false,
324 max_agents: config.max_parallel_tasks,
325 stall_report_after_ms: config.stall_report_after_ms,
326 stall_after_ms: config.stall_after_ms,
327 ticket_prefix: config.ticket_prefix.clone(),
328 project_prefix: config.project_prefix.clone(),
329 running_hours: config.running_hours.clone(),
330 agent: config.agent.clone(),
331 flows: config.flows.clone(),
332 default_flow: config.default_flow.clone(),
333 flow_test_cmd: config.flow_test_cmd.clone(),
334 root: repository.root.clone(),
335 project_dir: config.project_dir.clone(),
336 ticket_dir: config.ticket_dir.clone(),
337 ticket_source,
338 work_state_author_enabled,
339 worktree_dir: repository.root.join(&config.worktree_dir),
340 worktree_retention_ms: config.worktree_retention_ms,
341 state_dir: repository.state_dir.clone(),
342 runtime_dir: repository.runtime_dir.clone(),
343 socket: repository.operator_socket.clone(),
344 daemon_log: repository.daemon_log.clone(),
345 local_work_state,
346 work_state,
347 run_store,
348 storage_full: Cell::new(false),
349 reconciliation_blocked: false,
350 active: HashSet::new(),
351 supervised: HashSet::new(),
352 suspected_dead: HashSet::new(),
353 recovering: HashSet::new(),
354 cancelling: HashSet::new(),
355 stalling: HashSet::new(),
356 worker_tokens: HashMap::new(),
357 worker_listeners: HashMap::new(),
358 worker_socket_paths: HashMap::new(),
359 pending_exits: HashMap::new(),
360 reported_stalls: HashMap::new(),
361 output_notify: Arc::new(Notify::new()),
362 requests_tx: dispatcher_tx.clone(),
363 log: log.clone(),
364 clock,
365 run_ids: Arc::new(RandomRunIds) as Arc<dyn RunIdSource>,
366 classifier,
367 shutdown: shutdown_tx.clone(),
368 shutdown_flag: shutdown_flag.clone(),
369 };
370 restore_reported_output_stalls(&mut state);
371 reconcile_merged_ticket_triggers(&state, &log);
372 recover_inflight_runs(&mut state, &events_tx, &log).await?;
373 let dispatcher_task = tokio::spawn(run_dispatcher(
374 state,
375 dispatcher_rx,
376 events_rx,
377 events_tx,
378 log.clone(),
379 ));
380
381 loop {
382 tokio::select! {
383 accepted = listener.accept() => {
384 let (stream, _) = accepted.map_err(|source| DaemonError::Io {
385 path: repository.operator_socket.clone(),
386 source,
387 })?;
388 let dispatcher_tx = dispatcher_tx.clone();
389 let log = log.clone();
390 let shutdown = shutdown_tx.clone();
391 tokio::spawn(async move {
392 if let Err(error) = handle_connection(stream, dispatcher_tx, shutdown).await {
393 log.emit_with_fields(
394 LogLevel::Error,
395 "sloop::socket",
396 "connection_failed",
397 json!({"error": error.to_string()}),
398 );
399 }
400 });
401 }
402 control = shutdown_rx.recv() => {
403 let control = control.unwrap_or(DaemonControl::Stop);
404 shutdown_flag.store(true, Ordering::Release);
405 if control == DaemonControl::Stop {
406 log.emit(LogLevel::Info, "sloop::daemon", "daemon_stopped");
407 }
408 let _ = fs::remove_file(&repository.operator_socket);
409 dispatcher_task.abort();
410 let _ = dispatcher_task.await;
411 return Ok(control);
412 }
413 }
414 }
415}
416
417async fn handle_connection(
418 stream: UnixStream,
419 dispatcher: mpsc::Sender<DispatcherMessage>,
420 shutdown: mpsc::Sender<DaemonControl>,
421) -> io::Result<()> {
422 let reader = AsyncBufReader::new(stream);
423 let mut limited = reader.take(MAX_ENVELOPE_BYTES + 1);
424 let mut bytes = Vec::new();
425 let read = limited.read_until(b'\n', &mut bytes).await?;
426 if read == 0 {
427 return Ok(());
428 }
429
430 let mut stream = limited.into_inner().into_inner();
431 let envelope = if bytes.len() as u64 > MAX_ENVELOPE_BYTES {
432 Err(protocol_error("request envelope is too large"))
433 } else {
434 std::str::from_utf8(&bytes)
435 .map_err(|_| protocol_error("request envelope must be UTF-8"))
436 .and_then(|line| RequestEnvelope::decode(line.trim_end()).map_err(|error| error.body))
437 };
438
439 let is_stop = matches!(
440 &envelope,
441 Ok(envelope) if matches!(envelope.request, Request::Stop(_))
442 );
443 let is_restart = matches!(
444 &envelope,
445 Ok(envelope) if matches!(envelope.request, Request::Restart(_))
446 );
447 let response = match envelope {
448 Ok(envelope) if envelope.token.is_some() => ResponseEnvelope::failure(
449 Some(envelope.id),
450 unauthorized("operator socket does not accept worker tokens"),
451 ),
452 Ok(envelope)
453 if !matches!(
454 envelope.request.capability(),
455 Capability::Operator | Capability::Both
456 ) =>
457 {
458 ResponseEnvelope::failure(
459 Some(envelope.id),
460 unauthorized(
461 "worker verbs are not available on the operator socket; \
462 run `sloop` or `sloop show <ticket>` to inspect tickets from here",
463 ),
464 )
465 }
466 Ok(envelope) => dispatch_envelope(envelope, RequestOrigin::Operator, &dispatcher).await,
467 Err(error) => ResponseEnvelope::failure(None, error),
468 };
469
470 let stopping = is_stop && response.ok;
471 let encoded = serde_json::to_vec(&response).map_err(io::Error::other)?;
472 stream.write_all(&encoded).await?;
473 stream.write_all(b"\n").await?;
474 stream.shutdown().await?;
475 if stopping {
476 let _ = shutdown.send(DaemonControl::Stop).await;
477 } else if is_restart && response.ok {
478 let _ = dispatcher
479 .send(DispatcherMessage::RestartAcknowledged)
480 .await;
481 }
482 Ok(())
483}
484
485async fn handle_worker_connection(
489 stream: UnixStream,
490 run_id: String,
491 dispatcher: mpsc::Sender<DispatcherMessage>,
492) -> io::Result<()> {
493 let reader = AsyncBufReader::new(stream);
494 let mut limited = reader.take(MAX_ENVELOPE_BYTES + 1);
495 let mut bytes = Vec::new();
496 let read = limited.read_until(b'\n', &mut bytes).await?;
497 if read == 0 {
498 return Ok(());
499 }
500
501 let mut stream = limited.into_inner().into_inner();
502 let envelope = if bytes.len() as u64 > MAX_ENVELOPE_BYTES {
503 Err(protocol_error("request envelope is too large"))
504 } else {
505 std::str::from_utf8(&bytes)
506 .map_err(|_| protocol_error("request envelope must be UTF-8"))
507 .and_then(|line| RequestEnvelope::decode(line.trim_end()).map_err(|error| error.body))
508 };
509
510 let response = match envelope {
511 Ok(envelope)
512 if !matches!(
513 envelope.request.capability(),
514 Capability::Worker | Capability::Both
515 ) =>
516 {
517 ResponseEnvelope::failure(
518 Some(envelope.id),
519 unauthorized("operator verbs are not available on a worker socket"),
520 )
521 }
522 Ok(envelope) => {
523 let token = envelope.token.clone();
524 dispatch_envelope(
525 envelope,
526 RequestOrigin::Worker { run_id, token },
527 &dispatcher,
528 )
529 .await
530 }
531 Err(error) => ResponseEnvelope::failure(None, error),
532 };
533
534 let encoded = serde_json::to_vec(&response).map_err(io::Error::other)?;
535 stream.write_all(&encoded).await?;
536 stream.write_all(b"\n").await?;
537 stream.shutdown().await
538}
539
540async fn dispatch_envelope(
541 envelope: RequestEnvelope,
542 origin: RequestOrigin,
543 dispatcher: &mpsc::Sender<DispatcherMessage>,
544) -> ResponseEnvelope {
545 let (reply_tx, reply_rx) = oneshot::channel();
546 let id = envelope.id;
547 if dispatcher
548 .send(DispatcherMessage::Request {
549 id: id.clone(),
550 request: envelope.request,
551 origin,
552 reply: reply_tx,
553 })
554 .await
555 .is_err()
556 {
557 ResponseEnvelope::failure(Some(id), internal("dispatcher is unavailable"))
558 } else {
559 reply_rx.await.unwrap_or_else(|_| {
560 ResponseEnvelope::failure(Some(id), internal("dispatcher dropped response"))
561 })
562 }
563}
564
565pub(super) async fn serve_worker_socket(
569 listener: UnixListener,
570 run_id: String,
571 dispatcher: mpsc::Sender<DispatcherMessage>,
572 log: OperationalLog,
573) {
574 loop {
575 let Ok((stream, _)) = listener.accept().await else {
576 return;
577 };
578 let run_id = run_id.clone();
579 let dispatcher = dispatcher.clone();
580 let log = log.clone();
581 tokio::spawn(async move {
582 if let Err(error) = handle_worker_connection(stream, run_id.clone(), dispatcher).await {
583 log.emit_with_fields(
584 LogLevel::Error,
585 "sloop::socket",
586 "worker_connection_failed",
587 json!({"run_id": run_id, "error": error.to_string()}),
588 );
589 }
590 });
591 }
592}
593
594fn acquire_daemon_lock(path: &Path) -> Result<File, DaemonError> {
597 let lock = OpenOptions::new()
598 .create(true)
599 .truncate(false)
600 .read(true)
601 .write(true)
602 .open(path)
603 .map_err(|source| DaemonError::Io {
604 path: path.to_path_buf(),
605 source,
606 })?;
607 let deadline = Instant::now() + LOCK_GRACE;
608 loop {
609 match lock.try_lock_exclusive() {
610 Ok(()) => return Ok(lock),
611 Err(source) if source.kind() == io::ErrorKind::WouldBlock => {
612 if Instant::now() >= deadline {
613 return Err(DaemonError::AlreadyRunning);
614 }
615 std::thread::sleep(LOCK_POLL);
616 }
617 Err(source) => {
618 return Err(DaemonError::Io {
619 path: path.to_path_buf(),
620 source,
621 });
622 }
623 }
624 }
625}
626
627fn spawn_daemon(repository: &Repository) -> Result<(), DaemonError> {
628 let executable = std::env::current_exe().map_err(DaemonError::CurrentExecutable)?;
629 let mut command = Command::new(executable);
630 command
631 .args(["daemon", "--foreground"])
632 .current_dir(&repository.root)
633 .stdin(Stdio::null())
634 .stdout(Stdio::null())
635 .stderr(Stdio::null());
636 unsafe {
644 command.pre_exec(|| {
645 if libc::setsid() == -1 {
646 return Err(io::Error::last_os_error());
647 }
648 Ok(())
649 });
650 }
651 command.spawn().map(|_| ()).map_err(DaemonError::Spawn)
652}
653
654fn send_existing(
655 repository: &Repository,
656 request: Request,
657) -> Result<ResponseEnvelope, DaemonError> {
658 match send(&repository.operator_socket, request.clone()) {
659 Ok(response) => Ok(response),
660 Err(current_error) => {
661 let Some(identity) = read_lock_identity(&repository.lock_path) else {
662 return Err(current_error);
663 };
664 if !process_identity_matches(identity.pid, identity.started_at_ms) {
665 return Err(current_error);
666 }
667 let Some(socket) = identity.socket else {
668 return Err(current_error);
669 };
670 if socket == repository.operator_socket {
671 return Err(current_error);
672 }
673 send(&socket, request)
674 }
675 }
676}
677
678fn send(socket: &Path, request: Request) -> Result<ResponseEnvelope, DaemonError> {
679 let mut stream = StdUnixStream::connect(socket).map_err(DaemonError::Connect)?;
680 stream
681 .set_read_timeout(Some(CLIENT_TIMEOUT))
682 .map_err(DaemonError::Connect)?;
683 stream
684 .set_write_timeout(Some(CLIENT_TIMEOUT))
685 .map_err(DaemonError::Connect)?;
686
687 let sequence = NEXT_REQUEST_ID.fetch_add(1, Ordering::Relaxed);
688 let envelope = RequestEnvelope::new(
689 RequestId::new(format!("req-{}-{sequence}", std::process::id())),
690 request,
691 None,
692 );
693 serde_json::to_writer(&mut stream, &envelope).map_err(DaemonError::Encode)?;
694 stream.write_all(b"\n").map_err(DaemonError::Write)?;
695
696 let mut line = String::new();
697 let mut reader = BufReader::new(stream).take(MAX_ENVELOPE_BYTES + 1);
698 reader.read_line(&mut line).map_err(DaemonError::Read)?;
699 if line.len() as u64 > MAX_ENVELOPE_BYTES {
700 return Err(DaemonError::InvalidResponse(
701 "response envelope is too large".into(),
702 ));
703 }
704 serde_json::from_str(line.trim_end()).map_err(DaemonError::Decode)
705}
706
707#[derive(Debug, Clone, PartialEq, Eq)]
710pub struct LockIdentity {
711 pub pid: u32,
712 pub started_at_ms: Option<i64>,
713 pub socket: Option<PathBuf>,
714}
715
716pub fn read_lock_identity(path: &Path) -> Option<LockIdentity> {
717 let content = fs::read_to_string(path).ok()?;
718 let value: serde_json::Value = serde_json::from_str(content.trim()).ok()?;
719 Some(LockIdentity {
720 pid: u32::try_from(value["pid"].as_u64()?).ok()?,
721 started_at_ms: value["started_at_ms"].as_i64(),
722 socket: value["socket"].as_str().map(PathBuf::from),
723 })
724}
725
726#[derive(Debug)]
727pub enum DaemonError {
728 Config(ConfigError),
729 Catalog(CatalogError),
730 Store(StoreError),
731 WorkState(crate::work_state::SourceError),
732 CurrentDirectory(io::Error),
733 CurrentExecutable(io::Error),
734 Io {
735 path: PathBuf,
736 source: io::Error,
737 },
738 AlreadyRunning,
739 Runtime(io::Error),
740 Spawn(io::Error),
741 Connect(io::Error),
742 Write(io::Error),
743 Read(io::Error),
744 Encode(serde_json::Error),
745 Decode(serde_json::Error),
746 InvalidResponse(String),
747 Frontmatter {
748 path: PathBuf,
749 error: FrontmatterError,
750 },
751 IdAllocation(IdError),
752}
753
754impl DaemonError {
755 pub fn error_body(&self) -> ErrorBody {
756 let code = match self {
757 Self::Config(_) => ErrorCode::InvalidArguments,
758 _ => ErrorCode::DaemonUnavailable,
759 };
760 ErrorBody {
761 code,
762 message: self.to_string(),
763 details: json!({}),
764 }
765 }
766}
767
768impl From<ConfigError> for DaemonError {
769 fn from(error: ConfigError) -> Self {
770 Self::Config(error)
771 }
772}
773
774impl From<IdError> for DaemonError {
775 fn from(error: IdError) -> Self {
776 Self::IdAllocation(error)
777 }
778}
779
780impl std::fmt::Display for DaemonError {
781 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
782 match self {
783 Self::Config(error) => error.fmt(formatter),
784 Self::Catalog(error) => error.fmt(formatter),
785 Self::Store(error) => error.fmt(formatter),
786 Self::WorkState(error) => write!(formatter, "work source error: {error:?}"),
787 Self::CurrentDirectory(error) => {
788 write!(formatter, "cannot read current directory: {error}")
789 }
790 Self::CurrentExecutable(error) => {
791 write!(formatter, "cannot locate sloop executable: {error}")
792 }
793 Self::Io { path, source } => write!(formatter, "{}: {source}", path.display()),
794 Self::AlreadyRunning => formatter.write_str("another sloop daemon holds the lock"),
795 Self::Runtime(error) => write!(formatter, "cannot start async runtime: {error}"),
796 Self::Spawn(error) => write!(formatter, "cannot spawn daemon: {error}"),
797 Self::Connect(error) => write!(formatter, "cannot connect to daemon: {error}"),
798 Self::Write(error) => write!(formatter, "cannot write daemon request: {error}"),
799 Self::Read(error) => write!(formatter, "cannot read daemon response: {error}"),
800 Self::Encode(error) => write!(formatter, "cannot encode daemon request: {error}"),
801 Self::Decode(error) => write!(formatter, "cannot decode daemon response: {error}"),
802 Self::InvalidResponse(message) => formatter.write_str(message),
803 Self::Frontmatter { path, error } => write!(formatter, "{}: {error}", path.display()),
804 Self::IdAllocation(error) => error.fmt(formatter),
805 }
806 }
807}
808
809impl std::error::Error for DaemonError {}