Skip to main content

talos_agent/
session.rs

1//! AppServerSession actor — bridges SQ→Agent→EQ (ADR-005 L2 seam).
2//!
3//! The session actor owns an [`Agent`] and runs a message loop:
4//! - Receives [`SessionOp`] on the bounded SQ (cap=512)
5//! - Drives agent turns via [`Agent::run_streaming`]
6//! - Emits [`SessionEvent`] on the unbounded EQ
7
8use 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
67/// Session actor that owns an [`Agent`] and processes commands from the SQ.
68///
69/// Created via [`AppServerSession::new`], which returns a [`SessionHandle`]
70/// for the UI layer and the actor itself for spawning on a tokio task.
71pub 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    /// Creates a new generation-zero session actor with the given agent and configuration.
90    ///
91    /// Product composition roots that replace an Actor must call [`Self::set_generation`]
92    /// before spawning it. Returns a [`SessionHandle`] and the actor itself.
93    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    /// Assigns the authoritative generation before this Actor is spawned.
130    pub fn set_generation(&mut self, generation: u64) {
131        self.session_generation = generation;
132    }
133
134    /// Returns the authoritative generation assigned by the composition root.
135    #[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    /// Assigns the durable session that owns successful turn-message writes.
146    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    /// Assigns an atomic durable session used by the embedded runtime.
157    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    /// Runs the session actor until shutdown or SQ disconnect.
170    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}