bus

package
v0.6.1-rc.1 Latest Latest
Warning

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

Go to latest
Published: Oct 4, 2026 License: Apache-2.0 Imports: 2 Imported by: 0

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

Constants

This section is empty.

Variables

View Source
var Default = New(Context{})

Default is the package-level runtime used by the convenience functions.

View Source
var InstanceDisposed = Define("server.instance.disposed", struct {
	Directory string `json:"directory"`
}{})

InstanceDisposed is published to wildcard subscribers during Bus.Dispose.

Functions

func CreateID

func CreateID() string

CreateID creates an ascending evt identifier.

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 New

func New(context Context, options ...BusOption) *Bus

New constructs an instance 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

func (b *Bus) SubscribeAllCallback(callback func(Payload)) func()

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

func WithIDGenerator(createID func() string) BusOption

WithIDGenerator pins payload IDs.

type Context

type Context struct {
	Directory string
	Project   string
	Workspace string
}

Context is the instance metadata a bus is created for.

type Definition

type Definition struct {
	Type       string `json:"type"`
	Properties any    `json:"properties"`
}

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.

Jump to

Keyboard shortcuts

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