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