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:
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 itssource_rid_out_out-parameter. Pair it with a poller watchingZLINK_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_ridis auto-assigned per connection by the server, always fixed 4 bytes (uint32, big-endian).- To close one client, pass the
source_ridreceived from recv tozlink_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 asZLINK_EVENT_CONNECTION_READY/ZLINK_EVENT_DISCONNECTED, each carrying the 4-byterouting_id. A raw payload that happens to be a single0x00/0x01byte is delivered as ordinary data. IfZLINK_STREAM_OPT_NOTIFYis set to1before bind/connect in RAW mode,zlink_recv()additionally returns a zero-length record with the affectedsource_ridfor 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 returnsZLINK_SUBMIT_BACKPRESSUREDwithZLINK_DONTWAIT. For advanced backpressure patterns, see Performance Guide.
- The caller owns a successfully received
zlink_msg_tand closes it exactly once. - Copy the borrowed
source_ridbefore the next data receive when it must be retained. ZLINK_RECV_FLAGS_DONTWAITreturnsZLINK_RECV_NO_DATAwithEAGAINwhen 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_sizeorbody_sizeequal to zero is allowed; both sides are still delivered as validzlink_msg_tobjects.- Ownership of
headerandbodyis transferred to the caller. The caller must close or consume eachmsg_texactly once. - In PACKET mode, whole-message RAW receive (
zlink_recv()) fails withENOTSUP. 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_idframe is a protocol error - Messages larger than
MAXMSGSIZEare dropped and connection is closed (disconnect event)
8. Reference Tests¶
core/tests/integration/test_stream_socket.cppcore/tests/integration/test_stream_fastpath.cppcore/tests/integration/routing-id/test_connect_rid_string_alias.cppcore/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");