stackless_daemon/
server.rs1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
20pub enum DaemonRole {
21 Operator,
23 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
35pub 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
45pub 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 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 if role == DaemonRole::Operator {
79 crate::launchd::ensure_registered(paths);
80 }
81
82 let state = Arc::new(DaemonState::default());
83
84 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 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}