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}