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
8pub 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
24pub 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 NotAuthenticated,
69}
70
71pub 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 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 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 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 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 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
305pub 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}