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