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