refactor(10.2-01): nest festival module
- Move festival under modules\n- Update framework and example importers
This commit is contained in:
149
modules/festival/bus.go
Normal file
149
modules/festival/bus.go
Normal file
@@ -0,0 +1,149 @@
|
||||
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
|
||||
}
|
||||
}
|
||||
274
modules/festival/bus_test.go
Normal file
274
modules/festival/bus_test.go
Normal file
@@ -0,0 +1,274 @@
|
||||
package festival
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
)
|
||||
|
||||
type evt struct {
|
||||
data map[string]any
|
||||
handled bool
|
||||
}
|
||||
|
||||
func (e *evt) Collected() map[string]any { return e.data }
|
||||
func (e *evt) IsHandled() bool { return e.handled }
|
||||
|
||||
func TestFireRunsAllListenersAndJoinsErrors(t *testing.T) {
|
||||
bus := New()
|
||||
var order []string
|
||||
bus.Listen[*evt]("golem15.a", func(ctx context.Context, e *evt) error {
|
||||
order = append(order, "a")
|
||||
return errors.New("err-a")
|
||||
})
|
||||
bus.Listen[*evt]("golem15.b", func(ctx context.Context, e *evt) error {
|
||||
order = append(order, "b")
|
||||
return errors.New("err-b")
|
||||
})
|
||||
err := bus.Fire(context.Background(), &evt{})
|
||||
if err == nil {
|
||||
t.Fatal("expected joined errors")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "err-a") || !strings.Contains(err.Error(), "err-b") {
|
||||
t.Fatalf("joined error = %v", err)
|
||||
}
|
||||
if strings.Join(order, ",") != "a,b" {
|
||||
t.Fatalf("order = %v, want a,b", order)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCollectMergesPayloadsLaterWinsAndKeepsPartialOnError(t *testing.T) {
|
||||
bus := New()
|
||||
bus.Listen[*evt]("golem15.a", func(ctx context.Context, e *evt) error {
|
||||
e.data = map[string]any{"k": "a", "only": "a"}
|
||||
return nil
|
||||
})
|
||||
bus.Listen[*evt]("golem15.b", func(ctx context.Context, e *evt) error {
|
||||
e.data = map[string]any{"k": "b"}
|
||||
return errors.New("collect-b")
|
||||
})
|
||||
got, err := bus.Collect(context.Background(), &evt{})
|
||||
if err == nil || !strings.Contains(err.Error(), "collect-b") {
|
||||
t.Fatalf("expected collect error, got %v", err)
|
||||
}
|
||||
if got["k"] != "b" {
|
||||
t.Fatalf("later listener should win k, got %v", got)
|
||||
}
|
||||
if got["only"] != "a" {
|
||||
t.Fatalf("partial payload missing only=a: %v", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUntilHandledStopsOnFirstHandled(t *testing.T) {
|
||||
bus := New()
|
||||
var order []string
|
||||
bus.Listen[*evt]("golem15.a", func(ctx context.Context, e *evt) error {
|
||||
order = append(order, "a")
|
||||
e.handled = true
|
||||
return nil
|
||||
})
|
||||
bus.Listen[*evt]("golem15.b", func(ctx context.Context, e *evt) error {
|
||||
order = append(order, "b")
|
||||
return nil
|
||||
})
|
||||
ok, err := bus.UntilHandled(context.Background(), &evt{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !ok {
|
||||
t.Fatal("expected handled")
|
||||
}
|
||||
if strings.Join(order, ",") != "a" {
|
||||
t.Fatalf("order = %v, want only a", order)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUntilHandledStopsOnFirstError(t *testing.T) {
|
||||
bus := New()
|
||||
var order []string
|
||||
bus.Listen[*evt]("golem15.a", func(ctx context.Context, e *evt) error {
|
||||
order = append(order, "a")
|
||||
return errors.New("stop")
|
||||
})
|
||||
bus.Listen[*evt]("golem15.b", func(ctx context.Context, e *evt) error {
|
||||
order = append(order, "b")
|
||||
e.handled = true
|
||||
return nil
|
||||
})
|
||||
ok, err := bus.UntilHandled(context.Background(), &evt{})
|
||||
if err == nil || !strings.Contains(err.Error(), "stop") {
|
||||
t.Fatalf("expected first error, got ok=%v err=%v", ok, err)
|
||||
}
|
||||
if ok {
|
||||
t.Fatal("error should not report handled")
|
||||
}
|
||||
if strings.Join(order, ",") != "a" {
|
||||
t.Fatalf("order = %v, want only a", order)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPriorityDescendingStableTies(t *testing.T) {
|
||||
bus := New()
|
||||
var order []string
|
||||
bus.ListenPriority[*evt]("golem15.low", 1, func(ctx context.Context, e *evt) error {
|
||||
order = append(order, "low")
|
||||
return nil
|
||||
})
|
||||
bus.Listen[*evt]("golem15.default-a", func(ctx context.Context, e *evt) error {
|
||||
order = append(order, "default-a")
|
||||
return nil
|
||||
})
|
||||
bus.Listen[*evt]("golem15.default-b", func(ctx context.Context, e *evt) error {
|
||||
order = append(order, "default-b")
|
||||
return nil
|
||||
})
|
||||
bus.ListenPriority[*evt]("golem15.high", 10, func(ctx context.Context, e *evt) error {
|
||||
order = append(order, "high")
|
||||
return nil
|
||||
})
|
||||
if err := bus.Fire(context.Background(), &evt{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if strings.Join(order, ",") != "high,low,default-a,default-b" {
|
||||
t.Fatalf("order = %v", order)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPanicRecoveredNamesOwnerPlugin(t *testing.T) {
|
||||
bus := New()
|
||||
bus.Listen[*evt]("golem15.boom", func(ctx context.Context, e *evt) error {
|
||||
panic("kapow")
|
||||
})
|
||||
err := bus.Fire(context.Background(), &evt{})
|
||||
if err == nil {
|
||||
t.Fatal("expected panic error")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "golem15.boom") {
|
||||
t.Fatalf("error should name owner plugin, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBusesAreIsolated(t *testing.T) {
|
||||
a := New()
|
||||
b := New()
|
||||
called := false
|
||||
a.Listen[*evt]("golem15.a", func(ctx context.Context, e *evt) error {
|
||||
called = true
|
||||
return nil
|
||||
})
|
||||
if err := b.Fire(context.Background(), &evt{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if called {
|
||||
t.Fatal("bus B should not invoke bus A's listeners")
|
||||
}
|
||||
}
|
||||
|
||||
func TestConcurrentIndependentAppBuses(t *testing.T) {
|
||||
a := New()
|
||||
b := New()
|
||||
var aCount, bCount atomic.Int64
|
||||
a.Listen[*evt]("golem15.a", func(ctx context.Context, e *evt) error {
|
||||
aCount.Add(1)
|
||||
return nil
|
||||
})
|
||||
b.Listen[*evt]("golem15.b", func(ctx context.Context, e *evt) error {
|
||||
bCount.Add(1)
|
||||
return nil
|
||||
})
|
||||
var wg sync.WaitGroup
|
||||
const n = 40
|
||||
for i := 0; i < n; i++ {
|
||||
wg.Add(2)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
if err := a.Fire(context.Background(), &evt{}); err != nil {
|
||||
t.Errorf("bus A: %v", err)
|
||||
}
|
||||
}()
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
if err := b.Fire(context.Background(), &evt{}); err != nil {
|
||||
t.Errorf("bus B: %v", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
if aCount.Load() != n || bCount.Load() != n {
|
||||
t.Fatalf("counts a=%d b=%d, want %d each", aCount.Load(), bCount.Load(), n)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPanicRecoveryOmitsEventPayload(t *testing.T) {
|
||||
bus := New()
|
||||
bus.Listen[*evt]("golem15.boom", func(ctx context.Context, e *evt) error {
|
||||
panic("secret-payload-xyz")
|
||||
})
|
||||
event := &evt{data: map[string]any{"secret": "classified"}}
|
||||
err := bus.Fire(context.Background(), event)
|
||||
if err == nil {
|
||||
t.Fatal("expected panic error")
|
||||
}
|
||||
msg := err.Error()
|
||||
if !strings.Contains(msg, "golem15.boom") {
|
||||
t.Fatalf("error should name owner plugin, got %v", err)
|
||||
}
|
||||
if strings.Contains(msg, "secret-payload-xyz") || strings.Contains(msg, "classified") {
|
||||
t.Fatalf("panic recovery leaked payload: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestFireWithNoListeners(t *testing.T) {
|
||||
bus := New()
|
||||
if err := bus.Fire(context.Background(), &evt{}); err != nil {
|
||||
t.Fatalf("Fire with no listeners: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCollectWithoutCollectableReturnsEmptyMap(t *testing.T) {
|
||||
type plain struct{ n int }
|
||||
bus := New()
|
||||
bus.Listen[plain]("golem15.a", func(ctx context.Context, e plain) error { return nil })
|
||||
got, err := bus.Collect(context.Background(), plain{n: 1})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(got) != 0 {
|
||||
t.Fatalf("payload = %v, want empty", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUntilHandledNeverHandled(t *testing.T) {
|
||||
bus := New()
|
||||
var order []string
|
||||
bus.Listen[*evt]("golem15.a", func(ctx context.Context, e *evt) error {
|
||||
order = append(order, "a")
|
||||
return nil
|
||||
})
|
||||
ok, err := bus.UntilHandled(context.Background(), &evt{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if ok {
|
||||
t.Fatal("expected unhandled")
|
||||
}
|
||||
if strings.Join(order, ",") != "a" {
|
||||
t.Fatalf("order = %v, want a", order)
|
||||
}
|
||||
}
|
||||
|
||||
func TestListenNilBusOrFnIsNoop(t *testing.T) {
|
||||
var bus *Bus
|
||||
bus.Listen[*evt]("golem15.a", func(ctx context.Context, e *evt) error { return errors.New("should not run") })
|
||||
if err := bus.Fire(context.Background(), &evt{}); err != nil {
|
||||
t.Fatalf("nil bus Fire: %v", err)
|
||||
}
|
||||
live := New()
|
||||
live.Listen[*evt]("golem15.a", nil)
|
||||
if err := live.Fire(context.Background(), &evt{}); err != nil {
|
||||
t.Fatalf("nil fn Fire: %v", err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user