Weave documentation
Rust referenceweave-sdk

weave-sdk · replication

Source declarations, signatures and documentation for replication.

Source: sigil/weave/libs/weave-sdk/src/replication.rs. SHA-256: 29b2a49d34c219793056def49ac590d299e88d51d3477898d706a581a4316e4c.

This reference follows declared source modules, retains conditional attributes, and includes public declarations and implementation methods. Private-module re-exports and trait resolution require the compiler; this is a source reference, not a claim that every listed item is a root import. Function bodies and constant values are omitted.

replication::ReplicationMessage

Message types exchanged during strand replication.

Each message is length-prefixed and JSON-serialized on the wire.

#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type")]
pub enum ReplicationMessage {
    /// Anonymous resource preflight. Private repositories attach a one-time
    /// proof derived from the current repository epoch and Noise channel.
    /// No DID or authentication public key is disclosed at this stage.
    #[serde(rename = "hello")]
    Hello {
        version: u8,
        discovery_key: Vec<u8>,
        live: bool,
        #[serde(default)]
        access_proof: Vec<u8>,
    },

    /// Initial handshake: exchange public keys and discovery key.
    #[serde(rename = "handshake")]
    Handshake {
        version: u8,
        public_key: Vec<u8>,
        discovery_key: Vec<u8>,
        live: bool,
        #[serde(default)]
        peer_did: Option<String>,
        #[serde(default)]
        auth_public_key: Option<Vec<u8>>,
        #[serde(default)]
        auth_signature: Option<Vec<u8>>,
    },

    /// Announce available block range.
    #[serde(rename = "have")]
    Have { start: u64, length: u64 },

    /// Request the remote to announce its availability from `start`.
    #[serde(rename = "want")]
    Want { start: u64 },

    /// Request a range of blocks.
    #[serde(rename = "request")]
    Request { start: u64, end: u64 },

    /// Deliver a block.
    #[serde(rename = "data")]
    Data {
        seq: u64,
        #[serde(with = "base64_bytes")]
        data: Vec<u8>,
        header: WireHeader,
    },

    /// Cancel a previous request.
    #[serde(rename = "cancel")]
    Cancel { start: u64, end: u64 },

    /// Indicate we do not have a block.
    #[serde(rename = "unhave")]
    Unhave { seq: u64 },

    /// Synchronization complete (non-live sessions).
    #[serde(rename = "sync")]
    Sync,
}

Source line: 115.

replication::WireHeader

Minimal header representation for wire transfer.

#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WireHeader {
pub version: u8,
pub seq: u64,
pub length: u64,
pub byte_offset: u64,
#[serde(default, with = "base64_bytes")]
pub prev_hash: Vec<u8>,
pub writer_public_key: Vec<u8>,
pub timestamp: u64,
#[serde(with = "base64_bytes")]
pub signature: Vec<u8>
}

Source line: 179.

replication::base64_bytes::serialize

pub fn serialize<S>(bytes: &Vec<u8>, serializer: S) -> std::result::Result<S::Ok, S::Error>
    where
        S: Serializer,;

Source line: 241.

replication::base64_bytes::deserialize

pub fn deserialize<'de, D>(deserializer: D) -> std::result::Result<Vec<u8>, D::Error>
    where
        D: Deserializer<'de>,;

Source line: 250.

replication::SessionState

State of a replication session.

#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SessionState {
    /// Handshake in progress.
    Connecting,
    /// Actively replicating.
    Active,
    /// Paused (e.g., waiting for data).
    Idle,
    /// Session closed.
    Closed,
}

Source line: 334.

replication::ReplicationStats

Statistics for a replication session.

#[derive(Debug, Clone, Default)]
pub struct ReplicationStats {
pub blocks_uploaded: u64,
pub blocks_downloaded: u64,
pub bytes_uploaded: u64,
pub bytes_downloaded: u64
}

Source line: 347.

replication::AuthenticatedPeer

A remote peer whose handshake signature verified.

#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AuthenticatedPeer {
/// Claimed and verified DID.

pub did: String,
/// Ed25519 authentication public key used to verify the handshake.

pub auth_public_key: [u8; 32]
}

Source line: 356.

replication::ReplicationAdmissionPolicy

Admission callback for authenticated replication peers.

pub type ReplicationAdmissionPolicy =
    Arc<dyn Fn(&AuthenticatedPeer, &[u8; 32]) -> bool + Send + Sync>;

Source line: 364.

replication::ReplicationAccessProof

Channel-bound anonymous preflight for private replication.

The generator is used by an outbound peer before it discloses its DID. The verifier is used by the responder before it sends its own signed handshake. Public repositories leave this control unset.

#[derive(Clone)]
pub struct ReplicationAccessProof {

}

Source line: 373.

replication::ReplicationAccessProof::new

Construct a private replication preflight.

pub fn new(
        generator: impl Fn(&[u8; 32], &[u8; 32], bool) -> Result<[u8; 32]> + Send + Sync + 'static,
        verifier: impl Fn(&[u8; 32], &[u8; 32], bool, &[u8; 32]) -> bool + Send + Sync + 'static,
    ) -> Self;

Source line: 380.

replication::ReplicationAccessProof::dogfood_shared_secret

Dogfood / same-host LAN only — shared-secret AccessProof for WEAVE_PRIVATE_REPLICATION secondary prove.

This is not a public DHT PASS path. Both peers must use the same secret. Prefer a strong random secret; do not use in production.

pub fn dogfood_shared_secret(secret: &[u8]) -> Self;

Source line: 395.

replication::dogfood_allow_authenticated_peers

Dogfood / same-host LAN only — admission policy that accepts any post-handshake authenticated peer. Pair with [ReplicationAccessProof::dogfood_shared_secret]. Not a public DHT PASS.

pub fn dogfood_allow_authenticated_peers() -> ReplicationAdmissionPolicy;

Source line: 426.

replication::IncomingStrandResolver

Resolver used by inbound multi-resource replication. The responder reads the peer's requested discovery key from the handshake, then asks this resolver for the matching local strand.

pub type IncomingStrandResolver =
    Arc<dyn Fn([u8; 32]) -> Option<(String, Arc<tokio::sync::RwLock<Strand>>)> + Send + Sync>;

Source line: 433.

replication::LocusReplicationTargets

Strand handles required to replicate a complete Locus drive.

#[derive(Clone)]
pub struct LocusReplicationTargets {
/// Metadata B-tree strand.

pub meta_strand: Arc<RwLock<Strand>>,
/// Blob metadata strand.

pub blob_meta_strand: Arc<RwLock<Strand>>,
/// Blob payload strand.

pub blob_data_strand: Arc<RwLock<Strand>>,
/// Keep the sessions open for live updates.

pub live: bool
}

Source line: 438.

replication::LocusReplicationTargets::new

Create replication targets for the three strands backing a Locus drive.

pub fn new(
        meta_strand: Arc<RwLock<Strand>>,
        blob_meta_strand: Arc<RwLock<Strand>>,
        blob_data_strand: Arc<RwLock<Strand>>,
        live: bool,
    ) -> Self;

Source line: 451.

replication::ReplicationSession

A single replication session with a remote peer.

#[derive(Debug)]
pub struct ReplicationSession {
/// The remote peer's DID.

pub peer_did: String,
/// Remote peer's public key.

pub peer_public_key: [u8; 32],
/// Names of strands being replicated.

pub strands: Vec<String>,
/// Session state.

pub state: SessionState,
/// Accumulated stats.

pub stats: ReplicationStats
}

Source line: 468.

replication::ReplicationManager

Manages peer-to-peer replication of strands between nodes.

This integrates three layers:

  1. Secret Stream (zer0-secret-stream) -- Noise-encrypted transport
  2. Proto-Mux (zer0-proto-mp) -- Channel multiplexing with flow control
  3. Strand replication messages -- Block exchange protocol
pub struct ReplicationManager {

}

Source line: 498.

replication::ReplicationManager::new

Create a new replication manager.

pub fn new() -> Self;

Source line: 508.

replication::ReplicationManager::with_identity

Create a replication manager that authenticates handshakes with identity.

pub fn with_identity(identity: WeaveIdentity) -> Self;

Source line: 518.

replication::ReplicationManager::set_private_controls

Atomically configure or clear the two controls required for a private replication endpoint.

A private endpoint must never expose only the epoch proof or only the post-identity policy: the former admits any current-key holder, while the latter discloses responder identity before authorization. This API therefore rejects incomplete configurations.

pub async fn set_private_controls(
        &self,
        proof: Option<ReplicationAccessProof>,
        policy: Option<ReplicationAdmissionPolicy>,
    ) -> Result<()>;

Source line: 534.

replication::ReplicationManager::private_controls_enabled

Whether this manager is operating a complete private endpoint.

pub async fn private_controls_enabled(&self) -> Result<bool>;

Source line: 569.

replication::ReplicationManager::private_discovery_authenticator

Produce an anonymous, current-epoch authenticator for one leased discovery provider record. The payload digest is domain-separated before being supplied as the proof channel binding, so private DHT records disclose no member DID or authentication key.

pub async fn private_discovery_authenticator(
        &self,
        discovery_key: &[u8; 32],
        payload_digest: &[u8; 32],
    ) -> Result<Option<[u8; 32]>>;

Source line: 578.

replication::ReplicationManager::verify_private_discovery_authenticator

Verify the anonymous authenticator on a leased private discovery record using the current repository epoch.

pub async fn verify_private_discovery_authenticator(
        &self,
        discovery_key: &[u8; 32],
        payload_digest: &[u8; 32],
        authenticator: &[u8; 32],
    ) -> Result<bool>;

Source line: 593.

replication::ReplicationManager::connect_and_replicate

Connect to a remote peer and begin replicating a strand.

This performs the full connection lifecycle:

  1. TCP connect to addr
  2. Noise XX handshake via SecretStream
  3. Open a channel keyed by the strand's discovery key
  4. Run the Weave block exchange protocol

The replication loop runs in a background tokio::spawn task. Use [Self::stop_replication] to tear it down.

pub async fn connect_and_replicate(
        &self,
        _peer_did: &str,
        _addr: SocketAddr,
        _strand_name: &str,
        _strand: Arc<tokio::sync::RwLock<Strand>>,
        _live: bool,
    ) -> Result<()>;

Source line: 622.

replication::ReplicationManager::connect_and_replicate_pinned

Connect only when both the remote DID and its Ed25519 authentication key match a separately authenticated descriptor.

pub async fn connect_and_replicate_pinned(
        &self,
        peer_did: &str,
        expected_peer_auth_key: [u8; 32],
        addr: SocketAddr,
        strand_name: &str,
        strand: Arc<tokio::sync::RwLock<Strand>>,
        live: bool,
    ) -> Result<()>;

Source line: 638.

replication::ReplicationManager::handle_incoming_replication

Accept an incoming TCP connection and run replication for strand.

This is the responder-side counterpart of [Self::connect_and_replicate]. Typically called from a listener loop.

pub async fn handle_incoming_replication(
        &self,
        peer_did: &str,
        stream: TcpStream,
        strand_name: &str,
        strand: Arc<tokio::sync::RwLock<Strand>>,
        live: bool,
    ) -> Result<()>;

Source line: 813.

replication::ReplicationManager::handle_incoming_replication_dynamic

Accept an incoming TCP connection and select the replicated strand from the discovery key in the peer's handshake.

pub async fn handle_incoming_replication_dynamic(
        &self,
        peer_did: &str,
        stream: TcpStream,
        live: bool,
        resolver: IncomingStrandResolver,
    ) -> Result<()>;

Source line: 931.

replication::ReplicationManager::listen_for_replication

Start a TCP listener that accepts replication connections for the given strand. Each accepted connection spawns a responder session.

Returns the local address the listener is bound to and a handle to the listener task.

pub async fn listen_for_replication(
        &self,
        bind_addr: SocketAddr,
        strand: Arc<tokio::sync::RwLock<Strand>>,
        live: bool,
    ) -> Result<(SocketAddr, JoinHandle<()>)>;

Source line: 1037.

replication::ReplicationManager::stop_replication

Stop replicating with a peer. Signals the background task to shut down.

pub async fn stop_replication(&self, peer_did: &str) -> Result<()>;

Source line: 1103.

replication::ReplicationManager::sessions

List all active replication sessions.

pub async fn sessions(&self) -> Vec<(String, SessionState)>;

Source line: 1121.

replication::ReplicationManager::is_session_live

Return whether peer_did currently has a connecting or authenticated replication task. Discovery uses this to forget closed sessions and permit a clean retry.

pub async fn is_session_live(&self, peer_did: &str) -> bool;

Source line: 1132.

replication::ReplicationManager::stats

Get stats for a peer.

pub async fn stats(&self, peer_did: &str) -> ReplicationStats;

Source line: 1147.

replication::ReplicationManager::replicate_lens

Replicate a Lens database to a peer.

A Lens is backed by a single strand, so this discovers the underlying strand and replicates it.

pub async fn replicate_lens(
        &self,
        peer_did: &str,
        peer_auth_public_key: [u8; 32],
        addr: SocketAddr,
        lens_name: &str,
        lens_strand: Arc<tokio::sync::RwLock<Strand>>,
        live: bool,
    ) -> Result<()>;

Source line: 1163.

replication::ReplicationManager::replicate_locus

Replicate a Locus drive to a peer.

A Locus drive is backed by 3 strands: meta-strand, blob-meta, and blob-data. All three must be replicated for a complete Locus copy.

Returns the number of strands successfully started.

pub async fn replicate_locus(
        &self,
        peer_did: &str,
        peer_auth_public_key: [u8; 32],
        addr: SocketAddr,
        locus_name: &str,
        targets: LocusReplicationTargets,
    ) -> Result<u32>;

Source line: 1195.

replication::ReplicationManager::replicate_nexus

Replicate a Nexus multi-writer view to a peer.

A Nexus aggregates multiple input strands. This method replicates each of the provided input strands to the peer.

Returns the number of strands successfully started.

pub async fn replicate_nexus(
        &self,
        peer_did: &str,
        peer_auth_public_key: [u8; 32],
        addr: SocketAddr,
        nexus_name: &str,
        input_strands: Vec<(String, Arc<tokio::sync::RwLock<Strand>>)>,
        live: bool,
    ) -> Result<u32>;

Source line: 1311.

On this page

replication::ReplicationMessagereplication::WireHeaderreplication::base64_bytes::serializereplication::base64_bytes::deserializereplication::SessionStatereplication::ReplicationStatsreplication::AuthenticatedPeerreplication::ReplicationAdmissionPolicyreplication::ReplicationAccessProofreplication::ReplicationAccessProof::newreplication::ReplicationAccessProof::dogfood_shared_secretreplication::dogfood_allow_authenticated_peersreplication::IncomingStrandResolverreplication::LocusReplicationTargetsreplication::LocusReplicationTargets::newreplication::ReplicationSessionreplication::ReplicationManagerreplication::ReplicationManager::newreplication::ReplicationManager::with_identityreplication::ReplicationManager::set_private_controlsreplication::ReplicationManager::private_controls_enabledreplication::ReplicationManager::private_discovery_authenticatorreplication::ReplicationManager::verify_private_discovery_authenticatorreplication::ReplicationManager::connect_and_replicatereplication::ReplicationManager::connect_and_replicate_pinnedreplication::ReplicationManager::handle_incoming_replicationreplication::ReplicationManager::handle_incoming_replication_dynamicreplication::ReplicationManager::listen_for_replicationreplication::ReplicationManager::stop_replicationreplication::ReplicationManager::sessionsreplication::ReplicationManager::is_session_livereplication::ReplicationManager::statsreplication::ReplicationManager::replicate_lensreplication::ReplicationManager::replicate_locusreplication::ReplicationManager::replicate_nexus