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