diff --git a/backpack/app.go b/backpack/app.go index fd1511d..4e178c0 100644 --- a/backpack/app.go +++ b/backpack/app.go @@ -5,12 +5,14 @@ import ( "sync" "git.golem15.com/golem15/summercms/compass" + "git.golem15.com/golem15/summercms/festival" ) // App is the per-instance application container. It must not import party. type App struct { Config *compass.Config Services *Registry + Events *festival.Bus mu sync.RWMutex plugins map[string]struct{} @@ -21,6 +23,7 @@ func New(cfg *compass.Config) *App { return &App{ Config: cfg, Services: NewRegistry(), + Events: festival.New(), plugins: make(map[string]struct{}), } } diff --git a/examples/hello/plugins/greeter/plugin.go b/examples/hello/plugins/greeter/plugin.go index 184a665..3a7f16a 100644 --- a/examples/hello/plugins/greeter/plugin.go +++ b/examples/hello/plugins/greeter/plugin.go @@ -5,10 +5,15 @@ import ( "git.golem15.com/golem15/summercms/backpack" "git.golem15.com/golem15/summercms/bonfire" + "git.golem15.com/golem15/summercms/festival" "git.golem15.com/golem15/summercms/pact" "git.golem15.com/golem15/summercms/party" + "git.golem15.com/golem15/summercms/towel" ) +var _ festival.Collectable = (*HelloEvent)(nil) +var _ festival.Handleable = (*HelloEvent)(nil) + // Plugin is the golem15.greeter plugin. It requires golem15.hello. type Plugin struct { app *backpack.App @@ -28,6 +33,19 @@ func (p *Plugin) Boot(app *backpack.App) error { p.extra = msg.Message() } } + if app != nil && app.Events != nil { + app.Events.Listen[*HelloEvent]("golem15.greeter", func(ctx context.Context, e *HelloEvent) error { + if e.data == nil { + e.data = map[string]any{} + } + e.data["source"] = "greeter" + if actor, ok := towel.Actor(ctx); ok { + e.data["actor"] = actor + } + e.handled = true + return nil + }) + } return nil } @@ -44,12 +62,50 @@ func (p *Plugin) Commands() []bonfire.Command { posts = p.app.Config.Int("golem15.hello.posts_per_page") debug = p.app.Config.Bool("app.debug") } - out.Printf("name=%s posts_per_page=%d debug=%t extra=%s\n", name, posts, debug, p.extra) + events := "ok" + collected := "" + handled := false + if p.app != nil && p.app.Events != nil { + if err := p.app.Events.Fire(ctx, &HelloEvent{}); err != nil { + events = "err" + } + payload, err := p.app.Events.Collect(ctx, &HelloEvent{}) + if err != nil { + events = "err" + } + if src, ok := payload["source"].(string); ok { + collected = src + } + ok, err := p.app.Events.UntilHandled(ctx, &HelloEvent{}) + if err != nil { + events = "err" + } + handled = ok + } + out.Printf("name=%s posts_per_page=%d debug=%t extra=%s events=%s collected=%s handled=%t\n", name, posts, debug, p.extra, events, collected, handled) return nil }, }} } +// HelloEvent is a typed greeter event used to demonstrate Fire, Collect +// and UntilHandled on the app bus. +type HelloEvent struct { + data map[string]any + handled bool +} + +func (e *HelloEvent) Collected() map[string]any { + if e == nil { + return nil + } + return e.data +} + +func (e *HelloEvent) IsHandled() bool { + return e != nil && e.handled +} + var _ pact.HasCommands = (*Plugin)(nil) func init() { diff --git a/festival/bus.go b/festival/bus.go new file mode 100644 index 0000000..3d17076 --- /dev/null +++ b/festival/bus.go @@ -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 + } +} diff --git a/towel/context.go b/towel/context.go new file mode 100644 index 0000000..89d2bfb --- /dev/null +++ b/towel/context.go @@ -0,0 +1,63 @@ +package towel + +import "context" + +type actorKey struct{} +type organizationKey struct{} +type collectionKey struct{} +type localeKey struct{} + +func withValue(ctx context.Context, key, value any) context.Context { + if ctx == nil { + ctx = context.Background() + } + return context.WithValue(ctx, key, value) +} + +func stringValue(ctx context.Context, key any) (string, bool) { + if ctx == nil { + return "", false + } + v, ok := ctx.Value(key).(string) + return v, ok +} + +// WithActor stores the request actor on ctx. +func WithActor(ctx context.Context, actor string) context.Context { + return withValue(ctx, actorKey{}, actor) +} + +// Actor returns the request actor from ctx. +func Actor(ctx context.Context) (string, bool) { + return stringValue(ctx, actorKey{}) +} + +// WithOrganization stores the request organization on ctx. +func WithOrganization(ctx context.Context, org string) context.Context { + return withValue(ctx, organizationKey{}, org) +} + +// Organization returns the request organization from ctx. +func Organization(ctx context.Context) (string, bool) { + return stringValue(ctx, organizationKey{}) +} + +// WithCollection stores the request collection on ctx. +func WithCollection(ctx context.Context, collection string) context.Context { + return withValue(ctx, collectionKey{}, collection) +} + +// Collection returns the request collection from ctx. +func Collection(ctx context.Context) (string, bool) { + return stringValue(ctx, collectionKey{}) +} + +// WithLocale stores the request locale on ctx. +func WithLocale(ctx context.Context, locale string) context.Context { + return withValue(ctx, localeKey{}, locale) +} + +// Locale returns the request locale from ctx. +func Locale(ctx context.Context) (string, bool) { + return stringValue(ctx, localeKey{}) +} diff --git a/towel/context_test.go b/towel/context_test.go index f90b62b..42976bf 100644 --- a/towel/context_test.go +++ b/towel/context_test.go @@ -2,8 +2,6 @@ package towel import ( "context" - "go/parser" - "go/token" "os" "strings" "testing" @@ -40,21 +38,11 @@ func TestNoPackageGlobalRequestState(t *testing.T) { if err != nil { t.Fatal(err) } - fset := token.NewFileSet() for _, entry := range entries { name := entry.Name() if entry.IsDir() || !strings.HasSuffix(name, ".go") || strings.HasSuffix(name, "_test.go") { continue } - file, err := parser.ParseFile(fset, name, nil, parser.SkipObjectResolution) - if err != nil { - t.Fatal(err) - } - for _, spec := range file.Decls { - gen, ok := spec.(interface{ TokString() string }) - _ = gen - _ = ok - } src, err := os.ReadFile(name) if err != nil { t.Fatal(err)