onlyne_client/session/dispatch/projection.rs
1use super::*;
2
3use super::outbound::send_frame;
4use super::state::{DispatchInner, DispatchState};
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 reply_to,
59 ..
60 } => Report::Complete {
61 task_id,
62 outcome,
63 head,
64 reply_to,
65 cluster_ref: Some(cluster),
66 },
67 other => other,
68 }
69}
70
71/// The wire projection of one stored session row, given how its task ended.
72///
73/// The row and the task record answer different questions. The row holds the
74/// session's own tuple — agent, intent, resource, recovery line — and the task
75/// table holds the verdict the agent filed. Neither one says whether the session
76/// is over on its own, so nothing here is read as a lifecycle: `project` takes
77/// the tuple's dimensions plus `task_state` and derives the public view. A
78/// caller that wants a projection reads both rows.
79pub fn projection_of(row: &SessionRecord, task_state: TaskState) -> SessionProjection {
80 let observed: Option<serde_json::Value> = serde_json::from_str(&row.observed_json).ok();
81 // The dimension columns hold the reducer's own words, written from the same
82 // observation as `observed_json`, so the derivation reads them rather than
83 // the JSON bytes. A word that does not decode falls back to the freshly
84 // created dimension, which is the reading the wire columns below give too.
85 let lifecycle = wire_lifecycle(project(
86 phase(&row.agent_state, AgentState::Booting),
87 phase(&row.delivery_state, DeliveryState::None),
88 phase(&row.resource_state, ResourceState::Detached),
89 phase(&row.recovery_substate, RecoveryState::None),
90 task_state,
91 ));
92 SessionProjection {
93 lifecycle,
94 agent: phase(&row.agent_state, AgentPhase::Booting),
95 delivery: phase(&row.delivery_state, DeliveryPhase::NoIntent),
96 resource: phase(&row.resource_state, ResourcePhase::Detached),
97 recovery: phase(&row.recovery_substate, RecoveryPhase::NoRecovery),
98 outcome: task_outcome_of(task_state),
99 observed,
100 }
101}
102
103/// The wire word for a derived lifecycle. `PublicLifecycle` and `Lifecycle` are
104/// the same view named in two crates, and a match is what notices if one of them
105/// grows an arm the other does not have.
106fn wire_lifecycle(lifecycle: PublicLifecycle) -> Lifecycle {
107 match lifecycle {
108 PublicLifecycle::Created => Lifecycle::Created,
109 PublicLifecycle::Working => Lifecycle::Working,
110 PublicLifecycle::Idle => Lifecycle::Idle,
111 PublicLifecycle::Exited => Lifecycle::Exited,
112 }
113}
114
115/// The task state a reported outcome names. The wire has no word for work still
116/// in flight — a completion either says how it ended or is not a completion — so
117/// `pending` comes from the absence of a verdict, not from a report.
118pub fn task_state_of(outcome: Outcome) -> TaskState {
119 match outcome {
120 Outcome::Done => TaskState::Done,
121 Outcome::Failed => TaskState::Failed,
122 Outcome::Cancelled => TaskState::Cancelled,
123 }
124}
125
126/// The wire verdict for one task state, which is what a published projection
127/// carries beside its derived lifecycle. `pending` has no wire word.
128pub fn task_outcome_of(task_state: TaskState) -> Option<Outcome> {
129 match task_state {
130 TaskState::Pending => None,
131 TaskState::Done => Some(Outcome::Done),
132 TaskState::Failed => Some(Outcome::Failed),
133 TaskState::Cancelled => Some(Outcome::Cancelled),
134 }
135}
136
137/// How the task of one stored session ended, read from its own record.
138///
139/// A task with no record is `pending`: the open and the settle both write the
140/// row, so nothing having been written about a task means no delivery became a
141/// session for it and no verdict came in for it. Callers hold the dispatch guard,
142/// which this read needs to reach the store; the lock is not reentrant.
143pub(super) fn stored_task_state(inner: &DispatchInner, task_id: &str) -> TaskState {
144 inner
145 .store
146 .task(task_id)
147 .ok()
148 .flatten()
149 .map(|record| record.task_state)
150 .unwrap_or(TaskState::Pending)
151}
152
153/// Decode one stored enum word, falling back to the freshly created phase.
154pub(super) fn phase<T: serde::de::DeserializeOwned>(word: &str, fallback: T) -> T {
155 serde_json::from_value(serde_json::Value::String(word.to_string())).unwrap_or(fallback)
156}
157
158/// Publish the current projection of one session.
159///
160/// The frame is the client's heartbeat report carrying the whole projection:
161/// one frame per session activity, and the only one that puts a session's state
162/// on the wire. `session_id` is the row the server keys the mirror by, which for
163/// a client-held session is its task id.
164pub async fn sync_session(state: &DispatchState, task_id: &str) -> Result<()> {
165 // One section for both reads: a publish that took the session tuple before a
166 // settle and its verdict after would derive a lifecycle the pair never agreed
167 // to, and the store's lock is what keeps the two rows in step.
168 let (row, task_state) = {
169 let inner = state.inner.lock();
170 (
171 inner.store.get_session(task_id)?,
172 stored_task_state(&inner, task_id),
173 )
174 };
175 let Some(row) = row else { return Ok(()) };
176 let projection = projection_of(&row, task_state);
177 send_frame(
178 state,
179 ClientOp::Report(Report::Heartbeat {
180 task_id: row.task_id.clone(),
181 session_id: row.task_id.clone(),
182 generation: row.generation.max(0) as u64,
183 seq: row.seq.max(0) as u64,
184 observed: projection
185 .observed
186 .clone()
187 .unwrap_or(serde_json::Value::Null),
188 cluster_ref: None,
189 projection: Some(projection),
190 }),
191 )
192 .await
193}
194
195/// Log a reducer verdict and answer the version it advanced to.
196pub fn note_verdict(verdict: &Verdict, task_id: &str) -> Option<Version> {
197 match verdict {
198 Verdict::Applied(observation) => Some(observation.version),
199 Verdict::Ignored(reason) => {
200 tracing::debug!(task = %task_id, ?reason, "lifecycle event ignored");
201 None
202 }
203 Verdict::Rejected(reason) => {
204 tracing::warn!(task = %task_id, ?reason, "lifecycle event rejected; the ledger kept its state");
205 None
206 }
207 }
208}
209
210/// Feed the reducer the receipt the outbound queue observed for one session.
211///
212/// The durable queue is the only witness of the server's answer, so this is the
213/// door a receipt comes in by: the flusher learns that the completion was taken
214/// and the tuple records it as a fact it observed rather than as a body this
215/// client composed for itself. A session with no row, and a drain that is not
216/// open, are the reducer's own answers and leave the ledger alone — the verdict
217/// is logged like every other feed's and nothing is retried here.
218pub fn note_intent_receipt(state: &DispatchState, task_id: &str) {
219 let verdict = {
220 let inner = state.inner.lock();
221 feed_intent_receipt(&inner.bridge, &inner.store, task_id)
222 };
223 match verdict {
224 Ok(verdict) => {
225 note_verdict(&verdict, task_id);
226 }
227 Err(error) => tracing::debug!(
228 task = %task_id,
229 error = %error,
230 "intent receipt not recorded in the session ledger"
231 ),
232 }
233}