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)]
235pub struct ResolvedModel {
236 pub model: String,
238 pub provider_type: DriverId,
240 pub api_key: Option<String>,
242 pub base_url: Option<String>,
244 pub provider_metadata: Option<crate::driver_registry::ProviderMetadata>,
247}
248
249#[async_trait]
259pub trait ProviderStore: Send + Sync {
260 async fn get_resolved_model(&self, model_id: ModelId) -> Result<Option<ResolvedModel>>;
265
266 async fn get_default_model(&self) -> Result<Option<ResolvedModel>>;
270}
271
272#[async_trait]
273impl<T: ProviderStore + ?Sized> ProviderStore for std::sync::Arc<T> {
274 async fn get_resolved_model(&self, model_id: ModelId) -> Result<Option<ResolvedModel>> {
275 (**self).get_resolved_model(model_id).await
276 }
277
278 async fn get_default_model(&self) -> Result<Option<ResolvedModel>> {
279 (**self).get_default_model().await
280 }
281}
282
283#[derive(Debug, Clone)]
289pub struct StoredImageInfo {
290 pub id: ImageId,
291 pub filename: String,
292 pub content_type: String,
293 pub size_bytes: i64,
294 pub metadata: serde_json::Value,
295 pub created_at: DateTime<Utc>,
296}
297
298#[derive(Debug, Clone)]
300pub struct StoredImage {
301 pub info: StoredImageInfo,
302 pub data: Vec<u8>,
303}
304
305#[derive(Debug, Clone)]
307pub struct CreateStoredImage {
308 pub filename: String,
309 pub content_type: String,
310 pub data: Vec<u8>,
311 pub metadata: serde_json::Value,
312}
313
314#[async_trait]
315pub trait ImageArtifactStore: Send + Sync {
316 async fn create_image(&self, input: CreateStoredImage) -> Result<StoredImageInfo>;
318
319 async fn get_image(&self, image_id: ImageId) -> Result<Option<StoredImage>>;
321
322 async fn get_image_info(&self, image_id: ImageId) -> Result<Option<StoredImageInfo>>;
324}
325
326#[derive(Debug, Clone)]
332pub struct ProviderCredentials {
333 pub api_key: String,
334 pub base_url: Option<String>,
335}
336
337#[async_trait]
338pub trait ProviderCredentialStore: Send + Sync {
339 async fn get_default_provider_credentials(
344 &self,
345 provider_type: &str,
346 ) -> Result<Option<ProviderCredentials>>;
347}
348
349#[async_trait]
360pub trait ToolExecutor: Send + Sync {
361 async fn execute(&self, tool_call: &ToolCall, tool_def: &ToolDefinition) -> Result<ToolResult>;
366
367 async fn execute_with_context(
372 &self,
373 tool_call: &ToolCall,
374 tool_def: &ToolDefinition,
375 _context: &ToolContext,
376 ) -> Result<ToolResult> {
377 self.execute(tool_call, tool_def).await
379 }
380
381 async fn execute_batch(
383 &self,
384 tool_calls: &[ToolCall],
385 tool_defs: &[ToolDefinition],
386 ) -> Result<Vec<ToolResult>> {
387 let mut results = Vec::with_capacity(tool_calls.len());
388
389 let tool_map = build_tool_map(tool_defs);
390
391 for tool_call in tool_calls {
392 let tool_def = tool_map.get(tool_call.name.as_str()).ok_or_else(|| {
393 crate::error::AgentLoopError::tool(format!(
394 "Tool definition not found: {}",
395 tool_call.name
396 ))
397 })?;
398
399 results.push(self.execute(tool_call, tool_def).await?);
400 }
401
402 Ok(results)
403 }
404
405 async fn execute_parallel(
407 &self,
408 tool_calls: &[ToolCall],
409 tool_defs: &[ToolDefinition],
410 ) -> Result<Vec<ToolResult>>
411 where
412 Self: Sized,
413 {
414 use futures::future::join_all;
415
416 let tool_map = build_tool_map(tool_defs);
417
418 let futures: Vec<_> = tool_calls
419 .iter()
420 .map(|tool_call| async {
421 let tool_def = tool_map.get(tool_call.name.as_str()).ok_or_else(|| {
422 crate::error::AgentLoopError::tool(format!(
423 "Tool definition not found: {}",
424 tool_call.name
425 ))
426 })?;
427 self.execute(tool_call, tool_def).await
428 })
429 .collect();
430
431 let results = join_all(futures).await;
432 results.into_iter().collect()
433 }
434}
435
436#[async_trait]
440impl ToolExecutor for std::sync::Arc<dyn ToolExecutor> {
441 async fn execute(&self, tool_call: &ToolCall, tool_def: &ToolDefinition) -> Result<ToolResult> {
442 (**self).execute(tool_call, tool_def).await
443 }
444
445 async fn execute_with_context(
446 &self,
447 tool_call: &ToolCall,
448 tool_def: &ToolDefinition,
449 context: &ToolContext,
450 ) -> Result<ToolResult> {
451 (**self)
452 .execute_with_context(tool_call, tool_def, context)
453 .await
454 }
455
456 async fn execute_batch(
457 &self,
458 tool_calls: &[ToolCall],
459 tool_defs: &[ToolDefinition],
460 ) -> Result<Vec<ToolResult>> {
461 (**self).execute_batch(tool_calls, tool_defs).await
462 }
463}
464
465#[async_trait]
477pub trait SessionFileSystem: Send + Sync {
478 fn display_root(&self) -> String {
484 crate::session_path::WORKSPACE_PREFIX.to_string()
485 }
486
487 fn display_path(&self, path: &str) -> String {
493 crate::session_path::to_display_path(path)
494 }
495
496 fn resolve_path(&self, input: &str) -> String {
506 crate::session_path::to_session_path(input)
507 }
508
509 fn is_mount_resolver(&self) -> bool;
513
514 async fn read_file(&self, session_id: SessionId, path: &str) -> Result<Option<SessionFile>>;
516
517 async fn write_file(
519 &self,
520 session_id: SessionId,
521 path: &str,
522 content: &str,
523 encoding: &str,
524 ) -> Result<SessionFile>;
525
526 async fn write_file_if_content_matches(
531 &self,
532 session_id: SessionId,
533 path: &str,
534 expected_content: &str,
535 expected_encoding: &str,
536 content: &str,
537 encoding: &str,
538 ) -> Result<Option<SessionFile>> {
539 let Some(existing) = self.read_file(session_id, path).await? else {
540 return Ok(None);
541 };
542
543 if existing.is_directory {
544 return Ok(None);
545 }
546
547 let current_content = existing.content.unwrap_or_default();
548 if current_content != expected_content || existing.encoding != expected_encoding {
549 return Ok(None);
550 }
551
552 self.write_file(session_id, path, content, encoding)
553 .await
554 .map(Some)
555 }
556
557 async fn delete_file(&self, session_id: SessionId, path: &str, recursive: bool)
559 -> Result<bool>;
560
561 async fn list_directory(&self, session_id: SessionId, path: &str) -> Result<Vec<FileInfo>>;
563
564 async fn stat_file(&self, session_id: SessionId, path: &str) -> Result<Option<FileStat>>;
566
567 async fn grep_files(
574 &self,
575 session_id: SessionId,
576 pattern: &str,
577 path_pattern: Option<&str>,
578 ) -> Result<Vec<GrepMatch>>;
579
580 async fn grep_files_with_options(
586 &self,
587 session_id: SessionId,
588 pattern: &str,
589 options: &GrepOptions,
590 ) -> Result<GrepSearchResult> {
591 if options.before_context != 0 || options.after_context != 0 {
592 return Err(crate::error::AgentLoopError::tool(
593 "this file store does not support grep context",
594 ));
595 }
596 let all = self
597 .grep_files(session_id, pattern, options.path_pattern.as_deref())
598 .await?;
599 Ok(crate::session_file::bound_grep_matches(all, options))
600 }
601
602 async fn create_directory(&self, session_id: SessionId, path: &str) -> Result<FileInfo>;
604
605 async fn seed_initial_file(&self, session_id: SessionId, file: &InitialFile) -> Result<()> {
607 if file.is_readonly {
608 return Err(crate::error::AgentLoopError::store(
609 "read-only initial files require a SessionFileSystem-specific seed implementation",
610 ));
611 }
612 self.write_file(session_id, &file.path, &file.content, &file.encoding)
613 .await?;
614 Ok(())
615 }
616}
617
618pub struct WorkspaceScopedFileSystem {
628 inner: Arc<dyn SessionFileSystem>,
629 key: SessionId,
630}
631
632impl WorkspaceScopedFileSystem {
633 pub fn wrap(
635 inner: Arc<dyn SessionFileSystem>,
636 workspace_id: WorkspaceId,
637 ) -> Arc<dyn SessionFileSystem> {
638 Arc::new(Self {
639 inner,
640 key: SessionId::from_uuid(workspace_id.uuid()),
641 })
642 }
643}
644
645#[async_trait]
646impl SessionFileSystem for WorkspaceScopedFileSystem {
647 async fn read_file(&self, _session_id: SessionId, path: &str) -> Result<Option<SessionFile>> {
648 self.inner.read_file(self.key, path).await
649 }
650 async fn write_file(
651 &self,
652 _session_id: SessionId,
653 path: &str,
654 content: &str,
655 encoding: &str,
656 ) -> Result<SessionFile> {
657 self.inner
658 .write_file(self.key, path, content, encoding)
659 .await
660 }
661 async fn write_file_if_content_matches(
662 &self,
663 _session_id: SessionId,
664 path: &str,
665 expected_content: &str,
666 expected_encoding: &str,
667 content: &str,
668 encoding: &str,
669 ) -> Result<Option<SessionFile>> {
670 self.inner
671 .write_file_if_content_matches(
672 self.key,
673 path,
674 expected_content,
675 expected_encoding,
676 content,
677 encoding,
678 )
679 .await
680 }
681 async fn delete_file(
682 &self,
683 _session_id: SessionId,
684 path: &str,
685 recursive: bool,
686 ) -> Result<bool> {
687 self.inner.delete_file(self.key, path, recursive).await
688 }
689 async fn list_directory(&self, _session_id: SessionId, path: &str) -> Result<Vec<FileInfo>> {
690 self.inner.list_directory(self.key, path).await
691 }
692 async fn stat_file(&self, _session_id: SessionId, path: &str) -> Result<Option<FileStat>> {
693 self.inner.stat_file(self.key, path).await
694 }
695 async fn grep_files(
696 &self,
697 _session_id: SessionId,
698 pattern: &str,
699 path_pattern: Option<&str>,
700 ) -> Result<Vec<GrepMatch>> {
701 self.inner.grep_files(self.key, pattern, path_pattern).await
702 }
703 async fn grep_files_with_options(
704 &self,
705 _session_id: SessionId,
706 pattern: &str,
707 options: &GrepOptions,
708 ) -> Result<GrepSearchResult> {
709 self.inner
710 .grep_files_with_options(self.key, pattern, options)
711 .await
712 }
713 async fn create_directory(&self, _session_id: SessionId, path: &str) -> Result<FileInfo> {
714 self.inner.create_directory(self.key, path).await
715 }
716 async fn seed_initial_file(&self, _session_id: SessionId, file: &InitialFile) -> Result<()> {
717 self.inner.seed_initial_file(self.key, file).await
718 }
719
720 fn display_root(&self) -> String {
721 self.inner.display_root()
722 }
723
724 fn display_path(&self, path: &str) -> String {
725 self.inner.display_path(path)
726 }
727
728 fn resolve_path(&self, input: &str) -> String {
729 self.inner.resolve_path(input)
730 }
731
732 fn is_mount_resolver(&self) -> bool {
733 self.inner.is_mount_resolver()
734 }
735}
736
737#[async_trait]
738impl<T: SessionFileSystem + ?Sized> SessionFileSystem for std::sync::Arc<T> {
739 fn display_root(&self) -> String {
740 (**self).display_root()
741 }
742
743 fn display_path(&self, path: &str) -> String {
744 (**self).display_path(path)
745 }
746
747 fn resolve_path(&self, input: &str) -> String {
748 (**self).resolve_path(input)
749 }
750
751 fn is_mount_resolver(&self) -> bool {
752 (**self).is_mount_resolver()
753 }
754
755 async fn read_file(&self, session_id: SessionId, path: &str) -> Result<Option<SessionFile>> {
756 (**self).read_file(session_id, path).await
757 }
758
759 async fn write_file(
760 &self,
761 session_id: SessionId,
762 path: &str,
763 content: &str,
764 encoding: &str,
765 ) -> Result<SessionFile> {
766 (**self)
767 .write_file(session_id, path, content, encoding)
768 .await
769 }
770
771 async fn write_file_if_content_matches(
772 &self,
773 session_id: SessionId,
774 path: &str,
775 expected_content: &str,
776 expected_encoding: &str,
777 content: &str,
778 encoding: &str,
779 ) -> Result<Option<SessionFile>> {
780 (**self)
781 .write_file_if_content_matches(
782 session_id,
783 path,
784 expected_content,
785 expected_encoding,
786 content,
787 encoding,
788 )
789 .await
790 }
791
792 async fn delete_file(
793 &self,
794 session_id: SessionId,
795 path: &str,
796 recursive: bool,
797 ) -> Result<bool> {
798 (**self).delete_file(session_id, path, recursive).await
799 }
800
801 async fn list_directory(&self, session_id: SessionId, path: &str) -> Result<Vec<FileInfo>> {
802 (**self).list_directory(session_id, path).await
803 }
804
805 async fn stat_file(&self, session_id: SessionId, path: &str) -> Result<Option<FileStat>> {
806 (**self).stat_file(session_id, path).await
807 }
808
809 async fn grep_files(
810 &self,
811 session_id: SessionId,
812 pattern: &str,
813 path_pattern: Option<&str>,
814 ) -> Result<Vec<GrepMatch>> {
815 (**self).grep_files(session_id, pattern, path_pattern).await
816 }
817
818 async fn grep_files_with_options(
819 &self,
820 session_id: SessionId,
821 pattern: &str,
822 options: &GrepOptions,
823 ) -> Result<GrepSearchResult> {
824 (**self)
825 .grep_files_with_options(session_id, pattern, options)
826 .await
827 }
828
829 async fn create_directory(&self, session_id: SessionId, path: &str) -> Result<FileInfo> {
830 (**self).create_directory(session_id, path).await
831 }
832
833 async fn seed_initial_file(&self, session_id: SessionId, file: &InitialFile) -> Result<()> {
834 (**self).seed_initial_file(session_id, file).await
835 }
836}
837
838pub use SessionFileSystem as SessionFileStore;
840
841#[derive(Clone, Default)]
847pub struct SessionFileSystemFactoryContext {
848 values: Arc<HashMap<TypeId, Arc<dyn Any + Send + Sync>>>,
849}
850
851impl SessionFileSystemFactoryContext {
852 pub fn new() -> Self {
853 Self::default()
854 }
855
856 pub fn with<T: Any + Send + Sync>(mut self, value: Arc<T>) -> Self {
857 let values = Arc::make_mut(&mut self.values);
858 values.insert(TypeId::of::<T>(), value);
859 self
860 }
861
862 pub fn get<T: Any + Send + Sync>(&self) -> Option<Arc<T>> {
863 self.values
864 .get(&TypeId::of::<T>())
865 .and_then(|value| value.clone().downcast::<T>().ok())
866 }
867
868 pub fn with_workspace_roots(self, roots: Arc<crate::WorkspaceRootSet>) -> Self {
869 self.with(roots)
870 }
871
872 pub fn workspace_roots(&self) -> Option<Arc<crate::WorkspaceRootSet>> {
873 self.get::<crate::WorkspaceRootSet>()
874 }
875}
876
877#[async_trait]
879pub trait SessionFileSystemFactory: Send + Sync {
880 fn name(&self) -> &'static str {
882 "SessionFileSystemFactory"
883 }
884
885 fn is_disabled(&self) -> bool {
888 false
889 }
890
891 async fn create_session_file_system(
893 &self,
894 context: SessionFileSystemFactoryContext,
895 ) -> Result<Arc<dyn SessionFileSystem>>;
896}
897
898#[derive(Debug, Clone, Default)]
900pub struct DisabledSessionFileSystemFactory;
901
902#[async_trait]
903impl SessionFileSystemFactory for DisabledSessionFileSystemFactory {
904 fn name(&self) -> &'static str {
905 "DisabledSessionFileSystemFactory"
906 }
907
908 fn is_disabled(&self) -> bool {
909 true
910 }
911
912 async fn create_session_file_system(
913 &self,
914 _context: SessionFileSystemFactoryContext,
915 ) -> Result<Arc<dyn SessionFileSystem>> {
916 Err(crate::error::AgentLoopError::config(
917 "session filesystem is disabled",
918 ))
919 }
920}
921
922#[derive(Debug, Clone)]
928pub struct KeyInfo {
929 pub key: String,
930 pub created_at: chrono::DateTime<chrono::Utc>,
931 pub updated_at: chrono::DateTime<chrono::Utc>,
932}
933
934#[derive(Debug, Clone)]
936pub struct SecretInfo {
937 pub name: String,
938 pub created_at: chrono::DateTime<chrono::Utc>,
939 pub updated_at: chrono::DateTime<chrono::Utc>,
940}
941
942#[derive(Debug, Clone, serde::Serialize)]
952pub struct KnowledgeSearchHit {
953 pub id: String,
955 pub kb_id: String,
957 pub title: String,
958 pub kind: String,
959 pub tags: Vec<String>,
960 pub snippet: String,
962 pub resource: Option<String>,
964}
965
966#[async_trait]
970pub trait KnowledgeStore: Send + Sync {
971 async fn search_knowledge(
972 &self,
973 org_id: crate::typed_id::OrgId,
974 kb_public_ids: &[String],
975 query: &str,
976 kind: Option<&str>,
977 tags: &[String],
978 limit: usize,
979 ) -> Result<Vec<KnowledgeSearchHit>>;
980}
981
982#[async_trait]
987pub trait SessionStorageStore: Send + Sync {
988 async fn set_value(&self, session_id: SessionId, key: &str, value: &str) -> Result<()>;
992
993 async fn get_value(&self, session_id: SessionId, key: &str) -> Result<Option<String>>;
995
996 async fn delete_value(&self, session_id: SessionId, key: &str) -> Result<bool>;
998
999 async fn list_keys(&self, session_id: SessionId) -> Result<Vec<KeyInfo>>;
1001
1002 async fn set_secret(&self, session_id: SessionId, name: &str, value: &str) -> Result<()>;
1006
1007 async fn get_secret(&self, session_id: SessionId, name: &str) -> Result<Option<String>>;
1009
1010 async fn delete_secret(&self, session_id: SessionId, name: &str) -> Result<bool>;
1012
1013 async fn list_secrets(&self, session_id: SessionId) -> Result<Vec<SecretInfo>>;
1015}
1016
1017use crate::session_schedule::SessionSchedule;
1022use crate::typed_id::ScheduleId;
1023
1024#[async_trait]
1028pub trait SessionScheduleStore: Send + Sync {
1029 async fn create_schedule(
1031 &self,
1032 session_id: SessionId,
1033 description: String,
1034 cron_expression: Option<String>,
1035 scheduled_at: Option<chrono::DateTime<chrono::Utc>>,
1036 timezone: String,
1037 ) -> Result<SessionSchedule>;
1038
1039 async fn create_schedule_enforcing_limits(
1043 &self,
1044 session_id: SessionId,
1045 description: String,
1046 cron_expression: Option<String>,
1047 scheduled_at: Option<chrono::DateTime<chrono::Utc>>,
1048 timezone: String,
1049 ) -> std::result::Result<SessionSchedule, crate::session_schedule::ScheduleLimitError> {
1050 let per_session = self
1051 .count_active_schedules(session_id)
1052 .await
1053 .map_err(crate::session_schedule::ScheduleLimitError::Store)?;
1054 if per_session >= crate::session_schedule::MAX_ACTIVE_SCHEDULES_PER_SESSION {
1055 return Err(crate::session_schedule::ScheduleLimitError::Rejected(
1056 format!(
1057 "Maximum {} active schedules per session. Cancel an existing schedule first.",
1058 crate::session_schedule::MAX_ACTIVE_SCHEDULES_PER_SESSION
1059 ),
1060 ));
1061 }
1062
1063 let max_per_org = crate::session_schedule::max_active_schedules_per_org();
1064 let per_org = self
1065 .count_active_org_schedules()
1066 .await
1067 .map_err(crate::session_schedule::ScheduleLimitError::Store)?;
1068 if i64::from(per_org) >= max_per_org {
1069 return Err(crate::session_schedule::ScheduleLimitError::Rejected(
1070 format!(
1071 "Maximum {max_per_org} active schedules per org reached. Cancel an existing schedule first."
1072 ),
1073 ));
1074 }
1075
1076 if let Some(cron) = cron_expression.as_deref() {
1077 crate::session_schedule::validate_cron_min_interval(cron)
1078 .map_err(crate::session_schedule::ScheduleLimitError::Rejected)?;
1079 }
1080
1081 self.create_schedule(
1082 session_id,
1083 description,
1084 cron_expression,
1085 scheduled_at,
1086 timezone,
1087 )
1088 .await
1089 .map_err(crate::session_schedule::ScheduleLimitError::Store)
1090 }
1091
1092 async fn cancel_schedule(
1094 &self,
1095 session_id: SessionId,
1096 schedule_id: ScheduleId,
1097 ) -> Result<SessionSchedule>;
1098
1099 async fn list_schedules(&self, session_id: SessionId) -> Result<Vec<SessionSchedule>>;
1101
1102 async fn count_active_schedules(&self, session_id: SessionId) -> Result<u32>;
1104
1105 async fn count_active_org_schedules(&self) -> Result<u32>;
1110}
1111
1112#[async_trait]
1122pub trait SessionResourceRegistry: Send + Sync {
1123 async fn register(
1125 &self,
1126 entry: crate::session_resource::RegisterSessionResource,
1127 ) -> Result<crate::session_resource::SessionResourceEntry>;
1128
1129 async fn update_status(
1131 &self,
1132 session_id: SessionId,
1133 resource_id: &str,
1134 status: crate::session_resource::SessionResourceStatus,
1135 ) -> Result<Option<crate::session_resource::SessionResourceEntry>>;
1136
1137 async fn get(
1139 &self,
1140 session_id: SessionId,
1141 resource_id: &str,
1142 ) -> Result<Option<crate::session_resource::SessionResourceEntry>>;
1143
1144 async fn list(
1146 &self,
1147 session_id: SessionId,
1148 filter: Option<&crate::session_resource::SessionResourceFilter>,
1149 ) -> Result<Vec<crate::session_resource::SessionResourceEntry>>;
1150
1151 async fn deregister(&self, session_id: SessionId, resource_id: &str) -> Result<bool>;
1153}
1154
1155#[async_trait]
1165pub trait LeasedResourceStore: Send + Sync {
1166 async fn upsert_resource(&self, input: UpsertLeasedResource) -> Result<LeasedResource>;
1172
1173 async fn release_resource(
1179 &self,
1180 session_id: SessionId,
1181 provider: &str,
1182 resource_type: &str,
1183 external_id: &str,
1184 ) -> Result<Option<LeasedResource>>;
1185
1186 async fn list_resources(&self, session_id: SessionId) -> Result<Vec<LeasedResource>>;
1191}
1192
1193pub const DEFAULT_MAX_SUBAGENT_DEPTH: u32 = 2;
1200pub const DEFAULT_MAX_ACTIVE_DESCENDANT_SUBAGENT_TASKS: u32 = 16;
1201pub const DEFAULT_MAX_TOTAL_DESCENDANT_SUBAGENT_TASKS: u32 = 200;
1202pub const DEFAULT_MAX_ACTIVE_DETACHED_TASKS: u32 = 8;
1208pub const DEFAULT_MAX_TOTAL_DETACHED_TASKS: u32 = 50;
1209
1210#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1212pub struct SubagentNestingPolicy {
1213 pub platform_default: u32,
1214 pub org_override: Option<u32>,
1215 pub agent_override: Option<u32>,
1216 pub platform_default_max_active_descendant_tasks: u32,
1217 pub org_override_max_active_descendant_tasks: Option<u32>,
1218 pub agent_override_max_active_descendant_tasks: Option<u32>,
1219 pub platform_default_max_total_descendant_tasks: u32,
1220 pub org_override_max_total_descendant_tasks: Option<u32>,
1221 pub agent_override_max_total_descendant_tasks: Option<u32>,
1222 pub platform_default_max_active_detached_tasks: u32,
1223 pub org_override_max_active_detached_tasks: Option<u32>,
1224 pub agent_override_max_active_detached_tasks: Option<u32>,
1225 pub platform_default_max_total_detached_tasks: u32,
1226 pub org_override_max_total_detached_tasks: Option<u32>,
1227 pub agent_override_max_total_detached_tasks: Option<u32>,
1228}
1229
1230impl Default for SubagentNestingPolicy {
1231 fn default() -> Self {
1232 Self {
1233 platform_default: DEFAULT_MAX_SUBAGENT_DEPTH,
1234 org_override: None,
1235 agent_override: None,
1236 platform_default_max_active_descendant_tasks:
1237 DEFAULT_MAX_ACTIVE_DESCENDANT_SUBAGENT_TASKS,
1238 org_override_max_active_descendant_tasks: None,
1239 agent_override_max_active_descendant_tasks: None,
1240 platform_default_max_total_descendant_tasks:
1241 DEFAULT_MAX_TOTAL_DESCENDANT_SUBAGENT_TASKS,
1242 org_override_max_total_descendant_tasks: None,
1243 agent_override_max_total_descendant_tasks: None,
1244 platform_default_max_active_detached_tasks: DEFAULT_MAX_ACTIVE_DETACHED_TASKS,
1245 org_override_max_active_detached_tasks: None,
1246 agent_override_max_active_detached_tasks: None,
1247 platform_default_max_total_detached_tasks: DEFAULT_MAX_TOTAL_DETACHED_TASKS,
1248 org_override_max_total_detached_tasks: None,
1249 agent_override_max_total_detached_tasks: None,
1250 }
1251 }
1252}
1253
1254impl SubagentNestingPolicy {
1255 pub fn max_subagent_depth(self) -> u32 {
1256 self.agent_override
1257 .or(self.org_override)
1258 .unwrap_or(self.platform_default)
1259 }
1260
1261 pub fn max_active_descendant_tasks(self) -> u32 {
1262 self.agent_override_max_active_descendant_tasks
1263 .or(self.org_override_max_active_descendant_tasks)
1264 .unwrap_or(self.platform_default_max_active_descendant_tasks)
1265 }
1266
1267 pub fn max_total_descendant_tasks(self) -> u32 {
1268 self.agent_override_max_total_descendant_tasks
1269 .or(self.org_override_max_total_descendant_tasks)
1270 .unwrap_or(self.platform_default_max_total_descendant_tasks)
1271 }
1272
1273 pub fn max_active_detached_tasks(self) -> u32 {
1274 self.agent_override_max_active_detached_tasks
1275 .or(self.org_override_max_active_detached_tasks)
1276 .unwrap_or(self.platform_default_max_active_detached_tasks)
1277 }
1278
1279 pub fn max_total_detached_tasks(self) -> u32 {
1280 self.agent_override_max_total_detached_tasks
1281 .or(self.org_override_max_total_detached_tasks)
1282 .unwrap_or(self.platform_default_max_total_detached_tasks)
1283 }
1284
1285 pub fn with_platform_default(mut self, depth: u32) -> Self {
1286 self.platform_default = depth;
1287 self
1288 }
1289
1290 pub fn with_org_override(mut self, depth: Option<u32>) -> Self {
1291 self.org_override = depth;
1292 self
1293 }
1294
1295 pub fn with_agent_override(mut self, depth: Option<u32>) -> Self {
1296 self.agent_override = depth;
1297 self
1298 }
1299
1300 pub fn with_agent_task_caps_override(
1301 mut self,
1302 max_active: Option<u32>,
1303 max_total: Option<u32>,
1304 ) -> Self {
1305 self.agent_override_max_active_descendant_tasks = max_active;
1306 self.agent_override_max_total_descendant_tasks = max_total;
1307 self
1308 }
1309
1310 pub fn with_agent_detached_task_caps_override(
1311 mut self,
1312 max_active: Option<u32>,
1313 max_total: Option<u32>,
1314 ) -> Self {
1315 self.agent_override_max_active_detached_tasks = max_active;
1316 self.agent_override_max_total_detached_tasks = max_total;
1317 self
1318 }
1319}
1320
1321pub type SessionSqlDbStoreRef = Arc<dyn crate::session_sqldb::SessionSqlDbStore>;
1323
1324#[async_trait]
1329pub trait UserConnectionResolver: Send + Sync {
1330 async fn get_connection_token(
1333 &self,
1334 session_id: SessionId,
1335 provider: &str,
1336 ) -> Result<Option<String>>;
1337
1338 async fn get_connection_user(
1343 &self,
1344 _session_id: SessionId,
1345 _provider: &str,
1346 ) -> Result<Option<Uuid>> {
1347 Ok(None)
1348 }
1349
1350 async fn get_connection_token_for_user(
1355 &self,
1356 _user_id: Uuid,
1357 _provider: &str,
1358 ) -> Result<Option<String>> {
1359 Ok(None)
1360 }
1361
1362 async fn get_connection_metadata(
1365 &self,
1366 _session_id: SessionId,
1367 _provider: &str,
1368 ) -> Result<Option<serde_json::Value>> {
1369 Ok(None)
1370 }
1371}
1372
1373#[async_trait]
1383pub trait BudgetChecker: Send + Sync {
1384 async fn check_budgets(&self, session_id: &str) -> Result<crate::budget::BudgetToolResponse>;
1386}
1387
1388#[async_trait]
1397pub trait PaymentAuthority: Send + Sync {
1398 async fn execute_machine_payment(
1399 &self,
1400 session_id: SessionId,
1401 request: crate::payment::MachinePaymentRequest,
1402 ) -> Result<crate::payment::MachinePaymentResponse>;
1403}
1404
1405#[async_trait]
1415pub trait SessionCreationAuthority: Send + Sync {
1416 async fn authorize_session_creation(&self, session_id: SessionId) -> Result<SessionId>;
1420}
1421
1422#[async_trait]
1432pub trait OutboundToolRateLimiter: Send + Sync {
1433 async fn check_org(&self, org_id: &crate::typed_id::OrgId) -> bool;
1435}
1436
1437#[derive(Debug)]
1443pub enum ToolCallClaimResult {
1444 Claimed { claim_token: uuid::Uuid },
1447 AlreadySettled {
1449 result_json: serde_json::Value,
1450 args_fingerprint: String,
1451 },
1452 AlreadyRunning { args_fingerprint: String },
1457 DeterminismViolation {
1461 stored_fingerprint: String,
1462 current_fingerprint: String,
1463 },
1464}
1465
1466#[derive(Debug, Clone)]
1468pub enum DurableToolCallStatus {
1469 Settled { result_json: serde_json::Value },
1471 Interrupted {
1473 result_json: Option<serde_json::Value>,
1474 },
1475 Running,
1477}
1478
1479#[async_trait]
1484pub trait DurableToolResultStore: Send + Sync + 'static {
1485 async fn try_claim_tool_call(
1493 &self,
1494 turn_id: &str,
1495 tool_call_id: &str,
1496 tool_name: &str,
1497 args_fingerprint: &str,
1498 ) -> Result<ToolCallClaimResult>;
1499
1500 async fn settle_tool_call(
1506 &self,
1507 turn_id: &str,
1508 tool_call_id: &str,
1509 result_json: serde_json::Value,
1510 status: &str,
1511 claim_token: uuid::Uuid,
1512 ) -> Result<bool>;
1513
1514 async fn get_tool_call_status(
1519 &self,
1520 turn_id: &str,
1521 tool_call_id: &str,
1522 ) -> Result<Option<DurableToolCallStatus>>;
1523}
1524
1525pub struct NoopDurableToolResultStore;
1528
1529#[async_trait]
1530impl DurableToolResultStore for NoopDurableToolResultStore {
1531 async fn try_claim_tool_call(
1532 &self,
1533 _turn_id: &str,
1534 _tool_call_id: &str,
1535 _tool_name: &str,
1536 _args_fingerprint: &str,
1537 ) -> Result<ToolCallClaimResult> {
1538 Ok(ToolCallClaimResult::Claimed {
1539 claim_token: uuid::Uuid::new_v4(),
1540 })
1541 }
1542
1543 async fn settle_tool_call(
1544 &self,
1545 _turn_id: &str,
1546 _tool_call_id: &str,
1547 _result_json: serde_json::Value,
1548 _status: &str,
1549 _claim_token: uuid::Uuid,
1550 ) -> Result<bool> {
1551 Ok(true)
1552 }
1553
1554 async fn get_tool_call_status(
1555 &self,
1556 _turn_id: &str,
1557 _tool_call_id: &str,
1558 ) -> Result<Option<DurableToolCallStatus>> {
1559 Ok(None)
1560 }
1561}
1562
1563#[derive(Debug, Clone)]
1569pub struct StreamProgress {
1570 pub accumulated_len: usize,
1572 pub last_delta_at: u64,
1574}
1575
1576#[async_trait]
1582pub trait StreamHeartbeater: Send + Sync {
1583 async fn heartbeat(&self, progress: StreamProgress);
1589}
1590
1591pub struct NoopStreamHeartbeater;
1593
1594#[async_trait]
1595impl StreamHeartbeater for NoopStreamHeartbeater {
1596 async fn heartbeat(&self, _progress: StreamProgress) {}
1597}
1598
1599#[derive(Debug, Clone)]
1605pub struct PartialStreamState {
1606 pub message_id: MessageId,
1608
1609 pub accumulated: String,
1612}
1613
1614#[async_trait]
1622pub trait PartialStreamStore: Send + Sync {
1623 async fn get_partial_stream(
1626 &self,
1627 session_id: SessionId,
1628 turn_id: &str,
1629 ) -> Result<Option<PartialStreamState>>;
1630}
1631
1632pub struct NoopPartialStreamStore;
1634
1635#[async_trait]
1636impl PartialStreamStore for NoopPartialStreamStore {
1637 async fn get_partial_stream(
1638 &self,
1639 _session_id: SessionId,
1640 _turn_id: &str,
1641 ) -> Result<Option<PartialStreamState>> {
1642 Ok(None)
1643 }
1644}
1645
1646#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
1652pub enum ToolContextService {
1653 SessionFileSystem,
1654 SessionStorageStore,
1655 ImageArtifactStore,
1656 ProviderCredentialStore,
1657 UtilityLlmService,
1658 McpInvoker,
1659 EgressService,
1660 SessionSqlDbStore,
1661 MessageRetriever,
1662 SessionStore,
1663 SessionMutator,
1664 AgentStore,
1665 ConnectionResolver,
1666 ScheduleStore,
1667 PlatformStore,
1668 KnowledgeStore,
1669 KnowledgeIndexSearch,
1670 LeasedResourceStore,
1671 SessionResourceRegistry,
1672 SessionTaskRegistry,
1673 EventEmitter,
1674 CapabilityRegistry,
1675 ToolRegistry,
1676 OrgId,
1677 BudgetChecker,
1678 PaymentAuthority,
1679 SessionCreationAuthority,
1680 SubagentSpawnStore,
1681 ReasoningEffortHandle,
1682}
1683
1684impl ToolContextService {
1685 pub const fn name(self) -> &'static str {
1686 match self {
1687 Self::SessionFileSystem => "SessionFileSystem",
1688 Self::SessionStorageStore => "SessionStorageStore",
1689 Self::ImageArtifactStore => "ImageArtifactStore",
1690 Self::ProviderCredentialStore => "ProviderCredentialStore",
1691 Self::UtilityLlmService => "UtilityLlmService",
1692 Self::McpInvoker => "McpInvoker",
1693 Self::EgressService => "EgressService",
1694 Self::SessionSqlDbStore => "SessionSqlDbStore",
1695 Self::MessageRetriever => "MessageRetriever",
1696 Self::SessionStore => "SessionStore",
1697 Self::SessionMutator => "SessionMutator",
1698 Self::AgentStore => "AgentStore",
1699 Self::ConnectionResolver => "ConnectionResolver",
1700 Self::ScheduleStore => "SessionScheduleStore",
1701 Self::PlatformStore => "PlatformStore",
1702 Self::KnowledgeStore => "KnowledgeStore",
1703 Self::KnowledgeIndexSearch => "KnowledgeIndexSearch",
1704 Self::LeasedResourceStore => "LeasedResourceStore",
1705 Self::SessionResourceRegistry => "SessionResourceRegistry",
1706 Self::SessionTaskRegistry => "SessionTaskRegistry",
1707 Self::EventEmitter => "EventEmitter",
1708 Self::CapabilityRegistry => "CapabilityRegistry",
1709 Self::ToolRegistry => "ToolRegistry",
1710 Self::OrgId => "OrgId",
1711 Self::BudgetChecker => "BudgetChecker",
1712 Self::PaymentAuthority => "PaymentAuthority",
1713 Self::SessionCreationAuthority => "SessionCreationAuthority",
1714 Self::SubagentSpawnStore => "SubagentSpawnStore",
1715 Self::ReasoningEffortHandle => "ReasoningEffortHandle",
1716 }
1717 }
1718}
1719
1720#[derive(Clone, Default)]
1725pub struct ToolContextServices {
1726 pub file_store: Option<Arc<dyn SessionFileSystem>>,
1727 pub storage_store: Option<Arc<dyn SessionStorageStore>>,
1728 pub image_store: Option<Arc<dyn ImageArtifactStore>>,
1729 pub provider_credential_store: Option<Arc<dyn ProviderCredentialStore>>,
1730 pub utility_llm_service: Option<Arc<dyn crate::UtilityLlmService>>,
1731 pub mcp_invoker: Option<Arc<dyn crate::McpToolInvoker>>,
1732 pub egress_service: Option<Arc<dyn crate::EgressService>>,
1733 pub sqldb_store: Option<SessionSqlDbStoreRef>,
1734 pub message_retriever: Option<Arc<dyn crate::message_retriever::MessageRetriever>>,
1735 pub session_store: Option<Arc<dyn SessionStore>>,
1736 pub session_mutator: Option<Arc<dyn SessionMutator>>,
1737 pub agent_store: Option<Arc<dyn AgentStore>>,
1738 pub connection_resolver: Option<Arc<dyn UserConnectionResolver>>,
1739 pub schedule_store: Option<Arc<dyn SessionScheduleStore>>,
1740 pub platform_store: Option<Arc<dyn crate::platform_store::PlatformStore>>,
1741 pub knowledge_store: Option<Arc<dyn KnowledgeStore>>,
1742 pub knowledge_index_search: Option<Arc<dyn crate::vector_store::KnowledgeIndexSearch>>,
1743 pub leased_resource_store: Option<Arc<dyn LeasedResourceStore>>,
1744 pub session_resource_registry: Option<Arc<dyn SessionResourceRegistry>>,
1745 pub session_task_registry: Option<Arc<dyn crate::session_task::SessionTaskRegistry>>,
1746 pub event_emitter: Option<Arc<dyn EventEmitter>>,
1747 pub capability_registry: Option<crate::capabilities::CapabilityRegistry>,
1748 pub tool_registry: Option<Arc<crate::tools::ToolRegistry>>,
1749 pub org_id: Option<crate::typed_id::OrgId>,
1750 pub network_access: Option<crate::network_access::NetworkAccessList>,
1751 pub budget_checker: Option<Arc<dyn BudgetChecker>>,
1752 pub payment_authority: Option<Arc<dyn PaymentAuthority>>,
1753 pub session_creation_authority: Option<Arc<dyn SessionCreationAuthority>>,
1754 pub subagent_spawn_store: Option<Arc<dyn SubagentSpawnStore>>,
1755 pub subagent_nesting_policy: SubagentNestingPolicy,
1756 pub reasoning_effort_handle: Option<ReasoningEffortHandle>,
1757}
1758
1759impl ToolContextServices {
1760 pub fn provides(&self, service: ToolContextService) -> bool {
1761 match service {
1762 ToolContextService::SessionFileSystem => self.file_store.is_some(),
1763 ToolContextService::SessionStorageStore => self.storage_store.is_some(),
1764 ToolContextService::ImageArtifactStore => self.image_store.is_some(),
1765 ToolContextService::ProviderCredentialStore => self.provider_credential_store.is_some(),
1766 ToolContextService::UtilityLlmService => self.utility_llm_service.is_some(),
1767 ToolContextService::McpInvoker => self.mcp_invoker.is_some(),
1768 ToolContextService::EgressService => self.egress_service.is_some(),
1769 ToolContextService::SessionSqlDbStore => self.sqldb_store.is_some(),
1770 ToolContextService::MessageRetriever => self.message_retriever.is_some(),
1771 ToolContextService::SessionStore => self.session_store.is_some(),
1772 ToolContextService::SessionMutator => self.session_mutator.is_some(),
1773 ToolContextService::AgentStore => self.agent_store.is_some(),
1774 ToolContextService::ConnectionResolver => self.connection_resolver.is_some(),
1775 ToolContextService::ScheduleStore => self.schedule_store.is_some(),
1776 ToolContextService::PlatformStore => self.platform_store.is_some(),
1777 ToolContextService::KnowledgeStore => self.knowledge_store.is_some(),
1778 ToolContextService::KnowledgeIndexSearch => self.knowledge_index_search.is_some(),
1779 ToolContextService::LeasedResourceStore => self.leased_resource_store.is_some(),
1780 ToolContextService::SessionResourceRegistry => self.session_resource_registry.is_some(),
1781 ToolContextService::SessionTaskRegistry => self.session_task_registry.is_some(),
1782 ToolContextService::EventEmitter => self.event_emitter.is_some(),
1783 ToolContextService::CapabilityRegistry => self.capability_registry.is_some(),
1784 ToolContextService::ToolRegistry => self.tool_registry.is_some(),
1785 ToolContextService::OrgId => self.org_id.is_some(),
1786 ToolContextService::BudgetChecker => self.budget_checker.is_some(),
1787 ToolContextService::PaymentAuthority => self.payment_authority.is_some(),
1788 ToolContextService::SessionCreationAuthority => {
1789 self.session_creation_authority.is_some()
1790 }
1791 ToolContextService::SubagentSpawnStore => self.subagent_spawn_store.is_some(),
1792 ToolContextService::ReasoningEffortHandle => self.reasoning_effort_handle.is_some(),
1793 }
1794 }
1795}
1796
1797#[derive(Clone)]
1806pub struct ToolContext {
1807 pub session_id: SessionId,
1809 pub workspace_id: WorkspaceId,
1816
1817 pub file_store: Option<Arc<dyn SessionFileSystem>>,
1819
1820 pub storage_store: Option<Arc<dyn SessionStorageStore>>,
1822
1823 pub image_store: Option<Arc<dyn ImageArtifactStore>>,
1825
1826 pub provider_credential_store: Option<Arc<dyn ProviderCredentialStore>>,
1828
1829 pub utility_llm_service: Option<Arc<dyn crate::UtilityLlmService>>,
1831
1832 pub mcp_invoker: Option<Arc<dyn crate::McpToolInvoker>>,
1838
1839 pub egress_service: Option<Arc<dyn crate::EgressService>>,
1841
1842 pub sqldb_store: Option<SessionSqlDbStoreRef>,
1844
1845 pub message_retriever: Option<Arc<dyn crate::message_retriever::MessageRetriever>>,
1847
1848 pub session_store: Option<Arc<dyn SessionStore>>,
1850
1851 pub session_mutator: Option<Arc<dyn SessionMutator>>,
1853
1854 pub agent_store: Option<Arc<dyn AgentStore>>,
1856
1857 pub connection_resolver: Option<Arc<dyn UserConnectionResolver>>,
1859
1860 pub schedule_store: Option<Arc<dyn SessionScheduleStore>>,
1862
1863 pub platform_store: Option<Arc<dyn crate::platform_store::PlatformStore>>,
1865 pub knowledge_store: Option<Arc<dyn KnowledgeStore>>,
1867
1868 pub knowledge_index_search: Option<Arc<dyn crate::vector_store::KnowledgeIndexSearch>>,
1872
1873 pub leased_resource_store: Option<Arc<dyn LeasedResourceStore>>,
1875
1876 pub session_resource_registry: Option<Arc<dyn SessionResourceRegistry>>,
1878
1879 pub session_task_registry: Option<Arc<dyn crate::session_task::SessionTaskRegistry>>,
1882
1883 pub event_emitter: Option<Arc<dyn EventEmitter>>,
1886
1887 pub event_context: Option<crate::events::EventContext>,
1890
1891 pub tool_call_id: Option<String>,
1894 pub capability_registry: Option<crate::capabilities::CapabilityRegistry>,
1896
1897 pub tool_registry: Option<Arc<crate::tools::ToolRegistry>>,
1900
1901 pub visible_tool_names: Option<Arc<HashSet<String>>>,
1905
1906 pub org_id: Option<crate::typed_id::OrgId>,
1908
1909 pub network_access: Option<crate::network_access::NetworkAccessList>,
1912
1913 pub locale: Option<String>,
1917
1918 pub budget_checker: Option<Arc<dyn BudgetChecker>>,
1920
1921 pub payment_authority: Option<Arc<dyn PaymentAuthority>>,
1923
1924 pub session_creation_authority: Option<Arc<dyn SessionCreationAuthority>>,
1926
1927 pub subagent_spawn_store: Option<Arc<dyn SubagentSpawnStore>>,
1931
1932 pub subagent_nesting_policy: SubagentNestingPolicy,
1934
1935 pub reasoning_effort_handle: Option<ReasoningEffortHandle>,
1939
1940 pub cancellation: Option<tokio_util::sync::CancellationToken>,
1950}
1951
1952impl ToolContext {
1953 pub fn workspace_fs_key(&self) -> SessionId {
1958 SessionId::from_uuid(self.workspace_id.uuid())
1959 }
1960
1961 pub fn with_workspace_id(mut self, workspace_id: WorkspaceId) -> Self {
1963 self.workspace_id = workspace_id;
1964 self
1965 }
1966
1967 pub fn new(session_id: SessionId) -> Self {
1969 Self {
1970 session_id,
1971 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
1972 file_store: None,
1973 storage_store: None,
1974 image_store: None,
1975 provider_credential_store: None,
1976 utility_llm_service: None,
1977 mcp_invoker: None,
1978 egress_service: None,
1979 sqldb_store: None,
1980 message_retriever: None,
1981 session_store: None,
1982 session_mutator: None,
1983 agent_store: None,
1984 connection_resolver: None,
1985 schedule_store: None,
1986 platform_store: None,
1987 knowledge_store: None,
1988 knowledge_index_search: None,
1989 leased_resource_store: None,
1990 session_resource_registry: None,
1991 session_task_registry: None,
1992 event_emitter: None,
1993 event_context: None,
1994 tool_call_id: None,
1995 capability_registry: None,
1996 tool_registry: None,
1997 visible_tool_names: None,
1998 org_id: None,
1999 network_access: None,
2000 locale: None,
2001 budget_checker: None,
2002 payment_authority: None,
2003 session_creation_authority: None,
2004 subagent_spawn_store: None,
2005 subagent_nesting_policy: SubagentNestingPolicy::default(),
2006 reasoning_effort_handle: None,
2007 cancellation: None,
2008 }
2009 }
2010
2011 pub fn from_services(session_id: SessionId, services: &ToolContextServices) -> Self {
2013 Self {
2014 session_id,
2015 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
2016 file_store: services.file_store.clone(),
2017 storage_store: services.storage_store.clone(),
2018 image_store: services.image_store.clone(),
2019 provider_credential_store: services.provider_credential_store.clone(),
2020 utility_llm_service: services.utility_llm_service.clone(),
2021 mcp_invoker: services.mcp_invoker.clone(),
2022 egress_service: services.egress_service.clone(),
2023 sqldb_store: services.sqldb_store.clone(),
2024 message_retriever: services.message_retriever.clone(),
2025 session_store: services.session_store.clone(),
2026 session_mutator: services.session_mutator.clone(),
2027 agent_store: services.agent_store.clone(),
2028 connection_resolver: services.connection_resolver.clone(),
2029 schedule_store: services.schedule_store.clone(),
2030 platform_store: services.platform_store.clone(),
2031 knowledge_store: services.knowledge_store.clone(),
2032 knowledge_index_search: services.knowledge_index_search.clone(),
2033 leased_resource_store: services.leased_resource_store.clone(),
2034 session_resource_registry: services.session_resource_registry.clone(),
2035 session_task_registry: services.session_task_registry.clone(),
2036 event_emitter: services.event_emitter.clone(),
2037 event_context: None,
2038 tool_call_id: None,
2039 capability_registry: services.capability_registry.clone(),
2040 tool_registry: services.tool_registry.clone(),
2041 visible_tool_names: None,
2042 org_id: services.org_id,
2043 network_access: services.network_access.clone(),
2044 locale: None,
2045 budget_checker: services.budget_checker.clone(),
2046 payment_authority: services.payment_authority.clone(),
2047 session_creation_authority: services.session_creation_authority.clone(),
2048 subagent_spawn_store: services.subagent_spawn_store.clone(),
2049 subagent_nesting_policy: services.subagent_nesting_policy,
2050 reasoning_effort_handle: services.reasoning_effort_handle.clone(),
2051 cancellation: None,
2052 }
2053 }
2054
2055 pub fn with_file_store(session_id: SessionId, file_store: Arc<dyn SessionFileSystem>) -> Self {
2057 Self {
2058 session_id,
2059 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
2060 file_store: Some(file_store),
2061 storage_store: None,
2062 image_store: None,
2063 provider_credential_store: None,
2064 utility_llm_service: None,
2065 mcp_invoker: None,
2066 egress_service: None,
2067 sqldb_store: None,
2068 message_retriever: None,
2069 session_store: None,
2070 session_mutator: None,
2071 agent_store: None,
2072 connection_resolver: None,
2073 schedule_store: None,
2074 platform_store: None,
2075 knowledge_store: None,
2076 knowledge_index_search: None,
2077 leased_resource_store: None,
2078 session_resource_registry: None,
2079 session_task_registry: None,
2080 event_emitter: None,
2081 event_context: None,
2082 tool_call_id: None,
2083 capability_registry: None,
2084 tool_registry: None,
2085 visible_tool_names: None,
2086 org_id: None,
2087 network_access: None,
2088 locale: None,
2089 budget_checker: None,
2090 payment_authority: None,
2091 session_creation_authority: None,
2092 subagent_spawn_store: None,
2093 subagent_nesting_policy: SubagentNestingPolicy::default(),
2094 reasoning_effort_handle: None,
2095 cancellation: None,
2096 }
2097 }
2098
2099 pub fn with_storage_store(
2101 session_id: SessionId,
2102 storage_store: Arc<dyn SessionStorageStore>,
2103 ) -> Self {
2104 Self {
2105 session_id,
2106 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
2107 file_store: None,
2108 storage_store: Some(storage_store),
2109 image_store: None,
2110 provider_credential_store: None,
2111 utility_llm_service: None,
2112 mcp_invoker: None,
2113 egress_service: None,
2114 sqldb_store: None,
2115 message_retriever: None,
2116 session_store: None,
2117 session_mutator: None,
2118 agent_store: None,
2119 connection_resolver: None,
2120 schedule_store: None,
2121 platform_store: None,
2122 knowledge_store: None,
2123 knowledge_index_search: None,
2124 leased_resource_store: None,
2125 session_resource_registry: None,
2126 session_task_registry: None,
2127 event_emitter: None,
2128 event_context: None,
2129 tool_call_id: None,
2130 capability_registry: None,
2131 tool_registry: None,
2132 visible_tool_names: None,
2133 org_id: None,
2134 network_access: None,
2135 locale: None,
2136 budget_checker: None,
2137 payment_authority: None,
2138 session_creation_authority: None,
2139 subagent_spawn_store: None,
2140 subagent_nesting_policy: SubagentNestingPolicy::default(),
2141 reasoning_effort_handle: None,
2142 cancellation: None,
2143 }
2144 }
2145
2146 pub fn with_stores(
2148 session_id: SessionId,
2149 file_store: Arc<dyn SessionFileSystem>,
2150 storage_store: Arc<dyn SessionStorageStore>,
2151 ) -> Self {
2152 Self {
2153 session_id,
2154 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
2155 file_store: Some(file_store),
2156 storage_store: Some(storage_store),
2157 sqldb_store: None,
2158 image_store: None,
2159 provider_credential_store: None,
2160 utility_llm_service: None,
2161 mcp_invoker: None,
2162 egress_service: None,
2163 message_retriever: None,
2164 session_store: None,
2165 session_mutator: None,
2166 agent_store: None,
2167 connection_resolver: None,
2168 schedule_store: None,
2169 platform_store: None,
2170 knowledge_store: None,
2171 knowledge_index_search: None,
2172 leased_resource_store: None,
2173 session_resource_registry: None,
2174 session_task_registry: None,
2175 event_emitter: None,
2176 event_context: None,
2177 tool_call_id: None,
2178 capability_registry: None,
2179 tool_registry: None,
2180 visible_tool_names: None,
2181 org_id: None,
2182 network_access: None,
2183 locale: None,
2184 budget_checker: None,
2185 payment_authority: None,
2186 session_creation_authority: None,
2187 subagent_spawn_store: None,
2188 subagent_nesting_policy: SubagentNestingPolicy::default(),
2189 reasoning_effort_handle: None,
2190 cancellation: None,
2191 }
2192 }
2193
2194 pub fn with_sqldb_store(mut self, sqldb_store: SessionSqlDbStoreRef) -> Self {
2196 self.sqldb_store = Some(sqldb_store);
2197 self
2198 }
2199
2200 pub fn with_message_retriever(
2202 mut self,
2203 retriever: Arc<dyn crate::message_retriever::MessageRetriever>,
2204 ) -> Self {
2205 self.message_retriever = Some(retriever);
2206 self
2207 }
2208
2209 pub fn with_session_store(mut self, store: Arc<dyn SessionStore>) -> Self {
2211 self.session_store = Some(store);
2212 self
2213 }
2214
2215 pub fn with_session_mutator(mut self, mutator: Arc<dyn SessionMutator>) -> Self {
2217 self.session_mutator = Some(mutator);
2218 self
2219 }
2220
2221 pub fn with_cancellation(mut self, token: tokio_util::sync::CancellationToken) -> Self {
2227 self.cancellation = Some(token);
2228 self
2229 }
2230
2231 pub fn is_cancelled(&self) -> bool {
2234 self.cancellation
2235 .as_ref()
2236 .is_some_and(|token| token.is_cancelled())
2237 }
2238
2239 pub fn with_reasoning_effort_handle(mut self, handle: ReasoningEffortHandle) -> Self {
2240 self.reasoning_effort_handle = Some(handle);
2241 self
2242 }
2243
2244 pub fn with_agent_store(mut self, store: Arc<dyn AgentStore>) -> Self {
2246 self.agent_store = Some(store);
2247 self
2248 }
2249
2250 pub fn with_connection_resolver(mut self, resolver: Arc<dyn UserConnectionResolver>) -> Self {
2252 self.connection_resolver = Some(resolver);
2253 self
2254 }
2255
2256 pub fn with_image_store(
2258 session_id: SessionId,
2259 image_store: Arc<dyn ImageArtifactStore>,
2260 ) -> Self {
2261 Self {
2262 session_id,
2263 workspace_id: WorkspaceId::from_uuid(session_id.uuid()),
2264 file_store: None,
2265 storage_store: None,
2266 image_store: Some(image_store),
2267 provider_credential_store: None,
2268 utility_llm_service: None,
2269 mcp_invoker: None,
2270 egress_service: None,
2271 sqldb_store: None,
2272 message_retriever: None,
2273 session_store: None,
2274 session_mutator: None,
2275 agent_store: None,
2276 connection_resolver: None,
2277 schedule_store: None,
2278 platform_store: None,
2279 knowledge_store: None,
2280 knowledge_index_search: None,
2281 leased_resource_store: None,
2282 session_resource_registry: None,
2283 session_task_registry: None,
2284 event_emitter: None,
2285 event_context: None,
2286 tool_call_id: None,
2287 capability_registry: None,
2288 tool_registry: None,
2289 visible_tool_names: None,
2290 org_id: None,
2291 network_access: None,
2292 locale: None,
2293 budget_checker: None,
2294 payment_authority: None,
2295 session_creation_authority: None,
2296 subagent_spawn_store: None,
2297 subagent_nesting_policy: SubagentNestingPolicy::default(),
2298 reasoning_effort_handle: None,
2299 cancellation: None,
2300 }
2301 }
2302
2303 pub fn with_provider_credential_store(
2305 mut self,
2306 store: Arc<dyn ProviderCredentialStore>,
2307 ) -> Self {
2308 self.provider_credential_store = Some(store);
2309 self
2310 }
2311
2312 pub fn with_utility_llm_service(mut self, service: Arc<dyn crate::UtilityLlmService>) -> Self {
2314 self.utility_llm_service = Some(service);
2315 self
2316 }
2317
2318 pub fn with_mcp_invoker(mut self, invoker: Arc<dyn crate::McpToolInvoker>) -> Self {
2320 self.mcp_invoker = Some(invoker);
2321 self
2322 }
2323
2324 pub fn with_egress_service(mut self, service: Arc<dyn crate::EgressService>) -> Self {
2326 self.egress_service = Some(service);
2327 self
2328 }
2329
2330 pub fn with_egress_service_opt(
2333 mut self,
2334 service: Option<Arc<dyn crate::EgressService>>,
2335 ) -> Self {
2336 if let Some(service) = service {
2337 self.egress_service = Some(service);
2338 }
2339 self
2340 }
2341
2342 pub fn with_storage_store_arc(mut self, store: Arc<dyn SessionStorageStore>) -> Self {
2344 self.storage_store = Some(store);
2345 self
2346 }
2347
2348 pub fn with_schedule_store(mut self, store: Arc<dyn SessionScheduleStore>) -> Self {
2350 self.schedule_store = Some(store);
2351 self
2352 }
2353
2354 pub fn with_platform_store(
2356 mut self,
2357 store: Arc<dyn crate::platform_store::PlatformStore>,
2358 ) -> Self {
2359 self.platform_store = Some(store);
2360 self
2361 }
2362
2363 pub fn with_knowledge_index_search(
2365 mut self,
2366 search: Arc<dyn crate::vector_store::KnowledgeIndexSearch>,
2367 ) -> Self {
2368 self.knowledge_index_search = Some(search);
2369 self
2370 }
2371
2372 pub fn with_leased_resource_store(mut self, store: Arc<dyn LeasedResourceStore>) -> Self {
2374 self.leased_resource_store = Some(store);
2375 self
2376 }
2377
2378 pub fn with_session_resource_registry(
2380 mut self,
2381 registry: Arc<dyn SessionResourceRegistry>,
2382 ) -> Self {
2383 self.session_resource_registry = Some(registry);
2384 self
2385 }
2386
2387 pub fn with_session_task_registry(
2389 mut self,
2390 registry: Arc<dyn crate::session_task::SessionTaskRegistry>,
2391 ) -> Self {
2392 self.session_task_registry = Some(registry);
2393 self
2394 }
2395
2396 pub fn with_org_id(mut self, org_id: crate::typed_id::OrgId) -> Self {
2398 self.org_id = Some(org_id);
2399 self
2400 }
2401
2402 pub fn with_tool_registry(mut self, registry: Arc<crate::tools::ToolRegistry>) -> Self {
2404 self.tool_registry = Some(registry);
2405 self
2406 }
2407
2408 pub fn with_visible_tool_names(mut self, names: Arc<HashSet<String>>) -> Self {
2410 self.visible_tool_names = Some(names);
2411 self
2412 }
2413
2414 pub fn with_network_access(
2416 mut self,
2417 network_access: Option<crate::network_access::NetworkAccessList>,
2418 ) -> Self {
2419 self.network_access = network_access;
2420 self
2421 }
2422
2423 pub fn with_payment_authority(mut self, authority: Arc<dyn PaymentAuthority>) -> Self {
2425 self.payment_authority = Some(authority);
2426 self
2427 }
2428
2429 pub fn with_subagent_spawn_store(mut self, store: Arc<dyn SubagentSpawnStore>) -> Self {
2431 self.subagent_spawn_store = Some(store);
2432 self
2433 }
2434
2435 pub fn with_subagent_nesting_policy(mut self, policy: SubagentNestingPolicy) -> Self {
2437 self.subagent_nesting_policy = policy;
2438 self
2439 }
2440
2441 pub async fn emit_progress(&self, tool_name: &str, message: &str) {
2446 let (Some(emitter), Some(ctx), Some(call_id)) =
2447 (&self.event_emitter, &self.event_context, &self.tool_call_id)
2448 else {
2449 return;
2450 };
2451 if let Err(e) = emitter
2452 .emit(EventRequest::new(
2453 self.session_id,
2454 ctx.clone(),
2455 crate::events::ToolProgressData {
2456 tool_call_id: call_id.clone(),
2457 tool_name: tool_name.to_string(),
2458 message: message.to_string(),
2459 display_name: None,
2460 },
2461 ))
2462 .await
2463 {
2464 tracing::debug!(
2465 tool_call_id = call_id,
2466 tool_name,
2467 error = %e,
2468 "Failed to emit tool.progress event"
2469 );
2470 }
2471 }
2472
2473 pub async fn emit_tool_output(&self, tool_name: &str, delta: &str, stream: &str) {
2478 let (Some(emitter), Some(ctx), Some(call_id)) =
2479 (&self.event_emitter, &self.event_context, &self.tool_call_id)
2480 else {
2481 return;
2482 };
2483 if let Err(e) = emitter
2484 .emit(EventRequest::new(
2485 self.session_id,
2486 ctx.clone(),
2487 crate::events::ToolOutputDeltaData {
2488 tool_call_id: call_id.clone(),
2489 tool_name: tool_name.to_string(),
2490 delta: delta.to_string(),
2491 stream: stream.to_string(),
2492 },
2493 ))
2494 .await
2495 {
2496 tracing::debug!(
2497 tool_call_id = call_id,
2498 tool_name,
2499 error = %e,
2500 "Failed to emit tool.output.delta event"
2501 );
2502 }
2503 }
2504}
2505
2506impl std::fmt::Debug for ToolContext {
2507 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2508 f.debug_struct("ToolContext")
2509 .field("session_id", &self.session_id)
2510 .field("file_store", &self.file_store.is_some())
2511 .field("storage_store", &self.storage_store.is_some())
2512 .field("image_store", &self.image_store.is_some())
2513 .field(
2514 "provider_credential_store",
2515 &self.provider_credential_store.is_some(),
2516 )
2517 .field("utility_llm_service", &self.utility_llm_service.is_some())
2518 .field("egress_service", &self.egress_service.is_some())
2519 .field("sqldb_store", &self.sqldb_store.is_some())
2520 .field("message_retriever", &self.message_retriever.is_some())
2521 .field("session_store", &self.session_store.is_some())
2522 .field("session_mutator", &self.session_mutator.is_some())
2523 .field("agent_store", &self.agent_store.is_some())
2524 .field("connection_resolver", &self.connection_resolver.is_some())
2525 .field("schedule_store", &self.schedule_store.is_some())
2526 .field("platform_store", &self.platform_store.is_some())
2527 .field(
2528 "knowledge_index_search",
2529 &self.knowledge_index_search.is_some(),
2530 )
2531 .field(
2532 "leased_resource_store",
2533 &self.leased_resource_store.is_some(),
2534 )
2535 .field("event_emitter", &self.event_emitter.is_some())
2536 .field("tool_registry", &self.tool_registry.is_some())
2537 .field("payment_authority", &self.payment_authority.is_some())
2538 .field("subagent_spawn_store", &self.subagent_spawn_store.is_some())
2539 .field("subagent_nesting_policy", &self.subagent_nesting_policy)
2540 .field("org_id", &self.org_id)
2541 .finish()
2542 }
2543}
2544
2545use crate::events::{Event, EventRequest};
2550
2551#[async_trait]
2562pub trait EventEmitter: Send + Sync {
2563 async fn emit(&self, request: EventRequest) -> Result<Event>;
2568}
2569
2570#[async_trait]
2572impl<E: EventEmitter + ?Sized> EventEmitter for Arc<E> {
2573 async fn emit(&self, request: EventRequest) -> Result<Event> {
2574 (**self).emit(request).await
2575 }
2576}
2577
2578#[derive(Debug, Clone, Default)]
2582pub struct NoopEventEmitter;
2583
2584#[async_trait]
2585impl EventEmitter for NoopEventEmitter {
2586 async fn emit(&self, request: EventRequest) -> Result<Event> {
2587 Ok(request.into_event(crate::typed_id::EventId::new(), 0))
2589 }
2590}
2591
2592#[derive(Debug, Clone)]
2605pub struct ResolvedImage {
2606 pub base64: String,
2608 pub media_type: String,
2610}
2611
2612impl ResolvedImage {
2613 pub fn new(base64: impl Into<String>, media_type: impl Into<String>) -> Self {
2615 Self {
2616 base64: base64.into(),
2617 media_type: media_type.into(),
2618 }
2619 }
2620
2621 pub fn to_data_url(&self) -> String {
2625 format!("data:{};base64,{}", self.media_type, self.base64)
2626 }
2627}
2628
2629#[async_trait]
2662pub trait ImageResolver: Send + Sync {
2663 async fn resolve_image(&self, image_id: Uuid) -> Result<Option<ResolvedImage>>;
2667}
2668
2669#[derive(Debug)]
2675pub enum SpawnClaimResult {
2676 Claimed {
2679 spawn_handle_id: uuid::Uuid,
2680 claim_token: uuid::Uuid,
2681 },
2682 ClaimedPendingChild {
2686 spawn_handle_id: uuid::Uuid,
2687 claim_token: uuid::Uuid,
2688 },
2689 AlreadyRunning {
2692 child_session_id: crate::typed_id::SessionId,
2693 claim_token: uuid::Uuid,
2695 },
2696 AlreadySettled {
2699 child_session_id: crate::typed_id::SessionId,
2700 terminal_status: String,
2702 terminal_result: String,
2703 },
2704}
2705
2706#[async_trait]
2714pub trait SubagentSpawnStore: Send + Sync + 'static {
2715 async fn try_claim_spawn(
2720 &self,
2721 parent_session_id: crate::typed_id::SessionId,
2722 tool_call_id: &str,
2723 claim_token: uuid::Uuid,
2724 ) -> Result<SpawnClaimResult>;
2725
2726 async fn register_child_session(
2731 &self,
2732 spawn_handle_id: uuid::Uuid,
2733 claim_token: uuid::Uuid,
2734 child_session_id: crate::typed_id::SessionId,
2735 ) -> Result<()>;
2736
2737 async fn settle_spawn(
2743 &self,
2744 parent_session_id: crate::typed_id::SessionId,
2745 tool_call_id: &str,
2746 claim_token: uuid::Uuid,
2747 terminal_status: &str,
2748 terminal_result: &str,
2749 ) -> Result<()>;
2750}
2751
2752#[async_trait]
2754impl<S: SubagentSpawnStore + ?Sized> SubagentSpawnStore for Arc<S> {
2755 async fn try_claim_spawn(
2756 &self,
2757 parent_session_id: crate::typed_id::SessionId,
2758 tool_call_id: &str,
2759 claim_token: uuid::Uuid,
2760 ) -> Result<SpawnClaimResult> {
2761 (**self)
2762 .try_claim_spawn(parent_session_id, tool_call_id, claim_token)
2763 .await
2764 }
2765
2766 async fn register_child_session(
2767 &self,
2768 spawn_handle_id: uuid::Uuid,
2769 claim_token: uuid::Uuid,
2770 child_session_id: crate::typed_id::SessionId,
2771 ) -> Result<()> {
2772 (**self)
2773 .register_child_session(spawn_handle_id, claim_token, child_session_id)
2774 .await
2775 }
2776
2777 async fn settle_spawn(
2778 &self,
2779 parent_session_id: crate::typed_id::SessionId,
2780 tool_call_id: &str,
2781 claim_token: uuid::Uuid,
2782 terminal_status: &str,
2783 terminal_result: &str,
2784 ) -> Result<()> {
2785 (**self)
2786 .settle_spawn(
2787 parent_session_id,
2788 tool_call_id,
2789 claim_token,
2790 terminal_status,
2791 terminal_result,
2792 )
2793 .await
2794 }
2795}
2796
2797pub struct NoopSubagentSpawnStore;
2801
2802#[async_trait]
2803impl SubagentSpawnStore for NoopSubagentSpawnStore {
2804 async fn try_claim_spawn(
2805 &self,
2806 _parent_session_id: crate::typed_id::SessionId,
2807 _tool_call_id: &str,
2808 claim_token: uuid::Uuid,
2809 ) -> Result<SpawnClaimResult> {
2810 Ok(SpawnClaimResult::Claimed {
2811 spawn_handle_id: uuid::Uuid::new_v4(),
2812 claim_token,
2813 })
2814 }
2815
2816 async fn register_child_session(
2817 &self,
2818 _spawn_handle_id: uuid::Uuid,
2819 _claim_token: uuid::Uuid,
2820 _child_session_id: crate::typed_id::SessionId,
2821 ) -> Result<()> {
2822 Ok(())
2823 }
2824
2825 async fn settle_spawn(
2826 &self,
2827 _parent_session_id: crate::typed_id::SessionId,
2828 _tool_call_id: &str,
2829 _claim_token: uuid::Uuid,
2830 _terminal_status: &str,
2831 _terminal_result: &str,
2832 ) -> Result<()> {
2833 Ok(())
2834 }
2835}
2836
2837#[cfg(test)]
2842mod tests {
2843 use super::*;
2844
2845 #[test]
2846 fn test_resolved_image_new() {
2847 let image = ResolvedImage::new("SGVsbG8=", "image/png");
2848 assert_eq!(image.base64, "SGVsbG8=");
2849 assert_eq!(image.media_type, "image/png");
2850 }
2851
2852 #[test]
2853 fn test_resolved_image_to_data_url() {
2854 let image = ResolvedImage::new("SGVsbG8=", "image/png");
2855 let data_url = image.to_data_url();
2856 assert_eq!(data_url, "data:image/png;base64,SGVsbG8=");
2857 }
2858
2859 #[test]
2860 fn test_resolved_image_jpeg() {
2861 let image = ResolvedImage::new("base64data", "image/jpeg");
2862 let data_url = image.to_data_url();
2863 assert!(data_url.starts_with("data:image/jpeg;base64,"));
2864 }
2865}