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