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