1use std::collections::HashMap;
2
3use parking_lot::{Mutex, MutexGuard};
4
5use serde_json::json;
6
7use crate::backend::SessionRef;
8use onlyne_proto::lifecycle::{self, LifecycleEvent, Observation, Verdict, Version};
9
10use super::fault::{DEFAULT_ISOLATE_AFTER, DEFAULT_TERMINATE_AFTER};
11use onlyne_store::session::{SessionLedger, SessionRecord};
12
13use super::record::{backend_ref_json, short, stored_observation, to_versioned};
14
15#[derive(Debug, Default)]
18pub struct Bridge {
19 pub(super) live: Mutex<HashMap<String, SessionRef>>,
20 applying: Mutex<()>,
30}
31
32impl Bridge {
33 pub fn new() -> Self {
34 Self::default()
35 }
36
37 fn applying(&self) -> MutexGuard<'_, ()> {
41 self.applying.lock()
42 }
43
44 pub fn track_live(&self, session: SessionRef) {
47 self.live.lock().insert(session.task_id.clone(), session);
48 }
49
50 pub fn untrack_live(&self, task_id: &str) {
52 self.live.lock().remove(task_id);
53 }
54}
55
56pub(super) fn initial_observation() -> Observation {
57 Observation::initial(DEFAULT_ISOLATE_AFTER, DEFAULT_TERMINATE_AFTER)
58}
59
60fn version_of(event: &LifecycleEvent) -> Version {
62 match event {
63 LifecycleEvent::Created { v }
64 | LifecycleEvent::Ready { v }
65 | LifecycleEvent::TurnStarted { v }
66 | LifecycleEvent::TurnEnded { v }
67 | LifecycleEvent::Heartbeat { v, .. }
68 | LifecycleEvent::Complete { v }
69 | LifecycleEvent::IntentPending { v }
70 | LifecycleEvent::IntentRetry { v }
71 | LifecycleEvent::IntentReceipt { v }
72 | LifecycleEvent::IntentExhausted { v }
73 | LifecycleEvent::ResourceAttach { v }
74 | LifecycleEvent::ResourceCloseRequested { v }
75 | LifecycleEvent::ResourceClosed { v }
76 | LifecycleEvent::Suspend { v }
77 | LifecycleEvent::Resume { v }
78 | LifecycleEvent::AgentGone { v }
79 | LifecycleEvent::Cancel { v }
80 | LifecycleEvent::Fail { v }
81 | LifecycleEvent::ReconcileMismatch { v }
82 | LifecycleEvent::ReconcileOk { v }
83 | LifecycleEvent::AdoptNewGeneration { v }
84 | LifecycleEvent::Supersede { v, .. } => *v,
85 }
86}
87
88fn event_name(event: &LifecycleEvent) -> &'static str {
89 match event {
90 LifecycleEvent::Created { .. } => "created",
91 LifecycleEvent::Ready { .. } => "ready",
92 LifecycleEvent::TurnStarted { .. } => "turn_started",
93 LifecycleEvent::TurnEnded { .. } => "turn_ended",
94 LifecycleEvent::Heartbeat { .. } => "heartbeat",
95 LifecycleEvent::Complete { .. } => "complete",
96 LifecycleEvent::IntentPending { .. } => "intent_pending",
97 LifecycleEvent::IntentRetry { .. } => "intent_retry",
98 LifecycleEvent::IntentReceipt { .. } => "intent_receipt",
99 LifecycleEvent::IntentExhausted { .. } => "intent_exhausted",
100 LifecycleEvent::ResourceAttach { .. } => "resource_attach",
101 LifecycleEvent::ResourceCloseRequested { .. } => "resource_close_requested",
102 LifecycleEvent::ResourceClosed { .. } => "resource_closed",
103 LifecycleEvent::Suspend { .. } => "suspend",
104 LifecycleEvent::Resume { .. } => "resume",
105 LifecycleEvent::AgentGone { .. } => "agent_gone",
106 LifecycleEvent::Cancel { .. } => "cancel",
107 LifecycleEvent::Fail { .. } => "fail",
108 LifecycleEvent::ReconcileMismatch { .. } => "reconcile_mismatch",
109 LifecycleEvent::ReconcileOk { .. } => "reconcile_ok",
110 LifecycleEvent::AdoptNewGeneration { .. } => "adopt_new_generation",
111 LifecycleEvent::Supersede { .. } => "supersede",
112 }
113}
114
115pub fn apply_persist(
134 bridge: &Bridge,
135 ledger: &dyn SessionLedger,
136 task_id: &str,
137 event: &LifecycleEvent,
138) -> anyhow::Result<Verdict> {
139 let (verdict, landed) = apply_persist_reported(bridge, ledger, task_id, event)?;
140 if !landed {
141 report_lost_write(ledger, task_id, event, &verdict, 1);
142 }
143 Ok(verdict)
144}
145
146fn apply_persist_reported(
152 bridge: &Bridge,
153 ledger: &dyn SessionLedger,
154 task_id: &str,
155 event: &LifecycleEvent,
156) -> anyhow::Result<(Verdict, bool)> {
157 let _gate = bridge.applying();
158 let row = ledger.get_session(task_id)?;
159 apply_round(bridge, ledger, task_id, row.as_ref(), event)
160}
161
162fn apply_round(
166 bridge: &Bridge,
167 ledger: &dyn SessionLedger,
168 task_id: &str,
169 row: Option<&SessionRecord>,
170 event: &LifecycleEvent,
171) -> anyhow::Result<(Verdict, bool)> {
172 if row.is_none() {
173 let known = ledger.task_is_known(task_id)? || bridge.live.lock().contains_key(task_id);
174 if !known {
175 tracing::warn!(
176 task = %task_id,
177 event = event_name(event),
178 "lifecycle event for a session never tracked; nothing persisted"
179 );
180 anyhow::bail!("unknown session {task_id}");
181 }
182 }
183 let current = stored_observation(ledger, row);
184 if row.is_none() && matches!(event, LifecycleEvent::Created { .. }) {
185 let (seeded, landed) = seed_created(bridge, ledger, task_id, version_of(event))?;
186 return Ok((Verdict::Applied(seeded), landed));
187 }
188 let stamped = if row.is_some() {
192 stamp_adoption(event, ¤t)
193 } else {
194 None
195 };
196 if let Some(adopting) = stamped.as_ref() {
197 tracing::warn!(
198 task = %task_id,
199 event = event_name(event),
200 from = version_of(event).seq,
201 to = version_of(adopting).seq,
202 watermark = current.version.seq,
203 "adoption carried a sequence the stored watermark had passed; advanced it past the gate"
204 );
205 }
206 let event = stamped.as_ref().unwrap_or(event);
207 let verdict = lifecycle::apply(¤t, event);
208 let landed = record_verdict(bridge, ledger, task_id, row, event, ¤t, &verdict)?;
209 Ok((verdict, landed))
210}
211
212fn stamp_adoption(event: &LifecycleEvent, current: &Observation) -> Option<LifecycleEvent> {
225 let v = version_of(event);
226 let adoption = matches!(
227 event,
228 LifecycleEvent::AdoptNewGeneration { .. } | LifecycleEvent::Supersede { .. }
229 );
230 if !adoption || v.generation != current.version.generation || v.seq > current.version.seq {
231 return None;
232 }
233 let newer = Version::new(v.generation, current.version.seq.saturating_add(1));
234 Some(match event {
235 LifecycleEvent::Supersede {
236 old_generation_dead,
237 body,
238 ..
239 } => LifecycleEvent::Supersede {
240 v: newer,
241 old_generation_dead: *old_generation_dead,
242 body: body.clone(),
243 },
244 _ => LifecycleEvent::AdoptNewGeneration { v: newer },
245 })
246}
247
248fn report_lost_write(
254 ledger: &dyn SessionLedger,
255 task_id: &str,
256 event: &LifecycleEvent,
257 verdict: &Verdict,
258 attempts: usize,
259) {
260 let v = version_of(event);
261 tracing::warn!(
262 task = %task_id,
263 event = event_name(event),
264 generation = v.generation,
265 seq = v.seq,
266 attempts,
267 applied = matches!(verdict, Verdict::Applied(_)),
268 "session transition did not land; the row keeps the older tuple"
269 );
270 ledger.note_alert(format!(
271 "session {} lost its {} write to a newer watermark",
272 short(task_id),
273 event_name(event)
274 ));
275 ledger.emit(
276 "lifecycle_write_lost",
277 json!({
278 "task_id": task_id,
279 "event": event_name(event),
280 "generation": v.generation,
281 "seq": v.seq,
282 "attempts": attempts,
283 }),
284 );
285}
286
287fn seed_created(
291 bridge: &Bridge,
292 ledger: &dyn SessionLedger,
293 task_id: &str,
294 version: Version,
295) -> anyhow::Result<(Observation, bool)> {
296 let mut seeded = initial_observation();
297 seeded.version = version;
298 let backend_ref = backend_ref_json(bridge, task_id, None);
299 let desired = serde_json::to_string(&LifecycleEvent::Created { v: version })?;
300 let stored = to_versioned(&seeded, &backend_ref, &desired)?;
301 if ledger.upsert_session(task_id, &stored)? {
302 ledger.emit("lifecycle", transition_payload(task_id, &seeded, "created"));
303 Ok((seeded, true))
304 } else {
305 tracing::warn!(
306 task = %task_id,
307 "a newer session row appeared while seeding Created; kept the newer watermark"
308 );
309 Ok((seeded, false))
310 }
311}
312
313fn transition_payload(task_id: &str, to: &Observation, event: &str) -> serde_json::Value {
318 json!({
319 "task_id": task_id,
320 "agent": to.agent,
321 "delivery": to.delivery,
322 "resource": to.resource,
323 "recovery": to.recovery,
324 "generation": to.version.generation,
325 "seq": to.version.seq,
326 "event": event,
327 })
328}
329
330#[allow(clippy::too_many_arguments)]
331fn record_verdict(
332 bridge: &Bridge,
333 ledger: &dyn SessionLedger,
334 task_id: &str,
335 row: Option<&SessionRecord>,
336 event: &LifecycleEvent,
337 current: &Observation,
338 verdict: &Verdict,
339) -> anyhow::Result<bool> {
340 match verdict {
341 Verdict::Applied(next) => {
342 let backend_ref = backend_ref_json(bridge, task_id, row);
343 let desired = serde_json::to_string(event)?;
344 let stored = to_versioned(next, &backend_ref, &desired)?;
345 if !ledger.upsert_session(task_id, &stored)? {
346 tracing::warn!(
347 task = %task_id,
348 event = event_name(event),
349 generation = next.version.generation,
350 seq = next.version.seq,
351 "session write lost to a newer watermark; left the row alone"
352 );
353 return Ok(false);
354 }
355 tracing::debug!(
356 task = %task_id,
357 agent = ?next.agent,
358 delivery = ?next.delivery,
359 resource = ?next.resource,
360 recovery = ?next.recovery,
361 generation = next.version.generation,
362 seq = next.version.seq,
363 "lifecycle transition applied"
364 );
365 ledger.emit(
366 "lifecycle",
367 transition_payload(task_id, next, event_name(event)),
368 );
369 Ok(true)
370 }
371 Verdict::Ignored(reason) => {
372 tracing::debug!(
373 task = %task_id,
374 reason = ?reason,
375 event = event_name(event),
376 generation = current.version.generation,
377 seq = current.version.seq,
378 "lifecycle event ignored; ledger left as-is"
379 );
380 Ok(true)
381 }
382 Verdict::Rejected(reason) => {
383 tracing::warn!(
384 task = %task_id,
385 reason = ?reason,
386 event = event_name(event),
387 generation = current.version.generation,
388 seq = current.version.seq,
389 "lifecycle event rejected; ledger left as-is"
390 );
391 Ok(true)
392 }
393 }
394}
395
396pub fn next_version(ledger: &dyn SessionLedger, task_id: &str) -> anyhow::Result<Version> {
403 let row = ledger.get_session(task_id)?;
404 let current = stored_observation(ledger, row.as_ref());
405 Ok(next_after(¤t))
406}
407
408fn next_after(current: &Observation) -> Version {
410 Version::new(
411 current.version.generation,
412 current.version.seq.saturating_add(1),
413 )
414}
415
416pub fn apply_at_next(
420 bridge: &Bridge,
421 ledger: &dyn SessionLedger,
422 task_id: &str,
423 make: impl FnMut(Version) -> LifecycleEvent,
424) -> anyhow::Result<Verdict> {
425 let mut make = make;
426 apply_locally(bridge, ledger, task_id, |_, v| make(v))
427}
428
429pub fn apply_from_stored(
438 bridge: &Bridge,
439 ledger: &dyn SessionLedger,
440 task_id: &str,
441 make: impl FnMut(&Observation, Version) -> LifecycleEvent,
442) -> anyhow::Result<Verdict> {
443 apply_locally(bridge, ledger, task_id, make)
444}
445
446fn apply_locally(
456 bridge: &Bridge,
457 ledger: &dyn SessionLedger,
458 task_id: &str,
459 mut make: impl FnMut(&Observation, Version) -> LifecycleEvent,
460) -> anyhow::Result<Verdict> {
461 let _gate = bridge.applying();
462 let mut attempt = 1;
463 loop {
464 let row = ledger.get_session(task_id)?;
465 let current = stored_observation(ledger, row.as_ref());
466 let event = make(¤t, next_after(¤t));
467 let (verdict, landed) = apply_round(bridge, ledger, task_id, row.as_ref(), &event)?;
468 if landed {
469 return Ok(verdict);
470 }
471 if attempt == APPLY_ATTEMPTS {
472 report_lost_write(ledger, task_id, &event, &verdict, attempt);
473 return Ok(verdict);
474 }
475 tracing::debug!(
476 task = %task_id,
477 event = event_name(&event),
478 attempt,
479 "write lost to a competing watermark; replaying the round on the fresher row"
480 );
481 attempt += 1;
482 }
483}
484
485const APPLY_ATTEMPTS: usize = 4;
490
491pub fn try_feed(
493 bridge: &Bridge,
494 ledger: &dyn SessionLedger,
495 task_id: &str,
496 make: impl FnMut(Version) -> LifecycleEvent,
497) {
498 if let Err(err) = apply_at_next(bridge, ledger, task_id, make) {
499 tracing::warn!(
500 task = %task_id,
501 error = %err,
502 "could not record a lifecycle event in the session ledger"
503 );
504 }
505}