Skip to main content

onlyne_client/session/adapter_socket/
socket.rs

1use crate::session::dispatch::DispatchState;
2use anyhow::{Context, Result};
3use onlyne_adapter::AdapterIo;
4use onlyne_config::layout::RoleWorkspace;
5use onlyne_proto::{AdapterMsg, HelloArgs, HostOp, MountKind, PROTOCOL_VERSION, PluginOp};
6use onlyne_wire::socket::prelude::TokioListener;
7use onlyne_wire::socket::{
8    LocalListener, RegistrationFile, SocketEndpoint, bind_socket_v2, connect_local,
9    registration_path, remove_registration, write_registration,
10};
11use std::path::{Path, PathBuf};
12use std::time::Duration;
13
14#[derive(Clone)]
15pub struct AdapterSocket {
16    pub workspace: PathBuf,
17    pub role: String,
18    pub cluster: String,
19    pub server: String,
20    pub dispatch: DispatchState,
21}
22
23impl AdapterSocket {
24    /// The path this role's clients connect to, per
25    /// [`RoleWorkspace::socket_path`].
26    ///
27    /// One accessor answers for the whole tree, so the daemon that bound
28    /// `<runtime>/<digest>.sock` and a caller that derived it from the
29    /// workspace root reach the same socket.
30    pub fn path(&self) -> PathBuf {
31        RoleWorkspace::resolve(&self.workspace).socket_path()
32    }
33
34    /// The registration this client publishes beside its socket.
35    ///
36    /// The kind is the client's own: a role workspace's daemon serves the
37    /// client surface, and no reader infers that from what happens to sit
38    /// beside the tree. `role` names the one role this process serves,
39    /// `runtime` the session backend that hosts it, and `placement` where this
40    /// machine displays that runtime — the machine-owned half of the pair an
41    /// external runtime's plugin matches on to find the clients it serves.
42    pub fn registration(&self) -> RegistrationFile {
43        client_registration(
44            &self.workspace,
45            &self.role,
46            &self.dispatch.runtime_name(),
47            self.dispatch.placement_name(),
48        )
49    }
50
51    /// Bind the workspace socket and report the endpoint that was served.
52    ///
53    /// [`bind_socket_v2`] owns the whole sequence — resolve the runtime
54    /// directory, drop a stale name, and bind `<runtime>/<digest>.sock` — and
55    /// its failure message names the path it tried. The registration is
56    /// published after the bind succeeds, so a surface that never opened never
57    /// names itself.
58    ///
59    /// The socket is served off the canonical `<run_dir>/s` spelling, so the log
60    /// line names that spelling too: it is what an operator reading the tree
61    /// will look for and not find.
62    #[allow(clippy::unused_async)]
63    pub async fn bind(&self) -> Result<(LocalListener, SocketEndpoint)> {
64        let layout = RoleWorkspace::resolve(&self.workspace);
65        let endpoint = layout.socket_endpoint();
66        let listener = bind_socket_v2(layout.root()).with_context(|| {
67            format!(
68                "bind the workspace socket {}",
69                layout.socket_path().display(),
70            )
71        })?;
72        endpoint.publish(&self.registration()).with_context(|| {
73            format!(
74                "publish the client registration {}",
75                endpoint.registration().display(),
76            )
77        })?;
78        // Under v2 the socket always lives in the runtime directory, so serving
79        // off the `<run>/s` spelling is the normal case and not a move worth a
80        // warning. The canonical spelling is logged because it is what an
81        // operator reading the tree will look for and not find.
82        tracing::info!(
83            canonical = %endpoint.natural().display(),
84            served = %endpoint.actual().display(),
85            registration = %endpoint.registration().display(),
86            role = %self.role,
87            "adapter socket serving"
88        );
89        Ok((listener, endpoint))
90    }
91
92    /// Answer every connection on `listener` for the life of the process.
93    ///
94    /// The listener stays open across a failed `accept`: one bad handshake on one
95    /// socket is that client's problem, and every plugin already mounted would
96    /// lose its transport if the surface came down for it. The pause keeps a
97    /// persistent failure — a descriptor ceiling, a name pulled out from under
98    /// the listener — from spinning the loop at full speed.
99    ///
100    /// The registration published by [`bind`](Self::bind) belongs to this scope:
101    /// it names this process as the client for the workspace, and the guard here
102    /// takes the file with it when the surface stops serving, however it stops.
103    pub async fn accept_loop(&self, listener: LocalListener) -> Result<()> {
104        let _registration = RegistrationGuard {
105            root: RoleWorkspace::resolve(&self.workspace).root().to_path_buf(),
106        };
107        loop {
108            match listener.accept().await {
109                Ok(stream) => {
110                    let this = self.clone();
111                    tokio::spawn(async move {
112                        if let Err(err) = this.connection(stream).await {
113                            tracing::debug!(error = %err, "adapter connection closed");
114                        }
115                    });
116                }
117                Err(error) => {
118                    tracing::error!(
119                        error = %error,
120                        kind = ?error.kind(),
121                        "adapter socket accept failed; retrying"
122                    );
123                    tokio::time::sleep(ACCEPT_RETRY_PAUSE).await;
124                }
125            }
126        }
127    }
128
129    /// Bind and serve in one call, for a caller that owns no endpoint interest.
130    pub async fn serve(self) -> Result<()> {
131        let (listener, _) = self.bind().await?;
132        self.accept_loop(listener).await
133    }
134}
135
136/// The registration a client for `workspace` publishes beside its socket.
137///
138/// One builder answers for the bind and for every republish after it, because
139/// the file is the machine-level record of which surface serves this tree: a
140/// second spelling of that record would drift from the first.
141pub fn client_registration(
142    workspace: &Path,
143    role: &str,
144    runtime: &str,
145    placement: Option<&str>,
146) -> RegistrationFile {
147    let root = RoleWorkspace::resolve(workspace);
148    let registration = RegistrationFile::client(root.root())
149        .with_role(role)
150        .with_runtime(runtime);
151    match placement {
152        Some(name) => registration.with_placement(name),
153        None => registration,
154    }
155}
156
157/// Rewrite the registration for `workspace` with the facts this client has now.
158///
159/// The bind publishes a registration before the role's drive is known, so the
160/// `runtime` field it carries is the default drive's backend. Every `welcome`
161/// names the real drive and the client installs the backend it selects, and
162/// this is what puts that answer on the file an external runtime's plugin reads
163/// to find the clients it serves.
164pub fn republish_registration(
165    workspace: &Path,
166    role: &str,
167    runtime: &str,
168    placement: Option<&str>,
169) -> Result<()> {
170    let root = RoleWorkspace::resolve(workspace);
171    let path = registration_path(root.root());
172    write_registration(
173        root.root(),
174        &client_registration(workspace, role, runtime, placement),
175    )
176    .with_context(|| format!("republish the client registration {}", path.display()))
177}
178
179/// The registration a serving surface owns for as long as it serves.
180///
181/// The file is the one answer to "which client serves this workspace", and a
182/// registration outliving its surface is what an external runtime's plugin
183/// reads as a live client. Tying it to the scope means every end reaches the
184/// removal: a caller that drops the future, an acceptor the run aborts, and a
185/// surface that fails.
186struct RegistrationGuard {
187    root: PathBuf,
188}
189
190impl Drop for RegistrationGuard {
191    fn drop(&mut self) {
192        if let Err(error) = remove_registration(&self.root) {
193            tracing::warn!(
194                error = %error,
195                registration = %registration_path(&self.root).display(),
196                "could not remove the client registration"
197            );
198        } else {
199            tracing::info!(root = %self.root.display(), "removed the client registration");
200        }
201    }
202}
203
204/// Wait bound for the link probe's `hello` round trip.
205///
206/// The handshake runs over a local socket, and the bound matches the adapter
207/// protocol's own hello budget.
208pub const PROBE_TIMEOUT: Duration = Duration::from_secs(5);
209
210/// Pause before the accept loop asks the listener for a connection again after
211/// an `accept` failure.
212///
213/// The value is short enough that a transient failure costs one plugin one
214/// keystroke and long enough that a persistent one stays off the CPU.
215pub const ACCEPT_RETRY_PAUSE: Duration = Duration::from_millis(100);
216
217/// The link state of the client serving `socket`, and `None` when no client
218/// answers the probe.
219///
220/// The probe is an `admin` `hello` on the adapter surface (§7): it names no
221/// mount, binds nothing, and the host's `HelloAck` carries the link state. A
222/// socket nobody answers is a client that is not serving — `status` reads that
223/// as not running — while a client that answers with a down link is running
224/// and disconnected.
225pub async fn server_link_state(socket: &Path) -> Option<bool> {
226    let stream = connect_local(socket).await.ok()?;
227    let io = AdapterIo::new(stream, PROBE_TIMEOUT, PROBE_TIMEOUT);
228    let hello = HelloArgs {
229        protocol: PROTOCOL_VERSION,
230        plugin: "onlyne-client".to_string(),
231        version: env!("CARGO_PKG_VERSION").to_string(),
232        kind: MountKind::Admin,
233        capabilities: Vec::new(),
234        mount: None,
235    };
236    let body = io
237        .request(AdapterMsg::Plugin(PluginOp::Hello(hello)))
238        .await
239        .ok()?;
240    if !body.ok {
241        return Some(false);
242    }
243    let ack = body
244        .data
245        .and_then(|value| serde_json::from_value::<HostOp>(value).ok());
246    Some(matches!(
247        ack,
248        Some(HostOp::Welcome(ack)) if ack.server.connected
249    ))
250}
251
252/// Remove a socket file a previous run left behind.
253///
254/// The leaf is absent in the common case, and `NotFound` is that answer. A name
255/// that holds a live listener answers `connect` and keeps the bind of a client
256/// restarting behind it, so the readiness path clears it first, and the log line
257/// is the record that a surface was cleared.
258pub async fn stale_socket_removed(path: &Path) -> Result<()> {
259    match tokio::fs::remove_file(path).await {
260        Ok(()) => {
261            tracing::info!(socket = %path.display(), "removed stale adapter socket");
262            Ok(())
263        }
264        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
265        Err(error) => Err(error).with_context(|| format!("remove stale socket {}", path.display())),
266    }
267}