1use std::collections::VecDeque;
9use std::panic::AssertUnwindSafe;
10use std::path::{Path, PathBuf};
11use std::sync::Arc;
12use std::sync::atomic::{AtomicU64, Ordering};
13
14use futures_util::FutureExt;
15use tokio::sync::mpsc;
16use tokio::task::JoinHandle;
17use tokio_util::sync::CancellationToken;
18
19use talos_core::message::{AgentEvent, Message};
20use talos_core::session::{
21 MAX_STEERING_QUEUE_BYTES, MAX_STEERING_QUEUE_IMAGE_BYTES, MAX_STEERING_QUEUE_IMAGES,
22 MAX_STEERING_QUEUE_ITEMS, SessionConfig, SessionEvent, SessionHandle, SessionOp,
23 StructuredSubmission, SubmissionItem, SubmissionKind, SubmissionRejectionReason,
24 SubmissionSource, TurnCompletionStatus, TurnEventPayload,
25};
26#[cfg(test)]
27use talos_core::session::{
28 MAX_SUBMISSION_BATCH_ITEMS, MAX_SUBMISSION_IMAGE_COUNT, SubmissionReceiptDisposition,
29};
30use talos_session::PendingSubmissionStore;
31
32use crate::compaction::Compactor;
33use crate::token::TokenEstimator;
34use crate::{ActivatedSkillContext, Agent};
35
36mod custody;
37mod turn;
38
39#[cfg(test)]
40#[allow(warnings)]
41mod tests;
42
43use turn::{
44 DurableTurnPersistence, TurnForwarding, TurnPersistence, TurnRecord, TurnRecordStatus,
45 run_turn_with_forwarding,
46};
47
48static NEXT_RUNTIME_SESSION_ID: AtomicU64 = AtomicU64::new(1);
49const MAX_RECENT_SUBMISSION_IDS: usize = 1024;
50const MAX_RECENT_ITEM_IDS: usize = 4096;
51
52#[derive(Debug, Clone)]
53struct ActiveStructuredTurn {
54 submission_id: String,
55 receipt_id: String,
56 session_generation: u64,
57 source: SubmissionSource,
58 turn_id: String,
59}
60
61struct StartedTurn {
62 handle: JoinHandle<Option<TurnRecord>>,
63 token: CancellationToken,
64 structured: Option<ActiveStructuredTurn>,
65}
66
67pub struct AppServerSession {
72 agent: Arc<Agent>,
73 sq_rx: tokio::sync::mpsc::Receiver<SessionOp>,
74 eq_tx: mpsc::UnboundedSender<SessionEvent>,
75 history: Vec<Message>,
76 compactor: Compactor,
77 session_file: Option<PathBuf>,
78 session_dir: Option<PathBuf>,
79 persistence: Option<TurnPersistence>,
80 durable_persistence: Option<DurableTurnPersistence>,
81 pending_store: PendingSubmissionStore,
82 session_id: String,
83 session_generation: u64,
84 turn_prefix: String,
85 model_context_limit: u32,
86}
87
88impl AppServerSession {
89 pub fn new(agent: Agent, config: SessionConfig) -> (SessionHandle, Self) {
94 let (sq_tx, sq_rx) = tokio::sync::mpsc::channel(512);
95 let (eq_tx, eq_rx) = mpsc::unbounded_channel();
96
97 let handle = SessionHandle { sq_tx, eq_rx };
98 let compactor = Compactor::new(TokenEstimator::new(), config.model_context_limit);
99 let instance_id = NEXT_RUNTIME_SESSION_ID.fetch_add(1, Ordering::Relaxed);
100 let session_id = format!("runtime_{}_{}", std::process::id(), instance_id);
101 let pending_session_file = config
102 .workspace_root
103 .join(".talos")
104 .join("runtime")
105 .join(format!("{session_id}.tlog"));
106 let pending_store =
107 PendingSubmissionStore::for_session_file(&pending_session_file, &session_id);
108
109 let actor = Self {
110 agent: Arc::new(agent),
111 sq_rx,
112 eq_tx,
113 history: config.initial_history,
114 compactor,
115 session_file: None,
116 session_dir: None,
117 persistence: None,
118 durable_persistence: None,
119 pending_store,
120 session_id,
121 session_generation: 0,
122 turn_prefix: format!("turn_{}_{}", std::process::id(), instance_id),
123 model_context_limit: config.model_context_limit,
124 };
125
126 (handle, actor)
127 }
128
129 pub fn set_generation(&mut self, generation: u64) {
131 self.session_generation = generation;
132 }
133
134 #[must_use]
136 pub fn generation(&self) -> u64 {
137 self.session_generation
138 }
139
140 pub fn set_session_paths(&mut self, file: PathBuf, dir: PathBuf) {
141 self.session_file = Some(file);
142 self.session_dir = Some(dir);
143 }
144
145 pub fn set_persistence(
147 &mut self,
148 session: talos_session::Session,
149 metadata: talos_session::SessionMetadata,
150 ) {
151 self.pending_store = PendingSubmissionStore::for_session(&session);
152 self.session_id = session.id.to_string();
153 self.persistence = Some(TurnPersistence { session, metadata });
154 }
155
156 pub fn set_durable_persistence(
158 &mut self,
159 session: talos_session::DurableSession,
160 policy: talos_session::PersistencePolicy,
161 ) {
162 let session_id = session.id().to_string();
163 self.pending_store =
164 PendingSubmissionStore::for_session_file(session.file_path(), &session_id);
165 self.session_id = session_id;
166 self.durable_persistence = Some(DurableTurnPersistence { session, policy });
167 }
168
169 pub async fn run(&mut self) {
171 self.reconcile_running_submissions();
172
173 let mut turn_counter: u64 = 0;
174 let mut submission_counter: u64 = 0;
175 let mut pending = VecDeque::<StructuredSubmission>::new();
176 let mut pending_items = 0_usize;
177 let mut pending_bytes = 0_usize;
178 let mut pending_images = 0_usize;
179 let mut pending_image_bytes = 0_u64;
180 let mut recent_submission_ids = VecDeque::<String>::new();
181 let mut recent_item_ids = VecDeque::<String>::new();
182 let mut current_turn: Option<JoinHandle<Option<TurnRecord>>> = None;
183 let mut current_submission_size: Option<(usize, usize, usize, u64)> = None;
184 let mut current_structured: Option<ActiveStructuredTurn> = None;
185 let mut cancel_token: Option<CancellationToken> = None;
186 let mut paused = self.restore_pending_submissions(
187 &mut pending,
188 &mut pending_items,
189 &mut pending_bytes,
190 &mut pending_images,
191 &mut pending_image_bytes,
192 &mut recent_submission_ids,
193 &mut recent_item_ids,
194 );
195 let mut shutting_down = false;
196
197 loop {
198 if current_turn.is_none()
199 && !paused
200 && let Some(submission) = pending.pop_front()
201 {
202 let (images, image_bytes) = submission.image_totals();
203 let submission_size = (
204 submission.items.len(),
205 submission.total_text_bytes(),
206 images,
207 image_bytes,
208 );
209 turn_counter = turn_counter.saturating_add(1);
210 match self
211 .start_submission(submission.clone(), turn_counter)
212 .await
213 {
214 Some(started) => {
215 current_turn = Some(started.handle);
216 current_submission_size = Some(submission_size);
217 current_structured = started.structured;
218 cancel_token = Some(started.token);
219 }
220 None => {
221 pending.push_front(submission);
222 paused = true;
223 }
224 }
225 }
226
227 if shutting_down && current_turn.is_none() && pending.is_empty() {
228 break;
229 }
230
231 tokio::select! {
232 completed = async {
233 match current_turn.as_mut() {
234 Some(turn) => Some(turn.await),
235 None => None,
236 }
237 }, if current_turn.is_some() => {
238 current_turn = None;
239 cancel_token = None;
240 if let Some((items, bytes, images, image_bytes)) = current_submission_size.take() {
241 pending_items = pending_items.saturating_sub(items);
242 pending_bytes = pending_bytes.saturating_sub(bytes);
243 pending_images = pending_images.saturating_sub(images);
244 pending_image_bytes = pending_image_bytes.saturating_sub(image_bytes);
245 }
246 let structured = current_structured.take();
247 match completed.and_then(Result::ok).flatten() {
248 Some(record) => {
249 let status = record.status;
250 let completion = record.completion.clone();
251 self.commit_turn_record(record);
252 let custody_ok = structured.as_ref().is_none_or(|active| {
253 self.finish_structured_turn(active, &completion)
254 });
255 paused = status != TurnRecordStatus::Success || !custody_ok;
256 }
257 None => {
258 if let Some(active) = structured.as_ref() {
259 let completion = TurnCompletionStatus::Error {
260 message: "turn task ended without a completion record".into(),
261 };
262 let _ = self.finish_structured_turn(active, &completion);
263 }
264 paused = true;
265 }
266 }
267 }
268 op = self.sq_rx.recv(), if !shutting_down => {
269 let Some(op) = op else {
270 shutting_down = true;
271 if let Some(token) = cancel_token.take() {
272 token.cancel();
273 }
274 self.release_in_memory_pending_on_shutdown(&mut pending);
275 let _ = self.pending_store.pause_unstarted();
276 pending_items = current_submission_size.map_or(0, |size| size.0);
277 pending_bytes = current_submission_size.map_or(0, |size| size.1);
278 pending_images = current_submission_size.map_or(0, |size| size.2);
279 pending_image_bytes = current_submission_size.map_or(0, |size| size.3);
280 continue;
281 };
282 match op {
283 SessionOp::Submit { message } => {
284 submission_counter = submission_counter.saturating_add(1);
285 let submission = compatibility_submission(
286 submission_counter,
287 self.session_generation,
288 SubmissionKind::UserTurn,
289 message,
290 Vec::new(),
291 );
292 if self.accept_submission(
293 submission,
294 &mut pending,
295 &mut pending_items,
296 &mut pending_bytes,
297 &mut pending_images,
298 &mut pending_image_bytes,
299 &mut recent_submission_ids,
300 &mut recent_item_ids,
301 ) {
302 paused = false;
303 }
304 }
305 SessionOp::SubmitMultimodal { text, attachments } => {
306 submission_counter = submission_counter.saturating_add(1);
307 let submission = compatibility_submission(
308 submission_counter,
309 self.session_generation,
310 SubmissionKind::UserTurn,
311 text,
312 attachments,
313 );
314 if self.accept_submission(
315 submission,
316 &mut pending,
317 &mut pending_items,
318 &mut pending_bytes,
319 &mut pending_images,
320 &mut pending_image_bytes,
321 &mut recent_submission_ids,
322 &mut recent_item_ids,
323 ) {
324 paused = false;
325 }
326 }
327 SessionOp::PreviewRequest { message } => {
328 submission_counter = submission_counter.saturating_add(1);
329 let submission = compatibility_submission(
330 submission_counter,
331 self.session_generation,
332 SubmissionKind::PreviewRequest,
333 message,
334 Vec::new(),
335 );
336 if self.accept_submission(
337 submission,
338 &mut pending,
339 &mut pending_items,
340 &mut pending_bytes,
341 &mut pending_images,
342 &mut pending_image_bytes,
343 &mut recent_submission_ids,
344 &mut recent_item_ids,
345 ) {
346 paused = false;
347 }
348 }
349 SessionOp::SubmitStructured { submission } => {
350 let resumes = submission.source != SubmissionSource::Scheduler;
351 if self.accept_durable_submission(
352 submission,
353 &mut pending,
354 &mut pending_items,
355 &mut pending_bytes,
356 &mut pending_images,
357 &mut pending_image_bytes,
358 &mut recent_submission_ids,
359 &mut recent_item_ids,
360 None,
361 ) && resumes
362 {
363 paused = false;
364 }
365 }
366 SessionOp::SubmitStructuredTracked {
367 submission,
368 receipt_tx,
369 } => {
370 let resumes = submission.source != SubmissionSource::Scheduler;
371 if self.accept_durable_submission(
372 submission,
373 &mut pending,
374 &mut pending_items,
375 &mut pending_bytes,
376 &mut pending_images,
377 &mut pending_image_bytes,
378 &mut recent_submission_ids,
379 &mut recent_item_ids,
380 receipt_tx.as_ref(),
381 ) && resumes
382 {
383 paused = false;
384 }
385 }
386 SessionOp::ReconcileStructured { submission }
387 | SessionOp::SubmitStructuredReconcile { submission } => {
388 self.reconcile_submission(&submission, None);
389 }
390 SessionOp::ReconcileStructuredTracked {
391 submission,
392 receipt_tx,
393 }
394 | SessionOp::SubmitStructuredReconcileTracked {
395 submission,
396 receipt_tx,
397 } => {
398 self.reconcile_submission(&submission, receipt_tx.as_ref());
399 }
400 SessionOp::Interrupt => {
401 if let Some(token) = &cancel_token {
402 token.cancel();
403 } else if !pending.is_empty() {
404 if let Err(error) = self.pending_store.pause_unstarted() {
405 self.emit_custody_error("failed to pause pending work", &error);
406 }
407 paused = true;
408 }
409 }
410 SessionOp::InterruptTurn {
411 session_generation,
412 turn_id,
413 } => {
414 let matches = session_generation == self.session_generation
415 && current_structured
416 .as_ref()
417 .is_some_and(|active| active.turn_id == turn_id);
418 if matches && let Some(token) = &cancel_token {
419 token.cancel();
420 }
421 }
422 SessionOp::CancelPausedSubmission {
423 session_generation,
424 submission_id,
425 } => {
426 if current_turn.is_none()
427 && session_generation == self.session_generation
428 && self.cancel_paused_submission(
429 &submission_id,
430 &mut pending,
431 &mut pending_items,
432 &mut pending_bytes,
433 &mut pending_images,
434 &mut pending_image_bytes,
435 )
436 {
437 paused = false;
438 }
439 }
440 SessionOp::SetSkillContext { name, content } => {
441 if current_turn.is_some() || !pending.is_empty() {
442 let _ = self.eq_tx.send(SessionEvent::Error {
443 message: "cannot change active skill while a turn is active"
444 .into(),
445 });
446 continue;
447 }
448 let context = match (name, content) {
449 (Some(name), Some(content)) => {
450 Some(ActivatedSkillContext { name, content })
451 }
452 _ => None,
453 };
454 if let Some(agent_mut) = Arc::get_mut(&mut self.agent) {
455 agent_mut.set_activated_skill_context(context);
456 } else {
457 let _ = self.eq_tx.send(SessionEvent::Error {
458 message: "cannot change active skill while agent is busy"
459 .into(),
460 });
461 }
462 }
463 SessionOp::Shutdown => {
464 shutting_down = true;
465 if let Some(token) = &cancel_token {
466 token.cancel();
467 }
468 self.release_in_memory_pending_on_shutdown(&mut pending);
469 if let Err(error) = self.pending_store.pause_unstarted() {
470 self.emit_custody_error("failed to persist shutdown pause", &error);
471 }
472 pending_items = current_submission_size.map_or(0, |size| size.0);
473 pending_bytes = current_submission_size.map_or(0, |size| size.1);
474 pending_images = current_submission_size.map_or(0, |size| size.2);
475 pending_image_bytes = current_submission_size.map_or(0, |size| size.3);
476 }
477 }
478 }
479 }
480 }
481 }
482
483 fn commit_turn_record(&mut self, record: TurnRecord) {
484 for msg in record.new_messages {
485 self.history.push(msg);
486 }
487 }
488
489 async fn start_submission(
490 &mut self,
491 submission: StructuredSubmission,
492 turn_counter: u64,
493 ) -> Option<StartedTurn> {
494 if self.compactor.should_compact(&self.history) {
495 let compacted = self.compactor.apply_budget(self.history.clone());
496 let compacted = self.compactor.apply_trim(compacted);
497 let compacted = self.compactor.apply_microcompact(compacted);
498 self.history = match self
499 .compactor
500 .compact(compacted, self.agent.provider())
501 .await
502 {
503 Ok(history) => history,
504 Err(_) => self.compactor.compact_deterministic(self.history.clone()).0,
505 };
506 if let (Some(file), Some(dir)) = (&self.session_file, &self.session_dir) {
507 let _ = self.try_archive_session(file, dir, &self.history);
508 }
509 }
510
511 let prepared_turn = if submission.common_kind() == Some(SubmissionKind::PreviewRequest) {
512 None
513 } else {
514 match self
515 .agent
516 .prepare_session_turn(
517 &submission.items,
518 self.history.clone(),
519 self.model_context_limit,
520 )
521 .await
522 {
523 Ok(prepared) => Some(prepared),
524 Err(crate::AgentError::ContextBudgetExceeded { .. }) => {
525 self.pause_before_start(
526 &submission,
527 SubmissionRejectionReason::ContextBudgetExceeded,
528 );
529 return None;
530 }
531 Err(error) => {
532 let _ = self.eq_tx.send(SessionEvent::Error {
533 message: format!(
534 "failed to seal Provider request plan for {}: {error}",
535 submission.id
536 ),
537 });
538 self.pause_before_start(
539 &submission,
540 SubmissionRejectionReason::InvalidStructure,
541 );
542 return None;
543 }
544 }
545 };
546
547 let turn_id = format!("{}_{}", self.turn_prefix, turn_counter);
548 let structured = if submission.source == SubmissionSource::Compatibility {
549 None
550 } else {
551 let record = match self.pending_store.get(&submission.id) {
552 Ok(Some(record)) => record,
553 Ok(None) => {
554 self.emit_custody_error(
555 "accepted structured submission is missing from journal",
556 &submission.id,
557 );
558 return None;
559 }
560 Err(error) => {
561 self.emit_custody_error("failed to load accepted submission", &error);
562 return None;
563 }
564 };
565 if let Err(error) = self.pending_store.mark_running(&submission.id, &turn_id) {
566 self.emit_custody_error("failed to mark structured submission running", &error);
567 return None;
568 }
569 let active = ActiveStructuredTurn {
570 submission_id: submission.id.clone(),
571 receipt_id: record.receipt_id,
572 session_generation: self.session_generation,
573 source: submission.source,
574 turn_id: turn_id.clone(),
575 };
576 let _ = self.eq_tx.send(SessionEvent::StructuredSubmissionStarted {
577 session_id: self.session_id.clone(),
578 session_generation: active.session_generation,
579 submission: submission.clone(),
580 receipt_id: active.receipt_id.clone(),
581 turn_id: active.turn_id.clone(),
582 });
583 let _ = self.eq_tx.send(SessionEvent::StructuredTurnEvent {
584 session_id: self.session_id.clone(),
585 session_generation: active.session_generation,
586 source: active.source,
587 submission_id: active.submission_id.clone(),
588 receipt_id: active.receipt_id.clone(),
589 turn_id: active.turn_id.clone(),
590 sequence: 0,
591 payload: TurnEventPayload::Started,
592 });
593 Some(active)
594 };
595
596 let _ = self.eq_tx.send(SessionEvent::SubmissionStarted {
597 session_id: self.session_id.clone(),
598 submission_id: submission.id.clone(),
599 sender_generation: submission.sender_generation,
600 turn_id: turn_id.clone(),
601 });
602 let _ = self.eq_tx.send(SessionEvent::TurnEvent {
603 session_id: self.session_id.clone(),
604 turn_id: turn_id.clone(),
605 sequence: 0,
606 payload: TurnEventPayload::Started,
607 });
608
609 if submission.common_kind() == Some(SubmissionKind::PreviewRequest) {
610 let agent = self.agent.clone();
611 let eq_tx = self.eq_tx.clone();
612 let history = self.history.clone();
613 let session_id = self.session_id.clone();
614 let message = submission.items[0].text.clone();
615 let token = CancellationToken::new();
616 let preview_token = token.clone();
617 let handle = tokio::spawn(async move {
618 let result = tokio::select! {
619 () = preview_token.cancelled() => {
620 let completion = TurnCompletionStatus::Cancelled;
621 let _ = eq_tx.send(SessionEvent::TurnEvent {
622 session_id,
623 turn_id,
624 sequence: 1,
625 payload: TurnEventPayload::Completed {
626 status: completion.clone(),
627 },
628 });
629 return Some(TurnRecord {
630 new_messages: Vec::new(),
631 status: TurnRecordStatus::Cancelled,
632 completion,
633 });
634 }
635 result = agent.preview_request(message, history) => result,
636 };
637 let (completion, record_status) = match result {
638 Ok(Some(preview)) => {
639 let _ = eq_tx.send(SessionEvent::TurnEvent {
640 session_id: session_id.clone(),
641 turn_id: turn_id.clone(),
642 sequence: 1,
643 payload: TurnEventPayload::Progress {
644 event: AgentEvent::TurnStart,
645 },
646 });
647 let _ = eq_tx.send(SessionEvent::TurnEvent {
648 session_id: session_id.clone(),
649 turn_id: turn_id.clone(),
650 sequence: 2,
651 payload: TurnEventPayload::Progress {
652 event: AgentEvent::TextDelta {
653 delta: preview.clone(),
654 },
655 },
656 });
657 let _ = eq_tx.send(SessionEvent::TurnEvent {
658 session_id: session_id.clone(),
659 turn_id: turn_id.clone(),
660 sequence: 3,
661 payload: TurnEventPayload::Progress {
662 event: AgentEvent::TurnEnd {
663 stop_reason: talos_core::message::StopReason::EndTurn,
664 usage: talos_core::message::Usage::default(),
665 },
666 },
667 });
668 (
669 TurnCompletionStatus::Success {
670 final_text: preview,
671 new_messages: Vec::new(),
672 },
673 TurnRecordStatus::Success,
674 )
675 }
676 Ok(None) => (
677 TurnCompletionStatus::Error {
678 message: "request preview is unavailable for this provider".into(),
679 },
680 TurnRecordStatus::Error,
681 ),
682 Err(error) => (
683 TurnCompletionStatus::Error {
684 message: error.to_string(),
685 },
686 TurnRecordStatus::Error,
687 ),
688 };
689 let _ = eq_tx.send(SessionEvent::TurnEvent {
690 session_id,
691 turn_id,
692 sequence: if record_status == TurnRecordStatus::Success {
693 4
694 } else {
695 1
696 },
697 payload: TurnEventPayload::Completed {
698 status: completion.clone(),
699 },
700 });
701 Some(TurnRecord {
702 new_messages: Vec::new(),
703 status: record_status,
704 completion,
705 })
706 });
707 return Some(StartedTurn {
708 handle,
709 token,
710 structured,
711 });
712 }
713
714 if let Some(agent_mut) = Arc::get_mut(&mut self.agent) {
715 agent_mut.set_append_prompt_opt(None);
716 }
717
718 let sequence = Arc::new(AtomicU64::new(1));
719 let token = CancellationToken::new();
720 let token_clone = token.clone();
721 let agent = self.agent.clone();
722 let eq_tx = self.eq_tx.clone();
723 let prepared = prepared_turn.expect("non-preview submission must be prepared");
724 let persistence = self.persistence.clone();
725 let durable_persistence = self.durable_persistence.clone();
726 let session_id = self.session_id.clone();
727 let handle = tokio::spawn(async move {
728 let (event_tx, event_rx) = mpsc::unbounded_channel::<AgentEvent>();
729 let (result_tx, result_rx) = tokio::sync::oneshot::channel::<TurnRecord>();
730 let _ = AssertUnwindSafe(run_turn_with_forwarding(TurnForwarding {
731 agent,
732 prepared,
733 event_tx,
734 event_rx,
735 eq_tx,
736 cancel_token: token_clone,
737 turn_id,
738 session_id,
739 sequence,
740 persistence,
741 durable_persistence,
742 result_tx,
743 }))
744 .catch_unwind()
745 .await;
746 result_rx.await.ok()
747 });
748 Some(StartedTurn {
749 handle,
750 token,
751 structured,
752 })
753 }
754
755 fn try_archive_session(
756 &self,
757 file: &Path,
758 dir: &Path,
759 _compacted: &[Message],
760 ) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
761 use talos_session::CompactTextSessionStore;
762 use talos_session::compaction_engine::CompactionEngine;
763
764 let store = std::sync::Arc::new(CompactTextSessionStore);
765 let engine = CompactionEngine::new(store);
766
767 if !engine.should_compact(file, 0) {
768 return Ok(());
769 }
770
771 match engine.compact_segment(file, dir, 0)? {
772 talos_session::compaction_engine::CompactionResult::Compacted {
773 segment_id,
774 original_count,
775 ..
776 } => {
777 let _ = self.eq_tx.send(SessionEvent::Error {
778 message: format!(
779 "Session compacted: {original_count} entries archived to {segment_id}"
780 ),
781 });
782 }
783 talos_session::compaction_engine::CompactionResult::Skipped => {}
784 }
785
786 Ok(())
787 }
788}
789
790fn compatibility_submission(
791 sequence: u64,
792 sender_generation: u64,
793 kind: SubmissionKind,
794 text: String,
795 attachments: Vec<talos_core::message::ContentPart>,
796) -> StructuredSubmission {
797 StructuredSubmission {
798 id: format!("compatibility_{sequence}"),
799 source: SubmissionSource::Compatibility,
800 sender_generation,
801 items: vec![SubmissionItem {
802 id: format!("compatibility_item_{sequence}"),
803 enqueue_sequence: sequence,
804 kind,
805 text,
806 attachments,
807 }],
808 }
809}
810
811fn enqueue_by_source(
812 pending: &mut VecDeque<StructuredSubmission>,
813 submission: StructuredSubmission,
814) {
815 if submission.source == SubmissionSource::Scheduler {
816 pending.push_back(submission);
817 return;
818 }
819 let scheduler_boundary = pending
820 .iter()
821 .position(|queued| queued.source == SubmissionSource::Scheduler)
822 .unwrap_or(pending.len());
823 pending.insert(scheduler_boundary, submission);
824}
825
826fn record_recent_identity(identities: &mut VecDeque<String>, id: String, capacity: usize) {
827 while identities.len() >= capacity {
828 identities.pop_front();
829 }
830 identities.push_back(id);
831}
832
833fn validate_submission(submission: &StructuredSubmission) -> Result<(), SubmissionRejectionReason> {
834 submission.validate()
835}
836
837fn queue_has_capacity(
838 submission: &StructuredSubmission,
839 pending_items: usize,
840 pending_bytes: usize,
841 pending_images: usize,
842 pending_image_bytes: u64,
843) -> bool {
844 queue_totals_after(
845 submission,
846 pending_items,
847 pending_bytes,
848 pending_images,
849 pending_image_bytes,
850 )
851 .is_some()
852}
853
854fn queue_totals_after(
855 submission: &StructuredSubmission,
856 pending_items: usize,
857 pending_bytes: usize,
858 pending_images: usize,
859 pending_image_bytes: u64,
860) -> Option<(usize, usize, usize, u64)> {
861 let next_items = pending_items.checked_add(submission.items.len())?;
862 let next_bytes = pending_bytes.checked_add(submission.total_text_bytes())?;
863 let (submission_images, submission_image_bytes) = submission.image_totals();
864 let next_images = pending_images.checked_add(submission_images)?;
865 let next_image_bytes = pending_image_bytes.checked_add(submission_image_bytes)?;
866 (next_items <= MAX_STEERING_QUEUE_ITEMS
867 && next_bytes <= MAX_STEERING_QUEUE_BYTES
868 && next_images <= MAX_STEERING_QUEUE_IMAGES
869 && next_image_bytes <= MAX_STEERING_QUEUE_IMAGE_BYTES)
870 .then_some((next_items, next_bytes, next_images, next_image_bytes))
871}