1use std::collections::{HashMap, VecDeque};
4
5use crate::{
6 CoreError,
7 events::{
8 AcknowledgeResult, CommitReceipt, SessionEvent, SessionStatus, SnapshotRequiredReason,
9 TerminalEffect,
10 },
11 input::{InputId, InputOutcome, InputRequest},
12 types::{StateRevision, TerminalDelta, TerminalSnapshot, TerminalState},
13};
14
15#[derive(Clone, Copy, Debug, Eq, PartialEq)]
19pub struct SessionOptions {
20 pub event_capacity: usize,
21}
22
23impl Default for SessionOptions {
24 fn default() -> Self {
25 Self {
26 event_capacity: crate::types::MAX_EVENT_QUEUE,
27 }
28 }
29}
30
31#[derive(Clone, Copy, Debug, Eq, PartialEq)]
32enum ReceiptState {
33 Pending,
34 Committed,
35}
36
37#[derive(Debug)]
40pub struct Session {
41 options: SessionOptions,
42 status: SessionStatus,
43 state: Option<TerminalState>,
44 revision: StateRevision,
45 events: VecDeque<SessionEvent>,
46 receipts: HashMap<CommitReceipt, ReceiptState>,
47 next_receipt: u64,
48 next_input: u64,
49 pending_input: VecDeque<InputRequest>,
50 input_generation: u64,
51}
52
53impl Session {
54 pub fn new(options: SessionOptions) -> Self {
55 Self {
56 options: SessionOptions {
57 event_capacity: options.event_capacity.max(1),
58 },
59 status: SessionStatus::Connecting,
60 state: None,
61 revision: StateRevision::new(0),
62 events: VecDeque::new(),
63 receipts: HashMap::new(),
64 next_receipt: 0,
65 next_input: 0,
66 pending_input: VecDeque::new(),
67 input_generation: 0,
68 }
69 }
70
71 pub fn status(&self) -> SessionStatus {
72 self.status
73 }
74
75 pub fn revision(&self) -> StateRevision {
76 self.revision
77 }
78
79 pub fn state(&self) -> Option<&TerminalState> {
80 self.state.as_ref()
81 }
82
83 pub fn install_snapshot(&mut self, snapshot: TerminalSnapshot) -> Result<(), CoreError> {
84 if self.status == SessionStatus::Closed {
85 return Err(CoreError::NotLive);
86 }
87 if !self.can_enqueue() {
88 return Err(CoreError::ConsumerStalled);
89 }
90 self.state = Some(snapshot.state.clone());
91 self.revision = snapshot.revision;
92 self.status = SessionStatus::Live;
93 self.events.push_back(SessionEvent::Snapshot { snapshot });
94 if self.can_enqueue() {
95 self.events
96 .push_back(SessionEvent::StatusChanged(SessionStatus::Live));
97 }
98 Ok(())
99 }
100
101 pub fn apply_delta(&mut self, delta: &TerminalDelta) -> Result<(), CoreError> {
102 if self.status == SessionStatus::Closed {
103 return Err(CoreError::NotLive);
104 }
105 let Some(current) = self.state.as_ref() else {
106 self.enqueue_snapshot_required(
107 StateRevision::new(0),
108 delta.revision,
109 SnapshotRequiredReason::NoBase,
110 )?;
111 return Err(CoreError::NoSnapshot);
112 };
113 if delta.base_revision != self.revision
114 || delta.revision.get() != self.revision.get().saturating_add(1)
115 {
116 self.enqueue_snapshot_required(
117 self.revision.next(),
118 delta.revision,
119 SnapshotRequiredReason::RevisionGap,
120 )?;
121 return Err(CoreError::RevisionGap {
122 expected: self.revision.next(),
123 got: delta.revision,
124 });
125 }
126 if delta.state.dimensions != current.dimensions {
127 self.enqueue_snapshot_required(
128 self.revision.next(),
129 delta.revision,
130 SnapshotRequiredReason::InvalidDelta,
131 )?;
132 return Err(CoreError::DimensionsMismatch);
133 }
134 let mut next = current.clone();
135 delta.state.apply_to(&mut next)?;
136 let effects = effects_for(&delta.state, current);
137 let receipt = (!effects.is_empty()).then(|| self.issue_receipt());
138 if let Some(receipt) = receipt {
139 self.receipts.insert(receipt, ReceiptState::Pending);
140 }
141 if !self.can_enqueue() {
142 if let Some(receipt) = receipt {
143 self.receipts.remove(&receipt);
144 }
145 return Err(CoreError::ConsumerStalled);
146 }
147 self.state = Some(next);
148 self.revision = delta.revision;
149 self.events.push_back(SessionEvent::Delta {
150 delta: delta.clone(),
151 effects,
152 receipt,
153 });
154 Ok(())
155 }
156
157 pub fn acknowledge(&mut self, receipt: CommitReceipt) -> AcknowledgeResult {
158 match self.receipts.get_mut(&receipt) {
159 Some(state @ ReceiptState::Pending) => {
160 *state = ReceiptState::Committed;
161 AcknowledgeResult::Committed
162 }
163 Some(ReceiptState::Committed) => AcknowledgeResult::Duplicate,
164 None if receipt.0 <= self.next_receipt => AcknowledgeResult::RejectedStale,
165 None => AcknowledgeResult::RejectedUnknown,
166 }
167 }
168
169 pub fn acknowledge_raw(&mut self, receipt: u64) -> AcknowledgeResult {
172 self.acknowledge(CommitReceipt(receipt))
173 }
174
175 pub fn submit_input(&mut self, bytes: &[u8]) -> Result<InputId, CoreError> {
176 if self.status != SessionStatus::Live {
177 return Err(CoreError::NotLive);
178 }
179 let id = InputId(self.next_input);
180 self.next_input = self.next_input.saturating_add(1);
181 let request = InputRequest::new(id, bytes.to_vec())?;
182 if !self.can_enqueue() {
183 return Err(CoreError::ConsumerStalled);
184 }
185 self.pending_input.push_back(request.clone());
186 self.events.push_back(SessionEvent::InputQueued { request });
187 Ok(id)
188 }
189
190 pub fn complete_input(&mut self, id: InputId, outcome: InputOutcome) -> Result<(), CoreError> {
191 if self.status == SessionStatus::Closed {
192 return Err(CoreError::NotLive);
193 }
194 self.pending_input.retain(|request| request.id() != id);
195 if !self.can_enqueue() {
196 return Err(CoreError::ConsumerStalled);
197 }
198 self.events
199 .push_back(SessionEvent::InputCompleted { id, outcome });
200 Ok(())
201 }
202
203 pub fn complete_input_raw(&mut self, id: u64, outcome: InputOutcome) -> Result<(), CoreError> {
206 self.complete_input(InputId(id), outcome)
207 }
208
209 pub fn begin_recovery(&mut self) -> Result<(), CoreError> {
210 if self.status == SessionStatus::Closed {
211 return Err(CoreError::NotLive);
212 }
213 self.input_generation = self.input_generation.saturating_add(1);
214 let pending: Vec<_> = self.pending_input.drain(..).collect();
215 self.status = SessionStatus::Recovering;
216 for request in pending {
217 if !self.can_enqueue() {
218 return Err(CoreError::ConsumerStalled);
219 }
220 self.events.push_back(SessionEvent::InputCompleted {
221 id: request.id(),
222 outcome: InputOutcome::Superseded,
223 });
224 }
225 if self.can_enqueue() {
226 self.events
227 .push_back(SessionEvent::StatusChanged(SessionStatus::Recovering));
228 Ok(())
229 } else {
230 Err(CoreError::ConsumerStalled)
231 }
232 }
233
234 pub fn suspend(&mut self) -> Result<(), CoreError> {
235 if self.status == SessionStatus::Closed {
236 return Err(CoreError::NotLive);
237 }
238 self.input_generation = self.input_generation.saturating_add(1);
239 let pending: Vec<_> = self.pending_input.drain(..).collect();
240 self.status = SessionStatus::Suspended;
241 for request in pending {
242 if !self.can_enqueue() {
243 return Err(CoreError::ConsumerStalled);
244 }
245 self.events.push_back(SessionEvent::InputCompleted {
246 id: request.id(),
247 outcome: InputOutcome::Superseded,
248 });
249 }
250 self.enqueue_status(SessionStatus::Suspended)
251 }
252
253 pub fn close(&mut self) {
254 self.pending_input.clear();
255 self.receipts.clear();
256 self.state = None;
257 self.status = SessionStatus::Closed;
258 self.events.clear();
259 }
260
261 pub fn next_event(&mut self) -> Option<SessionEvent> {
262 self.events.pop_front()
263 }
264
265 pub fn drain_events(&mut self) {
266 self.events.clear();
267 }
268
269 pub fn pending_input_count(&self) -> usize {
270 self.pending_input.len()
271 }
272
273 fn issue_receipt(&mut self) -> CommitReceipt {
274 self.next_receipt = self.next_receipt.saturating_add(1);
275 CommitReceipt(self.next_receipt)
276 }
277
278 fn can_enqueue(&self) -> bool {
279 self.events.len() < self.options.event_capacity
280 }
281
282 fn enqueue_snapshot_required(
283 &mut self,
284 expected: StateRevision,
285 received: StateRevision,
286 reason: SnapshotRequiredReason,
287 ) -> Result<(), CoreError> {
288 if !self.can_enqueue() {
289 return Err(CoreError::ConsumerStalled);
290 }
291 self.events.push_back(SessionEvent::SnapshotRequired {
292 expected,
293 received,
294 reason,
295 });
296 Ok(())
297 }
298
299 fn enqueue_status(&mut self, status: SessionStatus) -> Result<(), CoreError> {
300 if !self.can_enqueue() {
301 return Err(CoreError::ConsumerStalled);
302 }
303 self.events.push_back(SessionEvent::StatusChanged(status));
304 Ok(())
305 }
306}
307
308fn effects_for(
309 delta: &crate::types::TerminalStateDelta,
310 current: &TerminalState,
311) -> Vec<TerminalEffect> {
312 let mut effects = Vec::new();
313 if !delta.primary_scroll.is_empty() {
314 effects.push(TerminalEffect::PrimaryScroll {
315 rows: delta.primary_scroll.clone(),
316 });
317 }
318 if delta.alternate != current.alternate {
319 effects.push(TerminalEffect::Screen {
320 alternate: delta.alternate,
321 });
322 }
323 if delta.modes != current.modes {
324 effects.push(TerminalEffect::InputModes { modes: delta.modes });
325 }
326 effects
327}