가이드 목록 | 이전: PUB/SUB | 다음: ROUTER
DEALER 소켓¶
이 장의 계약 소유 문서 — DEALER socket 스펙이 다룬다. 이 챕터는 그 계약을 언어별 예제로 보여준다.
1. 개요¶
DEALER 소켓은 비동기 요청 소켓이다. 여러 peer에 round-robin으로 송신하고, fair-queuing으로 수신한다. send/recv 순서 강제가 없어 자유로운 비동기 메시징이 가능하다.
핵심 특성: - 송신: round-robin — 연결된 peer에 차례로 분배 - 수신: fair-queuing — 모든 peer에서 공정하게 수신 - send/recv 순서 강제 없음 (비동기)
유효한 소켓 조합: DEALER ↔ ROUTER, DEALER ↔ DEALER
%%{init: {'flowchart': {'nodeSpacing': 32, 'rankSpacing': 40, 'padding': 8, 'wrappingWidth': 180}, 'themeVariables': {'fontSize': '18px'}}}%%
flowchart LR
D1[DEALER 1] -->|round-robin| R[ROUTER]
D2[DEALER 2] -->|round-robin| R
구체적 시나리오: 3개 DEALER가 1개 ROUTER로 전송¶
3개의 DEALER 클라이언트가 하나의 ROUTER 서버에 연결한다. 각 DEALER는
독립적으로 요청을 전송하며 ROUTER는 fair-queuing으로 수신하고
source_rid로 각 송신자를 구분한다.
| 송신자 | routing_id | 메시지 | ROUTER 수신 |
|---|---|---|---|
| DEALER 1 | D1 |
"buy AAPL 100" |
source_rid=D1, data="buy AAPL 100" |
| DEALER 2 | D2 |
"sell TSLA 50" |
source_rid=D2, data="sell TSLA 50" |
| DEALER 3 | D3 |
"buy MSFT 200" |
source_rid=D3, data="buy MSFT 200" |
ROUTER는 zlink_send_rid()에 해당 source_rid를 전달하여 각 DEALER에
응답한다. DEALER는 송신 연결에 round-robin을 사용하므로 하나의
DEALER가 여러 ROUTER에 연결하면 메시지가 round-robin으로 순환 분배된다
(msg1 -> ROUTER-A, msg2 -> ROUTER-B, ...).
2. 기본 사용법¶
생성 및 연결¶
void *dealer = zlink_socket(ctx, ZLINK_SOCKET_DEALER);
/* routing_id 설정(선택, ROUTER 쪽 식별용) */
zlink_set_routing_id(dealer, "client-1", 8);
/* 서버에 연결 */
zlink_connect(dealer, "tcp://127.0.0.1:5558");
메시지 송수신¶
/* 요청 전송 -- 각 호출이 단일 part record 전체를 제출하므로
순서 제약 없이 연속으로 보낼 수 있다. */
zlink_msg_t msg1, msg2, msg3;
zlink_msg_init_size(&msg1, 9);
memcpy(zlink_msg_data(&msg1), "request-1", 9);
zlink_send(dealer, &msg1, 1, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
zlink_msg_init_size(&msg2, 9);
memcpy(zlink_msg_data(&msg2), "request-2", 9);
zlink_send(dealer, &msg2, 1, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
zlink_msg_init_size(&msg3, 9);
memcpy(zlink_msg_data(&msg3), "request-3", 9);
zlink_send(dealer, &msg3, 1, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
/* 일반 DATA는 zlink_recv()로 소비한다.
zlink_request() 결과는 zlink_completion_recv()로 drain한다. */
수신 모드¶
DEALER는 zlink_recv()로 한 번에 record 전체를 동기 수신한다.
source_rid_out_은 DEALER에서는 선택 사항이다 — 어차피 DEALER는 이 자리에
NULL을 반환하므로, 필요 없으면 NULL을 넘긴다.
zlink_msg_t parts[8];
size_t part_count = 0;
zlink_recv_result_t rc = zlink_recv(
dealer, NULL, parts, 8, &part_count, ZLINK_RECV_FLAGS_NONE);
if (rc == ZLINK_RECV_OK) {
/* parts[0..part_count)를 처리한 뒤 배열 전체를 닫는다. */
zlink_multipart_close(parts, part_count);
}
/* 그 밖의 rc 값: ZLINK_RECV_NO_DATA (EAGAIN), TERMINATED, INVALID_HANDLE */
HWM(High-Water Mark, queue가 보관할 수 있는 accounted byte 상한) 도달 시
zlink_send()는 대기(기본) 또는ZLINK_SEND_FLAGS_DONTWAIT로ZLINK_SUBMIT_BACKPRESSURED를 반환한다. 고급 배압(backpressure) 패턴은 성능 가이드를 참고.
3. 사용 예제¶
/* DEALER → ROUTER send: 배열 하나에 담은 2-part record. */
zlink_msg_t parts[2];
zlink_msg_init_size(&parts[0], 6);
memcpy(zlink_msg_data(&parts[0]), "header", 6);
zlink_msg_init_size(&parts[1], 4);
memcpy(zlink_msg_data(&parts[1]), "body", 4);
zlink_submit_result_t rc = zlink_send(
dealer, parts, 2, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
4. 소켓 옵션¶
| 옵션 | 타입 | 기본값 | 설명 |
|---|---|---|---|
zlink_set_routing_id() |
binary | 자동(UUID) | ROUTER에서 식별할 ID (전용 함수) |
ZLINK_DEALER_OPT_PROBE |
int | 0 | 연결 시 빈 메시지 전송 (연결 알림) |
ZLINK_DEALER_OPT_REQUEST_TIMEOUT_MS |
int (ms) | 5000 |
request의 timeout_ms_ == 0일 때 사용되는 기본 timeout |
ZLINK_DEALER_OPT_WEIGHT |
int | 100 | 송신 round-robin의 peer별 load balancing 가중치 |
ZLINK_OPT_SNDHWM |
uint64_t bytes |
자동 | DEALER의 peer-queue 역할에 맞춰 산정된 자동 HWM. 수동 설정이 우선하며 0은 무제한 |
ZLINK_OPT_RCVHWM |
uint64_t bytes |
자동 | DEALER의 peer-queue 역할에 맞춰 산정된 자동 HWM. 수동 설정이 우선하며 0은 무제한 |
ZLINK_OPT_LINGER |
int | -1 | close 시 대기 시간 (ms) |
ZLINK_OPT_SNDTIMEO |
int | 1000 | 송신 타임아웃(ms). 무한 대기는 -1을 명시적으로 설정 |
ZLINK_OPT_RCVTIMEO |
int | 1000 | 수신 타임아웃(ms). 무한 대기는 -1을 명시적으로 설정 |
routing_id 설정¶
ROUTER가 DEALER를 식별하려면 명시적으로 routing_id를 설정한다.
/* bind/connect 전에 설정 */
zlink_set_routing_id(dealer, "D1", 2);
zlink_connect(dealer, "tcp://127.0.0.1:5558");
참고:
core/tests/integration/test_router_multiple_dealers.cpp—zlink_set_routing_id(dealer1, "D1", 2)
4.1 request-reply¶
DEALER가 상관된 응답을 기다리려면 일반 DATA send/recv 대신 zlink_request()를
사용한다. 이 함수는 ZMP(zlink 전용 메시지 프로토콜) request-reply envelope(요청-응답
식별용 header wrapper)를 붙이고, reply 또는 terminal 결과를 socket completion queue에 넣는다.
ZMP request-reply envelope의 와이어 프레임 형식은 ZMP 프로토콜을 참고.
int timeout_ms = 1000;
zlink_set_dealer_option(
dealer,
ZLINK_DEALER_OPT_REQUEST_TIMEOUT_MS,
&timeout_ms,
sizeof(timeout_ms));
zlink_msg_t req;
zlink_msg_init_size(&req, 4);
memcpy(zlink_msg_data(&req), "ping", 4);
zlink_completion_id_t id = 0;
zlink_submit_result_t rc = zlink_request(
dealer, NULL, &req, 1, ZLINK_SEND_FLAGS_NONE,
0 /* ZLINK_DEALER_OPT_REQUEST_TIMEOUT_MS 사용 */, NULL, &id);
if (rc == ZLINK_SUBMIT_OK) {
zlink_completion_t completion = {0};
completion.struct_size = sizeof(completion);
if (zlink_completion_recv(dealer, &completion, ZLINK_RECV_FLAGS_NONE)
== ZLINK_RECV_OK) {
/* completion.request_result는 OK, TIMED_OUT, NOT_FOUND 등이다. */
zlink_completion_close(&completion);
}
}
zlink_request()에 timeout_ms_ == 0을 전달하면 DEALER 소켓 기본값인
ZLINK_DEALER_OPT_REQUEST_TIMEOUT_MS(별도 설정이 없으면 5000ms)가
적용된다.
5. 사용 패턴¶
패턴 1: DEALER → ROUTER Request-Reply¶
가장 기본적인 패턴이다. DEALER가 일반 raw 메시지를 보내면, ROUTER는
source_rid로 송신자를 구분하여 같은 id로 응답한다.
void *router = zlink_socket(ctx, ZLINK_SOCKET_ROUTER);
zlink_bind(router, "tcp://*:5558");
void *dealer = zlink_socket(ctx, ZLINK_SOCKET_DEALER);
zlink_set_routing_id(dealer, "D1", 2);
zlink_connect(dealer, "tcp://127.0.0.1:5558");
/* 클라이언트 typed request */
zlink_msg_t req;
zlink_msg_init_size(&req, 5);
memcpy(zlink_msg_data(&req), "Hello", 5);
zlink_completion_id_t request_id = 0;
zlink_request(dealer, NULL, &req, 1, ZLINK_SEND_FLAGS_NONE,
0, NULL, &request_id);
/* 서버: source_rid + opaque reply_token으로 REQUEST 수신 */
const zlink_routing_id_t *source_rid = NULL;
zlink_reply_token_t reply_token = 0;
zlink_msg_t parts[8];
size_t part_count = 0;
zlink_router_recv(router, &source_rid, &reply_token,
parts, 8, &part_count, ZLINK_RECV_FLAGS_NONE);
printf("Received from [%.*s]: %.*s\n", (int)source_rid->size, source_rid->data,
(int)zlink_msg_size(&parts[0]), (char *)zlink_msg_data(&parts[0]));
/* 0이 아닌 token은 REQUEST이며 그대로 돌려준다. */
zlink_msg_t reply;
zlink_msg_init_size(&reply, 5);
memcpy(zlink_msg_data(&reply), "World", 5);
zlink_reply(router, source_rid, reply_token, &reply, 1);
zlink_multipart_close(parts, part_count);
/* 클라이언트: DATA가 아닌 REQUEST completion으로 응답 수신 */
zlink_completion_t completion = {0};
completion.struct_size = sizeof(completion);
zlink_completion_recv(dealer, &completion, ZLINK_RECV_FLAGS_NONE);
zlink_completion_close(&completion);
참고:
core/tests/integration/test_router_multiple_dealers.cpp— TCP/IPC/inproc 예제
패턴 2: 다중 DEALER Load Balancing¶
여러 DEALER가 하나의 ROUTER에 연결한다. ROUTER는 routing_id로 각 DEALER를 구분한다.
void *router = zlink_socket(ctx, ZLINK_SOCKET_ROUTER);
zlink_bind(router, "tcp://127.0.0.1:*");
char endpoint[256];
size_t len = sizeof(endpoint);
zlink_get_option(router, ZLINK_OPT_LAST_ENDPOINT, endpoint, &len);
void *dealer1 = zlink_socket(ctx, ZLINK_SOCKET_DEALER);
zlink_set_routing_id(dealer1, "D1", 2);
zlink_connect(dealer1, endpoint);
void *dealer2 = zlink_socket(ctx, ZLINK_SOCKET_DEALER);
zlink_set_routing_id(dealer2, "D2", 2);
zlink_connect(dealer2, endpoint);
/* 각 DEALER가 메시지를 보낸다 */
zlink_msg_t m1;
zlink_msg_init_size(&m1, 12);
memcpy(zlink_msg_data(&m1), "from_dealer1", 12);
zlink_send(dealer1, &m1, 1, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
zlink_msg_t m2;
zlink_msg_init_size(&m2, 12);
memcpy(zlink_msg_data(&m2), "from_dealer2", 12);
zlink_send(dealer2, &m2, 1, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
/* zlink_router_recv()는 각 DEALER의 메시지를 source_rid로 구분한다 */
참고:
core/tests/integration/test_router_multiple_dealers.cpp—test_router_multiple_dealers_tcp()
패턴 3: Proxy 패턴 (ROUTER-DEALER)¶
ROUTER(frontend) + DEALER(backend)로 멀티스레드 서버를 구성한다.
zlink_proxy()는 클라이언트의 routing-id envelope를 포함해 두 raw 소켓
사이의 multipart 프레임을 양방향으로 그대로 전달하므로, 워커 코드는
routing-id를 지정하는 send를 호출할 필요가 없다 — 받은 프레임을
그대로 되돌려 보내면 된다.
/* Frontend: 클라이언트가 여기에 연결한다 */
void *frontend = zlink_socket(ctx, ZLINK_SOCKET_ROUTER);
zlink_bind(frontend, "tcp://*:5558");
/* Backend: 워커 스레드가 여기에 연결한다 */
void *backend = zlink_socket(ctx, ZLINK_SOCKET_DEALER);
zlink_bind(backend, "inproc://backend");
/* 워커 스레드를 먼저 시작한 뒤 proxy 루프를 돌린다(호출 스레드를 블록) */
zlink_proxy(frontend, backend, NULL);
/* 워커 스레드: backend에 연결된 일반 DEALER다. proxy가 클라이언트의
envelope와 payload를 한 record로 넘겨주므로, 워커는 배열 전체를 한 번에
받고 응답 앞에 같은 envelope를 다시 실어 보낸다. Target routing id는
명시적으로 넘기지 않는다. */
void worker_thread(void *arg) {
void *worker = zlink_socket(ctx, ZLINK_SOCKET_DEALER);
zlink_connect(worker, "inproc://backend");
zlink_msg_t parts[2];
size_t part_count = 0;
zlink_recv(worker, NULL, parts, 2, &part_count, ZLINK_RECV_FLAGS_NONE);
/* parts[0]은 envelope, parts[1]은 요청 payload다. */
/* 처리 후 응답: envelope를 다시 보내고 이어서 응답 payload를 보낸다 */
zlink_msg_close(&parts[1]);
zlink_msg_init_size(&parts[1], 5);
memcpy(zlink_msg_data(&parts[1]), "World", 5);
zlink_send(worker, parts, 2, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
/* 워커는 소켓이 닫힐 때까지 살아 있는다 */
}
참고:
core/tests/integration/test_proxy.cpp— ROUTER(frontend) + DEALER(backend) + worker pool
패턴 4: DEALER ↔ DEALER 비동기 통신¶
양쪽 모두 DEALER를 사용하는 완전한 비동기 P2P 통신이다. PAIR와 유사하지만 HWM과 자동 재연결을 지원한다. 응답이 필요한 경우 반드시 1:1로 구성해야 한다 — routing_id가 없으므로 1:N 구성에서는 어떤 peer가 응답했는지 구분할 수 없다.
void *a = zlink_socket(ctx, ZLINK_SOCKET_DEALER);
zlink_bind(a, "tcp://*:5558");
void *b = zlink_socket(ctx, ZLINK_SOCKET_DEALER);
zlink_connect(b, "tcp://127.0.0.1:5558");
/* 양방향 자유 송신 */
zlink_msg_t ping;
zlink_msg_init_size(&ping, 4);
memcpy(zlink_msg_data(&ping), "ping", 4);
zlink_send(a, &ping, 1, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
zlink_msg_t pong;
zlink_msg_init_size(&pong, 4);
memcpy(zlink_msg_data(&pong), "pong", 4);
zlink_send(b, &pong, 1, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
/* b는 zlink_recv()로 "ping"을, a는 "pong"을 받는다 */
6. 주의사항¶
peer 없음 vs HWM 배압¶
둘은 별개의 결과다. 연결된 peer가 없으면(양수 가중치 pipe 없음) 송신은
ZLINK_SUBMIT_NOT_ADMITTED를 반환하며 메시지는 큐에 쌓이지 않는다. peer가
연결되어 있지만 그 송신 큐가 HWM에 도달하면 대기(기본) 또는
ZLINK_SEND_FLAGS_DONTWAIT 시 ZLINK_SUBMIT_BACKPRESSURED를 반환한다.
/* 연결된 peer 없이 송신 */
zlink_msg_t msg;
zlink_msg_init_size(&msg, 4);
memcpy(zlink_msg_data(&msg), "data", 4);
zlink_completion_id_t wait_token = 0; /* BACKPRESSURED일 때 WRITABLE record를 식별하는 token */
zlink_submit_result_t rc = zlink_send(
dealer, &msg, 1, ZLINK_SEND_FLAGS_DONTWAIT, NULL, &wait_token);
if (rc == ZLINK_SUBMIT_NOT_ADMITTED) {
/* 메시지를 받아줄 연결된 peer가 없음 */
} else if (rc == ZLINK_SUBMIT_BACKPRESSURED) {
/* peer는 연결됐지만 큐가 HWM에 도달 — part는 소비됐으므로 보관해 둔 사본을
ZLINK_POLLCOMPLETION 뒤 zlink_completion_recv()로 wait_token의 WRITABLE record를 받고 다시 제출 */
}
round-robin 분배¶
여러 peer가 연결된 경우 메시지는 round-robin으로 순환 분배된다. 특정 peer에게만 전송하려면 ROUTER를 사용한다.
가중치 기반 송신 대상 선택¶
원격 ROUTER는 자신의 peer 가중치(0..10000)를 함께 광고한다. DEALER는
가중치가 0인 ROUTER를 round-robin 후보 집합에서 자동으로 제외하고,
양수 가중치를 가진 ROUTER들 사이에서만 outbound 대상을 고른다. 연결
자체는 유지되므로, 가중치가 0이던 ROUTER가 다시 양수 값으로 돌아오면
재연결 없이 후보 집합에 복귀한다.
연결된 ROUTER가 모두 가중치 0이면 zlink_send()와
zlink_request()는 ZLINK_SUBMIT_NOT_ADMITTED를 반환한다.
연결이 끊긴 것이 아니라 보낼 대상이 일시적으로 없는 상태이므로 호출자는
최소 한 대의 ROUTER가 양수 가중치로 복귀할 때까지 기다렸다가
재시도해야 한다. NOT_ADMITTED를 영구 실패로 취급하면 유지보수가
끝나면 성공했을 메시지를 폐기하게 된다.
상세 규약은 DEALER spec §3 Outbound peer 선택을 참고.
routing_id는 connect 전에 설정¶
zlink_set_routing_id()는 zlink_connect() 호출 전에 호출해야 한다. 연결 후 변경은 적용되지 않는다.
언어별 완전한 예제¶
DEALER가 ROUTER로 보내고 응답을 받는 자립형 예제다(모든 바인딩, 빌드·실행 검증됨).
zlink::context_t ctx;
zlink::router_socket_t router (ctx);
zlink::dealer_socket_t dealer (ctx);
zlink::socket_monitor_t router_monitor = router.monitor_open ();
zlink::socket_monitor_t dealer_monitor = dealer.monitor_open ();
router.bind ("tcp://127.0.0.1:0");
const std::string endpoint = router.options ().last_endpoint ();
assert (!endpoint.empty ());
dealer.connect (endpoint);
assert (detail::wait_connected (router_monitor, dealer_monitor, 2000, &router));
const std::string sent = detail::k_dealer_router_request;
zlink::message_t outbound = detail::make_message (sent);
dealer.send ().message (std::move (outbound)).submit ();
zlink::received_t inbound;
assert (router.recv (inbound) == 0);
assert (inbound.routing_id ().has_value ());
assert (inbound.parts ().size () == 1);
assert (inbound.parts ()[0].to_string () == detail::k_dealer_router_request);
const std::string reply_payload = detail::k_dealer_router_reply;
zlink::message_t reply = detail::make_message (reply_payload);
// Reply on the ROUTER socket, addressed by the received envelope's routing id.
router.send (*inbound.routing_id ()).message (std::move (reply)).submit ();
inbound.close ();
zlink::received_t echoed;
assert (dealer.recv (echoed) == 0);
assert (echoed.parts ().size () == 1);
const std::string received = echoed.parts ()[0].to_string ();
assert (received == detail::k_dealer_router_reply);
echoed.close ();
std::printf ("[dealer-router/recv] send: \"%s\" → recv: \"%s\"\n", sent.c_str (),
received.c_str ());
return 0;
if (!SampleSupport.IsNativeAvailable())
return;
using var ctx = Zlink.CreateContext();
using var dealer = ctx.CreateDealerSocket();
using var router = ctx.CreateRouterSocket();
string endpoint = SampleSupport.NewEndpoint("tcp", "sample");
using var dealerMonitor = dealer.MonitorOpen(SocketEvent.ConnectionReady);
using var routerMonitor = router.MonitorOpen(SocketEvent.ConnectionReady);
router.Bind(endpoint);
dealer.Connect(endpoint);
SampleSupport.WaitConnected(routerMonitor, dealerMonitor);
using (Message request = Message.From("ping"))
await dealer.Send().Message(request).Async().Admitted;
using var received = Received.Create();
if (!router.Recv(received))
throw new InvalidOperationException("recv failed");
string requestPayload = received.Parts[0].GetString();
SampleSupport.EnsureEqual("ping", requestPayload, "request");
using var reply = Message.From("pong");
received.Send().Message(reply).Submit();
string payload = SampleSupport.ReceiveUtf8(dealer, 2000);
Console.WriteLine($"[dealer-router/recv] send: \"ping\" -> recv: \"{payload}\"");
SampleSupport.ensureNative();
String endpoint = SampleSupport.tcpEndpoint();
try (Context ctx = Zlink.createContext();
RouterSocket router = ctx.createRouterSocket();
DealerSocket dealer = ctx.createDealerSocket();
var routerMonitor = router.monitorOpen(
systems.zlink.contracts.eventing.MonitorEventType.CONNECTION_READY);
var dealerMonitor = dealer.monitorOpen(
systems.zlink.contracts.eventing.MonitorEventType.CONNECTION_READY)) {
router.bind(endpoint);
dealer.connect(endpoint);
SampleSupport.waitConnected(routerMonitor, dealerMonitor);
try (Message request = Message.from(SampleSupport.DEALER_REQUEST)) {
dealer.send().message(request).submit().admitted().toCompletableFuture().join();
}
try (systems.zlink.contracts.messaging.Received received = new systems.zlink.contracts.messaging.Received()) {
router.recv(received, systems.zlink.contracts.sockets.RecvFlags.NONE);
String value = SampleSupport.singleUtf8(received);
if (!SampleSupport.DEALER_REQUEST.equals(value)) {
throw new IllegalStateException("unexpected request: " + value);
}
try (Message reply = Message.from(SampleSupport.DEALER_REPLY)) {
received.send().message(reply).submit();
}
}
try (systems.zlink.contracts.messaging.Received received = new systems.zlink.contracts.messaging.Received()) {
dealer.recv(received, systems.zlink.contracts.sockets.RecvFlags.NONE);
String value = SampleSupport.singleUtf8(received);
if (!SampleSupport.DEALER_REPLY.equals(value)) {
throw new IllegalStateException("unexpected reply: " + value);
}
System.out.println("[dealer-router/recv] send: \""
+ SampleSupport.DEALER_REQUEST + "\" \u2192 recv: \"" + value + "\"");
}
}
SampleSupport.ensureNative()
val endpoint = SampleSupport.tcpEndpoint()
Zlink.createContext().use { ctx ->
ctx.createRouterSocket().use { router ->
ctx.createDealerSocket().use { dealer ->
router.monitorOpen(MonitorEventType.CONNECTION_READY).use { routerMonitor ->
dealer.monitorOpen(MonitorEventType.CONNECTION_READY).use { dealerMonitor ->
router.bind(endpoint)
dealer.connect(endpoint)
SampleSupport.waitConnected(routerMonitor, dealerMonitor)
Message.from(SampleSupport.DEALER_REQUEST).use { request ->
dealer.send().message(request).submit().await()
}
Received().use { received ->
router.recv(received, RecvFlags.NONE)
val value = SampleSupport.singleUtf8(received)
check(SampleSupport.DEALER_REQUEST == value) { "unexpected request: $value" }
Message.from(SampleSupport.DEALER_REPLY).use { reply ->
received.send().message(reply).submit()
}
}
Received().use { received ->
dealer.recv(received, RecvFlags.NONE)
val value = SampleSupport.singleUtf8(received)
check(SampleSupport.DEALER_REPLY == value) { "unexpected reply: $value" }
println(
"[dealer-router/recv] send: \"${SampleSupport.DEALER_REQUEST}\"" +
" → recv: \"$value\""
)
}
}
}
}
}
}
_, endpoint = tcp_endpoint()
with zlink.create_context() as ctx:
with zlink.create_router_socket(ctx) as router:
with zlink.create_dealer_socket(ctx) as dealer:
with router.monitor_open(zlink.MonitorEventMask.CONNECTION_READY) as router_monitor:
with dealer.monitor_open(zlink.MonitorEventMask.CONNECTION_READY) as dealer_monitor:
dealer.set_routing_id(b"CLIENT")
router.bind(endpoint)
dealer.connect(endpoint)
wait_connected(router_monitor, dealer_monitor)
send = dealer.send().message(b"ping").submit()
await send.admitted
request = zlink.create_received()
if not router.recv_into(request):
raise AssertionError("expected dealer-router request")
with request:
if request.routing_id != zlink.RoutingId(b"CLIENT"):
raise AssertionError(f"unexpected routing id: {request.routing_id!r}")
if request.to_bytes_list() != [b"ping"]:
raise AssertionError("unexpected dealer-router request payload")
send = request.send().message(b"pong").submit()
await send.admitted
reply = zlink.create_received()
if not dealer.recv_into(reply):
raise AssertionError("expected dealer-router reply")
with reply:
if reply.to_bytes_list() != [b"pong"]:
raise AssertionError("unexpected dealer-router reply payload")
print('[dealer-router/recv] send: "ping" → recv: "pong"')
const endpoint = await tcpEndpoint();
const ctx = zlink.createContext();
const router = zlink.createRouterSocket(ctx);
const dealer = zlink.createDealerSocket(ctx);
try {
const routerMonitor = router.monitorOpen([zlink.MonitorEventType.ConnectionReady]);
const dealerMonitor = dealer.monitorOpen([zlink.MonitorEventType.ConnectionReady]);
try {
router.bind(endpoint);
dealer.connect(endpoint);
await waitForConnectionReady(routerMonitor, zlink);
await waitForConnectionReady(dealerMonitor, zlink);
} finally {
routerMonitor.close();
dealerMonitor.close();
}
const sent = 'ping';
await dealer.send().message(Buffer.from(sent)).submit().admitted;
const reply = 'pong';
const request = new zlink.Received();
router.recv(request);
try {
const recvReq = request.parts[0].data().toString();
assert.equal(recvReq, sent);
assert.ok(request.routingId instanceof zlink.RoutingId);
request.send().message(Buffer.from(reply)).submit();
} finally {
request.close();
}
const response = new zlink.Received();
dealer.recv(response);
try {
const recv = response.parts[0].data().toString();
assert.equal(recv, reply);
console.log(`[dealer-router/recv] send: "${sent}" \u2192 recv: "${recv}"`);
} finally {
response.close();
}
} finally {
dealer.close();
router.close();
ctx.close();
}
const port = await reservePort();
const endpoint = `tcp://127.0.0.1:${port}`;
const ctx = zlink.createContext();
const router = zlink.createRouterSocket(ctx);
const dealer = zlink.createDealerSocket(ctx);
try {
const routerMonitor = router.monitorOpen([zlink.MonitorEventType.ConnectionReady]);
const dealerMonitor = dealer.monitorOpen([zlink.MonitorEventType.ConnectionReady]);
try {
router.bind(endpoint);
dealer.connect(endpoint);
await waitForConnectionReady(routerMonitor);
await waitForConnectionReady(dealerMonitor);
} finally {
routerMonitor.close();
dealerMonitor.close();
}
const sent = 'ping';
dealer.send().message(Buffer.from(sent)).submit();
const reply = 'pong';
const request = new zlink.Received();
router.recv(request);
try {
const recvReq = request.parts[0].data().toString();
assert.equal(recvReq, sent);
assert.ok(request.routingId instanceof zlink.RoutingId);
request.send().message(Buffer.from(reply)).submit();
} finally {
request.close();
}
const response = new zlink.Received();
dealer.recv(response);
try {
const recv = response.parts[0].data().toString();
assert.equal(recv, reply);
console.log(`[dealer-router/recv] send: "${sent}" → recv: "${recv}"`);
} finally {
response.close();
}
} finally {
dealer.close();
router.close();
ctx.close();
}
ctx, err := zlink.NewContext()
samplecommon.Must(err)
defer ctx.Close()
router, err := ctx.RouterSocket()
samplecommon.Must(err)
defer router.Close()
dealer, err := ctx.DealerSocket()
samplecommon.Must(err)
defer dealer.Close()
routerMon := samplecommon.OpenMonitor(router)
defer routerMon.Close()
dealerMon := samplecommon.OpenMonitor(dealer)
defer dealerMon.Close()
endpoint := samplecommon.UniqueTCP("dealer-router-recv")
rid := zlink.NewRoutingID([]byte("dealer-sample"))
samplecommon.Must(router.Bind(endpoint))
samplecommon.Must(dealer.SetRoutingID(rid))
samplecommon.Must(dealer.Connect(endpoint))
samplecommon.WaitConnected(routerMon, dealerMon)
submission, err := dealer.Send().Message(
samplecommon.Message("ping")).Submit(context.Background())
samplecommon.Must(err)
samplecommon.Must(submission.Admitted(context.Background()))
var request zlink.Received
_, err = router.Recv(&request, zlink.RecvFlagsNone)
samplecommon.Must(err)
defer request.Close()
submission, err = request.Send().Message(samplecommon.Message("pong")).Submit(context.Background())
if err == nil {
err = submission.Admitted(context.Background())
}
samplecommon.Must(err)
var reply zlink.Received
_, err = dealer.Recv(&reply, zlink.RecvFlagsNone)
samplecommon.Must(err)
defer reply.Close()
part, err := reply.SinglePartOrError()
samplecommon.Must(err)
if !bytes.Equal(part.Data(), []byte("pong")) {
samplecommon.Must(fmt.Errorf("unexpected reply %q", string(part.Data())))
}
fmt.Printf("[dealer-router/recv] send: %q -> recv: %q\n", "ping", string(part.Data()))
let ctx = Context::new().expect("context creation failed");
let endpoint = sample_support::tcp_endpoint();
let router = ctx.router_socket().expect("router socket failed");
let dealer = ctx.dealer_socket().expect("dealer socket failed");
let rid = RoutingId::from(b"dealer-node-7");
dealer.set_routing_id(&rid).expect("set routing id failed");
let router_mon = SocketMonitor::open(&router).expect("router monitor open failed");
let dealer_mon = SocketMonitor::open(&dealer).expect("dealer monitor open failed");
router.bind(&endpoint).expect("bind failed");
dealer.connect(&endpoint).expect("connect failed");
sample_support::wait_connected(&[&router_mon, &dealer_mon]);
drop(router_mon);
drop(dealer_mon);
let req = Message::try_from(b"ping").expect("message failed");
let submission = dealer.send().message(req).submit().expect("send failed");
sample_support::block_on(submission.admitted).expect("send admission failed");
let mut received = zlink::Received::empty();
router
.recv(&mut received, zlink::RecvFlags::NONE)
.expect("router recv failed");
assert!(received.routing_id().is_some());
assert_eq!(received.parts()[0].as_str().unwrap(), "ping");
let resp = Message::try_from(b"pong").expect("message failed");
let submission = received
.send()
.message(resp)
.submit()
.expect("received send failed");
sample_support::block_on(submission.admitted).expect("received send admission failed");
let mut response = zlink::Received::empty();
dealer
.recv(&mut response, zlink::RecvFlags::NONE)
.expect("dealer recv failed");
assert_eq!(response.parts()[0].as_str().unwrap(), "pong");
println!(
"[dealer-router/recv] send: \"ping\" → recv: \"{}\"",
response.parts()[0].as_str().unwrap()
);