1use crate::agent::Agent;
9use crate::harness::Harness;
10use crate::provider::DriverId;
11use crate::session_file::{
12 FileInfo, FileStat, GrepMatch, GrepOptions, GrepSearchResult, InitialFile, SessionFile,
13};
14use crate::tool_types::{ToolCall, ToolDefinition, ToolResult};
15use crate::typed_id::{AgentId, HarnessId, ImageId, MessageId, ModelId, SessionId, WorkspaceId};
16use async_trait::async_trait;
17use chrono::{DateTime, Utc};
18use std::any::{Any, TypeId};
19use std::collections::{HashMap, HashSet};
20use std::sync::Arc;
21use uuid::Uuid;
22
23fn build_tool_map(tool_defs: &[ToolDefinition]) -> HashMap<&str, &ToolDefinition> {
25 tool_defs.iter().map(|def| (def.name(), def)).collect()
26}
27
28use crate::error::Result;
29
30#[derive(Clone, Default)]
45pub struct ReasoningEffortHandle {
46 inner: Arc<std::sync::RwLock<Option<String>>>,
47}
48
49impl ReasoningEffortHandle {
50 pub fn new() -> Self {
52 Self::default()
53 }
54
55 pub fn with_effort(effort: impl Into<String>) -> Self {
57 Self {
58 inner: Arc::new(std::sync::RwLock::new(Some(effort.into()))),
59 }
60 }
61
62 pub fn set(&self, effort: Option<String>) {
66 let mut guard = self.inner.write().unwrap_or_else(|e| e.into_inner());
70 *guard = effort;
71 }
72
73 pub fn get(&self) -> Option<String> {
75 let guard = self.inner.read().unwrap_or_else(|e| e.into_inner());
78 guard.clone()
79 }
80}
81
82impl std::fmt::Debug for ReasoningEffortHandle {
83 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
84 f.debug_struct("ReasoningEffortHandle")
85 .field("effort", &self.get())
86 .finish()
87 }
88}
89
90#[async_trait]
101pub trait AgentStore: Send + Sync {
102 async fn get_agent(&self, agent_id: AgentId) -> Result<Option<Agent>>;
104}
105
106#[async_trait]
107impl<T: AgentStore + ?Sized> AgentStore for std::sync::Arc<T> {
108 async fn get_agent(&self, agent_id: AgentId) -> Result<Option<Agent>> {
109 (**self).get_agent(agent_id).await
110 }
111}
112
113#[async_trait]
128pub trait HarnessStore: Send + Sync {
129 async fn get_harness_chain(&self, harness_id: HarnessId) -> Result<Vec<Harness>>;
134}
135
136#[async_trait]
137impl<T: HarnessStore + ?Sized> HarnessStore for std::sync::Arc<T> {
138 async fn get_harness_chain(&self, harness_id: HarnessId) -> Result<Vec<Harness>> {
139 (**self).get_harness_chain(harness_id).await
140 }
141}
142
143use crate::capability_types::AgentCapabilityConfig;
148use crate::leased_resource::{LeasedResource, UpsertLeasedResource};
149use crate::session::Session;
150
151#[async_trait]
157pub trait SessionStore: Send + Sync {
158 async fn get_session(&self, session_id: SessionId) -> Result<Option<Session>>;
160}
161
162#[async_trait]
163impl<T: SessionStore + ?Sized> SessionStore for std::sync::Arc<T> {
164 async fn get_session(&self, session_id: SessionId) -> Result<Option<Session>> {
165 (**self).get_session(session_id).await
166 }
167}
168
169#[async_trait]
171pub trait SessionMutator: Send + Sync {
172 async fn update_session_title(&self, session_id: SessionId, title: String) -> Result<Session>;
174
175 async fn upsert_session_capability(
181 &self,
182 _session_id: SessionId,
183 _capability: AgentCapabilityConfig,
184 ) -> Result<Session> {
185 Err(crate::error::AgentLoopError::store(
186 "session backend does not support live capability reconfiguration",
187 ))
188 }
189
190 async fn remove_session_capability(
192 &self,
193 _session_id: SessionId,
194 _capability_id: &str,
195 ) -> Result<Session> {
196 Err(crate::error::AgentLoopError::store(
197 "session backend does not support live capability reconfiguration",
198 ))
199 }
200}
201
202#[async_trait]
203impl<T: SessionMutator + ?Sized> SessionMutator for std::sync::Arc<T> {
204 async fn update_session_title(&self, session_id: SessionId, title: String) -> Result<Session> {
205 (**self).update_session_title(session_id, title).await
206 }
207
208 async fn upsert_session_capability(
209 &self,
210 session_id: SessionId,
211 capability: AgentCapabilityConfig,
212 ) -> Result<Session> {
213 (**self)
214 .upsert_session_capability(session_id, capability)
215 .await
216 }
217
218 async fn remove_session_capability(
219 &self,
220 session_id: SessionId,
221 capability_id: &str,
222 ) -> Result<Session> {
223 (**self)
224 .remove_session_capability(session_id, capability_id)
225 .await
226 }
227}
228
229#[derive(Debug, Clone)]
239pub struct ResolvedModel {
240 pub model: String,
241 pub provider_type: DriverId,
242 pub api_key: Option<String>,
243 pub base_url: Option<String>,
244 pub provider_metadata: Option<crate::driver_registry::ProviderMetadata>,
245}
246
247impl ResolvedModel {
248 pub fn provider_key(&self) -> crate::ProviderKey {
249 self.provider_metadata
250 .as_ref()
251 .and_then(|metadata| metadata.extra.as_ref())
252 .and_then(|extra| extra.get("provider_id"))
253 .and_then(serde_json::Value::as_str)
254 .map(crate::ProviderKey::new)
255 .unwrap_or_else(|| crate::ProviderKey::new(self.provider_type.as_str()))
256 }
257
258 pub fn canonical_parts(&self) -> (crate::ModelSpec, crate::ProviderConfig) {
259 let provider = self.provider_key();
260 let spec = crate::ModelSpec::on(provider.clone(), self.model.clone());
261 let config = crate::ProviderConfig {
262 provider,
263 provider_type: self.provider_type.clone(),
264 api_key: self.api_key.clone(),
265 base_url: self.base_url.clone(),
266 metadata: self.provider_metadata.clone().unwrap_or_default(),
267 };
268 (spec, config)
269 }
270}
271
272#[async_trait]
283pub trait ProviderStore: Send + Sync {
284 async fn get_resolved_model(&self, model_id: ModelId) -> Result<Option<ResolvedModel>>;
290
291 async fn get_default_model(&self) -> Result<Option<ResolvedModel>>;
295
296 async fn get_provider_config(
301 &self,
302 _provider: &crate::ProviderKey,
303 ) -> Result<Option<crate::ProviderConfig>> {
304 Ok(None)
305 }
306}
307
308#[async_trait]
309impl<T: ProviderStore + ?Sized> ProviderStore for std::sync::Arc<T> {
310 async fn get_resolved_model(&self, model_id: ModelId) -> Result<Option<ResolvedModel>> {
311 (**self).get_resolved_model(model_id).await
312 }
313
314 async fn get_default_model(&self) -> Result<Option<ResolvedModel>> {
315 (**self).get_default_model().await
316 }
317
318 async fn get_provider_config(
319 &self,
320 provider: &crate::ProviderKey,
321 ) -> Result<Option<crate::ProviderConfig>> {
322 (**self).get_provider_config(provider).await
323 }
324}
325
326#[derive(Debug, Clone)]
332pub struct StoredImageInfo {
333 pub id: ImageId,
334 pub filename: String,
335 pub content_type: String,
336 pub size_bytes: i64,
337 pub metadata: serde_json::Value,
338 pub created_at: DateTime<Utc>,
339}
340
341#[derive(Debug, Clone)]
343pub struct StoredImage {
344 pub info: StoredImageInfo,
345 pub data: Vec<u8>,
346}
347
348#[derive(Debug, Clone)]
350pub struct CreateStoredImage {
351 pub filename: String,
352 pub content_type: String,
353 pub data: Vec<u8>,
354 pub metadata: serde_json::Value,
355}
356
357#[async_trait]
358pub trait ImageArtifactStore: Send + Sync {
359 async fn create_image(&self, input: CreateStoredImage) -> Result<StoredImageInfo>;
361
362 async fn get_image(&self, image_id: ImageId) -> Result<Option<StoredImage>>;
364
365 async fn get_image_info(&self, image_id: ImageId) -> Result<Option<StoredImageInfo>>;
367}
368
369#[derive(Debug, Clone)]
375pub struct ProviderCredentials {
376 pub api_key: String,
377 pub base_url: Option<String>,
378}
379
380#[async_trait]
381pub trait ProviderCredentialStore: Send + Sync {
382 async fn get_default_provider_credentials(
387 &self,
388 provider_type: &str,
389 ) -> Result<Option<ProviderCredentials>>;
390}
391
392#[async_trait]
403pub trait ToolExecutor: Send + Sync {
404 async fn execute(&self, tool_call: &ToolCall, tool_def: &ToolDefinition) -> Result<ToolResult>;
409
410 async fn execute_with_context(
415 &self,
416 tool_call: &ToolCall,
417 tool_def: &ToolDefinition,
418 _context: &ToolContext,
419 ) -> Result<ToolResult> {
420 self.execute(tool_call, tool_def).await
422 }
423
424 async fn execute_batch(
426 &self,
427 tool_calls: &[ToolCall],
428 tool_defs: &[ToolDefinition],
429 ) -> Result<Vec<ToolResult>> {
430 let mut results = Vec::with_capacity(tool_calls.len());
431
432 let tool_map = build_tool_map(tool_defs);
433
434 for tool_call in tool_calls {
435 let tool_def = tool_map.get(tool_call.name.as_str()).ok_or_else(|| {
436 crate::error::AgentLoopError::tool(format!(
437 "Tool definition not found: {}",
438 tool_call.name
439 ))
440 })?;
441
442 results.push(self.execute(tool_call, tool_def).await?);
443 }
444
445 Ok(results)
446 }
447
448 async fn execute_parallel(
450 &self,
451 tool_calls: &[ToolCall],
452 tool_defs: &[ToolDefinition],
453 ) -> Result<Vec<ToolResult>>
454 where
455 Self: Sized,
456 {
457 use futures::future::join_all;
458
459 let tool_map = build_tool_map(tool_defs);
460
461 let futures: Vec<_> = tool_calls
462 .iter()
463 .map(|tool_call| async {
464 let tool_def = tool_map.get(tool_call.name.as_str()).ok_or_else(|| {
465 crate::error::AgentLoopError::tool(format!(
466 "Tool definition not found: {}",
467 tool_call.name
468 ))
469 })?;
470 self.execute(tool_call, tool_def).await
471 })
472 .collect();
473
474 let results = join_all(futures).await;
475 results.into_iter().collect()
476 }
477}
478
479#[async_trait]
483impl ToolExecutor for std::sync::Arc<dyn ToolExecutor> {
484 async fn execute(&self, tool_call: &ToolCall, tool_def: &ToolDefinition) -> Result<ToolResult> {
485 (**self).execute(tool_call, tool_def).await
486 }
487
488 async fn execute_with_context(
489 &self,
490 tool_call: &ToolCall,
491 tool_def: &ToolDefinition,
492 context: &ToolContext,
493 ) -> Result<ToolResult> {
494 (**self)
495 .execute_with_context(tool_call, tool_def, context)
496 .await
497 }
498
499 async fn execute_batch(
500 &self,
501 tool_calls: &[ToolCall],
502 tool_defs: &[ToolDefinition],
503 ) -> Result<Vec<ToolResult>> {
504 (**self).execute_batch(tool_calls, tool_defs).await
505 }
506}
507
508#[async_trait]
520pub trait SessionFileSystem: Send + Sync {
521 fn display_root(&self) -> String {
527 crate::session_path::WORKSPACE_PREFIX.to_string()
528 }
529
530 fn display_path(&self, path: &str) -> String {
536 crate::session_path::to_display_path(path)
537 }
538
539 fn resolve_path(&self, input: &str) -> String {
553 crate::session_path::to_session_path(input)
554 }
555
556 fn is_mount_resolver(&self) -> bool;
560
561 async fn read_file(&self, session_id: SessionId, path: &str) -> Result<Option<SessionFile>>;
563
564 async fn write_file(
566 &self,
567 session_id: SessionId,
568 path: &str,
569 content: &str,
570 encoding: &str,
571 ) -> Result<SessionFile>;
572
573 async fn write_file_if_content_matches(
578 &self,
579 session_id: SessionId,
580 path: &str,
581 expected_content: &str,
582 expected_encoding: &str,
583 content: &str,
584 encoding: &str,
585 ) -> Result<Option<SessionFile>> {
586 let Some(existing) = self.read_file(session_id, path).await? else {
587 return Ok(None);
588 };
589
590 if existing.is_directory {
591 return Ok(None);
592 }
593
594 let current_content = existing.content.unwrap_or_default();
595 if current_content != expected_content || existing.encoding != expected_encoding {
596 return Ok(None);
597 }
598
599 self.write_file(session_id, path, content, encoding)
600 .await
601 .map(Some)
602 }
603
604 async fn delete_file(&self, session_id: SessionId, path: &str, recursive: bool)
606 -> Result<bool>;
607
608 async fn list_directory(&self, session_id: SessionId, path: &str) -> Result<Vec<FileInfo>>;
610
611 async fn stat_file(&self, session_id: SessionId, path: &str) -> Result<Option<FileStat>>;
613
614 async fn grep_files(
621 &self,
622 session_id: SessionId,
623 pattern: &str,
624 path_pattern: Option<&str>,
625 ) -> Result<Vec<GrepMatch>>;
626
627 async fn grep_files_with_options(
633 &self,
634 session_id: SessionId,
635 pattern: &str,
636 options: &GrepOptions,
637 ) -> Result<GrepSearchResult> {
638 if options.before_context != 0 || options.after_context != 0 {
639 return Err(crate::error::AgentLoopError::tool(
640 "this file store does not support grep context",
641 ));
642 }
643 let all = self
644 .grep_files(session_id, pattern, options.path_pattern.as_deref())
645 .await?;
646 Ok(crate::session_file::bound_grep_matches(all, options))
647 }
648
649 async fn create_directory(&self, session_id: SessionId, path: &str) -> Result<FileInfo>;
651
652 async fn seed_initial_file(&self, session_id: SessionId, file: &InitialFile) -> Result<()> {
654 if file.is_readonly {
655 return Err(crate::error::AgentLoopError::store(
656 "read-only initial files require a SessionFileSystem-specific seed implementation",
657 ));
658 }
659 self.write_file(session_id, &file.path, &file.content, &file.encoding)
660 .await?;
661 Ok(())
662 }
663}
664
665pub struct WorkspaceScopedFileSystem {
675 inner: Arc<dyn SessionFileSystem>,
676 key: SessionId,
677}
678
679impl WorkspaceScopedFileSystem {
680 pub fn wrap(
682 inner: Arc<dyn SessionFileSystem>,
683 workspace_id: WorkspaceId,
684 ) -> Arc<dyn SessionFileSystem> {
685 Arc::new(Self {
686 inner,
687 key: SessionId::from_uuid(workspace_id.uuid()),
688 })
689 }
690}
691
692#[async_trait]
693impl SessionFileSystem for WorkspaceScopedFileSystem {
694 async fn read_file(&self, _session_id: SessionId, path: &str) -> Result<Option<SessionFile>> {
695 self.inner.read_file(self.key, path).await
696 }
697 async fn write_file(
698 &self,
699 _session_id: SessionId,
700 path: &str,
701 content: &str,
702 encoding: &str,
703 ) -> Result<SessionFile> {
704 self.inner
705 .write_file(self.key, path, content, encoding)
706 .await
707 }
708 async fn write_file_if_content_matches(
709 &self,
710 _session_id: SessionId,
711 path: &str,
712 expected_content: &str,
713 expected_encoding: &str,
714 content: &str,
715 encoding: &str,
716 ) -> Result<Option<SessionFile>> {
717 self.inner
718 .write_file_if_content_matches(
719 self.key,
720 path,
721 expected_content,
722 expected_encoding,
723 content,
724 encoding,
725 )
726 .await
727 }
728 async fn delete_file(
729 &self,
730 _session_id: SessionId,
731 path: &str,
732 recursive: bool,
733 ) -> Result<bool> {
734 self.inner.delete_file(self.key, path, recursive).await
735 }
736 async fn list_directory(&self, _session_id: SessionId, path: &str) -> Result<Vec<FileInfo>> {
737 self.inner.list_directory(self.key, path).await
738 }
739 async fn stat_file(&self, _session_id: SessionId, path: &str) -> Result<Option<FileStat>> {
740 self.inner.stat_file(self.key, path).await
741 }
742 async fn grep_files(
743 &self,
744 _session_id: SessionId,
745 pattern: &str,
746 path_pattern: Option<&str>,
747 ) -> Result<Vec<GrepMatch>> {
748 self.inner.grep_files(self.key, pattern, path_pattern).await
749 }
750 async fn grep_files_with_options(
751 &self,
752 _session_id: SessionId,
753 pattern: &str,
754 options: &GrepOptions,
755 ) -> Result<GrepSearchResult> {
756 self.inner
757 .grep_files_with_options(self.key, pattern, options)
758 .await
759 }
760 async fn create_directory(&self, _session_id: SessionId, path: &str) -> Result<FileInfo> {
761 self.inner.create_directory(self.key, path).await
762 }
763 async fn seed_initial_file(&self, _session_id: SessionId, file: &InitialFile) -> Result<()> {
764 self.inner.seed_initial_file(self.key, file).await
765 }
766
767 fn display_root(&self) -> String {
768 self.inner.display_root()
769 }
770
771 fn display_path(&self, path: &str) -> String {
772 self.inner.display_path(path)
773 }
774
775 fn resolve_path(&self, input: &str) -> String {
776 self.inner.resolve_path(input)
777 }
778
779 fn is_mount_resolver(&self) -> bool {
780 self.inner.is_mount_resolver()
781 }
782}
783
784#[async_trait]
785impl<T: SessionFileSystem + ?Sized> SessionFileSystem for std::sync::Arc<T> {
786 fn display_root(&self) -> String {
787 (**self).display_root()
788 }
789
790 fn display_path(&self, path: &str) -> String {
791 (**self).display_path(path)
792 }
793
794 fn resolve_path(&self, input: &str) -> String {
795 (**self).resolve_path(input)
796 }
797
798 fn is_mount_resolver(&self) -> bool {
799 (**self).is_mount_resolver()
800 }
801
802 async fn read_file(&self, session_id: SessionId, path: &str) -> Result<Option<SessionFile>> {
803 (**self).read_file(session_id, path).await
804 }
805
806 async fn write_file(
807 &self,
808 session_id: SessionId,
809 path: &str,
810 content: &str,
811 encoding: &str,
812 ) -> Result<SessionFile> {
813 (**self)
814 .write_file(session_id, path, content, encoding)
815 .await
816 }
817
818 async fn write_file_if_content_matches(
819 &self,
820 session_id: SessionId,
821 path: &str,
822 expected_content: &str,
823 expected_encoding: &str,
824 content: &str,
825 encoding: &str,
826 ) -> Result<Option<SessionFile>> {
827 (**self)
828 .write_file_if_content_matches(
829 session_id,
830 path,
831 expected_content,
832 expected_encoding,
833 content,
834 encoding,
835 )
836 .await
837 }
838
839 async fn delete_file(
840 &self,
841 session_id: SessionId,
842 path: &str,
843 recursive: bool,
844 ) -> Result<bool> {
845 (**self).delete_file(session_id, path, recursive).await
846 }
847
848 async fn list_directory(&self, session_id: SessionId, path: &str) -> Result<Vec<FileInfo>> {
849 (**self).list_directory(session_id, path).await
850 }
851
852 async fn stat_file(&self, session_id: SessionId, path: &str) -> Result<Option<FileStat>> {
853 (**self).stat_file(session_id, path).await
854 }
855
856 async fn grep_files(
857 &self,
858 session_id: SessionId,
859 pattern: &str,
860 path_pattern: Option<&str>,
861 ) -> Result<Vec<GrepMatch>> {
862 (**self).grep_files(session_id, pattern, path_pattern).await
863 }
864
865 async fn grep_files_with_options(
866 &self,
867 session_id: SessionId,
868 pattern: &str,
869 options: &GrepOptions,
870 ) -> Result<GrepSearchResult> {
871 (**self)
872 .grep_files_with_options(session_id, pattern, options)
873 .await
874 }
875
876 async fn create_directory(&self, session_id: SessionId, path: &str) -> Result<FileInfo> {
877 (**self).create_directory(session_id, path).await
878 }
879
880 async fn seed_initial_file(&self, session_id: SessionId, file: &InitialFile) -> Result<()> {
881 (**self).seed_initial_file(session_id, file).await
882 }
883}
884
885pub use SessionFileSystem as SessionFileStore;
887
888#[derive(Clone, Default)]
894pub struct SessionFileSystemFactoryContext {
895 values: Arc<HashMap<TypeId, Arc<dyn Any + Send + Sync>>>,
896}
897
898impl SessionFileSystemFactoryContext {
899 pub fn new() -> Self {
900 Self::default()
901 }
902
903 pub fn with<T: Any + Send + Sync>(mut self, value: Arc<T>) -> Self {
904 let values = Arc::make_mut(&mut self.values);
905 values.insert(TypeId::of::<T>(), value);
906 self
907 }
908
909 pub fn get<T: Any + Send + Sync>(&self) -> Option<Arc<T>> {
910 self.values
911 .get(&TypeId::of::<T>())
912 .and_then(|value| value.clone().downcast::<T>().ok())
913 }
914
915 pub fn with_workspace_roots(self, roots: Arc<crate::WorkspaceRootSet>) -> Self {
916 self.with(roots)
917 }
918
919 pub fn workspace_roots(&self) -> Option<Arc<crate::WorkspaceRootSet>> {
920 self.get::<crate::WorkspaceRootSet>()
921 }
922}
923
924#[async_trait]
926pub trait SessionFileSystemFactory: Send + Sync {
927 fn name(&self) -> &'static str {
929 "SessionFileSystemFactory"
930 }
931
932 fn is_disabled(&self) -> bool {
935 false
936 }
937
938 async fn create_session_file_system(
940 &self,
941 context: SessionFileSystemFactoryContext,
942 ) -> Result<Arc<dyn SessionFileSystem>>;
943}
944
945#[derive(Debug, Clone, Default)]
947pub struct DisabledSessionFileSystemFactory;
948
949#[async_trait]
950impl SessionFileSystemFactory for DisabledSessionFileSystemFactory {
951 fn name(&self) -> &'static str {
952 "DisabledSessionFileSystemFactory"
953 }
954
955 fn is_disabled(&self) -> bool {
956 true
957 }
958
959 async fn create_session_file_system(
960 &self,
961 _context: SessionFileSystemFactoryContext,
962 ) -> Result<Arc<dyn SessionFileSystem>> {
963 Err(crate::error::AgentLoopError::config(
964 "session filesystem is disabled",
965 ))
966 }
967}
968
969#[derive(Debug, Clone)]
975pub struct KeyInfo {
976 pub key: String,
977 pub created_at: chrono::DateTime<chrono::Utc>,
978 pub updated_at: chrono::DateTime<chrono::Utc>,
979}
980
981#[derive(Debug, Clone)]
983pub struct SecretInfo {
984 pub name: String,
985 pub created_at: chrono::DateTime<chrono::Utc>,
986 pub updated_at: chrono::DateTime<chrono::Utc>,
987}
988
989#[derive(Debug, Clone, serde::Serialize)]
999pub struct KnowledgeSearchHit {
1000 pub id: String,
1002 pub kb_id: String,
1004 pub title: String,
1005 pub kind: String,
1006 pub tags: Vec<String>,
1007 pub snippet: String,
1009 pub resource: Option<String>,
1011}
1012
1013#[async_trait]
1017pub trait KnowledgeStore: Send + Sync {
1018 async fn search_knowledge(
1019 &self,
1020 org_id: crate::typed_id::OrgId,
1021 kb_public_ids: &[String],
1022 query: &str,
1023 kind: Option<&str>,
1024 tags: &[String],
1025 limit: usize,
1026 ) -> Result<Vec<KnowledgeSearchHit>>;
1027}
1028
1029#[async_trait]
1034pub trait SessionStorageStore: Send + Sync {
1035 async fn set_value(&self, session_id: SessionId, key: &str, value: &str) -> Result<()>;
1039
1040 async fn get_value(&self, session_id: SessionId, key: &str) -> Result<Option<String>>;
1042
1043 async fn delete_value(&self, session_id: SessionId, key: &str) -> Result<bool>;
1045
1046 async fn list_keys(&self, session_id: SessionId) -> Result<Vec<KeyInfo>>;
1048
1049 async fn set_secret(&self, session_id: SessionId, name: &str, value: &str) -> Result<()>;
1053
1054 async fn get_secret(&self, session_id: SessionId, name: &str) -> Result<Option<String>>;
1056
1057 async fn delete_secret(&self, session_id: SessionId, name: &str) -> Result<bool>;
1059
1060 async fn list_secrets(&self, session_id: SessionId) -> Result<Vec<SecretInfo>>;
1062}
1063
1064use crate::session_schedule::SessionSchedule;
1069use crate::typed_id::ScheduleId;
1070
1071#[async_trait]
1075pub trait SessionScheduleStore: Send + Sync {
1076 async fn create_schedule(
1078 &self,
1079 session_id: SessionId,
1080 description: String,
1081 cron_expression: Option<String>,
1082 scheduled_at: Option<chrono::DateTime<chrono::Utc>>,
1083 timezone: String,
1084 ) -> Result<SessionSchedule>;
1085
1086 async fn create_schedule_enforcing_limits(
1090 &self,
1091 session_id: SessionId,
1092 description: String,
1093 cron_expression: Option<String>,
1094 scheduled_at: Option<chrono::DateTime<chrono::Utc>>,
1095 timezone: String,
1096 ) -> std::result::Result<SessionSchedule, crate::session_schedule::ScheduleLimitError> {
1097 let per_session = self
1098 .count_active_schedules(session_id)
1099 .await
1100 .map_err(crate::session_schedule::ScheduleLimitError::Store)?;
1101 if per_session >= crate::session_schedule::MAX_ACTIVE_SCHEDULES_PER_SESSION {
1102 return Err(crate::session_schedule::ScheduleLimitError::Rejected(
1103 format!(
1104 "Maximum {} active schedules per session. Cancel an existing schedule first.",
1105 crate::session_schedule::MAX_ACTIVE_SCHEDULES_PER_SESSION
1106 ),
1107 ));
1108 }
1109
1110 let max_per_org = crate::session_schedule::max_active_schedules_per_org();
1111 let per_org = self
1112 .count_active_org_schedules()
1113 .await
1114 .map_err(crate::session_schedule::ScheduleLimitError::Store)?;
1115 if i64::from(per_org) >= max_per_org {
1116 return Err(crate::session_schedule::ScheduleLimitError::Rejected(
1117 format!(
1118 "Maximum {max_per_org} active schedules per org reached. Cancel an existing schedule first."
1119 ),
1120 ));
1121 }
1122
1123 if let Some(cron) = cron_expression.as_deref() {
1124 crate::session_schedule::validate_cron_min_interval(cron)
1125 .map_err(crate::session_schedule::ScheduleLimitError::Rejected)?;
1126 }
1127
1128 self.create_schedule(
1129 session_id,
1130 description,
1131 cron_expression,
1132 scheduled_at,
1133 timezone,
1134 )
1135 .await
1136 .map_err(crate::session_schedule::ScheduleLimitError::Store)
1137 }
1138
1139 async fn cancel_schedule(
1141 &self,
1142 session_id: SessionId,
1143 schedule_id: ScheduleId,
1144 ) -> Result<SessionSchedule>;
1145
1146 async fn list_schedules(&self, session_id: SessionId) -> Result<Vec<SessionSchedule>>;
1148
1149 async fn count_active_schedules(&self, session_id: SessionId) -> Result<u32>;
1151
1152 async fn count_active_org_schedules(&self) -> Result<u32>;
1157}
1158
1159#[async_trait]
1169pub trait SessionResourceRegistry: Send + Sync {
1170 async fn register(
1172 &self,
1173 entry: crate::session_resource::RegisterSessionResource,
1174 ) -> Result<crate::session_resource::SessionResourceEntry>;
1175
1176 async fn update_status(
1178 &self,
1179 session_id: SessionId,
1180 resource_id: &str,
1181 status: crate::session_resource::SessionResourceStatus,
1182 ) -> Result<Option<crate::session_resource::SessionResourceEntry>>;
1183
1184 async fn get(
1186 &self,
1187 session_id: SessionId,
1188 resource_id: &str,
1189 ) -> Result<Option<crate::session_resource::SessionResourceEntry>>;
1190
1191 async fn list(
1193 &self,
1194 session_id: SessionId,
1195 filter: Option<&crate::session_resource::SessionResourceFilter>,
1196 ) -> Result<Vec<crate::session_resource::SessionResourceEntry>>;
1197
1198 async fn deregister(&self, session_id: SessionId, resource_id: &str) -> Result<bool>;
1200}
1201
1202#[async_trait]
1212pub trait LeasedResourceStore: Send + Sync {
1213 async fn upsert_resource(&self, input: UpsertLeasedResource) -> Result<LeasedResource>;
1219
1220 async fn release_resource(
1226 &self,
1227 session_id: SessionId,
1228 provider: &str,
1229 resource_type: &str,
1230 external_id: &str,
1231 ) -> Result<Option<LeasedResource>>;
1232
1233 async fn list_resources(&self, session_id: SessionId) -> Result<Vec<LeasedResource>>;
1238}
1239
1240pub const DEFAULT_MAX_SUBAGENT_DEPTH: u32 = 2;
1247pub const DEFAULT_MAX_ACTIVE_DESCENDANT_SUBAGENT_TASKS: u32 = 16;
1248pub const DEFAULT_MAX_TOTAL_DESCENDANT_SUBAGENT_TASKS: u32 = 200;
1249pub const DEFAULT_MAX_ACTIVE_DETACHED_TASKS: u32 = 8;
1255pub const DEFAULT_MAX_TOTAL_DETACHED_TASKS: u32 = 50;
1256
1257#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1259pub struct SubagentNestingPolicy {
1260 pub platform_default: u32,
1261 pub org_override: Option<u32>,
1262 pub agent_override: Option<u32>,
1263 pub platform_default_max_active_descendant_tasks: u32,
1264 pub org_override_max_active_descendant_tasks: Option<u32>,
1265 pub agent_override_max_active_descendant_tasks: Option<u32>,
1266 pub platform_default_max_total_descendant_tasks: u32,
1267 pub org_override_max_total_descendant_tasks: Option<u32>,
1268 pub agent_override_max_total_descendant_tasks: Option<u32>,
1269 pub platform_default_max_active_detached_tasks: u32,
1270 pub org_override_max_active_detached_tasks: Option<u32>,
1271 pub agent_override_max_active_detached_tasks: Option<u32>,
1272 pub platform_default_max_total_detached_tasks: u32,
1273 pub org_override_max_total_detached_tasks: Option<u32>,
1274 pub agent_override_max_total_detached_tasks: Option<u32>,
1275}
1276
1277impl Default for SubagentNestingPolicy {
1278 fn default() -> Self {
1279 Self {
1280 platform_default: DEFAULT_MAX_SUBAGENT_DEPTH,
1281 org_override: None,
1282 agent_override: None,
1283 platform_default_max_active_descendant_tasks:
1284 DEFAULT_MAX_ACTIVE_DESCENDANT_SUBAGENT_TASKS,
1285 org_override_max_active_descendant_tasks: None,
1286 agent_override_max_active_descendant_tasks: None,
1287 platform_default_max_total_descendant_tasks:
1288 DEFAULT_MAX_TOTAL_DESCENDANT_SUBAGENT_TASKS,
1289 org_override_max_total_descendant_tasks: None,
1290 agent_override_max_total_descendant_tasks: None,
1291 platform_default_max_active_detached_tasks: DEFAULT_MAX_ACTIVE_DETACHED_TASKS,
1292 org_override_max_active_detached_tasks: None,
1293 agent_override_max_active_detached_tasks: None,
1294 platform_default_max_total_detached_tasks: DEFAULT_MAX_TOTAL_DETACHED_TASKS,
1295 org_override_max_total_detached_tasks: None,
1296 agent_override_max_total_detached_tasks: None,
1297 }
1298 }
1299}
1300
1301impl SubagentNestingPolicy {
1302 pub fn max_subagent_depth(self) -> u32 {
1303 self.agent_override
1304 .or(self.org_override)
1305 .unwrap_or(self.platform_default)
1306 }
1307
1308 pub fn max_active_descendant_tasks(self) -> u32 {
1309 self.agent_override_max_active_descendant_tasks
1310 .or(self.org_override_max_active_descendant_tasks)
1311 .unwrap_or(self.platform_default_max_active_descendant_tasks)
1312 }
1313
1314 pub fn max_total_descendant_tasks(self) -> u32 {
1315 self.agent_override_max_total_descendant_tasks
1316 .or(self.org_override_max_total_descendant_tasks)
1317 .unwrap_or(self.platform_default_max_total_descendant_tasks)
1318 }
1319
1320 pub fn max_active_detached_tasks(self) -> u32 {
1321 self.agent_override_max_active_detached_tasks
1322 .or(self.org_override_max_active_detached_tasks)
1323 .unwrap_or(self.platform_default_max_active_detached_tasks)
1324 }
1325
1326 pub fn max_total_detached_tasks(self) -> u32 {
1327 self.agent_override_max_total_detached_tasks
1328 .or(self.org_override_max_total_detached_tasks)
1329 .unwrap_or(self.platform_default_max_total_detached_tasks)
1330 }
1331
1332 pub fn with_platform_default(mut self, depth: u32) -> Self {
1333 self.platform_default = depth;
1334 self
1335 }
1336
1337 pub fn with_org_override(mut self, depth: Option<u32>) -> Self {
1338 self.org_override = depth;
1339 self
1340 }
1341
1342 pub fn with_agent_override(mut self, depth: Option<u32>) -> Self {
1343 self.agent_override = depth;
1344 self
1345 }
1346
1347 pub fn with_agent_task_caps_override(
1348 mut self,
1349 max_active: Option<u32>,
1350 max_total: Option<u32>,
1351 ) -> Self {
1352 self.agent_override_max_active_descendant_tasks = max_active;
1353 self.agent_override_max_total_descendant_tasks = max_total;
1354 self
1355 }
1356
1357 pub fn with_agent_detached_task_caps_override(
1358 mut self,
1359 max_active: Option<u32>,
1360 max_total: Option<u32>,
1361 ) -> Self {
1362 self.agent_override_max_active_detached_tasks = max_active;
1363 self.agent_override_max_total_detached_tasks = max_total;
1364 self
1365 }
1366}
1367
1368pub type SessionSqlDbStoreRef = Arc<dyn crate::session_sqldb::SessionSqlDbStore>;
1370
1371#[async_trait]
1376pub trait UserConnectionResolver: Send + Sync {
1377 async fn get_connection_token(
1380 &self,
1381 session_id: SessionId,
1382 provider: &str,
1383 ) -> Result<Option<String>>;
1384
1385 async fn get_connection_user(
1390 &self,
1391 _session_id: SessionId,
1392 _provider: &str,
1393 ) -> Result<Option<Uuid>> {
1394 Ok(None)
1395 }
1396
1397 async fn get_connection_token_for_user(
1402 &self,
1403 _user_id: Uuid,
1404 _provider: &str,
1405 ) -> Result<Option<String>> {
1406 Ok(None)
1407 }
1408
1409 async fn get_connection_metadata(
1412 &self,
1413 _session_id: SessionId,
1414 _provider: &str,
1415 ) -> Result<Option<serde_json::Value>> {
1416 Ok(None)
1417 }
1418}
1419
1420#[async_trait]
1430pub trait BudgetChecker: Send + Sync {
1431 async fn check_budgets(&self, session_id: &str) -> Result<crate::budget::BudgetToolResponse>;
1433}
1434
1435#[async_trait]
1444pub trait PaymentAuthority: Send + Sync {
1445 async fn execute_machine_payment(
1446 &self,
1447 session_id: SessionId,
1448 request: crate::payment::MachinePaymentRequest,
1449 ) -> Result<crate::payment::MachinePaymentResponse>;
1450}
1451
1452#[async_trait]
1462pub trait SessionCreationAuthority: Send + Sync {
1463 async fn authorize_session_creation(&self, session_id: SessionId) -> Result<SessionId>;
1467}
1468
1469#[async_trait]
1479pub trait OutboundToolRateLimiter: Send + Sync {
1480 async fn check_org(&self, org_id: &crate::typed_id::OrgId) -> bool;
1482}
1483
1484#[derive(Debug)]
1490pub enum ToolCallClaimResult {
1491 Claimed { claim_token: uuid::Uuid },
1494 AlreadySettled {
1496 result_json: serde_json::Value,
1497 args_fingerprint: String,
1498 },
1499 AlreadyRunning { args_fingerprint: String },
1504 DeterminismViolation {
1508 stored_fingerprint: String,
1509 current_fingerprint: String,
1510 },
1511}
1512
1513#[derive(Debug, Clone)]
1515pub enum DurableToolCallStatus {
1516 Settled { result_json: serde_json::Value },
1518 Interrupted {
1520 result_json: Option<serde_json::Value>,
1521 },
1522 Running,
1524}
1525
1526#[async_trait]
1531pub trait DurableToolResultStore: Send + Sync + 'static {
1532 async fn try_claim_tool_call(
1540 &self,
1541 turn_id: &str,
1542 tool_call_id: &str,
1543 tool_name: &str,
1544 args_fingerprint: &str,
1545 ) -> Result<ToolCallClaimResult>;
1546
1547 async fn settle_tool_call(
1553 &self,
1554 turn_id: &str,
1555 tool_call_id: &str,
1556 result_json: serde_json::Value,
1557 status: &str,
1558 claim_token: uuid::Uuid,
1559 ) -> Result<bool>;
1560
1561 async fn get_tool_call_status(
1566 &self,
1567 turn_id: &str,
1568 tool_call_id: &str,
1569 ) -> Result<Option<DurableToolCallStatus>>;
1570}
1571
1572pub struct NoopDurableToolResultStore;
1575
1576#[async_trait]
1577impl DurableToolResultStore for NoopDurableToolResultStore {
1578 async fn try_claim_tool_call(
1579 &self,
1580 _turn_id: &str,
1581 _tool_call_id: &str,
1582 _tool_name: &str,
1583 _args_fingerprint: &str,
1584 ) -> Result<ToolCallClaimResult> {
1585 Ok(ToolCallClaimResult::Claimed {
1586 claim_token: uuid::Uuid::new_v4(),
1587 })
1588 }
1589
1590 async fn settle_tool_call(
1591 &self,
1592 _turn_id: &str,
1593 _tool_call_id: &str,
1594 _result_json: serde_json::Value,
1595 _status: &str,
1596 _claim_token: uuid::Uuid,
1597 ) -> Result<bool> {
1598 Ok(true)
1599 }
1600
1601 async fn get_tool_call_status(
1602 &self,
1603 _turn_id: &str,
1604 _tool_call_id: &str,
1605 ) -> Result<Option<DurableToolCallStatus>> {
1606 Ok(None)
1607 }
1608}
1609
1610#[derive(Debug, Clone)]
1616pub struct StreamProgress {
1617 pub accumulated_len: usize,
1619 pub last_delta_at: u64,
1621}
1622
1623#[async_trait]
1629pub trait StreamHeartbeater: Send + Sync {
1630 async fn heartbeat(&self, progress: StreamProgress);
1636}
1637
1638pub struct NoopStreamHeartbeater;
1640
1641#[async_trait]
1642impl StreamHeartbeater for NoopStreamHeartbeater {
1643 async fn heartbeat(&self, _progress: StreamProgress) {}
1644}
1645
1646#[derive(Debug, Clone)]
1652pub struct PartialStreamState {
1653 pub message_id: MessageId,
1655
1656 pub accumulated: String,
1659}
1660
1661#[async_trait]
1669pub trait PartialStreamStore: Send + Sync {
1670 async fn get_partial_stream(
1673 &self,
1674 session_id: SessionId,
1675 turn_id: &str,
1676 ) -> Result<Option<PartialStreamState>>;
1677}
1678
1679pub struct NoopPartialStreamStore;
1681
1682#[async_trait]
1683impl PartialStreamStore for NoopPartialStreamStore {
1684 async fn get_partial_stream(
1685 &self,
1686 _session_id: SessionId,
1687 _turn_id: &str,
1688 ) -> Result<Option<PartialStreamState>> {
1689 Ok(None)
1690 }
1691}
1692
1693#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
1699pub enum ToolContextService {
1700 SessionFileSystem,
1701 SessionStorageStore,
1702 ImageArtifactStore,
1703 ProviderCredentialStore,
1704 UtilityLlmService,
1705 McpInvoker,
1706 EgressService,
1707 SessionSqlDbStore,
1708 MessageRetriever,
1709 SessionStore,
1710 SessionMutator,
1711 AgentStore,
1712 ConnectionResolver,
1713 ScheduleStore,
1714 SubagentSessionDelegate,
1715 KnowledgeStore,
1716 KnowledgeIndexSearch,
1717 LeasedResourceStore,
1718 SessionResourceRegistry,
1719 SessionTaskRegistry,
1720 EventEmitter,
1721 CapabilityRegistry,
1722 ToolRegistry,
1723 OrgId,
1724 BudgetChecker,
1725 PaymentAuthority,
1726 SessionCreationAuthority,
1727 SubagentSpawnStore,
1728 ReasoningEffortHandle,
1729}
1730
1731impl ToolContextService {
1732 pub const fn name(self) -> &'static str {
1733 match self {
1734 Self::SessionFileSystem => "SessionFileSystem",
1735 Self::SessionStorageStore => "SessionStorageStore",
1736 Self::ImageArtifactStore => "ImageArtifactStore",
1737 Self::ProviderCredentialStore => "ProviderCredentialStore",
1738 Self::UtilityLlmService => "UtilityLlmService",
1739 Self::McpInvoker => "McpInvoker",
1740 Self::EgressService => "EgressService",
1741 Self::SessionSqlDbStore => "SessionSqlDbStore",
1742 Self::MessageRetriever => "MessageRetriever",
1743 Self::SessionStore => "SessionStore",
1744 Self::SessionMutator => "SessionMutator",
1745 Self::AgentStore => "AgentStore",
1746 Self::ConnectionResolver => "ConnectionResolver",
1747 Self::ScheduleStore => "SessionScheduleStore",
1748 Self::SubagentSessionDelegate => "SubagentSessionDelegate",
1749 Self::KnowledgeStore => "KnowledgeStore",
1750 Self::KnowledgeIndexSearch => "KnowledgeIndexSearch",
1751 Self::LeasedResourceStore => "LeasedResourceStore",
1752 Self::SessionResourceRegistry => "SessionResourceRegistry",
1753 Self::SessionTaskRegistry => "SessionTaskRegistry",
1754 Self::EventEmitter => "EventEmitter",
1755 Self::CapabilityRegistry => "CapabilityRegistry",
1756 Self::ToolRegistry => "ToolRegistry",
1757 Self::OrgId => "OrgId",
1758 Self::BudgetChecker => "BudgetChecker",
1759 Self::PaymentAuthority => "PaymentAuthority",
1760 Self::SessionCreationAuthority => "SessionCreationAuthority",
1761 Self::SubagentSpawnStore => "SubagentSpawnStore",
1762 Self::ReasoningEffortHandle => "ReasoningEffortHandle",
1763 }
1764 }
1765}
1766
1767#[derive(Clone, Default)]
1772pub struct ToolContextServices {
1773 pub file_store: Option<Arc<dyn SessionFileSystem>>,
1774 pub storage_store: Option<Arc<dyn SessionStorageStore>>,
1775 pub image_store: Option<Arc<dyn ImageArtifactStore>>,
1776 pub provider_credential_store: Option<Arc<dyn ProviderCredentialStore>>,
1777 pub utility_llm_service: Option<Arc<dyn crate::UtilityLlmService>>,
1778 pub mcp_invoker: Option<Arc<dyn crate::McpToolInvoker>>,
1779 pub egress_service: Option<Arc<dyn crate::EgressService>>,
1780 pub sqldb_store: Option<SessionSqlDbStoreRef>,
1781 pub message_retriever: Option<Arc<dyn crate::message_retriever::MessageRetriever>>,
1782 pub session_store: Option<Arc<dyn SessionStore>>,
1783 pub session_mutator: Option<Arc<dyn SessionMutator>>,
1784 pub agent_store: Option<Arc<dyn AgentStore>>,
1785 pub connection_resolver: Option<Arc<dyn UserConnectionResolver>>,
1786 pub schedule_store: Option<Arc<dyn SessionScheduleStore>>,
1787 pub subagent_delegate: Option<Arc<dyn crate::subagent_delegation::SubagentSessionDelegate>>,
1788 pub extensions: ToolContextExtensions,
1789 pub knowledge_store: Option<Arc<dyn KnowledgeStore>>,
1790 pub knowledge_index_search: Option<Arc<dyn crate::vector_store::KnowledgeIndexSearch>>,
1791 pub leased_resource_store: Option<Arc<dyn LeasedResourceStore>>,
1792 pub session_resource_registry: Option<Arc<dyn SessionResourceRegistry>>,
1793 pub session_task_registry: Option<Arc<dyn crate::session_task::SessionTaskRegistry>>,
1794 pub event_emitter: Option<Arc<dyn EventEmitter>>,
1795 pub capability_registry: Option<crate::capabilities::CapabilityRegistry>,
1796 pub tool_registry: Option<Arc<crate::tools::ToolRegistry>>,
1797 pub org_id: Option<crate::typed_id::OrgId>,
1798 pub network_access: Option<crate::network_access::NetworkAccessList>,
1799 pub budget_checker: Option<Arc<dyn BudgetChecker>>,
1800 pub payment_authority: Option<Arc<dyn PaymentAuthority>>,
1801 pub session_creation_authority: Option<Arc<dyn SessionCreationAuthority>>,
1802 pub subagent_spawn_store: Option<Arc<dyn SubagentSpawnStore>>,
1803 pub subagent_nesting_policy: SubagentNestingPolicy,
1804 pub reasoning_effort_handle: Option<ReasoningEffortHandle>,
1805}
1806
1807impl ToolContextServices {
1808 pub fn provides(&self, service: ToolContextService) -> bool {
1809 match service {
1810 ToolContextService::SessionFileSystem => self.file_store.is_some(),
1811 ToolContextService::SessionStorageStore => self.storage_store.is_some(),
1812 ToolContextService::ImageArtifactStore => self.image_store.is_some(),
1813 ToolContextService::ProviderCredentialStore => self.provider_credential_store.is_some(),
1814 ToolContextService::UtilityLlmService => self.utility_llm_service.is_some(),
1815 ToolContextService::McpInvoker => self.mcp_invoker.is_some(),
1816 ToolContextService::EgressService => self.egress_service.is_some(),
1817 ToolContextService::SessionSqlDbStore => self.sqldb_store.is_some(),
1818 ToolContextService::MessageRetriever => self.message_retriever.is_some(),
1819 ToolContextService::SessionStore => self.session_store.is_some(),
1820 ToolContextService::SessionMutator => self.session_mutator.is_some(),
1821 ToolContextService::AgentStore => self.agent_store.is_some(),
1822 ToolContextService::ConnectionResolver => self.connection_resolver.is_some(),
1823 ToolContextService::ScheduleStore => self.schedule_store.is_some(),
1824 ToolContextService::SubagentSessionDelegate => self.subagent_delegate.is_some(),
1825 ToolContextService::KnowledgeStore => self.knowledge_store.is_some(),
1826 ToolContextService::KnowledgeIndexSearch => self.knowledge_index_search.is_some(),
1827 ToolContextService::LeasedResourceStore => self.leased_resource_store.is_some(),
1828 ToolContextService::SessionResourceRegistry => self.session_resource_registry.is_some(),
1829 ToolContextService::SessionTaskRegistry => self.session_task_registry.is_some(),
1830 ToolContextService::EventEmitter => self.event_emitter.is_some(),
1831 ToolContextService::CapabilityRegistry => self.capability_registry.is_some(),
1832 ToolContextService::ToolRegistry => self.tool_registry.is_some(),
1833 ToolContextService::OrgId => self.org_id.is_some(),
1834 ToolContextService::BudgetChecker => self.budget_checker.is_some(),
1835 ToolContextService::PaymentAuthority => self.payment_authority.is_some(),
1836 ToolContextService::SessionCreationAuthority => {
1837 self.session_creation_authority.is_some()
1838 }
1839 ToolContextService::SubagentSpawnStore => self.subagent_spawn_store.is_some(),
1840 ToolContextService::ReasoningEffortHandle => self.reasoning_effort_handle.is_some(),
1841 }
1842 }
1843}
1844
1845#[derive(Clone, Default)]
1850pub struct ToolContextExtensions {
1851 values: Arc<HashMap<TypeId, Arc<dyn Any + Send + Sync>>>,
1852}
1853
1854impl ToolContextExtensions {
1855 pub fn insert<T: Any + Send + Sync>(&mut self, value: Arc<T>) {
1857 Arc::make_mut(&mut self.values).insert(TypeId::of::<T>(), value);
1858 }
1859
1860 pub fn get<T: Any + Send + Sync>(&self) -> Option<Arc<T>> {
1862 self.values
1863 .get(&TypeId::of::<T>())
1864 .and_then(|value| value.clone().downcast::<T>().ok())
1865 }
1866
1867 pub fn is_empty(&self) -> bool {
1868 self.values.is_empty()
1869 }
1870}
1871
1872impl std::fmt::Debug for ToolContextExtensions {
1873 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1874 f.debug_struct("ToolContextExtensions")
1875 .field("len", &self.values.len())
1876 .finish()
1877 }
1878}
1879
1880#[derive(Clone)]
1889pub struct ToolContext {
1890 pub session_id: SessionId,
1892 pub workspace_id: WorkspaceId,
1899
1900 pub file_store: Option<Arc<dyn SessionFileSystem>>,
1902
1903 pub storage_store: Option<Arc<dyn SessionStorageStore>>,
1905
1906 pub image_store: Option<Arc<dyn ImageArtifactStore>>,
1908
1909 pub provider_credential_store: Option<Arc<dyn ProviderCredentialStore>>,
1911
1912 pub utility_llm_service: Option<Arc<dyn crate::UtilityLlmService>>,
1914
1915 pub mcp_invoker: Option<Arc<dyn crate::McpToolInvoker>>,
1921
1922 pub egress_service: Option<Arc<dyn crate::EgressService>>,
1924
1925 pub sqldb_store: Option<SessionSqlDbStoreRef>,
1927
1928 pub message_retriever: Option<Arc<dyn crate::message_retriever::MessageRetriever>>,
1930
1931 pub session_store: Option<Arc<dyn SessionStore>>,
1933
1934 pub session_mutator: Option<Arc<dyn SessionMutator>>,
1936
1937 pub agent_store: Option<Arc<dyn AgentStore>>,
1939
1940 pub connection_resolver: Option<Arc<dyn UserConnectionResolver>>,
1942
1943 pub schedule_store: Option<Arc<dyn SessionScheduleStore>>,
1945
1946 pub subagent_delegate: Option<Arc<dyn crate::subagent_delegation::SubagentSessionDelegate>>,
1951 pub extensions: ToolContextExtensions,
1955 pub knowledge_store: Option<Arc<dyn KnowledgeStore>>,
1957
1958 pub knowledge_index_search: Option<Arc<dyn crate::vector_store::KnowledgeIndexSearch>>,
1962
1963 pub leased_resource_store: Option<Arc<dyn LeasedResourceStore>>,
1965
1966 pub session_resource_registry: Option<Arc<dyn SessionResourceRegistry>>,
1968
1969 pub session_task_registry: Option<Arc<dyn crate::session_task::SessionTaskRegistry>>,
1972
1973 pub event_emitter: Option<Arc<dyn EventEmitter>>,
1976
1977 pub event_context: Option<crate::events::EventContext>,
1980
1981 pub tool_call_id: Option<String>,
1984 pub capability_registry: Option<crate::capabilities::CapabilityRegistry>,
1986
1987 pub tool_registry: Option<Arc<crate::tools::ToolRegistry>>,
1990
1991 pub visible_tool_names: Option<Arc<HashSet<String>>>,
1995
1996 pub org_id: Option<crate::typed_id::OrgId>,
1998
1999 pub network_access: Option<crate::network_access::NetworkAccessList>,
2002
2003 pub locale: Option<String>,
2007
2008 pub budget_checker: Option<Arc<dyn BudgetChecker>>,
2010
2011 pub payment_authority: Option<Arc<dyn PaymentAuthority>>,
2013
2014 pub session_creation_authority: Option<Arc<dyn SessionCreationAuthority>>,
2016
2017 pub subagent_spawn_store: Option<Arc<dyn SubagentSpawnStore>>,
2021
2022 pub subagent_nesting_policy: SubagentNestingPolicy,
2024
2025 pub reasoning_effort_handle: Option<ReasoningEffortHandle>,
2029
2030 pub cancellation: Option<tokio_util::sync::CancellationToken>,
2040}
2041
2042impl ToolContext {
2043 pub fn workspace_fs_key(&self) -> SessionId {
2048 SessionId::from_uuid(self.workspace_id.uuid())
2049 }
2050
2051 pub fn with_workspace_id(mut self, workspace_id: WorkspaceId) -> Self {
2053 self.workspace_id = workspace_id;
2054 self
2055 }
2056
2057 pub fn new(session_id: SessionId) -> Self {
2059 Self {
2060 session_id,
2061 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
2062 file_store: None,
2063 storage_store: None,
2064 image_store: None,
2065 provider_credential_store: None,
2066 utility_llm_service: None,
2067 mcp_invoker: None,
2068 egress_service: None,
2069 sqldb_store: None,
2070 message_retriever: None,
2071 session_store: None,
2072 session_mutator: None,
2073 agent_store: None,
2074 connection_resolver: None,
2075 schedule_store: None,
2076 subagent_delegate: None,
2077 extensions: ToolContextExtensions::default(),
2078 knowledge_store: None,
2079 knowledge_index_search: None,
2080 leased_resource_store: None,
2081 session_resource_registry: None,
2082 session_task_registry: None,
2083 event_emitter: None,
2084 event_context: None,
2085 tool_call_id: None,
2086 capability_registry: None,
2087 tool_registry: None,
2088 visible_tool_names: None,
2089 org_id: None,
2090 network_access: None,
2091 locale: None,
2092 budget_checker: None,
2093 payment_authority: None,
2094 session_creation_authority: None,
2095 subagent_spawn_store: None,
2096 subagent_nesting_policy: SubagentNestingPolicy::default(),
2097 reasoning_effort_handle: None,
2098 cancellation: None,
2099 }
2100 }
2101
2102 pub fn from_services(session_id: SessionId, services: &ToolContextServices) -> Self {
2104 Self {
2105 session_id,
2106 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
2107 file_store: services.file_store.clone(),
2108 storage_store: services.storage_store.clone(),
2109 image_store: services.image_store.clone(),
2110 provider_credential_store: services.provider_credential_store.clone(),
2111 utility_llm_service: services.utility_llm_service.clone(),
2112 mcp_invoker: services.mcp_invoker.clone(),
2113 egress_service: services.egress_service.clone(),
2114 sqldb_store: services.sqldb_store.clone(),
2115 message_retriever: services.message_retriever.clone(),
2116 session_store: services.session_store.clone(),
2117 session_mutator: services.session_mutator.clone(),
2118 agent_store: services.agent_store.clone(),
2119 connection_resolver: services.connection_resolver.clone(),
2120 schedule_store: services.schedule_store.clone(),
2121 subagent_delegate: services.subagent_delegate.clone(),
2122 extensions: services.extensions.clone(),
2123 knowledge_store: services.knowledge_store.clone(),
2124 knowledge_index_search: services.knowledge_index_search.clone(),
2125 leased_resource_store: services.leased_resource_store.clone(),
2126 session_resource_registry: services.session_resource_registry.clone(),
2127 session_task_registry: services.session_task_registry.clone(),
2128 event_emitter: services.event_emitter.clone(),
2129 event_context: None,
2130 tool_call_id: None,
2131 capability_registry: services.capability_registry.clone(),
2132 tool_registry: services.tool_registry.clone(),
2133 visible_tool_names: None,
2134 org_id: services.org_id,
2135 network_access: services.network_access.clone(),
2136 locale: None,
2137 budget_checker: services.budget_checker.clone(),
2138 payment_authority: services.payment_authority.clone(),
2139 session_creation_authority: services.session_creation_authority.clone(),
2140 subagent_spawn_store: services.subagent_spawn_store.clone(),
2141 subagent_nesting_policy: services.subagent_nesting_policy,
2142 reasoning_effort_handle: services.reasoning_effort_handle.clone(),
2143 cancellation: None,
2144 }
2145 }
2146
2147 pub fn with_file_store(session_id: SessionId, file_store: Arc<dyn SessionFileSystem>) -> Self {
2149 Self {
2150 session_id,
2151 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
2152 file_store: Some(file_store),
2153 storage_store: None,
2154 image_store: None,
2155 provider_credential_store: None,
2156 utility_llm_service: None,
2157 mcp_invoker: None,
2158 egress_service: None,
2159 sqldb_store: None,
2160 message_retriever: None,
2161 session_store: None,
2162 session_mutator: None,
2163 agent_store: None,
2164 connection_resolver: None,
2165 schedule_store: None,
2166 subagent_delegate: None,
2167 extensions: ToolContextExtensions::default(),
2168 knowledge_store: None,
2169 knowledge_index_search: None,
2170 leased_resource_store: None,
2171 session_resource_registry: None,
2172 session_task_registry: None,
2173 event_emitter: None,
2174 event_context: None,
2175 tool_call_id: None,
2176 capability_registry: None,
2177 tool_registry: None,
2178 visible_tool_names: None,
2179 org_id: None,
2180 network_access: None,
2181 locale: None,
2182 budget_checker: None,
2183 payment_authority: None,
2184 session_creation_authority: None,
2185 subagent_spawn_store: None,
2186 subagent_nesting_policy: SubagentNestingPolicy::default(),
2187 reasoning_effort_handle: None,
2188 cancellation: None,
2189 }
2190 }
2191
2192 pub fn with_storage_store(
2194 session_id: SessionId,
2195 storage_store: Arc<dyn SessionStorageStore>,
2196 ) -> Self {
2197 Self {
2198 session_id,
2199 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
2200 file_store: None,
2201 storage_store: Some(storage_store),
2202 image_store: None,
2203 provider_credential_store: None,
2204 utility_llm_service: None,
2205 mcp_invoker: None,
2206 egress_service: None,
2207 sqldb_store: None,
2208 message_retriever: None,
2209 session_store: None,
2210 session_mutator: None,
2211 agent_store: None,
2212 connection_resolver: None,
2213 schedule_store: None,
2214 subagent_delegate: None,
2215 extensions: ToolContextExtensions::default(),
2216 knowledge_store: None,
2217 knowledge_index_search: None,
2218 leased_resource_store: None,
2219 session_resource_registry: None,
2220 session_task_registry: None,
2221 event_emitter: None,
2222 event_context: None,
2223 tool_call_id: None,
2224 capability_registry: None,
2225 tool_registry: None,
2226 visible_tool_names: None,
2227 org_id: None,
2228 network_access: None,
2229 locale: None,
2230 budget_checker: None,
2231 payment_authority: None,
2232 session_creation_authority: None,
2233 subagent_spawn_store: None,
2234 subagent_nesting_policy: SubagentNestingPolicy::default(),
2235 reasoning_effort_handle: None,
2236 cancellation: None,
2237 }
2238 }
2239
2240 pub fn with_stores(
2242 session_id: SessionId,
2243 file_store: Arc<dyn SessionFileSystem>,
2244 storage_store: Arc<dyn SessionStorageStore>,
2245 ) -> Self {
2246 Self {
2247 session_id,
2248 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
2249 file_store: Some(file_store),
2250 storage_store: Some(storage_store),
2251 sqldb_store: None,
2252 image_store: None,
2253 provider_credential_store: None,
2254 utility_llm_service: None,
2255 mcp_invoker: None,
2256 egress_service: None,
2257 message_retriever: None,
2258 session_store: None,
2259 session_mutator: None,
2260 agent_store: None,
2261 connection_resolver: None,
2262 schedule_store: None,
2263 subagent_delegate: None,
2264 extensions: ToolContextExtensions::default(),
2265 knowledge_store: None,
2266 knowledge_index_search: None,
2267 leased_resource_store: None,
2268 session_resource_registry: None,
2269 session_task_registry: None,
2270 event_emitter: None,
2271 event_context: None,
2272 tool_call_id: None,
2273 capability_registry: None,
2274 tool_registry: None,
2275 visible_tool_names: None,
2276 org_id: None,
2277 network_access: None,
2278 locale: None,
2279 budget_checker: None,
2280 payment_authority: None,
2281 session_creation_authority: None,
2282 subagent_spawn_store: None,
2283 subagent_nesting_policy: SubagentNestingPolicy::default(),
2284 reasoning_effort_handle: None,
2285 cancellation: None,
2286 }
2287 }
2288
2289 pub fn with_sqldb_store(mut self, sqldb_store: SessionSqlDbStoreRef) -> Self {
2291 self.sqldb_store = Some(sqldb_store);
2292 self
2293 }
2294
2295 pub fn with_message_retriever(
2297 mut self,
2298 retriever: Arc<dyn crate::message_retriever::MessageRetriever>,
2299 ) -> Self {
2300 self.message_retriever = Some(retriever);
2301 self
2302 }
2303
2304 pub fn with_session_store(mut self, store: Arc<dyn SessionStore>) -> Self {
2306 self.session_store = Some(store);
2307 self
2308 }
2309
2310 pub fn with_session_mutator(mut self, mutator: Arc<dyn SessionMutator>) -> Self {
2312 self.session_mutator = Some(mutator);
2313 self
2314 }
2315
2316 pub fn with_cancellation(mut self, token: tokio_util::sync::CancellationToken) -> Self {
2322 self.cancellation = Some(token);
2323 self
2324 }
2325
2326 pub fn is_cancelled(&self) -> bool {
2329 self.cancellation
2330 .as_ref()
2331 .is_some_and(|token| token.is_cancelled())
2332 }
2333
2334 pub fn with_reasoning_effort_handle(mut self, handle: ReasoningEffortHandle) -> Self {
2335 self.reasoning_effort_handle = Some(handle);
2336 self
2337 }
2338
2339 pub fn with_agent_store(mut self, store: Arc<dyn AgentStore>) -> Self {
2341 self.agent_store = Some(store);
2342 self
2343 }
2344
2345 pub fn with_connection_resolver(mut self, resolver: Arc<dyn UserConnectionResolver>) -> Self {
2347 self.connection_resolver = Some(resolver);
2348 self
2349 }
2350
2351 pub fn with_image_store(
2353 session_id: SessionId,
2354 image_store: Arc<dyn ImageArtifactStore>,
2355 ) -> Self {
2356 Self {
2357 session_id,
2358 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
2359 file_store: None,
2360 storage_store: None,
2361 image_store: Some(image_store),
2362 provider_credential_store: None,
2363 utility_llm_service: None,
2364 mcp_invoker: None,
2365 egress_service: None,
2366 sqldb_store: None,
2367 message_retriever: None,
2368 session_store: None,
2369 session_mutator: None,
2370 agent_store: None,
2371 connection_resolver: None,
2372 schedule_store: None,
2373 subagent_delegate: None,
2374 extensions: ToolContextExtensions::default(),
2375 knowledge_store: None,
2376 knowledge_index_search: None,
2377 leased_resource_store: None,
2378 session_resource_registry: None,
2379 session_task_registry: None,
2380 event_emitter: None,
2381 event_context: None,
2382 tool_call_id: None,
2383 capability_registry: None,
2384 tool_registry: None,
2385 visible_tool_names: None,
2386 org_id: None,
2387 network_access: None,
2388 locale: None,
2389 budget_checker: None,
2390 payment_authority: None,
2391 session_creation_authority: None,
2392 subagent_spawn_store: None,
2393 subagent_nesting_policy: SubagentNestingPolicy::default(),
2394 reasoning_effort_handle: None,
2395 cancellation: None,
2396 }
2397 }
2398
2399 pub fn with_provider_credential_store(
2401 mut self,
2402 store: Arc<dyn ProviderCredentialStore>,
2403 ) -> Self {
2404 self.provider_credential_store = Some(store);
2405 self
2406 }
2407
2408 pub fn with_utility_llm_service(mut self, service: Arc<dyn crate::UtilityLlmService>) -> Self {
2410 self.utility_llm_service = Some(service);
2411 self
2412 }
2413
2414 pub fn with_mcp_invoker(mut self, invoker: Arc<dyn crate::McpToolInvoker>) -> Self {
2416 self.mcp_invoker = Some(invoker);
2417 self
2418 }
2419
2420 pub fn with_egress_service(mut self, service: Arc<dyn crate::EgressService>) -> Self {
2422 self.egress_service = Some(service);
2423 self
2424 }
2425
2426 pub fn with_egress_service_opt(
2429 mut self,
2430 service: Option<Arc<dyn crate::EgressService>>,
2431 ) -> Self {
2432 if let Some(service) = service {
2433 self.egress_service = Some(service);
2434 }
2435 self
2436 }
2437
2438 pub fn with_storage_store_arc(mut self, store: Arc<dyn SessionStorageStore>) -> Self {
2440 self.storage_store = Some(store);
2441 self
2442 }
2443
2444 pub fn with_schedule_store(mut self, store: Arc<dyn SessionScheduleStore>) -> Self {
2446 self.schedule_store = Some(store);
2447 self
2448 }
2449
2450 pub fn with_subagent_delegate(
2452 mut self,
2453 delegate: Arc<dyn crate::subagent_delegation::SubagentSessionDelegate>,
2454 ) -> Self {
2455 self.subagent_delegate = Some(delegate);
2456 self
2457 }
2458
2459 pub fn with_extension<T: std::any::Any + Send + Sync>(mut self, value: Arc<T>) -> Self {
2462 self.extensions.insert(value);
2463 self
2464 }
2465
2466 pub fn extension<T: std::any::Any + Send + Sync>(&self) -> Option<Arc<T>> {
2469 self.extensions.get::<T>()
2470 }
2471
2472 pub fn with_knowledge_index_search(
2474 mut self,
2475 search: Arc<dyn crate::vector_store::KnowledgeIndexSearch>,
2476 ) -> Self {
2477 self.knowledge_index_search = Some(search);
2478 self
2479 }
2480
2481 pub fn with_leased_resource_store(mut self, store: Arc<dyn LeasedResourceStore>) -> Self {
2483 self.leased_resource_store = Some(store);
2484 self
2485 }
2486
2487 pub fn with_session_resource_registry(
2489 mut self,
2490 registry: Arc<dyn SessionResourceRegistry>,
2491 ) -> Self {
2492 self.session_resource_registry = Some(registry);
2493 self
2494 }
2495
2496 pub fn with_session_task_registry(
2498 mut self,
2499 registry: Arc<dyn crate::session_task::SessionTaskRegistry>,
2500 ) -> Self {
2501 self.session_task_registry = Some(registry);
2502 self
2503 }
2504
2505 pub fn with_org_id(mut self, org_id: crate::typed_id::OrgId) -> Self {
2507 self.org_id = Some(org_id);
2508 self
2509 }
2510
2511 pub fn with_tool_registry(mut self, registry: Arc<crate::tools::ToolRegistry>) -> Self {
2513 self.tool_registry = Some(registry);
2514 self
2515 }
2516
2517 pub fn with_visible_tool_names(mut self, names: Arc<HashSet<String>>) -> Self {
2519 self.visible_tool_names = Some(names);
2520 self
2521 }
2522
2523 pub fn with_network_access(
2525 mut self,
2526 network_access: Option<crate::network_access::NetworkAccessList>,
2527 ) -> Self {
2528 self.network_access = network_access;
2529 self
2530 }
2531
2532 pub fn with_payment_authority(mut self, authority: Arc<dyn PaymentAuthority>) -> Self {
2534 self.payment_authority = Some(authority);
2535 self
2536 }
2537
2538 pub fn with_subagent_spawn_store(mut self, store: Arc<dyn SubagentSpawnStore>) -> Self {
2540 self.subagent_spawn_store = Some(store);
2541 self
2542 }
2543
2544 pub fn with_subagent_nesting_policy(mut self, policy: SubagentNestingPolicy) -> Self {
2546 self.subagent_nesting_policy = policy;
2547 self
2548 }
2549
2550 pub async fn emit_progress(&self, tool_name: &str, message: &str) {
2555 let (Some(emitter), Some(ctx), Some(call_id)) =
2556 (&self.event_emitter, &self.event_context, &self.tool_call_id)
2557 else {
2558 return;
2559 };
2560 if let Err(e) = emitter
2561 .emit(EventRequest::new(
2562 self.session_id,
2563 ctx.clone(),
2564 crate::events::ToolProgressData {
2565 tool_call_id: call_id.clone(),
2566 tool_name: tool_name.to_string(),
2567 message: message.to_string(),
2568 display_name: None,
2569 },
2570 ))
2571 .await
2572 {
2573 tracing::debug!(
2574 tool_call_id = call_id,
2575 tool_name,
2576 error = %e,
2577 "Failed to emit tool.progress event"
2578 );
2579 }
2580 }
2581
2582 pub async fn emit_tool_output(&self, tool_name: &str, delta: &str, stream: &str) {
2587 let (Some(emitter), Some(ctx), Some(call_id)) =
2588 (&self.event_emitter, &self.event_context, &self.tool_call_id)
2589 else {
2590 return;
2591 };
2592 if let Err(e) = emitter
2593 .emit(EventRequest::new(
2594 self.session_id,
2595 ctx.clone(),
2596 crate::events::ToolOutputDeltaData {
2597 tool_call_id: call_id.clone(),
2598 tool_name: tool_name.to_string(),
2599 delta: delta.to_string(),
2600 stream: stream.to_string(),
2601 },
2602 ))
2603 .await
2604 {
2605 tracing::debug!(
2606 tool_call_id = call_id,
2607 tool_name,
2608 error = %e,
2609 "Failed to emit tool.output.delta event"
2610 );
2611 }
2612 }
2613}
2614
2615impl std::fmt::Debug for ToolContext {
2616 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2617 f.debug_struct("ToolContext")
2618 .field("session_id", &self.session_id)
2619 .field("file_store", &self.file_store.is_some())
2620 .field("storage_store", &self.storage_store.is_some())
2621 .field("image_store", &self.image_store.is_some())
2622 .field(
2623 "provider_credential_store",
2624 &self.provider_credential_store.is_some(),
2625 )
2626 .field("utility_llm_service", &self.utility_llm_service.is_some())
2627 .field("egress_service", &self.egress_service.is_some())
2628 .field("sqldb_store", &self.sqldb_store.is_some())
2629 .field("message_retriever", &self.message_retriever.is_some())
2630 .field("session_store", &self.session_store.is_some())
2631 .field("session_mutator", &self.session_mutator.is_some())
2632 .field("agent_store", &self.agent_store.is_some())
2633 .field("connection_resolver", &self.connection_resolver.is_some())
2634 .field("schedule_store", &self.schedule_store.is_some())
2635 .field("subagent_delegate", &self.subagent_delegate.is_some())
2636 .field(
2637 "knowledge_index_search",
2638 &self.knowledge_index_search.is_some(),
2639 )
2640 .field(
2641 "leased_resource_store",
2642 &self.leased_resource_store.is_some(),
2643 )
2644 .field("event_emitter", &self.event_emitter.is_some())
2645 .field("tool_registry", &self.tool_registry.is_some())
2646 .field("payment_authority", &self.payment_authority.is_some())
2647 .field("subagent_spawn_store", &self.subagent_spawn_store.is_some())
2648 .field("subagent_nesting_policy", &self.subagent_nesting_policy)
2649 .field("org_id", &self.org_id)
2650 .finish()
2651 }
2652}
2653
2654use crate::events::{Event, EventRequest};
2659
2660#[async_trait]
2671pub trait EventEmitter: Send + Sync {
2672 async fn emit(&self, request: EventRequest) -> Result<Event>;
2677}
2678
2679#[async_trait]
2681impl<E: EventEmitter + ?Sized> EventEmitter for Arc<E> {
2682 async fn emit(&self, request: EventRequest) -> Result<Event> {
2683 (**self).emit(request).await
2684 }
2685}
2686
2687#[derive(Debug, Clone, Default)]
2691pub struct NoopEventEmitter;
2692
2693#[async_trait]
2694impl EventEmitter for NoopEventEmitter {
2695 async fn emit(&self, request: EventRequest) -> Result<Event> {
2696 Ok(request.into_event(crate::typed_id::EventId::new(), 0))
2698 }
2699}
2700
2701#[derive(Debug, Clone)]
2714pub struct ResolvedImage {
2715 pub base64: String,
2717 pub media_type: String,
2719}
2720
2721impl ResolvedImage {
2722 pub fn new(base64: impl Into<String>, media_type: impl Into<String>) -> Self {
2724 Self {
2725 base64: base64.into(),
2726 media_type: media_type.into(),
2727 }
2728 }
2729
2730 pub fn to_data_url(&self) -> String {
2734 format!("data:{};base64,{}", self.media_type, self.base64)
2735 }
2736}
2737
2738#[async_trait]
2771pub trait ImageResolver: Send + Sync {
2772 async fn resolve_image(&self, image_id: Uuid) -> Result<Option<ResolvedImage>>;
2776}
2777
2778#[derive(Debug)]
2784pub enum SpawnClaimResult {
2785 Claimed {
2788 spawn_handle_id: uuid::Uuid,
2789 claim_token: uuid::Uuid,
2790 },
2791 ClaimedPendingChild {
2795 spawn_handle_id: uuid::Uuid,
2796 claim_token: uuid::Uuid,
2797 },
2798 AlreadyRunning {
2801 child_session_id: crate::typed_id::SessionId,
2802 claim_token: uuid::Uuid,
2804 },
2805 AlreadySettled {
2808 child_session_id: crate::typed_id::SessionId,
2809 terminal_status: String,
2811 terminal_result: String,
2812 },
2813}
2814
2815#[async_trait]
2823pub trait SubagentSpawnStore: Send + Sync + 'static {
2824 async fn try_claim_spawn(
2829 &self,
2830 parent_session_id: crate::typed_id::SessionId,
2831 tool_call_id: &str,
2832 claim_token: uuid::Uuid,
2833 ) -> Result<SpawnClaimResult>;
2834
2835 async fn register_child_session(
2840 &self,
2841 spawn_handle_id: uuid::Uuid,
2842 claim_token: uuid::Uuid,
2843 child_session_id: crate::typed_id::SessionId,
2844 ) -> Result<()>;
2845
2846 async fn settle_spawn(
2852 &self,
2853 parent_session_id: crate::typed_id::SessionId,
2854 tool_call_id: &str,
2855 claim_token: uuid::Uuid,
2856 terminal_status: &str,
2857 terminal_result: &str,
2858 ) -> Result<()>;
2859}
2860
2861#[async_trait]
2863impl<S: SubagentSpawnStore + ?Sized> SubagentSpawnStore for Arc<S> {
2864 async fn try_claim_spawn(
2865 &self,
2866 parent_session_id: crate::typed_id::SessionId,
2867 tool_call_id: &str,
2868 claim_token: uuid::Uuid,
2869 ) -> Result<SpawnClaimResult> {
2870 (**self)
2871 .try_claim_spawn(parent_session_id, tool_call_id, claim_token)
2872 .await
2873 }
2874
2875 async fn register_child_session(
2876 &self,
2877 spawn_handle_id: uuid::Uuid,
2878 claim_token: uuid::Uuid,
2879 child_session_id: crate::typed_id::SessionId,
2880 ) -> Result<()> {
2881 (**self)
2882 .register_child_session(spawn_handle_id, claim_token, child_session_id)
2883 .await
2884 }
2885
2886 async fn settle_spawn(
2887 &self,
2888 parent_session_id: crate::typed_id::SessionId,
2889 tool_call_id: &str,
2890 claim_token: uuid::Uuid,
2891 terminal_status: &str,
2892 terminal_result: &str,
2893 ) -> Result<()> {
2894 (**self)
2895 .settle_spawn(
2896 parent_session_id,
2897 tool_call_id,
2898 claim_token,
2899 terminal_status,
2900 terminal_result,
2901 )
2902 .await
2903 }
2904}
2905
2906pub struct NoopSubagentSpawnStore;
2910
2911#[async_trait]
2912impl SubagentSpawnStore for NoopSubagentSpawnStore {
2913 async fn try_claim_spawn(
2914 &self,
2915 _parent_session_id: crate::typed_id::SessionId,
2916 _tool_call_id: &str,
2917 claim_token: uuid::Uuid,
2918 ) -> Result<SpawnClaimResult> {
2919 Ok(SpawnClaimResult::Claimed {
2920 spawn_handle_id: uuid::Uuid::new_v4(),
2921 claim_token,
2922 })
2923 }
2924
2925 async fn register_child_session(
2926 &self,
2927 _spawn_handle_id: uuid::Uuid,
2928 _claim_token: uuid::Uuid,
2929 _child_session_id: crate::typed_id::SessionId,
2930 ) -> Result<()> {
2931 Ok(())
2932 }
2933
2934 async fn settle_spawn(
2935 &self,
2936 _parent_session_id: crate::typed_id::SessionId,
2937 _tool_call_id: &str,
2938 _claim_token: uuid::Uuid,
2939 _terminal_status: &str,
2940 _terminal_result: &str,
2941 ) -> Result<()> {
2942 Ok(())
2943 }
2944}
2945
2946#[cfg(test)]
2951mod tests {
2952 use super::*;
2953
2954 #[test]
2955 fn test_resolved_image_new() {
2956 let image = ResolvedImage::new("SGVsbG8=", "image/png");
2957 assert_eq!(image.base64, "SGVsbG8=");
2958 assert_eq!(image.media_type, "image/png");
2959 }
2960
2961 #[test]
2962 fn test_resolved_image_to_data_url() {
2963 let image = ResolvedImage::new("SGVsbG8=", "image/png");
2964 let data_url = image.to_data_url();
2965 assert_eq!(data_url, "data:image/png;base64,SGVsbG8=");
2966 }
2967
2968 #[test]
2969 fn test_resolved_image_jpeg() {
2970 let image = ResolvedImage::new("base64data", "image/jpeg");
2971 let data_url = image.to_data_url();
2972 assert!(data_url.starts_with("data:image/jpeg;base64,"));
2973 }
2974}