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 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 let keep_from = summary_head.len() / 2;
141 summary_head = summary_head.split_off(keep_from);
142 continue;
143 }
144 Err(_) => {
145 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 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}