Skip to main content

sloop/daemon/
server.rs

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::{mpsc, oneshot};
19
20use crate::clock::{Clock, FileClock, SystemClock};
21use crate::config::{Config, ConfigError, Repository, TicketSourceConfig};
22use crate::frontmatter::FrontmatterError;
23use crate::ids::IdError;
24use crate::logging::{LogLevel, OperationalLog};
25use crate::protocol::{
26    Capability, ErrorBody, ErrorCode, Request, RequestEnvelope, RequestId, ResponseEnvelope,
27};
28use crate::run_ref::{RandomRunIds, RunIdSource};
29use crate::runner::local::{process_identity_matches, process_start_time};
30use crate::sources::TicketSource;
31use crate::sources::exec::ExecTicketSource;
32use crate::sources::markdown::MarkdownTicketSource;
33use crate::store::{Store, StoreError};
34use crate::vendor_error::{CatalogError, VendorErrorClassifier};
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::{index_projects, reconcile_tickets};
42
43const MAX_ENVELOPE_BYTES: u64 = 1024 * 1024;
44const STARTUP_TIMEOUT: Duration = Duration::from_secs(5);
45/// How long a starting daemon waits for a predecessor to drop its lock. A
46/// stopping daemon replies before its flock is released, so an immediate
47/// restart can arrive while the old process is still exiting. Shorter than
48/// `STARTUP_TIMEOUT` so a client that spawned this daemon still connects.
49const LOCK_GRACE: Duration = Duration::from_secs(2);
50const LOCK_POLL: Duration = Duration::from_millis(20);
51const CLIENT_TIMEOUT: Duration = Duration::from_secs(5);
52const DISPATCH_CHANNEL_CAPACITY: usize = 64;
53/// Activity-feed rows kept across daemon restarts; events are tiny, so this
54/// is weeks of history for a busy repository.
55const EVENT_RETENTION: i64 = 10_000;
56
57static NEXT_REQUEST_ID: AtomicU64 = AtomicU64::new(1);
58
59pub struct ClientResponse {
60    pub response: ResponseEnvelope,
61    pub started: bool,
62}
63
64pub fn request(request: Request) -> Result<ClientResponse, DaemonError> {
65    let cwd = std::env::current_dir().map_err(DaemonError::CurrentDirectory)?;
66    let repository = Repository::discover(&cwd)?;
67    Config::validate_client_essentials(&repository)?;
68
69    // Posting binds and snapshots a live flow definition. All other requests
70    // can use the configuration snapshot held by an existing daemon.
71    if matches!(&request, Request::Post(_)) {
72        Config::load(&repository)?;
73    }
74
75    if let Ok(response) = send_existing(&repository, request.clone()) {
76        return Ok(ClientResponse {
77            response,
78            started: false,
79        });
80    }
81
82    // Implicit startup has the same validation boundary as `sloop daemon`.
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
100/// Sends a request only if a daemon is already listening; `Ok(None)` means
101/// no daemon. Never spawns one.
102pub 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    // Hold the pre-v7 runtime lock as well during the lock-location
194    // transition, preventing an already-running older daemon in this runtime
195    // root from sharing the database with the new process.
196    let legacy_lock_path = repository.runtime_dir.join("daemon.lock");
197    let legacy_lock = acquire_daemon_lock(&legacy_lock_path)?;
198    // Identity is advisory; the flock is the guard, so write errors are
199    // ignored rather than fatal.
200    let identity = json!({
201        "pid": std::process::id(),
202        "started_at_ms": process_start_time(std::process::id()),
203        "socket": repository.operator_socket,
204    });
205    let _ = lock.set_len(0);
206    let _ = {
207        use std::io::Write as _;
208        writeln!(&lock, "{identity}")
209    };
210
211    let clock: Arc<dyn Clock> = match std::env::var_os("SLOOP_TEST_CLOCK_PATH") {
212        Some(path) => Arc::new(FileClock::new(path.into())),
213        None => Arc::new(SystemClock),
214    };
215    let store = Store::open(&repository.db_path, clock.now_ms()).map_err(DaemonError::Store)?;
216    // A fresh process satisfies any persisted restart intent, including after
217    // a crash during a drain or a successful self-exec.
218    store
219        .clear_restart_draining(clock.now_ms())
220        .map_err(DaemonError::Store)?;
221    // Bound the activity feed once per daemon lifetime; watchers page by
222    // sequence, so trimming old rows never invalidates a held cursor.
223    store
224        .trim_events(EVENT_RETENTION)
225        .map_err(DaemonError::Store)?;
226    if let Some(agent) = &config.agent {
227        store
228            .backfill_ticket_targets(&agent.default_target, clock.now_ms())
229            .map_err(DaemonError::Store)?;
230    }
231    let _ = index_projects(
232        &repository.root,
233        &config.project_dir,
234        &store,
235        clock.now_ms(),
236        &config.project_prefix,
237    )?;
238    reconcile_tickets(
239        &repository.root,
240        &store,
241        clock.now_ms(),
242        config.delete_missing_after_ms,
243    )?;
244
245    let runtime = tokio::runtime::Builder::new_multi_thread()
246        .enable_all()
247        .build()
248        .map_err(DaemonError::Runtime)?;
249    let root = repository.root.clone();
250    let daemon_log = repository.daemon_log.clone();
251    let control = runtime.block_on(serve(
252        repository,
253        config,
254        store,
255        lock,
256        legacy_lock,
257        clock,
258        classifier,
259    ))?;
260    drop(runtime);
261    Ok(match control {
262        DaemonControl::Stop => ServeExit::Stopped,
263        DaemonControl::Restart => ServeExit::Restart { root, daemon_log },
264    })
265}
266
267async fn serve(
268    repository: Repository,
269    config: Config,
270    store: Store,
271    _lock: fs::File,
272    _legacy_lock: fs::File,
273    clock: Arc<dyn Clock>,
274    classifier: Arc<VendorErrorClassifier>,
275) -> Result<DaemonControl, DaemonError> {
276    if repository.operator_socket.exists() {
277        fs::remove_file(&repository.operator_socket).map_err(|source| DaemonError::Io {
278            path: repository.operator_socket.clone(),
279            source,
280        })?;
281    }
282
283    let listener =
284        UnixListener::bind(&repository.operator_socket).map_err(|source| DaemonError::Io {
285            path: repository.operator_socket.clone(),
286            source,
287        })?;
288    fs::set_permissions(
289        &repository.operator_socket,
290        fs::Permissions::from_mode(0o600),
291    )
292    .map_err(|source| DaemonError::Io {
293        path: repository.operator_socket.clone(),
294        source,
295    })?;
296
297    let log = OperationalLog::open(&repository.daemon_log).map_err(|source| DaemonError::Io {
298        path: repository.daemon_log.clone(),
299        source,
300    })?;
301    log.emit(LogLevel::Info, "sloop::daemon", "daemon_started");
302
303    let paused = store.paused().map_err(DaemonError::Store)?;
304    let (dispatcher_tx, dispatcher_rx) = mpsc::channel(DISPATCH_CHANNEL_CAPACITY);
305    let (events_tx, events_rx) = mpsc::channel(DISPATCH_CHANNEL_CAPACITY);
306    let (shutdown_tx, mut shutdown_rx) = mpsc::channel::<DaemonControl>(1);
307    let shutdown_flag = Arc::new(AtomicBool::new(false));
308    let ticket_source: Arc<dyn TicketSource> = match &config.ticket_source {
309        TicketSourceConfig::Markdown => Arc::new(MarkdownTicketSource::new(
310            &repository.root,
311            &config.ticket_dir,
312        )),
313        TicketSourceConfig::Exec(argv) => {
314            Arc::new(ExecTicketSource::new(&repository.root, argv.clone()))
315        }
316    };
317    let mut state = DispatcherState {
318        pid: std::process::id(),
319        paused,
320        draining: false,
321        restart_acknowledged: false,
322        restart_signalled: false,
323        max_agents: config.max_parallel_tasks,
324        ticket_prefix: config.ticket_prefix.clone(),
325        project_prefix: config.project_prefix.clone(),
326        running_hours: config.running_hours.clone(),
327        agent: config.agent.clone(),
328        flows: config.flows.clone(),
329        default_flow: config.default_flow.clone(),
330        aftercare_test_cmd: config.aftercare_test_cmd.clone(),
331        root: repository.root.clone(),
332        project_dir: config.project_dir.clone(),
333        ticket_dir: config.ticket_dir.clone(),
334        ticket_source,
335        worktree_dir: repository.root.join(&config.worktree_dir),
336        worktree_retention_ms: config.worktree_retention_ms,
337        state_dir: repository.state_dir.clone(),
338        runtime_dir: repository.runtime_dir.clone(),
339        socket: repository.operator_socket.clone(),
340        daemon_log: repository.daemon_log.clone(),
341        store,
342        storage_full: Cell::new(false),
343        reconciliation_blocked: false,
344        active: HashSet::new(),
345        supervised: HashSet::new(),
346        suspected_dead: HashSet::new(),
347        recovering: HashSet::new(),
348        cancelling: HashSet::new(),
349        worker_tokens: HashMap::new(),
350        worker_listeners: HashMap::new(),
351        worker_socket_paths: HashMap::new(),
352        pending_exits: HashMap::new(),
353        requests_tx: dispatcher_tx.clone(),
354        log: log.clone(),
355        clock,
356        run_ids: Arc::new(RandomRunIds) as Arc<dyn RunIdSource>,
357        classifier,
358        shutdown: shutdown_tx.clone(),
359        shutdown_flag: shutdown_flag.clone(),
360    };
361    recover_inflight_runs(&mut state, &events_tx, &log)?;
362    let dispatcher_task = tokio::spawn(run_dispatcher(
363        state,
364        dispatcher_rx,
365        events_rx,
366        events_tx,
367        log.clone(),
368    ));
369
370    loop {
371        tokio::select! {
372            accepted = listener.accept() => {
373                let (stream, _) = accepted.map_err(|source| DaemonError::Io {
374                    path: repository.operator_socket.clone(),
375                    source,
376                })?;
377                let dispatcher_tx = dispatcher_tx.clone();
378                let log = log.clone();
379                let shutdown = shutdown_tx.clone();
380                tokio::spawn(async move {
381                    if let Err(error) = handle_connection(stream, dispatcher_tx, shutdown).await {
382                        log.emit_with_fields(
383                            LogLevel::Error,
384                            "sloop::socket",
385                            "connection_failed",
386                            json!({"error": error.to_string()}),
387                        );
388                    }
389                });
390            }
391            control = shutdown_rx.recv() => {
392                let control = control.unwrap_or(DaemonControl::Stop);
393                shutdown_flag.store(true, Ordering::Release);
394                if control == DaemonControl::Stop {
395                    log.emit(LogLevel::Info, "sloop::daemon", "daemon_stopped");
396                }
397                let _ = fs::remove_file(&repository.operator_socket);
398                dispatcher_task.abort();
399                let _ = dispatcher_task.await;
400                return Ok(control);
401            }
402        }
403    }
404}
405
406async fn handle_connection(
407    stream: UnixStream,
408    dispatcher: mpsc::Sender<DispatcherMessage>,
409    shutdown: mpsc::Sender<DaemonControl>,
410) -> io::Result<()> {
411    let reader = AsyncBufReader::new(stream);
412    let mut limited = reader.take(MAX_ENVELOPE_BYTES + 1);
413    let mut bytes = Vec::new();
414    let read = limited.read_until(b'\n', &mut bytes).await?;
415    if read == 0 {
416        return Ok(());
417    }
418
419    let mut stream = limited.into_inner().into_inner();
420    let envelope = if bytes.len() as u64 > MAX_ENVELOPE_BYTES {
421        Err(protocol_error("request envelope is too large"))
422    } else {
423        std::str::from_utf8(&bytes)
424            .map_err(|_| protocol_error("request envelope must be UTF-8"))
425            .and_then(|line| RequestEnvelope::decode(line.trim_end()).map_err(|error| error.body))
426    };
427
428    let is_stop = matches!(
429        &envelope,
430        Ok(envelope) if matches!(envelope.request, Request::Stop(_))
431    );
432    let is_restart = matches!(
433        &envelope,
434        Ok(envelope) if matches!(envelope.request, Request::Restart(_))
435    );
436    let response = match envelope {
437        Ok(envelope) if envelope.token.is_some() => ResponseEnvelope::failure(
438            Some(envelope.id),
439            unauthorized("operator socket does not accept worker tokens"),
440        ),
441        Ok(envelope)
442            if !matches!(
443                envelope.request.capability(),
444                Capability::Operator | Capability::Both
445            ) =>
446        {
447            ResponseEnvelope::failure(
448                Some(envelope.id),
449                unauthorized(
450                    "worker verbs are not available on the operator socket; \
451                     run `sloop list` or `sloop show <ticket>` to inspect tickets from here",
452                ),
453            )
454        }
455        Ok(envelope) => dispatch_envelope(envelope, RequestOrigin::Operator, &dispatcher).await,
456        Err(error) => ResponseEnvelope::failure(None, error),
457    };
458
459    // The reply must be flushed before the daemon exits, so the connection
460    // handler owns the shutdown signal for an accepted stop.
461    let stopping = is_stop && response.ok;
462    let encoded = serde_json::to_vec(&response).map_err(io::Error::other)?;
463    stream.write_all(&encoded).await?;
464    stream.write_all(b"\n").await?;
465    stream.shutdown().await?;
466    if stopping {
467        let _ = shutdown.send(DaemonControl::Stop).await;
468    } else if is_restart && response.ok {
469        let _ = dispatcher
470            .send(DispatcherMessage::RestartAcknowledged)
471            .await;
472    }
473    Ok(())
474}
475
476/// Reads one request line from a per-run worker socket, enforces the verb
477/// split at the boundary, and funnels the request through the dispatcher
478/// with the presented token for validation against the run's issued one.
479async fn handle_worker_connection(
480    stream: UnixStream,
481    run_id: String,
482    dispatcher: mpsc::Sender<DispatcherMessage>,
483) -> io::Result<()> {
484    let reader = AsyncBufReader::new(stream);
485    let mut limited = reader.take(MAX_ENVELOPE_BYTES + 1);
486    let mut bytes = Vec::new();
487    let read = limited.read_until(b'\n', &mut bytes).await?;
488    if read == 0 {
489        return Ok(());
490    }
491
492    let mut stream = limited.into_inner().into_inner();
493    let envelope = if bytes.len() as u64 > MAX_ENVELOPE_BYTES {
494        Err(protocol_error("request envelope is too large"))
495    } else {
496        std::str::from_utf8(&bytes)
497            .map_err(|_| protocol_error("request envelope must be UTF-8"))
498            .and_then(|line| RequestEnvelope::decode(line.trim_end()).map_err(|error| error.body))
499    };
500
501    let response = match envelope {
502        Ok(envelope)
503            if !matches!(
504                envelope.request.capability(),
505                Capability::Worker | Capability::Both
506            ) =>
507        {
508            ResponseEnvelope::failure(
509                Some(envelope.id),
510                unauthorized("operator verbs are not available on a worker socket"),
511            )
512        }
513        Ok(envelope) => {
514            let token = envelope.token.clone();
515            dispatch_envelope(
516                envelope,
517                RequestOrigin::Worker { run_id, token },
518                &dispatcher,
519            )
520            .await
521        }
522        Err(error) => ResponseEnvelope::failure(None, error),
523    };
524
525    let encoded = serde_json::to_vec(&response).map_err(io::Error::other)?;
526    stream.write_all(&encoded).await?;
527    stream.write_all(b"\n").await?;
528    stream.shutdown().await
529}
530
531async fn dispatch_envelope(
532    envelope: RequestEnvelope,
533    origin: RequestOrigin,
534    dispatcher: &mpsc::Sender<DispatcherMessage>,
535) -> ResponseEnvelope {
536    let (reply_tx, reply_rx) = oneshot::channel();
537    let id = envelope.id;
538    if dispatcher
539        .send(DispatcherMessage::Request {
540            id: id.clone(),
541            request: envelope.request,
542            origin,
543            reply: reply_tx,
544        })
545        .await
546        .is_err()
547    {
548        ResponseEnvelope::failure(Some(id), internal("dispatcher is unavailable"))
549    } else {
550        reply_rx.await.unwrap_or_else(|_| {
551            ResponseEnvelope::failure(Some(id), internal("dispatcher dropped response"))
552        })
553    }
554}
555
556/// Accepts connections on one run's worker socket until the settle path
557/// aborts this task. Each connection is one request, mirroring the
558/// operator socket.
559pub(super) async fn serve_worker_socket(
560    listener: UnixListener,
561    run_id: String,
562    dispatcher: mpsc::Sender<DispatcherMessage>,
563    log: OperationalLog,
564) {
565    loop {
566        let Ok((stream, _)) = listener.accept().await else {
567            return;
568        };
569        let run_id = run_id.clone();
570        let dispatcher = dispatcher.clone();
571        let log = log.clone();
572        tokio::spawn(async move {
573            if let Err(error) = handle_worker_connection(stream, run_id.clone(), dispatcher).await {
574                log.emit_with_fields(
575                    LogLevel::Error,
576                    "sloop::socket",
577                    "worker_connection_failed",
578                    json!({"run_id": run_id, "error": error.to_string()}),
579                );
580            }
581        });
582    }
583}
584
585/// Opens and exclusively locks a daemon lock file, waiting out `LOCK_GRACE`
586/// for a stopping predecessor before concluding another daemon is running.
587fn acquire_daemon_lock(path: &Path) -> Result<File, DaemonError> {
588    let lock = OpenOptions::new()
589        .create(true)
590        .truncate(false)
591        .read(true)
592        .write(true)
593        .open(path)
594        .map_err(|source| DaemonError::Io {
595            path: path.to_path_buf(),
596            source,
597        })?;
598    let deadline = Instant::now() + LOCK_GRACE;
599    loop {
600        match lock.try_lock_exclusive() {
601            Ok(()) => return Ok(lock),
602            Err(source) if source.kind() == io::ErrorKind::WouldBlock => {
603                if Instant::now() >= deadline {
604                    return Err(DaemonError::AlreadyRunning);
605                }
606                std::thread::sleep(LOCK_POLL);
607            }
608            Err(source) => {
609                return Err(DaemonError::Io {
610                    path: path.to_path_buf(),
611                    source,
612                });
613            }
614        }
615    }
616}
617
618fn spawn_daemon(repository: &Repository) -> Result<(), DaemonError> {
619    let executable = std::env::current_exe().map_err(DaemonError::CurrentExecutable)?;
620    Command::new(executable)
621        .args(["daemon", "--foreground"])
622        .current_dir(&repository.root)
623        .stdin(Stdio::null())
624        .stdout(Stdio::null())
625        .stderr(Stdio::null())
626        .spawn()
627        .map(|_| ())
628        .map_err(DaemonError::Spawn)
629}
630
631fn send_existing(
632    repository: &Repository,
633    request: Request,
634) -> Result<ResponseEnvelope, DaemonError> {
635    match send(&repository.operator_socket, request.clone()) {
636        Ok(response) => Ok(response),
637        Err(current_error) => {
638            let Some(identity) = read_lock_identity(&repository.lock_path) else {
639                return Err(current_error);
640            };
641            if !process_identity_matches(identity.pid, identity.started_at_ms) {
642                return Err(current_error);
643            }
644            let Some(socket) = identity.socket else {
645                return Err(current_error);
646            };
647            if socket == repository.operator_socket {
648                return Err(current_error);
649            }
650            send(&socket, request)
651        }
652    }
653}
654
655fn send(socket: &Path, request: Request) -> Result<ResponseEnvelope, DaemonError> {
656    let mut stream = StdUnixStream::connect(socket).map_err(DaemonError::Connect)?;
657    stream
658        .set_read_timeout(Some(CLIENT_TIMEOUT))
659        .map_err(DaemonError::Connect)?;
660    stream
661        .set_write_timeout(Some(CLIENT_TIMEOUT))
662        .map_err(DaemonError::Connect)?;
663
664    let sequence = NEXT_REQUEST_ID.fetch_add(1, Ordering::Relaxed);
665    let envelope = RequestEnvelope::new(
666        RequestId::new(format!("req-{}-{sequence}", std::process::id())),
667        request,
668        None,
669    );
670    serde_json::to_writer(&mut stream, &envelope).map_err(DaemonError::Encode)?;
671    stream.write_all(b"\n").map_err(DaemonError::Write)?;
672
673    let mut line = String::new();
674    let mut reader = BufReader::new(stream).take(MAX_ENVELOPE_BYTES + 1);
675    reader.read_line(&mut line).map_err(DaemonError::Read)?;
676    if line.len() as u64 > MAX_ENVELOPE_BYTES {
677        return Err(DaemonError::InvalidResponse(
678            "response envelope is too large".into(),
679        ));
680    }
681    serde_json::from_str(line.trim_end()).map_err(DaemonError::Decode)
682}
683
684/// Identity of the daemon owning a repository's lockfile: PID plus process start
685/// time, mirroring the identity rule used for supervised agents.
686#[derive(Debug, Clone, PartialEq, Eq)]
687pub struct LockIdentity {
688    pub pid: u32,
689    pub started_at_ms: Option<i64>,
690    pub socket: Option<PathBuf>,
691}
692
693pub fn read_lock_identity(path: &Path) -> Option<LockIdentity> {
694    let content = fs::read_to_string(path).ok()?;
695    let value: serde_json::Value = serde_json::from_str(content.trim()).ok()?;
696    Some(LockIdentity {
697        pid: u32::try_from(value["pid"].as_u64()?).ok()?,
698        started_at_ms: value["started_at_ms"].as_i64(),
699        socket: value["socket"].as_str().map(PathBuf::from),
700    })
701}
702
703#[derive(Debug)]
704pub enum DaemonError {
705    Config(ConfigError),
706    Catalog(CatalogError),
707    Store(StoreError),
708    CurrentDirectory(io::Error),
709    CurrentExecutable(io::Error),
710    Io {
711        path: PathBuf,
712        source: io::Error,
713    },
714    AlreadyRunning,
715    Runtime(io::Error),
716    Spawn(io::Error),
717    Connect(io::Error),
718    Write(io::Error),
719    Read(io::Error),
720    Encode(serde_json::Error),
721    Decode(serde_json::Error),
722    InvalidResponse(String),
723    Frontmatter {
724        path: PathBuf,
725        error: FrontmatterError,
726    },
727    IdAllocation(IdError),
728}
729
730impl DaemonError {
731    pub fn error_body(&self) -> ErrorBody {
732        let code = match self {
733            Self::Config(_) => ErrorCode::InvalidArguments,
734            _ => ErrorCode::DaemonUnavailable,
735        };
736        ErrorBody {
737            code,
738            message: self.to_string(),
739            details: json!({}),
740        }
741    }
742}
743
744impl From<ConfigError> for DaemonError {
745    fn from(error: ConfigError) -> Self {
746        Self::Config(error)
747    }
748}
749
750impl From<IdError> for DaemonError {
751    fn from(error: IdError) -> Self {
752        Self::IdAllocation(error)
753    }
754}
755
756impl std::fmt::Display for DaemonError {
757    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
758        match self {
759            Self::Config(error) => error.fmt(formatter),
760            Self::Catalog(error) => error.fmt(formatter),
761            Self::Store(error) => error.fmt(formatter),
762            Self::CurrentDirectory(error) => {
763                write!(formatter, "cannot read current directory: {error}")
764            }
765            Self::CurrentExecutable(error) => {
766                write!(formatter, "cannot locate sloop executable: {error}")
767            }
768            Self::Io { path, source } => write!(formatter, "{}: {source}", path.display()),
769            Self::AlreadyRunning => formatter.write_str("another sloop daemon holds the lock"),
770            Self::Runtime(error) => write!(formatter, "cannot start async runtime: {error}"),
771            Self::Spawn(error) => write!(formatter, "cannot spawn daemon: {error}"),
772            Self::Connect(error) => write!(formatter, "cannot connect to daemon: {error}"),
773            Self::Write(error) => write!(formatter, "cannot write daemon request: {error}"),
774            Self::Read(error) => write!(formatter, "cannot read daemon response: {error}"),
775            Self::Encode(error) => write!(formatter, "cannot encode daemon request: {error}"),
776            Self::Decode(error) => write!(formatter, "cannot decode daemon response: {error}"),
777            Self::InvalidResponse(message) => formatter.write_str(message),
778            Self::Frontmatter { path, error } => write!(formatter, "{}: {error}", path.display()),
779            Self::IdAllocation(error) => error.fmt(formatter),
780        }
781    }
782}
783
784impl std::error::Error for DaemonError {}