콘텐츠로 이동

한국어 | 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다(SendOpSendSubmitOp, RequestOpRequestSubmitOpReplyOpReplySubmitOp) — 단계에 맞는 메서드를 호출하면 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 / .BytesRequestSubmitOp 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.SubmitReplySubmitOp.Submiterror를 반환한다. 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 바인딩 스펙에서 전체 근거를 확인한다.