Files
synapbus/internal/api/broadcaster.go
T
QiuSWandClaude Sonnet 5.5 def4bb55a8 fix(#1): keep out-of-order live agent events, add store tests
Replay dedup now uses a fixed replayedUpTo that live events never advance.
Add ListEventMetaAfter store tests, document current-membership replay,
and skip member/subject lookups when no agent stream is connected.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
2026-10-05 16:15:13 +08:00

188 lines
5.8 KiB
Go

package api
import (
"context"
"log/slog"
"github.com/synapbus/synapbus/internal/agents"
"github.com/synapbus/synapbus/internal/channels"
"github.com/synapbus/synapbus/internal/messaging"
)
// NewMessageEvent is broadcast when a new message is sent.
type NewMessageEvent struct {
Channel string `json:"channel,omitempty"` // set for channel messages
FromAgent string `json:"from_agent,omitempty"` // set for DMs
ToAgent string `json:"to_agent,omitempty"` // set for DMs
MessageID int64 `json:"message_id"`
}
// UnreadUpdateEvent is broadcast when unread counts change (e.g. mark-read).
type UnreadUpdateEvent struct {
Channel string `json:"channel,omitempty"`
Agent string `json:"agent,omitempty"`
UnreadCount int `json:"unread_count"`
}
// EventBroadcaster broadcasts real-time events to connected SSE clients.
type EventBroadcaster interface {
BroadcastNewMessage(ctx context.Context, ownerID int64, event NewMessageEvent)
BroadcastUnreadUpdate(ctx context.Context, ownerID int64, event UnreadUpdateEvent)
}
// SSEBroadcaster implements EventBroadcaster using the SSEHub.
type SSEBroadcaster struct {
hub *SSEHub
agentService *agents.AgentService
channelService *channels.Service
msgService *messaging.MessagingService // optional: resolves conversation subjects
logger *slog.Logger
}
// NewSSEBroadcaster creates a broadcaster that sends events via SSE.
func NewSSEBroadcaster(hub *SSEHub, agentService *agents.AgentService, channelService *channels.Service) *SSEBroadcaster {
return &SSEBroadcaster{
hub: hub,
agentService: agentService,
channelService: channelService,
logger: slog.Default().With("component", "api.broadcaster"),
}
}
// SetMessageService sets the messaging service used to resolve conversation
// subjects for agent events. Optional; without it events carry no subject.
func (b *SSEBroadcaster) SetMessageService(svc *messaging.MessagingService) {
b.msgService = svc
}
// BroadcastNewMessage sends a new_message event to the given owner.
func (b *SSEBroadcaster) BroadcastNewMessage(_ context.Context, ownerID int64, event NewMessageEvent) {
b.hub.Broadcast(ownerID, SSEEvent{
Type: "new_message",
Data: event,
})
}
// BroadcastUnreadUpdate sends an unread_update event to the given owner.
func (b *SSEBroadcaster) BroadcastUnreadUpdate(_ context.Context, ownerID int64, event UnreadUpdateEvent) {
b.hub.Broadcast(ownerID, SSEEvent{
Type: "unread_update",
Data: event,
})
}
// BroadcastDM broadcasts a new_message event for a direct message.
// It resolves the recipient agent's owner and sends the event to them.
func (b *SSEBroadcaster) BroadcastDM(ctx context.Context, msg NewMessageEvent) {
if msg.ToAgent == "" {
return
}
agent, err := b.agentService.GetAgent(ctx, msg.ToAgent)
if err != nil {
b.logger.Debug("could not resolve recipient owner for SSE broadcast",
"to_agent", msg.ToAgent, "error", err)
return
}
b.BroadcastNewMessage(ctx, agent.OwnerID, msg)
}
// BroadcastChannelMessage broadcasts a new_message event for a channel message.
// It resolves all channel members' owners and sends the event to each unique owner.
func (b *SSEBroadcaster) BroadcastChannelMessage(ctx context.Context, channelID int64, msg NewMessageEvent) {
if b.channelService == nil {
return
}
members, err := b.channelService.GetMembers(ctx, channelID)
if err != nil {
b.logger.Debug("could not get channel members for SSE broadcast",
"channel_id", channelID, "error", err)
return
}
// Collect unique owner IDs to avoid duplicate broadcasts.
seen := make(map[int64]bool)
for _, m := range members {
agent, err := b.agentService.GetAgent(ctx, m.AgentName)
if err != nil {
continue
}
if !seen[agent.OwnerID] {
seen[agent.OwnerID] = true
b.BroadcastNewMessage(ctx, agent.OwnerID, msg)
}
}
}
// OnMessageSent implements messaging.MessageListener so the SSEBroadcaster
// can be wired directly into the messaging service. This ensures SSE events
// fire for messages sent via MCP (agents) as well as the REST API.
func (b *SSEBroadcaster) OnMessageSent(ctx context.Context, msg *messaging.Message) {
event := NewMessageEvent{
MessageID: msg.ID,
FromAgent: msg.FromAgent,
ToAgent: msg.ToAgent,
}
if msg.ChannelID != nil {
ch, err := b.channelService.GetChannel(ctx, *msg.ChannelID)
if err == nil {
event.Channel = ch.Name
}
b.BroadcastChannelMessage(ctx, *msg.ChannelID, event)
} else {
b.BroadcastDM(ctx, event)
}
b.broadcastToAgents(ctx, msg)
}
// broadcastToAgents pushes a body-free new_message event to the per-agent
// stream (GET /api/agent-events): the DM recipient, or the members of the
// channel at send time. The sender is never notified of its own message.
func (b *SSEBroadcaster) broadcastToAgents(ctx context.Context, msg *messaging.Message) {
if !b.hub.hasAgentSubs() {
return // avoid member/subject DB lookups when nobody is listening
}
var recipients []string
if msg.ChannelID != nil {
if b.channelService == nil {
return
}
members, err := b.channelService.GetMembers(ctx, *msg.ChannelID)
if err != nil {
b.logger.Debug("could not get channel members for agent SSE broadcast",
"channel_id", *msg.ChannelID, "error", err)
return
}
for _, m := range members {
if m.AgentName != msg.FromAgent {
recipients = append(recipients, m.AgentName)
}
}
} else if msg.ToAgent != "" && msg.ToAgent != msg.FromAgent {
recipients = []string{msg.ToAgent}
}
if len(recipients) == 0 {
return
}
ev := AgentMessageEvent{
MessageID: msg.ID,
FromAgent: msg.FromAgent,
ToAgent: msg.ToAgent,
}
if msg.ChannelID != nil {
ev.FromAgent = ""
if ch, err := b.channelService.GetChannel(ctx, *msg.ChannelID); err == nil {
ev.Channel = ch.Name
}
}
if b.msgService != nil {
ev.Subject = b.msgService.GetConversationSubject(ctx, msg.ConversationID)
}
b.hub.BroadcastAgentMessage(recipients, ev)
}