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 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 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 let keep_from = summary_head.len() / 2;
183 summary_head = summary_head.split_off(keep_from);
184 continue;
185 }
186 Err(_) => {
187 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 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}