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::{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);
48/// How long a starting daemon waits for a predecessor to drop its lock. A
49/// stopping daemon replies before its flock is released, so an immediate
50/// restart can arrive while the old process is still exiting. Shorter than
51/// `STARTUP_TIMEOUT` so a client that spawned this daemon still connects.
52const 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;
56/// Activity-feed rows kept across daemon restarts; events are tiny, so this
57/// is weeks of history for a busy repository.
58const 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    // Posting binds and snapshots a live flow definition. All other requests
73    // can use the configuration snapshot held by an existing daemon.
74    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    // Implicit startup has the same validation boundary as `sloop daemon`.
86    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
103/// Sends a request only if a daemon is already listening; `Ok(None)` means
104/// no daemon. Never spawns one.
105pub 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    // Hold the pre-v7 runtime lock as well during the lock-location
197    // transition, preventing an already-running older daemon in this runtime
198    // root from sharing the database with the new process.
199    let legacy_lock_path = repository.runtime_dir.join("daemon.lock");
200    let legacy_lock = acquire_daemon_lock(&legacy_lock_path)?;
201    // Identity is advisory; the flock is the guard, so write errors are
202    // ignored rather than fatal.
203    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    // A fresh process satisfies any persisted restart intent, including after
224    // a crash during a drain or a successful self-exec.
225    run_store
226        .clear_restart_draining(clock.now_ms())
227        .map_err(DaemonError::Store)?;
228    // Bound the activity feed once per daemon lifetime; watchers page by
229    // sequence, so trimming old rows never invalidates a held cursor.
230    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    // The reply must be flushed before the daemon exits, so the connection
483    // handler owns the shutdown signal for an accepted stop.
484    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
499/// Reads one request line from a per-run worker socket, enforces the verb
500/// split at the boundary, and funnels the request through the dispatcher
501/// with the presented token for validation against the run's issued one.
502async 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
579/// Accepts connections on one run's worker socket until the settle path
580/// aborts this task. Each connection is one request, mirroring the
581/// operator socket.
582pub(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
608/// Opens and exclusively locks a daemon lock file, waiting out `LOCK_GRACE`
609/// for a stopping predecessor before concluding another daemon is running.
610fn 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    // The spawning client exits as soon as its request is answered, so the
651    // daemon must not stay in that client's session or process group: a signal
652    // delivered to either -- a terminal hangup, a supervisor reaping a finished
653    // command's group -- would take the daemon down with an unrelated
654    // short-lived process, which reads as an unexplained crash. `setsid` leaves
655    // both in one call, and it cannot fail here because a freshly forked child
656    // is never already a process group leader.
657    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/// Identity of the daemon owning a repository's lockfile: PID plus process start
722/// time, mirroring the identity rule used for supervised agents.
723#[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 {}