Skip to main content

onlyne_client/reconcile/
bridge.rs

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/// A bridge instance: the reducer facts plus the ledger and the live session
16/// map it resolves probe targets from.
17#[derive(Debug, Default)]
18pub struct Bridge {
19    pub(super) live: Mutex<HashMap<String, SessionRef>>,
20    /// One session transaction at a time: the stored watermark is read, reduced
21    /// and written back while this is held, so the writers that share a bridge —
22    /// the plugin's beat, the stall clock, the settle, the sweep — cannot be
23    /// overtaken inside one another's window. The section spans two ledger
24    /// round-trips, which the store serializes anyway. A lock taken before any
25    /// other bridge lock and never taken from inside a ledger call, so the
26    /// ordering `applying` → `live` is the only one that can form. Unpoisoned by
27    /// construction: a panic inside a ledger call leaves the gate behind rather
28    /// than turning the next apply into a second panic.
29    applying: Mutex<()>,
30}
31
32impl Bridge {
33    pub fn new() -> Self {
34        Self::default()
35    }
36
37    /// Enter a reduce-and-persist transaction. Every public entry point in this
38    /// module takes the gate exactly once, at its own boundary, and the inner
39    /// rounds never re-enter it.
40    fn applying(&self) -> MutexGuard<'_, ()> {
41        self.applying.lock()
42    }
43
44    /// Remember a live session ref. The bridge prefers it over the stored
45    /// reference when it names the backend resource.
46    pub fn track_live(&self, session: SessionRef) {
47        self.live.lock().insert(session.task_id.clone(), session);
48    }
49
50    /// Forget a live session ref.
51    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
60/// The version an event carries.
61fn 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
115/// Reduce `event` against the stored tuple for `task_id` and persist the result.
116///
117/// The read, the reduce and the write are one transaction on the bridge's apply
118/// gate, so a writer that shares this bridge cannot take the watermark between
119/// them. A loss to a writer this bridge cannot see — a second client on the same
120/// store — is not swallowed: the caller handed over its own version, so there is
121/// no newer one to allocate on its behalf, and the loss is escalated instead
122/// (see [`report_lost_write`]).
123///
124/// `Applied` writes the row (unless a newer watermark already carries it) and emits one
125/// `lifecycle` event describing the transition.
126/// `Ignored` and `Rejected` leave the ledger untouched and are logged with the
127/// reducer's reason.
128///
129/// `Created` is the one event whose meaning at this boundary is the row itself:
130/// `Observation::initial` already is the post-`Created` tuple, so for a task
131/// with no stored row it seeds the row at the event's version. Use
132/// [`feed_created`] for the idempotent form.
133pub 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
146/// Reduce and persist, reporting alongside the verdict whether the write landed.
147///
148/// One implementation of the write, shared by every path. `false` means a
149/// transition the reducer accepted whose write lost the watermark; a verdict of
150/// `Ignored` or `Rejected` has nothing to write and reports `true`.
151fn 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
162/// One reduce-and-persist round against `row`, the watermark the caller read
163/// inside the same transaction. The gate is held by the caller, so this is the
164/// only place the write's answer is produced.
165fn 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    // An adoption is the one event the reducer exempts from the sequence gate,
189    // and the exemption does not reach the ledger: `upsert_session` refuses
190    // anything not strictly newer. Restamping is what makes the two gates agree.
191    let stamped = if row.is_some() {
192        stamp_adoption(event, &current)
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(&current, event);
208    let landed = record_verdict(bridge, ledger, task_id, row, event, &current, &verdict)?;
209    Ok((verdict, landed))
210}
211
212/// Re-version an adoption that the write gate would refuse.
213///
214/// `AdoptNewGeneration` and `Supersede` claim a generation, not a point in its
215/// sequence: the reducer says so by letting them through the same-generation
216/// sequence gate. A generation that is already the stored one is the case where
217/// that leaves the version at or behind the watermark the ledger compares, so
218/// the reducer answers `Applied` and `upsert_session` refuses the row — the
219/// adoption lands nowhere and the caller is told it did. Moving it one sequence
220/// past the stored watermark keeps the claim (same generation, newer sequence)
221/// and gives the gate an answer it can take. A generation strictly behind the
222/// stored one is not restamped: the reducer's own staleness gate is the right
223/// answer there.
224fn 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
248/// Escalate an accepted transition whose write the ledger refused, once retries
249/// cannot save it. The warn naming the loss is `record_verdict`'s; this is the
250/// part that reaches outside the process: an operator line on the client's alert
251/// surface and one event on the bus, so the server projecting this row learns
252/// the tuple it is showing is not the tuple the reducer settled on.
253fn 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
287/// Write the `Created` seed row. The reducer defines no transition into
288/// `Booting` from `Booting` — `initial()` is already the created tuple — so the
289/// seed is written directly at the event's version.
290fn 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
313/// One transition on the bus: the full session tuple the reducer settled on.
314/// No public view rides it — that needs the task state, which is not a session
315/// fact, so a subscriber that shows one derives it from these dimensions plus
316/// the task it is tracking.
317fn 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
396/// The next version for a session observed locally: same generation, one
397/// sequence past the stored watermark.
398///
399/// A standalone read, and the caller's write is its own transaction: a writer
400/// can take the watermark between the two. The bridge's local writers allocate
401/// inside their transaction instead of going through here.
402pub 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(&current))
406}
407
408/// One sequence past a watermark, in the watermark's own generation.
409fn next_after(current: &Observation) -> Version {
410    Version::new(
411        current.version.generation,
412        current.version.seq.saturating_add(1),
413    )
414}
415
416/// Reduce an event whose version is the stored watermark plus one. Every feed
417/// helper below goes through here, so an event that arrives twice is dropped by
418/// the reducer's watermark.
419pub 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
429/// Reduce an event the caller composes from the stored tuple.
430///
431/// For a caller whose body is derived from the row — a settlement replays the
432/// tuple it read with one dimension closed — reading outside the transaction is
433/// the race: the tuple it composed against can be the one the watermark has
434/// already left, and the reducer then answers a version its own reader never
435/// saw. Here the tuple and the version come from the same read that holds the
436/// apply gate, which is the only read that can promise either.
437pub 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
446/// The local transaction: read the row, let the caller compose its event from
447/// what that read shows, reduce, and write — all of it inside the apply gate, so
448/// the writers sharing this bridge cannot overtake each other between the read
449/// and the write. A write that still loses was taken by a writer this bridge
450/// cannot see (a second client on the same store); the round is replayed on the
451/// fresher row, bounded by [`APPLY_ATTEMPTS`], and the loss is escalated when the
452/// budget is spent. The reducer's own gate is untouched: a duplicate or an older
453/// event is still dropped, and a verdict of `Ignored` or `Rejected` has nothing
454/// to replay.
455fn 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(&current, next_after(&current));
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
485/// Attempts one local event gets to land its write. Four is the number that
486/// covers the writers that can share one store — the plugin's beat, the stall
487/// clock, the settle, and the sweep — one retry each, after which the loss is
488/// escalated rather than looped over.
489const APPLY_ATTEMPTS: usize = 4;
490
491/// Best-effort feed for local observation points.
492pub 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}