events

package
v0.0.0-...-51ec7a6 Latest Latest
Warning

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

Go to latest
Published: Sep 23, 2026 License: Apache-2.0 Imports: 8 Imported by: 0

Documentation

Overview

Package events provides typed, in-memory application events.

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrFrozen          = errors.New("event bus is frozen")
	ErrTypeMismatch    = errors.New("event definition type mismatch")
	ErrUndeclared      = errors.New("event topic is not declared")
	ErrReentrant       = errors.New("same-topic event publication is reentrant")
	ErrHandlerPanic    = errors.New("event handler panicked")
	ErrHandlerTimeout  = errors.New("event handler timed out")
	ErrSaturated       = errors.New("event handler concurrency limit reached")
	ErrSubscriberLimit = errors.New("event topic subscriber limit reached")
)
View Source
var ErrNoDeferral = errors.New("event deferral is unavailable")

ErrNoDeferral reports that no after-commit boundary is available, so the caller should dispatch immediately instead.

Functions

func Declare

func Declare[T any](bus *Bus, definition Definition[T]) error

Declare establishes a topic payload type before engine subscribers or publishers may use it. Direct event users may continue to rely on the historical lazy topic creation in Subscribe and Publish.

func Publish

func Publish[T any](ctx context.Context, bus *Bus, definition Definition[T], payload T) error

Publish invokes a topic's handlers serially for this publication. Separate publications may execute concurrently.

func PublishAfterCommit

func PublishAfterCommit[T any](ctx context.Context, bus *Bus, definition Definition[T], payload T) error

PublishAfterCommit dispatches subscribers once the caller's database transaction commits, keeping subscriber latency out of the transaction and off the pooled connection it holds.

Subscribers observe committed state, so a failure cannot roll the transaction back; it surfaces through the host's after-commit error path. When no transaction is active, this falls back to immediate publication and behaves exactly like Publish.

func RequireDeclared

func RequireDeclared[T any](bus *Bus, definition Definition[T]) error

RequireDeclared verifies that a topic was explicitly declared and that its payload type matches the caller's definition.

func Subscribe

func Subscribe[T any](bus *Bus, definition Definition[T], subscriber string, handler Handler[T]) error

Subscribe registers one named handler. Registration order is retained for all handlers on the same topic.

Types

type Bus

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

Bus stores in-memory event subscriptions. It does not persist, retry, or copy published payloads.

func New

func New(options ...Option) *Bus

func NewBus

func NewBus() *Bus

func (*Bus) Freeze

func (bus *Bus) Freeze() error

Freeze prevents new subscriptions. Existing subscriptions remain usable.

func (*Bus) Frozen

func (bus *Bus) Frozen() bool

type Deferrer

type Deferrer func(context.Context, func(context.Context) error) error

Deferrer schedules work to run after the caller's database transaction commits. It is injected so this package keeps no database dependency. Returning ErrNoDeferral means no transaction is active and the caller should dispatch immediately.

type Definition

type Definition[T any] struct {
	Name    string
	Version int
}

Definition identifies one version of a typed event. Name and Version form the runtime topic identity; T is checked when a topic is subscribed to or published.

func Define

func Define[T any](name string, version int) (Definition[T], error)

func MustDefine

func MustDefine[T any](name string, version int) Definition[T]

func (Definition[T]) Validate

func (definition Definition[T]) Validate() error

type Delivery

type Delivery struct {
	Topic      string
	Subscriber string
	Duration   time.Duration
	Err        error
}

type Handler

type Handler[T any] func(context.Context, T) error

Handler is invoked synchronously in subscription order.

type Observer

type Observer func(context.Context, Delivery)

type Option

type Option func(*busSettings)

func WithDeferrer

func WithDeferrer(deferrer Deferrer) Option

WithDeferrer supplies the after-commit scheduler used by PublishAfterCommit.

func WithHandlerTimeout

func WithHandlerTimeout(timeout time.Duration) Option

func WithMaxConcurrentHandlers

func WithMaxConcurrentHandlers(limit int) Option

func WithMaxSubscribersPerTopic

func WithMaxSubscribersPerTopic(limit int) Option

func WithObserver

func WithObserver(observer Observer) Option

Jump to

Keyboard shortcuts

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