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 ¶
- Constants
- Variables
- type Bus
- func (b *Bus) Close() error
- func (b *Bus) HandleMessage(_ context.Context, message ReceivedMessage) (bool, error)
- func (b *Bus) Handler(next Handler) Handler
- func (b *Bus) Publish(ctx context.Context, topic string, payload []byte, options PublishOptions) (PublishResult, error)
- func (b *Bus) Subscribe(topic string, handler EventHandler) (*Subscription, error)
- type BusConfig
- type DirectoryConfig
- type Event
- type EventHandler
- type EventID
- type Handler
- type Node
- func (n *Node) Authenticated(id PeerID) error
- func (n *Node) Close() error
- func (n *Node) Discovered(peer Peer) error
- func (n *Node) DiscoveredPeers() []Peer
- func (n *Node) Identity() PeerID
- func (n *Node) Peers() []Peer
- func (n *Node) Publish(ctx context.Context, topic string, payload []byte, options PublishOptions) (PublishResult, error)
- func (n *Node) Send(ctx context.Context, peer PeerID, payload []byte) error
- func (n *Node) SendMessage(ctx context.Context, peer PeerID, payload []byte) error
- func (n *Node) Start(ctx context.Context) error
- func (n *Node) Subscribe(topic string, handler EventHandler) (*Subscription, error)
- type NodeConfig
- type NodeTransport
- type Peer
- type PeerDirectory
- func (d *PeerDirectory) ExpireStale(olderThan time.Duration) int
- func (d *PeerDirectory) Get(id PeerID) (Peer, bool)
- func (d *PeerDirectory) IDs() []PeerID
- func (d *PeerDirectory) Len() int
- func (d *PeerDirectory) Peers() []Peer
- func (d *PeerDirectory) Remember(peer Peer) error
- func (d *PeerDirectory) Remove(id PeerID)
- type PeerID
- type PeerObserver
- type PeerSource
- type PeerSourceFunc
- type PublishOptions
- type PublishResult
- type ReceivedEvent
- type ReceivedMessage
- type Sender
- type Subscription
- type TransportFactory
Examples ¶
Constants ¶
const MaxPeerIDBytes = 256
MaxPeerIDBytes bounds identity material retained by the generic layer. Individual transports may impose a smaller bound.
const MaxTopicBytes = 128
Variables ¶
var ( ErrInvalidPeerID = errors.New("invalid peer identity") ErrInvalidMessage = errors.New("invalid direct message") )
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") )
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") )
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 (*Bus) HandleMessage ¶
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) 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.
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 (*Node) Discovered ¶
func (*Node) DiscoveredPeers ¶
DiscoveredPeers returns advisory presence candidates. They are usable by Send to initiate authentication but excluded from Peers and pub/sub fan-out.
func (*Node) Publish ¶
func (n *Node) Publish(ctx context.Context, topic string, payload []byte, options PublishOptions) (PublishResult, error)
func (*Node) SendMessage ¶
SendMessage implements Sender for the composed Bus. It is equivalent to Send and remains peer-addressed.
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.
type PeerObserver ¶
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 ¶
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 ¶
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. |