Forge documentation
Library referenceRust

forge-agent

ANVIL-compliant agent execution loop, workflows, and multi-agent orchestration for the Forge SDK

ANVIL-compliant agent execution loop, workflows, and multi-agent orchestration for the Forge SDK

Package contract

FieldValue
Languagerust
Source version0.2.0
Manifestforge-rs/crates/forge-agent/Cargo.toml
Source files12
EvidenceSource reference; registry publication and runtime conformance are separate checks

Import boundary

use forge_agent;

Use a source checkout or your verified private registry. Manifest coordinates identify the package; they do not establish that a public registry release exists.

Crate boundary

The following entries are taken from src/lib.rs. Feature conditions in the exact source still apply.

pub mod agent;

pub mod cambium;

pub mod context_manager;

pub mod error;

pub mod loop_control;

#[cfg(not(target_arch = "wasm32"))]
pub mod messaging;

pub mod observer;

pub mod streaming_tool_loop;

pub mod subagent;

pub mod tool_loop;

pub mod workflow;

pub mod prelude;

#[cfg(not(target_arch = "wasm32"))]
pub use crate::agent::AegisProviderAuthorityVerifier;

pub use crate::agent::{
        Agent, AgentCheckpointValidationError, AgentConfig, AgentEvent, AgentOutput,
        AgentRunCheckpoint, BoundaryContract, CodingProviderPreflight, ModelInferenceRecord,
        ProviderAuthorityVerifier, ProviderIdentityVerification, SubAgentDelegationRecord,
        ToolInvocationRecord, ToolInvocationStatus, AGENT_RUN_CHECKPOINT_SCHEMA,
    };

pub use crate::cambium::{
        CambiumAgentLoopObserver, CambiumEvent, CambiumJson, CambiumScope,
        ForgeCambiumObserverConfig,
    };

pub use crate::context_manager::{ContextWindowConfig, ContextWindowManager, PruningStrategy};

pub use crate::error::{ForgeAgentError, ForgeAgentResult};

pub use crate::loop_control::{
        AgentStopCondition, NoOpPrepare, PrepareStep, StopWhen, StopWhenCustom, StopWhenMaxSteps,
        StopWhenTextGenerated, StopWhenToolCalled,
    };

#[cfg(not(target_arch = "wasm32"))]
pub use crate::messaging::{
        AgentChannel, AgentChannelReceiver, AgentChannelSender, AgentMessage,
    };

pub use crate::observer::{AgentLoopObserver, NoOpObserver};

pub use crate::streaming_tool_loop::{StreamingLoopConfig, StreamingToolLoopAgent};

pub use crate::subagent::{
        create_subagent, create_subagent_with_observer, SubAgentConfig, SubAgentDelegationContext,
    };

pub use crate::tool_loop::ToolLoopAgent;

pub use crate::workflow::{
        ParallelWorkflow, RouterWorkflow, SequentialWorkflow, WorkflowOutput,
    };

Source reference

Download package reference JSON. Each original source file and generated declaration artifact has its own SHA-256 digest. Function bodies and constant values are omitted from downloads. These are source declaration inventories, not compiler-resolved rustdoc, TypeDoc, DocC, or Dokka output. Private modules can contain public declarations that are not reachable through the package boundary; consult the entry point before importing.

agent.rs

Read declaration text · 51 declaration entries

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
pub struct BoundaryContract {

}

pub fn new() -> Self;

pub fn with_network_access(mut self) -> Self;

pub fn with_delegation_access(mut self) -> Self;

pub fn allow_provider_namespace(mut self, namespace: impl Into<String>) -> Self;

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ProviderIdentityVerification {
/// The DID that was verified.

pub did: String,
/// Whether the DID document's signature validated successfully.

pub signature_valid: bool,
/// Whether the DID's lineage chain validated successfully.

pub lineage_valid: bool,
/// Forge-owned conformance level derived from the upstream verification pipeline.

pub conformance_level: u8
}

#[async_trait]
pub trait ProviderAuthorityVerifier: Send + Sync {
    /// Verifies the given DID and returns the normalized authority result.
    async fn verify_did(&self, did: &str) -> Result<ProviderIdentityVerification, ForgeAgentError>;
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CodingProviderPreflight {
/// The negotiated provider contract.

pub negotiation: ProviderNegotiationResult,
/// The DID that passed execution-authority verification.

pub verified_did: String,
/// The provider scopes that were validated against the active ACT.

pub required_scopes: Vec<String>
}

#[cfg(not(target_arch = "wasm32"))]
pub struct AegisProviderAuthorityVerifier {

}

#[cfg(not(target_arch = "wasm32"))]
pub fn new(
        registry: std::sync::Arc<aegis_core::PluginRegistry>,
        config: aegis_core::VerificationConfig,
    ) -> Self;

#[cfg(not(target_arch = "wasm32"))]
pub fn from_pipeline(pipeline: aegis_verify::VerificationPipeline) -> Self;

#[derive(Clone, Serialize, Deserialize)]
pub struct AgentConfig {

}

pub fn new(name: impl Into<String>, model: impl Into<String>) -> Self;

pub fn with_system_prompt(mut self, prompt: impl Into<String>) -> Self;

pub fn with_max_steps(mut self, max: u32) -> Self;

pub fn with_tool(mut self, tool: ToolDefinition) -> Self;

pub fn with_tools(mut self, tools: Vec<ToolDefinition>) -> Self;

pub fn with_tool_registry(mut self, registry: &forge_tool::registry::ToolRegistry) -> Self;

pub fn with_boundary_contract(mut self, boundary_contract: BoundaryContract) -> Self;

pub fn with_identity(mut self, identity: ForgeAgentIdentity) -> Self;

pub fn with_act(mut self, act: AgentCapabilityToken) -> Self;

pub fn name(&self) -> &str;

pub fn model(&self) -> &ProviderRef;

pub fn tools(&self) -> &[ToolDefinition];

pub fn max_steps(&self) -> u32;

pub fn system_prompt(&self) -> Option<&str>;

pub fn boundary_contract(&self) -> Option<&BoundaryContract>;

pub fn identity(&self) -> Option<&ForgeAgentIdentity>;

pub fn act(&self) -> Option<&AgentCapabilityToken>;

pub fn agent_did(&self) -> Option<&str>;

pub async fn preflight_coding_provider_execution<V: ProviderAuthorityVerifier>(
        &self,
        registry: &ProviderRegistry,
        request: &ProviderNegotiationRequest,
        verifier: &V,
    ) -> Result<CodingProviderPreflight, ForgeAgentError>;

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ToolInvocationRecord {
/// The tool call id that the model assigned.

pub id: String,
/// The tool's registered name.

pub name: String,
/// Wall-clock timestamp at which the agent began executing the tool.

///

/// Serialized as RFC 3339. Set from `chrono::Utc::now()` immediately

/// before the executor (or authorization gate / lookup) is invoked.

#[serde(with = "tool_invocation_started_at_serde")]
pub started_at: chrono::DateTime<chrono::Utc>,
/// Time elapsed between the start and the resolution of the invocation.

///

/// Captures the full latency: authorization check + executor + approval

/// handler. Serialized as a struct (`{ "secs": u64, "nanos": u32 }`).

pub duration: std::time::Duration,
/// Final status of the invocation.

pub status: ToolInvocationStatus
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ModelInferenceRecord {
/// Host or loop assigned id for this model call.

pub id: String,
/// Model identifier selected for the call.

pub model: String,
/// Provider identifier or namespace that served the call.

pub provider_id: String,
/// Optional model-router or Foundry route id.

pub route_id: Option<String>,
/// Governed artifact ref for the prompt payload.

pub prompt_ref: Option<String>,
/// Governed artifact ref for the completion payload.

pub completion_ref: Option<String>,
/// Prompt/input token count, when reported.

pub prompt_tokens: Option<u64>,
/// Completion/output token count, when reported.

pub completion_tokens: Option<u64>,
/// Wall-clock timestamp at which the model call began.

#[serde(with = "tool_invocation_started_at_serde")]
pub started_at: chrono::DateTime<chrono::Utc>,
/// Time elapsed between model call start and stream completion.

pub duration: std::time::Duration,
/// Canonical Cambium outcome label, such as `success` or `failed`.

pub outcome: String
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SubAgentDelegationRecord {
/// OAS/stable child agent id.

pub subagent_id: String,
/// Host-assigned run id for the delegated sub-agent.

pub subagent_run_id: String,
/// Optional summary of why the delegation happened.

pub delegation_reason: Option<String>,
/// Governed artifact ref for the delegated instruction.

pub instruction_ref: Option<String>,
/// Capability refs granted to the child.

pub capability_refs: Vec<String>,
/// Governed artifact ref for handoff state.

pub handoff_ref: Option<String>,
/// Wall-clock timestamp at which delegation started.

#[serde(with = "tool_invocation_started_at_serde")]
pub started_at: chrono::DateTime<chrono::Utc>
}

#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AgentEvent {
    /// A user message was added to the conversation.
    UserMessage(String),
    /// A system message was added to the conversation.
    SystemMessage(String),
    /// The assistant produced a text turn (no tool calls).
    AssistantText(String),
    /// The assistant called a tool. Pairs with a later [`AgentEvent::ToolResult`]
    /// (matched on `id`) and, if the run captured timing,
    /// [`AgentEvent::ToolInvocationCompleted`].
    AssistantToolCall {
        id: String,
        name: String,
        arguments: serde_json::Value,
    },
    /// The tool returned a result (already in the conversation as a `Role::Tool`
    /// message). `is_error` reflects the tool result's own flag.
    ToolResult {
        id: String,
        name: String,
        content: String,
        is_error: bool,
    },
    /// A timed tool invocation completed (lifted from
    /// [`AgentOutput::tool_invocations`]). Carries the typed status, the
    /// captured duration, and the wall-clock start.
    ToolInvocationCompleted(ToolInvocationRecord),
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ToolInvocationStatus {
    /// The executor returned `Ok` and the tool result is not flagged as an error.
    Succeeded,
    /// The executor returned `Ok` but the tool result is flagged as an error.
    /// Distinguished from `Failed` so callers can see "the tool ran fine but
    /// the model's request was wrong" vs. "the tool itself blew up."
    SucceededWithToolError,
    /// The tool name was not present in the registry.
    NotFound,
    /// Authorization (ACT scope check) denied the invocation.
    AuthorizationDenied,
    /// The executor returned `Err`.
    Failed,
}

pub fn serialize<S>(value: &DateTime<Utc>, serializer: S) -> Result<S::Ok, S::Error>
    where
        S: Serializer,;

pub fn deserialize<'de, D>(deserializer: D) -> Result<DateTime<Utc>, D::Error>
    where
        D: Deserializer<'de>,;

pub const AGENT_RUN_CHECKPOINT_SCHEMA: &str;

#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AgentRunCheckpoint {
/// Stable schema id. Must be [`AGENT_RUN_CHECKPOINT_SCHEMA`].

pub schema: String,
/// Host-assigned run identifier.

pub run_id: String,
/// Agent name from the originating [`AgentConfig`].

pub agent_name: String,
/// Provider model reference from the originating [`AgentConfig`].

pub model: String,
/// Number of tool-loop steps already consumed.

pub step_count: u32,
/// Maximum steps allowed for the originating run.

pub max_steps: u32,
/// Conversation state to pass back into [`Agent::run_with_messages`].

pub messages: Vec<ModelMessage>,
/// Usage accumulated before the checkpoint was written.

pub usage: Usage,
/// Final or latest assistant text available when the checkpoint was written.

pub final_text: String,
/// Tool invocation audit records accumulated before the checkpoint was

/// written.

#[serde(default)]
pub tool_invocations: Vec<ToolInvocationRecord>
}

pub fn validate_for(&self, config: &AgentConfig) -> Result<(), AgentCheckpointValidationError>;

pub fn resume_messages(&self) -> &[ModelMessage];

pub fn remaining_steps(&self) -> u32;

pub fn resume_config_for(
        &self,
        config: &AgentConfig,
    ) -> Result<AgentConfig, AgentCheckpointValidationError>;

pub fn digest_sha256(&self) -> Result<String, serde_json::Error>;

pub fn verify_digest_sha256(
        &self,
        expected: &str,
    ) -> Result<(), AgentCheckpointValidationError>;

#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AgentCheckpointValidationError {
    /// Checkpoint schema is not the supported Forge checkpoint schema.
    SchemaMismatch {
        /// Expected schema id.
        expected: String,
        /// Schema id found in the checkpoint.
        found: String,
    },
    /// Checkpoint run id is empty.
    EmptyRunId,
    /// Checkpoint belongs to a different agent name.
    AgentNameMismatch {
        /// Expected agent name.
        expected: String,
        /// Agent name found in the checkpoint.
        found: String,
    },
    /// Checkpoint belongs to a different model reference.
    ModelMismatch {
        /// Expected provider model reference.
        expected: String,
        /// Model reference found in the checkpoint.
        found: String,
    },
    /// Checkpoint max-step budget differs from the resume config.
    MaxStepsMismatch {
        /// Expected max-step budget.
        expected: u32,
        /// Max-step budget found in the checkpoint.
        found: u32,
    },
    /// Checkpoint has already consumed more steps than the budget allows.
    StepCountExceedsMax {
        /// Steps already consumed.
        step_count: u32,
        /// Maximum allowed steps.
        max_steps: u32,
    },
    /// Checkpoint has no remaining steps for a resumed run.
    StepBudgetExhausted {
        /// Original maximum allowed steps.
        max_steps: u32,
    },
    /// Expected checkpoint digest is not a SHA-256 hex string.
    InvalidDigest {
        /// Digest value that failed shape validation.
        found: String,
    },
    /// Checkpoint could not be serialized for digest verification.
    DigestSerializationFailed {
        /// Serialization failure message.
        message: String,
    },
    /// Checkpoint digest does not match expected evidence.
    DigestMismatch {
        /// Expected checkpoint digest.
        expected: String,
        /// Actual checkpoint digest.
        found: String,
    },
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentOutput {
/// The complete conversation history including all user, assistant,

/// and tool messages exchanged during execution.

pub messages: Vec<ModelMessage>,
/// Cumulative token usage across all LLM calls during execution.

pub usage: Usage,
/// The number of tool loop steps taken (each LLM call counts as one step).

pub steps_taken: u32,
/// The final text response produced by the agent.

///

/// This is the text content of the last assistant message that did not

/// contain tool calls (i.e., the terminal response).

pub final_text: String,
/// Per-tool invocation records collected during the run.

///

/// One entry per tool call attempted, in the order they were executed.

/// Includes name, wall-clock start, duration, and final status (succeeded /

/// failed / authorization-denied / not-found). Backfilled to an empty

/// `Vec` when constructing `AgentOutput` from older code paths or when

/// deserializing payloads that pre-date this field.

#[serde(default)]
pub tool_invocations: Vec<ToolInvocationRecord>
}

pub fn to_checkpoint(
        &self,
        run_id: impl Into<String>,
        config: &AgentConfig,
    ) -> AgentRunCheckpoint;

pub fn events(&self) -> impl Iterator<Item = AgentEvent> + '_;

#[async_trait]
pub trait Agent: Send + Sync {
    /// Returns the agent's configuration.
    ///
    /// # Returns
    ///
    /// A reference to the [`AgentConfig`] used to create this agent.
    fn config(&self) -> &AgentConfig;

    /// Returns a snapshot of the agent's current health profile.
    ///
    /// The health profile contains runtime metrics: uptime, error counts,
    /// tool invocations, inference statistics, and resource usage. Returns
    /// a clone because the profile is mutated concurrently during execution
    /// (behind interior mutability).
    ///
    /// # ANVIL Spec SS14.1
    ///
    /// Health profiles are part of the telemetry contract. They are exposed
    /// for monitoring and audit trail consumption.
    ///
    /// # Returns
    ///
    /// A cloned snapshot of the agent's [`HealthProfile`].
    fn health(&self) -> HealthProfile;

    /// Returns the agent's current lifecycle state.
    ///
    /// # ANVIL Spec SS13.2
    ///
    /// The lifecycle state machine has six states: Initializing, Ready,
    /// Running, Paused, Error, Terminated. See the spec for the complete
    /// transition table.
    ///
    /// # Returns
    ///
    /// The current [`LifecycleState`].
    fn lifecycle(&self) -> LifecycleState;

    /// Runs the agent with a text prompt.
    ///
    /// This is the primary entry point for agent execution. It creates an
    /// initial user message from the prompt and delegates to
    /// [`run_with_messages`](Self::run_with_messages).
    ///
    /// # ANVIL Spec SS7.1
    ///
    /// The agent execution loop:
    /// 1. Send messages to the LLM.
    /// 2. If the response contains tool calls, execute them and loop.
    /// 3. If the response is text only, return the final result.
    /// 4. Check stop conditions after each step.
    ///
    /// # Arguments
    ///
    /// * `prompt` - The user's text prompt to process.
    ///
    /// # Returns
    ///
    /// An [`AgentOutput`] containing the final response, conversation history,
    /// usage statistics, and step count.
    ///
    /// # Errors
    ///
    /// * [`ForgeAgentError::ToolLoopFailed`] -- if a tool loop step fails.
    /// * [`ForgeAgentError::LifecycleError`] -- if the agent is not in a runnable state.
    /// * [`ForgeAgentError::StopConditionReached`] -- if max steps or tokens are exceeded.
    /// * [`ForgeAgentError::Generate`] -- if the underlying model call fails.
    /// * [`ForgeAgentError::Tool`] -- if a tool execution fails.
    async fn run(&self, prompt: &str) -> Result<AgentOutput, ForgeAgentError>;

    /// Runs the agent with a pre-constructed message list.
    ///
    /// This is the lower-level entry point that allows callers to provide
    /// the complete conversation history, including system messages, tool
    /// results, and previous exchanges.
    ///
    /// # Arguments
    ///
    /// * `messages` - The initial message list to process.
    ///
    /// # Returns
    ///
    /// An [`AgentOutput`] containing the final response and metadata.
    ///
    /// # Errors
    ///
    /// Same error variants as [`run()`](Self::run).
    async fn run_with_messages(
        &self,
        messages: Vec<ModelMessage>,
    ) -> Result<AgentOutput, ForgeAgentError>;
}

cambium.rs

Read declaration text · 6 declaration entries

pub use cambium_sdk_rs::{CambiumEvent, CambiumJson, CambiumScope};

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ForgeCambiumObserverConfig {
/// Cambium tenant/product scope assigned by the embedding host.

pub scope: CambiumScope,
/// Host-assigned Forge run id.

pub run_id: String,
/// OAS agent id or stable agent identifier for lineage.

pub agent_id: String,
/// Principal that requested or authorized the Forge run.

pub principal_id: String,
/// Authentication authority for the principal, when known.

pub authenticated_by: Option<String>,
/// Stable event id prefix used with each Forge tool call id.

pub event_id_prefix: String,
/// Optional distributed trace id.

pub trace_id: Option<String>,
/// Optional span id paired with `trace_id`.

pub span_id: Option<String>,
/// Initial producer sequence. The first emitted event increments from this.

pub producer_sequence_start: u64
}

pub struct CambiumAgentLoopObserver {

}

pub fn try_new(config: ForgeCambiumObserverConfig) -> Result<Self, String>;

#[must_use]
pub fn events(&self) -> Vec<CambiumEvent>;

#[must_use]
pub fn drain_events(&self) -> Vec<CambiumEvent>;

context_manager.rs

Read declaration text · 7 declaration entries

#[derive(Debug, Clone, Copy)]
pub enum PruningStrategy {
    /// Remove the oldest non-system messages first.
    RemoveOldest,
    /// Keep system prompt + the last N messages (sliding window).
    SlidingWindow {
        /// Maximum number of recent messages to keep (in addition to system).
        keep_last: usize,
    },
    /// Summarize older messages into a single system message.
    /// (Future: will use the model to generate a summary.)
    Summarize,
}

#[derive(Debug, Clone)]
pub struct ContextWindowConfig {
/// Maximum number of tokens the model supports.

pub max_context_tokens: u64,
/// Threshold ratio (0.0–1.0) at which to start pruning.

/// Default: 0.85 (prune when 85% of context is used).

pub threshold_ratio: f64,
/// Estimated tokens per message for the simple estimator.

/// Default: 4 (roughly 4 tokens per word, ~100 words per message = 400).

pub tokens_per_char: f64,
/// Pruning strategy to use when threshold is exceeded.

pub strategy: PruningStrategy
}

pub struct ContextWindowManager {

}

pub fn new(config: ContextWindowConfig) -> Self;

pub fn for_context_size(max_tokens: u64) -> Self;

pub fn estimate_tokens(&self, messages: &[ModelMessage]) -> u64;

pub fn threshold_tokens(&self) -> u64;

error.rs

Read declaration text · 2 declaration entries

#[derive(Debug, Error)]
pub enum ForgeAgentError {
    /// The agent's tool loop failed during execution.
    ///
    /// This error indicates that a step within the agent's tool loop encountered
    /// an unrecoverable failure. The `step` field indicates which iteration
    /// failed, and `reason` explains what went wrong.
    ///
    /// See ANVIL Spec SS7.1 -- Agent Execution Loop.
    #[error("agent '{agent_name}' tool loop failed at step {step}: {reason}")]
    ToolLoopFailed {
        /// The name of the agent whose tool loop failed.
        agent_name: String,
        /// The step number (1-indexed) at which the failure occurred.
        step: u32,
        /// A human-readable description of what went wrong.
        reason: String,
    },

    /// A workflow orchestration step failed.
    ///
    /// This error wraps failures that occur during sequential, parallel, or
    /// router workflow execution. It identifies which agent within the workflow
    /// caused the failure.
    ///
    /// See ANVIL Spec SS7.2 -- Agent Workflows.
    #[error("workflow '{workflow_name}' failed at agent '{agent_name}': {reason}")]
    WorkflowFailed {
        /// The name of the workflow that failed.
        workflow_name: String,
        /// The name of the agent within the workflow that caused the failure.
        agent_name: String,
        /// A human-readable description of the failure.
        reason: String,
    },

    /// A sub-agent delegation failed.
    ///
    /// This error occurs when creating or executing a sub-agent that was
    /// delegated work from a parent agent. The parent-child relationship
    /// is captured in the error context.
    ///
    /// See ANVIL Spec SS11.2 -- Lineage Propagation.
    #[error(
        "sub-agent delegation failed: parent '{parent_name}' -> child '{child_name}': {reason}"
    )]
    SubAgentFailed {
        /// The parent agent that initiated the delegation.
        parent_name: String,
        /// The child agent that failed.
        child_name: String,
        /// A human-readable description of the failure.
        reason: String,
    },

    /// A lifecycle state machine violation occurred.
    ///
    /// This error wraps [`forge_health::error::ForgeHealthError`] when an agent
    /// attempts an invalid lifecycle transition (e.g., running a terminated agent).
    ///
    /// See ANVIL Spec SS5.1 -- Lifecycle State Machine.
    #[error("lifecycle error in agent '{agent_name}': {reason}")]
    LifecycleError {
        /// The agent whose lifecycle is in an invalid state.
        agent_name: String,
        /// A human-readable explanation of the lifecycle violation.
        reason: String,
    },

    /// An inter-agent messaging operation failed.
    ///
    /// This error occurs when sending or receiving messages between agents
    /// via agent channels. Common causes include channel closure and capacity
    /// exhaustion.
    #[error("agent messaging error: {reason}")]
    MessageError {
        /// A human-readable description of the messaging failure.
        reason: String,
    },

    /// The agent's stop condition was reached.
    ///
    /// This is a structured termination, not a failure. The agent stopped
    /// because its configured stop condition (max steps, max tokens, or a
    /// custom predicate) was satisfied.
    ///
    /// See ANVIL Spec SS7.1 -- Agent Execution Loop, stop conditions.
    #[error("agent '{agent_name}' stopped: {reason} (steps={steps_taken}, tokens={tokens_used})")]
    StopConditionReached {
        /// The agent that was stopped.
        agent_name: String,
        /// Why the agent stopped.
        reason: String,
        /// Number of tool loop steps completed.
        steps_taken: u32,
        /// Total tokens consumed across all steps.
        tokens_used: u64,
    },

    /// Tool invocation was denied by the authorization gate (ANVIL Spec SS8.7).
    ///
    /// This error occurs when a Host (Tier 2) tool is invoked but the agent's
    /// Arsenal ACT does not grant the required scope. Platform (Tier 1) and
    /// Embedded (Tier 3) tools are never denied.
    ///
    /// See ANVIL Spec SS8.7 -- Tool Authorization Gate.
    #[error("tool '{tool_name}' authorization denied for agent '{agent_did}': {reason}")]
    AuthorizationDenied {
        /// The DID of the agent that was denied.
        agent_did: String,
        /// The tool that was requested.
        tool_name: String,
        /// A human-readable explanation of the denial.
        reason: String,
    },

    /// Coding-provider execution was denied by the active boundary contract.
    #[error("coding-provider execution denied for agent '{agent_name}' on provider '{provider_ref}': {reason}")]
    BoundaryContractDenied {
        /// The human-readable agent name.
        agent_name: String,
        /// The provider reference that was denied.
        provider_ref: String,
        /// Why the boundary contract denied the execution.
        reason: String,
    },

    /// AEGIS-backed delegation or authority validation failed before provider execution.
    #[error("coding-provider delegation denied for agent '{agent_name}' on provider '{provider_ref}': {reason}")]
    DelegationDenied {
        /// The human-readable agent name.
        agent_name: String,
        /// The provider reference that was denied.
        provider_ref: String,
        /// Why the delegation or authority check failed.
        reason: String,
    },

    /// The agent's ACT does not grant the provider scope required for execution.
    #[error("coding-provider scope denied for agent '{agent_name}' on provider '{provider_ref}': missing scope '{required_scope}' in ACT {act_id}")]
    CredentialScopeDenied {
        /// The human-readable agent name.
        agent_name: String,
        /// The provider reference being executed.
        provider_ref: String,
        /// The scope required by the provider execution preflight.
        required_scope: String,
        /// The ACT that was checked, or `[missing]` when no ACT was configured.
        act_id: String,
    },

    /// Identity derivation failed during sub-agent creation.
    ///
    /// This error occurs when a parent agent attempts to derive a child agent
    /// identity, but the OAS key derivation or document construction fails.
    ///
    /// See ANVIL Spec SS11.2 -- Lineage Propagation.
    #[error("identity derivation failed for sub-agent '{child_name}' from parent '{parent_did}': {reason}")]
    IdentityDerivationFailed {
        /// The parent agent's OAS DID.
        parent_did: String,
        /// The child agent's intended name.
        child_name: String,
        /// A human-readable description of what went wrong.
        reason: String,
    },

    /// A `forge-core` error occurred during agent execution.
    #[error("core error: {0}")]
    Core(#[from] forge_core::error::ForgeError),

    /// A `forge-generate` error occurred during inference.
    #[error("generation error: {0}")]
    Generate(#[from] forge_generate::ForgeGenerateError),

    /// A `forge-tool` error occurred during tool execution.
    #[error("tool error: {0}")]
    Tool(#[from] forge_tool::error::ForgeToolError),

    /// A `forge-health` error occurred during lifecycle management.
    #[error("health error: {0}")]
    Health(#[from] forge_health::error::ForgeHealthError),

    /// A `forge-identity` error occurred during identity operations.
    #[error("identity error: {0}")]
    Identity(#[from] forge_identity::error::ForgeIdentityError),

    /// A `forge-auth` error occurred during authorization operations.
    #[error("auth error: {0}")]
    Auth(#[from] forge_auth::error::ForgeAuthError),
}

pub type ForgeAgentResult<T> = Result<T, ForgeAgentError>;

lib.rs

Read declaration text · 24 declaration entries

pub mod agent;

pub mod cambium;

pub mod context_manager;

pub mod error;

pub mod loop_control;

#[cfg(not(target_arch = "wasm32"))]
pub mod messaging;

pub mod observer;

pub mod streaming_tool_loop;

pub mod subagent;

pub mod tool_loop;

pub mod workflow;

pub mod prelude;

#[cfg(not(target_arch = "wasm32"))]
pub use crate::agent::AegisProviderAuthorityVerifier;

pub use crate::agent::{
        Agent, AgentCheckpointValidationError, AgentConfig, AgentEvent, AgentOutput,
        AgentRunCheckpoint, BoundaryContract, CodingProviderPreflight, ModelInferenceRecord,
        ProviderAuthorityVerifier, ProviderIdentityVerification, SubAgentDelegationRecord,
        ToolInvocationRecord, ToolInvocationStatus, AGENT_RUN_CHECKPOINT_SCHEMA,
    };

pub use crate::cambium::{
        CambiumAgentLoopObserver, CambiumEvent, CambiumJson, CambiumScope,
        ForgeCambiumObserverConfig,
    };

pub use crate::context_manager::{ContextWindowConfig, ContextWindowManager, PruningStrategy};

pub use crate::error::{ForgeAgentError, ForgeAgentResult};

pub use crate::loop_control::{
        AgentStopCondition, NoOpPrepare, PrepareStep, StopWhen, StopWhenCustom, StopWhenMaxSteps,
        StopWhenTextGenerated, StopWhenToolCalled,
    };

#[cfg(not(target_arch = "wasm32"))]
pub use crate::messaging::{
        AgentChannel, AgentChannelReceiver, AgentChannelSender, AgentMessage,
    };

pub use crate::observer::{AgentLoopObserver, NoOpObserver};

pub use crate::streaming_tool_loop::{StreamingLoopConfig, StreamingToolLoopAgent};

pub use crate::subagent::{
        create_subagent, create_subagent_with_observer, SubAgentConfig, SubAgentDelegationContext,
    };

pub use crate::tool_loop::ToolLoopAgent;

pub use crate::workflow::{
        ParallelWorkflow, RouterWorkflow, SequentialWorkflow, WorkflowOutput,
    };

loop_control.rs

Read declaration text · 12 declaration entries

pub enum AgentStopCondition {
    /// Stop when the step count reaches or exceeds this value.
    ///
    /// Steps are counted starting from 1 (the first LLM call is step 1).
    MaxSteps(u32),

    /// Stop when cumulative token usage reaches or exceeds this value.
    ///
    /// Token count includes both prompt and completion tokens across all steps.
    MaxTokens(u64),

    /// Stop when a custom predicate returns `true`.
    ///
    /// The predicate receives the current step count and total token usage.
    /// It must be `Send + Sync` for use in async contexts.
    ///
    /// # Examples
    ///
    /// ```
    /// use forge_agent::loop_control::AgentStopCondition;
    ///
    /// let stop = AgentStopCondition::Custom(Box::new(|steps, tokens| {
    ///     steps >= 5 && tokens >= 2000
    /// }));
    /// assert!(stop.should_stop(5, 2000));
    /// assert!(!stop.should_stop(4, 2000));
    /// assert!(!stop.should_stop(5, 1999));
    /// ```
    Custom(Box<dyn Fn(u32, u64) -> bool + Send + Sync>),
}

pub fn should_stop(&self, steps: u32, total_tokens: u64) -> bool;

pub fn description(&self) -> String;

pub trait StopWhen: Send + Sync {
    /// Determines whether the agent's tool loop should stop.
    ///
    /// # Arguments
    ///
    /// * `messages` - The current conversation history.
    /// * `step_count` - The number of tool loop steps completed (1-indexed).
    ///
    /// # Returns
    ///
    /// `true` if the loop should terminate, `false` to continue.
    fn should_stop(&self, messages: &[ModelMessage], step_count: u32) -> bool;
}

pub struct StopWhenTextGenerated;

pub struct StopWhenMaxSteps(pub u32);

pub struct StopWhenToolCalled {

}

pub fn new(tool_name: impl Into<String>) -> Self;

pub struct StopWhenCustom {

}

pub fn new<F>(predicate: F) -> Self
    where
        F: Fn(&[ModelMessage], u32) -> bool + Send + Sync + 'static,;

#[async_trait]
pub trait PrepareStep: Send + Sync {
    /// Modifies the message list before the next LLM call.
    ///
    /// # Arguments
    ///
    /// * `messages` - The mutable message list that will be sent to the model.
    ///   Implementors may add, remove, or modify messages.
    /// * `step` - The current step number (1-indexed). Step 1 is the first call.
    ///
    /// # Returns
    ///
    /// `Ok(())` on success, or a [`ForgeAgentError`](crate::error::ForgeAgentError)
    /// if preparation fails (which aborts the tool loop).
    ///
    /// # Errors
    ///
    /// Returning an error aborts the tool loop and propagates the error to
    /// the caller.
    async fn prepare(&self, messages: &mut Vec<ModelMessage>, step: u32) -> ForgeAgentResult<()>;
}

pub struct NoOpPrepare;

messaging.rs

Read declaration text · 17 declaration entries

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentMessage {

}

pub fn text(
        sender: impl Into<String>,
        recipient: impl Into<String>,
        payload: impl Into<String>,
    ) -> Self;

pub fn result(
        sender: impl Into<String>,
        recipient: impl Into<String>,
        payload: impl Into<String>,
    ) -> Self;

pub fn error(
        sender: impl Into<String>,
        recipient: impl Into<String>,
        payload: impl Into<String>,
    ) -> Self;

pub fn control(
        sender: impl Into<String>,
        recipient: impl Into<String>,
        payload: impl Into<String>,
    ) -> Self;

pub fn sender(&self) -> &str;

pub fn recipient(&self) -> &str;

pub fn payload(&self) -> &str;

pub fn kind(&self) -> &str;

#[derive(Debug, Clone)]
pub struct AgentChannelSender {

}

pub async fn send(&self, message: AgentMessage) -> ForgeAgentResult<()>;

pub fn try_send(&self, message: AgentMessage) -> ForgeAgentResult<()>;

#[derive(Debug)]
pub struct AgentChannelReceiver {

}

pub async fn recv(&mut self) -> Option<AgentMessage>;

pub fn try_recv(&mut self) -> Option<AgentMessage>;

pub struct AgentChannel;

#[allow(clippy::new_ret_no_self)]
pub fn new(capacity: usize) -> (AgentChannelSender, AgentChannelReceiver);

observer.rs

Read declaration text · 2 declaration entries

#[async_trait]
pub trait AgentLoopObserver: Send + Sync {
    /// Called at the start of each tool loop turn.
    ///
    /// # Arguments
    ///
    /// * `turn` - The turn number (1-indexed).
    async fn on_turn_start(&self, _turn: u32) ;

    /// Called at the end of each tool loop turn.
    ///
    /// # Arguments
    ///
    /// * `turn` - The turn number (1-indexed).
    /// * `usage` - Cumulative token usage after this turn.
    async fn on_turn_end(&self, _turn: u32, _usage: Usage) ;

    /// Called when the model produces a text delta (during streaming).
    ///
    /// # Arguments
    ///
    /// * `text` - The text fragment.
    async fn on_text_delta(&self, _text: &str) ;

    /// Called when the model starts a tool call.
    ///
    /// # Arguments
    ///
    /// * `tool_call` - The tool call being initiated.
    async fn on_tool_call_start(&self, _tool_call: &ToolCall) ;

    /// Called when a tool call completes.
    ///
    /// Legacy hook -- prefer [`on_tool_invocation`](Self::on_tool_invocation) for
    /// per-tool latency, status (succeeded / failed / authorization-denied /
    /// not-found), and the tool name. This hook is preserved for backward
    /// compatibility and continues to fire alongside `on_tool_invocation`.
    ///
    /// # Arguments
    ///
    /// * `tool_call_id` - The tool call ID.
    /// * `result` - The result content.
    /// * `is_error` - Whether the tool execution failed.
    async fn on_tool_call_end(&self, _tool_call_id: &str, _result: &str, _is_error: bool) ;

    /// Called when a tool call completes, with the full invocation record.
    ///
    /// Symmetric counterpart to [`on_tool_call_start`](Self::on_tool_call_start)
    /// -- this hook receives the tool name, wall-clock duration, and the typed
    /// `ToolInvocationStatus` so observers can build per-tool latency
    /// histograms, structured audit entries, or richer UI states without
    /// having to maintain their own start-time bookkeeping.
    ///
    /// Default implementation is a no-op so existing observers compile
    /// unchanged. The streaming tool loop fires both this and
    /// [`on_tool_call_end`](Self::on_tool_call_end) for every invocation.
    ///
    /// # Arguments
    ///
    /// * `record` - The complete invocation record.
    async fn on_tool_invocation(&self, _record: &ToolInvocationRecord) ;

    /// Called when a model inference call completes.
    ///
    /// The record carries model/provider ids, token counts, latency, outcome,
    /// and governed prompt/completion refs. It does not carry raw prompt or
    /// completion text.
    ///
    /// # Arguments
    ///
    /// * `record` - The complete model inference record.
    async fn on_model_inference_completed(&self, _record: &ModelInferenceRecord) ;

    /// Called when a parent agent delegates work to a sub-agent.
    ///
    /// The record carries child agent/run identifiers and governed refs for
    /// instruction or handoff context.
    ///
    /// # Arguments
    ///
    /// * `record` - The delegation record.
    async fn on_subagent_delegation_started(&self, _record: &SubAgentDelegationRecord) ;

    /// Called when a stream chunk is received.
    ///
    /// # Arguments
    ///
    /// * `chunk` - The stream chunk.
    async fn on_stream_chunk(&self, _chunk: &StreamChunk) ;

    /// Called when the agent loop completes successfully.
    ///
    /// # Arguments
    ///
    /// * `final_message` - The final assistant message.
    /// * `total_turns` - Total number of turns executed.
    /// * `total_usage` - Total token usage across all turns.
    async fn on_complete(
        &self,
        _final_message: &ModelMessage,
        _total_turns: u32,
        _total_usage: Usage,
    ) ;

    /// Called when the agent loop encounters an error.
    ///
    /// # Arguments
    ///
    /// * `error` - A human-readable error description.
    /// * `turn` - The turn in which the error occurred.
    async fn on_error(&self, _error: &str, _turn: u32) ;

    /// Called when the agent loop is stopped by a stop condition.
    ///
    /// # Arguments
    ///
    /// * `reason` - Why the loop stopped.
    /// * `turn` - The turn at which the loop stopped.
    async fn on_stopped(&self, _reason: FinishReason, _turn: u32) ;
}

pub struct NoOpObserver;

streaming_tool_loop.rs

Read declaration text · 11 declaration entries

#[derive(Debug, Clone)]
pub struct StreamingLoopConfig {
/// Maximum turns before forced termination.

pub max_turns: u32,
/// Context pruning threshold (fraction of max_context_tokens).

pub context_threshold: f64,
/// Context pruning strategy.

pub pruning_strategy: PruningStrategy,
/// Model's maximum context window size.

pub max_context_tokens: u64
}

pub struct StreamingToolLoopAgent {

}

pub fn new(
        config: AgentConfig,
        loop_config: StreamingLoopConfig,
        model: Arc<dyn LanguageModel>,
        tool_registry: ToolRegistry,
        approval: Arc<dyn ApprovalHandler>,
    ) -> Self;

pub fn with_observer(self, observer: Arc<dyn AgentLoopObserver>) -> Self;

pub fn set_observer(&self, observer: Arc<dyn AgentLoopObserver>);

pub fn clear_observer(&self);

pub fn observer(&self) -> Arc<dyn AgentLoopObserver>;

pub fn enqueue_interjection(&self, text: impl Into<String>);

pub fn enqueue_interjection_message(&self, message: ModelMessage);

pub fn model(&self) -> &Arc<dyn LanguageModel>;

pub fn with_generate_options(mut self, options: GenerateOptions) -> Self;

subagent.rs

Read declaration text · 12 declaration entries

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SubAgentConfig {

}

pub fn new(name: impl Into<String>) -> Self;

pub fn with_tool_names(mut self, names: Vec<String>) -> Self;

pub fn with_system_prompt(mut self, prompt: impl Into<String>) -> Self;

pub fn with_max_steps(mut self, max: u32) -> Self;

pub fn name(&self) -> &str;

pub fn tool_names(&self) -> &[String];

pub fn system_prompt(&self) -> Option<&str>;

pub fn max_steps(&self) -> Option<u32>;

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SubAgentDelegationContext {
/// Optional stable child agent id. Falls back to the child's configured DID

/// or child name when omitted.

pub subagent_id: Option<String>,
/// Host-assigned run id for the delegated child agent.

pub subagent_run_id: String,
/// Optional summary of why the delegation happened.

pub delegation_reason: Option<String>,
/// Governed artifact ref for delegated instructions.

pub instruction_ref: Option<String>,
/// Capability refs granted to the child.

pub capability_refs: Vec<String>,
/// Governed artifact ref for handoff state.

pub handoff_ref: Option<String>
}

pub fn create_subagent(
    parent_config: &AgentConfig,
    sub_config: SubAgentConfig,
    model: Arc<dyn LanguageModel>,
    tool_registry: ToolRegistry,
    approval: Arc<dyn ApprovalHandler>,
) -> ForgeAgentResult<ToolLoopAgent>;

pub async fn create_subagent_with_observer(
    parent_config: &AgentConfig,
    sub_config: SubAgentConfig,
    model: Arc<dyn LanguageModel>,
    tool_registry: ToolRegistry,
    approval: Arc<dyn ApprovalHandler>,
    observer: &dyn AgentLoopObserver,
    context: SubAgentDelegationContext,
) -> ForgeAgentResult<ToolLoopAgent>;

tool_loop.rs

Read declaration text · 5 declaration entries

pub struct ToolLoopAgent {

}

pub fn new(
        config: AgentConfig,
        model: Arc<dyn LanguageModel>,
        tool_registry: ToolRegistry,
        approval: Arc<dyn ApprovalHandler>,
    ) -> Self;

pub fn with_stop_condition(mut self, condition: AgentStopCondition) -> Self;

pub fn with_prepare_step(mut self, prepare: Arc<dyn PrepareStep>) -> Self;

pub fn with_generate_options(mut self, options: GenerateOptions) -> Self;

workflow.rs

Read declaration text · 10 declaration entries

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkflowOutput {
/// Individual outputs from each agent that executed.

pub agent_outputs: Vec<AgentOutput>,
/// Aggregate token usage across all agents in the workflow.

pub total_usage: Usage,
/// Total number of steps taken across all agents.

pub total_steps: u32,
/// The final text output of the workflow.

///

/// For sequential workflows, this is the last agent's output.

/// For parallel workflows, this is the concatenated outputs.

/// For router workflows, this is the selected agent's output.

pub final_text: String
}

pub struct SequentialWorkflow {

}

pub fn new(name: impl Into<String>, agents: Vec<Arc<dyn Agent>>) -> Self;

pub async fn execute(&self, prompt: &str) -> Result<WorkflowOutput, ForgeAgentError>;

pub struct ParallelWorkflow {

}

pub fn new(name: impl Into<String>, agents: Vec<Arc<dyn Agent>>) -> Self;

pub async fn execute(&self, prompt: &str) -> Result<WorkflowOutput, ForgeAgentError>;

pub struct RouterWorkflow {

}

pub fn new(
        name: impl Into<String>,
        classifier: Arc<dyn Agent>,
        agents: Vec<(String, Arc<dyn Agent>)>,
    ) -> Self;

pub async fn execute(&self, prompt: &str) -> Result<WorkflowOutput, ForgeAgentError>;

Continue

On this page