Forge documentation
Library referenceRust

forge-comm

ANVIL Communication Contract: message transport, envelopes, and protocol negotiation for the Forge SDK

ANVIL Communication Contract: message transport, envelopes, and protocol negotiation for the Forge SDK

Package contract

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

Import boundary

use forge_comm;

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.

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

pub mod error;

pub mod message;

pub mod noop;

pub mod protocol;

pub mod transport;

pub mod prelude;

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

pub use crate::error::{CommError, CommResult};

pub use crate::message::AgentMessage;

pub use crate::noop::NoopTransport;

pub use crate::protocol::{negotiate_protocol, ProtocolAccept, ProtocolOffer};

pub use crate::transport::MessageTransport;

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.

channel.rs

Read declaration text · 2 declaration entries

pub struct ChannelTransport {

}

pub fn new(capacity: usize) -> (Self, Self);

error.rs

Read declaration text · 2 declaration entries

#[derive(Debug, Error)]
pub enum CommError {
    /// Message serialization to JSON failed.
    ///
    /// This typically indicates that the message payload contains values that
    /// cannot be represented in JSON (e.g., NaN floats, circular references).
    ///
    /// See ANVIL Spec section 10.2 -- Agent Message Envelope.
    #[error("message serialization failed: {reason}")]
    SerializationFailed {
        /// A human-readable explanation of the serialization failure.
        reason: String,
    },

    /// Message deserialization from JSON failed.
    ///
    /// This indicates that the received bytes do not constitute a valid
    /// `AgentMessage` envelope. Common causes include malformed JSON,
    /// missing required fields, or incompatible schema versions.
    ///
    /// See ANVIL Spec section 10.2 -- Agent Message Envelope.
    #[error("message deserialization failed: {reason}")]
    DeserializationFailed {
        /// A human-readable explanation of the deserialization failure.
        reason: String,
    },

    /// The underlying transport mechanism failed.
    ///
    /// This covers network errors, I/O failures, and other transport-layer
    /// issues that prevent message delivery.
    ///
    /// See ANVIL Spec section 10.1 -- Communication Contract.
    #[error("transport failed: {reason}")]
    TransportFailed {
        /// A human-readable explanation of the transport failure.
        reason: String,
    },

    /// The transport is not connected and cannot send or receive messages.
    ///
    /// This is returned by transports that require an active connection
    /// (e.g., `NoopTransport::receive`).
    ///
    /// See ANVIL Spec section 10.1 -- Communication Contract.
    #[error("transport is not connected")]
    NotConnected,

    /// The communication channel has been closed.
    ///
    /// This is returned when the peer endpoint of a channel-based transport
    /// has been dropped, making further communication impossible.
    ///
    /// See ANVIL Spec section 10.3 -- Channel Lifecycle.
    #[error("communication channel is closed")]
    ChannelClosed,

    /// The message signature is invalid or cannot be verified.
    ///
    /// This indicates that the Ed25519 signature on the message does not
    /// match the claimed sender's public key, or the signature is malformed.
    ///
    /// See ANVIL Spec section 10.4 -- Message Integrity.
    #[error("invalid signature from sender '{sender}'")]
    SignatureInvalid {
        /// The OAS DID of the sender whose signature failed verification.
        sender: String,
    },

    /// Protocol negotiation between two agents failed.
    ///
    /// The offered protocol versions did not match any version supported
    /// by the receiving agent.
    ///
    /// See ANVIL Spec section 10.5 -- Protocol Negotiation.
    #[error("protocol negotiation failed for '{offered}': {reason}")]
    ProtocolNegotiationFailed {
        /// The protocol identifier that was offered.
        offered: String,
        /// A human-readable explanation of why negotiation failed.
        reason: String,
    },

    /// The message exceeds the maximum allowed size.
    ///
    /// Transports may enforce size limits to prevent resource exhaustion.
    /// The message must be split or the payload reduced.
    ///
    /// See ANVIL Spec section 10.2 -- Agent Message Envelope.
    #[error("message size {size} bytes exceeds maximum {max_size} bytes")]
    MessageTooLarge {
        /// The actual size of the message in bytes.
        size: usize,
        /// The maximum allowed size in bytes.
        max_size: usize,
    },

    /// A transport operation timed out.
    ///
    /// The operation did not complete within the specified duration.
    /// This may indicate network congestion, an unresponsive peer, or
    /// a misconfigured timeout value.
    ///
    /// See ANVIL Spec section 10.1 -- Communication Contract.
    #[error("transport operation timed out after {duration_ms}ms")]
    Timeout {
        /// The timeout duration in milliseconds.
        duration_ms: u64,
    },

    /// Replay protection rejected the received message.
    ///
    /// The envelope's nonce was already seen within the validity window, its
    /// timestamp was outside the configured clock-skew tolerance, or the
    /// envelope was a legacy wire format and legacy acceptance is disabled.
    /// The wrapped [`forge_core::replay::ReplayError`] carries the specific
    /// reason.
    ///
    /// See ANVIL Spec section 10.4 -- Message Integrity.
    #[error("replay protection rejected message: {0}")]
    ReplayDetected(#[from] forge_core::replay::ReplayError),
}

pub type CommResult<T> = Result<T, CommError>;

lib.rs

Read declaration text · 13 declaration entries

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

pub mod error;

pub mod message;

pub mod noop;

pub mod protocol;

pub mod transport;

pub mod prelude;

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

pub use crate::error::{CommError, CommResult};

pub use crate::message::AgentMessage;

pub use crate::noop::NoopTransport;

pub use crate::protocol::{negotiate_protocol, ProtocolAccept, ProtocolOffer};

pub use crate::transport::MessageTransport;

message.rs

Read declaration text · 12 declaration entries

pub const NONCE_LEN: usize;

#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct AgentMessage {
/// Unique message identifier (UUID v4).

pub id: String,
/// Correlation ID linking related messages in a conversation.

#[serde(skip_serializing_if = "Option::is_none")]
pub correlation_id: Option<String>,
/// Message ID this is a direct reply to.

#[serde(skip_serializing_if = "Option::is_none")]
pub reply_to: Option<String>,
/// Sender's OAS DID.

pub sender: String,
/// Recipient's OAS DID.

pub recipient: String,
/// Protocol identifier (dot-separated, e.g., `"anvil.task.v1"`).

pub protocol: String,
/// Message type within the protocol.

pub message_type: String,
/// JSON payload, opaque to the transport layer.

pub payload: serde_json::Value,
/// Ed25519 signature of the message (hex-encoded).

///

/// Computed over [`AgentMessage::signing_bytes`], which includes the

/// replay-protection `nonce` and `timestamp_ms` fields. Verification

/// uses the sender's public key resolved from their OAS DID document.

#[serde(skip_serializing_if = "Option::is_none")]
pub signature: Option<String>,
/// ISO 8601 creation timestamp (human-readable).

///

/// For replay protection use [`timestamp_ms`](Self::timestamp_ms); this

/// field is retained for logging and audit trails.

pub timestamp: String,
/// Unix-epoch milliseconds at send time. Covered by the signature.

///

/// `None` only for envelopes deserialized from a pre-replay-protection

/// wire format. Present on every message produced by [`AgentMessage::new`].

#[serde(skip_serializing_if = "Option::is_none")]
pub timestamp_ms: Option<i64>,
/// 16-byte random nonce, hex-encoded. Covered by the signature.

///

/// `None` only for envelopes deserialized from a pre-replay-protection

/// wire format. Present on every message produced by [`AgentMessage::new`].

#[serde(skip_serializing_if = "Option::is_none")]
pub nonce: Option<String>
}

pub fn new(
        sender: impl Into<String>,
        recipient: impl Into<String>,
        protocol: impl Into<String>,
        message_type: impl Into<String>,
        payload: serde_json::Value,
    ) -> Self;

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

pub fn is_signed(&self) -> bool;

pub fn has_replay_fields(&self) -> bool;

pub fn nonce_bytes(&self) -> Option<[u8; NONCE_LEN]>;

pub fn with_correlation_id(mut self, id: String) -> Self;

pub fn with_reply_to(mut self, id: String) -> Self;

pub fn with_replay_fields(mut self, timestamp_ms: i64, nonce: [u8; NONCE_LEN]) -> Self;

pub fn signing_bytes(&self) -> Vec<u8>;

pub fn validate_replay(
        &self,
        validator: &forge_core::replay::ReplayValidator,
    ) -> Result<(), forge_core::replay::ReplayError>;

noop.rs

Read declaration text · 1 declaration entries

#[derive(Debug, Clone, Copy, Default)]
pub struct NoopTransport;

protocol.rs

Read declaration text · 3 declaration entries

#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct ProtocolOffer {
/// The protocol identifier (e.g., `"anvil.task"`).

///

/// This identifies the protocol family without a version suffix.

pub protocol: String,
/// The versions the offering agent supports, ordered by preference.

///

/// Version strings follow the pattern `"v1"`, `"v2"`, etc. The first

/// version in the list is the most preferred by the offering agent.

pub versions: Vec<String>,
/// Optional extensions the offering agent supports.

///

/// Extensions are additional capabilities within the protocol, such as

/// `"streaming"`, `"compression"`, or `"batching"`. Both agents must

/// agree on extensions for them to be active.

pub extensions: Vec<String>
}

#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
pub struct ProtocolAccept {
/// The agreed-upon protocol identifier.

pub protocol: String,
/// The agreed-upon version string.

pub version: String,
/// The extensions active for this protocol session.

///

/// This is the full set of extensions from the offer, carried through

/// to the acceptance. In a full implementation, this would be intersected

/// with the accepting agent's supported extensions.

pub extensions: Vec<String>
}

pub fn negotiate_protocol(
    offered: &ProtocolOffer,
    supported_versions: &[&str],
) -> Result<ProtocolAccept, CommError>;

transport.rs

Read declaration text · 1 declaration entries

#[async_trait::async_trait]
pub trait MessageTransport: Send + Sync {
    /// Send a message to the recipient identified in the message envelope.
    ///
    /// The transport delivers the message to the recipient's receive queue.
    /// Delivery semantics (at-most-once, at-least-once, exactly-once) depend
    /// on the specific transport implementation.
    ///
    /// # Arguments
    ///
    /// * `message` - The agent message envelope to send.
    ///
    /// # Errors
    ///
    /// Returns [`CommError::TransportFailed`] if the underlying mechanism fails,
    /// [`CommError::ChannelClosed`] if the peer endpoint has been dropped, or
    /// [`CommError::MessageTooLarge`] if the message exceeds transport limits.
    async fn send(&self, message: AgentMessage) -> Result<(), CommError>;

    /// Receive the next available message.
    ///
    /// This method waits asynchronously until a message arrives or an error
    /// occurs. The specific blocking behavior depends on the transport
    /// implementation: channel transports suspend the current task, while
    /// network transports may perform I/O polling.
    ///
    /// # Errors
    ///
    /// Returns [`CommError::NotConnected`] if the transport is not active,
    /// [`CommError::ChannelClosed`] if the peer endpoint has been dropped, or
    /// [`CommError::TransportFailed`] for other transport-level failures.
    async fn receive(&self) -> Result<AgentMessage, CommError>;

    /// Receive the next available message and enforce replay protection.
    ///
    /// This is the preferred receive path for any caller that has already
    /// verified the sender's Ed25519 signature (or is about to, in a
    /// subsequent step). It delegates to [`receive`](Self::receive) to pull
    /// the next envelope off the wire, then calls
    /// [`AgentMessage::validate_replay`](crate::message::AgentMessage::validate_replay)
    /// with the provided `validator`.
    ///
    /// The replay validator enforces two invariants (see
    /// [`forge_core::replay::ReplayValidator`]):
    ///
    /// 1. The envelope's `timestamp_ms` must be within the validator's
    ///    configured clock-skew window.
    /// 2. The envelope's `nonce` must not have been seen before, within the
    ///    validity window.
    ///
    /// Legacy envelopes (no `nonce`/`timestamp_ms`) are rejected unless the
    /// validator has been configured with `accept_legacy = true` or the
    /// `FORGE_ACCEPT_LEGACY_MESSAGES=true` env flag is honored by the
    /// caller's [`forge_core::replay::ReplayConfig`].
    ///
    /// # Errors
    ///
    /// - Any error returned by [`receive`](Self::receive).
    /// - [`CommError::ReplayDetected`] if replay protection rejects the
    ///   message.
    async fn receive_validated(
        &self,
        validator: &ReplayValidator,
    ) -> Result<AgentMessage, CommError> ;
}

Continue

On this page