Skip to main content

roder_core/
compaction_runtime.rs

1use std::collections::HashMap;
2
3use futures::StreamExt;
4use roder_api::events::{ThreadId, TurnId};
5use roder_api::inference::{
6    AgentInferenceRequest, InferenceEvent, InferenceTurnContext, InstructionBundle, MessageDelta,
7    ModelSelection, OutputConfig, ReasoningConfig, RuntimeHints, RuntimeProfile,
8};
9use roder_api::tools::ToolChoice;
10use roder_api::transcript::{TranscriptItem, UserMessage};
11use std::sync::Mutex;
12
13use crate::compaction::{
14    CompactionOptions, accept_llm_compaction_summary, build_compaction_summary_prompt,
15    build_compaction_verify_prompt, estimate_prompt_tokens,
16};
17use crate::runtime::Runtime;
18
19#[derive(Default)]
20pub(crate) struct CompactionState {
21    tokens: u32,
22    generation: u64,
23}
24
25impl Runtime {
26    pub(crate) fn record_compaction_hysteresis(&self, thread_id: &ThreadId, trigger_tokens: u32) {
27        if let Ok(mut state) = self.compaction_hysteresis.lock() {
28            let state = state.entry(thread_id.clone()).or_default();
29            state.tokens = trigger_tokens;
30            state.generation = state.generation.wrapping_add(1);
31        }
32    }
33
34    pub(crate) fn compaction_hysteresis_baseline(&self, thread_id: &ThreadId) -> Option<u32> {
35        self.compaction_hysteresis
36            .lock()
37            .ok()
38            .and_then(|state| state.get(thread_id).map(|state| state.tokens))
39    }
40
41    pub(crate) fn compaction_generation(&self, thread_id: &ThreadId) -> u64 {
42        self.compaction_hysteresis
43            .lock()
44            .ok()
45            .and_then(|state| state.get(thread_id).map(|state| state.generation))
46            .unwrap_or(0)
47    }
48
49    pub(crate) fn compaction_options_for_turn(
50        &self,
51        thread_id: &ThreadId,
52        allow_repeat: bool,
53    ) -> CompactionOptions {
54        CompactionOptions {
55            allow_repeat,
56            force: false,
57            hysteresis_baseline: self.compaction_hysteresis_baseline(thread_id),
58            preserve_hint: None,
59        }
60    }
61
62    pub async fn force_compact_thread(
63        &self,
64        thread_id: &ThreadId,
65        turn_id: &TurnId,
66        preserve_hint: Option<String>,
67    ) -> anyhow::Result<ForceCompactOutcome> {
68        let _thread_admission = self.thread_admission(thread_id).await;
69        let cfg = self.status().await;
70        let selection = self
71            .parent_model_selection_for_subagents(thread_id, turn_id)
72            .await;
73        let provider = selection
74            .as_ref()
75            .map(|selection| selection.provider.clone())
76            .unwrap_or(cfg.default_provider.clone());
77        let model = selection
78            .as_ref()
79            .map(|selection| selection.model.clone())
80            .unwrap_or(cfg.default_model.clone());
81        let transcript = self.transcript_for_force_compact(thread_id).await?;
82        if transcript.is_empty() {
83            return Ok(ForceCompactOutcome {
84                compacted: false,
85                reason: Some("empty_transcript".to_string()),
86                estimated_tokens_before: 0,
87                estimated_tokens_after: 0,
88            });
89        }
90        let estimated_before = estimate_prompt_tokens(&transcript);
91        // Manual compaction runs only on an idle task. A second sampling
92        // request against a snapshot taken during tool execution can omit
93        // effects that complete before the boundary is committed.
94        if self.active_turn_for_thread(thread_id).await.is_some() {
95            return Ok(ForceCompactOutcome {
96                compacted: false,
97                reason: Some("turn_active".into()),
98                estimated_tokens_before: estimated_before,
99                estimated_tokens_after: estimated_before,
100            });
101        }
102
103        let compacted = self
104            .compact_transcript_if_needed(
105                thread_id,
106                turn_id,
107                &provider,
108                &model,
109                transcript,
110                CompactionOptions {
111                    allow_repeat: true,
112                    force: true,
113                    hysteresis_baseline: None,
114                    preserve_hint: preserve_hint.filter(|text| !text.trim().is_empty()),
115                },
116            )
117            .await?;
118        let estimated_after = estimate_prompt_tokens(&compacted);
119        Ok(ForceCompactOutcome {
120            compacted: estimated_after < estimated_before
121                || compacted
122                    .iter()
123                    .any(crate::compaction::is_compaction_boundary),
124            reason: None,
125            estimated_tokens_before: estimated_before,
126            estimated_tokens_after: estimated_after,
127        })
128    }
129
130    async fn transcript_for_force_compact(
131        &self,
132        thread_id: &ThreadId,
133    ) -> anyhow::Result<Vec<TranscriptItem>> {
134        let Some(store) = &self.thread_store else {
135            return Ok(Vec::new());
136        };
137        let Some(snapshot) = store.load_thread(thread_id).await? else {
138            return Ok(Vec::new());
139        };
140        let mut out = Vec::new();
141        for turn in snapshot.turns {
142            out.extend(turn.items);
143        }
144        Ok(crate::compaction::trim_to_last_compaction_boundary(out))
145    }
146
147    pub(crate) async fn summarize_compaction_head(
148        &self,
149        thread_id: &ThreadId,
150        turn_id: &TurnId,
151        provider: &str,
152        model: &str,
153        head: &[TranscriptItem],
154        preserve_hint: Option<&str>,
155    ) -> anyhow::Result<Option<String>> {
156        if head.is_empty() {
157            return Ok(None);
158        }
159        // Prefer a full-head LLM snapshot; if the summary request itself hits the
160        // provider prompt limit, shrink the head and retry once, then fall back
161        // to deterministic summarization.
162        let mut summary_head = head.to_vec();
163        let draft = loop {
164            match self
165                .run_compaction_summary_inference(
166                    thread_id,
167                    turn_id,
168                    provider,
169                    model,
170                    build_compaction_summary_prompt(&summary_head, preserve_hint),
171                )
172                .await
173            {
174                Ok(draft) => break draft,
175                Err(err)
176                    if crate::compaction::is_context_limit_failure_message(&err.to_string()) =>
177                {
178                    if summary_head.len() <= 1 {
179                        return Ok(None);
180                    }
181                    // Drop the oldest half so the summary prompt fits.
182                    let keep_from = summary_head.len() / 2;
183                    summary_head = summary_head.split_off(keep_from);
184                    continue;
185                }
186                Err(_) => {
187                    // Non-context failures already fall through to deterministic
188                    // summaries via Ok(None) below when the provider returns empty.
189                    return Ok(None);
190                }
191            }
192        };
193        let Some(draft) = draft else {
194            return Ok(None);
195        };
196        if !accept_llm_compaction_summary(head, &draft) {
197            return Ok(None);
198        }
199        let verified = match self
200            .run_compaction_summary_inference(
201                thread_id,
202                turn_id,
203                provider,
204                model,
205                build_compaction_verify_prompt(&draft),
206            )
207            .await
208        {
209            Ok(Some(text)) => text,
210            Ok(None) | Err(_) => draft.clone(),
211        };
212        if accept_llm_compaction_summary(head, &verified) {
213            Ok(Some(verified))
214        } else if accept_llm_compaction_summary(head, &draft) {
215            Ok(Some(draft))
216        } else {
217            Ok(None)
218        }
219    }
220
221    async fn run_compaction_summary_inference(
222        &self,
223        thread_id: &ThreadId,
224        turn_id: &TurnId,
225        provider: &str,
226        model: &str,
227        prompt: String,
228    ) -> anyhow::Result<Option<String>> {
229        let engine = self.engine_for(provider)?;
230        let request = AgentInferenceRequest {
231            model: ModelSelection {
232                provider: provider.to_string(),
233                model: model.to_string(),
234            },
235            instructions: InstructionBundle {
236                system: Some(
237                    "You compress conversation history into durable state snapshots.".to_string(),
238                ),
239                developer: None,
240                developer_context: None,
241            },
242            transcript: vec![TranscriptItem::UserMessage(UserMessage::text(prompt))],
243            tools: Vec::new(),
244            tool_choice: ToolChoice::None,
245            reasoning: ReasoningConfig::default(),
246            output: OutputConfig::default(),
247            runtime: RuntimeHints {
248                profile: RuntimeProfile::Interactive,
249                ..RuntimeHints::default()
250            },
251            metadata: serde_json::json!({ "roderCompactionSummary": true }),
252        };
253        let ctx = InferenceTurnContext {
254            thread_id,
255            turn_id,
256            tool_executor: None,
257        };
258        let mut stream = engine.stream_turn(ctx, request).await?;
259        let mut text = String::new();
260        let mut completed = false;
261        while let Some(event) = stream.next().await {
262            match event? {
263                InferenceEvent::MessageDelta(MessageDelta { text: delta, .. }) => {
264                    text.push_str(&delta)
265                }
266                InferenceEvent::Failed(failure) => {
267                    if crate::compaction::is_context_limit_failure_message(&failure.message) {
268                        anyhow::bail!("{}", failure.message);
269                    }
270                    // Non-context provider failures should not abort the parent
271                    // turn compaction path — fall back to deterministic summary.
272                    return Ok(None);
273                }
274                InferenceEvent::Completed(_) => {
275                    completed = true;
276                    break;
277                }
278                _ => {}
279            }
280        }
281        if !completed || text.trim().is_empty() {
282            Ok(None)
283        } else {
284            Ok(Some(text.trim().to_string()))
285        }
286    }
287}
288
289#[derive(Debug, Clone)]
290pub struct ForceCompactOutcome {
291    pub compacted: bool,
292    pub reason: Option<String>,
293    pub estimated_tokens_before: u32,
294    pub estimated_tokens_after: u32,
295}
296
297pub(crate) fn compaction_hysteresis_state() -> Mutex<HashMap<ThreadId, CompactionState>> {
298    Mutex::new(HashMap::new())
299}