한국어 | English
02. Messaging¶
이 category는 Message(frame 소유), 수신 envelope 타입(Received,
TopicMessage, SubscriptionEvent), 그리고 모든 socket type의 진입점이
반환하는 공유 send/request/reply operation-builder interface를 다룬다.
RoutingID는 messaging 전용이 아니라 package 전역 value type이므로 Core
category에서 다룬다. 정확한 signature는
internal/native/message.go,
received.go,
topic_message.go,
subscription_event.go,
operations.go가
소유하며,
contracts/messaging.go를
통해 alias로 re-export된다.
Message¶
하나의 native zlink frame을 소유한다 — 모든 send/request/reply/receive API가 옮기는 단위.
empty, err := contracts.NewMessage(nil)
sized, err := contracts.NewMessageWithSize(4096)
fromBytes, err := contracts.NewMessage(payload)
fromString, err := contracts.NewMessageString("payload")
Options.
| Member | 의미 |
|---|---|
NewMessage(data []byte) |
data를 복사; nil/빈 슬라이스는 길이 0 메시지를 만든다 |
NewMessageWithSize(size int) |
쓰기 가능한, 0으로 초기화된 storage |
NewMessageString(value string) |
UTF-8 인코딩 후 복사 |
Data() []byte |
Close까지만 유효한 view — 여러 goroutine에 걸쳐 또는 메시지 수명을 넘어 보관하지 말 것, 그럴 땐 Bytes()를 쓴다 |
Bytes() []byte |
payload의 스냅샷 복사 |
Size() |
payload 바이트 길이 |
IsEmpty() |
Size()가 0인지 |
Text() / String() |
둘 다 Data()를 UTF-8로 디코딩 |
CopyTo(destination []byte) (int, error) |
payload를 caller가 제공한 slice로 복사; destination이 너무 작으면 error |
TryCopyTo(destination []byte) bool |
error를 버리는 CopyTo의 non-erroring 버전 |
RefCount() int |
진단 전용; 실패 시 error를 전파하는 대신 0을 반환 |
Clone() / Copy() |
동등 — 둘 다 (*Message, error), 독립된 payload 복사본 |
Close() error |
message를 해제 |
Completion result. 모든 member는 동기다. Close()는 멱등이다(nil
receiver나 이미 닫힌 메시지는 nil을 반환); finalizer가 없으므로
caller가 명시적으로 호출해야 한다. MoveMessage 계열 builder 호출로
메시지를 보내면 native frame이 socket으로 이전되고 Go 쪽에서 closed로
표시된다(아래 operation-builder 형태 참고); 일반 Message 계열
호출로 보내면 caller의 Message는 여전히 소유 상태고 여전히
Close()가 필요하다.
선택 기준. outbound payload를 만들 땐 NewMessageWithSize/
NewMessage를 쓴다; MoveMessage로 소유권을 이전하는 대신 독립
복사본이 필요할 땐 Clone()을 쓴다. 다음 recv 호출 전에 슬라이스를
다 쓰는 hot path에선 Bytes()보다 Data()를 우선한다 —
Bytes()는 항상 복사를 할당한다.
Received¶
수신 메시지 envelope: routing metadata, ROUTER request의 opaque reply token,
message part. recv 호출마다
새로 만드는 대신 인스턴스 하나를 재사용해 매 수신당 할당을 피한다 —
socket의 receive 메서드가 이를 리셋하고 다시 채운다.
var received contracts.Received
ok, err := router.Recv(&received, contracts.RecvFlagsNone)
if _, replyable := received.ReplyToken(); ok && replyable {
received.Reply().Message(reply).Submit(ctx)
}
Options. 공개 생성자 없음 — caller가 zero-value Received{}를
선언하고 그 주소를 socket의 receive 메서드(Sockets category)에
넘기면 그것이 채운다.
| Member | 의미 |
|---|---|
RoutingID() RoutingID / HasRoutingID() bool |
peer routing id와 존재 여부 |
ReplyToken() (ReplyToken, bool) |
opaque reply capability와 존재 여부 |
IsSinglePart() bool |
Parts()가 정확히 하나인지 |
Parts() []*Message |
이 envelope이 담은 모든 message part |
FirstPart() (*Message, error) |
소유권 이전 없이 첫 part — part는 여전히 Received가 소유 |
Reply() ReplyOp |
공유 reply builder 시작; ReplyToken()이 있을 때만 유효 |
Send() SendOp |
공유 send builder를 시작, 이 envelope이 캡처한 source route로 향함 |
Close() error |
보관 중인 모든 part를 닫음; nil receiver에서 호출해도 안전 |
Completion result. 모든 member는 동기다. Close()(그리고 receive
메서드 자신의 reset 단계)는 이전에 보관 중이던 모든 part를 버리기
전에 닫는다 — 닫지 않고 Received를 재사용해도 native frame이 새는
일은 없다.
선택 기준. 수신 loop에서 매 메시지당 새로 할당하는 대신 Received
하나를 참조로 재사용한다. 목적지 route를 손으로 재구성하는 대신
Reply()를 쓴다 — routing id와 opaque reply token은 캡슐화돼 있어 이
builder 밖에선 접근할 수 없다.
TopicMessage¶
수신된 publish: topic, source routing id, message part. Received와
같은 방식으로 subscribe-receive 호출마다 인스턴스 하나를 재사용한다.
var published contracts.TopicMessage
ok, err := sub.Subscribe(&published, contracts.RecvFlagsNone)
topic := published.Topic()
Options. 공개 생성자 없음 — zero-value TopicMessage{}를 선언한다.
| Member | 의미 |
|---|---|
RoutingID() RoutingID / HasRoutingID() bool |
publisher의 routing id와 존재 여부 |
Topic() string |
이 publish가 전송된 topic |
IsSinglePart() bool / Parts() []*Message / FirstPart() (*Message, error) / Close() error |
Received와 같은 형태 |
Completion result. 동기다; Received와 동일한 close-후-재사용
계약.
선택 기준. Received와 같은 방식으로 subscribe-receive loop에서
인스턴스 하나를 재사용한다.
SubscriptionEvent¶
XPUB socket이 관찰한, 한 subscriber의 subscribe 또는 unsubscribe를 보고한다.
var evt contracts.SubscriptionEvent
ok, err := xpub.ReceiveSubscriptionEvent(&evt, contracts.RecvFlagsNone)
Options. 공개 생성자 없음 — zero-value SubscriptionEvent{}를
선언한다. 모든 member는 value receiver다, pointer가 아니다.
| Member | 의미 |
|---|---|
RoutingID() RoutingID / HasRoutingID() bool |
구독자의 routing id와 존재 여부 |
Subscribed() bool |
subscribe면 true, unsubscribe면 false |
Topic() string |
subscribe/unsubscribe된 topic |
Completion result. 동기다; Close() 없음 — 이 타입은 자신만의
native resource를 소유하지 않는다(message part를 소유하는
Received/TopicMessage와 다름).
선택 기준. subscriber 이탈을 관찰하려고 XPUB socket의 subscription-event receive 경로(Sockets category)에서 쓴다.
Send / request / reply operation-builder 형태¶
모든 socket type의 Send/Publish/Request/Reply 진입점(Sockets
category)이 반환하는, part와 terminal submit을 누적하는 fluent
builder interface. 각 단계는 별개의 Go interface다(SendOp →
SendSubmitOp, RequestOp → RequestSubmitOp →
ReplyOp → ReplySubmitOp) — 단계에
맞는 메서드를 호출하면 chain의 다음 interface를 반환하므로, caller는
part를 하나라도 추가하기 전엔 Submit을 호출할 수 없다.
err := dealer.Send().Message(part1).Message(part2).Submit(ctx)
parts, err := dealer.Request().
Message(payload).
Timeout(5 * time.Second).
Submit(ctx)
err := received.Reply().Message(reply).Submit(ctx)
Options. 모든 terminal Submit은 첫 인자로
context.Context를 받는다** — 취소되거나 deadline이 지난 context는
native submit이 실행되기 전에 그 context의 Err()로 호출을 즉시
중단시킨다, 다른 어떤 언어의 operation builder도 이러지 않는다.
| Stage | Member | 의미 |
|---|---|---|
SendOp |
.Message(*Message) / .MoveMessage(*Message) / .Bytes([]byte) → SendSubmitOp |
chain을 시작; MoveMessage는 submit이 성공하면 메시지 소유권을 socket으로 이전 |
SendSubmitOp |
같은 세 추가 메서드 + .Submit(ctx context.Context) error |
part 추가, completion-backed terminal |
RequestOp |
.Message / .Bytes → RequestSubmitOp |
request chain을 시작 |
RequestSubmitOp |
.Timeout(time.Duration) |
reply-wait timeout을 더함 |
RequestSubmitOp terminal |
.Submit(ctx context.Context) ([]*Message, error) |
completion-backed reply part를 직접 반환 |
context.Context |
send/request/reply Submit에 전달 |
Go waiter를 제한하며 successful native submit 뒤 Core cancel은 아님 |
ReplyOp |
.Message(*Message) → ReplySubmitOp |
reply chain을 시작 |
ReplySubmitOp |
.Message(...) / .Submit(ctx context.Context) error |
flag-free synchronous reply terminal |
Completion result. SendSubmitOp.Submit과 ReplySubmitOp.Submit은 error를
반환한다. RequestSubmitOp.Submit은 socket completion queue가 terminal을 낸 뒤
([]*Message, error)를 반환하고 caller가 모든 reply message를 닫는다. 모든 builder는
1회용이며 두 번째 Submit은 재제출하지 않고 error를 반환한다.
선택 기준. owning goroutine에서 Submit(ctx)의 직접 결과를 처리한다. 목적지
route를 손으로 재구성하는 대신
Received.Reply()/Send()를 쓴다 — routing id(그리고 Reply의
경우 reply token)는 builder에 캡슐화된다.
internal/native/message.go,
received.go,
operations.go,
Go 바인딩 스펙에서 전체 근거를 확인한다.