Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions .github/workflows/rust.yml
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
name: Rust

# Manual only. Nothing here runs on a push or a pull request: the gate is the
# one an author runs locally before asking for review, and a green tick that
# nobody asked for teaches people to stop reading it.
on:
push:
branches: [ "master" ]
pull_request:
branches: [ "master" ]
workflow_dispatch:

env:
CARGO_TERM_COLOR: always
Expand Down
7 changes: 7 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,13 @@ natively, and Claude Code loads it through the `@AGENTS.md` import in

## Invariants (don't break these)

- **No Python.** Not a script, not `python3 -c`, not a heredoc. Reaching for it is the
tell that a step is being solved by parsing when the tool that owns the answer could
just be asked. Do not swap it for another parser either, and do not assume `jq` is
present: it does not ship with macOS. A fixed-shape field is one `sed -nE` line;
anything needing real parsing belongs in this repo's own language, where it can be
tested. If a task seems to need Python, the approach is wrong.

- **Verify the whole chain, not just this repo.** endpoint-gen, honey_id-types,
endpoint-validator and the six backends all break silently when this crate moves.
`./scripts/check-chain.sh` verifies all of it; run it before calling a change done.
Expand Down
18 changes: 17 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "endpoint-libs"
version = "3.0.3"
version = "3.1.0"
edition = "2024"
authors = ["Veon <dveonloch@protonmail.com>"]
description = "Launch MCP services fast: describe endpoints once in RON, and endpoint-gen generates the Rust models, docs and MCP tool schemas that this crate serves over WebSocket RPC, with roles and typed errors built in."
Expand Down Expand Up @@ -67,6 +67,13 @@ framed-transport = [
"wire-core",
"dep:tokio-util",
]
nagoya-transport = [
# `framed_json_neutral` over a Nagoya socket. Off by default and additive: a
# consumer that does not run Nagoya never compiles it, and the tokio path is
# untouched either way.
"framed-transport",
"dep:nagoya",
]
agent-control = ["framed-transport"]
ws-client = [
# WS client (WsClient, WsClientBuilder) - standalone
Expand Down Expand Up @@ -155,6 +162,11 @@ tokio-rustls = { version = "0.26", optional = true, default-features = false, fe
rustls = { version = "0.23", optional = true, default-features = false, features = ["ring", "logging", "std"] }
tokio = { version = "1.39", features = ["full"] }
tokio-util = { version = "0.7", features = ["codec"], optional = true }
# The I/O driver is behind Nagoya's own non-default `reactor` feature; without it
# `nagoya::reactor` does not exist and the error reads like a missing module.
# 0.1.9 is the published release that already speaks AF_UNIX and carries `TaskSet`,
# so this resolves from the registry with no path dependency and no sibling checkout.
nagoya = { version = "0.1.9", default-features = false, features = ["reactor"], optional = true }
tokio-cron-scheduler = { version = "0.11", optional = true }
parking_lot = { version = "0.12", optional = true }
dashmap = { version = "6.0", optional = true }
Expand Down Expand Up @@ -191,6 +203,10 @@ tonic = { version = "0.14", optional = true }
[dev-dependencies]
tempfile = "3.19"
tokio = { version = "1.39", features = ["full", "test-util"] }
# `compat` only, and only for tests: it adapts a tokio duplex into the
# `futures-io` pair the neutral transport takes, which is what lets one test
# drive both framing paths and compare the bytes they produce.
tokio-util = { version = "0.7", features = ["codec", "compat"] }
tracing-throttle = { version = "0.4", features = ["async", "test-helpers"] }
uuid = { version = "1", features = ["v4", "serde"] }
rcgen = "0.14"
Expand Down
182 changes: 128 additions & 54 deletions src/libs/ws/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -171,31 +171,58 @@ impl WebsocketServer {
.upgrade_stream(stream, addr, &self.config, &cached_date)
.await?;

// Loop: spawn session task for each upgrade event
while let Ok(event) = rx.recv().await {
use futures::StreamExt;
use futures::future::FutureExt;
use futures::stream::FuturesUnordered;

// Poll sessions in place, the same one-thread model `serve_with` uses.
// The hyper upgrader still `spawn_local`s (TokioExecutor); that is why
// `run_shard` keeps a LocalSet. nago-wss is the replacement, not this loop.
let mut sessions = FuturesUnordered::new();
loop {
let event = if sessions.is_empty() {
match rx.recv().await {
Ok(event) => event,
Err(_) => break,
}
} else {
let recv = rx.recv();
futures::pin_mut!(recv);
match futures::future::select(recv, sessions.next()).await {
futures::future::Either::Left((Ok(event), _)) => event,
futures::future::Either::Left((Err(_), _)) => {
while sessions.next().await.is_some() {}
break;
}
futures::future::Either::Right(_) => continue,
}
};

let this = Arc::clone(&self);
let states = Arc::clone(&states);
let addr_clone = addr;
sessions.push(
async move {
let ws_stream = match create_ws_stream(event.on_upgrade).await {
Ok(s) => s,
Err(e) => {
error!(ws_server = true, ?addr_clone, "on_upgrade failed: {e}");
return;
}
};

tokio::task::spawn_local(async move {
let ws_stream = match create_ws_stream(event.on_upgrade).await {
Ok(s) => s,
Err(e) => {
error!(ws_server = true, ?addr_clone, "on_upgrade failed: {e}");
return;
}
};

debug!(
ws_server = true,
?addr_clone,
protocol = %event.protocol,
"WsServer: upgrade succeeded, protocol received"
);
debug!(
ws_server = true,
?addr_clone,
protocol = %event.protocol,
"WsServer: upgrade succeeded, protocol received"
);

this.post_upgrade_connection(addr_clone, states, ws_stream, event.protocol)
.await;
});
this.post_upgrade_connection(addr_clone, states, ws_stream, event.protocol)
.await;
}
.boxed_local(),
);
}

debug!(
Expand Down Expand Up @@ -230,8 +257,9 @@ impl WebsocketServer {
/// time — the WebSocket subprotocol string today, a handed-over token for local
/// transports. It is passed to [`AuthController::auth`] unchanged.
///
/// Must be called inside a `tokio::task::LocalSet`: [`MessageStream`]'s futures
/// are not `Send`.
/// [`MessageStream`]'s futures are not `Send`, so this must be polled on
/// the thread that owns the stream. [`Self::serve_with`] does that by
/// driving connections with `FuturesUnordered` rather than `spawn_local`.
pub async fn serve_connection(
self: Arc<Self>,
peer: PeerIdentity,
Expand Down Expand Up @@ -351,12 +379,19 @@ impl WebsocketServer {
/// runs on a single runtime — the shard-per-core model is a property of the TCP
/// path and buys nothing for a 1:1 sidecar channel.
///
/// Must be called inside a `tokio::task::LocalSet` (see
/// [`Self::serve_connection`]).
/// Connections are polled in place with `FuturesUnordered`. That is the
/// same one-thread model `spawn_local` had, without tying the method to a
/// tokio `LocalSet`, so a nagoya reactor can drive it. TCP `listen` now
/// polls its connections the same way; it still needs a `LocalSet` for the
/// hyper upgrader.
pub async fn serve_with<L>(self, listener: L) -> Result<()>
where
L: SessionListener + 'static,
{
use futures::StreamExt;
use futures::future::FutureExt;
use futures::stream::FuturesUnordered;

self.validate_protocol_mode()?;
let this = Arc::new(self);
let states = Arc::new(WebsocketStates::new());
Expand All @@ -366,8 +401,19 @@ impl WebsocketServer {
this.config.drop_conn_on_buffer_full,
);

let mut connections = FuturesUnordered::new();
loop {
let (stream, peer) = match listener.accept().await {
let accepted = if connections.is_empty() {
listener.accept().await
} else {
let accept = listener.accept();
futures::pin_mut!(accept);
match futures::future::select(accept, connections.next()).await {
futures::future::Either::Left((accepted, _)) => accepted,
futures::future::Either::Right(_) => continue,
}
};
let (stream, peer) = match accepted {
Ok(accepted) => accepted,
Err(err) => {
error!(ws_server = true, error = %err, "listener accept failed; stopping");
Expand All @@ -378,11 +424,14 @@ impl WebsocketServer {

let this = Arc::clone(&this);
let states = Arc::clone(&states);
tokio::task::spawn_local(async move {
// Local transports carry credentials out of band (an inherited fd is
// already a capability), so there is no subprotocol string to pass.
this.serve_connection(peer, states, stream, None).await;
});
connections.push(
async move {
// Local transports carry credentials out of band (an inherited fd is
// already a capability), so there is no subprotocol string to pass.
this.serve_connection(peer, states, stream, None).await;
}
.boxed_local(),
);
}
}

Expand Down Expand Up @@ -495,43 +544,68 @@ impl WebsocketServer {
.enable_all()
.build()
.expect("Failed to build shard runtime");
// LocalSet remains only because the hyper upgrader `spawn_local`s onto
// TokioExecutor. Connection tasks themselves are polled in place.
let local_set = LocalSet::new();
rt.block_on(local_set.run_until(async move {
use futures::StreamExt;
use futures::future::FutureExt;
use futures::stream::FuturesUnordered;

let mut connections = FuturesUnordered::new();
loop {
let Some((stream, addr)) = rx.recv().await else {
let received = if connections.is_empty() {
rx.recv().await
} else {
let recv = rx.recv();
futures::pin_mut!(recv);
match futures::future::select(recv, connections.next()).await {
futures::future::Either::Left((received, _)) => received,
futures::future::Either::Right(_) => continue,
}
};
let Some((stream, addr)) = received else {
while connections.next().await.is_some() {}
break;
};
let this = Arc::clone(&this);
let states = Arc::clone(&states);
let listener = Arc::clone(&listener);
tokio::task::spawn_local(async move {
let stream = match listener.handshake(stream).await {
Ok(channel) => {
debug!(ws_server = true, "Accepted stream from {}", addr);
channel
}
Err(err) => {
connections.push(
async move {
let stream = match listener.handshake(stream).await {
Ok(channel) => {
debug!(ws_server = true, "Accepted stream from {}", addr);
channel
}
Err(err) => {
error!(
ws_server = true,
"Error while handshaking stream: {:?}", err
);
return;
}
};
if let Err(err) = TOOLBOX
.scope(
this.toolbox.clone(),
this.handle_ws_handshake_and_connection(
addr,
states,
Box::new(stream),
),
)
.await
{
error!(
ws_server = true,
"Error while handshaking stream: {:?}", err
?addr,
"Failed to handle WS connection: {err}"
);
return;
}
};
if let Err(err) = TOOLBOX
.scope(
this.toolbox.clone(),
this.handle_ws_handshake_and_connection(addr, states, Box::new(stream)),
)
.await
{
error!(
ws_server = true,
?addr,
"Failed to handle WS connection: {err}"
);
}
});
.boxed_local(),
);
}
}));
}
Expand Down
Loading