Files
fips/docs/how-to/serve-many-peers-on-one-thread.md
Johnathan Corgan 3a789370b9 Add an experimental native datagram API addressed by public key
A client process opens a flow to a peer's public key on a chosen port and
sends and receives datagrams on a file descriptor the daemon hands it. No
IPv6 emulation, no TUN device, no DNS: a datagram travels from key to key.
The feature is off by default and is not a stable interface.

The wire needs no change and gets none. Every FSP data packet has carried a
port pair inside its AEAD envelope since v0.2.0, and port 256 is simply the
IPv6 shim. What was missing was a way for a program to ask for a port of its
own and be handed the traffic.

Addressing is the part worth reading twice, because the obvious design is
wrong. The x-only public key is the address. An npub is that key written in
bech32, so converting between them is a local encoding rather than a lookup
or a name service. The 16-byte node address that travels on the wire is the
first half of a SHA-256 of the key: it is a truncated hash, it does not
invert, and it appears nowhere a client can see. An earlier iteration of this
work reported a peer by that hash and could supply a key only sometimes,
which is what treating a wire identifier as an identity produces.

An accepted flow therefore always knows its peer. The key is captured where
the peer is authenticated rather than looked up when a report is rendered:
every inbound datagram passes one call site inside a handler that refuses
anything whose session is not established, and the responder has already
rejected the session unless the claimed address derives from the key it
proved. Reaching for the identity cache instead gives a best-effort answer
from a structure that evicts.

A listener is a descriptor. The daemon writes one message per arrival to it,
carrying the new flow's descriptor and the peer's address, so poll, select
and epoll work on a listener and accepting is a recvmsg. That is what lets
the API be used from a program that already has an event loop, which a
command-and-reply listener could not support: an arrival could not be waited
on beside anything else. There is no accept command and no reject command.
Refusing a flow is closing the descriptor you were handed.

The Rust surface mirrors std::net. FipsStream::connect, FipsListener::bind,
incoming, accept, io::Result and an errno mapping rather than a bespoke
error type. An address is given as an npub, as a key, or as a pair, through
one parameter, the way ToSocketAddrs takes several spellings of one thing.
Each type holds its descriptor and copies of what setup told it and nothing
else, so a stream that outlives its setup connection is not representable.

set_nonblocking, AsFd and the four deadline methods carry the names and
signatures std::net uses for the same jobs. They were asked for by a user
integrating the API with tokio: AsyncFd requires a non-blocking descriptor,
and anything receiving from a peer needs a bounded wait. AsFd is the better
of the two descriptor accessors, because the borrow cannot outlive the value
that owns the descriptor, so a reactor cannot hold a registration for a
descriptor that has since been closed and its number reused by the next
open. The non-blocking flag is read, modified and written back rather than
assigned, since the flag word carries more than that one bit and a caller may
have set O_ASYNC. A zero timeout is refused with EINVAL, because the kernel
reads a zero timeval as "wait for ever", which inverts what a caller passing
zero means; std::net refuses it for the same reason. The two directions are
separate options and stay that way. FipsListener gets no timeout methods,
matching TcpListener: bounding an accept is set_nonblocking plus the caller's
own poll, which the reactor how-to builds. A flow taken from accept is
blocking whatever the listener was set to, because the two are separate
sockets and the daemon hands over a fresh one.

One rule has no counterpart in Berkeley sockets and a client author must know
it: the v1 wire carries no half-close, so nothing peer-driven ever closes a
flow. A server written to read until the flow ends waits for a signal that
cannot arrive, holding a thread and a flow per peer until its process exits.
A program decides its own termination, and the example serves one datagram
per flow.

The tests reach a live daemon rather than a stand-in. Every public item had a
unit test against a hand-written stand-in with canned replies, and the five
entry points a program actually calls first, connect, connect_from,
connect_at, bind and the SOCKET constant, had no coverage of any kind,
because the tests that appear to cover them build a Wire over a socket pair
and hand it to the private open and hold, so nothing ever resolved a socket
path or mapped its errors. examples/native-surface.rs walks all thirty-eight
items against a running daemon and reports the number of assertions it made.
The count is read from the recorder rather than written as a literal, and the
harness asserts the exit status, the completion marker and the count
together, so deleting an assertion fails the check rather than quietly
shrinking it. Watchdogs turn a hang into a named failure, which several of
the walked behaviours would otherwise produce. The shared Docker image is
built once for every integration leg, so the new binary is staged at all ten
places the existing one is, the interop builder included, which gets a stub
because those images exercise the wire between daemon versions and older refs
do not carry the example. The platform gating was tested rather than reasoned
about: flipping all eleven gates so the native API is excluded leaves the
crate compiling clean across the workspace, every target and the profiling
feature.

The shipped docs tree gains what only the LaTeX manual under design/ had,
which is not published with the daemon. A reference entry covers the whole
surface: addressing and the port tiers, the Berkeley mapping, every method on
FipsAddr, FipsStream, FipsListener and Incoming, the errno table, the
ceilings, the four places data disappears with nothing reported, the line
protocol and the command reference. The errno table gives names rather than
numbers, since the client maps each name onto the libc constant for the
platform it was built for and the supported platforms disagree on the
numbers. A tutorial side trip stands up two throwaway nodes on one machine,
peered over loopback UDP with no TUN and no DNS, then writes a listening
program and a connecting program against them; it needs neither the public
mesh nor root, because the native path is the one that does not go through
the IPv6 adapter. The obligations a client in another language carries are a
how-to of their own, since they are a task rather than a description:
reading the setup connection with recvmsg, associating a descriptor with the
last complete line, telling an empty datagram from a close, and six others.
Serving many peers from one poll loop is another, with the whole program,
because the straightforward listener spawns a thread per flow and that is
wrong at the node's ceiling of 256. The drop causes are a table mapping each
of the seven texts DropReason::as_str produces to the counter it increments,
with drop_oversize called out as the ninth counter that is not in the table.
What a daemon restart costs is a section of its own: every flow and listener
ends, descriptors do not survive, there is no resumption, and datagrams sent
but not yet forwarded are lost through a window nothing bounds.

A stack comparison diagram places the interface against the stack a reader
already knows: the same application over HTTP, TLS, TCP, IP and Ethernet on
one side, and over its own format, FSP, FMP and a FIPS transport on the
other, aligned so each row is one concern. The two columns are not
alternatives and are not drawn as such. An unmodified IPv6 program's packets
reach fips0, and the adapter hands each one to FSP as a payload, so the left
stack runs inside the right one; the left column ends at a fork, eth0 for the
ordinary internet and fips0 for the mesh, and an arrow leaves fips0 and runs
back up into FSP's input. The row where TCP would be is empty on purpose and
names Reliable Object Delivery, which is where that capability is expected to
land. ROD is a v2 capability, the box is dashed because none of it exists
yet, and the design entry says the part a reader needs most: nothing on the
surface anticipates it, so a program written today should assume it does not
exist. Both endpoints carry a scheme and a worked port,
https://<npub>.fips:443 and fips://<npub>:443, with a footnote saying the two
ports are not the same kind of thing, a TCP port inside the tunnel on the
left and an FSP port on the right. The fips:// form is a coinage: nothing in
the tree parses it, nothing registers the scheme, and the API takes a key and
a port as separate arguments rather than a URL. The diagram also says where
the right column stops, since FIPS over UDP still rides IP and Ethernet
beneath. It appears in fips-concepts.md and fips-ipv6-adapter.md, which were
making its argument in prose without a picture, and deliberately not in
fips-architecture.md, which already carries the OSI mapping and makes the
same point about the transport row.

The gateway's control socket moves onto the same bind policy this API uses,
which is the one change here that touches deployed behaviour: fips-gateway
now tightens /run/fips to 0750. That is unreachable under the packaged
deployment, where fips.service has already created the directory at that
mode, and reachable for a source build or a container that starts the gateway
alone.

One changelog entry under Added, describing the released state: what a
client opens and reads, the addressing and why the node address is not it,
the listener being a descriptor, the std::net shape of the Rust surface,
and the one rule Berkeley sockets have no counterpart for. It says in as
many words that the wire is unchanged.
2026-08-21 05:48:23 +00:00

7.9 KiB

Serve Many Peers on One Thread

Goal: handle every native datagram API flow from a single poll loop, instead of dedicating a thread to each peer.

The straightforward listening program spawns a thread per flow. That is fine for a handful of peers and wrong at the node's ceiling of 256, where it costs 256 threads mostly parked in recv.

Every object on this surface is a descriptor, so there is nothing to integrate: both types implement AsFd and AsRawFd and go straight into a poll, select or epoll set. A listener is readable exactly when accept would not block, which is the property the whole shape rests on, and the crate asserts it as a test rather than claiming it.

Read use-the-native-datagram-api.md first if you have not opened a flow before.

The whole program

One dependency beyond the crate, libc, for poll itself:

[dependencies]
fips = { git = "https://github.com/jmcorgan/fips" }
libc = "0.2"
//! Serve many flows from one poll loop, with no thread per peer.
//!
//! ```text
//! eventloop /run/fips/api.sock 4242
//! ```

use fips::native::client::{FipsListener, FipsStream};
use std::env;
use std::error::Error;
use std::os::fd::{AsRawFd, RawFd};
use std::path::Path;
use std::process::ExitCode;
use std::time::{Duration, Instant};

/// How long a flow may go without a datagram before this program closes it.
///
/// Nothing else will end one. The wire carries no far-end close, so a peer that
/// has stopped sending is indistinguishable from one that is thinking, and a
/// reactor with no deadline holds every flow it ever accepted until it exits.
const IDLE: Duration = Duration::from_secs(30);

/// Hold the port named on the command line and serve every flow from one loop.
fn run() -> Result<(), Box<dyn Error>> {
    let mut args = env::args().skip(1);
    let (Some(socket), Some(port)) = (args.next(), args.next()) else {
        return Err("usage: eventloop <socket-path> <local-port>".into());
    };
    let port: u16 = port.parse()?;

    let listener = FipsListener::bind_at(Path::new(&socket), port)?;
    println!("holding {}", listener.local_addr());

    // Each flow with the time its last datagram arrived, which is what the
    // deadline is measured against.
    let mut flows: Vec<(FipsStream, Instant)> = Vec::new();
    loop {
        let mut fds = vec![watch(listener.as_raw_fd())];
        fds.extend(flows.iter().map(|(flow, _)| watch(flow.as_raw_fd())));

        let count = fds.len() as libc::nfds_t;
        // The timeout is what makes the deadline reachable: with no events at
        // all the loop must still wake to notice a flow that has gone quiet.
        let timeout = IDLE.as_millis() as libc::c_int;
        // SAFETY: `fds` is a live slice of `pollfd` for the whole call, and
        // `count` is its length.
        if unsafe { libc::poll(fds.as_mut_ptr(), count, timeout) } < 0 {
            return Err(std::io::Error::last_os_error().into());
        }

        // The established flows first, and backwards: removing a closed one
        // must not renumber one not yet examined, and accepting below would
        // otherwise push a flow this pass has no `revents` for. The listener
        // occupies slot 0, hence the offset.
        let now = Instant::now();
        for index in (0..flows.len()).rev() {
            if ready(&fds[index + 1]) {
                if echo(&flows[index].0) {
                    flows[index].1 = now;
                } else {
                    flows.remove(index); // dropping it releases the flow
                }
            } else if now.duration_since(flows[index].1) >= IDLE {
                println!("closing an idle flow from {}", flows[index].0.peer_addr());
                flows.remove(index);
            }
        }

        // Once per readiness rather than in a loop: the descriptor is blocking,
        // so a second `accept` with nothing queued would stall the whole loop.
        if ready(&fds[0]) {
            let (flow, peer) = listener.accept()?;
            println!("flow from {peer}");
            flows.push((flow, now));
        }
    }
}

/// A `pollfd` asking for readability on `fd`.
///
/// `POLLHUP` needs no asking for: it is reported in `revents` whether or not
/// it was requested, which is what lets one mask serve both cases.
fn watch(fd: RawFd) -> libc::pollfd {
    libc::pollfd {
        fd,
        events: libc::POLLIN,
        revents: 0,
    }
}

/// Whether this descriptor has something to read or has hung up.
fn ready(poll: &libc::pollfd) -> bool {
    poll.revents & (libc::POLLIN | libc::POLLHUP) != 0
}

/// Return one datagram, reporting whether the flow is still usable.
fn echo(flow: &FipsStream) -> bool {
    let mut buf = vec![0u8; flow.max_payload()];
    match flow.recv(&mut buf) {
        // `Ok(0)` is an empty datagram and not a close, so it is echoed like
        // any other. `EPIPE` is the daemon gone: no peer can close a flow.
        Ok(len) => flow.send(&buf[..len]).is_ok(),
        Err(_) => false,
    }
}

/// Report a failure on stderr and exit non-zero.
fn main() -> ExitCode {
    match run() {
        Ok(()) => ExitCode::SUCCESS,
        Err(error) => {
            eprintln!("eventloop: {error}");
            ExitCode::FAILURE
        }
    }
}

The four rules

Register the listener for readability only. There is nothing else to ask it for, and POLLHUP arrives in revents whether or not it was requested.

Accept once per readiness, not in a loop. The descriptor is blocking, so a second accept with nothing queued stalls the whole loop. Looping until WouldBlock is correct only after listener.set_nonblocking(true), which is also what an edge-triggered epoll requires.

Give every flow a deadline, and the poll a timeout that makes the deadline reachable. This is the most important of the four. Nothing will tell a reactor that a peer is finished, so a flow that goes quiet stays in the poll set for ever unless the program removes it — and with no events at all the loop must still wake in order to notice. A reactor with a deadline but no poll timeout has a deadline it can never reach.

A flow from accept arrives blocking, whatever the listener was set to. They are separate sockets and the daemon hands over a fresh one. The program above deliberately leaves them blocking and makes exactly one recv per readiness, which is safe on a blocking descriptor and is why it needs no flags at all. If you want them otherwise, call set_nonblocking on the flow.

Where this reaches past the client module

This program needs libc and an unsafe block, and it is the only one of the API's example programs that needs anything.

That is a narrower gap than it once was. Readiness and a bounded wait were both missing from the surface; set_nonblocking closed the first and set_read_timeout the second, and a program wanting an option on a flow now has a method for it. What is left is a different kind of thing: an option on a flow is something a surface can supply, and a reactor's own polling mechanism is not. No surface that stops at the descriptor can supply poll.

AsRawFd rather than AsFd here is deliberate: libc::poll takes a raw descriptor. Prefer AsFd anywhere you hold a registration, because its borrow cannot outlive the stream; this loop rebuilds its pollfd set from live references on every pass, so it holds nothing across an iteration.

See also