1#![warn(missing_docs)]
28
29mod envelope;
30pub mod events;
31pub mod flow;
32pub mod session;
33pub mod storage;
34
35pub use events::{AgentEvent, FlowStream, TurnStream};
36pub use flow::{assemble_registry, ExecutionResult, FlowClient, FlowClientBuilder};
37pub use session::{Fork, Session};
38pub use storage::Storage;
39
40pub use flux_flow::AgentSink;
44
45pub use tokio_util::sync::CancellationToken;
48
49pub use flux_provider::Provider;
54
55#[allow(deprecated)]
56pub use flux_agent::AgentExecutorConfig;
57pub use flux_agent::{
60 AdaptiveLoopPolicy, AgentLoopSpec, AgentSpec, AgentStagePolicy, BuiltinAgentLoop, Permissions,
61};
62
63pub use flux_core::Usage;
65
66pub use flux_core::PricingTable;
70
71pub use flux_flow::replay::ReplayReport;
74
75pub mod tools {
83 pub use flux_runtime::{tool_fn, FnTool, Tool, ToolContext, ToolRegistry, ToolResult};
84 pub use flux_spec::{Risk, ToolSpec};
85}
86
87pub fn stage_fn<I, O, F, Fut, E>(
93 name: impl Into<String>,
94 description: impl Into<String>,
95 handler: F,
96) -> Arc<dyn tools::Tool>
97where
98 I: serde::de::DeserializeOwned + schemars::JsonSchema + Send + 'static,
99 O: serde::Serialize + schemars::JsonSchema + Send + 'static,
100 F: Fn(I) -> Fut + Send + Sync + 'static,
101 Fut: Future<Output = std::result::Result<O, E>> + Send + 'static,
102 E: std::fmt::Display + Send + 'static,
103{
104 let spec = flux_spec::ToolSpec::read_only_typed::<I>(name, description)
105 .with_output_schema(flux_spec::tool_output_schema::<O>());
106 let handler = Arc::new(handler);
107 Arc::new(
108 flux_runtime::FnTool::new(spec, move |value| {
109 let handler = handler.clone();
110 async move {
111 let input = serde_json::from_value::<I>(value)
112 .map_err(|error| format!("invalid stage input: {error}"))?;
113 let output = handler(input).await.map_err(|error| error.to_string())?;
114 serde_json::to_value(output)
115 .map_err(|error| format!("stage output did not serialize: {error}"))
116 }
117 })
118 .with_staging_disposition(flux_spec::StagingDisposition::Gather),
119 )
120}
121
122pub mod approval {
127 pub use flux_runtime::{ApprovalChoice, Approver, RiskApprover};
128 pub use flux_spec::IntentSet;
129}
130
131pub mod authorization {
135 pub use flux_policy::{
136 default_local_grants, local_identity, AuthorizationPolicy, Caller, Trust,
137 };
138 pub use flux_runtime::{
139 ExecutionAuthorization, ExecutionEnvironment, IdentityCell, TurnIdentity,
140 };
141}
142
143pub mod observe {
157 pub use flux_core::Message;
158 pub use flux_events::EventStore;
159 pub use flux_events::{DiffRow, EfficiencySummary, ModelCost, RunDiff, TurnSummary};
160 pub use flux_evidence::{Observation, SignalMatch, ToolGroup, KIND_SIGNAL};
161 pub use flux_flow::state::FlowStore;
162 pub use flux_lang::ast::RunEvent;
163}
164
165pub mod voice {
174 pub use flux_flow::voice::{VoiceReply, VoiceSink};
175 pub use flux_provider::{RealtimeConfig, RealtimeProvider};
176}
177
178#[cfg(feature = "pricing")]
185pub mod pricing {
186 pub use flux_credentials::load_pricing_table;
187}
188
189#[cfg(feature = "providers")]
206pub mod providers {
207 pub use flux_providers::{anthropic, bedrock, codex, ollama, openai, openrouter, spec};
208
209 use flux_core::Result;
210 use flux_provider::Provider;
211
212 pub fn from_spec(spec: &str) -> Result<(Box<dyn Provider>, String)> {
219 let (native, _provider, model) = spec::build(spec)?;
220 Ok((Box::new(native), model))
221 }
222}
223
224#[cfg(feature = "plugins")]
237pub mod plugins {
238 pub use flux_plugin::{HostCapabilities, PluginDescriptor, PluginManifest, SystemHostCaps};
239
240 use std::sync::Arc;
241
242 use flux_core::Result;
243 use flux_runtime::Tool;
244 use flux_system::System;
245
246 pub async fn load_tools(
252 system: &Arc<System>,
253 name: &str,
254 descriptor: &PluginDescriptor,
255 ) -> Result<Vec<Arc<dyn Tool>>> {
256 let caps_system = system.clone();
257 let loaded = flux_plugin::load_plugin_tools(system, name, descriptor, move |m| {
258 Arc::new(SystemHostCaps::new(caps_system).with_manifest(m)) as Arc<dyn HostCapabilities>
259 })
260 .await?;
261 Ok(loaded.tools)
262 }
263}
264
265pub mod subagents {
276 #[allow(deprecated)]
277 pub use flux_orchestrate::{
278 parse_role, try_parse_role, ProviderFactory, Role, RoleRegistry, SpawnLimits, SubAgents,
279 };
280}
281
282pub use flux_system::sandbox::{Sandbox, SandboxSettings};
286
287pub use flux_lang::dsl;
292
293pub mod recipes;
294
295use std::future::Future;
296use std::path::PathBuf;
297use std::sync::Arc;
298
299use flux_cognition::CognitionPack;
301use flux_core::ContextBlock;
302use flux_core::Result;
303use flux_events::EventStore;
304use flux_flow::engine::FlowEngine;
305use flux_orchestrate::{SubAgents, TaskTool};
306#[cfg(test)]
307use flux_runtime::ToolContext;
308use flux_runtime::{Approver, ExecutionEnvironment, PermissionManager, Tool, ToolRegistry};
309use flux_secret::Redactor;
310use flux_system::{System, Workspace};
311
312#[derive(Debug, Default, Clone)]
318#[non_exhaustive]
319pub struct TurnOutput {
320 pub text: String,
322 pub tool_calls: Vec<String>,
324 pub usage: Option<Usage>,
326 pub suspended: bool,
331}
332
333type RegistryPack = Box<dyn FnOnce(&mut ToolRegistry) -> Result<()>>;
335
336pub struct ClientBuilder {
340 spec: AgentSpec,
341 envelope: envelope::Envelope,
342 storage: Option<Storage>,
343 cognition: bool,
344 ops: Vec<(String, Arc<dyn Tool>)>,
345 packs: Vec<RegistryPack>,
346 sub_agents: Option<SubAgents>,
347 sub_agent_adaptive_policy: Option<AdaptiveLoopPolicy>,
348}
349
350impl Default for ClientBuilder {
351 fn default() -> Self {
352 Self {
353 spec: AgentSpec::new("unknown"),
354 envelope: envelope::Envelope::with_default_allow(&["read"]),
356 storage: None,
358 cognition: false,
359 ops: Vec::new(),
360 packs: Vec::new(),
361 sub_agents: None,
363 sub_agent_adaptive_policy: None,
364 }
365 }
366}
367
368impl ClientBuilder {
369 pub fn from_spec(spec: AgentSpec) -> Self {
377 Self {
378 spec,
379 envelope: envelope::Envelope::bare(),
380 storage: None,
381 cognition: false,
382 ops: Vec::new(),
383 packs: Vec::new(),
384 sub_agents: None,
385 sub_agent_adaptive_policy: None,
386 }
387 }
388 pub fn model(mut self, m: impl Into<String>) -> Self {
390 self.spec.model = m.into();
391 self
392 }
393 pub fn system_prompt(mut self, s: impl Into<String>) -> Self {
395 self.spec.system_prompt = s.into();
396 self
397 }
398 pub fn max_tokens(mut self, n: u32) -> Self {
400 self.spec.max_tokens = n;
401 self
402 }
403 pub fn max_iterations(mut self, n: usize) -> Self {
407 self.spec.max_iterations = n;
408 self
409 }
410 pub fn adaptive_policy(mut self, policy: AdaptiveLoopPolicy) -> Self {
413 self.spec.adaptive_policy = policy;
414 self
415 }
416 pub fn agent_loop(mut self, agent_loop: AgentLoopSpec) -> Self {
418 self.spec.agent_loop = agent_loop;
419 self
420 }
421 pub fn allow(mut self, rule: impl Into<String>) -> Self {
423 self.envelope.allow.push(rule.into());
424 self
425 }
426 pub fn deny(mut self, rule: impl Into<String>) -> Self {
428 self.envelope.deny.push(rule.into());
429 self
430 }
431 pub fn auto_approve(mut self, yes: bool) -> Self {
433 self.envelope.auto_approve = yes;
434 self
435 }
436 pub fn approver(mut self, approver: Arc<dyn Approver>) -> Self {
442 self.envelope.approver = Some(approver);
443 self
444 }
445 pub fn with_authorization(
448 mut self,
449 policy: flux_policy::AuthorizationPolicy,
450 caller: flux_policy::Caller,
451 trust: flux_policy::Trust,
452 ) -> Self {
453 self.envelope.authorization =
454 flux_runtime::ExecutionAuthorization::new(policy, caller, trust);
455 self
456 }
457 pub fn with_redactor(mut self, redactor: Redactor) -> Self {
459 self.envelope.redactor = redactor;
460 self
461 }
462 pub fn register_op(mut self, tool: Arc<dyn Tool>) -> Self {
466 self.ops
467 .push(("sdk ClientBuilder::register_op".into(), tool));
468 self
469 }
470
471 pub fn register_op_from(mut self, source: impl Into<String>, tool: Arc<dyn Tool>) -> Self {
474 self.ops.push((source.into(), tool));
475 self
476 }
477 pub fn register_pack<F: FnOnce(&mut ToolRegistry) + 'static>(mut self, pack: F) -> Self {
480 self.packs.push(Box::new(move |registry| {
481 pack(registry);
482 Ok(())
483 }));
484 self
485 }
486 pub fn try_register_pack<F>(mut self, pack: F) -> Self
489 where
490 F: FnOnce(&mut ToolRegistry) -> Result<()> + 'static,
491 {
492 self.packs.push(Box::new(pack));
493 self
494 }
495 #[cfg(feature = "plugins")]
505 pub fn with_plugin_tools(mut self, tools: Vec<Arc<dyn Tool>>) -> Self {
506 self.ops.extend(tools.into_iter().map(|tool| {
507 (
508 "sdk ClientBuilder::with_plugin_tools (plugin source unspecified)".into(),
509 tool,
510 )
511 }));
512 self
513 }
514
515 #[cfg(feature = "plugins")]
517 pub fn with_plugin_tools_from(
518 mut self,
519 plugin: impl Into<String>,
520 tools: Vec<Arc<dyn Tool>>,
521 ) -> Self {
522 let source = format!("plugin:{}", plugin.into());
523 self.ops
524 .extend(tools.into_iter().map(|tool| (source.clone(), tool)));
525 self
526 }
527 pub fn with_sub_agents(mut self, mut sub_agents: SubAgents) -> Self {
543 if sub_agents.limits.wall_clock.is_none() {
544 sub_agents.limits.wall_clock = Some(std::time::Duration::from_secs(600));
545 }
546 self.sub_agents = Some(sub_agents);
547 self.sub_agent_adaptive_policy = None;
548 self
549 }
550 pub fn with_sub_agents_policy(
557 mut self,
558 mut sub_agents: SubAgents,
559 adaptive_policy: AdaptiveLoopPolicy,
560 ) -> Self {
561 if sub_agents.limits.wall_clock.is_none() {
562 sub_agents.limits.wall_clock = Some(std::time::Duration::from_secs(600));
563 }
564 self.sub_agents = Some(sub_agents);
565 self.sub_agent_adaptive_policy = Some(adaptive_policy);
566 self
567 }
568 pub fn tools<I, S>(mut self, subset: I) -> Self
571 where
572 I: IntoIterator<Item = S>,
573 S: Into<String>,
574 {
575 self.spec.tools = Some(subset.into_iter().map(Into::into).collect());
576 self
577 }
578 pub fn with_cognition(mut self, yes: bool) -> Self {
582 self.cognition = yes;
583 self
584 }
585 pub fn with_sandbox(mut self, sandbox: Sandbox) -> Self {
591 self.envelope.sandbox = Some(sandbox);
592 self
593 }
594 pub fn storage(mut self, storage: Storage) -> Self {
598 self.storage = Some(storage);
599 self
600 }
601 pub fn add_context(
604 mut self,
605 id: impl Into<String>,
606 title: impl Into<String>,
607 body: impl Into<String>,
608 ) -> Self {
609 self.spec.context.push(ContextBlock::new(id, title, body));
610 self
611 }
612 pub fn groups<I>(mut self, groups: I) -> Self
619 where
620 I: IntoIterator<Item = flux_evidence::ToolGroup>,
621 {
622 self.spec.groups = groups.into_iter().collect();
623 self
624 }
625 pub fn ambient_signals<I, S>(mut self, signals: I) -> Self
630 where
631 I: IntoIterator<Item = S>,
632 S: Into<String>,
633 {
634 self.spec.ambient_signals = signals.into_iter().map(Into::into).collect();
635 self
636 }
637 pub fn with_compaction(mut self, threshold_chars: usize) -> Self {
642 self.spec.compact_threshold_chars = threshold_chars;
643 self
644 }
645 pub fn context_budget(mut self, bytes: usize) -> Self {
649 self.spec.context_budget = bytes;
650 self
651 }
652
653 pub fn build(self, provider: Box<dyn Provider>, root: impl Into<PathBuf>) -> Result<Client> {
657 let root = root.into();
658 let provider: Arc<dyn Provider> = Arc::from(provider);
659 let sandbox = self.envelope.resolve_sandbox();
663 let system = Arc::new(System::new(Workspace::new(root.clone())?).with_sandbox(sandbox));
664 let mut registry = ToolRegistry::new();
665 flux_tools::try_register_builtins(&mut registry)?;
666 if self.cognition {
667 CognitionPack::new(provider.clone(), self.spec.model.clone())
668 .with_reasoning(self.spec.thinking, self.spec.effort)
669 .try_register_from("sdk ClientBuilder cognition pack", &mut registry)?;
670 }
671 let base_names: std::collections::HashSet<String> = registry.names().into_iter().collect();
675 for (source, tool) in self.ops {
678 registry.try_register_from(source, tool)?;
679 }
680 for pack in self.packs {
681 pack(&mut registry)?;
682 }
683 if self.sub_agents.is_some() {
688 registry.try_register_from(
689 "sdk ClientBuilder sub-agent task operation",
690 Arc::new(TaskTool),
691 )?;
692 }
693 let custom_names: Vec<String> = registry
694 .names()
695 .into_iter()
696 .filter(|n| !base_names.contains(n))
697 .collect();
698 let approver = self.envelope.resolve_approver();
699
700 let (events, flow) = self.storage.unwrap_or_default().resolve()?;
701
702 let mut spec = self.spec;
707 spec.permissions
708 .allow
709 .extend(self.envelope.allow.iter().cloned());
710 spec.permissions
711 .deny
712 .extend(self.envelope.deny.iter().cloned());
713 if let Some(tools) = spec.tools.as_mut() {
716 for name in custom_names {
717 if !tools.contains(&name) {
718 tools.push(name);
719 }
720 }
721 }
722 spec.cwd = root;
723 let authorization = self.envelope.authorization.clone();
724 let mut environment = ExecutionEnvironment::new(
725 system.clone(),
726 registry,
727 PermissionManager::new(),
728 approver,
729 authorization.clone(),
730 )
731 .with_redactor(self.envelope.redactor.clone());
732 if let Some(sub_agents) = self.sub_agents {
736 let sub_agents = sub_agents
737 .with_reasoning(spec.thinking, spec.effort)
738 .with_authorization_cell(authorization.policy().clone(), authorization.identity());
739 let spawner = match self.sub_agent_adaptive_policy {
740 Some(policy) => {
741 sub_agents.into_spawner_with_adaptive_policy(system.clone(), policy)
742 }
743 None => sub_agents.into_spawner(system.clone()),
744 };
745 environment = environment.with_spawner(spawner);
746 }
747 let model = spec.model.clone();
748 let engine = spec.assemble_in(provider, environment, events, flow)?;
749 Ok(Client {
750 engine: Arc::new(engine),
751 model,
752 default_session: std::sync::Mutex::new(None),
756 turn_guard: Arc::new(tokio::sync::Mutex::new(())),
757 })
758 }
759}
760
761pub struct Client {
767 engine: Arc<FlowEngine>,
768 model: String,
769 default_session: std::sync::Mutex<Option<String>>,
772 turn_guard: Arc<tokio::sync::Mutex<()>>,
776}
777
778impl Client {
779 pub fn builder() -> ClientBuilder {
781 ClientBuilder::default()
782 }
783
784 fn default_id(&self) -> Result<String> {
787 let mut slot = self.default_session.lock().unwrap();
788 if let Some(id) = slot.as_ref() {
789 return Ok(id.clone());
790 }
791 let id = self.engine.events.create_session(&self.model)?;
792 *slot = Some(id.clone());
793 Ok(id)
794 }
795
796 pub fn session_id(&self) -> Result<String> {
799 self.default_id()
800 }
801
802 pub async fn run(&self, input: &str) -> Result<TurnOutput> {
805 let id = self.default_id()?;
806 self.session(id).send(input).await
807 }
808
809 pub fn default_session(&self) -> Result<Session> {
811 Ok(self.session(self.default_id()?))
812 }
813
814 pub fn create_session(&self) -> Result<Session> {
816 let id = self.engine.events.create_session(&self.model)?;
817 Ok(self.session(id))
818 }
819
820 pub fn open_session(&self, id: &str) -> Result<Session> {
824 self.engine.events.info(id)?;
825 Ok(self.session(id.to_string()))
826 }
827
828 pub fn latest_session(&self) -> Result<Option<Session>> {
834 Ok(self
835 .engine
836 .events
837 .latest_session()?
838 .map(|id| self.session(id)))
839 }
840
841 pub fn event_store(&self) -> Arc<EventStore> {
844 self.engine.events.clone()
845 }
846
847 pub fn engine(&self) -> &Arc<FlowEngine> {
851 &self.engine
852 }
853
854 fn session(&self, id: String) -> Session {
855 Session {
856 engine: self.engine.clone(),
857 id,
858 turn_guard: self.turn_guard.clone(),
859 }
860 }
861}
862
863#[cfg(test)]
864mod tests {
865 use super::*;
866 use async_trait::async_trait;
867 use flux_core::{Chunk, ContentBlock, StopReason, Usage};
868 use flux_provider::{ChunkStream, Request};
869 use std::sync::Mutex;
870
871 fn parse_role(content: &str, name_fallback: &str) -> crate::subagents::Role {
872 crate::subagents::try_parse_role(content, name_fallback).unwrap()
873 }
874
875 fn request_has_tool(request: &Request, name: &str) -> bool {
876 request.tools.iter().any(|tool| tool.name == name)
877 }
878
879 fn intent_chunks(intent: &str, families: &[&str]) -> Vec<Chunk> {
880 vec![
881 Chunk::Block(ContentBlock::ToolUse {
882 id: "intent".into(),
883 name: "declare_intent".into(),
884 input: serde_json::json!({
885 "intent": intent,
886 "capability_families": families,
887 }),
888 }),
889 Chunk::Done {
890 stop_reason: Some(StopReason::ToolUse),
891 },
892 ]
893 }
894
895 fn native_call(id: &str, name: &str, input: serde_json::Value) -> Vec<Chunk> {
896 vec![
897 Chunk::Block(ContentBlock::ToolUse {
898 id: id.into(),
899 name: name.into(),
900 input,
901 }),
902 Chunk::Done {
903 stop_reason: Some(StopReason::ToolUse),
904 },
905 ]
906 }
907
908 struct OneShotMock {
909 chunks: Mutex<Option<Vec<Chunk>>>,
910 }
911 #[async_trait]
912 impl Provider for OneShotMock {
913 fn name(&self) -> &str {
914 "mock"
915 }
916 async fn stream(&self, req: Request) -> Result<ChunkStream> {
917 let chunks = if request_has_tool(&req, "declare_intent") {
918 intent_chunks("answer the user", &[])
919 } else {
920 self.chunks.lock().unwrap().take().unwrap_or_default()
921 };
922 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
923 }
924 }
925
926 #[tokio::test]
927 async fn client_runs_a_text_turn() {
928 let dir = std::env::temp_dir().join(format!("flux-sdk-test-{}", std::process::id()));
929 std::fs::create_dir_all(&dir).unwrap();
930 let provider = Box::new(OneShotMock {
933 chunks: Mutex::new(Some(vec![
934 Chunk::TextDelta("hello from sdk".into()),
935 Chunk::Block(ContentBlock::Text {
936 text: "hello from sdk".into(),
937 }),
938 Chunk::Usage(Usage {
939 input_tokens: 64,
940 output_tokens: 8,
941 cache_read_input_tokens: 16,
942 ..Default::default()
943 }),
944 Chunk::Done {
945 stop_reason: Some(StopReason::EndTurn),
946 },
947 ])),
948 });
949 let client = Client::builder()
950 .model("mock")
951 .build(provider, &dir)
952 .unwrap();
953 let out = client.run("hi").await.unwrap();
954 assert_eq!(out.text, "hello from sdk");
955 assert!(out.tool_calls.is_empty());
956 let usage = out
959 .usage
960 .expect("usage surfaced through the FlowEngine loop");
961 assert_eq!(usage.input_tokens, 64);
962 assert_eq!(usage.output_tokens, 8);
963 assert_eq!(usage.cache_read_input_tokens, 16);
964 std::fs::remove_dir_all(&dir).ok();
965 }
966
967 struct SystemCaptureMock {
970 systems: Arc<Mutex<Vec<String>>>,
971 }
972 #[async_trait]
973 impl Provider for SystemCaptureMock {
974 fn name(&self) -> &str {
975 "mock"
976 }
977 async fn stream(&self, req: Request) -> Result<ChunkStream> {
978 let mut sys = String::new();
979 for seg in &req.system_segments {
980 sys.push_str(&seg.text);
981 sys.push('\n');
982 }
983 if let Some(s) = &req.system {
984 sys.push_str(s);
985 }
986 self.systems.lock().unwrap().push(sys);
987 let chunks = if request_has_tool(&req, "declare_intent") {
988 intent_chunks("answer the user", &[])
989 } else {
990 vec![
991 Chunk::Block(ContentBlock::Text { text: "ok".into() }),
992 Chunk::Done {
993 stop_reason: Some(StopReason::EndTurn),
994 },
995 ]
996 };
997 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
998 }
999 }
1000
1001 #[tokio::test]
1004 async fn sdk_skills_require_an_explicit_agent_spec() {
1005 let dir = std::env::temp_dir().join(format!("flux-sdk-skills-{}", std::process::id()));
1006 let skills = dir.join(".flux").join("skills");
1007 std::fs::create_dir_all(&skills).unwrap();
1008 std::fs::write(
1009 skills.join("greeting.md"),
1010 "---\nname: greeting\ndescription: how to greet\ntriggers: [zorblefrazz]\n---\nAlways greet with ahoy.",
1011 )
1012 .unwrap();
1013
1014 let systems = Arc::new(Mutex::new(Vec::new()));
1015 let provider = Box::new(SystemCaptureMock {
1016 systems: systems.clone(),
1017 });
1018 let client = Client::builder()
1019 .model("mock")
1020 .build(provider, &dir)
1021 .unwrap();
1022 client.run("please zorblefrazz me").await.unwrap();
1023
1024 let sys = systems.lock().unwrap().join("\n---\n");
1025 assert!(
1026 !sys.contains("Always greet with ahoy."),
1027 "workspace discovery must not activate a skill implicitly; got:\n{sys}"
1028 );
1029
1030 let systems = Arc::new(Mutex::new(Vec::new()));
1031 let provider = Box::new(SystemCaptureMock {
1032 systems: systems.clone(),
1033 });
1034 let spec = AgentSpec {
1035 cwd: dir.clone(),
1036 ..AgentSpec::new("mock")
1037 }
1038 .try_with_default_skills()
1039 .unwrap();
1040 let client = ClientBuilder::from_spec(spec)
1041 .build(provider, &dir)
1042 .unwrap();
1043 client.run("an unrelated turn").await.unwrap();
1044 let sys = systems.lock().unwrap().join("\n---\n");
1045 assert!(
1046 sys.contains("<skill name=\"greeting\">") && sys.contains("Always greet with ahoy."),
1047 "an explicitly populated AgentSpec injects the skill; got:\n{sys}"
1048 );
1049 std::fs::remove_dir_all(&dir).ok();
1050 }
1051
1052 struct PlanThenProseMock {
1056 calls: std::sync::atomic::AtomicUsize,
1057 }
1058 #[async_trait]
1059 impl Provider for PlanThenProseMock {
1060 fn name(&self) -> &str {
1061 "mock"
1062 }
1063 async fn stream(&self, req: Request) -> Result<ChunkStream> {
1064 if request_has_tool(&req, "declare_intent") {
1065 return Ok(Box::pin(futures::stream::iter(
1066 intent_chunks("write a file", &["workspace.write"])
1067 .into_iter()
1068 .map(Ok),
1069 )));
1070 }
1071 let n = self
1072 .calls
1073 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1074 let chunks = if n == 0 {
1075 native_call(
1076 "write-1",
1077 "write",
1078 serde_json::json!({
1079 "path": "sdk-plan.txt",
1080 "content": "from the sdk action batch\n"
1081 }),
1082 )
1083 } else if n == 1 {
1084 native_call(
1085 "finalize-1",
1086 "finalize_plan",
1087 serde_json::json!({
1088 "instructions": "Report whether the file was written."
1089 }),
1090 )
1091 } else {
1092 vec![
1093 Chunk::Block(ContentBlock::Text {
1094 text: "Wrote the file.".into(),
1095 }),
1096 Chunk::Done {
1097 stop_reason: Some(StopReason::EndTurn),
1098 },
1099 ]
1100 };
1101 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
1102 }
1103 }
1104
1105 struct ProseMock {
1108 text: &'static str,
1109 }
1110 #[async_trait]
1111 impl Provider for ProseMock {
1112 fn name(&self) -> &str {
1113 "mock"
1114 }
1115 async fn stream(&self, req: Request) -> Result<ChunkStream> {
1116 let chunks = if request_has_tool(&req, "declare_intent") {
1117 intent_chunks("answer the user", &[])
1118 } else {
1119 vec![
1120 Chunk::TextDelta(self.text.into()),
1121 Chunk::Block(ContentBlock::Text {
1122 text: self.text.into(),
1123 }),
1124 Chunk::Done {
1125 stop_reason: Some(StopReason::EndTurn),
1126 },
1127 ]
1128 };
1129 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
1130 }
1131 }
1132
1133 #[tokio::test]
1136 async fn storage_dir_persists_and_resumes_a_session() {
1137 let dir = std::env::temp_dir().join(format!("flux-sdk-store-{}", std::process::id()));
1138 std::fs::remove_dir_all(&dir).ok();
1139 std::fs::create_dir_all(&dir).unwrap();
1140 let store_dir = dir.join("state");
1141
1142 let client = Client::builder()
1143 .model("mock")
1144 .storage(Storage::dir(&store_dir))
1145 .build(Box::new(ProseMock { text: "first" }), &dir)
1146 .unwrap();
1147 let out = client.run("hello").await.unwrap();
1148 assert_eq!(out.text, "first");
1149 let id = client.session_id().unwrap();
1150 drop(client);
1151
1152 let client = Client::builder()
1154 .model("mock")
1155 .storage(Storage::dir(&store_dir))
1156 .build(Box::new(ProseMock { text: "second" }), &dir)
1157 .unwrap();
1158 let session = client.open_session(&id).unwrap();
1159 let history = session.history().unwrap();
1160 assert!(
1161 history.len() >= 2,
1162 "expected the prior turn's user+assistant messages, got {}",
1163 history.len()
1164 );
1165 let out = session.send("again").await.unwrap();
1166 assert_eq!(out.text, "second");
1167 assert!(session.history().unwrap().len() > history.len());
1168 std::fs::remove_dir_all(&dir).ok();
1169 }
1170
1171 struct PlanOpMock {
1173 op: &'static str,
1174 calls: std::sync::atomic::AtomicUsize,
1175 }
1176 #[async_trait]
1177 impl Provider for PlanOpMock {
1178 fn name(&self) -> &str {
1179 "mock"
1180 }
1181 async fn stream(&self, req: Request) -> Result<ChunkStream> {
1182 if request_has_tool(&req, "declare_intent") {
1183 return Ok(Box::pin(futures::stream::iter(
1184 intent_chunks("call the requested operation", &["core"])
1185 .into_iter()
1186 .map(Ok),
1187 )));
1188 }
1189 let n = self
1190 .calls
1191 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1192 let chunks = if n == 0 {
1193 let input = if self.op == "park" {
1194 serde_json::json!({})
1195 } else {
1196 serde_json::json!({"name": "flux"})
1197 };
1198 native_call("op-1", self.op, input)
1199 } else {
1200 vec![
1201 Chunk::Block(ContentBlock::Text {
1202 text: "done".into(),
1203 }),
1204 Chunk::Done {
1205 stop_reason: Some(StopReason::EndTurn),
1206 },
1207 ]
1208 };
1209 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
1210 }
1211 }
1212
1213 fn greet_tool(hits: Arc<std::sync::atomic::AtomicUsize>) -> Arc<dyn flux_runtime::Tool> {
1214 flux_runtime::tool_fn(
1215 flux_spec::ToolSpec::read_only(
1216 "greet",
1217 "Greets by name",
1218 serde_json::json!({
1219 "type": "object",
1220 "properties": { "name": { "type": "string" } },
1221 "required": ["name"]
1222 }),
1223 ),
1224 move |input| {
1225 let hits = hits.clone();
1226 async move {
1227 hits.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1228 Ok(serde_json::json!(format!(
1229 "hello {}",
1230 input["name"].as_str().unwrap_or("?")
1231 )))
1232 }
1233 },
1234 )
1235 }
1236
1237 fn guarded_read_tool(hits: Arc<std::sync::atomic::AtomicUsize>) -> Arc<dyn flux_runtime::Tool> {
1238 flux_runtime::tool_fn(
1239 flux_spec::ToolSpec::read_only(
1240 "guarded_read",
1241 "A filesystem-scoped read used to exercise authorization",
1242 serde_json::json!({
1243 "type": "object",
1244 "properties": { "name": { "type": "string" } }
1245 }),
1246 )
1247 .with_access(vec![flux_spec::AccessKind::Filesystem]),
1248 move |_input| {
1249 let hits = hits.clone();
1250 async move {
1251 hits.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1252 Ok(serde_json::json!("ran"))
1253 }
1254 },
1255 )
1256 }
1257
1258 #[tokio::test]
1261 async fn a_registered_custom_tool_dispatches_through_a_planned_turn() {
1262 let dir = std::env::temp_dir().join(format!("flux-sdk-fntool-{}", std::process::id()));
1263 std::fs::create_dir_all(&dir).unwrap();
1264 let hits = Arc::new(std::sync::atomic::AtomicUsize::new(0));
1265 let client = Client::builder()
1266 .model("mock")
1267 .auto_approve(true)
1268 .register_op(greet_tool(hits.clone()))
1269 .build(
1270 Box::new(PlanOpMock {
1271 op: "greet",
1272 calls: std::sync::atomic::AtomicUsize::new(0),
1273 }),
1274 &dir,
1275 )
1276 .unwrap();
1277 let out = client.run("greet flux").await.unwrap();
1278 assert_eq!(hits.load(std::sync::atomic::Ordering::Relaxed), 1);
1279 assert_eq!(out.tool_calls, vec!["greet"]);
1280 std::fs::remove_dir_all(&dir).ok();
1281 }
1282
1283 #[tokio::test]
1285 async fn client_auto_approval_cannot_widen_authorization() {
1286 let dir =
1287 std::env::temp_dir().join(format!("flux-sdk-policy-floor-{}", std::process::id()));
1288 std::fs::create_dir_all(&dir).unwrap();
1289 let hits = Arc::new(std::sync::atomic::AtomicUsize::new(0));
1290 let mut policy = flux_policy::default_local_grants();
1291 policy.grants.retain(|grant| {
1292 !grant
1293 .actions
1294 .iter()
1295 .any(|action| action.0 == "workspace.read")
1296 });
1297 let (caller, trust) = flux_policy::local_identity("sdk-test");
1298 let client = Client::builder()
1299 .model("mock")
1300 .auto_approve(true)
1301 .with_authorization(policy, caller, trust)
1302 .register_op(guarded_read_tool(hits.clone()))
1303 .build(
1304 Box::new(PlanOpMock {
1305 op: "guarded_read",
1306 calls: std::sync::atomic::AtomicUsize::new(0),
1307 }),
1308 &dir,
1309 )
1310 .unwrap();
1311
1312 let _ = client.run("greet flux").await;
1313 assert_eq!(
1314 hits.load(std::sync::atomic::Ordering::SeqCst),
1315 0,
1316 "auto approval must not execute a policy-denied operation"
1317 );
1318 std::fs::remove_dir_all(&dir).ok();
1319 }
1320
1321 #[tokio::test]
1325 async fn an_injected_approver_gates_a_registered_custom_tool() {
1326 struct DenyGreet;
1327 #[async_trait]
1328 impl flux_runtime::Approver for DenyGreet {
1329 async fn request(
1330 &self,
1331 tool: &str,
1332 _subjects: &[String],
1333 _intents: &flux_spec::IntentSet,
1334 ) -> flux_runtime::ApprovalChoice {
1335 if tool == "greet" {
1336 flux_runtime::ApprovalChoice::Deny
1337 } else {
1338 flux_runtime::ApprovalChoice::Allow
1339 }
1340 }
1341 }
1342
1343 let dir = std::env::temp_dir().join(format!("flux-sdk-fngate-{}", std::process::id()));
1344 std::fs::create_dir_all(&dir).unwrap();
1345 let hits = Arc::new(std::sync::atomic::AtomicUsize::new(0));
1346 let client = Client::builder()
1347 .model("mock")
1348 .approver(Arc::new(DenyGreet))
1349 .register_op(greet_tool(hits.clone()))
1350 .build(
1351 Box::new(PlanOpMock {
1352 op: "greet",
1353 calls: std::sync::atomic::AtomicUsize::new(0),
1354 }),
1355 &dir,
1356 )
1357 .unwrap();
1358 let _ = client.run("greet flux").await;
1359 assert_eq!(
1360 hits.load(std::sync::atomic::Ordering::Relaxed),
1361 0,
1362 "the injected approver must gate the registered tool"
1363 );
1364 std::fs::remove_dir_all(&dir).ok();
1365 }
1366
1367 #[tokio::test]
1370 async fn tools_subset_removes_ops_from_the_registry() {
1371 let dir = std::env::temp_dir().join(format!("flux-sdk-subset-{}", std::process::id()));
1372 std::fs::create_dir_all(&dir).unwrap();
1373 let client = Client::builder()
1374 .model("mock")
1375 .auto_approve(true)
1376 .tools(["read"])
1377 .build(
1378 Box::new(PlanThenProseMock {
1379 calls: std::sync::atomic::AtomicUsize::new(0),
1380 }),
1381 &dir,
1382 )
1383 .unwrap();
1384 let _ = client.run("write a file").await;
1387 assert!(
1388 !dir.join("sdk-plan.txt").exists(),
1389 "an out-of-subset op must not execute"
1390 );
1391 std::fs::remove_dir_all(&dir).ok();
1392 }
1393
1394 #[tokio::test]
1397 async fn send_with_streams_deltas_and_tool_results_to_a_consumer_sink() {
1398 #[derive(Default)]
1399 struct Recording {
1400 deltas: Vec<String>,
1401 tool_results: Vec<String>,
1402 }
1403 impl AgentSink for Recording {
1404 fn text_delta(&mut self, t: &str) {
1405 self.deltas.push(t.to_string());
1406 }
1407 fn tool_result(&mut self, name: &str, _result: &flux_runtime::ToolResult) {
1408 self.tool_results.push(name.to_string());
1409 }
1410 }
1411
1412 let dir = std::env::temp_dir().join(format!("flux-sdk-sendwith-{}", std::process::id()));
1413 std::fs::create_dir_all(&dir).unwrap();
1414 let client = Client::builder()
1415 .model("mock")
1416 .auto_approve(true)
1417 .build(
1418 Box::new(PlanThenProseMock {
1419 calls: std::sync::atomic::AtomicUsize::new(0),
1420 }),
1421 &dir,
1422 )
1423 .unwrap();
1424 let session = client.default_session().unwrap();
1425 let mut sink = Recording::default();
1426 let out = session
1427 .send_with("write a file", &mut sink, &CancellationToken::new())
1428 .await
1429 .unwrap();
1430 assert_eq!(out.text, "Wrote the file.");
1431 assert!(
1435 sink.tool_results.contains(&"write".to_string()),
1436 "the consumer sink must receive tool_result events, got {:?}",
1437 sink.tool_results
1438 );
1439 std::fs::remove_dir_all(&dir).ok();
1440 }
1441
1442 struct TwoDeltaMock;
1444 #[async_trait]
1445 impl Provider for TwoDeltaMock {
1446 fn name(&self) -> &str {
1447 "mock"
1448 }
1449 async fn stream(&self, req: Request) -> Result<ChunkStream> {
1450 let chunks = if request_has_tool(&req, "declare_intent") {
1451 intent_chunks("answer the user", &[])
1452 } else {
1453 vec![
1454 Chunk::TextDelta("first ".into()),
1455 Chunk::TextDelta("second".into()),
1456 Chunk::Block(ContentBlock::Text {
1457 text: "first second".into(),
1458 }),
1459 Chunk::Done {
1460 stop_reason: Some(StopReason::EndTurn),
1461 },
1462 ]
1463 };
1464 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
1465 }
1466 }
1467
1468 #[tokio::test]
1473 async fn stream_yields_events_and_finish_collects() {
1474 let dir = std::env::temp_dir().join(format!("flux-sdk-stream-{}", std::process::id()));
1475 std::fs::create_dir_all(&dir).unwrap();
1476 let client = Client::builder()
1477 .model("mock")
1478 .build(Box::new(TwoDeltaMock), &dir)
1479 .unwrap();
1480 let mut stream = client.default_session().unwrap().stream("hi");
1481
1482 let mut deltas = String::new();
1483 let mut saw_turn_end = false;
1484 while let Some(event) = stream.next().await {
1485 match event {
1486 AgentEvent::TextDelta(t) => deltas.push_str(&t),
1487 AgentEvent::TurnEnd { .. } => saw_turn_end = true,
1488 _ => {}
1489 }
1490 }
1491 assert_eq!(deltas, "first second");
1492 assert!(saw_turn_end, "the stream must emit a TurnEnd event");
1493 let out = stream.finish().await.unwrap();
1494 assert_eq!(out.text, "first second");
1495 std::fs::remove_dir_all(&dir).ok();
1496 }
1497
1498 #[tokio::test]
1501 async fn cancelling_a_streamed_turn_keeps_the_session_shape_valid() {
1502 let dir = std::env::temp_dir().join(format!("flux-sdk-cancel-{}", std::process::id()));
1503 std::fs::create_dir_all(&dir).unwrap();
1504 let parked = flux_runtime::tool_fn(
1506 flux_spec::ToolSpec::read_only(
1507 "park",
1508 "Blocks until cancelled",
1509 serde_json::json!({ "type": "object", "properties": {} }),
1510 ),
1511 |_input| async {
1512 tokio::time::sleep(std::time::Duration::from_secs(300)).await;
1513 Ok(serde_json::json!("unreachable"))
1514 },
1515 );
1516 let client = Client::builder()
1517 .model("mock")
1518 .auto_approve(true)
1519 .register_op(parked)
1520 .build(
1521 Box::new(PlanOpMock {
1522 op: "park",
1523 calls: std::sync::atomic::AtomicUsize::new(0),
1524 }),
1525 &dir,
1526 )
1527 .unwrap();
1528 let session = client.default_session().unwrap();
1529 let mut stream = session.stream("park it");
1530 loop {
1532 match stream.next().await {
1533 Some(AgentEvent::ToolCall { name, .. }) if name == "park" => break,
1534 Some(_) => continue,
1535 None => panic!("stream ended before the tool call"),
1536 }
1537 }
1538 stream.cancel();
1539 let _ = stream.finish().await;
1540
1541 let history = session.history().unwrap();
1542 assert!(!history.is_empty());
1543 for pair in history.windows(2) {
1544 assert_ne!(
1545 pair[0].role, pair[1].role,
1546 "roles must alternate after a cancelled turn"
1547 );
1548 }
1549 assert!(
1550 matches!(history.last().unwrap().role, flux_core::Role::Assistant),
1551 "a cancelled turn must still persist exactly one closing assistant message"
1552 );
1553 std::fs::remove_dir_all(&dir).ok();
1554 }
1555
1556 #[tokio::test]
1558 async fn open_session_unknown_id_errors() {
1559 let dir = std::env::temp_dir().join(format!("flux-sdk-open-{}", std::process::id()));
1560 std::fs::create_dir_all(&dir).unwrap();
1561 let client = Client::builder()
1562 .model("mock")
1563 .build(Box::new(ProseMock { text: "x" }), &dir)
1564 .unwrap();
1565 assert!(client.open_session("no-such-session").is_err());
1566 std::fs::remove_dir_all(&dir).ok();
1567 }
1568
1569 struct SlowRecordingMock {
1572 calls: Arc<Mutex<Vec<(std::time::Instant, std::time::Instant)>>>,
1573 }
1574 #[async_trait]
1575 impl Provider for SlowRecordingMock {
1576 fn name(&self) -> &str {
1577 "mock"
1578 }
1579 async fn stream(&self, req: Request) -> Result<ChunkStream> {
1580 let start = std::time::Instant::now();
1581 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1582 self.calls
1583 .lock()
1584 .unwrap()
1585 .push((start, std::time::Instant::now()));
1586 let chunks = if request_has_tool(&req, "declare_intent") {
1587 intent_chunks("answer the user", &[])
1588 } else {
1589 vec![
1590 Chunk::Block(ContentBlock::Text { text: "ok".into() }),
1591 Chunk::Done {
1592 stop_reason: Some(StopReason::EndTurn),
1593 },
1594 ]
1595 };
1596 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
1597 }
1598 }
1599
1600 #[tokio::test]
1603 async fn concurrent_sends_serialize_on_the_turn_guard() {
1604 let dir = std::env::temp_dir().join(format!("flux-sdk-guard-{}", std::process::id()));
1605 std::fs::create_dir_all(&dir).unwrap();
1606 let calls = Arc::new(Mutex::new(Vec::new()));
1607 let client = Client::builder()
1608 .model("mock")
1609 .build(
1610 Box::new(SlowRecordingMock {
1611 calls: calls.clone(),
1612 }),
1613 &dir,
1614 )
1615 .unwrap();
1616 let a = client.create_session().unwrap();
1617 let b = client.create_session().unwrap();
1618 let (ra, rb) = tokio::join!(a.send("one"), b.send("two"));
1619 ra.unwrap();
1620 rb.unwrap();
1621
1622 let mut intervals = calls.lock().unwrap().clone();
1623 intervals.sort_by_key(|(s, _)| *s);
1624 assert_eq!(intervals.len(), 4, "two adaptive stages per chat turn");
1625 assert!(
1626 intervals.windows(2).all(|pair| pair[1].0 >= pair[0].1),
1627 "provider calls overlapped: the turn guard failed to serialize the turns"
1628 );
1629 std::fs::remove_dir_all(&dir).ok();
1630 }
1631
1632 #[tokio::test]
1633 async fn client_runs_an_action_batch_then_answers() {
1634 let dir = std::env::temp_dir().join(format!("flux-sdk-plan-{}", std::process::id()));
1635 std::fs::create_dir_all(&dir).unwrap();
1636 let provider = Box::new(PlanThenProseMock {
1637 calls: std::sync::atomic::AtomicUsize::new(0),
1638 });
1639 let client = Client::builder()
1640 .model("mock")
1641 .auto_approve(true) .build(provider, &dir)
1643 .unwrap();
1644 let out = client.run("write a file").await.unwrap();
1645 assert_eq!(out.text, "Wrote the file.");
1646 assert_eq!(out.tool_calls, vec!["write"]);
1648 assert!(dir.join("sdk-plan.txt").exists(), "the batch's write ran");
1650 std::fs::remove_dir_all(&dir).ok();
1651 }
1652
1653 #[tokio::test]
1656 async fn tools_subset_preserves_registered_custom_ops() {
1657 let dir = std::env::temp_dir().join(format!("flux-sdk-subkeep-{}", std::process::id()));
1658 std::fs::create_dir_all(&dir).unwrap();
1659 let hits = Arc::new(std::sync::atomic::AtomicUsize::new(0));
1660 let client = Client::builder()
1661 .model("mock")
1662 .auto_approve(true)
1663 .tools(["read"]) .register_op(greet_tool(hits.clone()))
1665 .build(
1666 Box::new(PlanOpMock {
1667 op: "greet",
1668 calls: std::sync::atomic::AtomicUsize::new(0),
1669 }),
1670 &dir,
1671 )
1672 .unwrap();
1673 let out = client.run("greet flux").await.unwrap();
1674 assert_eq!(
1675 hits.load(std::sync::atomic::Ordering::Relaxed),
1676 1,
1677 "the registered custom op must survive the tools() subset"
1678 );
1679 assert_eq!(out.tool_calls, vec!["greet"]);
1680 std::fs::remove_dir_all(&dir).ok();
1681 }
1682
1683 #[tokio::test]
1687 async fn lazy_default_session_does_not_shadow_the_prior_conversation() {
1688 let dir = std::env::temp_dir().join(format!("flux-sdk-lazy-{}", std::process::id()));
1689 std::fs::remove_dir_all(&dir).ok();
1690 std::fs::create_dir_all(&dir).unwrap();
1691 let store = dir.join("state");
1692
1693 let real_id = {
1695 let client = Client::builder()
1696 .model("mock")
1697 .storage(Storage::dir(&store))
1698 .build(Box::new(ProseMock { text: "hi" }), &dir)
1699 .unwrap();
1700 client.run("remember this").await.unwrap();
1701 let id = client.session_id().unwrap();
1702 drop(client);
1703 id
1704 };
1705
1706 let client = Client::builder()
1709 .model("mock")
1710 .storage(Storage::dir(&store))
1711 .build(Box::new(ProseMock { text: "hi" }), &dir)
1712 .unwrap();
1713 let latest = client
1714 .latest_session()
1715 .unwrap()
1716 .expect("a prior session exists");
1717 assert_eq!(
1718 latest.id(),
1719 real_id,
1720 "latest_session must return the real prior conversation, not a fresh empty default"
1721 );
1722 std::fs::remove_dir_all(&dir).ok();
1723 }
1724
1725 #[tokio::test]
1729 async fn dropping_a_turn_stream_cancels_the_turn() {
1730 let dir = std::env::temp_dir().join(format!("flux-sdk-dropcancel-{}", std::process::id()));
1731 std::fs::create_dir_all(&dir).unwrap();
1732 let parked = flux_runtime::tool_fn(
1733 flux_spec::ToolSpec::read_only(
1734 "park",
1735 "Blocks until cancelled",
1736 serde_json::json!({ "type": "object", "properties": {} }),
1737 ),
1738 |_input| async {
1739 tokio::time::sleep(std::time::Duration::from_secs(300)).await;
1740 Ok(serde_json::json!("unreachable"))
1741 },
1742 );
1743 let client = Client::builder()
1744 .model("mock")
1745 .auto_approve(true)
1746 .register_op(parked)
1747 .build(
1748 Box::new(PlanOpMock {
1749 op: "park",
1750 calls: std::sync::atomic::AtomicUsize::new(0),
1751 }),
1752 &dir,
1753 )
1754 .unwrap();
1755 {
1756 let session = client.default_session().unwrap();
1757 let mut stream = session.stream("park it");
1758 loop {
1760 match stream.next().await {
1761 Some(AgentEvent::ToolCall { name, .. }) if name == "park" => break,
1762 Some(_) => continue,
1763 None => panic!("stream ended before the tool call"),
1764 }
1765 }
1766 drop(stream);
1767 }
1768 let out = tokio::time::timeout(
1771 std::time::Duration::from_secs(20),
1772 client.default_session().unwrap().send("are you there"),
1773 )
1774 .await
1775 .expect("run after drop must not hang — the dropped stream should have cancelled the turn")
1776 .unwrap();
1777 assert_eq!(out.text, "done");
1778 std::fs::remove_dir_all(&dir).ok();
1779 }
1780
1781 struct NeverMock;
1784 #[async_trait]
1785 impl Provider for NeverMock {
1786 fn name(&self) -> &str {
1787 "mock"
1788 }
1789 async fn stream(&self, _req: Request) -> Result<ChunkStream> {
1790 panic!("a flow-driven session must not invoke a model stage");
1791 }
1792 }
1793
1794 struct EchoTool;
1797 #[async_trait]
1798 impl Tool for EchoTool {
1799 fn spec(&self) -> flux_spec::ToolSpec {
1800 flux_spec::ToolSpec::read_only(
1801 "echo",
1802 "echo text",
1803 serde_json::json!({
1804 "type": "object",
1805 "properties": { "text": { "type": "string" } },
1806 "required": ["text"]
1807 }),
1808 )
1809 }
1810 async fn execute(
1811 &self,
1812 _c: &ToolContext,
1813 params: serde_json::Value,
1814 ) -> Result<flux_runtime::ToolResult> {
1815 Ok(flux_runtime::ToolResult::ok(
1816 params
1817 .get("text")
1818 .and_then(|v| v.as_str())
1819 .unwrap_or("")
1820 .to_string(),
1821 ))
1822 }
1823 }
1824
1825 #[test]
1826 fn client_builder_reports_source_aware_custom_operation_collisions() {
1827 let error = Client::builder()
1828 .register_op_from("custom-pack:alpha", Arc::new(EchoTool))
1829 .register_op_from("custom-pack:beta", Arc::new(EchoTool))
1830 .build(Box::new(NeverMock), ".")
1831 .err()
1832 .expect("duplicate operation must fail client assembly")
1833 .to_string();
1834
1835 assert!(error.contains("duplicate operation `echo`"));
1836 assert!(error.contains("custom-pack:alpha"));
1837 assert!(error.contains("custom-pack:beta"));
1838 }
1839
1840 #[test]
1841 fn client_builder_rejects_custom_operation_shadowing_a_builtin() {
1842 let shadow = flux_runtime::tool_fn(
1843 flux_spec::ToolSpec::read_only(
1844 "read",
1845 "shadow the workspace reader",
1846 serde_json::json!({"type": "object"}),
1847 ),
1848 |_params| async { Ok(serde_json::Value::Null) },
1849 );
1850 let error = Client::builder()
1851 .register_op_from("custom-pack:shadow", shadow)
1852 .build(Box::new(NeverMock), ".")
1853 .err()
1854 .expect("custom operation must not replace a built-in")
1855 .to_string();
1856
1857 assert!(error.contains("duplicate operation `read`"), "{error}");
1858 assert!(error.contains("custom-pack:shadow"), "{error}");
1859 assert!(error.contains("flux-tools core coding pack"), "{error}");
1860 }
1861
1862 #[cfg(feature = "plugins")]
1863 #[test]
1864 fn client_builder_rejects_two_installed_plugins_with_the_same_public_operation() {
1865 let error = Client::builder()
1866 .with_plugin_tools_from(
1867 "alpha (/plugins/alpha.toml)",
1868 vec![Arc::new(EchoTool) as Arc<dyn Tool>],
1869 )
1870 .with_plugin_tools_from(
1871 "beta (/plugins/beta.toml)",
1872 vec![Arc::new(EchoTool) as Arc<dyn Tool>],
1873 )
1874 .build(Box::new(NeverMock), ".")
1875 .err()
1876 .expect("installed plugin public-name collision must fail assembly")
1877 .to_string();
1878
1879 assert!(error.contains("duplicate operation `echo`"), "{error}");
1880 assert!(
1881 error.contains("plugin:alpha (/plugins/alpha.toml)"),
1882 "{error}"
1883 );
1884 assert!(
1885 error.contains("plugin:beta (/plugins/beta.toml)"),
1886 "{error}"
1887 );
1888 }
1889
1890 fn interview_flow() -> flux_lang::ast::DraftAst {
1893 use flux_lang::ast::{Node, SymbolName};
1894 let prompt = |t: &str| Node::Call {
1895 op: "echo".into(),
1896 args: vec![Node::Lit {
1897 value: serde_json::json!(t),
1898 }],
1899 };
1900 let await_reply = |name: &str| Node::Await {
1901 binding: Some(SymbolName(name.into())),
1902 source: "user_input".into(),
1903 as_type: None,
1904 condition: None,
1905 };
1906 flux_lang::ast::DraftAst {
1907 body: vec![
1908 prompt("What is your name?"),
1909 await_reply("name"),
1910 prompt("Nice to meet you. Favorite color?"),
1911 await_reply("color"),
1912 prompt("All done — thanks!"),
1913 ],
1914 ..Default::default()
1915 }
1916 }
1917
1918 #[tokio::test]
1923 async fn start_flow_suspends_surfaces_prompt_and_send_resumes() {
1924 let dir = std::env::temp_dir().join(format!("flux-sdk-startflow-{}", std::process::id()));
1925 std::fs::create_dir_all(&dir).unwrap();
1926 let client = Client::builder()
1927 .model("mock")
1928 .auto_approve(true)
1929 .register_op(Arc::new(EchoTool))
1930 .build(Box::new(NeverMock), &dir)
1931 .unwrap();
1932 let session = client.create_session().unwrap();
1933
1934 let out = session.start_flow(&interview_flow()).await.unwrap();
1936 assert!(
1937 out.text.contains("What is your name?"),
1938 "start_flow surfaces the first authored prompt: {:?}",
1939 out.text
1940 );
1941 assert!(out.suspended, "the flow parked on its first `await`");
1942 assert!(
1943 session.suspended().unwrap(),
1944 "the session reports suspended"
1945 );
1946
1947 let out = session.send("Timo").await.unwrap();
1949 assert!(
1950 out.text.contains("Favorite color?"),
1951 "send resumes to the second authored prompt: {:?}",
1952 out.text
1953 );
1954 assert!(out.suspended, "still parked on the second `await`");
1955
1956 let out = session.send("blue").await.unwrap();
1958 assert!(
1959 out.text.contains("All done"),
1960 "the final send completes the flow: {:?}",
1961 out.text
1962 );
1963 assert!(!out.suspended, "a completed flow is no longer suspended");
1964 assert!(
1965 !session.suspended().unwrap(),
1966 "the session reports not suspended"
1967 );
1968
1969 std::fs::remove_dir_all(&dir).ok();
1970 }
1971
1972 #[tokio::test]
1976 async fn suspended_flow_survives_a_process_restart() {
1977 let dir =
1978 std::env::temp_dir().join(format!("flux-sdk-startflow-restart-{}", std::process::id()));
1979 std::fs::create_dir_all(&dir).unwrap();
1980
1981 let session_id = {
1982 let client = Client::builder()
1983 .model("mock")
1984 .auto_approve(true)
1985 .register_op(Arc::new(EchoTool))
1986 .storage(Storage::dir(&dir))
1987 .build(Box::new(NeverMock), &dir)
1988 .unwrap();
1989 let session = client.create_session().unwrap();
1990 let out = session.start_flow(&interview_flow()).await.unwrap();
1991 assert!(out.suspended, "parked on await #1 before the restart");
1992 session.id().to_string()
1993 };
1995
1996 let client = Client::builder()
1998 .model("mock")
1999 .auto_approve(true)
2000 .register_op(Arc::new(EchoTool))
2001 .storage(Storage::dir(&dir))
2002 .build(Box::new(NeverMock), &dir)
2003 .unwrap();
2004 let session = client.open_session(&session_id).unwrap();
2005 assert!(
2006 session.suspended().unwrap(),
2007 "the persisted suspension is visible after the restart"
2008 );
2009 let out = session.send("Timo").await.unwrap();
2010 assert!(
2011 out.text.contains("Favorite color?"),
2012 "the parked flow resumes across the restart: {:?}",
2013 out.text
2014 );
2015 assert!(out.suspended, "re-parked on await #2");
2016
2017 std::fs::remove_dir_all(&dir).ok();
2018 }
2019
2020 struct WorkerMock(Usage);
2023 #[async_trait]
2024 impl Provider for WorkerMock {
2025 fn name(&self) -> &str {
2026 "mock"
2027 }
2028 async fn stream(&self, req: Request) -> Result<ChunkStream> {
2029 let chunks = if request_has_tool(&req, "declare_intent") {
2030 intent_chunks("complete the delegated task", &[])
2031 } else {
2032 vec![
2033 Chunk::Block(ContentBlock::Text {
2034 text: "did the subtask".into(),
2035 }),
2036 Chunk::Usage(self.0.clone()),
2037 Chunk::Done {
2038 stop_reason: Some(StopReason::EndTurn),
2039 },
2040 ]
2041 };
2042 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
2043 }
2044 }
2045
2046 struct DelegatingMock {
2049 calls: std::sync::atomic::AtomicUsize,
2050 }
2051 #[async_trait]
2052 impl Provider for DelegatingMock {
2053 fn name(&self) -> &str {
2054 "mock"
2055 }
2056 async fn stream(&self, req: Request) -> Result<ChunkStream> {
2057 if request_has_tool(&req, "declare_intent") {
2058 return Ok(Box::pin(futures::stream::iter(
2059 intent_chunks("delegate the task", &["process"])
2060 .into_iter()
2061 .map(Ok),
2062 )));
2063 }
2064 let n = self
2065 .calls
2066 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
2067 let chunks = if n == 0 {
2068 native_call(
2069 "task-1",
2070 "task",
2071 serde_json::json!({"role": "worker", "task": "do it"}),
2072 )
2073 } else if n == 1 {
2074 native_call(
2075 "finalize-1",
2076 "finalize_plan",
2077 serde_json::json!({
2078 "instructions": "Report the delegated task's actual result."
2079 }),
2080 )
2081 } else {
2082 vec![
2083 Chunk::Block(ContentBlock::Text {
2084 text: "delegated to the worker".into(),
2085 }),
2086 Chunk::Done {
2087 stop_reason: Some(StopReason::EndTurn),
2088 },
2089 ]
2090 };
2091 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
2092 }
2093 }
2094
2095 #[tokio::test]
2100 async fn with_sub_agents_runs_a_delegated_task_and_records_child_usage() {
2101 use crate::subagents::{RoleRegistry, SubAgents};
2102
2103 let dir = std::env::temp_dir().join(format!("flux-sdk-subagents-{}", std::process::id()));
2104 std::fs::create_dir_all(&dir).unwrap();
2105
2106 let mut roles = RoleRegistry::default();
2108 roles.insert(parse_role("---\n---\nworker prompt", "worker"));
2109
2110 let child_usage = Usage {
2112 input_tokens: 1000,
2113 output_tokens: 200,
2114 ..Default::default()
2115 };
2116 let factory: crate::subagents::ProviderFactory = Arc::new({
2117 let u = child_usage.clone();
2118 move || Ok(Box::new(WorkerMock(u.clone())) as Box<dyn Provider>)
2119 });
2120 let sub_agents = SubAgents::new(roles, ToolRegistry::new(), factory, "mock", 1024);
2121
2122 let client = Client::builder()
2123 .model("mock")
2124 .auto_approve(true)
2125 .with_sub_agents(sub_agents)
2126 .build(
2127 Box::new(DelegatingMock {
2128 calls: std::sync::atomic::AtomicUsize::new(0),
2129 }),
2130 &dir,
2131 )
2132 .unwrap();
2133
2134 let out = client.run("delegate this").await.unwrap();
2135 assert!(
2136 out.tool_calls.contains(&"task".to_string()),
2137 "the adaptive turn delegated via `task`: {:?}",
2138 out.tool_calls
2139 );
2140
2141 let sid = client.session_id().unwrap();
2144 let events = client.event_store();
2145 let turns = events.turns(&sid).unwrap();
2146 let usage = turns
2147 .last()
2148 .and_then(|t| t.usage.as_ref())
2149 .expect("the parent turn's usage must be Some — the sub-agent billed tokens");
2150 assert_eq!(
2151 usage.input_tokens, 1000,
2152 "the sub-agent's input tokens landed in the session's run trace"
2153 );
2154 assert_eq!(
2155 usage.output_tokens, 200,
2156 "the sub-agent's output tokens landed in the session's run trace"
2157 );
2158
2159 std::fs::remove_dir_all(&dir).ok();
2160 }
2161
2162 struct PolicyCaptureWorker {
2163 requests: Arc<Mutex<Vec<Request>>>,
2164 }
2165
2166 #[async_trait]
2167 impl Provider for PolicyCaptureWorker {
2168 fn name(&self) -> &str {
2169 "mock"
2170 }
2171
2172 async fn stream(&self, req: Request) -> Result<ChunkStream> {
2173 self.requests.lock().unwrap().push(req.clone());
2174 let chunks = if request_has_tool(&req, "declare_intent") {
2175 intent_chunks("complete the delegated task", &[])
2176 } else {
2177 vec![
2178 Chunk::Block(ContentBlock::Text {
2179 text: "policy child done".into(),
2180 }),
2181 Chunk::Done {
2182 stop_reason: Some(StopReason::EndTurn),
2183 },
2184 ]
2185 };
2186 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
2187 }
2188 }
2189
2190 #[tokio::test]
2194 async fn with_sub_agents_policy_reaches_conversational_children() {
2195 use crate::subagents::{RoleRegistry, SubAgents};
2196
2197 let dir =
2198 std::env::temp_dir().join(format!("flux-sdk-subagent-policy-{}", std::process::id()));
2199 std::fs::create_dir_all(&dir).unwrap();
2200 let mut roles = RoleRegistry::default();
2201 roles.insert(parse_role("---\n---\nworker prompt", "worker"));
2202 let requests = Arc::new(Mutex::new(Vec::new()));
2203 let factory: crate::subagents::ProviderFactory = Arc::new({
2204 let requests = requests.clone();
2205 move || {
2206 Ok(Box::new(PolicyCaptureWorker {
2207 requests: requests.clone(),
2208 }) as Box<dyn Provider>)
2209 }
2210 });
2211 let sub_agents = SubAgents::new(roles, ToolRegistry::new(), factory, "child-default", 1024);
2212 let policy = AdaptiveLoopPolicy {
2213 max_model_calls: 2,
2214 intent: AgentStagePolicy {
2215 model: Some("intent-fast".into()),
2216 effort: Some(flux_provider::Effort::Low),
2217 max_tokens: Some(111),
2218 max_calls: Some(1),
2219 },
2220 explore: AgentStagePolicy {
2221 model: Some("explore-deep".into()),
2222 effort: Some(flux_provider::Effort::High),
2223 max_tokens: Some(222),
2224 max_calls: Some(1),
2225 },
2226 };
2227
2228 let client = Client::builder()
2229 .model("mock")
2230 .auto_approve(true)
2231 .with_sub_agents_policy(sub_agents, policy)
2232 .build(
2233 Box::new(DelegatingMock {
2234 calls: std::sync::atomic::AtomicUsize::new(0),
2235 }),
2236 &dir,
2237 )
2238 .unwrap();
2239 let out = client.run("delegate this").await.unwrap();
2240 assert!(out.tool_calls.contains(&"task".to_string()));
2241
2242 let requests = requests.lock().unwrap();
2243 assert_eq!(requests.len(), 2);
2244 assert_eq!(requests[0].trace.as_ref().unwrap().stage, "intent");
2245 assert_eq!(requests[0].model, "intent-fast");
2246 assert_eq!(requests[0].effort, Some(flux_provider::Effort::Low));
2247 assert_eq!(requests[0].max_tokens, 111);
2248 assert_eq!(requests[1].trace.as_ref().unwrap().stage, "explore");
2249 assert_eq!(requests[1].model, "explore-deep");
2250 assert_eq!(requests[1].effort, Some(flux_provider::Effort::High));
2251 assert_eq!(requests[1].max_tokens, 222);
2252 drop(requests);
2253
2254 std::fs::remove_dir_all(&dir).ok();
2255 }
2256
2257 struct NestedTaskAdapter {
2261 audit: Arc<EventStore>,
2262 child_entered: tokio::sync::mpsc::UnboundedSender<()>,
2263 child_dropped: tokio::sync::mpsc::UnboundedSender<()>,
2264 }
2265
2266 struct DropNotice(tokio::sync::mpsc::UnboundedSender<()>);
2267
2268 impl Drop for DropNotice {
2269 fn drop(&mut self) {
2270 let _ = self.0.send(());
2271 }
2272 }
2273
2274 struct HangingWorker {
2275 entered: tokio::sync::mpsc::UnboundedSender<()>,
2276 dropped: tokio::sync::mpsc::UnboundedSender<()>,
2277 }
2278
2279 #[async_trait]
2280 impl Provider for HangingWorker {
2281 fn name(&self) -> &str {
2282 "mock"
2283 }
2284
2285 async fn stream(&self, request: Request) -> Result<ChunkStream> {
2286 if request_has_tool(&request, "declare_intent") {
2287 return Ok(Box::pin(futures::stream::iter(
2288 intent_chunks("wait for the parent request", &[])
2289 .into_iter()
2290 .map(Ok),
2291 )));
2292 }
2293 let _notice = DropNotice(self.dropped.clone());
2294 let _ = self.entered.send(());
2295 futures::future::pending::<Result<ChunkStream>>().await
2296 }
2297 }
2298
2299 #[async_trait]
2300 impl Tool for NestedTaskAdapter {
2301 fn spec(&self) -> flux_spec::ToolSpec {
2302 flux_spec::ToolSpec::read_only(
2303 "nested_task_adapter",
2304 "delegate through a nested one-shot runtime",
2305 serde_json::json!({ "type": "object", "properties": {} }),
2306 )
2307 }
2308
2309 async fn execute(
2310 &self,
2311 _ctx: &ToolContext,
2312 _params: serde_json::Value,
2313 ) -> Result<flux_runtime::ToolResult> {
2314 use crate::subagents::{RoleRegistry, SpawnLimits, SubAgents};
2315
2316 let mut roles = RoleRegistry::default();
2317 roles.insert(parse_role(
2318 "---\n---\nWait until the parent request is cancelled.",
2319 "worker",
2320 ));
2321 let factory: crate::subagents::ProviderFactory = Arc::new({
2322 let entered = self.child_entered.clone();
2323 let dropped = self.child_dropped.clone();
2324 move || {
2325 Ok(Box::new(HangingWorker {
2326 entered: entered.clone(),
2327 dropped: dropped.clone(),
2328 }) as Box<dyn Provider>)
2329 }
2330 });
2331 let mut limits = SpawnLimits::new(1024);
2332 limits.wall_clock = Some(std::time::Duration::from_secs(30));
2334 let sub_agents = SubAgents::new(roles, ToolRegistry::new(), factory, "mock", 1024)
2335 .with_limits(limits)
2336 .with_audit(self.audit.clone());
2337
2338 let root = std::env::temp_dir().join(format!(
2339 "flux-sdk-nested-turn-context-{}",
2340 std::process::id()
2341 ));
2342 std::fs::create_dir_all(&root)?;
2343 let mut nested = FlowClient::builder()
2344 .model("mock")
2345 .auto_approve(true)
2346 .build(Arc::new(ProseMock { text: "unused" }), &root)?;
2347 nested.with_sub_agents(sub_agents);
2348 let flow: flux_flow::ast::DraftAst = serde_json::from_value(serde_json::json!({
2349 "body": [{
2350 "kind": "call",
2351 "op": "task",
2352 "args": [{
2353 "kind": "lit",
2354 "value": { "role": "worker", "task": "wait" }
2355 }]
2356 }]
2357 }))
2358 .map_err(|error| flux_core::Error::Other(error.to_string()))?;
2359 let outcome = nested.execute_streamed(&flow).finish().await?;
2360 Ok(flux_runtime::ToolResult::ok(outcome.result))
2361 }
2362 }
2363
2364 struct NestedAdapterParent {
2365 calls: std::sync::atomic::AtomicUsize,
2366 }
2367
2368 #[async_trait]
2369 impl Provider for NestedAdapterParent {
2370 fn name(&self) -> &str {
2371 "mock"
2372 }
2373
2374 async fn stream(&self, request: Request) -> Result<ChunkStream> {
2375 if request_has_tool(&request, "declare_intent") {
2376 return Ok(Box::pin(futures::stream::iter(
2377 intent_chunks("delegate through the adapter", &["core"])
2378 .into_iter()
2379 .map(Ok),
2380 )));
2381 }
2382 let chunks = match self
2383 .calls
2384 .fetch_add(1, std::sync::atomic::Ordering::Relaxed)
2385 {
2386 0 => native_call("adapter-1", "nested_task_adapter", serde_json::json!({})),
2387 _ => vec![
2388 Chunk::Block(ContentBlock::Text {
2389 text: "unexpected completion".into(),
2390 }),
2391 Chunk::Done {
2392 stop_reason: Some(StopReason::EndTurn),
2393 },
2394 ],
2395 };
2396 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
2397 }
2398 }
2399
2400 #[tokio::test]
2401 async fn nested_streamed_task_inherits_parent_cancel_and_session_lineage() {
2402 let dir = std::env::temp_dir().join(format!(
2403 "flux-sdk-parent-turn-context-{}",
2404 std::process::id()
2405 ));
2406 std::fs::create_dir_all(&dir).unwrap();
2407 let audit = Arc::new(EventStore::in_memory().unwrap());
2408 let (entered_tx, mut entered_rx) = tokio::sync::mpsc::unbounded_channel();
2409 let (dropped_tx, mut dropped_rx) = tokio::sync::mpsc::unbounded_channel();
2410 let client = Client::builder()
2411 .model("mock")
2412 .auto_approve(true)
2413 .register_op(Arc::new(NestedTaskAdapter {
2414 audit: audit.clone(),
2415 child_entered: entered_tx,
2416 child_dropped: dropped_tx,
2417 }))
2418 .build(
2419 Box::new(NestedAdapterParent {
2420 calls: std::sync::atomic::AtomicUsize::new(0),
2421 }),
2422 &dir,
2423 )
2424 .unwrap();
2425 let session = client.default_session().unwrap();
2426 let parent_session = session.id().to_string();
2427 let turn = session.stream("delegate through the adapter");
2428
2429 tokio::time::timeout(std::time::Duration::from_secs(5), entered_rx.recv())
2430 .await
2431 .expect("the nested child must enter its parked provider")
2432 .expect("the nested child entry channel closed");
2433 turn.cancel();
2434 tokio::time::timeout(std::time::Duration::from_secs(2), dropped_rx.recv())
2435 .await
2436 .expect("parent cancellation must reach the nested child before its 30s deadline")
2437 .expect("the nested child drop channel closed");
2438 turn.finish().await.unwrap();
2439
2440 let children = audit.children_of(&parent_session).unwrap();
2441 assert_eq!(
2442 children.len(),
2443 1,
2444 "the child audit stream must be parent-linked"
2445 );
2446 let child = audit.info(&children[0]).unwrap();
2447 assert_eq!(
2448 child.context.correlation_id.as_deref(),
2449 Some(parent_session.as_str())
2450 );
2451 std::fs::remove_dir_all(&dir).ok();
2452 }
2453
2454 struct WidgetTool;
2457 #[async_trait]
2458 impl Tool for WidgetTool {
2459 fn spec(&self) -> flux_spec::ToolSpec {
2460 flux_spec::ToolSpec::read_only(
2461 "zzquux_probe",
2462 "a gated probe op",
2463 serde_json::json!({ "type": "object", "properties": {} }),
2464 )
2465 }
2466 async fn execute(
2467 &self,
2468 _c: &ToolContext,
2469 _params: serde_json::Value,
2470 ) -> Result<flux_runtime::ToolResult> {
2471 Ok(flux_runtime::ToolResult::ok("ok"))
2472 }
2473 }
2474
2475 #[tokio::test]
2480 async fn groups_gate_an_op_until_its_ambient_signal_surfaces() {
2481 use crate::observe::{SignalMatch, ToolGroup, KIND_SIGNAL};
2482
2483 let dir = std::env::temp_dir().join(format!("flux-sdk-groups-{}", std::process::id()));
2484 std::fs::create_dir_all(&dir).unwrap();
2485 let group = ToolGroup {
2486 name: "widgets".into(),
2487 description: String::new(),
2488 tools: vec!["zzquux_probe".into()],
2489 surface_when: vec![SignalMatch {
2490 kind: KIND_SIGNAL.to_string(),
2491 signal: Some("widgets_on".into()),
2492 }],
2493 };
2494
2495 let systems_gated = Arc::new(Mutex::new(Vec::new()));
2497 let client = Client::builder()
2498 .model("mock")
2499 .register_op(Arc::new(WidgetTool))
2500 .groups([group.clone()])
2501 .build(
2502 Box::new(SystemCaptureMock {
2503 systems: systems_gated.clone(),
2504 }),
2505 &dir,
2506 )
2507 .unwrap();
2508 client.run("hi").await.unwrap();
2509 let gated = systems_gated.lock().unwrap().join("\n");
2510 assert!(
2511 !gated.contains("zzquux_probe"),
2512 "the gated op must be absent from the catalog until its signal fires"
2513 );
2514
2515 let systems_on = Arc::new(Mutex::new(Vec::new()));
2517 let client = Client::builder()
2518 .model("mock")
2519 .register_op(Arc::new(WidgetTool))
2520 .groups([group])
2521 .ambient_signals(["widgets_on"])
2522 .build(
2523 Box::new(SystemCaptureMock {
2524 systems: systems_on.clone(),
2525 }),
2526 &dir,
2527 )
2528 .unwrap();
2529 client.run("hi").await.unwrap();
2530 let surfaced = systems_on.lock().unwrap().join("\n");
2531 assert!(
2532 surfaced.contains("zzquux_probe"),
2533 "the op must be advertised once its group's signal surfaces:\n{surfaced}"
2534 );
2535
2536 std::fs::remove_dir_all(&dir).ok();
2537 }
2538
2539 #[tokio::test]
2543 async fn with_compaction_trips_and_records_a_context_compacted_observation() {
2544 let dir = std::env::temp_dir().join(format!("flux-sdk-compact-{}", std::process::id()));
2545 std::fs::create_dir_all(&dir).unwrap();
2546 let client = Client::builder()
2549 .model("mock")
2550 .with_compaction(10)
2551 .build(Box::new(ProseMock { text: "ok" }), &dir)
2552 .unwrap();
2553
2554 for _ in 0..3 {
2557 client.run("tell me something").await.unwrap();
2558 }
2559
2560 let sid = client.session_id().unwrap();
2561 let obs = client.event_store().observations(&sid).unwrap();
2562 let compacted: Option<&crate::observe::Observation> =
2563 obs.iter().find(|o| o.kind == "context.compacted");
2564 let compacted = compacted.expect("a context.compacted observation must be recorded");
2565 assert!(
2566 compacted.data["from_messages"].as_u64().unwrap()
2567 > compacted.data["to_messages"].as_u64().unwrap(),
2568 "compaction shrank the message count: {:?}",
2569 compacted.data
2570 );
2571
2572 std::fs::remove_dir_all(&dir).ok();
2573 }
2574
2575 struct PricedMock(Usage);
2578 #[async_trait]
2579 impl Provider for PricedMock {
2580 fn name(&self) -> &str {
2581 "mock"
2582 }
2583 async fn stream(&self, req: Request) -> Result<ChunkStream> {
2584 let chunks = if request_has_tool(&req, "declare_intent") {
2585 intent_chunks("answer the user", &[])
2586 } else {
2587 vec![
2588 Chunk::TextDelta("ok".into()),
2589 Chunk::Block(ContentBlock::Text { text: "ok".into() }),
2590 Chunk::Usage(self.0.clone()),
2591 Chunk::Done {
2592 stop_reason: Some(StopReason::EndTurn),
2593 },
2594 ]
2595 };
2596 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
2597 }
2598 }
2599
2600 #[tokio::test]
2605 async fn session_projections_report_turns_history_and_cost() {
2606 use crate::observe::{ModelCost, TurnSummary};
2607 use flux_core::Role;
2608
2609 let dir = std::env::temp_dir().join(format!("flux-sdk-proj-{}", std::process::id()));
2610 std::fs::create_dir_all(&dir).unwrap();
2611 let per_call = Usage {
2612 input_tokens: 1000,
2613 output_tokens: 500,
2614 ..Default::default()
2615 };
2616 let client = Client::builder()
2617 .model("priced-mock")
2618 .build(Box::new(PricedMock(per_call)), &dir)
2619 .unwrap();
2620 let session = client.default_session().unwrap();
2621 session.send("first").await.unwrap();
2622 session.send("second").await.unwrap();
2623
2624 let turns: Vec<TurnSummary> = session.turns().unwrap();
2626 assert_eq!(turns.len(), 2, "one TurnSummary per turn: {turns:?}");
2627
2628 let history = session.history().unwrap();
2630 assert_eq!(history.len(), 4, "two turns = four messages");
2631 assert_eq!(history[0].role, Role::User);
2632 assert_eq!(history[1].role, Role::Assistant);
2633 assert_eq!(history[2].role, Role::User);
2634 assert_eq!(history[3].role, Role::Assistant);
2635
2636 let mut pricing = PricingTable::builtin();
2638 pricing.set(
2639 "priced-mock",
2640 flux_core::Rates {
2641 input: 1000.0,
2642 output: 1000.0,
2643 ..Default::default()
2644 },
2645 );
2646 let cost: Vec<ModelCost> = session.cost(&pricing).unwrap();
2647 let usd: f64 = cost
2648 .iter()
2649 .filter_map(|c| c.cost.as_ref())
2650 .map(|m| m.usd)
2651 .sum();
2652 assert!(usd > 0.0, "the priced model reports non-zero USD: {cost:?}");
2653
2654 let _ = session.run_trace().unwrap();
2656 let _ = session.efficiency().unwrap();
2657
2658 std::fs::remove_dir_all(&dir).ok();
2659 }
2660
2661 #[cfg(feature = "pricing")]
2666 #[test]
2667 fn pricing_feature_exposes_the_loader() {
2668 let table = crate::pricing::load_pricing_table();
2669 assert!(
2671 table.rates_for("claude-sonnet-4.6").is_some() || !format!("{table:?}").is_empty(),
2672 "the loaded table carries the built-in rates"
2673 );
2674 }
2675
2676 #[cfg(feature = "providers")]
2680 #[test]
2681 fn providers_from_spec_builds_a_credential_free_provider() {
2682 let (provider, model) =
2683 crate::providers::from_spec("ollama/qwen3").expect("ollama needs no credential");
2684 assert_eq!(model, "qwen3", "the resolved model id rides back");
2685 let _name = provider.name();
2687 }
2688
2689 #[test]
2695 fn default_build_pulls_no_optional_provider_batteries() {
2696 let manifest = include_str!("../Cargo.toml");
2697 assert!(
2698 manifest.contains("default = []"),
2699 "default features must be empty (provider-agnostic default build)"
2700 );
2701 assert!(
2702 manifest.contains("flux-providers = { workspace = true, optional = true }"),
2703 "flux-providers must be optional (the `providers` feature only)"
2704 );
2705 assert!(
2706 manifest.contains("flux-credentials = { workspace = true, optional = true }"),
2707 "flux-credentials must be optional (the `pricing` feature only)"
2708 );
2709 assert!(
2710 manifest.contains("flux-plugin = { workspace = true, optional = true }"),
2711 "flux-plugin must be optional (the `plugins` feature only)"
2712 );
2713 }
2714
2715 #[tokio::test]
2720 async fn session_replays_a_recorded_plan_hermetically() {
2721 let dir = std::env::temp_dir().join(format!("flux-sdk-replay-{}", std::process::id()));
2722 std::fs::remove_dir_all(&dir).ok();
2723 std::fs::create_dir_all(&dir).unwrap();
2724 let store = dir.join("state");
2725
2726 let sid = {
2728 let client = Client::builder()
2729 .model("mock")
2730 .auto_approve(true)
2731 .storage(Storage::dir(&store))
2732 .build(
2733 Box::new(PlanThenProseMock {
2734 calls: std::sync::atomic::AtomicUsize::new(0),
2735 }),
2736 &dir,
2737 )
2738 .unwrap();
2739 client.run("write a file").await.unwrap();
2740 client.session_id().unwrap()
2741 };
2742
2743 let client = Client::builder()
2746 .model("mock")
2747 .storage(Storage::dir(&store))
2748 .build(Box::new(NeverMock), &dir)
2749 .unwrap();
2750 let session = client.open_session(&sid).unwrap();
2751 struct NullSink;
2752 impl AgentSink for NullSink {}
2753 let mut sink = NullSink;
2754 let report = session.replay(None, &mut sink).await.unwrap();
2755 assert!(
2756 !report.plans.is_empty(),
2757 "the recorded plan replayed: {report:?}"
2758 );
2759 assert!(
2760 report.diverged.is_none(),
2761 "a faithful replay does not diverge: {report:?}"
2762 );
2763
2764 std::fs::remove_dir_all(&dir).ok();
2765 }
2766
2767 #[tokio::test]
2769 async fn replay_of_a_non_recorded_session_errors_honestly() {
2770 let dir = std::env::temp_dir().join(format!("flux-sdk-replay-none-{}", std::process::id()));
2771 std::fs::create_dir_all(&dir).unwrap();
2772 let client = Client::builder()
2773 .model("mock")
2774 .storage(Storage::dir(dir.join("state")))
2775 .build(Box::new(ProseMock { text: "hello" }), &dir)
2776 .unwrap();
2777 let session = client.default_session().unwrap();
2778 session.send("hi").await.unwrap(); struct NullSink;
2781 impl AgentSink for NullSink {}
2782 let mut sink = NullSink;
2783 let err = session
2784 .replay(None, &mut sink)
2785 .await
2786 .expect_err("a chat-only session is not replayable");
2787 assert!(
2788 err.to_string().contains("not replayable"),
2789 "the error is honest about why: {err}"
2790 );
2791
2792 std::fs::remove_dir_all(&dir).ok();
2793 }
2794
2795 struct BindPlanMock {
2798 calls: std::sync::atomic::AtomicUsize,
2799 }
2800 #[async_trait]
2801 impl Provider for BindPlanMock {
2802 fn name(&self) -> &str {
2803 "mock"
2804 }
2805 async fn stream(&self, req: Request) -> Result<ChunkStream> {
2806 if request_has_tool(&req, "declare_intent") {
2807 return Ok(Box::pin(futures::stream::iter(
2808 intent_chunks("read the note", &["workspace.read"])
2809 .into_iter()
2810 .map(Ok),
2811 )));
2812 }
2813 let n = self
2814 .calls
2815 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
2816 let chunks = if n == 0 {
2817 native_call("read-1", "read", serde_json::json!({"path": "note.txt"}))
2818 } else {
2819 vec![
2820 Chunk::Block(ContentBlock::Text {
2821 text: "done".into(),
2822 }),
2823 Chunk::Done {
2824 stop_reason: Some(StopReason::EndTurn),
2825 },
2826 ]
2827 };
2828 Ok(Box::pin(futures::stream::iter(chunks.into_iter().map(Ok))))
2829 }
2830 }
2831
2832 async fn record_bind_session(tag: &str) -> (std::path::PathBuf, std::path::PathBuf, String) {
2835 let dir = std::env::temp_dir().join(format!("flux-sdk-fork-{tag}-{}", std::process::id()));
2836 std::fs::remove_dir_all(&dir).ok();
2837 std::fs::create_dir_all(&dir).unwrap();
2838 std::fs::write(dir.join("note.txt"), "original").unwrap();
2839 let store = dir.join("state");
2840 let client = Client::builder()
2841 .model("mock")
2842 .auto_approve(true)
2843 .storage(Storage::dir(store.clone()))
2844 .build(
2845 Box::new(BindPlanMock {
2846 calls: std::sync::atomic::AtomicUsize::new(0),
2847 }),
2848 &dir,
2849 )
2850 .unwrap();
2851 client.run("read the note").await.unwrap();
2852 let sid = client.session_id().unwrap();
2853 (dir, store, sid)
2854 }
2855
2856 struct NullSink;
2857 impl AgentSink for NullSink {}
2858
2859 #[tokio::test]
2862 async fn fork_inject_diverges_and_leaves_the_original_untouched() {
2863 let (dir, store, sid) = record_bind_session("inject").await;
2864
2865 let client = Client::builder()
2866 .model("mock")
2867 .auto_approve(true)
2868 .storage(Storage::dir(store))
2869 .build(Box::new(NeverMock), &dir)
2870 .unwrap();
2871 let events = client.event_store();
2872 let head_before = events.head_seq(&sid).unwrap();
2873
2874 let session = client.open_session(&sid).unwrap();
2875 let fork = session.fork(0).await.unwrap();
2876 let mut sink = NullSink;
2877 fork.inject(&serde_json::json!("injected"), &mut sink)
2878 .await
2879 .unwrap();
2880
2881 assert_eq!(
2883 events.head_seq(&sid).unwrap(),
2884 head_before,
2885 "forking must not touch the original session's log"
2886 );
2887
2888 let diff = fork.diff(&session).unwrap();
2890 assert!(
2891 !diff.identical && !diff.rows.is_empty(),
2892 "the injected fork diverges from the original: {diff:?}"
2893 );
2894
2895 std::fs::remove_dir_all(&dir).ok();
2896 }
2897
2898 #[tokio::test]
2901 async fn fork_edit_diverges_on_a_bound_value() {
2902 let (dir, store, sid) = record_bind_session("edit").await;
2903
2904 let client = Client::builder()
2905 .model("mock")
2906 .auto_approve(true)
2907 .storage(Storage::dir(store))
2908 .build(Box::new(NeverMock), &dir)
2909 .unwrap();
2910 let session = client.open_session(&sid).unwrap();
2911 let fork = session.fork(0).await.unwrap();
2912
2913 let edited: crate::flow::DraftAst = serde_json::from_value(serde_json::json!({ "body": [
2915 { "kind": "bind", "name": "x", "value": { "kind": "lit", "value": "edited-value" } },
2916 { "kind": "return", "value": { "kind": "var", "name": "x" } }
2917 ]}))
2918 .unwrap();
2919 let mut sink = NullSink;
2920 fork.edit(&edited, &mut sink).await.unwrap();
2921
2922 let diff = fork.diff(&session).unwrap();
2923 assert!(
2924 !diff.identical,
2925 "the edited fork diverges from the original: {diff:?}"
2926 );
2927
2928 std::fs::remove_dir_all(&dir).ok();
2929 }
2930}