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