mirror of
https://github.com/tiennm99/goclaw.git
synced 2026-10-03 07:12:50 +00:00
Bug fixes: - bus: add panic recovery in Broadcast() — panicking subscriber no longer crashes entire event bus goroutine - cron: fix deadlock in RunJob() — recordRun called while holding mutex, extracted recordRunLocked for callers already holding lock - cron: add run log recording to executeJobByID — automatic scheduler was not populating run log, only manual RunJob did New tests (~170 cases across 9 packages): - crypto: roundtrip, key derivation (hex/b64/raw), nonce uniqueness, backward compat, wrong key, corruption - permissions: role hierarchy, RoleFromScopes, CanAccess, scope precedence - bus: pub/sub delivery, buffer full, panic recovery, concurrent safety - providers/retry: IsRetryableError, backoff, jitter, Retry-After, context cancellation, hook callback - sessions: idempotency, concurrent writes, defensive copy, save/load roundtrip, metadata accumulation - scheduler: draining, drop policies, adaptive throttle, stale completion, lane concurrency, debounce, interrupt mode - config: JSON5 parsing, env overrides, FlexibleStringSlice, owner IDs - cron: schedule validation, computeNextRun, CRUD, job execution, failure tracking, persistence roundtrip - agent: history sanitization edge cases — all-tool history, partial results, cross-turn dedup, orphaned tools, 1000-msg performance
144 lines
3.7 KiB
Go
144 lines
3.7 KiB
Go
package bus
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log/slog"
|
|
"sync"
|
|
)
|
|
|
|
// MessageBus routes messages between channels and the agent runtime,
|
|
// and broadcasts events to WebSocket subscribers.
|
|
type MessageBus struct {
|
|
inbound chan InboundMessage
|
|
outbound chan OutboundMessage
|
|
|
|
// Channel message handlers (channel name → handler)
|
|
handlers map[string]MessageHandler
|
|
handlerMu sync.RWMutex
|
|
|
|
// Event subscribers (subscriber ID → handler)
|
|
subscribers map[string]EventHandler
|
|
subMu sync.RWMutex
|
|
}
|
|
|
|
func New() *MessageBus {
|
|
return &MessageBus{
|
|
inbound: make(chan InboundMessage, 1000),
|
|
outbound: make(chan OutboundMessage, 1000),
|
|
handlers: make(map[string]MessageHandler),
|
|
subscribers: make(map[string]EventHandler),
|
|
}
|
|
}
|
|
|
|
// PublishInbound queues an inbound message from a channel.
|
|
// Blocks if the inbound buffer is full.
|
|
func (mb *MessageBus) PublishInbound(msg InboundMessage) {
|
|
mb.inbound <- msg
|
|
}
|
|
|
|
// TryPublishInbound attempts to queue an inbound message without blocking.
|
|
// Returns false if the inbound buffer is full (message dropped).
|
|
func (mb *MessageBus) TryPublishInbound(msg InboundMessage) bool {
|
|
select {
|
|
case mb.inbound <- msg:
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
// ConsumeInbound blocks until an inbound message is available or ctx is cancelled.
|
|
func (mb *MessageBus) ConsumeInbound(ctx context.Context) (InboundMessage, bool) {
|
|
select {
|
|
case msg := <-mb.inbound:
|
|
return msg, true
|
|
case <-ctx.Done():
|
|
return InboundMessage{}, false
|
|
}
|
|
}
|
|
|
|
// PublishOutbound queues an outbound message to a channel.
|
|
// Blocks if the outbound buffer is full.
|
|
func (mb *MessageBus) PublishOutbound(msg OutboundMessage) {
|
|
mb.outbound <- msg
|
|
}
|
|
|
|
// TryPublishOutbound attempts to queue an outbound message without blocking.
|
|
// Returns false if the outbound buffer is full (message dropped).
|
|
func (mb *MessageBus) TryPublishOutbound(msg OutboundMessage) bool {
|
|
select {
|
|
case mb.outbound <- msg:
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
// SubscribeOutbound blocks until an outbound message is available or ctx is cancelled.
|
|
func (mb *MessageBus) SubscribeOutbound(ctx context.Context) (OutboundMessage, bool) {
|
|
select {
|
|
case msg := <-mb.outbound:
|
|
return msg, true
|
|
case <-ctx.Done():
|
|
return OutboundMessage{}, false
|
|
}
|
|
}
|
|
|
|
// RegisterHandler registers a message handler for a channel.
|
|
func (mb *MessageBus) RegisterHandler(channel string, handler MessageHandler) {
|
|
mb.handlerMu.Lock()
|
|
defer mb.handlerMu.Unlock()
|
|
mb.handlers[channel] = handler
|
|
}
|
|
|
|
// GetHandler returns the message handler for a channel.
|
|
func (mb *MessageBus) GetHandler(channel string) (MessageHandler, bool) {
|
|
mb.handlerMu.RLock()
|
|
defer mb.handlerMu.RUnlock()
|
|
handler, ok := mb.handlers[channel]
|
|
return handler, ok
|
|
}
|
|
|
|
// Subscribe registers an event subscriber. Returns the subscriber ID for unsubscribe.
|
|
func (mb *MessageBus) Subscribe(id string, handler EventHandler) {
|
|
mb.subMu.Lock()
|
|
defer mb.subMu.Unlock()
|
|
mb.subscribers[id] = handler
|
|
}
|
|
|
|
// Unsubscribe removes an event subscriber.
|
|
func (mb *MessageBus) Unsubscribe(id string) {
|
|
mb.subMu.Lock()
|
|
defer mb.subMu.Unlock()
|
|
delete(mb.subscribers, id)
|
|
}
|
|
|
|
// Broadcast sends an event to all subscribers (non-blocking per subscriber).
|
|
// Panicking handlers are caught and logged to prevent one bad subscriber
|
|
// from crashing the entire event bus.
|
|
func (mb *MessageBus) Broadcast(event Event) {
|
|
mb.subMu.RLock()
|
|
defer mb.subMu.RUnlock()
|
|
for id, handler := range mb.subscribers {
|
|
func() {
|
|
defer func() {
|
|
if r := recover(); r != nil {
|
|
slog.Error("bus: subscriber panicked",
|
|
"subscriber", id,
|
|
"event", event.Name,
|
|
"panic", fmt.Sprint(r),
|
|
)
|
|
}
|
|
}()
|
|
handler(event)
|
|
}()
|
|
}
|
|
}
|
|
|
|
// Close shuts down the message bus.
|
|
func (mb *MessageBus) Close() {
|
|
close(mb.inbound)
|
|
close(mb.outbound)
|
|
}
|