Skip to main content

ironflow_core/
provider.rs

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