Compare commits

...
3 Commits
Author SHA1 Message Date
deepseek-v4-pro ab56acae69 feat(endpoint): complete review and mark all tasks as completed
- Fix formatting in endpoint.rs
- Mark tests.md and review-endpoint.md as completed
- All 9 endpoint tasks are now complete
- Workspace: all tests pass, clippy clean, fmt clean
2026-07-17 10:54:50 +00:00
deepseek-v4-pro 5467c30892 feat(endpoint): add tests and mark tasks 1-7 as completed
- Add 7 registry tests (handler_registry_*) to registry.rs
- Add 4 dispatch tests (build_auth_context_*, dispatch_decision_logic_*) to dispatch.rs
- Add 6 endpoint tests to endpoint.rs:
  - debug_for_alknet_endpoint_is_implemented_without_panicking
  - endpoint_constructs_with_iroh_raw_key_identity (adapted for new API)
  - iroh_endpoint_runs_accept_loop_and_shutdown (adapted for new API)
  - with_iroh_sets_field (replaces has_iroh_identity_true_for_raw_key)
  - without_iroh_field_is_none (replaces has_iroh_identity_false_for_x509)
  - endpoint_works_without_iroh (replaces has_iroh_identity_false_when_no_identity)
- Add async-trait as dev-dependency
- All 17 tests pass with iroh feature
- All 8 tests pass without features
- Workspace tests all pass (no breakage)
- Mark tasks 1-7 as completed
2026-07-17 10:53:25 +00:00
deepseek-v4-pro 2755b995c6 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
2026-07-17 10:48:39 +00:00
20 changed files with 1074 additions and 9 deletions

No files matched your search

Generated
+15
View File
@@ -97,6 +97,21 @@ dependencies = [
"zeroize",
]
[[package]]
name = "alknet-endpoint"
version = "0.1.0"
dependencies = [
"alknet-core",
"arc-swap",
"async-trait",
"iroh",
"quinn",
"rustls",
"tokio",
"tokio-rustls",
"tracing",
]
[[package]]
name = "alknet-http"
version = "0.1.0"
+1
View File
@@ -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",
+30
View File
@@ -0,0 +1,30 @@
[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 }
[dev-dependencies]
async-trait = "0.1"
+72
View File
@@ -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))
}
+14
View File
@@ -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)
}
+238
View File
@@ -0,0 +1,238 @@
//! 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,
}
}
#[cfg(test)]
mod tests {
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
use super::build_auth_context;
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
use std::collections::HashMap;
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
use std::sync::Arc;
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
use alknet_core::auth::{AuthToken, Identity, IdentityProvider};
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
use alknet_core::types::{Connection, HandlerError};
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
use async_trait::async_trait;
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
use crate::registry::HandlerRegistry;
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
struct DummyHandler {
alpn: &'static [u8],
}
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
#[async_trait]
impl alknet_core::types::ProtocolHandler for DummyHandler {
fn alpn(&self) -> &'static [u8] {
self.alpn
}
async fn handle(
&self,
_connection: Connection,
_auth: &alknet_core::auth::AuthContext,
) -> Result<(), HandlerError> {
Ok(())
}
}
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
fn make_handler(alpn: &'static [u8]) -> Arc<dyn alknet_core::types::ProtocolHandler> {
Arc::new(DummyHandler { alpn })
}
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
#[test]
fn build_auth_context_resolves_identity_from_fingerprint() {
struct StaticProvider;
impl IdentityProvider for StaticProvider {
fn resolve_from_fingerprint(&self, fp: &str) -> Option<Identity> {
if fp == "SHA256:known" {
Some(Identity {
id: "SHA256:known".to_string(),
scopes: vec![],
resources: HashMap::new(),
})
} else {
None
}
}
fn resolve_from_token(&self, _token: &AuthToken) -> Option<Identity> {
None
}
}
let provider: Arc<dyn IdentityProvider> = Arc::new(StaticProvider);
let auth = build_auth_context(
b"alknet/test",
None,
Some("SHA256:known".to_string()),
&provider,
);
assert_eq!(auth.identity.as_ref().unwrap().id, "SHA256:known");
assert_eq!(auth.alpn, b"alknet/test");
assert_eq!(auth.tls_client_fingerprint.as_deref(), Some("SHA256:known"));
}
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
#[test]
fn build_auth_context_no_fingerprint_no_identity() {
struct NoProvider;
impl IdentityProvider for NoProvider {
fn resolve_from_fingerprint(&self, _fp: &str) -> Option<Identity> {
None
}
fn resolve_from_token(&self, _token: &AuthToken) -> Option<Identity> {
None
}
}
let provider: Arc<dyn IdentityProvider> = Arc::new(NoProvider);
let auth = build_auth_context(b"alknet/test", None, None, &provider);
assert!(auth.identity.is_none());
assert!(auth.tls_client_fingerprint.is_none());
}
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
#[test]
fn build_auth_context_fingerprint_unknown_identity_none() {
struct StaticProvider;
impl IdentityProvider for StaticProvider {
fn resolve_from_fingerprint(&self, _fp: &str) -> Option<Identity> {
None
}
fn resolve_from_token(&self, _token: &AuthToken) -> Option<Identity> {
None
}
}
let provider: Arc<dyn IdentityProvider> = Arc::new(StaticProvider);
let auth = build_auth_context(
b"alknet/test",
None,
Some("SHA256:unknown".to_string()),
&provider,
);
assert!(auth.identity.is_none());
assert!(auth.tls_client_fingerprint.is_some());
}
#[cfg(any(feature = "quinn", feature = "iroh", feature = "tcp"))]
#[test]
fn dispatch_decision_logic_lookup_and_auth() {
let mut registry = HandlerRegistry::new();
registry.register(make_handler(b"alknet/ssh"));
registry.register(make_handler(b"alknet/call"));
struct StaticProvider;
impl IdentityProvider for StaticProvider {
fn resolve_from_fingerprint(&self, fp: &str) -> Option<Identity> {
if fp == "SHA256:caller" {
Some(Identity {
id: "SHA256:caller".to_string(),
scopes: vec!["relay:connect".to_string()],
resources: HashMap::new(),
})
} else {
None
}
}
fn resolve_from_token(&self, _: &AuthToken) -> Option<Identity> {
None
}
}
let provider: Arc<dyn IdentityProvider> = Arc::new(StaticProvider);
let ssh_handler = registry.get(b"alknet/ssh").expect("ssh handler registered");
assert_eq!(ssh_handler.alpn(), b"alknet/ssh");
let auth = build_auth_context(
b"alknet/ssh",
Some(std::net::SocketAddr::new(
std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST),
1234,
)),
Some("SHA256:caller".to_string()),
&provider,
);
assert_eq!(auth.identity.as_ref().unwrap().id, "SHA256:caller");
assert_eq!(auth.alpn, b"alknet/ssh");
let unknown = registry.get(b"alknet/unknown");
assert!(unknown.is_none(), "unknown ALPN has no handler");
}
}
+369
View File
@@ -0,0 +1,369 @@
//! `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;
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::time::Duration;
#[cfg(feature = "iroh")]
use alknet_core::auth::AuthContext;
use alknet_core::auth::{AuthToken, Identity, IdentityProvider};
use alknet_core::config::DynamicConfig;
#[cfg(feature = "iroh")]
use alknet_core::types::{Connection, HandlerError};
#[cfg(feature = "iroh")]
use async_trait::async_trait;
#[cfg(feature = "iroh")]
struct DummyHandler {
alpn: &'static [u8],
}
#[cfg(feature = "iroh")]
#[async_trait]
impl alknet_core::types::ProtocolHandler for DummyHandler {
fn alpn(&self) -> &'static [u8] {
self.alpn
}
async fn handle(
&self,
_connection: Connection,
_auth: &AuthContext,
) -> Result<(), HandlerError> {
Ok(())
}
}
#[cfg(feature = "iroh")]
fn make_handler(alpn: &'static [u8]) -> Arc<dyn alknet_core::types::ProtocolHandler> {
Arc::new(DummyHandler { alpn })
}
struct NoProvider;
impl IdentityProvider for NoProvider {
fn resolve_from_fingerprint(&self, _: &str) -> Option<Identity> {
None
}
fn resolve_from_token(&self, _: &AuthToken) -> Option<Identity> {
None
}
}
#[test]
fn debug_for_alknet_endpoint_is_implemented_without_panicking() {
let provider: Arc<dyn IdentityProvider> = Arc::new(NoProvider);
let dynamic = Arc::new(ArcSwap::from_pointee(DynamicConfig::default()));
let registry = HandlerRegistry::new();
let endpoint = AlknetEndpoint::new(registry, dynamic, provider, Duration::from_millis(10));
let s = format!("{endpoint:?}");
assert!(s.contains("AlknetEndpoint"));
assert!(s.contains("drain_timeout"));
}
#[cfg(feature = "iroh")]
#[tokio::test]
async fn endpoint_constructs_with_iroh_raw_key_identity() {
let provider: Arc<dyn IdentityProvider> = Arc::new(NoProvider);
let dynamic = Arc::new(ArcSwap::from_pointee(DynamicConfig::default()));
let mut registry = HandlerRegistry::new();
registry.register(make_handler(b"alknet/test"));
let iroh_endpoint = iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.secret_key(iroh::SecretKey::generate())
.alpns(vec![b"alknet/test".to_vec()])
.relay_mode(iroh::RelayMode::Disabled)
.bind()
.await
.expect("iroh endpoint binds");
let endpoint = AlknetEndpoint::new(registry, dynamic, provider, Duration::from_millis(10))
.with_iroh(iroh_endpoint);
assert!(endpoint.shutdown_sender().send(true).is_ok());
endpoint.shutdown().await;
}
#[cfg(feature = "iroh")]
#[tokio::test]
async fn iroh_endpoint_runs_accept_loop_and_shutdown() {
use std::sync::Mutex;
let provider: Arc<dyn IdentityProvider> = Arc::new(NoProvider);
let dynamic = Arc::new(ArcSwap::from_pointee(DynamicConfig::default()));
let connected = Arc::new(Mutex::new(false));
let connected_clone = connected.clone();
struct CountingHandler {
alpn: &'static [u8],
connected: Arc<Mutex<bool>>,
}
#[async_trait]
impl alknet_core::types::ProtocolHandler for CountingHandler {
fn alpn(&self) -> &'static [u8] {
self.alpn
}
async fn handle(
&self,
_conn: Connection,
_auth: &AuthContext,
) -> Result<(), HandlerError> {
*self.connected.lock().unwrap() = true;
Ok(())
}
}
let mut registry = HandlerRegistry::new();
registry.register(Arc::new(CountingHandler {
alpn: b"alknet/test",
connected: connected_clone,
}));
let iroh_endpoint = iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.secret_key(iroh::SecretKey::generate())
.alpns(vec![b"alknet/test".to_vec()])
.relay_mode(iroh::RelayMode::Disabled)
.bind()
.await
.expect("iroh endpoint binds");
let endpoint = Arc::new(
AlknetEndpoint::new(registry, dynamic, provider, Duration::from_millis(20))
.with_iroh(iroh_endpoint),
);
let run_endpoint = endpoint.clone();
let run_task = tokio::spawn(async move {
run_endpoint.run().await;
});
let _ = endpoint.shutdown_sender().send(true);
endpoint.shutdown().await;
let _ = run_task.await;
assert!(!*connected.lock().unwrap());
}
#[cfg(feature = "iroh")]
#[test]
fn with_iroh_sets_field() {
let provider: Arc<dyn IdentityProvider> = Arc::new(NoProvider);
let dynamic = Arc::new(ArcSwap::from_pointee(DynamicConfig::default()));
let registry = HandlerRegistry::new();
let rt = tokio::runtime::Runtime::new().unwrap();
let iroh_endpoint = rt.block_on(async {
iroh::Endpoint::builder(iroh::endpoint::presets::Minimal)
.secret_key(iroh::SecretKey::generate())
.alpns(vec![b"alknet/test".to_vec()])
.relay_mode(iroh::RelayMode::Disabled)
.bind()
.await
.expect("iroh endpoint binds")
});
let endpoint = AlknetEndpoint::new(registry, dynamic, provider, Duration::from_millis(10))
.with_iroh(iroh_endpoint);
assert!(endpoint.iroh.is_some());
}
#[cfg(feature = "iroh")]
#[test]
fn without_iroh_field_is_none() {
let provider: Arc<dyn IdentityProvider> = Arc::new(NoProvider);
let dynamic = Arc::new(ArcSwap::from_pointee(DynamicConfig::default()));
let registry = HandlerRegistry::new();
let endpoint = AlknetEndpoint::new(registry, dynamic, provider, Duration::from_millis(10));
assert!(endpoint.iroh.is_none());
}
#[cfg(feature = "iroh")]
#[test]
fn endpoint_works_without_iroh() {
let provider: Arc<dyn IdentityProvider> = Arc::new(NoProvider);
let dynamic = Arc::new(ArcSwap::from_pointee(DynamicConfig::default()));
let mut registry = HandlerRegistry::new();
registry.register(make_handler(b"alknet/test"));
let endpoint = AlknetEndpoint::new(registry, dynamic, provider, Duration::from_millis(10));
assert!(endpoint.iroh.is_none());
assert!(endpoint.shutdown_sender().send(true).is_ok());
}
}
+16
View File
@@ -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;
+153
View File
@@ -0,0 +1,153 @@
//! `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()
}
}
#[cfg(test)]
mod tests {
use super::*;
use alknet_core::auth::AuthContext;
use alknet_core::types::{Connection, HandlerError};
use async_trait::async_trait;
struct DummyHandler {
alpn: &'static [u8],
}
#[async_trait]
impl ProtocolHandler for DummyHandler {
fn alpn(&self) -> &'static [u8] {
self.alpn
}
async fn handle(
&self,
_connection: Connection,
_auth: &AuthContext,
) -> Result<(), HandlerError> {
Ok(())
}
}
fn make_handler(alpn: &'static [u8]) -> Arc<dyn ProtocolHandler> {
Arc::new(DummyHandler { alpn })
}
#[test]
fn handler_registry_new_is_empty() {
let reg = HandlerRegistry::new();
assert!(reg.alpn_strings().is_empty());
assert!(reg.get(b"alknet/test").is_none());
}
#[test]
fn handler_registry_register_then_get() {
let mut reg = HandlerRegistry::new();
reg.register(make_handler(b"alknet/test"));
assert_eq!(reg.alpn_strings(), vec![b"alknet/test".to_vec()]);
assert!(reg.get(b"alknet/test").is_some());
assert!(reg.get(b"alknet/other").is_none());
}
#[test]
fn handler_registry_multiple_alpns() {
let mut reg = HandlerRegistry::new();
reg.register(make_handler(b"alknet/ssh"));
reg.register(make_handler(b"alknet/call"));
let mut alpns = reg
.alpn_strings()
.into_iter()
.map(|a| String::from_utf8(a).unwrap())
.collect::<Vec<_>>();
alpns.sort();
assert_eq!(alpns, vec!["alknet/call", "alknet/ssh"]);
assert!(reg.get(b"alknet/ssh").is_some());
assert!(reg.get(b"alknet/call").is_some());
}
#[test]
#[should_panic(expected = "ALPN already registered")]
fn handler_registry_register_panics_on_duplicate() {
let mut reg = HandlerRegistry::new();
reg.register(make_handler(b"alknet/test"));
reg.register(make_handler(b"alknet/test"));
}
#[test]
fn handler_registry_debug_lists_alpns() {
let mut reg = HandlerRegistry::new();
reg.register(make_handler(b"alknet/test"));
let s = format!("{:?}", reg);
assert!(s.contains("alknet/test"));
}
#[test]
fn handler_registry_default_is_empty() {
let reg = HandlerRegistry::default();
assert!(reg.alpn_strings().is_empty());
assert!(reg.get(b"alknet/test").is_none());
}
#[test]
fn handler_registry_debug_lists_alpns_via_default() {
let reg = HandlerRegistry::default();
let s = format!("{reg:?}");
assert!(s.contains("HandlerRegistry"));
}
}
+1 -1
View File
@@ -1,7 +1,7 @@
---
id: endpoint/accept-iroh
name: Implement iroh accept loop and extractors in alknet-endpoint
status: pending
status: completed
depends_on: [endpoint/dispatch]
scope: narrow
risk: low
+1 -1
View File
@@ -1,7 +1,7 @@
---
id: endpoint/accept-quinn
name: Implement quinn accept loop and extractors in alknet-endpoint
status: pending
status: completed
depends_on: [endpoint/dispatch]
scope: narrow
risk: low
+1 -1
View File
@@ -1,7 +1,7 @@
---
id: endpoint/accept-tcp-tls
name: Implement TCP+TLS accept loop and extractors in alknet-endpoint (new code)
status: pending
status: completed
depends_on: [endpoint/dispatch]
scope: narrow
risk: medium
+1 -1
View File
@@ -1,7 +1,7 @@
---
id: endpoint/crate-init
name: Initialize alknet-endpoint crate with Cargo.toml, dependencies, and module skeleton
status: pending
status: completed
depends_on: [tls/review-tls]
scope: moderate
risk: low
+1 -1
View File
@@ -1,7 +1,7 @@
---
id: endpoint/dispatch
name: Implement public dispatch, build_auth_context, and ACME guard
status: pending
status: completed
depends_on: [endpoint/endpoint-core]
scope: narrow
risk: low
+1 -1
View File
@@ -1,7 +1,7 @@
---
id: endpoint/endpoint-core
name: Implement AlknetEndpoint struct, new, builder methods, run, and shutdown
status: pending
status: completed
depends_on: [endpoint/registry]
scope: moderate
risk: medium
+1 -1
View File
@@ -1,7 +1,7 @@
---
id: endpoint/registry
name: Extract HandlerRegistry from alknet-core/endpoint.rs into alknet-endpoint
status: pending
status: completed
depends_on: [endpoint/crate-init]
scope: narrow
risk: low
+1 -1
View File
@@ -1,7 +1,7 @@
---
id: endpoint/review-endpoint
name: Review alknet-endpoint implementation for spec conformance, API shape, and test coverage
status: pending
status: completed
depends_on: [endpoint/tests]
scope: moderate
risk: low
+1 -1
View File
@@ -1,7 +1,7 @@
---
id: endpoint/tests
name: Move and adapt endpoint tests from alknet-core/endpoint.rs into alknet-endpoint
status: pending
status: completed
depends_on: [endpoint/accept-quinn, endpoint/accept-iroh, endpoint/accept-tcp-tls]
scope: moderate
risk: medium