event

package module
v1.0.2 Latest Latest
Warning

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

Go to latest
Published: Aug 23, 2026 License: MIT Imports: 14 Imported by: 0

README

github.com/adnilis/event

基于 github.com/adnilis/actor 的进程内通用事件派发库。它把事件正文保存在库自有的固定容量分区队列中,Actor 邮箱只接收可合并的唤醒 token,从而让背压、内存上界、分区并行和失败语义可验证。

快速开始

dispatcher, err := event.New(event.Config{
    QueueCapacity: 64,
    Partitions:    4,
    PublishPolicy: event.QueueBlock,
    MaxAttempts:   3,
})
if err != nil {
    return err
}
defer dispatcher.Shutdown(context.Background())

_, err = dispatcher.Subscribe("orders", event.HandlerFunc(func(ctx context.Context, e event.Event) error {
    // Handler 必须把 Payload 当作只读数据;需要独立副本时配置 Codec。
    return handleOrder(ctx, e)
}))
if err != nil {
    return err
}

result, err := dispatcher.Publish(ctx, event.Event{
    ID: "order-1", Topic: "orders", Key: "account-1", Payload: payload,
})
_ = result // Accepted 只表示已进入订阅队列,不表示 Handler 已成功。
return err

可运行示例位于 examples/basic/main.go。

核心语义

  • Topic 用于订阅过滤;非空 Key 使用固定 FNV-1a 分区,同一订阅同一 Key FIFO,不同分区可以并行;空 Key 使用轮询。订阅创建后分区数固定。
  • 每个订阅分区都有 QueueCapacity 硬上限。默认 QueueReject,满队列返回 ErrQueueFull;显式 QueueBlock 会等待空间,但必须传入可取消/有截止时间的 Context。
  • 近似内存预算是“订阅数 × 分区数 × QueueCapacity”个队列槽位,重试任务与普通事件共享该分区预算;Actor 邮箱不会承载事件正文。
  • Publish 的 Accepted 只表示入队成功。处理结果可能是成功、失败、再次重试、死信或关闭取消。系统提供至少一次投递语义,重试可能重复调用同一个 Event.ID,业务必须用 Event.ID 或业务幂等键保护副作用。
  • Handler panic 会被隔离;可重试错误按 MaxAttempts 和退避策略调度,最终失败进入 DeadLetterSink。死信包含原事件、最后错误、尝试次数、panic 元数据和时间。

生命周期、指标和健康

状态为 Running → Draining → Closed。Shutdown(ctx) 进入 Draining 后拒绝新发布/订阅,等待已接受事件排空;Context 超时会取消内部处理 Context、停止可停止资源并返回 Context 错误。Shutdown 可安全重复调用。

Metrics() 返回原子值快照:Published、Accepted、QueueFull、Processed、Failed、Retried、DeadLettered、Panicked、QueueDepth、InFlight 和 MaxLatency。Health() 返回 HealthReady、HealthDraining 或 HealthNotReady 及原因,不暴露内部 Actor、队列或可变注册表。

Codec 与 Store 扩展边界

配置 Config.Codec 后,库在发布前编码 Payload,并为每个 Handler 独立解码;Encode 错误返回 ErrCodec,Decode 错误进入死信。未配置 Codec 时不做隐式深拷贝,调用方必须在 Publish 返回后把 Payload 当作只读。

Store/StoredEvent 只定义未来外部持久化适配边界,首期不会调用 Store,也不提供 WAL、跨进程恢复或崩溃恢复。

错误与首期非目标

常见错误包括 ErrInvalidEvent、ErrNoSubscriber、ErrQueueFull、ErrClosed、ErrActorUnavailable 和 ErrCodec。

首期只保证单进程生命周期内的并发安全、有界内存、失败隔离、重试、死信和优雅关闭;不实现 Kafka/Redis/数据库后端、WAL、跨进程路由、节点发现、集群一致性、全局 exactly-once 或业务幂等存储。

验证

go test ./...
CGO_ENABLED=1 go test -race ./...
go vet ./...
go test -bench=. -benchmem ./...

Documentation

Overview

Package event 提供基于 Actor 的进程内事件派发能力。

Package event 提供基于 Actor 的进程内事件派发能力。

Package event 提供基于 Actor 的进程内事件派发能力。

Package event 提供基于 Actor 的进程内事件派发能力。

Package event 提供基于 Actor 的进程内事件派发能力。

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrInvalidConfig 表示配置为空、越界或相互冲突。
	ErrInvalidConfig = errors.New("event: invalid config")
	// ErrInvalidEvent 表示事件缺少必需的 ID 或 Topic。
	ErrInvalidEvent = errors.New("event: invalid event")
	// ErrNoSubscriber 表示目标 Topic 当前没有订阅者。
	ErrNoSubscriber = errors.New("event: no subscriber")
	// ErrQueueFull 表示订阅队列已达到容量上限。
	ErrQueueFull = errors.New("event: queue full")
	// ErrClosed 表示派发器或订阅已经关闭。
	ErrClosed = errors.New("event: dispatcher closed")
	// ErrNotReady 表示运行时资源尚未准备好。
	ErrNotReady = errors.New("event: dispatcher not ready")
	// ErrActorUnavailable 表示 Actor 运行时无法接收唤醒消息。
	ErrActorUnavailable = errors.New("event: actor unavailable")
	// ErrCodec 表示编码或解码过程失败。
	ErrCodec = errors.New("event: codec failure")
)

Functions

This section is empty.

Types

type Codec

type Codec interface {
	// Encode 在事件进入订阅队列前编码 Payload。
	Encode(context.Context, any) ([]byte, error)
	// Decode 为一次 Handler 调用生成独立的 Payload。
	Decode(context.Context, []byte) (any, error)
}

Codec 负责序列化和反序列化 Event.Payload。

派发器在入队前只编码一次,并为每次 Handler 调用独立解码。 Encode 返回后,调用方必须将返回的字节切片视为只读数据。

type Config

type Config struct {
	// QueueCapacity 是每个订阅分区的队列容量上限。
	QueueCapacity int
	// Partitions 是每个订阅创建的固定分区数。
	Partitions int
	// PublishPolicy 定义队列满时的发布行为。
	PublishPolicy QueuePolicy
	// MaxAttempts 是单个事件允许的最大处理次数。
	MaxAttempts int
	// HandlerTimeout 限制单次 Handler 调用时长,零表示不限制。
	HandlerTimeout time.Duration
	// RetryInitialBackoff 是第一次重试的退避时长。
	RetryInitialBackoff time.Duration
	// RetryMaxBackoff 是重试退避的最大时长。
	RetryMaxBackoff time.Duration
	// ErrorClassifier 判断 Handler 错误是否可重试。
	ErrorClassifier ErrorClassifier
	// DeadLetterSink 接收最终失败或不可重试的事件。
	DeadLetterSink DeadLetterSink
	// ActorSystemName 是底层 Actor 系统的名称前缀。
	ActorSystemName string
	// Codec 定义发布编码和 Handler 解码策略。
	Codec Codec
	// Store 预留外部持久化适配边界,首期不会主动调用。
	Store Store
}

Config 是 EventBus 的默认运行配置。

func DefaultConfig

func DefaultConfig() Config

DefaultConfig 返回一份可直接使用的默认配置。

type DeadLetter

type DeadLetter struct {
	// Event 是失败事件的副本。
	Event Event
	// LastError 是最后一次处理错误。
	LastError error
	// Attempts 是已经执行的处理次数。
	Attempts int
	// Panicked 表示失败是否由 Handler panic 引起。
	Panicked bool
	// PanicValue 保存捕获到的 panic 值。
	PanicValue any
	// FailedAt 是进入死信时的时间。
	FailedAt time.Time
}

DeadLetter 保存无法成功处理的事件及其最终失败信息。

type DeadLetterSink

type DeadLetterSink interface {
	// Handle 处理一条死信记录。
	Handle(context.Context, DeadLetter) error
}

DeadLetterSink 在重试耗尽或错误不可重试时接收死信。

type DeadLetterSinkFunc

type DeadLetterSinkFunc func(context.Context, DeadLetter) error

DeadLetterSinkFunc 将普通函数适配为 DeadLetterSink。

func (DeadLetterSinkFunc) Handle

func (f DeadLetterSinkFunc) Handle(ctx context.Context, dead DeadLetter) error

Handle 调用底层函数处理死信。

type DeliveryStatus

type DeliveryStatus uint8

DeliveryStatus 表示一次发布的总体投递状态。

const (
	// DeliveryAccepted 表示事件已被所有目标订阅接收。
	DeliveryAccepted DeliveryStatus = iota + 1
	// DeliveryRejected 表示至少一个目标订阅拒绝接收。
	DeliveryRejected
	// DeliveryCompleted 表示事件已完成处理。
	DeliveryCompleted
	// DeliveryFailed 表示事件处理失败但尚未进入最终状态。
	DeliveryFailed
	// DeliveryDeadLettered 表示事件已进入死信流程。
	DeliveryDeadLettered
	// DeliveryCancelled 表示发布因关闭或取消而终止。
	DeliveryCancelled
)

type Dispatcher

type Dispatcher interface {
	Publisher
	Subscribe(string, Handler, ...SubscribeOption) (Subscription, error)
	Shutdown(context.Context) error
	State() State
	Metrics() MetricsSnapshot
	Health() HealthSnapshot
}

Dispatcher 是事件发布、订阅、生命周期和观测能力的完整接口。

type ErrorClassifier

type ErrorClassifier func(error) bool

ErrorClassifier 判断 Handler 错误是否允许重试。

type Event

type Event struct {
	// ID 是业务侧用于幂等和追踪的事件标识。
	ID string
	// Topic 是订阅匹配使用的主题。
	Topic string
	// Key 决定分区;相同 Key 在同一订阅内保持 FIFO。
	Key string
	// Payload 是业务负载;配置 Codec 后会按 Handler 独立解码。
	Payload any
	// Metadata 保存可选的字符串元数据。
	Metadata map[string]string
	// CreatedAt 是事件创建时间,用于延迟指标计算。
	CreatedAt time.Time
}

Event 是传递给 Handler 的不可变事件信封。 未配置 Codec 时,Payload 由调用方保证在 Publish 返回后只读。

func (Event) Validate

func (e Event) Validate() error

Validate 校验事件是否包含必需的 ID 和 Topic。

type EventBus

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

EventBus 是事件派发器的默认实现。

func New

func New(config Config, opts ...Option) (*EventBus, error)

New 创建并启动一个 EventBus。

func (*EventBus) Health

func (d *EventBus) Health() HealthSnapshot

Health 返回当前时刻的就绪状态快照。

func (*EventBus) Metrics

func (d *EventBus) Metrics() MetricsSnapshot

Metrics 返回派发器计数器和水位的原子副本。

func (*EventBus) Publish

func (d *EventBus) Publish(ctx context.Context, e Event) (PublishResult, error)

Publish 校验事件并将其投递到 Topic 的所有订阅队列。 返回成功只代表入队,不代表 Handler 已经完成处理。

func (*EventBus) Shutdown

func (d *EventBus) Shutdown(ctx context.Context) error

Shutdown 停止接收新工作,等待已接收工作排空,并关闭底层运行时。 若 Context 超时,则取消内部处理并返回 Context 错误。

func (*EventBus) State

func (d *EventBus) State() State

State 返回派发器当前生命周期状态。

func (*EventBus) Subscribe

func (d *EventBus) Subscribe(topic string, handler Handler, opts ...SubscribeOption) (Subscription, error)

Subscribe 为 Topic 注册 Handler,并启动该订阅的所有分区。

type Handler

type Handler interface {
	// Handle 执行业务处理逻辑;返回错误时由派发器决定重试或进入死信。
	Handle(context.Context, Event) error
}

Handler 消费一个事件并返回处理结果。

type HandlerFunc

type HandlerFunc func(context.Context, Event) error

HandlerFunc 将普通函数适配为 Handler。

func (HandlerFunc) Handle

func (f HandlerFunc) Handle(ctx context.Context, e Event) error

Handle 调用底层函数完成事件处理。

type HealthSnapshot

type HealthSnapshot struct {
	// Status 是当前健康状态。
	Status HealthStatus
	// Reason 是状态的简要原因。
	Reason string
}

HealthSnapshot 是 Health 返回的健康状态快照。

type HealthStatus

type HealthStatus uint8

HealthStatus 表示派发器对外提供服务的健康状态。

const (
	// HealthReady 表示派发器可以接收新工作。
	HealthReady HealthStatus = iota + 1
	// HealthDraining 表示派发器正在排空已接收工作。
	HealthDraining
	// HealthNotReady 表示派发器不可接收新工作。
	HealthNotReady
)

type MetricsSnapshot

type MetricsSnapshot struct {
	// Published 是通过事件校验并进入发布流程的次数。
	Published uint64
	// Accepted 是成功进入订阅队列的投递次数。
	Accepted uint64
	// QueueFull 是因队列满而拒绝的投递次数。
	QueueFull uint64
	// Processed 是 Handler 成功完成的次数。
	Processed uint64
	// Failed 是 Handler 返回错误或 panic 的次数。
	Failed uint64
	// Retried 是已安排重试的次数。
	Retried uint64
	// DeadLettered 是进入死信流程的次数。
	DeadLettered uint64
	// Panicked 是捕获到 Handler panic 的次数。
	Panicked uint64
	// QueueDepth 是当前所有订阅队列中的事件数量。
	QueueDepth int
	// InFlight 是当前正在执行 Handler 的事件数量。
	InFlight int
	// MaxLatency 是观测到的最大处理延迟。
	MaxLatency time.Duration
}

MetricsSnapshot 是 Metrics 返回的原子指标快照。

type Option

type Option func(*Config) error

Option 修改 EventBus 配置。

type PublishResult

type PublishResult struct {
	// Accepted 是成功进入订阅队列的数量。
	Accepted int
	// Rejected 是未能进入订阅队列的数量。
	Rejected int
	// Status 是本次发布的总体状态。
	Status DeliveryStatus
}

PublishResult 描述一次 Publish 对各订阅的入队结果。

type Publisher

type Publisher interface {
	// Publish 将事件投递到匹配的订阅队列。
	Publish(context.Context, Event) (PublishResult, error)
}

Publisher 定义发布事件的最小接口。

type QueuePolicy

type QueuePolicy uint8

QueuePolicy 定义队列满时 Publish 的处理策略。

const (
	// QueueReject 立即返回 ErrQueueFull。
	QueueReject QueuePolicy = iota
	// QueueBlock 等待队列空间,并受 Context 取消或截止时间约束。
	QueueBlock
)

type State

type State uint8

State 表示派发器生命周期状态。

const (
	// StateRunning 表示正常接收发布和订阅请求。
	StateRunning State = iota + 1
	// StateDraining 表示拒绝新工作并等待已接收工作排空。
	StateDraining
	// StateClosed 表示派发器已经关闭。
	StateClosed
)

func (State) String

func (s State) String() string

String 返回生命周期状态的可读名称。

type Store

type Store interface {
	// Save 持久化一个事件记录。
	Save(context.Context, StoredEvent) error
	// Load 根据事件 ID 加载一个事件记录。
	Load(context.Context, string) (StoredEvent, error)
}

Store 描述可选的外部持久化集成边界。 首期派发器不会主动调用该接口。

type StoredEvent

type StoredEvent struct {
	// Event 保存事件信封。
	Event Event
	// EncodedPayload 保存 Codec 编码后的负载。
	EncodedPayload []byte
}

StoredEvent 是编码事件未来接入持久化存储时使用的边界结构。

首期实现不会调用 Store,也不提供崩溃恢复;未来接入外部可靠适配器时 可以复用该结构。

type SubscribeConfig

type SubscribeConfig struct {
	// QueueCapacity 是该订阅每个分区的队列容量。
	QueueCapacity int
	// Partitions 是该订阅的固定分区数。
	Partitions int
	// PublishPolicy 定义该订阅队列满时的行为。
	PublishPolicy QueuePolicy
	// HandlerTimeout 限制该订阅 Handler 的单次处理时长。
	HandlerTimeout time.Duration
	// MaxAttempts 是该订阅的最大处理次数。
	MaxAttempts int
	// RetryInitialBackoff 是该订阅第一次重试的退避时长。
	RetryInitialBackoff time.Duration
	// RetryMaxBackoff 是该订阅重试退避的上限。
	RetryMaxBackoff time.Duration
	// ErrorClassifier 判断该订阅错误是否可重试。
	ErrorClassifier ErrorClassifier
	// DeadLetterSink 接收该订阅的最终失败事件。
	DeadLetterSink DeadLetterSink
	// Codec 定义该订阅的 Payload 解码方式。
	Codec Codec
}

SubscribeConfig 是单个订阅可覆盖的运行配置。

type SubscribeOption

type SubscribeOption func(*SubscribeConfig) error

SubscribeOption 修改单个订阅配置。

type Subscription

type Subscription interface {
	// ID 返回订阅唯一标识。
	ID() string
	// Topic 返回订阅主题。
	Topic() string
	// Close 停止订阅并排空或取消其队列中的工作。
	Close() error
}

Subscription 表示一个 Topic 订阅及其生命周期控制句柄。

Directories

Path Synopsis
examples
basic command
basic 示例展示最小可运行的事件发布、订阅和优雅关闭流程。
basic 示例展示最小可运行的事件发布、订阅和优雅关闭流程。
internal
actorbridge
Package actorbridge 隔离外部 Actor 运行时与事件派发器实现。
Package actorbridge 隔离外部 Actor 运行时与事件派发器实现。
handler
Package handler 提供 Handler 调用、panic 隔离和错误分类工具。
Package handler 提供 Handler 调用、panic 隔离和错误分类工具。
queue
Package queue 提供固定容量、并发安全的泛型队列。
Package queue 提供固定容量、并发安全的泛型队列。
retry
Package retry 提供重试退避和延迟任务调度能力。
Package retry 提供重试退避和延迟任务调度能力。
wakeup
Package wakeup 提供合并 Actor 唤醒请求的轻量状态门。
Package wakeup 提供合并 Actor 唤醒请求的轻量状态门。

Jump to

Keyboard shortcuts

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