Skip to content

← ROUTER | Proxy →

STREAM Socket

1. Overview

STREAM is a server-only socket for communicating with external raw clients.

Core rules: - Set ZLINK_STREAM_OPT_RECV_MODE to RAW or PACKET before the first zlink_bind() or zlink_connect(). Bind or connect without a mode fails with EINVAL. - ZLINK_SOCKET_STREAM supports both bind and connect (zlink_connect() returns ZLINK_CONNECT_OK on success). The peer can be an OS/Asio/WebSocket raw client or another zlink STREAM socket with a receive mode set. - RAW mode has no zlink-level wire framing — it is a transparent byte stream (the encoder/decoder pass bytes through unchanged). For length-delimited packets, use PACKET mode, which frames as 2-byte BE header size + 4-byte BE body size + header + body. - At the zlink API level, RAW zlink_recv() returns a one-part byte record and exposes the source client's 4-byte routing_id through its own source_rid_out_ out-parameter, and packet zlink_stream_recv_packet() exposes the same Core-owned borrowed view through source_rid_out_.

Valid combination:

external raw client  <---- RAW byte stream (no framing) ---->  STREAM(server)

STREAM is not directly compatible with zlink internal sockets (PAIR/PUB/SUB/DEALER/ROUTER).


2. Server Create/Bind

void *stream = zlink_socket(ctx, ZLINK_SOCKET_STREAM);
zlink_stream_recv_mode_t mode = ZLINK_STREAM_RECV_MODE_RAW;   /* required before bind */
zlink_set_stream_option(stream, ZLINK_STREAM_OPT_RECV_MODE, &mode, sizeof(mode));
int linger = 0;
zlink_set_option(stream, ZLINK_OPT_LINGER, &linger, sizeof(linger));
zlink_bind(stream, "tcp://0.0.0.0:8080");

Supported transports (bind and connect): - tcp:// - tls:// - ws:// - wss://


3. STREAM-Specific Behavior

STREAM is the only exception type in the raw socket family. Exactly one of two receive modes must be selected before the first successful bind.

  • RAW: zlink_recv() pulls one transport byte record. The record has one part, with the source routing id returned through its source_rid_out_ out-parameter. Pair it with a poller watching ZLINK_POLLIN.
  • PACKET: zlink_stream_recv_packet() pulls packets assembled from a fixed framing convention (2B header size + 4B body size + header + body, all big-endian) as header/body pairs.

Set ZLINK_STREAM_OPT_RECV_MODE to ZLINK_STREAM_RECV_MODE_RAW or ZLINK_STREAM_RECV_MODE_PACKET before bind. The mode becomes immutable after the first successful bind; the receive API for the other mode returns ENOTSUP.

STREAM-specific behavior:

  • source_rid is auto-assigned per connection by the server, always fixed 4 bytes (uint32, big-endian).
  • To close one client, pass the source_rid received from recv to zlink_disconnect_rid(). STREAM target routing ids must be 4 bytes.
  • With the default ZLINK_STREAM_OPT_NOTIFY=0, connect/disconnect are not in-band data markers. They are reported through the socket monitor as ZLINK_EVENT_CONNECTION_READY / ZLINK_EVENT_DISCONNECTED, each carrying the 4-byte routing_id. A raw payload that happens to be a single 0x00/0x01 byte is delivered as ordinary data. If ZLINK_STREAM_OPT_NOTIFY is set to 1 before bind/connect in RAW mode, zlink_recv() additionally returns a zero-length record with the affected source_rid for every connect and disconnect, so a zero-length part must then be treated as a notification rather than data.

4. RAW Pull Example

With ZLINK_STREAM_OPT_NOTIFY=0 (the default), every part pulled in STREAM RAW mode is application data; observe connect/disconnect on the socket monitor (see Monitoring). With NOTIFY=1, zero-length records arrive interleaved as connect/disconnect notifications.

zlink_stream_recv_mode_t mode = ZLINK_STREAM_RECV_MODE_RAW;
zlink_set_stream_option(stream, ZLINK_STREAM_OPT_RECV_MODE,
                        &mode, sizeof(mode));
zlink_bind(stream, "tcp://0.0.0.0:8080");

const zlink_routing_id_t *source_rid = NULL;
zlink_msg_t parts[1];
size_t part_count = 0;
if (zlink_recv(stream, &source_rid, parts, 1, &part_count,
               ZLINK_RECV_FLAGS_NONE) == ZLINK_RECV_OK) {
    zlink_msg_t reply;
    zlink_msg_init_size(&reply, zlink_msg_size(&parts[0]));
    memcpy(zlink_msg_data(&reply), zlink_msg_data(&parts[0]), zlink_msg_size(&parts[0]));
    zlink_send_rid(stream, source_rid, &reply, 1, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
    zlink_multipart_close(parts, part_count);
}

Key Points

Item Description
Receive API zlink_recv()
Readiness A poller reports ZLINK_POLLIN; the application then drains receives
Lifetime source_rid remains valid until the same socket's next data-recv entry or close
Framing Raw bytes as received from the transport
Send zlink_send_rid()

When the send queue is full (HWM), zlink_send_rid() blocks (default) or returns ZLINK_SUBMIT_BACKPRESSURED with ZLINK_DONTWAIT. For advanced backpressure patterns, see Performance Guide.

  • The caller owns a successfully received zlink_msg_t and closes it exactly once.
  • Copy the borrowed source_rid before the next data receive when it must be retained.
  • ZLINK_RECV_FLAGS_DONTWAIT returns ZLINK_RECV_NO_DATA with EAGAIN when empty.

4.1 PACKET Pull Mode

When the upstream protocol uses the fixed framing convention (2-byte big-endian header size + 4-byte big-endian body size + header payload + body payload), select PACKET mode and pull with zlink_stream_recv_packet(). Core handles fragment accumulation and length parsing, so the application receives assembled header/body pairs directly.

zlink_stream_recv_mode_t mode = ZLINK_STREAM_RECV_MODE_PACKET;
zlink_set_stream_option(stream, ZLINK_STREAM_OPT_RECV_MODE,
                        &mode, sizeof(mode));
zlink_bind(stream, "tcp://0.0.0.0:8080");

const zlink_routing_id_t *source_rid = NULL;
zlink_msg_t header;
zlink_msg_t body;
zlink_msg_init(&header);
zlink_msg_init(&body);
if (zlink_stream_recv_packet(stream, &source_rid, &header, &body,
                             ZLINK_RECV_FLAGS_NONE) == ZLINK_RECV_OK) {
    /* process header/body; zero-length messages are valid */
    zlink_msg_close(&header);
    zlink_msg_close(&body);
}

Rules for PACKET mode:

  • header_size or body_size equal to zero is allowed; both sides are still delivered as valid zlink_msg_t objects.
  • Ownership of header and body is transferred to the caller. The caller must close or consume each msg_t exactly once.
  • In PACKET mode, whole-message RAW receive (zlink_recv()) fails with ENOTSUP. In RAW mode, zlink_stream_recv_packet() fails the same way.
  • Malformed packets (length exceeding implementation limits, assembly failure, premature close, etc.) result in the connection being closed as the default policy. Observe such events via the socket monitor.

This mode relieves the application from re-implementing fragment accumulation, but it does not change the fact that transport fragment boundaries differ from packet boundaries.


5. Client Implementation Rule

Clients can be raw socket/websocket clients, or another zlink STREAM socket that sets a receive mode and calls zlink_connect().

Conceptual POSIX TCP example (RAW mode — no zlink framing, just bytes):

// RAW mode: send/recv raw bytes; message boundaries are application-defined
send(fd, body, body_len, 0);

char buf[4096];
ssize_t n = recv(fd, buf, sizeof(buf), 0);

If the server side uses PACKET mode (zlink_stream_recv_packet), the client must frame each packet as 2-byte BE header size + 4-byte BE body size + header + body:

// packet mode: [2B header_size BE][4B body_size BE][header][body]
uint16_t hsz_be = htons(header_len);
uint32_t bsz_be = htonl(body_len);
send(fd, &hsz_be, 2, 0);
send(fd, &bsz_be, 4, 0);
send(fd, header, header_len, 0);
send(fd, body, body_len, 0);

6. Option and Runtime Policy

Main supported options: - ZLINK_OPT_MAXMSGSIZE, ZLINK_OPT_SNDHWM, ZLINK_OPT_RCVHWM, ZLINK_OPT_SNDBUF, ZLINK_OPT_RCVBUF, ZLINK_OPT_BACKLOG, ZLINK_OPT_LINGER - ZLINK_STREAM_OPT_RECV_MODE (via zlink_set_stream_option() / zlink_get_stream_option()): select RAW or PACKET before the first bind or connect - ZLINK_STREAM_OPT_NOTIFY: enable zero-length connect/disconnect records in RAW mode; it cannot be combined with PACKET mode - TLS/WSS server: zlink_set_tls_server() / TLS client: zlink_set_tls_client()

STREAM listeners often receive bytes from raw TCP peers. If the peer is not fully trusted, set ZLINK_OPT_MAXMSGSIZE to the largest application message you intend to accept before calling zlink_bind. Without that setting the compatibility default is unlimited.

Unsupported/changed: - Setting ZLINK_ROUTER_OPT_CONNECT_ROUTING_ID on STREAM returns EOPNOTSUPP.

6.1 Default STREAM runtime profile

Defaults a STREAM socket applies to its public options: - ZLINK_OPT_BACKLOG: 65536 - ZLINK_OPT_SNDHWM / ZLINK_OPT_RCVHWM: STREAM profile byte value from the default balanced auto-HWM policy, or the manual byte default if context auto-HWM is disabled - ZLINK_OPT_SNDBUF / ZLINK_OPT_RCVBUF: default -1, leaving OS buffer defaults and TCP autotuning in control

The exact contract of each option is owned by the STREAM socket spec.


7. Errors and Constraints

  • zlink_connect(stream, ...) -> EOPNOTSUPP
  • On STREAM, non-4-byte routing_id frame is a protocol error
  • Messages larger than MAXMSGSIZE are dropped and connection is closed (disconnect event)

8. Reference Tests

  • core/tests/integration/test_stream_socket.cpp
  • core/tests/integration/test_stream_fastpath.cpp
  • core/tests/integration/routing-id/test_connect_rid_string_alias.cpp
  • core/tests/scenario/stream/zlink/test_scenario_stream_zlink.cpp

These tests use STREAM server + raw client paths.


← ROUTER | Proxy → | Transport →

Full language examples

zlink::context_t ctx;
zlink::stream_socket_t server (ctx);
zlink::socket_monitor_t server_monitor = server.monitor_open ();
server.options ().notify (false);
server.options ().recv_mode (zlink::stream_recv_mode_t::raw);

server.bind ("tcp://127.0.0.1:0");
const std::string endpoint = server.options ().last_endpoint ();
assert (!endpoint.empty ());

detail::raw_tcp_client_t client (endpoint);
assert (detail::wait_stream_connected (server_monitor));

const char *request = detail::k_stream_payload;
const size_t request_size = std::strlen (request);
client.send_all (request, request_size);

zlink::received_t inbound;
assert (server.recv (inbound) == 0);
assert (inbound.routing_id ().has_value ());
assert (inbound.parts ().size () == 1);
assert (inbound.parts ()[0].to_string () == detail::k_stream_payload);

zlink::message_t reply = detail::make_message (detail::k_stream_payload);
// Reply on the STREAM socket, addressed by the received envelope's routing id.
server.send (*inbound.routing_id ()).message (reply).submit ();
inbound.close ();

char response[64];
const int received = client.recv_exact (response, request_size);
assert (received == static_cast<int> (std::strlen (detail::k_stream_payload)));
assert (std::memcmp (response, detail::k_stream_payload, received) == 0);
std::printf ("[stream/recv] send: \"%s\" → recv: \"%.*s\"\n", request, received, response);

client.close ();
return 0;
if (!SampleSupport.IsNativeAvailable())
    return;

using var ctx = Zlink.CreateContext();
using var stream = ctx.CreateStreamSocket();
stream.Options.ReceiveMode = StreamReceiveMode.Raw;
stream.Options.Linger = TimeSpan.Zero;
string endpoint = SampleSupport.NewEndpoint("tcp", "sample");
int port = SampleSupport.ExtractPort(endpoint);
using var monitor = stream.MonitorOpen(SocketEvent.Accepted);
stream.Bind(endpoint);

using var client = SampleSupport.ConnectRawClient(port);
SampleSupport.WaitMonitorEvent(monitor, 5000, SocketEvent.Accepted);
NetworkStream network = client.GetStream();
byte[] request = "hello-stream"u8.ToArray();
SampleSupport.SendAll(network, request);

using var received = Received.Create();
if (!stream.Recv(received))
    throw new InvalidOperationException("recv failed");
if (received.RoutingId == null)
    throw new InvalidOperationException("missing routing id");
string payload = received.Parts[0].GetString();
SampleSupport.EnsureEqual("hello-stream", payload, "payload");

using var reply = Message.From("hello-stream");
received.Send().Message(reply).Submit();
string echoed = System.Text.Encoding.UTF8.GetString(
    SampleSupport.ReceiveExact(network, "hello-stream".Length));
Console.WriteLine(
    $"[stream/recv] send: \"hello-stream\" -> recv: \"{echoed}\"");
SampleSupport.ensureNative();
String endpoint = SampleSupport.tcpEndpoint();

try (Context ctx = Zlink.createContext();
     StreamSocket server = ctx.createStreamSocket();
     var monitor = server.monitorOpen(
         systems.zlink.contracts.eventing.MonitorEventType.ACCEPTED,
         systems.zlink.contracts.eventing.MonitorEventType.CONNECTION_READY)) {
    server.options().recvMode(StreamRecvMode.RAW);
    server.bind(endpoint);
    try (var rawClient = SampleSupport.connectRawTcp(endpoint)) {
        SampleSupport.waitStreamConnected(monitor);

        byte[] payload =
            SampleSupport.STREAM_PAYLOAD.getBytes(StandardCharsets.UTF_8);
        SampleSupport.sendRawTcp(rawClient, payload);

        try (systems.zlink.contracts.messaging.Received received = new systems.zlink.contracts.messaging.Received()) {
            server.recv(received, systems.zlink.contracts.sockets.RecvFlags.NONE);
            String value = SampleSupport.singleUtf8(received);
            if (!SampleSupport.STREAM_PAYLOAD.equals(value)
                || received.getRoutingId().isEmpty()) {
                throw new IllegalStateException("unexpected stream delivery");
            }

            try (Message reply = Message.from(
                     SampleSupport.STREAM_PAYLOAD)) {
                received.send().message(reply).submit();
            }
        }

        String echoed = new String(
            SampleSupport.recvExactRawTcp(rawClient, payload.length),
            StandardCharsets.UTF_8);
        if (!SampleSupport.STREAM_PAYLOAD.equals(echoed)) {
            throw new IllegalStateException("unexpected stream echo: " + echoed);
        }
        System.out.println("[stream/recv] send: \"" +
            SampleSupport.STREAM_PAYLOAD + "\" \u2192 recv: \"" +
            echoed + "\"");
    }
}
SampleSupport.ensureNative()
val endpoint = SampleSupport.tcpEndpoint()

Zlink.createContext().use { ctx ->
    ctx.createStreamSocket().use { server ->
        server.monitorOpen(
            MonitorEventType.ACCEPTED,
            MonitorEventType.CONNECTION_READY
        ).use { monitor ->
            server.bind(endpoint)
            SampleSupport.connectRawTcp(endpoint).use { rawClient ->
                SampleSupport.waitStreamConnected(monitor)

                val payload = SampleSupport.STREAM_PAYLOAD.toByteArray(StandardCharsets.UTF_8)
                SampleSupport.sendRawTcp(rawClient, payload)

                Received().use { received ->
                    server.recv(received, RecvFlags.NONE)
                    val value = SampleSupport.singleUtf8(received)
                    check(SampleSupport.STREAM_PAYLOAD == value && received.routingId.isPresent) {
                        "unexpected stream delivery"
                    }
                    Message.from(SampleSupport.STREAM_PAYLOAD).use { reply ->
                        received.send().message(reply).submit()
                    }
                }

                val echoed = String(
                    SampleSupport.recvExactRawTcp(rawClient, payload.size),
                    StandardCharsets.UTF_8
                )
                check(SampleSupport.STREAM_PAYLOAD == echoed) { "unexpected stream echo: $echoed" }
                println(
                    "[stream/recv] send: \"${SampleSupport.STREAM_PAYLOAD}\"" +
                        " → recv: \"$echoed\""
                )
            }
        }
    }
}
port, endpoint = tcp_endpoint()

with zlink.create_context() as ctx:
    with zlink.create_stream_socket(ctx) as server:
        server.stream_options.recv_mode = zlink.StreamRecvMode.RAW
        with server.monitor_open(zlink.MonitorEventMask.ACCEPTED) as server_monitor:
            server.bind(endpoint)
            with socket.create_connection(("127.0.0.1", port), timeout=3.0) as client:
                wait_socket_monitor_event(server_monitor, zlink.MonitorEventMask.ACCEPTED)
                client.sendall(b"hello-stream")
                received = zlink.create_received()
                if not server.recv_into(received):
                    raise AssertionError("expected stream payload")
                with received:
                    if received.to_bytes_list() != [b"hello-stream"]:
                        raise AssertionError("unexpected stream payload")
                    if not received.routing_id:
                        raise AssertionError("stream sample expected a routing id")
                    server.send(received.routing_id).message(b"hello-stream").submit_sync()
                reply = client.recv(64)
                if reply != b"hello-stream":
                    raise AssertionError(f"unexpected stream reply: {reply!r}")
        print('[stream/recv] send: "hello-stream" → recv: "hello-stream"')
const port = await reservePort();
const endpoint = `tcp://127.0.0.1:${port}`;
const ctx = zlink.createContext();
const stream = zlink.createStreamSocket(ctx);
let client;

try {
  stream.options.recvMode = zlink.StreamRecvMode.Raw;
  stream.bind(endpoint);
  client = net.createConnection({ host: '127.0.0.1', port });
  await once(client, 'connect');
  const sent = 'hello-stream';
  client.write(Buffer.from(sent));

  const received = new zlink.Received();
  assert.equal(stream.recv(received), true);
  try {
    assert.ok(received.routingId instanceof zlink.RoutingId);
    const recv = received.parts[0].data().toString();
    assert.equal(recv, sent);

    await received.send().message(Buffer.from(sent)).submit().admitted;
    const [reply] = await once(client, 'data');
    assert.equal(reply.toString(), sent);
    console.log(`[stream/recv] send: "${sent}" \u2192 recv: "${recv}"`);
  } finally {
    received.close();
  }
} finally {
  if (client) {
    client.destroy();
  }
  stream.close();
  ctx.close();
}
const port = await reservePort();
const endpoint = `tcp://127.0.0.1:${port}`;
const ctx = zlink.createContext();
const stream = zlink.createStreamSocket(ctx);
let client;

try {
  stream.bind(endpoint);
  client = net.createConnection({ host: '127.0.0.1', port });
  await once(client, 'connect');
  const sent = 'hello-stream';
  client.write(Buffer.from(sent));

  const received = new zlink.Received();
  assert.equal(stream.recv(received), true);
  try {
    assert.ok(received.routingId instanceof zlink.RoutingId);
    const recv = received.parts[0].data().toString();
    assert.equal(recv, sent);

    received.send().message(Buffer.from(sent)).submit();
    const [reply] = await once(client, 'data');
    assert.equal(reply.toString(), sent);
    console.log(`[stream/recv] send: "${sent}" → recv: "${recv}"`);
  } finally {
    received.close();
  }
} finally {
  if (client) {
    client.destroy();
  }
  stream.close();
  ctx.close();
}
ctx, err := zlink.NewContext()
samplecommon.Must(err)
defer ctx.Close()

server, err := ctx.StreamSocket()
samplecommon.Must(err)
defer server.Close()
samplecommon.Must(server.SetReceiveMode(zlink.StreamReceiveRaw))

endpoint := samplecommon.UniqueTCP("stream-recv")
samplecommon.Must(server.Bind(endpoint))

conn := samplecommon.DialEndpoint(endpoint)
defer conn.Close()

sent := "hello-stream"
_, err = conn.Write([]byte(sent))
samplecommon.Must(err)

var received zlink.Received
_, err = server.Recv(&received, zlink.RecvFlagsNone)
samplecommon.Must(err)
defer received.Close()
part, err := received.SinglePartOrError()
samplecommon.Must(err)
if !bytes.Equal(part.Data(), []byte(sent)) {
    samplecommon.Must(fmt.Errorf("unexpected payload %q", string(part.Data())))
}

submission, err := received.Send().Message(samplecommon.Message(sent)).Submit(context.Background())
if err == nil {
    err = submission.Admitted(context.Background())
}
samplecommon.Must(err)

buffer := make([]byte, len(sent))
_, err = io.ReadFull(conn, buffer)
samplecommon.Must(err)
fmt.Printf("[stream/recv] send: %q -> recv: %q\n", sent, string(buffer))
let ctx = Context::new().expect("context creation failed");

let stream = ctx.stream_socket().expect("stream socket failed");
stream
    .stream_options()
    .set_recv_mode(zlink::StreamRecvMode::Raw)
    .expect("set raw receive mode failed");
stream.bind("tcp://127.0.0.1:0").expect("bind failed");
let endpoint = stream.last_endpoint().expect("last_endpoint failed");
let stream_mon = SocketMonitor::open(&stream).expect("stream monitor open failed");

let tcp_addr = endpoint.strip_prefix("tcp://").unwrap();
let mut tcp_client = TcpStream::connect(tcp_addr).expect("tcp connect failed");
tcp_client.set_nodelay(true).expect("set_nodelay failed");
sample_support::wait_stream_connected(&stream_mon);
drop(stream_mon);

tcp_client
    .write_all(b"hello-stream")
    .expect("tcp write failed");
tcp_client.flush().expect("tcp flush failed");

let mut received = zlink::Received::empty();
stream
    .recv(&mut received, zlink::RecvFlags::NONE)
    .expect("server recv failed");
assert_eq!(received.parts()[0].as_bytes(), b"hello-stream");
let submission = received
    .send()
    .message(zlink::Message::try_from(b"hello-stream").expect("reply message failed"))
    .submit()
    .expect("stream reply failed");
sample_support::block_on(submission.admitted).expect("stream reply admission failed");
let mut response = [0u8; 12];
tcp_client
    .read_exact(&mut response)
    .expect("tcp read failed");
assert_eq!(&response, b"hello-stream");
println!(
    "[stream/recv] send: \"hello-stream\" → recv: \"{}\"",
    received.parts()[0].as_str().unwrap()
);

PACKET pull examples

These variants pull packets that use the fixed framing convention.

zlink::context_t ctx;
zlink::stream_socket_t server (ctx);
zlink::socket_monitor_t server_monitor = server.monitor_open ();
server.options ().notify (false);
server.options ().recv_mode (zlink::stream_recv_mode_t::packet);

server.bind ("tcp://127.0.0.1:0");
const std::string endpoint = server.options ().last_endpoint ();
assert (!endpoint.empty ());

detail::raw_tcp_client_t client (endpoint);
assert (detail::wait_stream_connected (server_monitor));

const std::string request = detail::k_stream_payload;
const std::vector<unsigned char> request_frame = encode_packet_frame (request);
client.send_all (reinterpret_cast<const char *> (request_frame.data ()), request_frame.size ());

zlink::stream_packet_t packet;
const auto deadline = std::chrono::steady_clock::now () + std::chrono::seconds (2);
while (!server.recv_packet (packet, zlink::recv_flags_t::dontwait)) {
    assert (std::chrono::steady_clock::now () < deadline);
    std::this_thread::sleep_for (std::chrono::milliseconds (1));
}
assert (packet_payload (packet.header (), packet.body ()) == detail::k_stream_payload);

zlink::message_t reply = detail::make_message (detail::k_stream_payload);
assert (packet.routing_id ().has_value ());
server.send (*packet.routing_id ()).message (reply).submit ();

char response[64];
const int received = client.recv_exact (response, request.size ());
assert (received == static_cast<int> (std::strlen (detail::k_stream_payload)));
assert (std::memcmp (response, detail::k_stream_payload, received) == 0);
std::printf ("[stream/packet-pull] send: \"%s\" → recv: \"%.*s\"\n", request.c_str (),
             received, response);

client.close ();
return 0;
if (!SampleSupport.IsNativeAvailable())
    return;

using var ctx = Zlink.CreateContext();
using var stream = ctx.CreateStreamSocket();
stream.Options.Linger = TimeSpan.Zero;
stream.Options.ReceiveMode = StreamReceiveMode.Packet;
string endpoint = SampleSupport.NewEndpoint("tcp", "sample");
int port = SampleSupport.ExtractPort(endpoint);
using var monitor = stream.MonitorOpen(SocketEvent.Accepted);
stream.Bind(endpoint);

using var client = SampleSupport.ConnectRawClient(port);
SampleSupport.WaitMonitorEvent(monitor, 5000, SocketEvent.Accepted);
SampleSupport.SendStreamPacket(client.GetStream(), "hello-stream"u8);
using var packet = StreamPacket.Create();
if (!stream.RecvPacket(packet))
    throw new InvalidOperationException("packet receive failed");
string callbackPayload = packet.Body?.GetString()
    ?? throw new InvalidOperationException("missing packet body");
Console.WriteLine(
    $"[stream/packet-pull] send: \"hello-stream\" -> recv: \"{callbackPayload}\"");
SampleSupport.ensureNative();
String endpoint = SampleSupport.tcpEndpoint();
try (Context ctx = Zlink.createContext();
     StreamSocket server = ctx.createStreamSocket();
     var monitor = server.monitorOpen(
         systems.zlink.contracts.eventing.MonitorEventType.ACCEPTED,
         systems.zlink.contracts.eventing.MonitorEventType.CONNECTION_READY)) {
    server.options().recvMode(StreamRecvMode.PACKET);
    server.bind(endpoint);
    try (var rawClient = SampleSupport.connectRawTcp(endpoint)) {
        SampleSupport.waitStreamConnected(monitor);

        try (StreamPacket packet = new StreamPacket()) {
            byte[] payload = SampleSupport.STREAM_PAYLOAD
                .getBytes(StandardCharsets.UTF_8);
            byte[] frame = frameBytes(payload);
            SampleSupport.sendRawTcp(rawClient, frame);
            if (!server.recvPacket(packet, RecvFlags.NONE)) {
                throw new IllegalStateException("missing stream packet");
            }
            var routingId = packet.routingId().orElseThrow();
            String value = packet.body().toUtf8String();
            if (!SampleSupport.STREAM_PAYLOAD.equals(value)) {
                throw new IllegalStateException("unexpected stream payload");
            }
            try (Message replyHeader = Message.from(new byte[0]);
                 Message replyBody = Message.from(
                     SampleSupport.STREAM_PAYLOAD);
                 Message reply = frame(replyHeader, replyBody)) {
                server.send(routingId)
                    .message(reply)
                    .submit().admitted().toCompletableFuture().join();
            }

            byte[] echoedFrame = SampleSupport.recvExactRawTcp(
                rawClient, frame.length);
        String echoed = new String(
            java.util.Arrays.copyOfRange(echoedFrame, 6,
                echoedFrame.length),
            StandardCharsets.UTF_8);
        if (!SampleSupport.STREAM_PAYLOAD.equals(echoed)) {
            throw new IllegalStateException("unexpected stream echo: " + echoed);
        }
        System.out.println("[stream/packet-callback] send: \"" +
            SampleSupport.STREAM_PAYLOAD + "\" -> recv: \"" +
            echoed + "\"");
        }
    }
}

A current pull-based Kotlin PACKET sample is not yet available in this source tree.

port, endpoint = tcp_endpoint()

with zlink.create_context() as ctx:
    with zlink.create_stream_socket(ctx) as server:
        server.stream_options.recv_mode = zlink.StreamRecvMode.PACKET
        with server.monitor_open(zlink.MonitorEventMask.ACCEPTED) as server_monitor:
            server.bind(endpoint)
            with socket.create_connection(("127.0.0.1", port), timeout=3.0) as client:
                wait_socket_monitor_event(server_monitor, zlink.MonitorEventMask.ACCEPTED)
                client.sendall(
                    struct.pack("!HI", 0, len(b"hello-stream")) + b"hello-stream"
                )
                packet = zlink.StreamPacket()
                if not server.recv_packet_into(packet):
                    raise AssertionError("expected stream packet")
                with packet:
                    if packet.header.to_bytes() != b"":
                        raise AssertionError("unexpected stream packet header")
                    if packet.body.to_bytes() != b"hello-stream":
                        raise AssertionError("unexpected stream packet body")
                    if not packet.routing_id:
                        raise AssertionError("stream packet expected a routing id")
                    server.send(packet.routing_id).message(b"hello-stream").submit_sync()
                reply = client.recv(64)
                if reply != b"hello-stream":
                    raise AssertionError(f"unexpected stream packet reply: {reply!r}")
        print('[stream/packet-recv] send: "hello-stream" -> recv: "hello-stream"')
const port = await reservePort();
const endpoint = `tcp://127.0.0.1:${port}`;
const ctx = zlink.createContext();
const stream = zlink.createStreamSocket(ctx);
const packet = new zlink.StreamPacket();
let client;

try {
  stream.options.recvMode = zlink.StreamRecvMode.Packet;
  stream.bind(endpoint);
  client = net.createConnection({ host: '127.0.0.1', port });
  await once(client, 'connect');
  client.write(frame(Buffer.from('hello-stream')));

  for (let attempt = 0; attempt < 200; attempt += 1) {
    if (stream.recvPacket(packet, zlink.RecvFlags.DontWait)) break;
    await new Promise((resolve) => setTimeout(resolve, 2));
  }
  assert.ok(packet.routingId instanceof zlink.RoutingId);
  assert.equal(packet.header.data().length, 0);
  assert.equal(packet.body.data().toString(), 'hello-stream');
  console.log('[stream/packet] send: "hello-stream" -> recv: "hello-stream"');
} finally {
  packet.close();
  if (client) client.destroy();
  stream.close();
  ctx.close();
}

A current pull-based JavaScript PACKET sample is not yet available in this source tree.

ctx, err := zlink.NewContext()
samplecommon.Must(err)
defer ctx.Close()

server, err := ctx.StreamSocket()
samplecommon.Must(err)
defer server.Close()
samplecommon.Must(server.SetReceiveMode(zlink.StreamReceivePacket))

endpoint := samplecommon.UniqueTCP("stream-packet-callback")
samplecommon.Must(server.Bind(endpoint))

conn := samplecommon.DialEndpoint(endpoint)
defer conn.Close()

sent := "hello-stream"
samplecommon.WriteStreamPacket(conn, []byte(sent))
var received zlink.StreamPacket
ok, err := server.RecvPacket(&received, zlink.RecvFlagsNone)
samplecommon.Must(err)
if !ok || string(received.Body().Data()) != sent {
    samplecommon.Must(fmt.Errorf("unexpected packet body"))
}
packet := samplecommon.FrameStreamPacketMessage(received.Header(), received.Body())
submission, err := server.SendTo(received.RoutingID()).Message(packet).Submit(context.Background())
samplecommon.Must(err)
samplecommon.Must(submission.Admitted(context.Background()))

buffer := samplecommon.ReadStreamPacketBody(conn)

fmt.Printf("[stream/packet-callback] send: %q -> recv: %q\n", sent, string(buffer))
let ctx = Context::new().expect("context creation failed");

let stream = ctx.stream_socket().expect("stream socket failed");
stream
    .stream_options()
    .set_recv_mode(StreamRecvMode::Packet)
    .expect("set packet mode failed");
let endpoint = sample_support::tcp_endpoint();
stream.bind(&endpoint).expect("bind failed");
let stream_mon = SocketMonitor::open(&stream).expect("stream monitor open failed");

let tcp_addr = endpoint.strip_prefix("tcp://").unwrap();
let mut tcp_client = TcpStream::connect(tcp_addr).expect("tcp connect failed");
tcp_client.set_nodelay(true).expect("set_nodelay failed");
sample_support::wait_stream_connected(&stream_mon);
drop(stream_mon);

write_stream_packet(&mut tcp_client, b"hello-stream");
let mut packet = StreamPacket::empty();
assert!(
    stream
        .recv_packet(&mut packet, RecvFlags::NONE)
        .expect("packet receive failed")
);
assert!(
    !packet
        .routing_id()
        .expect("packet routing id")
        .as_bytes()
        .is_empty()
);
assert!(
    packet
        .header()
        .expect("packet header")
        .as_bytes()
        .is_empty()
);
let payload = packet.body().expect("packet body").as_bytes();
assert_eq!(payload, b"hello-stream");
let recv_str = std::str::from_utf8(payload).unwrap();
println!(
    "[stream/packet-recv] send: \"hello-stream\" → recv: \"{}\"",
    recv_str
);
packet.close().expect("packet close failed");

← ROUTER | Proxy →