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