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
19impl Runtime {
20    pub(crate) fn record_compaction_hysteresis(&self, thread_id: &ThreadId, trigger_tokens: u32) {
21        if let Ok(mut state) = self.compaction_hysteresis.lock() {
22            state.insert(thread_id.clone(), trigger_tokens);
23        }
24    }
25
26    pub(crate) fn compaction_hysteresis_baseline(&self, thread_id: &ThreadId) -> Option<u32> {
27        self.compaction_hysteresis
28            .lock()
29            .ok()
30            .and_then(|state| state.get(thread_id).copied())
31    }
32
33    pub(crate) fn compaction_options_for_turn(
34        &self,
35        thread_id: &ThreadId,
36        allow_repeat: bool,
37    ) -> CompactionOptions {
38        CompactionOptions {
39            allow_repeat,
40            force: false,
41            hysteresis_baseline: self.compaction_hysteresis_baseline(thread_id),
42            preserve_hint: None,
43        }
44    }
45
46    pub async fn force_compact_thread(
47        &self,
48        thread_id: &ThreadId,
49        turn_id: &TurnId,
50        preserve_hint: Option<String>,
51    ) -> anyhow::Result<ForceCompactOutcome> {
52        let cfg = self.status().await;
53        let provider = cfg.default_provider.clone();
54        let model = cfg.default_model.clone();
55        let transcript = self.transcript_for_force_compact(thread_id).await?;
56        if transcript.is_empty() {
57            return Ok(ForceCompactOutcome {
58                compacted: false,
59                reason: Some("empty_transcript".to_string()),
60                estimated_tokens_before: 0,
61                estimated_tokens_after: 0,
62            });
63        }
64        let estimated_before = estimate_prompt_tokens(&transcript);
65        let compacted = self
66            .compact_transcript_if_needed(
67                thread_id,
68                turn_id,
69                &provider,
70                &model,
71                transcript,
72                CompactionOptions {
73                    allow_repeat: true,
74                    force: true,
75                    hysteresis_baseline: None,
76                    preserve_hint: preserve_hint.filter(|text| !text.trim().is_empty()),
77                },
78            )
79            .await?;
80        let estimated_after = estimate_prompt_tokens(&compacted);
81        Ok(ForceCompactOutcome {
82            compacted: estimated_after < estimated_before
83                || compacted
84                    .iter()
85                    .any(|item| matches!(item, TranscriptItem::ContextCompaction(_))),
86            reason: None,
87            estimated_tokens_before: estimated_before,
88            estimated_tokens_after: estimated_after,
89        })
90    }
91
92    async fn transcript_for_force_compact(
93        &self,
94        thread_id: &ThreadId,
95    ) -> anyhow::Result<Vec<TranscriptItem>> {
96        let Some(store) = &self.thread_store else {
97            return Ok(Vec::new());
98        };
99        let Some(snapshot) = store.load_thread(thread_id).await? else {
100            return Ok(Vec::new());
101        };
102        let mut out = Vec::new();
103        for turn in snapshot.turns {
104            out.extend(turn.items);
105        }
106        Ok(crate::compaction::trim_to_last_compaction_boundary(out))
107    }
108
109    pub(crate) async fn summarize_compaction_head(
110        &self,
111        provider: &str,
112        model: &str,
113        head: &[TranscriptItem],
114        preserve_hint: Option<&str>,
115    ) -> anyhow::Result<Option<String>> {
116        if head.is_empty() {
117            return Ok(None);
118        }
119        // Prefer a full-head LLM snapshot; if the summary request itself hits the
120        // provider prompt limit, shrink the head and retry once, then fall back
121        // to deterministic summarization.
122        let mut summary_head = head.to_vec();
123        let draft = loop {
124            match self
125                .run_compaction_summary_inference(
126                    provider,
127                    model,
128                    build_compaction_summary_prompt(&summary_head, preserve_hint),
129                )
130                .await
131            {
132                Ok(draft) => break draft,
133                Err(err)
134                    if crate::compaction::is_context_limit_failure_message(&err.to_string()) =>
135                {
136                    if summary_head.len() <= 1 {
137                        return Ok(None);
138                    }
139                    // Drop the oldest half so the summary prompt fits.
140                    let keep_from = summary_head.len() / 2;
141                    summary_head = summary_head.split_off(keep_from);
142                    continue;
143                }
144                Err(_) => {
145                    // Non-context failures already fall through to deterministic
146                    // summaries via Ok(None) below when the provider returns empty.
147                    return Ok(None);
148                }
149            }
150        };
151        let Some(draft) = draft else {
152            return Ok(None);
153        };
154        if !accept_llm_compaction_summary(head, &draft) {
155            return Ok(None);
156        }
157        let verified = match self
158            .run_compaction_summary_inference(
159                provider,
160                model,
161                build_compaction_verify_prompt(&draft),
162            )
163            .await
164        {
165            Ok(Some(text)) => text,
166            Ok(None) | Err(_) => draft.clone(),
167        };
168        if accept_llm_compaction_summary(head, &verified) {
169            Ok(Some(verified))
170        } else if accept_llm_compaction_summary(head, &draft) {
171            Ok(Some(draft))
172        } else {
173            Ok(None)
174        }
175    }
176
177    async fn run_compaction_summary_inference(
178        &self,
179        provider: &str,
180        model: &str,
181        prompt: String,
182    ) -> anyhow::Result<Option<String>> {
183        let engine = self.engine_for(provider)?;
184        let request = AgentInferenceRequest {
185            model: ModelSelection {
186                provider: provider.to_string(),
187                model: model.to_string(),
188            },
189            instructions: InstructionBundle {
190                system: Some(
191                    "You compress conversation history into durable state snapshots.".to_string(),
192                ),
193                developer: None,
194                developer_context: None,
195            },
196            transcript: vec![TranscriptItem::UserMessage(UserMessage::text(prompt))],
197            tools: Vec::new(),
198            tool_choice: ToolChoice::None,
199            reasoning: ReasoningConfig::default(),
200            output: OutputConfig::default(),
201            runtime: RuntimeHints {
202                profile: RuntimeProfile::Interactive,
203                ..RuntimeHints::default()
204            },
205            metadata: serde_json::json!({ "roderCompactionSummary": true }),
206        };
207        let ctx = InferenceTurnContext {
208            thread_id: &"compaction-summary".to_string(),
209            turn_id: &"compaction-summary".to_string(),
210            tool_executor: None,
211        };
212        let mut stream = engine.stream_turn(ctx, request).await?;
213        let mut text = String::new();
214        while let Some(event) = stream.next().await {
215            match event? {
216                InferenceEvent::MessageDelta(MessageDelta { text: delta, .. }) => {
217                    text.push_str(&delta)
218                }
219                InferenceEvent::Failed(failure) => {
220                    if crate::compaction::is_context_limit_failure_message(&failure.message) {
221                        anyhow::bail!("{}", failure.message);
222                    }
223                    // Non-context provider failures should not abort the parent
224                    // turn compaction path — fall back to deterministic summary.
225                    return Ok(None);
226                }
227                InferenceEvent::Completed(_) => break,
228                _ => {}
229            }
230        }
231        if text.trim().is_empty() {
232            Ok(None)
233        } else {
234            Ok(Some(text.trim().to_string()))
235        }
236    }
237}
238
239#[derive(Debug, Clone)]
240pub struct ForceCompactOutcome {
241    pub compacted: bool,
242    pub reason: Option<String>,
243    pub estimated_tokens_before: u32,
244    pub estimated_tokens_after: u32,
245}
246
247pub(crate) fn compaction_hysteresis_state() -> Mutex<HashMap<ThreadId, u32>> {
248    Mutex::new(HashMap::new())
249}