Skip to main content

onlyne_client/session/dispatch/
outbound.rs

1use super::*;
2
3use super::env::{AGENT, REQUEST_TIMEOUT};
4use super::projection::with_cluster;
5use super::state::{DispatchInner, DispatchState};
6
7/// Write one ack into the durable intent queue.
8pub(super) fn store_ack(inner: &DispatchInner, mut ack: AckArgs) {
9    if ack.op_id.is_none() {
10        ack.op_id = Some(onlyne_proto::new_op_id());
11    }
12    let Some(op_id) = ack.op_id.clone() else {
13        return;
14    };
15    match serde_json::to_value(ClientOp::Ack(ack)) {
16        Ok(value) => {
17            if let Err(error) = inner.store.enqueue_intent(&op_id, &value) {
18                tracing::warn!(error = %error, "settled ack was not stored");
19            }
20        }
21        Err(error) => tracing::warn!(error = %error, "settled ack did not serialize"),
22    }
23}
24
25/// Queue one envelope into the durable intent table while the dispatch lock is
26/// already held. The `op_id` rule and the validation are `enqueue_outbound`'s.
27pub(super) fn queue_outbound_locked(
28    inner: &mut DispatchInner,
29    envelope: &Envelope,
30) -> Result<String> {
31    let mut stamped = envelope.clone();
32    let op_id = stamp_op_id(&mut stamped);
33    stamped
34        .validate()
35        .map_err(|error| anyhow!(error.to_string()))?;
36    inner
37        .store
38        .enqueue_intent(&op_id, &serde_json::to_value(&stamped)?)?;
39    Ok(op_id)
40}
41
42/// Hand one envelope to the live link, or to the intent queue when it is down.
43pub(super) async fn transport_envelope(state: &DispatchState, envelope: &Envelope) -> Result<()> {
44    let op = ClientOp::Send(Box::new(envelope.clone()));
45    let outbox = { state.inner.lock().outbox.clone() };
46    match outbox {
47        Some(outbox) => match outbox.send(op).await {
48            Ok(()) => Ok(()),
49            Err(error) => {
50                tracing::warn!(error = %error, "completion fell back to the intent queue");
51                state.enqueue_outbound(envelope).map(|_| ())
52            }
53        },
54        None => state.enqueue_outbound(envelope).map(|_| ()),
55    }
56}
57
58/// One authenticated server link.
59///
60/// The first five methods are the whole transport surface the runloop uses, so
61/// a change of transport touches this struct alone. The remaining two read the
62/// supervision state of the connection the handle keeps across redials.
63#[derive(Clone)]
64pub struct ClientLink {
65    handle: ClientConn,
66    welcome: Arc<Welcome>,
67    /// Skeleton of the routed `hello`. Its `live_sessions` is deliberately
68    /// never filled in: every send — `connect` and `authenticate` both — stamps
69    /// a freshly-read claim through [`hello_with_live_sessions`], so no caller
70    /// reads the stored list and a write-back here would only go stale.
71    hello: HandshakeArgs,
72}
73
74impl ClientLink {
75    /// Dial, verify the certificate pin, sign the server challenge, then read
76    /// the role slice with `hello`.
77    pub async fn connect(
78        init: &ClientInit,
79        live_sessions: Vec<LiveSession>,
80    ) -> Result<Self, NetError> {
81        let keypair = KeyPair::load(&init.key_path)?;
82        let settings = ConnSettings {
83            agent: AGENT.to_string(),
84            version: env!("CARGO_PKG_VERSION").to_string(),
85            ..ConnSettings::new(PROTOCOL_VERSION)
86        };
87        let handle: ClientConn =
88            dial(&init.server, &keypair, &init.cert_pin, &init.role, settings).await?;
89        let hello = HandshakeArgs {
90            protocol: PROTOCOL_VERSION,
91            role: init.role.clone(),
92            key: keypair.public_str(),
93            signature: String::new(),
94            agent: AGENT.to_string(),
95            version: env!("CARGO_PKG_VERSION").to_string(),
96            aggregate: false,
97            live_sessions: Vec::new(),
98        };
99        let body = handle
100            .request(
101                Frame::req(
102                    String::new(),
103                    ClientOp::Hello(hello_with_live_sessions(&hello, live_sessions)),
104                ),
105                REQUEST_TIMEOUT,
106            )
107            .await?;
108        if !body.ok {
109            let error = body.error.clone().unwrap_or(onlyne_proto::ErrorPayload {
110                code: onlyne_proto::ErrorCode::Internal,
111                message: "hello refused".to_string(),
112                field: None,
113            });
114            return Err(NetError::Rejected {
115                code: wire_code(error.code),
116                message: error.message,
117            });
118        }
119        let data = body.data().cloned().ok_or(NetError::BadFrame)?;
120        let welcome: Welcome = serde_json::from_value(data).map_err(|_| NetError::BadFrame)?;
121        Ok(Self {
122            handle,
123            welcome: Arc::new(welcome),
124            hello,
125        })
126    }
127
128    /// Send the routed `hello` again on the connection this link now holds.
129    ///
130    /// The net layer redials on its own, and a fresh connection carries no role
131    /// binding until this frame lands, so a caller replays it whenever readiness
132    /// returns (plan §7 line 310).
133    pub async fn authenticate(&self, live_sessions: Vec<LiveSession>) -> Result<(), NetError> {
134        let body = self
135            .handle
136            .request(
137                Frame::req(
138                    String::new(),
139                    ClientOp::Hello(hello_with_live_sessions(&self.hello, live_sessions)),
140                ),
141                REQUEST_TIMEOUT,
142            )
143            .await?;
144        if !body.ok {
145            let error = body.error.clone().unwrap_or(onlyne_proto::ErrorPayload {
146                code: onlyne_proto::ErrorCode::Internal,
147                message: "hello refused".to_string(),
148                field: None,
149            });
150            return Err(NetError::Rejected {
151                code: wire_code(error.code),
152                message: error.message,
153            });
154        }
155        Ok(())
156    }
157
158    /// Role slice the server bound this link to.
159    pub fn welcome(&self) -> &Welcome {
160        &self.welcome
161    }
162
163    /// Send one request frame and answer with its body. A refusal from the
164    /// server arrives inside the body, which keeps the retry decision in the
165    /// intent machine.
166    pub async fn request(&self, op: ClientOp) -> Result<ResBody, NetError> {
167        self.handle
168            .request(Frame::req(String::new(), op), REQUEST_TIMEOUT)
169            .await
170    }
171
172    /// Clone the server observation stream.
173    pub fn events(&self) -> broadcast::Receiver<Frame<ClientOp>> {
174        self.handle.events()
175    }
176
177    /// Send `bye` and drain in-flight requests.
178    pub async fn close(&self) -> Result<(), NetError> {
179        self.handle.close().await
180    }
181
182    /// Liveness of the connection behind this link.
183    pub fn readiness(&self) -> ConnReadiness {
184        self.handle.readiness()
185    }
186    /// Reason the supervisor stopped redialing.
187    pub async fn failure(&self) -> Option<NetError> {
188        self.handle.failure().await
189    }
190}
191
192/// Stamp the live sessions this client holds onto a hello skeleton.
193///
194/// Slots exist from assign until release. A fresh process holds none, so hello
195/// sends an empty list and the server requeues. A live client whose link flaps
196/// still holds its slots, so the deliveries those sessions are bound to stay
197/// in_flight.
198pub(crate) fn hello_with_live_sessions(
199    hello: &HandshakeArgs,
200    live_sessions: Vec<LiveSession>,
201) -> HandshakeArgs {
202    let mut hello = hello.clone();
203    hello.live_sessions = live_sessions;
204    hello
205}
206
207/// Wire code of a refusal, in the snake_case spelling both sides share.
208fn wire_code(code: onlyne_proto::ErrorCode) -> String {
209    serde_json::to_value(code)
210        .ok()
211        .and_then(|value| value.as_str().map(str::to_string))
212        .unwrap_or_else(|| "internal".to_string())
213}
214
215/// Where lifecycle frames leave the dispatcher.
216///
217/// `send` completes once the frame is on the wire. The ready report awaits this
218/// before the payload reaches the agent, which is the causal order §6 fixes.
219pub trait Outbox: Send + Sync {
220    fn send(&self, op: ClientOp)
221    -> Pin<Box<dyn Future<Output = Result<(), NetError>> + Send + '_>>;
222
223    /// One request round trip, for a caller that needs the server's answer
224    /// rather than a queued frame: the local CLI reports the verdict a send got.
225    fn request(
226        &self,
227        op: ClientOp,
228    ) -> Pin<Box<dyn Future<Output = Result<ResBody, NetError>> + Send + '_>>;
229}
230
231impl Outbox for ClientLink {
232    fn send(
233        &self,
234        op: ClientOp,
235    ) -> Pin<Box<dyn Future<Output = Result<(), NetError>> + Send + '_>> {
236        Box::pin(async move {
237            // A refusal is not a send. The server answered, which says the frame
238            // arrived, and it took nothing: the callers of this method treat `Ok`
239            // as "the server has it" — `send_frame` skips the durable queue on
240            // `Ok`, and `transport_envelope` skips it for a completion — so a
241            // refusal answered here would drop the frame outright. The two
242            // refusals this path meets while the queue is warming are the
243            // pre-handshake one (a window of the connection) and a transient
244            // internal one, and both are the intent machine's to judge, which is
245            // where the caller hands the frame instead.
246            let body = self.request(op).await?;
247            if !body.ok {
248                let error = body.error.clone().unwrap_or(onlyne_proto::ErrorPayload {
249                    code: onlyne_proto::ErrorCode::Internal,
250                    message: "the server took nothing and named no reason".to_string(),
251                    field: None,
252                });
253                return Err(NetError::Rejected {
254                    code: wire_code(error.code),
255                    message: error.message,
256                });
257            }
258            Ok(())
259        })
260    }
261
262    fn request(
263        &self,
264        op: ClientOp,
265    ) -> Pin<Box<dyn Future<Output = Result<ResBody, NetError>> + Send + '_>> {
266        Box::pin(async move { ClientLink::request(self, op).await })
267    }
268}
269
270/// Deliver one lifecycle frame, and queue it durably when the link is down.
271///
272/// §6 line 289: a running session reaches its terminal state while the outbound
273/// work waits in `client.db` intents for the flusher.
274///
275/// The frame that could not leave says nothing about the link, so the shared
276/// accept gate is left exactly where the runloop put it. A send fails with the
277/// link still `Ready` — the request deadline belongs to the caller, and
278/// `onlyne_net::conn` records that "the silent peer keeps the link up; only this
279/// call gave up" — and dropping the flag here latched it: `watch_readiness` only
280/// re-arms `accept_new` on a readiness transition, so the pull loop stopped
281/// draining the role's inbox for the life of that link while the work sat queued
282/// on the server, and the next delivery to arrive found a client that refused
283/// new work. The connection's own state is the flag's only author
284/// (`runloop::link`).
285pub async fn send_frame(state: &DispatchState, op: ClientOp) -> Result<()> {
286    let op = match op {
287        ClientOp::Report(report) => ClientOp::Report(with_cluster(state, report)),
288        other => other,
289    };
290    if let Some(outbox) = state.outbox() {
291        if outbox.send(op.clone()).await.is_ok() {
292            return Ok(());
293        }
294    }
295    state.enqueue_op(&op)?;
296    Ok(())
297}
298
299impl DispatchState {
300    /// Queue an outbound envelope before its first write and answer its op_id.
301    ///
302    /// The queue keys every row by an `op_id`, and the proto requires that key
303    /// only for the non-note kinds, so a note that arrives without one gets a
304    /// fresh client-minted id here: the row is keyed and what it replays is the
305    /// whole stamped envelope. A non-note keeps the id it brought, so a
306    /// re-delivered task still dedups on its original one.
307    pub fn enqueue_outbound(&self, envelope: &Envelope) -> Result<String> {
308        queue_outbound_locked(&mut self.inner.lock(), envelope)
309    }
310
311    /// The flag the runloop and the dispatcher share.
312    pub fn accept_new(&self) -> Arc<AtomicBool> {
313        self.inner.lock().accept_new.clone()
314    }
315
316    /// Whether the role holds a ready server link.
317    pub fn link_up(&self) -> bool {
318        self.inner.lock().link_up.load(Ordering::SeqCst)
319    }
320
321    /// The name of the session runtime hosting this role's sessions.
322    ///
323    /// The registration beside the client socket carries it, and an external
324    /// runtime's plugin matches on it to find the client its own session
325    /// belongs to, so the answer is the backend's own name rather than a
326    /// spelling invented here.
327    pub fn runtime_name(&self) -> String {
328        self.inner.lock().backend.name().to_string()
329    }
330
331    /// Record that the server link came up or went down.
332    pub fn set_link_up(&self, up: bool) {
333        self.inner.lock().link_up.store(up, Ordering::SeqCst);
334    }
335
336    /// Aggregate name this role supervises, empty for a plain role.
337    pub fn cluster_ref(&self) -> String {
338        self.inner.lock().cluster_ref.clone()
339    }
340
341    /// Record the server's topology name, read from `welcome.cluster`.
342    ///
343    /// Each spawned session carries it as `ONLYNE_CLUSTER`, which is how a host
344    /// backend addresses the tree it splits panes into. The runloop calls
345    /// this on every welcome, so a server that reloads under a new name is
346    /// followed by the sessions spawned after that point.
347    pub fn set_topology(&self, cluster: &str) {
348        self.inner.lock().topology = cluster.trim().to_string();
349    }
350
351    /// The topology name recorded from `welcome`, empty before the first welcome.
352    pub fn topology(&self) -> String {
353        self.inner.lock().topology.clone()
354    }
355
356    /// Record the aggregate name once, so every report keeps the same value
357    /// across a reconnect.
358    pub fn set_cluster_ref(&self, aggregate: impl Into<String>) {
359        self.inner.lock().cluster_ref = aggregate.into();
360    }
361
362    /// Ask the server one question over the live link.
363    ///
364    /// Err means the link is down, never a refusal: a refusal arrives as an
365    /// `Ok` body carrying `ok: false`, which is what the local CLI shows.
366    pub async fn request(&self, op: ClientOp) -> Result<ResBody, NetError> {
367        let outbox = { self.inner.lock().outbox.clone() };
368        let Some(outbox) = outbox else {
369            return Err(NetError::NotReady);
370        };
371        outbox.request(op).await
372    }
373
374    /// Install the live link as the outbound path.
375    pub fn attach_outbox(&self, outbox: Arc<dyn Outbox>) {
376        self.inner.lock().outbox = Some(outbox);
377    }
378
379    /// Remove the outbound path; lifecycle frames then queue as intents.
380    pub fn detach_outbox(&self) {
381        self.inner.lock().outbox = None;
382    }
383
384    fn outbox(&self) -> Option<Arc<dyn Outbox>> {
385        self.inner.lock().outbox.clone()
386    }
387
388    /// Queue one client op in the durable intent table and answer its op_id.
389    pub fn enqueue_op(&self, op: &ClientOp) -> Result<String> {
390        let op_id = onlyne_proto::new_id();
391        self.inner
392            .lock()
393            .store
394            .enqueue_intent(&op_id, &serde_json::to_value(op)?)?;
395        Ok(op_id)
396    }
397}