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