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#[derive(Clone)]
52pub struct ApprovalResolver {
53 pending_approvals: PendingApprovals,
54 runtime_events_tx: broadcast::Sender<RuntimeEvent>,
55}
56
57#[derive(Clone)]
60pub struct QuestionResolver {
61 pending_questions: PendingQuestions,
62 runtime_events_tx: broadcast::Sender<RuntimeEvent>,
63}
64
65#[derive(Clone)]
67pub struct PlanReviewResolver {
68 pending_reviews: PendingPlanReviews,
69 runtime_events_tx: broadcast::Sender<RuntimeEvent>,
70}
71
72#[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 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 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 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 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 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#[derive(Clone)]
272pub struct TurnCanceller {
273 inner: CancelToken,
274}
275
276impl TurnCanceller {
277 pub fn cancel(&self) {
279 self.inner.cancel();
280 }
281
282 pub fn is_cancelled(&self) -> bool {
284 self.inner.is_requested()
285 }
286}
287
288pub struct AgentRuntimeOptions {
290 pub loaded_config: LoadedConfig,
292 pub model_provider: Arc<dyn ModelProvider>,
294 pub project_dir: PathBuf,
296 pub tool_executor: Option<Arc<ToolExecutor>>,
298 pub context_packets: Vec<ContextPacket>,
300 pub active_skills: Vec<String>,
302 pub initial_messages: Vec<ModelMessage>,
304 pub initial_events: Vec<AgentEvent>,
306 pub initial_created_at: Option<u64>,
308 pub initial_updated_at: Option<u64>,
310 pub initial_goal: Option<crate::goal::types::SessionGoal>,
312 pub session_id: Option<SessionId>,
314 pub event_tx: Option<tokio::sync::mpsc::UnboundedSender<AgentEvent>>,
316 pub runtime_components: Option<RuntimeComponents>,
318 pub session_title_handle: Option<SessionTitleHandle>,
321 pub memory_extraction_model: Option<MemoryExtractionModel>,
324}
325
326#[derive(Clone)]
328pub struct MemoryExtractionModel {
329 pub provider: Arc<dyn ModelProvider>,
330 pub model_name: String,
331}
332
333pub 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: Arc<GoalRuntimeHandle>,
361 goal_extension: GoalExtension,
363 turn_used_memory_write: bool,
366 last_user_task: String,
368 session_title_handle: SessionTitleHandle,
370 memory_extraction_model: Option<MemoryExtractionModel>,
372 pending_user_input: std::sync::atomic::AtomicBool,
375 agent_mode: std::sync::RwLock<crate::plan_mode::AgentMode>,
378 plan_parser: std::sync::Mutex<crate::plan_mode::ProposedPlanParser>,
380 memory_manager: Arc<std::sync::Mutex<Option<Arc<crate::memory::MemoryManager>>>>,
382}
383
384impl AgentRuntime {
385 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 pub fn get_goal(&self) -> Option<crate::goal::types::SessionGoal> {
458 self.goal_runtime.get_goal()
459 }
460
461 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 pub fn clear_goal(&self) {
472 self.goal_runtime.clear_goal();
473 }
474
475 pub fn update_goal(&self, goal: crate::goal::types::SessionGoal) {
477 self.goal_runtime.update_goal(goal);
478 }
479
480 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 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 pub fn goal_idle_prompt(&self) -> Option<String> {
499 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 pub fn goals_config(&self) -> crate::config::GoalsConfig {
516 self.loaded_config.config.goals.clone()
517 }
518
519 pub fn goal_runtime(&self) -> &Arc<GoalRuntimeHandle> {
521 &self.goal_runtime
522 }
523
524 pub fn agent_mode(&self) -> crate::plan_mode::AgentMode {
528 *self.agent_mode.read().unwrap_or_else(|e| e.into_inner())
529 }
530
531 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 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 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 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 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 pub fn has_pending_user_input(&self) -> bool {
581 self.pending_user_input
582 .load(std::sync::atomic::Ordering::SeqCst)
583 }
584
585 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 pub fn events(&self) -> &[AgentEvent] {
593 self.session.events()
594 }
595
596 pub fn session_id(&self) -> &SessionId {
598 self.session.id()
599 }
600
601 pub fn session_title(&self) -> Option<&str> {
603 self.session.title()
604 }
605
606 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 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 pub fn context_packets(&self) -> &[ContextPacket] {
628 &self.context_packets
629 }
630
631 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 pub fn list_models(&self) -> Vec<ModelOption> {
644 available_model_options(&self.loaded_config.config)
645 }
646
647 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 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 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 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 pub fn stream_events(&self) -> broadcast::Receiver<RuntimeEvent> {
706 self.event_bus.stream_events()
707 }
708
709 pub fn cancel_turn(&self) {
711 self.turn_canceller().cancel();
712 }
713
714 pub fn resolve_approval(&self, decision: ApprovalDecision) -> bool {
716 self.approval_resolver().resolve(decision)
717 }
718
719 pub fn resolve_question(&self, response: QuestionResponse) -> bool {
721 self.question_resolver().resolve(response)
722 }
723
724 pub fn resolve_plan_review(&self, response: PlanReviewResponse) -> bool {
726 self.plan_review_resolver().resolve(response)
727 }
728
729 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 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 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 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 pub fn turn_canceller(&self) -> TurnCanceller {
766 TurnCanceller {
767 inner: self.cancel_token.clone(),
768 }
769 }
770
771 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 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 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 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 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 if let Some(thinking) = thinking_override {
832 let level_str = thinking.as_config_str();
833 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 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 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 let model_wrote_memory = self.turn_used_memory_write;
942
943 if !model_wrote_memory {
944 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 self.turn_used_memory_write = false;
956
957 self.try_auto_dream();
959
960 self.try_auto_distill();
962 }
963 Err(err) => {
964 self.goal_extension.on_turn_error(&err.to_string());
965 self.flush_partial_model_output_from_events();
968 self.record_event(AgentEvent::Error {
969 message: err.to_string(),
970 });
971 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 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 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 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 self.session.truncate_events_to_user_turns(keep_user_turns);
1049 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 crate::session::truncate_messages_to_user_turns(
1079 &mut self.initial_messages,
1080 keep_user_turns,
1081 );
1082
1083 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 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 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 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 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 pub fn tool_executor(&self) -> Option<Arc<ToolExecutor>> {
1243 self.tool_executor.clone()
1244 }
1245
1246 pub fn session_title_handle(&self) -> SessionTitleHandle {
1249 self.session_title_handle.clone()
1250 }
1251
1252 pub fn set_tool_executor(&mut self, executor: Arc<ToolExecutor>) {
1254 self.tool_executor = Some(executor);
1255 }
1256
1257 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 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 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 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 fn persist_submitted_session(&mut self) {
1331 if let Err(err) = self.snapshot_session() {
1334 tracing::warn!(error = %err, "early session snapshot after user message failed");
1335 }
1336 }
1337
1338 fn try_extract_memories(&self, _session_id: &str, conversation: &str) {
1342 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 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 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 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 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 *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 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 match &event {
1603 AgentEvent::ToolCompleted(_result) => {
1604 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 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 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 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 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 AgentEvent::AgentModeChanged { .. } => None,
1907 AgentEvent::SessionRecap { .. } => None,
1908 AgentEvent::StreamResuming { .. } => None,
1910 AgentEvent::NotificationRequested { .. } | AgentEvent::UpdateAvailable { .. } => None,
1912 }
1913}
1914
1915async 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}