net — the network as a stream
This page specifies how a Flux application reaches the network: the five verbs that describe a
connection, the arrival envelope every message wears, the codec and backpressure fields of a
NetSpec, the protocol presets, and the seam that promotes a live feed into causal analysis. The
capability model it obeys — resource handles, default-deny, consent — is owned by
host-services; the other side of the wire, a server you run, is
server.
New here? Start with Guide §10 — Building an application →
The network pillar is sealed in design, rolling out after the v1 core; its capability face
(net:fetch, net:stream) is part of the APP-plane contract.
A network connection is a stream on the arrival axis. That is not a metaphor — it is the same statement the language already makes about a chart’s x axis, applied to a different clock. And it is what lets one word, stream, cover request/response, server push, polling, a stream-of-streams, and pagination: they differ only in their arrival regime.
The script never opens a socket. It describes a connection; the host — the only holder of the file descriptor, the TLS session and the token — opens it, and delivers what arrives as journaled messages.
The placement theorem
Where may a network stream live? The answer is forced, not chosen:
Arrival is non-deterministic with respect to the world — the network arrives when it arrives. So it lives entirely in the APP plane, where non-determinism enters as a journaled message and contaminates nothing. Reading the network from analysis is
[ErrFirewall].Promotion into analysis happens by exactly one seam: a causal
scanthat folds arrivals into closed bars, followed bytoSource, which the host ingests append-only — never revising a closed bar. Only then does the ordinary resample operator apply.
A pane can therefore read an exchange feed, a chat gateway, an RSS river. An indicator
cannot — it only ever sees the causal series that toSource produced. The firewall is not
weakened by the network; the network is routed around it.
The five verbs
A connection is described by a spec and consumed through the two frozen doors of the APP plane —
commands out, subscriptions in. No verb ever carries a file descriptor, a token, or a forged
URL: the url is an allowlisted string key the host resolves, at a consented domain.
| Verb | Kind | What it is |
|---|---|---|
request(spec, On) |
Cmd |
one-shot request/response; the result re-enters as a message through the constructor it carries |
subscribe(spec, On) |
Sub |
an inbound push stream (server-sent events, a socket subscription, a polled feed) |
connect(spec, On) |
Sub |
a persistent bidirectional channel; the companion sink(connKey).send(v) emits an outbound command |
datagrams(spec, On) |
Sub |
unordered, unreliable datagrams |
paginate(spec, On, next, maxPages) |
Sub |
unrolls a bounded loop, yielding a stream of pages |
A spec is a NetSpec record. You rarely build one field by field: a preset returns one, and
with overrides the fields you care about.
variant Msg { Got(f: Recv) | Ping }
app feed {
capabilities: [ net:stream ]
init(p) = { last: na, missed: 0 }
update(m, msg) = match msg {
Got(f) -> match f {
Data(t) -> { model: m with { last: t.price }, cmds: [] }
Dropped(n) -> { model: m with { missed: m.missed + n }, cmds: [] }
_ -> { model: m, cmds: [] }
}
Ping -> { model: m, cmds: [ sink("exchange.ws").send(Ping) ] }
}
view(m) = text("last {fmt.price(m.last)} · missed {m.missed}")
subs(m) = [ connect(ws("exchange.ws") with { codec: Json, schema: Trade }, Got) ]
}Note what the _ arm is doing: the envelope has seven arms, and match is exhaustive, so the
compiler will not let you forget that a socket can drop, lag, or be revoked. Note also
Dropped(n) — the model counts what it never saw.
The arrival envelope
Everything that arrives is wrapped in one variant, declared once, and eliminated by match.
It is shown here in spec notation — the <κ> is the metalanguage for “the schema you
declared”, not surface syntax: a v1 record is monomorphic, and the host supplies Recv already
specialised to your schema.
variant Recv<κ> {
Data(κ) // a decoded unit — κ is the schema you declared
| DecodeError(field: string, reason: string) // a broken required field — never a silent `na`
| Dropped(count: num) // backpressure evicted this many
| Lagged // the consumer fell behind
| NetErr(class: NetErrClass, reason: string) // a TYPED failure — branch on the class
| SchemaMismatch(expected: string, got: string)// the provider's version left the accepted window
| Revoked(capRef: string) // the capability was revoked mid-session
}
variant NetErrClass {
RateLimited(retryAfter: duration | na) // pace — never hammer
| Unauthorized // re-authenticate; do not blind-retry
| Http(status: num) | Timeout | Dns | Tls | Cors | Refused | Closed(code: num | na)
} // the WS close code drives resume vs re-identifyWhy the error class is a variant and not a string. A reconnection policy that branches on the wording of a host error message is not portable, not testable, and — because the wording is not byte-stable across engines — not replayable. The class is host-authoritative and byte-identical, so
RateLimited(retryAfter)→ pace,Unauthorized→ re-auth,Http(404)→ give up,Timeout→ retry, is a total function you can write once and trust.
A one-shot request never delivers Dropped or Lagged (there is no backpressure on a single
response). A stream can deliver all seven arms — and match makes you handle them.
The lifecycle is declarative
A subscription that is present means “this connection is wanted”; removing it means “close it”.
The runtime diffs subscriptions by their connKey, exactly as it diffs everything else.
Each field of a spec carries a class:
| Class | Fields | Effect of a change |
|---|---|---|
| Hot | backpressure, heartbeat, timeout, subscription filters | applied without reopening the connection |
| Reopen | url, protocol, auth, codec, framing | reconnects |
Changing the connKey itself is an explicit reconnection. A send on a key whose subscription
is gone is a no-op that surfaces as a message — never a crash, never a write to a closed socket.
Codecs
A codec is a projection of a kind — not a parser. You declare the shape you expect, and the host decodes into it and kind-checks at the boundary:
record Trade { price: price(BTC, USD) ; qty: volume ; ts: time }
trades = subscribe(ws("exchange.ws") with { codec: Json, schema: Trade }, Got)codec and schema are two fields, not one. The codec says how the bytes are shaped; the
schema says what kind they decode into. Leave schema off and it is inferred from the codec —
but naming it is what lets the boundary check hold.
This is what makes the ban on regular expressions and ad-hoc parsing satisfiable rather than
merely restrictive: you never parse a payload in the script, because the payload arrives already
typed. A required field that is broken gives you DecodeError(field, reason) — never a silent
na that poisons a computation three hops later.
The decoder is incremental and arena-backed: it fills a bounded buffer as bytes arrive,
zero-copy where the layout permits, with no allocation in the steady state. Framing (how a byte
stream is cut into messages) is a separate, composable layer, so a codec can be reused across
transports. The catalogue is closed — extensible only as a vetted list, never as a runtime
grammar. In v1 it is Json, XmlFeed (which is what the RSS/Atom preset decodes with), Utf8
and Raw, plus the compact binaries Cbor and MsgPack, the tabular Csv and Tsv, and
Url. Md follows the text pillar, with the protocols that carry it.
Each codec carries an accepted version window. A provider that drifts outside it yields
SchemaMismatch(expected, got) at connect, before a single datum is delivered — an
“update required”, rather than a corrupted model.
Backpressure is declared, never implicit
A fast producer and a slow consumer is not an edge case; it is Tuesday. So the policy is part of the spec, and every option is total:
BackPressure is a closed variant of five arms, and it is the back: field of the spec:
| Policy | Behaviour |
|---|---|
Latest(n) |
a bounded ring of the n most recent arrivals — the default for a price feed |
DropOldest(n) |
a bounded queue of n; when it is full, evict the oldest |
DropNewest(n) |
a bounded queue of n; when it is full, refuse the newcomer |
Sample(clk) |
keep-last at each tick of a bar clock — coalescing, for a feed you only sample |
Block |
true backpressure, toward the host: the host stops reading the socket |
A hot stream that declares no back takes Latest. There is no unbounded queue, so there is no
way to write the classic memory leak where an app quietly buffers a firehose until it dies. And
because eviction is counted — Dropped(n) is an arm of the envelope — the model always knows
what it did not see.
The one that is easy to get wrong:
Block. A fold from ticks into OHLCV bars must see every tick, or the bar it builds is not the bar that happened.Samplewould silently lose ticks between clock edges, andLatest(n)would lose them under a burst. So thebars(tf)fold that feeds the seam into analysis declaresback: Block— and the cost of that choice is explicit and local: the host stops reading the socket, rather than the script quietly aggregating a lie.
Rate-shaping is a different axis and lives elsewhere: throttle / debounce / sample /
dedup / merge are combinators of the stream algebra, applied to a stream you already have;
and pace on the spec token-buckets the pipeline’s own outbound requests. Backpressure is what
happens when arrivals outrun the consumer — not how fast you choose to ask.
Pagination follows the same discipline: a bounded loop with a declared maximum page count, or an incremental pull driven by the model — never a free loop over an unknown number of pages.
Asynchrony without callbacks
There is no await, no promise, no callback in the language. A pipeline is described, the
host drives it, and every result re-enters as a message through the constructor the command
carried. That is the whole story — and it is the reason an application’s entire behaviour is
reconstructible from its journal.
The cost is honest and worth stating: a long asynchronous chain becomes several message
constructors and several arms of update, rather than three lines of await. What you buy is
that every intermediate step is in the journal — so time travel, replay and server-side
verification all work, which they cannot if the intermediates are hidden inside a task.
Protocols
The protocol catalogue decomposes into five orthogonal axes (transport, framing, encoding,
session, delivery), which is why one spec type covers everything rather than one API per
protocol. A preset is a plain function returning a NetSpec with the sane defaults for its
protocol already set — a reconnection backoff, a heartbeat, a backpressure ring — and with
overrides whatever you need. That collapses the common case to one line:
quotes = request(rest("api.example.com") with { codec: Json, schema: Quote }, GotQuote)
events = subscribe(sse("events.example.com") with { codec: Json, schema: Event }, GotEvent)
ticks = connect(ws("stream.example.com") with { codec: Cbor, schema: Tick, back: Latest(256) }, GotTick)
posts = subscribe(rss("blog.example.com/feed", 300s), GotPost)The presets are rest · ws · sse · rss · wt. rss takes its poll interval as an argument,
because a poll without a declared cadence is not a poll; the others read theirs from the transport.
Browser-possible: HTTP, server-sent events, WebSocket, WebTransport, and the polling feeds. Deferred to a relay: protocols that need a raw TCP socket — the browser cannot open one, and pretending otherwise would be a lie in the documentation rather than a feature.
Capabilities, auth and governance
The whole pillar sits behind the net:* family, default-deny — the net: face of the capability
model host-services specifies in full:
net:fetch(domain)andnet:stream(domain)are granted per domain, with user consent, and enforced by the host’s content-security policy. The script never holds the socket.- Authentication is host-held. The token never enters the script — the app receives an opaque handle, and the host attaches the credential. A leaked script cannot leak a credential it never had.
- Egress is governable. A supervisor can narrow (never widen) what an application may reach; a revocation is journaled as a bound, so a re-fold reproduces it deterministically and a command issued after it fails closed.
- Offline is a first-class state. A per-grant cache policy, a bounded offline command queue
replayed under an idempotency key on reconnection, and
OnConnectivityedges — never a silent multi-writer merge, which is a named non-goal.
The seam into analysis
// APP plane: every tick (`back: Block`), folded into closed bars, handed to the host
feed = connect(ws("exchange.ws") with { codec: Cbor, schema: Trade, back: Block }, Got)
src = feed.bars(tf("1m")).toSource("BTC-ext") // the stream is the receiver — UFCS, no pipe operator
// ANALYSIS plane: it is now an ordinary, causal series
plot ema(series("BTC-ext").close, 20) @ tf("1h")The host ingests append-only: a closed bar is never revised. So a series fed by an exchange socket carries exactly the same no-repaint guarantee as one fed by the first-party pipeline — not because the network is trusted, but because the seam refuses to let it rewrite history.
Figure — ticks arrive as journaled APP-plane messages, a causal scan folds them into closed bars, and toSource crosses into analysis append-only: the one seam by which the network becomes causal data.
See also
- App plane — commands, subscriptions, the journal, capabilities.
- Host services — the resource-handle doctrine this pillar follows.
- server — the other side of the wire, when it is ours.
- text — why there is no regular expression, and what replaces it.
- Host integration —
data:sourceand the causal ingestion contract. - The four planes — the firewall this pillar routes around.