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
| Field | Value |
|---|---|
| Language | rust |
| Source version | 0.2.0 |
| Manifest | forge-rs/crates/forge-agent/Cargo.toml |
| Source files | 12 |
| Evidence | Source 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>;