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