Skip to main content

navi_core/runtime/
mod.rs

1mod event_bus;
2mod session_state;
3
4#[cfg(test)]
5mod tests;
6
7use std::collections::HashMap;
8use std::path::PathBuf;
9use std::sync::{Arc, RwLock};
10use tokio::sync::{broadcast, mpsc, oneshot};
11
12use crate::cancel::CancelToken;
13use crate::config::{LoadedConfig, PermissionMode, SecurityConfig};
14use crate::context::ContextPacket;
15use crate::event::{
16    AgentEvent, ApprovalDecision, PlanReviewResponse, QuestionResponse, RuntimeEvent,
17    RuntimeEventKind, SudoPasswordResponse,
18};
19use crate::goal::{
20    CreateGoalTool, GetGoalTool, GoalExtension, GoalRuntimeHandle, GoalService, UpdateGoalTool,
21};
22use crate::harness::select_harness_policy;
23use crate::model::{ModelMessage, ModelProvider, ModelResponse};
24use crate::runtime_components::RuntimeComponents;
25use crate::security::SecurityPolicy;
26use crate::session::{SessionId, SessionStore, current_unix_timestamp};
27use crate::session_title::SessionTitleHandle;
28use crate::skills::{
29    SkillManifest, SkillPool, active_skills, discover_catalog_entries, discover_configured_skills,
30};
31use crate::tool::builtin::{RepoExploreTool, SubagentTool};
32use crate::tool::{Tool, ToolExecutor};
33use crate::trace::{TraceStore, turn_traces_from_events};
34use crate::{
35    ModelOption, SessionSnapshot, available_model_options, canonical_provider_id,
36    provider_request_model_name,
37};
38use anyhow::Result;
39
40pub use event_bus::EventBus;
41pub use session_state::SessionState;
42
43type PendingApprovals = Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<ApprovalDecision>>>>;
44type PendingQuestions = Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<QuestionResponse>>>>;
45type PendingPlanReviews =
46    Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<PlanReviewResponse>>>>;
47type PendingSudoPasswords =
48    Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<SudoPasswordResponse>>>>;
49
50/// Resolves pending tool approvals by matching decision ids to waiting
51/// receivers. Cloneable so it can be handed to the UI layer.
52#[derive(Clone)]
53pub struct ApprovalResolver {
54    pending_approvals: PendingApprovals,
55    runtime_events_tx: broadcast::Sender<RuntimeEvent>,
56}
57
58/// Resolves pending interactive questions by matching response ids to waiting
59/// receivers. Cloneable so it can be handed to the UI layer.
60#[derive(Clone)]
61pub struct QuestionResolver {
62    pending_questions: PendingQuestions,
63    runtime_events_tx: broadcast::Sender<RuntimeEvent>,
64}
65
66/// Resolves pending plan reviews (blocks `plan` create until the user finishes).
67#[derive(Clone)]
68pub struct PlanReviewResolver {
69    pending_reviews: PendingPlanReviews,
70    runtime_events_tx: broadcast::Sender<RuntimeEvent>,
71}
72
73/// Resolves sudo password prompts without putting secrets on the event bus.
74#[derive(Clone)]
75pub struct SudoPasswordResolver {
76    pending: PendingSudoPasswords,
77}
78
79impl SudoPasswordResolver {
80    #[cfg(test)]
81    pub fn new_for_test() -> Self {
82        Self {
83            pending: Arc::new(std::sync::Mutex::new(HashMap::new())),
84        }
85    }
86
87    pub(crate) fn new_standalone() -> Self {
88        Self {
89            pending: Arc::new(std::sync::Mutex::new(HashMap::new())),
90        }
91    }
92
93    pub fn register(&self, id: String) -> oneshot::Receiver<SudoPasswordResponse> {
94        let (tx, rx) = oneshot::channel();
95        self.pending
96            .lock()
97            .unwrap_or_else(|e| e.into_inner())
98            .insert(id, tx);
99        rx
100    }
101
102    /// Deliver the user's response. The password (if any) never goes on the runtime event bus.
103    pub fn resolve(&self, response: SudoPasswordResponse) -> bool {
104        let id = response.id().to_string();
105        if let Some(tx) = self
106            .pending
107            .lock()
108            .unwrap_or_else(|e| e.into_inner())
109            .remove(&id)
110        {
111            let _ = tx.send(response);
112            true
113        } else {
114            false
115        }
116    }
117}
118
119impl PlanReviewResolver {
120    #[cfg(test)]
121    pub fn new_for_test() -> Self {
122        let (tx, _) = broadcast::channel(16);
123        Self {
124            pending_reviews: Arc::new(std::sync::Mutex::new(HashMap::new())),
125            runtime_events_tx: tx,
126        }
127    }
128
129    pub(crate) fn new_standalone() -> Self {
130        let (tx, _) = broadcast::channel(16);
131        Self {
132            pending_reviews: Arc::new(std::sync::Mutex::new(HashMap::new())),
133            runtime_events_tx: tx,
134        }
135    }
136
137    pub fn register(&self, id: String) -> oneshot::Receiver<PlanReviewResponse> {
138        let (tx, rx) = oneshot::channel();
139        self.pending_reviews
140            .lock()
141            .unwrap_or_else(|e| e.into_inner())
142            .insert(id, tx);
143        rx
144    }
145
146    pub fn resolve(&self, response: PlanReviewResponse) -> bool {
147        let id = response.id.clone();
148        if let Some(tx) = self
149            .pending_reviews
150            .lock()
151            .unwrap_or_else(|e| e.into_inner())
152            .remove(&id)
153        {
154            let _ = tx.send(response.clone());
155            let _ = self.runtime_events_tx.send(RuntimeEvent::new(
156                RuntimeEventKind::PlanReviewResolved(response),
157            ));
158            true
159        } else {
160            false
161        }
162    }
163}
164
165impl QuestionResolver {
166    #[cfg(test)]
167    pub fn new_for_test() -> Self {
168        let (tx, _) = broadcast::channel(16);
169        Self {
170            pending_questions: Arc::new(std::sync::Mutex::new(HashMap::new())),
171            runtime_events_tx: tx,
172        }
173    }
174
175    pub(crate) fn new_standalone() -> Self {
176        let (tx, _) = broadcast::channel(16);
177        Self {
178            pending_questions: Arc::new(std::sync::Mutex::new(HashMap::new())),
179            runtime_events_tx: tx,
180        }
181    }
182
183    /// Register a pending question, returning the receiver for the response.
184    pub fn register(&self, id: String) -> oneshot::Receiver<QuestionResponse> {
185        let (tx, rx) = oneshot::channel();
186        self.pending_questions
187            .lock()
188            .unwrap_or_else(|e| e.into_inner())
189            .insert(id, tx);
190        rx
191    }
192
193    /// Resolves a pending question by id. Returns `true` if a matching request
194    /// was found and resolved.
195    pub fn resolve(&self, response: QuestionResponse) -> bool {
196        let id = response.id().to_string();
197        if let Some(tx) = self
198            .pending_questions
199            .lock()
200            .unwrap_or_else(|e| e.into_inner())
201            .remove(&id)
202        {
203            let _ = tx.send(response.clone());
204            let _ =
205                self.runtime_events_tx
206                    .send(RuntimeEvent::new(RuntimeEventKind::QuestionResolved(
207                        response,
208                    )));
209            true
210        } else {
211            false
212        }
213    }
214}
215
216impl ApprovalResolver {
217    #[cfg(test)]
218    pub fn new_for_test() -> Self {
219        let (tx, _) = broadcast::channel(16);
220        Self {
221            pending_approvals: Arc::new(std::sync::Mutex::new(HashMap::new())),
222            runtime_events_tx: tx,
223        }
224    }
225
226    pub(crate) fn new_standalone() -> Self {
227        let (tx, _) = broadcast::channel(16);
228        Self {
229            pending_approvals: Arc::new(std::sync::Mutex::new(HashMap::new())),
230            runtime_events_tx: tx,
231        }
232    }
233
234    /// Register a pending approval, returning the receiver for the decision.
235    pub fn register(&self, id: String) -> oneshot::Receiver<ApprovalDecision> {
236        let (tx, rx) = oneshot::channel();
237        self.pending_approvals
238            .lock()
239            .unwrap_or_else(|e| e.into_inner())
240            .insert(id, tx);
241        rx
242    }
243
244    /// Resolves a pending approval by id. Returns `true` if a matching
245    /// pending request was found and resolved.
246    pub fn resolve(&self, decision: ApprovalDecision) -> bool {
247        let id = match &decision {
248            ApprovalDecision::Approved { id } => id,
249            ApprovalDecision::Denied { id } => id,
250        };
251        if let Some(tx) = self
252            .pending_approvals
253            .lock()
254            .unwrap_or_else(|e| e.into_inner())
255            .remove(id)
256        {
257            let _ = tx.send(decision.clone());
258            let _ =
259                self.runtime_events_tx
260                    .send(RuntimeEvent::new(RuntimeEventKind::ApprovalResolved(
261                        decision,
262                    )));
263            true
264        } else {
265            false
266        }
267    }
268}
269
270/// A lightweight handle that cancels the current turn when dropped or called.
271/// Cloneable so it can be handed to the UI layer.
272#[derive(Clone)]
273pub struct TurnCanceller {
274    inner: CancelToken,
275}
276
277impl TurnCanceller {
278    /// Cancels the current turn.
279    pub fn cancel(&self) {
280        self.inner.cancel();
281    }
282
283    /// Returns `true` if cancellation has been requested.
284    pub fn is_cancelled(&self) -> bool {
285        self.inner.is_requested()
286    }
287}
288
289/// Options for constructing an [`AgentRuntime`].
290pub struct AgentRuntimeOptions {
291    /// Loaded and merged configuration.
292    pub loaded_config: LoadedConfig,
293    /// The model provider implementation.
294    pub model_provider: Arc<dyn ModelProvider>,
295    /// Project root directory.
296    pub project_dir: PathBuf,
297    /// Optional custom tool executor (defaults to built-in tools).
298    pub tool_executor: Option<Arc<ToolExecutor>>,
299    /// Context packets to inject into the session.
300    pub context_packets: Vec<ContextPacket>,
301    /// Active skill names for this session.
302    pub active_skills: Vec<String>,
303    /// Seed messages for restoring a session.
304    pub initial_messages: Vec<ModelMessage>,
305    /// Seed events for restoring a persisted session without losing history.
306    pub initial_events: Vec<AgentEvent>,
307    /// Original creation timestamp for restored sessions.
308    pub initial_created_at: Option<u64>,
309    /// Original update timestamp for restored sessions.
310    pub initial_updated_at: Option<u64>,
311    /// Goal restored from a persisted session snapshot.
312    pub initial_goal: Option<crate::goal::types::SessionGoal>,
313    /// Session id for restoring an existing session.
314    pub session_id: Option<SessionId>,
315    /// Optional channel for forwarding agent events outside the runtime.
316    pub event_tx: Option<tokio::sync::mpsc::UnboundedSender<AgentEvent>>,
317    /// Replaceable runtime components. Defaults preserve NAVI's code-agent behavior.
318    pub runtime_components: Option<RuntimeComponents>,
319    /// Session-local title state written by the `set_session_title` tool.
320    /// When omitted, direct runtime users simply do not get agent-managed titles.
321    pub session_title_handle: Option<SessionTitleHandle>,
322    /// Explicit model for automatic durable-memory extraction. `None` means
323    /// extraction is disabled rather than silently billing the chat model.
324    pub memory_extraction_model: Option<MemoryExtractionModel>,
325    /// When true, skip auto-registering goal tools and skill loaders on the
326    /// provided tool executor. Used by host tool profiles (`chat_only`,
327    /// `host_tools_only`) so the host's filtered tool set is preserved.
328    pub skip_auto_tool_bootstrap: bool,
329}
330
331/// Provider selection for the opt-in per-turn memory extractor.
332#[derive(Clone)]
333pub struct MemoryExtractionModel {
334    pub provider: Arc<dyn ModelProvider>,
335    pub model_name: String,
336}
337
338/// The core agent runtime that manages sessions, turns, approvals, and events.
339pub struct AgentRuntime {
340    loaded_config: LoadedConfig,
341    model_provider: Arc<dyn ModelProvider>,
342    shared_model_provider: Arc<RwLock<Arc<dyn ModelProvider>>>,
343    shared_model_name: Arc<RwLock<String>>,
344    shared_config: Arc<RwLock<crate::config::NaviConfig>>,
345    project_dir: PathBuf,
346    tool_executor: Option<Arc<ToolExecutor>>,
347    session_store: SessionStore,
348    context_packets: Vec<ContextPacket>,
349    shared_context_packets: Arc<std::sync::Mutex<Vec<ContextPacket>>>,
350    active_skills: Vec<String>,
351    shared_available_skills: Arc<std::sync::Mutex<Vec<crate::skills::SkillManifest>>>,
352    shared_skill_pools: Arc<std::sync::Mutex<Vec<crate::skills::SkillPool>>>,
353    shared_active_skills: Arc<std::sync::Mutex<Vec<crate::skills::SkillManifest>>>,
354    prompt_cache: Arc<crate::prompt::PromptCache>,
355    runtime_components: RuntimeComponents,
356    initial_messages: Vec<ModelMessage>,
357    event_tx: Option<mpsc::UnboundedSender<AgentEvent>>,
358    cancel_token: CancelToken,
359    pending_approvals: PendingApprovals,
360    pending_questions: PendingQuestions,
361    pending_plan_reviews: PendingPlanReviews,
362    pending_sudo_passwords: PendingSudoPasswords,
363    event_bus: EventBus,
364    session: SessionState,
365    /// Goal runtime handle for the current session.
366    goal_runtime: Arc<GoalRuntimeHandle>,
367    /// Goal extension providing lifecycle hooks.
368    goal_extension: GoalExtension,
369    /// Whether the model used the `memory` tool with `write` action during the current turn.
370    /// Used for mutual exclusion with background extractMemories.
371    turn_used_memory_write: bool,
372    /// Last user task text — used for extractMemories context.
373    last_user_task: String,
374    /// Session title assigned by the chat model through `set_session_title`.
375    session_title_handle: SessionTitleHandle,
376    /// User-selected model for asynchronous automatic memory extraction.
377    memory_extraction_model: Option<MemoryExtractionModel>,
378    /// When true, do not auto-register goal/skill tools on the provided executor.
379    skip_auto_tool_bootstrap: bool,
380    /// Whether a user message is pending (set when send_turn is called
381    /// while the runtime is busy with auto-continuation).
382    pending_user_input: std::sync::atomic::AtomicBool,
383    /// Current agent mode (Default or Plan). In Plan mode, only read-only
384    /// tools are available and the model is instructed to propose a plan.
385    agent_mode: std::sync::RwLock<crate::plan_mode::AgentMode>,
386    /// Parser for `<proposed_plan>` tags in streaming text.
387    plan_parser: std::sync::Mutex<crate::plan_mode::ProposedPlanParser>,
388    /// Session-scoped memory manager shared with TurnContext (open once).
389    memory_manager: Arc<std::sync::Mutex<Option<Arc<crate::memory::MemoryManager>>>>,
390    /// Per-turn file snapshots for rewind / Revert.
391    rewind_store: Arc<std::sync::Mutex<crate::rewind::RewindStore>>,
392    /// From active harness packs: override for goals.max_auto_continue_turns.
393    harness_max_auto_continue: Option<u32>,
394    /// From active harness packs: preferred token budget when setting a goal.
395    harness_token_budget: Option<i64>,
396    /// Developer-context harness card for active packs.
397    harness_card: Option<String>,
398    /// Soft graph + skill allowlist merged (None = no extra lock).
399    harness_allow_tools: Option<Vec<String>>,
400}
401
402impl AgentRuntime {
403    /// Creates a new runtime from the given options.
404    pub fn new(options: AgentRuntimeOptions) -> Self {
405        let session_store = SessionStore::with_redaction(
406            options.loaded_config.data_dir.clone(),
407            options
408                .loaded_config
409                .config
410                .security
411                .redact_secrets_in_sessions,
412        );
413
414        let shared_context_packets =
415            Arc::new(std::sync::Mutex::new(options.context_packets.clone()));
416        let shared_available_skills = Arc::new(std::sync::Mutex::new(Vec::new()));
417        let shared_skill_pools = Arc::new(std::sync::Mutex::new(Vec::new()));
418        let shared_active_skills = Arc::new(std::sync::Mutex::new(Vec::new()));
419        let shared_model_provider = Arc::new(RwLock::new(options.model_provider.clone()));
420        let shared_model_name = Arc::new(RwLock::new(provider_request_model_name(
421            &options.loaded_config.config.model.provider,
422            &options.loaded_config.config.model.name,
423        )));
424        let shared_config = Arc::new(RwLock::new(options.loaded_config.config.clone()));
425        let prompt_cache = Arc::new(crate::prompt::PromptCache::new());
426        let runtime_components = options.runtime_components.unwrap_or_default();
427        let goal_service = Arc::new(GoalService::new());
428        let goal_runtime = Arc::new(GoalRuntimeHandle::new(options.initial_goal.clone()));
429        let goal_extension = GoalExtension::new(goal_service.clone(), goal_runtime.clone());
430
431        let rewind_sid = options
432            .session_id
433            .as_ref()
434            .map(|s| s.as_str().to_string())
435            .unwrap_or_else(|| "pending".to_string());
436        let rewind_store = Arc::new(std::sync::Mutex::new(crate::rewind::RewindStore::new(
437            &options.loaded_config.data_dir,
438            &rewind_sid,
439            &options.project_dir,
440        )));
441
442        Self {
443            loaded_config: options.loaded_config,
444            model_provider: options.model_provider,
445            shared_model_provider,
446            shared_model_name,
447            shared_config,
448            project_dir: options.project_dir,
449            tool_executor: options.tool_executor,
450            session_store,
451            context_packets: options.context_packets,
452            shared_context_packets,
453            active_skills: options.active_skills,
454            shared_available_skills,
455            shared_skill_pools,
456            shared_active_skills,
457            prompt_cache,
458            runtime_components,
459            initial_messages: options.initial_messages,
460            event_tx: options.event_tx,
461            cancel_token: CancelToken::new(),
462            pending_approvals: Arc::new(std::sync::Mutex::new(HashMap::new())),
463            pending_questions: Arc::new(std::sync::Mutex::new(HashMap::new())),
464            pending_plan_reviews: Arc::new(std::sync::Mutex::new(HashMap::new())),
465            pending_sudo_passwords: Arc::new(std::sync::Mutex::new(HashMap::new())),
466            event_bus: EventBus::new(),
467            session: SessionState::new_with_history(
468                options.session_id,
469                options.initial_events,
470                options.initial_created_at,
471                options.initial_updated_at,
472            ),
473            goal_runtime,
474            goal_extension,
475            turn_used_memory_write: false,
476            last_user_task: String::new(),
477            session_title_handle: options.session_title_handle.unwrap_or_default(),
478            memory_extraction_model: options.memory_extraction_model,
479            skip_auto_tool_bootstrap: options.skip_auto_tool_bootstrap,
480            pending_user_input: std::sync::atomic::AtomicBool::new(false),
481            agent_mode: std::sync::RwLock::new(crate::plan_mode::AgentMode::Default),
482            plan_parser: std::sync::Mutex::new(crate::plan_mode::ProposedPlanParser::new()),
483            memory_manager: Arc::new(std::sync::Mutex::new(None)),
484            rewind_store,
485            harness_max_auto_continue: None,
486            harness_token_budget: None,
487            harness_card: None,
488            harness_allow_tools: None,
489        }
490    }
491
492    /// Shared rewind store (write tools note dirty paths through this handle).
493    pub fn rewind_store_handle(&self) -> Arc<std::sync::Mutex<crate::rewind::RewindStore>> {
494        self.rewind_store.clone()
495    }
496
497    /// Rebind rewind store to the live session id (after start_session).
498    pub fn rebind_rewind_store(&self) {
499        let sid = self.session.id().as_str().to_string();
500        let mut store = self.rewind_store.lock().unwrap_or_else(|e| e.into_inner());
501        *store =
502            crate::rewind::RewindStore::new(&self.loaded_config.data_dir, &sid, &self.project_dir);
503    }
504
505    /// Returns all agent events recorded so far.
506    /// Returns the current session goal, if any.
507    pub fn get_goal(&self) -> Option<crate::goal::types::SessionGoal> {
508        self.goal_runtime.get_goal()
509    }
510
511    /// Sets or updates the session goal and notifies live clients.
512    pub fn set_goal(
513        &self,
514        objective: String,
515        token_budget: Option<i64>,
516    ) -> crate::goal::types::SessionGoal {
517        let goal = self.goal_runtime.set_objective(objective, token_budget);
518        self.goal_runtime.set_auto_continue(true);
519        self.publish_goal_updated(&goal);
520        goal
521    }
522
523    /// Sets or updates the session goal with a compact UI label.
524    pub fn set_goal_with_short_description(
525        &self,
526        objective: String,
527        short_description: Option<String>,
528        token_budget: Option<i64>,
529    ) -> crate::goal::types::SessionGoal {
530        let goal = self.goal_runtime.set_objective_with_short_description(
531            objective,
532            short_description,
533            token_budget,
534        );
535        self.goal_runtime.set_auto_continue(true);
536        self.publish_goal_updated(&goal);
537        goal
538    }
539
540    /// Clears the current session goal.
541    ///
542    /// Live UIs typically clear their chip optimistically when the user
543    /// requests clear; no terminal GoalUpdated is published (that would look
544    /// like a successful complete).
545    pub fn clear_goal(&self) {
546        self.goal_runtime.clear_goal();
547    }
548
549    /// Updates the stored goal (used after status transitions) and notifies clients.
550    pub fn update_goal(&self, goal: crate::goal::types::SessionGoal) {
551        self.goal_runtime.update_goal(goal.clone());
552        self.publish_goal_updated(&goal);
553    }
554
555    fn publish_goal_updated(&self, goal: &crate::goal::types::SessionGoal) {
556        self.event_bus.publish(RuntimeEventKind::GoalUpdated {
557            session_id: goal.session_id.clone(),
558            goal_id: goal.goal_id.as_str().to_string(),
559            objective: goal.objective.clone(),
560            short_description: goal.short_description.clone(),
561            status: goal.status,
562            tokens_used: goal.tokens_used,
563            token_budget: goal.token_budget,
564        });
565    }
566
567    /// Updates the goal checklist (replaces all tasks).
568    pub fn update_goal_checklist(
569        &self,
570        tasks: Vec<crate::goal::types::GoalTask>,
571    ) -> Option<crate::goal::types::SessionGoal> {
572        self.goal_runtime.update_checklist(tasks)
573    }
574
575    /// Updates a single task's status in the goal checklist.
576    pub fn update_goal_task_status(
577        &self,
578        task_id: usize,
579        status: crate::goal::types::TaskStatus,
580    ) -> Option<crate::goal::types::SessionGoal> {
581        self.goal_runtime.update_task_status(task_id, status)
582    }
583
584    /// Returns a continuation steering prompt if the goal is active and should auto-continue.
585    pub fn goal_idle_prompt(&self) -> Option<String> {
586        // In Plan mode, don't auto-continue — the user needs to confirm the plan first.
587        if self
588            .agent_mode
589            .read()
590            .unwrap_or_else(|e| e.into_inner())
591            .restricts_tools()
592        {
593            return None;
594        }
595        if !self.loaded_config.config.goals.enabled {
596            return None;
597        }
598        self.goal_extension.on_idle()
599    }
600
601    /// Returns the goals configuration.
602    pub fn goals_config(&self) -> crate::config::GoalsConfig {
603        self.loaded_config.config.goals.clone()
604    }
605
606    /// Returns a reference to the goal runtime handle.
607    pub fn goal_runtime(&self) -> &Arc<GoalRuntimeHandle> {
608        &self.goal_runtime
609    }
610
611    // ── Plan Mode ──────────────────────────────────────────────
612
613    /// Returns the current agent mode.
614    pub fn agent_mode(&self) -> crate::plan_mode::AgentMode {
615        *self.agent_mode.read().unwrap_or_else(|e| e.into_inner())
616    }
617
618    /// Enters Plan mode. Exploration tools + the session plan markdown file are
619    /// available; project writes and commands are blocked.
620    pub fn enter_plan_mode(&self) {
621        *self.agent_mode.write().unwrap_or_else(|e| e.into_inner()) =
622            crate::plan_mode::AgentMode::Plan;
623        *self.plan_parser.lock().unwrap_or_else(|e| e.into_inner()) =
624            crate::plan_mode::ProposedPlanParser::new();
625
626        let plan_path = crate::plan_store::session_plan_file_path(
627            &self.loaded_config.data_dir,
628            self.session.id().as_str(),
629        );
630        if let Some(parent) = plan_path.parent() {
631            let _ = std::fs::create_dir_all(parent);
632        }
633        if let Some(exec) = self.tool_executor.as_ref() {
634            exec.policy().set_plan_mode(true, Some(plan_path.clone()));
635        }
636
637        self.event_bus.publish(RuntimeEventKind::AgentModeChanged {
638            session_id: self.session.id().as_str().to_string(),
639            mode: crate::plan_mode::AgentMode::Plan,
640        });
641    }
642
643    /// Exits Plan mode and returns to normal execution.
644    pub fn exit_plan_mode(&self) {
645        *self.agent_mode.write().unwrap_or_else(|e| e.into_inner()) =
646            crate::plan_mode::AgentMode::Default;
647        if let Some(exec) = self.tool_executor.as_ref() {
648            exec.policy().set_plan_mode(false, None);
649        }
650        self.event_bus.publish(RuntimeEventKind::AgentModeChanged {
651            session_id: self.session.id().as_str().to_string(),
652            mode: crate::plan_mode::AgentMode::Default,
653        });
654    }
655
656    /// Feeds a text delta into the plan parser. Returns any completed plans.
657    pub fn feed_plan_text(&self, text: &str) -> Vec<crate::plan_mode::ProposedPlan> {
658        self.plan_parser
659            .lock()
660            .unwrap_or_else(|e| e.into_inner())
661            .push_text(text)
662    }
663
664    /// Drains any pending plans from the parser (call at end of turn).
665    pub fn drain_plans(&self) -> Vec<crate::plan_mode::ProposedPlan> {
666        self.plan_parser
667            .lock()
668            .unwrap_or_else(|e| e.into_inner())
669            .drain()
670    }
671
672    /// Returns true if the parser is currently inside a `<proposed_plan>` block.
673    pub fn is_parsing_plan(&self) -> bool {
674        self.plan_parser
675            .lock()
676            .unwrap_or_else(|e| e.into_inner())
677            .is_in_plan()
678    }
679
680    /// Returns `true` if a user message is pending (set when `send_turn` is
681    /// called while the runtime is busy with auto-continuation).
682    pub fn has_pending_user_input(&self) -> bool {
683        self.pending_user_input
684            .load(std::sync::atomic::Ordering::SeqCst)
685    }
686
687    /// Marks that a user message is pending.
688    pub fn set_pending_user_input(&self, pending: bool) {
689        self.pending_user_input
690            .store(pending, std::sync::atomic::Ordering::SeqCst);
691    }
692
693    /// Returns all agent events recorded so far.
694    pub fn events(&self) -> &[AgentEvent] {
695        self.session.events()
696    }
697
698    /// Returns the current session id.
699    pub fn session_id(&self) -> &SessionId {
700        self.session.id()
701    }
702
703    /// Returns the session title, if one has been derived.
704    pub fn session_title(&self) -> Option<&str> {
705        self.session.title()
706    }
707
708    /// Adds a context packet to the session and emits a `ContextUpdated` event.
709    pub fn add_context_packet(&mut self, packet: ContextPacket) {
710        self.context_packets.push(packet.clone());
711        self.shared_context_packets
712            .lock()
713            .unwrap_or_else(|e| e.into_inner())
714            .push(packet);
715        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
716    }
717
718    /// Clears all context packets and emits a `ContextUpdated` event.
719    pub fn clear_context_packets(&mut self) {
720        self.context_packets.clear();
721        self.shared_context_packets
722            .lock()
723            .unwrap_or_else(|e| e.into_inner())
724            .clear();
725        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
726    }
727
728    /// Returns the current context packets.
729    pub fn context_packets(&self) -> &[ContextPacket] {
730        &self.context_packets
731    }
732
733    /// Sets the session catalog filter for skills.
734    ///
735    /// Empty `skills` means **all** discovered skills are catalog-active (default).
736    /// Non-empty restricts the Available Skills catalog to those ids/names.
737    /// Never injects skill instruction bodies into the prompt.
738    ///
739    /// When a skill has a harness pack under `{data_dir}/harnesses/<id>/`, applies
740    /// soft graph allowlists, loop max_turns / token_budget, and harness card text.
741    pub fn set_active_skills(&mut self, skills: Vec<String>) {
742        self.active_skills = skills;
743        let catalog = self.load_catalog_skills();
744        let pools = self.load_catalog_pools();
745        *self
746            .shared_available_skills
747            .lock()
748            .unwrap_or_else(|e| e.into_inner()) = catalog.clone();
749        *self
750            .shared_skill_pools
751            .lock()
752            .unwrap_or_else(|e| e.into_inner()) = pools;
753        // Soft harness applies only to session-active skills, never the full catalog.
754        let harness_skills = self.session_active_skill_manifests();
755        *self
756            .shared_active_skills
757            .lock()
758            .unwrap_or_else(|e| e.into_inner()) = harness_skills.clone();
759        self.apply_harness_packs_for_active(&harness_skills);
760        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
761    }
762
763    /// Skill manifests selected for **session** activation (CLI `--skill`, host
764    /// `set_active_skills`, config `skills.active` when session list is empty).
765    /// Empty selection → no soft harness allowlist.
766    fn session_active_skill_manifests(&self) -> Vec<crate::skills::SkillManifest> {
767        let available = self.discover_available_skills().unwrap_or_default();
768        let session = &self.active_skills;
769        let configured = &self.loaded_config.config.skills.active;
770        // Unlike catalog `active_skills`, an empty selection means **no**
771        // harness soft-apply (do not default to every discovered skill).
772        if session.is_empty() && configured.is_empty() {
773            return Vec::new();
774        }
775        crate::skills::active_skills(&available, configured, session)
776    }
777
778    /// Recompute harness soft-apply state from currently active skill manifests.
779    fn apply_harness_packs_for_active(&mut self, manifests: &[crate::skills::SkillManifest]) {
780        let applied =
781            crate::harness_pack::apply_harness_for_skills(&self.loaded_config.data_dir, manifests);
782        self.harness_max_auto_continue = applied.max_auto_continue_turns;
783        self.harness_token_budget = applied.token_budget;
784        self.harness_card = if applied.harness_card.is_empty() {
785            None
786        } else {
787            Some(applied.harness_card)
788        };
789        self.harness_allow_tools = applied.allow_tools;
790        // If a harness sets a token budget and no goal exists yet, leave budget
791        // for host/model create_goal; if a goal is Active, optionally annotate budget.
792        if let Some(budget) = applied.token_budget {
793            if let Some(mut goal) = self.goal_runtime.get_goal() {
794                if goal.token_budget.is_none() && goal.status.should_auto_continue() {
795                    goal.token_budget = Some(budget);
796                    self.goal_runtime.update_goal(goal);
797                }
798            }
799        }
800    }
801
802    /// Effective max auto-continue turns (harness pack overrides goals config when set).
803    pub fn effective_max_auto_continue_turns(&self) -> u32 {
804        let cfg = self.goals_config().max_auto_continue_turns;
805        match self.harness_max_auto_continue {
806            Some(h) if h > 0 => {
807                if cfg > 0 {
808                    cfg.min(h)
809                } else {
810                    h
811                }
812            }
813            _ => cfg,
814        }
815    }
816
817    /// Developer harness card for the active pack(s), if any.
818    pub fn harness_card(&self) -> Option<&str> {
819        self.harness_card.as_deref()
820    }
821
822    /// Lists available model options from the loaded configuration.
823    pub fn list_models(&self) -> Vec<ModelOption> {
824        available_model_options(&self.loaded_config.config)
825    }
826
827    /// Current `(provider_id, model_name)` selection for this runtime.
828    pub fn model_selection(&self) -> (&str, &str) {
829        (
830            self.loaded_config.config.model.provider.as_str(),
831            self.loaded_config.config.model.name.as_str(),
832        )
833    }
834
835    /// Changes the selected model and emits a `ContextUpdated` event.
836    pub fn set_model(&mut self, provider: impl Into<String>, model: impl Into<String>) {
837        self.loaded_config.config.model.provider =
838            canonical_provider_id(&provider.into()).to_string();
839        self.loaded_config.config.model.name = model.into();
840        self.update_shared_model_state();
841        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
842    }
843
844    /// Replaces the runtime configuration and provider used by subsequent turns.
845    pub fn set_model_provider(
846        &mut self,
847        loaded_config: LoadedConfig,
848        model_provider: Arc<dyn ModelProvider>,
849    ) {
850        self.loaded_config = loaded_config;
851        self.model_provider = model_provider;
852        self.update_shared_model_state();
853        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
854    }
855
856    /// Registers a host-provided tool with the runtime's tool executor.
857    /// Creates a default executor if none exists yet.
858    pub fn register_host_tool(&mut self, tool: Arc<dyn Tool>) -> Result<()> {
859        if self.tool_executor.is_none() {
860            let security_policy = SecurityPolicy::new(
861                self.project_dir.clone(),
862                self.loaded_config.data_dir.clone(),
863                self.loaded_config.config.effective_security_config(),
864            )?;
865            self.tool_executor = Some(Arc::new(ToolExecutor::with_security_policy(
866                security_policy,
867                self.runtime_components.security.clone(),
868            )));
869        }
870
871        let Some(executor) = self.tool_executor.as_mut() else {
872            return Err(anyhow::anyhow!("tool executor unavailable"));
873        };
874        let Some(executor) = Arc::get_mut(executor) else {
875            return Err(anyhow::anyhow!(
876                "cannot register host tool while tool executor is shared"
877            ));
878        };
879        executor.register_tool(tool);
880        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
881        Ok(())
882    }
883
884    /// Returns a broadcast receiver for [`RuntimeEvent`]s.
885    pub fn stream_events(&self) -> broadcast::Receiver<RuntimeEvent> {
886        self.event_bus.stream_events()
887    }
888
889    /// Cancels the currently running turn.
890    pub fn cancel_turn(&self) {
891        self.turn_canceller().cancel();
892    }
893
894    /// Resolves a pending approval by id. Returns `true` if found.
895    pub fn resolve_approval(&self, decision: ApprovalDecision) -> bool {
896        self.approval_resolver().resolve(decision)
897    }
898
899    /// Resolves a pending interactive question by id. Returns `true` if found.
900    pub fn resolve_question(&self, response: QuestionResponse) -> bool {
901        self.question_resolver().resolve(response)
902    }
903
904    /// Resolves a pending plan review by invocation id. Returns `true` if found.
905    pub fn resolve_plan_review(&self, response: PlanReviewResponse) -> bool {
906        self.plan_review_resolver().resolve(response)
907    }
908
909    /// Returns an [`ApprovalResolver`] handle for external approval resolution.
910    pub fn approval_resolver(&self) -> ApprovalResolver {
911        ApprovalResolver {
912            pending_approvals: self.pending_approvals.clone(),
913            runtime_events_tx: self.event_bus.sender(),
914        }
915    }
916
917    /// Returns a [`QuestionResolver`] handle for external question resolution.
918    pub fn question_resolver(&self) -> QuestionResolver {
919        QuestionResolver {
920            pending_questions: self.pending_questions.clone(),
921            runtime_events_tx: self.event_bus.sender(),
922        }
923    }
924
925    /// Returns a [`PlanReviewResolver`] handle for external plan-review resolution.
926    pub fn plan_review_resolver(&self) -> PlanReviewResolver {
927        PlanReviewResolver {
928            pending_reviews: self.pending_plan_reviews.clone(),
929            runtime_events_tx: self.event_bus.sender(),
930        }
931    }
932
933    /// Returns a [`SudoPasswordResolver`] for interactive sudo prompts.
934    pub fn sudo_password_resolver(&self) -> SudoPasswordResolver {
935        SudoPasswordResolver {
936            pending: self.pending_sudo_passwords.clone(),
937        }
938    }
939
940    pub fn resolve_sudo_password(&self, response: SudoPasswordResponse) -> bool {
941        self.sudo_password_resolver().resolve(response)
942    }
943
944    /// Returns a [`TurnCanceller`] handle for external cancellation.
945    pub fn turn_canceller(&self) -> TurnCanceller {
946        TurnCanceller {
947            inner: self.cancel_token.clone(),
948        }
949    }
950
951    /// Starts a new session (or restarts if one is already active).
952    /// Returns the session id.
953    pub fn start_session(&mut self) -> Result<SessionId> {
954        if self.session.started() {
955            self.goal_extension
956                .on_session_end(self.session.id().as_str());
957            self.runtime_components
958                .hooks
959                .on_session_end(self.session.id().as_str());
960
961            // Light auto-memory consolidation on session end (stale + dedup, no model needed)
962            let _ = self.consolidate_auto_memory();
963
964            self.event_bus.publish(RuntimeEventKind::SessionFinished {
965                session_id: self.session.id().as_str().to_string(),
966            });
967        }
968        self.cancel_token.reset();
969        self.pending_approvals = Arc::new(std::sync::Mutex::new(HashMap::new()));
970        self.pending_questions = Arc::new(std::sync::Mutex::new(HashMap::new()));
971        self.pending_plan_reviews = Arc::new(std::sync::Mutex::new(HashMap::new()));
972        self.pending_sudo_passwords = Arc::new(std::sync::Mutex::new(HashMap::new()));
973        self.session.start();
974
975        let (session_runtime, event_rx) = self.build_session_runtime()?;
976        self.session.set_runtime(session_runtime, event_rx);
977
978        let id = self.session.id().clone();
979        // Bind rewind store to this session id and attach it to the tool executor.
980        self.rebind_rewind_store();
981        if let Ok(_) = self.ensure_tool_executor() {
982            if let Some(stored) = self.tool_executor.as_mut() {
983                if let Some(exec) = Arc::get_mut(stored) {
984                    exec.set_rewind_store(Some(self.rewind_store.clone()));
985                }
986            }
987        }
988        // Goal lifecycle: session start + register runtime
989        self.goal_extension.on_session_start(id.as_str());
990        self.runtime_components.hooks.on_session_start(id.as_str());
991        self.event_bus.publish(RuntimeEventKind::SessionStarted {
992            session_id: id.as_str().to_string(),
993        });
994
995        Ok(id)
996    }
997
998    /// Sends a user task to the agent and waits for the full response.
999    /// Starts a session automatically if one is not active.
1000    /// Sends a user turn with optional multimodal content parts.
1001    ///
1002    /// When `content_parts` is non-empty, the message is created as a
1003    /// multimodal user message containing both text and images.
1004    ///
1005    /// After a successful turn, if a thread goal is still active, runs automatic
1006    /// continuation turns (idle lifecycle) until the goal stops, the auto-continue
1007    /// limit is hit, plan mode is on, or the user has pending input.
1008    pub async fn send_turn_with_parts(
1009        &mut self,
1010        task: String,
1011        content_parts: Vec<crate::model::ContentPart>,
1012        thinking_override: Option<crate::model::ThinkingConfig>,
1013    ) -> Result<ModelResponse> {
1014        let mut response = self
1015            .send_turn_once(task, content_parts, thinking_override)
1016            .await?;
1017
1018        let goals_config = self.goals_config();
1019        let max_auto = self.effective_max_auto_continue_turns();
1020        let mut auto_continue_count = 0u32;
1021        loop {
1022            if !goals_config.enabled {
1023                break;
1024            }
1025            if max_auto > 0 && auto_continue_count >= max_auto {
1026                self.goal_runtime
1027                    .record_blocked_turn("auto-continuation limit reached");
1028                if let Some(goal) = self.goal_runtime.get_goal() {
1029                    self.publish_goal_updated(&goal);
1030                }
1031                break;
1032            }
1033            if self.has_pending_user_input() {
1034                break;
1035            }
1036            let Some(prompt) = self.goal_idle_prompt() else {
1037                break;
1038            };
1039            auto_continue_count += 1;
1040            tracing::info!(
1041                auto_continue_count,
1042                "starting automatic goal continuation turn"
1043            );
1044            // Continuation is injected as a model-visible steering message (not a
1045            // host/user chat turn). Content matches the goal steering template.
1046            response = self.send_turn_once(prompt, Vec::new(), None).await?;
1047        }
1048
1049        Ok(response)
1050    }
1051
1052    /// Single turn without goal auto-continuation.
1053    async fn send_turn_once(
1054        &mut self,
1055        task: String,
1056        content_parts: Vec<crate::model::ContentPart>,
1057        thinking_override: Option<crate::model::ThinkingConfig>,
1058    ) -> Result<ModelResponse> {
1059        // Mark that user input is being processed.
1060        self.pending_user_input
1061            .store(false, std::sync::atomic::Ordering::SeqCst);
1062
1063        if !self.session.started() || self.session.runtime().is_none() {
1064            self.start_session()?;
1065        }
1066
1067        // Apply per-turn thinking override before the turn runs so
1068        // build_model_request picks it up from the shared config.
1069        if let Some(thinking) = thinking_override {
1070            let level_str = thinking.as_config_str();
1071            // NOTE: we only update shared_config, NOT loaded_config. Mutating
1072            // loaded_config would permanently corrupt the original config and
1073            // leak the last override into future turns that pass thinking: None.
1074            self.shared_config
1075                .write()
1076                .unwrap_or_else(|e| e.into_inner())
1077                .tui
1078                .thinking_level = level_str.to_string();
1079        }
1080
1081        let submission_tx = self
1082            .session
1083            .runtime()
1084            .ok_or_else(|| anyhow::anyhow!("session not started"))?
1085            .submission_tx
1086            .clone();
1087
1088        let mut event_rx = self
1089            .session
1090            .take_event_rx()
1091            .ok_or_else(|| anyhow::anyhow!("session event stream unavailable"))?;
1092
1093        self.cancel_token.reset();
1094
1095        let turn_id = self.session.next_turn_id();
1096        tracing::info!(
1097            project = %self.project_dir.display(),
1098            provider = %self.loaded_config.config.model.provider,
1099            model = %self.loaded_config.config.model.name,
1100            "agent task submitted"
1101        );
1102        let session_id = self.session.id().as_str().to_string();
1103        self.runtime_components
1104            .hooks
1105            .on_turn_start(&session_id, &task);
1106        self.goal_extension.on_turn_start(&session_id, &task);
1107        self.record_event(AgentEvent::UserTaskSubmitted {
1108            text: task.clone(),
1109            content_parts: content_parts.clone(),
1110            submitted_at: Some(crate::session::current_unix_timestamp()),
1111        });
1112        self.last_user_task = task.clone();
1113
1114        // Pre-turn file snapshot (dirty paths only).
1115        let prompt_index = self
1116            .session
1117            .events()
1118            .iter()
1119            .filter(|e| matches!(e, AgentEvent::UserTaskSubmitted { .. }))
1120            .count()
1121            .saturating_sub(1);
1122        let captured_index = {
1123            let mut store = self.rewind_store.lock().unwrap_or_else(|e| e.into_inner());
1124            // Ensure store is keyed to live session id.
1125            if store.root().to_string_lossy().contains("pending") {
1126                *store = crate::rewind::RewindStore::new(
1127                    &self.loaded_config.data_dir,
1128                    self.session.id().as_str(),
1129                    &self.project_dir,
1130                );
1131            }
1132            // New turn: only count writes that happen after this capture.
1133            store.clear_turn_written();
1134            match store.capture_point(
1135                prompt_index,
1136                &task,
1137                crate::session::current_unix_timestamp(),
1138            ) {
1139                Ok(p) => Some(p.prompt_index),
1140                Err(e) => {
1141                    tracing::warn!(error = %e, "rewind: failed to capture pre-turn snapshot");
1142                    None
1143                }
1144            }
1145        };
1146
1147        // Persist immediately so sidebars list the session while the turn is
1148        // still running. The chat model supplies its title through the tool.
1149        self.persist_submitted_session();
1150        self.event_bus.publish(RuntimeEventKind::TurnStarted {
1151            turn_id: turn_id.clone(),
1152        });
1153
1154        let (response_tx, response_rx) = tokio::sync::oneshot::channel();
1155        if let Err(e) = submission_tx.send(crate::session::SessionCommand::Turn(
1156            crate::session::Submission {
1157                task,
1158                content_parts,
1159                response_tx,
1160            },
1161        )) {
1162            return Err(anyhow::anyhow!("failed to send submission: {}", e));
1163        }
1164
1165        let mut response_rx = response_rx;
1166        let result: Result<String> = loop {
1167            tokio::select! {
1168                res = &mut response_rx => {
1169                    break match res {
1170                        Ok(Ok(text)) => Ok(text),
1171                        Ok(Err(err)) => Err(anyhow::anyhow!("turn failed: {err}")),
1172                        Err(_) => Err(anyhow::anyhow!("turn cancelled or panicked")),
1173                    };
1174                }
1175                Some(event) = event_rx.recv() => {
1176                    self.record_event(event);
1177                    // Apply + persist as soon as `set_session_title` completes so
1178                    // live UIs and session lists do not wait for the whole turn.
1179                    if self.apply_pending_session_title() {
1180                        if let Err(err) = self.session.snapshot(
1181                            &self.project_dir,
1182                            &self.session_store,
1183                            &self.event_bus,
1184                            self.goal_runtime.get_goal(),
1185                        ) {
1186                            tracing::debug!(error = %err, "early title snapshot failed");
1187                        }
1188                    }
1189                }
1190            }
1191        };
1192
1193        while let Ok(event) = event_rx.try_recv() {
1194            self.record_event(event);
1195        }
1196        drop(event_rx);
1197        self.session.set_updated_at(current_unix_timestamp());
1198        let _ = self.apply_pending_session_title();
1199
1200        // Record paths created this turn so restore can delete them.
1201        if let Some(idx) = captured_index {
1202            let mut store = self.rewind_store.lock().unwrap_or_else(|e| e.into_inner());
1203            if let Err(e) = store.finalize_turn_created_paths(idx) {
1204                tracing::debug!(error = %e, prompt_index = idx, "rewind: finalize created_paths failed");
1205            }
1206        }
1207
1208        match &result {
1209            Ok(text) => {
1210                self.goal_extension.on_turn_end(&session_id);
1211                self.runtime_components
1212                    .hooks
1213                    .on_turn_end(self.session.id().as_str(), text);
1214                self.event_bus.publish(RuntimeEventKind::TurnCompleted {
1215                    turn_id,
1216                    text: text.clone(),
1217                });
1218
1219                // extractMemories: background extraction per turn (fire-and-forget)
1220                // Skip if the model already wrote memories during this turn
1221                let model_wrote_memory = self.turn_used_memory_write;
1222
1223                if !model_wrote_memory {
1224                    // Build conversation snippet from user task + assistant response
1225                    let user_task = self.last_user_task.clone();
1226                    let conversation = if user_task.is_empty() {
1227                        format!("Assistant: {}", text)
1228                    } else {
1229                        format!("User: {}\n\nAssistant: {}", user_task, text)
1230                    };
1231                    self.try_extract_memories(&session_id, &conversation);
1232                }
1233
1234                // Reset per-turn flag
1235                self.turn_used_memory_write = false;
1236
1237                // Auto-dream: fire-and-forget check after each turn
1238                self.try_auto_dream();
1239
1240                // Auto-distill: fire-and-forget check after each turn
1241                self.try_auto_distill();
1242            }
1243            Err(err) => {
1244                self.goal_extension.on_turn_error(&err.to_string());
1245                // Coalesce streamed deltas into a durable ModelOutput so a mid-turn
1246                // failure / crash does not erase what the model already produced.
1247                self.flush_partial_model_output_from_events();
1248                self.record_event(AgentEvent::Error {
1249                    message: err.to_string(),
1250                });
1251                // Best-effort persist immediately — Desktop/TUI also snapshot after.
1252                if let Err(snap_err) = self.snapshot_session() {
1253                    tracing::warn!(
1254                        error = %snap_err,
1255                        "failed to snapshot session after turn error"
1256                    );
1257                }
1258            }
1259        }
1260
1261        result.map(|text| {
1262            tracing::info!(chars = text.len(), "agent task completed");
1263            ModelResponse { text }
1264        })
1265    }
1266
1267    /// If the current turn streamed text/thinking but never emitted `ModelOutput`
1268    /// (e.g. provider error mid-stream), write a `ModelOutput` from the deltas so
1269    /// session JSON reload keeps the partial answer.
1270    fn flush_partial_model_output_from_events(&mut self) {
1271        let events = self.session.events();
1272        let mut last_user = None;
1273        for (i, event) in events.iter().enumerate() {
1274            if matches!(event, AgentEvent::UserTaskSubmitted { .. }) {
1275                last_user = Some(i);
1276            }
1277        }
1278        let start = last_user.map(|i| i + 1).unwrap_or(0);
1279        let mut text = String::new();
1280        let mut thinking = String::new();
1281        let mut saw_output = false;
1282        for event in &events[start..] {
1283            match event {
1284                AgentEvent::ModelOutput { .. } => {
1285                    saw_output = true;
1286                    break;
1287                }
1288                AgentEvent::ModelDelta { text: delta } => text.push_str(delta),
1289                AgentEvent::ModelThinkingDelta { text: delta } => thinking.push_str(delta),
1290                _ => {}
1291            }
1292        }
1293        if saw_output {
1294            return;
1295        }
1296        if text.is_empty() && thinking.is_empty() {
1297            return;
1298        }
1299        self.record_event(AgentEvent::ModelOutput {
1300            text,
1301            thinking: if thinking.is_empty() {
1302                None
1303            } else {
1304                Some(thinking)
1305            },
1306        });
1307    }
1308
1309    pub async fn submit_task(&mut self, task: String) -> Result<ModelResponse> {
1310        self.send_turn_with_parts(task, Vec::new(), None).await
1311    }
1312
1313    /// Sends a plain text user turn (no images).
1314    pub async fn send_turn(&mut self, task: String) -> Result<ModelResponse> {
1315        self.send_turn_with_parts(task, Vec::new(), None).await
1316    }
1317
1318    /// Force-compact live conversation history using the session's own model.
1319    ///
1320    /// Summarizes older turns, replaces the live message list, and emits
1321    /// `AutoCompactStarted` / `AutoCompactCompleted` (or Failed) events.
1322    pub async fn compact_now(&mut self) -> Result<crate::compact::CompactOutcome> {
1323        if !self.session.started() || self.session.runtime().is_none() {
1324            // No live session loop yet — compact the seed history in place.
1325            self.record_event(AgentEvent::AutoCompactStarted);
1326            let provider = self
1327                .shared_model_provider
1328                .read()
1329                .unwrap_or_else(|e| e.into_inner())
1330                .clone();
1331            let model = self
1332                .shared_model_name
1333                .read()
1334                .unwrap_or_else(|e| e.into_inner())
1335                .clone();
1336            let harness = self.loaded_config.config.harness.clone();
1337            let mut state = crate::compact::CompactState::new(
1338                crate::config::effective_context_window(&self.loaded_config.config),
1339            );
1340            match state
1341                .force_compact(
1342                    &mut self.initial_messages,
1343                    provider.as_ref(),
1344                    &model,
1345                    &harness,
1346                )
1347                .await
1348            {
1349                Ok(Some(outcome)) => {
1350                    self.collapse_events_after_compact(&outcome);
1351                    // collapse already includes AutoCompactCompleted; also
1352                    // publish for live subscribers if session wasn't recording.
1353                    self.event_bus
1354                        .publish(RuntimeEventKind::AutoCompactCompleted {
1355                            tokens_saved: outcome.tokens_saved,
1356                            summary: outcome.summary.clone(),
1357                            kept_recent_messages: outcome.kept_recent_messages,
1358                        });
1359                    Ok(outcome)
1360                }
1361                Ok(None) => {
1362                    let reason = "nothing to compact".to_string();
1363                    self.record_event(AgentEvent::AutoCompactFailed {
1364                        reason: reason.clone(),
1365                    });
1366                    Err(anyhow::anyhow!(reason))
1367                }
1368                Err(e) => {
1369                    self.record_event(AgentEvent::AutoCompactFailed {
1370                        reason: e.to_string(),
1371                    });
1372                    Err(e)
1373                }
1374            }
1375        } else {
1376            let submission_tx = self
1377                .session
1378                .runtime()
1379                .ok_or_else(|| anyhow::anyhow!("session not started"))?
1380                .submission_tx
1381                .clone();
1382
1383            // Drain compact events from the session loop through the normal
1384            // event channel so subscribers (TUI) see them live.
1385            let mut event_rx = self
1386                .session
1387                .take_event_rx()
1388                .ok_or_else(|| anyhow::anyhow!("session event stream unavailable"))?;
1389
1390            let (response_tx, response_rx) = tokio::sync::oneshot::channel();
1391            if let Err(e) =
1392                submission_tx.send(crate::session::SessionCommand::Compact { response_tx })
1393            {
1394                return Err(anyhow::anyhow!("failed to send compact command: {}", e));
1395            }
1396
1397            let mut response_rx = response_rx;
1398            let result: Result<crate::compact::CompactOutcome> = loop {
1399                tokio::select! {
1400                    res = &mut response_rx => {
1401                        break match res {
1402                            Ok(Ok(outcome)) => Ok(outcome),
1403                            Ok(Err(err)) => Err(anyhow::anyhow!("compact failed: {err}")),
1404                            Err(_) => Err(anyhow::anyhow!("compact cancelled or session loop exited")),
1405                        };
1406                    }
1407                    Some(event) = event_rx.recv() => {
1408                        self.record_event(event);
1409                    }
1410                }
1411            };
1412
1413            while let Ok(event) = event_rx.try_recv() {
1414                self.record_event(event);
1415            }
1416            // SessionEventReceiver::Drop returns the channel to the session slot.
1417
1418            if let Ok(ref outcome) = result {
1419                self.collapse_events_after_compact(outcome);
1420                self.sync_initial_messages_from_compact(outcome);
1421                if let Err(err) = self.snapshot_session() {
1422                    tracing::warn!(error = %err, "failed to snapshot session after compact");
1423                }
1424            }
1425
1426            result
1427        }
1428    }
1429
1430    fn collapse_events_after_compact(&mut self, outcome: &crate::compact::CompactOutcome) {
1431        // Rebuild a minimal event log: summary as the sole prior user turn.
1432        // Keeps session reloads from resurrecting the pre-compact history.
1433        let summary_event = AgentEvent::UserTaskSubmitted {
1434            text: format!(
1435                "Here is a summary of the conversation so far:\n\n{}",
1436                outcome.summary
1437            ),
1438            content_parts: Vec::new(),
1439            submitted_at: Some(crate::session::current_unix_timestamp()),
1440        };
1441        let completed = AgentEvent::AutoCompactCompleted {
1442            tokens_saved: outcome.tokens_saved,
1443            summary: outcome.summary.clone(),
1444            // Event log is fully replaced — recent turns are not reconstructed here.
1445            kept_recent_messages: 0,
1446        };
1447        self.session.replace_events(vec![summary_event, completed]);
1448    }
1449
1450    fn sync_initial_messages_from_compact(&mut self, outcome: &crate::compact::CompactOutcome) {
1451        // Preserve any system/developer prefix already in initial_messages.
1452        let prefix: Vec<_> = self
1453            .initial_messages
1454            .iter()
1455            .filter(|m| {
1456                matches!(
1457                    m.role,
1458                    crate::model::ModelRole::System | crate::model::ModelRole::Developer
1459                )
1460            })
1461            .cloned()
1462            .collect();
1463        self.initial_messages = prefix;
1464        self.initial_messages
1465            .push(crate::compact::compact_summary_user_message(
1466                &outcome.summary,
1467            ));
1468    }
1469
1470    /// Rewind live conversation history for an edited past user message.
1471    ///
1472    /// Keeps the first `keep_user_turns` user turns (and their assistant/tool
1473    /// follow-ups), drops everything after, and truncates recorded session
1474    /// events the same way. Also restores project files to the pre-turn state
1475    /// of the first dropped user prompt (session rewind).
1476    ///
1477    /// Caller should then `send_turn` with the new text (or leave the prompt in
1478    /// the input for the user to edit).
1479    ///
1480    /// `keep_user_turns = 0` keeps only system/developer preamble.
1481    pub async fn rewind_to_user_turns(&mut self, keep_user_turns: usize) -> Result<usize> {
1482        // Restore filesystem to state before user turn `keep_user_turns` (and
1483        // drop checkpoints at/after that index). Best-effort; history still rewinds.
1484        let fs_summary = {
1485            let mut store = self.rewind_store.lock().unwrap_or_else(|e| e.into_inner());
1486            store.restore_to(keep_user_turns)
1487        };
1488        if !fs_summary.errors.is_empty() {
1489            for err in &fs_summary.errors {
1490                tracing::warn!(error = %err, "rewind: filesystem restore issue");
1491            }
1492        }
1493        if fs_summary.total_changes() > 0 {
1494            tracing::info!(
1495                restored = fs_summary.restored,
1496                deleted = fs_summary.deleted,
1497                keep_user_turns,
1498                "rewind: restored project files"
1499            );
1500        }
1501
1502        if !self.session.started() || self.session.runtime().is_none() {
1503            // Nothing live yet — just reset seed messages/events for a clean start.
1504            self.session.truncate_events_to_user_turns(keep_user_turns);
1505            // Seed messages for next start_session: drop user turns past keep.
1506            crate::session::truncate_messages_to_user_turns(
1507                &mut self.initial_messages,
1508                keep_user_turns,
1509            );
1510            return Ok(self.initial_messages.len());
1511        }
1512
1513        let submission_tx = self
1514            .session
1515            .runtime()
1516            .ok_or_else(|| anyhow::anyhow!("session not started"))?
1517            .submission_tx
1518            .clone();
1519
1520        let (response_tx, response_rx) = tokio::sync::oneshot::channel();
1521        if let Err(e) = submission_tx.send(crate::session::SessionCommand::TruncateToUserTurns {
1522            keep_user_turns,
1523            response_tx,
1524        }) {
1525            return Err(anyhow::anyhow!("failed to send rewind command: {}", e));
1526        }
1527
1528        let remaining = response_rx
1529            .await
1530            .map_err(|_| anyhow::anyhow!("rewind cancelled or session loop exited"))??;
1531
1532        self.session.truncate_events_to_user_turns(keep_user_turns);
1533        // Keep initial_messages in sync if the session is later restarted.
1534        crate::session::truncate_messages_to_user_turns(
1535            &mut self.initial_messages,
1536            keep_user_turns,
1537        );
1538
1539        // Persist truncated history so reload does not resurrect dropped turns.
1540        if let Err(err) = self.snapshot_session() {
1541            tracing::warn!(error = %err, "failed to snapshot session after rewind");
1542        }
1543
1544        Ok(remaining)
1545    }
1546
1547    /// List rewind checkpoints (user-turn snapshots) for the live session.
1548    pub fn list_rewind_points(&self) -> Vec<crate::rewind::RewindPointMeta> {
1549        self.rewind_store
1550            .lock()
1551            .unwrap_or_else(|e| e.into_inner())
1552            .load_points()
1553    }
1554
1555    /// Filesystem restore summary from the last `restore_to` is not retained;
1556    /// use [`Self::rewind_to_user_turns`] which restores as a side effect.
1557    pub fn rewind_store_points_path(&self) -> std::path::PathBuf {
1558        self.rewind_store
1559            .lock()
1560            .unwrap_or_else(|e| e.into_inner())
1561            .root()
1562            .to_path_buf()
1563    }
1564
1565    /// Applies a title set by the current chat model and informs live clients.
1566    ///
1567    /// Returns `true` when a new title was applied so callers can decide whether
1568    /// to persist immediately (mid-turn) or fold the change into a later snapshot.
1569    fn apply_pending_session_title(&mut self) -> bool {
1570        let Some(title) = self.session_title_handle.take() else {
1571            return false;
1572        };
1573        self.session.set_title(Some(title.clone()));
1574        self.event_bus
1575            .publish(RuntimeEventKind::SessionTitleUpdated {
1576                session_id: self.session.id().as_str().to_string(),
1577                title,
1578            });
1579        true
1580    }
1581
1582    /// Creates a [`SessionSnapshot`] of the current session state for persistence.
1583    pub fn snapshot_session(&mut self) -> Result<SessionSnapshot> {
1584        self.session.set_updated_at(current_unix_timestamp());
1585        let _ = self.apply_pending_session_title();
1586        let snapshot = self.session.snapshot(
1587            &self.project_dir,
1588            &self.session_store,
1589            &self.event_bus,
1590            self.goal_runtime.get_goal(),
1591        )?;
1592        self.save_trace_snapshot(snapshot.id.as_str());
1593        Ok(snapshot)
1594    }
1595
1596    /// Creates and persists a session snapshot without blocking the async runtime.
1597    pub async fn snapshot_session_async(&mut self) -> Result<SessionSnapshot> {
1598        self.session.set_updated_at(current_unix_timestamp());
1599        let _ = self.apply_pending_session_title();
1600        let snapshot = self
1601            .session
1602            .snapshot_async(
1603                &self.project_dir,
1604                &self.session_store,
1605                &self.event_bus,
1606                self.goal_runtime.get_goal(),
1607            )
1608            .await?;
1609        self.save_trace_snapshot(snapshot.id.as_str());
1610        Ok(snapshot)
1611    }
1612
1613    fn save_trace_snapshot(&self, session_id: &str) {
1614        let traces = turn_traces_from_events(
1615            session_id,
1616            &self.loaded_config.config.model.provider,
1617            &self.loaded_config.config.model.name,
1618            self.session.events(),
1619        );
1620        if traces.is_empty() {
1621            return;
1622        }
1623        let store = TraceStore::new(&self.loaded_config.data_dir);
1624        if let Err(err) = store.save_session_traces(session_id, &traces) {
1625            tracing::warn!(error = %err, session_id, "failed to save turn traces");
1626        }
1627    }
1628
1629    /// Test-only accessor; production code uses `self.session_store` directly.
1630    #[cfg(test)]
1631    pub(crate) fn session_store(&self) -> &SessionStore {
1632        &self.session_store
1633    }
1634
1635    fn ensure_tool_executor(&mut self) -> Result<Arc<ToolExecutor>> {
1636        if let Some(executor) = self.tool_executor.as_mut() {
1637            if let Some(executor) = Arc::get_mut(executor) {
1638                if !self.skip_auto_tool_bootstrap {
1639                    Self::register_goal_tools_on_executor(
1640                        executor,
1641                        self.goal_runtime.clone(),
1642                        self.loaded_config.config.goals.enabled,
1643                    );
1644                    executor.register_skill_loader(
1645                        self.project_dir.clone(),
1646                        self.loaded_config.data_dir.clone(),
1647                        self.shared_config.clone(),
1648                    );
1649                }
1650                executor.set_rewind_store(Some(self.rewind_store.clone()));
1651            } else {
1652                tracing::warn!(
1653                    "tool executor Arc is shared; cannot install goal tools in place \
1654                     (strong_count={})",
1655                    Arc::strong_count(executor)
1656                );
1657            }
1658            return Ok(executor.clone());
1659        }
1660
1661        let security_policy = SecurityPolicy::new(
1662            self.project_dir.clone(),
1663            self.loaded_config.data_dir.clone(),
1664            self.loaded_config.config.effective_security_config(),
1665        )?;
1666        let harness_policy = crate::harness::select_harness_policy(&self.loaded_config.config);
1667        let profile_name = format!("{:?}", harness_policy.profile).to_lowercase();
1668        let mut executor = ToolExecutor::with_security_policy(
1669            security_policy,
1670            self.runtime_components.security.clone(),
1671        );
1672        executor.set_harness_profile(profile_name);
1673        Self::register_goal_tools_on_executor(
1674            &mut executor,
1675            self.goal_runtime.clone(),
1676            self.loaded_config.config.goals.enabled,
1677        );
1678        executor.register_skill_loader(
1679            self.project_dir.clone(),
1680            self.loaded_config.data_dir.clone(),
1681            self.shared_config.clone(),
1682        );
1683        executor.set_rewind_store(Some(self.rewind_store.clone()));
1684
1685        let workflow_config = self.loaded_config.config.workflow.clone();
1686        let workflow_policy = SecurityPolicy::new(
1687            self.project_dir.clone(),
1688            self.loaded_config.data_dir.clone(),
1689            self.loaded_config.config.security.clone(),
1690        )
1691        .unwrap_or_else(|_| {
1692            SecurityPolicy::new(
1693                self.project_dir.clone(),
1694                self.loaded_config.data_dir.clone(),
1695                crate::config::SecurityConfig::default(),
1696            )
1697            .expect("default security policy")
1698        });
1699        let executor = Arc::new_cyclic(|executor_weak| {
1700            let subagent = SubagentTool::new(
1701                executor_weak.clone(),
1702                self.shared_model_provider.clone(),
1703                self.project_dir.clone(),
1704                self.loaded_config.data_dir.clone(),
1705                self.shared_model_name.clone(),
1706                self.loaded_config.config.harness.clone(),
1707                self.shared_config.clone(),
1708                self.prompt_cache.clone(),
1709                self.runtime_components.clone(),
1710            );
1711            executor.register_tool(Arc::new(subagent));
1712            // BM25 + symbol index — cheap, no nested agent.
1713            executor.register_tool(Arc::new(RepoExploreTool::new(self.project_dir.clone())));
1714            // Production workflow: each agent() is a real nested subagent turn
1715            // with policy-intersected tool allowlists (own max_parallel semaphore).
1716            executor.register_tool(Arc::new(crate::tool::WorkflowTool::with_subagent_bridge(
1717                workflow_policy.clone(),
1718                workflow_config.clone(),
1719                executor_weak.clone(),
1720            )));
1721            executor
1722        });
1723        // If plan mode was entered before the executor existed, re-apply the gate.
1724        if self.agent_mode() == crate::plan_mode::AgentMode::Plan {
1725            let plan_path = crate::plan_store::session_plan_file_path(
1726                &self.loaded_config.data_dir,
1727                self.session.id().as_str(),
1728            );
1729            executor.policy().set_plan_mode(true, Some(plan_path));
1730        }
1731
1732        self.tool_executor = Some(executor.clone());
1733        Ok(executor)
1734    }
1735
1736    fn register_goal_tools_on_executor(
1737        executor: &mut ToolExecutor,
1738        goal_runtime: Arc<GoalRuntimeHandle>,
1739        goals_enabled: bool,
1740    ) {
1741        if !goals_enabled {
1742            return;
1743        }
1744        // Model-facing goal surface: get / create / update only.
1745        // Host checklist mutations stay on the SDK/server API, not the model schema.
1746        executor.register_tool(Arc::new(GetGoalTool::new(goal_runtime.clone())));
1747        executor.register_tool(Arc::new(CreateGoalTool::new(goal_runtime.clone())));
1748        executor.register_tool(Arc::new(UpdateGoalTool::new(goal_runtime)));
1749    }
1750
1751    /// Returns the configured tool executor, if any.
1752    pub fn tool_executor(&self) -> Option<Arc<ToolExecutor>> {
1753        self.tool_executor.clone()
1754    }
1755
1756    /// Keep only tools whose names satisfy `pred` on the live session executor.
1757    ///
1758    /// Used by host tool profiles after start (goal/skill tools may re-register).
1759    /// No-ops when the executor Arc is shared.
1760    pub fn retain_tools<F>(&mut self, pred: F)
1761    where
1762        F: FnMut(&str) -> bool,
1763    {
1764        if let Some(exec) = self.tool_executor.as_mut() {
1765            if let Some(e) = Arc::get_mut(exec) {
1766                e.retain_tools(pred);
1767            } else {
1768                tracing::warn!(
1769                    strong_count = Arc::strong_count(exec),
1770                    "cannot retain tools; tool executor Arc is shared"
1771                );
1772            }
1773        }
1774    }
1775
1776    /// Returns the session-local handle used by the title tool. SDK tooling
1777    /// reloads use this to keep the built-in title capability installed.
1778    pub fn session_title_handle(&self) -> SessionTitleHandle {
1779        self.session_title_handle.clone()
1780    }
1781
1782    /// Replaces the session tool executor (e.g. after installing WASM plugins).
1783    pub fn set_tool_executor(&mut self, executor: Arc<ToolExecutor>) {
1784        self.tool_executor = Some(executor);
1785    }
1786
1787    /// Updates the security configuration used by the tool executor and the
1788    /// runtime's loaded config. No-op when no executor has been created yet.
1789    pub fn set_security_config(&mut self, security: SecurityConfig) -> Result<()> {
1790        self.loaded_config.config.security = security.clone();
1791        self.loaded_config.config.tui.yolo_mode =
1792            matches!(security.permission_mode, PermissionMode::Yolo);
1793
1794        let Some(executor) = self.tool_executor.as_mut() else {
1795            return Ok(());
1796        };
1797        let Some(executor) = Arc::get_mut(executor) else {
1798            return Err(anyhow::anyhow!(
1799                "cannot update security policy while tool executor is shared"
1800            ));
1801        };
1802        let mut policy = executor.security_policy().clone();
1803        policy.set_config(security);
1804        executor.set_security_policy(policy);
1805        Ok(())
1806    }
1807
1808    /// Returns a shared [`MemoryManager`], opening SQLite stores at most once.
1809    fn get_or_init_memory_manager(&self) -> Result<Option<Arc<crate::memory::MemoryManager>>> {
1810        let memory_config = &self.loaded_config.config.memory;
1811        if !memory_config.enabled {
1812            return Ok(None);
1813        }
1814        let mut guard = self
1815            .memory_manager
1816            .lock()
1817            .unwrap_or_else(|e| e.into_inner());
1818        if let Some(manager) = guard.as_ref() {
1819            return Ok(Some(manager.clone()));
1820        }
1821        let manager = Arc::new(crate::memory::MemoryManager::new(
1822            self.project_dir.clone(),
1823            self.loaded_config.data_dir.clone(),
1824            memory_config,
1825        )?);
1826        *guard = Some(manager.clone());
1827        Ok(Some(manager))
1828    }
1829
1830    /// Runs a lightweight auto-memory consolidation (stale detection + dedup).
1831    /// Called on session end. Does not require a model provider.
1832    fn consolidate_auto_memory(&self) -> Result<()> {
1833        let Some(manager) = self.get_or_init_memory_manager()? else {
1834            return Ok(());
1835        };
1836
1837        let db_path = manager.auto_memory.db_path.clone();
1838        if !db_path.exists() {
1839            return Ok(());
1840        }
1841
1842        // Reuse the already-open store rather than reopening the DB.
1843        let report = manager.auto_memory.consolidate(30)?;
1844
1845        if report.marked_stale > 0 || report.duplicates_merged > 0 {
1846            tracing::info!(
1847                "auto-memory consolidation on session end: {} stale, {} duplicates, {} active",
1848                report.marked_stale,
1849                report.duplicates_merged,
1850                report.remaining_active
1851            );
1852        }
1853
1854        Ok(())
1855    }
1856
1857    /// Persist a newly submitted session so it appears in saved-session lists.
1858    /// Naming is deliberately not done here: a second completion would both cost
1859    /// credits and create an unrelated request that harms provider cache affinity.
1860    fn persist_submitted_session(&mut self) {
1861        // Always snapshot on user submit so live sessions appear in list_saved
1862        // while the agent is still working.
1863        if let Err(err) = self.snapshot_session() {
1864            tracing::warn!(error = %err, "early session snapshot after user message failed");
1865        }
1866    }
1867
1868    /// extractMemories: background per-turn memory extraction.
1869    /// Spawns a tokio task that calls the model to extract durable memories
1870    /// from the completed turn. Fire-and-forget — does not block the agent loop.
1871    fn try_extract_memories(&self, _session_id: &str, conversation: &str) {
1872        // Extraction never borrows the interactive chat model. It only runs
1873        // after the user explicitly configures a dedicated background model.
1874        let Some(extraction_model) = &self.memory_extraction_model else {
1875            tracing::debug!("extract-memories skipped: no memory extraction model configured");
1876            return;
1877        };
1878        let provider = extraction_model.provider.clone();
1879        let model_name = extraction_model.model_name.clone();
1880        let conversation = conversation.to_string();
1881
1882        let manager = match self.get_or_init_memory_manager() {
1883            Ok(Some(m)) => m,
1884            Ok(None) => return,
1885            Err(e) => {
1886                tracing::debug!("extract-memories: failed to init memory manager: {}", e);
1887                return;
1888            }
1889        };
1890
1891        let db_path = manager.auto_memory.db_path.clone();
1892        if !db_path.exists() {
1893            return;
1894        }
1895
1896        let store = manager.auto_memory.clone();
1897
1898        // Fire-and-forget
1899        tokio::spawn(async move {
1900            match crate::memory::extract::extract_memories(
1901                &conversation,
1902                provider.as_ref(),
1903                &model_name,
1904                &store,
1905            )
1906            .await
1907            {
1908                Ok(n) => {
1909                    if n > 0 {
1910                        tracing::info!("extract-memories: saved {} memories from turn", n);
1911                    }
1912                }
1913                Err(e) => {
1914                    tracing::debug!("extract-memories failed: {}", e);
1915                }
1916            }
1917        });
1918    }
1919
1920    /// Auto-dream: checks 3 gates after each turn and spawns consolidation if all pass.
1921    /// Fire-and-forget — does not block the agent loop.
1922    fn try_auto_dream(&self) {
1923        let memory_config = &self.loaded_config.config.memory;
1924        let manager = match self.get_or_init_memory_manager() {
1925            Ok(Some(m)) => m,
1926            Ok(None) => return,
1927            Err(e) => {
1928                tracing::debug!("auto-dream: failed to init memory manager: {}", e);
1929                return;
1930            }
1931        };
1932
1933        let interval_hours = memory_config.dream_interval_days * 24;
1934        let state =
1935            crate::memory::auto_dream::AutoDreamState::new(manager.store.memory_root.clone())
1936                .with_interval(interval_hours.max(1));
1937
1938        let history = manager.history.clone();
1939
1940        let last_dream = state.read_last_dream_at();
1941        if !state.should_dream(&history) {
1942            return;
1943        }
1944
1945        let db_path = manager.auto_memory.db_path.clone();
1946        let hours_since = if last_dream > 0 {
1947            let now = std::time::SystemTime::now()
1948                .duration_since(std::time::UNIX_EPOCH)
1949                .map(|d| d.as_secs())
1950                .unwrap_or(0);
1951            (now.saturating_sub(last_dream)) / 3600
1952        } else {
1953            0
1954        };
1955        let sessions_count = history.list_sessions().map(|s| s.len()).unwrap_or(0);
1956
1957        tracing::info!(
1958            "auto-dream triggered: {}h since last, {} sessions",
1959            hours_since,
1960            sessions_count
1961        );
1962
1963        self.event_bus.publish(RuntimeEventKind::AutoDreamStarted {
1964            hours_since_last: hours_since,
1965            sessions_reviewed: sessions_count,
1966        });
1967
1968        let memory_root = manager.store.memory_root.clone();
1969        let event_sender = self.event_bus.sender();
1970        tokio::spawn(async move {
1971            let result = run_auto_dream_consolidation(&db_path).await;
1972
1973            let dream_state = crate::memory::auto_dream::AutoDreamState::new(memory_root);
1974            match result {
1975                Ok(report) => {
1976                    tracing::info!(
1977                        "auto-dream completed: {} stale, {} duplicates, {} active",
1978                        report.marked_stale,
1979                        report.duplicates_merged,
1980                        report.remaining_active
1981                    );
1982                    dream_state.mark_completed();
1983                    let _ = event_sender.send(RuntimeEvent::new(
1984                        RuntimeEventKind::AutoDreamCompleted {
1985                            marked_stale: report.marked_stale,
1986                            duplicates_merged: report.duplicates_merged,
1987                            active_count: report.remaining_active,
1988                        },
1989                    ));
1990                }
1991                Err(e) => {
1992                    tracing::warn!("auto-dream failed: {}", e);
1993                    dream_state.release();
1994                    let _ =
1995                        event_sender.send(RuntimeEvent::new(RuntimeEventKind::AutoDreamFailed {
1996                            reason: e.to_string(),
1997                        }));
1998                }
1999            }
2000        });
2001    }
2002
2003    /// Auto-distill: checks time gate after each turn and spawns distill if enough time passed.
2004    /// Fire-and-forget — does not block the agent loop.
2005    fn try_auto_distill(&self) {
2006        let memory_config = &self.loaded_config.config.memory;
2007        if memory_config.distill_interval_days == 0 {
2008            return;
2009        }
2010
2011        let manager = match self.get_or_init_memory_manager() {
2012            Ok(Some(m)) => m,
2013            _ => return,
2014        };
2015
2016        let interval_hours = memory_config.distill_interval_days * 24;
2017        let state = crate::memory::auto_dream::AutoDreamState::new(
2018            manager.store.memory_root.join("distill"),
2019        )
2020        .with_interval(interval_hours)
2021        .with_min_sessions(3);
2022
2023        let history = manager.history.clone();
2024        if !state.should_dream(&history) {
2025            return;
2026        }
2027
2028        tracing::info!("auto-distill triggered");
2029
2030        let auto_memory = manager.auto_memory.clone();
2031        tokio::spawn(async move {
2032            // Distill only does stale + dedup (no model-based SOP extraction in auto mode)
2033            match auto_memory.consolidate(60) {
2034                Ok(report) => {
2035                    tracing::info!(
2036                        "auto-distill completed: {} stale, {} duplicates, {} active",
2037                        report.marked_stale,
2038                        report.duplicates_merged,
2039                        report.remaining_active
2040                    );
2041                    state.mark_completed();
2042                }
2043                Err(e) => {
2044                    tracing::warn!("auto-distill failed: {}", e);
2045                    state.release();
2046                }
2047            }
2048        });
2049    }
2050
2051    fn build_session_runtime(
2052        &mut self,
2053    ) -> Result<(
2054        crate::session::SessionRuntime,
2055        mpsc::UnboundedReceiver<AgentEvent>,
2056    )> {
2057        let tool_executor = self.ensure_tool_executor()?;
2058        let (event_tx, event_rx) = mpsc::unbounded_channel();
2059        let memory_injection =
2060            self.session_store
2061                .load_memory(&self.project_dir)
2062                .and_then(|memory| {
2063                    memory.format_injection(self.loaded_config.config.memory.max_memory_entries)
2064                });
2065
2066        // Catalog: root skills + pools (metadata only). Bodies via load_skill.
2067        let root = self.load_catalog_skills();
2068        let pools = self.load_catalog_pools();
2069        *self
2070            .shared_available_skills
2071            .lock()
2072            .unwrap_or_else(|e| e.into_inner()) = root.clone();
2073        *self
2074            .shared_skill_pools
2075            .lock()
2076            .unwrap_or_else(|e| e.into_inner()) = pools;
2077        // Soft harness only for session-active skills (not every discovered skill).
2078        let harness_skills = self.session_active_skill_manifests();
2079        *self
2080            .shared_active_skills
2081            .lock()
2082            .unwrap_or_else(|e| e.into_inner()) = harness_skills.clone();
2083        self.apply_harness_packs_for_active(&harness_skills);
2084
2085        let ctx = Arc::new(crate::turn::TurnContext {
2086            model_provider: self.shared_model_provider.clone(),
2087            tool_executor,
2088            project_dir: self.project_dir.clone(),
2089            data_dir: self.loaded_config.data_dir.clone(),
2090            model_name: self.shared_model_name.clone(),
2091            event_tx: Some(event_tx),
2092            approval_resolver: self.approval_resolver(),
2093            question_resolver: self.question_resolver(),
2094            plan_review_resolver: self.plan_review_resolver(),
2095            sudo_password_resolver: self.sudo_password_resolver(),
2096            compact_state: Arc::new(tokio::sync::Mutex::new(crate::compact::CompactState::new(
2097                crate::config::effective_context_window(&self.loaded_config.config),
2098            ))),
2099            harness_config: self.loaded_config.config.harness.clone(),
2100            include_tool_prompt_manifest: crate::config::effective_tool_prompt_manifest(
2101                &self.loaded_config.config,
2102            ),
2103            context_packets: self.shared_context_packets.clone(),
2104            available_skills: self.shared_available_skills.clone(),
2105            skill_pools: self.shared_skill_pools.clone(),
2106            // Prompt path uses available_skills + skill_pools; no instruction bodies.
2107            active_skills: Arc::new(std::sync::Mutex::new(Vec::new())),
2108            prompt_cache: self.prompt_cache.clone(),
2109            instructions: std::sync::Arc::new(std::sync::RwLock::new(None)),
2110            prompt_prefix: std::sync::Arc::new(std::sync::Mutex::new(None)),
2111            components: self.runtime_components.clone(),
2112            cancel_token: self.cancel_token.clone(),
2113            config: self.shared_config.clone(),
2114            memory_injection: memory_injection.clone(),
2115            compaction_provider: None,
2116            agent_mode: self.agent_mode(),
2117            compaction_model_name: None,
2118            session_id: self.session.id().as_str().to_string(),
2119            // Catalog-active skills do not lock tools. Only session-active harness
2120            // packs / harness-flagged skills may contribute an allowlist.
2121            allowed_tool_names: self.harness_allow_tools.clone(),
2122            is_subagent: false,
2123            memory_manager: self.memory_manager.clone(),
2124            harness_card: self.harness_card.clone(),
2125        });
2126
2127        let policy = select_harness_policy(&self.loaded_config.config);
2128        let session_runtime = crate::session::SessionRuntime::spawn(
2129            ctx,
2130            policy,
2131            self.initial_messages.clone(),
2132            memory_injection,
2133        );
2134
2135        Ok((session_runtime, event_rx))
2136    }
2137
2138    fn record_event(&mut self, event: AgentEvent) {
2139        // ── Goal accounting driven by agent events ────────────────
2140        match &event {
2141            AgentEvent::ToolCompleted(_result) => {
2142                // Track if the model used the memory tool with write action during this turn
2143                if _result.ok {
2144                    if let Some(output) = _result.output.as_object() {
2145                        if output.get("tool_name").and_then(|v| v.as_str()) == Some("memory") {
2146                            self.turn_used_memory_write = true;
2147                        }
2148                    }
2149                }
2150                let budget_prompt = self.goal_extension.on_tool_complete();
2151                if budget_prompt.is_some() {
2152                    if let Some(goal) = self.goal_runtime.get_goal() {
2153                        self.event_bus.publish(RuntimeEventKind::GoalUpdated {
2154                            session_id: goal.session_id.clone(),
2155                            goal_id: goal.goal_id.as_str().to_string(),
2156                            objective: goal.objective.clone(),
2157                            short_description: goal.short_description.clone(),
2158                            status: goal.status,
2159                            tokens_used: goal.tokens_used,
2160                            token_budget: goal.token_budget,
2161                        });
2162                    }
2163                }
2164            }
2165            AgentEvent::UsageReported {
2166                input_tokens,
2167                output_tokens,
2168                ..
2169            } => {
2170                let exceeded = self
2171                    .goal_extension
2172                    .on_token_usage(*input_tokens, *output_tokens);
2173                if exceeded {
2174                    if let Some(goal) = self.goal_runtime.get_goal() {
2175                        self.event_bus.publish(RuntimeEventKind::GoalUpdated {
2176                            session_id: goal.session_id.clone(),
2177                            goal_id: goal.goal_id.as_str().to_string(),
2178                            objective: goal.objective.clone(),
2179                            short_description: goal.short_description.clone(),
2180                            status: goal.status,
2181                            tokens_used: goal.tokens_used,
2182                            token_budget: goal.token_budget,
2183                        });
2184                    }
2185                }
2186            }
2187            AgentEvent::SetGoalRequested {
2188                objective,
2189                short_description,
2190                token_budget,
2191            } => {
2192                let goal = self.goal_runtime.set_objective_with_short_description(
2193                    objective.clone(),
2194                    short_description.clone(),
2195                    *token_budget,
2196                );
2197                self.goal_runtime.set_auto_continue(true);
2198                self.publish_goal_updated(&goal);
2199            }
2200            AgentEvent::GoalUpdated {
2201                session_id,
2202                goal_id,
2203                objective,
2204                short_description,
2205                status,
2206                tokens_used,
2207                token_budget,
2208            } => {
2209                // Tool-side goal tools already mutated GoalRuntimeHandle; fan out
2210                // to the runtime event bus so TUI/SDK subscribers stay in sync.
2211                self.event_bus.publish(RuntimeEventKind::GoalUpdated {
2212                    session_id: session_id.clone(),
2213                    goal_id: goal_id.clone(),
2214                    objective: objective.clone(),
2215                    short_description: short_description.clone(),
2216                    status: *status,
2217                    tokens_used: *tokens_used,
2218                    token_budget: *token_budget,
2219                });
2220            }
2221            AgentEvent::ModelDelta { text } => {
2222                // Feed text deltas to plan parser when in Plan mode.
2223                if self.agent_mode() == crate::plan_mode::AgentMode::Plan {
2224                    let plans = self.feed_plan_text(&text);
2225                    for plan in plans {
2226                        self.event_bus.publish(RuntimeEventKind::PlanProposed {
2227                            session_id: self.session.id().as_str().to_string(),
2228                            title: plan.title,
2229                            steps: plan.steps,
2230                        });
2231                    }
2232                }
2233            }
2234            AgentEvent::ModelOutput { text, .. } => {
2235                // Feed final model output to plan parser (catches plans at end of turn).
2236                if self.agent_mode() == crate::plan_mode::AgentMode::Plan {
2237                    let plans = self.feed_plan_text(text);
2238                    for plan in plans {
2239                        self.event_bus.publish(RuntimeEventKind::PlanProposed {
2240                            session_id: self.session.id().as_str().to_string(),
2241                            title: plan.title,
2242                            steps: plan.steps,
2243                        });
2244                    }
2245                    // Also drain any unclosed plans at end of turn.
2246                    let remaining = self.drain_plans();
2247                    for plan in remaining {
2248                        self.event_bus.publish(RuntimeEventKind::PlanProposed {
2249                            session_id: self.session.id().as_str().to_string(),
2250                            title: plan.title,
2251                            steps: plan.steps,
2252                        });
2253                    }
2254                }
2255            }
2256            AgentEvent::PlanProposed { .. } => {}
2257            AgentEvent::PlanReviewRequested(_) | AgentEvent::PlanReviewResolved(_) => {}
2258            AgentEvent::SudoPasswordRequested(_) => {}
2259            AgentEvent::AgentModeChanged { .. } => {}
2260            _ => {}
2261        }
2262
2263        if let Some(tx) = &self.event_tx {
2264            let _ = tx.send(event.clone());
2265        }
2266        // Stream deltas are live-only: keep publishing to subscribers / event bus
2267        // but do not persist them in the session event log (RAM + disk blowup).
2268        let transient = matches!(
2269            event,
2270            AgentEvent::SubagentActivity { .. }
2271                | AgentEvent::SubagentTranscript { .. }
2272                | AgentEvent::ModelDelta { .. }
2273                | AgentEvent::ModelThinkingDelta { .. }
2274                | AgentEvent::ToolCallStreaming { .. }
2275                | AgentEvent::StreamResuming { .. }
2276        );
2277        if let Some(kind) = runtime_event_kind_from_agent_event(&event) {
2278            self.event_bus.publish(kind);
2279        }
2280        if !transient {
2281            self.session.push_event(event);
2282        }
2283    }
2284
2285    fn update_shared_model_state(&self) {
2286        *self
2287            .shared_model_provider
2288            .write()
2289            .unwrap_or_else(|e| e.into_inner()) = self.model_provider.clone();
2290        *self
2291            .shared_model_name
2292            .write()
2293            .unwrap_or_else(|e| e.into_inner()) = provider_request_model_name(
2294            &self.loaded_config.config.model.provider,
2295            &self.loaded_config.config.model.name,
2296        );
2297        *self
2298            .shared_config
2299            .write()
2300            .unwrap_or_else(|e| e.into_inner()) = self.loaded_config.config.clone();
2301    }
2302
2303    /// Root-level catalog skills (metadata only). Pool members are not included.
2304    fn load_catalog_skills(&self) -> Vec<SkillManifest> {
2305        match discover_catalog_entries(
2306            &self.loaded_config.config.skills,
2307            &self.project_dir,
2308            &self.loaded_config.data_dir,
2309        ) {
2310            Ok(catalog) => active_skills(
2311                &catalog.root_skills,
2312                &self.loaded_config.config.skills.active,
2313                &self.active_skills,
2314            ),
2315            Err(err) => {
2316                tracing::warn!(error = %err, "failed to load catalog skills");
2317                Vec::new()
2318            }
2319        }
2320    }
2321
2322    fn load_catalog_pools(&self) -> Vec<SkillPool> {
2323        match discover_catalog_entries(
2324            &self.loaded_config.config.skills,
2325            &self.project_dir,
2326            &self.loaded_config.data_dir,
2327        ) {
2328            Ok(catalog) => catalog.pools,
2329            Err(err) => {
2330                tracing::warn!(error = %err, "failed to load skill pools");
2331                Vec::new()
2332            }
2333        }
2334    }
2335
2336    fn load_active_skills(&self) -> Vec<SkillManifest> {
2337        self.load_catalog_skills()
2338    }
2339
2340    fn load_available_skills(&self) -> Vec<SkillManifest> {
2341        self.load_catalog_skills()
2342    }
2343
2344    fn discover_available_skills(&self) -> Result<Vec<SkillManifest>> {
2345        // Full discovery (root + pool members) for load_skill / harness.
2346        discover_configured_skills(
2347            &self.loaded_config.config.skills,
2348            &self.project_dir,
2349            &self.loaded_config.data_dir,
2350        )
2351    }
2352}
2353
2354fn runtime_event_kind_from_agent_event(event: &AgentEvent) -> Option<RuntimeEventKind> {
2355    match event {
2356        AgentEvent::ModelDelta { text } => {
2357            Some(RuntimeEventKind::AssistantDelta { text: text.clone() })
2358        }
2359        AgentEvent::ModelThinkingDelta { text } => {
2360            Some(RuntimeEventKind::AssistantThinkingDelta { text: text.clone() })
2361        }
2362        AgentEvent::ToolRequested(invocation) => {
2363            Some(RuntimeEventKind::ToolRequested(invocation.clone()))
2364        }
2365        AgentEvent::ToolCompleted(result) => Some(RuntimeEventKind::ToolCompleted(result.clone())),
2366        AgentEvent::SubagentActivity {
2367            invocation_id,
2368            message,
2369        } => Some(RuntimeEventKind::SubagentActivity {
2370            invocation_id: invocation_id.clone(),
2371            message: message.clone(),
2372        }),
2373        AgentEvent::SubagentTranscript {
2374            invocation_id,
2375            item,
2376        } => Some(RuntimeEventKind::SubagentTranscript {
2377            invocation_id: invocation_id.clone(),
2378            item: item.clone(),
2379        }),
2380        AgentEvent::ApprovalRequested(request) => {
2381            Some(RuntimeEventKind::ApprovalRequired(request.clone()))
2382        }
2383        AgentEvent::UsageReported {
2384            input_tokens,
2385            output_tokens,
2386            cache_creation_tokens,
2387            cache_read_tokens,
2388        } => Some(RuntimeEventKind::TokensUpdated {
2389            input_tokens: *input_tokens,
2390            output_tokens: *output_tokens,
2391            cache_creation_tokens: *cache_creation_tokens,
2392            cache_read_tokens: *cache_read_tokens,
2393        }),
2394        AgentEvent::Error { message } => Some(RuntimeEventKind::Error {
2395            message: message.clone(),
2396        }),
2397        AgentEvent::ApprovalResolved(decision) => {
2398            Some(RuntimeEventKind::ApprovalResolved(decision.clone()))
2399        }
2400        AgentEvent::CapabilityRecorded(entry) => {
2401            Some(RuntimeEventKind::CapabilityRecorded(entry.clone()))
2402        }
2403        AgentEvent::QuestionRequested(request) => {
2404            Some(RuntimeEventKind::QuestionRequired(request.clone()))
2405        }
2406        AgentEvent::QuestionResolved(response) => {
2407            Some(RuntimeEventKind::QuestionResolved(response.clone()))
2408        }
2409        AgentEvent::PlanReviewRequested(request) => {
2410            Some(RuntimeEventKind::PlanReviewRequired(request.clone()))
2411        }
2412        AgentEvent::PlanReviewResolved(response) => {
2413            Some(RuntimeEventKind::PlanReviewResolved(response.clone()))
2414        }
2415        AgentEvent::SudoPasswordRequested(request) => {
2416            Some(RuntimeEventKind::SudoPasswordRequired(request.clone()))
2417        }
2418        AgentEvent::HarnessTrace(value) => Some(RuntimeEventKind::HarnessTrace(value.clone())),
2419        AgentEvent::HarnessStopped {
2420            reason,
2421            message,
2422            tool_name,
2423        } => Some(RuntimeEventKind::HarnessStopped {
2424            reason: reason.clone(),
2425            message: message.clone(),
2426            tool_name: tool_name.clone(),
2427        }),
2428        AgentEvent::PatchProposed(patch) => Some(RuntimeEventKind::PatchProposed(patch.clone())),
2429        AgentEvent::MicroCompactApplied { messages_cleared } => {
2430            Some(RuntimeEventKind::MicroCompactApplied {
2431                messages_cleared: *messages_cleared,
2432            })
2433        }
2434        AgentEvent::AutoCompactStarted => Some(RuntimeEventKind::AutoCompactStarted),
2435        AgentEvent::AutoCompactCompleted {
2436            tokens_saved,
2437            summary,
2438            kept_recent_messages,
2439        } => Some(RuntimeEventKind::AutoCompactCompleted {
2440            tokens_saved: *tokens_saved,
2441            summary: summary.clone(),
2442            kept_recent_messages: *kept_recent_messages,
2443        }),
2444        AgentEvent::AutoCompactFailed { reason } => Some(RuntimeEventKind::AutoCompactFailed {
2445            reason: reason.clone(),
2446        }),
2447        AgentEvent::UserTaskSubmitted { .. } | AgentEvent::ModelOutput { .. } => None,
2448        AgentEvent::RepeatedToolCallWarning { .. } => None,
2449        AgentEvent::RepetitionDetected { .. } => None,
2450        AgentEvent::GoalUpdated { .. } => None,
2451        AgentEvent::SetGoalRequested {
2452            objective,
2453            short_description,
2454            token_budget,
2455        } => Some(RuntimeEventKind::SetGoalRequested {
2456            objective: objective.clone(),
2457            short_description: short_description.clone(),
2458            token_budget: *token_budget,
2459        }),
2460        AgentEvent::AutoDreamStarted {
2461            hours_since_last,
2462            sessions_reviewed,
2463        } => Some(RuntimeEventKind::AutoDreamStarted {
2464            hours_since_last: *hours_since_last,
2465            sessions_reviewed: *sessions_reviewed,
2466        }),
2467        AgentEvent::AutoDreamCompleted {
2468            marked_stale,
2469            duplicates_merged,
2470            active_count,
2471        } => Some(RuntimeEventKind::AutoDreamCompleted {
2472            marked_stale: *marked_stale,
2473            duplicates_merged: *duplicates_merged,
2474            active_count: *active_count,
2475        }),
2476        AgentEvent::AutoDreamFailed { reason } => Some(RuntimeEventKind::AutoDreamFailed {
2477            reason: reason.clone(),
2478        }),
2479        AgentEvent::PlanProposed { .. } => None,
2480        // Also mapped above via dedicated arms; keep fallthrough safe.
2481        AgentEvent::AgentModeChanged { .. } => None,
2482        AgentEvent::SessionRecap { .. } => None,
2483        // Transient mid-stream resume signal; live UI only, not a runtime kind.
2484        AgentEvent::StreamResuming { .. } => None,
2485        // Live-only tool-call argument stream; UI progress, not a runtime kind.
2486        AgentEvent::ToolCallStreaming { .. } => None,
2487        // Host-facing UI events (not replayed as runtime kinds).
2488        AgentEvent::NotificationRequested { .. } | AgentEvent::UpdateAvailable { .. } => None,
2489    }
2490}
2491
2492/// Runs the auto-dream consolidation pass (stale + dedup) on the auto-memory SQLite store.
2493/// Called from a tokio::spawn background task — must not access AgentRuntime state.
2494async fn run_auto_dream_consolidation(
2495    db_path: &std::path::Path,
2496) -> anyhow::Result<crate::memory::ConsolidationReport> {
2497    let store = crate::memory::AutoMemoryStore::open(db_path)?;
2498    store.consolidate(30)
2499}