libbus

package module
v0.0.0-...-6216d24 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 30, 2026 License: Apache-2.0 Imports: 13 Imported by: 0

README

libbus

libbus is a small, backend-agnostic publish-subscribe abstraction for Go: fire-and-forget publish, streaming subscriptions, and request-reply, on top of three pluggable backends.

It was extracted from contenox/contenox.

What it provides

The whole public surface is one interface plus a handler type:

type Handler func(ctx context.Context, data []byte) ([]byte, error)

type Messenger interface {
	Publish(ctx context.Context, subject string, data []byte) error
	Stream(ctx context.Context, subject string, ch chan<- []byte) (Subscription, error)
	Request(ctx context.Context, subject string, data []byte) ([]byte, error)
	Serve(ctx context.Context, subject string, handler Handler) (Subscription, error)
	Close() error
}

type Subscription interface {
	Unsubscribe() error
}

Every backend guarantees the same contract (enforced by a shared conformance test suite): Publish to a subject with no subscribers is a no-op and never blocks; Stream delivers in publish order until Unsubscribe or context cancellation; a handler error still produces a reply — Request returns a non-nil error only on transport failure, never on handler failure; and after Close, every method returns ErrConnectionClosed.

Backends differ in ways callers must tolerate: NATS and InMem are at-most-once under backpressure (they drop once a subscriber's buffer fills), while SQLiteBus is durable; NATS and InMem require Serve to be registered before Request is called, SQLiteBus does not (it polls, so a late Serve still picks up a pending request); handler concurrency and delivery latency (SQLiteBus is poll-driven) vary per backend. Always give Request a deadline.

The three backends
  • InMem (NewInMem()) — single-process, no external dependencies. Reproduces NATS's observable behavior (at-most-once, drop-under-backpressure, immediate failure on Request with no registered handler) so it's a faithful stand-in for tests or single-instance deployments that don't need durability or cross-process delivery.

  • SQLiteBus (NewSQLite(exec)) — durable, poll-driven, backed by a database/sql connection satisfying a minimal ExecContext/QueryContext interface. Requires three tables to already exist in the database: bus_events, bus_requests, bus_replies (see schema below). Pick this when you already have a SQLite database in the process and want durable pub-sub/request-reply without standing up a broker — messages survive a crash between publish and delivery, at the cost of poll-interval latency (tunable via NewSQLiteWithOptions).

  • NATS (NewPubSub(ctx, cfg)) — real broker, for multi-process/networked deployments needing low-latency delivery and standard pub-sub semantics. Not durable by default (matches InMem's at-most-once-under-backpressure behavior). local_nats.go also provides NewTestPubSub() / SetupNatsInstance(ctx), which spin up a real NATS server via testcontainers — useful in your own integration tests, but it requires Docker.

Install

go get github.com/contenox/libbus

Usage

Minimal in-memory example — publish/subscribe and request/reply, no external services required:

package main

import (
	"context"
	"fmt"
	"time"

	"github.com/contenox/libbus"
)

func main() {
	ctx := context.Background()
	bus := libbus.NewInMem()
	defer bus.Close()

	// Publish/Subscribe
	ch := make(chan []byte, 1)
	sub, err := bus.Stream(ctx, "greetings", ch)
	if err != nil {
		panic(err)
	}
	defer sub.Unsubscribe()

	if err := bus.Publish(ctx, "greetings", []byte("hello")); err != nil {
		panic(err)
	}
	fmt.Println(string(<-ch)) // "hello"

	// Request/Reply
	echoSub, err := bus.Serve(ctx, "echo", func(ctx context.Context, data []byte) ([]byte, error) {
		return data, nil
	})
	if err != nil {
		panic(err)
	}
	defer echoSub.Unsubscribe()

	reqCtx, cancel := context.WithTimeout(ctx, 2*time.Second)
	defer cancel()
	reply, err := bus.Request(reqCtx, "echo", []byte("ping"))
	if err != nil {
		panic(err)
	}
	fmt.Println(string(reply)) // "ping"
}

Swapping in NATS only changes construction:

bus, err := libbus.NewPubSub(ctx, &libbus.Config{
	NATSURL: "nats://localhost:4222",
})

Swapping in SQLite requires the schema to exist first:

const schema = `
CREATE TABLE IF NOT EXISTS bus_events (
    id         INTEGER PRIMARY KEY AUTOINCREMENT,
    subject    TEXT    NOT NULL,
    data       BLOB    NOT NULL,
    created_at INTEGER NOT NULL DEFAULT (unixepoch('now'))
);
CREATE TABLE IF NOT EXISTS bus_requests (
    id         TEXT    PRIMARY KEY,
    subject    TEXT    NOT NULL,
    data       BLOB    NOT NULL,
    created_at INTEGER NOT NULL DEFAULT (unixepoch('now'))
);
CREATE TABLE IF NOT EXISTS bus_replies (
    request_id TEXT    PRIMARY KEY,
    data       BLOB    NOT NULL,
    created_at INTEGER NOT NULL DEFAULT (unixepoch('now'))
);
`

db, _ := sql.Open("sqlite", "app.db?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)")
db.Exec(schema)
bus := libbus.NewSQLite(db) // db satisfies the minimal ExecContext/QueryContext interface

Testing

The package's own test suite (conformance_test.go) runs the same behavioral test matrix against all three backends. The NATS backend is skipped automatically when Docker isn't available; set LIBBUS_REQUIRE_NATS=1 to turn that skip into a hard failure (e.g. in CI where Docker is expected).

License

Apache-2.0, see LICENSE.

Documentation

Overview

Package libbus is a high-level publish-subscribe abstraction over a message broker, offering fire-and-forget publish, streaming subscriptions, and request-reply (Serve/Request) on top of pluggable backends (NATS, SQLite, in-memory).

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrConnectionClosed is returned when an operation is attempted on a closed connection.
	ErrConnectionClosed = errors.New("connection closed")
	// ErrStreamSubscriptionFail is returned when a stream subscription fails.
	ErrStreamSubscriptionFail = errors.New("stream subscription failed")
	// ErrMessagePublish is returned when publishing a message fails for reasons other than a closed connection.
	ErrMessagePublish = errors.New("message publishing failed")
	// ErrRequestTimeout is returned when a request-reply operation times out.
	ErrRequestTimeout = errors.New("request timed out")
)

Functions

func SetupNatsInstance

func SetupNatsInstance(ctx context.Context) (string, testcontainers.Container, func(), error)

Types

type Config

type Config struct {
	NATSURL      string
	NATSPassword string
	NATSUser     string
}

type Handler

type Handler func(ctx context.Context, data []byte) ([]byte, error)

Handler is a function that processes a request and returns a response. It is used by the Serve method to handle incoming requests.

type InMem

type InMem struct {
	// contains filtered or unexported fields
}

InMem is an in-memory Messenger for single-process use: no NATS, no network. It intentionally reproduces the NATS backend's observable contract (at-most-once delivery, Request failing immediately with no handler registered) — see the Messenger interface docs for the full matrix.

func NewInMem

func NewInMem() *InMem

NewInMem returns a new in-memory Messenger. Use for local single-process mode (no NATS).

func (*InMem) Close

func (p *InMem) Close() error

Close marks the messenger closed and releases resources.

func (*InMem) Publish

func (p *InMem) Publish(ctx context.Context, subject string, data []byte) error

Publish hands the message to every Stream subscriber's queue and returns. It never blocks on a consumer: a full queue drops the message (see inmemStreamBuffer) so that a wedged subscriber cannot take the publisher with it.

func (*InMem) Request

func (p *InMem) Request(ctx context.Context, subject string, data []byte) ([]byte, error)

Request invokes the Serve handler registered for the subject, in the caller's goroutine. Like the NATS backend it does NOT wait for a handler to appear: a missing handler fails immediately rather than after the context deadline.

func (*InMem) Serve

func (p *InMem) Serve(ctx context.Context, subject string, handler Handler) (Subscription, error)

Serve registers a handler for the subject. Request calls will invoke this handler.

func (*InMem) Stream

func (p *InMem) Stream(ctx context.Context, subject string, ch chan<- []byte) (Subscription, error)

Stream creates a subscription to a subject; messages are delivered to ch.

type Messenger

type Messenger interface {
	// Publish sends a fire-and-forget message to a given subject.
	Publish(ctx context.Context, subject string, data []byte) error

	// Stream creates a subscription to a subject and delivers messages asynchronously
	// to the provided channel. The subscription is automatically managed and will
	// be closed when the provided context is canceled.
	Stream(ctx context.Context, subject string, ch chan<- []byte) (Subscription, error)

	// Request sends a request message and waits for a reply. The context can be
	// used to set a timeout or to cancel the request.
	Request(ctx context.Context, subject string, data []byte) ([]byte, error)

	// Serve registers a handler for a given subject to respond to requests.
	// It starts a worker that listens for requests and executes the handler.
	// The returned Subscription can be used to stop serving.
	Serve(ctx context.Context, subject string, handler Handler) (Subscription, error)

	// Close disconnects from the messaging server and cleans up any underlying resources.
	Close() error
}

Messenger is a high-level pub-sub/request-reply interface for real-time notifications and distributing lightweight messages between services.

Guaranteed by every backend (enforced by conformance_test.go): Publish to a subject with no subscribers is a no-op and never blocks; Stream delivers in publish order until Unsubscribe or context cancel; a handler error still yields a reply — Request returns a non-nil error only on transport failure, never on handler failure; and after Close, every method returns ErrConnectionClosed.

Backends differ in ways callers must tolerate: NATS/InMem are at-most-once under backpressure (drop once a subscriber's buffer fills) while SQLiteBus is durable; NATS/InMem require Serve to return before Request is called, SQLiteBus does not; handler concurrency and delivery latency (SQLiteBus is poll-driven) vary per backend. Always give Request a deadline.

func NewPubSub

func NewPubSub(ctx context.Context, cfg *Config) (Messenger, error)

func NewTestPubSub

func NewTestPubSub() (Messenger, func(), error)

NewTestPubSub starts a NATS container using SetupNatsInstance, creates a new PubSub instance, and returns it along with a cleanup function.

type SQLiteBus

type SQLiteBus struct {
	// contains filtered or unexported fields
}

SQLiteBus implements Messenger over a SQLite database.

Schema tables (bus_events, bus_requests, bus_replies) must exist before use. They are part of runtimetypes.SchemaSQLite and are created automatically when the CLI database is opened.

Usage:

bus := libbus.NewSQLite(dbManager.WithoutTransaction())
defer bus.Close()

func NewSQLite

func NewSQLite(exec sqlExec) *SQLiteBus

NewSQLite creates a SQLite-backed Messenger. exec must be the result of dbManager.WithoutTransaction() — it satisfies sqlExec.

func NewSQLiteWithOptions

func NewSQLiteWithOptions(exec sqlExec, opt SQLiteBusOptions) *SQLiteBus

NewSQLiteWithOptions is like NewSQLite but allows tuning poll intervals for tests.

func (*SQLiteBus) Close

func (b *SQLiteBus) Close() error

Close stops all background goroutines. The underlying database is NOT closed (it is owned by the caller who provided the sqlExec).

func (*SQLiteBus) Publish

func (b *SQLiteBus) Publish(ctx context.Context, subject string, data []byte) error

Publish inserts a row into bus_events so Stream subscribers can pick it up.

func (*SQLiteBus) Request

func (b *SQLiteBus) Request(ctx context.Context, subject string, data []byte) ([]byte, error)

Request inserts a request row and polls for the reply until ctx deadline or 10s timeout.

func (*SQLiteBus) Serve

func (b *SQLiteBus) Serve(ctx context.Context, subject string, handler Handler) (Subscription, error)

Serve registers a handler for subject. A polling goroutine picks up rows from bus_requests, calls the handler, and writes the reply to bus_replies.

func (*SQLiteBus) Stream

func (b *SQLiteBus) Stream(ctx context.Context, subject string, ch chan<- []byte) (Subscription, error)

Stream starts a polling goroutine that delivers new bus_events for subject to ch. The subscription goroutine stops when ctx is cancelled.

type SQLiteBusOptions

type SQLiteBusOptions struct {
	EventPoll   time.Duration
	RequestPoll time.Duration
}

SQLiteBusOptions overrides poll intervals (e.g. tests use 1ms so request/reply is deterministic).

type Subscription

type Subscription interface {
	// Unsubscribe removes the subscription, stopping the delivery of messages.
	Unsubscribe() error
}

Subscription represents an active subscription to a subject.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL