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,
21    UpdateGoalChecklistTool, UpdateGoalTool,
22};
23use crate::harness::select_harness_policy;
24use crate::model::{ModelMessage, ModelProvider, ModelResponse};
25use crate::runtime_components::RuntimeComponents;
26use crate::security::SecurityPolicy;
27use crate::session::{SessionId, SessionStore, current_unix_timestamp};
28use crate::session_title::SessionTitleHandle;
29use crate::skills::{SkillManifest, active_skills, discover_configured_skills};
30use crate::tool::builtin::{RepoExploreTool, SubagentTool};
31use crate::tool::{Tool, ToolExecutor};
32use crate::trace::{TraceStore, turn_traces_from_events};
33use crate::{
34    ModelOption, SessionSnapshot, available_model_options, canonical_provider_id,
35    provider_request_model_name,
36};
37use anyhow::Result;
38
39pub use event_bus::EventBus;
40pub use session_state::SessionState;
41
42type PendingApprovals = Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<ApprovalDecision>>>>;
43type PendingQuestions = Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<QuestionResponse>>>>;
44type PendingPlanReviews =
45    Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<PlanReviewResponse>>>>;
46type PendingSudoPasswords =
47    Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<SudoPasswordResponse>>>>;
48
49/// Resolves pending tool approvals by matching decision ids to waiting
50/// receivers. Cloneable so it can be handed to the UI layer.
51#[derive(Clone)]
52pub struct ApprovalResolver {
53    pending_approvals: PendingApprovals,
54    runtime_events_tx: broadcast::Sender<RuntimeEvent>,
55}
56
57/// Resolves pending interactive questions by matching response ids to waiting
58/// receivers. Cloneable so it can be handed to the UI layer.
59#[derive(Clone)]
60pub struct QuestionResolver {
61    pending_questions: PendingQuestions,
62    runtime_events_tx: broadcast::Sender<RuntimeEvent>,
63}
64
65/// Resolves pending plan reviews (blocks `plan` create until the user finishes).
66#[derive(Clone)]
67pub struct PlanReviewResolver {
68    pending_reviews: PendingPlanReviews,
69    runtime_events_tx: broadcast::Sender<RuntimeEvent>,
70}
71
72/// Resolves sudo password prompts without putting secrets on the event bus.
73#[derive(Clone)]
74pub struct SudoPasswordResolver {
75    pending: PendingSudoPasswords,
76}
77
78impl SudoPasswordResolver {
79    #[cfg(test)]
80    pub fn new_for_test() -> Self {
81        Self {
82            pending: Arc::new(std::sync::Mutex::new(HashMap::new())),
83        }
84    }
85
86    pub(crate) fn new_standalone() -> Self {
87        Self {
88            pending: Arc::new(std::sync::Mutex::new(HashMap::new())),
89        }
90    }
91
92    pub fn register(&self, id: String) -> oneshot::Receiver<SudoPasswordResponse> {
93        let (tx, rx) = oneshot::channel();
94        self.pending
95            .lock()
96            .unwrap_or_else(|e| e.into_inner())
97            .insert(id, tx);
98        rx
99    }
100
101    /// Deliver the user's response. The password (if any) never goes on the runtime event bus.
102    pub fn resolve(&self, response: SudoPasswordResponse) -> bool {
103        let id = response.id().to_string();
104        if let Some(tx) = self
105            .pending
106            .lock()
107            .unwrap_or_else(|e| e.into_inner())
108            .remove(&id)
109        {
110            let _ = tx.send(response);
111            true
112        } else {
113            false
114        }
115    }
116}
117
118impl PlanReviewResolver {
119    #[cfg(test)]
120    pub fn new_for_test() -> Self {
121        let (tx, _) = broadcast::channel(16);
122        Self {
123            pending_reviews: Arc::new(std::sync::Mutex::new(HashMap::new())),
124            runtime_events_tx: tx,
125        }
126    }
127
128    pub(crate) fn new_standalone() -> Self {
129        let (tx, _) = broadcast::channel(16);
130        Self {
131            pending_reviews: Arc::new(std::sync::Mutex::new(HashMap::new())),
132            runtime_events_tx: tx,
133        }
134    }
135
136    pub fn register(&self, id: String) -> oneshot::Receiver<PlanReviewResponse> {
137        let (tx, rx) = oneshot::channel();
138        self.pending_reviews
139            .lock()
140            .unwrap_or_else(|e| e.into_inner())
141            .insert(id, tx);
142        rx
143    }
144
145    pub fn resolve(&self, response: PlanReviewResponse) -> bool {
146        let id = response.id.clone();
147        if let Some(tx) = self
148            .pending_reviews
149            .lock()
150            .unwrap_or_else(|e| e.into_inner())
151            .remove(&id)
152        {
153            let _ = tx.send(response.clone());
154            let _ = self.runtime_events_tx.send(RuntimeEvent::new(
155                RuntimeEventKind::PlanReviewResolved(response),
156            ));
157            true
158        } else {
159            false
160        }
161    }
162}
163
164impl QuestionResolver {
165    #[cfg(test)]
166    pub fn new_for_test() -> Self {
167        let (tx, _) = broadcast::channel(16);
168        Self {
169            pending_questions: Arc::new(std::sync::Mutex::new(HashMap::new())),
170            runtime_events_tx: tx,
171        }
172    }
173
174    pub(crate) fn new_standalone() -> Self {
175        let (tx, _) = broadcast::channel(16);
176        Self {
177            pending_questions: Arc::new(std::sync::Mutex::new(HashMap::new())),
178            runtime_events_tx: tx,
179        }
180    }
181
182    /// Register a pending question, returning the receiver for the response.
183    pub fn register(&self, id: String) -> oneshot::Receiver<QuestionResponse> {
184        let (tx, rx) = oneshot::channel();
185        self.pending_questions
186            .lock()
187            .unwrap_or_else(|e| e.into_inner())
188            .insert(id, tx);
189        rx
190    }
191
192    /// Resolves a pending question by id. Returns `true` if a matching request
193    /// was found and resolved.
194    pub fn resolve(&self, response: QuestionResponse) -> bool {
195        let id = response.id().to_string();
196        if let Some(tx) = self
197            .pending_questions
198            .lock()
199            .unwrap_or_else(|e| e.into_inner())
200            .remove(&id)
201        {
202            let _ = tx.send(response.clone());
203            let _ =
204                self.runtime_events_tx
205                    .send(RuntimeEvent::new(RuntimeEventKind::QuestionResolved(
206                        response,
207                    )));
208            true
209        } else {
210            false
211        }
212    }
213}
214
215impl ApprovalResolver {
216    #[cfg(test)]
217    pub fn new_for_test() -> Self {
218        let (tx, _) = broadcast::channel(16);
219        Self {
220            pending_approvals: Arc::new(std::sync::Mutex::new(HashMap::new())),
221            runtime_events_tx: tx,
222        }
223    }
224
225    pub(crate) fn new_standalone() -> Self {
226        let (tx, _) = broadcast::channel(16);
227        Self {
228            pending_approvals: Arc::new(std::sync::Mutex::new(HashMap::new())),
229            runtime_events_tx: tx,
230        }
231    }
232
233    /// Register a pending approval, returning the receiver for the decision.
234    pub fn register(&self, id: String) -> oneshot::Receiver<ApprovalDecision> {
235        let (tx, rx) = oneshot::channel();
236        self.pending_approvals
237            .lock()
238            .unwrap_or_else(|e| e.into_inner())
239            .insert(id, tx);
240        rx
241    }
242
243    /// Resolves a pending approval by id. Returns `true` if a matching
244    /// pending request was found and resolved.
245    pub fn resolve(&self, decision: ApprovalDecision) -> bool {
246        let id = match &decision {
247            ApprovalDecision::Approved { id } => id,
248            ApprovalDecision::Denied { id } => id,
249        };
250        if let Some(tx) = self
251            .pending_approvals
252            .lock()
253            .unwrap_or_else(|e| e.into_inner())
254            .remove(id)
255        {
256            let _ = tx.send(decision.clone());
257            let _ =
258                self.runtime_events_tx
259                    .send(RuntimeEvent::new(RuntimeEventKind::ApprovalResolved(
260                        decision,
261                    )));
262            true
263        } else {
264            false
265        }
266    }
267}
268
269/// A lightweight handle that cancels the current turn when dropped or called.
270/// Cloneable so it can be handed to the UI layer.
271#[derive(Clone)]
272pub struct TurnCanceller {
273    inner: CancelToken,
274}
275
276impl TurnCanceller {
277    /// Cancels the current turn.
278    pub fn cancel(&self) {
279        self.inner.cancel();
280    }
281
282    /// Returns `true` if cancellation has been requested.
283    pub fn is_cancelled(&self) -> bool {
284        self.inner.is_requested()
285    }
286}
287
288/// Options for constructing an [`AgentRuntime`].
289pub struct AgentRuntimeOptions {
290    /// Loaded and merged configuration.
291    pub loaded_config: LoadedConfig,
292    /// The model provider implementation.
293    pub model_provider: Arc<dyn ModelProvider>,
294    /// Project root directory.
295    pub project_dir: PathBuf,
296    /// Optional custom tool executor (defaults to built-in tools).
297    pub tool_executor: Option<Arc<ToolExecutor>>,
298    /// Context packets to inject into the session.
299    pub context_packets: Vec<ContextPacket>,
300    /// Active skill names for this session.
301    pub active_skills: Vec<String>,
302    /// Seed messages for restoring a session.
303    pub initial_messages: Vec<ModelMessage>,
304    /// Seed events for restoring a persisted session without losing history.
305    pub initial_events: Vec<AgentEvent>,
306    /// Original creation timestamp for restored sessions.
307    pub initial_created_at: Option<u64>,
308    /// Original update timestamp for restored sessions.
309    pub initial_updated_at: Option<u64>,
310    /// Goal restored from a persisted session snapshot.
311    pub initial_goal: Option<crate::goal::types::SessionGoal>,
312    /// Session id for restoring an existing session.
313    pub session_id: Option<SessionId>,
314    /// Optional channel for forwarding agent events outside the runtime.
315    pub event_tx: Option<tokio::sync::mpsc::UnboundedSender<AgentEvent>>,
316    /// Replaceable runtime components. Defaults preserve NAVI's code-agent behavior.
317    pub runtime_components: Option<RuntimeComponents>,
318    /// Session-local title state written by the `set_session_title` tool.
319    /// When omitted, direct runtime users simply do not get agent-managed titles.
320    pub session_title_handle: Option<SessionTitleHandle>,
321    /// Explicit model for automatic durable-memory extraction. `None` means
322    /// extraction is disabled rather than silently billing the chat model.
323    pub memory_extraction_model: Option<MemoryExtractionModel>,
324}
325
326/// Provider selection for the opt-in per-turn memory extractor.
327#[derive(Clone)]
328pub struct MemoryExtractionModel {
329    pub provider: Arc<dyn ModelProvider>,
330    pub model_name: String,
331}
332
333/// The core agent runtime that manages sessions, turns, approvals, and events.
334pub struct AgentRuntime {
335    loaded_config: LoadedConfig,
336    model_provider: Arc<dyn ModelProvider>,
337    shared_model_provider: Arc<RwLock<Arc<dyn ModelProvider>>>,
338    shared_model_name: Arc<RwLock<String>>,
339    shared_config: Arc<RwLock<crate::config::NaviConfig>>,
340    project_dir: PathBuf,
341    tool_executor: Option<Arc<ToolExecutor>>,
342    session_store: SessionStore,
343    context_packets: Vec<ContextPacket>,
344    shared_context_packets: Arc<std::sync::Mutex<Vec<ContextPacket>>>,
345    active_skills: Vec<String>,
346    shared_available_skills: Arc<std::sync::Mutex<Vec<crate::skills::SkillManifest>>>,
347    shared_active_skills: Arc<std::sync::Mutex<Vec<crate::skills::SkillManifest>>>,
348    prompt_cache: Arc<crate::prompt::PromptCache>,
349    runtime_components: RuntimeComponents,
350    initial_messages: Vec<ModelMessage>,
351    event_tx: Option<mpsc::UnboundedSender<AgentEvent>>,
352    cancel_token: CancelToken,
353    pending_approvals: PendingApprovals,
354    pending_questions: PendingQuestions,
355    pending_plan_reviews: PendingPlanReviews,
356    pending_sudo_passwords: PendingSudoPasswords,
357    event_bus: EventBus,
358    session: SessionState,
359    /// Goal runtime handle for the current session.
360    goal_runtime: Arc<GoalRuntimeHandle>,
361    /// Goal extension providing lifecycle hooks.
362    goal_extension: GoalExtension,
363    /// Whether the model used the `memory` tool with `write` action during the current turn.
364    /// Used for mutual exclusion with background extractMemories.
365    turn_used_memory_write: bool,
366    /// Last user task text — used for extractMemories context.
367    last_user_task: String,
368    /// Session title assigned by the chat model through `set_session_title`.
369    session_title_handle: SessionTitleHandle,
370    /// User-selected model for asynchronous automatic memory extraction.
371    memory_extraction_model: Option<MemoryExtractionModel>,
372    /// Whether a user message is pending (set when send_turn is called
373    /// while the runtime is busy with auto-continuation).
374    pending_user_input: std::sync::atomic::AtomicBool,
375    /// Current agent mode (Default or Plan). In Plan mode, only read-only
376    /// tools are available and the model is instructed to propose a plan.
377    agent_mode: std::sync::RwLock<crate::plan_mode::AgentMode>,
378    /// Parser for `<proposed_plan>` tags in streaming text.
379    plan_parser: std::sync::Mutex<crate::plan_mode::ProposedPlanParser>,
380    /// Session-scoped memory manager shared with TurnContext (open once).
381    memory_manager: Arc<std::sync::Mutex<Option<Arc<crate::memory::MemoryManager>>>>,
382}
383
384impl AgentRuntime {
385    /// Creates a new runtime from the given options.
386    pub fn new(options: AgentRuntimeOptions) -> Self {
387        let session_store = SessionStore::with_redaction(
388            options.loaded_config.data_dir.clone(),
389            options
390                .loaded_config
391                .config
392                .security
393                .redact_secrets_in_sessions,
394        );
395
396        let shared_context_packets =
397            Arc::new(std::sync::Mutex::new(options.context_packets.clone()));
398        let shared_available_skills = Arc::new(std::sync::Mutex::new(Vec::new()));
399        let shared_active_skills = Arc::new(std::sync::Mutex::new(Vec::new()));
400        let shared_model_provider = Arc::new(RwLock::new(options.model_provider.clone()));
401        let shared_model_name = Arc::new(RwLock::new(provider_request_model_name(
402            &options.loaded_config.config.model.provider,
403            &options.loaded_config.config.model.name,
404        )));
405        let shared_config = Arc::new(RwLock::new(options.loaded_config.config.clone()));
406        let prompt_cache = Arc::new(crate::prompt::PromptCache::new());
407        let runtime_components = options.runtime_components.unwrap_or_default();
408        let goal_service = Arc::new(GoalService::new());
409        let goal_runtime = Arc::new(GoalRuntimeHandle::new(options.initial_goal.clone()));
410        let goal_extension = GoalExtension::new(goal_service.clone(), goal_runtime.clone());
411
412        Self {
413            loaded_config: options.loaded_config,
414            model_provider: options.model_provider,
415            shared_model_provider,
416            shared_model_name,
417            shared_config,
418            project_dir: options.project_dir,
419            tool_executor: options.tool_executor,
420            session_store,
421            context_packets: options.context_packets,
422            shared_context_packets,
423            active_skills: options.active_skills,
424            shared_available_skills,
425            shared_active_skills,
426            prompt_cache,
427            runtime_components,
428            initial_messages: options.initial_messages,
429            event_tx: options.event_tx,
430            cancel_token: CancelToken::new(),
431            pending_approvals: Arc::new(std::sync::Mutex::new(HashMap::new())),
432            pending_questions: Arc::new(std::sync::Mutex::new(HashMap::new())),
433            pending_plan_reviews: Arc::new(std::sync::Mutex::new(HashMap::new())),
434            pending_sudo_passwords: Arc::new(std::sync::Mutex::new(HashMap::new())),
435            event_bus: EventBus::new(),
436            session: SessionState::new_with_history(
437                options.session_id,
438                options.initial_events,
439                options.initial_created_at,
440                options.initial_updated_at,
441            ),
442            goal_runtime,
443            goal_extension,
444            turn_used_memory_write: false,
445            last_user_task: String::new(),
446            session_title_handle: options.session_title_handle.unwrap_or_default(),
447            memory_extraction_model: options.memory_extraction_model,
448            pending_user_input: std::sync::atomic::AtomicBool::new(false),
449            agent_mode: std::sync::RwLock::new(crate::plan_mode::AgentMode::Default),
450            plan_parser: std::sync::Mutex::new(crate::plan_mode::ProposedPlanParser::new()),
451            memory_manager: Arc::new(std::sync::Mutex::new(None)),
452        }
453    }
454
455    /// Returns all agent events recorded so far.
456    /// Returns the current session goal, if any.
457    pub fn get_goal(&self) -> Option<crate::goal::types::SessionGoal> {
458        self.goal_runtime.get_goal()
459    }
460
461    /// Sets or updates the session goal.
462    pub fn set_goal(
463        &self,
464        objective: String,
465        token_budget: Option<i64>,
466    ) -> crate::goal::types::SessionGoal {
467        self.goal_runtime.set_objective(objective, token_budget)
468    }
469
470    /// Clears the current session goal.
471    pub fn clear_goal(&self) {
472        self.goal_runtime.clear_goal();
473    }
474
475    /// Updates the stored goal (used after status transitions).
476    pub fn update_goal(&self, goal: crate::goal::types::SessionGoal) {
477        self.goal_runtime.update_goal(goal);
478    }
479
480    /// Updates the goal checklist (replaces all tasks).
481    pub fn update_goal_checklist(
482        &self,
483        tasks: Vec<crate::goal::types::GoalTask>,
484    ) -> Option<crate::goal::types::SessionGoal> {
485        self.goal_runtime.update_checklist(tasks)
486    }
487
488    /// Updates a single task's status in the goal checklist.
489    pub fn update_goal_task_status(
490        &self,
491        task_id: usize,
492        status: crate::goal::types::TaskStatus,
493    ) -> Option<crate::goal::types::SessionGoal> {
494        self.goal_runtime.update_task_status(task_id, status)
495    }
496
497    /// Returns a continuation steering prompt if the goal is active and should auto-continue.
498    pub fn goal_idle_prompt(&self) -> Option<String> {
499        // In Plan mode, don't auto-continue — the user needs to confirm the plan first.
500        if self
501            .agent_mode
502            .read()
503            .unwrap_or_else(|e| e.into_inner())
504            .restricts_tools()
505        {
506            return None;
507        }
508        if !self.loaded_config.config.goals.enabled {
509            return None;
510        }
511        self.goal_extension.on_idle()
512    }
513
514    /// Returns the goals configuration.
515    pub fn goals_config(&self) -> crate::config::GoalsConfig {
516        self.loaded_config.config.goals.clone()
517    }
518
519    /// Returns a reference to the goal runtime handle.
520    pub fn goal_runtime(&self) -> &Arc<GoalRuntimeHandle> {
521        &self.goal_runtime
522    }
523
524    // ── Plan Mode ──────────────────────────────────────────────
525
526    /// Returns the current agent mode.
527    pub fn agent_mode(&self) -> crate::plan_mode::AgentMode {
528        *self.agent_mode.read().unwrap_or_else(|e| e.into_inner())
529    }
530
531    /// Enters Plan mode. Only read-only tools will be available, and the
532    /// model is instructed to propose a plan via `<proposed_plan>` tags.
533    pub fn enter_plan_mode(&self) {
534        *self.agent_mode.write().unwrap_or_else(|e| e.into_inner()) =
535            crate::plan_mode::AgentMode::Plan;
536        *self.plan_parser.lock().unwrap_or_else(|e| e.into_inner()) =
537            crate::plan_mode::ProposedPlanParser::new();
538        self.event_bus.publish(RuntimeEventKind::AgentModeChanged {
539            session_id: self.session.id().as_str().to_string(),
540            mode: crate::plan_mode::AgentMode::Plan,
541        });
542    }
543
544    /// Exits Plan mode and returns to normal execution.
545    pub fn exit_plan_mode(&self) {
546        *self.agent_mode.write().unwrap_or_else(|e| e.into_inner()) =
547            crate::plan_mode::AgentMode::Default;
548        self.event_bus.publish(RuntimeEventKind::AgentModeChanged {
549            session_id: self.session.id().as_str().to_string(),
550            mode: crate::plan_mode::AgentMode::Default,
551        });
552    }
553
554    /// Feeds a text delta into the plan parser. Returns any completed plans.
555    pub fn feed_plan_text(&self, text: &str) -> Vec<crate::plan_mode::ProposedPlan> {
556        self.plan_parser
557            .lock()
558            .unwrap_or_else(|e| e.into_inner())
559            .push_text(text)
560    }
561
562    /// Drains any pending plans from the parser (call at end of turn).
563    pub fn drain_plans(&self) -> Vec<crate::plan_mode::ProposedPlan> {
564        self.plan_parser
565            .lock()
566            .unwrap_or_else(|e| e.into_inner())
567            .drain()
568    }
569
570    /// Returns true if the parser is currently inside a `<proposed_plan>` block.
571    pub fn is_parsing_plan(&self) -> bool {
572        self.plan_parser
573            .lock()
574            .unwrap_or_else(|e| e.into_inner())
575            .is_in_plan()
576    }
577
578    /// Returns `true` if a user message is pending (set when `send_turn` is
579    /// called while the runtime is busy with auto-continuation).
580    pub fn has_pending_user_input(&self) -> bool {
581        self.pending_user_input
582            .load(std::sync::atomic::Ordering::SeqCst)
583    }
584
585    /// Marks that a user message is pending.
586    pub fn set_pending_user_input(&self, pending: bool) {
587        self.pending_user_input
588            .store(pending, std::sync::atomic::Ordering::SeqCst);
589    }
590
591    /// Returns all agent events recorded so far.
592    pub fn events(&self) -> &[AgentEvent] {
593        self.session.events()
594    }
595
596    /// Returns the current session id.
597    pub fn session_id(&self) -> &SessionId {
598        self.session.id()
599    }
600
601    /// Returns the session title, if one has been derived.
602    pub fn session_title(&self) -> Option<&str> {
603        self.session.title()
604    }
605
606    /// Adds a context packet to the session and emits a `ContextUpdated` event.
607    pub fn add_context_packet(&mut self, packet: ContextPacket) {
608        self.context_packets.push(packet.clone());
609        self.shared_context_packets
610            .lock()
611            .unwrap_or_else(|e| e.into_inner())
612            .push(packet);
613        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
614    }
615
616    /// Clears all context packets and emits a `ContextUpdated` event.
617    pub fn clear_context_packets(&mut self) {
618        self.context_packets.clear();
619        self.shared_context_packets
620            .lock()
621            .unwrap_or_else(|e| e.into_inner())
622            .clear();
623        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
624    }
625
626    /// Returns the current context packets.
627    pub fn context_packets(&self) -> &[ContextPacket] {
628        &self.context_packets
629    }
630
631    /// Sets the active skills for this session and emits a `ContextUpdated` event.
632    pub fn set_active_skills(&mut self, skills: Vec<String>) {
633        self.active_skills = skills;
634        let manifests = self.load_active_skills();
635        *self
636            .shared_active_skills
637            .lock()
638            .unwrap_or_else(|e| e.into_inner()) = manifests;
639        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
640    }
641
642    /// Lists available model options from the loaded configuration.
643    pub fn list_models(&self) -> Vec<ModelOption> {
644        available_model_options(&self.loaded_config.config)
645    }
646
647    /// Current `(provider_id, model_name)` selection for this runtime.
648    pub fn model_selection(&self) -> (&str, &str) {
649        (
650            self.loaded_config.config.model.provider.as_str(),
651            self.loaded_config.config.model.name.as_str(),
652        )
653    }
654
655    /// Changes the selected model and emits a `ContextUpdated` event.
656    pub fn set_model(&mut self, provider: impl Into<String>, model: impl Into<String>) {
657        self.loaded_config.config.model.provider =
658            canonical_provider_id(&provider.into()).to_string();
659        self.loaded_config.config.model.name = model.into();
660        self.update_shared_model_state();
661        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
662    }
663
664    /// Replaces the runtime configuration and provider used by subsequent turns.
665    pub fn set_model_provider(
666        &mut self,
667        loaded_config: LoadedConfig,
668        model_provider: Arc<dyn ModelProvider>,
669    ) {
670        self.loaded_config = loaded_config;
671        self.model_provider = model_provider;
672        self.update_shared_model_state();
673        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
674    }
675
676    /// Registers a host-provided tool with the runtime's tool executor.
677    /// Creates a default executor if none exists yet.
678    pub fn register_host_tool(&mut self, tool: Arc<dyn Tool>) -> Result<()> {
679        if self.tool_executor.is_none() {
680            let security_policy = SecurityPolicy::new(
681                self.project_dir.clone(),
682                self.loaded_config.data_dir.clone(),
683                self.loaded_config.config.effective_security_config(),
684            )?;
685            self.tool_executor = Some(Arc::new(ToolExecutor::with_security_policy(
686                security_policy,
687                self.runtime_components.security.clone(),
688            )));
689        }
690
691        let Some(executor) = self.tool_executor.as_mut() else {
692            return Err(anyhow::anyhow!("tool executor unavailable"));
693        };
694        let Some(executor) = Arc::get_mut(executor) else {
695            return Err(anyhow::anyhow!(
696                "cannot register host tool while tool executor is shared"
697            ));
698        };
699        executor.register_tool(tool);
700        self.event_bus.publish(RuntimeEventKind::ContextUpdated);
701        Ok(())
702    }
703
704    /// Returns a broadcast receiver for [`RuntimeEvent`]s.
705    pub fn stream_events(&self) -> broadcast::Receiver<RuntimeEvent> {
706        self.event_bus.stream_events()
707    }
708
709    /// Cancels the currently running turn.
710    pub fn cancel_turn(&self) {
711        self.turn_canceller().cancel();
712    }
713
714    /// Resolves a pending approval by id. Returns `true` if found.
715    pub fn resolve_approval(&self, decision: ApprovalDecision) -> bool {
716        self.approval_resolver().resolve(decision)
717    }
718
719    /// Resolves a pending interactive question by id. Returns `true` if found.
720    pub fn resolve_question(&self, response: QuestionResponse) -> bool {
721        self.question_resolver().resolve(response)
722    }
723
724    /// Resolves a pending plan review by invocation id. Returns `true` if found.
725    pub fn resolve_plan_review(&self, response: PlanReviewResponse) -> bool {
726        self.plan_review_resolver().resolve(response)
727    }
728
729    /// Returns an [`ApprovalResolver`] handle for external approval resolution.
730    pub fn approval_resolver(&self) -> ApprovalResolver {
731        ApprovalResolver {
732            pending_approvals: self.pending_approvals.clone(),
733            runtime_events_tx: self.event_bus.sender(),
734        }
735    }
736
737    /// Returns a [`QuestionResolver`] handle for external question resolution.
738    pub fn question_resolver(&self) -> QuestionResolver {
739        QuestionResolver {
740            pending_questions: self.pending_questions.clone(),
741            runtime_events_tx: self.event_bus.sender(),
742        }
743    }
744
745    /// Returns a [`PlanReviewResolver`] handle for external plan-review resolution.
746    pub fn plan_review_resolver(&self) -> PlanReviewResolver {
747        PlanReviewResolver {
748            pending_reviews: self.pending_plan_reviews.clone(),
749            runtime_events_tx: self.event_bus.sender(),
750        }
751    }
752
753    /// Returns a [`SudoPasswordResolver`] for interactive sudo prompts.
754    pub fn sudo_password_resolver(&self) -> SudoPasswordResolver {
755        SudoPasswordResolver {
756            pending: self.pending_sudo_passwords.clone(),
757        }
758    }
759
760    pub fn resolve_sudo_password(&self, response: SudoPasswordResponse) -> bool {
761        self.sudo_password_resolver().resolve(response)
762    }
763
764    /// Returns a [`TurnCanceller`] handle for external cancellation.
765    pub fn turn_canceller(&self) -> TurnCanceller {
766        TurnCanceller {
767            inner: self.cancel_token.clone(),
768        }
769    }
770
771    /// Starts a new session (or restarts if one is already active).
772    /// Returns the session id.
773    pub fn start_session(&mut self) -> Result<SessionId> {
774        if self.session.started() {
775            self.goal_extension
776                .on_session_end(self.session.id().as_str());
777            self.runtime_components
778                .hooks
779                .on_session_end(self.session.id().as_str());
780
781            // Light auto-memory consolidation on session end (stale + dedup, no model needed)
782            let _ = self.consolidate_auto_memory();
783
784            self.event_bus.publish(RuntimeEventKind::SessionFinished {
785                session_id: self.session.id().as_str().to_string(),
786            });
787        }
788        self.cancel_token.reset();
789        self.pending_approvals = Arc::new(std::sync::Mutex::new(HashMap::new()));
790        self.pending_questions = Arc::new(std::sync::Mutex::new(HashMap::new()));
791        self.pending_plan_reviews = Arc::new(std::sync::Mutex::new(HashMap::new()));
792        self.pending_sudo_passwords = Arc::new(std::sync::Mutex::new(HashMap::new()));
793        self.session.start();
794
795        let (session_runtime, event_rx) = self.build_session_runtime()?;
796        self.session.set_runtime(session_runtime, event_rx);
797
798        let id = self.session.id().clone();
799        // Goal lifecycle: session start + register runtime
800        self.goal_extension.on_session_start(id.as_str());
801        self.runtime_components.hooks.on_session_start(id.as_str());
802        self.event_bus.publish(RuntimeEventKind::SessionStarted {
803            session_id: id.as_str().to_string(),
804        });
805
806        Ok(id)
807    }
808
809    /// Sends a user task to the agent and waits for the full response.
810    /// Starts a session automatically if one is not active.
811    /// Sends a user turn with optional multimodal content parts.
812    ///
813    /// When `content_parts` is non-empty, the message is created as a
814    /// multimodal user message containing both text and images.
815    pub async fn send_turn_with_parts(
816        &mut self,
817        task: String,
818        content_parts: Vec<crate::model::ContentPart>,
819        thinking_override: Option<crate::model::ThinkingConfig>,
820    ) -> Result<ModelResponse> {
821        // Mark that user input is being processed.
822        self.pending_user_input
823            .store(false, std::sync::atomic::Ordering::SeqCst);
824
825        if !self.session.started() || self.session.runtime().is_none() {
826            self.start_session()?;
827        }
828
829        // Apply per-turn thinking override before the turn runs so
830        // build_model_request picks it up from the shared config.
831        if let Some(thinking) = thinking_override {
832            let level_str = thinking.as_config_str();
833            // NOTE: we only update shared_config, NOT loaded_config. Mutating
834            // loaded_config would permanently corrupt the original config and
835            // leak the last override into future turns that pass thinking: None.
836            self.shared_config
837                .write()
838                .unwrap_or_else(|e| e.into_inner())
839                .tui
840                .thinking_level = level_str.to_string();
841        }
842
843        let submission_tx = self
844            .session
845            .runtime()
846            .ok_or_else(|| anyhow::anyhow!("session not started"))?
847            .submission_tx
848            .clone();
849
850        let mut event_rx = self
851            .session
852            .take_event_rx()
853            .ok_or_else(|| anyhow::anyhow!("session event stream unavailable"))?;
854
855        self.cancel_token.reset();
856
857        let turn_id = self.session.next_turn_id();
858        tracing::info!(
859            project = %self.project_dir.display(),
860            provider = %self.loaded_config.config.model.provider,
861            model = %self.loaded_config.config.model.name,
862            "agent task submitted"
863        );
864        let session_id = self.session.id().as_str().to_string();
865        self.runtime_components
866            .hooks
867            .on_turn_start(&session_id, &task);
868        self.goal_extension.on_turn_start(&session_id, &task);
869        self.record_event(AgentEvent::UserTaskSubmitted {
870            text: task.clone(),
871            content_parts: content_parts.clone(),
872            submitted_at: Some(crate::session::current_unix_timestamp()),
873        });
874        self.last_user_task = task.clone();
875        // Persist immediately so sidebars list the session while the turn is
876        // still running. The chat model supplies its title through the tool.
877        self.persist_submitted_session();
878        self.event_bus.publish(RuntimeEventKind::TurnStarted {
879            turn_id: turn_id.clone(),
880        });
881
882        let (response_tx, response_rx) = tokio::sync::oneshot::channel();
883        if let Err(e) = submission_tx.send(crate::session::SessionCommand::Turn(
884            crate::session::Submission {
885                task,
886                content_parts,
887                response_tx,
888            },
889        )) {
890            return Err(anyhow::anyhow!("failed to send submission: {}", e));
891        }
892
893        let mut response_rx = response_rx;
894        let result: Result<String> = loop {
895            tokio::select! {
896                res = &mut response_rx => {
897                    break match res {
898                        Ok(Ok(text)) => Ok(text),
899                        Ok(Err(err)) => Err(anyhow::anyhow!(err)),
900                        Err(_) => Err(anyhow::anyhow!("turn cancelled or panicked")),
901                    };
902                }
903                Some(event) = event_rx.recv() => {
904                    self.record_event(event);
905                    // Apply + persist as soon as `set_session_title` completes so
906                    // live UIs and session lists do not wait for the whole turn.
907                    if self.apply_pending_session_title() {
908                        if let Err(err) = self.session.snapshot(
909                            &self.project_dir,
910                            &self.session_store,
911                            &self.event_bus,
912                            self.goal_runtime.get_goal(),
913                        ) {
914                            tracing::debug!(error = %err, "early title snapshot failed");
915                        }
916                    }
917                }
918            }
919        };
920
921        while let Ok(event) = event_rx.try_recv() {
922            self.record_event(event);
923        }
924        drop(event_rx);
925        self.session.set_updated_at(current_unix_timestamp());
926        let _ = self.apply_pending_session_title();
927
928        match &result {
929            Ok(text) => {
930                self.goal_extension.on_turn_end(&session_id);
931                self.runtime_components
932                    .hooks
933                    .on_turn_end(self.session.id().as_str(), text);
934                self.event_bus.publish(RuntimeEventKind::TurnCompleted {
935                    turn_id,
936                    text: text.clone(),
937                });
938
939                // extractMemories: background extraction per turn (fire-and-forget)
940                // Skip if the model already wrote memories during this turn
941                let model_wrote_memory = self.turn_used_memory_write;
942
943                if !model_wrote_memory {
944                    // Build conversation snippet from user task + assistant response
945                    let user_task = self.last_user_task.clone();
946                    let conversation = if user_task.is_empty() {
947                        format!("Assistant: {}", text)
948                    } else {
949                        format!("User: {}\n\nAssistant: {}", user_task, text)
950                    };
951                    self.try_extract_memories(&session_id, &conversation);
952                }
953
954                // Reset per-turn flag
955                self.turn_used_memory_write = false;
956
957                // Auto-dream: fire-and-forget check after each turn
958                self.try_auto_dream();
959
960                // Auto-distill: fire-and-forget check after each turn
961                self.try_auto_distill();
962            }
963            Err(err) => {
964                self.goal_extension.on_turn_error(&err.to_string());
965                // Coalesce streamed deltas into a durable ModelOutput so a mid-turn
966                // failure / crash does not erase what the model already produced.
967                self.flush_partial_model_output_from_events();
968                self.record_event(AgentEvent::Error {
969                    message: err.to_string(),
970                });
971                // Best-effort persist immediately — Desktop/TUI also snapshot after.
972                if let Err(snap_err) = self.snapshot_session() {
973                    tracing::warn!(
974                        error = %snap_err,
975                        "failed to snapshot session after turn error"
976                    );
977                }
978            }
979        }
980
981        result.map(|text| {
982            tracing::info!(chars = text.len(), "agent task completed");
983            ModelResponse { text }
984        })
985    }
986
987    /// If the current turn streamed text/thinking but never emitted `ModelOutput`
988    /// (e.g. provider error mid-stream), write a `ModelOutput` from the deltas so
989    /// session JSON reload keeps the partial answer.
990    fn flush_partial_model_output_from_events(&mut self) {
991        let events = self.session.events();
992        let mut last_user = None;
993        for (i, event) in events.iter().enumerate() {
994            if matches!(event, AgentEvent::UserTaskSubmitted { .. }) {
995                last_user = Some(i);
996            }
997        }
998        let start = last_user.map(|i| i + 1).unwrap_or(0);
999        let mut text = String::new();
1000        let mut thinking = String::new();
1001        let mut saw_output = false;
1002        for event in &events[start..] {
1003            match event {
1004                AgentEvent::ModelOutput { .. } => {
1005                    saw_output = true;
1006                    break;
1007                }
1008                AgentEvent::ModelDelta { text: delta } => text.push_str(delta),
1009                AgentEvent::ModelThinkingDelta { text: delta } => thinking.push_str(delta),
1010                _ => {}
1011            }
1012        }
1013        if saw_output {
1014            return;
1015        }
1016        if text.is_empty() && thinking.is_empty() {
1017            return;
1018        }
1019        self.record_event(AgentEvent::ModelOutput {
1020            text,
1021            thinking: if thinking.is_empty() {
1022                None
1023            } else {
1024                Some(thinking)
1025            },
1026        });
1027    }
1028
1029    pub async fn submit_task(&mut self, task: String) -> Result<ModelResponse> {
1030        self.send_turn_with_parts(task, Vec::new(), None).await
1031    }
1032
1033    /// Sends a plain text user turn (no images).
1034    pub async fn send_turn(&mut self, task: String) -> Result<ModelResponse> {
1035        self.send_turn_with_parts(task, Vec::new(), None).await
1036    }
1037
1038    /// Rewind live conversation history for an edited past user message.
1039    ///
1040    /// Keeps the first `keep_user_turns` user turns (and their assistant/tool
1041    /// follow-ups), drops everything after, and truncates recorded session
1042    /// events the same way. Caller should then `send_turn` with the new text.
1043    ///
1044    /// `keep_user_turns = 0` keeps only system/developer preamble.
1045    pub async fn rewind_to_user_turns(&mut self, keep_user_turns: usize) -> Result<usize> {
1046        if !self.session.started() || self.session.runtime().is_none() {
1047            // Nothing live yet — just reset seed messages/events for a clean start.
1048            self.session.truncate_events_to_user_turns(keep_user_turns);
1049            // Seed messages for next start_session: drop user turns past keep.
1050            crate::session::truncate_messages_to_user_turns(
1051                &mut self.initial_messages,
1052                keep_user_turns,
1053            );
1054            return Ok(self.initial_messages.len());
1055        }
1056
1057        let submission_tx = self
1058            .session
1059            .runtime()
1060            .ok_or_else(|| anyhow::anyhow!("session not started"))?
1061            .submission_tx
1062            .clone();
1063
1064        let (response_tx, response_rx) = tokio::sync::oneshot::channel();
1065        if let Err(e) = submission_tx.send(crate::session::SessionCommand::TruncateToUserTurns {
1066            keep_user_turns,
1067            response_tx,
1068        }) {
1069            return Err(anyhow::anyhow!("failed to send rewind command: {}", e));
1070        }
1071
1072        let remaining = response_rx
1073            .await
1074            .map_err(|_| anyhow::anyhow!("rewind cancelled or session loop exited"))??;
1075
1076        self.session.truncate_events_to_user_turns(keep_user_turns);
1077        // Keep initial_messages in sync if the session is later restarted.
1078        crate::session::truncate_messages_to_user_turns(
1079            &mut self.initial_messages,
1080            keep_user_turns,
1081        );
1082
1083        // Persist truncated history so reload does not resurrect dropped turns.
1084        if let Err(err) = self.snapshot_session() {
1085            tracing::warn!(error = %err, "failed to snapshot session after rewind");
1086        }
1087
1088        Ok(remaining)
1089    }
1090
1091    /// Applies a title set by the current chat model and informs live clients.
1092    ///
1093    /// Returns `true` when a new title was applied so callers can decide whether
1094    /// to persist immediately (mid-turn) or fold the change into a later snapshot.
1095    fn apply_pending_session_title(&mut self) -> bool {
1096        let Some(title) = self.session_title_handle.take() else {
1097            return false;
1098        };
1099        self.session.set_title(Some(title.clone()));
1100        self.event_bus
1101            .publish(RuntimeEventKind::SessionTitleUpdated {
1102                session_id: self.session.id().as_str().to_string(),
1103                title,
1104            });
1105        true
1106    }
1107
1108    /// Creates a [`SessionSnapshot`] of the current session state for persistence.
1109    pub fn snapshot_session(&mut self) -> Result<SessionSnapshot> {
1110        self.session.set_updated_at(current_unix_timestamp());
1111        let _ = self.apply_pending_session_title();
1112        let snapshot = self.session.snapshot(
1113            &self.project_dir,
1114            &self.session_store,
1115            &self.event_bus,
1116            self.goal_runtime.get_goal(),
1117        )?;
1118        self.save_trace_snapshot(snapshot.id.as_str());
1119        Ok(snapshot)
1120    }
1121
1122    /// Creates and persists a session snapshot without blocking the async runtime.
1123    pub async fn snapshot_session_async(&mut self) -> Result<SessionSnapshot> {
1124        self.session.set_updated_at(current_unix_timestamp());
1125        let _ = self.apply_pending_session_title();
1126        let snapshot = self
1127            .session
1128            .snapshot_async(
1129                &self.project_dir,
1130                &self.session_store,
1131                &self.event_bus,
1132                self.goal_runtime.get_goal(),
1133            )
1134            .await?;
1135        self.save_trace_snapshot(snapshot.id.as_str());
1136        Ok(snapshot)
1137    }
1138
1139    fn save_trace_snapshot(&self, session_id: &str) {
1140        let traces = turn_traces_from_events(
1141            session_id,
1142            &self.loaded_config.config.model.provider,
1143            &self.loaded_config.config.model.name,
1144            self.session.events(),
1145        );
1146        if traces.is_empty() {
1147            return;
1148        }
1149        let store = TraceStore::new(&self.loaded_config.data_dir);
1150        if let Err(err) = store.save_session_traces(session_id, &traces) {
1151            tracing::warn!(error = %err, session_id, "failed to save turn traces");
1152        }
1153    }
1154
1155    #[cfg_attr(not(test), allow(dead_code))]
1156    pub(crate) fn session_store(&self) -> &SessionStore {
1157        &self.session_store
1158    }
1159
1160    fn ensure_tool_executor(&mut self) -> Result<Arc<ToolExecutor>> {
1161        if let Some(executor) = self.tool_executor.as_mut() {
1162            if let Some(executor) = Arc::get_mut(executor) {
1163                Self::register_goal_tools_on_executor(
1164                    executor,
1165                    self.goal_runtime.clone(),
1166                    self.loaded_config.config.goals.enabled,
1167                );
1168                executor.register_skill_loader(
1169                    self.project_dir.clone(),
1170                    self.loaded_config.data_dir.clone(),
1171                    self.shared_config.clone(),
1172                );
1173            } else {
1174                tracing::warn!(
1175                    "tool executor Arc is shared; cannot install goal tools in place \
1176                     (strong_count={})",
1177                    Arc::strong_count(executor)
1178                );
1179            }
1180            return Ok(executor.clone());
1181        }
1182
1183        let security_policy = SecurityPolicy::new(
1184            self.project_dir.clone(),
1185            self.loaded_config.data_dir.clone(),
1186            self.loaded_config.config.effective_security_config(),
1187        )?;
1188        let harness_policy = crate::harness::select_harness_policy(&self.loaded_config.config);
1189        let profile_name = format!("{:?}", harness_policy.profile).to_lowercase();
1190        let mut executor = ToolExecutor::with_security_policy(
1191            security_policy,
1192            self.runtime_components.security.clone(),
1193        );
1194        executor.set_harness_profile(profile_name);
1195        Self::register_goal_tools_on_executor(
1196            &mut executor,
1197            self.goal_runtime.clone(),
1198            self.loaded_config.config.goals.enabled,
1199        );
1200        executor.register_skill_loader(
1201            self.project_dir.clone(),
1202            self.loaded_config.data_dir.clone(),
1203            self.shared_config.clone(),
1204        );
1205
1206        let executor = Arc::new_cyclic(|executor_weak| {
1207            let subagent = SubagentTool::new(
1208                executor_weak.clone(),
1209                self.shared_model_provider.clone(),
1210                self.project_dir.clone(),
1211                self.loaded_config.data_dir.clone(),
1212                self.shared_model_name.clone(),
1213                self.loaded_config.config.harness.clone(),
1214                self.shared_config.clone(),
1215                self.prompt_cache.clone(),
1216                self.runtime_components.clone(),
1217            );
1218            executor.register_tool(Arc::new(subagent));
1219            // BM25 + symbol index — cheap, no nested agent.
1220            executor.register_tool(Arc::new(RepoExploreTool::new(self.project_dir.clone())));
1221            executor
1222        });
1223        self.tool_executor = Some(executor.clone());
1224        Ok(executor)
1225    }
1226
1227    fn register_goal_tools_on_executor(
1228        executor: &mut ToolExecutor,
1229        goal_runtime: Arc<GoalRuntimeHandle>,
1230        goals_enabled: bool,
1231    ) {
1232        if !goals_enabled {
1233            return;
1234        }
1235        executor.register_tool(Arc::new(GetGoalTool::new(goal_runtime.clone())));
1236        executor.register_tool(Arc::new(CreateGoalTool::new(goal_runtime.clone())));
1237        executor.register_tool(Arc::new(UpdateGoalTool::new(goal_runtime.clone())));
1238        executor.register_tool(Arc::new(UpdateGoalChecklistTool::new(goal_runtime)));
1239    }
1240
1241    /// Returns the configured tool executor, if any.
1242    pub fn tool_executor(&self) -> Option<Arc<ToolExecutor>> {
1243        self.tool_executor.clone()
1244    }
1245
1246    /// Returns the session-local handle used by the title tool. SDK tooling
1247    /// reloads use this to keep the built-in title capability installed.
1248    pub fn session_title_handle(&self) -> SessionTitleHandle {
1249        self.session_title_handle.clone()
1250    }
1251
1252    /// Replaces the session tool executor (e.g. after installing WASM plugins).
1253    pub fn set_tool_executor(&mut self, executor: Arc<ToolExecutor>) {
1254        self.tool_executor = Some(executor);
1255    }
1256
1257    /// Updates the security configuration used by the tool executor and the
1258    /// runtime's loaded config. No-op when no executor has been created yet.
1259    pub fn set_security_config(&mut self, security: SecurityConfig) -> Result<()> {
1260        self.loaded_config.config.security = security.clone();
1261        self.loaded_config.config.tui.yolo_mode =
1262            matches!(security.permission_mode, PermissionMode::Yolo);
1263
1264        let Some(executor) = self.tool_executor.as_mut() else {
1265            return Ok(());
1266        };
1267        let Some(executor) = Arc::get_mut(executor) else {
1268            return Err(anyhow::anyhow!(
1269                "cannot update security policy while tool executor is shared"
1270            ));
1271        };
1272        let mut policy = executor.security_policy().clone();
1273        policy.set_config(security);
1274        executor.set_security_policy(policy);
1275        Ok(())
1276    }
1277
1278    /// Returns a shared [`MemoryManager`], opening SQLite stores at most once.
1279    fn get_or_init_memory_manager(&self) -> Result<Option<Arc<crate::memory::MemoryManager>>> {
1280        let memory_config = &self.loaded_config.config.memory;
1281        if !memory_config.enabled {
1282            return Ok(None);
1283        }
1284        let mut guard = self
1285            .memory_manager
1286            .lock()
1287            .unwrap_or_else(|e| e.into_inner());
1288        if let Some(manager) = guard.as_ref() {
1289            return Ok(Some(manager.clone()));
1290        }
1291        let manager = Arc::new(crate::memory::MemoryManager::new(
1292            self.project_dir.clone(),
1293            self.loaded_config.data_dir.clone(),
1294            memory_config,
1295        )?);
1296        *guard = Some(manager.clone());
1297        Ok(Some(manager))
1298    }
1299
1300    /// Runs a lightweight auto-memory consolidation (stale detection + dedup).
1301    /// Called on session end. Does not require a model provider.
1302    fn consolidate_auto_memory(&self) -> Result<()> {
1303        let Some(manager) = self.get_or_init_memory_manager()? else {
1304            return Ok(());
1305        };
1306
1307        let db_path = manager.auto_memory.db_path.clone();
1308        if !db_path.exists() {
1309            return Ok(());
1310        }
1311
1312        // Reuse the already-open store rather than reopening the DB.
1313        let report = manager.auto_memory.consolidate(30)?;
1314
1315        if report.marked_stale > 0 || report.duplicates_merged > 0 {
1316            tracing::info!(
1317                "auto-memory consolidation on session end: {} stale, {} duplicates, {} active",
1318                report.marked_stale,
1319                report.duplicates_merged,
1320                report.remaining_active
1321            );
1322        }
1323
1324        Ok(())
1325    }
1326
1327    /// Persist a newly submitted session so it appears in saved-session lists.
1328    /// Naming is deliberately not done here: a second completion would both cost
1329    /// credits and create an unrelated request that harms provider cache affinity.
1330    fn persist_submitted_session(&mut self) {
1331        // Always snapshot on user submit so live sessions appear in list_saved
1332        // while the agent is still working.
1333        if let Err(err) = self.snapshot_session() {
1334            tracing::warn!(error = %err, "early session snapshot after user message failed");
1335        }
1336    }
1337
1338    /// extractMemories: background per-turn memory extraction.
1339    /// Spawns a tokio task that calls the model to extract durable memories
1340    /// from the completed turn. Fire-and-forget — does not block the agent loop.
1341    fn try_extract_memories(&self, _session_id: &str, conversation: &str) {
1342        // Extraction never borrows the interactive chat model. It only runs
1343        // after the user explicitly configures a dedicated background model.
1344        let Some(extraction_model) = &self.memory_extraction_model else {
1345            tracing::debug!("extract-memories skipped: no memory extraction model configured");
1346            return;
1347        };
1348        let provider = extraction_model.provider.clone();
1349        let model_name = extraction_model.model_name.clone();
1350        let conversation = conversation.to_string();
1351
1352        let manager = match self.get_or_init_memory_manager() {
1353            Ok(Some(m)) => m,
1354            Ok(None) => return,
1355            Err(e) => {
1356                tracing::debug!("extract-memories: failed to init memory manager: {}", e);
1357                return;
1358            }
1359        };
1360
1361        let db_path = manager.auto_memory.db_path.clone();
1362        if !db_path.exists() {
1363            return;
1364        }
1365
1366        let store = manager.auto_memory.clone();
1367
1368        // Fire-and-forget
1369        tokio::spawn(async move {
1370            match crate::memory::extract::extract_memories(
1371                &conversation,
1372                provider.as_ref(),
1373                &model_name,
1374                &store,
1375            )
1376            .await
1377            {
1378                Ok(n) => {
1379                    if n > 0 {
1380                        tracing::info!("extract-memories: saved {} memories from turn", n);
1381                    }
1382                }
1383                Err(e) => {
1384                    tracing::debug!("extract-memories failed: {}", e);
1385                }
1386            }
1387        });
1388    }
1389
1390    /// Auto-dream: checks 3 gates after each turn and spawns consolidation if all pass.
1391    /// Fire-and-forget — does not block the agent loop.
1392    fn try_auto_dream(&self) {
1393        let memory_config = &self.loaded_config.config.memory;
1394        let manager = match self.get_or_init_memory_manager() {
1395            Ok(Some(m)) => m,
1396            Ok(None) => return,
1397            Err(e) => {
1398                tracing::debug!("auto-dream: failed to init memory manager: {}", e);
1399                return;
1400            }
1401        };
1402
1403        let interval_hours = memory_config.dream_interval_days * 24;
1404        let state =
1405            crate::memory::auto_dream::AutoDreamState::new(manager.store.memory_root.clone())
1406                .with_interval(interval_hours.max(1));
1407
1408        let history = manager.history.clone();
1409
1410        let last_dream = state.read_last_dream_at();
1411        if !state.should_dream(&history) {
1412            return;
1413        }
1414
1415        let db_path = manager.auto_memory.db_path.clone();
1416        let hours_since = if last_dream > 0 {
1417            let now = std::time::SystemTime::now()
1418                .duration_since(std::time::UNIX_EPOCH)
1419                .map(|d| d.as_secs())
1420                .unwrap_or(0);
1421            (now.saturating_sub(last_dream)) / 3600
1422        } else {
1423            0
1424        };
1425        let sessions_count = history.list_sessions().map(|s| s.len()).unwrap_or(0);
1426
1427        tracing::info!(
1428            "auto-dream triggered: {}h since last, {} sessions",
1429            hours_since,
1430            sessions_count
1431        );
1432
1433        self.event_bus.publish(RuntimeEventKind::AutoDreamStarted {
1434            hours_since_last: hours_since,
1435            sessions_reviewed: sessions_count,
1436        });
1437
1438        let memory_root = manager.store.memory_root.clone();
1439        let event_sender = self.event_bus.sender();
1440        tokio::spawn(async move {
1441            let result = run_auto_dream_consolidation(&db_path).await;
1442
1443            let dream_state = crate::memory::auto_dream::AutoDreamState::new(memory_root);
1444            match result {
1445                Ok(report) => {
1446                    tracing::info!(
1447                        "auto-dream completed: {} stale, {} duplicates, {} active",
1448                        report.marked_stale,
1449                        report.duplicates_merged,
1450                        report.remaining_active
1451                    );
1452                    dream_state.mark_completed();
1453                    let _ = event_sender.send(RuntimeEvent::new(
1454                        RuntimeEventKind::AutoDreamCompleted {
1455                            marked_stale: report.marked_stale,
1456                            duplicates_merged: report.duplicates_merged,
1457                            active_count: report.remaining_active,
1458                        },
1459                    ));
1460                }
1461                Err(e) => {
1462                    tracing::warn!("auto-dream failed: {}", e);
1463                    dream_state.release();
1464                    let _ =
1465                        event_sender.send(RuntimeEvent::new(RuntimeEventKind::AutoDreamFailed {
1466                            reason: e.to_string(),
1467                        }));
1468                }
1469            }
1470        });
1471    }
1472
1473    /// Auto-distill: checks time gate after each turn and spawns distill if enough time passed.
1474    /// Fire-and-forget — does not block the agent loop.
1475    fn try_auto_distill(&self) {
1476        let memory_config = &self.loaded_config.config.memory;
1477        if memory_config.distill_interval_days == 0 {
1478            return;
1479        }
1480
1481        let manager = match self.get_or_init_memory_manager() {
1482            Ok(Some(m)) => m,
1483            _ => return,
1484        };
1485
1486        let interval_hours = memory_config.distill_interval_days * 24;
1487        let state = crate::memory::auto_dream::AutoDreamState::new(
1488            manager.store.memory_root.join("distill"),
1489        )
1490        .with_interval(interval_hours)
1491        .with_min_sessions(3);
1492
1493        let history = manager.history.clone();
1494        if !state.should_dream(&history) {
1495            return;
1496        }
1497
1498        tracing::info!("auto-distill triggered");
1499
1500        let auto_memory = manager.auto_memory.clone();
1501        tokio::spawn(async move {
1502            // Distill only does stale + dedup (no model-based SOP extraction in auto mode)
1503            match auto_memory.consolidate(60) {
1504                Ok(report) => {
1505                    tracing::info!(
1506                        "auto-distill completed: {} stale, {} duplicates, {} active",
1507                        report.marked_stale,
1508                        report.duplicates_merged,
1509                        report.remaining_active
1510                    );
1511                    state.mark_completed();
1512                }
1513                Err(e) => {
1514                    tracing::warn!("auto-distill failed: {}", e);
1515                    state.release();
1516                }
1517            }
1518        });
1519    }
1520
1521    fn build_session_runtime(
1522        &mut self,
1523    ) -> Result<(
1524        crate::session::SessionRuntime,
1525        mpsc::UnboundedReceiver<AgentEvent>,
1526    )> {
1527        let tool_executor = self.ensure_tool_executor()?;
1528        let (event_tx, event_rx) = mpsc::unbounded_channel();
1529        let memory_injection =
1530            self.session_store
1531                .load_memory(&self.project_dir)
1532                .and_then(|memory| {
1533                    memory.format_injection(self.loaded_config.config.memory.max_memory_entries)
1534                });
1535
1536        // Initialize skill snapshots for prompt rendering.
1537        *self
1538            .shared_available_skills
1539            .lock()
1540            .unwrap_or_else(|e| e.into_inner()) = self.load_available_skills();
1541        *self
1542            .shared_active_skills
1543            .lock()
1544            .unwrap_or_else(|e| e.into_inner()) = self.load_active_skills();
1545
1546        let ctx = Arc::new(crate::turn::TurnContext {
1547            model_provider: self.shared_model_provider.clone(),
1548            tool_executor,
1549            project_dir: self.project_dir.clone(),
1550            data_dir: self.loaded_config.data_dir.clone(),
1551            model_name: self.shared_model_name.clone(),
1552            event_tx: Some(event_tx),
1553            approval_resolver: self.approval_resolver(),
1554            question_resolver: self.question_resolver(),
1555            plan_review_resolver: self.plan_review_resolver(),
1556            sudo_password_resolver: self.sudo_password_resolver(),
1557            compact_state: Arc::new(tokio::sync::Mutex::new(crate::compact::CompactState::new(
1558                crate::config::effective_context_window(&self.loaded_config.config),
1559            ))),
1560            harness_config: self.loaded_config.config.harness.clone(),
1561            include_tool_prompt_manifest: crate::config::effective_tool_prompt_manifest(
1562                &self.loaded_config.config,
1563            ),
1564            context_packets: self.shared_context_packets.clone(),
1565            available_skills: self.shared_available_skills.clone(),
1566            active_skills: self.shared_active_skills.clone(),
1567            prompt_cache: self.prompt_cache.clone(),
1568            instructions: std::sync::Arc::new(std::sync::RwLock::new(None)),
1569            prompt_prefix: std::sync::Arc::new(std::sync::Mutex::new(None)),
1570            components: self.runtime_components.clone(),
1571            cancel_token: self.cancel_token.clone(),
1572            config: self.shared_config.clone(),
1573            memory_injection: memory_injection.clone(),
1574            compaction_provider: None,
1575            agent_mode: self.agent_mode(),
1576            compaction_model_name: None,
1577            session_id: self.session.id().as_str().to_string(),
1578            // When active skills declare allow_tools, lock the turn to that set.
1579            allowed_tool_names: {
1580                let active = self
1581                    .shared_active_skills
1582                    .lock()
1583                    .unwrap_or_else(|e| e.into_inner());
1584                crate::skills::skill_tool_allowlist(&active)
1585            },
1586            memory_manager: self.memory_manager.clone(),
1587        });
1588
1589        let policy = select_harness_policy(&self.loaded_config.config);
1590        let session_runtime = crate::session::SessionRuntime::spawn(
1591            ctx,
1592            policy,
1593            self.initial_messages.clone(),
1594            memory_injection,
1595        );
1596
1597        Ok((session_runtime, event_rx))
1598    }
1599
1600    fn record_event(&mut self, event: AgentEvent) {
1601        // ── Goal accounting driven by agent events ────────────────
1602        match &event {
1603            AgentEvent::ToolCompleted(_result) => {
1604                // Track if the model used the memory tool with write action during this turn
1605                if _result.ok {
1606                    if let Some(output) = _result.output.as_object() {
1607                        if output.get("tool_name").and_then(|v| v.as_str()) == Some("memory") {
1608                            self.turn_used_memory_write = true;
1609                        }
1610                    }
1611                }
1612                let budget_prompt = self.goal_extension.on_tool_complete();
1613                if budget_prompt.is_some() {
1614                    if let Some(goal) = self.goal_runtime.get_goal() {
1615                        self.event_bus.publish(RuntimeEventKind::GoalUpdated {
1616                            session_id: goal.session_id.clone(),
1617                            goal_id: goal.goal_id.as_str().to_string(),
1618                            objective: goal.objective.clone(),
1619                            short_description: goal.short_description.clone(),
1620                            status: goal.status,
1621                            tokens_used: goal.tokens_used,
1622                            token_budget: goal.token_budget,
1623                        });
1624                    }
1625                }
1626            }
1627            AgentEvent::UsageReported {
1628                input_tokens,
1629                output_tokens,
1630                ..
1631            } => {
1632                let exceeded = self
1633                    .goal_extension
1634                    .on_token_usage(*input_tokens, *output_tokens);
1635                if exceeded {
1636                    if let Some(goal) = self.goal_runtime.get_goal() {
1637                        self.event_bus.publish(RuntimeEventKind::GoalUpdated {
1638                            session_id: goal.session_id.clone(),
1639                            goal_id: goal.goal_id.as_str().to_string(),
1640                            objective: goal.objective.clone(),
1641                            short_description: goal.short_description.clone(),
1642                            status: goal.status,
1643                            tokens_used: goal.tokens_used,
1644                            token_budget: goal.token_budget,
1645                        });
1646                    }
1647                }
1648            }
1649            AgentEvent::SetGoalRequested {
1650                objective,
1651                short_description,
1652                token_budget,
1653            } => {
1654                let goal = self.goal_runtime.set_objective_with_short_description(
1655                    objective.clone(),
1656                    short_description.clone(),
1657                    *token_budget,
1658                );
1659                self.event_bus.publish(RuntimeEventKind::GoalUpdated {
1660                    session_id: goal.session_id.clone(),
1661                    goal_id: goal.goal_id.as_str().to_string(),
1662                    objective: goal.objective.clone(),
1663                    short_description: goal.short_description.clone(),
1664                    status: goal.status,
1665                    tokens_used: goal.tokens_used,
1666                    token_budget: goal.token_budget,
1667                });
1668            }
1669            AgentEvent::ModelDelta { text } => {
1670                // Feed text deltas to plan parser when in Plan mode.
1671                if self.agent_mode() == crate::plan_mode::AgentMode::Plan {
1672                    let plans = self.feed_plan_text(&text);
1673                    for plan in plans {
1674                        self.event_bus.publish(RuntimeEventKind::PlanProposed {
1675                            session_id: self.session.id().as_str().to_string(),
1676                            title: plan.title,
1677                            steps: plan.steps,
1678                        });
1679                    }
1680                }
1681            }
1682            AgentEvent::ModelOutput { text, .. } => {
1683                // Feed final model output to plan parser (catches plans at end of turn).
1684                if self.agent_mode() == crate::plan_mode::AgentMode::Plan {
1685                    let plans = self.feed_plan_text(text);
1686                    for plan in plans {
1687                        self.event_bus.publish(RuntimeEventKind::PlanProposed {
1688                            session_id: self.session.id().as_str().to_string(),
1689                            title: plan.title,
1690                            steps: plan.steps,
1691                        });
1692                    }
1693                    // Also drain any unclosed plans at end of turn.
1694                    let remaining = self.drain_plans();
1695                    for plan in remaining {
1696                        self.event_bus.publish(RuntimeEventKind::PlanProposed {
1697                            session_id: self.session.id().as_str().to_string(),
1698                            title: plan.title,
1699                            steps: plan.steps,
1700                        });
1701                    }
1702                }
1703            }
1704            AgentEvent::PlanProposed { .. } => {}
1705            AgentEvent::PlanReviewRequested(_) | AgentEvent::PlanReviewResolved(_) => {}
1706            AgentEvent::SudoPasswordRequested(_) => {}
1707            AgentEvent::AgentModeChanged { .. } => {}
1708            _ => {}
1709        }
1710
1711        if let Some(tx) = &self.event_tx {
1712            let _ = tx.send(event.clone());
1713        }
1714        // Stream deltas are live-only: keep publishing to subscribers / event bus
1715        // but do not persist them in the session event log (RAM + disk blowup).
1716        let transient = matches!(
1717            event,
1718            AgentEvent::SubagentActivity { .. }
1719                | AgentEvent::SubagentTranscript { .. }
1720                | AgentEvent::ModelDelta { .. }
1721                | AgentEvent::ModelThinkingDelta { .. }
1722                | AgentEvent::StreamResuming { .. }
1723        );
1724        if let Some(kind) = runtime_event_kind_from_agent_event(&event) {
1725            self.event_bus.publish(kind);
1726        }
1727        if !transient {
1728            self.session.push_event(event);
1729        }
1730    }
1731
1732    fn update_shared_model_state(&self) {
1733        *self
1734            .shared_model_provider
1735            .write()
1736            .unwrap_or_else(|e| e.into_inner()) = self.model_provider.clone();
1737        *self
1738            .shared_model_name
1739            .write()
1740            .unwrap_or_else(|e| e.into_inner()) = provider_request_model_name(
1741            &self.loaded_config.config.model.provider,
1742            &self.loaded_config.config.model.name,
1743        );
1744        *self
1745            .shared_config
1746            .write()
1747            .unwrap_or_else(|e| e.into_inner()) = self.loaded_config.config.clone();
1748    }
1749
1750    fn load_active_skills(&self) -> Vec<SkillManifest> {
1751        match self.discover_available_skills() {
1752            Ok(skills) => active_skills(
1753                &skills,
1754                &self.loaded_config.config.skills.active,
1755                &self.active_skills,
1756            ),
1757            Err(err) => {
1758                tracing::warn!(error = %err, "failed to load configured skills");
1759                Vec::new()
1760            }
1761        }
1762    }
1763
1764    fn load_available_skills(&self) -> Vec<SkillManifest> {
1765        match self.discover_available_skills() {
1766            Ok(skills) => skills,
1767            Err(err) => {
1768                tracing::warn!(error = %err, "failed to discover configured skills");
1769                Vec::new()
1770            }
1771        }
1772    }
1773
1774    fn discover_available_skills(&self) -> Result<Vec<SkillManifest>> {
1775        discover_configured_skills(
1776            &self.loaded_config.config.skills,
1777            &self.project_dir,
1778            &self.loaded_config.data_dir,
1779        )
1780    }
1781}
1782
1783fn runtime_event_kind_from_agent_event(event: &AgentEvent) -> Option<RuntimeEventKind> {
1784    match event {
1785        AgentEvent::ModelDelta { text } => {
1786            Some(RuntimeEventKind::AssistantDelta { text: text.clone() })
1787        }
1788        AgentEvent::ModelThinkingDelta { text } => {
1789            Some(RuntimeEventKind::AssistantThinkingDelta { text: text.clone() })
1790        }
1791        AgentEvent::ToolRequested(invocation) => {
1792            Some(RuntimeEventKind::ToolRequested(invocation.clone()))
1793        }
1794        AgentEvent::ToolCompleted(result) => Some(RuntimeEventKind::ToolCompleted(result.clone())),
1795        AgentEvent::SubagentActivity {
1796            invocation_id,
1797            message,
1798        } => Some(RuntimeEventKind::SubagentActivity {
1799            invocation_id: invocation_id.clone(),
1800            message: message.clone(),
1801        }),
1802        AgentEvent::SubagentTranscript {
1803            invocation_id,
1804            item,
1805        } => Some(RuntimeEventKind::SubagentTranscript {
1806            invocation_id: invocation_id.clone(),
1807            item: item.clone(),
1808        }),
1809        AgentEvent::ApprovalRequested(request) => {
1810            Some(RuntimeEventKind::ApprovalRequired(request.clone()))
1811        }
1812        AgentEvent::UsageReported {
1813            input_tokens,
1814            output_tokens,
1815            cache_creation_tokens,
1816            cache_read_tokens,
1817        } => Some(RuntimeEventKind::TokensUpdated {
1818            input_tokens: *input_tokens,
1819            output_tokens: *output_tokens,
1820            cache_creation_tokens: *cache_creation_tokens,
1821            cache_read_tokens: *cache_read_tokens,
1822        }),
1823        AgentEvent::Error { message } => Some(RuntimeEventKind::Error {
1824            message: message.clone(),
1825        }),
1826        AgentEvent::ApprovalResolved(decision) => {
1827            Some(RuntimeEventKind::ApprovalResolved(decision.clone()))
1828        }
1829        AgentEvent::CapabilityRecorded(entry) => {
1830            Some(RuntimeEventKind::CapabilityRecorded(entry.clone()))
1831        }
1832        AgentEvent::QuestionRequested(request) => {
1833            Some(RuntimeEventKind::QuestionRequired(request.clone()))
1834        }
1835        AgentEvent::QuestionResolved(response) => {
1836            Some(RuntimeEventKind::QuestionResolved(response.clone()))
1837        }
1838        AgentEvent::PlanReviewRequested(request) => {
1839            Some(RuntimeEventKind::PlanReviewRequired(request.clone()))
1840        }
1841        AgentEvent::PlanReviewResolved(response) => {
1842            Some(RuntimeEventKind::PlanReviewResolved(response.clone()))
1843        }
1844        AgentEvent::SudoPasswordRequested(request) => {
1845            Some(RuntimeEventKind::SudoPasswordRequired(request.clone()))
1846        }
1847        AgentEvent::HarnessTrace(value) => Some(RuntimeEventKind::HarnessTrace(value.clone())),
1848        AgentEvent::HarnessStopped {
1849            reason,
1850            message,
1851            tool_name,
1852        } => Some(RuntimeEventKind::HarnessStopped {
1853            reason: reason.clone(),
1854            message: message.clone(),
1855            tool_name: tool_name.clone(),
1856        }),
1857        AgentEvent::PatchProposed(patch) => Some(RuntimeEventKind::PatchProposed(patch.clone())),
1858        AgentEvent::MicroCompactApplied { messages_cleared } => {
1859            Some(RuntimeEventKind::MicroCompactApplied {
1860                messages_cleared: *messages_cleared,
1861            })
1862        }
1863        AgentEvent::AutoCompactStarted => Some(RuntimeEventKind::AutoCompactStarted),
1864        AgentEvent::AutoCompactCompleted { tokens_saved } => {
1865            Some(RuntimeEventKind::AutoCompactCompleted {
1866                tokens_saved: *tokens_saved,
1867            })
1868        }
1869        AgentEvent::AutoCompactFailed { reason } => Some(RuntimeEventKind::AutoCompactFailed {
1870            reason: reason.clone(),
1871        }),
1872        AgentEvent::UserTaskSubmitted { .. } | AgentEvent::ModelOutput { .. } => None,
1873        AgentEvent::RepeatedToolCallWarning { .. } => None,
1874        AgentEvent::RepetitionDetected { .. } => None,
1875        AgentEvent::GoalUpdated { .. } => None,
1876        AgentEvent::SetGoalRequested {
1877            objective,
1878            short_description,
1879            token_budget,
1880        } => Some(RuntimeEventKind::SetGoalRequested {
1881            objective: objective.clone(),
1882            short_description: short_description.clone(),
1883            token_budget: *token_budget,
1884        }),
1885        AgentEvent::AutoDreamStarted {
1886            hours_since_last,
1887            sessions_reviewed,
1888        } => Some(RuntimeEventKind::AutoDreamStarted {
1889            hours_since_last: *hours_since_last,
1890            sessions_reviewed: *sessions_reviewed,
1891        }),
1892        AgentEvent::AutoDreamCompleted {
1893            marked_stale,
1894            duplicates_merged,
1895            active_count,
1896        } => Some(RuntimeEventKind::AutoDreamCompleted {
1897            marked_stale: *marked_stale,
1898            duplicates_merged: *duplicates_merged,
1899            active_count: *active_count,
1900        }),
1901        AgentEvent::AutoDreamFailed { reason } => Some(RuntimeEventKind::AutoDreamFailed {
1902            reason: reason.clone(),
1903        }),
1904        AgentEvent::PlanProposed { .. } => None,
1905        // Also mapped above via dedicated arms; keep fallthrough safe.
1906        AgentEvent::AgentModeChanged { .. } => None,
1907        AgentEvent::SessionRecap { .. } => None,
1908        // Transient mid-stream resume signal; live UI only, not a runtime kind.
1909        AgentEvent::StreamResuming { .. } => None,
1910        // Host-facing UI events (not replayed as runtime kinds).
1911        AgentEvent::NotificationRequested { .. } | AgentEvent::UpdateAvailable { .. } => None,
1912    }
1913}
1914
1915/// Runs the auto-dream consolidation pass (stale + dedup) on the auto-memory SQLite store.
1916/// Called from a tokio::spawn background task — must not access AgentRuntime state.
1917async fn run_auto_dream_consolidation(
1918    db_path: &std::path::Path,
1919) -> anyhow::Result<crate::memory::ConsolidationReport> {
1920    let store = crate::memory::AutoMemoryStore::open(db_path)?;
1921    store.consolidate(30)
1922}