Skip to content

Repository files navigation

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 ./...

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages