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