返回 CodeWhale
purge.rs
根目录 / crates / tui / src / purge.rs
1 //! Agent-driven context purging.
2 //!
3 //! Unlike compaction (which summarises old messages via LLM), purge lets the
4 //! agent analyse the conversation history and surgically remove or rewrite
5 //! individual messages that are no longer needed. The agent uses the
6 //! `purge_context` tool to submit a list of operations; the engine validates
7 //! and executes them.
8
9 use regex::Regex;
10 use std::fmt::Write;
11 use tokio::sync::mpsc::Sender;
12
13 use crate::config::ProviderKind;
14 use crate::core::events::Event;
15 use crate::fast_hash::{FastHashMap, FastHashSet};
16 use crate::llm_client::LlmClient;
17 use crate::regex_cache::compile_user_regex;
18 use codewhale_models::Role;
19 use codewhale_models::{ContentBlock, Message, MessageRequest, Tool};
20
21 // ── Prompt‑building constants ──────────────────────────────────────────────
22
23 const TEXT_SNIPPET_CHARS: usize = 60;
24 const TOOL_RESULT_SNIPPET_CHARS: usize = 80;
25 const TOOL_USE_ARGS_CHARS: usize = 120;
26
27 // ── Prompt instruction template ─────────────────────────────────────────────
28
29 const PURGE_INSTRUCTIONS: &str = "\
30 ## Context Purge
31
32 Free space in the conversation's context window. Below is the current history with stable numeric IDs.\
33 Identify content that is clearly no longer needed for the ongoing work.
34
35 ### Operations
36
37 remove — Delete an entire message by its ID. Example:
38 {\"op\": \"remove\", \"msg\": 3}
39
40 replace — Rewrite part of a specific content block using regex substitution.
41 pattern uses Rust regex syntax. Must specify both `block` and
42 `pattern` and `with`. Example:
43 {\"op\": \"replace\", \"msg\": 7, \"block\": 0,
44 \"pattern\": \"read \\\\d+ files\", \"with\": \"read files\"}
45
46 offload — Preserve a complete message in durable session storage and replace it
47 with a compact retrieval handle. Use for long material that may be needed
48 again; retrieve_tool_result can recover exact content later. Example:
49 {\"op\": \"offload\", \"msg\": 3}
50
51 ### Pairing rule
52
53 Every ToolUse block is paired with its ToolResult. If you remove a message
54 containing a tool call, its result will be removed too — and vice versa. You
55 do not need to list both. Offload similarly archives the entire paired group,
56 including thinking, signatures, media, and tool data; no original blocks are discarded.
57
58 ### What to keep
59
60 - Important decisions, architectural choices
61 - File paths that are still relevant
62 - Tool outputs that contain information not yet acted upon
63
64 ### What to prune
65
66 - Verbose tool outputs whose information has been fully consumed
67 - Redundant confirmations (\"done\", \"ok\", \"that worked\")
68 - Superseded file reads (the file was later written/modified)
69 - Boilerplate that the model already incorporated into later work
70
71 Be conservative. When in doubt, keep the message.
72
73 ### Conversation
74 ";
75
76 // ── Purge operation types ───────────────────────────────────────────────────
77
78 /// A single purge operation submitted by the agent.
79 #[derive(Debug, Clone)]
80 pub enum PurgeOp {
81 /// Remove an entire message (plus its tool-call/result counterpart).
82 Remove { msg_id: usize },
83 /// Archive an entire message and its paired tool messages before replacing them with a handle.
84 Offload { msg_id: usize },
85 /// Regex-replace within a specific content block.
86 Replace {
87 msg_id: usize,
88 block_idx: usize,
89 pattern: Regex,
90 with: String,
91 },
92 }
93
94 /// Result of executing purge operations.
95 #[derive(Debug, Clone)]
96 pub struct PurgeResult {
97 /// The remaining messages after all operations.
98 pub messages: Vec<Message>,
99 /// How many messages were removed.
100 pub removed_count: usize,
101 /// How many replace operations were applied.
102 pub replaced_count: usize,
103 /// Messages preserved in a durable session artifact instead of active context.
104 pub offloaded_count: usize,
105 }
106
107 // ── Event emission helpers ──────────────────────────────────────────────────
108
109 /// Emit a `PurgeStarted` event to the UI.
110 pub async fn emit_purge_started(tx: &Sender<Event>, message: String) {
111 let _ = tx.send(Event::PurgeStarted { message }).await;
112 }
113
114 /// Emit a `PurgeCompleted` event to the UI.
115 pub async fn emit_purge_completed(
116 tx: &Sender<Event>,
117 messages_before: usize,
118 messages_after: usize,
119 removed_count: usize,
120 replaced_count: usize,
121 message: String,
122 ) {
123 let _ = tx
124 .send(Event::PurgeCompleted {
125 messages_before,
126 messages_after,
127 removed_count,
128 replaced_count,
129 message,
130 })
131 .await;
132 }
133
134 /// Emit a `PurgeFailed` event to the UI.
135 pub async fn emit_purge_failed(tx: &Sender<Event>, message: String) {
136 let _ = tx.send(Event::PurgeFailed { message }).await;
137 }
138
139 // ── Prompt builder ──────────────────────────────────────────────────────────
140
141 /// Build the purge request user message — a formatted listing of the current
142 /// conversation with ephemeral sequential IDs.
143 pub fn build_purge_prompt(messages: &[Message]) -> String {
144 let mut buf = String::with_capacity(messages.len().saturating_mul(256));
145 buf.push_str(PURGE_INSTRUCTIONS);
146
147 for (idx, msg) in messages.iter().enumerate() {
148 let msg_id = idx + 1; // 1‑based for the agent
149 if msg.role == "user" {
150 // User messages: always a single block — omit block index.
151 format_user_message(&mut buf, msg_id, msg);
152 } else {
153 // Assistant messages: may be multi‑block — show block indices.
154 let _ = writeln!(buf, "[{msg_id}] {role}", role = msg.role);
155 for (blk_idx, block) in msg.content.iter().enumerate() {
156 format_content_block(&mut buf, blk_idx, block);
157 }
158 buf.push('\n');
159 }
160 }
161
162 buf
163 }
164
165 fn format_user_message(buf: &mut String, msg_id: usize, msg: &Message) {
166 let block = msg.content.first();
167 match block {
168 Some(ContentBlock::Text { text, .. }) => {
169 let snippet = truncate_str(text, TEXT_SNIPPET_CHARS);
170 let _ = writeln!(
171 buf,
172 "[{msg_id}] user Text ({len} chars): \"{snippet}\"",
173 len = text.len()
174 );
175 }
176 Some(ContentBlock::ToolResult {
177 content,
178 tool_use_id,
179 ..
180 }) => {
181 let snippet = truncate_str(content, TOOL_RESULT_SNIPPET_CHARS);
182 let _ = writeln!(
183 buf,
184 "[{msg_id}] user ToolResult (id={tool_use_id}, {len} chars): \"{snippet}\"",
185 len = content.len(),
186 );
187 }
188 _ => {
189 let _ = writeln!(buf, "[{msg_id}] user (non‑text block)");
190 }
191 }
192 }
193
194 fn format_content_block(buf: &mut String, blk_idx: usize, block: &ContentBlock) {
195 match block {
196 ContentBlock::Text { text, .. } => {
197 let snippet = truncate_str(text, TEXT_SNIPPET_CHARS);
198 let _ = writeln!(
199 buf,
200 " [{blk_idx}] Text ({len} chars): \"{snippet}\"",
201 len = text.len(),
202 );
203 }
204 ContentBlock::Thinking { .. } => {
205 // Omit thinking blocks — API-mandated on tool-call messages;
206 // the agent cannot remove them, so listing them only adds noise.
207 }
208 ContentBlock::ToolUse {
209 name, input, id, ..
210 } => {
211 let args = serde_json::to_string(input).unwrap_or_default();
212 let args_preview = truncate_str(&args, TOOL_USE_ARGS_CHARS);
213 let _ = writeln!(
214 buf,
215 " [{blk_idx}] ToolUse ({name}, id={id}, args={args_preview})"
216 );
217 }
218 ContentBlock::ToolResult {
219 content,
220 tool_use_id,
221 ..
222 } => {
223 let snippet = truncate_str(content, TOOL_RESULT_SNIPPET_CHARS);
224 let _ = writeln!(
225 buf,
226 " [{blk_idx}] ToolResult (id={tool_use_id}, {len} chars): \"{snippet}\"",
227 len = content.len(),
228 );
229 }
230 ContentBlock::ServerToolUse {
231 name, input, id, ..
232 } => {
233 let args = serde_json::to_string(input).unwrap_or_default();
234 let args_preview = truncate_str(&args, TOOL_USE_ARGS_CHARS);
235 let _ = writeln!(
236 buf,
237 " [{blk_idx}] ServerToolUse ({name}, id={id}, args={args_preview})"
238 );
239 }
240 ContentBlock::ToolSearchToolResult {
241 tool_use_id,
242 content,
243 ..
244 } => {
245 let snippet = truncate_str(&content.to_string(), TOOL_RESULT_SNIPPET_CHARS);
246 let _ = writeln!(
247 buf,
248 " [{blk_idx}] ToolSearchToolResult (id={tool_use_id}, content={snippet})"
249 );
250 }
251 ContentBlock::CodeExecutionToolResult {
252 tool_use_id,
253 content,
254 ..
255 } => {
256 let snippet = truncate_str(&content.to_string(), TOOL_RESULT_SNIPPET_CHARS);
257 let _ = writeln!(
258 buf,
259 " [{blk_idx}] CodeExecutionToolResult (id={tool_use_id}, content={snippet})"
260 );
261 }
262 ContentBlock::ImageUrl { .. } => {}
263 }
264 }
265
266 fn truncate_str(text: &str, max_chars: usize) -> String {
267 if text.chars().count() <= max_chars {
268 return text.to_string();
269 }
270 let take = max_chars.saturating_sub(3);
271 let mut out: String = text.chars().take(take).collect();
272 out.push_str("...");
273 out
274 }
275
276 // ── Operation parser ────────────────────────────────────────────────────────
277
278 /// Parse the `purge_context` tool input JSON into a list of validated
279 /// `PurgeOp`s. Returns an error string on invalid input.
280 pub fn parse_purge_operations(
281 input: &serde_json::Value,
282 message_count: usize,
283 ) -> Result<Vec<PurgeOp>, String> {
284 let ops = input
285 .get("operations")
286 .and_then(|v| v.as_array())
287 .ok_or_else(|| "missing or invalid 'operations' array".to_string())?;
288
289 let mut parsed = Vec::with_capacity(ops.len());
290
291 for (i, op) in ops.iter().enumerate() {
292 let op_type = op
293 .get("op")
294 .and_then(|v| v.as_str())
295 .ok_or_else(|| format!("operation[{i}]: missing 'op' field"))?;
296
297 let msg = op
298 .get("msg")
299 .and_then(|v| v.as_u64())
300 .ok_or_else(|| format!("operation[{i}]: missing or invalid 'msg'"))?;
301
302 let msg_id = usize::try_from(msg).unwrap_or(usize::MAX);
303 if msg_id == 0 || msg_id > message_count {
304 return Err(format!(
305 "operation[{i}]: msg {msg} out of range (1–{message_count})"
306 ));
307 }
308
309 match op_type {
310 "remove" => {
311 parsed.push(PurgeOp::Remove { msg_id });
312 }
313 "offload" => parsed.push(PurgeOp::Offload { msg_id }),
314 "replace" => {
315 let block_idx = op
316 .get("block")
317 .and_then(|v| v.as_u64())
318 .map(|v| v as usize)
319 .ok_or_else(|| format!("operation[{i}]: 'replace' requires 'block'"))?;
320
321 let pattern_str = op
322 .get("pattern")
323 .and_then(|v| v.as_str())
324 .ok_or_else(|| format!("operation[{i}]: 'replace' requires 'pattern'"))?;
325
326 let with = op
327 .get("with")
328 .and_then(|v| v.as_str())
329 .unwrap_or("")
330 .to_string();
331
332 let pattern = compile_user_regex(pattern_str)
333 .map_err(|e| format!("operation[{i}]: invalid regex pattern: {e}"))?;
334
335 parsed.push(PurgeOp::Replace {
336 msg_id,
337 block_idx,
338 pattern,
339 with,
340 });
341 }
342 other => {
343 return Err(format!(
344 "operation[{i}]: unknown op '{other}' (expected 'remove', 'replace' or 'offload')"
345 ));
346 }
347 }
348 }
349
350 Ok(parsed)
351 }
352
353 // ── Operation executor ──────────────────────────────────────────────────────
354
355 /// Execute a list of purge operations against the message history.
356 ///
357 /// Operations are processed in the order given but effective removal runs
358 /// from highest index to lowest to keep earlier indices stable. After all
359 /// user-requested operations, tool‑call/result pair cascading runs to
360 /// prevent orphaned blocks.
361 pub fn execute_purge_operations(
362 messages: &[Message],
363 ops: &[PurgeOp],
364 session_id: &str,
365 ) -> Result<PurgeResult, String> {
366 let mut offloaded: FastHashSet<usize> = ops
367 .iter()
368 .filter_map(|op| {
369 if let PurgeOp::Offload { msg_id } = op {
370 msg_id.checked_sub(1).filter(|idx| *idx < messages.len())
371 } else {
372 None
373 }
374 })
375 .collect();
376 cascade_tool_pair_removals(messages, &mut offloaded);
377 // Publish original, full-fidelity messages before any destructive operation.
378 // Failure leaves the caller's conversation intact, including mixed operations.
379 let mut pointer = if offloaded.is_empty() {
380 None
381 } else {
382 Some(publish_offloaded_context(session_id, messages, &offloaded)?)
383 };
384 let mut msgs = messages.to_vec();
385 let mut msg_indices_to_remove: FastHashSet<usize> = FastHashSet::default();
386 let mut replaced_count = 0usize;
387
388 // Phase 1: collect removes and apply replaces.
389 for op in ops {
390 match op {
391 PurgeOp::Remove { msg_id } => {
392 let idx = msg_id.saturating_sub(1);
393 if idx < msgs.len() {
394 msg_indices_to_remove.insert(idx);
395 }
396 }
397 PurgeOp::Offload { .. } => {}
398 PurgeOp::Replace {
399 msg_id,
400 block_idx,
401 pattern,
402 with,
403 } => {
404 let idx = msg_id.saturating_sub(1);
405 if idx >= msgs.len() || offloaded.contains(&idx) {
406 continue;
407 }
408 if let Some(block) = msgs[idx].content.get_mut(*block_idx) {
409 let old_text = block_content_text(block).to_string();
410 let new_text = pattern.replace_all(&old_text, with.as_str()).to_string();
411 apply_block_replacement(block, &new_text);
412 replaced_count = replaced_count.saturating_add(1);
413 }
414 }
415 }
416 }
417
418 // Phase 2: cascade removal to tool-call/result counterparts.
419 cascade_tool_pair_removals(&msgs, &mut msg_indices_to_remove);
420
421 // Archival takes precedence over remove/replace when a paired group overlaps.
422 let removed_count = msg_indices_to_remove.difference(&offloaded).count();
423 let first_offloaded = offloaded.iter().min().copied();
424 let mut retained = Vec::with_capacity(msgs.len());
425 for (idx, msg) in msgs.into_iter().enumerate() {
426 if Some(idx) == first_offloaded
427 && let Some(text) = pointer.take()
428 {
429 retained.push(Message {
430 role: messages[idx].role.clone(),
431 content: vec![ContentBlock::Text {
432 text,
433 cache_control: None,
434 }],
435 });
436 }
437 if !offloaded.contains(&idx) && !msg_indices_to_remove.contains(&idx) {
438 retained.push(msg);
439 }
440 }
441
442 Ok(PurgeResult {
443 messages: retained,
444 removed_count,
445 replaced_count,
446 offloaded_count: offloaded.len(),
447 })
448 }
449
450 fn publish_offloaded_context(
451 session_id: &str,
452 messages: &[Message],
453 selected: &FastHashSet<usize>,
454 ) -> Result<String, String> {
455 use crate::tools::large_output_router::{
456 EvidenceArtifact, EvidenceRetentionState, publish_evidence_metadata, unix_millis_now,
457 };
458 let archived: Vec<_> = messages
459 .iter()
460 .enumerate()
461 .filter(|(idx, _)| selected.contains(idx))
462 .map(|(idx, message)| serde_json::json!({"message_id": idx + 1, "message": message}))
463 .collect();
464 let bytes = serde_json::to_vec_pretty(&serde_json::json!({
465 "schema_version": 1,
466 "messages": archived,
467 }))
468 .map_err(|err| format!("Could not encode offloaded context; history unchanged: {err}"))?;
469 let call_id = format!("purge_{}", uuid::Uuid::new_v4());
470 let handle = crate::artifacts::artifact_id_for_tool_call(&call_id);
471 let metadata = EvidenceArtifact {
472 handle: handle.clone(),
473 digest: crate::hashing::sha256_hex(&bytes),
474 size_bytes: bytes.len().try_into().unwrap_or(u64::MAX),
475 content_type: "application/json".to_string(),
476 tool_name: "purge_context".to_string(),
477 call_id,
478 origin_session: session_id.to_string(),
479 generation: 1,
480 redacted: false,
481 encoding: "utf-8".to_string(),
482 retention_state: EvidenceRetentionState::Live,
483 created_at_unix_ms: unix_millis_now(),
484 // This is the only full copy after offload: retain it with the session.
485 retain_until_unix_ms: u64::MAX,
486 storage_path: crate::artifacts::session_artifact_relative_path(&handle),
487 };
488 publish_evidence_metadata(session_id, &metadata)
489 .and_then(|_| {
490 crate::artifacts::write_session_artifact_immutable(session_id, &handle, &bytes)
491 })
492 .map_err(|err| format!("Could not store offloaded context; history unchanged: {err}"))?;
493 Ok(format!(
494 "[Offloaded context: {} messages preserved exactly in session artifact {handle}, generation 1. \
495 Use retrieve_tool_result with ref={handle}, mode=query/lines to inspect, or mode=bytes for exact recovery. \
496 The archive includes original message IDs and complete content blocks. Retained until this session is deleted.]",
497 selected.len()
498 ))
499 }
500
501 /// When a message containing a ToolUse or ToolResult is marked for removal,
502 /// cascade that removal to its unambiguous counterpart. Corrupt or ambiguous
503 /// identities cannot select another exchange for deletion. Runs a fixpoint
504 /// loop until the remove set is closed under proven pairing.
505 fn cascade_tool_pair_removals(messages: &[Message], remove_set: &mut FastHashSet<usize>) {
506 if remove_set.is_empty() {
507 return;
508 }
509
510 // A reused provider ID is not a unique host identity. Ambiguous entries
511 // cannot select another exchange's counterpart for deletion.
512 let mut call_id_to_idx = FastHashMap::default();
513 let mut result_id_to_idx = FastHashMap::default();
514 for (idx, msg) in messages.iter().enumerate() {
515 for block in &msg.content {
516 let Some(key) = block
517 .tool_call_key()
518 .filter(|key| !key.as_str().trim().is_empty())
519 else {
520 continue;
521 };
522 let (map, provider_id) = match block {
523 ContentBlock::ToolUse { id, .. } | ContentBlock::ServerToolUse { id, .. } => {
524 (&mut call_id_to_idx, id.as_str())
525 }
526 ContentBlock::ToolResult { tool_use_id, .. }
527 | ContentBlock::ToolSearchToolResult { tool_use_id, .. }
528 | ContentBlock::CodeExecutionToolResult { tool_use_id, .. } => {
529 (&mut result_id_to_idx, tool_use_id.as_str())
530 }
531 _ => continue,
532 };
533 map.entry(key)
534 .and_modify(|entry| *entry = None)
535 .or_insert(Some((idx, provider_id)));
536 }
537 }
538
539 // Fixpoint: when a tool-call is removed, also remove its result (and vice versa).
540 let max_iters = messages.len().max(10);
541 for _ in 0..max_iters {
542 let snapshot: Vec<usize> = remove_set.iter().copied().collect();
543 let mut changed = false;
544
545 for idx in snapshot {
546 let msg = &messages[idx];
547 for block in &msg.content {
548 match block {
549 ContentBlock::ToolUse { id, .. } | ContentBlock::ServerToolUse { id, .. } => {
550 if let Some(&(result_idx, provider_id)) = block
551 .tool_call_key()
552 .and_then(|key| result_id_to_idx.get(&key))
553 .and_then(Option::as_ref)
554 && provider_id == id.as_str()
555 && block
556 .tool_call_key()
557 .and_then(|key| call_id_to_idx.get(&key))
558 .and_then(Option::as_ref)
559 .is_some_and(|(call_idx, _)| *call_idx == idx)
560 && remove_set.insert(result_idx)
561 {
562 changed = true;
563 }
564 }
565 ContentBlock::ToolResult { tool_use_id, .. }
566 | ContentBlock::ToolSearchToolResult { tool_use_id, .. }
567 | ContentBlock::CodeExecutionToolResult { tool_use_id, .. } => {
568 if let Some(&(call_idx, provider_id)) = block
569 .tool_call_key()
570 .and_then(|key| call_id_to_idx.get(&key))
571 .and_then(Option::as_ref)
572 && provider_id == tool_use_id.as_str()
573 && block
574 .tool_call_key()
575 .and_then(|key| result_id_to_idx.get(&key))
576 .and_then(Option::as_ref)
577 .is_some_and(|(result_idx, _)| *result_idx == idx)
578 && remove_set.insert(call_idx)
579 {
580 changed = true;
581 }
582 }
583 _ => {}
584 }
585 }
586 }
587
588 if !changed {
589 break;
590 }
591 }
592 }
593
594 fn block_content_text(block: &ContentBlock) -> &str {
595 match block {
596 ContentBlock::Text { text, .. } => text,
597 ContentBlock::ToolResult { content, .. } => content,
598 _ => "",
599 }
600 }
601
602 fn apply_block_replacement(block: &mut ContentBlock, new_text: &str) {
603 match block {
604 ContentBlock::Text { text, .. } => {
605 *text = new_text.to_string();
606 }
607 ContentBlock::ToolResult { content, .. } => {
608 *content = new_text.to_string();
609 }
610 _ => {}
611 }
612 }
613
614 // ── Tool definition builder ──────────────────────────────────────────────────
615
616 /// Build the `purge_context` tool definition sent to the model during a purge
617 /// turn. This tool is ad-hoc — it is not registered in the normal tool catalog
618 /// and has no dispatch handler.
619 pub fn build_purge_tool() -> Tool {
620 Tool {
621 tool_type: None,
622 name: "purge_context".to_string(),
623 description: "Remove, condense, or durably offload conversation history to free context window space."
624 .to_string(),
625 input_schema: serde_json::json!({
626 "type": "object",
627 "properties": {
628 "operations": {
629 "type": "array",
630 "items": {
631 "type": "object",
632 "properties": {
633 "op": {"type": "string", "enum": ["remove", "replace", "offload"]},
634 "msg": {"type": "integer"},
635 "block": {"type": "integer"},
636 "pattern": {"type": "string"},
637 "with": {"type": "string"}
638 },
639 "required": ["op", "msg"]
640 }
641 }
642 },
643 "required": ["operations"]
644 }),
645 allowed_callers: None,
646 defer_loading: None,
647 input_examples: None,
648 strict: Some(true),
649 cache_control: None,
650 }
651 }
652
653 // ── Orchestration ────────────────────────────────────────────────────────────
654
655 /// Run a full purge cycle: build the prompt, call the model with the
656 /// `purge_context` tool, parse the response, and execute the operations.
657 ///
658 /// Returns the `PurgeResult` with the modified message list on success,
659 /// or a human-readable error string on failure.
660 ///
661 /// Cost reporting is handled internally as a side-effect of the API call.
662 /// The caller is responsible for emitting start/completed/failed events
663 /// and for replacing the session message list with `PurgeResult.messages`.
664 pub async fn run_purge(
665 client: &impl LlmClient,
666 _provider: ProviderKind,
667 session_id: &str,
668 messages: &[Message],
669 model: &str,
670 reasoning_effort: Option<String>,
671 max_tokens: u32,
672 ) -> Result<PurgeResult, String> {
673 // 1. Build the purge prompt from the current conversation.
674 let prompt = build_purge_prompt(messages);
675
676 // 2. Clone messages and inject the prompt as a user message.
677 let mut request_messages = messages.to_vec();
678 request_messages.push(Message {
679 role: Role::User,
680 content: vec![ContentBlock::Text {
681 text: prompt,
682 cache_control: None,
683 }],
684 });
685
686 // 3. Build the tool definition and the request.
687 let purge_tool = build_purge_tool();
688 let request = MessageRequest {
689 model: model.to_string(),
690 messages: request_messages,
691 max_tokens,
692 system: None,
693 tools: Some(vec![purge_tool]),
694 tool_choice: None,
695 metadata: None,
696 thinking: None,
697 reasoning_effort,
698 stream: Some(false),
699 temperature: None,
700 top_p: None,
701 };
702
703 // 4. Send to the model. Capture the session scope before awaiting so a
704 // late response cannot accrue into a subsequently loaded/new session.
705 let cost_scope = crate::cost_status::scope_token();
706 let cost_route = client.effective_route_envelope(model, chrono::Utc::now());
707 let response = client
708 .create_message(request)
709 .await
710 .map_err(|e| format!("Purge API error: {e}"))?;
711
712 // Report the route, not just the provider name: the endpoint decides
713 // whether this is a metered public API, a plan quota, or a local runtime.
714 // Purge currently has TUI admission only. Freeze that session origin rather
715 // than borrowing a previous Runtime turn's owner from ambient config.
716 let source_id = format!(
717 "purge:{}:{}",
718 cost_route
719 .dispatched_at
720 .timestamp_nanos_opt()
721 .unwrap_or_default(),
722 response.id
723 );
724 crate::cost_status::report_effective_route_for_interactive_origin(
725 cost_scope,
726 session_id,
727 &source_id,
728 &source_id,
729 &cost_route,
730 &response.usage,
731 );
732
733 // A truncated response can still carry a complete-looking `purge_context`
734 // call; executing it would mutate the session from incomplete output.
735 if codewhale_models::is_incomplete_stop_reason(response.stop_reason.as_deref()) {
736 return Err(format!(
737 "Purge model response incomplete: provider stop reason `{}`; no purge was applied.",
738 codewhale_models::stop_reason_detail(response.stop_reason.as_deref())
739 ));
740 }
741
742 // 5. Find the `purge_context` tool call in the response.
743 let tool_input = response.content.iter().find_map(|block| {
744 if let ContentBlock::ToolUse { name, input, .. } = block
745 && name == "purge_context"
746 {
747 return Some(input.clone());
748 }
749 None
750 });
751
752 match tool_input {
753 Some(input) => {
754 let ops = parse_purge_operations(&input, messages.len())
755 .map_err(|e| format!("Purge parse error: {e}"))?;
756 execute_purge_operations(messages, &ops, session_id)
757 }
758 None => Err("Purge: model did not call purge_context tool".to_string()),
759 }
760 }
761
762 // ── Tests ───────────────────────────────────────────────────────────────────
763
764 #[cfg(test)]
765 mod tests {
766 use super::*;
767 use serde_json::json;
768
769 fn msg_text(role: &str, text: &str) -> Message {
770 Message {
771 role: Role::from(role),
772 content: vec![ContentBlock::Text {
773 text: text.to_string(),
774 cache_control: None,
775 }],
776 }
777 }
778
779 fn msg_tool_use(id: &str, name: &str, input: serde_json::Value) -> Message {
780 Message {
781 role: Role::Assistant,
782 content: vec![ContentBlock::ToolUse {
783 execution_id: None,
784 id: id.to_string(),
785 name: name.to_string(),
786 input,
787 caller: None,
788 thought_signature: None,
789 }],
790 }
791 }
792
793 fn msg_tool_result(id: &str, content: &str) -> Message {
794 Message {
795 role: Role::User,
796 content: vec![ContentBlock::ToolResult {
797 execution_id: None,
798 tool_use_id: id.to_string(),
799 content: content.to_string(),
800 is_error: None,
801 content_blocks: None,
802 }],
803 }
804 }
805
806 #[test]
807 fn purge_counterparts_use_execution_identity_without_legacy_fallback() {
808 let messages: Vec<Message> = serde_json::from_value(json!([
809 {"role":"assistant","content":[{"type":"tool_use","id":"wire","execution_id":"first","name":"read","input":{}}]},
810 {"role":"user","content":[{"type":"tool_result","tool_use_id":"wire","execution_id":"first","content":"first"}]},
811 {"role":"assistant","content":[{"type":"tool_use","id":"wire","execution_id":"second","name":"read","input":{}}]},
812 {"role":"user","content":[{"type":"tool_result","tool_use_id":"wire","execution_id":"second","content":"second"}]},
813 {"role":"assistant","content":[{"type":"tool_use","id":"first","name":"read","input":{}}]},
814 {"role":"user","content":[{"type":"tool_result","tool_use_id":"first","content":"legacy"}]}
815 ])).unwrap();
816 for (selected, expected) in [(0, vec![0, 1]), (3, vec![2, 3]), (4, vec![4, 5])] {
817 let mut remove = FastHashSet::from_iter([selected]);
818 cascade_tool_pair_removals(&messages, &mut remove);
819 assert_eq!(remove, FastHashSet::from_iter(expected));
820 }
821 let mut ambiguous = messages.clone();
822 ambiguous.push(messages[0].clone());
823 let mut remove = FastHashSet::from_iter([0]);
824 cascade_tool_pair_removals(&ambiguous, &mut remove);
825 assert_eq!(remove, FastHashSet::from_iter([0]));
826 }
827
828 #[test]
829 fn offload_round_trips_full_paired_context_through_the_existing_retrieval_tool() {
830 use crate::test_support::{EnvVarGuard, lock_test_env};
831 use crate::tools::spec::{ToolContext, ToolSpec};
832 use crate::tools::tool_result_retrieval::RetrieveToolResultTool;
833 use base64::Engine as _;
834
835 let _env = lock_test_env();
836 let _cost_guard = crate::cost_status::test_scope();
837 let home = tempfile::tempdir().unwrap();
838 let state = home.path().join("explicit-state");
839 let _state = EnvVarGuard::set("CODEWHALE_HOME", &state);
840 let session_id = "purge-owned-session";
841 let messages: Vec<Message> = serde_json::from_value(json!([
842 {"role":"user", "content":[{"type":"text", "text":"Keep the task"}]},
843 {"role":"assistant", "content":[
844 {"type":"thinking", "thinking":"retained reasoning", "signature":"signed-exact"},
845 {"type":"tool_use", "id":"read-original", "name":"read_file", "input":{"path":"old.rs"}, "thought_signature":"google-exact"}
846 ]},
847 {"role":"user", "content":[
848 {"type":"tool_result", "tool_use_id":"read-original", "content":"original-file-data\n".repeat(1000), "content_blocks":[{"type":"image","source":{"data":"Zml4dHVyZQ=="}}]},
849 {"type":"image_url", "image_url":{"url":"data:image/png;base64,Zml4dHVyZQ=="}}
850 ]},
851 {"role":"assistant", "content":[{"type":"text", "text":"Keep the current answer"}]}
852 ])).unwrap();
853 let original = messages.clone();
854 let mock = MockLlmClient::new(vec![]);
855 mock.push_message_response(msg_response_with_tool_call(json!([
856 {"op":"offload", "msg":3}
857 ])));
858 let runtime = tokio::runtime::Builder::new_current_thread()
859 .enable_all()
860 .build()
861 .unwrap();
862 let result = runtime
863 .block_on(run_purge(
864 &mock,
865 ProviderKind::Deepseek,
866 session_id,
867 &messages,
868 "mock",
869 None,
870 4096,
871 ))
872 .unwrap();
873 assert_eq!(messages, original);
874 assert_eq!(result.offloaded_count, 2);
875 assert_eq!(result.removed_count, 0);
876 assert_eq!(result.messages.len(), 3);
877 assert_eq!(result.messages[0], original[0]);
878 assert_eq!(result.messages[2], original[3]);
879 assert!(
880 serde_json::to_vec(&result.messages).unwrap().len()
881 < serde_json::to_vec(&original).unwrap().len()
882 );
883 let pointer = block_content_text(&result.messages[1].content[0]);
884 let handle = pointer
885 .split_whitespace()
886 .find(|part| part.starts_with("art_purge_"))
887 .unwrap()
888 .trim_end_matches(',');
889 let context = ToolContext::new(home.path()).with_state_namespace(session_id);
890 let retrieved = runtime
891 .block_on(RetrieveToolResultTool.execute(
892 json!({"ref":handle,"mode":"bytes","generation":1,"max_bytes":131072}),
893 &context,
894 ))
895 .unwrap();
896 let payload: serde_json::Value = serde_json::from_str(&retrieved.content).unwrap();
897 let bytes = base64::engine::general_purpose::STANDARD
898 .decode(payload["data"].as_str().unwrap())
899 .unwrap();
900 let archive: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
901 let restored: Vec<Message> = archive["messages"]
902 .as_array()
903 .unwrap()
904 .iter()
905 .map(|entry| serde_json::from_value(entry["message"].clone()).unwrap())
906 .collect();
907 assert_eq!(restored, original[1..3]);
908 assert_eq!(archive["messages"][0]["message_id"], 2);
909 assert_eq!(archive["messages"][1]["message_id"], 3);
910 let metadata =
911 crate::tools::large_output_router::read_evidence_metadata(session_id, handle).unwrap();
912 assert_eq!(metadata.retain_until_unix_ms, u64::MAX);
913 let path = state
914 .join("sessions")
915 .join(session_id)
916 .join(&metadata.storage_path);
917 assert_eq!(std::fs::read(&path).unwrap(), bytes);
918 let other = ToolContext::new(home.path()).with_state_namespace("another-session");
919 assert!(
920 runtime
921 .block_on(
922 RetrieveToolResultTool.execute(json!({"ref":handle,"mode":"bytes"}), &other)
923 )
924 .is_err()
925 );
926 std::fs::write(path, b"changed fixture").unwrap();
927 assert!(
928 runtime
929 .block_on(
930 RetrieveToolResultTool.execute(json!({"ref":handle,"mode":"bytes"}), &context)
931 )
932 .is_err()
933 );
934 }
935
936 #[test]
937 fn offload_closes_server_tool_pairs_and_wins_over_overlapping_destructive_operations() {
938 use crate::test_support::{EnvVarGuard, lock_test_env};
939 let _env = lock_test_env();
940 let home = tempfile::tempdir().unwrap();
941 let _state = EnvVarGuard::set("CODEWHALE_HOME", home.path());
942 let messages: Vec<Message> = serde_json::from_value(json!([
943 {"role":"assistant", "content":[
944 {"type":"server_tool_use", "id":"search", "name":"tool_search", "input":{}},
945 {"type":"server_tool_use", "id":"execute", "name":"code_execution", "input":{}}
946 ]},
947 {"role":"assistant", "content":[{"type":"tool_search_tool_result", "tool_use_id":"search", "content":{"tools":["read"]}}]},
948 {"role":"assistant", "content":[{"type":"code_execution_tool_result", "tool_use_id":"execute", "content":{"output":"exact"}}]},
949 {"role":"user", "content":[{"type":"text", "text":"keep"}]}
950 ])).unwrap();
951 let ops = parse_purge_operations(
952 &json!({"operations":[
953 {"op":"remove", "msg":1}, {"op":"offload", "msg":2},
954 {"op":"replace", "msg":3, "block":0, "pattern":"exact", "with":"lost"}
955 ]}),
956 messages.len(),
957 )
958 .unwrap();
959 let result = execute_purge_operations(&messages, &ops, "server-pairs").unwrap();
960 assert_eq!(result.offloaded_count, 3);
961 assert_eq!(result.removed_count, 0);
962 assert_eq!(result.replaced_count, 0);
963 assert_eq!(result.messages.len(), 2);
964 assert_eq!(result.messages[1], messages[3]);
965 }
966
967 #[test]
968 fn offload_storage_failure_leaves_mixed_purge_history_unchanged() {
969 use crate::test_support::{EnvVarGuard, lock_test_env};
970 let _env = lock_test_env();
971 let home = tempfile::tempdir().unwrap();
972 let _state = EnvVarGuard::set("CODEWHALE_HOME", home.path());
973 std::fs::write(home.path().join("sessions"), b"not a directory").unwrap();
974 let messages = vec![msg_text("user", "keep"), msg_text("assistant", "original")];
975 let original = messages.clone();
976 let ops = parse_purge_operations(&json!({"operations":[
977 {"op":"remove", "msg":1}, {"op":"replace", "msg":2, "block":0, "pattern":"original", "with":"lost"}, {"op":"offload", "msg":2}
978 ]}), messages.len()).unwrap();
979 let error = execute_purge_operations(&messages, &ops, "storage-failure").unwrap_err();
980 assert!(error.contains("history unchanged"));
981 assert_eq!(messages, original);
982 }
983
984 #[test]
985 fn parse_remove_operations() {
986 let input = json!({
987 "operations": [
988 {"op": "remove", "msg": 1},
989 {"op": "remove", "msg": 3}
990 ]
991 });
992 let ops = parse_purge_operations(&input, 5).unwrap();
993 assert_eq!(ops.len(), 2);
994 assert!(matches!(ops[0], PurgeOp::Remove { msg_id: 1 }));
995 assert!(matches!(ops[1], PurgeOp::Remove { msg_id: 3 }));
996 }
997
998 #[test]
999 fn parse_replace_operation() {
1000 let input = json!({
1001 "operations": [
1002 {"op": "replace", "msg": 2, "block": 0, "pattern": "hello", "with": "hi"}
1003 ]
1004 });
1005 let ops = parse_purge_operations(&input, 5).unwrap();
1006 assert_eq!(ops.len(), 1);
1007 assert!(matches!(ops[0], PurgeOp::Replace { msg_id: 2, .. }));
1008 }
1009
1010 #[test]
1011 fn parse_rejects_out_of_range_msg() {
1012 let input = json!({"operations": [{"op": "remove", "msg": 10}]});
1013 assert!(parse_purge_operations(&input, 5).is_err());
1014 }
1015
1016 #[test]
1017 fn parse_rejects_invalid_regex() {
1018 let input = json!({
1019 "operations": [{"op": "replace", "msg": 1, "block": 0, "pattern": "[", "with": "x"}]
1020 });
1021 assert!(parse_purge_operations(&input, 5).is_err());
1022 }
1023
1024 #[test]
1025 fn execute_remove_works() {
1026 let msgs = vec![
1027 msg_text("user", "hello"),
1028 msg_text("assistant", "hi there"),
1029 msg_text("user", "bye"),
1030 ];
1031 let ops = vec![PurgeOp::Remove { msg_id: 2 }];
1032 let result = execute_purge_operations(&msgs, &ops, "purge-test").unwrap();
1033 assert_eq!(result.removed_count, 1);
1034 assert_eq!(result.messages.len(), 2);
1035 }
1036
1037 #[test]
1038 fn execute_replace_text_block() {
1039 let msgs = vec![msg_text("assistant", "Hello world! Hello again!")];
1040 let pattern = Regex::new("Hello").unwrap();
1041 let ops = vec![PurgeOp::Replace {
1042 msg_id: 1,
1043 block_idx: 0,
1044 pattern,
1045 with: "Hi".to_string(),
1046 }];
1047 let result = execute_purge_operations(&msgs, &ops, "purge-test").unwrap();
1048 assert_eq!(result.replaced_count, 1);
1049
1050 if let ContentBlock::Text { text, .. } = &result.messages[0].content[0] {
1051 assert_eq!(text, "Hi world! Hi again!");
1052 } else {
1053 panic!("expected text block");
1054 }
1055 }
1056
1057 #[test]
1058 fn tool_call_result_pairing_cascaded() {
1059 // Message 2 (idx 1) is a tool call. Message 3 (idx 2) is its result.
1060 // Removing the tool call should cascade to remove the result too.
1061 let msgs = vec![
1062 msg_text("user", "read a file"),
1063 msg_tool_use("call_01", "read_file", json!({"path": "x.rs"})),
1064 msg_tool_result("call_01", "fn main() {}"),
1065 ];
1066 let ops = vec![PurgeOp::Remove { msg_id: 2 }]; // remove tool call only
1067 let result = execute_purge_operations(&msgs, &ops, "purge-test").unwrap();
1068 // Both tool call and its result should be gone (cascaded).
1069 assert_eq!(
1070 result.removed_count, 2,
1071 "tool call + its result should both be removed"
1072 );
1073 assert_eq!(result.messages.len(), 1);
1074 }
1075
1076 #[test]
1077 fn tool_result_removal_cascades_to_call() {
1078 // Removing the result should cascade to remove the call.
1079 let msgs = vec![
1080 msg_text("user", "read a file"),
1081 msg_tool_use("call_01", "read_file", json!({"path": "x.rs"})),
1082 msg_tool_result("call_01", "fn main() {}"),
1083 ];
1084 let ops = vec![PurgeOp::Remove { msg_id: 3 }]; // remove result only
1085 let result = execute_purge_operations(&msgs, &ops, "purge-test").unwrap();
1086 assert_eq!(
1087 result.removed_count, 2,
1088 "tool result + its call should both be removed"
1089 );
1090 assert_eq!(result.messages.len(), 1);
1091 }
1092
1093 #[test]
1094 fn prompt_truncates_long_content() {
1095 let long_text = "x".repeat(200);
1096 let msgs = vec![msg_text("user", &long_text)];
1097 let prompt = build_purge_prompt(&msgs);
1098 assert!(prompt.contains("(200 chars)"));
1099 assert!(prompt.contains("xxx...")); // truncated
1100 assert!(!prompt.contains(&long_text));
1101 }
1102
1103 #[test]
1104 fn prompt_shows_full_short_content() {
1105 let msgs = vec![msg_text("user", "hi")];
1106 let prompt = build_purge_prompt(&msgs);
1107 assert!(prompt.contains("\"hi\""));
1108 assert!(!prompt.contains("..."));
1109 }
1110
1111 #[test]
1112 fn prompt_omits_thinking_blocks() {
1113 let msgs = vec![Message {
1114 role: Role::Assistant,
1115 content: vec![
1116 ContentBlock::Thinking {
1117 signature: None,
1118 state: None,
1119 thinking: "let me think...".to_string(),
1120 },
1121 ContentBlock::Text {
1122 text: "done".to_string(),
1123 cache_control: None,
1124 },
1125 ],
1126 }];
1127 let prompt = build_purge_prompt(&msgs);
1128 assert!(!prompt.contains("let me think"));
1129 assert!(prompt.contains("Text (4 chars)"));
1130 }
1131
1132 #[test]
1133 fn build_purge_tool_has_correct_shape() {
1134 let tool = build_purge_tool();
1135 assert_eq!(tool.name, "purge_context");
1136 let schema = &tool.input_schema;
1137 assert_eq!(schema["type"], "object");
1138 assert!(schema["properties"]["operations"]["type"] == "array");
1139 let ops_item = &schema["properties"]["operations"]["items"];
1140 assert_eq!(ops_item["type"], "object");
1141 let required = ops_item["required"].as_array().unwrap();
1142 assert!(required.contains(&json!("op")));
1143 assert!(required.contains(&json!("msg")));
1144 }
1145
1146 use crate::llm_client::mock::MockLlmClient;
1147 use codewhale_models::{MessageResponse, Usage};
1148
1149 fn msg_response_with_tool_call(operations: serde_json::Value) -> MessageResponse {
1150 MessageResponse {
1151 id: "resp_test".to_string(),
1152 r#type: "message".to_string(),
1153 role: "assistant".to_string(),
1154 content: vec![ContentBlock::ToolUse {
1155 execution_id: None,
1156 id: "call_purge".to_string(),
1157 name: "purge_context".to_string(),
1158 input: json!({"operations": operations}),
1159 caller: None,
1160 thought_signature: None,
1161 }],
1162 model: "mock-model".to_string(),
1163 stop_reason: None,
1164 stop_sequence: None,
1165 container: None,
1166 usage: Usage::default(),
1167 }
1168 }
1169
1170 fn msg_response_without_tool_call(text: &str) -> MessageResponse {
1171 MessageResponse {
1172 id: "resp_plain".to_string(),
1173 r#type: "message".to_string(),
1174 role: "assistant".to_string(),
1175 content: vec![ContentBlock::Text {
1176 text: text.to_string(),
1177 cache_control: None,
1178 }],
1179 model: "mock".to_string(),
1180 stop_reason: None,
1181 stop_sequence: None,
1182 container: None,
1183 usage: Usage::default(),
1184 }
1185 }
1186
1187 #[tokio::test]
1188 async fn run_purge_removes_message() {
1189 let _env = crate::test_support::lock_test_env();
1190 let home = tempfile::tempdir().unwrap();
1191 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", home.path());
1192 let _cost_guard = crate::cost_status::test_scope();
1193 let mock = MockLlmClient::new(vec![]);
1194 mock.push_message_response(msg_response_with_tool_call(json!([
1195 {"op": "remove", "msg": 2}
1196 ])));
1197
1198 let messages = vec![
1199 msg_text("user", "hello"),
1200 msg_text("assistant", "remove me"),
1201 msg_text("user", "bye"),
1202 ];
1203
1204 let result = run_purge(
1205 &mock,
1206 ProviderKind::Deepseek,
1207 "purge-test",
1208 &messages,
1209 "mock",
1210 None,
1211 4096,
1212 )
1213 .await
1214 .unwrap();
1215 assert_eq!(result.removed_count, 1);
1216 assert_eq!(result.replaced_count, 0);
1217 assert_eq!(result.messages.len(), 2);
1218
1219 if let ContentBlock::Text { text, .. } = &result.messages[0].content[0] {
1220 assert_eq!(text, "hello");
1221 } else {
1222 panic!(
1223 "expected text block, got {:?}",
1224 result.messages[0].content[0]
1225 );
1226 }
1227 if let ContentBlock::Text { text, .. } = &result.messages[1].content[0] {
1228 assert_eq!(text, "bye");
1229 } else {
1230 panic!(
1231 "expected text block, got {:?}",
1232 result.messages[1].content[0]
1233 );
1234 }
1235 }
1236
1237 #[tokio::test]
1238 async fn run_purge_replace_condenses_text() {
1239 let _env = crate::test_support::lock_test_env();
1240 let home = tempfile::tempdir().unwrap();
1241 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", home.path());
1242 let _cost_guard = crate::cost_status::test_scope();
1243 let mock = MockLlmClient::new(vec![]);
1244 mock.push_message_response(msg_response_with_tool_call(json!([
1245 {"op": "replace", "msg": 1, "block": 0, "pattern": "very long and verbose", "with": "short"}
1246 ])));
1247
1248 let messages = vec![msg_text("assistant", "this is very long and verbose text")];
1249
1250 let result = run_purge(
1251 &mock,
1252 ProviderKind::Deepseek,
1253 "purge-test",
1254 &messages,
1255 "mock",
1256 None,
1257 4096,
1258 )
1259 .await
1260 .unwrap();
1261 assert_eq!(result.removed_count, 0);
1262 assert_eq!(result.replaced_count, 1);
1263
1264 if let ContentBlock::Text { text, .. } = &result.messages[0].content[0] {
1265 assert_eq!(text, "this is short text");
1266 } else {
1267 panic!(
1268 "expected text block, got {:?}",
1269 result.messages[0].content[0]
1270 );
1271 }
1272 }
1273
1274 #[tokio::test]
1275 async fn run_purge_errors_when_no_tool_call() {
1276 let _env = crate::test_support::lock_test_env();
1277 let home = tempfile::tempdir().unwrap();
1278 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", home.path());
1279 let _cost_guard = crate::cost_status::test_scope();
1280 let mock = MockLlmClient::new(vec![]);
1281 mock.push_message_response(msg_response_without_tool_call("nothing to clean up"));
1282
1283 let messages = vec![msg_text("user", "hi")];
1284 let err = run_purge(
1285 &mock,
1286 ProviderKind::Deepseek,
1287 "purge-test",
1288 &messages,
1289 "mock",
1290 None,
1291 4096,
1292 )
1293 .await
1294 .unwrap_err();
1295 assert!(err.contains("did not call purge_context"));
1296 }
1297
1298 #[tokio::test]
1299 async fn run_purge_errors_on_api_failure() {
1300 let _env = crate::test_support::lock_test_env();
1301 let home = tempfile::tempdir().unwrap();
1302 let _home = crate::test_support::EnvVarGuard::set("CODEWHALE_HOME", home.path());
1303 let _cost_guard = crate::cost_status::test_scope();
1304 // No canned response — MockLlmClient returns an error.
1305 let mock = MockLlmClient::new(vec![]);
1306 let messages = vec![msg_text("user", "hi")];
1307 let err = run_purge(
1308 &mock,
1309 ProviderKind::Deepseek,
1310 "purge-test",
1311 &messages,
1312 "mock",
1313 None,
1314 4096,
1315 )
1316 .await
1317 .unwrap_err();
1318 assert!(err.contains("Purge API error"));
1319 }
1320 }
1321
1321 lines RUST