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}