feat(01-02): dispatch typed events with request context
- App-owned festival bus with Fire, Collect and UntilHandled - Recover listener panics with owner plugin IDs; isolate buses per app - towel context accessors and hello command demo of all three modes
This commit is contained in:
@@ -5,12 +5,14 @@ import (
|
|||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
"git.golem15.com/golem15/summercms/compass"
|
"git.golem15.com/golem15/summercms/compass"
|
||||||
|
"git.golem15.com/golem15/summercms/festival"
|
||||||
)
|
)
|
||||||
|
|
||||||
// App is the per-instance application container. It must not import party.
|
// App is the per-instance application container. It must not import party.
|
||||||
type App struct {
|
type App struct {
|
||||||
Config *compass.Config
|
Config *compass.Config
|
||||||
Services *Registry
|
Services *Registry
|
||||||
|
Events *festival.Bus
|
||||||
|
|
||||||
mu sync.RWMutex
|
mu sync.RWMutex
|
||||||
plugins map[string]struct{}
|
plugins map[string]struct{}
|
||||||
@@ -21,6 +23,7 @@ func New(cfg *compass.Config) *App {
|
|||||||
return &App{
|
return &App{
|
||||||
Config: cfg,
|
Config: cfg,
|
||||||
Services: NewRegistry(),
|
Services: NewRegistry(),
|
||||||
|
Events: festival.New(),
|
||||||
plugins: make(map[string]struct{}),
|
plugins: make(map[string]struct{}),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -5,10 +5,15 @@ import (
|
|||||||
|
|
||||||
"git.golem15.com/golem15/summercms/backpack"
|
"git.golem15.com/golem15/summercms/backpack"
|
||||||
"git.golem15.com/golem15/summercms/bonfire"
|
"git.golem15.com/golem15/summercms/bonfire"
|
||||||
|
"git.golem15.com/golem15/summercms/festival"
|
||||||
"git.golem15.com/golem15/summercms/pact"
|
"git.golem15.com/golem15/summercms/pact"
|
||||||
"git.golem15.com/golem15/summercms/party"
|
"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.
|
// Plugin is the golem15.greeter plugin. It requires golem15.hello.
|
||||||
type Plugin struct {
|
type Plugin struct {
|
||||||
app *backpack.App
|
app *backpack.App
|
||||||
@@ -28,6 +33,19 @@ func (p *Plugin) Boot(app *backpack.App) error {
|
|||||||
p.extra = msg.Message()
|
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
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -44,12 +62,50 @@ func (p *Plugin) Commands() []bonfire.Command {
|
|||||||
posts = p.app.Config.Int("golem15.hello.posts_per_page")
|
posts = p.app.Config.Int("golem15.hello.posts_per_page")
|
||||||
debug = p.app.Config.Bool("app.debug")
|
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
|
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)
|
var _ pact.HasCommands = (*Plugin)(nil)
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
|
|||||||
149
festival/bus.go
Normal file
149
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
|
||||||
|
}
|
||||||
|
}
|
||||||
63
towel/context.go
Normal file
63
towel/context.go
Normal file
@@ -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{})
|
||||||
|
}
|
||||||
@@ -2,8 +2,6 @@ package towel
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"go/parser"
|
|
||||||
"go/token"
|
|
||||||
"os"
|
"os"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
@@ -40,21 +38,11 @@ func TestNoPackageGlobalRequestState(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
}
|
}
|
||||||
fset := token.NewFileSet()
|
|
||||||
for _, entry := range entries {
|
for _, entry := range entries {
|
||||||
name := entry.Name()
|
name := entry.Name()
|
||||||
if entry.IsDir() || !strings.HasSuffix(name, ".go") || strings.HasSuffix(name, "_test.go") {
|
if entry.IsDir() || !strings.HasSuffix(name, ".go") || strings.HasSuffix(name, "_test.go") {
|
||||||
continue
|
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)
|
src, err := os.ReadFile(name)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatal(err)
|
t.Fatal(err)
|
||||||
|
|||||||
Reference in New Issue
Block a user