don't use this until i actually make it good
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305//! 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<Self, NodeError> { 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<PeerSession, NodeError> { 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<PeerSession, NodeError> { 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<IncomingConnection, NodeError> { 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<PeerSession, NodeError> { 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<Mutex<SendStream>>, recv: Arc<Mutex<RecvStream>>,}
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<Mutex<Stream>> 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<Mutex<SendStream>>, 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<Mutex<RecvStream>>) -> Result<Vec<u8>, 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<Mutex<SendStream>>, Arc<Mutex<RecvStream>>) { 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<u8>| { 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<Mutex<SendStream>>, Arc<Mutex<RecvStream>>), NodeError> { let (send_fn, recv_fn, send, recv) = framed(send, recv); let mut hs: Handshake<SendFn, RecvFn> = 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<Mutex<SendStream>>, Arc<Mutex<RecvStream>>), NodeError> { let (send_fn, recv_fn, send, recv) = framed(send, recv); let mut hs: Handshake<SendFn, RecvFn> = 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<Message, NodeError> { let frame = read_frame(&self.recv).await?; Message::from_drisl(&frame).map_err(|e| NodeError::Codec(e.to_string())) }}