콘텐츠로 이동

가이드 목록 | 이전: ROUTER | 다음: 프록시 패턴

STREAM 소켓

이 장의 계약 소유 문서STREAM socket 스펙이 다룬다. 이 챕터는 그 계약을 언어별 예제로 보여준다.

1. 개요

STREAM 소켓은 zlink framing 없이 raw byte를 주고받는 소켓이다. 반대편은 OS socket·WebSocket 같은 외부 raw client일 수도 있고, receive mode를 설정한 다른 zlink STREAM 소켓일 수도 있다.

핵심 규칙: - 첫 zlink_bind() 또는 zlink_connect() 전에 ZLINK_STREAM_OPT_RECV_MODE를 RAW 또는 PACKET으로 설정한다. mode 없이 bind/connect하면 EINVAL로 실패한다. - ZLINK_SOCKET_STREAM은 bind와 connect를 모두 지원한다(zlink_connect() 성공 시 ZLINK_CONNECT_OK). - RAW 모드는 raw 바이트 스트림을 그대로 전달한다. PACKET 모드는 2바이트 BE header 크기 + 4바이트 BE body 크기 + header + body framing을 사용한다. - zlink API 수준에서: RAW zlink_recv()는 한 part로 된 byte record와 발신 클라이언트의 4바이트 routing_id를 자체 source_rid_out_ out-parameter로 노출하며, packet zlink_stream_recv_packet()도 같은 Core-owned borrowed view를 source_rid_out_로 노출한다.

유효 조합:

external raw client  <---- RAW 바이트 스트림 (framing 없음) ---->  STREAM(server)

STREAM은 zlink 내부 소켓(PAIR/PUB/SUB/DEALER/ROUTER)과 직접 호환되지 않는다.


2. 서버 생성/바인드

void *stream = zlink_socket(ctx, ZLINK_SOCKET_STREAM);
zlink_stream_recv_mode_t mode = ZLINK_STREAM_RECV_MODE_RAW;   /* 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");

지원 transport(bind·connect 공통): - tcp:// - tls:// - ws:// - wss://


3. STREAM 고유 동작

STREAM은 기반 소켓 계열(raw socket family)에서 유일한 예외 타입이다. 첫 bind 성공 전에 두 가지 수신 모드 중 정확히 하나를 고른다.

  • RAW: zlink_recv()로 transport byte record 하나를 가져온다. Record는 part 하나로 반환되고, 소스 routing id는 source_rid_out_ out-parameter로 받는다. poller의 ZLINK_POLLIN과 함께 사용한다.
  • PACKET: zlink_stream_recv_packet()으로 고정 framing(framing, 패킷 경계를 구분하는 방식) 규약(2B header size + 4B body size + header + body, big-endian)을 따르는 패킷을 조립된 header/body 형태로 받는다.

bind 전에 ZLINK_STREAM_OPT_RECV_MODEZLINK_STREAM_RECV_MODE_RAW 또는 ZLINK_STREAM_RECV_MODE_PACKET으로 설정한다. 첫 bind 성공 뒤에는 모드를 바꿀 수 없고, 다른 모드의 수신 API는 ENOTSUP로 실패한다.

STREAM만의 고유 동작은 다음과 같다.

  • source_rid는 서버가 연결별로 자동 할당하며, 고정 4바이트(uint32, big-endian)이다.
  • 특정 클라이언트를 끊어야 하면 recv에서 받은 source_ridzlink_disconnect_rid()에 넘긴다. STREAM의 대상 rid는 반드시 4바이트다.
  • 기본값(ZLINK_STREAM_OPT_NOTIFY=0)에서 connect/disconnect는 in-band 데이터 마커가 아니다. 소켓 monitor의 ZLINK_EVENT_CONNECTION_READY / ZLINK_EVENT_DISCONNECTED 이벤트로 보고되며 각각 4바이트 routing_id를 담는다. 우연히 1바이트 0x00/0x01인 raw payload는 일반 데이터로 전달된다. RAW 모드에서 bind/connect 전에 ZLINK_STREAM_OPT_NOTIFY1로 설정하면 zlink_recv()가 연결·해제마다 해당 source_rid와 함께 길이 0 record를 추가로 돌려주므로, 그 경우 길이 0 part는 데이터가 아니라 알림으로 해석한다.

4. RAW pull 예시

ZLINK_STREAM_OPT_NOTIFY=0(기본)이면 STREAM RAW 모드에서 pull한 모든 part는 애플리케이션 데이터이고, connect/disconnect는 소켓 monitor로 관찰한다(Monitoring 참고). NOTIFY=1이면 길이 0 record가 연결·해제 알림으로 섞여 들어온다.

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);
}

주요 사항

항목 설명
수신 API zlink_recv()
readiness poller가 ZLINK_POLLIN을 알리면 애플리케이션이 recv를 drain
수명 source_rid는 같은 소켓의 다음 data-recv 진입 또는 close 전까지 유효
framing transport에서 수신된 raw 바이트
전송 zlink_send_rid()

송신 큐가 가득 차면(HWM, 고수위 표시) zlink_send_rid()는 블록(기본) 또는 ZLINK_DONTWAITZLINK_SUBMIT_BACKPRESSURED 를 반환한다. 배압(backpressure) 패턴은 성능 가이드를 참고.

  • 성공적으로 받은 zlink_msg_t는 호출자가 소유하며 정확히 한 번 close한다.
  • 다음 data receive 뒤에도 source_rid가 필요하면 borrowed view를 미리 복사한다.
  • 빈 큐에서 ZLINK_RECV_FLAGS_DONTWAIT을 쓰면 EAGAINZLINK_RECV_NO_DATA를 반환한다.

4.1 PACKET pull 모드

고정 framing 규약(2바이트 big-endian header size + 4바이트 big-endian body size + header payload + body payload)을 사용하는 상위 프로토콜에서는 PACKET 모드를 선택하고 zlink_stream_recv_packet()으로 pull한다. Core가 조각(fragment) 누적과 길이 해석을 직접 처리하므로 응용은 header/body를 그대로 받아 처리한다.

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) {
    /* header/body 처리. 길이 0인 메시지도 유효하다. */
    zlink_msg_close(&header);
    zlink_msg_close(&body);
}

PACKET 모드의 규칙은 다음과 같다.

  • header_sizebody_size 는 각각 0 도 허용된다. 길이가 0 이어도 msg_t 는 유효한 객체로 전달된다.
  • headerbody 의 소유권은 호출자로 이전된다. 호출자는 두 msg_t 를 각각 정확히 한 번 close 하거나 소비해야 한다.
  • PACKET 모드에서 whole-message RAW receive(zlink_recv())는 ENOTSUP로 실패한다. RAW 모드에서 zlink_stream_recv_packet()도 같은 방식으로 실패한다.
  • framing 규약을 지키지 않는 비정형 패킷(malformed packet)(길이 제한 초과, 조립 실패, 불완전 상태 연결 종료 등)은 연결을 닫는 기본 동작으로 이어진다. 이 이벤트는 소켓 모니터(socket monitor) 경로로 관찰한다.

이 모드는 조각 누적을 응용 쪽에서 다시 구현하지 않아도 되는 편의를 주지만 transport 조각 경계와 패킷 경계가 다르다는 점 자체를 바꾸지는 않는다.


5. 클라이언트 구현 원칙

클라이언트는 raw socket/websocket로 구현하거나, receive mode를 설정한 zlink STREAM 소켓으로 zlink_connect()해도 된다. STREAM은 raw 바이트를 그대로 전달하므로 패킷 경계(framing)는 애플리케이션이 정의해야 한다.

개념적 POSIX TCP 예시 (RAW 모드 — zlink framing 없음, 바이트 그대로):

// RAW 모드: raw 바이트 송수신. 메시지 경계는 애플리케이션이 정의
send(fd, body, body_len, 0);

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

서버가 PACKET 모드(zlink_stream_recv_packet)를 쓰면 클라이언트는 각 패킷을 2바이트 BE header size + 4바이트 BE body size + header + body로 framing해야 한다:

// 패킷 모드: [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. 옵션 및 런타임 정책

주요 옵션:

  • 지원:
  • 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 (zlink_set_stream_option() / zlink_get_stream_option()): 첫 bind 또는 connect 전에 RAW 또는 PACKET 선택
  • ZLINK_STREAM_OPT_NOTIFY: RAW 모드의 길이 0 connect/disconnect record 활성화. PACKET과 결합 불가
  • TLS/WSS 서버: zlink_set_tls_server()
  • TLS 클라이언트: zlink_set_tls_client()

STREAM listener는 raw TCP 피어가 보낸 바이트를 직접 받을 수 있다. 피어를 완전히 신뢰할 수 없다면 zlink_bind를 호출하기 전에 애플리케이션이 받아들일 최대 메시지 크기로 ZLINK_OPT_MAXMSGSIZE를 설정한다. 이 값을 설정하지 않으면 호환성을 위해 기본값은 무제한이다.

비지원/변경: - ZLINK_ROUTER_OPT_CONNECT_ROUTING_ID를 STREAM에 설정하면 EOPNOTSUPP

6.1 STREAM 기본 런타임 프로파일

STREAM socket이 public option에 두는 기본값: - ZLINK_OPT_BACKLOG: 65536 - ZLINK_OPT_SNDHWM / ZLINK_OPT_RCVHWM: 기본 balanced auto-HWM 정책의 STREAM profile byte 값. context auto-HWM을 끄면 수동 byte 기본값 사용 - ZLINK_OPT_SNDBUF / ZLINK_OPT_RCVBUF: 기본값 -1. OS 기본 버퍼와 TCP 자동 조정에 맡김

각 option의 정확한 계약은 STREAM socket 스펙이 소유한다.


7. 에러/제약

  • zlink_connect(stream, ...) -> EOPNOTSUPP
  • STREAM 대상 routing_id는 4바이트여야 하며, 크기가 다르면 호출자 인자 오류(EINVAL)
  • MAXMSGSIZE 초과 메시지는 연결 종료(disconnect 이벤트)

8. 테스트 기준 구현

참고 파일: - 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

위 테스트들은 STREAM 서버 + raw client 경로를 기준으로 동작한다.


← ROUTER | Proxy → | Transport →

언어별 완전한 예제

STREAM 소켓으로 원시 바이트를 주고받는 자립형 예제다(모든 바인딩, 빌드·실행 검증됨).

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 방식

고정 framing의 수신 패킷을 애플리케이션이 pull하는 변형이다.

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 + "\"");
        }
    }
}

현재 source tree에는 pull 기반 Kotlin PACKET sample이 아직 없다.

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();
}

현재 source tree에는 pull 기반 JavaScript PACKET sample이 아직 없다.

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 | 프록시 →