Forge documentation
Library referenceRust

forge-contracts

Formal interface contracts between Forge (agent substrate) and Aut0 (organization platform)

Formal interface contracts between Forge (agent substrate) and Aut0 (organization platform)

Package contract

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

Import boundary

use forge_contracts;

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 auth;

pub mod brew;

pub mod comm;

pub mod error;

pub mod flowers;

pub mod identity;

pub mod mcp;

pub mod memory;

pub mod provider;

pub mod runtime;

pub mod telemetry;

pub mod prelude;

pub use crate::auth::{
        AuthContract, AuthDecision, CapabilityNarrowingRequest, DelegationChainEntry,
    };

pub use crate::brew::{BrewContract, PlanExecutionResult, PlanHandle};

pub use crate::comm::{ChannelConfig, ChannelHandle, CommContract, SessionConfig};

pub use crate::error::ContractError;

pub use crate::flowers::{CheckpointData, FlowersBridgeContract, FlowersExecutionHandle};

pub use crate::identity::{
        DerivedIdentityRequest, IdentityContract, IdentityHandle, LineageInfo,
    };

pub use crate::mcp::{McpContract, McpServerHandle, McpToolDescriptor};

pub use crate::memory::{
        MemoryContract, MemoryQuery, MemoryRecord, MemoryScope, MemoryWritePolicy,
    };

pub use crate::provider::{
        OrgProviderPolicy, ProviderContract, ProviderFallbackStrategy, ProviderSessionHandle,
    };

pub use crate::runtime::{
        AgentCreateRequest, AgentHandle, AgentLifecycleCommand, AgentRuntimeContract, AgentStatus,
    };

pub use crate::telemetry::{
        HealthSummary, OrgHealthContract, OrgTelemetryContract, SpanFilter,
    };

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.

auth.rs

Read declaration text · 8 declaration entries

pub const CONTRACT_VERSION: &str;

#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum AuthDecision {
    /// Authorization was granted.
    Allowed {
        /// The specific scope that was matched.
        scope: String,
        /// When this authorization expires, if applicable.
        expires_at: Option<DateTime<Utc>>,
    },
    /// Authorization was denied.
    Denied {
        /// The scope that was requested.
        requested_scope: String,
        /// The reason for denial.
        reason: String,
    },
}

pub fn is_allowed(&self) -> bool;

pub fn is_denied(&self) -> bool;

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CapabilityNarrowingRequest {
/// The parent agent's DID (must have a valid ACT).

pub parent_did: String,
/// The child agent's DID (will receive the narrowed ACT).

pub child_did: String,
/// The scopes to grant to the child. Must be a subset of parent's scopes.

pub requested_scopes: Vec<String>,
/// Optional time-to-live in seconds for the child's token.

/// If `None`, inherits the parent's expiration.

pub ttl_seconds: Option<u64>,
/// Maximum delegation chain depth. If the parent is already at

/// this depth, delegation fails.

pub max_delegation_depth: u32
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DelegationChainEntry {
/// The delegator's DID.

pub delegator_did: String,
/// The delegatee's DID.

pub delegatee_did: String,
/// The scopes that were delegated.

pub scopes: Vec<String>,
/// When the delegation was created.

pub created_at: DateTime<Utc>,
/// When the delegation expires.

pub expires_at: Option<DateTime<Utc>>,
/// Depth in the delegation chain (0 = root grant).

pub depth: u32
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CapabilityTokenHandle {
/// Unique identifier for this token.

pub token_id: String,
/// The agent DID this token belongs to.

pub agent_did: String,
/// The scopes granted by this token.

pub scopes: Vec<String>,
/// When this token expires, if applicable.

pub expires_at: Option<DateTime<Utc>>,
/// Depth in the delegation chain.

pub delegation_depth: u32
}

#[async_trait]
pub trait AuthContract: Send + Sync {
    /// Checks whether an agent is authorized for a specific scope.
    ///
    /// # Arguments
    ///
    /// * `agent_did` - The agent requesting authorization.
    /// * `scope` - The scope to check (e.g., "tool:web_search",
    ///   "memory:write:department:engineering").
    ///
    /// # Returns
    ///
    /// An `AuthDecision` indicating whether access is granted or denied.
    ///
    /// # Errors
    ///
    /// - `ContractError::DidResolutionFailed` if the agent's DID cannot
    ///   be resolved to find its ACT.
    async fn check_authorization(
        &self,
        agent_did: &str,
        scope: &str,
    ) -> ContractResult<AuthDecision>;

    /// Creates a narrowed capability token for a child agent.
    ///
    /// # Arguments
    ///
    /// * `request` - The narrowing parameters.
    ///
    /// # Returns
    ///
    /// A handle to the newly created child token.
    ///
    /// # Errors
    ///
    /// - `ContractError::CapabilityEscalation` if any requested scope
    ///   exceeds the parent's grants.
    /// - `ContractError::TokenExpired` if the parent's token has expired.
    async fn delegate_capabilities(
        &self,
        request: CapabilityNarrowingRequest,
    ) -> ContractResult<CapabilityTokenHandle>;

    /// Revokes a capability token, immediately invalidating it.
    ///
    /// Revocation cascades: revoking a parent token also revokes all
    /// tokens derived from it.
    ///
    /// # Arguments
    ///
    /// * `token_id` - The token to revoke.
    ///
    /// # Errors
    ///
    /// - `ContractError::AuthorizationDenied` if the caller does not
    ///   have authority to revoke this token.
    async fn revoke_token(&self, token_id: &str) -> ContractResult<()>;

    /// Returns the full delegation chain for an agent's current token.
    ///
    /// # Arguments
    ///
    /// * `agent_did` - The agent whose delegation chain to retrieve.
    ///
    /// # Returns
    ///
    /// The chain of delegations from root to the agent, ordered by depth.
    ///
    /// # Errors
    ///
    /// - `ContractError::DidResolutionFailed` if the agent DID cannot
    ///   be resolved.
    async fn get_delegation_chain(
        &self,
        agent_did: &str,
    ) -> ContractResult<Vec<DelegationChainEntry>>;

    /// Lists all scopes currently granted to an agent.
    ///
    /// # Arguments
    ///
    /// * `agent_did` - The agent to query.
    ///
    /// # Returns
    ///
    /// The list of scope strings currently active for this agent.
    async fn list_scopes(&self, agent_did: &str) -> ContractResult<Vec<String>>;
}

brew.rs

Read declaration text · 9 declaration entries

pub const CONTRACT_VERSION: &str;

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PlanHandle {
/// Unique plan identifier.

pub plan_id: String,
/// The Brew graph identifier this plan was frozen from.

pub brew_id: String,
/// Number of nodes in the plan.

pub node_count: u32,
/// Number of edges in the plan.

pub edge_count: u32,
/// Whether this plan is currently executing.

pub executing: bool
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BrewNodeSpec {
/// Unique node identifier within the brew.

pub node_id: String,
/// The type of node.

pub node_type: BrewNodeType,
/// Input mapping: key is parameter name, value is source expression.

pub inputs: BTreeMap<String, String>,
/// Configuration specific to the node type.

pub config: serde_json::Value
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum BrewNodeType {
    /// An agent execution node (runs a tool loop).
    Agent,
    /// A tool invocation node (calls a single tool).
    Tool,
    /// A provider call node (single LLM inference).
    Inference,
    /// A conditional branch node.
    Condition,
    /// A data transformation node.
    Transform,
    /// A sub-brew reference node.
    SubBrew,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BrewEdgeSpec {
/// Source node ID.

pub from_node: String,
/// Target node ID.

pub to_node: String,
/// The type of edge.

pub edge_type: BrewEdgeType,
/// Optional condition expression for conditional edges.

pub condition: Option<String>
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum BrewEdgeType {
    /// Data flows from source output to target input.
    Data,
    /// Control flow: target executes after source completes.
    Control,
    /// Error flow: target executes if source fails.
    Error,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PlanExecutionResult {
/// The plan that was executed.

pub plan_id: String,
/// Whether the plan completed successfully.

pub success: bool,
/// Results from each node, keyed by node ID.

pub node_results: BTreeMap<String, NodeResult>,
/// Total execution time in milliseconds.

pub duration_ms: u64,
/// Total tokens consumed across all nodes.

pub total_tokens: u64
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct NodeResult {
/// The node identifier.

pub node_id: String,
/// Whether this node succeeded.

pub success: bool,
/// The node's output as JSON.

pub output: Option<serde_json::Value>,
/// Error message if the node failed.

pub error: Option<String>,
/// Execution time for this node in milliseconds.

pub duration_ms: u64
}

#[async_trait]
pub trait BrewContract: Send + Sync {
    /// Creates a new empty Brew graph.
    ///
    /// # Arguments
    ///
    /// * `brew_id` - Unique identifier for this brew.
    /// * `description` - Human-readable description of the plan's purpose.
    ///
    /// # Returns
    ///
    /// A handle to the unfrozen plan.
    async fn create_plan(&self, brew_id: &str, description: &str) -> ContractResult<PlanHandle>;

    /// Adds a node to an unfrozen plan.
    ///
    /// # Arguments
    ///
    /// * `plan_id` - The plan to modify.
    /// * `node` - The node specification.
    ///
    /// # Errors
    ///
    /// - `ContractError::PlanResolutionFailed` if the plan is already frozen.
    async fn add_node(&self, plan_id: &str, node: BrewNodeSpec) -> ContractResult<()>;

    /// Adds an edge to an unfrozen plan.
    ///
    /// # Arguments
    ///
    /// * `plan_id` - The plan to modify.
    /// * `edge` - The edge specification.
    ///
    /// # Errors
    ///
    /// - `ContractError::PlanResolutionFailed` if the plan is already frozen
    ///   or if referenced nodes do not exist.
    async fn add_edge(&self, plan_id: &str, edge: BrewEdgeSpec) -> ContractResult<()>;

    /// Freezes a plan, validating and resolving all symbols.
    ///
    /// After freezing, no modifications are allowed. The plan is ready
    /// for execution.
    ///
    /// # Arguments
    ///
    /// * `plan_id` - The plan to freeze.
    ///
    /// # Returns
    ///
    /// The updated plan handle with `executing: false`.
    ///
    /// # Errors
    ///
    /// - `ContractError::PlanResolutionFailed` if validation fails.
    async fn freeze_plan(&self, plan_id: &str) -> ContractResult<PlanHandle>;

    /// Executes a frozen plan.
    ///
    /// # Arguments
    ///
    /// * `plan_id` - The frozen plan to execute.
    /// * `inputs` - Input values keyed by parameter name.
    ///
    /// # Returns
    ///
    /// The execution result with outputs from all nodes.
    ///
    /// # Errors
    ///
    /// - `ContractError::PlanExecutionFailed` if execution fails.
    async fn execute_plan(
        &self,
        plan_id: &str,
        inputs: BTreeMap<String, serde_json::Value>,
    ) -> ContractResult<PlanExecutionResult>;

    /// Queries the current status of a plan.
    ///
    /// # Arguments
    ///
    /// * `plan_id` - The plan to query.
    async fn get_plan_status(&self, plan_id: &str) -> ContractResult<PlanHandle>;
}

comm.rs

Read declaration text · 9 declaration entries

pub const CONTRACT_VERSION: &str;

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ChannelConfig {
/// Human-readable channel name.

pub channel_name: String,
/// The type of channel (maps to different session semantics).

pub channel_type: ChannelType,
/// The organization this channel belongs to.

pub org_id: String,
/// Optional department scope.

pub department_id: Option<String>,
/// DIDs of agents allowed in this channel.

pub participants: Vec<String>,
/// Maximum message payload size in bytes.

pub max_message_size_bytes: u32,
/// How long to retain message history, in hours. `None` = forever.

pub history_retention_hours: Option<u32>
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum ChannelType {
    /// Organization-wide broadcast channel.
    OrgWide,
    /// Department-scoped channel.
    Department,
    /// Direct message between two agents.
    DirectMessage,
    /// Executive channel (Company Director + root-holder).
    Executive,
    /// Root-holder secure channel.
    RootHolder,
    /// Custom channel type.
    Custom(String),
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ChannelHandle {
/// Unique channel identifier.

pub channel_id: String,
/// The underlying Forge session identifier.

pub session_id: String,
/// The channel name.

pub name: String,
/// The channel type.

pub channel_type: ChannelType,
/// Number of participants.

pub participant_count: u32
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionConfig {
/// Human-readable session name.

pub session_name: String,
/// The coordinator agent's DID.

pub coordinator_did: String,
/// Worker agent DIDs.

pub worker_dids: Vec<String>,
/// Optional timeout for the session in seconds.

pub timeout_seconds: Option<u64>,
/// Whether to enable shared context for this session.

pub shared_context_enabled: bool
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ChannelMessage {
/// The sender's DID.

pub sender_did: String,
/// The message content type.

pub content_type: MessageContentType,
/// The message payload as JSON.

pub payload: serde_json::Value,
/// Optional correlation ID for threading.

pub correlation_id: Option<String>,
/// Optional reply-to message ID.

pub reply_to: Option<String>
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum MessageContentType {
    /// Plain text message.
    Text,
    /// Structured data message.
    Structured,
    /// Task delegation message.
    TaskDelegation,
    /// Task result message.
    TaskResult,
    /// Interrupt/signal message.
    Interrupt,
    /// Status update message.
    StatusUpdate,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ReceivedMessage {
/// Unique message identifier.

pub message_id: String,
/// The sender's DID.

pub sender_did: String,
/// The message content type.

pub content_type: MessageContentType,
/// The message payload.

pub payload: serde_json::Value,
/// When the message was sent.

pub sent_at: DateTime<Utc>,
/// Ed25519 signature from the sender (base64-encoded).

pub signature: Option<String>,
/// Correlation ID for threading.

pub correlation_id: Option<String>
}

#[async_trait]
pub trait CommContract: Send + Sync {
    /// Creates a new communication channel.
    ///
    /// # Arguments
    ///
    /// * `config` - The channel configuration.
    ///
    /// # Returns
    ///
    /// A handle to the created channel.
    ///
    /// # Errors
    ///
    /// - `ContractError::ChannelError` if creation fails.
    /// - `ContractError::DidResolutionFailed` if any participant DID
    ///   cannot be resolved.
    async fn create_channel(&self, config: ChannelConfig) -> ContractResult<ChannelHandle>;

    /// Sends a message to a channel.
    ///
    /// The message is automatically wrapped in an `AgentMessage` envelope,
    /// signed by the sender, and delivered to all channel participants.
    ///
    /// # Arguments
    ///
    /// * `channel_id` - The target channel.
    /// * `message` - The message to send.
    ///
    /// # Returns
    ///
    /// The message ID assigned by the transport.
    ///
    /// # Errors
    ///
    /// - `ContractError::ChannelError` if the channel does not exist.
    /// - `ContractError::AuthorizationDenied` if the sender is not a
    ///   participant.
    async fn send_message(
        &self,
        channel_id: &str,
        message: ChannelMessage,
    ) -> ContractResult<String>;

    /// Receives the next message from a channel.
    ///
    /// This is a pull-based interface. For push-based delivery, use
    /// `subscribe`.
    ///
    /// # Arguments
    ///
    /// * `channel_id` - The channel to receive from.
    /// * `agent_did` - The receiving agent's DID.
    ///
    /// # Returns
    ///
    /// The next unread message, or `None` if no messages are pending.
    async fn receive_message(
        &self,
        channel_id: &str,
        agent_did: &str,
    ) -> ContractResult<Option<ReceivedMessage>>;

    /// Closes a channel, cleaning up resources.
    ///
    /// # Arguments
    ///
    /// * `channel_id` - The channel to close.
    async fn close_channel(&self, channel_id: &str) -> ContractResult<()>;

    /// Lists all channels for an organization.
    ///
    /// # Arguments
    ///
    /// * `org_id` - The organization to query.
    /// * `department_id` - Optional department filter.
    async fn list_channels(
        &self,
        org_id: &str,
        department_id: Option<&str>,
    ) -> ContractResult<Vec<ChannelHandle>>;
}

error.rs

Read declaration text · 2 declaration entries

#[derive(Debug, Error)]
pub enum ContractError {
    // -----------------------------------------------------------------------
    // S-01: Agent Runtime
    // -----------------------------------------------------------------------
    /// The requested agent does not exist or has been terminated.
    #[error("agent '{agent_id}' not found in runtime")]
    AgentNotFound {
        /// The agent identifier that was not found.
        agent_id: String,
    },

    /// Agent creation failed due to invalid configuration.
    #[error("agent creation failed: {reason}")]
    AgentCreationFailed {
        /// Human-readable explanation of why creation failed.
        reason: String,
    },

    /// An invalid lifecycle transition was attempted.
    #[error("invalid lifecycle transition from {from} to {to} for agent '{agent_id}'")]
    InvalidLifecycleTransition {
        /// The agent that owns the lifecycle.
        agent_id: String,
        /// The current state name.
        from: String,
        /// The requested target state name.
        to: String,
    },

    // -----------------------------------------------------------------------
    // S-02: Identity & Lineage
    // -----------------------------------------------------------------------
    /// Identity derivation failed.
    #[error("identity derivation failed for path '{path}' from parent '{parent_did}': {reason}")]
    IdentityDerivationFailed {
        /// The parent DID from which derivation was attempted.
        parent_did: String,
        /// The derivation path that failed.
        path: String,
        /// Explanation of the failure.
        reason: String,
    },

    /// Lineage verification failed at the specified depth.
    #[error("lineage verification failed at depth {depth} for '{did}': {reason}")]
    LineageVerificationFailed {
        /// The DID whose lineage failed verification.
        did: String,
        /// The depth at which verification failed.
        depth: u32,
        /// Explanation of the failure.
        reason: String,
    },

    /// DID resolution failed.
    #[error("DID resolution failed for '{did}': {reason}")]
    DidResolutionFailed {
        /// The DID that could not be resolved.
        did: String,
        /// Explanation of the failure.
        reason: String,
    },

    // -----------------------------------------------------------------------
    // S-03: Auth & Delegation
    // -----------------------------------------------------------------------
    /// Capability escalation denied: child requested more than parent has.
    #[error("capability escalation denied: agent '{agent_did}' requested scope '{requested}' but parent ACT only grants {available:?}")]
    CapabilityEscalation {
        /// The agent DID that attempted escalation.
        agent_did: String,
        /// The scope that was requested.
        requested: String,
        /// The scopes available to the parent.
        available: Vec<String>,
    },

    /// A capability token has expired.
    #[error("capability token '{token_id}' for agent '{agent_did}' expired at {expired_at}")]
    TokenExpired {
        /// The token identifier.
        token_id: String,
        /// The agent DID that owns the token.
        agent_did: String,
        /// ISO 8601 timestamp when the token expired.
        expired_at: String,
    },

    /// Authorization check failed.
    #[error("authorization denied for agent '{agent_did}' on scope '{scope}': {reason}")]
    AuthorizationDenied {
        /// The agent DID that was denied.
        agent_did: String,
        /// The scope that was requested.
        scope: String,
        /// Explanation of why authorization was denied.
        reason: String,
    },

    // -----------------------------------------------------------------------
    // S-04: Provider Routing
    // -----------------------------------------------------------------------
    /// No provider matched the routing policy.
    #[error("no provider matched routing policy for org '{org_id}': {reason}")]
    NoProviderAvailable {
        /// The organization identifier.
        org_id: String,
        /// Explanation of why no provider matched.
        reason: String,
    },

    /// Provider session creation or management failed.
    #[error("provider session error for '{provider_ref}': {reason}")]
    ProviderSessionError {
        /// The provider reference string.
        provider_ref: String,
        /// Explanation of the failure.
        reason: String,
    },

    // -----------------------------------------------------------------------
    // S-05: Brew Plans
    // -----------------------------------------------------------------------
    /// Plan resolution failed (unresolved symbols, cycles, etc.).
    #[error("brew plan '{plan_id}' resolution failed: {reason}")]
    PlanResolutionFailed {
        /// The plan identifier.
        plan_id: String,
        /// Explanation of the failure.
        reason: String,
    },

    /// Plan execution failed at a specific node.
    #[error("brew plan '{plan_id}' failed at node '{node_id}': {reason}")]
    PlanExecutionFailed {
        /// The plan identifier.
        plan_id: String,
        /// The node that failed.
        node_id: String,
        /// Explanation of the failure.
        reason: String,
    },

    // -----------------------------------------------------------------------
    // S-06: Comm/Collab
    // -----------------------------------------------------------------------
    /// Channel creation or operation failed.
    #[error("channel '{channel_id}' error: {reason}")]
    ChannelError {
        /// The channel identifier.
        channel_id: String,
        /// Explanation of the failure.
        reason: String,
    },

    // -----------------------------------------------------------------------
    // S-07: MCP
    // -----------------------------------------------------------------------
    /// MCP server connection or operation failed.
    #[error("MCP server '{server_name}' error: {reason}")]
    McpError {
        /// The MCP server name.
        server_name: String,
        /// Explanation of the failure.
        reason: String,
    },

    // -----------------------------------------------------------------------
    // S-08: Telemetry & Health
    // -----------------------------------------------------------------------
    /// Telemetry collection or aggregation failed.
    #[error("telemetry error: {reason}")]
    TelemetryError {
        /// Explanation of the failure.
        reason: String,
    },

    // -----------------------------------------------------------------------
    // S-09: Flowers Bridge
    // -----------------------------------------------------------------------
    /// Durable execution checkpoint or resume failed.
    #[error("flowers execution '{execution_id}' error: {reason}")]
    FlowersError {
        /// The Flowers execution identifier.
        execution_id: String,
        /// Explanation of the failure.
        reason: String,
    },

    // -----------------------------------------------------------------------
    // S-10: Memory
    // -----------------------------------------------------------------------
    /// Memory read or write operation failed.
    #[error("memory error in scope '{scope}': {reason}")]
    MemoryError {
        /// The memory scope where the error occurred.
        scope: String,
        /// Explanation of the failure.
        reason: String,
    },

    /// Memory access denied by policy.
    #[error("memory access denied for agent '{agent_did}' in scope '{scope}': {reason}")]
    MemoryAccessDenied {
        /// The agent DID that was denied.
        agent_did: String,
        /// The memory scope.
        scope: String,
        /// Explanation of the denial.
        reason: String,
    },

    // -----------------------------------------------------------------------
    // Cross-cutting
    // -----------------------------------------------------------------------
    /// Contract version mismatch between Forge and Aut0.
    #[error("contract version mismatch: Forge has {forge_version}, Aut0 expects {aut0_version}")]
    VersionMismatch {
        /// The version Forge compiled against.
        forge_version: String,
        /// The version Aut0 compiled against.
        aut0_version: String,
    },

    /// A required configuration field was missing or invalid.
    #[error("configuration error: {reason}")]
    ConfigurationError {
        /// Explanation of the configuration problem.
        reason: String,
    },

    /// An internal Forge error that Aut0 should not need to handle in detail.
    #[error("internal error: {reason}")]
    Internal {
        /// Explanation for diagnostic purposes.
        reason: String,
    },
}

pub type ContractResult<T> = Result<T, ContractError>;

flowers.rs

Read declaration text · 9 declaration entries

pub const CONTRACT_VERSION: &str;

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct FlowersExecutionHandle {
/// Unique execution identifier.

pub execution_id: String,
/// The workflow definition this execution instantiates.

pub workflow_id: String,
/// Current execution state.

pub state: FlowersExecutionState,
/// When the execution was created.

pub created_at: DateTime<Utc>,
/// When the execution last changed state.

pub updated_at: DateTime<Utc>
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum FlowersExecutionState {
    /// Execution is queued but not yet started.
    Queued,
    /// Execution is actively running.
    Running,
    /// Execution is paused (awaiting signal or timer).
    Suspended,
    /// Execution completed successfully.
    Completed,
    /// Execution failed with an error.
    Failed,
    /// Execution was cancelled.
    Cancelled,
    /// Execution is being compensated (saga rollback).
    Compensating,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct FlowersJournalEntry {
/// Sequential event index within the execution.

pub index: u64,
/// The event type.

pub event_type: JournalEntryType,
/// The operation name (e.g., "agent.invoke", "tool.execute").

pub operation: String,
/// Input data for this operation (JSON-serialized).

pub input: serde_json::Value,
/// Output data from this operation (JSON-serialized), if completed.

pub output: Option<serde_json::Value>,
/// When this event was recorded.

pub timestamp: DateTime<Utc>,
/// Hash for integrity verification.

pub hash: String
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum JournalEntryType {
    /// A provider (LLM) call.
    ProviderCall,
    /// A tool execution.
    ToolExecution,
    /// A timer event.
    Timer,
    /// A signal received.
    Signal,
    /// A checkpoint created.
    Checkpoint,
    /// A compensation (rollback) action.
    Compensation,
    /// An arbitrary side effect.
    SideEffect,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CheckpointData {
/// Unique checkpoint identifier.

pub checkpoint_id: String,
/// The execution this checkpoint belongs to.

pub execution_id: String,
/// The journal index at which this checkpoint was taken.

pub journal_index: u64,
/// Serialized execution state.

pub state: serde_json::Value,
/// When the checkpoint was created.

pub created_at: DateTime<Utc>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct FlowersSignal {
/// Signal name.

pub name: String,
/// Signal payload.

pub payload: serde_json::Value
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct FlowersSubmitRequest {
/// The agent to run (referenced by handle from S-01).

pub agent_id: String,
/// The input to pass to the agent.

pub input: String,
/// Optional workflow ID override (for resume scenarios).

pub workflow_id: Option<String>,
/// Optional checkpoint to resume from.

pub resume_from_checkpoint: Option<String>,
/// Maximum execution time in seconds.

pub timeout_seconds: Option<u64>,
/// Organization context.

pub org_id: Option<String>,
/// Department context.

pub department_id: Option<String>
}

#[async_trait]
pub trait FlowersBridgeContract: Send + Sync {
    /// Submits an agent for durable execution.
    ///
    /// The agent's LLM and tool calls will be journaled for crash
    /// recovery. On process restart, the execution resumes from the
    /// last committed journal entry.
    ///
    /// # Arguments
    ///
    /// * `request` - The submission parameters.
    ///
    /// # Returns
    ///
    /// A handle to the created execution.
    ///
    /// # Errors
    ///
    /// - `ContractError::FlowersError` if submission fails.
    /// - `ContractError::AgentNotFound` if the agent_id is invalid.
    async fn submit_agent(
        &self,
        request: FlowersSubmitRequest,
    ) -> ContractResult<FlowersExecutionHandle>;

    /// Sends a signal to a running or suspended execution.
    ///
    /// # Arguments
    ///
    /// * `execution_id` - The target execution.
    /// * `signal` - The signal to deliver.
    ///
    /// # Errors
    ///
    /// - `ContractError::FlowersError` if the execution does not exist
    ///   or cannot receive signals.
    async fn send_signal(&self, execution_id: &str, signal: FlowersSignal) -> ContractResult<()>;

    /// Cancels a running or suspended execution.
    ///
    /// Cancellation triggers compensation if a `CompensationStack` is
    /// registered for the execution.
    ///
    /// # Arguments
    ///
    /// * `execution_id` - The execution to cancel.
    /// * `reason` - Optional cancellation reason.
    async fn cancel_execution(
        &self,
        execution_id: &str,
        reason: Option<&str>,
    ) -> ContractResult<()>;

    /// Returns the journal entries for an execution.
    ///
    /// # Arguments
    ///
    /// * `execution_id` - The execution to query.
    /// * `from_index` - Start reading from this index.
    /// * `limit` - Maximum entries to return.
    ///
    /// # Returns
    ///
    /// Journal entries in sequential order.
    async fn get_journal(
        &self,
        execution_id: &str,
        from_index: u64,
        limit: u32,
    ) -> ContractResult<Vec<FlowersJournalEntry>>;

    /// Creates a checkpoint of the current execution state.
    ///
    /// # Arguments
    ///
    /// * `execution_id` - The execution to checkpoint.
    ///
    /// # Returns
    ///
    /// The checkpoint data.
    async fn create_checkpoint(&self, execution_id: &str) -> ContractResult<CheckpointData>;

    /// Returns the current status of an execution.
    ///
    /// # Arguments
    ///
    /// * `execution_id` - The execution to query.
    async fn get_execution_status(
        &self,
        execution_id: &str,
    ) -> ContractResult<FlowersExecutionHandle>;

    /// Lists all executions for an organization.
    ///
    /// # Arguments
    ///
    /// * `org_id` - The organization to query.
    /// * `state_filter` - Optional state filter.
    /// * `limit` - Maximum results.
    async fn list_executions(
        &self,
        org_id: &str,
        state_filter: Option<FlowersExecutionState>,
        limit: u32,
    ) -> ContractResult<Vec<FlowersExecutionHandle>>;
}

identity.rs

Read declaration text · 6 declaration entries

pub const CONTRACT_VERSION: &str;

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct IdentityHandle {
/// The OAS DID string (e.g., "did:oas:l1fe:agent:code-reviewer").

pub did: String,
/// Depth in the lineage chain (0 = HMR root, 1 = direct child, etc.).

pub lineage_depth: u32,
/// The OAS namespace (e.g., "l1fe").

pub namespace: String,
/// The entity name within the DID (e.g., "code-reviewer").

pub entity_name: String
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DerivedIdentityRequest {
/// The parent agent's DID from which to derive.

pub parent_did: String,
/// The name for the child agent identity.

pub child_name: String,
/// The OAS namespace for the child DID.

pub namespace: String,
/// Optional org ID for organizational context.

pub org_id: Option<String>,
/// Optional department ID for organizational context.

pub department_id: Option<String>,
/// Maximum allowed lineage depth. Derivation fails if this would

/// be exceeded. Default: 16.

pub max_lineage_depth: u32
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct LineageInfo {
/// The DID of the identity being verified.

pub did: String,
/// Depth in the lineage chain (0 = root).

pub depth: u32,
/// The root DID at the top of the chain (usually an HMR/MHR).

pub root_did: String,
/// Whether the full chain from root to this identity verifies.

pub chain_valid: bool,
/// Each hop in the chain from root to this identity.

pub chain: Vec<LineageHop>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct LineageHop {
/// The parent DID at this hop.

pub parent_did: String,
/// The child DID derived at this hop.

pub child_did: String,
/// The derivation path used.

pub derivation_path: String,
/// Whether this individual hop's signature verifies.

pub signature_valid: bool
}

#[async_trait]
pub trait IdentityContract: Send + Sync {
    /// Creates a new root identity (HMR) for an organization.
    ///
    /// This is called once per organization founding. The HMR is the
    /// cryptographic root from which all agent identities are derived.
    ///
    /// # Arguments
    ///
    /// * `namespace` - The OAS namespace (e.g., "l1fe").
    /// * `root_name` - The root identity name (e.g., "root-holder").
    ///
    /// # Returns
    ///
    /// A handle to the created root identity.
    async fn create_root_identity(
        &self,
        namespace: &str,
        root_name: &str,
    ) -> ContractResult<IdentityHandle>;

    /// Derives a child identity from an existing parent.
    ///
    /// The child's Ed25519 keypair is deterministically derived via
    /// HKDF-SHA256 from the parent's keypair and the derivation path.
    ///
    /// # Arguments
    ///
    /// * `request` - The derivation parameters.
    ///
    /// # Returns
    ///
    /// A handle to the derived child identity.
    ///
    /// # Errors
    ///
    /// - `ContractError::IdentityDerivationFailed` if derivation fails.
    /// - `ContractError::IdentityDerivationFailed` if max lineage depth
    ///   would be exceeded.
    async fn derive_identity(
        &self,
        request: DerivedIdentityRequest,
    ) -> ContractResult<IdentityHandle>;

    /// Verifies the lineage chain of an identity.
    ///
    /// Walks the chain from the given DID back to its root and verifies
    /// every hop's cryptographic proof.
    ///
    /// # Arguments
    ///
    /// * `did` - The DID to verify.
    ///
    /// # Returns
    ///
    /// Lineage information including chain validity.
    ///
    /// # Errors
    ///
    /// - `ContractError::DidResolutionFailed` if the DID cannot be resolved.
    /// - `ContractError::LineageVerificationFailed` if any hop fails.
    async fn verify_lineage(&self, did: &str) -> ContractResult<LineageInfo>;

    /// Resolves a DID to its public identity information.
    ///
    /// For local identities, this is immediate. For remote identities,
    /// this may involve network resolution.
    ///
    /// # Arguments
    ///
    /// * `did` - The DID string to resolve.
    ///
    /// # Returns
    ///
    /// The identity handle with public information.
    ///
    /// # Errors
    ///
    /// - `ContractError::DidResolutionFailed` if resolution fails.
    async fn resolve_did(&self, did: &str) -> ContractResult<IdentityHandle>;

    /// Signs arbitrary data with the specified identity's private key.
    ///
    /// # Arguments
    ///
    /// * `did` - The DID of the signing identity.
    /// * `data` - The data to sign.
    ///
    /// # Returns
    ///
    /// The Ed25519 signature bytes.
    ///
    /// # Errors
    ///
    /// - `ContractError::DidResolutionFailed` if the DID is not local.
    async fn sign(&self, did: &str, data: &[u8]) -> ContractResult<Vec<u8>>;

    /// Verifies a signature against a DID's public key.
    ///
    /// # Arguments
    ///
    /// * `did` - The DID of the alleged signer.
    /// * `data` - The data that was signed.
    /// * `signature` - The signature to verify.
    ///
    /// # Returns
    ///
    /// `true` if the signature is valid.
    ///
    /// # Errors
    ///
    /// - `ContractError::DidResolutionFailed` if the DID cannot be resolved.
    async fn verify(&self, did: &str, data: &[u8], signature: &[u8]) -> ContractResult<bool>;
}

lib.rs

Read declaration text · 25 declaration entries

pub mod auth;

pub mod brew;

pub mod comm;

pub mod error;

pub mod flowers;

pub mod identity;

pub mod mcp;

pub mod memory;

pub mod provider;

pub mod runtime;

pub mod telemetry;

pub const CONTRACTS_VERSION: &str;

pub fn contracts_major_version() -> u32;

pub mod prelude;

pub use crate::auth::{
        AuthContract, AuthDecision, CapabilityNarrowingRequest, DelegationChainEntry,
    };

pub use crate::brew::{BrewContract, PlanExecutionResult, PlanHandle};

pub use crate::comm::{ChannelConfig, ChannelHandle, CommContract, SessionConfig};

pub use crate::error::ContractError;

pub use crate::flowers::{CheckpointData, FlowersBridgeContract, FlowersExecutionHandle};

pub use crate::identity::{
        DerivedIdentityRequest, IdentityContract, IdentityHandle, LineageInfo,
    };

pub use crate::mcp::{McpContract, McpServerHandle, McpToolDescriptor};

pub use crate::memory::{
        MemoryContract, MemoryQuery, MemoryRecord, MemoryScope, MemoryWritePolicy,
    };

pub use crate::provider::{
        OrgProviderPolicy, ProviderContract, ProviderFallbackStrategy, ProviderSessionHandle,
    };

pub use crate::runtime::{
        AgentCreateRequest, AgentHandle, AgentLifecycleCommand, AgentRuntimeContract, AgentStatus,
    };

pub use crate::telemetry::{
        HealthSummary, OrgHealthContract, OrgTelemetryContract, SpanFilter,
    };

mcp.rs

Read declaration text · 9 declaration entries

pub const CONTRACT_VERSION: &str;

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct McpServerConfig {
/// Unique name for this MCP server connection.

pub server_name: String,
/// The transport type for connecting to the server.

pub transport: McpTransportType,
/// Optional authentication configuration.

pub auth: Option<McpAuthConfig>,
/// Organization that owns this connection.

pub org_id: String,
/// Optional department scope (only agents in this department can use).

pub department_id: Option<String>
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum McpTransportType {
    /// Standard I/O transport (subprocess).
    Stdio {
        /// Command to execute.
        command: String,
        /// Command arguments.
        args: Vec<String>,
    },
    /// Server-Sent Events over HTTP.
    Sse {
        /// The SSE endpoint URL.
        url: String,
    },
    /// Streamable HTTP transport.
    Http {
        /// The HTTP endpoint URL.
        url: String,
    },
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct McpAuthConfig {
/// The authentication method.

pub method: McpAuthMethod
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum McpAuthMethod {
    /// No authentication.
    None,
    /// API key authentication.
    ApiKey {
        /// Header name for the API key.
        header: String,
    },
    /// OAuth 2.0 with PKCE.
    OAuth {
        /// Authorization endpoint URL.
        auth_url: String,
        /// Token endpoint URL.
        token_url: String,
        /// Client ID.
        client_id: String,
        /// Scopes to request.
        scopes: Vec<String>,
    },
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct McpServerHandle {
/// Unique connection identifier.

pub connection_id: String,
/// The server name.

pub server_name: String,
/// Whether the connection is currently alive.

pub connected: bool,
/// Number of tools available from this server.

pub tool_count: u32,
/// Number of resources available from this server.

pub resource_count: u32
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct McpToolDescriptor {
/// The tool name.

pub name: String,
/// Human-readable description.

pub description: String,
/// JSON Schema for the tool's input parameters.

pub input_schema: serde_json::Value,
/// The MCP server that provides this tool.

pub server_name: String
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct McpResourceDescriptor {
/// The resource URI.

pub uri: String,
/// Human-readable name.

pub name: String,
/// Resource description.

pub description: Option<String>,
/// MIME type of the resource.

pub mime_type: Option<String>,
/// The MCP server that provides this resource.

pub server_name: String
}

#[async_trait]
pub trait McpContract: Send + Sync {
    /// Connects to an external MCP server.
    ///
    /// # Arguments
    ///
    /// * `config` - The server connection configuration.
    ///
    /// # Returns
    ///
    /// A handle to the connected server.
    ///
    /// # Errors
    ///
    /// - `ContractError::McpError` if the connection or handshake fails.
    async fn connect_server(&self, config: McpServerConfig) -> ContractResult<McpServerHandle>;

    /// Disconnects from an MCP server.
    ///
    /// # Arguments
    ///
    /// * `connection_id` - The connection to close.
    async fn disconnect_server(&self, connection_id: &str) -> ContractResult<()>;

    /// Discovers all tools available from a connected server.
    ///
    /// # Arguments
    ///
    /// * `connection_id` - The server to query.
    ///
    /// # Returns
    ///
    /// Tool descriptors from the server.
    async fn discover_tools(&self, connection_id: &str) -> ContractResult<Vec<McpToolDescriptor>>;

    /// Discovers all resources available from a connected server.
    ///
    /// # Arguments
    ///
    /// * `connection_id` - The server to query.
    ///
    /// # Returns
    ///
    /// Resource descriptors from the server.
    async fn discover_resources(
        &self,
        connection_id: &str,
    ) -> ContractResult<Vec<McpResourceDescriptor>>;

    /// Invokes a tool on a connected MCP server.
    ///
    /// # Arguments
    ///
    /// * `connection_id` - The server hosting the tool.
    /// * `tool_name` - The tool to invoke.
    /// * `arguments` - The tool's input arguments as JSON.
    /// * `agent_did` - The agent invoking the tool (for ACT checks).
    ///
    /// # Returns
    ///
    /// The tool's output as JSON.
    ///
    /// # Errors
    ///
    /// - `ContractError::McpError` if the tool invocation fails.
    /// - `ContractError::AuthorizationDenied` if the agent lacks the
    ///   required tool scope.
    async fn invoke_tool(
        &self,
        connection_id: &str,
        tool_name: &str,
        arguments: serde_json::Value,
        agent_did: &str,
    ) -> ContractResult<serde_json::Value>;

    /// Reads a resource from a connected MCP server.
    ///
    /// # Arguments
    ///
    /// * `connection_id` - The server hosting the resource.
    /// * `uri` - The resource URI.
    ///
    /// # Returns
    ///
    /// The resource content as JSON.
    async fn read_resource(
        &self,
        connection_id: &str,
        uri: &str,
    ) -> ContractResult<serde_json::Value>;

    /// Lists all connected MCP servers for an organization.
    ///
    /// # Arguments
    ///
    /// * `org_id` - The organization to query.
    async fn list_servers(&self, org_id: &str) -> ContractResult<Vec<McpServerHandle>>;
}

memory.rs

Read declaration text · 10 declaration entries

pub const CONTRACT_VERSION: &str;

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum MemoryScope {
    /// Organization-wide scope (all agents can read).
    Organization {
        /// The organization identifier.
        org_id: String,
    },
    /// Department-level scope.
    Department {
        /// The organization identifier.
        org_id: String,
        /// The department identifier.
        department_id: String,
    },
    /// Team-level scope.
    Team {
        /// The organization identifier.
        org_id: String,
        /// The department identifier.
        department_id: String,
        /// The team identifier.
        team_id: String,
    },
    /// Agent-private scope.
    Agent {
        /// The agent's DID.
        agent_did: String,
    },
}

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

pub fn display_scope(&self) -> String;

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum MemoryType {
    /// Temporal events and interaction histories.
    Episodic,
    /// Facts, relationships, and knowledge.
    Semantic,
    /// Workflows, runbooks, and procedures.
    Procedural,
    /// Documents, artifacts, and code snippets.
    Resource,
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MemoryRecord {
/// The scope where this record should be stored.

pub scope: MemoryScope,
/// The type of memory.

pub memory_type: MemoryType,
/// A unique key within the scope (for retrieval and updates).

pub key: String,
/// The memory content as JSON.

pub content: serde_json::Value,
/// Tags for discovery and filtering.

pub tags: Vec<String>,
/// The DID of the agent that authored this record.

pub author_did: String
}

#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct MemoryQuery {
/// The scope to query. The query also includes all parent scopes

/// (unless `include_parents` is `false`).

pub scope: Option<MemoryScope>,
/// Filter by memory type.

pub memory_type: Option<MemoryType>,
/// Filter by key prefix.

pub key_prefix: Option<String>,
/// Filter by tags (records must have ALL specified tags).

pub tags: Vec<String>,
/// Semantic search query (uses vector similarity).

pub semantic_query: Option<String>,
/// Maximum results to return.

pub limit: Option<u32>,
/// Whether to include records from parent scopes. Default: true.

pub include_parents: bool
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MemoryResult {
/// The record's unique key.

pub key: String,
/// The scope this record belongs to.

pub scope: MemoryScope,
/// The memory type.

pub memory_type: MemoryType,
/// The record content.

pub content: serde_json::Value,
/// Tags on this record.

pub tags: Vec<String>,
/// Who authored this record.

pub author_did: String,
/// When this record was created.

pub created_at: DateTime<Utc>,
/// When this record was last updated.

pub updated_at: DateTime<Utc>,
/// Similarity score (0.0 to 1.0) when using semantic search.

pub similarity: Option<f64>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MemoryWritePolicy {
/// The scope this policy applies to.

pub scope: MemoryScope,
/// Whether writes are allowed to this scope.

pub writes_allowed: bool,
/// DIDs allowed to write (empty = all agents in scope can write).

pub allowed_writers: Vec<String>,
/// Memory types that are writable (empty = all types).

pub writable_types: Vec<MemoryType>,
/// Whether writes require approval from a parent scope agent.

pub requires_approval: bool,
/// Maximum record size in bytes.

pub max_record_size_bytes: Option<u64>
}

#[async_trait]
pub trait MemoryContract: Send + Sync {
    /// Stores a memory record.
    ///
    /// # Arguments
    ///
    /// * `record` - The record to store.
    ///
    /// # Errors
    ///
    /// - `ContractError::MemoryAccessDenied` if the author lacks write
    ///   access to the specified scope.
    /// - `ContractError::MemoryError` if storage fails.
    async fn store(&self, record: MemoryRecord) -> ContractResult<()>;

    /// Retrieves memory records matching a query.
    ///
    /// # Arguments
    ///
    /// * `query` - The query criteria.
    /// * `requester_did` - The agent requesting the records (for access
    ///   control).
    ///
    /// # Returns
    ///
    /// Matching records ordered by relevance (semantic search) or
    /// recency (non-semantic queries).
    ///
    /// # Errors
    ///
    /// - `ContractError::MemoryAccessDenied` if the requester lacks
    ///   read access.
    async fn retrieve(
        &self,
        query: MemoryQuery,
        requester_did: &str,
    ) -> ContractResult<Vec<MemoryResult>>;

    /// Deletes a memory record.
    ///
    /// # Arguments
    ///
    /// * `scope` - The scope containing the record.
    /// * `key` - The record key.
    /// * `requester_did` - The agent requesting deletion (for access
    ///   control).
    ///
    /// # Errors
    ///
    /// - `ContractError::MemoryAccessDenied` if the requester lacks
    ///   write access.
    async fn delete(
        &self,
        scope: &MemoryScope,
        key: &str,
        requester_did: &str,
    ) -> ContractResult<()>;

    /// Sets the write policy for a memory scope.
    ///
    /// # Arguments
    ///
    /// * `policy` - The write policy to set.
    ///
    /// # Errors
    ///
    /// - `ContractError::AuthorizationDenied` if the caller lacks
    ///   authority to set policies for this scope.
    async fn set_write_policy(&self, policy: MemoryWritePolicy) -> ContractResult<()>;

    /// Returns the current write policy for a scope.
    ///
    /// # Arguments
    ///
    /// * `scope` - The scope to query.
    async fn get_write_policy(
        &self,
        scope: &MemoryScope,
    ) -> ContractResult<Option<MemoryWritePolicy>>;
}

provider.rs

Read declaration text · 8 declaration entries

pub const CONTRACT_VERSION: &str;

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct OrgProviderPolicy {
/// The organization this policy applies to.

pub org_id: String,
/// Ordered list of available providers. Lower priority number = preferred.

pub providers: Vec<ProviderEntry>,
/// Strategy for handling provider failures.

pub fallback_strategy: ProviderFallbackStrategy,
/// Department-level overrides. Key is department ID.

/// Overrides merge with (not replace) the org-level policy.

pub department_overrides: BTreeMap<String, DepartmentProviderOverride>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ProviderEntry {
/// Provider namespace (e.g., "openai", "anthropic", "local").

pub namespace: String,
/// Whether this provider is currently enabled.

pub enabled: bool,
/// Priority for routing (lower = preferred).

pub priority: u32,
/// Optional rate limit: max tokens per minute across all agents.

pub max_tokens_per_minute: Option<u64>,
/// Optional cost limit: max USD per hour.

pub max_cost_per_hour_usd: Option<f64>,
/// Allowed model names within this provider. Empty = all models.

pub allowed_models: Vec<String>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DepartmentProviderOverride {
/// Department identifier.

pub department_id: String,
/// Provider namespaces explicitly allowed for this department.

/// Empty means "inherit org policy."

pub allowed_providers: Vec<String>,
/// Provider namespaces explicitly blocked for this department.

pub blocked_providers: Vec<String>,
/// Optional department-level cost cap (USD per hour).

pub max_cost_per_hour_usd: Option<f64>
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum ProviderFallbackStrategy {
    /// Try the next provider in priority order.
    NextPriority,
    /// Fail immediately without trying alternatives.
    FailFast,
    /// Retry the same provider up to N times, then fail.
    RetryThenFail {
        /// Maximum number of retries.
        max_retries: u32,
    },
    /// Retry the same provider, then fall back to next priority.
    RetryThenFallback {
        /// Maximum retries before fallback.
        max_retries: u32,
    },
}

#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ProviderSessionHandle {
/// Unique session identifier.

pub session_id: String,
/// The provider namespace for this session.

pub provider_namespace: String,
/// The specific model in use.

pub model: String,
/// Tokens consumed in this session so far.

pub tokens_consumed: u64,
/// Estimated cost in USD so far.

pub estimated_cost_usd: f64
}

#[async_trait]
pub trait ProviderContract: Send + Sync {
    /// Sets or updates the organization-level provider policy.
    ///
    /// This replaces the entire policy for the given organization.
    ///
    /// # Arguments
    ///
    /// * `policy` - The complete provider policy.
    ///
    /// # Errors
    ///
    /// - `ContractError::ConfigurationError` if the policy is invalid.
    async fn set_org_policy(&self, policy: OrgProviderPolicy) -> ContractResult<()>;

    /// Resolves the best provider for a given request.
    ///
    /// Considers org policy, department overrides, agent capabilities,
    /// current rate limits, and cost budgets.
    ///
    /// # Arguments
    ///
    /// * `org_id` - The organization making the request.
    /// * `department_id` - Optional department for override lookup.
    /// * `agent_did` - The requesting agent (for ACT scope checks).
    /// * `requested_provider` - Optional provider preference (e.g., "anthropic:claude-sonnet-4-5-20250929").
    ///
    /// # Returns
    ///
    /// A handle to the created provider session.
    ///
    /// # Errors
    ///
    /// - `ContractError::NoProviderAvailable` if no provider matches.
    /// - `ContractError::AuthorizationDenied` if the agent lacks provider scopes.
    async fn resolve_provider(
        &self,
        org_id: &str,
        department_id: Option<&str>,
        agent_did: &str,
        requested_provider: Option<&str>,
    ) -> ContractResult<ProviderSessionHandle>;

    /// Returns the current status of a provider session.
    ///
    /// # Arguments
    ///
    /// * `session_id` - The session to query.
    ///
    /// # Returns
    ///
    /// The session handle with current usage statistics.
    async fn get_session(&self, session_id: &str) -> ContractResult<ProviderSessionHandle>;

    /// Closes a provider session, releasing resources.
    ///
    /// # Arguments
    ///
    /// * `session_id` - The session to close.
    async fn close_session(&self, session_id: &str) -> ContractResult<()>;

    /// Returns usage statistics for an organization.
    ///
    /// # Arguments
    ///
    /// * `org_id` - The organization to query.
    ///
    /// # Returns
    ///
    /// Aggregate usage per provider namespace.
    async fn get_org_usage(
        &self,
        org_id: &str,
    ) -> ContractResult<BTreeMap<String, ProviderUsageStats>>;
}

#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct ProviderUsageStats {
/// Total tokens consumed.

pub total_tokens: u64,
/// Total estimated cost in USD.

pub total_cost_usd: f64,
/// Number of active sessions.

pub active_sessions: u32,
/// Number of requests in the current rate limit window.

pub requests_this_window: u64,
/// Number of requests that were rate-limited.

pub rate_limited_count: u64,
/// Number of requests that failed.

pub failure_count: u64
}

runtime.rs

Read declaration text · 10 declaration entries

pub const CONTRACT_VERSION: &str;

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentCreateRequest {
/// Human-readable name for the agent (e.g., "code-reviewer").

pub agent_name: String,
/// The provider:model reference (e.g., "anthropic:claude-sonnet-4-5-20250929").

pub provider_ref: String,
/// Optional system prompt prepended to every LLM call.

pub system_prompt: Option<String>,
/// Tool definitions available to this agent.

pub tools: Vec<ToolDefinition>,
/// Maximum tool loop steps before forced termination.

pub max_steps: u32,
/// Aut0 organization ID (org-level context).

pub org_id: Option<String>,
/// Aut0 department ID (department-level context).

pub department_id: Option<String>,
/// Role within the organization (e.g., "senior-reviewer", "pm").

pub role: Option<String>,
/// Parent agent DID for sub-agent derivation. When `None`, the agent

/// is a root-level agent derived from an HMR/MHR.

pub parent_agent_did: Option<String>,
/// Arbitrary key-value metadata for Aut0-specific context.

pub metadata: std::collections::BTreeMap<String, String>
}

#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct AgentHandle {

}

pub fn new(id: impl Into<String>, did: impl Into<String>) -> Self;

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

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

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum AgentLifecycleCommand {
    /// Transition from Initializing to Ready, then to Running.
    Start,
    /// Transition from Running to Paused.
    Pause,
    /// Transition from Paused to Running.
    Resume,
    /// Transition from any state to Terminated.
    Terminate {
        /// Optional reason for termination.
        reason: Option<String>,
    },
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentStatus {
/// The agent's OAS DID.

pub did: String,
/// Current ANVIL lifecycle state.

pub lifecycle_state: LifecycleState,
/// Current health profile snapshot.

pub health: HealthProfile,
/// Number of tool loop steps completed.

pub steps_completed: u32,
/// Number of tool loop steps remaining before max_steps.

pub steps_remaining: u32,
/// Whether the agent is currently executing a tool call.

pub executing_tool: bool,
/// The provider:model reference the agent is using.

pub provider_ref: String,
/// Aut0 metadata passed at creation.

pub metadata: std::collections::BTreeMap<String, String>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AgentRunResult {
/// The final text output from the agent, if any.

pub output_text: Option<String>,
/// Structured output as JSON, if the agent produced structured output.

pub output_json: Option<serde_json::Value>,
/// Total tool invocations during the run.

pub tool_invocations: u32,
/// Total LLM inference calls during the run.

pub inference_calls: u32,
/// Total tokens consumed (input + output).

pub total_tokens: u64,
/// Whether the agent terminated normally or was force-stopped.

pub terminated_normally: bool,
/// The final lifecycle state.

pub final_state: LifecycleState
}

#[async_trait]
pub trait AgentRuntimeContract: Send + Sync {
    /// Creates a new agent in the Forge runtime.
    ///
    /// The agent starts in `Initializing` state. Call `lifecycle_command`
    /// with `Start` to advance it to `Running`.
    ///
    /// # Arguments
    ///
    /// * `request` - The agent creation parameters including identity,
    ///   provider, tools, and Aut0-specific metadata.
    ///
    /// # Returns
    ///
    /// An opaque handle to the created agent.
    ///
    /// # Errors
    ///
    /// - `ContractError::AgentCreationFailed` if the configuration is invalid.
    /// - `ContractError::IdentityDerivationFailed` if identity cannot be derived.
    /// - `ContractError::NoProviderAvailable` if the provider ref cannot be resolved.
    async fn create_agent(&self, request: AgentCreateRequest) -> ContractResult<AgentHandle>;

    /// Issues a lifecycle command to an existing agent.
    ///
    /// # Arguments
    ///
    /// * `handle` - The agent to command.
    /// * `command` - The lifecycle transition to perform.
    ///
    /// # Returns
    ///
    /// The new status after the transition.
    ///
    /// # Errors
    ///
    /// - `ContractError::AgentNotFound` if the handle is invalid.
    /// - `ContractError::InvalidLifecycleTransition` if the transition
    ///   violates the ANVIL state machine.
    async fn lifecycle_command(
        &self,
        handle: &AgentHandle,
        command: AgentLifecycleCommand,
    ) -> ContractResult<AgentStatus>;

    /// Queries the current status of an agent.
    ///
    /// # Arguments
    ///
    /// * `handle` - The agent to query.
    ///
    /// # Returns
    ///
    /// A snapshot of the agent's lifecycle, health, and execution state.
    ///
    /// # Errors
    ///
    /// - `ContractError::AgentNotFound` if the handle is invalid.
    async fn get_status(&self, handle: &AgentHandle) -> ContractResult<AgentStatus>;

    /// Runs an agent to completion with the given input.
    ///
    /// This is a convenience method that starts the agent (if not already
    /// running), sends the input through the tool loop, and blocks until
    /// termination or max_steps.
    ///
    /// # Arguments
    ///
    /// * `handle` - The agent to run.
    /// * `input` - The user/task input string.
    ///
    /// # Returns
    ///
    /// The execution result including output, tool counts, and token usage.
    ///
    /// # Errors
    ///
    /// - `ContractError::AgentNotFound` if the handle is invalid.
    /// - `ContractError::InvalidLifecycleTransition` if the agent is in
    ///   a state that cannot transition to Running.
    async fn run_to_completion(
        &self,
        handle: &AgentHandle,
        input: &str,
    ) -> ContractResult<AgentRunResult>;

    /// Lists all active agents matching an optional filter.
    ///
    /// # Arguments
    ///
    /// * `org_id` - Filter by organization. `None` returns all.
    /// * `department_id` - Filter by department. `None` returns all in org.
    ///
    /// # Returns
    ///
    /// Handles for all matching agents.
    async fn list_agents(
        &self,
        org_id: Option<&str>,
        department_id: Option<&str>,
    ) -> ContractResult<Vec<AgentHandle>>;
}

telemetry.rs

Read declaration text · 9 declaration entries

pub const CONTRACT_VERSION: &str;

#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct SpanFilter {
/// Filter by agent DID.

pub agent_did: Option<String>,
/// Filter by span name prefix (e.g., "anvil.tool.").

pub name_prefix: Option<String>,
/// Filter by time range start (inclusive).

pub from: Option<DateTime<Utc>>,
/// Filter by time range end (exclusive).

pub to: Option<DateTime<Utc>>,
/// Filter by organization.

pub org_id: Option<String>,
/// Filter by department.

pub department_id: Option<String>,
/// Maximum number of spans to return.

pub limit: Option<u32>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct CollectedSpan {
/// Unique span identifier.

pub span_id: String,
/// Optional parent span ID.

pub parent_span_id: Option<String>,
/// The span name (e.g., "anvil.generate", "anvil.tool.invoke").

pub name: String,
/// The agent DID that emitted this span.

pub agent_did: String,
/// Start timestamp.

pub started_at: DateTime<Utc>,
/// End timestamp.

pub ended_at: Option<DateTime<Utc>>,
/// Duration in microseconds.

pub duration_us: Option<u64>,
/// Span attributes as key-value pairs.

pub attributes: BTreeMap<String, String>
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct AuditEntry {
/// Sequential entry index.

pub index: u64,
/// The agent DID that created this entry.

pub agent_did: String,
/// The event kind.

pub event_kind: String,
/// The event payload as JSON.

pub payload: serde_json::Value,
/// ISO 8601 timestamp.

pub timestamp: DateTime<Utc>,
/// Ed25519 signature (base64-encoded).

pub signature: String
}

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HealthSummary {
/// The scope identifier (org ID, department ID, or team ID).

pub scope_id: String,
/// The scope type.

pub scope_type: HealthScopeType,
/// Overall health status for this scope.

pub overall_status: HealthStatus,
/// Number of agents in each lifecycle state.

pub lifecycle_counts: BTreeMap<String, u32>,
/// Number of agents at each health level.

pub health_counts: HealthCounts,
/// Total active agents in this scope.

pub total_agents: u32,
/// Total tool invocations across all agents.

pub total_tool_invocations: u64,
/// Total inference calls across all agents.

pub total_inference_calls: u64,
/// Total tokens consumed across all agents.

pub total_tokens: u64,
/// Estimated total cost in USD.

pub estimated_cost_usd: f64,
/// When this summary was last computed.

pub computed_at: DateTime<Utc>
}

#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum HealthScopeType {
    /// Organization-wide summary.
    Organization,
    /// Department-level summary.
    Department,
    /// Team-level summary.
    Team,
}

#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct HealthCounts {
/// Number of agents in Healthy state.

pub healthy: u32,
/// Number of agents in Degraded state.

pub degraded: u32,
/// Number of agents in Critical state.

pub critical: u32
}

#[async_trait]
pub trait OrgTelemetryContract: Send + Sync {
    /// Queries collected spans with filtering.
    ///
    /// # Arguments
    ///
    /// * `filter` - Criteria for filtering spans.
    ///
    /// # Returns
    ///
    /// Matching spans ordered by start time.
    async fn query_spans(&self, filter: SpanFilter) -> ContractResult<Vec<CollectedSpan>>;

    /// Returns the audit trail for an agent.
    ///
    /// # Arguments
    ///
    /// * `agent_did` - The agent whose audit trail to retrieve.
    /// * `from_index` - Start reading from this entry index.
    /// * `limit` - Maximum entries to return.
    ///
    /// # Returns
    ///
    /// Audit entries in sequential order.
    async fn get_audit_trail(
        &self,
        agent_did: &str,
        from_index: u64,
        limit: u32,
    ) -> ContractResult<Vec<AuditEntry>>;

    /// Returns the audit trail for an entire organization.
    ///
    /// # Arguments
    ///
    /// * `org_id` - The organization to query.
    /// * `from` - Start time (inclusive).
    /// * `to` - End time (exclusive).
    /// * `limit` - Maximum entries to return.
    async fn get_org_audit_trail(
        &self,
        org_id: &str,
        from: DateTime<Utc>,
        to: DateTime<Utc>,
        limit: u32,
    ) -> ContractResult<Vec<AuditEntry>>;
}

#[async_trait]
pub trait OrgHealthContract: Send + Sync {
    /// Returns a health summary for an organizational scope.
    ///
    /// # Arguments
    ///
    /// * `org_id` - The organization.
    /// * `scope_type` - The scope level (org, department, or team).
    /// * `scope_id` - The scope identifier. For `Organization`, this
    ///   is the org ID. For `Department`, the department ID, etc.
    ///
    /// # Returns
    ///
    /// Aggregated health summary for all agents in the scope.
    async fn get_health_summary(
        &self,
        org_id: &str,
        scope_type: HealthScopeType,
        scope_id: &str,
    ) -> ContractResult<HealthSummary>;

    /// Returns health summaries for all departments in an organization.
    ///
    /// # Arguments
    ///
    /// * `org_id` - The organization to query.
    ///
    /// # Returns
    ///
    /// One summary per department.
    async fn get_all_department_health(&self, org_id: &str) -> ContractResult<Vec<HealthSummary>>;

    /// Returns the health status of a specific agent.
    ///
    /// # Arguments
    ///
    /// * `agent_did` - The agent to query.
    async fn get_agent_health(&self, agent_did: &str) -> ContractResult<HealthSummary>;
}

Continue

On this page