Skip to main content

onlyne_client/session/dispatch/
projection.rs

1use super::*;
2
3use super::outbound::send_frame;
4use super::state::{DispatchInner, DispatchState, binding_task_state, slot_key_named};
5
6/// Stamp the origin cluster on the state-carrying report kinds.
7///
8/// `cluster_ref` names the cluster whose supervisor observed this projection
9/// (`docs/v1-PLAN.md` line 248, the `aggregate` annotation of the role's own
10/// spec entry, which is `Principal::Cluster` on the wire at line 122). A plain
11/// role leaves the field unset, which is the `skip_serializing_if` shape the
12/// byte-identical replay rule at line 502 depends on.
13///
14/// A heartbeat that carries a projection is the one state-carrying report this
15/// client builds rather than relays, and it stays unstamped: the server mirrors
16/// such a projection verbatim, and the publish has never named an origin
17/// cluster.
18pub fn with_cluster(state: &DispatchState, report: Report) -> Report {
19    let cluster = state.cluster_ref();
20    if cluster.is_empty() {
21        return report;
22    }
23    match report {
24        Report::Ready {
25            task_id,
26            session_id,
27            generation,
28            seq,
29            ..
30        } => Report::Ready {
31            task_id,
32            session_id,
33            generation,
34            seq,
35            cluster_ref: Some(cluster),
36        },
37        Report::Heartbeat {
38            projection: None,
39            task_id,
40            session_id,
41            generation,
42            seq,
43            observed,
44            ..
45        } => Report::Heartbeat {
46            task_id,
47            session_id,
48            generation,
49            seq,
50            observed,
51            projection: None,
52            cluster_ref: Some(cluster),
53        },
54        Report::Complete {
55            task_id,
56            outcome,
57            head,
58            details,
59            files,
60            reply_to,
61            cluster_ref: _,
62        } => Report::Complete {
63            task_id,
64            outcome,
65            head,
66            details,
67            files,
68            reply_to,
69            cluster_ref: Some(cluster),
70        },
71        other => other,
72    }
73}
74
75/// The wire projection of one stored session row, given how its task ended.
76///
77/// The row and the task record answer different questions. The row holds the
78/// session's own tuple — agent, intent, resource, recovery line — and the task
79/// table holds the verdict the agent filed. Neither one says whether the session
80/// is over on its own, so nothing here is read as a lifecycle: `project` takes
81/// the tuple's dimensions plus `task_state` and derives the public view. A
82/// caller that wants a projection reads both rows.
83pub fn projection_of(row: &SessionRecord, task_state: TaskState) -> SessionProjection {
84    let observed: Option<serde_json::Value> = serde_json::from_str(&row.observed_json).ok();
85    // The dimension columns hold the reducer's own words, written from the same
86    // observation as `observed_json`, so the derivation reads them rather than
87    // the JSON bytes. A word that does not decode falls back to the freshly
88    // created dimension, which is the reading the wire columns below give too.
89    let lifecycle = project(
90        phase(&row.agent_state, AgentPhase::Booting),
91        phase(&row.delivery_state, DeliveryPhase::NoIntent),
92        phase(&row.resource_state, ResourcePhase::Detached),
93        phase(&row.recovery_substate, RecoveryPhase::NoRecovery),
94        task_state,
95    );
96    SessionProjection {
97        lifecycle,
98        agent: phase(&row.agent_state, AgentPhase::Booting),
99        delivery: phase(&row.delivery_state, DeliveryPhase::NoIntent),
100        resource: phase(&row.resource_state, ResourcePhase::Detached),
101        recovery: phase(&row.recovery_substate, RecoveryPhase::NoRecovery),
102        outcome: task_outcome_of(task_state),
103        observed,
104    }
105}
106
107/// The task state a reported outcome names. The wire has no word for work still
108/// in flight — a completion either says how it ended or is not a completion — so
109/// `pending` comes from the absence of a verdict, not from a report.
110pub fn task_state_of(outcome: Outcome) -> TaskState {
111    match outcome {
112        Outcome::Done => TaskState::Done,
113        Outcome::Failed => TaskState::Failed,
114        Outcome::Cancelled => TaskState::Cancelled,
115        // A blocked delivery is a settled delivery whose work waits on
116        // something outside it, which is the task record's own blocked state.
117        Outcome::Blocked => TaskState::Blocked,
118    }
119}
120
121/// The wire verdict for one task state, which is what a published projection
122/// carries beside its derived lifecycle. `pending` has no wire word.
123pub fn task_outcome_of(task_state: TaskState) -> Option<Outcome> {
124    match task_state {
125        TaskState::Pending => None,
126        TaskState::Done => Some(Outcome::Done),
127        TaskState::Failed => Some(Outcome::Failed),
128        TaskState::Cancelled => Some(Outcome::Cancelled),
129        TaskState::Blocked => Some(Outcome::Blocked),
130    }
131}
132
133/// How the task of one stored session ended, read from its own record.
134///
135/// A task with no record is `pending`: the open and the settle both write the
136/// row, so nothing having been written about a task means no delivery became a
137/// session for it and no verdict came in for it. Callers hold the dispatch guard,
138/// which this read needs to reach the store; the lock is not reentrant.
139pub(super) fn stored_task_state(inner: &DispatchInner, task_id: &str) -> TaskState {
140    inner
141        .store
142        .task(task_id)
143        .ok()
144        .flatten()
145        .map(|record| record.task_state)
146        .unwrap_or(TaskState::Pending)
147}
148
149/// Decode one stored enum word, falling back to the freshly created phase.
150pub(super) fn phase<T: serde::de::DeserializeOwned>(word: &str, fallback: T) -> T {
151    serde_json::from_value(serde_json::Value::String(word.to_string())).unwrap_or(fallback)
152}
153
154/// Publish the current projection of one session.
155///
156/// The frame is the client's heartbeat report carrying the whole projection:
157/// one frame per session activity, and the only one that puts a session's state
158/// on the wire. The argument names the session, or the delivery a caller has in
159/// hand: a client-held session takes its id from the delivery that opened it, so
160/// both spellings reach the same row and the slot the client holds decides which
161/// session that is.
162pub async fn sync_session(state: &DispatchState, session_id: &str) -> Result<()> {
163    let Some(op) = sync_frame(state, session_id)? else {
164        return Ok(());
165    };
166    send_frame(state, op).await
167}
168
169/// The durable frame one session's current projection publishes, or `None`
170/// when the store holds no row for it.
171///
172/// `sync_session` sends this through `send_frame`, which already queues it when
173/// the link is down; the frame is exposed separately for the caller whose send
174/// failed for some other reason and owes the exit a second attempt through
175/// [`DispatchState::enqueue_op`].
176///
177/// The session's own id is what the server keys the mirror by, and it stays put
178/// while a scope hands the session delivery after delivery. The delivery travels
179/// beside it as `task_id`: the one the session is bound to now, or the last one
180/// it served once the binding is released. Whether the session is between
181/// deliveries is the projection's own delivery dimension, not the id's — so the
182/// mirror keeps reading the delivery a row belongs to while the pair says the
183/// session is live.
184///
185/// A session whose row was never given a delivery sends an empty `task_id`,
186/// which is the wire's only spelling for "this session is on no delivery"; a
187/// client-held session always has the delivery that opened it.
188pub fn sync_frame(state: &DispatchState, session_id: &str) -> Result<Option<ClientOp>> {
189    // One section for both reads: a publish that took the session tuple before a
190    // settle and its verdict after would derive a lifecycle the pair never agreed
191    // to, and the store's lock is what keeps the two rows in step.
192    let (key, row, task_state) = {
193        let inner = state.inner.lock();
194        let key = slot_key_named(&inner, session_id);
195        let task_state = match key.as_deref().and_then(|key| inner.sessions.get(key)) {
196            // A session this client holds answers from its slot: the delivery it
197            // serves now is the binding, and a session serving nothing reads as
198            // `pending`, which is a live session rather than an exit.
199            Some(slot) => binding_task_state(&inner, slot),
200            // A session whose slot is gone is read from the delivery the caller
201            // named, which is how a retired session's row stays publishable.
202            None => match inner.store.get_session(session_id)? {
203                Some(row) => stored_task_state(&inner, &row.task_id),
204                None => TaskState::Pending,
205            },
206        };
207        (key, inner.store.get_session(session_id)?, task_state)
208    };
209    let Some(row) = row else { return Ok(None) };
210    let session_id = key.unwrap_or_else(|| row.session_id.clone());
211    let projection = projection_of(&row, task_state);
212    Ok(Some(ClientOp::Report(Report::Heartbeat {
213        task_id: row.task_id.clone(),
214        session_id,
215        generation: row.generation.max(0) as u64,
216        seq: row.seq.max(0) as u64,
217        observed: projection
218            .observed
219            .clone()
220            .unwrap_or(serde_json::Value::Null),
221        cluster_ref: None,
222        projection: Some(projection),
223    })))
224}
225
226/// Log a reducer verdict and answer the version it advanced to.
227pub fn note_verdict(verdict: &Verdict, task_id: &str) -> Option<Version> {
228    match verdict {
229        Verdict::Applied(observation) => Some(observation.version),
230        Verdict::Ignored(reason) => {
231            tracing::debug!(task = %task_id, ?reason, "lifecycle event ignored");
232            None
233        }
234        Verdict::Rejected(reason) => {
235            tracing::warn!(task = %task_id, ?reason, "lifecycle event rejected; the ledger kept its state");
236            None
237        }
238    }
239}
240
241/// Feed the reducer the receipt the outbound queue observed for one session.
242///
243/// The durable queue is the only witness of the server's answer, so this is the
244/// door a receipt comes in by: the flusher learns that the completion was taken
245/// and the tuple records it as a fact it observed rather than as a body this
246/// client composed for itself. A session with no row, and a drain that is not
247/// open, are the reducer's own answers and leave the ledger alone — the verdict
248/// is logged like every other feed's and nothing is retried here.
249pub fn note_intent_receipt(state: &DispatchState, task_id: &str) {
250    let verdict = {
251        let inner = state.inner.lock();
252        feed_intent_receipt(&inner.bridge, &inner.store, task_id)
253    };
254    match verdict {
255        Ok(verdict) => {
256            note_verdict(&verdict, task_id);
257        }
258        Err(error) => tracing::debug!(
259            task = %task_id,
260            error = %error,
261            "intent receipt not recorded in the session ledger"
262        ),
263    }
264}