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