miden-validator 0.17.2

Miden validator
use std::future::IntoFuture;
use std::time::Duration;

use anyhow::{Context, ensure};
use iroh::endpoint::{Connection, Incoming, Side};
use iroh::{Endpoint, EndpointAddr, EndpointId};
use miden_node_tracing::warn;
use miden_node_utils::retry::{self, Retryable};
use miden_protocol::crypto::dsa::ecdsa_k256_keccak::PublicKey;
use miden_validator::ValidatorSigner;
use rand_core_06::OsRng;

use super::super::wire::{RecvStream, SendStream};
use super::Ceremony;
use super::ceremony_config::CeremonyConfig;
use super::challenge::{Challenge, ChallengeResponse};
use super::completion::Completion;
use super::dkg::confirmation::DkgDealingsCommitment;
use super::dkg::{DealerMessages, DkgPublicKey, DkgRegistryRoot};
use super::session::{CeremonyNonce, SessionId};

/// An established Iroh connection with a verified endpoint identity but no validator identity.
pub struct ConnectedPeer {
    connection: Connection,
}

impl ConnectedPeer {
    /// Retries connection attempts with backoff until the peer connects or the endpoint closes.
    /// The caller must apply the ceremony timeout.
    ///
    /// Peers can start at different times. Authentication runs after this method returns and
    /// must not be retried.
    pub async fn connect(endpoint: &Endpoint, peer_addr: EndpointAddr) -> anyhow::Result<Self> {
        let peer_endpoint = peer_addr.id;
        let connection = (|| endpoint.connect(peer_addr.clone(), Ceremony::ALPN))
            .retry(retry::exponential(Duration::from_secs(1), Duration::from_secs(10)))
            .when(|_| !endpoint.is_closed())
            .notify(|error, delay| {
                warn!(
                    error,
                    target: miden_validator::LOG_TARGET,
                    "Retrying DKG peer connection",
                    dkg.peer.endpoint = peer_endpoint.to_string() #[nonstandard],
                    retry.delay_ms = delay.as_millis() as u64
                );
            })
            .await
            .with_context(|| format!("failed to connect to peer endpoint {peer_endpoint}"))?;
        Ok(Self { connection })
    }

    /// Establishes one incoming connection within ten seconds.
    ///
    /// The timeout releases capacity when a remote endpoint does not complete establishment.
    pub async fn accept(incoming: Incoming) -> anyhow::Result<Self> {
        let connection = tokio::time::timeout(Duration::from_secs(10), incoming.into_future())
            .await
            .context("incoming peer connection timed out")?
            .context("failed to establish incoming peer connection")?;
        Ok(Self { connection })
    }

    async fn bi_stream(&self) -> anyhow::Result<(SendStream, RecvStream)> {
        let (send, receive) = match self.connection.side() {
            Side::Client => {
                self.connection.open_bi().await.context("failed to open bidirectional stream")?
            },
            Side::Server => self
                .connection
                .accept_bi()
                .await
                .context("failed to accept bidirectional stream")?,
        };
        Ok((SendStream::from(send), RecvStream::from(receive)))
    }

    pub fn endpoint_id(&self) -> EndpointId {
        self.connection.remote_id()
    }

    pub fn close(self, reason: &[u8]) {
        self.connection.close(0u8.into(), reason);
    }

    /// Exchanges signed challenges and binds this connection to its configured validator key.
    ///
    /// Both sides send before reading at each exchange. The returned peer retains the same stream
    /// for subsequent ceremony messages, but no ceremony configuration has been exchanged yet.
    pub async fn authenticate(
        self,
        expected_validator_key: &PublicKey,
        signer: &ValidatorSigner,
    ) -> anyhow::Result<AuthenticatedPeer> {
        // Derive the tls-exporter channel binding defined in RFC 9266, section 2.
        let mut channel_binding = [0; 32];
        self.connection
            .export_keying_material(&mut channel_binding, b"EXPORTER-Channel-Binding", b"")
            .map_err(|error| {
                anyhow::anyhow!("failed to derive authentication channel binding: {error:?}")
            })?;
        let (mut send, mut receive) =
            self.bi_stream().await.context("failed to establish authentication stream")?;
        let challenge = Challenge::random(&mut OsRng);

        send.write(&challenge)
            .await
            .context("failed to send authentication challenge")?;

        let peer_challenge = receive
            .read_exact::<Challenge>(Challenge::BYTES)
            .await
            .context("failed to read authentication challenge")?;

        let response = peer_challenge.sign(signer, &channel_binding).await?;
        send.write(&response).await.context("failed to send authentication response")?;

        let response = receive
            .read_exact::<ChallengeResponse>(ChallengeResponse::BYTES)
            .await
            .context("failed to read challenge response")?;
        let validator_public_key = response.verify_against(&challenge, &channel_binding)?;
        ensure!(
            &validator_public_key == expected_validator_key,
            "peer validator key does not match the key configured for endpoint {}",
            self.endpoint_id(),
        );
        Ok(AuthenticatedPeer {
            validator_public_key,
            connection: self,
            send,
            receive,
        })
    }
}

/// A connection authenticated as its configured validator, with one ordered stream for the ceremony.
///
/// Both ends must call the exchange methods in protocol order. Buffered stream data allows one
/// peer to reach the next step while the other is still waiting for its remaining peers.
pub struct AuthenticatedPeer {
    validator_public_key: PublicKey,
    connection: ConnectedPeer,
    send: SendStream,
    receive: RecvStream,
}

impl AuthenticatedPeer {
    pub async fn exchange_ceremony_config(&mut self, local: &CeremonyConfig) -> anyhow::Result<()> {
        self.send.write(local).await.context("failed to send ceremony config")?;

        let peer_config = self
            .receive
            .read_exact::<CeremonyConfig>(CeremonyConfig::BYTES)
            .await
            .context("failed to read ceremony config")?;
        ensure!(
            &peer_config == local,
            "peer ceremony config does not match: local {local:?}, peer {peer_config:?}",
        );
        Ok(())
    }

    pub async fn exchange_ceremony_nonce(
        &mut self,
        local: &CeremonyNonce,
    ) -> anyhow::Result<CeremonyNonce> {
        self.send.write(local).await.context("failed to send ceremony nonce")?;
        let peer_nonce = self
            .receive
            .read_exact::<CeremonyNonce>(CeremonyNonce::BYTES)
            .await
            .context("failed to read ceremony nonce")?;
        Ok(peer_nonce)
    }

    pub async fn exchange_session_id(&mut self, local: &SessionId) -> anyhow::Result<SessionId> {
        self.send.write(local).await.context("failed to send session ID")?;
        let peer_session_id = self
            .receive
            .read_exact::<SessionId>(SessionId::BYTES)
            .await
            .context("failed to read session ID")?;
        Ok(peer_session_id)
    }

    pub async fn exchange_dkg_public_key(
        &mut self,
        local: &DkgPublicKey,
    ) -> anyhow::Result<DkgPublicKey> {
        self.send.write(local).await.context("failed to send DKG public key")?;
        let peer_dkg_public_key = self
            .receive
            .read_exact::<DkgPublicKey>(DkgPublicKey::BYTES)
            .await
            .context("failed to read DKG public key")?;
        Ok(peer_dkg_public_key)
    }

    pub async fn exchange_dkg_registry_root(
        &mut self,
        local: &DkgRegistryRoot,
    ) -> anyhow::Result<DkgRegistryRoot> {
        self.send.write(local).await.context("failed to send DKG registry root")?;
        let peer_registry_root = self
            .receive
            .read_exact::<DkgRegistryRoot>(DkgRegistryRoot::BYTES)
            .await
            .context("failed to read DKG registry root")?;
        Ok(peer_registry_root)
    }

    /// Exchanges both dealer messages using the local encoding length as the receive bound.
    ///
    /// The agreed registry and threshold give every dealer the same encoded message size. The
    /// peer's internal length prefixes therefore cannot request an unbounded receive allocation.
    pub async fn exchange_dealer_messages(
        &mut self,
        local: &DealerMessages,
    ) -> anyhow::Result<DealerMessages> {
        let message_bytes =
            self.send.write(local).await.context("failed to send local dealer messages")?;
        let peer_messages = self
            .receive
            .read_exact::<DealerMessages>(message_bytes)
            .await
            .context("failed to read peer dealer messages")?;
        Ok(peer_messages)
    }

    pub async fn exchange_dealings_commitment(
        &mut self,
        local: &DkgDealingsCommitment,
    ) -> anyhow::Result<DkgDealingsCommitment> {
        self.send.write(local).await.context("failed to send DKG dealings commitment")?;
        let peer_commitment = self
            .receive
            .read_exact::<DkgDealingsCommitment>(DkgDealingsCommitment::BYTES)
            .await
            .context("failed to read DKG dealings commitment")?;
        Ok(peer_commitment)
    }

    /// Exchanges the final completion message and waits for transport acknowledgement of our send.
    ///
    /// Finishing and draining the send stream keeps endpoint shutdown from discarding the local
    /// completion message. The ceremony separately checks that the received completion matches.
    pub async fn exchange_completion(&mut self, local: &Completion) -> anyhow::Result<Completion> {
        self.send.write(local).await.context("failed to send ceremony completion")?;
        self.send.finish()?;
        let peer_completion = self
            .receive
            .read_exact::<Completion>(Completion::BYTES)
            .await
            .context("failed to read peer ceremony completion")?;
        self.send.wait_for_delivery().await.with_context(|| {
            format!("failed to deliver ceremony completion to {}", self.connection.endpoint_id())
        })?;
        Ok(peer_completion)
    }

    pub fn validator_public_key(&self) -> &PublicKey {
        &self.validator_public_key
    }

    #[cfg(test)]
    pub fn connection(&self) -> &Connection {
        &self.connection.connection
    }

    #[cfg(test)]
    pub fn into_streams(
        self,
    ) -> (Connection, iroh::endpoint::SendStream, iroh::endpoint::RecvStream) {
        (self.connection.connection, self.send.into_inner(), self.receive.into_inner())
    }
}