1use crate::cancel::CancelToken;
2use crate::compact::CompactState;
3use crate::context::ContextPacket;
4use crate::event::{AgentEvent, RepetitionWarningKind};
5use crate::harness::{
6 AgentRunState, HarnessPolicy, HarnessStop, HarnessStopReason, ToolLoopDecision,
7 tool_error_result, trace_request_summary,
8};
9use crate::model::{
10 AttachmentKind, ContentPart, ModelMessage, ModelProvider, ModelRequest, ModelRole,
11 ModelStreamEvent, ThinkingConfig,
12};
13use crate::prompt::{PromptCache, SystemPromptInput};
14use crate::runtime_components::RuntimeComponents;
15use crate::security::SecurityDecision;
16use crate::skills::SkillManifest;
17use crate::tool::{ToolExecutor, ToolParallelism, take_tool_content_parts};
18use anyhow::Result;
19use futures_util::StreamExt;
20use serde_json::{Value, json};
21use std::hash::{Hash, Hasher};
22use std::path::PathBuf;
23use std::sync::{Arc, RwLock};
24
25const QUESTION_TOOL_NAME: &str = "question";
26const PLAN_TOOL_NAME: &str = "plan";
27
28struct ModelTurnOutput {
29 text: String,
30 thinking: String,
31 tool_calls: Vec<crate::tool::ToolInvocation>,
32 harness_stop: Option<HarnessStop>,
33}
34
35type ToolExecutionResult = (
36 crate::tool::ToolInvocation,
37 crate::tool::ToolResult,
38 String,
39 Vec<ContentPart>,
40);
41
42pub struct TurnContext {
43 pub model_provider: Arc<RwLock<Arc<dyn ModelProvider>>>,
44 pub tool_executor: Arc<ToolExecutor>,
45 pub project_dir: PathBuf,
46 pub data_dir: PathBuf,
47 pub model_name: Arc<RwLock<String>>,
48 pub event_tx: Option<tokio::sync::mpsc::UnboundedSender<AgentEvent>>,
49 pub approval_resolver: crate::runtime::ApprovalResolver,
50 pub question_resolver: crate::runtime::QuestionResolver,
51 pub plan_review_resolver: crate::runtime::PlanReviewResolver,
52 pub sudo_password_resolver: crate::runtime::SudoPasswordResolver,
53 pub compact_state: Arc<tokio::sync::Mutex<CompactState>>,
54 pub harness_config: crate::config::HarnessConfig,
55 pub include_tool_prompt_manifest: bool,
56 pub context_packets: Arc<std::sync::Mutex<Vec<ContextPacket>>>,
57 pub available_skills: Arc<std::sync::Mutex<Vec<SkillManifest>>>,
58 pub skill_pools: Arc<std::sync::Mutex<Vec<crate::skills::SkillPool>>>,
60 pub active_skills: Arc<std::sync::Mutex<Vec<SkillManifest>>>,
61 pub prompt_cache: Arc<PromptCache>,
62 pub components: RuntimeComponents,
63 pub cancel_token: CancelToken,
64 pub instructions: Arc<RwLock<Option<String>>>,
68 pub prompt_prefix: Arc<std::sync::Mutex<Option<Vec<ModelMessage>>>>,
73 pub config: Arc<RwLock<crate::config::NaviConfig>>,
77 pub memory_injection: Option<String>,
79 pub compaction_provider: Option<Arc<dyn ModelProvider>>,
82 pub agent_mode: crate::plan_mode::AgentMode,
85 pub compaction_model_name: Option<String>,
87 pub session_id: String,
88 pub allowed_tool_names: Option<Vec<String>>,
91 pub is_subagent: bool,
94 pub memory_manager: Arc<std::sync::Mutex<Option<Arc<crate::memory::MemoryManager>>>>,
98 pub harness_card: Option<String>,
100}
101
102impl TurnContext {
103 pub fn active_model_provider(&self) -> Arc<dyn ModelProvider> {
104 self.model_provider
105 .read()
106 .unwrap_or_else(|e| e.into_inner())
107 .clone()
108 }
109
110 pub fn active_model_name(&self) -> String {
111 self.model_name
112 .read()
113 .unwrap_or_else(|e| e.into_inner())
114 .clone()
115 }
116
117 pub fn active_config(&self) -> crate::config::NaviConfig {
118 self.config
119 .read()
120 .unwrap_or_else(|e| e.into_inner())
121 .clone()
122 }
123
124 pub fn get_or_init_memory_manager(&self) -> Result<Option<Arc<crate::memory::MemoryManager>>> {
127 let memory_config = self.active_config().memory;
128 if !memory_config.enabled {
129 return Ok(None);
130 }
131 let mut guard = self
132 .memory_manager
133 .lock()
134 .unwrap_or_else(|e| e.into_inner());
135 if let Some(manager) = guard.as_ref() {
136 return Ok(Some(manager.clone()));
137 }
138 let manager = Arc::new(crate::memory::MemoryManager::new(
139 self.project_dir.clone(),
140 self.data_dir.clone(),
141 &memory_config,
142 )?);
143 *guard = Some(manager.clone());
144 Ok(Some(manager))
145 }
146
147 pub fn cancellation_requested(&self) -> bool {
148 self.cancel_token.is_requested()
149 }
150
151 pub fn resolve_approval(&self, decision: crate::event::ApprovalDecision) -> bool {
152 self.approval_resolver.resolve(decision)
153 }
154}
155
156pub struct Prompt {
157 pub input: Vec<ModelMessage>,
158 pub tools: Vec<crate::tool::ToolDefinition>,
159 pub base_instructions: String,
160}
161
162pub async fn run_turn(
163 ctx: &TurnContext,
164 messages: &mut Vec<ModelMessage>,
165 policy: HarnessPolicy,
166) -> Result<String> {
167 ensure_not_cancelled(ctx)?;
168 ensure_system_prompt(ctx, messages).await;
169
170 let mut run_state = AgentRunState::default();
171 let final_text = loop {
172 ensure_not_cancelled(ctx)?;
173 maintain_context_budget(ctx, messages).await;
174 ensure_not_cancelled(ctx)?;
175
176 let request = build_model_request(ctx, messages);
177 emit_request_trace(ctx, &request, policy);
178
179 let output = collect_model_output(ctx, request).await?;
180 ensure_not_cancelled(ctx)?;
181
182 if let Some(stop) = output.harness_stop.clone() {
183 let text = finalize_harness_stop(ctx, messages, stop);
184 break text;
185 }
186
187 if !output.tool_calls.is_empty() {
188 if let Some(text) =
189 handle_tool_calls(ctx, messages, &mut run_state, policy, output).await
190 {
191 break text;
192 }
193 continue;
194 }
195
196 persist_final_model_output(ctx, messages, &output);
197 break output.text;
198 };
199
200 let _ = sync_messages_to_history(ctx, messages).await;
201 Ok(final_text)
202}
203
204fn ensure_not_cancelled(ctx: &TurnContext) -> Result<()> {
205 if ctx.cancellation_requested() {
206 Err(anyhow::anyhow!("turn cancelled"))
207 } else {
208 Ok(())
209 }
210}
211
212async fn ensure_system_prompt(ctx: &TurnContext, messages: &mut Vec<ModelMessage>) {
213 if let Some(prefix) = ctx
214 .prompt_prefix
215 .lock()
216 .unwrap_or_else(|e| e.into_inner())
217 .clone()
218 {
219 replace_prompt_prefix(messages, prefix);
220 return;
221 }
222
223 let context_packets = ctx
224 .context_packets
225 .lock()
226 .unwrap_or_else(|e| e.into_inner())
227 .clone();
228 let active_skills = ctx
229 .active_skills
230 .lock()
231 .unwrap_or_else(|e| e.into_inner())
232 .clone();
233 let available_skills = ctx
234 .available_skills
235 .lock()
236 .unwrap_or_else(|e| e.into_inner())
237 .clone();
238 let skill_pools = ctx
239 .skill_pools
240 .lock()
241 .unwrap_or_else(|e| e.into_inner())
242 .clone();
243 let memory_injection = combined_memory_injection(ctx).await;
244 let mut tools = ctx.tool_executor.definitions();
245
246 if ctx.agent_mode.restricts_tools() {
248 tools.retain(|t| crate::plan_mode::is_tool_allowed_in_plan_mode_named(&t.name, t.kind));
249 }
250
251 let input = SystemPromptInput {
252 config: ctx.active_config(),
253 project_dir: ctx.project_dir.clone(),
254 memory_injection,
255 tools: ctx
256 .components
257 .harness
258 .filter_tools(tools, ctx.allowed_tool_names.as_deref()),
259 include_tool_prompt_manifest: ctx.include_tool_prompt_manifest,
260 context_packets,
261 available_skills,
262 active_skills,
263 skill_pools,
264 harness_card: ctx.harness_card.clone(),
265 };
266 let prompt = ctx.components.prompt.clone();
267 let prompt_cache = ctx.prompt_cache.clone();
268 let rendered = tokio::task::spawn_blocking(move || prompt.build(input, prompt_cache))
269 .await
270 .unwrap_or_else(|_| crate::prompt::RenderedPrompt {
271 instructions: "Default NAVI base instructions".to_string(),
272 developer_messages: Vec::new(),
273 });
274
275 *ctx.instructions.write().unwrap_or_else(|e| e.into_inner()) =
279 Some(rendered.instructions.clone());
280
281 let mut prefix = Vec::with_capacity(2 + rendered.developer_messages.len());
285 prefix.push(ModelMessage::system(rendered.instructions));
286 if ctx.agent_mode.restricts_tools() {
288 let plan_path = ctx
289 .tool_executor
290 .policy()
291 .plan_file_path()
292 .unwrap_or_else(|| {
293 crate::plan_store::session_plan_file_path(&ctx.data_dir, &ctx.session_id)
294 });
295 let plan_exists = plan_path.is_file();
296 let file_info = if plan_exists {
297 format!(
298 "A plan file already exists at `{}`. Read it and edit incrementally with `edit` or rewrite with `write_file` / plan(action='write').",
299 plan_path.display()
300 )
301 } else {
302 format!(
303 "No plan file yet. Create it at `{}` using write_file or plan(action='write', plan='...markdown...').",
304 plan_path.display()
305 )
306 };
307 prefix.push(ModelMessage::developer(format!(
308 "Plan mode is active. The user does not want execution yet — do NOT edit project files, \
309run non-readonly tools, change configs, or make commits. This supersedes other instructions.\n\
310\n\
311## Plan file\n\
312{file_info}\n\
313This is the ONLY path you may write. Build the plan as a **markdown design document** (not JSON).\n\
314\n\
315## Workflow\n\
3161. Explore with read-only tools (search, read_file, …).\n\
3172. Draft/update the plan file incrementally (Context, Approach, Files to modify with paths, Verification).\n\
3183. Use `question` for requirements/approach clarifications — not to ask \"is this plan okay?\".\n\
3194. When the plan is ready for approval, call plan(action='submit') (reads the plan file; empty args are fine).\n\
320\n\
321Your turn should end with either `question` or plan(action='submit'). After approval the host exits plan mode and you implement."
322 )));
323 }
324 prefix.extend(rendered.developer_messages);
325 *ctx.prompt_prefix.lock().unwrap_or_else(|e| e.into_inner()) = Some(prefix.clone());
326 replace_prompt_prefix(messages, prefix);
327}
328
329fn replace_prompt_prefix(messages: &mut Vec<ModelMessage>, prefix: Vec<ModelMessage>) {
330 while matches!(
331 messages.first(),
332 Some(m) if m.role == ModelRole::System || m.role == ModelRole::Developer
333 ) {
334 messages.remove(0);
335 }
336 messages.splice(0..0, prefix);
337}
338
339async fn maintain_context_budget(ctx: &TurnContext, messages: &mut Vec<ModelMessage>) {
340 let _ = sync_messages_to_history(ctx, messages).await;
341
342 let cleared = ctx
347 .components
348 .compaction
349 .micro_compact(messages, ctx.harness_config.micro_compact_gap_minutes);
350 if cleared > 0 {
351 tracing::info!(cleared, "micro-compact applied");
352 if let Some(ref tx) = ctx.event_tx {
353 let _ = tx.send(AgentEvent::MicroCompactApplied {
354 messages_cleared: cleared,
355 });
356 }
357 }
358
359 let should_autocompact = {
360 let state = ctx.compact_state.lock().await;
361 state.should_autocompact(ctx.harness_config.autocompact_buffer_tokens)
362 };
363 if should_autocompact {
364 if let Some(ref tx) = ctx.event_tx {
365 let _ = tx.send(AgentEvent::AutoCompactStarted);
366 }
367 let provider = ctx.active_model_provider();
369 let model = ctx.active_model_name();
370 let mut state = ctx.compact_state.lock().await;
371 match ctx
372 .components
373 .compaction
374 .auto_compact(
375 &mut state,
376 messages,
377 provider.as_ref(),
378 &model,
379 &ctx.harness_config,
380 )
381 .await
382 {
383 Ok(Some(outcome)) => {
384 if let Some(ref tx) = ctx.event_tx {
385 let _ = tx.send(AgentEvent::AutoCompactCompleted {
386 tokens_saved: outcome.tokens_saved,
387 summary: outcome.summary,
388 kept_recent_messages: outcome.kept_recent_messages,
389 });
390 }
391 }
392 Ok(None) => {}
393 Err(e) => {
394 if let Some(ref tx) = ctx.event_tx {
395 let _ = tx.send(AgentEvent::AutoCompactFailed {
396 reason: e.to_string(),
397 });
398 }
399 }
400 }
401 }
402
403 let _ = evaluate_memory_triggers(ctx, messages).await;
407}
408
409fn build_model_request(ctx: &TurnContext, messages: &[ModelMessage]) -> ModelRequest {
410 let config = ctx.active_config();
411 let mut thinking = ThinkingConfig::from_config_str(&config.tui.thinking_level);
414
415 let model_name = ctx.active_model_name();
418 let provider_id = config.model.provider.clone();
419 if let Some(provider) = crate::config::resolve_provider_config(&config, &provider_id) {
420 if let Some(model) = provider
421 .models
422 .iter()
423 .find(|m| m.name == model_name || m.name.eq_ignore_ascii_case(&model_name))
424 {
425 thinking = crate::resolve_model_thinking_level(
426 thinking,
427 model.supports_thinking,
428 &model.reasoning_levels,
429 model.default_reasoning_effort.as_deref(),
430 );
431 }
432 }
433
434 ModelRequest {
435 model: model_name,
436 instructions: ctx
437 .instructions
438 .read()
439 .unwrap_or_else(|e| e.into_inner())
440 .clone(),
441 messages: rewrite_unsupported_attachments(ctx, messages),
442 thinking,
443 tools: match crate::config::effective_tool_calling_mode(&config) {
444 crate::config::ToolCallingMode::Native => {
445 let all_tools = ctx.tool_executor.definitions();
446 let mut tools = ctx
447 .components
448 .harness
449 .filter_tools(all_tools, ctx.allowed_tool_names.as_deref());
450 if tools.windows(2).any(|w| w[0].name > w[1].name) {
454 tools.sort_by(|a, b| a.name.cmp(&b.name));
455 }
456 tools
457 }
458 crate::config::ToolCallingMode::TextExtracted
459 | crate::config::ToolCallingMode::ManifestOnly
460 | crate::config::ToolCallingMode::Disabled => Vec::new(),
461 },
462 session_id: Some(ctx.session_id.clone()),
464 }
465}
466
467fn rewrite_unsupported_attachments(
468 ctx: &TurnContext,
469 messages: &[ModelMessage],
470) -> Vec<ModelMessage> {
471 let has_attachments = messages.iter().any(|m| {
474 matches!(m.role, ModelRole::User | ModelRole::Tool) && !m.content_parts.is_empty()
475 });
476 if !has_attachments {
477 return messages.to_vec();
478 }
479
480 let config = ctx.active_config();
481 let provider_id = config.model.provider.clone();
482 let model_name = config.model.name.clone();
483
484 messages
485 .iter()
486 .cloned()
487 .map(|mut message| {
488 if !matches!(message.role, ModelRole::User | ModelRole::Tool)
489 || message.content_parts.is_empty()
490 {
491 return message;
492 }
493
494 let mut rewritten = Vec::with_capacity(message.content_parts.len());
495 for part in message.content_parts {
496 let Some(kind) = part.attachment_kind() else {
497 rewritten.push(part);
498 continue;
499 };
500
501 if crate::config::model_supports_attachment(
502 &config,
503 &provider_id,
504 &model_name,
505 kind,
506 ) {
507 rewritten.push(part);
508 } else {
509 rewritten.push(ContentPart::Text {
511 text: unsupported_attachment_tool_instruction(kind, &part),
512 });
513 }
514 }
515 message.content_parts = rewritten;
516 message
517 })
518 .collect()
519}
520
521fn unsupported_attachment_tool_instruction(kind: AttachmentKind, part: &ContentPart) -> String {
522 let media_type = part.media_type().unwrap_or("application/octet-stream");
523 let byte_len = part
527 .data()
528 .map(|data| approx_decoded_attachment_bytes(data))
529 .unwrap_or(0);
530 let name = part
531 .name()
532 .map(|n| format!(" name={n:?}"))
533 .unwrap_or_default();
534 format!(
535 "[NAVI attachment unavailable to this chat model]\n\
536 kind={kind} media_type={media_type} approx_bytes={byte_len}{name}\n\
537 This model cannot view {kind} attachments directly. Tell the user to:\n\
538 1) switch to a vision-capable model (registry supports_images), or\n\
539 2) configure attachment_models.{kind} so analyze_attachment can run on a specialized model.\n\
540 Do not invent the contents of the attachment.",
541 kind = kind.as_str(),
542 )
543}
544
545fn approx_decoded_attachment_bytes(b64: &str) -> usize {
547 let trimmed = b64.trim();
548 if trimmed.is_empty() {
549 return 0;
550 }
551 let padding = trimmed.chars().rev().take_while(|c| *c == '=').count();
552 trimmed.len().saturating_mul(3) / 4 - padding.min(2)
553}
554
555fn emit_request_trace(ctx: &TurnContext, request: &ModelRequest, policy: HarnessPolicy) {
556 if let Some(ref tx) = ctx.event_tx {
557 let _ = tx.send(AgentEvent::HarnessTrace(trace_request_summary(
558 request, policy,
559 )));
560 }
561
562 tracing::info!(
563 model = %request.model,
564 messages = request.messages.len(),
565 tools = request.tools.len(),
566 "turn request started"
567 );
568}
569
570fn finalize_harness_stop(
571 ctx: &TurnContext,
572 messages: &mut Vec<ModelMessage>,
573 stop: HarnessStop,
574) -> String {
575 emit_harness_stop(ctx, &stop);
576 let text = persist_harness_stop_output(messages, &stop);
577 if let Some(ref tx) = ctx.event_tx {
578 let _ = tx.send(AgentEvent::ModelOutput {
579 text: text.clone(),
580 thinking: None,
581 });
582 }
583 text
584}
585
586fn emit_harness_stop(ctx: &TurnContext, stop: &HarnessStop) {
587 if let Some(ref tx) = ctx.event_tx {
588 let _ = tx.send(AgentEvent::HarnessStopped {
589 reason: stop.reason.as_str().to_string(),
590 message: stop.message.clone(),
591 tool_name: stop.tool_name.clone(),
592 });
593 }
594}
595
596fn persist_harness_stop_output(messages: &mut Vec<ModelMessage>, stop: &HarnessStop) -> String {
597 let mut text = format!(
598 "Stopped the run because the harness detected `{}`.\n\n{}",
599 stop.reason.as_str(),
600 stop.message
601 );
602 if let Some(tool_name) = &stop.tool_name {
603 text.push_str(&format!("\n\nLast tool: `{tool_name}`."));
604 }
605 text.push_str(
606 "\n\nTry again with a smaller instruction, or switch to a model/provider with more stable tool calling.",
607 );
608 messages.push(ModelMessage::assistant(text.clone()));
609 text
610}
611
612async fn collect_model_output(ctx: &TurnContext, request: ModelRequest) -> Result<ModelTurnOutput> {
613 let provider = ctx.active_model_provider();
614 let mut stream = provider.stream(request);
615 let mut output = ModelTurnOutput {
616 text: String::new(),
617 thinking: String::new(),
618 tool_calls: Vec::new(),
619 harness_stop: None,
620 };
621 let mut think_tags = ThinkTagSplitter::default();
622 let mut repetition_detector = crate::repetition::RepetitionDetector::default();
623
624 loop {
629 let event = tokio::select! {
630 biased;
631 _ = ctx.cancel_token.notified() => {
632 return Err(anyhow::anyhow!("turn cancelled"));
633 }
634 event = stream.next() => event,
635 };
636
637 let Some(event) = event else {
638 break;
639 };
640 ensure_not_cancelled(ctx)?;
641 match event? {
642 ModelStreamEvent::TextDelta { text } => {
643 let warning = repetition_detector.feed_text(&text);
644 emit_split_text(ctx, &mut output, think_tags.push(&text));
645 if let Some(warning) = warning {
646 output.harness_stop = Some(stop_for_repetition(ctx, warning));
647 break;
648 }
649 }
650 ModelStreamEvent::ThinkingDelta { text } => {
651 let warning = repetition_detector.feed_thinking(&text);
652 output.thinking.push_str(&text);
653 if let Some(ref tx) = ctx.event_tx {
654 let _ = tx.send(AgentEvent::ModelThinkingDelta { text });
655 }
656 if let Some(warning) = warning {
657 output.harness_stop = Some(stop_for_repetition(ctx, warning));
658 break;
659 }
660 }
661 ModelStreamEvent::ToolCall(invocation) => {
662 if invocation.tool_name.is_empty() {
663 tracing::warn!(
664 invocation_id = %invocation.id,
665 "skipping tool call with empty tool name from model"
666 );
667 continue;
668 }
669 tracing::info!(
670 tool = %invocation.tool_name,
671 invocation_id = %invocation.id,
672 "turn requested tool call"
673 );
674 if let Some(ref tx) = ctx.event_tx {
675 let _ = tx.send(AgentEvent::ToolRequested(invocation.clone()));
676 }
677 output.tool_calls.push(invocation);
678 }
679 ModelStreamEvent::ToolCallProgress {
680 id,
681 tool_name,
682 arguments_chars,
683 } => {
684 if tool_name.is_empty() && arguments_chars == 0 {
685 continue;
686 }
687 if let Some(ref tx) = ctx.event_tx {
688 let _ = tx.send(AgentEvent::ToolCallStreaming {
689 id,
690 tool_name,
691 arguments_chars,
692 });
693 }
694 }
695 ModelStreamEvent::Usage {
696 input_tokens,
697 output_tokens,
698 cache_creation_tokens,
699 cache_read_tokens,
700 } => {
701 let out_tok = output_tokens.unwrap_or(0);
702 let cache_create = cache_creation_tokens.unwrap_or(0);
703 let cache_read = cache_read_tokens.unwrap_or(0);
704 let context_in = crate::compact::context_tokens_for_meter(
708 input_tokens,
709 cache_create,
710 cache_read,
711 );
712 if let Some(ref tx) = ctx.event_tx {
713 let _ = tx.send(AgentEvent::UsageReported {
714 input_tokens: context_in.unwrap_or(input_tokens.unwrap_or(0)),
716 output_tokens: out_tok,
717 cache_creation_tokens: cache_create,
718 cache_read_tokens: cache_read,
719 });
720 }
721 if let Some(in_tok) = context_in {
722 if in_tok > 0 {
724 let mut state = ctx.compact_state.lock().await;
725 state.update_usage_full(in_tok, out_tok);
726 }
727 }
728 }
729 ModelStreamEvent::Done => {
730 emit_split_text(ctx, &mut output, think_tags.drain_pending());
731 break;
732 }
733 ModelStreamEvent::Status { label } => {
734 if label == "resuming" {
735 if let Some(ref tx) = ctx.event_tx {
736 let _ = tx.send(AgentEvent::StreamResuming {
737 accumulated_chars: output.text.len(),
738 attempt: 0,
739 });
740 }
741 } else if label == "thinking" {
742 if output.thinking.is_empty()
746 && let Some(ref tx) = ctx.event_tx
747 {
748 let _ = tx.send(AgentEvent::ModelThinkingDelta {
749 text: String::new(),
750 });
751 }
752 }
753 }
754 }
755 }
756
757 ensure_not_cancelled(ctx)?;
758 Ok(output)
759}
760
761fn emit_split_text(ctx: &TurnContext, output: &mut ModelTurnOutput, parts: Vec<SplitTextPart>) {
762 for part in parts {
763 match part {
764 SplitTextPart::Text(text) => {
765 output.text.push_str(&text);
766 if let Some(ref tx) = ctx.event_tx {
767 let _ = tx.send(AgentEvent::ModelDelta { text });
768 }
769 }
770 SplitTextPart::Thinking(text) => {
771 output.thinking.push_str(&text);
772 if let Some(ref tx) = ctx.event_tx {
773 let _ = tx.send(AgentEvent::ModelThinkingDelta { text });
774 }
775 }
776 }
777 }
778}
779
780#[derive(Debug, PartialEq, Eq)]
781enum SplitTextPart {
782 Text(String),
783 Thinking(String),
784}
785
786#[derive(Default)]
787struct ThinkTagSplitter {
788 in_think: bool,
789 pending: String,
790}
791
792impl ThinkTagSplitter {
793 fn push(&mut self, content: &str) -> Vec<SplitTextPart> {
794 let mut input = std::mem::take(&mut self.pending);
795 input.push_str(content);
796 self.split(&input, false)
797 }
798
799 fn drain_pending(&mut self) -> Vec<SplitTextPart> {
800 let pending = std::mem::take(&mut self.pending);
801 let tag = if self.in_think { "</think>" } else { "<think>" };
802 if is_partial_tag_prefix(&pending, tag) {
803 return Vec::new();
804 }
805 self.split(&pending, true)
806 }
807
808 fn split(&mut self, input: &str, final_chunk: bool) -> Vec<SplitTextPart> {
809 let mut parts = Vec::new();
810 let mut remaining = input;
811
812 while !remaining.is_empty() {
813 let tag = if self.in_think { "</think>" } else { "<think>" };
814 if let Some(pos) = find_ascii_case_insensitive(remaining, tag) {
815 self.push_segment(&mut parts, &remaining[..pos]);
816 remaining = &remaining[pos + tag.len()..];
817 self.in_think = !self.in_think;
818 continue;
819 }
820
821 let keep = if final_chunk {
822 0
823 } else {
824 partial_tag_suffix_len(remaining, tag)
825 };
826 let emit_len = remaining.len().saturating_sub(keep);
827 self.push_segment(&mut parts, &remaining[..emit_len]);
828 self.pending.push_str(&remaining[emit_len..]);
829 break;
830 }
831
832 parts
833 }
834
835 fn push_segment(&self, parts: &mut Vec<SplitTextPart>, text: &str) {
836 if text.is_empty() {
837 return;
838 }
839 if self.in_think {
840 parts.push(SplitTextPart::Thinking(text.to_string()));
841 } else {
842 parts.push(SplitTextPart::Text(text.to_string()));
843 }
844 }
845}
846
847fn find_ascii_case_insensitive(haystack: &str, needle: &str) -> Option<usize> {
848 haystack
849 .as_bytes()
850 .windows(needle.len())
851 .position(|window| window.eq_ignore_ascii_case(needle.as_bytes()))
852}
853
854fn partial_tag_suffix_len(text: &str, tag: &str) -> usize {
855 let bytes = text.as_bytes();
856 let tag_bytes = tag.as_bytes();
857 let max_len = bytes.len().min(tag_bytes.len().saturating_sub(1));
858 for len in (1..=max_len).rev() {
859 if bytes[bytes.len() - len..].eq_ignore_ascii_case(&tag_bytes[..len]) {
860 return len;
861 }
862 }
863 0
864}
865
866fn is_partial_tag_prefix(text: &str, tag: &str) -> bool {
867 !text.is_empty()
868 && text.len() < tag.len()
869 && tag.as_bytes()[..text.len()].eq_ignore_ascii_case(text.as_bytes())
870}
871
872async fn handle_tool_calls(
873 ctx: &TurnContext,
874 messages: &mut Vec<ModelMessage>,
875 run_state: &mut AgentRunState,
876 policy: HarnessPolicy,
877 mut output: ModelTurnOutput,
878) -> Option<String> {
879 let tool_call_content = std::mem::take(&mut output.text);
880 let tool_call_thinking =
881 (!output.thinking.is_empty()).then(|| std::mem::take(&mut output.thinking));
882 messages.push(ModelMessage::assistant_tool_calls_with_context(
883 output.tool_calls.clone(),
884 tool_call_content,
885 tool_call_thinking,
886 ));
887
888 let mut executable_calls = Vec::new();
889 let mut immediate_results = Vec::new();
890 for invocation in output.tool_calls {
891 ctx.components.hooks.on_tool_call(&invocation);
892 match ctx
893 .components
894 .harness
895 .record_tool_call(run_state, policy, &invocation)
896 {
897 ToolLoopDecision::Continue => executable_calls.push(invocation),
898 ToolLoopDecision::Stop(stop) => {
899 let result = tool_error_result(&invocation, &stop.message);
900 let observation =
901 ctx.components
902 .harness
903 .compact_tool_observation(&invocation, &result, policy);
904 immediate_results.push((invocation, result, observation, Vec::new()));
905 for (invocation, _result, observation, content_parts) in immediate_results {
906 messages.push(ModelMessage::tool_result_with_parts(
907 invocation.id,
908 invocation.tool_name,
909 observation,
910 content_parts,
911 ));
912 }
913 let text = finalize_harness_stop(ctx, messages, stop);
914 return Some(text);
915 }
916 }
917 }
918
919 let mut all_results = immediate_results;
920 let execution_lock = Arc::new(tokio::sync::RwLock::new(()));
921 for chunk in executable_calls.chunks(policy.max_parallel_tool_calls.max(1)) {
922 let tool_futures = chunk.iter().cloned().map(|invocation| {
923 execute_tool_call_with_parallelism(ctx, policy, invocation, execution_lock.clone())
924 });
925 all_results.extend(futures_util::future::join_all(tool_futures).await);
926 }
927
928 for (invocation, result, observation, content_parts) in all_results {
929 ctx.components.hooks.on_tool_result(&result);
930 let stop =
931 match ctx
932 .components
933 .harness
934 .record_tool_result(run_state, policy, &invocation, &result)
935 {
936 ToolLoopDecision::Continue => None,
937 ToolLoopDecision::Stop(stop) => Some(stop),
938 };
939 messages.push(ModelMessage::tool_result_with_parts(
940 invocation.id,
941 invocation.tool_name,
942 observation,
943 content_parts,
944 ));
945 if let Some(summary) = manual_context_summary(&result) {
946 if let Some(ref tx) = ctx.event_tx {
947 let _ = tx.send(AgentEvent::AutoCompactStarted);
948 }
949 let outcome = {
950 let mut state = ctx.compact_state.lock().await;
951 state.apply_manual_summary(messages, summary)
952 };
953 if let Some(ref tx) = ctx.event_tx {
954 let _ = tx.send(AgentEvent::AutoCompactCompleted {
955 tokens_saved: outcome.tokens_saved,
956 summary: outcome.summary,
957 kept_recent_messages: outcome.kept_recent_messages,
958 });
959 }
960 return Some("Context compacted.".to_string());
961 }
962 if let Some(stop) = stop {
963 let text = finalize_harness_stop(ctx, messages, stop);
964 return Some(text);
965 }
966 }
967
968 None
969}
970
971fn manual_context_summary(result: &crate::tool::ToolResult) -> Option<String> {
972 if !result.ok || result.output.get("new_context_requested")?.as_bool()? != true {
973 return None;
974 }
975 let summary = result.output.get("summary")?.as_str()?.trim();
976 (!summary.is_empty()).then(|| summary.to_string())
977}
978
979async fn execute_tool_call_with_parallelism(
980 ctx: &TurnContext,
981 policy: HarnessPolicy,
982 invocation: crate::tool::ToolInvocation,
983 execution_lock: Arc<tokio::sync::RwLock<()>>,
984) -> ToolExecutionResult {
985 match ctx.tool_executor.parallelism_for(&invocation.tool_name) {
986 ToolParallelism::Shared => {
987 let _guard = execution_lock.read().await;
988 execute_tool_call(ctx, policy, invocation).await
989 }
990 ToolParallelism::Exclusive => {
991 let _guard = execution_lock.write().await;
992 execute_tool_call(ctx, policy, invocation).await
993 }
994 }
995}
996
997pub fn tool_allowlist_deny_message(is_subagent: bool, tool_name: &str) -> String {
1002 let scope = if is_subagent {
1003 "for this subagent"
1004 } else {
1005 "for the active harness"
1006 };
1007 format!("tool `{tool_name}` is not in the allowed tool set {scope}")
1008}
1009
1010async fn execute_tool_call(
1011 ctx: &TurnContext,
1012 policy: HarnessPolicy,
1013 invocation: crate::tool::ToolInvocation,
1014) -> ToolExecutionResult {
1015 if ctx.cancellation_requested() {
1016 let result = tool_error_result(&invocation, "turn cancelled");
1017 let observation =
1018 ctx.components
1019 .harness
1020 .compact_tool_observation(&invocation, &result, policy);
1021 return (invocation, result, observation, Vec::new());
1022 }
1023
1024 if let Some(ref allowed) = ctx.allowed_tool_names {
1026 if !allowed.contains(&invocation.tool_name) {
1027 let result = tool_error_result(
1028 &invocation,
1029 tool_allowlist_deny_message(ctx.is_subagent, &invocation.tool_name),
1030 );
1031 if let Some(ref tx) = ctx.event_tx {
1032 let _ = tx.send(AgentEvent::ToolCompleted(result.clone()));
1033 }
1034 let observation =
1035 ctx.components
1036 .harness
1037 .compact_tool_observation(&invocation, &result, policy);
1038 return (invocation, result, observation, Vec::new());
1039 }
1040 }
1041
1042 if let Err(invalid) = ctx.tool_executor.validate_arguments(&invocation) {
1043 let result = ctx.tool_executor.invalid_tool_result(&invocation, invalid);
1044 if let Some(ref tx) = ctx.event_tx {
1045 let _ = tx.send(AgentEvent::ToolCompleted(result.clone()));
1046 }
1047 let observation =
1048 ctx.components
1049 .harness
1050 .compact_tool_observation(&invocation, &result, policy);
1051 return (invocation, result, observation, Vec::new());
1052 }
1053
1054 if invocation.tool_name == QUESTION_TOOL_NAME {
1055 let result = ask_user_question(ctx, &invocation).await;
1056 if let Some(ref tx) = ctx.event_tx {
1057 let _ = tx.send(AgentEvent::ToolCompleted(result.clone()));
1058 }
1059 let observation =
1060 ctx.components
1061 .harness
1062 .compact_tool_observation(&invocation, &result, policy);
1063 return (invocation, result, observation, Vec::new());
1064 }
1065
1066 let tool_ctx = crate::tool::ToolInvocationContext {
1067 event_tx: ctx.event_tx.clone(),
1068 sudo_password_resolver: Some(ctx.sudo_password_resolver.clone()),
1069 cancel_token: Some(ctx.cancel_token.clone()),
1070 };
1071
1072 let mut result = match ctx.tool_executor.validate(&invocation) {
1073 SecurityDecision::Allow => {
1074 ctx.tool_executor
1075 .invoke_with_full_context(invocation.clone(), tool_ctx, false)
1076 .await
1077 }
1078 SecurityDecision::NeedsApproval(risk) => {
1079 approve_and_invoke_tool(ctx, &invocation, risk).await
1080 }
1081 SecurityDecision::Deny(reason) => tool_error_result(&invocation, reason),
1082 };
1083
1084 if result.ok
1086 && invocation.tool_name == PLAN_TOOL_NAME
1087 && result.output.get("needs_review").and_then(|v| v.as_bool()) == Some(true)
1088 {
1089 result = wait_for_plan_review(ctx, &invocation, result).await;
1090 }
1091
1092 if ctx.cancellation_requested() {
1093 let result = tool_error_result(&invocation, "turn cancelled");
1094 let observation =
1095 ctx.components
1096 .harness
1097 .compact_tool_observation(&invocation, &result, policy);
1098 return (invocation, result, observation, Vec::new());
1099 }
1100
1101 let content_parts = take_tool_content_parts(&mut result);
1104
1105 if let Some(ref tx) = ctx.event_tx {
1106 let _ = tx.send(AgentEvent::ToolCompleted(result.clone()));
1107 }
1108
1109 let observation = ctx
1110 .components
1111 .harness
1112 .compact_tool_observation(&invocation, &result, policy);
1113 (invocation, result, observation, content_parts)
1114}
1115
1116async fn wait_for_plan_review(
1118 ctx: &TurnContext,
1119 invocation: &crate::tool::ToolInvocation,
1120 created: crate::tool::ToolResult,
1121) -> crate::tool::ToolResult {
1122 let Some(ref tx) = ctx.event_tx else {
1123 return created;
1125 };
1126
1127 let plan_id = created
1128 .output
1129 .get("plan_id")
1130 .and_then(|v| v.as_str())
1131 .unwrap_or("")
1132 .to_string();
1133 let title = created
1134 .output
1135 .get("title")
1136 .and_then(|v| v.as_str())
1137 .unwrap_or("Plan")
1138 .to_string();
1139 let description = created
1140 .output
1141 .get("description")
1142 .and_then(|v| v.as_str())
1143 .unwrap_or("")
1144 .to_string();
1145 let steps: Vec<String> = created
1146 .output
1147 .get("steps")
1148 .and_then(|v| v.as_array())
1149 .map(|arr| {
1150 arr.iter()
1151 .filter_map(|s| {
1152 s.get("description")
1153 .and_then(|d| d.as_str())
1154 .map(str::to_string)
1155 })
1156 .collect()
1157 })
1158 .unwrap_or_default();
1159
1160 let body_markdown = created
1161 .output
1162 .get("body_markdown")
1163 .and_then(|v| v.as_str())
1164 .unwrap_or("")
1165 .to_string();
1166 let plan_file_path = created
1167 .output
1168 .get("plan_file_path")
1169 .and_then(|v| v.as_str())
1170 .unwrap_or("")
1171 .to_string();
1172
1173 let request = crate::event::PlanReviewRequest {
1174 id: invocation.id.clone(),
1175 plan_id: plan_id.clone(),
1176 title,
1177 description,
1178 steps,
1179 body_markdown,
1180 plan_file_path,
1181 };
1182
1183 let answer_rx = ctx.plan_review_resolver.register(invocation.id.clone());
1184 let _ = tx.send(AgentEvent::PlanReviewRequested(request));
1185
1186 let response = tokio::select! {
1187 response = answer_rx => response.ok(),
1188 _ = ctx.cancel_token.notified() => None,
1189 };
1190
1191 match response {
1192 Some(resp) => {
1193 let _ = tx.send(AgentEvent::PlanReviewResolved(resp.clone()));
1194 let decision = match resp.decision {
1195 crate::event::PlanReviewDecision::Approve => "approve",
1196 crate::event::PlanReviewDecision::RequestChanges => "request_changes",
1197 crate::event::PlanReviewDecision::Quit => "quit",
1198 };
1199 let comments_json: Vec<Value> = resp
1200 .comments
1201 .iter()
1202 .map(|c| {
1203 json!({
1204 "start_line": c.start_line,
1205 "end_line": c.end_line,
1206 "text": c.text,
1207 })
1208 })
1209 .collect();
1210 let ok = !matches!(resp.decision, crate::event::PlanReviewDecision::Quit);
1211 let mut output = created.output;
1212 if let Some(obj) = output.as_object_mut() {
1213 obj.insert("decision".into(), json!(decision));
1214 obj.insert("comments".into(), json!(comments_json));
1215 obj.insert("freeform".into(), json!(resp.freeform));
1216 if matches!(resp.decision, crate::event::PlanReviewDecision::Approve) {
1218 if let Some(path) = obj
1219 .get("plan_file_path")
1220 .and_then(|v| v.as_str())
1221 .map(std::path::PathBuf::from)
1222 {
1223 if let Some(md) = crate::plan_store::read_plan_file(&path) {
1224 obj.insert("body_markdown".into(), json!(md));
1225 }
1226 }
1227 }
1228 obj.insert(
1229 "message".into(),
1230 json!(match resp.decision {
1231 crate::event::PlanReviewDecision::Approve =>
1232 "User approved the plan. You are now in normal mode — implement the plan. Use the body_markdown / plan file as the source of truth.",
1233 crate::event::PlanReviewDecision::RequestChanges =>
1234 "User requested changes to the plan. Revise the plan file based on comments, then submit again.",
1235 crate::event::PlanReviewDecision::Quit =>
1236 "User abandoned the plan. Do not implement it.",
1237 }),
1238 );
1239 obj.insert("needs_review".into(), json!(false));
1241 obj.insert("review_complete".into(), json!(true));
1242 }
1243 crate::tool::ToolResult {
1244 invocation_id: invocation.id.clone(),
1245 ok,
1246 output,
1247 }
1248 }
1249 None => {
1250 let mut output = created.output;
1251 if let Some(obj) = output.as_object_mut() {
1252 obj.insert("decision".into(), json!("cancelled"));
1253 obj.insert("needs_review".into(), json!(false));
1254 obj.insert(
1255 "message".into(),
1256 json!("Plan review cancelled (turn cancelled)."),
1257 );
1258 }
1259 crate::tool::ToolResult {
1260 invocation_id: invocation.id.clone(),
1261 ok: false,
1262 output,
1263 }
1264 }
1265 }
1266}
1267
1268async fn ask_user_question(
1269 ctx: &TurnContext,
1270 invocation: &crate::tool::ToolInvocation,
1271) -> crate::tool::ToolResult {
1272 let Some(ref tx) = ctx.event_tx else {
1273 return tool_error_result(invocation, "question requires an interactive client");
1274 };
1275
1276 let request = match question_request_from_invocation(invocation) {
1277 Ok(request) => request,
1278 Err(message) => return tool_error_result(invocation, message),
1279 };
1280
1281 let answer_rx = ctx.question_resolver.register(invocation.id.clone());
1282 let _ = tx.send(AgentEvent::QuestionRequested(request));
1283
1284 let response = tokio::select! {
1285 response = answer_rx => response.ok(),
1286 _ = ctx.cancel_token.notified() => None,
1287 };
1288
1289 match response {
1290 Some(crate::event::QuestionResponse::Answered { id, answers }) => {
1291 let response = crate::event::QuestionResponse::Answered {
1292 id,
1293 answers: answers.clone(),
1294 };
1295 let _ = tx.send(AgentEvent::QuestionResolved(response));
1296 crate::tool::ToolResult {
1297 invocation_id: invocation.id.clone(),
1298 ok: true,
1299 output: json!({
1300 "schema_version": 1,
1301 "answers": answers,
1302 "answer": answers.join("\n"),
1303 }),
1304 }
1305 }
1306 Some(response @ crate::event::QuestionResponse::Dismissed { .. }) => {
1307 let _ = tx.send(AgentEvent::QuestionResolved(response));
1308 tool_error_result(invocation, "user dismissed question")
1309 }
1310 None => {
1311 let response = crate::event::QuestionResponse::Dismissed {
1312 id: invocation.id.clone(),
1313 };
1314 let _ = tx.send(AgentEvent::QuestionResolved(response));
1315 tool_error_result(invocation, "turn cancelled")
1316 }
1317 }
1318}
1319
1320fn question_request_from_invocation(
1321 invocation: &crate::tool::ToolInvocation,
1322) -> std::result::Result<crate::event::QuestionRequest, String> {
1323 let question = invocation
1324 .input
1325 .get("question")
1326 .and_then(Value::as_str)
1327 .filter(|value| !value.trim().is_empty())
1328 .ok_or_else(|| "question must include a non-empty `question` string".to_string())?
1329 .trim()
1330 .to_string();
1331 let options_value = invocation
1332 .input
1333 .get("options")
1334 .and_then(Value::as_array)
1335 .ok_or_else(|| "question must include an `options` array".to_string())?;
1336 let mut options = Vec::new();
1337 for option in options_value {
1338 if let Some(label) = option.as_str() {
1339 options.push(crate::event::QuestionOption {
1340 label: label.to_string(),
1341 description: None,
1342 });
1343 continue;
1344 }
1345 let Some(object) = option.as_object() else {
1346 return Err("question options must be strings or objects".to_string());
1347 };
1348 let label = object
1349 .get("label")
1350 .and_then(Value::as_str)
1351 .filter(|value| !value.trim().is_empty())
1352 .ok_or_else(|| "question option objects need a non-empty `label`".to_string())?;
1353 let description = object
1354 .get("description")
1355 .and_then(Value::as_str)
1356 .filter(|value| !value.trim().is_empty())
1357 .map(str::to_string);
1358 options.push(crate::event::QuestionOption {
1359 label: label.to_string(),
1360 description,
1361 });
1362 }
1363 if options.is_empty() {
1364 return Err("question must include at least one option".to_string());
1365 }
1366 Ok(crate::event::QuestionRequest {
1367 id: invocation.id.clone(),
1368 question,
1369 options,
1370 multiple: invocation
1371 .input
1372 .get("multiple")
1373 .and_then(Value::as_bool)
1374 .unwrap_or(false),
1375 allow_custom: invocation
1376 .input
1377 .get("custom")
1378 .or_else(|| invocation.input.get("allow_custom"))
1379 .and_then(Value::as_bool)
1380 .unwrap_or(false),
1381 })
1382}
1383
1384async fn approve_and_invoke_tool(
1385 ctx: &TurnContext,
1386 invocation: &crate::tool::ToolInvocation,
1387 risk: crate::security::SecurityRisk,
1388) -> crate::tool::ToolResult {
1389 let Some(ref tx) = ctx.event_tx else {
1390 return tool_error_result(
1391 invocation,
1392 "approval required in headless mode; rerun in TUI",
1393 );
1394 };
1395
1396 let approval_risk = match risk {
1397 crate::security::SecurityRisk::Tool => crate::event::ApprovalRisk::Tool,
1398 crate::security::SecurityRisk::Write => crate::event::ApprovalRisk::Write,
1399 crate::security::SecurityRisk::Command => crate::event::ApprovalRisk::Command,
1400 crate::security::SecurityRisk::GuardedCommand => crate::event::ApprovalRisk::Guarded,
1401 crate::security::SecurityRisk::ExternalPlugin => crate::event::ApprovalRisk::ExternalPlugin,
1402 };
1403 let approve_rx = ctx.approval_resolver.register(invocation.id.clone());
1404
1405 let _ = tx.send(AgentEvent::ApprovalRequested(
1406 crate::event::ApprovalRequest {
1407 id: invocation.id.clone(),
1408 summary: format!("Run tool `{}`", invocation.tool_name),
1409 risk: approval_risk,
1410 },
1411 ));
1412
1413 let approved = tokio::select! {
1414 decision = approve_rx => decision.ok(),
1415 _ = ctx.cancel_token.notified() => None,
1416 };
1417
1418 match approved {
1419 Some(decision) => {
1420 let is_approved = matches!(decision, crate::event::ApprovalDecision::Approved { .. });
1421 if is_approved {
1422 ctx.tool_executor
1423 .invoke_with_full_context(
1424 invocation.clone(),
1425 crate::tool::ToolInvocationContext {
1426 event_tx: ctx.event_tx.clone(),
1427 sudo_password_resolver: Some(ctx.sudo_password_resolver.clone()),
1428 cancel_token: Some(ctx.cancel_token.clone()),
1429 },
1430 true,
1431 )
1432 .await
1433 } else {
1434 tool_error_result(invocation, "user denied tool execution")
1435 }
1436 }
1437 None => {
1438 let _ = tx.send(AgentEvent::ApprovalResolved(
1439 crate::event::ApprovalDecision::Denied {
1440 id: invocation.id.clone(),
1441 },
1442 ));
1443 tool_error_result(invocation, "turn cancelled")
1444 }
1445 }
1446}
1447
1448fn persist_final_model_output(
1449 ctx: &TurnContext,
1450 messages: &mut Vec<ModelMessage>,
1451 output: &ModelTurnOutput,
1452) {
1453 if let Some(ref tx) = ctx.event_tx {
1454 let _ = tx.send(AgentEvent::ModelOutput {
1455 text: output.text.clone(),
1456 thinking: (!output.thinking.is_empty()).then(|| output.thinking.clone()),
1457 });
1458 }
1459
1460 if !output.text.trim().is_empty() || !output.thinking.is_empty() {
1461 messages.push(ModelMessage::assistant_with_thinking(
1462 output.text.clone(),
1463 (!output.thinking.is_empty()).then(|| output.thinking.clone()),
1464 ));
1465 }
1466}
1467
1468fn stop_for_repetition(
1469 ctx: &TurnContext,
1470 warning: crate::repetition::RepetitionWarning,
1471) -> HarnessStop {
1472 if let Some(ref tx) = ctx.event_tx {
1473 let _ = tx.send(AgentEvent::RepetitionDetected {
1474 kind: map_repetition_kind(&warning.kind),
1475 message: warning.message.clone(),
1476 });
1477 }
1478 HarnessStop {
1479 reason: HarnessStopReason::DegenerateModelOutput,
1480 message: warning.message,
1481 tool_name: None,
1482 }
1483}
1484
1485fn map_repetition_kind(kind: &crate::repetition::RepetitionKind) -> RepetitionWarningKind {
1486 match kind {
1487 crate::repetition::RepetitionKind::CharRun { ch, count } => {
1488 RepetitionWarningKind::CharRun {
1489 ch: *ch,
1490 count: *count,
1491 }
1492 }
1493 crate::repetition::RepetitionKind::AlternatingPattern { pattern, cycles } => {
1494 RepetitionWarningKind::AlternatingPattern {
1495 pattern: pattern.clone(),
1496 cycles: *cycles,
1497 }
1498 }
1499 }
1500}
1501
1502async fn combined_memory_injection(ctx: &TurnContext) -> Option<String> {
1504 let rebuild_context = {
1505 let state = ctx.compact_state.lock().await;
1506 state.rebuild_context.clone()
1507 };
1508
1509 let auto_memory_index = load_auto_memory_index(ctx);
1510
1511 let parts: Vec<String> = Vec::new();
1512 let mut parts = parts;
1513
1514 if let Some(ref idx) = auto_memory_index {
1515 if !idx.trim().is_empty() {
1516 parts.push(format!("=== AUTO-MEMORY INDEX ===\n{}", idx));
1517 }
1518 }
1519
1520 match (ctx.memory_injection.clone(), rebuild_context) {
1521 (Some(session_memory), Some(rebuild_context)) => {
1522 parts.push(session_memory);
1523 parts.push(format!("Rebuilt session context:\n\n{rebuild_context}"));
1524 }
1525 (Some(session_memory), None) => {
1526 parts.push(session_memory);
1527 }
1528 (None, Some(rebuild_context)) => {
1529 parts.push(format!("Rebuilt session context:\n\n{rebuild_context}"));
1530 }
1531 (None, None) => {}
1532 }
1533
1534 if parts.is_empty() {
1535 None
1536 } else {
1537 Some(parts.join("\n\n"))
1538 }
1539}
1540
1541fn load_auto_memory_index(ctx: &TurnContext) -> Option<String> {
1543 let manager = ctx.get_or_init_memory_manager().ok().flatten()?;
1544 let store = manager.auto_memory.clone();
1545 let index = store.build_prompt_context(2000);
1546 if index.trim().is_empty() {
1547 None
1548 } else {
1549 Some(index)
1550 }
1551}
1552
1553pub async fn sync_messages_to_history(ctx: &TurnContext, messages: &[ModelMessage]) -> Result<()> {
1554 let Some(manager) = ctx.get_or_init_memory_manager()? else {
1555 return Ok(());
1556 };
1557 manager
1558 .history
1559 .record_session_start(&ctx.session_id, &ctx.project_dir.to_string_lossy())?;
1560
1561 let pending: Vec<(u64, &ModelMessage)> = {
1562 let state = ctx.compact_state.lock().await;
1563 messages
1564 .iter()
1565 .filter_map(|msg| {
1566 let key = history_message_key(msg);
1567 (!state.history_synced_message_keys.contains(&key)).then_some((key, msg))
1568 })
1569 .collect()
1570 };
1571
1572 let mut recorded_keys = Vec::new();
1573 for (key, msg) in pending {
1574 let role_str = match msg.role {
1575 crate::model::ModelRole::User => "user",
1576 crate::model::ModelRole::Assistant => "assistant",
1577 crate::model::ModelRole::Tool => "tool",
1578 crate::model::ModelRole::System => "system",
1579 crate::model::ModelRole::Developer => "developer",
1580 };
1581
1582 let tool_name = msg.tool_name.clone();
1583 let tool_input: Option<String> = None;
1584 let mut tool_output = None;
1585
1586 if msg.role == crate::model::ModelRole::Tool {
1587 tool_output = Some(msg.content.clone());
1588 }
1589
1590 manager.history.record_event(
1591 &ctx.session_id,
1592 "message",
1593 Some(role_str),
1594 Some(&msg.content),
1595 tool_name.as_deref(),
1596 tool_input.as_deref(),
1597 tool_output.as_deref(),
1598 None,
1599 None,
1600 )?;
1601 recorded_keys.push(key);
1602 }
1603
1604 if !recorded_keys.is_empty() {
1605 let mut state = ctx.compact_state.lock().await;
1606 state.history_synced_message_keys.extend(recorded_keys);
1607 }
1608 Ok(())
1609}
1610
1611fn history_message_key(msg: &ModelMessage) -> u64 {
1612 let mut hasher = std::collections::hash_map::DefaultHasher::new();
1613 std::mem::discriminant(&msg.role).hash(&mut hasher);
1614 msg.content.hash(&mut hasher);
1615 msg.tool_call_id.hash(&mut hasher);
1616 msg.tool_name.hash(&mut hasher);
1617 msg.created_at.hash(&mut hasher);
1618 serde_json::to_string(&msg.content_parts)
1619 .unwrap_or_default()
1620 .hash(&mut hasher);
1621 serde_json::to_string(&msg.tool_calls)
1622 .unwrap_or_default()
1623 .hash(&mut hasher);
1624 msg.thinking_content.hash(&mut hasher);
1625 hasher.finish()
1626}
1627
1628pub(crate) async fn evaluate_memory_triggers(
1630 ctx: &TurnContext,
1631 messages: &mut Vec<ModelMessage>,
1632) -> Result<bool> {
1633 let memory_config = ctx.active_config().memory;
1634 let Some(manager) = ctx.get_or_init_memory_manager()? else {
1635 return Ok(false);
1636 };
1637
1638 let (percentage, _total_tokens) = {
1639 let state = ctx.compact_state.lock().await;
1640 (
1641 state.context_percentage(0) as f64 / 100.0,
1642 state.total_estimated_tokens(0),
1643 )
1644 };
1645
1646 let mut thresholds_to_trigger = Vec::new();
1648 {
1649 let state = ctx.compact_state.lock().await;
1650 for &t in &memory_config.checkpoint_thresholds {
1651 if percentage >= t && !state.crossed_thresholds.contains(&t) {
1652 thresholds_to_trigger.push(t);
1653 }
1654 }
1655 }
1656
1657 if !thresholds_to_trigger.is_empty() {
1658 manager
1659 .history
1660 .record_session_start(&ctx.session_id, &ctx.project_dir.to_string_lossy())?;
1661
1662 let provider = ctx.active_model_provider();
1663 let model_name = ctx.active_model_name();
1664
1665 crate::memory::run_checkpoint_writer(
1666 &ctx.session_id,
1667 messages,
1668 &manager.auto_memory,
1669 provider.as_ref(),
1670 &model_name,
1671 )
1672 .await?;
1673
1674 {
1676 let mut state = ctx.compact_state.lock().await;
1677 for t in thresholds_to_trigger {
1678 state.crossed_thresholds.push(t);
1679 let cp_path = manager.auto_memory.db_path.to_string_lossy().to_string();
1680 manager.history.record_checkpoint(
1681 &ctx.session_id,
1682 state.crossed_thresholds.len() as i64,
1683 percentage,
1684 &cp_path,
1685 )?;
1686 }
1687 }
1688 }
1689
1690 if percentage >= memory_config.rebuild_threshold {
1695 tracing::info!(
1696 "Rebuild threshold reached ({}% >= {}%) — applying long-horizon context rebuild fallback",
1697 percentage * 100.0,
1698 memory_config.rebuild_threshold * 100.0
1699 );
1700
1701 let context_window = {
1702 let state = ctx.compact_state.lock().await;
1703 state.context_window
1704 };
1705
1706 let boot_context = crate::memory::build_rebuild_context(
1708 messages,
1709 &manager.auto_memory,
1710 &manager.global_memory,
1711 context_window,
1712 memory_config.injected_context_token_budget,
1713 );
1714
1715 let cycle_num = {
1717 let mut state = ctx.compact_state.lock().await;
1718 state.crossed_thresholds.clear(); state.rebuild_context = Some(boot_context.clone());
1720 1 };
1722 manager
1723 .history
1724 .record_rebuild(&ctx.session_id, cycle_num, cycle_num + 1, &boot_context)?;
1725
1726 messages.clear();
1728 ensure_system_prompt(ctx, messages).await;
1729
1730 {
1732 let mut state = ctx.compact_state.lock().await;
1733 state.last_input_tokens = None;
1734 state.clear_unsent_bytes();
1735 }
1736
1737 if let Some(ref tx) = ctx.event_tx {
1740 let _ = tx.send(AgentEvent::AutoCompactCompleted {
1741 tokens_saved: 1,
1742 summary: format!(
1743 "Context was near the model limit. Session history was rebuilt from long-horizon memory.\n\n{}",
1744 boot_context
1745 ),
1746 kept_recent_messages: 0,
1747 });
1748 }
1749
1750 return Ok(true);
1751 }
1752
1753 Ok(false)
1754}
1755
1756#[cfg(test)]
1757mod tests;