streambus API

streambus

package

API reference for the streambus package.

F
function

BenchmarkPublishFanout

Parameters

streambus/benchmark_test.go:9-37
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)
				}
			}
		})
	}
}
S
struct

prefixNode

streambus/inmemory.go:10-13
type prefixNode struct

Fields

Name Type Description
subs map[uint64]*Subscription
children map[byte]*prefixNode
F
function

newPrefixNode

Returns

streambus/inmemory.go:15-20
func newPrefixNode() *prefixNode

{
	return &prefixNode{
		subs:     make(map[uint64]*Subscription),
		children: make(map[byte]*prefixNode),
	}
}
S
struct

replayBuffer

streambus/inmemory.go:22-26
type replayBuffer struct

Methods

add
Method

Parameters

frame Frame
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++
	}
}
last
Method

Parameters

count int

Returns

[]Frame
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
}
after
Method

Parameters

sequence uint64

Returns

[]Frame
bool
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
F
function

newReplayBuffer

Parameters

capacity
int

Returns

streambus/inmemory.go:28-30
func newReplayBuffer(capacity int) *replayBuffer

{
	return &replayBuffer{frames: make([]Frame, capacity)}
}
S
struct

InMemory

InMemory is a bounded, transport-independent StreamBus implementation.

streambus/inmemory.go:75-84
type InMemory struct

Methods

Publish
Method

Publish fans a frame out to all exact and prefix subscribers.

Parameters

Returns

uint64
error
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
Method

Subscribe creates a bounded exact-topic or prefix subscription.

Parameters

Returns

error
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
}
remove
Method

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
Method

Close stops every subscription and rejects future operations.

Returns

error
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
F
function

NewInMemory

NewInMemory creates a StreamBus using process-local memory.

Parameters

config

Returns

streambus/inmemory.go:87-94
func NewInMemory(config Config) *InMemory

{
	return &InMemory{
		config: config.normalized(),
		exact:  make(map[string]map[uint64]*Subscription),
		prefix: newPrefixNode(),
		replay: make(map[string]*replayBuffer),
	}
}
F
function

collectPrefixSubscriptions

Parameters

node
dst
map[uint64]*Subscription
streambus/inmemory.go:289-296
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)
	}
}
F
function

TestPublishExactAndPrefix

Parameters

streambus/inmemory_test.go:10-37
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):
	}
}
F
function

TestReplayPreservesOrder

Parameters

streambus/inmemory_test.go:39-54
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)
	}
}
F
function

TestResumeFromSequence

Parameters

streambus/inmemory_test.go:56-74
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)
	}
}
F
function

TestResumeReportsKnownGap

Parameters

streambus/inmemory_test.go:76-88
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)
	}
}
F
function

TestLatestOnlyCollapsesQueuedFrames

Parameters

streambus/inmemory_test.go:90-109
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")
	}
}
F
function

TestDropNewestPreservesQueuedFrame

Parameters

streambus/inmemory_test.go:111-130
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)
	}
}
F
function

TestBlockAppliesBackpressure

Parameters

streambus/inmemory_test.go:132-160
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")
	}
}
F
function

TestPriorityBypassesQueuedBulkFrames

Parameters

streambus/inmemory_test.go:162-187
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)
	}
}
F
function

TestCopyPayload

Parameters

streambus/inmemory_test.go:189-201
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)
	}
}
F
function

TestSubscriptionContextAndBusClose

Parameters

streambus/inmemory_test.go:203-219
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)
	}
}
F
function

TestValidation

Parameters

streambus/inmemory_test.go:221-235
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)
	}
}
F
function

mustSubscribe

Parameters

bus

Returns

streambus/inmemory_test.go:237-240
func mustSubscribe(t *testing.T, bus Bus, options SubscribeOptions) *Subscription

{
	t.Helper()
	return mustSubscribeContext(t, bus, context.Background(), options)
}
F
function

mustSubscribeContext

Parameters

Returns

streambus/inmemory_test.go:242-250
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
}
F
function

publishString

Parameters

bus
topic
string
payload
string
streambus/inmemory_test.go:252-257
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)
	}
}
F
function

receiveFrame

Parameters

subscription

Returns

streambus/inmemory_test.go:259-271
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{}
	}
}
F
function

waitFor

Parameters

condition
func() bool
streambus/inmemory_test.go:273-283
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")
}
S
struct

ringQueue

streambus/queue.go:3-7
type ringQueue struct

Methods

push
Method

Parameters

frame Frame
func (*ringQueue) push(frame Frame)
{
	idx := (q.head + q.len) % len(q.items)
	q.items[idx] = frame
	q.len++
}
peek
Method

Returns

bool
func (*ringQueue) peek() (Frame, bool)
{
	if q.len == 0 {
		return Frame{}, false
	}
	return q.items[q.head], true
}
pop
Method

Returns

bool
func (*ringQueue) pop() (Frame, bool)
{
	if q.len == 0 {
		return Frame{}, false
	}
	frame := q.items[q.head]
	q.items[q.head] = Frame{}
	q.head = (q.head + 1) % len(q.items)
	q.len--
	return frame, true
}

Fields

Name Type Description
items []Frame
head int
len int
F
function

newRingQueue

Parameters

capacity
int

Returns

streambus/queue.go:9-11
func newRingQueue(capacity int) ringQueue

{
	return ringQueue{items: make([]Frame, capacity)}
}
S
struct

frameQueue

streambus/queue.go:37-41
type frameQueue struct

Methods

push
Method

Parameters

frame Frame

Returns

bool
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
Method

pop returns the highest-priority queued frame while retaining FIFO order within a priority class.

Returns

bool
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
Method

popOldest removes the earliest sequence regardless of priority. It is used by DropOldest, whose name describes age rather than scheduling priority.

Returns

bool
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
}
clear
Method

Returns

int
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
F
function

newFrameQueue

Parameters

capacity
int

Returns

streambus/queue.go:43-49
func newFrameQueue(capacity int) frameQueue

{
	queue := frameQueue{capacity: capacity}
	for i := range queue.queues {
		queue.queues[i] = newRingQueue(capacity)
	}
	return queue
}
T
type

MessageKind

MessageKind identifies StreamBus control and data messages on a transport.

streambus/wire.go:21-21
type MessageKind uint8
S
struct

Message

Message is the binary protocol envelope shared by transport adapters.

streambus/wire.go:36-44
type Message struct

Fields

Name Type Description
Kind MessageKind
ID uint64
SubscriptionID uint64
Ack uint64
Frame Frame
Options SubscribeOptions
Error string
F
function

EncodeMessage

EncodeMessage encodes one message without stream-length framing. It is
suitable for a WebTransport datagram.

Parameters

message

Returns

[]byte
error
streambus/wire.go:48-83
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
}
F
function

DecodeMessage

DecodeMessage decodes one unframed message.

Parameters

data
[]byte

Returns

error
streambus/wire.go:86-184
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
}
F
function

WriteMessage

WriteMessage writes a length-delimited message to a reliable stream.

Parameters

writer
message

Returns

error
streambus/wire.go:187-198
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)
}
S
struct

Reader

Reader decodes length-delimited StreamBus messages.

streambus/wire.go:201-204
type Reader struct

Methods

ReadMessage
Method

ReadMessage reads one length-delimited message.

Returns

error
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
F
function

NewReader

NewReader creates a stream decoder. A non-positive maximum uses 64 MiB.

Parameters

reader
maximum
int

Returns

streambus/wire.go:207-212
func NewReader(reader io.Reader, maximum int) *Reader

{
	if maximum <= 0 {
		maximum = 64 << 20
	}
	return &Reader{reader: bufio.NewReader(reader), max: maximum}
}
F
function

putUvarint

Parameters

value
uint64
streambus/wire.go:230-234
func putUvarint(out *bytes.Buffer, value uint64)

{
	var data [binary.MaxVarintLen64]byte
	n := binary.PutUvarint(data[:], value)
	out.Write(data[:n])
}
F
function

putVarint

Parameters

value
int64
streambus/wire.go:236-240
func putVarint(out *bytes.Buffer, value int64)

{
	var data [binary.MaxVarintLen64]byte
	n := binary.PutVarint(data[:], value)
	out.Write(data[:n])
}
F
function

putBytes

Parameters

data
[]byte
streambus/wire.go:242-245
func putBytes(out *bytes.Buffer, data []byte)

{
	putUvarint(out, uint64(len(data)))
	out.Write(data)
}
F
function

readBytes

Parameters

reader
maximum
uint64

Returns

[]byte
error
streambus/wire.go:247-260
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
}
F
function

writeAll

Parameters

writer
data
[]byte

Returns

error
streambus/wire.go:262-274
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
}
F
function

TestMessageRoundTrip

Parameters

streambus/wire_test.go:12-47
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)
	}
}
F
function

TestFramedMessages

Parameters

streambus/wire_test.go:49-70
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])
		}
	}
}
F
function

TestInvalidWireData

Parameters

streambus/wire_test.go:72-84
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)
	}
}
F
function

TestWriteMessageHandlesShortWrites

Parameters

streambus/wire_test.go:86-99
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)
	}
}
S
struct

shortWriter

streambus/wire_test.go:101-104
type shortWriter struct

Methods

Write
Method

Parameters

data []byte

Returns

int
error
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
S
struct

Subscription

Subscription exposes a bounded stream of frames.

streambus/subscription.go:10-26
type Subscription struct

Methods

ID
Method

ID is stable for the lifetime of the subscription.

Returns

uint64
func (*Subscription) ID() uint64
{ return s.id }
Frames
Method

Frames returns the delivery channel. It is closed when the subscription is closed or its parent context is canceled.

Returns

<-chan Frame
func (*Subscription) Frames() <-chan Frame
{ return s.frames }
Done
Method

Done is closed when the subscription stops.

Returns

<-chan struct{}
func (*Subscription) Done() <-chan struct{}
{ return s.done }
Stats
Method

Stats returns current counters and queue depth.

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
Method

Close releases the subscription. It is safe to call more than once.

Returns

error
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
}
preload
Method

Parameters

frames []Frame
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)
	}
}
enqueue
Method

Parameters

Returns

error
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:
			}
		}
	}
}
run
Method
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)
				}
			}
		}
	}
}
signal
Method

Parameters

ch chan struct{}
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
F
function

newSubscription

Parameters

id
uint64
unregister
func()

Returns

streambus/subscription.go:28-41
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
}
T
type

Reliability

Reliability describes whether a transport may discard a frame in flight.
The in-memory bus preserves this value for transport adapters.

streambus/types.go:20-20
type Reliability uint8
T
type

Priority

Priority allows transports to isolate latency-sensitive traffic from bulk
streams. Higher values represent more urgent traffic.

streambus/types.go:29-29
type Priority uint8
T
type

OverflowPolicy

OverflowPolicy defines what happens when a subscriber cannot keep up.

streambus/types.go:39-39
type OverflowPolicy uint8
S
struct

Frame

Frame is an immutable unit delivered by StreamBus. Payload must not be
modified after Publish unless Config.CopyPayload is enabled.

streambus/types.go:54-61
type Frame struct

Fields

Name Type Description
Sequence uint64
Timestamp time.Time
Topic string
Payload []byte
Reliability Reliability
Priority Priority
S
struct

SubscribeOptions

SubscribeOptions configure an exact-topic or prefix subscription.

streambus/types.go:64-75
type SubscribeOptions struct

Fields

Name Type Description
Topic string
Prefix bool
Buffer int
Overflow OverflowPolicy
Replay int
Since uint64
S
struct

Config

Config controls an in-memory StreamBus.

streambus/types.go:78-84
type Config struct

Methods

normalized
Method

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
I
interface

Bus

Bus is the transport-independent StreamBus contract.

streambus/types.go:103-107
type Bus interface

Methods

Publish
Method

Parameters

Returns

uint64
error
func Publish(...)
Subscribe
Method

Returns

error
func Subscribe(...)
Close
Method

Returns

error
func Close(...)
S
struct

SubscriptionStats

SubscriptionStats is a point-in-time view of subscriber pressure.

streambus/types.go:110-115
type SubscriptionStats struct

Fields

Name Type Description
Enqueued uint64
Delivered uint64
Dropped uint64
Queued int