feat(wasi): run EasyTier core on Cloudflare Workers and browsers (#2548)

* fix(core): normalize secure keys for TOML instances

* feat(wasi): run core behind Cloudflare WebSockets

Introduce the Cloudflare Worker WASI host that runs the EasyTier core
behind host-upgraded WebSockets.

- Worker package scaffold (wrangler Durable Object, build-wasm script,
  vitest config) and core-runtime/websocket-host/data-plane runtime.
- WASI host WebSocket tunnel ABI (imports, adapter, runtime exports)
  with bounded receive memory and bounded admission queue.
- Route host sockets through the portable listener plan
  (HostListenerRegistration, listener queue, admission handler split).
- Build the WASM guest with the aes-gcm feature so secure peer
  sessions have their cipher available.

* feat(wasi): add outbound browser client runtime

Add the outbound-only WASI runtime and browser connector host so
browser pages can dial EasyTier peers through WebSocket relays.

- CoreConnectivityMode::{OutboundOnly, InboundOnly} gating for
  listeners, discovery, and direct connectivity modules.
- ExternalTunnelConnector plumbing through composite/connector_host/
  manual for browser WebSocket dials.
- Browser/Node smoke entries with shared helpers
  (smoke-shared.ts).

* feat(wasi): extend browser data plane with TCP half-close

Add the data-plane pieces the browser runtime needs for full-duplex
TCP streams behind host WebSockets:

- Guest TCP shutdown_write operation with submit/take ABI pair
  (DATA_PLANE_ABI_VERSION 3 -> 4) and smoltcp half-close support.
- Worker data-plane TCP listener/stream plumbing and core-runtime
  listener registration.
- Unit coverage for the new session ops and listener wiring.

* refactor(wasi): make host tunnel ABI transport-neutral

Replace WebSocket-specific core and WASI boundaries with a
message-oriented Host Tunnel interface. Keep WebSocket framing and text
rejection in the Cloudflare host while preserving payload boundaries,
ownership, cancellation, backpressure, and EOF behavior.

Rename feature flags and guest imports and exports to the Host Tunnel
ABI. Update both Worker profiles, tests, and architecture documentation.

* feat(web): split WASI hosts into publishable npm packages

Extract the shared JSPI, WASI, Host Tunnel, and data-plane runtime
into @easytier/runtime. Keep ABI handles, guest memory, TOML, and
operation broker details behind its adapter entry point.

Add typed, auto-starting @easytier/browser and factory-based
@easytier/cloudflare packages. Ship a matching Wasm profile with
each platform package and validate its capabilities before packing.

Persist Cloudflare instance identity in Durable Object storage,
centralize WebSocket admission ownership, and add package-level
coverage for the public interfaces.

* fix(web): make public packages portable

Embed the browser Wasm artifact in the published JavaScript entry
point. This lets esbuild consumers bundle the package without an asset
loader or a copied file.

Return Cloudflare's nominal Durable Object base type and document the
named subclass export required by generated Wrangler bindings.

* docs(web): add public package walkthrough

Expand both package READMEs with installation, configuration, local
validation, health checks, and deployment instructions.

Add a standalone Vite and Wrangler example that imports only the
public Browser and Cloudflare entries. Generate Worker bindings from
configuration and keep local secrets outside version control.

* chore(go): import EasyTier Go host

Add the standalone Go host runtime as a monorepo subtree without
carrying its development branch ancestry.

Preserve its API, tests, examples, generated protobuf bindings, and
embedded WASI artifacts.

* refactor(hosts): colocate Go and JavaScript runtimes

Move the browser, Cloudflare, shared runtime, and web example into
the easytier-js subtree. Update workspace metadata, build paths, and
documentation for the new layout.

Adopt github.com/EasyTier/EasyTier/easytier-go as the Go module path.
Resolve artifact and protobuf generation from the enclosing monorepo.

* build(web): isolate JavaScript host workspace

Keep public browser and Cloudflare packages outside the legacy frontend
workspace so root installs and cross-platform builds do not pull workerd.

Make each package build generate its required WASI artifact from a clean
checkout. Add a dedicated workflow that runs the same install and check
commands documented for contributors.

Move JavaScript dependencies into a scoped lockfile and restore the root
workspace lockfile to its pre-host state.
This commit is contained in:
KKRainbow authored and GitHub committed 2026-09-06 13:35:02 +08:00
1 parent 56be71c7f9
commit f26c2aa147
215 files changed
+45277 -144

No files matched your search

+2
View File
@@ -146,6 +146,8 @@ proxy-smoltcp-stack = [
"smoltcp/proto-ipv6",
"smoltcp/async",
]
wasm-host-tunnel = []
wasm-host-tunnel-outbound = ["wasm-host-tunnel"]
test-utils = []
tracing-log = ["tracing/log"]
zstd = ["dep:zstd"]
@@ -103,6 +103,17 @@ impl InterfaceAddrCache {
/// Mechanical connector operations supplied by one process-wide runtime.
#[async_trait]
pub trait ConnectorRuntime: VirtualTcpSocketFactory + Send + Sync + 'static {
fn supports_external_tunnel(&self, _scheme: &str) -> bool {
false
}
async fn connect_external_tunnel(
&self,
_url: &Url,
) -> anyhow::Result<Option<Box<dyn crate::tunnel::Tunnel>>> {
Ok(None)
}
async fn connect_byte_stream(
&self,
url: &Url,
@@ -207,6 +218,17 @@ where
S: ConnectorRuntime + VirtualUdpSocketFactory,
E: ConnectorEnvironment,
{
fn supports_external_tunnel(&self, scheme: &str) -> bool {
self.sockets.supports_external_tunnel(scheme)
}
async fn connect_external_tunnel(
&self,
url: &Url,
) -> anyhow::Result<Option<Box<dyn crate::tunnel::Tunnel>>> {
self.sockets.connect_external_tunnel(url).await
}
async fn local_addr_for_remote(
&self,
remote_addr: SocketAddr,
@@ -24,6 +24,7 @@ use url::Url;
use crate::{
connectivity::{
composite::{ConnectorEnvironment, ConnectorHostAdapter, ConnectorRuntime},
manual::ExternalTunnelConnector,
transport::ConnectedByteStream,
},
host::environment::{HostConnectorEnvironmentIo, local_addr_for_remote},
@@ -113,6 +114,7 @@ where
listeners: HostTcpListenerFactory<B>,
environment: Arc<HostConnectorEnvironmentSnapshot>,
environment_io: Arc<E>,
external_tunnel_connector: Option<Arc<dyn ExternalTunnelConnector>>,
}
impl<B, E> HostConnectorRuntime<B, E>
@@ -132,8 +134,17 @@ where
listeners: HostTcpListenerFactory::new(runtime, backend),
environment: Arc::new(environment),
environment_io,
external_tunnel_connector: None,
}
}
pub fn with_external_tunnel_connector(
mut self,
connector: Arc<dyn ExternalTunnelConnector>,
) -> Self {
self.external_tunnel_connector = Some(connector);
self
}
}
#[async_trait]
@@ -181,6 +192,25 @@ where
B: ConnectorHostSocketBackend,
E: HostConnectorEnvironmentIo,
{
fn supports_external_tunnel(&self, scheme: &str) -> bool {
self.external_tunnel_connector
.as_ref()
.is_some_and(|connector| connector.supports_scheme(scheme))
}
async fn connect_external_tunnel(
&self,
url: &Url,
) -> anyhow::Result<Option<Box<dyn crate::tunnel::Tunnel>>> {
let Some(connector) = &self.external_tunnel_connector else {
return Ok(None);
};
if !connector.supports_scheme(url.scheme()) {
return Ok(None);
}
Ok(Some(connector.connect(url).await?))
}
async fn connect_byte_stream(
&self,
url: &Url,
@@ -257,6 +287,24 @@ where
ConnectorHostAdapter::new(runtime.clone(), runtime)
}
pub fn new_connector_host_with_external_tunnel<B, E>(
socket_runtime: HostSocketRuntime,
backend: Arc<B>,
environment: HostConnectorEnvironmentSnapshot,
environment_io: Arc<E>,
connector: Arc<dyn ExternalTunnelConnector>,
) -> ConnectorHost<B, E>
where
B: ConnectorHostSocketBackend,
E: HostConnectorEnvironmentIo,
{
let runtime = Arc::new(
HostConnectorRuntime::new(socket_runtime, backend, environment, environment_io)
.with_external_tunnel_connector(connector),
);
ConnectorHostAdapter::new(runtime.clone(), runtime)
}
#[cfg(test)]
mod tests {
use std::{
@@ -299,6 +347,19 @@ mod tests {
struct FixedStunProvider;
struct ExternalFailingConnector;
#[async_trait]
impl ExternalTunnelConnector for ExternalFailingConnector {
fn supports_scheme(&self, scheme: &str) -> bool {
scheme == "ws"
}
async fn connect(&self, _url: &Url) -> anyhow::Result<Box<dyn crate::tunnel::Tunnel>> {
anyhow::bail!("external connector called")
}
}
#[async_trait]
impl StunInfoProvider for FixedStunProvider {
fn get_stun_info(&self) -> StunInfo {
@@ -579,6 +640,34 @@ mod tests {
);
}
#[tokio::test]
async fn delegates_external_tunnel_connections() {
let host = new_connector_host_with_external_tunnel(
HostSocketRuntime::new(),
Arc::new(UnsupportedBackend::default()),
test_environment_snapshot(),
Arc::new(TestEnvironmentIo::default()),
Arc::new(ExternalFailingConnector),
);
let error = ManualConnectorHost::connect_external_tunnel(
&host,
&"ws://relay.example/".parse().unwrap(),
)
.await
.unwrap_err();
assert_eq!(error.to_string(), "external connector called");
assert!(
ManualConnectorHost::connect_external_tunnel(
&host,
&"tcp://relay.example:11010".parse().unwrap(),
)
.await
.unwrap()
.is_none()
);
}
#[tokio::test]
async fn direct_rpc_projects_host_observations_without_instance_policy() {
let host = Arc::new(new_connector_host(
+52 -13
View File
@@ -94,6 +94,14 @@ pub struct ManualInterfaceAddrs {
#[async_trait]
pub trait ManualConnectorHost: VirtualTcpSocketFactory + VirtualUdpSocketFactory {
fn supports_external_tunnel(&self, _scheme: &str) -> bool {
false
}
async fn connect_external_tunnel(&self, _url: &Url) -> anyhow::Result<Option<Box<dyn Tunnel>>> {
Ok(None)
}
async fn local_addr_for_remote(
&self,
remote_addr: SocketAddr,
@@ -110,6 +118,13 @@ pub trait ManualConnectorHost: VirtualTcpSocketFactory + VirtualUdpSocketFactory
}
}
#[async_trait]
pub trait ExternalTunnelConnector: Send + Sync + 'static {
fn supports_scheme(&self, scheme: &str) -> bool;
async fn connect(&self, url: &Url) -> anyhow::Result<Box<dyn Tunnel>>;
}
#[async_trait]
pub(crate) trait ManualEndpointResolver: Send + Sync + 'static {
async fn resolve_endpoint(&self, url: &Url) -> anyhow::Result<Url>;
@@ -182,6 +197,14 @@ where
));
}
if let Some(tunnel) = self.host.connect_external_tunnel(&endpoint.url).await? {
return Ok(apply_resolved_endpoint_info(
tunnel,
requested_url,
endpoint.tunnel_prefixes,
));
}
if !self.protocol.supports_scheme(endpoint.url.scheme()) {
anyhow::bail!(
"unsupported client protocol upgrader: {}",
@@ -710,17 +733,20 @@ where
return Err(error);
}
};
let ip_versions = match resolve_reconnect_ip_versions(
&normalized_url,
connect_timeout,
ManualTransport::from_url(&normalized_url)
.ok()
.map(|transport| data.options.socket_context(transport, IpVersion::Both))
.unwrap_or_default(),
data.dns.as_ref(),
)
.await
{
let ip_versions = match if data.host.supports_external_tunnel(normalized_url.scheme()) {
Ok(vec![IpVersion::Both])
} else {
resolve_reconnect_ip_versions(
&normalized_url,
connect_timeout,
ManualTransport::from_url(&normalized_url)
.ok()
.map(|transport| data.options.socket_context(transport, IpVersion::Both))
.unwrap_or_default(),
data.dns.as_ref(),
)
.await
} {
Ok(ip_versions) => ip_versions,
Err(error) => {
emit_connect_error(&data, &url, IpVersion::Both, &error);
@@ -770,13 +796,17 @@ where
),
)
.await?;
if endpoint.url.scheme() != "ring" && !data.protocol.supports_scheme(endpoint.url.scheme()) {
let uses_external_tunnel = data.host.supports_external_tunnel(endpoint.url.scheme());
if endpoint.url.scheme() != "ring"
&& !uses_external_tunnel
&& !data.protocol.supports_scheme(endpoint.url.scheme())
{
anyhow::bail!(
"unsupported client protocol upgrader: {}",
endpoint.url.scheme()
);
}
let transport = (endpoint.url.scheme() != "ring")
let transport = (endpoint.url.scheme() != "ring" && !uses_external_tunnel)
.then(|| ManualTransport::from_url(&endpoint.url))
.transpose()?;
let resolved = match transport {
@@ -822,6 +852,15 @@ where
if endpoint.url.scheme() == "ring" {
return connect_ring_tunnel(&data.ring_registry, &endpoint.url);
}
if uses_external_tunnel {
return data
.host
.connect_external_tunnel(&endpoint.url)
.await?
.ok_or_else(|| {
anyhow::anyhow!("host did not provide external tunnel for {}", endpoint.url)
});
}
let transport = transport.expect("non-Ring endpoint should have a transport");
let connected = match resolved {
Some((remote_addr, bind_addrs)) => {
@@ -53,6 +53,7 @@ pub enum DataPlaneOperationKind {
UdpBind = 6,
UdpReceive = 7,
UdpSend = 8,
TcpShutdownWrite = 9,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -100,6 +101,7 @@ pub enum DataPlaneOperationResult {
TcpWritten {
len: usize,
},
TcpWriteShutdown,
UdpBound {
socket: DataPlaneResourceId,
local_addr: SocketAddr,
@@ -136,6 +138,7 @@ impl DataPlaneOperationResult {
Self::UdpBound { socket, .. } => Some(*socket),
Self::TcpRead { .. }
| Self::TcpWritten { .. }
| Self::TcpWriteShutdown
| Self::UdpReceived { .. }
| Self::UdpSent { .. } => None,
}
@@ -168,6 +168,7 @@ enum PendingOperationResult {
eof: bool,
},
TcpWritten(usize),
TcpWriteShutdown,
UdpBound(DataPlaneUdpSocket),
UdpReceived {
data: Vec<u8>,
@@ -624,6 +625,36 @@ where
Ok(operation_id)
}
pub fn submit_tcp_shutdown_write(
self: &Arc<Self>,
stream_id: DataPlaneResourceId,
) -> DataPlaneResult<DataPlaneOperationId> {
Self::ensure_executor()?;
let (stream, operation_id, cancel) = {
let mut state = self.lock_state();
let stream = Self::require_tcp(&state, stream_id)?;
let (operation_id, cancel) = self.admit_locked(
&mut state,
DataPlaneOperationKind::TcpShutdownWrite,
Some(stream_id),
0,
false,
)?;
(stream, operation_id, cancel)
};
self.spawn_operation(operation_id, async move {
stream
.write_deadline
.run(cancel, async {
stream.write.lock().await.shutdown().await?;
Ok::<_, std::io::Error>(())
})
.await?;
Ok(PendingOperationResult::TcpWriteShutdown)
});
Ok(operation_id)
}
pub fn submit_udp_bind(
self: &Arc<Self>,
local_port: u16,
@@ -880,6 +911,7 @@ where
DataPlaneOperationResult::TcpRead { data, eof }
}
PendingOperationResult::TcpWritten(len) => DataPlaneOperationResult::TcpWritten { len },
PendingOperationResult::TcpWriteShutdown => DataPlaneOperationResult::TcpWriteShutdown,
PendingOperationResult::UdpBound(socket) => {
let local_addr = socket.local_addr();
let socket = Self::insert_udp_resource_locked(resources, socket)?;
@@ -327,6 +327,55 @@ async fn data_plane_sessions_complete_tcp_operations_end_to_end() {
assert_eq!(written, 4);
assert_eq!(received, b"ping");
let eof_read = session_b.submit_tcp_read(server, 16).unwrap();
let shutdown = session_a.submit_tcp_shutdown_write(client).unwrap();
let shutdown_completion = wait_for_session_completion(&session_a).await;
let eof_completion = wait_for_session_completion(&session_b).await;
assert_eq!(shutdown_completion.operation_id, shutdown);
assert_eq!(eof_completion.operation_id, eof_read);
session_a
.take_result_with(shutdown, |outcome| match outcome {
Ok(DataPlaneOperationResult::TcpWriteShutdown) => Some(()),
_ => None,
})
.unwrap()
.unwrap();
let eof = session_b
.take_result_with(eof_read, |outcome| match outcome {
Ok(DataPlaneOperationResult::TcpRead { data, eof }) => Some((data.clone(), *eof)),
_ => None,
})
.unwrap()
.unwrap();
assert_eq!(eof, (Vec::new(), true));
let response_read = session_a.submit_tcp_read(client, 16).unwrap();
let response_write = session_b
.submit_tcp_write(server, b"pong".to_vec())
.unwrap();
let (response_read_completion, response_write_completion) = tokio::join!(
wait_for_session_completion(&session_a),
wait_for_session_completion(&session_b),
);
assert_eq!(response_read_completion.operation_id, response_read);
assert_eq!(response_write_completion.operation_id, response_write);
let response = session_a
.take_result_with(response_read, |outcome| match outcome {
Ok(DataPlaneOperationResult::TcpRead { data, eof }) if !eof => Some(data.clone()),
_ => None,
})
.unwrap()
.unwrap();
let response_len = session_b
.take_result_with(response_write, |outcome| match outcome {
Ok(DataPlaneOperationResult::TcpWritten { len }) => Some(*len),
_ => None,
})
.unwrap()
.unwrap();
assert_eq!(response, b"pong");
assert_eq!(response_len, 4);
let blocked_read = session_b.submit_tcp_read(server, 16).unwrap();
session_b.close_resource(server);
let close_completion = wait_for_session_completion(&session_b).await;
@@ -199,11 +199,14 @@ impl AsyncWrite for TcpStream {
fn poll_shutdown(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Result<(), io::Error>> {
let mut socket = self.reactor.get_socket::<tcp::Socket>(*self.handle);
if socket.is_open() {
if socket.may_send() {
socket.close();
self.reactor.notify();
}
if socket.state() == tcp::State::Closed {
if matches!(
socket.state(),
tcp::State::FinWait2 | tcp::State::TimeWait | tcp::State::Closed
) {
return Poll::Ready(Ok(()));
}
+1
View File
@@ -15,3 +15,4 @@ pub mod packet;
pub mod socket;
#[cfg(test)]
pub(crate) mod testkit;
pub mod tunnel;
+1 -1
View File
@@ -184,7 +184,7 @@ impl HostSocketRuntime {
}
}
pub(in crate::host) async fn run_operation<I, T>(
pub(crate) async fn run_operation<I, T>(
&self,
io: Arc<I>,
submit: impl FnOnce(&I, HostOperationId) -> io::Result<()>,
+36
View File
@@ -0,0 +1,36 @@
//! Host-backed message tunnels.
//!
//! The host owns the concrete transport and exposes complete EasyTier tunnel
//! payloads. Core owns packet semantics and the peer lifecycle above them.
use std::{io, task::Poll};
use super::socket::{HostOperationId, HostSocketHandle, HostSocketIo};
/// Maximum message payload accepted from a host tunnel adapter.
pub const MAX_HOST_TUNNEL_PAYLOAD_LEN: usize = 1024 * 1024;
/// Mechanical host I/O below EasyTier's message-tunnel seam.
///
/// Submit methods copy their complete input before returning. A receive keeps
/// the host transport's message boundary intact. A clean remote close is
/// reported as [`io::ErrorKind::UnexpectedEof`] by `take_receive`.
pub trait HostTunnelIo: HostSocketIo {
fn submit_receive(
&self,
handle: HostSocketHandle,
operation: HostOperationId,
capacity: usize,
) -> io::Result<()>;
fn take_receive(&self, operation: HostOperationId) -> Poll<io::Result<Vec<u8>>>;
fn submit_send(
&self,
handle: HostSocketHandle,
operation: HostOperationId,
source: &[u8],
) -> io::Result<()>;
fn take_send(&self, operation: HostOperationId) -> Poll<io::Result<()>>;
}
+11 -1
View File
@@ -28,7 +28,7 @@ use crate::{
use easytier_proto::common::CompressionAlgoPb;
use super::{CoreConnectivityConfig, CoreInstanceConfig};
use super::{CoreConnectivityConfig, CoreConnectivityMode, CoreInstanceConfig};
const OSPF_UPDATE_MY_FOREIGN_NETWORK_INTERVAL_SEC: u64 = 10;
const MAX_DIRECT_CONNS_PER_PEER_IN_FOREIGN_NETWORK: usize = 3;
@@ -58,6 +58,7 @@ pub struct CoreInstanceHostConfig {
pub upnp_enabled: bool,
pub tcp_hole_punching_enabled: bool,
pub ignore_unsupported_config: bool,
pub connectivity: CoreConnectivityMode,
pub easytier_version: String,
pub endpoint_protocols: Vec<String>,
}
@@ -83,6 +84,7 @@ impl Default for CoreInstanceHostConfig {
upnp_enabled: true,
tcp_hole_punching_enabled: true,
ignore_unsupported_config: false,
connectivity: CoreConnectivityMode::Full,
easytier_version: env!("CARGO_PKG_VERSION").to_owned(),
endpoint_protocols: ManualEndpointDiscoveryConfig::default().srv_protocols,
}
@@ -352,6 +354,8 @@ impl CoreInstanceConfig {
runtime,
startup_plan: super::CoreInstanceStartupPlan {
gateway: host.gateway_enabled,
packet_proxy: host.proxy_enabled,
connectivity: host.connectivity,
},
stun: StunServerConfig {
udp_servers: stun_servers
@@ -588,6 +592,7 @@ disable_p2p = true
icmp_failure_is_fatal: true,
public_ipv6_provider_supported: true,
gateway_enabled: false,
connectivity: CoreConnectivityMode::InboundOnly,
easytier_version: "host-version".to_owned(),
endpoint_protocols: vec!["host-protocol".to_owned()],
..Default::default()
@@ -606,6 +611,10 @@ disable_p2p = true
.as_deref(),
Some("host-fallback")
);
assert_eq!(
normalized.connectivity.startup_plan.connectivity,
CoreConnectivityMode::InboundOnly
);
assert!(
normalized
.peer
@@ -677,6 +686,7 @@ data_compress_algo = "Zstd"
let normalized = CoreInstanceConfig::from_toml_with_host(&config, &host).unwrap();
assert_eq!(config.dump(), before);
assert!(!normalized.connectivity.startup_plan.packet_proxy);
assert_eq!(normalized.connectivity.initial_peers.len(), 1);
assert_eq!(
normalized
+27 -9
View File
@@ -111,10 +111,16 @@ where
.ok_or_else(|| anyhow::anyhow!("packet egress is one-shot and already started"))?;
self.packet_egress.start(packet_receiver).await?;
self.peer_manager.run().await.map_err(anyhow::Error::from)?;
self.direct.run();
if let Some(direct) = &self.direct {
direct.run();
}
#[cfg(feature = "tcp-hole-punch")]
self.tcp_hole_punch.run();
self.manual.start();
if let Some(tcp_hole_punch) = &self.tcp_hole_punch {
tcp_hole_punch.run();
}
if let Some(manual) = &self.manual {
manual.start();
}
#[cfg(feature = "public-ipv6-provider")]
self.public_ipv6_provider.start().await;
@@ -128,8 +134,12 @@ where
wrapped_transport.start().await?;
}
#[cfg(feature = "proxy-packet")]
self.start_packet_proxy().await?;
self.udp_hole_punch.start().await?;
if self.startup_plan.packet_proxy {
self.start_packet_proxy().await?;
}
if let Some(udp_hole_punch) = &self.udp_hole_punch {
udp_hole_punch.start().await?;
}
self.peer_center.init().await;
#[cfg(feature = "proxy-cidr-monitor")]
self.proxy_cidr_monitor
@@ -170,7 +180,9 @@ where
if let Some(listener) = &self.listener {
listener.stop().await;
}
self.udp_hole_punch.stop().await;
if let Some(udp_hole_punch) = &self.udp_hole_punch {
udp_hole_punch.stop().await;
}
#[cfg(feature = "proxy-smoltcp-stack")]
self.port_forward_adapter.stop().await;
#[cfg(feature = "proxy-smoltcp-stack")]
@@ -185,10 +197,16 @@ where
}
#[cfg(feature = "proxy-packet")]
self.packet_proxy.stop().await;
self.manual.stop().await;
if let Some(manual) = &self.manual {
manual.stop().await;
}
#[cfg(feature = "tcp-hole-punch")]
self.tcp_hole_punch.stop().await;
self.direct.stop().await;
if let Some(tcp_hole_punch) = &self.tcp_hole_punch {
tcp_hole_punch.stop().await;
}
if let Some(direct) = &self.direct {
direct.stop().await;
}
self.peer_center.stop().await;
// Host packet tasks can still call the packet plane, so stop them
+19 -8
View File
@@ -30,19 +30,29 @@ where
}
pub fn add_connector(&self, url: Url) -> anyhow::Result<()> {
self.manual.add_connector(url)
self.manual
.as_ref()
.ok_or_else(|| anyhow::anyhow!("inbound-only connectivity has no outbound connectors"))?
.add_connector(url)
}
pub fn remove_connector(&self, url: &Url) -> bool {
self.manual.remove_connector(url)
self.manual
.as_ref()
.is_some_and(|manual| manual.remove_connector(url))
}
pub fn clear_connectors(&self) {
self.manual.clear_connectors();
if let Some(manual) = &self.manual {
manual.clear_connectors();
}
}
pub fn list_connectors(&self) -> Vec<ManualConnectorSnapshot> {
self.manual.list_connectors()
self.manual
.as_ref()
.map(|manual| manual.list_connectors())
.unwrap_or_default()
}
pub fn running_listeners(&self) -> Vec<Url> {
@@ -92,10 +102,11 @@ where
.peer_manager
.node_snapshot(self.running_listeners())
.await;
snapshot.ip_list = self
.direct
.local_address_observations_with_stun(&snapshot.stun_info)
.await;
if let Some(direct) = &self.direct {
snapshot.ip_list = direct
.local_address_observations_with_stun(&snapshot.stun_info)
.await;
}
snapshot
}
+170 -83
View File
@@ -64,8 +64,8 @@ use crate::{
packet::{HostPacketReceiver, PacketSink, host_packet_channel},
},
listener::{
AcceptedSocketHandler, ExternalListenerFactory, ExternalListenerRequest, ListenerFactory,
RunningListenerRegistry,
AcceptedSocketHandler, ExternalListenerFactory, ExternalListenerRequest,
HostListenerRegistration, ListenerFactory, RunningListenerRegistry,
plan::{ListenerRuntimeConfig, PreparedListenerPlan, prepare_listener_plan},
transport::{
AcceptedTransport, CoreListenerRuntime, HostAcceptedTcpSocket,
@@ -146,9 +146,18 @@ impl CoreInstanceState {
}
}
fn default_true() -> bool {
true
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct CoreInstanceStartupPlan {
#[serde(default = "default_true")]
pub gateway: bool,
#[serde(default = "default_true")]
pub packet_proxy: bool,
#[serde(default)]
pub connectivity: CoreConnectivityMode,
}
impl CoreInstanceStartupPlan {
@@ -159,10 +168,26 @@ impl CoreInstanceStartupPlan {
impl Default for CoreInstanceStartupPlan {
fn default() -> Self {
Self { gateway: true }
Self {
gateway: true,
packet_proxy: true,
connectivity: CoreConnectivityMode::Full,
}
}
}
/// Selects which portable connectivity Modules participate in one instance.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "snake_case")]
pub enum CoreConnectivityMode {
#[default]
Full,
/// Dial configured peers without listeners, discovery, or direct connectivity.
OutboundOnly,
/// Accept Host-registered listeners without constructing outbound socket Modules.
InboundOnly,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct CoreConnectivityConfig {
pub initial_peers: Vec<Url>,
@@ -279,6 +304,7 @@ where
pub protocol: Option<Arc<dyn ClientProtocolUpgrader<<H as VirtualTcpSocketFactory>::Socket>>>,
pub external_listener_factory:
Option<Arc<dyn ExternalListenerFactory<AcceptedTransport<HostAcceptedTcpSocket<H>>>>>,
pub host_listener_registrations: Vec<HostListenerRegistration>,
pub server_protocol: Option<Arc<dyn ServerProtocolUpgrader<HostAcceptedTcpSocket<H>>>>,
/// Optional OS port-mapping adapter. STUN-only hole punching remains
/// available when the host does not provide one.
@@ -338,6 +364,7 @@ where
wrapped_transports: WrappedTransportEngines::default(),
protocol: None,
external_listener_factory: None,
host_listener_registrations: Vec::new(),
server_protocol: None,
udp_hole_punch_platform: None,
#[cfg(feature = "proxy-packet")]
@@ -380,13 +407,13 @@ where
pub(super) cancel: CancellationToken,
pub(super) peer_manager: Arc<PeerManagerCore>,
packet_plane: Arc<CorePacketPlane>,
pub(super) manual: ManualConnectorManager<H>,
pub(super) direct: DirectConnectorManager<H>,
pub(super) manual: Option<ManualConnectorManager<H>>,
pub(super) direct: Option<DirectConnectorManager<H>>,
#[cfg(feature = "tcp-hole-punch")]
tcp_hole_punch: TcpHolePunchConnector<H, PeerManagerCore>,
tcp_hole_punch: Option<TcpHolePunchConnector<H, PeerManagerCore>>,
pub(super) listener: Option<Arc<CoreListenerRuntime<H>>>,
running_listeners: Arc<RunningListenerRegistry>,
pub(super) udp_hole_punch: CoreUdpHolePunchService<H, PeerManagerCore>,
pub(super) udp_hole_punch: Option<CoreUdpHolePunchService<H, PeerManagerCore>>,
#[cfg(feature = "wrapped-transport")]
wrapped_transport: Option<Arc<WrappedTransportProxyModule>>,
#[cfg(feature = "proxy-smoltcp-stack")]
@@ -411,7 +438,7 @@ where
public_ipv6_provider: PublicIpv6ProviderRuntime,
#[cfg(feature = "vpn-portal")]
vpn_portal: Arc<PortalModule>,
#[cfg(feature = "proxy-smoltcp-stack")]
#[cfg(feature = "proxy-packet")]
pub(super) startup_plan: CoreInstanceStartupPlan,
pub(super) runtime_config: CoreRuntimeConfigStore,
#[cfg(feature = "test-utils")]
@@ -472,6 +499,7 @@ where
mut adapters: CoreHostAdapters<H>,
) -> anyhow::Result<Arc<Self>> {
build_capabilities::validate(&config)?;
let connectivity_mode = config.connectivity.startup_plan.connectivity;
let instance_name = config.instance_name;
#[cfg(feature = "vpn-portal")]
let vpn_portal_config = config.vpn_portal.clone();
@@ -497,17 +525,26 @@ where
public_ipv6_host,
public_ipv6_events,
);
let stun = Self::prepare_stun(&adapters, &config.connectivity);
let peer_stun: Arc<dyn StunInfoProvider> = stun.clone();
let foreign_rpc_registrar = Arc::new(ForeignDirectConnectorRpcRegistrar::new(
adapters.host.clone(),
stun.clone(),
));
let stun = (connectivity_mode == CoreConnectivityMode::Full)
.then(|| Self::prepare_stun(&adapters, &config.connectivity));
let (peer_stun, foreign_rpc_registrar): (
Arc<dyn PeerStunInfoSource>,
Arc<dyn crate::peers::foreign_network::ForeignNetworkRpcRegistrar>,
) = match &stun {
Some(stun) => (
Arc::new(CoreStunPeerInfoSource(stun.clone())),
Arc::new(ForeignDirectConnectorRpcRegistrar::new(
adapters.host.clone(),
stun.clone(),
)),
),
None => (Arc::new(()), Arc::new(())),
};
let peer_manager = Arc::new(PeerManagerCore::new(
config.peer,
config.managed_credentials,
runtime_config.clone(),
Arc::new(CoreStunPeerInfoSource(peer_stun)),
peer_stun,
packet_tx,
public_ipv6_runtime.clone(),
events.clone(),
@@ -515,8 +552,11 @@ where
foreign_rpc_registrar,
)?);
let config = config.connectivity;
let configured_listeners = (connectivity_mode != CoreConnectivityMode::OutboundOnly)
.then_some(config.listeners.as_ref())
.flatten();
let listener_plan = prepare_listener_plan(
config.listeners.as_ref(),
configured_listeners,
peer_manager.instance_id(),
adapters.server_protocol.as_deref(),
adapters.external_listener_factory.as_deref(),
@@ -536,6 +576,7 @@ where
wrapped_transports,
protocol,
external_listener_factory,
host_listener_registrations,
server_protocol,
udp_hole_punch_platform,
#[cfg(feature = "proxy-packet")]
@@ -549,6 +590,12 @@ where
#[cfg(feature = "vpn-portal")]
vpn_portal,
} = adapters;
let host_listener_registrations = if connectivity_mode == CoreConnectivityMode::OutboundOnly
{
Vec::new()
} else {
host_listener_registrations
};
let dns_records: Arc<dyn DnsRecordResolver> = dns.clone();
let dns: Arc<dyn DnsResolver> = dns;
let ring_registry = process_runtime.ring_registry();
@@ -563,19 +610,22 @@ where
manual: manual_options,
direct: direct_options,
} = config;
#[cfg(not(feature = "proxy-smoltcp-stack"))]
#[cfg(not(feature = "proxy-packet"))]
let _ = startup_plan;
if connectivity_mode == CoreConnectivityMode::InboundOnly && !initial_peers.is_empty() {
anyhow::bail!("inbound-only connectivity does not support outbound peers");
}
let accepted_tunnel_handler = PeerAcceptedTunnelHandler::new(&peer_manager, events.clone());
let accepted_transport_handler: Arc<
dyn AcceptedSocketHandler<AcceptedTransport<HostAcceptedTcpSocket<H>>>,
> = match server_protocol {
Some(server_protocol) => {
let tunnel_handler = PeerAcceptedTunnelHandler::new(&peer_manager, events.clone());
Arc::new(ProtocolAcceptedTransportHandler::new(
&tunnel_handler,
server_protocol,
))
}
None => Arc::new(RawAcceptedTransportHandler::new(&peer_manager)),
Some(server_protocol) => Arc::new(ProtocolAcceptedTransportHandler::new(
&accepted_tunnel_handler,
server_protocol,
)),
None => Arc::new(RawAcceptedTransportHandler::new(
accepted_tunnel_handler.clone(),
)),
};
let running_listeners = Arc::new(RunningListenerRegistry::default());
let PreparedListenerPlan {
@@ -583,8 +633,11 @@ where
external,
failures,
} = listener_plan;
let mut external_factories = Vec::with_capacity(external.len());
if !external.is_empty() && external_listener_factory.is_none() {
let mut external_factories =
Vec::with_capacity(external.len() + host_listener_registrations.len());
if (!external.is_empty() || !host_listener_registrations.is_empty())
&& external_listener_factory.is_none()
{
anyhow::bail!("listener plan requires an external listener factory");
}
for (listener, socket_context) in external {
@@ -598,6 +651,19 @@ where
listener.must_succeed,
));
}
for request in host_listener_registrations {
let factory = external_listener_factory.clone().unwrap();
if !factory.supports_scheme(request.url.scheme()) {
anyhow::bail!(
"external listener factory does not support Host listener scheme {}",
request.url.scheme()
);
}
external_factories.push(ListenerFactory::new(
move || factory.create(request.clone()),
true,
));
}
let has_listener_work =
!transports.is_empty() || !external_factories.is_empty() || !failures.is_empty();
let listener = has_listener_work.then(|| {
@@ -613,40 +679,51 @@ where
running_listeners.clone(),
))
});
let protocol = protocol.unwrap_or_else(|| {
Arc::new(CoreClientProtocolUpgrader::new(
CoreClientProtocolConfig::default(),
))
let protocol = (connectivity_mode != CoreConnectivityMode::InboundOnly).then(|| {
protocol.unwrap_or_else(|| {
Arc::new(CoreClientProtocolUpgrader::new(
CoreClientProtocolConfig::default(),
))
})
});
let endpoint_resolver = Arc::new(CoreManualEndpointResolver::new(
host.clone(),
dns.clone(),
dns_records,
endpoint_discovery,
));
let manual = ManualConnectorManager::new(
peer_manager.clone(),
host.clone(),
dns.clone(),
endpoint_resolver,
protocol.clone(),
ring_registry,
manual_options,
events.clone(),
);
for url in initial_peers {
manual.add_connector(url)?;
}
let udp_hole_punch_socket_context = direct_options.udp_bind.context.clone();
let udp_hole_punch = CoreUdpHolePunchService::new(
peer_manager.clone(),
host.clone(),
stun.clone(),
udp_hole_punch_platform,
events.clone(),
udp_hole_punch_socket_context,
protocol.clone(),
);
let manual = if let Some(protocol) = &protocol {
let endpoint_resolver = Arc::new(CoreManualEndpointResolver::new(
host.clone(),
dns.clone(),
dns_records,
endpoint_discovery,
));
let manual = ManualConnectorManager::new(
peer_manager.clone(),
host.clone(),
dns.clone(),
endpoint_resolver,
protocol.clone(),
ring_registry,
manual_options,
events.clone(),
);
for url in initial_peers {
manual.add_connector(url)?;
}
Some(manual)
} else {
None
};
let udp_hole_punch = stun
.as_ref()
.zip(protocol.as_ref())
.map(|(stun, protocol)| {
CoreUdpHolePunchService::new(
peer_manager.clone(),
host.clone(),
stun.clone(),
udp_hole_punch_platform,
events.clone(),
direct_options.udp_bind.context.clone(),
protocol.clone(),
)
});
let proxy_cidr_table = Arc::new(ProxyCidrTable::from_snapshot(proxy_cidr_snapshot(
runtime_config.snapshot().as_ref(),
)));
@@ -708,28 +785,38 @@ where
events.clone(),
);
#[cfg(feature = "tcp-hole-punch")]
let tcp_hole_punch = TcpHolePunchConnector::new(
peer_manager.clone(),
host.clone(),
stun.clone(),
direct_options.tcp_bind.context.clone(),
protocol.clone(),
Arc::new(crate::connectivity::protocol::CoreServerProtocolUpgrader::<
HostAcceptedTcpSocket<H>,
>::new(
crate::connectivity::protocol::CoreServerProtocolConfig::default(),
)),
);
let direct = DirectConnectorManager::new_with_running_listeners(
peer_manager.clone(),
host.clone(),
protected_tcp_ports,
stun.clone(),
running_listeners.clone(),
dns,
protocol,
direct_options,
);
let tcp_hole_punch = stun
.as_ref()
.zip(protocol.as_ref())
.map(|(stun, protocol)| {
TcpHolePunchConnector::new(
peer_manager.clone(),
host.clone(),
stun.clone(),
direct_options.tcp_bind.context.clone(),
protocol.clone(),
Arc::new(crate::connectivity::protocol::CoreServerProtocolUpgrader::<
HostAcceptedTcpSocket<H>,
>::new(
crate::connectivity::protocol::CoreServerProtocolConfig::default(),
)),
)
});
let direct = match (stun, protocol) {
(Some(stun), Some(protocol)) => {
Some(DirectConnectorManager::new_with_running_listeners(
peer_manager.clone(),
host.clone(),
protected_tcp_ports,
stun,
running_listeners.clone(),
dns,
protocol,
direct_options,
))
}
_ => None,
};
let peer_center = Arc::new(PeerCenterInstance::new(peer_manager.clone()));
#[cfg(feature = "public-ipv6-provider")]
let public_ipv6_provider = PublicIpv6ProviderRuntime::new(
@@ -799,7 +886,7 @@ where
public_ipv6_provider,
#[cfg(feature = "vpn-portal")]
vpn_portal,
#[cfg(feature = "proxy-smoltcp-stack")]
#[cfg(feature = "proxy-packet")]
startup_plan,
runtime_config,
#[cfg(feature = "test-utils")]
+88
View File
@@ -185,6 +185,10 @@ fn core_instance_config_round_trips_as_normalized_json() {
assert!(!decoded.connectivity.direct.testing);
assert!(decoded.connectivity.startup_plan.gateway);
assert_eq!(
decoded.connectivity.startup_plan.connectivity,
CoreConnectivityMode::Full
);
assert_eq!(serde_json::to_value(&decoded).unwrap(), encoded);
let mut legacy = encoded;
@@ -3016,4 +3020,88 @@ virtual_ip = "10.82.0.2/24"
instance.stop().await;
assert!(instance.running_listeners().is_empty());
}
#[tokio::test]
async fn inbound_only_uses_host_registered_listener_lifecycle() {
let external_url: Url = "unix:///tmp/easytier-host-listener-test".parse().unwrap();
let mut config = test_config("host-listener");
config.connectivity.startup_plan.connectivity = CoreConnectivityMode::InboundOnly;
let (packet_sink, _packet_receiver) = tokio::sync::mpsc::channel(16);
let mut adapters = adapters(
Some(Arc::new(ReadyExternalListenerFactory)),
Arc::new(packet_sink),
);
adapters
.host_listener_registrations
.push(ExternalListenerRequest {
url: external_url.clone(),
socket_context: SocketContext::default(),
});
let instance = CoreInstance::new(config, adapters).unwrap();
assert!(
instance
.add_connector("tcp://127.0.0.1:11010".parse().unwrap())
.is_err()
);
instance.start().await.unwrap();
assert!(instance.running_listeners().contains(&external_url));
instance.stop().await;
assert!(instance.running_listeners().is_empty());
}
#[tokio::test]
async fn inbound_only_rejects_initial_peers() {
let mut config = test_config("inbound-only-peer");
config.connectivity.startup_plan.connectivity = CoreConnectivityMode::InboundOnly;
config
.connectivity
.initial_peers
.push("tcp://127.0.0.1:11010".parse().unwrap());
let Err(error) = build_instance(config) else {
panic!("inbound-only instance accepted an outbound peer");
};
assert!(
error
.to_string()
.contains("inbound-only connectivity does not support outbound peers"),
"unexpected construction error: {error:#}"
);
}
#[tokio::test]
async fn outbound_only_ignores_configured_and_host_registered_listeners() {
let configured_url: Url = "unix:///tmp/easytier-outbound-configured-listener"
.parse()
.unwrap();
let host_url: Url = "unix:///tmp/easytier-outbound-host-listener"
.parse()
.unwrap();
let mut config = test_config("outbound-only-listeners");
config.connectivity.startup_plan.connectivity = CoreConnectivityMode::OutboundOnly;
config.connectivity.listeners = Some(ListenerRuntimeConfig::new(
vec![configured_url],
false,
SocketContext::default(),
));
let (packet_sink, _packet_receiver) = tokio::sync::mpsc::channel(16);
let mut adapters = adapters(
Some(Arc::new(ReadyExternalListenerFactory)),
Arc::new(packet_sink),
);
adapters
.host_listener_registrations
.push(ExternalListenerRequest {
url: host_url,
socket_context: SocketContext::default(),
});
let instance = CoreInstance::new(config, adapters).unwrap();
instance.start().await.unwrap();
assert!(instance.running_listeners().is_empty());
instance.stop().await;
}
}
+12
View File
@@ -15,6 +15,14 @@ use crate::{
};
pub mod plan;
#[cfg(any(
test,
all(
feature = "wasm-host-tunnel",
not(feature = "wasm-host-tunnel-outbound")
)
))]
pub(crate) mod queue;
pub mod transport;
pub trait ExternalListenerFactory<Accepted>: Send + Sync + 'static
@@ -35,6 +43,10 @@ pub struct ExternalListenerRequest {
pub socket_context: SocketContext,
}
/// One listener supplied by the Host rather than the portable TOML model.
/// Host listeners are mandatory: startup fails when one cannot bind.
pub type HostListenerRegistration = ExternalListenerRequest;
#[async_trait]
pub trait AcceptedSocketHandler<Accepted>: Send + Sync {
async fn handle_accepted_socket(&self, accepted: Accepted) -> anyhow::Result<()>;
+149
View File
@@ -0,0 +1,149 @@
use std::{collections::VecDeque, sync::Mutex};
use tokio::sync::Notify;
struct HostListenerQueueState<T> {
closed: bool,
listeners: usize,
pending: VecDeque<T>,
}
/// Bounded handoff from a synchronous Host callback to async listeners.
pub(crate) struct HostListenerQueue<T> {
capacity: usize,
state: Mutex<HostListenerQueueState<T>>,
changed: Notify,
}
impl<T> HostListenerQueue<T> {
pub(crate) fn new(capacity: usize) -> Self {
Self {
capacity,
state: Mutex::new(HostListenerQueueState {
closed: false,
listeners: 0,
pending: VecDeque::new(),
}),
changed: Notify::new(),
}
}
pub(crate) fn register_listener(&self) -> bool {
let mut state = self.state.lock().unwrap();
if state.closed {
return false;
}
state.listeners += 1;
true
}
pub(crate) fn unregister_listener(&self) {
let pending = {
let mut state = self.state.lock().unwrap();
debug_assert!(state.listeners > 0);
state.listeners -= 1;
if state.listeners != 0 {
return;
}
state.closed = true;
std::mem::take(&mut state.pending)
};
drop(pending);
self.changed.notify_waiters();
}
/// Constructs `T` only after the queue accepts Host-to-guest ownership.
pub(crate) fn enqueue_with(&self, create: impl FnOnce() -> T) -> anyhow::Result<()> {
{
let mut state = self.state.lock().unwrap();
if state.closed || state.listeners == 0 {
anyhow::bail!("Host listener queue is closed");
}
if state.pending.len() >= self.capacity {
anyhow::bail!("Host listener admission queue is full");
}
state.pending.push_back(create());
}
self.changed.notify_one();
Ok(())
}
pub(crate) async fn accept(&self) -> Option<T> {
loop {
let changed = self.changed.notified();
{
let mut state = self.state.lock().unwrap();
if let Some(item) = state.pending.pop_front() {
return Some(item);
}
if state.closed {
return None;
}
}
changed.await;
}
}
}
#[cfg(test)]
mod tests {
use std::sync::{
Arc,
atomic::{AtomicUsize, Ordering},
};
use super::HostListenerQueue;
struct DropCounter(Arc<AtomicUsize>);
impl Drop for DropCounter {
fn drop(&mut self) {
self.0.fetch_add(1, Ordering::Relaxed);
}
}
#[tokio::test]
async fn constructs_only_after_listener_accepts_ownership() {
let queue = HostListenerQueue::new(1);
let constructed = AtomicUsize::new(0);
assert!(
queue
.enqueue_with(|| {
constructed.fetch_add(1, Ordering::Relaxed);
})
.is_err()
);
assert_eq!(constructed.load(Ordering::Relaxed), 0);
assert!(queue.register_listener());
queue
.enqueue_with(|| {
constructed.fetch_add(1, Ordering::Relaxed);
})
.unwrap();
assert!(
queue
.enqueue_with(|| {
constructed.fetch_add(1, Ordering::Relaxed);
})
.is_err()
);
assert_eq!(constructed.load(Ordering::Relaxed), 1);
queue.accept().await.unwrap();
queue.unregister_listener();
}
#[test]
fn last_listener_closes_and_drains_pending_items() {
let queue = HostListenerQueue::new(1);
let drops = Arc::new(AtomicUsize::new(0));
assert!(queue.register_listener());
queue.enqueue_with(|| DropCounter(drops.clone())).unwrap();
queue.unregister_listener();
assert_eq!(drops.load(Ordering::Relaxed), 1);
assert!(!queue.register_listener());
}
}
+46 -11
View File
@@ -77,14 +77,12 @@ impl AcceptedTunnelHandler for PeerAcceptedTunnelHandler {
}
pub(crate) struct RawAcceptedTransportHandler {
peer_manager: Weak<PeerManagerCore>,
tunnel_handler: Arc<dyn AcceptedTunnelHandler>,
}
impl RawAcceptedTransportHandler {
pub(crate) fn new(peer_manager: &Arc<PeerManagerCore>) -> Self {
Self {
peer_manager: Arc::downgrade(peer_manager),
}
pub(crate) fn new(tunnel_handler: Arc<dyn AcceptedTunnelHandler>) -> Self {
Self { tunnel_handler }
}
}
@@ -97,10 +95,6 @@ where
&self,
accepted: AcceptedTransport<TcpSocket>,
) -> anyhow::Result<()> {
let peer_manager = self
.peer_manager
.upgrade()
.ok_or_else(|| anyhow::anyhow!("peer manager is gone"))?;
let tunnel = match accepted {
AcceptedTransport::Tunnel { tunnel, .. } => tunnel,
AcceptedTransport::Tcp {
@@ -125,7 +119,48 @@ where
remote_url,
} => raw::upgrade_accepted_byte_stream(socket, local_url, remote_url)?,
};
peer_manager.add_tunnel_as_server(tunnel, true).await?;
Ok(())
self.tunnel_handler.handle_tunnel(tunnel).await
}
}
#[cfg(test)]
mod tests {
use std::sync::{
Arc,
atomic::{AtomicUsize, Ordering},
};
use super::*;
use crate::{host::testkit::TestTcpSocket, tunnel::ring::RingTunnelRegistry};
struct RecordingTunnelHandler(AtomicUsize);
#[async_trait]
impl AcceptedTunnelHandler for RecordingTunnelHandler {
async fn handle_tunnel(&self, _tunnel: Box<dyn Tunnel>) -> anyhow::Result<()> {
self.0.fetch_add(1, Ordering::Relaxed);
Ok(())
}
}
#[tokio::test]
async fn raw_transport_delegates_tunnel_to_shared_admission_handler() {
let registry = Arc::new(RingTunnelRegistry::default());
let local_id = uuid::Uuid::new_v4();
let mut listener = registry.bind(local_id).unwrap();
let _client = registry.connect(local_id).unwrap();
let tunnel = listener.accept().await.unwrap().into_tunnel();
let recorder = Arc::new(RecordingTunnelHandler(AtomicUsize::new(0)));
let handler = RawAcceptedTransportHandler::new(recorder.clone());
handler
.handle_accepted_socket(AcceptedTransport::<TestTcpSocket>::Tunnel {
tunnel,
local_url: format!("ring://{local_id}").parse().unwrap(),
})
.await
.unwrap();
assert_eq!(recorder.0.load(Ordering::Relaxed), 1);
}
}
+221
View File
@@ -0,0 +1,221 @@
//! EasyTier message tunnel over a host-owned transport.
use std::{io, sync::Arc};
use futures::{sink, stream};
use url::Url;
use crate::{
host::{
socket::{HostSocketHandle, HostSocketRuntime},
tunnel::{HostTunnelIo, MAX_HOST_TUNNEL_PAYLOAD_LEN},
},
packet::{ZCPacket, ZCPacketType},
proto::common::TunnelInfo,
};
use super::{Tunnel, TunnelError, wrapper::TunnelWrapper};
struct HostTunnelResource {
io: Arc<dyn HostTunnelIo>,
handle: HostSocketHandle,
}
impl Drop for HostTunnelResource {
fn drop(&mut self) {
let _ = self.io.close(self.handle);
}
}
/// Builds a message-preserving EasyTier tunnel around a host-owned transport.
pub fn new_host_tunnel(
runtime: HostSocketRuntime,
io: Arc<dyn HostTunnelIo>,
handle: HostSocketHandle,
local_url: Url,
remote_url: Url,
resolved_remote_url: Option<Url>,
) -> Box<dyn Tunnel> {
let resource = Arc::new(HostTunnelResource { io, handle });
let reader_resource = resource.clone();
let reader_runtime = runtime.clone();
let reader = stream::unfold(
(reader_runtime, reader_resource),
|(runtime, resource)| async move {
let result = runtime
.run_operation(
resource.io.clone(),
|io, operation| {
io.submit_receive(resource.handle, operation, MAX_HOST_TUNNEL_PAYLOAD_LEN)
},
|io, operation| io.take_receive(operation),
|io, operation| io.cancel_operation(operation),
)
.await;
let item = match result {
Ok(message) => Ok(ZCPacket::new_from_buf(
bytes::BytesMut::from(message.as_slice()),
ZCPacketType::DummyTunnel,
)),
Err(error) if error.kind() == io::ErrorKind::UnexpectedEof => return None,
Err(error) => Err(TunnelError::IOError(error)),
};
Some((item, (runtime, resource)))
},
);
let writer = sink::unfold(
(runtime, resource),
|(runtime, resource), packet: ZCPacket| async move {
let payload = packet.tunnel_payload_bytes();
runtime
.run_operation(
resource.io.clone(),
|io, operation| io.submit_send(resource.handle, operation, &payload),
|io, operation| io.take_send(operation),
|io, operation| io.cancel_operation(operation),
)
.await
.map_err(TunnelError::IOError)?;
Ok((runtime, resource))
},
);
let remote_addr = remote_url.clone().into();
let resolved_remote_addr = resolved_remote_url
.unwrap_or_else(|| remote_url.clone())
.into();
let info = TunnelInfo {
tunnel_type: local_url.scheme().to_owned(),
local_addr: Some(local_url.into()),
remote_addr: Some(remote_addr),
resolved_remote_addr: Some(resolved_remote_addr),
};
Box::new(TunnelWrapper::new(reader, writer, Some(info)))
}
#[cfg(test)]
mod tests {
use std::{
collections::{HashMap, VecDeque},
sync::{
Mutex,
atomic::{AtomicUsize, Ordering},
},
task::Poll,
};
use futures::{SinkExt as _, StreamExt as _};
use super::*;
use crate::host::socket::{HostOperationId, HostSocketIo};
#[derive(Default)]
struct MockTunnelIo {
incoming: Mutex<VecDeque<io::Result<Vec<u8>>>>,
receives: Mutex<HashMap<HostOperationId, io::Result<Vec<u8>>>>,
sends: Mutex<HashMap<HostOperationId, Vec<u8>>>,
sent: Mutex<Vec<Vec<u8>>>,
closes: AtomicUsize,
}
impl HostSocketIo for MockTunnelIo {
fn cancel_operation(&self, operation: HostOperationId) -> io::Result<()> {
self.receives.lock().unwrap().remove(&operation);
self.sends.lock().unwrap().remove(&operation);
Ok(())
}
fn close(&self, _handle: HostSocketHandle) -> io::Result<()> {
self.closes.fetch_add(1, Ordering::SeqCst);
Ok(())
}
}
impl HostTunnelIo for MockTunnelIo {
fn submit_receive(
&self,
_handle: HostSocketHandle,
operation: HostOperationId,
_capacity: usize,
) -> io::Result<()> {
let result = self.incoming.lock().unwrap().pop_front().unwrap();
self.receives.lock().unwrap().insert(operation, result);
Ok(())
}
fn take_receive(&self, operation: HostOperationId) -> Poll<io::Result<Vec<u8>>> {
Poll::Ready(self.receives.lock().unwrap().remove(&operation).unwrap())
}
fn submit_send(
&self,
_handle: HostSocketHandle,
operation: HostOperationId,
source: &[u8],
) -> io::Result<()> {
self.sends
.lock()
.unwrap()
.insert(operation, source.to_vec());
Ok(())
}
fn take_send(&self, operation: HostOperationId) -> Poll<io::Result<()>> {
let message = self.sends.lock().unwrap().remove(&operation).unwrap();
self.sent.lock().unwrap().push(message);
Poll::Ready(Ok(()))
}
}
fn tunnel(io: Arc<MockTunnelIo>) -> Box<dyn Tunnel> {
new_host_tunnel(
HostSocketRuntime::new(),
io,
HostSocketHandle(7),
Url::parse("test-tunnel://relay.example/").unwrap(),
Url::parse("test-tunnel://client.example/").unwrap(),
None,
)
}
#[test]
fn preserves_message_boundaries_and_closes_once() {
let io = Arc::new(MockTunnelIo::default());
io.incoming.lock().unwrap().push_back(Ok(vec![1, 2, 3]));
io.incoming.lock().unwrap().push_back(Ok(vec![8, 9]));
let tunnel = tunnel(io.clone());
assert_eq!(tunnel.info().unwrap().tunnel_type, "test-tunnel");
let (mut reader, mut writer) = tunnel.split();
let first = futures::executor::block_on(reader.next()).unwrap().unwrap();
let second = futures::executor::block_on(reader.next()).unwrap().unwrap();
assert_eq!(first.tunnel_payload(), &[1, 2, 3]);
assert_eq!(second.tunnel_payload(), &[8, 9]);
let packet = ZCPacket::new_from_buf(
bytes::BytesMut::from(&[4, 5, 6][..]),
ZCPacketType::DummyTunnel,
);
futures::executor::block_on(writer.send(packet)).unwrap();
assert_eq!(*io.sent.lock().unwrap(), vec![vec![4, 5, 6]]);
drop(tunnel);
drop(reader);
drop(writer);
assert_eq!(io.closes.load(Ordering::SeqCst), 1);
}
#[test]
fn maps_clean_remote_close_to_stream_eof() {
let io = Arc::new(MockTunnelIo::default());
io.incoming
.lock()
.unwrap()
.push_back(Err(io::Error::new(io::ErrorKind::UnexpectedEof, "closed")));
let tunnel = tunnel(io);
let (mut reader, _writer) = tunnel.split();
assert!(futures::executor::block_on(reader.next()).is_none());
}
}
+1
View File
@@ -18,6 +18,7 @@ pub use crate::socket::IpVersion;
pub(crate) mod encrypt;
pub mod filter;
pub mod framed;
pub mod host_tunnel;
pub mod mpsc;
pub mod ring;
pub(crate) mod secure_datagram;
+15 -1
View File
@@ -34,7 +34,7 @@ pub const CORE_INSTANCE_CONFIG_VERSION: u32 = 14;
pub const WEB_CLIENT_CONFIG_VERSION: u32 = 1;
/// Version of the public data-plane guest export contract.
pub const DATA_PLANE_ABI_VERSION: u32 = 3;
pub const DATA_PLANE_ABI_VERSION: u32 = 4;
/// Version of the protobuf RPC guest export contract.
#[cfg(feature = "management-rpc")]
@@ -54,6 +54,8 @@ pub const DATA_PLANE_UDP_CAPABILITY: u64 = 1 << 2;
pub const DATA_PLANE_DEADLINE_READ: u32 = 1 << 0;
/// Update the write deadline in `easytier_data_plane_resource_deadline_set`.
pub const DATA_PLANE_DEADLINE_WRITE: u32 = 1 << 1;
/// Version of the host tunnel metadata and export contract.
pub const HOST_TUNNEL_ABI_VERSION: u32 = 1;
/// Guest exports a WASI runtime calls to manage a core instance.
///
@@ -62,6 +64,9 @@ pub const DATA_PLANE_DEADLINE_WRITE: u32 = 1 << 1;
/// all asynchronous guest work through `easytier_instance_drive` and host
/// completion notifications.
pub const GUEST_EXPORTS: &[&str] = &[
// WASI command initialization. This only binds the runtime; core lifecycle
// still starts through the instance exports below.
"_start",
// Guest-memory buffers.
"easytier_buffer_alloc",
"easytier_buffer_free",
@@ -112,6 +117,13 @@ pub const RPC_GUEST_EXPORTS: &[&str] = &[
"easytier_rpc_operation_free",
];
/// Guest exports present with the host tunnel feature.
#[cfg(feature = "wasm-host-tunnel")]
pub const HOST_TUNNEL_GUEST_EXPORTS: &[&str] = &[
"easytier_host_tunnel_abi_version",
"easytier_instance_accept_tunnel",
];
/// Guest exports present when the core is built with the smoltcp data plane.
#[cfg(feature = "proxy-smoltcp-stack")]
pub const DATA_PLANE_GUEST_EXPORTS: &[&str] = &[
@@ -124,6 +136,7 @@ pub const DATA_PLANE_GUEST_EXPORTS: &[&str] = &[
"easytier_data_plane_tcp_accept_submit",
"easytier_data_plane_tcp_read_submit",
"easytier_data_plane_tcp_write_submit",
"easytier_data_plane_tcp_shutdown_write_submit",
"easytier_data_plane_udp_bind_submit",
"easytier_data_plane_udp_receive_submit",
"easytier_data_plane_udp_send_submit",
@@ -136,6 +149,7 @@ pub const DATA_PLANE_GUEST_EXPORTS: &[&str] = &[
"easytier_data_plane_tcp_accept_result_take",
"easytier_data_plane_tcp_read_result_take",
"easytier_data_plane_tcp_write_result_take",
"easytier_data_plane_tcp_shutdown_write_result_take",
"easytier_data_plane_udp_bind_result_take",
"easytier_data_plane_udp_receive_result_take",
"easytier_data_plane_udp_send_result_take",
+2
View File
@@ -7,3 +7,5 @@ pub mod event;
pub mod management;
pub mod packet;
pub mod socket;
#[cfg(feature = "wasm-host-tunnel")]
pub mod tunnel;
+397
View File
@@ -0,0 +1,397 @@
//! WASI imports for host-owned message tunnels.
use std::{
collections::HashMap,
io,
sync::{Arc, Mutex},
task::Poll,
};
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
use std::{fmt, marker::PhantomData};
use crate::{
host::{
socket::{HostOperationId, HostSocketHandle, HostSocketIo},
tunnel::HostTunnelIo,
},
tunnel::Tunnel,
wasi::{
imports::{
HOST_PENDING, HOST_TUNNEL_CLOSED, cancel_operation, close, start_tunnel_receive,
start_tunnel_send, take_tunnel_receive, take_tunnel_send,
},
wire::common::{host_error, status},
},
};
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
use crate::{
listener::{
ExternalListenerFactory, ExternalListenerRequest, queue::HostListenerQueue,
transport::AcceptedTransport,
},
socket::{SocketListener, tcp::VirtualTcpSocket},
};
#[cfg(feature = "wasm-host-tunnel-outbound")]
use crate::connectivity::manual::ExternalTunnelConnector;
#[cfg(feature = "wasm-host-tunnel-outbound")]
use crate::wasi::imports::{start_tunnel_connect, take_tunnel_connect};
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
const MAX_PENDING_HOST_TUNNELS: usize = 256;
#[derive(Default)]
pub struct WasiHostTunnelIo {
receive_capacities: Mutex<HashMap<HostOperationId, usize>>,
}
impl WasiHostTunnelIo {
fn forget_operation(&self, operation: HostOperationId) {
self.receive_capacities.lock().unwrap().remove(&operation);
}
}
impl HostSocketIo for WasiHostTunnelIo {
fn cancel_operation(&self, operation: HostOperationId) -> io::Result<()> {
self.forget_operation(operation);
status("cancel_operation", unsafe { cancel_operation(operation.0) })
}
fn close(&self, handle: HostSocketHandle) -> io::Result<()> {
status("close", unsafe { close(handle.0) })
}
}
impl HostTunnelIo for WasiHostTunnelIo {
fn submit_receive(
&self,
handle: HostSocketHandle,
operation: HostOperationId,
capacity: usize,
) -> io::Result<()> {
let capacity = u32::try_from(capacity).map_err(|_| {
io::Error::new(
io::ErrorKind::InvalidInput,
"host tunnel receive buffer is too large",
)
})?;
status("start_tunnel_receive", unsafe {
start_tunnel_receive(handle.0, operation.0, capacity)
})?;
self.receive_capacities
.lock()
.unwrap()
.insert(operation, capacity as usize);
Ok(())
}
fn take_receive(&self, operation: HostOperationId) -> Poll<io::Result<Vec<u8>>> {
let mut capacities = self.receive_capacities.lock().unwrap();
let Some(&capacity) = capacities.get(&operation) else {
return Poll::Ready(Err(io::Error::new(
io::ErrorKind::NotFound,
"WASI host tunnel receive capacity is missing",
)));
};
let result = unsafe { take_tunnel_receive(operation.0, 0, 0) };
match result {
HOST_PENDING => Poll::Pending,
HOST_TUNNEL_CLOSED => {
capacities.remove(&operation);
Poll::Ready(Err(io::Error::new(
io::ErrorKind::UnexpectedEof,
"host tunnel closed",
)))
}
length if length >= 0 => {
let length = length as usize;
if length > capacity {
capacities.remove(&operation);
let _ = unsafe { cancel_operation(operation.0) };
return Poll::Ready(Err(io::Error::new(
io::ErrorKind::InvalidData,
"host tunnel payload exceeds submitted capacity",
)));
}
let mut buffer = vec![0; length];
let copied = unsafe {
take_tunnel_receive(
operation.0,
buffer.as_mut_ptr() as u32,
buffer.len() as u32,
)
};
capacities.remove(&operation);
if copied == length as i32 {
Poll::Ready(Ok(buffer))
} else {
let _ = unsafe { cancel_operation(operation.0) };
Poll::Ready(Err(host_error("take_tunnel_receive copy", copied)))
}
}
code => {
capacities.remove(&operation);
Poll::Ready(Err(host_error("take_tunnel_receive", code)))
}
}
}
fn submit_send(
&self,
handle: HostSocketHandle,
operation: HostOperationId,
source: &[u8],
) -> io::Result<()> {
let length = u32::try_from(source.len()).map_err(|_| {
io::Error::new(
io::ErrorKind::InvalidInput,
"host tunnel send buffer is too large",
)
})?;
status("start_tunnel_send", unsafe {
start_tunnel_send(handle.0, operation.0, source.as_ptr() as u32, length)
})
}
fn take_send(&self, operation: HostOperationId) -> Poll<io::Result<()>> {
match unsafe { take_tunnel_send(operation.0) } {
HOST_PENDING => Poll::Pending,
0 => Poll::Ready(Ok(())),
code => Poll::Ready(Err(host_error("take_tunnel_send", code))),
}
}
}
#[cfg(feature = "wasm-host-tunnel-outbound")]
pub(crate) struct WasiHostTunnelConnector {
runtime: crate::host::socket::HostSocketRuntime,
io: Arc<WasiHostTunnelIo>,
supported_schemes: Arc<[String]>,
}
#[cfg(feature = "wasm-host-tunnel-outbound")]
impl WasiHostTunnelConnector {
pub(crate) fn new(
runtime: crate::host::socket::HostSocketRuntime,
io: Arc<WasiHostTunnelIo>,
supported_schemes: Arc<[String]>,
) -> Self {
Self {
runtime,
io,
supported_schemes,
}
}
}
#[cfg(feature = "wasm-host-tunnel-outbound")]
#[async_trait::async_trait]
impl ExternalTunnelConnector for WasiHostTunnelConnector {
fn supports_scheme(&self, scheme: &str) -> bool {
self.supported_schemes
.iter()
.any(|supported| supported == scheme)
}
async fn connect(&self, url: &url::Url) -> anyhow::Result<Box<dyn Tunnel>> {
let encoded = url.as_str().as_bytes();
let encoded_len = u32::try_from(encoded.len())
.map_err(|_| anyhow::anyhow!("host tunnel URL exceeds WASI guest memory"))?;
let handle = self
.runtime
.run_operation(
self.io.clone(),
|_, operation| {
status("start_tunnel_connect", unsafe {
start_tunnel_connect(operation.0, encoded.as_ptr() as u32, encoded_len)
})
},
|_, operation| match unsafe { take_tunnel_connect(operation.0) } {
value if value == i64::from(HOST_PENDING) => Poll::Pending,
value if value > 0 => Poll::Ready(Ok(HostSocketHandle(value as u64))),
value => Poll::Ready(Err(host_error(
"take_tunnel_connect",
i32::try_from(value).unwrap_or(i32::MIN),
))),
},
|io, operation| io.cancel_operation(operation),
)
.await?;
let local_url = url::Url::parse(&format!("{}://0.0.0.0:0", url.scheme()))?;
Ok(crate::tunnel::host_tunnel::new_host_tunnel(
self.runtime.clone(),
self.io.clone(),
handle,
local_url,
url.clone(),
None,
))
}
}
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
type HostTunnelQueue = HostListenerQueue<Box<dyn Tunnel>>;
/// Owns the Host Tunnel I/O Adapter and its listener admission queue.
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
pub(crate) struct WasiHostTunnelIngress {
runtime: crate::host::socket::HostSocketRuntime,
io: Arc<WasiHostTunnelIo>,
queue: Arc<HostTunnelQueue>,
supported_schemes: Arc<[String]>,
}
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
impl WasiHostTunnelIngress {
pub(crate) fn new(
runtime: crate::host::socket::HostSocketRuntime,
supported_schemes: Arc<[String]>,
) -> Self {
Self::with_io(
runtime,
Arc::new(WasiHostTunnelIo::default()),
supported_schemes,
)
}
pub(crate) fn with_io(
runtime: crate::host::socket::HostSocketRuntime,
io: Arc<WasiHostTunnelIo>,
supported_schemes: Arc<[String]>,
) -> Self {
Self {
runtime,
io,
queue: Arc::new(HostListenerQueue::new(MAX_PENDING_HOST_TUNNELS)),
supported_schemes,
}
}
pub(crate) fn listener_factory<TcpSocket>(
&self,
) -> Arc<dyn ExternalListenerFactory<AcceptedTransport<TcpSocket>>>
where
TcpSocket: VirtualTcpSocket,
{
Arc::new(WasiHostTunnelListenerFactory {
queue: self.queue.clone(),
supported_schemes: self.supported_schemes.clone(),
})
}
pub(crate) fn accept(
&self,
handle: HostSocketHandle,
metadata: crate::wasi::schema::WasiHostTunnelMetadata,
) -> anyhow::Result<()> {
let crate::wasi::schema::WasiHostTunnelMetadata {
local_url,
remote_url,
resolved_remote_url,
..
} = metadata;
self.queue.enqueue_with(|| {
crate::tunnel::host_tunnel::new_host_tunnel(
self.runtime.clone(),
self.io.clone(),
handle,
local_url,
remote_url,
resolved_remote_url,
)
})
}
}
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
struct WasiHostTunnelListenerFactory {
queue: Arc<HostTunnelQueue>,
supported_schemes: Arc<[String]>,
}
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
impl<TcpSocket> ExternalListenerFactory<AcceptedTransport<TcpSocket>>
for WasiHostTunnelListenerFactory
where
TcpSocket: VirtualTcpSocket,
{
fn supports_scheme(&self, scheme: &str) -> bool {
self.supported_schemes
.iter()
.any(|supported| supported == scheme)
}
fn create(
&self,
request: ExternalListenerRequest,
) -> Box<dyn SocketListener<Accepted = AcceptedTransport<TcpSocket>>> {
Box::new(WasiHostTunnelListener {
registered: self.queue.register_listener(),
queue: self.queue.clone(),
local_url: request.url,
tcp_socket: PhantomData,
})
}
}
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
struct WasiHostTunnelListener<TcpSocket> {
registered: bool,
queue: Arc<HostTunnelQueue>,
local_url: url::Url,
tcp_socket: PhantomData<fn() -> TcpSocket>,
}
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
impl<TcpSocket> fmt::Debug for WasiHostTunnelListener<TcpSocket> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("WasiHostTunnelListener")
.field("local_url", &self.local_url)
.field("registered", &self.registered)
.finish()
}
}
#[async_trait::async_trait]
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
impl<TcpSocket> SocketListener for WasiHostTunnelListener<TcpSocket>
where
TcpSocket: VirtualTcpSocket,
{
type Accepted = AcceptedTransport<TcpSocket>;
async fn listen(&mut self) -> anyhow::Result<()> {
if !self.registered {
anyhow::bail!("Host Tunnel listener queue is closed");
}
Ok(())
}
async fn accept(&mut self) -> anyhow::Result<Self::Accepted> {
let tunnel = self
.queue
.accept()
.await
.ok_or_else(|| anyhow::anyhow!("Host Tunnel listener queue is closed"))?;
Ok(AcceptedTransport::Tunnel {
tunnel,
local_url: self.local_url.clone(),
})
}
fn local_url(&self) -> url::Url {
self.local_url.clone()
}
}
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
impl<TcpSocket> Drop for WasiHostTunnelListener<TcpSocket> {
fn drop(&mut self) {
if self.registered {
self.queue.unregister_listener();
}
}
}
+29
View File
@@ -6,6 +6,8 @@
pub(crate) const HOST_PENDING: i32 = -1;
pub(crate) const HOST_WOULD_BLOCK: i32 = -5;
#[cfg(feature = "wasm-host-tunnel")]
pub(crate) const HOST_TUNNEL_CLOSED: i32 = -10;
#[link(wasm_import_module = "easytier_host")]
unsafe extern "C" {
@@ -188,6 +190,33 @@ unsafe extern "C" {
/// Reports packet-sink write readiness; it never accepts a packet itself.
pub(crate) fn take_packet_write_ready(operation: u64) -> i32;
#[cfg(feature = "wasm-host-tunnel")]
/// Starts receiving one complete host tunnel payload.
pub(crate) fn start_tunnel_receive(handle: u64, operation: u64, capacity: u32) -> i32;
#[cfg(feature = "wasm-host-tunnel")]
/// Probes or copies one payload, or returns the closed sentinel.
///
/// A null destination with zero capacity returns the payload length without
/// consuming it. A second call copies and consumes the payload.
pub(crate) fn take_tunnel_receive(operation: u64, destination: u32, capacity: u32) -> i32;
#[cfg(feature = "wasm-host-tunnel")]
/// Starts sending one complete host tunnel payload.
pub(crate) fn start_tunnel_send(handle: u64, operation: u64, source: u32, length: u32) -> i32;
#[cfg(feature = "wasm-host-tunnel")]
/// Reports completion of one host tunnel send.
pub(crate) fn take_tunnel_send(operation: u64) -> i32;
#[cfg(feature = "wasm-host-tunnel-outbound")]
/// Starts opening one outbound host tunnel for the requested URL.
pub(crate) fn start_tunnel_connect(operation: u64, url: u32, url_len: u32) -> i32;
#[cfg(feature = "wasm-host-tunnel-outbound")]
/// Returns the connected tunnel handle, or a negative host status.
pub(crate) fn take_tunnel_connect(operation: u64) -> i64;
/// Cancels a pending or completed-but-unread operation and releases host state.
///
/// Cancellation must be idempotent when the operation is already absent.
+164 -3
View File
@@ -45,6 +45,11 @@ pub(super) type WasiCore = crate::instance::CoreInstance<
pub(super) struct WasiCoreRuntime {
socket_runtime: crate::host::socket::HostSocketRuntime,
#[cfg(all(
feature = "wasm-host-tunnel",
not(feature = "wasm-host-tunnel-outbound")
))]
tunnel_ingress: crate::wasi::adapter::tunnel::WasiHostTunnelIngress,
core: std::sync::Arc<WasiCore>,
}
@@ -56,6 +61,27 @@ impl WasiCoreRuntime {
pub(super) fn notify_host_completions(&self) {
self.socket_runtime.notify_completions();
}
#[cfg(all(
feature = "wasm-host-tunnel",
not(feature = "wasm-host-tunnel-outbound")
))]
pub(super) fn accept_tunnel(
&self,
handle: crate::host::socket::HostSocketHandle,
metadata: crate::wasi::schema::WasiHostTunnelMetadata,
) -> anyhow::Result<()> {
self.tunnel_ingress.accept(handle, metadata)
}
#[cfg(feature = "wasm-host-tunnel-outbound")]
pub(super) fn accept_tunnel(
&self,
_handle: crate::host::socket::HostSocketHandle,
_metadata: crate::wasi::schema::WasiHostTunnelMetadata,
) -> anyhow::Result<()> {
anyhow::bail!("Host Tunnel admission is unavailable in outbound-only WASI")
}
}
pub(super) fn new_wasi_core_runtime(
@@ -69,7 +95,6 @@ pub(super) fn new_wasi_core_runtime(
use crate::host::{dns::HostDnsResolver, packet::HostPacketSink, socket::HostSocketRuntime};
use crate::{
connectivity::connector_host::new_connector_host,
instance::{CoreHostAdapters, CoreInstance},
wasi::adapter::{
dns::WasiHostDnsIo, environment::WasiHostConnectorEnvironmentIo,
@@ -78,12 +103,36 @@ pub(super) fn new_wasi_core_runtime(
},
};
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
use crate::connectivity::connector_host::new_connector_host;
#[cfg(feature = "wasm-host-tunnel-outbound")]
use crate::connectivity::connector_host::new_connector_host_with_external_tunnel;
let socket_runtime = HostSocketRuntime::new();
let socket_backend = Arc::new(WasiHostSocketBackend::default());
let environment_io = Arc::new(WasiHostConnectorEnvironmentIo);
#[cfg(feature = "wasm-host-tunnel")]
let host_tunnel_schemes: Arc<[String]> = Arc::from(["ws".to_owned(), "wss".to_owned()]);
#[cfg(feature = "wasm-host-tunnel-outbound")]
let tunnel_io = Arc::new(crate::wasi::adapter::tunnel::WasiHostTunnelIo::default());
#[cfg(feature = "wasm-host-tunnel-outbound")]
let host = Arc::new(new_connector_host_with_external_tunnel(
socket_runtime.clone(),
socket_backend,
environment_snapshot,
environment_io,
Arc::new(crate::wasi::adapter::tunnel::WasiHostTunnelConnector::new(
socket_runtime.clone(),
tunnel_io,
host_tunnel_schemes.clone(),
)),
));
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
let host = Arc::new(new_connector_host(
socket_runtime.clone(),
Arc::new(WasiHostSocketBackend::default()),
socket_backend,
environment_snapshot,
Arc::new(WasiHostConnectorEnvironmentIo),
environment_io,
));
let dns = Arc::new(HostDnsResolver::new(
socket_runtime.clone(),
@@ -97,10 +146,45 @@ pub(super) fn new_wasi_core_runtime(
let mut adapters = CoreHostAdapters::new(host, dns, packet_sink, process_runtime);
adapters.instance_runtime = Arc::new(WasiInstanceRuntimeHost);
adapters.events = Arc::new(WasiHostEventSink::new(event_sink));
#[cfg(feature = "wasm-host-tunnel-outbound")]
{
adapters.config.connectivity = crate::instance::CoreConnectivityMode::OutboundOnly;
adapters.config.smoltcp_available = true;
adapters.config.requires_smoltcp = true;
adapters.config.gateway_enabled = false;
adapters.config.proxy_enabled = false;
adapters.config.ignore_unsupported_config = true;
adapters.config.endpoint_protocols = host_tunnel_schemes.to_vec();
}
#[cfg(feature = "wasm-host-tunnel")]
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
let tunnel_ingress = crate::wasi::adapter::tunnel::WasiHostTunnelIngress::new(
socket_runtime.clone(),
host_tunnel_schemes,
);
#[cfg(feature = "wasm-host-tunnel")]
#[cfg(not(feature = "wasm-host-tunnel-outbound"))]
{
adapters.config.connectivity = crate::instance::CoreConnectivityMode::InboundOnly;
adapters.external_listener_factory = Some(tunnel_ingress.listener_factory());
adapters
.host_listener_registrations
.push(crate::listener::ExternalListenerRequest {
url: "wss://0.0.0.0:443"
.parse()
.expect("Host Tunnel listener URL must be valid"),
socket_context: crate::socket::SocketContext::default(),
});
}
let core = CoreInstance::from_toml(config, adapters)?;
Ok(WasiCoreRuntime {
socket_runtime,
#[cfg(all(
feature = "wasm-host-tunnel",
not(feature = "wasm-host-tunnel-outbound")
))]
tunnel_ingress,
core,
})
}
@@ -128,6 +212,8 @@ mod abi {
use super::{WasiCoreRuntime, new_wasi_core_runtime};
use crate::wasi::schema::WasiCoreInstanceCreateConfig;
#[cfg(feature = "wasm-host-tunnel")]
use crate::wasi::schema::WasiHostTunnelMetadata;
#[cfg(feature = "proxy-smoltcp-stack")]
mod data_plane;
@@ -137,6 +223,8 @@ mod abi {
mod web_client;
const MAX_CREATE_CONFIG_LEN: usize = 16 * 1024 * 1024;
#[cfg(feature = "wasm-host-tunnel")]
const MAX_HOST_TUNNEL_METADATA_LEN: usize = 16 * 1024;
const MAX_GUEST_BUFFER_LEN: usize = MAX_CREATE_CONFIG_LEN;
#[cfg(feature = "management-rpc")]
const MAX_RPC_MESSAGE_LEN: usize = 16 * 1024 * 1024;
@@ -480,6 +568,18 @@ mod abi {
}
});
}
#[cfg(feature = "wasm-host-tunnel")]
fn accept_tunnel(
&self,
tunnel_handle: crate::host::socket::HostSocketHandle,
metadata: WasiHostTunnelMetadata,
) -> anyhow::Result<()> {
if self.core.core().state() != CoreInstanceState::Running {
anyhow::bail!("core instance is not running");
}
self.core.accept_tunnel(tunnel_handle, metadata)
}
}
fn decode_create_config(encoded: &[u8]) -> anyhow::Result<WasiCoreInstanceCreateConfig> {
@@ -583,6 +683,10 @@ mod abi {
with_abi_state(|state| state.read_buffer(pointer, length))
}
#[unsafe(no_mangle)]
/// WASI command entrypoint used only to initialize the host runtime.
pub extern "C" fn _start() {}
#[unsafe(no_mangle)]
/// Allocates a guest-owned ABI buffer and returns its linear-memory offset.
///
@@ -795,6 +899,63 @@ mod abi {
})
}
#[cfg(feature = "wasm-host-tunnel")]
#[unsafe(no_mangle)]
/// Returns the host tunnel ABI version implemented by this guest.
pub extern "C" fn easytier_host_tunnel_abi_version() -> u32 {
crate::wasi::abi::HOST_TUNNEL_ABI_VERSION
}
#[cfg(feature = "wasm-host-tunnel")]
#[unsafe(no_mangle)]
/// Transfers one host-owned transport into server tunnel admission.
///
/// `metadata` is a versioned JSON document. A zero return transfers
/// ownership of `tunnel_handle` to the guest; on failure the host keeps
/// ownership and must close it. Admission is scheduled and completed by
/// later drive calls so the export never waits for the peer handshake.
pub extern "C" fn easytier_instance_accept_tunnel(
handle: u64,
tunnel_handle: u64,
metadata_pointer: u32,
metadata_length: u32,
) -> i32 {
if tunnel_handle == 0 {
set_instance_error(handle, "host tunnel handle must be non-zero");
return INVALID_INPUT;
}
let encoded = match read_guest_buffer(
metadata_pointer,
metadata_length,
MAX_HOST_TUNNEL_METADATA_LEN,
) {
Ok(encoded) => encoded,
Err(error) => {
set_instance_error(handle, error);
return INVALID_INPUT;
}
};
let metadata: WasiHostTunnelMetadata = match serde_json::from_slice(&encoded) {
Ok(metadata) => metadata,
Err(error) => {
set_instance_error(handle, error);
return INVALID_INPUT;
}
};
if let Err(error) = metadata.validate() {
set_instance_error(handle, error);
return INVALID_INPUT;
}
with_instance(handle, |instance| {
instance.accept_tunnel(
crate::host::socket::HostSocketHandle(tunnel_handle),
metadata,
)?;
Ok(0)
})
}
#[unsafe(no_mangle)]
/// Destroys an instance and releases its lifecycle, timer, and runtime state.
pub extern "C" fn easytier_instance_drop(handle: u64) -> i32 {
@@ -233,7 +233,7 @@ fn require_ipv4(address: SocketAddr) -> Result<SocketAddr, DataPlaneError> {
address.is_ipv4().then_some(address).ok_or_else(|| {
error(
DataPlaneErrorKind::AddressFamilyUnsupported,
"data-plane ABI v2 supports IPv4 only",
"data-plane ABI supports IPv4 only",
)
})
}
@@ -355,6 +355,24 @@ pub extern "C" fn easytier_data_plane_tcp_write_submit(
})
}
#[unsafe(no_mangle)]
pub extern "C" fn easytier_data_plane_tcp_shutdown_write_submit(
handle: u64,
stream: u64,
output_operation: u32,
) -> i32 {
let stream = match resource_id(stream) {
Ok(stream) => stream,
Err(error) => {
set_instance_error(handle, error.message());
return error_status(error.kind());
}
};
submit_operation(handle, output_operation, |instance| {
instance.submit_data_plane(|session| session.submit_tcp_shutdown_write(stream))
})
}
#[unsafe(no_mangle)]
pub extern "C" fn easytier_data_plane_udp_bind_submit(
handle: u64,
@@ -660,6 +678,26 @@ pub extern "C" fn easytier_data_plane_tcp_write_result_take(handle: u64, operati
})
}
#[unsafe(no_mangle)]
pub extern "C" fn easytier_data_plane_tcp_shutdown_write_result_take(
handle: u64,
operation: u64,
) -> i32 {
data_plane_call(handle, |instance| {
let session = instance.data_plane_session();
take_result(
&session,
operation_id(operation)?,
DataPlaneOperationKind::TcpShutdownWrite,
|result| match result {
DataPlaneOperationResult::TcpWriteShutdown => Ok(()),
_ => Err(invalid_input("TCP write shutdown result variant mismatch")),
},
)?;
Ok(0)
})
}
#[unsafe(no_mangle)]
pub extern "C" fn easytier_data_plane_udp_bind_result_take(
handle: u64,
+5 -1
View File
@@ -12,6 +12,10 @@ use std::{
use tokio::runtime::Runtime;
// A zero-duration timer can win before the executor's park hook observes
// quiescence, especially when JSPI resumes the guest on a slower embedder.
const DRIVE_BUDGET: Duration = Duration::from_millis(10);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum RuntimeDriveOutcome {
Quiescent,
@@ -52,7 +56,7 @@ impl RuntimeDriver {
let _active = RuntimeDriverGuard::activate(self.state.as_ref());
runtime.block_on(async {
let budget = tokio::time::sleep(Duration::ZERO);
let budget = tokio::time::sleep(DRIVE_BUDGET);
tokio::pin!(budget);
poll_fn(|context| {
if self.state.poll_quiescent(context.waker()) {
+37
View File
@@ -60,3 +60,40 @@ impl WasiWebClientCreateConfig {
Ok(())
}
}
#[cfg(feature = "wasm-host-tunnel")]
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct WasiHostTunnelMetadata {
pub version: u32,
pub local_url: url::Url,
pub remote_url: url::Url,
#[serde(default)]
pub resolved_remote_url: Option<url::Url>,
}
#[cfg(feature = "wasm-host-tunnel")]
impl WasiHostTunnelMetadata {
pub(crate) fn validate(&self) -> anyhow::Result<()> {
if self.version != crate::wasi::abi::HOST_TUNNEL_ABI_VERSION {
anyhow::bail!("unsupported host tunnel ABI version: {}", self.version);
}
Ok(())
}
}
#[cfg(all(test, feature = "wasm-host-tunnel"))]
mod tests {
use super::*;
#[test]
fn host_tunnel_metadata_accepts_transport_neutral_urls() {
let metadata = WasiHostTunnelMetadata {
version: crate::wasi::abi::HOST_TUNNEL_ABI_VERSION,
local_url: "test-tunnel://listener.example/".parse().unwrap(),
remote_url: "test-tunnel://peer.example/".parse().unwrap(),
resolved_remote_url: None,
};
metadata.validate().unwrap();
}
}
+1 -1
View File
@@ -22,7 +22,7 @@ pub(crate) fn decode_ipv4_socket_address(wire: &[u8]) -> io::Result<SocketAddr>
if !address.is_ipv4() {
return Err(io::Error::new(
io::ErrorKind::Unsupported,
"data-plane ABI v2 supports IPv4 only",
"data-plane ABI supports IPv4 only",
));
}
Ok(address)