Skip to main content

navi_core/turn/
mod.rs

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    /// Skill pools (folders) for the top-level catalog.
59    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    /// Stable base instructions for the `instructions` field of the provider
65    /// request, separated from developer messages for prompt cache efficiency.
66    /// It is populated once, when the session prompt prefix is frozen.
67    pub instructions: Arc<RwLock<Option<String>>>,
68    /// Immutable prompt prefix for this session. Context packets, skills and
69    /// memory must not rewrite the provider-visible system/developer prefix in
70    /// the middle of a conversation: doing so fragments prompt caches and can
71    /// subtly change the model's operating instructions.
72    pub prompt_prefix: Arc<std::sync::Mutex<Option<Vec<ModelMessage>>>>,
73    /// Snapshot of the active `NaviConfig` taken at turn start. Used by
74    /// `ensure_system_prompt` so the model sees the user-configured harness
75    /// profile, model and provider rather than the defaults.
76    pub config: Arc<RwLock<crate::config::NaviConfig>>,
77    /// Optional previous-session memory loaded at session startup.
78    pub memory_injection: Option<String>,
79    /// Optional separate provider for compaction/summarization. When set,
80    /// auto_compact uses this instead of the main model provider.
81    pub compaction_provider: Option<Arc<dyn ModelProvider>>,
82    /// Current agent mode (Default or Plan). In Plan mode, only read-only
83    /// tools are available to the model.
84    pub agent_mode: crate::plan_mode::AgentMode,
85    /// Model name for the compaction provider.
86    pub compaction_model_name: Option<String>,
87    pub session_id: String,
88    /// Optional set of tool names this turn may call (subagent or active harness).
89    /// `None` means unrestricted (beyond security policy / agent mode).
90    pub allowed_tool_names: Option<Vec<String>>,
91    /// When true, allowlist denials use subagent-specific wording. Root/harness
92    /// turns set this false so denials do not say "for this subagent".
93    pub is_subagent: bool,
94    /// Session-scoped memory manager. Lazily opened once and reused for history
95    /// sync, auto-memory index, checkpoints, and rebuilds (avoids reopening
96    /// three SQLite DBs every tool-loop iteration).
97    pub memory_manager: Arc<std::sync::Mutex<Option<Arc<crate::memory::MemoryManager>>>>,
98    /// Optional harness pack card injected into developer context.
99    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    /// Returns a shared [`MemoryManager`] for this session when memory is enabled.
125    /// Opens the underlying SQLite stores at most once per slot.
126    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    // In Plan mode, allow explore + plan-file writes + plan/question tools.
247    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    // Store the stable base instructions on the turn context so
276    // `build_model_request` can place them in the provider's
277    // `instructions` field.
278    *ctx.instructions.write().unwrap_or_else(|e| e.into_inner()) =
279        Some(rendered.instructions.clone());
280
281    // Build the prefix once. It deliberately remains unchanged for the rest
282    // of the session so provider prompt-cache keys stay stable after the first
283    // request, even if memory, skill or external-context state changes later.
284    let mut prefix = Vec::with_capacity(2 + rendered.developer_messages.len());
285    prefix.push(ModelMessage::system(rendered.instructions));
286    // In Plan mode, inject instructions for markdown plan-file workflow.
287    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    // Micro-compact and auto-compact run *before* long-horizon memory rebuild.
343    // Previously rebuild short-circuited this path at ~85% usage, emitted a hard
344    // Error event ("physical context rebuild"), and skipped model summarization —
345    // so the turn never recovered and the compact text never appeared in chat.
346    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        // Always use the session's own model — not a background/subagent provider.
368        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    // Checkpoints + rebuild fallback. After a successful auto-compact,
404    // token usage is reset so rebuild will not fire. Rebuild only acts when
405    // context is still critically full (e.g. auto-compact failed).
406    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    // Fixed effort from config/session preference — never re-scored mid-turn so
412    // provider prefix/KV cache stays stable across tool-loop iterations.
413    let mut thinking = ThinkingConfig::from_config_str(&config.tui.thinking_level);
414
415    // Clamp to registry reasoning_levels for the active model when available.
416    // Models without reasoning support are forced to Off automatically.
417    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                // definitions() already returns name-sorted tools; re-sort only
451                // when a filter may have disordered a filtered subset (no-op
452                // when already sorted — keeps prefix-cache order stable).
453                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        // Pin provider KV-cache affinity to this agent session (Charm Hyper, etc.).
463        session_id: Some(ctx.session_id.clone()),
464    }
465}
466
467fn rewrite_unsupported_attachments(
468    ctx: &TurnContext,
469    messages: &[ModelMessage],
470) -> Vec<ModelMessage> {
471    // Fast path: no multimodal attachments on user or tool messages.
472    // ModelRequest still needs an owned list, so we clone once without mapping.
473    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                    // Never dump base64 into the prompt for non-vision models.
510                    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    // Never inline base64 attachment bytes into the prompt. Free/small models
524    // hang or rate-limit when the rewritten text carries multi-MB payloads, and
525    // models cannot reliably re-emit base64 into analyze_attachment anyway.
526    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
545/// Base64 payload size estimate (decoded). Avoids allocating a decoded buffer.
546fn 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    // Race the provider stream against cancel. Checking only *after* each
625    // `stream.next()` leaves the session loop parked on a hung/slow HTTP body
626    // after Esc-cancel; the next user turn then queues forever and the TUI
627    // shows "Waiting for model" until process restart.
628    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                // Context meter must count cached prompt tokens. Some aggregators
705                // report only non-cached prompt (e.g. 430) with a large cache_read
706                // (e.g. 63k) — without summing, the UI shows a bogus ~430 / 1M.
707                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                        // Prefer full context size for session accounting.
715                        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                    // Never clobber a real prior reading with a zero/empty partial.
723                    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                    // Provider signalled reasoning without a text delta (e.g.
743                    // Anthropic signature_delta). Mark the stream as live so the
744                    // UI does not escalate to "Still waiting for model".
745                    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
997/// Deny message when `allowed_tool_names` rejects a tool.
998///
999/// Root/harness turns must not say "subagent" — that misled the TUI when catalog
1000/// skills incorrectly locked the session.
1001pub 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    // Optional tool allowlist (nested subagent or soft harness on the root turn).
1025    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    // Plan create: block the turn until the user finishes the review modal.
1085    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    // Lift multimodal parts (e.g. view_image) before ToolCompleted/observation
1102    // so base64 never enters the text transcript or event payload.
1103    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
1116/// Block after `plan(create)` until the TUI resolves the review modal.
1117async 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        // Headless: return create result without blocking.
1124        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                // On approve, re-surface the markdown so the model has the plan after mode exit.
1217                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                // Still true that review finished; model should not open another modal.
1240                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
1502/// Synchronizes new messages in the session conversation history to SQLite.
1503async 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
1541/// Loads the auto-memory index for system prompt injection.
1542fn 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
1628/// Evaluates memory system checkpoint and rebuild thresholds based on context utilization.
1629pub(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    // 1. Checkpoint thresholds
1647    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        // Mark thresholds as crossed
1675        {
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    // 2. Rebuild threshold — last-resort fallback when auto-compact did not
1691    // reclaim enough budget (or the circuit breaker is open). Never emit a hard
1692    // Error: rebuild is recovery, not a turn failure, and the agent loop should
1693    // continue the active task with the rebuilt prompt.
1694    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        // Rebuild context!
1707        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        // Record rebuild in SQLite
1716        let cycle_num = {
1717            let mut state = ctx.compact_state.lock().await;
1718            state.crossed_thresholds.clear(); // reset thresholds for the new cycle!
1719            state.rebuild_context = Some(boot_context.clone());
1720            1 // Default cycle sequence number
1721        };
1722        manager
1723            .history
1724            .record_rebuild(&ctx.session_id, cycle_num, cycle_num + 1, &boot_context)?;
1725
1726        // Re-assemble conversation messages
1727        messages.clear();
1728        ensure_system_prompt(ctx, messages).await;
1729
1730        // Reset compaction state token usage
1731        {
1732            let mut state = ctx.compact_state.lock().await;
1733            state.last_input_tokens = None;
1734            state.clear_unsent_bytes();
1735        }
1736
1737        // Surface as a compact-style recovery notice so the TUI/chat show
1738        // continuity instead of a red error that looks like a hard failure.
1739        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;