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