가이드 목록 | 이전: DEALER | 다음: STREAM
ROUTER 소켓¶
이 장의 계약 소유 문서 — ROUTER socket 스펙이 다룬다. 이 챕터는 그 계약을 언어별 예제로 보여준다.
1. 개요¶
ROUTER는 하나의 socket에서 여러 peer와의 연결(pipe)을 관리하는 비동기 raw socket이다. 모든 수신 메시지에는 송신자의 routing id가 함께 오고, 모든 송신 메시지는 target routing id를 지정해야 한다. 하나의 socket이 여러 DEALER 또는 ROUTER peer를 (DEALER처럼 round-robin이 아니라) 개별적으로 지정해 통신해야 할 때 사용한다.
핵심 특성:
- 수신: 모든 record가 송신자의 routing id와 불투명 reply token을 함께 반환
- 송신: directed 전용 — 호출자가 routing id로 peer를 지정
- 하나의 socket에 두 가지 트래픽 형태 공존: 일반 DATA(reply token 0)와
reply를 기대하는 REQUEST record(0이 아닌 token)
유효한 소켓 조합: ROUTER ↔ DEALER, ROUTER ↔ ROUTER
%%{init: {'flowchart': {'nodeSpacing': 32, 'rankSpacing': 40, 'padding': 8, 'wrappingWidth': 180}, 'themeVariables': {'fontSize': '18px'}}}%%
flowchart LR
R[ROUTER] -->|routing id 지정| D1[DEALER 1]
R -->|routing id 지정| D2[DEALER 2]
D1 -->|fair-queue| R
D2 -->|fair-queue| R
2. 기본 사용법¶
생성 및 바인딩¶
메시지 수신¶
zlink_router_recv()는 payload record 전체를 caller가 제공한 배열에 반환한다. Routing id view는 같은
socket에서 다음 data receive 함수에 진입하기 전까지만 유효하다. 그 이후에도 사용해야
하면 복사한다.
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_recv_result_t rc = zlink_router_recv(
router, &source_rid, &reply_token, parts, 8, &part_count, ZLINK_RECV_FLAGS_NONE);
if (rc == ZLINK_RECV_OK) {
/* source_rid는 peer를 식별한다. 배열 전체를 처리한 뒤 닫는다. */
zlink_multipart_close(parts, part_count);
}
/* 그 밖의 rc 값: ZLINK_RECV_NO_DATA (EAGAIN), TERMINATED, INVALID_HANDLE */
일반 routed DATA에서는 reply_token이 0이다. 0이 아닌 token은 zlink_send_rid()가
아니라 zlink_reply()(§4 참고)로 응답해야 하는 REQUEST다.
Application은 token을 해석하지 않는다.
Routed message 송신¶
zlink_send_rid()는 target_rid_가 지정하는 peer에게 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_rid(
router, source_rid, parts, 2, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
3. 옵션¶
| 옵션 | 타입 | 기본값 | 설명 |
|---|---|---|---|
ZLINK_ROUTER_OPT_MANDATORY |
int | 1 |
0=off, 양수=on. on이면 연결되지 않은 routing id로의 directed submit이 조용히 버려지지 않고 ZLINK_SUBMIT_NOT_CONNECTED로 실패 |
ZLINK_ROUTER_OPT_PROBE |
int | 0 |
0=off, 양수=on. 연결 설정 시 빈 raw message를 보내 peer가 연결과 이 ROUTER의 routing id를 관찰하게 함 |
ZLINK_ROUTER_OPT_CONNECT_ROUTING_ID |
binary, set 전용 | — | 다음 zlink_connect()로 만들 pipe의 local alias. 각 connect 전에 설정 |
ZLINK_ROUTER_OPT_REQUEST_TIMEOUT_MS |
int (ms) | 5000 |
request의 timeout_ms_ == 0일 때 사용하는 기본 timeout |
ZLINK_ROUTER_OPT_WEIGHT |
int | 100, 범위 0..10000 |
연결된 peer에 알리는 이 ROUTER의 가중치 |
ZLINK_OPT_SNDHWM |
uint64_t bytes |
자동 | 수동 설정이 우선하며 0은 무제한 |
ZLINK_OPT_RCVHWM |
uint64_t bytes |
자동 | 수동 설정이 우선하며 0은 무제한 |
ZLINK_OPT_LINGER |
int | -1 |
close 시 대기 시간 (ms) |
ZLINK_OPT_SNDTIMEO |
int | 1000 |
송신 타임아웃(ms). 무한 대기는 -1을 명시적으로 설정 |
ZLINK_OPT_RCVTIMEO |
int | 1000 |
수신 타임아웃(ms). 무한 대기는 -1을 명시적으로 설정 |
ROUTER 전용 옵션은 typed accessor로 설정·조회한다.
ZLINK_EXPORT zlink_config_result_t zlink_set_router_option(
void *handle_, zlink_router_option_t option_, const void *optval_, size_t optvallen_);
ZLINK_EXPORT zlink_config_result_t zlink_get_router_option(
void *handle_, zlink_router_option_t option_, void *optval_, size_t *optvallen_);
zlink_get_router_option()을 호출할 때 *optvallen_은 optval_의 입력 용량이다. 성공하면
실제로 쓴 byte 수로 갱신된다.
ZLINK_ROUTER_OPT_MANDATORY¶
void *router = zlink_socket(ctx, ZLINK_SOCKET_ROUTER);
int mandatory = 1;
zlink_set_router_option(router, ZLINK_ROUTER_OPT_MANDATORY, &mandatory, sizeof(mandatory));
/* target_rid가 연결된 pipe가 없는 routing id를 가리킴 */
zlink_msg_t part;
zlink_msg_init_size(&part, 4);
memcpy(zlink_msg_data(&part), "data", 4);
zlink_submit_result_t rc = zlink_send_rid(
router, target_rid, &part, 1, ZLINK_SEND_FLAGS_NONE, NULL, NULL);
/* MANDATORY가 켜져 있으므로 rc == ZLINK_SUBMIT_NOT_CONNECTED */
참고:
core/tests/integration/test_router_mandatory.cpp
ZLINK_ROUTER_OPT_CONNECT_ROUTING_ID¶
ROUTER가 inbound 연결만 받지 않고 peer로 직접 connect할 때, 각 zlink_connect() 호출
전에 설정하면 그 호출이 만드는 pipe의 local alias를 고를 수 있다.
zlink_set_router_option(
router, ZLINK_ROUTER_OPT_CONNECT_ROUTING_ID, "peer-a", 6);
zlink_connect(router, "tcp://127.0.0.1:5559");
4. Request와 reply¶
zlink_request()는 routed request를 제출하고 0이 아닌 completion ID를 반환한다. Reply
또는 terminal 결과는 일반 DATA receive가 아니라 zlink_completion_recv()로 pull한다.
수신한 REQUEST(0이 아닌 reply token)는 receive 결과가 반환한 source RID와 token을 사용해
zlink_reply()로 응답한다.
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(
router, peer_rid, &req, 1, ZLINK_SEND_FLAGS_NONE,
0 /* ZLINK_ROUTER_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(router, &completion, ZLINK_RECV_FLAGS_NONE)
== ZLINK_RECV_OK)
zlink_completion_close(&completion);
}
수신 측은 receive 결과가 반환한 routing id와 불투명 token으로 응답한다.
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);
if (reply_token != 0) {
/* 이 record는 directed send가 아니라 reply를 기대한다. */
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);
DEALER peer로 보내는 reply는 DEALER-ROUTER Application connection의 FIFO, HWM과 PAUSED
state를 공유하므로 ZLINK_SUBMIT_BACKPRESSURED가 될 수 있다. ROUTER peer로 보내는 reply는
ROUTER-ROUTER Completion lane을 사용한다. 성공한 reply 제출만 token을 소비하며 request
lifecycle이 유효하면 실패한 완전한 시도를 재시도할 수 있다.
참고:
core/tests/integration/test_zmp_request_reply.cpp,core/tests/integration/test_zmp_request_reply_router_recv_surface.cpp
5. 사용 패턴¶
패턴 1: ROUTER ← 여러 DEALER¶
가장 흔한 형태다. 각 DEALER가 routing id를 갖고 연결하면 ROUTER는 source_rid로 송신자를
구분하고 같은 id로 응답한다.
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);
/* router의 recv는 source_rid로 "D1"과 "D2"를 구분하고,
zlink_send_rid(router, source_rid, ...)로 해당 peer에만 응답한다. */
참고:
core/tests/integration/test_router_multiple_dealers.cpp
패턴 2: 상관관계가 있는 request-reply¶
호출자가 자유 형식 send/recv 대신 전달 확인과 상관된 응답이 필요할 때
zlink_request() / zlink_reply()(§4 참고)를 사용한다.
Completion ID가 origin 결과를 상관시키고, 불투명한 0이 아닌 reply token이 responder에게
REQUEST 하나를 응답할 권한을 준다. 일반 DATA의 token은 0이다.
패턴 3: MANDATORY로 도달 가능성 강제¶
기본값에서는 연결된 pipe가 없는 routing id로의 directed send가 오류 없이 버려진다.
ZLINK_ROUTER_OPT_MANDATORY를 설정하면 이를 ZLINK_SUBMIT_NOT_CONNECTED로 드러내
호출자가 메시지를 조용히 잃는 대신 오래된 routing id를 감지할 수 있다.
int mandatory = 1;
zlink_set_router_option(router, ZLINK_ROUTER_OPT_MANDATORY, &mandatory, sizeof(mandatory));
참고:
core/tests/integration/test_router_mandatory.cpp,core/tests/integration/test_router_mandatory_hwm.cpp
패턴 4: 프록시(ROUTER-DEALER)¶
ROUTER를 frontend로, DEALER를 backend로 써서 멀티스레드 서버를 구성한다. 전체 프록시
예제는 DEALER §5 패턴 3을 참고한다.
그 예제의 ROUTER 쪽은 frontend로 바인딩한 평범한 zlink_socket(ctx, ZLINK_SOCKET_ROUTER)다.
6. 주의사항¶
Routing ID 수명¶
zlink_router_recv()가 반환하는 source_rid는 socket-owned view다. 같은 socket의
다음 data receive 진입 전까지만 유효하며 성공 여부와 관계없이 무효화되므로, 그 이후에도
id가 필요하면 byte를 복사한다. Record 전체와 함께 routing id와 reply token 하나를 반환한다.
전체 수명·복사 규칙은 Routing ID를 참고한다.
peer 없음 vs HWM 배압¶
DEALER와 마찬가지로 둘은 별개의 결과다. ZLINK_ROUTER_OPT_MANDATORY가 켜져 있으면
연결되지 않은 routing id로의 송신은 ZLINK_SUBMIT_NOT_CONNECTED를 반환하며 아무것도
큐에 쌓이지 않는다. 연결된 peer의 큐가 HWM에 도달하면 대기(기본) 또는
ZLINK_SEND_FLAGS_DONTWAIT 시 ZLINK_SUBMIT_BACKPRESSURED를 반환한다.
Logical RID 지정¶
zlink_send_rid()와 zlink_request()는 logical routing id만 받는다. Physical pair
ID와 generation은 public send selector가 아니다. Core가 DONTWAIT record를 admission 전에
보관하면 transient reconnect 동안 같은 logical RID를 유지하고 completion ID로 terminal을
보고한다. Local admission 뒤에는 새 connection에 payload를 replay하지 않는다.
동시성¶
ROUTER의 public handle은 Thread Safety가 설명하는 계층적 동시성 계약을 따른다. send/publish 경로는 같은 handle의 동시 사용을 허용하지만, 옵션 변경과 close는 정확성을 위해 직렬화된다. 각 send 호출은 완성된 record 하나를 원자적으로 제출하므로 여러 thread가 같은 handle에 독립된 record를 제출할 수 있다.
언어별 완전한 예제¶
DEALER가 ROUTER로 보내고 응답을 받는 자립형 예제다(모든 바인딩, 빌드·실행 검증됨).
ROUTER 쪽 처리는 위 예제의 router.recv / router.send에 해당한다.
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()
);
Routing id 수명과 복사 규칙은 Routing ID, 같은 handle의 동시 사용 조건은 Thread Safety를 참고한다.