feat(core): Connection::from_stream — generic single-stream connections

Add ConnectionKind::Stream: a single read/write pair behind a
Mutex<Option<...>> that accept_bi yields once (then ConnectionClosed).
Add Connection::from_stream(send, recv, alpn, remote_addr) and
Connection::from_bidi(stream, alpn, remote_addr) constructors.

Rename the stream-level SendStreamKind::Mock / RecvStreamKind::Mock to
::Stream (they were already generic Box<dyn AsyncRead/Write> — the
wrong name). Rename from_mock to from_stream.

Remove MockConnection trait + ConnectionKind::Mock entirely. Migrate
all alknet-call test StubConnections to use from_stream with
tokio::io::sink() + tokio::io::empty() (immediate EOF on the read side
causes handle_stream to exit cleanly, then accept_bi returns
ConnectionClosed and the run_loop exits).

The yield-once accept_bi contract: QUIC yields many streams, everything
else yields one. Handlers that loop (TtyAdapter) get one iteration per
single-stream connection; handlers that call once (HttpAdapter) get the
stream directly. Both correct, no branching on transport.

This unblocks TCP+TLS listeners, SSH channel dispatch, WebTransport
streams, and wasm — all via from_stream, all through the same
HandlerRegistry, zero handler code changes.

See docs/research/transport-generalization/findings.md for the full
trace and the yield-once contract.
This commit is contained in:
glm-5.2 committed 2026-07-09 07:06:12 +00:00
1 parent acd049e8b2
commit 865fef6210
6 files changed
+127 -192

No files matched your search

+7 -25
View File
@@ -576,34 +576,16 @@ mod tests {
};
use crate::registry::spec::{AccessControl, OperationSpec, OperationType, Visibility};
use alknet_core::auth::Identity;
use alknet_core::types::{Capabilities, MockConnection};
use alknet_core::types::Capabilities;
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::sync::Mutex as StdMutex;
struct StubConnection {
alpn: &'static [u8],
addr: Option<SocketAddr>,
closed: StdMutex<Option<(u32, String)>>,
}
impl MockConnection for StubConnection {
fn remote_alpn(&self) -> &[u8] {
self.alpn
}
fn remote_addr(&self) -> Option<SocketAddr> {
self.addr
}
fn close(&self, code: u32, reason: &str) {
*self.closed.lock().unwrap() = Some((code, reason.to_string()));
}
}
fn stub_connection() -> Connection {
Connection::from_mock(Arc::new(StubConnection {
alpn: b"alknet/call",
addr: Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)),
closed: StdMutex::new(None),
}))
Connection::from_stream(
tokio::io::sink(),
tokio::io::empty(),
b"alknet/call".to_vec(),
Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)),
)
}
fn external_spec(name: &str) -> OperationSpec {
+7 -24
View File
@@ -434,35 +434,18 @@ mod tests {
use crate::registry::registration::{make_handler, make_streaming_handler};
use crate::registry::spec::OperationType;
use alknet_core::auth::Identity;
use alknet_core::types::{Capabilities, MockConnection};
use alknet_core::types::Capabilities;
use std::collections::HashMap;
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::sync::Mutex as StdMutex;
struct StubConnection {
alpn: &'static [u8],
addr: Option<SocketAddr>,
closed: StdMutex<Option<(u32, String)>>,
}
impl MockConnection for StubConnection {
fn remote_alpn(&self) -> &[u8] {
self.alpn
}
fn remote_addr(&self) -> Option<SocketAddr> {
self.addr
}
fn close(&self, code: u32, reason: &str) {
*self.closed.lock().unwrap() = Some((code, reason.to_string()));
}
}
fn stub_connection() -> alknet_core::types::Connection {
alknet_core::types::Connection::from_mock(Arc::new(StubConnection {
alpn: b"alknet/call",
addr: Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)),
closed: StdMutex::new(None),
}))
alknet_core::types::Connection::from_stream(
tokio::io::sink(),
tokio::io::empty(),
b"alknet/call".to_vec(),
Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)),
)
}
fn sample_schema_json(name: &str, op_type: &str) -> Value {
+10 -27
View File
@@ -290,30 +290,13 @@ mod tests {
})
}
struct StubConnection {
alpn: &'static [u8],
addr: Option<SocketAddr>,
closed: StdMutex<Option<(u32, String)>>,
}
impl alknet_core::types::MockConnection for StubConnection {
fn remote_alpn(&self) -> &[u8] {
self.alpn
}
fn remote_addr(&self) -> Option<SocketAddr> {
self.addr
}
fn close(&self, code: u32, reason: &str) {
*self.closed.lock().unwrap() = Some((code, reason.to_string()));
}
}
fn stub_connection() -> Connection {
Connection::from_mock(Arc::new(StubConnection {
alpn: b"alknet/call",
addr: Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)),
closed: StdMutex::new(None),
}))
Connection::from_stream(
tokio::io::sink(),
tokio::io::empty(),
b"alknet/call".to_vec(),
Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)),
)
}
#[test]
@@ -1210,8 +1193,8 @@ mod tests {
let frame = encode_frame(&EventEnvelope::aborted("parent-1"));
let recv = tokio::io::BufReader::new(std::io::Cursor::new(frame));
let (send, _recv_sink) = tokio::io::duplex(64);
let send = alknet_core::types::SendStream::from_mock(send);
let recv = alknet_core::types::RecvStream::from_mock(recv);
let send = alknet_core::types::SendStream::from_stream(send);
let recv = alknet_core::types::RecvStream::from_stream(recv);
adapter.handle_stream(conn.clone(), send, recv).await;
@@ -1250,8 +1233,8 @@ mod tests {
let frame = encode_frame(&EventEnvelope::aborted("does-not-exist"));
let recv = tokio::io::BufReader::new(std::io::Cursor::new(frame));
let (send, _recv_sink) = tokio::io::duplex(64);
let send = alknet_core::types::SendStream::from_mock(send);
let recv = alknet_core::types::RecvStream::from_mock(recv);
let send = alknet_core::types::SendStream::from_stream(send);
let recv = alknet_core::types::RecvStream::from_stream(recv);
adapter.handle_stream(conn.clone(), send, recv).await;
+7 -25
View File
@@ -456,36 +456,18 @@ mod tests {
use crate::registry::context::CompositionAuthority;
use crate::registry::registration::{make_handler, Handler, HandlerKind, OperationProvenance};
use crate::registry::spec::{AccessControl, OperationSpec, OperationType, Visibility};
use alknet_core::types::{Capabilities, MockConnection};
use alknet_core::types::Capabilities;
use std::collections::HashMap;
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::sync::Mutex as StdMutex;
use std::time::{Duration, Instant};
struct StubConnection {
alpn: &'static [u8],
addr: Option<SocketAddr>,
closed: StdMutex<Option<(u32, String)>>,
}
impl MockConnection for StubConnection {
fn remote_alpn(&self) -> &[u8] {
self.alpn
}
fn remote_addr(&self) -> Option<SocketAddr> {
self.addr
}
fn close(&self, code: u32, reason: &str) {
*self.closed.lock().unwrap() = Some((code, reason.to_string()));
}
}
fn stub_connection() -> Connection {
Connection::from_mock(Arc::new(StubConnection {
alpn: b"alknet/call",
addr: Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)),
closed: StdMutex::new(None),
}))
Connection::from_stream(
tokio::io::sink(),
tokio::io::empty(),
b"alknet/call".to_vec(),
Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)),
)
}
fn external_spec(name: &str) -> OperationSpec {
+17 -34
View File
@@ -456,35 +456,18 @@ mod tests {
};
use crate::registry::spec::{AccessControl, OperationSpec, OperationType, Visibility};
use alknet_core::auth::{AuthToken, Identity, IdentityProvider};
use alknet_core::types::{Capabilities, MockConnection};
use alknet_core::types::Capabilities;
use std::collections::HashMap;
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::sync::Mutex as StdMutex;
struct StubConnection {
alpn: &'static [u8],
addr: Option<SocketAddr>,
closed: StdMutex<Option<(u32, String)>>,
}
impl MockConnection for StubConnection {
fn remote_alpn(&self) -> &[u8] {
self.alpn
}
fn remote_addr(&self) -> Option<SocketAddr> {
self.addr
}
fn close(&self, code: u32, reason: &str) {
*self.closed.lock().unwrap() = Some((code, reason.to_string()));
}
}
fn stub_connection() -> alknet_core::types::Connection {
alknet_core::types::Connection::from_mock(Arc::new(StubConnection {
alpn: b"alknet/call",
addr: Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)),
closed: StdMutex::new(None),
}))
alknet_core::types::Connection::from_stream(
tokio::io::sink(),
tokio::io::empty(),
b"alknet/call".to_vec(),
Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)),
)
}
struct StaticIdentityProvider {
@@ -1197,8 +1180,8 @@ mod tests {
);
let recv = tokio::io::BufReader::new(std::io::Cursor::new(encode_frame(&request)));
let (send, mut sink) = tokio::io::duplex(8 * 1024);
let send = alknet_core::types::SendStream::from_mock(send);
let recv = alknet_core::types::RecvStream::from_mock(recv);
let send = alknet_core::types::SendStream::from_stream(send);
let recv = alknet_core::types::RecvStream::from_stream(recv);
dp.handle_stream(conn, send, recv).await;
@@ -1236,8 +1219,8 @@ mod tests {
);
let recv = tokio::io::BufReader::new(std::io::Cursor::new(encode_frame(&request)));
let (send, mut sink) = tokio::io::duplex(8 * 1024);
let send = alknet_core::types::SendStream::from_mock(send);
let recv = alknet_core::types::RecvStream::from_mock(recv);
let send = alknet_core::types::SendStream::from_stream(send);
let recv = alknet_core::types::RecvStream::from_stream(recv);
dp.handle_stream(conn, send, recv).await;
@@ -1272,8 +1255,8 @@ mod tests {
);
let recv = tokio::io::BufReader::new(std::io::Cursor::new(encode_frame(&request)));
let (send, mut sink) = tokio::io::duplex(8 * 1024);
let send = alknet_core::types::SendStream::from_mock(send);
let recv = alknet_core::types::RecvStream::from_mock(recv);
let send = alknet_core::types::SendStream::from_stream(send);
let recv = alknet_core::types::RecvStream::from_stream(recv);
dp.handle_stream(conn, send, recv).await;
@@ -1303,8 +1286,8 @@ mod tests {
);
let recv = tokio::io::BufReader::new(std::io::Cursor::new(encode_frame(&request)));
let (send, mut sink) = tokio::io::duplex(8 * 1024);
let send = alknet_core::types::SendStream::from_mock(send);
let recv = alknet_core::types::RecvStream::from_mock(recv);
let send = alknet_core::types::SendStream::from_stream(send);
let recv = alknet_core::types::RecvStream::from_stream(recv);
dp.handle_stream(conn, send, recv).await;
@@ -1363,8 +1346,8 @@ mod tests {
);
let recv = tokio::io::BufReader::new(std::io::Cursor::new(encode_frame(&request)));
let (send, _sink) = tokio::io::duplex(8 * 1024);
let send = alknet_core::types::SendStream::from_mock(send);
let recv = alknet_core::types::RecvStream::from_mock(recv);
let send = alknet_core::types::SendStream::from_stream(send);
let recv = alknet_core::types::RecvStream::from_stream(recv);
let conn_clone = Arc::clone(&conn);
let dp_clone = dp.clone();
+79 -57
View File
@@ -6,7 +6,7 @@
use std::collections::HashMap;
use std::io;
use std::net::SocketAddr;
use std::sync::{Arc, OnceLock};
use std::sync::{Mutex, OnceLock};
use async_trait::async_trait;
use tokio::io::{AsyncRead, AsyncWrite};
@@ -230,7 +230,7 @@ enum SendStreamKind {
Quinn(quinn::SendStream),
#[cfg(feature = "iroh")]
Iroh(iroh::endpoint::SendStream),
Mock(Box<dyn AsyncWrite + Send + Unpin>),
Stream(Box<dyn AsyncWrite + Send + Unpin>),
}
enum RecvStreamKind {
@@ -238,7 +238,7 @@ enum RecvStreamKind {
Quinn(quinn::RecvStream),
#[cfg(feature = "iroh")]
Iroh(iroh::endpoint::RecvStream),
Mock(Box<dyn AsyncRead + Send + Unpin>),
Stream(Box<dyn AsyncRead + Send + Unpin>),
}
pub struct SendStream {
@@ -264,10 +264,9 @@ impl SendStream {
}
}
#[allow(dead_code)]
pub fn from_mock(stream: impl AsyncWrite + Send + Unpin + 'static) -> Self {
pub fn from_stream(stream: impl AsyncWrite + Send + Unpin + 'static) -> Self {
Self {
kind: SendStreamKind::Mock(Box::new(stream)),
kind: SendStreamKind::Stream(Box::new(stream)),
}
}
}
@@ -287,10 +286,9 @@ impl RecvStream {
}
}
#[allow(dead_code)]
pub fn from_mock(stream: impl AsyncRead + Send + Unpin + 'static) -> Self {
pub fn from_stream(stream: impl AsyncRead + Send + Unpin + 'static) -> Self {
Self {
kind: RecvStreamKind::Mock(Box::new(stream)),
kind: RecvStreamKind::Stream(Box::new(stream)),
}
}
}
@@ -306,7 +304,7 @@ impl AsyncWrite for SendStream {
SendStreamKind::Quinn(s) => AsyncWrite::poll_write(std::pin::Pin::new(s), cx, buf),
#[cfg(feature = "iroh")]
SendStreamKind::Iroh(s) => AsyncWrite::poll_write(std::pin::Pin::new(s), cx, buf),
SendStreamKind::Mock(s) => {
SendStreamKind::Stream(s) => {
AsyncWrite::poll_write(std::pin::Pin::new(s.as_mut()), cx, buf)
}
}
@@ -321,7 +319,7 @@ impl AsyncWrite for SendStream {
SendStreamKind::Quinn(s) => AsyncWrite::poll_flush(std::pin::Pin::new(s), cx),
#[cfg(feature = "iroh")]
SendStreamKind::Iroh(s) => AsyncWrite::poll_flush(std::pin::Pin::new(s), cx),
SendStreamKind::Mock(s) => AsyncWrite::poll_flush(std::pin::Pin::new(s.as_mut()), cx),
SendStreamKind::Stream(s) => AsyncWrite::poll_flush(std::pin::Pin::new(s.as_mut()), cx),
}
}
@@ -334,8 +332,8 @@ impl AsyncWrite for SendStream {
SendStreamKind::Quinn(s) => AsyncWrite::poll_shutdown(std::pin::Pin::new(s), cx),
#[cfg(feature = "iroh")]
SendStreamKind::Iroh(s) => AsyncWrite::poll_shutdown(std::pin::Pin::new(s), cx),
SendStreamKind::Mock(s) => {
AsyncWrite::poll_shutdown(std::pin::Pin::new(s.as_mut()), cx)
SendStreamKind::Stream(s) => {
AsyncWrite::poll_shutdown(std::pin::Pin::new(s), cx)
}
}
}
@@ -352,7 +350,7 @@ impl AsyncRead for RecvStream {
RecvStreamKind::Quinn(s) => AsyncRead::poll_read(std::pin::Pin::new(s), cx, buf),
#[cfg(feature = "iroh")]
RecvStreamKind::Iroh(s) => AsyncRead::poll_read(std::pin::Pin::new(s), cx, buf),
RecvStreamKind::Mock(s) => {
RecvStreamKind::Stream(s) => {
AsyncRead::poll_read(std::pin::Pin::new(s.as_mut()), cx, buf)
}
}
@@ -364,14 +362,12 @@ enum ConnectionKind {
Quinn(quinn::Connection),
#[cfg(feature = "iroh")]
Iroh(iroh::endpoint::Connection),
Mock(Arc<dyn MockConnection + Send + Sync>),
Stream(StreamConn),
}
#[allow(dead_code)]
pub trait MockConnection: Send + Sync {
fn remote_alpn(&self) -> &[u8];
fn remote_addr(&self) -> Option<SocketAddr>;
fn close(&self, code: u32, reason: &str);
struct StreamConn {
stream: Mutex<Option<(SendStream, RecvStream)>>,
remote_addr: Option<SocketAddr>,
}
pub struct Connection {
@@ -405,16 +401,52 @@ impl Connection {
}
}
#[allow(dead_code)]
pub fn from_mock(mock: Arc<dyn MockConnection + Send + Sync>) -> Self {
let alpn = mock.remote_alpn().to_vec();
/// Construct a `Connection` from a pre-split read/write pair.
/// `accept_bi()` yields this pair once, then returns `ConnectionClosed`.
/// `open_bi()` returns `StreamClosed` (a single stream can't open new streams).
pub fn from_stream(
send: impl AsyncWrite + Send + Unpin + 'static,
recv: impl AsyncRead + Send + Unpin + 'static,
alpn: Vec<u8>,
remote_addr: Option<SocketAddr>,
) -> Self {
Self {
kind: ConnectionKind::Mock(mock),
kind: ConnectionKind::Stream(StreamConn {
stream: Mutex::new(Some((
SendStream::from_stream(send),
RecvStream::from_stream(recv),
))),
remote_addr,
}),
alpn,
identity: OnceLock::new(),
}
}
/// Convenience for a single bidirectional stream (e.g. `TlsStream<TcpStream>`).
/// Splits internally via `tokio::io::split`.
pub fn from_bidi(
stream: impl AsyncRead + AsyncWrite + Send + Unpin + 'static,
alpn: Vec<u8>,
remote_addr: Option<SocketAddr>,
) -> Self {
let (recv, send) = tokio::io::split(stream);
Self::from_stream(send, recv, alpn, remote_addr)
}
/// Yield the next bidirectional stream this connection provides.
///
/// # Transport semantics
///
/// - **QUIC (quinn/iroh)**: returns a new bidi stream on each call.
/// `ConnectionClosed` when the underlying connection closes.
/// - **TCP+TLS / single-stream**: yields the underlying stream on the
/// first call, then `ConnectionClosed` on all subsequent calls.
/// A single transport stream cannot open new application streams.
///
/// Handlers that loop `accept_bi` (e.g. `TtyAdapter`) get one session
/// per single-stream connection; handlers that call once (e.g.
/// `HttpAdapter`) get the stream directly. Both are correct.
pub async fn accept_bi(&self) -> Result<(SendStream, RecvStream), StreamError> {
match &self.kind {
#[cfg(feature = "quinn")]
@@ -427,7 +459,13 @@ impl Connection {
let (send, recv) = c.accept_bi().await.map_err(map_iroh_connection_error)?;
Ok((SendStream::from_iroh(send), RecvStream::from_iroh(recv)))
}
ConnectionKind::Mock(_) => Err(StreamError::StreamClosed),
ConnectionKind::Stream(sc) => {
let mut guard = sc.stream.lock().expect("stream mutex poisoned");
match guard.take() {
Some(pair) => Ok(pair),
None => Err(StreamError::ConnectionClosed),
}
}
}
}
@@ -443,7 +481,7 @@ impl Connection {
let (send, recv) = c.open_bi().await.map_err(map_iroh_connection_error)?;
Ok((SendStream::from_iroh(send), RecvStream::from_iroh(recv)))
}
ConnectionKind::Mock(_) => Err(StreamError::StreamClosed),
ConnectionKind::Stream(_) => Err(StreamError::StreamClosed),
}
}
@@ -457,7 +495,7 @@ impl Connection {
ConnectionKind::Quinn(c) => Some(c.remote_address()),
#[cfg(feature = "iroh")]
ConnectionKind::Iroh(_) => None,
ConnectionKind::Mock(m) => m.remote_addr(),
ConnectionKind::Stream(sc) => sc.remote_addr,
}
}
@@ -473,7 +511,9 @@ impl Connection {
let code = iroh::endpoint::VarInt::from(code);
c.close(code, reason.as_bytes());
}
ConnectionKind::Mock(m) => m.close(code, reason),
ConnectionKind::Stream(sc) => {
let _ = sc.stream.lock().expect("stream mutex poisoned").take();
}
}
}
@@ -517,31 +557,13 @@ mod tests {
use super::*;
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
struct MockConn {
alpn: &'static [u8],
addr: Option<SocketAddr>,
closed: std::sync::Mutex<Option<(u32, String)>>,
}
#[allow(dead_code)]
impl MockConnection for MockConn {
fn remote_alpn(&self) -> &[u8] {
self.alpn
}
fn remote_addr(&self) -> Option<SocketAddr> {
self.addr
}
fn close(&self, code: u32, reason: &str) {
*self.closed.lock().unwrap() = Some((code, reason.to_string()));
}
}
fn mock_connection() -> Connection {
Connection::from_mock(Arc::new(MockConn {
alpn: b"alknet/test",
addr: Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 1234)),
closed: std::sync::Mutex::new(None),
}))
fn test_connection() -> Connection {
Connection::from_stream(
tokio::io::sink(),
tokio::io::empty(),
b"alknet/test".to_vec(),
Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 1234)),
)
}
#[test]
@@ -607,7 +629,7 @@ mod tests {
#[test]
fn set_identity_once_succeeds_twice_errors() {
let conn = mock_connection();
let conn = test_connection();
let id = Identity {
id: "alk_test".to_string(),
scopes: vec!["relay:connect".to_string()],
@@ -622,7 +644,7 @@ mod tests {
#[test]
fn identity_get_returns_set_value() {
let conn = mock_connection();
let conn = test_connection();
assert!(conn.identity().is_none());
let id = Identity {
id: "alk_test".to_string(),
@@ -634,8 +656,8 @@ mod tests {
}
#[test]
fn connection_remote_alpn_and_addr_from_mock() {
let conn = mock_connection();
fn connection_remote_alpn_and_addr_from_stream() {
let conn = test_connection();
assert_eq!(conn.remote_alpn(), b"alknet/test");
assert_eq!(
conn.remote_addr(),