Files
viettranx 90b7396d74 fix(bus,cron): panic recovery + deadlock fix + unit tests for 9 packages
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
2026-03-28 18:15:46 +07:00

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)
}