Wymcp. Transport. Stream
(Wymcp v0.6.2)
View Source
Legacy-only
The SSE stream serves one session, and sessions exist only because the legacy era does. This module is removed wholesale at the legacy decommission.
The chunked SSE connection for one session, run as a receive loop by the GET request process that owns its socket; one active stream per session — a new GET replaces the old.
Design decisions
Every chunk write happens in the request process, because real adapters
enforce socket ownership: Bandit raises when any process other than the
request's own calls Plug.Conn.chunk/2 (Plug.Test has no such check,
which is how a separate writer process ever appeared to work). The stream
is still a separate process from Wymcp.Session — the session must keep
serving POSTs while the stream blocks — but that process is the GET
request process itself, so the conn is never handed over.
serve/3 runs the stream's whole lifetime: register with the session,
commit the 200, send the priming event, then serve pushes, keepalives,
and replacement from the receive loop until the stream ends — the client
disconnects (a chunk write fails), the session dies (monitored), or a new
GET replaces this one. The final conn travels back up the plug stack like
any other response. A registration failure answers {:error, :session_gone} before anything is committed, so the router still sends a
clean 404; a priming write that fails ends the stream right there, before
the loop and before a keepalive timer is armed.
A disconnect is the only ending the session cannot see for itself. Ending
the stream is not ending this process: the request process belongs to the
adapter, which may keep it alive for the connection's next request, so no
:DOWN reaches Wymcp.Session's stream monitor and a registration left
behind would make every later push wait out its caller's own call timeout
before answering. Every failed write therefore goes through one funnel that
tells the session (Wymcp.Session.unregister_stream/2, a cast — the
loop must never wait on the session)
and then closes. Replacement and session death need no such message: the
session already holds the successor's pid in the first case and is gone
in the second.
The session reaches the loop only through this module's public API.
push/3 is a stream-answered push: the session hands the pre-encoded
payload and a reply address to the loop and never blocks on the write —
the loop answers the address with the push ack (answer_push/2): the
pusher directly on a plain push, the session on a server-request round
trip's push leg, nobody for a list-changed notification. replace/1 is a
plain send: asking, never killing, because the loop's process is serving
a committed 200. The loop ends with a catch-all clause that drops unknown
messages at debug level — it must never crash on a stray, and strays
exist by construction (Plug.Test's adapter notifies the conn owner;
ThousandIsland's read timer can fire mid-callback).
sequenceDiagram
participant Client
participant Router as Router (GET /)
participant Stream as Stream loop (request process)
participant Session
Client->>Router: GET / (Mcp-Session-Id)
Router->>Session: lookup + touch
Router->>Stream: serve(conn, session_pid)
Stream->>Session: register_stream(self())
Note over Session: asks a replaced old stream to stop
alt registration fails (session died or timed out)
Stream-->>Router: {:error, :session_gone}
Router-->>Client: 404 (pre-commit)
else registered
Stream-->>Client: 200 chunked, priming event, keepalives
Session->>Stream: push(json, reply_to) / replace
Stream-->>Client: SSE events
Note over Stream: answers reply_to with the push ack
Stream-->>Router: final conn (stream ended)
endEvent IDs
Each SSE event gets a monotonically increasing integer ID. Clients use
Last-Event-ID on reconnection to indicate the last event they received.
Full replay is out of scope — no missed events are re-sent; the counter
resumes after the client's last seen event (Last-Event-ID: evt-7 makes
the priming event evt-8).
The header is raw client input: an id that does not read as
evt-<non-negative integer> is not an error — the resumption point is
discarded with an :info log naming the raw value, and the counter
resumes from 0. evt-0 is a legitimate resumption point, not a discard.
flowchart TD
subgraph Stream
ST[Transport.Stream] --> L["receive loop (push / replace / keepalive)"]
ST --> PR["priming event"]
ST --> K["keepalive timer"]
end
subgraph External
ST -->|"register_stream / unregister_stream, monitors"| S[Session]
L --> SSE[Transport.SSE]
end
Summary
Functions
Hands one pre-encoded push to the stream loop — a stream-answered push:
the caller never blocks here; the loop answers reply_to with the push
ack after the chunk write.
Asks the stream loop to end so a new GET can take over the session.
Runs the SSE stream for session_pid in the calling process — the GET
request process — and blocks until the stream ends.
Types
@type reply_to() :: {:caller, GenServer.from()} | {:ack, pid(), term()} | :none
Who awaits one push — the address the loop answers with the push ack.
@type serve_opts() :: %{optional(:keepalive_interval) => pos_integer()}
Functions
Hands one pre-encoded push to the stream loop — a stream-answered push:
the caller never blocks here; the loop answers reply_to with the push
ack after the chunk write.
{:caller, from} answers a GenServer.caller directly via
GenServer.reply/2 (Wymcp.Session forwards the pusher's own from, so
a late ack is dropped by the caller's reply alias); {:ack, session_pid, tag} sends {:push_ack, tag, result} to the session — the server-request
round trip's push leg; :none expects no answer (list-changed
notifications). The ack vocabulary: :ok — written; {:error, :disconnected} — the chunk write failed and the stream is closing, or
the push was still queued when the stream closed (the close drain answers
it). A push sent in the gap between registration and loop entry queues in
the loop's mailbox and is answered after the priming event.
Asks the stream loop to end so a new GET can take over the session.
A plain send — never blocks the caller, never kills: the loop's process is a request process serving a committed 200, and it finishes normally by returning its conn up the plug stack.
Runs the SSE stream for session_pid in the calling process — the GET
request process — and blocks until the stream ends.
Returns {:ok, conn} with the final conn once the stream ends (client
disconnect, session death, or replacement by a new GET). Returns
{:error, :session_gone} before anything is committed when registering
with the session fails — the session died between the router's lookup and
registration, or the registration call timed out; the log distinguishes
the two. The resumption point comes from Wymcp.Plugs.SingletonHeaders'
:wymcp_last_event_id assign, not from the header directly: a single
well-formed value resumes the counter, and a missing, malformed, or
duplicated one starts fresh from 0 (the last two log the discard at
:info). Calling this function on a conn the check never touched
raises — the assign is the contract.
opts may carry :keepalive_interval in milliseconds (default 15 000);
only tests override it.