Skip to main content

scv_client/
lib.rs

1//! Shared local transport interfaces, without server policy or bridge dependencies.
2
3use anyhow::{Context, Result, bail};
4use scv_protocol::{
5    ClientMessage, DaemonCommand, DaemonStatus, PROTOCOL_VERSION, PeerInfo, ServerEvent,
6};
7use std::{
8    path::{Path, PathBuf},
9    time::Duration,
10};
11use tokio::{
12    io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
13    net::UnixStream,
14};
15
16pub fn default_socket_path() -> Result<PathBuf> {
17    let root = std::env::var_os("SCV_HOME")
18        .map(PathBuf::from)
19        .or_else(|| dirs::home_dir().map(|path| path.join(".scv")))
20        .context("cannot determine SCV_HOME")?;
21    Ok(root.join("server.sock"))
22}
23
24/// A bounded management exchange. Never retries mutations on ambiguous failure.
25pub async fn control(path: &Path, command: DaemonCommand) -> Result<DaemonStatus> {
26    tokio::time::timeout(Duration::from_secs(20), async {
27        let stream = UnixStream::connect(path)
28            .await
29            .context("SCV daemon unavailable; start it with `scv start` or `scv run`")?;
30        let (reader, mut writer) = stream.into_split();
31        let mut reader = BufReader::new(reader);
32        for message in [
33            ClientMessage::Initialize {
34                request_id: "init".into(),
35                protocol_version: PROTOCOL_VERSION,
36                client: PeerInfo {
37                    name: "scv-control".into(),
38                    version: env!("CARGO_PKG_VERSION").into(),
39                },
40            },
41            ClientMessage::DaemonControl {
42                request_id: "control".into(),
43                command,
44            },
45        ] {
46            let mut frame = serde_json::to_vec(&message)?;
47            frame.push(b'\n');
48            writer.write_all(&frame).await?;
49            let mut bytes = Vec::new();
50            loop {
51                let buf = reader.fill_buf().await?;
52                if buf.is_empty() {
53                    bail!("SCV daemon closed the management connection");
54                }
55                let take = buf
56                    .iter()
57                    .position(|b| *b == b'\n')
58                    .map_or(buf.len(), |n| n + 1);
59                if bytes.len() + take > 1024 * 1024 {
60                    bail!("SCV status exceeds frame limit");
61                }
62                bytes.extend_from_slice(&buf[..take]);
63                reader.consume(take);
64                if bytes.last() == Some(&b'\n') {
65                    break;
66                }
67            }
68            match serde_json::from_slice::<ServerEvent>(&bytes)? {
69                ServerEvent::Initialized {
70                    protocol_version: PROTOCOL_VERSION,
71                    ..
72                } if matches!(message, ClientMessage::Initialize { .. }) => {}
73                ServerEvent::DaemonStatus { status, .. } => return Ok(status),
74                ServerEvent::Error { message, .. } => bail!("{message}"),
75                _ => bail!("unexpected SCV management response; upgrade/restart the daemon"),
76            }
77        }
78        bail!("SCV daemon omitted status")
79    })
80    .await
81    .context("SCV management request timed out; query status before retrying")?
82}