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