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::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);
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    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
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    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
485/// Reads one request line from a per-run worker socket, enforces the verb
486/// split at the boundary, and funnels the request through the dispatcher
487/// with the presented token for validation against the run's issued one.
488async 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
565/// Accepts connections on one run's worker socket until the settle path
566/// aborts this task. Each connection is one request, mirroring the
567/// operator socket.
568pub(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
594/// Opens and exclusively locks a daemon lock file, waiting out `LOCK_GRACE`
595/// for a stopping predecessor before concluding another daemon is running.
596fn 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    // The spawning client exits as soon as its request is answered, so the
637    // daemon must not stay in that client's session or process group: a signal
638    // delivered to either -- a terminal hangup, a supervisor reaping a finished
639    // command's group -- would take the daemon down with an unrelated
640    // short-lived process, which reads as an unexplained crash. `setsid` leaves
641    // both in one call, and it cannot fail here because a freshly forked child
642    // is never already a process group leader.
643    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/// Identity of the daemon owning a repository's lockfile: PID plus process start
708/// time, mirroring the identity rule used for supervised agents.
709#[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 {}