Skip to main content

onlyne_client/runtime/
intent.rs

1use anyhow::{Context, Result};
2use chrono::{DateTime, Utc};
3use onlyne_proto::{ClientOp, Envelope, ErrorCode, MsgKind, Receipt, Report, ResBody, new_op_id};
4use onlyne_store::{ClientStore, IntentRow};
5use serde_json::Value;
6use std::time::Duration;
7
8/// Refusals that end an intent.
9///
10/// `Invalid` stays out: the server answers it for several conditions a retry
11/// clears, and exhaustion keeps the row visible as a fault, so a transient
12/// refusal costs retries rather than the queued message (plan §6 line 289).
13pub const PERMANENT_ERRORS: &[ErrorCode] = &[
14    ErrorCode::AclDenied,
15    ErrorCode::Conflict,
16    ErrorCode::Forbidden,
17    ErrorCode::UnknownRole,
18    ErrorCode::NotAdmin,
19    ErrorCode::BadFrame,
20    ErrorCode::FrameTooLarge,
21    ErrorCode::ProtocolVersion,
22];
23
24/// Passes a row this process cannot act on may be skipped before it retires.
25///
26/// The role's own ceiling (`attempts`) cannot bound this: it counts the answers
27/// the server gave, and a row whose stored payload no longer decodes never
28/// reaches the server to be answered. Without a bound of its own such a row sits
29/// at its deadline, first in the flush batch, taking a slot on every pass for
30/// the life of the queue (round-2 audit C1).
31pub const MAX_LOCAL_FAILURES: u32 = 5;
32
33#[derive(Debug, Clone, Copy, PartialEq, Eq)]
34pub enum IntentState {
35    Pending,
36    Retrying,
37    Accepted,
38    Exhausted,
39}
40
41impl IntentState {
42    pub fn as_str(self) -> &'static str {
43        match self {
44            Self::Pending => "pending",
45            Self::Retrying => "retrying",
46            Self::Accepted => "accepted",
47            Self::Exhausted => "exhausted",
48        }
49    }
50}
51
52#[derive(Debug, Clone, PartialEq)]
53pub enum IntentResult {
54    Accepted(Option<Receipt>),
55    Retryable(ErrorCode, String),
56    Dropped(ErrorCode, String),
57    Exhausted,
58    /// The frame reached a connection whose routed `hello` has not landed yet.
59    ///
60    /// The net layer redials behind the runloop's back, so a connection reads
61    /// `Ready` while it still carries no role binding, and the server refuses
62    /// every frame that arrives in that window. The refusal is a fact about the
63    /// connection, not about the row: charging it a backoff rung takes the row
64    /// out of the due window for a whole delay, and the flush the runloop runs
65    /// behind the replayed `hello` then finds nothing to send — which is how a
66    /// reconnecting client's own state stayed off the server's mirror for a
67    /// second (e2e case 15). The row keeps its place in the queue instead.
68    NotAuthenticated,
69}
70
71/// Give one outbound envelope the `op_id` its intent row is keyed by.
72///
73/// The proto requires the key for every kind but `Note` (`Envelope::validate`
74/// in `onlyne-proto`), so a plugin's note legitimately arrives without one
75/// while the queue still keys every row by an id. Minting it here keeps the
76/// stamped envelope and the row's key one value: what the row replays is what
77/// it stored. A non-note envelope keeps the key it brought, which is what
78/// makes a re-delivered task dedup on its original id.
79pub fn stamp_op_id(envelope: &mut Envelope) -> String {
80    if let Some(op_id) = &envelope.op_id {
81        return op_id.clone();
82    }
83    let op_id = new_op_id();
84    envelope.op_id = Some(op_id.clone());
85    op_id
86}
87
88#[derive(Clone)]
89pub struct IntentMachine {
90    pub store: ClientStore,
91    pub attempts: u32,
92    pub backoff_ms: Vec<u64>,
93}
94
95impl IntentMachine {
96    pub fn new(store: ClientStore, attempts: u32, backoff_ms: Vec<u64>) -> Self {
97        Self {
98            store,
99            attempts,
100            backoff_ms,
101        }
102    }
103
104    /// Queue one envelope under the id its row is keyed by.
105    ///
106    /// A kind that may arrive without a key (a note) gets a fresh one from
107    /// [`stamp_op_id`], and the stamped envelope is what the row stores and
108    /// replays, so the row's `op_id` and its `env_json` never disagree.
109    pub fn enqueue(&self, envelope: &Envelope) -> Result<bool> {
110        let mut stamped = envelope.clone();
111        let op_id = stamp_op_id(&mut stamped);
112        Ok(self
113            .store
114            .enqueue_intent(&op_id, &serde_json::to_value(&stamped)?)?)
115    }
116
117    pub fn enqueue_value(&self, op_id: &str, envelope: &Value) -> Result<bool> {
118        Ok(self.store.enqueue_intent(op_id, envelope)?)
119    }
120
121    /// The rows this pass may send: due, oldest deadline first, one batch.
122    ///
123    /// Not the whole queue — [`ClientStore::flush_order`] bounds what is due and
124    /// how many rows one pass takes, and a row still inside its backoff is asked
125    /// for again on the pass its deadline arrives.
126    pub fn pending(&self) -> Result<Vec<IntentRow>> {
127        Ok(self.store.flush_order()?)
128    }
129
130    pub fn next_delay(&self, attempt: u32) -> Duration {
131        let idx = attempt.saturating_sub(1) as usize;
132        Duration::from_millis(
133            self.backoff_ms
134                .get(idx)
135                .copied()
136                .or_else(|| self.backoff_ms.last().copied())
137                .unwrap_or(1_000),
138        )
139    }
140
141    /// Count one pass in which this process could not act on a row at all.
142    ///
143    /// A payload that no longer decodes will not decode on a later pass, and a
144    /// row left at its old deadline holds the head of the flush batch forever:
145    /// every pass spends a slot on a frame that can never be built, and behind a
146    /// batch cap that slot is one the queue's oldest sendable row did not get.
147    /// Charging the pass retires the row at [`MAX_LOCAL_FAILURES`] the way a
148    /// refused row retires at `attempts` — `exhausted`, with the fault that names
149    /// the reason — and the backoff meanwhile takes it out of the due window, so
150    /// the head moves on the same pass.
151    ///
152    /// A transport failure is not this path: the link being down says nothing
153    /// about the row, so [`IntentMachine::defer`] keeps its budget intact and the
154    /// reconnect flushes it (plan §6 line 289).
155    pub fn fail_local(&self, row: &IntentRow, reason: &str) -> Result<IntentResult> {
156        let failures = row.attempt.max(0) as u32 + 1;
157        if failures >= MAX_LOCAL_FAILURES {
158            self.store.exhaust_intent(&row.op_id, reason)?;
159            self.record_exhausted(&row.env_json, row.attempt, reason)?;
160            return Ok(IntentResult::Exhausted);
161        }
162        let due = Utc::now() + self.next_delay(failures);
163        self.store.bump_intent(&row.op_id, due, reason)?;
164        Ok(IntentResult::Retryable(
165            ErrorCode::Internal,
166            reason.to_string(),
167        ))
168    }
169
170    pub fn attempt(&self, row: &IntentRow, response: Option<&ResBody>) -> Result<IntentResult> {
171        let op_id = row.op_id.as_str();
172        let Some(body) = response else {
173            let next = row.attempt.saturating_add(1) as u32;
174            if next >= self.attempts {
175                self.store
176                    .exhaust_intent(op_id, "intent attempts exhausted")?;
177                self.record_exhausted(&row.env_json, row.attempt, "intent attempts exhausted")?;
178                return Ok(IntentResult::Exhausted);
179            }
180            let due = Utc::now() + self.next_delay(next);
181            self.store
182                .bump_intent(op_id, due, "connection unavailable")?;
183            return Ok(IntentResult::Retryable(
184                ErrorCode::Internal,
185                "connection unavailable".into(),
186            ));
187        };
188        self.apply_response(row, body)
189    }
190
191    /// Hold an intent whose send never reached the server.
192    ///
193    /// The link being down says nothing about the message, so the queue keeps
194    /// the row and the attempt counter stays where it was (plan §6 line 289).
195    pub fn defer(&self, row: &IntentRow, reason: &str) -> Result<IntentResult> {
196        let due = Utc::now() + self.next_delay(row.attempt.max(0) as u32 + 1);
197        self.store.defer_intent(&row.op_id, due, reason)?;
198        Ok(IntentResult::Retryable(
199            ErrorCode::Internal,
200            reason.to_string(),
201        ))
202    }
203
204    fn apply_response(
205        &self,
206        row: &IntentRow,
207        body: &onlyne_proto::ResBody,
208    ) -> Result<IntentResult> {
209        if body.ok {
210            let receipt = body
211                .data
212                .as_ref()
213                .and_then(|v| serde_json::from_value::<Receipt>(v.clone()).ok());
214            self.store
215                .accept_intent(&row.op_id, &body.data.clone().unwrap_or(Value::Null))?;
216            return Ok(IntentResult::Accepted(receipt));
217        }
218        let error = body
219            .error
220            .as_ref()
221            .context("error response missing payload")?;
222        if error.code == ErrorCode::Duplicate {
223            let receipt = body
224                .data
225                .as_ref()
226                .and_then(|v| serde_json::from_value::<Receipt>(v.clone()).ok());
227            self.store
228                .accept_intent(&row.op_id, &body.data.clone().unwrap_or(Value::Null))?;
229            return Ok(IntentResult::Accepted(receipt));
230        }
231        // A frame answered before the server session finished its `hello` names a
232        // window of the connection, so the row waits for the handshake instead of
233        // leaving the queue (plan §7 line 310's refusal). It waits with its
234        // deadline where it is rather than on a backoff rung: the rows behind it
235        // in the same batch are refused the same way, and the reconnect's own
236        // flush has to be able to send them (plan §7 line 310).
237        if error.message == onlyne_proto::HELLO_REQUIRED_MESSAGE {
238            return Ok(IntentResult::NotAuthenticated);
239        }
240        if PERMANENT_ERRORS.contains(&error.code) {
241            self.delete_intent(&row.op_id)?;
242            return Ok(IntentResult::Dropped(error.code, error.message.clone()));
243        }
244        let next = row.attempt.saturating_add(1) as u32;
245        if next >= self.attempts {
246            self.store.exhaust_intent(&row.op_id, &error.message)?;
247            self.record_exhausted(&row.env_json, row.attempt, &error.message)?;
248            Ok(IntentResult::Exhausted)
249        } else {
250            let due = Utc::now() + self.next_delay(next);
251            self.store.bump_intent(&row.op_id, due, &error.message)?;
252            Ok(IntentResult::Retryable(error.code, error.message.clone()))
253        }
254    }
255
256    fn delete_intent(&self, op_id: &str) -> Result<()> {
257        self.store.delete_intent(op_id)?;
258        Ok(())
259    }
260
261    fn record_exhausted(&self, env: &Value, attempt: i64, reason: &str) -> Result<()> {
262        let task = env
263            .get("causality")
264            .and_then(|v| v.get("task"))
265            .and_then(Value::as_str)
266            .unwrap_or("");
267        let _ = crate::reconcile::record_fault(
268            &self.store,
269            task,
270            "intent_exhausted",
271            "intent",
272            reason,
273        )?;
274        let _ = attempt;
275        let report = Report::Fault {
276            task_id: Some(task.to_string()),
277            session_id: None,
278            generation: None,
279            seq: None,
280            kind: "intent_exhausted".into(),
281            reason: reason.into(),
282            desired: None,
283            observed: None,
284        };
285        let _ = self.store.append_event(
286            "report_fault",
287            &serde_json::to_value(&report).unwrap_or(Value::Null),
288        );
289        Ok(())
290    }
291}
292
293pub fn op_for_intent(row: &IntentRow) -> Result<ClientOp> {
294    if let Ok(op) = serde_json::from_value::<ClientOp>(row.env_json.clone()) {
295        return Ok(op);
296    }
297    let envelope: Envelope = serde_json::from_value(row.env_json.clone())?;
298    Ok(ClientOp::Send(Box::new(envelope)))
299}
300
301pub fn due(row: &IntentRow) -> Result<DateTime<Utc>> {
302    Ok(DateTime::parse_from_rfc3339(&row.next_attempt_at)?.with_timezone(&Utc))
303}
304
305/// The session whose completion intent one queue row carries, when the row is
306/// one.
307///
308/// Two shapes and no others, and the restriction is a soundness rule rather
309/// than tidiness. A receipt closes the drain the completion is riding on, and
310/// the reducer only admits one from a drain that is open — so feeding a receipt
311/// for an accepted op that merely *names* a task would close the completion's
312/// drain on the strength of something else entirely. An ack settles a delivery,
313/// a ready or heartbeat report publishes a projection, and a note wakes an
314/// agent: none of them is the completion, and none of them may end it. The two
315/// that are, are the `Completion` envelope the settlement sends and the
316/// `complete` report the residual reconcile sends when it holds no envelope.
317///
318/// The session is keyed by the task the intent answers, which is what a
319/// client-held session's own row is keyed by. A row whose payload names no task
320/// — a note, or an ack — answers `None`, and the caller leaves the reducer
321/// alone rather than guessing a session.
322pub fn completion_task_id(row: &IntentRow) -> Option<String> {
323    if let Ok(envelope) = serde_json::from_value::<Envelope>(row.env_json.clone()) {
324        if envelope.kind != MsgKind::Completion {
325            return None;
326        }
327        return envelope.task_id().map(str::to_string);
328    }
329    match serde_json::from_value::<ClientOp>(row.env_json.clone()) {
330        Ok(ClientOp::Report(Report::Complete { task_id, .. })) => Some(task_id),
331        _ => None,
332    }
333}
334
335pub fn permanent_error(code: ErrorCode) -> bool {
336    PERMANENT_ERRORS.contains(&code)
337}