Skip to main content

stackless_daemon/
server.rs

1//! The resident daemon (§3): unix-socket RPC + the reverse proxy.
2//! Spun up on demand by the CLI; same binary, `daemon run` subcommand.
3
4use std::path::PathBuf;
5use std::sync::Arc;
6
7use stackless_core::paths::Paths;
8use stackless_core::process::ProcessStamp;
9use stackless_core::types::{ProtocolVersion, TcpPort};
10use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
11use tokio::net::{UnixListener, UnixStream};
12
13use crate::proxy;
14
15use crate::rpc::{Envelope, Request, Response, ResponseBody, build_version};
16use crate::state::DaemonState;
17
18/// Whether this daemon is the operator process or an embedded/test instance.
19#[derive(Debug, Clone, Copy, PartialEq, Eq)]
20pub enum DaemonRole {
21    /// Operator daemon: register launchd and run the lease reaper.
22    Operator,
23    /// Embedded/test daemon: skip launchd registration and the reaper.
24    Embedded,
25}
26
27pub fn socket_path() -> PathBuf {
28    socket_path_for(&Paths::from_env())
29}
30
31pub fn socket_path_for(paths: &Paths) -> PathBuf {
32    paths.socket_path()
33}
34
35/// Run the daemon until told to shut down. Returns once drained.
36pub async fn run() -> std::io::Result<()> {
37    run_with(
38        &Paths::from_env(),
39        proxy::proxy_port(),
40        DaemonRole::Operator,
41    )
42    .await
43}
44
45/// Like [`run`], but binds the socket under an injectable state layout,
46/// listens on an injectable proxy port, and selects operator vs embedded
47/// behavior via [`DaemonRole`].
48///
49/// Operator mode requires [`crate::mark_cli_process`]: launchd registration
50/// and the lease reaper shell out via `current_exe`, which must be the CLI.
51pub async fn run_with(paths: &Paths, proxy_port: TcpPort, role: DaemonRole) -> std::io::Result<()> {
52    if role == DaemonRole::Operator && !crate::is_cli_process() {
53        return Err(std::io::Error::new(
54            std::io::ErrorKind::PermissionDenied,
55            "operator daemon requires the stackless CLI process \
56             (mark_cli_process); use DaemonRole::Embedded for in-process tests \
57             or spawn via the resolved CLI binary",
58        ));
59    }
60    let path = socket_path_for(paths);
61    if let Some(dir) = path.parent() {
62        std::fs::create_dir_all(dir)?;
63    }
64    // A live daemon answers on the socket; a dead one leaves a stale
65    // file behind. Probe before stealing the path.
66    if UnixStream::connect(&path).await.is_ok() {
67        return Err(std::io::Error::new(
68            std::io::ErrorKind::AddrInUse,
69            "another stackless daemon is already serving this socket",
70        ));
71    }
72    let _ = std::fs::remove_file(&path);
73    let listener = UnixListener::bind(&path)?;
74
75    // Boot persistence (§3): register as a launchd user agent so leases
76    // survive reboots/crashes. Refusal degrades loudly, never aborts.
77    // Skip for embedded test daemons.
78    if role == DaemonRole::Operator {
79        crate::launchd::ensure_registered(paths);
80    }
81
82    let state = Arc::new(DaemonState::default());
83
84    // Re-adopt before serving (§3: upgrade = restart + re-adopt). Routes
85    // and supervision live only in memory, so they died with the prior
86    // daemon — rebuild them from the journal before the proxy or socket
87    // can field a request, so the first proxied call already routes.
88    let summary = crate::adopt::readopt(&state, paths);
89    if !summary.adopted.is_empty() || !summary.dead.is_empty() {
90        eprintln!(
91            "stackless daemon: re-adopted {} live process(es), noted {} dead",
92            summary.adopted.len(),
93            summary.dead.len()
94        );
95    }
96
97    let proxy_state = state.clone();
98    tokio::spawn(async move {
99        if let Err(err) = proxy::serve(proxy_state, proxy_port).await {
100            eprintln!(
101                "stackless daemon: proxy failed to bind port {}: {err}",
102                proxy_port.get()
103            );
104        }
105    });
106
107    // The reaper (§6): one immediate pass reaps leases overdue while the
108    // daemon was down (start/wake), then a tick every minute. Operator only.
109    if role == DaemonRole::Operator {
110        let reaper_paths = paths.clone();
111        let reaper_port = proxy_port;
112        tokio::spawn(async move {
113            crate::reaper::tick_once(&reaper_paths, reaper_port).await;
114            crate::reaper::run(reaper_paths, reaper_port).await;
115        });
116    }
117
118    let (shutdown_tx, mut shutdown_rx) = tokio::sync::mpsc::channel::<()>(1);
119    loop {
120        tokio::select! {
121            accepted = listener.accept() => {
122                let Ok((stream, _)) = accepted else { continue };
123                let state = state.clone();
124                let shutdown = shutdown_tx.clone();
125                tokio::spawn(async move {
126                    let _ = handle_connection(stream, state, shutdown).await;
127                });
128            }
129            _ = shutdown_rx.recv() => break,
130        }
131    }
132    let _ = std::fs::remove_file(&path);
133    Ok(())
134}
135
136async fn handle_connection(
137    stream: UnixStream,
138    state: Arc<DaemonState>,
139    shutdown: tokio::sync::mpsc::Sender<()>,
140) -> std::io::Result<()> {
141    let (read_half, mut write_half) = stream.into_split();
142    let mut lines = BufReader::new(read_half).lines();
143    while let Some(line) = lines.next_line().await? {
144        if line.trim().is_empty() {
145            continue;
146        }
147        let response = match serde_json::from_str::<Envelope<Request>>(&line) {
148            Ok(envelope) => dispatch(envelope.body, &state, &shutdown).await,
149            Err(err) => Response::Err {
150                error: format!("unparseable request: {err}"),
151            },
152        };
153        let envelope = Envelope {
154            protocol: ProtocolVersion::V1,
155            version: build_version().to_owned(),
156            body: response,
157        };
158        let mut serialized = serde_json::to_string(&envelope)
159            .unwrap_or_else(|_| r#"{"error":"response serialization failed"}"#.to_owned());
160        serialized.push('\n');
161        write_half.write_all(serialized.as_bytes()).await?;
162    }
163    Ok(())
164}
165
166async fn dispatch(
167    request: Request,
168    state: &Arc<DaemonState>,
169    shutdown: &tokio::sync::mpsc::Sender<()>,
170) -> Response {
171    match request {
172        Request::Ping => Response::Ok(ResponseBody::Pong),
173        Request::RouteSet { host, port } => {
174            state.route_set(host, port);
175            Response::Ok(ResponseBody::Done)
176        }
177        Request::RouteDelete { host } => {
178            state.route_delete(&host);
179            Response::Ok(ResponseBody::Done)
180        }
181        Request::Routes => Response::Ok(ResponseBody::Routes {
182            routes: state.routes(),
183        }),
184        Request::Supervise {
185            instance,
186            service,
187            pid,
188            start_time,
189        } => {
190            state.supervise(instance, service, ProcessStamp { pid, start_time });
191            Response::Ok(ResponseBody::Done)
192        }
193        Request::Forget { instance } => {
194            state.forget(instance.as_str());
195            Response::Ok(ResponseBody::Done)
196        }
197        Request::InstanceProcesses { instance } => Response::Ok(ResponseBody::Processes {
198            processes: state.instance_processes(instance.as_str()),
199        }),
200        Request::Shutdown => {
201            let _ = shutdown.send(()).await;
202            Response::Ok(ResponseBody::Done)
203        }
204    }
205}