返回 CodeWhale
daemon_client.rs
根目录 / crates / app-server / src / daemon_client.rs
1 //! Authenticated local control attachment to the already held Runtime owner.
2 //! The private receipt routes; the connected kernel process authenticates it.
3 //! This forwards the existing bounded dispatcher protocol, never HTTP bearers.
4
5 use anyhow::Result;
6 use std::path::PathBuf;
7
8 #[cfg(any(unix, windows))]
9 mod platform {
10 use super::*;
11 use crate::daemon_socket::AuthorizedPeer;
12 use crate::daemon_socket::{default_socket_path, owner_work};
13 use crate::{BoundedLines, ParsedStdioLine, parse_stdio_line, write_stdio_line};
14 use anyhow::{Context as _, bail};
15 use codewhale_config::private_directory::PrivateDirectory;
16 use codewhale_protocol::RuntimeOwnerReceipt;
17 use serde_json::{Value, json};
18 use std::sync::Arc;
19 use std::time::Duration;
20 use tokio::io::{AsyncWriteExt as _, BufReader};
21 #[cfg(unix)]
22 use tokio::net::UnixStream;
23 #[cfg(unix)]
24 type Connection = UnixStream;
25 #[cfg(windows)]
26 type Connection = tokio::net::windows::named_pipe::NamedPipeClient;
27
28 pub struct OwnerClient {
29 read: BoundedLines<BufReader<tokio::io::ReadHalf<Connection>>>,
30 write: tokio::io::WriteHalf<Connection>,
31 peer: Arc<AuthorizedPeer>,
32 receipt: RuntimeOwnerReceipt,
33 routing: Option<crate::RuntimeOwnerRouting>,
34 frontend: crate::daemon_socket::AttachFrontend,
35 }
36
37 pub async fn connect(
38 config_path: Option<PathBuf>,
39 selected: Option<PathBuf>,
40 ) -> Result<OwnerClient> {
41 connect_if_published(config_path, selected).await?.context(
42 "the selected Runtime has no authenticated live owner; start its canonical host",
43 )
44 }
45
46 pub async fn connect_if_published(
47 config_path: Option<PathBuf>,
48 selected: Option<PathBuf>,
49 ) -> Result<Option<OwnerClient>> {
50 connect_frontend(
51 config_path,
52 selected,
53 crate::daemon_socket::AttachFrontend::Control,
54 None,
55 None,
56 None,
57 None,
58 )
59 .await
60 }
61 pub async fn connect_acp_if_published(
62 config_path: Option<PathBuf>,
63 selected: Option<PathBuf>,
64 ) -> Result<Option<OwnerClient>> {
65 connect_frontend(
66 config_path,
67 selected,
68 crate::daemon_socket::AttachFrontend::Acp,
69 None,
70 None,
71 None,
72 None,
73 )
74 .await
75 }
76 pub async fn connect_listener_if_published(
77 config_path: Option<PathBuf>,
78 selected: Option<PathBuf>,
79 listener: crate::RuntimeListenerSelection,
80 expected_owner: RuntimeOwnerReceipt,
81 ) -> Result<Option<OwnerClient>> {
82 listener.validate_bounds()?;
83 connect_frontend(
84 config_path,
85 selected,
86 crate::daemon_socket::AttachFrontend::Listener,
87 Some(listener),
88 Some(expected_owner),
89 None,
90 None,
91 )
92 .await
93 }
94 pub async fn connect_scoped_control_if_published(
95 config_path: Option<PathBuf>,
96 selected: Option<PathBuf>,
97 scope: crate::RuntimeFrontendScope,
98 owner: RuntimeOwnerReceipt,
99 ) -> Result<Option<OwnerClient>> {
100 scope.validate_bounds()?;
101 connect_frontend(
102 config_path,
103 selected,
104 crate::daemon_socket::AttachFrontend::Control,
105 None,
106 Some(owner),
107 Some(scope),
108 None,
109 )
110 .await
111 }
112 pub async fn connect_selected_acp_if_published(
113 config_path: Option<PathBuf>,
114 selected: Option<PathBuf>,
115 scope: crate::RuntimeFrontendScope,
116 model: String,
117 owner: RuntimeOwnerReceipt,
118 ) -> Result<Option<OwnerClient>> {
119 scope.validate_bounds()?;
120 anyhow::ensure!(
121 !model.trim().is_empty() && model.len() <= 1024,
122 "invalid selected ACP model"
123 );
124 connect_frontend(
125 config_path,
126 selected,
127 crate::daemon_socket::AttachFrontend::Acp,
128 None,
129 Some(owner),
130 Some(scope),
131 Some(model),
132 )
133 .await
134 }
135 async fn connect_frontend(
136 config_path: Option<PathBuf>,
137 selected: Option<PathBuf>,
138 frontend: crate::daemon_socket::AttachFrontend,
139 listener: Option<crate::RuntimeListenerSelection>,
140 expected_owner: Option<RuntimeOwnerReceipt>,
141 scope: Option<crate::RuntimeFrontendScope>,
142 acp_model: Option<String>,
143 ) -> Result<Option<OwnerClient>> {
144 let socket_path = owner_work(move || match selected {
145 Some(path) => Ok(path),
146 None => default_socket_path().map_err(Into::into),
147 })
148 .await?;
149 #[cfg(unix)]
150 let directory = socket_path
151 .parent()
152 .context("owner endpoint has no private parent")?
153 .to_path_buf();
154 #[cfg(windows)]
155 let directory = owner_work(crate::daemon_windows::owner_directory).await?;
156 #[cfg(unix)]
157 let name = format!(
158 "{}.owner.json",
159 socket_path
160 .file_name()
161 .and_then(|name| name.to_str())
162 .context("invalid owner endpoint basename")?
163 );
164 #[cfg(windows)]
165 let name = "daemon.owner.json".to_string();
166 let receipt_name = name.clone();
167 let found = owner_work(move || {
168 let parent = match PrivateDirectory::inspect(&directory) {
169 Ok(parent) => Arc::new(parent),
170 Err(error)
171 if error
172 .downcast_ref::<std::io::Error>()
173 .is_some_and(|error| error.kind() == std::io::ErrorKind::NotFound) =>
174 {
175 return Ok(None);
176 }
177 Err(error) => return Err(error),
178 };
179 let Some((bytes, file)) = parent.read_private_receipt(&name, 16384)? else {
180 return Ok(None);
181 };
182 let receipt: RuntimeOwnerReceipt = serde_json::from_slice(&bytes)?;
183 Ok(Some((parent, receipt, file)))
184 })
185 .await?;
186 let Some((parent, receipt, file)) = found else {
187 return Ok(None);
188 };
189 anyhow::ensure!(
190 receipt.version == 1
191 && receipt.pid > 0
192 && receipt.socket_path == socket_path
193 && !receipt.lease_generation.is_empty(),
194 "selected owner receipt is invalid"
195 );
196 anyhow::ensure!(
197 receipt.config_path == config_path,
198 "selected owner config does not match; refusing another store"
199 );
200 #[cfg(unix)]
201 let (stream, peer) = {
202 let stream =
203 tokio::time::timeout(Duration::from_secs(5), UnixStream::connect(&socket_path))
204 .await
205 .context("owner connect deadline expired")??;
206 let credential = stream
207 .peer_cred()
208 .context("owner peer credentials unavailable")?;
209 let pid = credential
210 .pid()
211 .and_then(|pid| u32::try_from(pid).ok())
212 .context("owner peer PID unavailable")?;
213 anyhow::ensure!(
214 credential.uid() == PrivateDirectory::current_user_id()
215 && receipt.principal == credential.uid().to_string()
216 && pid == receipt.pid,
217 "connected process is not the selected Runtime owner"
218 );
219 (
220 stream,
221 Arc::new(AuthorizedPeer {
222 pid,
223 start: receipt.process_start.clone(),
224 }),
225 )
226 };
227 #[cfg(windows)]
228 let (stream, peer) = {
229 let (stream, process) = crate::daemon_windows::connect_owner(&receipt).await?;
230 let peer = AuthorizedPeer {
231 pid: process.pid(),
232 start: process.start().to_string(),
233 process,
234 };
235 (stream, Arc::new(peer))
236 };
237 let held_parent = parent.clone();
238 let held_receipt = receipt.clone();
239 let check = peer.clone();
240 owner_work(move || {
241 check.check()?;
242 anyhow::ensure!(
243 held_parent.is_at_selected_path()?,
244 "selected owner parent changed"
245 );
246 let (bytes, current) = held_parent
247 .read_private_receipt(&receipt_name, 16384)?
248 .context("selected owner receipt withdrawn")?;
249 anyhow::ensure!(
250 serde_json::from_slice::<RuntimeOwnerReceipt>(&bytes)? == held_receipt
251 && PrivateDirectory::same_file_identity(&file, &current)?,
252 "selected owner receipt changed before attach"
253 );
254 Ok(())
255 })
256 .await?;
257 anyhow::ensure!(
258 expected_owner
259 .as_ref()
260 .is_none_or(|expected| expected == &receipt),
261 "selected owner changed before frontend admission"
262 );
263 let (read, mut write) = tokio::io::split(stream);
264 let mut replies = BoundedLines::new(BufReader::new(read));
265 let attach_id = uuid::Uuid::new_v4().to_string();
266 let attach = json!({"jsonrpc":"2.0","id":attach_id,"method":"daemon/attach","params":{
267 "client":{"name":"codewhale-stdio","version":env!("CARGO_PKG_VERSION"),"pid":std::process::id()},
268 "mode":"attach","frontend":frontend,"listener":listener,"scope":scope,"acp_model":acp_model,"expect_daemon_version":env!("CARGO_PKG_VERSION"),"expect_owner":receipt
269 }});
270 tokio::time::timeout(
271 Duration::from_secs(5),
272 write_stdio_line(&mut write, &attach),
273 )
274 .await
275 .context("owner attach write deadline expired; outcome uncertain")??;
276 let line = tokio::time::timeout(Duration::from_secs(5), replies.next_line())
277 .await
278 .context("owner attach response deadline expired")??
279 .context("owner closed before attach")?;
280 anyhow::ensure!(line.len() <= 16384, "oversized owner attach response");
281 let response: Value = serde_json::from_str(&line)?;
282 anyhow::ensure!(
283 response["id"] == attach_id
284 && response["error"].is_null()
285 && response["result"]["role"] == "attached",
286 "selected owner refused guest attachment"
287 );
288 let returned: RuntimeOwnerReceipt =
289 serde_json::from_value(response["result"]["owner_receipt"].clone())
290 .context("owner attach receipt unavailable")?;
291 anyhow::ensure!(
292 returned == receipt,
293 "connected owner did not confirm the exact captured generation/store"
294 );
295 let routing = response["result"]
296 .get("runtime_routing")
297 .filter(|value| !value.is_null())
298 .map(|value| serde_json::from_value(value.clone()))
299 .transpose()?;
300 Ok(Some(OwnerClient {
301 read: replies,
302 write,
303 peer,
304 receipt,
305 routing,
306 frontend,
307 }))
308 }
309
310 impl OwnerClient {
311 pub fn routing(&self) -> Option<&crate::RuntimeOwnerRouting> {
312 self.routing.as_ref()
313 }
314
315 pub fn receipt(&self) -> &RuntimeOwnerReceipt {
316 &self.receipt
317 }
318 pub async fn send(&mut self, id: Value, method: &str, params: Value) -> Result<()> {
319 let frame = json!({"jsonrpc":"2.0","id":id,"method":method,"params":params});
320 let bytes = serde_json::to_vec(&frame)?;
321 anyhow::ensure!(
322 bytes.len() <= crate::MAX_RUNTIME_IMAGE_BODY_BYTES,
323 "owner request exceeds transport bound"
324 );
325 let check = self.peer.clone();
326 owner_work(move || check.check()).await?;
327 tokio::time::timeout(
328 Duration::from_secs(5),
329 write_stdio_line(&mut self.write, &frame),
330 )
331 .await
332 .context("owner write deadline expired; outcome uncertain, not replayed")??;
333 Ok(())
334 }
335 pub async fn recv(&mut self) -> Result<Option<Value>> {
336 let Some(line) = self.read.next_line().await? else {
337 return Ok(None);
338 };
339 let check = self.peer.clone();
340 owner_work(move || check.check()).await?;
341 serde_json::from_str(&line).map(Some).map_err(Into::into)
342 }
343 pub async fn forward<I, O>(mut self, input: I, output: O) -> Result<()>
344 where
345 I: tokio::io::AsyncRead + Unpin,
346 O: tokio::io::AsyncWrite + Unpin,
347 {
348 let mut input = BoundedLines::new(BufReader::new(input));
349 let mut output = tokio::io::BufWriter::new(output);
350 let mut input_open = true;
351 loop {
352 tokio::select! {
353 line = input.next_line(), if input_open => match line? {
354 None => {
355 input_open=false;
356 tokio::time::timeout(Duration::from_secs(5),write_stdio_line(&mut self.write,&json!({"jsonrpc":"2.0","method":"daemon/input_closed","params":{}}))).await.context("owner input-close write outcome uncertain, not replayed")??;
357 self.write.shutdown().await?;
358 }
359 Some(line) if self.frontend==crate::daemon_socket::AttachFrontend::Acp => {
360 let check=self.peer.clone();
361 owner_work(move || check.check()).await?;
362 let value:Value=serde_json::from_str(&line)?;
363 anyhow::ensure!(value.is_object() && value["jsonrpc"]=="2.0","ACP frame must be a JSON-RPC object");
364 // Preserve response IDs/permission replies verbatim.
365 // Actual ACP authority/parser stays in AcpServer.
366 tokio::time::timeout(Duration::from_secs(5),async {self.write.write_all(line.as_bytes()).await?;self.write.write_all(b"\n").await?;self.write.flush().await}).await.context("ACP write outcome uncertain; not replayed")??;
367 }
368 Some(line) => match parse_stdio_line(&line) {
369 ParsedStdioLine::Blank => {},
370 ParsedStdioLine::Rejected(response) => write_stdio_line(&mut output, &response).await?,
371 ParsedStdioLine::Request(request) => {
372 self.send(request.id.unwrap_or(Value::Null),&request.method,request.params).await?;
373 }
374 }
375 },
376 line = self.read.next_line() => match line? {
377 None if !input_open => return Ok(()),
378 None => bail!("selected owner connection closed; pending operation outcomes may be uncertain"),
379 Some(line) => {
380 let _: Value = serde_json::from_str(&line).context("invalid owner response")?;
381 output.write_all(line.as_bytes()).await?;
382 output.write_all(b"\n").await?;
383 output.flush().await?;
384 }
385 }
386 }
387 }
388 }
389 }
390
391 pub async fn forward_stdio(config_path: Option<PathBuf>) -> Result<()> {
392 connect(config_path, None)
393 .await?
394 .forward(tokio::io::stdin(), tokio::io::stdout())
395 .await
396 }
397 }
398
399 #[cfg(any(unix, windows))]
400 pub use platform::{
401 OwnerClient, connect, connect_acp_if_published, connect_if_published,
402 connect_listener_if_published, connect_scoped_control_if_published,
403 connect_selected_acp_if_published, forward_stdio,
404 };
405
406 #[cfg(not(any(unix, windows)))]
407 pub async fn forward_stdio(_config_path: Option<PathBuf>) -> Result<()> {
408 anyhow::bail!("authenticated live-owner control is not implemented on this platform")
409 }
410
410 lines RUST