1mod event_bus;
2mod session_state;
3
4#[cfg(test)]
5mod tests;
6
7use std::collections::HashMap;
8use std::path::PathBuf;
9use std::sync::{Arc, RwLock};
10use tokio::sync::{broadcast, mpsc, oneshot};
11
12use crate::cancel::CancelToken;
13use crate::config::{LoadedConfig, PermissionMode, SecurityConfig};
14use crate::context::ContextPacket;
15use crate::event::{
16 AgentEvent, ApprovalDecision, PlanReviewResponse, QuestionResponse, RuntimeEvent,
17 RuntimeEventKind, SudoPasswordResponse,
18};
19use crate::goal::{
20 CreateGoalTool, GetGoalTool, GoalExtension, GoalRuntimeHandle, GoalService, UpdateGoalTool,
21};
22use crate::harness::select_harness_policy;
23use crate::model::{ModelMessage, ModelProvider, ModelResponse};
24use crate::runtime_components::RuntimeComponents;
25use crate::security::SecurityPolicy;
26use crate::session::{SessionId, SessionStore, current_unix_timestamp};
27use crate::session_title::SessionTitleHandle;
28use crate::skills::{
29 SkillManifest, SkillPool, active_skills, discover_catalog_entries, discover_configured_skills,
30};
31use crate::tool::builtin::{RepoExploreTool, SubagentTool};
32use crate::tool::{Tool, ToolExecutor};
33use crate::trace::{TraceStore, turn_traces_from_events};
34use crate::{
35 ModelOption, SessionSnapshot, available_model_options, canonical_provider_id,
36 provider_request_model_name,
37};
38use anyhow::Result;
39
40pub use event_bus::EventBus;
41pub use session_state::SessionState;
42
43type PendingApprovals = Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<ApprovalDecision>>>>;
44type PendingQuestions = Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<QuestionResponse>>>>;
45type PendingPlanReviews =
46 Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<PlanReviewResponse>>>>;
47type PendingSudoPasswords =
48 Arc<std::sync::Mutex<HashMap<String, oneshot::Sender<SudoPasswordResponse>>>>;
49
50#[derive(Clone)]
53pub struct ApprovalResolver {
54 pending_approvals: PendingApprovals,
55 runtime_events_tx: broadcast::Sender<RuntimeEvent>,
56}
57
58#[derive(Clone)]
61pub struct QuestionResolver {
62 pending_questions: PendingQuestions,
63 runtime_events_tx: broadcast::Sender<RuntimeEvent>,
64}
65
66#[derive(Clone)]
68pub struct PlanReviewResolver {
69 pending_reviews: PendingPlanReviews,
70 runtime_events_tx: broadcast::Sender<RuntimeEvent>,
71}
72
73#[derive(Clone)]
75pub struct SudoPasswordResolver {
76 pending: PendingSudoPasswords,
77}
78
79impl SudoPasswordResolver {
80 #[cfg(test)]
81 pub fn new_for_test() -> Self {
82 Self {
83 pending: Arc::new(std::sync::Mutex::new(HashMap::new())),
84 }
85 }
86
87 pub(crate) fn new_standalone() -> Self {
88 Self {
89 pending: Arc::new(std::sync::Mutex::new(HashMap::new())),
90 }
91 }
92
93 pub fn register(&self, id: String) -> oneshot::Receiver<SudoPasswordResponse> {
94 let (tx, rx) = oneshot::channel();
95 self.pending
96 .lock()
97 .unwrap_or_else(|e| e.into_inner())
98 .insert(id, tx);
99 rx
100 }
101
102 pub fn resolve(&self, response: SudoPasswordResponse) -> bool {
104 let id = response.id().to_string();
105 if let Some(tx) = self
106 .pending
107 .lock()
108 .unwrap_or_else(|e| e.into_inner())
109 .remove(&id)
110 {
111 let _ = tx.send(response);
112 true
113 } else {
114 false
115 }
116 }
117}
118
119impl PlanReviewResolver {
120 #[cfg(test)]
121 pub fn new_for_test() -> Self {
122 let (tx, _) = broadcast::channel(16);
123 Self {
124 pending_reviews: Arc::new(std::sync::Mutex::new(HashMap::new())),
125 runtime_events_tx: tx,
126 }
127 }
128
129 pub(crate) fn new_standalone() -> Self {
130 let (tx, _) = broadcast::channel(16);
131 Self {
132 pending_reviews: Arc::new(std::sync::Mutex::new(HashMap::new())),
133 runtime_events_tx: tx,
134 }
135 }
136
137 pub fn register(&self, id: String) -> oneshot::Receiver<PlanReviewResponse> {
138 let (tx, rx) = oneshot::channel();
139 self.pending_reviews
140 .lock()
141 .unwrap_or_else(|e| e.into_inner())
142 .insert(id, tx);
143 rx
144 }
145
146 pub fn resolve(&self, response: PlanReviewResponse) -> bool {
147 let id = response.id.clone();
148 if let Some(tx) = self
149 .pending_reviews
150 .lock()
151 .unwrap_or_else(|e| e.into_inner())
152 .remove(&id)
153 {
154 let _ = tx.send(response.clone());
155 let _ = self.runtime_events_tx.send(RuntimeEvent::new(
156 RuntimeEventKind::PlanReviewResolved(response),
157 ));
158 true
159 } else {
160 false
161 }
162 }
163}
164
165impl QuestionResolver {
166 #[cfg(test)]
167 pub fn new_for_test() -> Self {
168 let (tx, _) = broadcast::channel(16);
169 Self {
170 pending_questions: Arc::new(std::sync::Mutex::new(HashMap::new())),
171 runtime_events_tx: tx,
172 }
173 }
174
175 pub(crate) fn new_standalone() -> Self {
176 let (tx, _) = broadcast::channel(16);
177 Self {
178 pending_questions: Arc::new(std::sync::Mutex::new(HashMap::new())),
179 runtime_events_tx: tx,
180 }
181 }
182
183 pub fn register(&self, id: String) -> oneshot::Receiver<QuestionResponse> {
185 let (tx, rx) = oneshot::channel();
186 self.pending_questions
187 .lock()
188 .unwrap_or_else(|e| e.into_inner())
189 .insert(id, tx);
190 rx
191 }
192
193 pub fn resolve(&self, response: QuestionResponse) -> bool {
196 let id = response.id().to_string();
197 if let Some(tx) = self
198 .pending_questions
199 .lock()
200 .unwrap_or_else(|e| e.into_inner())
201 .remove(&id)
202 {
203 let _ = tx.send(response.clone());
204 let _ =
205 self.runtime_events_tx
206 .send(RuntimeEvent::new(RuntimeEventKind::QuestionResolved(
207 response,
208 )));
209 true
210 } else {
211 false
212 }
213 }
214}
215
216impl ApprovalResolver {
217 #[cfg(test)]
218 pub fn new_for_test() -> Self {
219 let (tx, _) = broadcast::channel(16);
220 Self {
221 pending_approvals: Arc::new(std::sync::Mutex::new(HashMap::new())),
222 runtime_events_tx: tx,
223 }
224 }
225
226 pub(crate) fn new_standalone() -> Self {
227 let (tx, _) = broadcast::channel(16);
228 Self {
229 pending_approvals: Arc::new(std::sync::Mutex::new(HashMap::new())),
230 runtime_events_tx: tx,
231 }
232 }
233
234 pub fn register(&self, id: String) -> oneshot::Receiver<ApprovalDecision> {
236 let (tx, rx) = oneshot::channel();
237 self.pending_approvals
238 .lock()
239 .unwrap_or_else(|e| e.into_inner())
240 .insert(id, tx);
241 rx
242 }
243
244 pub fn resolve(&self, decision: ApprovalDecision) -> bool {
247 let id = match &decision {
248 ApprovalDecision::Approved { id } => id,
249 ApprovalDecision::Denied { id } => id,
250 };
251 if let Some(tx) = self
252 .pending_approvals
253 .lock()
254 .unwrap_or_else(|e| e.into_inner())
255 .remove(id)
256 {
257 let _ = tx.send(decision.clone());
258 let _ =
259 self.runtime_events_tx
260 .send(RuntimeEvent::new(RuntimeEventKind::ApprovalResolved(
261 decision,
262 )));
263 true
264 } else {
265 false
266 }
267 }
268}
269
270#[derive(Clone)]
273pub struct TurnCanceller {
274 inner: CancelToken,
275}
276
277impl TurnCanceller {
278 pub fn cancel(&self) {
280 self.inner.cancel();
281 }
282
283 pub fn is_cancelled(&self) -> bool {
285 self.inner.is_requested()
286 }
287}
288
289pub struct AgentRuntimeOptions {
291 pub loaded_config: LoadedConfig,
293 pub model_provider: Arc<dyn ModelProvider>,
295 pub project_dir: PathBuf,
297 pub tool_executor: Option<Arc<ToolExecutor>>,
299 pub context_packets: Vec<ContextPacket>,
301 pub active_skills: Vec<String>,
303 pub initial_messages: Vec<ModelMessage>,
305 pub initial_events: Vec<AgentEvent>,
307 pub initial_created_at: Option<u64>,
309 pub initial_updated_at: Option<u64>,
311 pub initial_goal: Option<crate::goal::types::SessionGoal>,
313 pub session_id: Option<SessionId>,
315 pub event_tx: Option<tokio::sync::mpsc::UnboundedSender<AgentEvent>>,
317 pub runtime_components: Option<RuntimeComponents>,
319 pub session_title_handle: Option<SessionTitleHandle>,
322 pub memory_extraction_model: Option<MemoryExtractionModel>,
325 pub skip_auto_tool_bootstrap: bool,
329}
330
331#[derive(Clone)]
333pub struct MemoryExtractionModel {
334 pub provider: Arc<dyn ModelProvider>,
335 pub model_name: String,
336}
337
338pub struct AgentRuntime {
340 loaded_config: LoadedConfig,
341 model_provider: Arc<dyn ModelProvider>,
342 shared_model_provider: Arc<RwLock<Arc<dyn ModelProvider>>>,
343 shared_model_name: Arc<RwLock<String>>,
344 shared_config: Arc<RwLock<crate::config::NaviConfig>>,
345 project_dir: PathBuf,
346 tool_executor: Option<Arc<ToolExecutor>>,
347 session_store: SessionStore,
348 context_packets: Vec<ContextPacket>,
349 shared_context_packets: Arc<std::sync::Mutex<Vec<ContextPacket>>>,
350 active_skills: Vec<String>,
351 shared_available_skills: Arc<std::sync::Mutex<Vec<crate::skills::SkillManifest>>>,
352 shared_skill_pools: Arc<std::sync::Mutex<Vec<crate::skills::SkillPool>>>,
353 shared_active_skills: Arc<std::sync::Mutex<Vec<crate::skills::SkillManifest>>>,
354 prompt_cache: Arc<crate::prompt::PromptCache>,
355 runtime_components: RuntimeComponents,
356 initial_messages: Vec<ModelMessage>,
357 event_tx: Option<mpsc::UnboundedSender<AgentEvent>>,
358 cancel_token: CancelToken,
359 pending_approvals: PendingApprovals,
360 pending_questions: PendingQuestions,
361 pending_plan_reviews: PendingPlanReviews,
362 pending_sudo_passwords: PendingSudoPasswords,
363 event_bus: EventBus,
364 session: SessionState,
365 goal_runtime: Arc<GoalRuntimeHandle>,
367 goal_extension: GoalExtension,
369 turn_used_memory_write: bool,
372 last_user_task: String,
374 session_title_handle: SessionTitleHandle,
376 memory_extraction_model: Option<MemoryExtractionModel>,
378 skip_auto_tool_bootstrap: bool,
380 pending_user_input: std::sync::atomic::AtomicBool,
383 agent_mode: std::sync::RwLock<crate::plan_mode::AgentMode>,
386 plan_parser: std::sync::Mutex<crate::plan_mode::ProposedPlanParser>,
388 memory_manager: Arc<std::sync::Mutex<Option<Arc<crate::memory::MemoryManager>>>>,
390 rewind_store: Arc<std::sync::Mutex<crate::rewind::RewindStore>>,
392 harness_max_auto_continue: Option<u32>,
394 harness_token_budget: Option<i64>,
396 harness_card: Option<String>,
398 harness_allow_tools: Option<Vec<String>>,
400}
401
402impl AgentRuntime {
403 pub fn new(options: AgentRuntimeOptions) -> Self {
405 let session_store = SessionStore::with_redaction(
406 options.loaded_config.data_dir.clone(),
407 options
408 .loaded_config
409 .config
410 .security
411 .redact_secrets_in_sessions,
412 );
413
414 let shared_context_packets =
415 Arc::new(std::sync::Mutex::new(options.context_packets.clone()));
416 let shared_available_skills = Arc::new(std::sync::Mutex::new(Vec::new()));
417 let shared_skill_pools = Arc::new(std::sync::Mutex::new(Vec::new()));
418 let shared_active_skills = Arc::new(std::sync::Mutex::new(Vec::new()));
419 let shared_model_provider = Arc::new(RwLock::new(options.model_provider.clone()));
420 let shared_model_name = Arc::new(RwLock::new(provider_request_model_name(
421 &options.loaded_config.config.model.provider,
422 &options.loaded_config.config.model.name,
423 )));
424 let shared_config = Arc::new(RwLock::new(options.loaded_config.config.clone()));
425 let prompt_cache = Arc::new(crate::prompt::PromptCache::new());
426 let runtime_components = options.runtime_components.unwrap_or_default();
427 let goal_service = Arc::new(GoalService::new());
428 let goal_runtime = Arc::new(GoalRuntimeHandle::new(options.initial_goal.clone()));
429 let goal_extension = GoalExtension::new(goal_service.clone(), goal_runtime.clone());
430
431 let rewind_sid = options
432 .session_id
433 .as_ref()
434 .map(|s| s.as_str().to_string())
435 .unwrap_or_else(|| "pending".to_string());
436 let rewind_store = Arc::new(std::sync::Mutex::new(crate::rewind::RewindStore::new(
437 &options.loaded_config.data_dir,
438 &rewind_sid,
439 &options.project_dir,
440 )));
441
442 Self {
443 loaded_config: options.loaded_config,
444 model_provider: options.model_provider,
445 shared_model_provider,
446 shared_model_name,
447 shared_config,
448 project_dir: options.project_dir,
449 tool_executor: options.tool_executor,
450 session_store,
451 context_packets: options.context_packets,
452 shared_context_packets,
453 active_skills: options.active_skills,
454 shared_available_skills,
455 shared_skill_pools,
456 shared_active_skills,
457 prompt_cache,
458 runtime_components,
459 initial_messages: options.initial_messages,
460 event_tx: options.event_tx,
461 cancel_token: CancelToken::new(),
462 pending_approvals: Arc::new(std::sync::Mutex::new(HashMap::new())),
463 pending_questions: Arc::new(std::sync::Mutex::new(HashMap::new())),
464 pending_plan_reviews: Arc::new(std::sync::Mutex::new(HashMap::new())),
465 pending_sudo_passwords: Arc::new(std::sync::Mutex::new(HashMap::new())),
466 event_bus: EventBus::new(),
467 session: SessionState::new_with_history(
468 options.session_id,
469 options.initial_events,
470 options.initial_created_at,
471 options.initial_updated_at,
472 ),
473 goal_runtime,
474 goal_extension,
475 turn_used_memory_write: false,
476 last_user_task: String::new(),
477 session_title_handle: options.session_title_handle.unwrap_or_default(),
478 memory_extraction_model: options.memory_extraction_model,
479 skip_auto_tool_bootstrap: options.skip_auto_tool_bootstrap,
480 pending_user_input: std::sync::atomic::AtomicBool::new(false),
481 agent_mode: std::sync::RwLock::new(crate::plan_mode::AgentMode::Default),
482 plan_parser: std::sync::Mutex::new(crate::plan_mode::ProposedPlanParser::new()),
483 memory_manager: Arc::new(std::sync::Mutex::new(None)),
484 rewind_store,
485 harness_max_auto_continue: None,
486 harness_token_budget: None,
487 harness_card: None,
488 harness_allow_tools: None,
489 }
490 }
491
492 pub fn rewind_store_handle(&self) -> Arc<std::sync::Mutex<crate::rewind::RewindStore>> {
494 self.rewind_store.clone()
495 }
496
497 pub fn rebind_rewind_store(&self) {
499 let sid = self.session.id().as_str().to_string();
500 let mut store = self.rewind_store.lock().unwrap_or_else(|e| e.into_inner());
501 *store =
502 crate::rewind::RewindStore::new(&self.loaded_config.data_dir, &sid, &self.project_dir);
503 }
504
505 pub fn get_goal(&self) -> Option<crate::goal::types::SessionGoal> {
508 self.goal_runtime.get_goal()
509 }
510
511 pub fn set_goal(
513 &self,
514 objective: String,
515 token_budget: Option<i64>,
516 ) -> crate::goal::types::SessionGoal {
517 let goal = self.goal_runtime.set_objective(objective, token_budget);
518 self.goal_runtime.set_auto_continue(true);
519 self.publish_goal_updated(&goal);
520 goal
521 }
522
523 pub fn set_goal_with_short_description(
525 &self,
526 objective: String,
527 short_description: Option<String>,
528 token_budget: Option<i64>,
529 ) -> crate::goal::types::SessionGoal {
530 let goal = self.goal_runtime.set_objective_with_short_description(
531 objective,
532 short_description,
533 token_budget,
534 );
535 self.goal_runtime.set_auto_continue(true);
536 self.publish_goal_updated(&goal);
537 goal
538 }
539
540 pub fn clear_goal(&self) {
546 self.goal_runtime.clear_goal();
547 }
548
549 pub fn update_goal(&self, goal: crate::goal::types::SessionGoal) {
551 self.goal_runtime.update_goal(goal.clone());
552 self.publish_goal_updated(&goal);
553 }
554
555 fn publish_goal_updated(&self, goal: &crate::goal::types::SessionGoal) {
556 self.event_bus.publish(RuntimeEventKind::GoalUpdated {
557 session_id: goal.session_id.clone(),
558 goal_id: goal.goal_id.as_str().to_string(),
559 objective: goal.objective.clone(),
560 short_description: goal.short_description.clone(),
561 status: goal.status,
562 tokens_used: goal.tokens_used,
563 token_budget: goal.token_budget,
564 });
565 }
566
567 pub fn update_goal_checklist(
569 &self,
570 tasks: Vec<crate::goal::types::GoalTask>,
571 ) -> Option<crate::goal::types::SessionGoal> {
572 self.goal_runtime.update_checklist(tasks)
573 }
574
575 pub fn update_goal_task_status(
577 &self,
578 task_id: usize,
579 status: crate::goal::types::TaskStatus,
580 ) -> Option<crate::goal::types::SessionGoal> {
581 self.goal_runtime.update_task_status(task_id, status)
582 }
583
584 pub fn goal_idle_prompt(&self) -> Option<String> {
586 if self
588 .agent_mode
589 .read()
590 .unwrap_or_else(|e| e.into_inner())
591 .restricts_tools()
592 {
593 return None;
594 }
595 if !self.loaded_config.config.goals.enabled {
596 return None;
597 }
598 self.goal_extension.on_idle()
599 }
600
601 pub fn goals_config(&self) -> crate::config::GoalsConfig {
603 self.loaded_config.config.goals.clone()
604 }
605
606 pub fn goal_runtime(&self) -> &Arc<GoalRuntimeHandle> {
608 &self.goal_runtime
609 }
610
611 pub fn agent_mode(&self) -> crate::plan_mode::AgentMode {
615 *self.agent_mode.read().unwrap_or_else(|e| e.into_inner())
616 }
617
618 pub fn enter_plan_mode(&self) {
621 *self.agent_mode.write().unwrap_or_else(|e| e.into_inner()) =
622 crate::plan_mode::AgentMode::Plan;
623 *self.plan_parser.lock().unwrap_or_else(|e| e.into_inner()) =
624 crate::plan_mode::ProposedPlanParser::new();
625
626 let plan_path = crate::plan_store::session_plan_file_path(
627 &self.loaded_config.data_dir,
628 self.session.id().as_str(),
629 );
630 if let Some(parent) = plan_path.parent() {
631 let _ = std::fs::create_dir_all(parent);
632 }
633 if let Some(exec) = self.tool_executor.as_ref() {
634 exec.policy().set_plan_mode(true, Some(plan_path.clone()));
635 }
636
637 self.event_bus.publish(RuntimeEventKind::AgentModeChanged {
638 session_id: self.session.id().as_str().to_string(),
639 mode: crate::plan_mode::AgentMode::Plan,
640 });
641 }
642
643 pub fn exit_plan_mode(&self) {
645 *self.agent_mode.write().unwrap_or_else(|e| e.into_inner()) =
646 crate::plan_mode::AgentMode::Default;
647 if let Some(exec) = self.tool_executor.as_ref() {
648 exec.policy().set_plan_mode(false, None);
649 }
650 self.event_bus.publish(RuntimeEventKind::AgentModeChanged {
651 session_id: self.session.id().as_str().to_string(),
652 mode: crate::plan_mode::AgentMode::Default,
653 });
654 }
655
656 pub fn feed_plan_text(&self, text: &str) -> Vec<crate::plan_mode::ProposedPlan> {
658 self.plan_parser
659 .lock()
660 .unwrap_or_else(|e| e.into_inner())
661 .push_text(text)
662 }
663
664 pub fn drain_plans(&self) -> Vec<crate::plan_mode::ProposedPlan> {
666 self.plan_parser
667 .lock()
668 .unwrap_or_else(|e| e.into_inner())
669 .drain()
670 }
671
672 pub fn is_parsing_plan(&self) -> bool {
674 self.plan_parser
675 .lock()
676 .unwrap_or_else(|e| e.into_inner())
677 .is_in_plan()
678 }
679
680 pub fn has_pending_user_input(&self) -> bool {
683 self.pending_user_input
684 .load(std::sync::atomic::Ordering::SeqCst)
685 }
686
687 pub fn set_pending_user_input(&self, pending: bool) {
689 self.pending_user_input
690 .store(pending, std::sync::atomic::Ordering::SeqCst);
691 }
692
693 pub fn events(&self) -> &[AgentEvent] {
695 self.session.events()
696 }
697
698 pub fn session_id(&self) -> &SessionId {
700 self.session.id()
701 }
702
703 pub fn session_title(&self) -> Option<&str> {
705 self.session.title()
706 }
707
708 pub fn add_context_packet(&mut self, packet: ContextPacket) {
710 self.context_packets.push(packet.clone());
711 self.shared_context_packets
712 .lock()
713 .unwrap_or_else(|e| e.into_inner())
714 .push(packet);
715 self.event_bus.publish(RuntimeEventKind::ContextUpdated);
716 }
717
718 pub fn clear_context_packets(&mut self) {
720 self.context_packets.clear();
721 self.shared_context_packets
722 .lock()
723 .unwrap_or_else(|e| e.into_inner())
724 .clear();
725 self.event_bus.publish(RuntimeEventKind::ContextUpdated);
726 }
727
728 pub fn context_packets(&self) -> &[ContextPacket] {
730 &self.context_packets
731 }
732
733 pub fn set_active_skills(&mut self, skills: Vec<String>) {
742 self.active_skills = skills;
743 let catalog = self.load_catalog_skills();
744 let pools = self.load_catalog_pools();
745 *self
746 .shared_available_skills
747 .lock()
748 .unwrap_or_else(|e| e.into_inner()) = catalog.clone();
749 *self
750 .shared_skill_pools
751 .lock()
752 .unwrap_or_else(|e| e.into_inner()) = pools;
753 let harness_skills = self.session_active_skill_manifests();
755 *self
756 .shared_active_skills
757 .lock()
758 .unwrap_or_else(|e| e.into_inner()) = harness_skills.clone();
759 self.apply_harness_packs_for_active(&harness_skills);
760 self.event_bus.publish(RuntimeEventKind::ContextUpdated);
761 }
762
763 fn session_active_skill_manifests(&self) -> Vec<crate::skills::SkillManifest> {
767 let available = self.discover_available_skills().unwrap_or_default();
768 let session = &self.active_skills;
769 let configured = &self.loaded_config.config.skills.active;
770 if session.is_empty() && configured.is_empty() {
773 return Vec::new();
774 }
775 crate::skills::active_skills(&available, configured, session)
776 }
777
778 fn apply_harness_packs_for_active(&mut self, manifests: &[crate::skills::SkillManifest]) {
780 let applied =
781 crate::harness_pack::apply_harness_for_skills(&self.loaded_config.data_dir, manifests);
782 self.harness_max_auto_continue = applied.max_auto_continue_turns;
783 self.harness_token_budget = applied.token_budget;
784 self.harness_card = if applied.harness_card.is_empty() {
785 None
786 } else {
787 Some(applied.harness_card)
788 };
789 self.harness_allow_tools = applied.allow_tools;
790 if let Some(budget) = applied.token_budget {
793 if let Some(mut goal) = self.goal_runtime.get_goal() {
794 if goal.token_budget.is_none() && goal.status.should_auto_continue() {
795 goal.token_budget = Some(budget);
796 self.goal_runtime.update_goal(goal);
797 }
798 }
799 }
800 }
801
802 pub fn effective_max_auto_continue_turns(&self) -> u32 {
804 let cfg = self.goals_config().max_auto_continue_turns;
805 match self.harness_max_auto_continue {
806 Some(h) if h > 0 => {
807 if cfg > 0 {
808 cfg.min(h)
809 } else {
810 h
811 }
812 }
813 _ => cfg,
814 }
815 }
816
817 pub fn harness_card(&self) -> Option<&str> {
819 self.harness_card.as_deref()
820 }
821
822 pub fn list_models(&self) -> Vec<ModelOption> {
824 available_model_options(&self.loaded_config.config)
825 }
826
827 pub fn model_selection(&self) -> (&str, &str) {
829 (
830 self.loaded_config.config.model.provider.as_str(),
831 self.loaded_config.config.model.name.as_str(),
832 )
833 }
834
835 pub fn set_model(&mut self, provider: impl Into<String>, model: impl Into<String>) {
837 self.loaded_config.config.model.provider =
838 canonical_provider_id(&provider.into()).to_string();
839 self.loaded_config.config.model.name = model.into();
840 self.update_shared_model_state();
841 self.event_bus.publish(RuntimeEventKind::ContextUpdated);
842 }
843
844 pub fn set_model_provider(
846 &mut self,
847 loaded_config: LoadedConfig,
848 model_provider: Arc<dyn ModelProvider>,
849 ) {
850 self.loaded_config = loaded_config;
851 self.model_provider = model_provider;
852 self.update_shared_model_state();
853 self.event_bus.publish(RuntimeEventKind::ContextUpdated);
854 }
855
856 pub fn register_host_tool(&mut self, tool: Arc<dyn Tool>) -> Result<()> {
859 if self.tool_executor.is_none() {
860 let security_policy = SecurityPolicy::new(
861 self.project_dir.clone(),
862 self.loaded_config.data_dir.clone(),
863 self.loaded_config.config.effective_security_config(),
864 )?;
865 self.tool_executor = Some(Arc::new(ToolExecutor::with_security_policy(
866 security_policy,
867 self.runtime_components.security.clone(),
868 )));
869 }
870
871 let Some(executor) = self.tool_executor.as_mut() else {
872 return Err(anyhow::anyhow!("tool executor unavailable"));
873 };
874 let Some(executor) = Arc::get_mut(executor) else {
875 return Err(anyhow::anyhow!(
876 "cannot register host tool while tool executor is shared"
877 ));
878 };
879 executor.register_tool(tool);
880 self.event_bus.publish(RuntimeEventKind::ContextUpdated);
881 Ok(())
882 }
883
884 pub fn stream_events(&self) -> broadcast::Receiver<RuntimeEvent> {
886 self.event_bus.stream_events()
887 }
888
889 pub fn cancel_turn(&self) {
891 self.turn_canceller().cancel();
892 }
893
894 pub fn resolve_approval(&self, decision: ApprovalDecision) -> bool {
896 self.approval_resolver().resolve(decision)
897 }
898
899 pub fn resolve_question(&self, response: QuestionResponse) -> bool {
901 self.question_resolver().resolve(response)
902 }
903
904 pub fn resolve_plan_review(&self, response: PlanReviewResponse) -> bool {
906 self.plan_review_resolver().resolve(response)
907 }
908
909 pub fn approval_resolver(&self) -> ApprovalResolver {
911 ApprovalResolver {
912 pending_approvals: self.pending_approvals.clone(),
913 runtime_events_tx: self.event_bus.sender(),
914 }
915 }
916
917 pub fn question_resolver(&self) -> QuestionResolver {
919 QuestionResolver {
920 pending_questions: self.pending_questions.clone(),
921 runtime_events_tx: self.event_bus.sender(),
922 }
923 }
924
925 pub fn plan_review_resolver(&self) -> PlanReviewResolver {
927 PlanReviewResolver {
928 pending_reviews: self.pending_plan_reviews.clone(),
929 runtime_events_tx: self.event_bus.sender(),
930 }
931 }
932
933 pub fn sudo_password_resolver(&self) -> SudoPasswordResolver {
935 SudoPasswordResolver {
936 pending: self.pending_sudo_passwords.clone(),
937 }
938 }
939
940 pub fn resolve_sudo_password(&self, response: SudoPasswordResponse) -> bool {
941 self.sudo_password_resolver().resolve(response)
942 }
943
944 pub fn turn_canceller(&self) -> TurnCanceller {
946 TurnCanceller {
947 inner: self.cancel_token.clone(),
948 }
949 }
950
951 pub fn start_session(&mut self) -> Result<SessionId> {
954 if self.session.started() {
955 self.goal_extension
956 .on_session_end(self.session.id().as_str());
957 self.runtime_components
958 .hooks
959 .on_session_end(self.session.id().as_str());
960
961 let _ = self.consolidate_auto_memory();
963
964 self.event_bus.publish(RuntimeEventKind::SessionFinished {
965 session_id: self.session.id().as_str().to_string(),
966 });
967 }
968 self.cancel_token.reset();
969 self.pending_approvals = Arc::new(std::sync::Mutex::new(HashMap::new()));
970 self.pending_questions = Arc::new(std::sync::Mutex::new(HashMap::new()));
971 self.pending_plan_reviews = Arc::new(std::sync::Mutex::new(HashMap::new()));
972 self.pending_sudo_passwords = Arc::new(std::sync::Mutex::new(HashMap::new()));
973 self.session.start();
974
975 let (session_runtime, event_rx) = self.build_session_runtime()?;
976 self.session.set_runtime(session_runtime, event_rx);
977
978 let id = self.session.id().clone();
979 self.rebind_rewind_store();
981 if let Ok(_) = self.ensure_tool_executor() {
982 if let Some(stored) = self.tool_executor.as_mut() {
983 if let Some(exec) = Arc::get_mut(stored) {
984 exec.set_rewind_store(Some(self.rewind_store.clone()));
985 }
986 }
987 }
988 self.goal_extension.on_session_start(id.as_str());
990 self.runtime_components.hooks.on_session_start(id.as_str());
991 self.event_bus.publish(RuntimeEventKind::SessionStarted {
992 session_id: id.as_str().to_string(),
993 });
994
995 Ok(id)
996 }
997
998 pub async fn send_turn_with_parts(
1009 &mut self,
1010 task: String,
1011 content_parts: Vec<crate::model::ContentPart>,
1012 thinking_override: Option<crate::model::ThinkingConfig>,
1013 ) -> Result<ModelResponse> {
1014 let mut response = self
1015 .send_turn_once(task, content_parts, thinking_override)
1016 .await?;
1017
1018 let goals_config = self.goals_config();
1019 let max_auto = self.effective_max_auto_continue_turns();
1020 let mut auto_continue_count = 0u32;
1021 loop {
1022 if !goals_config.enabled {
1023 break;
1024 }
1025 if max_auto > 0 && auto_continue_count >= max_auto {
1026 self.goal_runtime
1027 .record_blocked_turn("auto-continuation limit reached");
1028 if let Some(goal) = self.goal_runtime.get_goal() {
1029 self.publish_goal_updated(&goal);
1030 }
1031 break;
1032 }
1033 if self.has_pending_user_input() {
1034 break;
1035 }
1036 let Some(prompt) = self.goal_idle_prompt() else {
1037 break;
1038 };
1039 auto_continue_count += 1;
1040 tracing::info!(
1041 auto_continue_count,
1042 "starting automatic goal continuation turn"
1043 );
1044 response = self.send_turn_once(prompt, Vec::new(), None).await?;
1047 }
1048
1049 Ok(response)
1050 }
1051
1052 async fn send_turn_once(
1054 &mut self,
1055 task: String,
1056 content_parts: Vec<crate::model::ContentPart>,
1057 thinking_override: Option<crate::model::ThinkingConfig>,
1058 ) -> Result<ModelResponse> {
1059 self.pending_user_input
1061 .store(false, std::sync::atomic::Ordering::SeqCst);
1062
1063 if !self.session.started() || self.session.runtime().is_none() {
1064 self.start_session()?;
1065 }
1066
1067 if let Some(thinking) = thinking_override {
1070 let level_str = thinking.as_config_str();
1071 self.shared_config
1075 .write()
1076 .unwrap_or_else(|e| e.into_inner())
1077 .tui
1078 .thinking_level = level_str.to_string();
1079 }
1080
1081 let submission_tx = self
1082 .session
1083 .runtime()
1084 .ok_or_else(|| anyhow::anyhow!("session not started"))?
1085 .submission_tx
1086 .clone();
1087
1088 let mut event_rx = self
1089 .session
1090 .take_event_rx()
1091 .ok_or_else(|| anyhow::anyhow!("session event stream unavailable"))?;
1092
1093 self.cancel_token.reset();
1094
1095 let turn_id = self.session.next_turn_id();
1096 tracing::info!(
1097 project = %self.project_dir.display(),
1098 provider = %self.loaded_config.config.model.provider,
1099 model = %self.loaded_config.config.model.name,
1100 "agent task submitted"
1101 );
1102 let session_id = self.session.id().as_str().to_string();
1103 self.runtime_components
1104 .hooks
1105 .on_turn_start(&session_id, &task);
1106 self.goal_extension.on_turn_start(&session_id, &task);
1107 self.record_event(AgentEvent::UserTaskSubmitted {
1108 text: task.clone(),
1109 content_parts: content_parts.clone(),
1110 submitted_at: Some(crate::session::current_unix_timestamp()),
1111 });
1112 self.last_user_task = task.clone();
1113
1114 let prompt_index = self
1116 .session
1117 .events()
1118 .iter()
1119 .filter(|e| matches!(e, AgentEvent::UserTaskSubmitted { .. }))
1120 .count()
1121 .saturating_sub(1);
1122 let captured_index = {
1123 let mut store = self.rewind_store.lock().unwrap_or_else(|e| e.into_inner());
1124 if store.root().to_string_lossy().contains("pending") {
1126 *store = crate::rewind::RewindStore::new(
1127 &self.loaded_config.data_dir,
1128 self.session.id().as_str(),
1129 &self.project_dir,
1130 );
1131 }
1132 store.clear_turn_written();
1134 match store.capture_point(
1135 prompt_index,
1136 &task,
1137 crate::session::current_unix_timestamp(),
1138 ) {
1139 Ok(p) => Some(p.prompt_index),
1140 Err(e) => {
1141 tracing::warn!(error = %e, "rewind: failed to capture pre-turn snapshot");
1142 None
1143 }
1144 }
1145 };
1146
1147 self.persist_submitted_session();
1150 self.event_bus.publish(RuntimeEventKind::TurnStarted {
1151 turn_id: turn_id.clone(),
1152 });
1153
1154 let (response_tx, response_rx) = tokio::sync::oneshot::channel();
1155 if let Err(e) = submission_tx.send(crate::session::SessionCommand::Turn(
1156 crate::session::Submission {
1157 task,
1158 content_parts,
1159 response_tx,
1160 },
1161 )) {
1162 return Err(anyhow::anyhow!("failed to send submission: {}", e));
1163 }
1164
1165 let mut response_rx = response_rx;
1166 let result: Result<String> = loop {
1167 tokio::select! {
1168 res = &mut response_rx => {
1169 break match res {
1170 Ok(Ok(text)) => Ok(text),
1171 Ok(Err(err)) => Err(anyhow::anyhow!("turn failed: {err}")),
1172 Err(_) => Err(anyhow::anyhow!("turn cancelled or panicked")),
1173 };
1174 }
1175 Some(event) = event_rx.recv() => {
1176 self.record_event(event);
1177 if self.apply_pending_session_title() {
1180 if let Err(err) = self.session.snapshot(
1181 &self.project_dir,
1182 &self.session_store,
1183 &self.event_bus,
1184 self.goal_runtime.get_goal(),
1185 ) {
1186 tracing::debug!(error = %err, "early title snapshot failed");
1187 }
1188 }
1189 }
1190 }
1191 };
1192
1193 while let Ok(event) = event_rx.try_recv() {
1194 self.record_event(event);
1195 }
1196 drop(event_rx);
1197 self.session.set_updated_at(current_unix_timestamp());
1198 let _ = self.apply_pending_session_title();
1199
1200 if let Some(idx) = captured_index {
1202 let mut store = self.rewind_store.lock().unwrap_or_else(|e| e.into_inner());
1203 if let Err(e) = store.finalize_turn_created_paths(idx) {
1204 tracing::debug!(error = %e, prompt_index = idx, "rewind: finalize created_paths failed");
1205 }
1206 }
1207
1208 match &result {
1209 Ok(text) => {
1210 self.goal_extension.on_turn_end(&session_id);
1211 self.runtime_components
1212 .hooks
1213 .on_turn_end(self.session.id().as_str(), text);
1214 self.event_bus.publish(RuntimeEventKind::TurnCompleted {
1215 turn_id,
1216 text: text.clone(),
1217 });
1218
1219 let model_wrote_memory = self.turn_used_memory_write;
1222
1223 if !model_wrote_memory {
1224 let user_task = self.last_user_task.clone();
1226 let conversation = if user_task.is_empty() {
1227 format!("Assistant: {}", text)
1228 } else {
1229 format!("User: {}\n\nAssistant: {}", user_task, text)
1230 };
1231 self.try_extract_memories(&session_id, &conversation);
1232 }
1233
1234 self.turn_used_memory_write = false;
1236
1237 self.try_auto_dream();
1239
1240 self.try_auto_distill();
1242 }
1243 Err(err) => {
1244 self.goal_extension.on_turn_error(&err.to_string());
1245 self.flush_partial_model_output_from_events();
1248 self.record_event(AgentEvent::Error {
1249 message: err.to_string(),
1250 });
1251 if let Err(snap_err) = self.snapshot_session() {
1253 tracing::warn!(
1254 error = %snap_err,
1255 "failed to snapshot session after turn error"
1256 );
1257 }
1258 }
1259 }
1260
1261 result.map(|text| {
1262 tracing::info!(chars = text.len(), "agent task completed");
1263 ModelResponse { text }
1264 })
1265 }
1266
1267 fn flush_partial_model_output_from_events(&mut self) {
1271 let events = self.session.events();
1272 let mut last_user = None;
1273 for (i, event) in events.iter().enumerate() {
1274 if matches!(event, AgentEvent::UserTaskSubmitted { .. }) {
1275 last_user = Some(i);
1276 }
1277 }
1278 let start = last_user.map(|i| i + 1).unwrap_or(0);
1279 let mut text = String::new();
1280 let mut thinking = String::new();
1281 let mut saw_output = false;
1282 for event in &events[start..] {
1283 match event {
1284 AgentEvent::ModelOutput { .. } => {
1285 saw_output = true;
1286 break;
1287 }
1288 AgentEvent::ModelDelta { text: delta } => text.push_str(delta),
1289 AgentEvent::ModelThinkingDelta { text: delta } => thinking.push_str(delta),
1290 _ => {}
1291 }
1292 }
1293 if saw_output {
1294 return;
1295 }
1296 if text.is_empty() && thinking.is_empty() {
1297 return;
1298 }
1299 self.record_event(AgentEvent::ModelOutput {
1300 text,
1301 thinking: if thinking.is_empty() {
1302 None
1303 } else {
1304 Some(thinking)
1305 },
1306 });
1307 }
1308
1309 pub async fn submit_task(&mut self, task: String) -> Result<ModelResponse> {
1310 self.send_turn_with_parts(task, Vec::new(), None).await
1311 }
1312
1313 pub async fn send_turn(&mut self, task: String) -> Result<ModelResponse> {
1315 self.send_turn_with_parts(task, Vec::new(), None).await
1316 }
1317
1318 pub async fn compact_now(&mut self) -> Result<crate::compact::CompactOutcome> {
1323 if !self.session.started() || self.session.runtime().is_none() {
1324 self.record_event(AgentEvent::AutoCompactStarted);
1326 let provider = self
1327 .shared_model_provider
1328 .read()
1329 .unwrap_or_else(|e| e.into_inner())
1330 .clone();
1331 let model = self
1332 .shared_model_name
1333 .read()
1334 .unwrap_or_else(|e| e.into_inner())
1335 .clone();
1336 let harness = self.loaded_config.config.harness.clone();
1337 let mut state = crate::compact::CompactState::new(
1338 crate::config::effective_context_window(&self.loaded_config.config),
1339 );
1340 match state
1341 .force_compact(
1342 &mut self.initial_messages,
1343 provider.as_ref(),
1344 &model,
1345 &harness,
1346 )
1347 .await
1348 {
1349 Ok(Some(outcome)) => {
1350 self.collapse_events_after_compact(&outcome);
1351 self.event_bus
1354 .publish(RuntimeEventKind::AutoCompactCompleted {
1355 tokens_saved: outcome.tokens_saved,
1356 summary: outcome.summary.clone(),
1357 kept_recent_messages: outcome.kept_recent_messages,
1358 });
1359 Ok(outcome)
1360 }
1361 Ok(None) => {
1362 let reason = "nothing to compact".to_string();
1363 self.record_event(AgentEvent::AutoCompactFailed {
1364 reason: reason.clone(),
1365 });
1366 Err(anyhow::anyhow!(reason))
1367 }
1368 Err(e) => {
1369 self.record_event(AgentEvent::AutoCompactFailed {
1370 reason: e.to_string(),
1371 });
1372 Err(e)
1373 }
1374 }
1375 } else {
1376 let submission_tx = self
1377 .session
1378 .runtime()
1379 .ok_or_else(|| anyhow::anyhow!("session not started"))?
1380 .submission_tx
1381 .clone();
1382
1383 let mut event_rx = self
1386 .session
1387 .take_event_rx()
1388 .ok_or_else(|| anyhow::anyhow!("session event stream unavailable"))?;
1389
1390 let (response_tx, response_rx) = tokio::sync::oneshot::channel();
1391 if let Err(e) =
1392 submission_tx.send(crate::session::SessionCommand::Compact { response_tx })
1393 {
1394 return Err(anyhow::anyhow!("failed to send compact command: {}", e));
1395 }
1396
1397 let mut response_rx = response_rx;
1398 let result: Result<crate::compact::CompactOutcome> = loop {
1399 tokio::select! {
1400 res = &mut response_rx => {
1401 break match res {
1402 Ok(Ok(outcome)) => Ok(outcome),
1403 Ok(Err(err)) => Err(anyhow::anyhow!("compact failed: {err}")),
1404 Err(_) => Err(anyhow::anyhow!("compact cancelled or session loop exited")),
1405 };
1406 }
1407 Some(event) = event_rx.recv() => {
1408 self.record_event(event);
1409 }
1410 }
1411 };
1412
1413 while let Ok(event) = event_rx.try_recv() {
1414 self.record_event(event);
1415 }
1416 if let Ok(ref outcome) = result {
1419 self.collapse_events_after_compact(outcome);
1420 self.sync_initial_messages_from_compact(outcome);
1421 if let Err(err) = self.snapshot_session() {
1422 tracing::warn!(error = %err, "failed to snapshot session after compact");
1423 }
1424 }
1425
1426 result
1427 }
1428 }
1429
1430 fn collapse_events_after_compact(&mut self, outcome: &crate::compact::CompactOutcome) {
1431 let summary_event = AgentEvent::UserTaskSubmitted {
1434 text: format!(
1435 "Here is a summary of the conversation so far:\n\n{}",
1436 outcome.summary
1437 ),
1438 content_parts: Vec::new(),
1439 submitted_at: Some(crate::session::current_unix_timestamp()),
1440 };
1441 let completed = AgentEvent::AutoCompactCompleted {
1442 tokens_saved: outcome.tokens_saved,
1443 summary: outcome.summary.clone(),
1444 kept_recent_messages: 0,
1446 };
1447 self.session.replace_events(vec![summary_event, completed]);
1448 }
1449
1450 fn sync_initial_messages_from_compact(&mut self, outcome: &crate::compact::CompactOutcome) {
1451 let prefix: Vec<_> = self
1453 .initial_messages
1454 .iter()
1455 .filter(|m| {
1456 matches!(
1457 m.role,
1458 crate::model::ModelRole::System | crate::model::ModelRole::Developer
1459 )
1460 })
1461 .cloned()
1462 .collect();
1463 self.initial_messages = prefix;
1464 self.initial_messages
1465 .push(crate::compact::compact_summary_user_message(
1466 &outcome.summary,
1467 ));
1468 }
1469
1470 pub async fn rewind_to_user_turns(&mut self, keep_user_turns: usize) -> Result<usize> {
1482 let fs_summary = {
1485 let mut store = self.rewind_store.lock().unwrap_or_else(|e| e.into_inner());
1486 store.restore_to(keep_user_turns)
1487 };
1488 if !fs_summary.errors.is_empty() {
1489 for err in &fs_summary.errors {
1490 tracing::warn!(error = %err, "rewind: filesystem restore issue");
1491 }
1492 }
1493 if fs_summary.total_changes() > 0 {
1494 tracing::info!(
1495 restored = fs_summary.restored,
1496 deleted = fs_summary.deleted,
1497 keep_user_turns,
1498 "rewind: restored project files"
1499 );
1500 }
1501
1502 if !self.session.started() || self.session.runtime().is_none() {
1503 self.session.truncate_events_to_user_turns(keep_user_turns);
1505 crate::session::truncate_messages_to_user_turns(
1507 &mut self.initial_messages,
1508 keep_user_turns,
1509 );
1510 return Ok(self.initial_messages.len());
1511 }
1512
1513 let submission_tx = self
1514 .session
1515 .runtime()
1516 .ok_or_else(|| anyhow::anyhow!("session not started"))?
1517 .submission_tx
1518 .clone();
1519
1520 let (response_tx, response_rx) = tokio::sync::oneshot::channel();
1521 if let Err(e) = submission_tx.send(crate::session::SessionCommand::TruncateToUserTurns {
1522 keep_user_turns,
1523 response_tx,
1524 }) {
1525 return Err(anyhow::anyhow!("failed to send rewind command: {}", e));
1526 }
1527
1528 let remaining = response_rx
1529 .await
1530 .map_err(|_| anyhow::anyhow!("rewind cancelled or session loop exited"))??;
1531
1532 self.session.truncate_events_to_user_turns(keep_user_turns);
1533 crate::session::truncate_messages_to_user_turns(
1535 &mut self.initial_messages,
1536 keep_user_turns,
1537 );
1538
1539 if let Err(err) = self.snapshot_session() {
1541 tracing::warn!(error = %err, "failed to snapshot session after rewind");
1542 }
1543
1544 Ok(remaining)
1545 }
1546
1547 pub fn list_rewind_points(&self) -> Vec<crate::rewind::RewindPointMeta> {
1549 self.rewind_store
1550 .lock()
1551 .unwrap_or_else(|e| e.into_inner())
1552 .load_points()
1553 }
1554
1555 pub fn rewind_store_points_path(&self) -> std::path::PathBuf {
1558 self.rewind_store
1559 .lock()
1560 .unwrap_or_else(|e| e.into_inner())
1561 .root()
1562 .to_path_buf()
1563 }
1564
1565 fn apply_pending_session_title(&mut self) -> bool {
1570 let Some(title) = self.session_title_handle.take() else {
1571 return false;
1572 };
1573 self.session.set_title(Some(title.clone()));
1574 self.event_bus
1575 .publish(RuntimeEventKind::SessionTitleUpdated {
1576 session_id: self.session.id().as_str().to_string(),
1577 title,
1578 });
1579 true
1580 }
1581
1582 pub fn snapshot_session(&mut self) -> Result<SessionSnapshot> {
1584 self.session.set_updated_at(current_unix_timestamp());
1585 let _ = self.apply_pending_session_title();
1586 let snapshot = self.session.snapshot(
1587 &self.project_dir,
1588 &self.session_store,
1589 &self.event_bus,
1590 self.goal_runtime.get_goal(),
1591 )?;
1592 self.save_trace_snapshot(snapshot.id.as_str());
1593 Ok(snapshot)
1594 }
1595
1596 pub async fn snapshot_session_async(&mut self) -> Result<SessionSnapshot> {
1598 self.session.set_updated_at(current_unix_timestamp());
1599 let _ = self.apply_pending_session_title();
1600 let snapshot = self
1601 .session
1602 .snapshot_async(
1603 &self.project_dir,
1604 &self.session_store,
1605 &self.event_bus,
1606 self.goal_runtime.get_goal(),
1607 )
1608 .await?;
1609 self.save_trace_snapshot(snapshot.id.as_str());
1610 Ok(snapshot)
1611 }
1612
1613 fn save_trace_snapshot(&self, session_id: &str) {
1614 let traces = turn_traces_from_events(
1615 session_id,
1616 &self.loaded_config.config.model.provider,
1617 &self.loaded_config.config.model.name,
1618 self.session.events(),
1619 );
1620 if traces.is_empty() {
1621 return;
1622 }
1623 let store = TraceStore::new(&self.loaded_config.data_dir);
1624 if let Err(err) = store.save_session_traces(session_id, &traces) {
1625 tracing::warn!(error = %err, session_id, "failed to save turn traces");
1626 }
1627 }
1628
1629 #[cfg(test)]
1631 pub(crate) fn session_store(&self) -> &SessionStore {
1632 &self.session_store
1633 }
1634
1635 fn ensure_tool_executor(&mut self) -> Result<Arc<ToolExecutor>> {
1636 if let Some(executor) = self.tool_executor.as_mut() {
1637 if let Some(executor) = Arc::get_mut(executor) {
1638 if !self.skip_auto_tool_bootstrap {
1639 Self::register_goal_tools_on_executor(
1640 executor,
1641 self.goal_runtime.clone(),
1642 self.loaded_config.config.goals.enabled,
1643 );
1644 executor.register_skill_loader(
1645 self.project_dir.clone(),
1646 self.loaded_config.data_dir.clone(),
1647 self.shared_config.clone(),
1648 );
1649 }
1650 executor.set_rewind_store(Some(self.rewind_store.clone()));
1651 } else {
1652 tracing::warn!(
1653 "tool executor Arc is shared; cannot install goal tools in place \
1654 (strong_count={})",
1655 Arc::strong_count(executor)
1656 );
1657 }
1658 return Ok(executor.clone());
1659 }
1660
1661 let security_policy = SecurityPolicy::new(
1662 self.project_dir.clone(),
1663 self.loaded_config.data_dir.clone(),
1664 self.loaded_config.config.effective_security_config(),
1665 )?;
1666 let harness_policy = crate::harness::select_harness_policy(&self.loaded_config.config);
1667 let profile_name = format!("{:?}", harness_policy.profile).to_lowercase();
1668 let mut executor = ToolExecutor::with_security_policy(
1669 security_policy,
1670 self.runtime_components.security.clone(),
1671 );
1672 executor.set_harness_profile(profile_name);
1673 Self::register_goal_tools_on_executor(
1674 &mut executor,
1675 self.goal_runtime.clone(),
1676 self.loaded_config.config.goals.enabled,
1677 );
1678 executor.register_skill_loader(
1679 self.project_dir.clone(),
1680 self.loaded_config.data_dir.clone(),
1681 self.shared_config.clone(),
1682 );
1683 executor.set_rewind_store(Some(self.rewind_store.clone()));
1684
1685 let workflow_config = self.loaded_config.config.workflow.clone();
1686 let workflow_policy = SecurityPolicy::new(
1687 self.project_dir.clone(),
1688 self.loaded_config.data_dir.clone(),
1689 self.loaded_config.config.security.clone(),
1690 )
1691 .unwrap_or_else(|_| {
1692 SecurityPolicy::new(
1693 self.project_dir.clone(),
1694 self.loaded_config.data_dir.clone(),
1695 crate::config::SecurityConfig::default(),
1696 )
1697 .expect("default security policy")
1698 });
1699 let executor = Arc::new_cyclic(|executor_weak| {
1700 let subagent = SubagentTool::new(
1701 executor_weak.clone(),
1702 self.shared_model_provider.clone(),
1703 self.project_dir.clone(),
1704 self.loaded_config.data_dir.clone(),
1705 self.shared_model_name.clone(),
1706 self.loaded_config.config.harness.clone(),
1707 self.shared_config.clone(),
1708 self.prompt_cache.clone(),
1709 self.runtime_components.clone(),
1710 );
1711 executor.register_tool(Arc::new(subagent));
1712 executor.register_tool(Arc::new(RepoExploreTool::new(self.project_dir.clone())));
1714 executor.register_tool(Arc::new(crate::tool::WorkflowTool::with_subagent_bridge(
1717 workflow_policy.clone(),
1718 workflow_config.clone(),
1719 executor_weak.clone(),
1720 )));
1721 executor
1722 });
1723 if self.agent_mode() == crate::plan_mode::AgentMode::Plan {
1725 let plan_path = crate::plan_store::session_plan_file_path(
1726 &self.loaded_config.data_dir,
1727 self.session.id().as_str(),
1728 );
1729 executor.policy().set_plan_mode(true, Some(plan_path));
1730 }
1731
1732 self.tool_executor = Some(executor.clone());
1733 Ok(executor)
1734 }
1735
1736 fn register_goal_tools_on_executor(
1737 executor: &mut ToolExecutor,
1738 goal_runtime: Arc<GoalRuntimeHandle>,
1739 goals_enabled: bool,
1740 ) {
1741 if !goals_enabled {
1742 return;
1743 }
1744 executor.register_tool(Arc::new(GetGoalTool::new(goal_runtime.clone())));
1747 executor.register_tool(Arc::new(CreateGoalTool::new(goal_runtime.clone())));
1748 executor.register_tool(Arc::new(UpdateGoalTool::new(goal_runtime)));
1749 }
1750
1751 pub fn tool_executor(&self) -> Option<Arc<ToolExecutor>> {
1753 self.tool_executor.clone()
1754 }
1755
1756 pub fn retain_tools<F>(&mut self, pred: F)
1761 where
1762 F: FnMut(&str) -> bool,
1763 {
1764 if let Some(exec) = self.tool_executor.as_mut() {
1765 if let Some(e) = Arc::get_mut(exec) {
1766 e.retain_tools(pred);
1767 } else {
1768 tracing::warn!(
1769 strong_count = Arc::strong_count(exec),
1770 "cannot retain tools; tool executor Arc is shared"
1771 );
1772 }
1773 }
1774 }
1775
1776 pub fn session_title_handle(&self) -> SessionTitleHandle {
1779 self.session_title_handle.clone()
1780 }
1781
1782 pub fn set_tool_executor(&mut self, executor: Arc<ToolExecutor>) {
1784 self.tool_executor = Some(executor);
1785 }
1786
1787 pub fn set_security_config(&mut self, security: SecurityConfig) -> Result<()> {
1790 self.loaded_config.config.security = security.clone();
1791 self.loaded_config.config.tui.yolo_mode =
1792 matches!(security.permission_mode, PermissionMode::Yolo);
1793
1794 let Some(executor) = self.tool_executor.as_mut() else {
1795 return Ok(());
1796 };
1797 let Some(executor) = Arc::get_mut(executor) else {
1798 return Err(anyhow::anyhow!(
1799 "cannot update security policy while tool executor is shared"
1800 ));
1801 };
1802 let mut policy = executor.security_policy().clone();
1803 policy.set_config(security);
1804 executor.set_security_policy(policy);
1805 Ok(())
1806 }
1807
1808 fn get_or_init_memory_manager(&self) -> Result<Option<Arc<crate::memory::MemoryManager>>> {
1810 let memory_config = &self.loaded_config.config.memory;
1811 if !memory_config.enabled {
1812 return Ok(None);
1813 }
1814 let mut guard = self
1815 .memory_manager
1816 .lock()
1817 .unwrap_or_else(|e| e.into_inner());
1818 if let Some(manager) = guard.as_ref() {
1819 return Ok(Some(manager.clone()));
1820 }
1821 let manager = Arc::new(crate::memory::MemoryManager::new(
1822 self.project_dir.clone(),
1823 self.loaded_config.data_dir.clone(),
1824 memory_config,
1825 )?);
1826 *guard = Some(manager.clone());
1827 Ok(Some(manager))
1828 }
1829
1830 fn consolidate_auto_memory(&self) -> Result<()> {
1833 let Some(manager) = self.get_or_init_memory_manager()? else {
1834 return Ok(());
1835 };
1836
1837 let db_path = manager.auto_memory.db_path.clone();
1838 if !db_path.exists() {
1839 return Ok(());
1840 }
1841
1842 let report = manager.auto_memory.consolidate(30)?;
1844
1845 if report.marked_stale > 0 || report.duplicates_merged > 0 {
1846 tracing::info!(
1847 "auto-memory consolidation on session end: {} stale, {} duplicates, {} active",
1848 report.marked_stale,
1849 report.duplicates_merged,
1850 report.remaining_active
1851 );
1852 }
1853
1854 Ok(())
1855 }
1856
1857 fn persist_submitted_session(&mut self) {
1861 if let Err(err) = self.snapshot_session() {
1864 tracing::warn!(error = %err, "early session snapshot after user message failed");
1865 }
1866 }
1867
1868 fn try_extract_memories(&self, _session_id: &str, conversation: &str) {
1872 let Some(extraction_model) = &self.memory_extraction_model else {
1875 tracing::debug!("extract-memories skipped: no memory extraction model configured");
1876 return;
1877 };
1878 let provider = extraction_model.provider.clone();
1879 let model_name = extraction_model.model_name.clone();
1880 let conversation = conversation.to_string();
1881
1882 let manager = match self.get_or_init_memory_manager() {
1883 Ok(Some(m)) => m,
1884 Ok(None) => return,
1885 Err(e) => {
1886 tracing::debug!("extract-memories: failed to init memory manager: {}", e);
1887 return;
1888 }
1889 };
1890
1891 let db_path = manager.auto_memory.db_path.clone();
1892 if !db_path.exists() {
1893 return;
1894 }
1895
1896 let store = manager.auto_memory.clone();
1897
1898 tokio::spawn(async move {
1900 match crate::memory::extract::extract_memories(
1901 &conversation,
1902 provider.as_ref(),
1903 &model_name,
1904 &store,
1905 )
1906 .await
1907 {
1908 Ok(n) => {
1909 if n > 0 {
1910 tracing::info!("extract-memories: saved {} memories from turn", n);
1911 }
1912 }
1913 Err(e) => {
1914 tracing::debug!("extract-memories failed: {}", e);
1915 }
1916 }
1917 });
1918 }
1919
1920 fn try_auto_dream(&self) {
1923 let memory_config = &self.loaded_config.config.memory;
1924 let manager = match self.get_or_init_memory_manager() {
1925 Ok(Some(m)) => m,
1926 Ok(None) => return,
1927 Err(e) => {
1928 tracing::debug!("auto-dream: failed to init memory manager: {}", e);
1929 return;
1930 }
1931 };
1932
1933 let interval_hours = memory_config.dream_interval_days * 24;
1934 let state =
1935 crate::memory::auto_dream::AutoDreamState::new(manager.store.memory_root.clone())
1936 .with_interval(interval_hours.max(1));
1937
1938 let history = manager.history.clone();
1939
1940 let last_dream = state.read_last_dream_at();
1941 if !state.should_dream(&history) {
1942 return;
1943 }
1944
1945 let db_path = manager.auto_memory.db_path.clone();
1946 let hours_since = if last_dream > 0 {
1947 let now = std::time::SystemTime::now()
1948 .duration_since(std::time::UNIX_EPOCH)
1949 .map(|d| d.as_secs())
1950 .unwrap_or(0);
1951 (now.saturating_sub(last_dream)) / 3600
1952 } else {
1953 0
1954 };
1955 let sessions_count = history.list_sessions().map(|s| s.len()).unwrap_or(0);
1956
1957 tracing::info!(
1958 "auto-dream triggered: {}h since last, {} sessions",
1959 hours_since,
1960 sessions_count
1961 );
1962
1963 self.event_bus.publish(RuntimeEventKind::AutoDreamStarted {
1964 hours_since_last: hours_since,
1965 sessions_reviewed: sessions_count,
1966 });
1967
1968 let memory_root = manager.store.memory_root.clone();
1969 let event_sender = self.event_bus.sender();
1970 tokio::spawn(async move {
1971 let result = run_auto_dream_consolidation(&db_path).await;
1972
1973 let dream_state = crate::memory::auto_dream::AutoDreamState::new(memory_root);
1974 match result {
1975 Ok(report) => {
1976 tracing::info!(
1977 "auto-dream completed: {} stale, {} duplicates, {} active",
1978 report.marked_stale,
1979 report.duplicates_merged,
1980 report.remaining_active
1981 );
1982 dream_state.mark_completed();
1983 let _ = event_sender.send(RuntimeEvent::new(
1984 RuntimeEventKind::AutoDreamCompleted {
1985 marked_stale: report.marked_stale,
1986 duplicates_merged: report.duplicates_merged,
1987 active_count: report.remaining_active,
1988 },
1989 ));
1990 }
1991 Err(e) => {
1992 tracing::warn!("auto-dream failed: {}", e);
1993 dream_state.release();
1994 let _ =
1995 event_sender.send(RuntimeEvent::new(RuntimeEventKind::AutoDreamFailed {
1996 reason: e.to_string(),
1997 }));
1998 }
1999 }
2000 });
2001 }
2002
2003 fn try_auto_distill(&self) {
2006 let memory_config = &self.loaded_config.config.memory;
2007 if memory_config.distill_interval_days == 0 {
2008 return;
2009 }
2010
2011 let manager = match self.get_or_init_memory_manager() {
2012 Ok(Some(m)) => m,
2013 _ => return,
2014 };
2015
2016 let interval_hours = memory_config.distill_interval_days * 24;
2017 let state = crate::memory::auto_dream::AutoDreamState::new(
2018 manager.store.memory_root.join("distill"),
2019 )
2020 .with_interval(interval_hours)
2021 .with_min_sessions(3);
2022
2023 let history = manager.history.clone();
2024 if !state.should_dream(&history) {
2025 return;
2026 }
2027
2028 tracing::info!("auto-distill triggered");
2029
2030 let auto_memory = manager.auto_memory.clone();
2031 tokio::spawn(async move {
2032 match auto_memory.consolidate(60) {
2034 Ok(report) => {
2035 tracing::info!(
2036 "auto-distill completed: {} stale, {} duplicates, {} active",
2037 report.marked_stale,
2038 report.duplicates_merged,
2039 report.remaining_active
2040 );
2041 state.mark_completed();
2042 }
2043 Err(e) => {
2044 tracing::warn!("auto-distill failed: {}", e);
2045 state.release();
2046 }
2047 }
2048 });
2049 }
2050
2051 fn build_session_runtime(
2052 &mut self,
2053 ) -> Result<(
2054 crate::session::SessionRuntime,
2055 mpsc::UnboundedReceiver<AgentEvent>,
2056 )> {
2057 let tool_executor = self.ensure_tool_executor()?;
2058 let (event_tx, event_rx) = mpsc::unbounded_channel();
2059 let memory_injection =
2060 self.session_store
2061 .load_memory(&self.project_dir)
2062 .and_then(|memory| {
2063 memory.format_injection(self.loaded_config.config.memory.max_memory_entries)
2064 });
2065
2066 let root = self.load_catalog_skills();
2068 let pools = self.load_catalog_pools();
2069 *self
2070 .shared_available_skills
2071 .lock()
2072 .unwrap_or_else(|e| e.into_inner()) = root.clone();
2073 *self
2074 .shared_skill_pools
2075 .lock()
2076 .unwrap_or_else(|e| e.into_inner()) = pools;
2077 let harness_skills = self.session_active_skill_manifests();
2079 *self
2080 .shared_active_skills
2081 .lock()
2082 .unwrap_or_else(|e| e.into_inner()) = harness_skills.clone();
2083 self.apply_harness_packs_for_active(&harness_skills);
2084
2085 let ctx = Arc::new(crate::turn::TurnContext {
2086 model_provider: self.shared_model_provider.clone(),
2087 tool_executor,
2088 project_dir: self.project_dir.clone(),
2089 data_dir: self.loaded_config.data_dir.clone(),
2090 model_name: self.shared_model_name.clone(),
2091 event_tx: Some(event_tx),
2092 approval_resolver: self.approval_resolver(),
2093 question_resolver: self.question_resolver(),
2094 plan_review_resolver: self.plan_review_resolver(),
2095 sudo_password_resolver: self.sudo_password_resolver(),
2096 compact_state: Arc::new(tokio::sync::Mutex::new(crate::compact::CompactState::new(
2097 crate::config::effective_context_window(&self.loaded_config.config),
2098 ))),
2099 harness_config: self.loaded_config.config.harness.clone(),
2100 include_tool_prompt_manifest: crate::config::effective_tool_prompt_manifest(
2101 &self.loaded_config.config,
2102 ),
2103 context_packets: self.shared_context_packets.clone(),
2104 available_skills: self.shared_available_skills.clone(),
2105 skill_pools: self.shared_skill_pools.clone(),
2106 active_skills: Arc::new(std::sync::Mutex::new(Vec::new())),
2108 prompt_cache: self.prompt_cache.clone(),
2109 instructions: std::sync::Arc::new(std::sync::RwLock::new(None)),
2110 prompt_prefix: std::sync::Arc::new(std::sync::Mutex::new(None)),
2111 components: self.runtime_components.clone(),
2112 cancel_token: self.cancel_token.clone(),
2113 config: self.shared_config.clone(),
2114 memory_injection: memory_injection.clone(),
2115 compaction_provider: None,
2116 agent_mode: self.agent_mode(),
2117 compaction_model_name: None,
2118 session_id: self.session.id().as_str().to_string(),
2119 allowed_tool_names: self.harness_allow_tools.clone(),
2122 is_subagent: false,
2123 memory_manager: self.memory_manager.clone(),
2124 harness_card: self.harness_card.clone(),
2125 });
2126
2127 let policy = select_harness_policy(&self.loaded_config.config);
2128 let session_runtime = crate::session::SessionRuntime::spawn(
2129 ctx,
2130 policy,
2131 self.initial_messages.clone(),
2132 memory_injection,
2133 );
2134
2135 Ok((session_runtime, event_rx))
2136 }
2137
2138 fn record_event(&mut self, event: AgentEvent) {
2139 match &event {
2141 AgentEvent::ToolCompleted(_result) => {
2142 if _result.ok {
2144 if let Some(output) = _result.output.as_object() {
2145 if output.get("tool_name").and_then(|v| v.as_str()) == Some("memory") {
2146 self.turn_used_memory_write = true;
2147 }
2148 }
2149 }
2150 let budget_prompt = self.goal_extension.on_tool_complete();
2151 if budget_prompt.is_some() {
2152 if let Some(goal) = self.goal_runtime.get_goal() {
2153 self.event_bus.publish(RuntimeEventKind::GoalUpdated {
2154 session_id: goal.session_id.clone(),
2155 goal_id: goal.goal_id.as_str().to_string(),
2156 objective: goal.objective.clone(),
2157 short_description: goal.short_description.clone(),
2158 status: goal.status,
2159 tokens_used: goal.tokens_used,
2160 token_budget: goal.token_budget,
2161 });
2162 }
2163 }
2164 }
2165 AgentEvent::UsageReported {
2166 input_tokens,
2167 output_tokens,
2168 ..
2169 } => {
2170 let exceeded = self
2171 .goal_extension
2172 .on_token_usage(*input_tokens, *output_tokens);
2173 if exceeded {
2174 if let Some(goal) = self.goal_runtime.get_goal() {
2175 self.event_bus.publish(RuntimeEventKind::GoalUpdated {
2176 session_id: goal.session_id.clone(),
2177 goal_id: goal.goal_id.as_str().to_string(),
2178 objective: goal.objective.clone(),
2179 short_description: goal.short_description.clone(),
2180 status: goal.status,
2181 tokens_used: goal.tokens_used,
2182 token_budget: goal.token_budget,
2183 });
2184 }
2185 }
2186 }
2187 AgentEvent::SetGoalRequested {
2188 objective,
2189 short_description,
2190 token_budget,
2191 } => {
2192 let goal = self.goal_runtime.set_objective_with_short_description(
2193 objective.clone(),
2194 short_description.clone(),
2195 *token_budget,
2196 );
2197 self.goal_runtime.set_auto_continue(true);
2198 self.publish_goal_updated(&goal);
2199 }
2200 AgentEvent::GoalUpdated {
2201 session_id,
2202 goal_id,
2203 objective,
2204 short_description,
2205 status,
2206 tokens_used,
2207 token_budget,
2208 } => {
2209 self.event_bus.publish(RuntimeEventKind::GoalUpdated {
2212 session_id: session_id.clone(),
2213 goal_id: goal_id.clone(),
2214 objective: objective.clone(),
2215 short_description: short_description.clone(),
2216 status: *status,
2217 tokens_used: *tokens_used,
2218 token_budget: *token_budget,
2219 });
2220 }
2221 AgentEvent::ModelDelta { text } => {
2222 if self.agent_mode() == crate::plan_mode::AgentMode::Plan {
2224 let plans = self.feed_plan_text(&text);
2225 for plan in plans {
2226 self.event_bus.publish(RuntimeEventKind::PlanProposed {
2227 session_id: self.session.id().as_str().to_string(),
2228 title: plan.title,
2229 steps: plan.steps,
2230 });
2231 }
2232 }
2233 }
2234 AgentEvent::ModelOutput { text, .. } => {
2235 if self.agent_mode() == crate::plan_mode::AgentMode::Plan {
2237 let plans = self.feed_plan_text(text);
2238 for plan in plans {
2239 self.event_bus.publish(RuntimeEventKind::PlanProposed {
2240 session_id: self.session.id().as_str().to_string(),
2241 title: plan.title,
2242 steps: plan.steps,
2243 });
2244 }
2245 let remaining = self.drain_plans();
2247 for plan in remaining {
2248 self.event_bus.publish(RuntimeEventKind::PlanProposed {
2249 session_id: self.session.id().as_str().to_string(),
2250 title: plan.title,
2251 steps: plan.steps,
2252 });
2253 }
2254 }
2255 }
2256 AgentEvent::PlanProposed { .. } => {}
2257 AgentEvent::PlanReviewRequested(_) | AgentEvent::PlanReviewResolved(_) => {}
2258 AgentEvent::SudoPasswordRequested(_) => {}
2259 AgentEvent::AgentModeChanged { .. } => {}
2260 _ => {}
2261 }
2262
2263 if let Some(tx) = &self.event_tx {
2264 let _ = tx.send(event.clone());
2265 }
2266 let transient = matches!(
2269 event,
2270 AgentEvent::SubagentActivity { .. }
2271 | AgentEvent::SubagentTranscript { .. }
2272 | AgentEvent::ModelDelta { .. }
2273 | AgentEvent::ModelThinkingDelta { .. }
2274 | AgentEvent::ToolCallStreaming { .. }
2275 | AgentEvent::StreamResuming { .. }
2276 );
2277 if let Some(kind) = runtime_event_kind_from_agent_event(&event) {
2278 self.event_bus.publish(kind);
2279 }
2280 if !transient {
2281 self.session.push_event(event);
2282 }
2283 }
2284
2285 fn update_shared_model_state(&self) {
2286 *self
2287 .shared_model_provider
2288 .write()
2289 .unwrap_or_else(|e| e.into_inner()) = self.model_provider.clone();
2290 *self
2291 .shared_model_name
2292 .write()
2293 .unwrap_or_else(|e| e.into_inner()) = provider_request_model_name(
2294 &self.loaded_config.config.model.provider,
2295 &self.loaded_config.config.model.name,
2296 );
2297 *self
2298 .shared_config
2299 .write()
2300 .unwrap_or_else(|e| e.into_inner()) = self.loaded_config.config.clone();
2301 }
2302
2303 fn load_catalog_skills(&self) -> Vec<SkillManifest> {
2305 match discover_catalog_entries(
2306 &self.loaded_config.config.skills,
2307 &self.project_dir,
2308 &self.loaded_config.data_dir,
2309 ) {
2310 Ok(catalog) => active_skills(
2311 &catalog.root_skills,
2312 &self.loaded_config.config.skills.active,
2313 &self.active_skills,
2314 ),
2315 Err(err) => {
2316 tracing::warn!(error = %err, "failed to load catalog skills");
2317 Vec::new()
2318 }
2319 }
2320 }
2321
2322 fn load_catalog_pools(&self) -> Vec<SkillPool> {
2323 match discover_catalog_entries(
2324 &self.loaded_config.config.skills,
2325 &self.project_dir,
2326 &self.loaded_config.data_dir,
2327 ) {
2328 Ok(catalog) => catalog.pools,
2329 Err(err) => {
2330 tracing::warn!(error = %err, "failed to load skill pools");
2331 Vec::new()
2332 }
2333 }
2334 }
2335
2336 fn load_active_skills(&self) -> Vec<SkillManifest> {
2337 self.load_catalog_skills()
2338 }
2339
2340 fn load_available_skills(&self) -> Vec<SkillManifest> {
2341 self.load_catalog_skills()
2342 }
2343
2344 fn discover_available_skills(&self) -> Result<Vec<SkillManifest>> {
2345 discover_configured_skills(
2347 &self.loaded_config.config.skills,
2348 &self.project_dir,
2349 &self.loaded_config.data_dir,
2350 )
2351 }
2352}
2353
2354fn runtime_event_kind_from_agent_event(event: &AgentEvent) -> Option<RuntimeEventKind> {
2355 match event {
2356 AgentEvent::ModelDelta { text } => {
2357 Some(RuntimeEventKind::AssistantDelta { text: text.clone() })
2358 }
2359 AgentEvent::ModelThinkingDelta { text } => {
2360 Some(RuntimeEventKind::AssistantThinkingDelta { text: text.clone() })
2361 }
2362 AgentEvent::ToolRequested(invocation) => {
2363 Some(RuntimeEventKind::ToolRequested(invocation.clone()))
2364 }
2365 AgentEvent::ToolCompleted(result) => Some(RuntimeEventKind::ToolCompleted(result.clone())),
2366 AgentEvent::SubagentActivity {
2367 invocation_id,
2368 message,
2369 } => Some(RuntimeEventKind::SubagentActivity {
2370 invocation_id: invocation_id.clone(),
2371 message: message.clone(),
2372 }),
2373 AgentEvent::SubagentTranscript {
2374 invocation_id,
2375 item,
2376 } => Some(RuntimeEventKind::SubagentTranscript {
2377 invocation_id: invocation_id.clone(),
2378 item: item.clone(),
2379 }),
2380 AgentEvent::ApprovalRequested(request) => {
2381 Some(RuntimeEventKind::ApprovalRequired(request.clone()))
2382 }
2383 AgentEvent::UsageReported {
2384 input_tokens,
2385 output_tokens,
2386 cache_creation_tokens,
2387 cache_read_tokens,
2388 } => Some(RuntimeEventKind::TokensUpdated {
2389 input_tokens: *input_tokens,
2390 output_tokens: *output_tokens,
2391 cache_creation_tokens: *cache_creation_tokens,
2392 cache_read_tokens: *cache_read_tokens,
2393 }),
2394 AgentEvent::Error { message } => Some(RuntimeEventKind::Error {
2395 message: message.clone(),
2396 }),
2397 AgentEvent::ApprovalResolved(decision) => {
2398 Some(RuntimeEventKind::ApprovalResolved(decision.clone()))
2399 }
2400 AgentEvent::CapabilityRecorded(entry) => {
2401 Some(RuntimeEventKind::CapabilityRecorded(entry.clone()))
2402 }
2403 AgentEvent::QuestionRequested(request) => {
2404 Some(RuntimeEventKind::QuestionRequired(request.clone()))
2405 }
2406 AgentEvent::QuestionResolved(response) => {
2407 Some(RuntimeEventKind::QuestionResolved(response.clone()))
2408 }
2409 AgentEvent::PlanReviewRequested(request) => {
2410 Some(RuntimeEventKind::PlanReviewRequired(request.clone()))
2411 }
2412 AgentEvent::PlanReviewResolved(response) => {
2413 Some(RuntimeEventKind::PlanReviewResolved(response.clone()))
2414 }
2415 AgentEvent::SudoPasswordRequested(request) => {
2416 Some(RuntimeEventKind::SudoPasswordRequired(request.clone()))
2417 }
2418 AgentEvent::HarnessTrace(value) => Some(RuntimeEventKind::HarnessTrace(value.clone())),
2419 AgentEvent::HarnessStopped {
2420 reason,
2421 message,
2422 tool_name,
2423 } => Some(RuntimeEventKind::HarnessStopped {
2424 reason: reason.clone(),
2425 message: message.clone(),
2426 tool_name: tool_name.clone(),
2427 }),
2428 AgentEvent::PatchProposed(patch) => Some(RuntimeEventKind::PatchProposed(patch.clone())),
2429 AgentEvent::MicroCompactApplied { messages_cleared } => {
2430 Some(RuntimeEventKind::MicroCompactApplied {
2431 messages_cleared: *messages_cleared,
2432 })
2433 }
2434 AgentEvent::AutoCompactStarted => Some(RuntimeEventKind::AutoCompactStarted),
2435 AgentEvent::AutoCompactCompleted {
2436 tokens_saved,
2437 summary,
2438 kept_recent_messages,
2439 } => Some(RuntimeEventKind::AutoCompactCompleted {
2440 tokens_saved: *tokens_saved,
2441 summary: summary.clone(),
2442 kept_recent_messages: *kept_recent_messages,
2443 }),
2444 AgentEvent::AutoCompactFailed { reason } => Some(RuntimeEventKind::AutoCompactFailed {
2445 reason: reason.clone(),
2446 }),
2447 AgentEvent::UserTaskSubmitted { .. } | AgentEvent::ModelOutput { .. } => None,
2448 AgentEvent::RepeatedToolCallWarning { .. } => None,
2449 AgentEvent::RepetitionDetected { .. } => None,
2450 AgentEvent::GoalUpdated { .. } => None,
2451 AgentEvent::SetGoalRequested {
2452 objective,
2453 short_description,
2454 token_budget,
2455 } => Some(RuntimeEventKind::SetGoalRequested {
2456 objective: objective.clone(),
2457 short_description: short_description.clone(),
2458 token_budget: *token_budget,
2459 }),
2460 AgentEvent::AutoDreamStarted {
2461 hours_since_last,
2462 sessions_reviewed,
2463 } => Some(RuntimeEventKind::AutoDreamStarted {
2464 hours_since_last: *hours_since_last,
2465 sessions_reviewed: *sessions_reviewed,
2466 }),
2467 AgentEvent::AutoDreamCompleted {
2468 marked_stale,
2469 duplicates_merged,
2470 active_count,
2471 } => Some(RuntimeEventKind::AutoDreamCompleted {
2472 marked_stale: *marked_stale,
2473 duplicates_merged: *duplicates_merged,
2474 active_count: *active_count,
2475 }),
2476 AgentEvent::AutoDreamFailed { reason } => Some(RuntimeEventKind::AutoDreamFailed {
2477 reason: reason.clone(),
2478 }),
2479 AgentEvent::PlanProposed { .. } => None,
2480 AgentEvent::AgentModeChanged { .. } => None,
2482 AgentEvent::SessionRecap { .. } => None,
2483 AgentEvent::StreamResuming { .. } => None,
2485 AgentEvent::ToolCallStreaming { .. } => None,
2487 AgentEvent::NotificationRequested { .. } | AgentEvent::UpdateAvailable { .. } => None,
2489 }
2490}
2491
2492async fn run_auto_dream_consolidation(
2495 db_path: &std::path::Path,
2496) -> anyhow::Result<crate::memory::ConsolidationReport> {
2497 let store = crate::memory::AutoMemoryStore::open(db_path)?;
2498 store.consolidate(30)
2499}