Documentation
¶
Overview ¶
Package bus is the in-process event bus. Subscriber snapshots are invoked synchronously in registration order. All mutable state is protected for concurrent publishers/subscribers.
Event-definition registry. Registry iteration preserves first-definition order.
Index ¶
- Variables
- func CreateID() string
- func Publish(def Definition, properties any, options ...PublishOptions)
- func SubscribeAllCallback(callback func(Payload)) func()
- func SubscribeCallback(def Definition, callback func(Payload)) func()
- type Bus
- func (b *Bus) Dispose()
- func (b *Bus) Publish(def Definition, properties any, options ...PublishOptions)
- func (b *Bus) Subscribe(def Definition) *Subscription
- func (b *Bus) SubscribeAll() *Subscription
- func (b *Bus) SubscribeAllCallback(callback func(Payload)) func()
- func (b *Bus) SubscribeCallback(def Definition, callback func(Payload)) func()
- type BusOption
- type Context
- type Definition
- type Payload
- type PayloadDefinition
- type PublishOptions
- type Subscription
Constants ¶
This section is empty.
Variables ¶
var Default = New(Context{})
Default is the package-level runtime used by the convenience functions.
var InstanceDisposed = Define("server.instance.disposed", struct { Directory string `json:"directory"` }{})
InstanceDisposed is published to wildcard subscribers during Bus.Dispose.
Functions ¶
func Publish ¶
func Publish(def Definition, properties any, options ...PublishOptions)
Publish emits on Default.
func SubscribeAllCallback ¶
func SubscribeAllCallback(callback func(Payload)) func()
SubscribeAllCallback subscribes on Default.
func SubscribeCallback ¶
func SubscribeCallback(def Definition, callback func(Payload)) func()
SubscribeCallback subscribes on Default.
Types ¶
type Bus ¶
type Bus struct {
// contains filtered or unexported fields
}
Bus is an instance-scoped pub/sub bus.
func (*Bus) Dispose ¶
func (b *Bus) Dispose()
Dispose publishes InstanceDisposed to wildcard subscribers only, then closes streams and makes later publishes/subscriptions inert.
func (*Bus) Publish ¶
func (b *Bus) Publish(def Definition, properties any, options ...PublishOptions)
Publish delivers to typed subscribers, then to wildcard subscribers, in that order.
func (*Bus) Subscribe ¶
func (b *Bus) Subscribe(def Definition) *Subscription
Subscribe returns a typed stream. Call Close when finished.
func (*Bus) SubscribeAll ¶
func (b *Bus) SubscribeAll() *Subscription
SubscribeAll returns a wildcard stream.
func (*Bus) SubscribeAllCallback ¶
SubscribeAllCallback subscribes to every event.
func (*Bus) SubscribeCallback ¶
func (b *Bus) SubscribeCallback(def Definition, callback func(Payload)) func()
SubscribeCallback subscribes to one event definition.
type BusOption ¶
type BusOption func(*Bus)
BusOption configures New.
func WithIDGenerator ¶
WithIDGenerator pins payload IDs.
type Definition ¶
Definition identifies an event type and carries its consumer-supplied property schema/descriptor.
func Define ¶
func Define(eventType string, properties any) Definition
Define registers and returns an event definition. Redefining a type updates its schema without changing its original insertion position.
type Payload ¶
type Payload struct {
ID string `json:"id"`
Type string `json:"type"`
Properties any `json:"properties"`
}
Payload is the wire event delivered to subscribers.
type PayloadDefinition ¶
type PayloadDefinition struct {
Type string `json:"type"`
Properties any `json:"properties"`
Identifier string `json:"identifier"`
}
PayloadDefinition describes one registered event type and its property schema.
func Payloads ¶
func Payloads() []PayloadDefinition
Payloads returns the payload descriptors in registry order.
type PublishOptions ¶
type PublishOptions struct {
ID string
}
PublishOptions lets a publisher pin the payload ID.
type Subscription ¶
type Subscription struct {
C <-chan Payload
// contains filtered or unexported fields
}
Subscription is an unbounded ordered stream subscription.
func (*Subscription) Close ¶
func (s *Subscription) Close()
Close unsubscribes and closes C after already queued events are delivered.