feat(endpoint): initialize alknet-endpoint crate with all modules
- Create crates/alknet-endpoint/ with Cargo.toml, feature flags (quinn, iroh, tcp, acme) - Implement HandlerRegistry (registry.rs) — extracted from core/endpoint.rs - Implement AlknetEndpoint (endpoint.rs) — fresh build against ADR-083 shape - new() takes no StaticConfig, no TLS config - with_quinn/with_iroh/with_tcp_tls builder methods - run() spawns accept loops for each active transport - shutdown() is infallible (no Result) - No EndpointError type - Implement dispatch (dispatch.rs) — shared dispatch_connection free function - Transport-agnostic: takes pre-extracted alpn, fingerprint, remote_addr - ACME guard (acme-tls/1) behind acme feature - build_auth_context resolves identity from fingerprint - Implement accept loops: - accept/quinn.rs — quinn accept loop with ALPN + fingerprint extraction - accept/iroh.rs — iroh accept loop with ALPN negotiation + fingerprint - accept/tcp_tls.rs — TCP+TLS accept loop (new code, not in old endpoint.rs) - Add alknet-endpoint to workspace members - All feature combos compile cleanly (no features, quinn, iroh, tcp, all three) - Workspace cargo check passes
This commit is contained in:
1 parent
67d2f4affb
commit
2755b995c6
11 files changed
+616
No files matched your search
Generated
+14
@@ -97,6 +97,20 @@ dependencies = [
|
||||
"zeroize",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "alknet-endpoint"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"alknet-core",
|
||||
"arc-swap",
|
||||
"iroh",
|
||||
"quinn",
|
||||
"rustls",
|
||||
"tokio",
|
||||
"tokio-rustls",
|
||||
"tracing",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "alknet-http"
|
||||
version = "0.1.0"
|
||||
|
||||
@@ -3,6 +3,7 @@ members = [
|
||||
"crates/alknet-vault",
|
||||
"crates/alknet-core",
|
||||
"crates/alknet-call",
|
||||
"crates/alknet-endpoint",
|
||||
"crates/alknet-http",
|
||||
"crates/alknet-tls",
|
||||
"crates/alknet-tty",
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
[package]
|
||||
name = "alknet-endpoint"
|
||||
version.workspace = true
|
||||
edition.workspace = true
|
||||
license.workspace = true
|
||||
description = "Server-side multi-transport accept-loop runner — dispatches incoming connections by ALPN"
|
||||
repository.workspace = true
|
||||
|
||||
[lib]
|
||||
name = "alknet_endpoint"
|
||||
|
||||
[features]
|
||||
default = []
|
||||
quinn = ["dep:quinn", "alknet-core/quinn"]
|
||||
iroh = ["dep:iroh", "alknet-core/iroh"]
|
||||
tcp = ["dep:tokio-rustls"]
|
||||
acme = []
|
||||
|
||||
[dependencies]
|
||||
alknet-core = { path = "../alknet-core" }
|
||||
tokio = { version = "1", features = ["full"] }
|
||||
arc-swap = "1"
|
||||
tracing = "0.1"
|
||||
rustls = "0.23"
|
||||
quinn = { version = "0.11", optional = true }
|
||||
iroh = { version = "1.0", optional = true, default-features = false, features = ["tls-aws-lc-rs"] }
|
||||
tokio-rustls = { version = "0.26", optional = true }
|
||||
@@ -0,0 +1,72 @@
|
||||
//! Iroh accept loop.
|
||||
//!
|
||||
//! Accepts iroh connections, negotiates ALPN, extracts the client
|
||||
//! fingerprint (NodeId), converts to `Connection`, and dispatches.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use tokio::sync::watch;
|
||||
use tracing::{debug, warn};
|
||||
|
||||
use alknet_core::auth::IdentityProvider;
|
||||
use alknet_core::types::Connection;
|
||||
|
||||
use crate::registry::HandlerRegistry;
|
||||
|
||||
pub(crate) async fn run_accept_loop(
|
||||
iroh: iroh::Endpoint,
|
||||
handlers: Arc<HandlerRegistry>,
|
||||
identity_provider: Arc<dyn IdentityProvider>,
|
||||
shutdown_rx: &mut watch::Receiver<bool>,
|
||||
) {
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = shutdown_rx.changed() => {
|
||||
debug!("iroh accept loop: shutdown signaled");
|
||||
break;
|
||||
}
|
||||
incoming = iroh.accept() => {
|
||||
let Some(incoming) = incoming else {
|
||||
debug!("iroh accept loop: endpoint closed");
|
||||
break;
|
||||
};
|
||||
let handlers = handlers.clone();
|
||||
let identity_provider = identity_provider.clone();
|
||||
tokio::spawn(async move {
|
||||
let mut connecting = match incoming.accept() {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
warn!("iroh accept failed: {e}");
|
||||
return;
|
||||
}
|
||||
};
|
||||
let alpn = match connecting.alpn().await {
|
||||
Ok(alpn) => alpn,
|
||||
Err(e) => {
|
||||
warn!("iroh ALPN negotiation failed: {e}");
|
||||
return;
|
||||
}
|
||||
};
|
||||
let connection = match connecting.await {
|
||||
Ok(conn) => conn,
|
||||
Err(e) => {
|
||||
warn!("iroh handshake completion failed: {e}");
|
||||
return;
|
||||
}
|
||||
};
|
||||
let fingerprint = extract_client_fingerprint(&connection);
|
||||
let conn = Connection::from_iroh(connection);
|
||||
crate::dispatch::dispatch_connection(
|
||||
conn, alpn, fingerprint, None,
|
||||
&handlers, &identity_provider,
|
||||
);
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn extract_client_fingerprint(connection: &iroh::endpoint::Connection) -> Option<String> {
|
||||
let node_id = connection.remote_id();
|
||||
Some(format!("ed25519:{}", node_id))
|
||||
}
|
||||
@@ -0,0 +1,14 @@
|
||||
//! Transport-specific accept loops.
|
||||
//!
|
||||
//! Each module contains a `run_accept_loop` function that accepts
|
||||
//! connections on its transport, extracts ALPN + fingerprint, and
|
||||
//! calls `crate::dispatch::dispatch_connection`.
|
||||
|
||||
#[cfg(feature = "quinn")]
|
||||
pub(crate) mod quinn;
|
||||
|
||||
#[cfg(feature = "iroh")]
|
||||
pub(crate) mod iroh;
|
||||
|
||||
#[cfg(feature = "tcp")]
|
||||
pub(crate) mod tcp_tls;
|
||||
@@ -0,0 +1,83 @@
|
||||
//! Quinn (QUIC) accept loop.
|
||||
//!
|
||||
//! Accepts QUIC connections, performs the TLS handshake, extracts
|
||||
//! ALPN + client fingerprint, converts to `Connection`, and dispatches.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use tokio::sync::watch;
|
||||
use tracing::{debug, warn};
|
||||
|
||||
use alknet_core::auth::IdentityProvider;
|
||||
use alknet_core::types::Connection;
|
||||
|
||||
use crate::registry::HandlerRegistry;
|
||||
|
||||
pub(crate) async fn run_accept_loop(
|
||||
quinn: quinn::Endpoint,
|
||||
handlers: Arc<HandlerRegistry>,
|
||||
identity_provider: Arc<dyn IdentityProvider>,
|
||||
shutdown_rx: &mut watch::Receiver<bool>,
|
||||
) {
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = shutdown_rx.changed() => {
|
||||
debug!("quinn accept loop: shutdown signaled");
|
||||
break;
|
||||
}
|
||||
incoming = quinn.accept() => {
|
||||
let Some(incoming) = incoming else {
|
||||
debug!("quinn accept loop: endpoint closed");
|
||||
break;
|
||||
};
|
||||
let connecting = match incoming.accept() {
|
||||
Ok(c) => c,
|
||||
Err(e) => {
|
||||
warn!("quinn accept failed: {e}");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let handlers = handlers.clone();
|
||||
let identity_provider = identity_provider.clone();
|
||||
tokio::spawn(async move {
|
||||
let connection = match connecting.await {
|
||||
Ok(conn) => conn,
|
||||
Err(e) => {
|
||||
warn!("quinn TLS handshake failure: {e}");
|
||||
return;
|
||||
}
|
||||
};
|
||||
let alpn = extract_alpn(&connection);
|
||||
let remote_addr = Some(connection.remote_address());
|
||||
let fingerprint = extract_client_fingerprint(&connection);
|
||||
let conn = Connection::from_quinn_with_alpn(connection, alpn.clone());
|
||||
crate::dispatch::dispatch_connection(
|
||||
conn, alpn, fingerprint, remote_addr,
|
||||
&handlers, &identity_provider,
|
||||
);
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn extract_alpn(connection: &quinn::Connection) -> Vec<u8> {
|
||||
use quinn::crypto::rustls::HandshakeData;
|
||||
if let Some(data) = connection.handshake_data() {
|
||||
if let Ok(hs) = data.downcast::<HandshakeData>() {
|
||||
if let Some(protocol) = hs.protocol {
|
||||
return protocol;
|
||||
}
|
||||
}
|
||||
}
|
||||
Vec::new()
|
||||
}
|
||||
|
||||
fn extract_client_fingerprint(connection: &quinn::Connection) -> Option<String> {
|
||||
let identity = connection.peer_identity()?;
|
||||
let certs = identity
|
||||
.downcast::<Vec<rustls::pki_types::CertificateDer>>()
|
||||
.ok()?;
|
||||
let leaf = certs.first()?;
|
||||
alknet_core::fingerprint::fingerprint_from_cert_der(leaf.as_ref())
|
||||
}
|
||||
@@ -0,0 +1,74 @@
|
||||
//! TCP+TLS accept loop.
|
||||
//!
|
||||
//! Accepts TCP connections, performs a TLS handshake, extracts
|
||||
//! ALPN + client fingerprint, converts to `Connection::from_bidi`,
|
||||
//! and dispatches.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use tokio::sync::watch;
|
||||
use tracing::{debug, warn};
|
||||
|
||||
use alknet_core::auth::IdentityProvider;
|
||||
use alknet_core::types::Connection;
|
||||
|
||||
use crate::registry::HandlerRegistry;
|
||||
|
||||
pub(crate) async fn run_accept_loop(
|
||||
listener: tokio::net::TcpListener,
|
||||
acceptor: tokio_rustls::TlsAcceptor,
|
||||
handlers: Arc<HandlerRegistry>,
|
||||
identity_provider: Arc<dyn IdentityProvider>,
|
||||
shutdown_rx: &mut watch::Receiver<bool>,
|
||||
) {
|
||||
loop {
|
||||
tokio::select! {
|
||||
_ = shutdown_rx.changed() => {
|
||||
debug!("tcp+tls accept loop: shutdown signaled");
|
||||
break;
|
||||
}
|
||||
result = listener.accept() => {
|
||||
let (tcp_stream, remote_addr) = match result {
|
||||
Ok(r) => r,
|
||||
Err(e) => {
|
||||
warn!("tcp+tls accept failed: {e}");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let acceptor = acceptor.clone();
|
||||
let handlers = handlers.clone();
|
||||
let identity_provider = identity_provider.clone();
|
||||
tokio::spawn(async move {
|
||||
let tls_stream = match acceptor.accept(tcp_stream).await {
|
||||
Ok(s) => s,
|
||||
Err(e) => {
|
||||
warn!("tcp+tls TLS handshake failure: {e}");
|
||||
return;
|
||||
}
|
||||
};
|
||||
let (alpn, fingerprint) = extract_tls_session_info(&tls_stream);
|
||||
let conn = Connection::from_bidi(tls_stream, alpn.clone(), Some(remote_addr));
|
||||
crate::dispatch::dispatch_connection(
|
||||
conn, alpn, fingerprint, Some(remote_addr),
|
||||
&handlers, &identity_provider,
|
||||
);
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn extract_tls_session_info(
|
||||
tls_stream: &tokio_rustls::server::TlsStream<tokio::net::TcpStream>,
|
||||
) -> (Vec<u8>, Option<String>) {
|
||||
let (_, session) = tls_stream.get_ref();
|
||||
let alpn = session
|
||||
.alpn_protocol()
|
||||
.map(|a| a.to_vec())
|
||||
.unwrap_or_default();
|
||||
let fingerprint = session
|
||||
.peer_certificates()
|
||||
.and_then(|certs| certs.first())
|
||||
.and_then(|cert| alknet_core::fingerprint::fingerprint_from_cert_der(cert.as_ref()));
|
||||
(alpn, fingerprint)
|
||||
}
|
||||
@@ -0,0 +1,76 @@
|
||||
//! Shared dispatch path for all transports.
|
||||
//!
|
||||
//! `dispatch_connection` is the free function called by every accept loop
|
||||
//! after transport-specific extraction. `AlknetEndpoint::dispatch` delegates
|
||||
//! to it. `build_auth_context` resolves the caller's identity from the
|
||||
//! TLS fingerprint.
|
||||
|
||||
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
|
||||
use std::net::SocketAddr;
|
||||
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
|
||||
use std::sync::Arc;
|
||||
|
||||
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
|
||||
use tracing::{error, warn};
|
||||
|
||||
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
|
||||
use alknet_core::auth::{AuthContext, IdentityProvider};
|
||||
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
|
||||
use alknet_core::types::Connection;
|
||||
|
||||
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
|
||||
use crate::registry::HandlerRegistry;
|
||||
|
||||
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
|
||||
pub(crate) fn dispatch_connection(
|
||||
connection: Connection,
|
||||
alpn: Vec<u8>,
|
||||
fingerprint: Option<String>,
|
||||
remote_addr: Option<SocketAddr>,
|
||||
handlers: &HandlerRegistry,
|
||||
identity_provider: &Arc<dyn IdentityProvider>,
|
||||
) {
|
||||
#[cfg(feature = "acme")]
|
||||
if alpn == b"acme-tls/1" {
|
||||
tracing::debug!("acme-tls/1 challenge connection; closing");
|
||||
connection.close(0, "acme done");
|
||||
return;
|
||||
}
|
||||
|
||||
let handler = match handlers.get(&alpn) {
|
||||
Some(h) => h.clone(),
|
||||
None => {
|
||||
connection.close(0, "no handler");
|
||||
warn!(
|
||||
"dispatch: no handler for ALPN {:?}",
|
||||
String::from_utf8_lossy(&alpn)
|
||||
);
|
||||
return;
|
||||
}
|
||||
};
|
||||
|
||||
let auth = build_auth_context(&alpn, remote_addr, fingerprint, identity_provider);
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = handler.handle(connection, &auth).await {
|
||||
error!("handler returned error: {e}");
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
|
||||
pub(crate) fn build_auth_context(
|
||||
alpn: &[u8],
|
||||
remote_addr: Option<SocketAddr>,
|
||||
tls_client_fingerprint: Option<String>,
|
||||
identity_provider: &Arc<dyn IdentityProvider>,
|
||||
) -> AuthContext {
|
||||
let identity = tls_client_fingerprint
|
||||
.as_ref()
|
||||
.and_then(|fp| identity_provider.resolve_from_fingerprint(fp));
|
||||
AuthContext {
|
||||
identity,
|
||||
alpn: alpn.to_vec(),
|
||||
remote_addr,
|
||||
tls_client_fingerprint,
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,178 @@
|
||||
//! `AlknetEndpoint` — the central runtime type for accepting inbound connections.
|
||||
//!
|
||||
//! Takes pre-built transports via builder methods, runs their accept loops
|
||||
//! inside `run()`, and dispatches each accepted connection to the registered
|
||||
//! `ProtocolHandler` by ALPN.
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use arc_swap::ArcSwap;
|
||||
use tokio::sync::watch;
|
||||
|
||||
use alknet_core::auth::IdentityProvider;
|
||||
use alknet_core::config::DynamicConfig;
|
||||
|
||||
use crate::registry::HandlerRegistry;
|
||||
|
||||
#[cfg(feature = "tcp")]
|
||||
pub(crate) type TcpTlsListener = (tokio::net::TcpListener, tokio_rustls::TlsAcceptor);
|
||||
|
||||
pub struct AlknetEndpoint {
|
||||
#[cfg(feature = "quinn")]
|
||||
quinn: Option<quinn::Endpoint>,
|
||||
#[cfg(feature = "iroh")]
|
||||
iroh: Option<iroh::Endpoint>,
|
||||
#[cfg(feature = "tcp")]
|
||||
tcp_tls: std::sync::Mutex<Option<TcpTlsListener>>,
|
||||
handlers: Arc<HandlerRegistry>,
|
||||
#[allow(dead_code)]
|
||||
dynamic: Arc<ArcSwap<DynamicConfig>>,
|
||||
#[allow(dead_code)]
|
||||
identity_provider: Arc<dyn IdentityProvider>,
|
||||
shutdown_tx: watch::Sender<bool>,
|
||||
#[allow(dead_code)]
|
||||
shutdown_rx: watch::Receiver<bool>,
|
||||
drain_timeout: Duration,
|
||||
}
|
||||
|
||||
impl std::fmt::Debug for AlknetEndpoint {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
f.debug_struct("AlknetEndpoint")
|
||||
.field("handlers", &self.handlers)
|
||||
.field("drain_timeout", &self.drain_timeout)
|
||||
.finish()
|
||||
}
|
||||
}
|
||||
|
||||
impl AlknetEndpoint {
|
||||
pub fn new(
|
||||
handlers: HandlerRegistry,
|
||||
dynamic: Arc<ArcSwap<DynamicConfig>>,
|
||||
identity_provider: Arc<dyn IdentityProvider>,
|
||||
drain_timeout: Duration,
|
||||
) -> Self {
|
||||
let (shutdown_tx, shutdown_rx) = watch::channel(false);
|
||||
Self {
|
||||
#[cfg(feature = "quinn")]
|
||||
quinn: None,
|
||||
#[cfg(feature = "iroh")]
|
||||
iroh: None,
|
||||
#[cfg(feature = "tcp")]
|
||||
tcp_tls: std::sync::Mutex::new(None),
|
||||
handlers: Arc::new(handlers),
|
||||
dynamic,
|
||||
identity_provider,
|
||||
shutdown_tx,
|
||||
shutdown_rx,
|
||||
drain_timeout,
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "quinn")]
|
||||
pub fn with_quinn(mut self, endpoint: quinn::Endpoint) -> Self {
|
||||
self.quinn = Some(endpoint);
|
||||
self
|
||||
}
|
||||
|
||||
#[cfg(feature = "iroh")]
|
||||
pub fn with_iroh(mut self, endpoint: iroh::Endpoint) -> Self {
|
||||
self.iroh = Some(endpoint);
|
||||
self
|
||||
}
|
||||
|
||||
#[cfg(feature = "tcp")]
|
||||
pub fn with_tcp_tls(
|
||||
self,
|
||||
listener: tokio::net::TcpListener,
|
||||
acceptor: tokio_rustls::TlsAcceptor,
|
||||
) -> Self {
|
||||
*self.tcp_tls.lock().unwrap_or_else(|e| e.into_inner()) = Some((listener, acceptor));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn shutdown_sender(&self) -> watch::Sender<bool> {
|
||||
self.shutdown_tx.clone()
|
||||
}
|
||||
|
||||
pub async fn run(self: Arc<Self>) {
|
||||
#[allow(unused_mut)]
|
||||
let mut tasks: Vec<tokio::task::JoinHandle<()>> = Vec::new();
|
||||
|
||||
#[cfg(feature = "quinn")]
|
||||
if let Some(quinn) = &self.quinn {
|
||||
let quinn = quinn.clone();
|
||||
let handlers = self.handlers.clone();
|
||||
let identity_provider = self.identity_provider.clone();
|
||||
let mut shutdown_rx = self.shutdown_rx.clone();
|
||||
tasks.push(tokio::spawn(async move {
|
||||
crate::accept::quinn::run_accept_loop(
|
||||
quinn,
|
||||
handlers,
|
||||
identity_provider,
|
||||
&mut shutdown_rx,
|
||||
)
|
||||
.await;
|
||||
}));
|
||||
}
|
||||
|
||||
#[cfg(feature = "iroh")]
|
||||
if let Some(iroh) = &self.iroh {
|
||||
let iroh = iroh.clone();
|
||||
let handlers = self.handlers.clone();
|
||||
let identity_provider = self.identity_provider.clone();
|
||||
let mut shutdown_rx = self.shutdown_rx.clone();
|
||||
tasks.push(tokio::spawn(async move {
|
||||
crate::accept::iroh::run_accept_loop(
|
||||
iroh,
|
||||
handlers,
|
||||
identity_provider,
|
||||
&mut shutdown_rx,
|
||||
)
|
||||
.await;
|
||||
}));
|
||||
}
|
||||
|
||||
#[cfg(feature = "tcp")]
|
||||
if let Some((listener, acceptor)) = self.tcp_tls.lock().unwrap_or_else(|e| e.into_inner()).take() {
|
||||
let handlers = self.handlers.clone();
|
||||
let identity_provider = self.identity_provider.clone();
|
||||
let mut shutdown_rx = self.shutdown_rx.clone();
|
||||
tasks.push(tokio::spawn(async move {
|
||||
crate::accept::tcp_tls::run_accept_loop(
|
||||
listener,
|
||||
acceptor,
|
||||
handlers,
|
||||
identity_provider,
|
||||
&mut shutdown_rx,
|
||||
)
|
||||
.await;
|
||||
}));
|
||||
}
|
||||
|
||||
for task in tasks {
|
||||
let _ = task.await;
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn shutdown(&self) {
|
||||
let _ = self.shutdown_tx.send(true);
|
||||
|
||||
#[cfg(feature = "quinn")]
|
||||
if let Some(quinn) = &self.quinn {
|
||||
quinn.close(0u32.into(), b"shutdown");
|
||||
}
|
||||
|
||||
#[cfg(feature = "iroh")]
|
||||
if let Some(iroh) = &self.iroh {
|
||||
iroh.close().await;
|
||||
}
|
||||
|
||||
tokio::time::sleep(self.drain_timeout).await;
|
||||
|
||||
#[cfg(feature = "quinn")]
|
||||
if let Some(quinn) = &self.quinn {
|
||||
quinn.wait_idle().await;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
//! alknet-endpoint: Server-side multi-transport accept-loop runner.
|
||||
//!
|
||||
//! `AlknetEndpoint` takes pre-built transports (quinn, iroh, TCP+TLS) via
|
||||
//! builder methods, runs their accept loops inside `run()`, and dispatches
|
||||
//! each accepted connection to the registered `ProtocolHandler` by ALPN.
|
||||
//!
|
||||
//! The endpoint does not build transports and does not depend on
|
||||
//! `alknet-tls` — transport construction is the assembly layer's concern.
|
||||
|
||||
pub mod accept;
|
||||
pub mod dispatch;
|
||||
pub mod endpoint;
|
||||
pub mod registry;
|
||||
|
||||
pub use endpoint::AlknetEndpoint;
|
||||
pub use registry::HandlerRegistry;
|
||||
@@ -0,0 +1,61 @@
|
||||
//! `HandlerRegistry` — maps ALPN byte strings to `ProtocolHandler` instances.
|
||||
//!
|
||||
//! Registered statically at startup by the assembly layer; the endpoint
|
||||
//! dispatches by looking up the negotiated ALPN.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
|
||||
use alknet_core::types::ProtocolHandler;
|
||||
|
||||
pub struct HandlerRegistry {
|
||||
handlers: HashMap<&'static [u8], Arc<dyn ProtocolHandler>>,
|
||||
}
|
||||
|
||||
impl HandlerRegistry {
|
||||
pub fn new() -> Self {
|
||||
Self {
|
||||
handlers: HashMap::new(),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn register(&mut self, handler: Arc<dyn ProtocolHandler>) {
|
||||
let alpn = handler.alpn();
|
||||
if self.handlers.contains_key(alpn) {
|
||||
panic!(
|
||||
"HandlerRegistry: ALPN already registered: {:?}",
|
||||
String::from_utf8_lossy(alpn)
|
||||
);
|
||||
}
|
||||
self.handlers.insert(alpn, handler);
|
||||
}
|
||||
|
||||
pub fn get(&self, alpn: &[u8]) -> Option<&Arc<dyn ProtocolHandler>> {
|
||||
self.handlers.get(alpn)
|
||||
}
|
||||
|
||||
pub fn alpn_strings(&self) -> Vec<Vec<u8>> {
|
||||
self.handlers.keys().map(|k| k.to_vec()).collect()
|
||||
}
|
||||
}
|
||||
|
||||
impl Default for HandlerRegistry {
|
||||
fn default() -> Self {
|
||||
Self::new()
|
||||
}
|
||||
}
|
||||
|
||||
impl std::fmt::Debug for HandlerRegistry {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
f.debug_struct("HandlerRegistry")
|
||||
.field(
|
||||
"alpns",
|
||||
&self
|
||||
.handlers
|
||||
.keys()
|
||||
.map(|k| String::from_utf8_lossy(k).to_string())
|
||||
.collect::<Vec<_>>(),
|
||||
)
|
||||
.finish()
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user