streambus
packageAPI reference for the streambus
package.
Imports
(12)BenchmarkPublishFanout
Parameters
func BenchmarkPublishFanout(b *testing.B)
{
for _, subscribers := range []int{1, 10, 100} {
b.Run(fmt.Sprintf("subscribers-%d", subscribers), func(b *testing.B) {
bus := NewInMemory(Config{DefaultBuffer: 1024, MaxBuffer: 1024})
defer bus.Close()
for i := 0; i < subscribers; i++ {
subscription, err := bus.Subscribe(context.Background(), SubscribeOptions{
Topic: "benchmark", Buffer: 1024, Overflow: LatestOnly,
})
if err != nil {
b.Fatal(err)
}
defer subscription.Close()
go func() {
for range subscription.Frames() {
}
}()
}
frame := Frame{Topic: "benchmark", Payload: []byte("payload")}
b.ReportAllocs()
b.ResetTimer()
for i := 0; i < b.N; i++ {
if _, err := bus.Publish(context.Background(), frame); err != nil {
b.Fatal(err)
}
}
})
}
}
prefixNode
type prefixNode struct
Fields
| Name | Type | Description |
|---|---|---|
| subs | map[uint64]*Subscription | |
| children | map[byte]*prefixNode |
newPrefixNode
Returns
func newPrefixNode() *prefixNode
{
return &prefixNode{
subs: make(map[uint64]*Subscription),
children: make(map[byte]*prefixNode),
}
}
replayBuffer
type replayBuffer struct
Methods
Parameters
func (*replayBuffer) add(frame Frame)
{
if len(r.frames) == 0 {
return
}
r.frames[r.next] = frame
r.next = (r.next + 1) % len(r.frames)
if r.len < len(r.frames) {
r.len++
}
}
Parameters
Returns
func (*replayBuffer) last(count int) []Frame
{
if count > r.len {
count = r.len
}
out := make([]Frame, 0, count)
start := (r.next - count + len(r.frames)) % len(r.frames)
for i := 0; i < count; i++ {
out = append(out, r.frames[(start+i)%len(r.frames)])
}
return out
}
Parameters
Returns
func (*replayBuffer) after(sequence uint64) ([]Frame, bool)
{
if r.len == 0 {
return nil, false
}
start := (r.next - r.len + len(r.frames)) % len(r.frames)
oldest := r.frames[start].Sequence
if r.len == len(r.frames) && sequence != 0 && sequence < oldest {
return nil, true
}
out := make([]Frame, 0, r.len)
for i := 0; i < r.len; i++ {
frame := r.frames[(start+i)%len(r.frames)]
if frame.Sequence > sequence {
out = append(out, frame)
}
}
return out, false
}
Fields
| Name | Type | Description |
|---|---|---|
| frames | []Frame | |
| next | int | |
| len | int |
newReplayBuffer
Parameters
Returns
func newReplayBuffer(capacity int) *replayBuffer
{
return &replayBuffer{frames: make([]Frame, capacity)}
}
InMemory
InMemory is a bounded, transport-independent StreamBus implementation.
type InMemory struct
Methods
Publish fans a frame out to all exact and prefix subscribers.
Parameters
Returns
func (*InMemory) Publish(ctx context.Context, frame Frame) (uint64, error)
{
if frame.Topic == "" {
return 0, ErrInvalidTopic
}
if len(frame.Payload) > b.config.MaxPayloadBytes {
return 0, ErrPayloadTooLarge
}
select {
case <-ctx.Done():
return 0, ctx.Err()
default:
}
if b.config.CopyPayload {
frame.Payload = append([]byte(nil), frame.Payload...)
}
// Sequence ownership stays with the bus. Transport clients cannot inject or
// reuse sequence numbers from another session.
frame.Sequence = b.seq.Add(1)
if frame.Timestamp.IsZero() {
frame.Timestamp = time.Now().UTC()
}
b.mu.Lock()
if b.closed {
b.mu.Unlock()
return 0, ErrClosed
}
targets := make(map[uint64]*Subscription)
for id, sub := range b.exact[frame.Topic] {
targets[id] = sub
}
node := b.prefix
for i := 0; i < len(frame.Topic); i++ {
next := node.children[frame.Topic[i]]
if next == nil {
break
}
node = next
for id, sub := range node.subs {
targets[id] = sub
}
}
if b.config.ReplayCapacity > 0 {
journal := b.replay[frame.Topic]
if journal == nil {
journal = newReplayBuffer(b.config.ReplayCapacity)
b.replay[frame.Topic] = journal
}
journal.add(frame)
}
b.mu.Unlock()
for _, sub := range targets {
if err := sub.enqueue(ctx, frame); err != nil && err != ErrSubscriptionDone {
return frame.Sequence, err
}
}
return frame.Sequence, nil
}
Subscribe creates a bounded exact-topic or prefix subscription.
Parameters
Returns
func (*InMemory) Subscribe(ctx context.Context, options SubscribeOptions) (*Subscription, error)
{
if options.Topic == "" {
return nil, ErrInvalidTopic
}
if options.Buffer == 0 {
options.Buffer = b.config.DefaultBuffer
}
if options.Buffer < 1 || options.Buffer > b.config.MaxBuffer {
return nil, ErrInvalidBuffer
}
if options.Replay < 0 {
options.Replay = 0
}
if options.Prefix {
options.Replay = 0
options.Since = 0
}
select {
case <-ctx.Done():
return nil, ctx.Err()
default:
}
id := b.subSeq.Add(1)
var sub *Subscription
sub = newSubscription(id, options, func() { b.remove(sub) })
b.mu.Lock()
if b.closed {
b.mu.Unlock()
_ = sub.Close()
return nil, ErrClosed
}
var replay []Frame
if journal := b.replay[options.Topic]; journal != nil {
if options.Since > 0 {
var gap bool
replay, gap = journal.after(options.Since)
if gap || len(replay) > options.Buffer {
b.mu.Unlock()
_ = sub.Close()
return nil, ErrReplayUnavailable
}
} else if options.Replay > 0 {
replay = journal.last(options.Replay)
}
} else if options.Since > 0 && b.config.ReplayCapacity == 0 {
b.mu.Unlock()
_ = sub.Close()
return nil, ErrReplayUnavailable
}
if len(replay) > 0 {
sub.preload(replay)
}
if options.Prefix {
node := b.prefix
for i := 0; i < len(options.Topic); i++ {
child := node.children[options.Topic[i]]
if child == nil {
child = newPrefixNode()
node.children[options.Topic[i]] = child
}
node = child
}
node.subs[id] = sub
} else {
if b.exact[options.Topic] == nil {
b.exact[options.Topic] = make(map[uint64]*Subscription)
}
b.exact[options.Topic][id] = sub
}
b.mu.Unlock()
go func() {
select {
case <-ctx.Done():
_ = sub.Close()
case <-sub.Done():
}
}()
return sub, nil
}
Parameters
func (*InMemory) remove(sub *Subscription)
{
b.mu.Lock()
defer b.mu.Unlock()
options := sub.options
if options.Prefix {
node := b.prefix
for i := 0; i < len(options.Topic); i++ {
node = node.children[options.Topic[i]]
if node == nil {
return
}
}
delete(node.subs, sub.id)
return
}
subs := b.exact[options.Topic]
delete(subs, sub.id)
if len(subs) == 0 {
delete(b.exact, options.Topic)
}
}
Close stops every subscription and rejects future operations.
Returns
func (*InMemory) Close() error
{
b.mu.Lock()
if b.closed {
b.mu.Unlock()
return nil
}
b.closed = true
unique := make(map[uint64]*Subscription)
for _, subscribers := range b.exact {
for id, sub := range subscribers {
unique[id] = sub
}
}
collectPrefixSubscriptions(b.prefix, unique)
b.exact = make(map[string]map[uint64]*Subscription)
b.prefix = newPrefixNode()
b.mu.Unlock()
for _, sub := range unique {
_ = sub.Close()
}
return nil
}
Fields
| Name | Type | Description |
|---|---|---|
| config | Config | |
| mu | sync.RWMutex | |
| exact | map[string]map[uint64]*Subscription | |
| prefix | *prefixNode | |
| replay | map[string]*replayBuffer | |
| seq | atomic.Uint64 | |
| subSeq | atomic.Uint64 | |
| closed | bool |
Uses
NewInMemory
NewInMemory creates a StreamBus using process-local memory.
func NewInMemory(config Config) *InMemory
{
return &InMemory{
config: config.normalized(),
exact: make(map[string]map[uint64]*Subscription),
prefix: newPrefixNode(),
replay: make(map[string]*replayBuffer),
}
}
Uses
collectPrefixSubscriptions
Parameters
func collectPrefixSubscriptions(node *prefixNode, dst map[uint64]*Subscription)
{
for id, sub := range node.subs {
dst[id] = sub
}
for _, child := range node.children {
collectPrefixSubscriptions(child, dst)
}
}
TestPublishExactAndPrefix
Parameters
func TestPublishExactAndPrefix(t *testing.T)
{
bus := NewInMemory(Config{})
t.Cleanup(func() { _ = bus.Close() })
exact := mustSubscribe(t, bus, SubscribeOptions{Topic: "ui:button", Buffer: 4})
prefix := mustSubscribe(t, bus, SubscribeOptions{Topic: "ui:", Prefix: true, Buffer: 4})
other := mustSubscribe(t, bus, SubscribeOptions{Topic: "data:", Prefix: true, Buffer: 4})
sequence, err := bus.Publish(context.Background(), Frame{Topic: "ui:button", Payload: []byte("clicked")})
if err != nil {
t.Fatal(err)
}
if sequence == 0 {
t.Fatal("expected an assigned sequence")
}
for name, subscription := range map[string]*Subscription{"exact": exact, "prefix": prefix} {
frame := receiveFrame(t, subscription)
if frame.Sequence != sequence || string(frame.Payload) != "clicked" {
t.Fatalf("%s received %#v", name, frame)
}
}
select {
case frame := <-other.Frames():
t.Fatalf("unrelated prefix received %#v", frame)
case <-time.After(20 * time.Millisecond):
}
}
TestReplayPreservesOrder
Parameters
func TestReplayPreservesOrder(t *testing.T)
{
bus := NewInMemory(Config{ReplayCapacity: 4})
t.Cleanup(func() { _ = bus.Close() })
for _, payload := range []string{"one", "two", "three"} {
if _, err := bus.Publish(context.Background(), Frame{Topic: "history", Payload: []byte(payload)}); err != nil {
t.Fatal(err)
}
}
subscription := mustSubscribe(t, bus, SubscribeOptions{Topic: "history", Buffer: 4, Replay: 2})
if got := string(receiveFrame(t, subscription).Payload); got != "two" {
t.Fatalf("first replay frame = %q", got)
}
if got := string(receiveFrame(t, subscription).Payload); got != "three" {
t.Fatalf("second replay frame = %q", got)
}
}
TestResumeFromSequence
Parameters
func TestResumeFromSequence(t *testing.T)
{
bus := NewInMemory(Config{ReplayCapacity: 4})
t.Cleanup(func() { _ = bus.Close() })
sequences := make([]uint64, 0, 3)
for _, payload := range []string{"one", "two", "three"} {
sequence, err := bus.Publish(context.Background(), Frame{Topic: "resume", Payload: []byte(payload)})
if err != nil {
t.Fatal(err)
}
sequences = append(sequences, sequence)
}
subscription := mustSubscribe(t, bus, SubscribeOptions{Topic: "resume", Buffer: 4, Since: sequences[0]})
if got := string(receiveFrame(t, subscription).Payload); got != "two" {
t.Fatalf("first resumed frame = %q", got)
}
if got := string(receiveFrame(t, subscription).Payload); got != "three" {
t.Fatalf("second resumed frame = %q", got)
}
}
TestResumeReportsKnownGap
Parameters
func TestResumeReportsKnownGap(t *testing.T)
{
bus := NewInMemory(Config{ReplayCapacity: 2})
t.Cleanup(func() { _ = bus.Close() })
first, err := bus.Publish(context.Background(), Frame{Topic: "resume", Payload: []byte("one")})
if err != nil {
t.Fatal(err)
}
publishString(t, bus, "resume", "two")
publishString(t, bus, "resume", "three")
if _, err := bus.Subscribe(context.Background(), SubscribeOptions{Topic: "resume", Buffer: 2, Since: first}); !errors.Is(err, ErrReplayUnavailable) {
t.Fatalf("resume gap = %v", err)
}
}
TestLatestOnlyCollapsesQueuedFrames
Parameters
func TestLatestOnlyCollapsesQueuedFrames(t *testing.T)
{
bus := NewInMemory(Config{})
t.Cleanup(func() { _ = bus.Close() })
subscription := mustSubscribe(t, bus, SubscribeOptions{Topic: "cursor", Buffer: 2, Overflow: LatestOnly})
publishString(t, bus, "cursor", "one")
waitFor(t, func() bool { return subscription.Stats().Queued == 0 })
for _, value := range []string{"two", "three", "four"} {
publishString(t, bus, "cursor", value)
}
if got := string(receiveFrame(t, subscription).Payload); got != "one" {
t.Fatalf("in-flight frame = %q", got)
}
if got := string(receiveFrame(t, subscription).Payload); got != "four" {
t.Fatalf("latest frame = %q", got)
}
if subscription.Stats().Dropped == 0 {
t.Fatal("expected collapsed frames to be counted")
}
}
TestDropNewestPreservesQueuedFrame
Parameters
func TestDropNewestPreservesQueuedFrame(t *testing.T)
{
bus := NewInMemory(Config{})
t.Cleanup(func() { _ = bus.Close() })
subscription := mustSubscribe(t, bus, SubscribeOptions{Topic: "events", Buffer: 1, Overflow: DropNewest})
publishString(t, bus, "events", "one")
waitFor(t, func() bool { return subscription.Stats().Queued == 0 })
publishString(t, bus, "events", "two")
publishString(t, bus, "events", "three")
if got := string(receiveFrame(t, subscription).Payload); got != "one" {
t.Fatalf("first frame = %q", got)
}
if got := string(receiveFrame(t, subscription).Payload); got != "two" {
t.Fatalf("queued frame = %q", got)
}
if got := subscription.Stats().Dropped; got != 1 {
t.Fatalf("dropped = %d", got)
}
}
TestBlockAppliesBackpressure
Parameters
func TestBlockAppliesBackpressure(t *testing.T)
{
bus := NewInMemory(Config{})
t.Cleanup(func() { _ = bus.Close() })
subscription := mustSubscribe(t, bus, SubscribeOptions{Topic: "bulk", Buffer: 1, Overflow: Block})
publishString(t, bus, "bulk", "one")
waitFor(t, func() bool { return subscription.Stats().Queued == 0 })
publishString(t, bus, "bulk", "two")
completed := make(chan error, 1)
go func() {
_, err := bus.Publish(context.Background(), Frame{Topic: "bulk", Payload: []byte("three")})
completed <- err
}()
select {
case err := <-completed:
t.Fatalf("publish completed before capacity was available: %v", err)
case <-time.After(20 * time.Millisecond):
}
_ = receiveFrame(t, subscription)
select {
case err := <-completed:
if err != nil {
t.Fatal(err)
}
case <-time.After(time.Second):
t.Fatal("blocked publish did not resume")
}
}
TestPriorityBypassesQueuedBulkFrames
Parameters
func TestPriorityBypassesQueuedBulkFrames(t *testing.T)
{
bus := NewInMemory(Config{})
t.Cleanup(func() { _ = bus.Close() })
subscription := mustSubscribe(t, bus, SubscribeOptions{Topic: "mixed", Buffer: 3, Overflow: Block})
if _, err := bus.Publish(context.Background(), Frame{Topic: "mixed", Payload: []byte("in-flight"), Priority: PriorityBulk}); err != nil {
t.Fatal(err)
}
waitFor(t, func() bool { return subscription.Stats().Queued == 0 })
if _, err := bus.Publish(context.Background(), Frame{Topic: "mixed", Payload: []byte("bulk"), Priority: PriorityBulk}); err != nil {
t.Fatal(err)
}
if _, err := bus.Publish(context.Background(), Frame{Topic: "mixed", Payload: []byte("critical"), Priority: PriorityCritical}); err != nil {
t.Fatal(err)
}
if got := string(receiveFrame(t, subscription).Payload); got != "in-flight" {
t.Fatalf("first frame = %q", got)
}
if got := string(receiveFrame(t, subscription).Payload); got != "critical" {
t.Fatalf("prioritized frame = %q", got)
}
if got := string(receiveFrame(t, subscription).Payload); got != "bulk" {
t.Fatalf("bulk frame = %q", got)
}
}
TestCopyPayload
Parameters
func TestCopyPayload(t *testing.T)
{
bus := NewInMemory(Config{CopyPayload: true})
t.Cleanup(func() { _ = bus.Close() })
subscription := mustSubscribe(t, bus, SubscribeOptions{Topic: "safe", Buffer: 1})
payload := []byte("before")
if _, err := bus.Publish(context.Background(), Frame{Topic: "safe", Payload: payload}); err != nil {
t.Fatal(err)
}
copy(payload, "after!")
if got := string(receiveFrame(t, subscription).Payload); got != "before" {
t.Fatalf("copied payload = %q", got)
}
}
TestSubscriptionContextAndBusClose
Parameters
func TestSubscriptionContextAndBusClose(t *testing.T)
{
bus := NewInMemory(Config{})
ctx, cancel := context.WithCancel(context.Background())
subscription := mustSubscribeContext(t, bus, ctx, SubscribeOptions{Topic: "done", Buffer: 1})
cancel()
select {
case <-subscription.Done():
case <-time.After(time.Second):
t.Fatal("subscription did not follow context cancellation")
}
if err := bus.Close(); err != nil {
t.Fatal(err)
}
if _, err := bus.Publish(context.Background(), Frame{Topic: "done"}); !errors.Is(err, ErrClosed) {
t.Fatalf("publish after close = %v", err)
}
}
TestValidation
Parameters
func TestValidation(t *testing.T)
{
bus := NewInMemory(Config{MaxBuffer: 2, MaxPayloadBytes: 3})
t.Cleanup(func() { _ = bus.Close() })
if _, err := bus.Subscribe(context.Background(), SubscribeOptions{Topic: "x", Buffer: 3}); !errors.Is(err, ErrInvalidBuffer) {
t.Fatalf("oversized buffer = %v", err)
}
if _, err := bus.Publish(context.Background(), Frame{Topic: "x", Payload: []byte("four")}); !errors.Is(err, ErrPayloadTooLarge) {
t.Fatalf("oversized payload = %v", err)
}
withoutReplay := NewInMemory(Config{})
t.Cleanup(func() { _ = withoutReplay.Close() })
if _, err := withoutReplay.Subscribe(context.Background(), SubscribeOptions{Topic: "x", Since: 1}); !errors.Is(err, ErrReplayUnavailable) {
t.Fatalf("resume without replay = %v", err)
}
}
mustSubscribe
Parameters
Returns
func mustSubscribe(t *testing.T, bus Bus, options SubscribeOptions) *Subscription
{
t.Helper()
return mustSubscribeContext(t, bus, context.Background(), options)
}
mustSubscribeContext
Parameters
Returns
func mustSubscribeContext(t *testing.T, bus Bus, ctx context.Context, options SubscribeOptions) *Subscription
{
t.Helper()
subscription, err := bus.Subscribe(ctx, options)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = subscription.Close() })
return subscription
}
publishString
func publishString(t *testing.T, bus Bus, topic, payload string)
{
t.Helper()
if _, err := bus.Publish(context.Background(), Frame{Topic: topic, Payload: []byte(payload)}); err != nil {
t.Fatal(err)
}
}
Uses
receiveFrame
Parameters
Returns
func receiveFrame(t *testing.T, subscription *Subscription) Frame
{
t.Helper()
select {
case frame, ok := <-subscription.Frames():
if !ok {
t.Fatal("subscription closed")
}
return frame
case <-time.After(time.Second):
t.Fatal("timed out waiting for frame")
return Frame{}
}
}
Uses
waitFor
Parameters
func waitFor(t *testing.T, condition func() bool)
{
t.Helper()
deadline := time.Now().Add(time.Second)
for time.Now().Before(deadline) {
if condition() {
return
}
time.Sleep(time.Millisecond)
}
t.Fatal("condition was not met")
}
ringQueue
type ringQueue struct
Methods
Fields
| Name | Type | Description |
|---|---|---|
| items | []Frame | |
| head | int | |
| len | int |
newRingQueue
Parameters
Returns
func newRingQueue(capacity int) ringQueue
{
return ringQueue{items: make([]Frame, capacity)}
}
frameQueue
type frameQueue struct
Methods
Parameters
Returns
func (*frameQueue) push(frame Frame) bool
{
if q.len == q.capacity {
return false
}
priority := int(frame.Priority)
if priority < int(PriorityBulk) || priority > int(PriorityCritical) {
priority = int(PriorityNormal)
}
q.queues[priority].push(frame)
q.len++
return true
}
pop returns the highest-priority queued frame while retaining FIFO order within a priority class.
Returns
func (*frameQueue) pop() (Frame, bool)
{
for priority := int(PriorityCritical); priority >= int(PriorityBulk); priority-- {
if frame, ok := q.queues[priority].pop(); ok {
q.len--
return frame, true
}
}
return Frame{}, false
}
popOldest removes the earliest sequence regardless of priority. It is used by DropOldest, whose name describes age rather than scheduling priority.
Returns
func (*frameQueue) popOldest() (Frame, bool)
{
selected := -1
var oldest Frame
for priority := int(PriorityBulk); priority <= int(PriorityCritical); priority++ {
frame, ok := q.queues[priority].peek()
if !ok {
continue
}
if selected == -1 || frame.Sequence < oldest.Sequence {
selected = priority
oldest = frame
}
}
if selected == -1 {
return Frame{}, false
}
frame, _ := q.queues[selected].pop()
q.len--
return frame, true
}
Returns
func (*frameQueue) clear() int
{
dropped := q.len
for priority := range q.queues {
for q.queues[priority].len > 0 {
_, _ = q.queues[priority].pop()
}
}
q.len = 0
return dropped
}
Fields
| Name | Type | Description |
|---|---|---|
| queues | [4]ringQueue | |
| capacity | int | |
| len | int |
newFrameQueue
Parameters
Returns
func newFrameQueue(capacity int) frameQueue
{
queue := frameQueue{capacity: capacity}
for i := range queue.queues {
queue.queues[i] = newRingQueue(capacity)
}
return queue
}
MessageKind
MessageKind identifies StreamBus control and data messages on a transport.
type MessageKind uint8
Message
Message is the binary protocol envelope shared by transport adapters.
type Message struct
Fields
| Name | Type | Description |
|---|---|---|
| Kind | MessageKind | |
| ID | uint64 | |
| SubscriptionID | uint64 | |
| Ack | uint64 | |
| Frame | Frame | |
| Options | SubscribeOptions | |
| Error | string |
EncodeMessage
EncodeMessage encodes one message without stream-length framing. It is
suitable for a WebTransport datagram.
Parameters
Returns
func EncodeMessage(message Message) ([]byte, error)
{
topic := message.Frame.Topic
if message.Kind == MessageSubscribe && topic == "" {
topic = message.Options.Topic
}
if len(topic) > 1<<20 || len(message.Frame.Payload) > 64<<20 || len(message.Error) > 1<<20 {
return nil, ErrInvalidWireMessage
}
var out bytes.Buffer
out.WriteByte(WireVersion)
out.WriteByte(byte(message.Kind))
putUvarint(&out, message.ID)
putUvarint(&out, message.SubscriptionID)
putUvarint(&out, message.Ack)
putUvarint(&out, message.Frame.Sequence)
timestamp := int64(0)
if !message.Frame.Timestamp.IsZero() {
timestamp = message.Frame.Timestamp.UnixNano()
}
putVarint(&out, timestamp)
out.WriteByte(byte(message.Frame.Reliability))
out.WriteByte(byte(message.Frame.Priority))
putBytes(&out, []byte(topic))
putBytes(&out, message.Frame.Payload)
if message.Options.Prefix {
out.WriteByte(1)
} else {
out.WriteByte(0)
}
putUvarint(&out, uint64(message.Options.Buffer))
out.WriteByte(byte(message.Options.Overflow))
putUvarint(&out, uint64(message.Options.Replay))
putUvarint(&out, message.Options.Since)
putBytes(&out, []byte(message.Error))
return out.Bytes(), nil
}
Uses
DecodeMessage
DecodeMessage decodes one unframed message.
Parameters
Returns
func DecodeMessage(data []byte) (Message, error)
{
reader := bytes.NewReader(data)
version, err := reader.ReadByte()
if err != nil {
return Message{}, ErrInvalidWireMessage
}
if version != WireVersion {
return Message{}, ErrInvalidWireVersion
}
kind, err := reader.ReadByte()
if err != nil {
return Message{}, ErrInvalidWireMessage
}
message := Message{Kind: MessageKind(kind)}
if message.Kind < MessageData || message.Kind > MessagePong {
return Message{}, ErrInvalidWireMessage
}
if message.ID, err = binary.ReadUvarint(reader); err != nil {
return Message{}, ErrInvalidWireMessage
}
if message.SubscriptionID, err = binary.ReadUvarint(reader); err != nil {
return Message{}, ErrInvalidWireMessage
}
if message.Ack, err = binary.ReadUvarint(reader); err != nil {
return Message{}, ErrInvalidWireMessage
}
if message.Frame.Sequence, err = binary.ReadUvarint(reader); err != nil {
return Message{}, ErrInvalidWireMessage
}
timestamp, err := binary.ReadVarint(reader)
if err != nil {
return Message{}, ErrInvalidWireMessage
}
if timestamp != 0 {
message.Frame.Timestamp = time.Unix(0, timestamp).UTC()
}
reliability, err := reader.ReadByte()
if err != nil {
return Message{}, ErrInvalidWireMessage
}
message.Frame.Reliability = Reliability(reliability)
if message.Frame.Reliability > Unreliable {
return Message{}, ErrInvalidWireMessage
}
priority, err := reader.ReadByte()
if err != nil {
return Message{}, ErrInvalidWireMessage
}
message.Frame.Priority = Priority(priority)
if message.Frame.Priority > PriorityCritical {
return Message{}, ErrInvalidWireMessage
}
topic, err := readBytes(reader, 1<<20)
if err != nil {
return Message{}, err
}
message.Frame.Topic = string(topic)
if message.Frame.Payload, err = readBytes(reader, 64<<20); err != nil {
return Message{}, err
}
prefix, err := reader.ReadByte()
if err != nil {
return Message{}, ErrInvalidWireMessage
}
if message.Kind == MessageSubscribe {
message.Options.Topic = message.Frame.Topic
}
message.Options.Prefix = prefix == 1
buffer, err := binary.ReadUvarint(reader)
if err != nil || buffer > uint64(^uint(0)>>1) {
return Message{}, ErrInvalidWireMessage
}
message.Options.Buffer = int(buffer)
overflow, err := reader.ReadByte()
if err != nil {
return Message{}, ErrInvalidWireMessage
}
message.Options.Overflow = OverflowPolicy(overflow)
if message.Options.Overflow > LatestOnly {
return Message{}, ErrInvalidWireMessage
}
replay, err := binary.ReadUvarint(reader)
if err != nil || replay > uint64(^uint(0)>>1) {
return Message{}, ErrInvalidWireMessage
}
message.Options.Replay = int(replay)
if message.Options.Since, err = binary.ReadUvarint(reader); err != nil {
return Message{}, ErrInvalidWireMessage
}
errorText, err := readBytes(reader, 1<<20)
if err != nil {
return Message{}, err
}
message.Error = string(errorText)
if reader.Len() != 0 {
return Message{}, ErrInvalidWireMessage
}
return message, nil
}
Uses
WriteMessage
WriteMessage writes a length-delimited message to a reliable stream.
func WriteMessage(writer io.Writer, message Message) error
{
data, err := EncodeMessage(message)
if err != nil {
return err
}
var prefix [binary.MaxVarintLen64]byte
n := binary.PutUvarint(prefix[:], uint64(len(data)))
if err := writeAll(writer, prefix[:n]); err != nil {
return err
}
return writeAll(writer, data)
}
Uses
Reader
Reader decodes length-delimited StreamBus messages.
type Reader struct
Methods
ReadMessage reads one length-delimited message.
Returns
func (*Reader) ReadMessage() (Message, error)
{
size, err := binary.ReadUvarint(r.reader)
if err != nil {
return Message{}, err
}
if size > uint64(r.max) {
return Message{}, fmt.Errorf("%w: frame length %d exceeds %d", ErrInvalidWireMessage, size, r.max)
}
data := make([]byte, int(size))
if _, err := io.ReadFull(r.reader, data); err != nil {
return Message{}, err
}
return DecodeMessage(data)
}
Fields
| Name | Type | Description |
|---|---|---|
| reader | *bufio.Reader | |
| max | int |
NewReader
NewReader creates a stream decoder. A non-positive maximum uses 64 MiB.
func NewReader(reader io.Reader, maximum int) *Reader
{
if maximum <= 0 {
maximum = 64 << 20
}
return &Reader{reader: bufio.NewReader(reader), max: maximum}
}
putUvarint
Parameters
func putUvarint(out *bytes.Buffer, value uint64)
{
var data [binary.MaxVarintLen64]byte
n := binary.PutUvarint(data[:], value)
out.Write(data[:n])
}
putVarint
Parameters
func putVarint(out *bytes.Buffer, value int64)
{
var data [binary.MaxVarintLen64]byte
n := binary.PutVarint(data[:], value)
out.Write(data[:n])
}
putBytes
Parameters
func putBytes(out *bytes.Buffer, data []byte)
{
putUvarint(out, uint64(len(data)))
out.Write(data)
}
readBytes
Parameters
Returns
func readBytes(reader *bytes.Reader, maximum uint64) ([]byte, error)
{
size, err := binary.ReadUvarint(reader)
if err != nil || size > maximum || size > uint64(reader.Len()) {
return nil, ErrInvalidWireMessage
}
if size == 0 {
return nil, nil
}
data := make([]byte, int(size))
if _, err := io.ReadFull(reader, data); err != nil {
return nil, ErrInvalidWireMessage
}
return data, nil
}
writeAll
Parameters
Returns
func writeAll(writer io.Writer, data []byte) error
{
for len(data) > 0 {
n, err := writer.Write(data)
if err != nil {
return err
}
if n <= 0 {
return io.ErrShortWrite
}
data = data[n:]
}
return nil
}
TestMessageRoundTrip
Parameters
func TestMessageRoundTrip(t *testing.T)
{
want := Message{
Kind: MessageSubscribe,
ID: 41,
SubscriptionID: 9,
Ack: 20,
Frame: Frame{
Sequence: 21,
Timestamp: time.Unix(123, 456).UTC(),
Topic: "component:chart",
Payload: []byte{0, 1, 2, 255},
Reliability: Unreliable,
Priority: PriorityInteractive,
},
Options: SubscribeOptions{
Topic: "component:chart",
Prefix: true,
Buffer: 128,
Overflow: LatestOnly,
Replay: 3,
Since: 99,
},
Error: "example",
}
encoded, err := EncodeMessage(want)
if err != nil {
t.Fatal(err)
}
got, err := DecodeMessage(encoded)
if err != nil {
t.Fatal(err)
}
if !reflect.DeepEqual(got, want) {
t.Fatalf("round trip mismatch\n got: %#v\nwant: %#v", got, want)
}
}
TestFramedMessages
Parameters
func TestFramedMessages(t *testing.T)
{
var stream bytes.Buffer
want := []Message{
{Kind: MessagePing, ID: 1},
{Kind: MessageData, ID: 2, Frame: Frame{Topic: "x", Payload: []byte("payload")}},
}
for _, message := range want {
if err := WriteMessage(&stream, message); err != nil {
t.Fatal(err)
}
}
reader := NewReader(&stream, 1024)
for i := range want {
got, err := reader.ReadMessage()
if err != nil {
t.Fatal(err)
}
if !reflect.DeepEqual(got, want[i]) {
t.Fatalf("message %d = %#v, want %#v", i, got, want[i])
}
}
}
TestInvalidWireData
Parameters
func TestInvalidWireData(t *testing.T)
{
if _, err := DecodeMessage(nil); !errors.Is(err, ErrInvalidWireMessage) {
t.Fatalf("empty message = %v", err)
}
data, err := EncodeMessage(Message{Kind: MessagePing})
if err != nil {
t.Fatal(err)
}
data[0]++
if _, err := DecodeMessage(data); !errors.Is(err, ErrInvalidWireVersion) {
t.Fatalf("wrong version = %v", err)
}
}
TestWriteMessageHandlesShortWrites
Parameters
func TestWriteMessageHandlesShortWrites(t *testing.T)
{
writer := &shortWriter{limit: 2}
want := Message{Kind: MessageData, Frame: Frame{Topic: "topic", Payload: []byte("payload")}}
if err := WriteMessage(writer, want); err != nil {
t.Fatal(err)
}
got, err := NewReader(bytes.NewReader(writer.data), 1024).ReadMessage()
if err != nil {
t.Fatal(err)
}
if !reflect.DeepEqual(got, want) {
t.Fatalf("message = %#v, want %#v", got, want)
}
}
shortWriter
type shortWriter struct
Methods
Parameters
Returns
func (*shortWriter) Write(data []byte) (int, error)
{
if len(data) == 0 {
return 0, io.ErrShortWrite
}
if len(data) > w.limit {
data = data[:w.limit]
}
w.data = append(w.data, data...)
return len(data), nil
}
Fields
| Name | Type | Description |
|---|---|---|
| data | []byte | |
| limit | int |
Subscription
Subscription exposes a bounded stream of frames.
type Subscription struct
Methods
ID is stable for the lifetime of the subscription.
Returns
func (*Subscription) ID() uint64
{ return s.id }
Frames returns the delivery channel. It is closed when the subscription is closed or its parent context is canceled.
Returns
func (*Subscription) Frames() <-chan Frame
{ return s.frames }
Done is closed when the subscription stops.
Returns
func (*Subscription) Done() <-chan struct{}
{ return s.done }
Stats returns current counters and queue depth.
Returns
func (*Subscription) Stats() SubscriptionStats
{
s.mu.Lock()
queued := s.queue.len
s.mu.Unlock()
return SubscriptionStats{
Enqueued: s.enqueued.Load(),
Delivered: s.delivered.Load(),
Dropped: s.dropped.Load(),
Queued: queued,
}
}
Close releases the subscription. It is safe to call more than once.
Returns
func (*Subscription) Close() error
{
s.closeOnce.Do(func() {
s.mu.Lock()
s.closed = true
s.mu.Unlock()
close(s.done)
if s.unregister != nil {
s.unregister()
}
})
return nil
}
Parameters
func (*Subscription) preload(frames []Frame)
{
s.mu.Lock()
defer s.mu.Unlock()
for _, frame := range frames {
if !s.queue.push(frame) {
_, _ = s.queue.popOldest()
_ = s.queue.push(frame)
s.dropped.Add(1)
}
s.enqueued.Add(1)
}
if s.queue.len > 0 {
s.signal(s.notify)
}
}
Parameters
Returns
func (*Subscription) enqueue(ctx context.Context, frame Frame) error
{
for {
s.mu.Lock()
if s.closed {
s.mu.Unlock()
return ErrSubscriptionDone
}
if s.queue.push(frame) {
s.enqueued.Add(1)
s.mu.Unlock()
s.signal(s.notify)
return nil
}
switch s.options.Overflow {
case DropNewest:
s.dropped.Add(1)
s.mu.Unlock()
return nil
case DropOldest:
_, _ = s.queue.popOldest()
_ = s.queue.push(frame)
s.dropped.Add(1)
s.enqueued.Add(1)
s.mu.Unlock()
s.signal(s.notify)
return nil
case LatestOnly:
dropped := s.queue.clear()
_ = s.queue.push(frame)
s.dropped.Add(uint64(dropped))
s.enqueued.Add(1)
s.mu.Unlock()
s.signal(s.notify)
return nil
default:
s.mu.Unlock()
select {
case <-ctx.Done():
return ctx.Err()
case <-s.done:
return ErrSubscriptionDone
case <-s.space:
}
}
}
}
func (*Subscription) run()
{
defer close(s.frames)
for {
select {
case <-s.done:
return
case <-s.notify:
for {
s.mu.Lock()
frame, ok := s.queue.pop()
s.mu.Unlock()
if !ok {
break
}
s.signal(s.space)
select {
case <-s.done:
return
case s.frames <- frame:
s.delivered.Add(1)
}
}
}
}
}
Parameters
func (*Subscription) signal(ch chan struct{})
{
select {
case ch <- struct{}{}:
default:
}
}
Fields
| Name | Type | Description |
|---|---|---|
| id | uint64 | |
| options | SubscribeOptions | |
| frames | chan Frame | |
| done | chan struct{} | |
| notify | chan struct{} | |
| space | chan struct{} | |
| queue | frameQueue | |
| mu | sync.Mutex | |
| closed | bool | |
| closeOnce | sync.Once | |
| unregister | func() | |
| enqueued | atomic.Uint64 | |
| delivered | atomic.Uint64 | |
| dropped | atomic.Uint64 |
newSubscription
Parameters
Returns
func newSubscription(id uint64, options SubscribeOptions, unregister func()) *Subscription
{
s := &Subscription{
id: id,
options: options,
frames: make(chan Frame),
done: make(chan struct{}),
notify: make(chan struct{}, 1),
space: make(chan struct{}, 1),
queue: newFrameQueue(options.Buffer),
unregister: unregister,
}
go s.run()
return s
}
Reliability
Reliability describes whether a transport may discard a frame in flight.
The in-memory bus preserves this value for transport adapters.
type Reliability uint8
Priority
Priority allows transports to isolate latency-sensitive traffic from bulk
streams. Higher values represent more urgent traffic.
type Priority uint8
OverflowPolicy
OverflowPolicy defines what happens when a subscriber cannot keep up.
type OverflowPolicy uint8
Frame
Frame is an immutable unit delivered by StreamBus. Payload must not be
modified after Publish unless Config.CopyPayload is enabled.
type Frame struct
Fields
| Name | Type | Description |
|---|---|---|
| Sequence | uint64 | |
| Timestamp | time.Time | |
| Topic | string | |
| Payload | []byte | |
| Reliability | Reliability | |
| Priority | Priority |
SubscribeOptions
SubscribeOptions configure an exact-topic or prefix subscription.
type SubscribeOptions struct
Fields
| Name | Type | Description |
|---|---|---|
| Topic | string | |
| Prefix | bool | |
| Buffer | int | |
| Overflow | OverflowPolicy | |
| Replay | int | |
| Since | uint64 |
Config
Config controls an in-memory StreamBus.
type Config struct
Methods
Returns
func (Config) normalized() Config
{
if c.DefaultBuffer <= 0 {
c.DefaultBuffer = 64
}
if c.MaxBuffer <= 0 {
c.MaxBuffer = 4096
}
if c.DefaultBuffer > c.MaxBuffer {
c.DefaultBuffer = c.MaxBuffer
}
if c.MaxPayloadBytes <= 0 {
c.MaxPayloadBytes = 16 << 20
}
return c
}
Fields
| Name | Type | Description |
|---|---|---|
| DefaultBuffer | int | |
| MaxBuffer | int | |
| ReplayCapacity | int | |
| MaxPayloadBytes | int | |
| CopyPayload | bool |
Bus
Bus is the transport-independent StreamBus contract.
type Bus interface
Methods
SubscriptionStats
SubscriptionStats is a point-in-time view of subscriber pressure.
type SubscriptionStats struct
Fields
| Name | Type | Description |
|---|---|---|
| Enqueued | uint64 | |
| Delivered | uint64 | |
| Dropped | uint64 | |
| Queued | int |