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