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