Files
summercms/modules/festival/bus.go
Jakub Zych ac1f6d14f4 refactor(10.2-01): nest festival module
- Move festival under modules\n- Update framework and example importers
2026-09-28 02:16:36 +02:00

150 lines
3.4 KiB
Go

package festival
import (
"context"
"errors"
"fmt"
"reflect"
"sort"
"sync"
)
// Collectable is implemented by events that expose a mergeable Collect payload.
type Collectable interface {
Collected() map[string]any
}
// Handleable is implemented by events that can stop UntilHandled.
type Handleable interface {
IsHandled() bool
}
type listener struct {
owner string
priority int
order int
fn func(context.Context, any) error
}
// Bus is an app-owned typed event dispatcher. Dispatch is synchronous on
// the caller goroutine.
type Bus struct {
mu sync.Mutex
seq int
byType map[reflect.Type][]listener
}
// New returns an empty bus.
func New() *Bus {
return &Bus{byType: make(map[reflect.Type][]listener)}
}
// Listen registers fn at priority 0.
func (b *Bus) Listen[T any](owner string, fn func(context.Context, T) error) {
b.ListenPriority(owner, 0, fn)
}
// ListenPriority registers fn. Higher priority runs first; equal priority
// keeps registration order.
func (b *Bus) ListenPriority[T any](owner string, priority int, fn func(context.Context, T) error) {
if b == nil || fn == nil {
return
}
key := reflect.TypeFor[T]()
b.mu.Lock()
defer b.mu.Unlock()
if b.byType == nil {
b.byType = make(map[reflect.Type][]listener)
}
b.seq++
b.byType[key] = append(b.byType[key], listener{
owner: owner,
priority: priority,
order: b.seq,
fn: func(ctx context.Context, event any) error {
v, ok := event.(T)
if !ok {
return nil
}
return fn(ctx, v)
},
})
}
// Fire invokes every listener and returns errors.Join of their errors.
func (b *Bus) Fire[T any](ctx context.Context, event T) error {
var errs []error
for _, l := range b.snapshot[T]() {
if err := invoke(l, ctx, event); err != nil {
errs = append(errs, err)
}
}
return errors.Join(errs...)
}
// Collect runs every listener, merges Collected() maps (later wins on
// duplicate keys), and returns the partial payload with joined errors.
func (b *Bus) Collect[T any](ctx context.Context, event T) (map[string]any, error) {
merged := map[string]any{}
var errs []error
for _, l := range b.snapshot[T]() {
if err := invoke(l, ctx, event); err != nil {
errs = append(errs, err)
}
mergeCollected(merged, event)
}
return merged, errors.Join(errs...)
}
// UntilHandled stops at the first handled event or the first error.
func (b *Bus) UntilHandled[T any](ctx context.Context, event T) (bool, error) {
for _, l := range b.snapshot[T]() {
if err := invoke(l, ctx, event); err != nil {
return false, err
}
if h, ok := any(event).(Handleable); ok && h.IsHandled() {
return true, nil
}
}
return false, nil
}
func (b *Bus) snapshot[T any]() []listener {
if b == nil {
return nil
}
key := reflect.TypeFor[T]()
b.mu.Lock()
list := append([]listener(nil), b.byType[key]...)
b.mu.Unlock()
sort.SliceStable(list, func(i, j int) bool {
if list[i].priority != list[j].priority {
return list[i].priority > list[j].priority
}
return list[i].order < list[j].order
})
return list
}
func invoke(l listener, ctx context.Context, event any) (err error) {
defer func() {
if rec := recover(); rec != nil {
err = fmt.Errorf("festival: plugin %s panicked", l.owner)
}
}()
if ctx == nil {
ctx = context.Background()
}
return l.fn(ctx, event)
}
func mergeCollected(dst map[string]any, event any) {
c, ok := any(event).(Collectable)
if !ok {
return
}
for k, v := range c.Collected() {
dst[k] = v
}
}