Skip to main content

orchestral_runtime/session_context/
pressure.rs

1use super::*;
2
3impl AgentSessionCompactor {
4    /// Compact against the same full-input meter and policy used by context
5    /// projection. Candidate summaries are measured before they become durable
6    /// facts; a failed candidate cannot enlarge the next model request.
7    pub async fn compact_active_run_for_context(
8        &self,
9        engine: &AgentSessionContextEngine,
10        request: SessionContextRequest,
11        policy: ContextTokenPolicy,
12    ) -> Result<Option<AgentSessionRecord>, SessionContextError> {
13        validate_context_request(&request)?;
14        if request.through_session_seq.is_some() {
15            return Err(SessionContextError::InvalidRequest(
16                "cannot compact a historical Context cursor".to_owned(),
17            ));
18        }
19        let records = self.journal.load_session(&request.session_id).await?;
20        validate_session_trace(&request.session_id, &records)?;
21        let groups = replay_groups(
22            &records,
23            &request.current_run_id,
24            &request.allowed_skill_digests,
25        )?;
26        let measure = |groups: &BTreeMap<u64, MessageGroup>| {
27            let selected = groups
28                .values()
29                .filter(|group| group.pinned)
30                .map(|group| group.key)
31                .collect();
32            engine.context_input_tokens(
33                &assemble_messages(&request.system_message, groups, &selected),
34                &request.tools,
35                policy,
36            )
37        };
38        let before = measure(&groups)?;
39        let budget = request.max_context_tokens - request.reserved_output_tokens;
40        if before <= budget {
41            return Ok(None);
42        }
43        let immutable = groups
44            .iter()
45            .filter(|(_, group)| !is_current_compactable(group, &records, &request.current_run_id))
46            .map(|(key, group)| (*key, group.clone()))
47            .collect();
48        if policy == ContextTokenPolicy::UpperBound && measure(&immutable)? > budget {
49            // Task, Skills, retained artifacts and safety facts cannot be
50            // removed to manufacture space, even if tool history is available.
51            return Ok(None);
52        }
53        // Planning estimates may use a measured prefix. The immutable subset
54        // can lose that prefix and fall back to a larger raw estimate, while a
55        // candidate retaining the prefix still fits. Only full candidates can
56        // establish progress under Planning; the subset is not a lower bound.
57
58        let mut sources = Vec::new();
59        if let Some(preferred) =
60            select_active_run_compaction_source(&records, &request.current_run_id, &self.policy)
61        {
62            sources.push(preferred);
63        }
64        // Keeping recent exchanges is a quality preference. If the preferred
65        // source cannot fit, consider each entire live segment, including its
66        // newest exchange or a summary whose budget has since become smaller.
67        for source in live_source::compactable_segments(&groups, &records, &request.current_run_id)
68        {
69            if !sources.contains(&source) {
70                sources.push(source);
71            }
72        }
73        // A late summary producer may be adjacent to exchanges on the other
74        // side of a surviving fact. Single groups remain useful safe candidates
75        // when their combined logical extent cannot be replaced in place.
76        for group in groups
77            .values()
78            .filter(|group| is_current_compactable(group, &records, &request.current_run_id))
79        {
80            let source = single_range(group.producer_seq);
81            if !sources.contains(&source) {
82                sources.push(source);
83            }
84        }
85        let mut best = None;
86        let mut best_tokens = before;
87        'sources: for source in sources {
88            let originals = original_compaction_groups(&records, &source)?;
89            let current_messages = groups
90                .values()
91                .filter(|group| source.contains(group.producer_seq))
92                .flat_map(|group| group.messages.clone())
93                .collect::<Vec<_>>();
94            let serialized = serde_json::to_string(&current_messages)
95                .map_err(|error| SessionContextError::Compaction(error.to_string()))?;
96            // This is a generation hint, not token accounting. Full-input
97            // metering below is authoritative for the selected token policy.
98            let mut max_chars = (serialized.chars().count() / 2).max(256);
99            let focus_messages = groups
100                .values()
101                .filter(|group| group.pinned && !source.contains(group.producer_seq))
102                .flat_map(|group| group.messages.clone())
103                .collect::<Vec<_>>();
104            let mut previous_summary = None;
105            loop {
106                let summary = self
107                    .summarizer
108                    .summarize_with_char_budget(
109                        SessionCompactionInput {
110                            session_id: request.session_id.clone(),
111                            source: source.clone(),
112                            groups: originals.clone(),
113                            focus_messages: focus_messages.clone(),
114                        },
115                        max_chars,
116                    )
117                    .await?;
118                summary.validate().map_err(|error| {
119                    SessionContextError::Compaction(format!("invalid active-Run summary: {error}"))
120                })?;
121                if previous_summary.as_ref() == Some(&summary) {
122                    break;
123                }
124                if summary.role == ModelRole::Assistant
125                    && !placement::can_replace_source(&groups, &records, &source)
126                {
127                    continue 'sources;
128                }
129                let next_seq = records.len() as u64 + 1;
130                let replacement =
131                    placement::summary_group(&groups, &source, next_seq, &summary, true, true)?;
132                let mut candidate = groups.clone();
133                candidate.retain(|_, group| !source.contains(group.producer_seq));
134                candidate.insert(next_seq, replacement);
135                let used = measure(&candidate)?;
136                if used < best_tokens {
137                    best_tokens = used;
138                    best = Some((source.clone(), summary.clone()));
139                    if used <= budget {
140                        break 'sources;
141                    }
142                }
143                if max_chars == 256 {
144                    break;
145                }
146                previous_summary = Some(summary);
147                max_chars = (max_chars / 2).max(256);
148            }
149        }
150        let Some((source, summary)) = best else {
151            return Ok(None);
152        };
153        // Separate live segments can straddle protected facts. A strictly
154        // smaller projection may need another pass; never persist no progress.
155        self.commit_active_run_summary(&records, &request.current_run_id, source, summary)
156            .await
157    }
158}
159
160fn is_current_compactable(
161    group: &MessageGroup,
162    records: &[AgentSessionRecord],
163    current_run_id: &RunId,
164) -> bool {
165    group.active_compactable && records[(group.producer_seq - 1) as usize].run_id == *current_run_id
166}
167
168#[cfg(test)]
169mod tests;