//! the iroh-facing node: endpoint bring-up, DRISL frame io over QUIC streams, //! and wiring the abstract [`Handshake`] onto a real connection. //! //! iroh's `Endpoint` is bound with the iwakura ALPN. outgoing connections run //! the initiator handshake; the accept loop runs the responder handshake. a //! completed handshake yields a [`PeerSession`]: the established connection //! plus the peer's verified identity, ready for message exchange. use std::pin::Pin; use std::sync::Arc; use ed25519_dalek::SigningKey; use iroh::endpoint::{presets, Connection}; use iroh::{Endpoint, EndpointAddr, EndpointId, SecretKey}; use thiserror::Error; use super::handshake::{Handshake, HandshakeError, PeerInfo, RecvFn, RecvFut, SendFn, SendFut, VerifyFn}; use super::{IWAKURA_ALPN, MAX_FRAME_LEN}; use crate::lexicon::{DidKey, EndpointRecord, Message}; #[derive(Error, Debug)] pub enum NodeError { #[error("failed to bind iroh endpoint: {0}")] Bind(String), #[error("failed to connect: {0}")] Connect(String), #[error("stream error: {0}")] Stream(String), #[error("handshake failed: {0}")] Handshake(#[from] HandshakeError), #[error("frame exceeds maximum length")] FrameTooLong, #[error("message codec error: {0}")] Codec(String), } /// an iroh endpoint bound for iwakura, ready to connect and accept. pub struct IwakuraNode { endpoint: Endpoint, } impl IwakuraNode { /// bind an iroh endpoint with the iwakura ALPN, using the n0 preset /// (default relay + hole-punching for NAT traversal). /// /// `device_key` is this node's iroh identity. it MUST be the same key whose /// public half is named in our endpoint record's `sub` — otherwise the /// handshake presents a record authorizing a different endpoint than the /// one actually connected, and every peer rejects it. this is why bind /// takes the key explicitly rather than generating one. pub async fn bind(device_key: &SigningKey) -> Result { let secret = SecretKey::from_bytes(&device_key.to_bytes()); let endpoint = Endpoint::builder(presets::N0) .secret_key(secret) .alpns(vec![IWAKURA_ALPN.to_vec()]) .bind() .await .map_err(|e| NodeError::Bind(e.to_string()))?; Ok(Self { endpoint }) } /// this node's endpoint id (its iroh ed25519 public key). pub fn endpoint_id(&self) -> EndpointId { self.endpoint.id() } /// this node's endpoint id as a did:key. pub fn endpoint_did_key(&self) -> DidKey { DidKey::from_public_key_bytes(*self.endpoint_id().as_bytes()) } /// connect to a peer and run the initiator handshake. /// /// `our_record` authorizes our own endpoint. `verify_presence` is the repo /// presence check used to validate the peer's record. on success returns a /// [`PeerSession`] with the peer's verified identity. pub async fn connect( &self, peer: EndpointAddr, our_record: &EndpointRecord, verify_presence: &VerifyFn, ) -> Result { let conn = self .endpoint .connect(peer, IWAKURA_ALPN) .await .map_err(|e| NodeError::Connect(e.to_string()))?; let (send, recv) = conn.open_bi().await.map_err(|e| NodeError::Stream(e.to_string()))?; let peer_endpoint = DidKey::from_public_key_bytes(*conn.remote_id().as_bytes()); let (peer, send, recv) = run_handshake_initiator(send, recv, our_record, &peer_endpoint, verify_presence).await?; Ok(PeerSession { conn, peer, send, recv }) } /// accept the next incoming connection and run the responder handshake. /// /// blocks until a peer connects. `our_record` authorizes our own endpoint; /// `verify_presence` validates the initiator's record. pub async fn accept( &self, our_record: &EndpointRecord, verify_presence: &VerifyFn, ) -> Result { let incoming = self.accept_connection().await?; self.handshake_responder(incoming, our_record, verify_presence).await } /// accept the next incoming connection WITHOUT running the handshake. /// /// split out so a persistent accept loop can hand the connection to a /// spawned task for the handshake — a slow or failing handshake then never /// blocks the next accept. pair with [`IwakuraNode::handshake_responder`]. pub async fn accept_connection(&self) -> Result { let conn = self .endpoint .accept() .await .ok_or_else(|| NodeError::Connect("endpoint closed".into()))? .await .map_err(|e| NodeError::Connect(e.to_string()))?; let (send, recv) = conn.accept_bi().await.map_err(|e| NodeError::Stream(e.to_string()))?; Ok(IncomingConnection { conn, send, recv }) } /// run the responder handshake on an accepted connection, producing a /// session. see [`IwakuraNode::accept_connection`]. pub async fn handshake_responder( &self, incoming: IncomingConnection, our_record: &EndpointRecord, verify_presence: &VerifyFn, ) -> Result { let IncomingConnection { conn, send, recv } = incoming; let peer_endpoint = DidKey::from_public_key_bytes(*conn.remote_id().as_bytes()); let (peer, send, recv) = run_handshake_responder(send, recv, our_record, &peer_endpoint, verify_presence).await?; Ok(PeerSession { conn, peer, send, recv }) } /// gracefully close the endpoint. takes &self so a shared (Arc'd) node /// can be closed without consuming it. pub async fn close(&self) { self.endpoint.close().await; } } /// an accepted incoming connection, pre-handshake. see /// [`IwakuraNode::accept_connection`]. pub struct IncomingConnection { conn: Connection, send: SendStream, recv: RecvStream, } /// an established, mutually-authenticated session with a peer. /// /// the handshake's bidirectional stream is kept open and reused for all /// messages, so delivery is ordered and there's no per-message stream churn. /// `send`/`recv` are the two halves of that single stream, arc+mutex shared so /// a background recv loop and a foreground send loop can run concurrently. pub struct PeerSession { conn: Connection, /// the peer's verified identity from the handshake. pub peer: PeerInfo, send: Arc>, recv: Arc>, } impl PeerSession { /// the peer's iroh endpoint id. pub fn remote_id(&self) -> EndpointId { self.conn.remote_id() } /// wait until the peer closes the connection. pub async fn closed(&self) { self.conn.closed().await; } /// close the connection from our side. /// /// used when we decide the session should end (e.g. the peer's endpoint /// record lapsed during a ttl re-check), rather than waiting for the peer. pub async fn close(&self) { // error code 0 = no error; an orderly close. self.conn.close(0u32.into(), b"done"); } } // --- frame io over a QUIC stream ------------------------------------------- // // the handshake and message layers speak in opaque DRISL frames. these // adapters wrap an iroh send/recv stream pair into the closure shape the // abstract `Handshake` expects, applying the 4-byte LE length-prefix framing. use iroh::endpoint::{RecvStream, SendStream}; use tokio::sync::Mutex; // the streams are not Clone and the boxed futures must be 'static, so each // closure holds an Arc> and clones the Arc into every future. // that gives the future owned, 'static access to the stream under a lock. the // same Arc'd halves are handed to the PeerSession afterwards, so the handshake // stream is reused for messages. /// write one length-prefixed frame to a stream half. async fn write_frame(send: &Arc>, frame: &[u8]) -> Result<(), NodeError> { let len = u32::try_from(frame.len()).map_err(|_| NodeError::FrameTooLong)?; let mut guard = send.lock().await; guard.write_all(&len.to_le_bytes()).await.map_err(|e| NodeError::Stream(e.to_string()))?; guard.write_all(frame).await.map_err(|e| NodeError::Stream(e.to_string()))?; Ok(()) } /// read one length-prefixed frame from a stream half. async fn read_frame(recv: &Arc>) -> Result, NodeError> { let mut guard = recv.lock().await; let mut len_buf = [0u8; 4]; guard.read_exact(&mut len_buf).await.map_err(|e| NodeError::Stream(e.to_string()))?; let len = u32::from_le_bytes(len_buf); if len > MAX_FRAME_LEN { return Err(NodeError::FrameTooLong); } let mut buf = vec![0u8; len as usize]; guard.read_exact(&mut buf).await.map_err(|e| NodeError::Stream(e.to_string()))?; Ok(buf) } /// the framed send/recv closures over a stream pair, plus the Arc'd halves to /// keep for message io afterwards. fn framed( send: SendStream, recv: RecvStream, ) -> (SendFn, RecvFn, Arc>, Arc>) { let send = Arc::new(Mutex::new(send)); let recv = Arc::new(Mutex::new(recv)); let send_fn = { let send = Arc::clone(&send); Box::new(move |frame: Vec| { let send = Arc::clone(&send); Box::pin(async move { write_frame(&send, &frame).await.map_err(|e| e.to_string()) }) as SendFut }) as SendFn }; let recv_fn = { let recv = Arc::clone(&recv); Box::new(move || { let recv = Arc::clone(&recv); Box::pin(async move { read_frame(&recv).await.map_err(|e| e.to_string()) }) as RecvFut }) as RecvFn }; (send_fn, recv_fn, send, recv) } async fn run_handshake_initiator( send: SendStream, recv: RecvStream, our_record: &EndpointRecord, peer_endpoint: &DidKey, verify_presence: &VerifyFn, ) -> Result<(PeerInfo, Arc>, Arc>), NodeError> { let (send_fn, recv_fn, send, recv) = framed(send, recv); let mut hs: Handshake = Handshake::new(send_fn, recv_fn); let peer = hs.run_initiator(our_record, peer_endpoint, verify_presence).await?; Ok((peer, send, recv)) } async fn run_handshake_responder( send: SendStream, recv: RecvStream, our_record: &EndpointRecord, peer_endpoint: &DidKey, verify_presence: &VerifyFn, ) -> Result<(PeerInfo, Arc>, Arc>), NodeError> { let (send_fn, recv_fn, send, recv) = framed(send, recv); let mut hs: Handshake = Handshake::new(send_fn, recv_fn); let peer = hs.run_responder(our_record, peer_endpoint, verify_presence).await?; Ok((peer, send, recv)) } // --- message io -------------------------------------------------------------- // // messages flow over the SAME long-lived bi stream the handshake used, with // the same length-prefix framing. one stream, both directions, ordered. impl PeerSession { /// send a message to the peer on the session's long-lived stream. pub async fn send_message(&self, message: &Message) -> Result<(), NodeError> { let frame = message.to_drisl().map_err(|e| NodeError::Codec(e.to_string()))?; write_frame(&self.send, &frame).await } /// read the next message from the session's long-lived stream. /// /// returns an error when the peer closes the stream (or on a stream /// error); callers treat that as end-of-conversation. pub async fn recv_message(&self) -> Result { let frame = read_frame(&self.recv).await?; Message::from_drisl(&frame).map_err(|e| NodeError::Codec(e.to_string())) } }