Skip to main content

ironflow_core/
provider.rs

1//! Provider trait and configuration types for agent invocations.
2//!
3//! The [`AgentProvider`] trait is the primary extension point in ironflow: implement it
4//! to plug in any AI backend (local model, HTTP API, mock, etc.) without changing
5//! your workflow code.
6//!
7//! Built-in implementations:
8//!
9//! * [`ClaudeCodeProvider`](crate::providers::claude::ClaudeCodeProvider) - local `claude` CLI.
10//! * `SshProvider` - remote via SSH (requires `transport-ssh` feature).
11//! * `DockerProvider` - Docker container (requires `transport-docker` feature).
12//! * `K8sEphemeralProvider` - ephemeral K8s pod (requires `transport-k8s` feature).
13//! * `K8sPersistentProvider` - persistent K8s pod (requires `transport-k8s` feature).
14//! * [`RecordReplayProvider`](crate::providers::record_replay::RecordReplayProvider) -
15//!   records and replays fixtures for deterministic testing.
16
17use std::collections::BTreeMap;
18use std::fmt;
19use std::future::Future;
20use std::marker::PhantomData;
21use std::pin::Pin;
22use std::sync::Arc;
23
24use schemars::JsonSchema;
25use serde::{Deserialize, Serialize};
26use serde_json::Value;
27
28use crate::error::AgentError;
29use crate::operations::agent::{Model, PermissionMode};
30use crate::retry::RetryPolicy;
31use crate::trace_context::WorkflowTraceContext;
32
33mod tool;
34mod tool_profile;
35
36pub use tool::Tool;
37pub use tool_profile::ToolProfile;
38
39/// Boxed future returned by [`AgentProvider::invoke`].
40pub type InvokeFuture<'a> =
41    Pin<Box<dyn Future<Output = Result<AgentOutput, AgentError>> + Send + 'a>>;
42
43// ── Typestate markers ──────────────────────────────────────────────
44
45/// Marker: no tools have been added via the builder.
46#[derive(Debug, Clone, Copy)]
47pub struct NoTools;
48
49/// Marker: at least one tool has been added via [`AgentConfig::allow_tool`],
50/// or a tool profile selected via [`AgentConfig::tool_profile`].
51#[derive(Debug, Clone, Copy)]
52pub struct WithTools;
53
54/// Marker: no JSON schema has been set via the builder.
55#[derive(Debug, Clone, Copy)]
56pub struct NoSchema;
57
58/// Marker: a JSON schema derived from `T` has been set via
59/// [`AgentConfig::output`]. The step answers with a `T`.
60pub struct WithSchema<T>(PhantomData<fn() -> T>);
61
62// Written by hand: a derive would require `T: Debug + Clone + Copy` for a
63// marker that never holds a `T`.
64impl<T> fmt::Debug for WithSchema<T> {
65    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
66        f.write_str("WithSchema")
67    }
68}
69
70impl<T> Clone for WithSchema<T> {
71    fn clone(&self) -> Self {
72        *self
73    }
74}
75
76impl<T> Copy for WithSchema<T> {}
77
78/// Marker: a pre-serialized JSON schema has been set via
79/// [`AgentConfig::output_schema_raw`]. The answer is not typed.
80#[derive(Debug, Clone, Copy)]
81pub struct RawSchema;
82
83// ── AgentInput ─────────────────────────────────────────────────────
84
85/// Declarative external input fetched into the agent's filesystem before invocation.
86///
87/// Each input is a URL that the provider must download and materialize at
88/// `mount_path` so the agent can read it via the `Read` tool.
89///
90/// Provider behavior:
91///
92/// * [`ClaudeCodeProvider`](crate::providers::claude::ClaudeCodeProvider) (local) -
93///   downloads via reqwest into a per-invocation temp directory and rewrites
94///   `mount_path` to the resolved local path.
95/// * `K8sEphemeralProvider` - injects a `curlimages/curl` initContainer that
96///   downloads each URL into a shared `emptyDir`, mounted on the main container
97///   at the parent directory of `mount_path`.
98///
99/// The `mount_path` must be an absolute path. Intermediate directories are
100/// created automatically.
101#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
102pub struct AgentInput {
103    /// Source URL to download (HTTP/HTTPS, including signed S3/R2 URLs).
104    pub url: String,
105
106    /// Absolute filesystem path where the file must be available inside the
107    /// agent's filesystem.
108    pub mount_path: String,
109}
110
111impl AgentInput {
112    /// Create a new input descriptor.
113    pub fn new(url: &str, mount_path: &str) -> Self {
114        Self {
115            url: url.to_string(),
116            mount_path: mount_path.to_string(),
117        }
118    }
119}
120
121// ── AgentConfig ────────────────────────────────────────────────────
122
123/// Serializable configuration passed to an [`AgentProvider`] for a single invocation.
124///
125/// Built by [`Agent::run`](crate::operations::agent::Agent::run) from the builder state.
126/// Provider implementations translate these fields into whatever format the underlying
127/// backend expects.
128///
129/// # Typestate: tools vs structured output
130///
131/// Claude CLI has a [known bug](https://github.com/anthropics/claude-code/issues/18536)
132/// where combining `--json-schema` with `--allowedTools` always returns
133/// `structured_output: null`. To prevent this at compile time, [`allow_tool`](Self::allow_tool)
134/// and [`output`](Self::output) / [`output_schema_raw`](Self::output_schema_raw) are mutually
135/// exclusive: using one removes the other from the available API.
136///
137/// ```
138/// use ironflow_core::provider::{AgentConfig, Tool};
139///
140/// // OK: tools only
141/// let _ = AgentConfig::new("search").allow_tool(Tool::WebSearch);
142///
143/// // OK: structured output only
144/// let _ = AgentConfig::new("classify").output_schema_raw(r#"{"type":"object"}"#);
145/// ```
146///
147/// ```compile_fail,E0599
148/// use ironflow_core::provider::{AgentConfig, Tool};
149/// // COMPILE ERROR: cannot add tools after setting structured output
150/// let _ = AgentConfig::new("x").output_schema_raw("{}").allow_tool(Tool::Read);
151/// ```
152///
153/// ```compile_fail,E0599
154/// use ironflow_core::provider::{AgentConfig, Tool};
155/// // COMPILE ERROR: cannot set structured output after adding tools
156/// let _ = AgentConfig::new("x").allow_tool(Tool::Read).output_schema_raw("{}");
157/// ```
158///
159/// ```compile_fail,E0599
160/// use ironflow_core::provider::{AgentConfig, ToolProfile};
161/// // COMPILE ERROR: a tool profile counts as tools
162/// let _ = AgentConfig::new("x").tool_profile(ToolProfile::new("bug")).output_schema_raw("{}");
163/// ```
164///
165/// **Workaround**: split the work into two steps -- one agent with tools to
166/// gather data, then a second agent with `.output::<T>()` to structure the result.
167#[derive(Debug, Clone, Serialize, Deserialize)]
168#[serde(bound(serialize = "", deserialize = ""))]
169#[non_exhaustive]
170pub struct AgentConfig<Tools = NoTools, Schema = NoSchema> {
171    /// Optional system prompt that sets the agent's persona or constraints.
172    pub system_prompt: Option<String>,
173
174    /// The user prompt - the main instruction to the agent.
175    pub prompt: String,
176
177    /// Which model to use for this invocation.
178    ///
179    /// Accepts any string. Use [`Model`] constants for well-known Claude models
180    /// (e.g. `Model::SONNET`), or pass a custom identifier for other providers.
181    #[serde(default = "default_model")]
182    pub model: String,
183
184    /// Allowlist of tool names the agent may invoke (empty = provider default).
185    #[serde(default)]
186    pub allowed_tools: Vec<String>,
187
188    /// Denylist of tool names the agent MUST NOT invoke.
189    ///
190    /// Maps to `--disallowedTools` on the Claude CLI. Unlike
191    /// [`allowed_tools`](Self::allowed_tools), this does **not** activate any
192    /// tools; it only filters out tools that would otherwise be loaded by
193    /// default. As such, it is safe to combine with structured output
194    /// ([`output`](Self::output)) without triggering the Claude CLI bug that
195    /// affects `--json-schema` + `--allowedTools`.
196    #[serde(default)]
197    pub disallowed_tools: Vec<String>,
198
199    /// Named tool profile the provider exposes to this step.
200    ///
201    /// Set it with [`AgentConfig::tool_profile`]. `None` means the provider's
202    /// default tools only (none unless it has some).
203    #[serde(default, skip_serializing_if = "Option::is_none")]
204    pub tool_profile: Option<ToolProfile>,
205
206    /// Maximum number of agentic turns before the provider should stop.
207    pub max_turns: Option<u32>,
208
209    /// Maximum number of tool calls executed concurrently within a single
210    /// turn's group of consecutive read-only calls (default 4). `1` restores
211    /// fully sequential tool execution.
212    #[serde(default = "default_max_parallel_tools")]
213    pub max_parallel_tools: usize,
214
215    /// Maximum spend in USD for this single invocation.
216    pub max_budget_usd: Option<f64>,
217
218    /// Working directory for the agent process.
219    pub working_dir: Option<String>,
220
221    /// Path to an MCP server configuration file.
222    pub mcp_config: Option<String>,
223
224    /// When `true`, pass `--strict-mcp-config` to the Claude CLI so it only
225    /// loads MCP servers from [`mcp_config`](Self::mcp_config) and ignores
226    /// any global/user MCP configuration (e.g. `~/.claude.json`).
227    ///
228    /// Useful to prevent global MCP servers from leaking tools into steps
229    /// that request `structured_output`, which triggers the Claude CLI bug
230    /// where `--json-schema` combined with any active tool returns
231    /// `structured_output: null`. See
232    /// <https://github.com/anthropics/claude-code/issues/18536>.
233    ///
234    /// Combine with `mcp_config` set to a file containing
235    /// `{"mcpServers":{}}` to disable every MCP server for the invocation.
236    #[serde(default)]
237    pub strict_mcp_config: bool,
238
239    /// When `true`, pass `--bare` to Claude CLI. Bare mode disables:
240    /// - auto-memory (automatic creation of `~/.claude/.../memory/*.md` files)
241    /// - `CLAUDE.md` auto-discovery (no global/project `CLAUDE.md` loaded)
242    /// - hooks, LSP, plugin sync, attribution, background prefetches
243    ///
244    /// Recommended for orchestrator agents that should not have any implicit
245    /// side effects on the user's filesystem or inherit user-level context.
246    ///
247    /// # Authentication requirement
248    ///
249    /// `--bare` is **only compatible with an Anthropic API key**
250    /// (`ANTHROPIC_API_KEY` environment variable). It does **not** work with
251    /// OAuth authentication (`claude /login` / keychain-stored credentials),
252    /// because bare mode disables keychain reads.
253    #[serde(default)]
254    pub bare: bool,
255
256    /// Permission mode controlling how the agent handles tool-use approvals.
257    #[serde(default)]
258    pub permission_mode: PermissionMode,
259
260    /// Optional JSON Schema string. When set, the provider should request
261    /// structured (typed) output from the model.
262    #[serde(alias = "output_schema")]
263    pub json_schema: Option<String>,
264
265    /// Optional session ID to resume a previous conversation.
266    ///
267    /// When set, the provider should continue the conversation from the
268    /// specified session rather than starting a new one.
269    pub resume_session_id: Option<String>,
270
271    /// Enable verbose/debug mode to capture the full conversation trace.
272    ///
273    /// When `true`, the provider uses streaming output (`stream-json`) to
274    /// record every assistant message and tool call. The resulting
275    /// [`AgentOutput::debug_messages`] field will contain the conversation
276    /// trace for inspection.
277    #[serde(default)]
278    pub verbose: bool,
279
280    /// Custom labels applied to the pod (K8s providers only).
281    ///
282    /// Non-K8s providers ignore this field. Labels are merged with the
283    /// provider-level pod labels and the hardcoded ironflow labels. In case
284    /// of conflict, hardcoded labels always win, then invocation-level labels,
285    /// then provider-level defaults.
286    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
287    pub pod_labels: BTreeMap<String, String>,
288
289    /// External inputs to materialize on the agent's filesystem before invocation.
290    ///
291    /// See [`AgentInput`] for the semantics. The provider is responsible for
292    /// fetching each URL and placing it at `mount_path` before the agent runs.
293    /// Add inputs with [`AgentConfig::input_file`].
294    #[serde(default, skip_serializing_if = "Vec::is_empty")]
295    pub inputs: Vec<AgentInput>,
296
297    /// When `true`, a failure of this step does not fail the run.
298    #[serde(default)]
299    pub allow_failure: bool,
300
301    /// Optional step-level retry policy.
302    #[serde(default, skip_serializing_if = "Option::is_none")]
303    pub retry: Option<RetryPolicy>,
304
305    /// Optional W3C trace context for distributed tracing propagation.
306    ///
307    /// When set, providers can inject the `traceparent` header into
308    /// outgoing HTTP requests (LLM APIs, MCP servers) to correlate
309    /// workflow spans with downstream service spans.
310    #[serde(default, skip_serializing_if = "Option::is_none")]
311    pub trace_context: Option<WorkflowTraceContext>,
312
313    /// Zero-sized typestate marker (not serialized).
314    #[serde(skip)]
315    pub(crate) _marker: PhantomData<(Tools, Schema)>,
316}
317
318fn default_model() -> String {
319    Model::SONNET.to_string()
320}
321
322fn default_max_parallel_tools() -> usize {
323    4
324}
325
326// ── Constructor (base type only) ───────────────────────────────────
327
328impl AgentConfig {
329    /// Create an `AgentConfig` with required fields and defaults for the rest.
330    pub fn new(prompt: &str) -> Self {
331        Self {
332            system_prompt: None,
333            prompt: prompt.to_string(),
334            model: Model::SONNET.to_string(),
335            allowed_tools: Vec::new(),
336            disallowed_tools: Vec::new(),
337            tool_profile: None,
338            max_turns: None,
339            max_parallel_tools: 4,
340            max_budget_usd: None,
341            working_dir: None,
342            mcp_config: None,
343            strict_mcp_config: false,
344            bare: false,
345            permission_mode: PermissionMode::Default,
346            json_schema: None,
347
348            resume_session_id: None,
349            verbose: false,
350            pod_labels: BTreeMap::new(),
351            inputs: Vec::new(),
352            allow_failure: false,
353            retry: None,
354            trace_context: None,
355            _marker: PhantomData,
356        }
357    }
358}
359
360// ── Methods available on ALL typestate variants ────────────────────
361
362impl<Tools, Schema> AgentConfig<Tools, Schema> {
363    /// Set the system prompt.
364    pub fn system_prompt(mut self, prompt: &str) -> Self {
365        self.system_prompt = Some(prompt.to_string());
366        self
367    }
368
369    /// Set the model name.
370    pub fn model(mut self, model: &str) -> Self {
371        self.model = model.to_string();
372        self
373    }
374
375    /// Set the maximum budget in USD.
376    pub fn max_budget_usd(mut self, budget: f64) -> Self {
377        self.max_budget_usd = Some(budget);
378        self
379    }
380
381    /// Set the maximum number of turns.
382    pub fn max_turns(mut self, turns: u32) -> Self {
383        self.max_turns = Some(turns);
384        self
385    }
386
387    /// Set the maximum number of tool calls executed concurrently within a
388    /// turn's group of consecutive read-only calls.
389    ///
390    /// # Panics
391    ///
392    /// Panics if `n` is `0`.
393    ///
394    /// # Examples
395    ///
396    /// ```
397    /// use ironflow_core::provider::AgentConfig;
398    ///
399    /// let config = AgentConfig::new("summarize the repo").max_parallel_tools(2);
400    /// assert_eq!(config.max_parallel_tools, 2);
401    /// ```
402    pub fn max_parallel_tools(mut self, n: usize) -> Self {
403        assert!(n > 0, "max_parallel_tools must be greater than 0");
404        self.max_parallel_tools = n;
405        self
406    }
407
408    /// Set the working directory.
409    pub fn working_dir(mut self, dir: &str) -> Self {
410        self.working_dir = Some(dir.to_string());
411        self
412    }
413
414    /// Set the permission mode.
415    pub fn permission_mode(mut self, mode: PermissionMode) -> Self {
416        self.permission_mode = mode;
417        self
418    }
419
420    /// Enable verbose/debug mode.
421    pub fn verbose(mut self, enabled: bool) -> Self {
422        self.verbose = enabled;
423        self
424    }
425
426    /// Set the MCP server configuration file path.
427    pub fn mcp_config(mut self, config: &str) -> Self {
428        self.mcp_config = Some(config.to_string());
429        self
430    }
431
432    /// Enable strict MCP config mode.
433    ///
434    /// When `true`, the Claude CLI is invoked with `--strict-mcp-config`,
435    /// which disables loading of any MCP server defined outside the
436    /// [`mcp_config`](Self::mcp_config) file (the global `~/.claude.json`
437    /// and user-level configs are ignored).
438    ///
439    /// This is the recommended way to prevent global MCP servers from
440    /// silently injecting tools into a structured-output step and
441    /// triggering the Claude CLI bug that returns `structured_output: null`
442    /// whenever any tool is active. See
443    /// <https://github.com/anthropics/claude-code/issues/18536>.
444    ///
445    /// # Examples
446    ///
447    /// ```
448    /// use ironflow_core::provider::AgentConfig;
449    /// use schemars::JsonSchema;
450    ///
451    /// #[derive(serde::Deserialize, JsonSchema)]
452    /// struct Out { ok: bool }
453    ///
454    /// // Isolate the step from any global MCP server so structured output works.
455    /// let config = AgentConfig::new("classify this")
456    ///     .strict_mcp_config(true)
457    ///     .mcp_config(r#"{"mcpServers":{}}"#)
458    ///     .output::<Out>();
459    /// ```
460    pub fn strict_mcp_config(mut self, strict: bool) -> Self {
461        self.strict_mcp_config = strict;
462        self
463    }
464
465    /// Enable bare mode (minimal Claude Code environment, see `--bare`).
466    ///
467    /// When `true`, the Claude CLI is invoked with `--bare`, which disables:
468    /// - auto-memory (no automatic `~/.claude/.../memory/*.md` file creation)
469    /// - `CLAUDE.md` auto-discovery (neither global nor project-level)
470    /// - hooks, LSP, plugin sync, attribution, background prefetches,
471    ///   keychain reads
472    ///
473    /// Sets `CLAUDE_CODE_SIMPLE=1` in the child process.
474    ///
475    /// Recommended for orchestrator steps that should not have any implicit
476    /// side effects on the user's filesystem or inherit user-level context
477    /// (email, preferences, etc.).
478    ///
479    /// # Authentication requirement
480    ///
481    /// `--bare` is **only compatible with an Anthropic API key**
482    /// (`ANTHROPIC_API_KEY` environment variable). It does **not** work with
483    /// OAuth authentication (`claude /login` / keychain-stored credentials),
484    /// because bare mode disables keychain reads. Invoking a bare agent on an
485    /// OAuth-only host will fail with an authentication error.
486    ///
487    /// # Examples
488    ///
489    /// ```
490    /// use ironflow_core::provider::AgentConfig;
491    ///
492    /// let config = AgentConfig::new("classify this")
493    ///     .bare(true);
494    /// ```
495    pub fn bare(mut self, enabled: bool) -> Self {
496        self.bare = enabled;
497        self
498    }
499
500    /// Mark this step as allowed to fail without stopping the run.
501    ///
502    /// # Examples
503    ///
504    /// ```
505    /// use ironflow_core::provider::AgentConfig;
506    ///
507    /// let config = AgentConfig::new("lint the code").allow_failure();
508    /// assert!(config.allow_failure);
509    /// ```
510    pub fn allow_failure(mut self) -> Self {
511        self.allow_failure = true;
512        self
513    }
514
515    /// Replace the entire disallowed-tools list.
516    ///
517    /// Maps to `--disallowedTools` on the Claude CLI. This method is available
518    /// on **every** typestate variant (including
519    /// [`AgentConfig<NoTools, WithSchema<T>>`]) because, unlike
520    /// [`allow_tool`](AgentConfig::allow_tool), `disallowed_tools` does not
521    /// activate any tool -- it only filters out tools that would otherwise be
522    /// loaded by default.
523    ///
524    /// As such, it is safe to combine with structured output:
525    ///
526    /// # Examples
527    ///
528    /// ```
529    /// use ironflow_core::provider::{AgentConfig, Tool};
530    /// use schemars::JsonSchema;
531    ///
532    /// #[derive(serde::Deserialize, JsonSchema)]
533    /// struct Out { ok: bool }
534    ///
535    /// let config = AgentConfig::new("classify this")
536    ///     .disallowed_tools([Tool::Write, Tool::Edit])
537    ///     .output::<Out>();
538    /// assert_eq!(config.disallowed_tools, vec!["Write", "Edit"]);
539    /// ```
540    pub fn disallowed_tools<I>(mut self, tools: I) -> Self
541    where
542        I: IntoIterator<Item = Tool>,
543    {
544        self.disallowed_tools = tools.into_iter().map(|tool| tool.to_string()).collect();
545        self
546    }
547
548    /// Add a single custom pod label (K8s providers only).
549    ///
550    /// Can be called multiple times. Non-K8s providers ignore this field.
551    ///
552    /// # Examples
553    ///
554    /// ```
555    /// use ironflow_core::provider::AgentConfig;
556    ///
557    /// let config = AgentConfig::new("analyze")
558    ///     .pod_label("ironflow.io/network-profile", "grafana-only")
559    ///     .pod_label("team", "observability");
560    /// ```
561    pub fn pod_label(mut self, key: &str, value: &str) -> Self {
562        self.pod_labels.insert(key.to_string(), value.to_string());
563        self
564    }
565
566    /// Replace the entire custom pod labels map (K8s providers only).
567    ///
568    /// Non-K8s providers ignore this field.
569    ///
570    /// # Examples
571    ///
572    /// ```
573    /// use std::collections::BTreeMap;
574    /// use ironflow_core::provider::AgentConfig;
575    ///
576    /// let mut labels = BTreeMap::new();
577    /// labels.insert("env".to_string(), "staging".to_string());
578    /// let config = AgentConfig::new("deploy").pod_labels(labels);
579    /// ```
580    pub fn pod_labels(mut self, labels: BTreeMap<String, String>) -> Self {
581        self.pod_labels = labels;
582        self
583    }
584
585    /// Set a session ID to resume a previous conversation.
586    pub fn resume(mut self, session_id: &str) -> Self {
587        self.resume_session_id = Some(session_id.to_string());
588        self
589    }
590
591    /// Set a step-level retry policy.
592    ///
593    /// # Examples
594    ///
595    /// ```
596    /// use ironflow_core::provider::AgentConfig;
597    /// use ironflow_core::retry::RetryPolicy;
598    ///
599    /// let config = AgentConfig::new("Summarize this document")
600    ///     .retry_policy(RetryPolicy::new(3));
601    /// assert!(config.retry.is_some());
602    /// ```
603    pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
604        self.retry = Some(policy);
605        self
606    }
607
608    /// Attach a [`WorkflowTraceContext`] for distributed tracing.
609    ///
610    /// When set, providers can inject the `traceparent` header into
611    /// outgoing HTTP requests to correlate workflow spans with
612    /// downstream service spans.
613    ///
614    /// # Examples
615    ///
616    /// ```
617    /// use ironflow_core::provider::AgentConfig;
618    /// use ironflow_core::trace_context::WorkflowTraceContext;
619    ///
620    /// let ctx = WorkflowTraceContext::new_root();
621    /// let config = AgentConfig::new("classify this")
622    ///     .trace_context(ctx);
623    /// assert!(config.trace_context.is_some());
624    /// ```
625    pub fn trace_context(mut self, ctx: WorkflowTraceContext) -> Self {
626        self.trace_context = Some(ctx);
627        self
628    }
629
630    /// Declare an external input that the provider must materialize on the
631    /// agent's filesystem before invocation.
632    ///
633    /// `url` is fetched (HTTP/HTTPS) and written to `mount_path` (absolute
634    /// path) inside the agent's runtime. Each provider materializes inputs
635    /// in its own way:
636    ///
637    /// * Local provider: downloads to a temp dir on the host.
638    /// * K8s providers: spawn a `curlimages/curl` initContainer that downloads
639    ///   into a shared `emptyDir` mounted on the main container.
640    ///
641    /// Can be called multiple times to declare several inputs.
642    ///
643    /// # Examples
644    ///
645    /// ```
646    /// use ironflow_core::provider::{AgentConfig, Tool};
647    ///
648    /// let config = AgentConfig::new("Read /work/dossier.pdf and summarize")
649    ///     .allow_tool(Tool::Read)
650    ///     .input_file("https://r2.example.com/dossier.pdf", "/work/dossier.pdf");
651    /// ```
652    pub fn input_file(mut self, url: &str, mount_path: &str) -> Self {
653        self.inputs.push(AgentInput::new(url, mount_path));
654        self
655    }
656
657    /// Convert to a different typestate by moving all fields.
658    ///
659    /// Safe because the marker is a zero-sized [`PhantomData`] -- no
660    /// runtime data changes.
661    fn change_state<T2, S2>(self) -> AgentConfig<T2, S2> {
662        AgentConfig {
663            system_prompt: self.system_prompt,
664            prompt: self.prompt,
665            model: self.model,
666            allowed_tools: self.allowed_tools,
667            disallowed_tools: self.disallowed_tools,
668            tool_profile: self.tool_profile,
669            max_turns: self.max_turns,
670            max_parallel_tools: self.max_parallel_tools,
671            max_budget_usd: self.max_budget_usd,
672            working_dir: self.working_dir,
673            mcp_config: self.mcp_config,
674            strict_mcp_config: self.strict_mcp_config,
675            bare: self.bare,
676            permission_mode: self.permission_mode,
677            json_schema: self.json_schema,
678            resume_session_id: self.resume_session_id,
679            verbose: self.verbose,
680            pod_labels: self.pod_labels,
681            inputs: self.inputs,
682            allow_failure: self.allow_failure,
683            retry: self.retry,
684            trace_context: self.trace_context,
685            _marker: PhantomData,
686        }
687    }
688}
689
690// ── allow_tool: only when no schema is set ─────────────────────────
691
692impl<Tools> AgentConfig<Tools, NoSchema> {
693    /// Add an allowed tool.
694    ///
695    /// Can be called multiple times to allow several tools. Returns an
696    /// [`AgentConfig<WithTools, NoSchema>`], which **cannot** call
697    /// [`output`](AgentConfig::output) or [`output_schema_raw`](AgentConfig::output_schema_raw).
698    ///
699    /// This restriction exists because Claude CLI has a
700    /// [known bug](https://github.com/anthropics/claude-code/issues/18536)
701    /// where `--json-schema` combined with `--allowedTools` always returns
702    /// `structured_output: null`.
703    ///
704    /// **Workaround**: use two sequential agent steps -- one with tools to
705    /// gather data, then one with `.output::<T>()` to structure the result.
706    ///
707    /// # Examples
708    ///
709    /// ```
710    /// use ironflow_core::provider::{AgentConfig, Tool};
711    ///
712    /// let config = AgentConfig::new("search the web")
713    ///     .allow_tool(Tool::WebSearch)
714    ///     .allow_tool(Tool::Custom("mcp__docs__lookup".to_string()));
715    /// assert_eq!(config.allowed_tools, vec!["WebSearch", "mcp__docs__lookup"]);
716    /// ```
717    ///
718    /// ```compile_fail,E0599
719    /// use ironflow_core::provider::{AgentConfig, Tool};
720    /// // ERROR: cannot set structured output after adding tools
721    /// let _ = AgentConfig::new("x")
722    ///     .allow_tool(Tool::Read)
723    ///     .output_schema_raw(r#"{"type":"object"}"#);
724    /// ```
725    pub fn allow_tool(mut self, tool: Tool) -> AgentConfig<WithTools, NoSchema> {
726        self.allowed_tools.push(tool.to_string());
727        self.change_state()
728    }
729
730    /// Select the named tool profile the provider exposes to this step.
731    ///
732    /// Profiles are registered with
733    /// [`HttpAgentProvider::with_tool_profile`](crate::providers::http::HttpAgentProvider::with_tool_profile);
734    /// declare each [`ToolProfile`] once as a constant and share it. A profile
735    /// the provider does not have fails the step with
736    /// [`AgentError::UnknownToolProfile`], never falling back to other tools.
737    /// Claude CLI providers fail with [`AgentError::ToolProfileUnsupported`].
738    /// Like [`allow_tool`](Self::allow_tool), it rules out
739    /// [`output`](AgentConfig::output): split into two steps.
740    ///
741    /// # Examples
742    ///
743    /// ```
744    /// use ironflow_core::provider::{AgentConfig, ToolProfile};
745    ///
746    /// const BUG: ToolProfile = ToolProfile::new("bug");
747    ///
748    /// let config = AgentConfig::new("Find the root cause").tool_profile(BUG);
749    /// assert_eq!(config.tool_profile, Some(BUG));
750    /// ```
751    ///
752    /// ```compile_fail,E0308
753    /// use ironflow_core::provider::AgentConfig;
754    /// // COMPILE ERROR: a profile is a `ToolProfile`, not a string
755    /// let _ = AgentConfig::new("x").tool_profile("bug");
756    /// ```
757    pub fn tool_profile(mut self, profile: ToolProfile) -> AgentConfig<WithTools, NoSchema> {
758        self.tool_profile = Some(profile);
759        self.change_state()
760    }
761}
762
763// ── output: only when no tools are set ─────────────────────────────
764
765impl<Schema> AgentConfig<NoTools, Schema> {
766    /// Set structured output from a Rust type implementing [`JsonSchema`].
767    ///
768    /// The schema is serialized once at build time. When set, the provider
769    /// will request typed output conforming to this schema.
770    ///
771    /// **Important:** structured output requires `max_turns >= 2`.
772    ///
773    /// Returns an [`AgentConfig<NoTools, WithSchema<T>>`], which remembers
774    /// `T` (a workflow step built from it answers with a `T`) and **cannot**
775    /// call [`allow_tool`](AgentConfig::allow_tool).
776    ///
777    /// This restriction exists because Claude CLI has a
778    /// [known bug](https://github.com/anthropics/claude-code/issues/18536)
779    /// where `--json-schema` combined with `--allowedTools` always returns
780    /// `structured_output: null`.
781    ///
782    /// **Workaround**: use two sequential agent steps -- one with tools to
783    /// gather data, then one with `.output::<T>()` to structure the result.
784    ///
785    /// # Known limitations of Claude CLI structured output
786    ///
787    /// The Claude CLI does not guarantee strict schema conformance for
788    /// structured output. The following upstream bugs affect the behavior:
789    ///
790    /// - **Schema flattening** ([anthropics/claude-agent-sdk-python#502]):
791    ///   a schema like `{"type":"object","properties":{"items":{"type":"array",...}}}`
792    ///   may return a bare array instead of the wrapper object. The CLI
793    ///   non-deterministically flattens schemas with a single array field.
794    /// - **Non-deterministic wrapping** ([anthropics/claude-agent-sdk-python#374]):
795    ///   the same prompt can produce differently wrapped output across runs.
796    /// - **No conformance guarantee** ([anthropics/claude-code#9058]):
797    ///   the CLI does not validate output against the provided JSON schema.
798    ///
799    /// Because of these bugs, ironflow's provider layer applies multiple
800    /// fallback strategies when extracting the structured value (see
801    /// [`extract_structured_value`](crate::providers::claude::common::extract_structured_value)).
802    ///
803    /// [anthropics/claude-agent-sdk-python#502]: https://github.com/anthropics/claude-agent-sdk-python/issues/502
804    /// [anthropics/claude-agent-sdk-python#374]: https://github.com/anthropics/claude-agent-sdk-python/issues/374
805    /// [anthropics/claude-code#9058]: https://github.com/anthropics/claude-code/issues/9058
806    ///
807    /// # Examples
808    ///
809    /// ```
810    /// use ironflow_core::provider::AgentConfig;
811    /// use schemars::JsonSchema;
812    ///
813    /// #[derive(serde::Deserialize, JsonSchema)]
814    /// struct Labels { labels: Vec<String> }
815    ///
816    /// let config = AgentConfig::new("classify this text")
817    ///     .output::<Labels>();
818    /// ```
819    ///
820    /// ```compile_fail,E0599
821    /// use ironflow_core::provider::{AgentConfig, Tool};
822    /// use schemars::JsonSchema;
823    /// #[derive(serde::Deserialize, JsonSchema)]
824    /// struct Out { x: i32 }
825    /// // ERROR: cannot add tools after setting structured output
826    /// let _ = AgentConfig::new("x").output::<Out>().allow_tool(Tool::Read);
827    /// ```
828    /// # Panics
829    ///
830    /// Panics if the schema generated by `schemars` cannot be serialized
831    /// to JSON. This indicates a bug in the type's `JsonSchema` derive,
832    /// not a recoverable runtime error.
833    pub fn output<T: JsonSchema>(mut self) -> AgentConfig<NoTools, WithSchema<T>> {
834        let schema = schemars::schema_for!(T);
835        let serialized = serde_json::to_string(&schema).unwrap_or_else(|e| {
836            panic!(
837                "failed to serialize JSON schema for {}: {e}",
838                std::any::type_name::<T>()
839            )
840        });
841        self.json_schema = Some(serialized);
842        self.change_state()
843    }
844
845    /// Set structured output from a pre-serialized JSON Schema string.
846    ///
847    /// Returns an [`AgentConfig<NoTools, RawSchema>`], whose answer is not
848    /// typed, and which **cannot** call [`allow_tool`](AgentConfig::allow_tool).
849    /// Prefer [`output`](Self::output). See it for the rationale and
850    /// workaround.
851    pub fn output_schema_raw(mut self, schema: &str) -> AgentConfig<NoTools, RawSchema> {
852        self.json_schema = Some(schema.to_string());
853        self.change_state()
854    }
855}
856
857// ── From conversions to base type ──────────────────────────────────
858
859impl From<AgentConfig<WithTools, NoSchema>> for AgentConfig {
860    fn from(config: AgentConfig<WithTools, NoSchema>) -> Self {
861        config.change_state()
862    }
863}
864
865impl<T> From<AgentConfig<NoTools, WithSchema<T>>> for AgentConfig {
866    fn from(config: AgentConfig<NoTools, WithSchema<T>>) -> Self {
867        config.change_state()
868    }
869}
870
871impl From<AgentConfig<NoTools, RawSchema>> for AgentConfig {
872    fn from(config: AgentConfig<NoTools, RawSchema>) -> Self {
873        config.change_state()
874    }
875}
876
877// ── AgentOutput ────────────────────────────────────────────────────
878
879/// Raw output returned by an [`AgentProvider`] after a successful invocation.
880///
881/// Carries the agent's response value together with usage and billing metadata.
882#[derive(Clone, Debug, Serialize, Deserialize)]
883#[non_exhaustive]
884pub struct AgentOutput {
885    /// The agent's response. A plain [`Value::String`] for text mode, or an
886    /// arbitrary JSON value when a JSON schema was requested.
887    pub value: Value,
888
889    /// Provider-assigned session identifier, useful for resuming conversations.
890    pub session_id: Option<String>,
891
892    /// Total cost in USD for this invocation, if reported by the provider.
893    pub cost_usd: Option<f64>,
894
895    /// Uncached input tokens (excludes cache reads and writes), if reported.
896    pub input_tokens: Option<u64>,
897
898    /// Input tokens served from the prompt cache, if reported.
899    #[serde(default)]
900    pub cache_read_input_tokens: Option<u64>,
901
902    /// Input tokens written to the prompt cache, if reported.
903    #[serde(default)]
904    pub cache_creation_input_tokens: Option<u64>,
905
906    /// Number of output tokens generated, if reported.
907    pub output_tokens: Option<u64>,
908
909    /// The concrete model identifier used (e.g. `"claude-sonnet-4-20250514"`).
910    pub model: Option<String>,
911
912    /// Wall-clock duration of the invocation in milliseconds.
913    pub duration_ms: u64,
914
915    /// Conversation trace captured when [`AgentConfig::verbose`] is `true`.
916    ///
917    /// Contains every assistant message and tool call made during the
918    /// invocation, in chronological order. `None` when verbose mode is off.
919    pub debug_messages: Option<Vec<DebugMessage>>,
920}
921
922/// A single assistant turn captured during a verbose invocation.
923///
924/// Each `DebugMessage` represents one assistant response, which may contain
925/// free-form text, tool calls, or both.
926///
927/// # Examples
928///
929/// ```no_run
930/// use ironflow_core::prelude::*;
931///
932/// # async fn example() -> Result<(), OperationError> {
933/// let provider = ClaudeCodeProvider::new();
934/// let result = Agent::new()
935///     .prompt("List files in src/")
936///     .verbose()
937///     .run(&provider)
938///     .await?;
939///
940/// if let Some(messages) = result.debug_messages() {
941///     for msg in messages {
942///         println!("{msg}");
943///     }
944/// }
945/// # Ok(())
946/// # }
947/// ```
948#[derive(Debug, Clone, Serialize, Deserialize)]
949#[non_exhaustive]
950pub struct DebugMessage {
951    /// Free-form text produced by the assistant in this turn, if any.
952    pub text: Option<String>,
953
954    /// Extended thinking blocks produced by the model in this turn.
955    ///
956    /// Available only when the model emits `thinking` content blocks
957    /// (Opus 4.7 adaptive thinking, Claude 3.7+ extended thinking, etc.).
958    /// The blocks are joined in arrival order.
959    #[serde(default, skip_serializing_if = "Option::is_none")]
960    pub thinking: Option<String>,
961
962    /// `true` when the model emitted a `thinking` content block but the
963    /// text was redacted (only a signature is provided).
964    ///
965    /// Opus 4.7 adaptive thinking and the `display: "omitted"` setting both
966    /// produce signature-only thinking blocks: the model proves it reasoned
967    /// without exposing the chain of thought. The UI should still show a
968    /// badge so the user knows thinking happened.
969    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
970    pub thinking_redacted: bool,
971
972    /// Tool calls made by the assistant in this turn.
973    pub tool_calls: Vec<DebugToolCall>,
974
975    /// Tool results received from the user/runtime for the preceding tool calls.
976    ///
977    /// In the Claude stream-json format, tool results come as `"type":"user"`
978    /// messages whose content is a list of `tool_result` blocks. We attach
979    /// them to the turn that emitted the matching `tool_use` so the timeline
980    /// stays compact.
981    #[serde(default, skip_serializing_if = "Vec::is_empty")]
982    pub tool_results: Vec<DebugToolResult>,
983
984    /// The model's stop reason for this turn (e.g. `"end_turn"`, `"tool_use"`).
985    pub stop_reason: Option<String>,
986
987    /// Input tokens consumed by this turn, if reported.
988    #[serde(default, skip_serializing_if = "Option::is_none")]
989    pub input_tokens: Option<u64>,
990
991    /// Output tokens generated by this turn, if reported.
992    #[serde(default, skip_serializing_if = "Option::is_none")]
993    pub output_tokens: Option<u64>,
994}
995
996impl fmt::Display for DebugMessage {
997    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
998        if let Some(ref thinking) = self.thinking {
999            writeln!(f, "[thinking] {thinking}")?;
1000        } else if self.thinking_redacted {
1001            writeln!(f, "[thinking redacted]")?;
1002        }
1003        if let Some(ref text) = self.text {
1004            writeln!(f, "[assistant] {text}")?;
1005        }
1006        for tc in &self.tool_calls {
1007            write!(f, "{tc}")?;
1008        }
1009        for tr in &self.tool_results {
1010            write!(f, "{tr}")?;
1011        }
1012        Ok(())
1013    }
1014}
1015
1016/// A single tool call captured during a verbose invocation.
1017///
1018/// Records the tool name and its input arguments as a raw JSON value.
1019#[derive(Debug, Clone, Serialize, Deserialize)]
1020#[non_exhaustive]
1021pub struct DebugToolCall {
1022    /// Stable identifier assigned by the model (`tool_use_id`).
1023    ///
1024    /// Used to correlate a call with its subsequent [`DebugToolResult`].
1025    #[serde(default, skip_serializing_if = "Option::is_none")]
1026    pub id: Option<String>,
1027
1028    /// Name of the tool invoked (e.g. `"Read"`, `"Bash"`, `"Grep"`).
1029    pub name: String,
1030
1031    /// Input arguments passed to the tool, as raw JSON.
1032    pub input: Value,
1033}
1034
1035impl fmt::Display for DebugToolCall {
1036    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1037        writeln!(f, "  [tool_use] {} -> {}", self.name, self.input)
1038    }
1039}
1040
1041/// A tool result returned to the model after a tool call.
1042///
1043/// Carries the tool output (any JSON value: string, object, array) and
1044/// an error flag if the tool failed.
1045#[derive(Debug, Clone, Serialize, Deserialize)]
1046#[non_exhaustive]
1047pub struct DebugToolResult {
1048    /// The `tool_use_id` this result answers, matching [`DebugToolCall::id`].
1049    #[serde(default, skip_serializing_if = "Option::is_none")]
1050    pub tool_use_id: Option<String>,
1051
1052    /// Raw content returned by the tool.
1053    pub content: Value,
1054
1055    /// Whether the tool reported an error.
1056    #[serde(default)]
1057    pub is_error: bool,
1058}
1059
1060impl fmt::Display for DebugToolResult {
1061    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1062        let kind = if self.is_error {
1063            "tool_error"
1064        } else {
1065            "tool_result"
1066        };
1067        writeln!(f, "  [{kind}] {}", self.content)
1068    }
1069}
1070
1071impl AgentOutput {
1072    /// Create an `AgentOutput` with the given value and sensible defaults.
1073    pub fn new(value: Value) -> Self {
1074        Self {
1075            value,
1076            session_id: None,
1077            cost_usd: None,
1078            input_tokens: None,
1079            cache_read_input_tokens: None,
1080            cache_creation_input_tokens: None,
1081            output_tokens: None,
1082            model: None,
1083            duration_ms: 0,
1084            debug_messages: None,
1085        }
1086    }
1087}
1088
1089// ── Log sink ──────────────────────────────────────────────────────
1090
1091/// Sink for streaming log lines from provider invocations in real time.
1092///
1093/// Providers that support live log streaming (e.g. K8s ephemeral) call
1094/// [`log`](LogSink::log) for each output line as it is produced, enabling
1095/// downstream consumers (SSE endpoints, log pushers) to display progress
1096/// before the invocation completes.
1097///
1098/// This trait lives in `ironflow-core` so providers can emit logs without
1099/// depending on higher-level crates.
1100///
1101/// # Examples
1102///
1103/// ```
1104/// use std::sync::{Arc, Mutex};
1105/// use ironflow_core::provider::LogSink;
1106///
1107/// struct VecSink(Mutex<Vec<(String, String)>>);
1108///
1109/// impl LogSink for VecSink {
1110///     fn log(&self, stream: &str, line: &str) {
1111///         self.0.lock().unwrap().push((stream.to_string(), line.to_string()));
1112///     }
1113/// }
1114///
1115/// let sink = Arc::new(VecSink(Mutex::new(Vec::new())));
1116/// sink.log("stdout", "hello world");
1117/// assert_eq!(sink.0.lock().unwrap().len(), 1);
1118/// ```
1119pub trait LogSink: Send + Sync {
1120    /// Emit a single log line on the given stream.
1121    ///
1122    /// `stream` is one of `"stdout"`, `"stderr"`, or `"system"`.
1123    /// Implementations should silently drop lines if the receiver is closed.
1124    fn log(&self, stream: &str, line: &str);
1125}
1126
1127// ── Provider trait ─────────────────────────────────────────────────
1128
1129/// Trait for AI agent backends.
1130///
1131/// Implement this trait to provide a custom AI backend for [`Agent`](crate::operations::agent::Agent).
1132/// The only required method is [`invoke`](AgentProvider::invoke), which takes an
1133/// [`AgentConfig`] and returns an [`AgentOutput`] (or an [`AgentError`]).
1134///
1135/// # Examples
1136///
1137/// ```no_run
1138/// use ironflow_core::provider::{AgentConfig, AgentOutput, AgentProvider, InvokeFuture};
1139///
1140/// struct MyProvider;
1141///
1142/// impl AgentProvider for MyProvider {
1143///     fn invoke<'a>(&'a self, config: &'a AgentConfig) -> InvokeFuture<'a> {
1144///         Box::pin(async move {
1145///             // Call your custom backend here...
1146///             todo!()
1147///         })
1148///     }
1149/// }
1150/// ```
1151pub trait AgentProvider: Send + Sync {
1152    /// Execute a single agent invocation with the given configuration.
1153    ///
1154    /// # Errors
1155    ///
1156    /// Returns [`AgentError`] if the underlying backend process fails,
1157    /// times out, or produces output that does not match the requested schema.
1158    fn invoke<'a>(&'a self, config: &'a AgentConfig) -> InvokeFuture<'a>;
1159
1160    /// Execute an agent invocation with real-time log streaming.
1161    ///
1162    /// Providers that support live output streaming should override this
1163    /// method to pipe each output line to the [`LogSink`] as it arrives.
1164    /// The default implementation ignores the sink and delegates to
1165    /// [`invoke`](AgentProvider::invoke).
1166    ///
1167    /// # Errors
1168    ///
1169    /// Returns [`AgentError`] if the underlying backend process fails,
1170    /// times out, or produces output that does not match the requested schema.
1171    fn invoke_with_logs<'a>(
1172        &'a self,
1173        config: &'a AgentConfig,
1174        log_sink: Arc<dyn LogSink>,
1175    ) -> InvokeFuture<'a> {
1176        let _ = log_sink;
1177        self.invoke(config)
1178    }
1179}
1180
1181// The decision abstraction lives beside `AgentProvider`: re-exported here so
1182// `ironflow_core::provider::DecisionProvider` resolves alongside it, while the
1183// types themselves live in the `decision` module.
1184pub use crate::decision::{
1185    ChoiceAnswer, DecideFuture, DecisionAnswer, DecisionOutput, DecisionProvider, DecisionQuestion,
1186    DecisionRequest, DecisionUsage, NoulAnswer, NoulCriteria, ScoreAnswer,
1187};
1188
1189#[cfg(test)]
1190mod tests {
1191    use super::*;
1192    use serde_json::json;
1193
1194    fn full_config() -> AgentConfig {
1195        AgentConfig {
1196            system_prompt: Some("you are helpful".to_string()),
1197            prompt: "do stuff".to_string(),
1198            model: Model::OPUS.to_string(),
1199            allowed_tools: vec!["Read".to_string(), "Write".to_string()],
1200            disallowed_tools: vec!["Bash".to_string()],
1201            tool_profile: None,
1202            max_turns: Some(10),
1203            max_parallel_tools: 2,
1204            max_budget_usd: Some(2.5),
1205            working_dir: Some("/tmp".to_string()),
1206            mcp_config: Some("{}".to_string()),
1207            strict_mcp_config: true,
1208            bare: true,
1209            permission_mode: PermissionMode::Auto,
1210            json_schema: Some(r#"{"type":"object"}"#.to_string()),
1211
1212            resume_session_id: None,
1213            verbose: false,
1214            pod_labels: BTreeMap::new(),
1215            inputs: Vec::new(),
1216            allow_failure: false,
1217            retry: None,
1218            trace_context: None,
1219            _marker: PhantomData,
1220        }
1221    }
1222
1223    #[test]
1224    fn agent_config_serialize_deserialize_roundtrip() {
1225        let config = full_config();
1226        let json = serde_json::to_string(&config).unwrap();
1227        let back: AgentConfig = serde_json::from_str(&json).unwrap();
1228
1229        assert_eq!(back.system_prompt, Some("you are helpful".to_string()));
1230        assert_eq!(back.prompt, "do stuff");
1231        assert_eq!(back.allowed_tools, vec!["Read", "Write"]);
1232        assert_eq!(back.max_turns, Some(10));
1233        assert_eq!(back.max_parallel_tools, 2);
1234        assert_eq!(back.max_budget_usd, Some(2.5));
1235        assert_eq!(back.working_dir, Some("/tmp".to_string()));
1236        assert_eq!(back.mcp_config, Some("{}".to_string()));
1237        assert_eq!(back.json_schema, Some(r#"{"type":"object"}"#.to_string()));
1238    }
1239
1240    #[test]
1241    fn agent_config_with_all_optional_fields_none() {
1242        let config: AgentConfig = AgentConfig {
1243            system_prompt: None,
1244            prompt: "hello".to_string(),
1245            model: Model::HAIKU.to_string(),
1246            allowed_tools: vec![],
1247            disallowed_tools: vec![],
1248            tool_profile: None,
1249            max_turns: None,
1250            max_parallel_tools: 4,
1251            max_budget_usd: None,
1252            working_dir: None,
1253            mcp_config: None,
1254            strict_mcp_config: false,
1255            bare: false,
1256            permission_mode: PermissionMode::Default,
1257            json_schema: None,
1258
1259            resume_session_id: None,
1260            verbose: false,
1261            pod_labels: BTreeMap::new(),
1262            inputs: Vec::new(),
1263            allow_failure: false,
1264            retry: None,
1265            trace_context: None,
1266            _marker: PhantomData,
1267        };
1268        let json = serde_json::to_string(&config).unwrap();
1269        let back: AgentConfig = serde_json::from_str(&json).unwrap();
1270
1271        assert_eq!(back.system_prompt, None);
1272        assert_eq!(back.prompt, "hello");
1273        assert!(back.allowed_tools.is_empty());
1274        assert_eq!(back.max_turns, None);
1275        assert_eq!(back.max_budget_usd, None);
1276        assert_eq!(back.working_dir, None);
1277        assert_eq!(back.mcp_config, None);
1278        assert_eq!(back.json_schema, None);
1279    }
1280
1281    #[test]
1282    fn agent_output_serialize_deserialize_roundtrip() {
1283        let output = AgentOutput {
1284            value: json!({"key": "value"}),
1285            session_id: Some("sess-abc".to_string()),
1286            cost_usd: Some(0.01),
1287            input_tokens: Some(500),
1288            cache_read_input_tokens: Some(4000),
1289            cache_creation_input_tokens: Some(120),
1290            output_tokens: Some(200),
1291            model: Some("claude-sonnet".to_string()),
1292            duration_ms: 3000,
1293            debug_messages: None,
1294        };
1295        let json = serde_json::to_string(&output).unwrap();
1296        let back: AgentOutput = serde_json::from_str(&json).unwrap();
1297
1298        assert_eq!(back.value, json!({"key": "value"}));
1299        assert_eq!(back.session_id, Some("sess-abc".to_string()));
1300        assert_eq!(back.cost_usd, Some(0.01));
1301        assert_eq!(back.input_tokens, Some(500));
1302        assert_eq!(back.cache_read_input_tokens, Some(4000));
1303        assert_eq!(back.cache_creation_input_tokens, Some(120));
1304        assert_eq!(back.output_tokens, Some(200));
1305        assert_eq!(back.model, Some("claude-sonnet".to_string()));
1306        assert_eq!(back.duration_ms, 3000);
1307    }
1308
1309    #[test]
1310    fn agent_output_deserializes_without_cache_fields() {
1311        let raw = json!({
1312            "value": "ok",
1313            "session_id": null,
1314            "cost_usd": 0.01,
1315            "input_tokens": 10,
1316            "output_tokens": 5,
1317            "model": null,
1318            "duration_ms": 100,
1319            "debug_messages": null
1320        });
1321        let back: AgentOutput = serde_json::from_value(raw).unwrap();
1322        assert_eq!(back.input_tokens, Some(10));
1323        assert_eq!(back.cache_read_input_tokens, None);
1324        assert_eq!(back.cache_creation_input_tokens, None);
1325    }
1326
1327    #[test]
1328    fn agent_config_new_has_correct_defaults() {
1329        let config = AgentConfig::new("test prompt");
1330        assert_eq!(config.prompt, "test prompt");
1331        assert_eq!(config.system_prompt, None);
1332        assert_eq!(config.model, Model::SONNET);
1333        assert!(config.allowed_tools.is_empty());
1334        assert_eq!(config.max_turns, None);
1335        assert_eq!(config.max_budget_usd, None);
1336        assert_eq!(config.working_dir, None);
1337        assert_eq!(config.mcp_config, None);
1338        assert!(matches!(config.permission_mode, PermissionMode::Default));
1339        assert_eq!(config.json_schema, None);
1340        assert_eq!(config.resume_session_id, None);
1341        assert!(!config.verbose);
1342    }
1343
1344    #[test]
1345    fn agent_output_new_has_correct_defaults() {
1346        let output = AgentOutput::new(json!("test"));
1347        assert_eq!(output.value, json!("test"));
1348        assert_eq!(output.session_id, None);
1349        assert_eq!(output.cost_usd, None);
1350        assert_eq!(output.input_tokens, None);
1351        assert_eq!(output.cache_read_input_tokens, None);
1352        assert_eq!(output.cache_creation_input_tokens, None);
1353        assert_eq!(output.output_tokens, None);
1354        assert_eq!(output.model, None);
1355        assert_eq!(output.duration_ms, 0);
1356        assert!(output.debug_messages.is_none());
1357    }
1358
1359    #[test]
1360    fn agent_config_resume_session_roundtrip() {
1361        let mut config = AgentConfig::new("test");
1362        config.resume_session_id = Some("sess-xyz".to_string());
1363        let json = serde_json::to_string(&config).unwrap();
1364        let back: AgentConfig = serde_json::from_str(&json).unwrap();
1365        assert_eq!(back.resume_session_id, Some("sess-xyz".to_string()));
1366    }
1367
1368    #[test]
1369    fn agent_output_debug_does_not_panic() {
1370        let output = AgentOutput {
1371            value: json!(null),
1372            session_id: None,
1373            cost_usd: None,
1374            input_tokens: None,
1375            cache_read_input_tokens: None,
1376            cache_creation_input_tokens: None,
1377            output_tokens: None,
1378            model: None,
1379            duration_ms: 0,
1380            debug_messages: None,
1381        };
1382        let debug_str = format!("{:?}", output);
1383        assert!(!debug_str.is_empty());
1384    }
1385
1386    #[test]
1387    fn allow_tool_transitions_to_with_tools() {
1388        let config = AgentConfig::new("test").allow_tool(Tool::Read);
1389        assert_eq!(config.allowed_tools, vec!["Read"]);
1390
1391        // Can add more tools, known or custom.
1392        let config = config
1393            .allow_tool(Tool::Write)
1394            .allow_tool(Tool::Custom("mcp__github__search".to_string()));
1395        assert_eq!(
1396            config.allowed_tools,
1397            vec!["Read", "Write", "mcp__github__search"]
1398        );
1399    }
1400
1401    #[test]
1402    fn output_carries_the_output_type_in_the_typestate() {
1403        #[derive(serde::Deserialize, JsonSchema)]
1404        #[allow(dead_code)]
1405        struct Verdict {
1406            approved: bool,
1407        }
1408
1409        let config: AgentConfig<NoTools, WithSchema<Verdict>> =
1410            AgentConfig::new("review").output::<Verdict>();
1411        assert!(
1412            config
1413                .json_schema
1414                .as_deref()
1415                .is_some_and(|s| s.contains("approved"))
1416        );
1417    }
1418
1419    #[test]
1420    fn output_schema_raw_transitions_to_with_schema() {
1421        let config = AgentConfig::new("test").output_schema_raw(r#"{"type":"object"}"#);
1422        assert_eq!(config.json_schema.as_deref(), Some(r#"{"type":"object"}"#));
1423    }
1424
1425    #[test]
1426    fn with_tools_converts_to_base_type() {
1427        let typed = AgentConfig::new("test").allow_tool(Tool::Read);
1428        let base: AgentConfig = typed.into();
1429        assert_eq!(base.allowed_tools, vec!["Read"]);
1430    }
1431
1432    #[test]
1433    fn with_schema_converts_to_base_type() {
1434        let typed = AgentConfig::new("test").output_schema_raw(r#"{"type":"object"}"#);
1435        let base: AgentConfig = typed.into();
1436        assert_eq!(base.json_schema.as_deref(), Some(r#"{"type":"object"}"#));
1437    }
1438
1439    #[test]
1440    fn serde_roundtrip_ignores_marker() {
1441        let config = AgentConfig::new("test").allow_tool(Tool::Read);
1442        let json = serde_json::to_string(&config).unwrap();
1443        assert!(!json.contains("marker"));
1444
1445        let back: AgentConfig = serde_json::from_str(&json).unwrap();
1446        assert_eq!(back.allowed_tools, vec!["Read"]);
1447    }
1448
1449    #[test]
1450    fn bare_defaults_to_false() {
1451        let config = AgentConfig::new("hello");
1452        assert!(!config.bare, "bare must default to false");
1453    }
1454
1455    #[test]
1456    fn bare_builder_sets_flag() {
1457        let config = AgentConfig::new("hello").bare(true);
1458        assert!(config.bare, "bare(true) must enable the flag");
1459
1460        let config = config.bare(false);
1461        assert!(!config.bare, "bare(false) must disable the flag");
1462    }
1463
1464    #[test]
1465    fn bare_serde_default_when_missing() {
1466        let raw = r#"{"prompt":"hello","model":"sonnet"}"#;
1467        let config: AgentConfig = serde_json::from_str(raw).unwrap();
1468        assert!(
1469            !config.bare,
1470            "bare must default to false when absent from serialized payload"
1471        );
1472    }
1473
1474    #[test]
1475    fn bare_serde_roundtrip() {
1476        let mut config = AgentConfig::new("hello");
1477        config.bare = true;
1478        let json = serde_json::to_string(&config).unwrap();
1479        assert!(
1480            json.contains("\"bare\":true"),
1481            "serialized form must contain bare:true, got: {json}"
1482        );
1483
1484        let back: AgentConfig = serde_json::from_str(&json).unwrap();
1485        assert!(back.bare, "bare must survive a serde roundtrip");
1486    }
1487
1488    #[test]
1489    fn disallowed_tools_defaults_to_empty() {
1490        let config = AgentConfig::new("hello");
1491        assert!(
1492            config.disallowed_tools.is_empty(),
1493            "disallowed_tools must default to empty"
1494        );
1495    }
1496
1497    #[test]
1498    fn disallowed_tools_builder_replaces_list() {
1499        let config = AgentConfig::new("hello").disallowed_tools([Tool::Write, Tool::Edit]);
1500        assert_eq!(config.disallowed_tools, vec!["Write", "Edit"]);
1501
1502        // Subsequent call fully replaces the list.
1503        let config = config.disallowed_tools([Tool::Bash]);
1504        assert_eq!(config.disallowed_tools, vec!["Bash"]);
1505
1506        // Empty input clears the list.
1507        let config = config.disallowed_tools([]);
1508        assert!(config.disallowed_tools.is_empty());
1509    }
1510
1511    #[test]
1512    fn disallowed_tools_compatible_with_output() {
1513        #[derive(serde::Deserialize, JsonSchema)]
1514        #[allow(dead_code)]
1515        struct Out {
1516            ok: bool,
1517        }
1518
1519        // Typestate compile check: .disallowed_tools(...) must be callable
1520        // before AND after .output::<T>() because it lives on
1521        // impl<Tools, Schema>, not impl<Tools, NoSchema>.
1522        let before: AgentConfig<NoTools, WithSchema<Out>> = AgentConfig::new("classify")
1523            .disallowed_tools([Tool::Write, Tool::Edit])
1524            .output::<Out>();
1525        assert_eq!(before.disallowed_tools, vec!["Write", "Edit"]);
1526        assert!(before.json_schema.is_some());
1527
1528        let after: AgentConfig<NoTools, WithSchema<Out>> = AgentConfig::new("classify")
1529            .output::<Out>()
1530            .disallowed_tools([Tool::Write]);
1531        assert_eq!(after.disallowed_tools, vec!["Write"]);
1532        assert!(after.json_schema.is_some());
1533    }
1534
1535    #[test]
1536    fn disallowed_tools_serde_default_when_missing() {
1537        let raw = r#"{"prompt":"hello","model":"sonnet"}"#;
1538        let config: AgentConfig = serde_json::from_str(raw).unwrap();
1539        assert!(
1540            config.disallowed_tools.is_empty(),
1541            "disallowed_tools must default to empty when absent from serialized payload"
1542        );
1543    }
1544
1545    #[test]
1546    fn disallowed_tools_serde_roundtrip() {
1547        let config = AgentConfig::new("hello").disallowed_tools([Tool::Write, Tool::Edit]);
1548        let json = serde_json::to_string(&config).unwrap();
1549        assert!(
1550            json.contains("\"disallowed_tools\":[\"Write\",\"Edit\"]"),
1551            "serialized form must contain the disallowed_tools array, got: {json}"
1552        );
1553
1554        let back: AgentConfig = serde_json::from_str(&json).unwrap();
1555        assert_eq!(back.disallowed_tools, vec!["Write", "Edit"]);
1556    }
1557
1558    #[test]
1559    fn pod_labels_defaults_to_empty() {
1560        let config = AgentConfig::new("test");
1561        assert!(config.pod_labels.is_empty());
1562    }
1563
1564    #[test]
1565    fn pod_label_builder_adds_entry() {
1566        let config = AgentConfig::new("test").pod_label("k", "v");
1567        assert_eq!(config.pod_labels.len(), 1);
1568        assert_eq!(config.pod_labels["k"], "v");
1569    }
1570
1571    #[test]
1572    fn pod_labels_builder_replaces_map() {
1573        let config = AgentConfig::new("test").pod_label("old", "value");
1574        let mut new_map = BTreeMap::new();
1575        new_map.insert("new".to_string(), "value".to_string());
1576        let config = config.pod_labels(new_map);
1577        assert_eq!(config.pod_labels.len(), 1);
1578        assert_eq!(config.pod_labels["new"], "value");
1579        assert!(!config.pod_labels.contains_key("old"));
1580    }
1581
1582    #[test]
1583    fn pod_labels_serde_default_when_missing() {
1584        let raw = r#"{"prompt":"hello","model":"sonnet"}"#;
1585        let config: AgentConfig = serde_json::from_str(raw).unwrap();
1586        assert!(
1587            config.pod_labels.is_empty(),
1588            "pod_labels must default to empty when absent from serialized payload"
1589        );
1590    }
1591
1592    #[test]
1593    fn pod_labels_serde_skip_when_empty() {
1594        let config = AgentConfig::new("hello");
1595        let json = serde_json::to_string(&config).unwrap();
1596        assert!(
1597            !json.contains("pod_labels"),
1598            "empty pod_labels must be skipped during serialization, got: {json}"
1599        );
1600    }
1601
1602    #[test]
1603    fn pod_labels_serde_roundtrip() {
1604        let config = AgentConfig::new("hello")
1605            .pod_label("ironflow.io/network-profile", "grafana-only")
1606            .pod_label("team", "observability");
1607        let json = serde_json::to_string(&config).unwrap();
1608        assert!(
1609            json.contains("pod_labels"),
1610            "non-empty pod_labels must be present in serialized form, got: {json}"
1611        );
1612
1613        let back: AgentConfig = serde_json::from_str(&json).unwrap();
1614        assert_eq!(back.pod_labels.len(), 2);
1615        assert_eq!(
1616            back.pod_labels["ironflow.io/network-profile"],
1617            "grafana-only"
1618        );
1619        assert_eq!(back.pod_labels["team"], "observability");
1620    }
1621
1622    // ── LogSink tests ─────────────────────────────────────────────
1623
1624    use crate::test_support::VecSink;
1625
1626    #[test]
1627    fn log_sink_collects_lines() {
1628        let sink = VecSink::new();
1629        sink.log("stdout", "line 1");
1630        sink.log("stderr", "err!");
1631        sink.log("system", "done");
1632
1633        let lines = sink.0.lock().unwrap();
1634        assert_eq!(lines.len(), 3);
1635        assert_eq!(lines[0], ("stdout".to_string(), "line 1".to_string()));
1636        assert_eq!(lines[1], ("stderr".to_string(), "err!".to_string()));
1637        assert_eq!(lines[2], ("system".to_string(), "done".to_string()));
1638    }
1639
1640    #[test]
1641    fn log_sink_arc_is_clone_and_send() {
1642        let sink: Arc<dyn LogSink> = VecSink::new();
1643        let cloned = sink.clone();
1644        sink.log("stdout", "from original");
1645        cloned.log("stdout", "from clone");
1646    }
1647
1648    // ── invoke_with_logs default impl ─────────────────────────────
1649
1650    struct FixedProvider {
1651        output: AgentOutput,
1652    }
1653
1654    impl AgentProvider for FixedProvider {
1655        fn invoke<'a>(&'a self, _config: &'a AgentConfig) -> InvokeFuture<'a> {
1656            Box::pin(async {
1657                Ok(AgentOutput {
1658                    value: self.output.value.clone(),
1659                    session_id: self.output.session_id.clone(),
1660                    cost_usd: self.output.cost_usd,
1661                    input_tokens: self.output.input_tokens,
1662                    cache_read_input_tokens: self.output.cache_read_input_tokens,
1663                    cache_creation_input_tokens: self.output.cache_creation_input_tokens,
1664                    output_tokens: self.output.output_tokens,
1665                    model: self.output.model.clone(),
1666                    duration_ms: self.output.duration_ms,
1667                    debug_messages: None,
1668                })
1669            })
1670        }
1671    }
1672
1673    #[tokio::test]
1674    async fn invoke_with_logs_default_delegates_to_invoke() {
1675        let provider = FixedProvider {
1676            output: AgentOutput::new(json!("ok")),
1677        };
1678        let config = AgentConfig::new("test");
1679        let sink: Arc<dyn LogSink> = VecSink::new();
1680
1681        let result = provider.invoke_with_logs(&config, sink.clone()).await;
1682        assert!(result.is_ok());
1683        assert_eq!(result.unwrap().value, json!("ok"));
1684    }
1685
1686    #[tokio::test]
1687    async fn invoke_with_logs_default_ignores_sink() {
1688        let provider = FixedProvider {
1689            output: AgentOutput::new(json!("ok")),
1690        };
1691        let config = AgentConfig::new("test");
1692        let sink = VecSink::new();
1693
1694        let _ = provider
1695            .invoke_with_logs(&config, sink.clone() as Arc<dyn LogSink>)
1696            .await;
1697
1698        let lines = sink.0.lock().unwrap();
1699        assert!(lines.is_empty(), "default impl should not emit any logs");
1700    }
1701
1702    #[test]
1703    fn max_parallel_tools_defaults_to_four() {
1704        assert_eq!(AgentConfig::new("hi").max_parallel_tools, 4);
1705    }
1706
1707    #[test]
1708    fn max_parallel_tools_builder_sets_value() {
1709        let config = AgentConfig::new("hi").max_parallel_tools(2);
1710        assert_eq!(config.max_parallel_tools, 2);
1711    }
1712
1713    #[test]
1714    #[should_panic(expected = "max_parallel_tools must be greater than 0")]
1715    fn max_parallel_tools_zero_panics() {
1716        let _ = AgentConfig::new("hi").max_parallel_tools(0);
1717    }
1718
1719    #[test]
1720    fn max_parallel_tools_missing_from_json_defaults_to_four() {
1721        let json = serde_json::to_value(AgentConfig::new("hi")).unwrap();
1722        let mut obj = json.as_object().unwrap().clone();
1723        obj.remove("max_parallel_tools");
1724        let back: AgentConfig = serde_json::from_value(Value::Object(obj)).unwrap();
1725        assert_eq!(back.max_parallel_tools, 4);
1726    }
1727}