meshbus

package module
v0.1.0 Latest Latest
Warning

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

Go to latest
Published: Sep 28, 2026 License: MIT Imports: 12 Imported by: 0

README

meshbus

meshbus is a small brokerless messaging layer for authenticated peers. It provides shared-secret realm membership, bounded peer discovery, opaque direct messages, and best-effort one-hop pub/sub.

The core package and realm use only the Go standard library. rns is the Reticulum adapter and is the only package that depends on Reticulum-Go.

Guarantees

  • The sender exposed to an application always comes from the authenticated transport session.
  • Serialized payload identities are never authoritative.
  • Discovery metadata is advisory and bounded.
  • Pub/sub is in-memory, TTL-bounded, deduplicated, and best-effort; it has no replay, forwarding, offsets, consumer groups, or exactly-once guarantee.
  • Transport routes remain adapter-private; applications address PeerID values.

See WIRE.md for the interoperable wire formats and compatibility policy.

Verification

go test ./...
go test -race ./...
go vet ./...

Status

The module is pre-v1 while its public API is exercised by external consumers. Published wire version markers and domain separators are changed only by introducing a new version.

Documentation

Overview

Package meshbus defines transport-independent authenticated peer messaging, bounded peer discovery, and best-effort pub/sub composition.

Generic, transport-independent realm peer discovery.

The entity here is a PeerDirectory, not a Cluster: the realm is a security boundary, whereas the directory is only the observable state of the network. A discovered peer carries a transport-authenticated public identity and bounded advisory metadata. Routes remain private to the transport adapter. Discovery is advisory only: an announce does not prove realm-key possession. Node promotes a candidate only after the transport reports successful realm authentication.

Index

Examples

Constants

View Source
const MaxPeerIDBytes = 256

MaxPeerIDBytes bounds identity material retained by the generic layer. Individual transports may impose a smaller bound.

View Source
const MaxTopicBytes = 128

Variables

View Source
var (
	ErrInvalidPeerID  = errors.New("invalid peer identity")
	ErrInvalidMessage = errors.New("invalid direct message")
)
View Source
var (
	ErrInvalidEvent      = errors.New("invalid event")
	ErrInvalidTopic      = errors.New("invalid event topic")
	ErrEventBackpressure = errors.New("event subscription queue is full")
	ErrSubscriptionLimit = errors.New("event subscription limit reached")
	ErrFanoutLimit       = errors.New("event fan-out peer limit exceeded")
	ErrBusClosed         = errors.New("event bus is closed")
)
View Source
var (
	ErrInvalidNode        = errors.New("invalid meshbus node")
	ErrUnknownPeer        = errors.New("unknown meshbus peer")
	ErrNodeNotStarted     = errors.New("meshbus node is not started")
	ErrNodeAlreadyStarted = errors.New("meshbus node is already started")
	ErrNodeClosed         = errors.New("meshbus node is closed")
)
View Source
var (
	// ErrInvalidPeer is returned when a peer record fails validation.
	ErrInvalidPeer = errors.New("invalid discovered peer")
	// ErrPeerLimit is returned when the directory is at capacity.
	ErrPeerLimit = errors.New("peer directory capacity reached")
	// ErrMetadataLimit is returned when application metadata exceeds bounds.
	ErrMetadataLimit = errors.New("peer metadata exceeds bound")
	// ErrInvalidDirectoryConfig is returned for negative resource bounds.
	ErrInvalidDirectoryConfig = errors.New("invalid peer directory configuration")
)

Functions

This section is empty.

Types

type Bus

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

Bus distributes best-effort events over authenticated direct messages. It owns no durable log, replay cursor, or consumer group state.

func NewBus

func NewBus(config BusConfig) (*Bus, error)

NewBus creates a stopped-empty but immediately usable event bus.

func (*Bus) Close

func (b *Bus) Close() error

Close stops subscriptions and rejects future publication or registration.

func (*Bus) HandleMessage

func (b *Bus) HandleMessage(_ context.Context, message ReceivedMessage) (bool, error)

HandleMessage consumes meshbus event frames and leaves other direct messages to the caller. A true result means the frame belonged to pub/sub even when it was expired, duplicated, invalid, or backpressured.

func (*Bus) Handler

func (b *Bus) Handler(next Handler) Handler

Handler composes pub/sub with a fallback direct-message handler.

func (*Bus) Publish

func (b *Bus) Publish(ctx context.Context, topic string, payload []byte, options PublishOptions) (PublishResult, error)

Publish creates one event and fans it out once to the current unique peer snapshot. All peers are attempted; partial failures are joined in the result.

func (*Bus) Subscribe

func (b *Bus) Subscribe(topic string, handler EventHandler) (*Subscription, error)

Subscribe registers one exact topic. Each subscription has one worker and a bounded queue, so handler concurrency and memory remain finite.

type BusConfig

type BusConfig struct {
	Sender            Sender
	Peers             PeerSource
	DefaultTTL        time.Duration
	MaxTTL            time.Duration
	MaxPayloadBytes   int
	DedupCapacity     int
	QueueCapacity     int
	MaxSubscriptions  int
	MaxFanoutPeers    int
	FanoutConcurrency int
	OnHandlerError    func(error)
	// contains filtered or unexported fields
}

BusConfig sets finite resource bounds for one in-memory event bus.

type DirectoryConfig

type DirectoryConfig struct {
	// MaxPeers bounds the number of remembered peers.
	MaxPeers int
	// MaxMetadataBytes bounds the total application metadata carried on each
	// peer record. Zero uses the default bound.
	MaxMetadataBytes int
	// contains filtered or unexported fields
}

DirectoryConfig sets the finite resource bounds for one peer directory.

type Event

type Event struct {
	ID          EventID
	Topic       string
	PublishedAt time.Time
	TTL         time.Duration
	ContentType string
	Payload     []byte
}

Event is the transport-neutral event data. PublishedAt and other serialized fields are metadata; ReceivedEvent.Sender is the authoritative peer.

type EventHandler

type EventHandler func(context.Context, ReceivedEvent) error

EventHandler handles one event. Duplicate delivery remains possible across process restarts and must be safe at the application boundary.

type EventID

type EventID [16]byte

EventID is a publisher-generated 128-bit identifier used for bounded duplicate suppression. It is not an authority token.

func (EventID) String

func (id EventID) String() string

type Handler

type Handler func(context.Context, ReceivedMessage) error

Handler receives an authenticated direct message.

type Node

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

Node owns authenticated peer state and composes it with direct messaging and bounded pub/sub. Transport discovery remains advisory until a realm proof promotes the candidate into Peers.

Example
package main

import (
	"context"
	"errors"
	"fmt"
	"sync"

	"github.com/mytecor/meshbus"
)

type externalNetwork struct {
	mu    sync.RWMutex
	nodes map[meshbus.PeerID]*externalTransport
}

type externalTransport struct {
	network  *externalNetwork
	id       meshbus.PeerID
	handler  meshbus.Handler
	observer meshbus.PeerObserver
	started  bool
}

func newExternalNetwork() *externalNetwork {
	return &externalNetwork{nodes: make(map[meshbus.PeerID]*externalTransport)}
}

func (n *externalNetwork) factory(identity byte, captured **externalTransport) meshbus.TransportFactory {
	return func(handler meshbus.Handler) (meshbus.NodeTransport, error) {
		id, err := meshbus.NewPeerID([]byte{identity})
		if err != nil {
			return nil, err
		}
		transport := &externalTransport{network: n, id: id, handler: handler}
		n.mu.Lock()
		n.nodes[id] = transport
		n.mu.Unlock()
		*captured = transport
		return transport, nil
	}
}

func (t *externalTransport) Identity() meshbus.PeerID { return t.id }

func (t *externalTransport) SetPeerObserver(observer meshbus.PeerObserver) error {
	t.observer = observer
	return nil
}

func (t *externalTransport) Start(context.Context) error {
	t.started = true
	return nil
}

func (t *externalTransport) Close() error {
	t.started = false
	return nil
}

func (t *externalTransport) SendMessage(ctx context.Context, peer meshbus.PeerID, payload []byte) error {
	t.network.mu.RLock()
	target := t.network.nodes[peer]
	t.network.mu.RUnlock()
	if !t.started || target == nil || !target.started {
		return errors.New("peer unavailable")
	}
	if err := t.observer.Authenticated(target.id); err != nil {
		return err
	}
	if err := target.observer.Authenticated(t.id); err != nil {
		return err
	}
	message, err := meshbus.NewReceivedMessage(t.id.Bytes(), payload)
	if err != nil {
		return err
	}
	return target.handler(ctx, message)
}

func (t *externalTransport) discover(peer *externalTransport) error {
	return t.observer.Discovered(meshbus.Peer{ID: peer.id})
}

func main() {
	network := newExternalNetwork()
	var transportA, transportB *externalTransport
	nodeA, _ := meshbus.NewNode(meshbus.NodeConfig{Transport: network.factory(1, &transportA)})
	nodeB, _ := meshbus.NewNode(meshbus.NodeConfig{Transport: network.factory(2, &transportB)})
	defer nodeA.Close()
	defer nodeB.Close()
	ctx := context.Background()
	_ = nodeA.Start(ctx)
	_ = nodeB.Start(ctx)
	_ = transportA.discover(transportB)
	_ = transportB.discover(transportA)
	_ = nodeA.Send(ctx, nodeB.Identity(), []byte("authenticate"))

	received := make(chan meshbus.ReceivedEvent, 1)
	_, _ = nodeB.Subscribe("example.event", func(_ context.Context, event meshbus.ReceivedEvent) error {
		received <- event
		return nil
	})
	_, _ = nodeA.Publish(ctx, "example.event", []byte("hello"), meshbus.PublishOptions{})
	event := <-received
	fmt.Printf("%s: %s\n", event.Sender, event.Payload)
}
Output:
01: hello

func NewNode

func NewNode(config NodeConfig) (*Node, error)

func (*Node) Authenticated

func (n *Node) Authenticated(id PeerID) error

func (*Node) Close

func (n *Node) Close() error

func (*Node) Discovered

func (n *Node) Discovered(peer Peer) error

func (*Node) DiscoveredPeers

func (n *Node) DiscoveredPeers() []Peer

DiscoveredPeers returns advisory presence candidates. They are usable by Send to initiate authentication but excluded from Peers and pub/sub fan-out.

func (*Node) Identity

func (n *Node) Identity() PeerID

Identity returns this node's transport-authenticated identity.

func (*Node) Peers

func (n *Node) Peers() []Peer

Peers returns only identities that completed realm authentication.

func (*Node) Publish

func (n *Node) Publish(ctx context.Context, topic string, payload []byte, options PublishOptions) (PublishResult, error)

func (*Node) Send

func (n *Node) Send(ctx context.Context, peer PeerID, payload []byte) error

func (*Node) SendMessage

func (n *Node) SendMessage(ctx context.Context, peer PeerID, payload []byte) error

SendMessage implements Sender for the composed Bus. It is equivalent to Send and remains peer-addressed.

func (*Node) Start

func (n *Node) Start(ctx context.Context) error

func (*Node) Subscribe

func (n *Node) Subscribe(topic string, handler EventHandler) (*Subscription, error)

type NodeConfig

type NodeConfig struct {
	Transport     TransportFactory
	DirectHandler Handler
	Bus           BusConfig
	Directory     DirectoryConfig
	PeerTTL       time.Duration
	SweepInterval time.Duration
	OnPeerError   func(error)
}

type NodeTransport

type NodeTransport interface {
	Sender
	Identity() PeerID
	SetPeerObserver(PeerObserver) error
	Start(context.Context) error
	Close() error
}

NodeTransport supplies authenticated direct delivery and reports discovery separately from successful realm authentication.

type Peer

type Peer struct {
	// ID is the authenticated public identity established by the transport.
	// It never comes from application metadata on the wire.
	ID PeerID
	// Metadata holds optional bounded application metadata (for example
	// advisory os/arch hints). Presence and metadata are advisory only.
	Metadata map[string]string
	// Hops is an optional advisory path metric such as hop count.
	Hops uint8
	// LastSeen is the local time the peer was last (re)discovered.
	LastSeen time.Time
}

Peer describes one authenticated realm peer learned through discovery.

type PeerDirectory

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

PeerDirectory is a bounded, observable record of realm peers learned through discovery. It remembers peers keyed by authenticated PeerID, updates an existing peer, resolves a PeerID to its current transport route, returns a bounded immutable snapshot for pub/sub fan-out, and removes or expires stale entries. It is safe for concurrent use.

func NewPeerDirectory

func NewPeerDirectory(config DirectoryConfig) (*PeerDirectory, error)

NewPeerDirectory creates an empty directory with the given resource bounds.

func (*PeerDirectory) ExpireStale

func (d *PeerDirectory) ExpireStale(olderThan time.Duration) int

ExpireStale removes peers whose LastSeen is older than the given age. With a zero age only zero-valued LastSeen records age out. It updates LastSeen on every Remember, so idle peers are what age out.

func (*PeerDirectory) Get

func (d *PeerDirectory) Get(id PeerID) (Peer, bool)

Get returns an immutable copy of one discovered peer.

func (*PeerDirectory) IDs

func (d *PeerDirectory) IDs() []PeerID

IDs returns the bounded snapshot of peer identities used for pub/sub fan-out, ordered by identity.

func (*PeerDirectory) Len

func (d *PeerDirectory) Len() int

Len returns the number of remembered peers.

func (*PeerDirectory) Peers

func (d *PeerDirectory) Peers() []Peer

Peers returns an immutable, copy-safe snapshot of the known realm peers, ordered by identity for deterministic fan-out. The result is capped at the configured MaxPeers. Discovery is advisory only.

func (*PeerDirectory) Remember

func (d *PeerDirectory) Remember(peer Peer) error

Remember adds or updates one discovered peer. An existing peer with the same authenticated PeerID is updated in place instead of creating a duplicate. The route may change while the identity stays stable. Returns ErrPeerLimit when the directory is at capacity and the identity is new, and ErrMetadataLimit when application metadata exceeds the configured bound.

func (*PeerDirectory) Remove

func (d *PeerDirectory) Remove(id PeerID)

Remove forgets one peer. It is a no-op when the peer is unknown.

type PeerID

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

PeerID is an opaque identity established by the transport. Its binary form is kept immutable so application payloads cannot replace sender authority.

func NewPeerID

func NewPeerID(value []byte) (PeerID, error)

NewPeerID copies a non-empty transport-authenticated identity.

func (PeerID) Bytes

func (p PeerID) Bytes() []byte

Bytes returns a copy of the transport identity.

func (PeerID) IsZero

func (p PeerID) IsZero() bool

IsZero reports whether the identity was not constructed from authenticated transport bytes.

func (PeerID) String

func (p PeerID) String() string

String returns the identity in a stable, non-secret hexadecimal form.

type PeerObserver

type PeerObserver interface {
	Discovered(Peer) error
	Authenticated(PeerID) error
}

PeerObserver receives transport discovery and successful realm authentication separately. Discovery is advisory; only Authenticated proves that a transport identity holds the realm key.

type PeerSource

type PeerSource interface {
	Peers() []PeerID
}

PeerSource returns a snapshot of authenticated peers. Publish never forwards an event beyond this one-hop snapshot.

type PeerSourceFunc

type PeerSourceFunc func() []PeerID

PeerSourceFunc adapts a function to PeerSource.

func (PeerSourceFunc) Peers

func (f PeerSourceFunc) Peers() []PeerID

type PublishOptions

type PublishOptions struct {
	TTL         time.Duration
	ContentType string
	// RemoteOnly suppresses delivery to local subscriptions when publishing
	// through Node. Direct Bus users have no local identity and remain remote-only.
	RemoteOnly bool
	// contains filtered or unexported fields
}

PublishOptions controls bounded event lifetime and payload interpretation.

type PublishResult

type PublishResult struct {
	ID             EventID
	Attempted      int
	Delivered      int
	Failed         map[PeerID]error
	LocalDelivered bool
	LocalError     error
}

PublishResult reports best-effort delivery without claiming remote acknowledgement or exactly-once semantics.

type ReceivedEvent

type ReceivedEvent struct {
	Event
	Sender     PeerID
	ReceivedAt time.Time
}

ReceivedEvent pairs an event with transport-authenticated delivery metadata.

type ReceivedMessage

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

ReceivedMessage is an opaque payload paired with its transport-authenticated sender. The message does not interpret or trust identities inside Payload.

func NewReceivedMessage

func NewReceivedMessage(sender []byte, payload []byte) (ReceivedMessage, error)

NewReceivedMessage constructs an isolated inbound message.

func (ReceivedMessage) Payload

func (m ReceivedMessage) Payload() []byte

Payload returns a copy of the opaque application bytes.

func (ReceivedMessage) Sender

func (m ReceivedMessage) Sender() PeerID

Sender returns the authenticated transport peer.

type Sender

type Sender interface {
	SendMessage(context.Context, PeerID, []byte) error
}

Sender sends opaque application bytes to an authenticated peer. Transport route resolution stays behind this boundary.

type Subscription

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

Subscription is one exact-topic handler with a bounded private queue.

func (*Subscription) Close

func (s *Subscription) Close() error

Close removes the subscription and cancels its handler context.

func (*Subscription) Done

func (s *Subscription) Done() <-chan struct{}

Done closes after the subscription worker stops.

type TransportFactory

type TransportFactory func(Handler) (NodeTransport, error)

Directories

Path Synopsis
Package realm defines a transport-independent shared-secret membership realm.
Package realm defines a transport-independent shared-secret membership realm.
Package rns implements a reusable, authenticated Reticulum adapter for meshbus.
Package rns implements a reusable, authenticated Reticulum adapter for meshbus.

Jump to

Keyboard shortcuts

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