Skip to main content

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}