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