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