mengstack-api/internal/kernel/eventbus/bus.go
MengStack Dev df809f1045
Some checks failed
CI / Build & Test (push) Failing after 1m31s
feat: user CRUD API + dashboard stats + kernel infrastructure
- Add user management endpoints (list/create/get/update/delete) with pagination and search
- Add dashboard stats endpoint with tenant/user/online counts and growth metrics
- Add tenant resolver middleware for multi-tenant request scoping
- Add i18n kernel with zh/en message files and AcceptLanguage middleware
- Add WebSocket hub/handler for real-time communication
- Add job scheduler kernel with cron support
- Add plugin sandbox for isolated execution
- Add storage kernel (local filesystem)
- Add event bus kernel for pub/sub
- Add cache kernel abstraction
- Add database migration runner and version upgrade checker
- Add rate limiting middleware with Redis backend
- Add SQL migrations for rbac, audit_logs, settings, notifications, examples
- Extend user repository with list/delete/count operations
- Register all module routes with tenant resolver
2026-10-03 03:42:58 +08:00

74 lines
1.2 KiB
Go

package eventbus
import (
"context"
"sync"
)
type Event interface {
Name() string
}
type Handler func(ctx context.Context, event Event) error
type Bus struct {
mu sync.RWMutex
handlers map[string][]Handler
async bool
queue chan asyncEvent
}
type asyncEvent struct {
ctx context.Context
event Event
}
func New() *Bus {
return &Bus{
handlers: make(map[string][]Handler),
queue: make(chan asyncEvent, 1024),
}
}
func (b *Bus) Subscribe(eventName string, handler Handler) {
b.mu.Lock()
defer b.mu.Unlock()
b.handlers[eventName] = append(b.handlers[eventName], handler)
}
func (b *Bus) Publish(ctx context.Context, event Event) error {
b.mu.RLock()
handlers := b.handlers[event.Name()]
b.mu.RUnlock()
for _, h := range handlers {
if err := h(ctx, event); err != nil {
return err
}
}
return nil
}
func (b *Bus) PublishAsync(ctx context.Context, event Event) {
b.queue <- asyncEvent{ctx: ctx, event: event}
}
func (b *Bus) StartWorkers(n int) {
for i := 0; i < n; i++ {
go func() {
for ae := range b.queue {
b.mu.RLock()
handlers := b.handlers[ae.event.Name()]
b.mu.RUnlock()
for _, h := range handlers {
_ = h(ae.ctx, ae.event)
}
}
}()
}
}
func (b *Bus) Stop() {
close(b.queue)
}