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}