The stream registry represents multiple virtual application streams inside one established MixTransport session. It defines TransportStream in libp2p_mix_transport/streams.nim and makes each TransportSession responsible for registering and removing its streams.
TransportStream inherits from libp2p’s BufferStream, and therefore from Connection. The session registry, frame router, application protocol handler and buffered read/write path all refer to the same stream object; there is no separate application wrapper.
Stream Identity Is Scoped by the Session
A streamId is unique within one transport session. The complete identity of a stream is therefore the pair:
(sessionId, streamId)Two different sessions may both contain stream 1 without ambiguity because their sessionId values differ.
The endpoint opening a stream selects its streamId. It will send that ID in an OpenStream frame, and the remote endpoint will register an inbound stream with the same ID. One logical stream never receives separate local and remote identifiers:
session initiator session recipient
outbound stream 1
|
| OpenStream(streamId = 1)
v
inbound stream 1Preventing Concurrent Allocation Collisions
Both endpoints may open streams within the same established session. If both started independent counters at 1, simultaneous opens could assign 1 to two different streams before either OpenStream frame arrived.
The registry avoids coordination by partitioning the positive uint64 identifier space:
- The endpoint that initiated the transport session allocates odd identifiers:
1,3,5, and so on. - The endpoint that received the transport session allocates even identifiers:
2,4,6, and so on.
This partition determines only which endpoint may originate an identifier. Both endpoints still use the same ID for the resulting stream. For example, stream 2 is outbound at the session recipient and inbound at the session initiator.
The IDs are sequential because they need to be unique only within one session. Small sequential values also have a compact Protobuf varint representation. The available space is large enough that exhaustion is not a practical concern during one session’s lifetime.
Outbound Stream Registration
addOutboundStream(codec) may be called only on an established session and requires a non-empty application codec. It takes the session’s next locally assigned identifier, creates a pending outbound TransportStream, stores it by streamId, and advances the local counter by two.
let stream = newTransportStream(
session.sessionId, session.peerId, session.nextOutboundStreamId, codec,
StreamDirection.Outbound,
)
session.streams[stream.streamId] = stream
session.nextOutboundStreamId += 2The stream remains pending because local registration alone does not show that the remote endpoint accepted the application codec. MixTransport.dial sends OpenStream and waits for this stream to become established after StreamAck returns.
Inbound Stream Registration
addInboundStream(streamId, codec) also requires an established session and a non-empty codec. Before changing the registry, it verifies two independent conditions:
- The identifier belongs to the remote endpoint’s odd or even allocation space. Identifier zero is never valid.
- No stream with that exact identifier is already registered in the session.
The first check rejects a remote frame that attempts to use an identifier reserved for locally opened streams. The second check detects duplicate OpenStream frames and prevents them from replacing the existing stream.
The role-based validation is deliberately local: it answers whether an identifier received from the other endpoint belongs to that endpoint’s allocation space.
func isValidInboundStreamId(session: TransportSession, streamId: uint64): bool =
if streamId == 0:
return false
case session.role
of SessionRole.Initiator:
streamId mod 2 == 0
of SessionRole.Recipient:
streamId mod 2 == 1An initiator allocates odd outbound IDs itself, so it accepts even IDs in incoming OpenStream frames. A recipient allocates even outbound IDs itself, so it accepts odd incoming IDs. The modulo test does not claim that every identifier with the expected parity already exists; addInboundStream creates the stream only after also checking the session state, codec and duplicate table entry.
After validation, the session creates a pending inbound stream with the supplied ID. The accepting endpoint configures and establishes that stream before sending StreamAck. For a forward OpenStream, the response uses its dedicated SURBs; for a reverse OpenStream, it uses the forward path. If response submission fails completely, the handler removes and shuts down the stream. Before registration, acceptInboundStreamOpening records the attempt in a bounded history, preventing a duplicate or late opening from recreating a removed stream. Mix Transport Implementation Walk Through - Recipient-Originated Streams explains that history and the symmetric opening paths.
if not session.isValidInboundStreamId(streamId):
return err("stream identifier was not allocated by the remote endpoint")
if session.streams.hasKey(streamId):
return err("stream identifier is already registered")
let stream = newTransportStream(
session.sessionId, session.peerId, streamId, codec, StreamDirection.Inbound
)
session.streams[streamId] = stream
ok(stream)Immediately after this registration step, before the handler configures and establishes it, the stream is still Pending. Registration lets subsequent control code refer to the same object, while the pending state prevents application use before the remote opener has been sent an acknowledgement.
Direction and State
StreamDirection is local to an endpoint:
opener: Outbound
remote endpoint: InboundIt is independent of SessionRole. A session initiator can receive an inbound stream opened later by the session recipient, and a session recipient can hold an outbound stream that it opened itself.
Every stream starts in Pending. establish changes the state to Established, while reject changes it to Rejected. Both operations fire the stream’s resolution AsyncEvent. On the accepting endpoint, the OpenStream handler establishes its inbound stream before submitting StreamAck. On the opener, the reply handler establishes the pending outbound stream after StreamAck, or rejects it after StreamReject. dial waits on the resolution event instead of polling and then returns either the established stream or an error.
Construction Boundary
newTransportStream is exported from the internal streams.nim module because sessions.nim is a separate Nim module and must be able to call it. The package facade deliberately excludes that constructor from its exports. Normal consumers create registered streams through addOutboundStream and addInboundStream, which enforce the session and identifier rules described above.
Stream and Session Lifetimes
The session owns its stream table. Removing one stream deletes only that stream’s entry and leaves the established transport session registered. This preserves the session pseudonym and allows later streams to reuse the same anonymous peer relationship.
Removing the session also shuts down its streams and releases session-owned state. Stream closure wakes waiters and requests cancellation of stream-owned tasks; asynchronous shutdown awaits their completion. CloseStream, ResetStream, Disconnect, and ResetSession coordinate remote teardown on a best-effort basis, as described in Mix Transport Implementation Walk Through - Remote Teardown.