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>
This commit is contained in:
co-authored by
Claude Sonnet 5.5
parent
2e8a49ca5a
commit
def4bb55a8
@@ -54,6 +54,13 @@ func (h *SSEHub) SetAgentEventBacklog(b AgentEventBacklog) {
|
||||
h.agentBacklog = b
|
||||
}
|
||||
|
||||
// hasAgentSubs reports whether any /api/agent-events connection is open.
|
||||
func (h *SSEHub) hasAgentSubs() bool {
|
||||
h.mu.RLock()
|
||||
defer h.mu.RUnlock()
|
||||
return len(h.agentSubs) > 0
|
||||
}
|
||||
|
||||
// BroadcastAgentMessage delivers a new_message event to every connection of the
|
||||
// named agents. Connections that cannot keep up are dropped.
|
||||
func (h *SSEHub) BroadcastAgentMessage(agentNames []string, ev AgentMessageEvent) {
|
||||
@@ -180,9 +187,10 @@ func (h *SSEHub) HandleAgentEvents(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
// sentUpTo is the highest message id already delivered; live events at or
|
||||
// below it are duplicates of the replay.
|
||||
sentUpTo := lastID
|
||||
// replayedUpTo is fixed once the replay ends: live events at or below it
|
||||
// are duplicates of the replay. It is never advanced by live events, because
|
||||
// concurrent senders may broadcast ids out of order (11 before 10).
|
||||
replayedUpTo := lastID
|
||||
if hasLast {
|
||||
h.mu.RLock()
|
||||
backlog := h.agentBacklog
|
||||
@@ -211,7 +219,7 @@ func (h *SSEHub) HandleAgentEvents(w http.ResponseWriter, r *http.Request) {
|
||||
if !write(m.MessageID, "new_message", ev) {
|
||||
return
|
||||
}
|
||||
sentUpTo = m.MessageID
|
||||
replayedUpTo = m.MessageID
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -233,13 +241,12 @@ func (h *SSEHub) HandleAgentEvents(w http.ResponseWriter, r *http.Request) {
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
if ev.MessageID <= sentUpTo {
|
||||
if ev.MessageID <= replayedUpTo {
|
||||
continue
|
||||
}
|
||||
if !write(ev.MessageID, "new_message", ev) {
|
||||
return
|
||||
}
|
||||
sentUpTo = ev.MessageID
|
||||
case <-ticker.C:
|
||||
if !write(0, "heartbeat", map[string]any{
|
||||
"timestamp": time.Now().Format(time.RFC3339),
|
||||
|
||||
@@ -487,3 +487,56 @@ func TestAgentEvents_Heartbeat(t *testing.T) {
|
||||
t.Fatalf("event = %q, want heartbeat", f.Event)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentEvents_OutOfOrderLiveEventsNotDropped(t *testing.T) {
|
||||
env := newAgentEventsEnv(t, "alice")
|
||||
c := env.connect(env.keys["alice"], "")
|
||||
c.expectConnected("alice")
|
||||
|
||||
// Concurrent senders can broadcast ids out of order; both must arrive.
|
||||
env.hub.BroadcastAgentMessage([]string{"alice"}, AgentMessageEvent{MessageID: 11, FromAgent: "x", ToAgent: "alice"})
|
||||
env.hub.BroadcastAgentMessage([]string{"alice"}, AgentMessageEvent{MessageID: 10, FromAgent: "y", ToAgent: "alice"})
|
||||
|
||||
for _, want := range []string{"11", "10"} {
|
||||
if f, _ := c.nextMessage(); f.ID != want {
|
||||
t.Fatalf("got id %s, want %s", f.ID, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentEvents_OutOfOrderLiveAfterReplay(t *testing.T) {
|
||||
env := newAgentEventsEnv(t, "alice", "bob")
|
||||
first := env.dm("bob", "alice", "old", "")
|
||||
missed := env.dm("bob", "alice", "missed", "")
|
||||
|
||||
c := env.connect(env.keys["alice"], strconv.FormatInt(first, 10))
|
||||
c.expectConnected("alice")
|
||||
if f, _ := c.nextMessage(); f.ID != strconv.FormatInt(missed, 10) {
|
||||
t.Fatalf("replay id = %s, want %d", f.ID, missed)
|
||||
}
|
||||
|
||||
// A duplicate of the replayed message is skipped; later out-of-order ids pass.
|
||||
env.hub.BroadcastAgentMessage([]string{"alice"}, AgentMessageEvent{MessageID: missed})
|
||||
env.hub.BroadcastAgentMessage([]string{"alice"}, AgentMessageEvent{MessageID: missed + 20})
|
||||
env.hub.BroadcastAgentMessage([]string{"alice"}, AgentMessageEvent{MessageID: missed + 10})
|
||||
for _, want := range []int64{missed + 20, missed + 10} {
|
||||
if f, _ := c.nextMessage(); f.ID != strconv.FormatInt(want, 10) {
|
||||
t.Fatalf("got id %s, want %d", f.ID, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestSSEHub_HasAgentSubs(t *testing.T) {
|
||||
hub := NewSSEHub()
|
||||
if hub.hasAgentSubs() {
|
||||
t.Fatal("empty hub reports subscriptions")
|
||||
}
|
||||
sub := hub.addAgentSub("alice")
|
||||
if !hub.hasAgentSubs() {
|
||||
t.Fatal("hub with a subscription reports none")
|
||||
}
|
||||
hub.removeAgentSub("alice", sub)
|
||||
if hub.hasAgentSubs() {
|
||||
t.Fatal("hub still reports subscriptions after removal")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -143,6 +143,9 @@ func (b *SSEBroadcaster) OnMessageSent(ctx context.Context, msg *messaging.Messa
|
||||
// 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 {
|
||||
|
||||
@@ -22,6 +22,10 @@ type EventMeta struct {
|
||||
// the agent is a member of (same scope as SearchMessages). Messages sent by the
|
||||
// agent itself are excluded.
|
||||
//
|
||||
// Membership is evaluated at query time (current membership), so messages
|
||||
// posted in a channel before the agent joined it are also returned when their
|
||||
// id > afterID. Channel members can read that history anyway.
|
||||
//
|
||||
// Note: channel messages are stored with to_agent NULL (see InsertMessage), so
|
||||
// the channel branch matches both NULL and ''.
|
||||
func (s *SQLiteMessageStore) ListEventMetaAfter(ctx context.Context, agentName string, afterID int64, limit int) ([]*EventMeta, error) {
|
||||
|
||||
@@ -0,0 +1,97 @@
|
||||
package messaging
|
||||
|
||||
import (
|
||||
"context"
|
||||
"reflect"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestSQLiteMessageStore_ListEventMetaAfter(t *testing.T) {
|
||||
db := newTestDB(t)
|
||||
store := NewSQLiteMessageStore(db)
|
||||
ctx := context.Background()
|
||||
|
||||
for _, a := range []string{"a", "b", "c"} {
|
||||
seedAgent(t, db, a)
|
||||
}
|
||||
for _, q := range []string{
|
||||
`INSERT INTO channels (id, name, description, topic, type, is_private, is_system, created_by, created_at, updated_at)
|
||||
VALUES (1, 'room', '', '', 'standard', 0, 0, 'a', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`,
|
||||
`INSERT INTO channels (id, name, description, topic, type, is_private, is_system, created_by, created_at, updated_at)
|
||||
VALUES (2, 'hidden', '', '', 'standard', 0, 0, 'c', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)`,
|
||||
`INSERT INTO channel_members (channel_id, agent_name, role, joined_at) VALUES (1, 'a', 'member', CURRENT_TIMESTAMP)`,
|
||||
`INSERT INTO channel_members (channel_id, agent_name, role, joined_at) VALUES (1, 'b', 'member', CURRENT_TIMESTAMP)`,
|
||||
`INSERT INTO channel_members (channel_id, agent_name, role, joined_at) VALUES (2, 'c', 'member', CURRENT_TIMESTAMP)`,
|
||||
} {
|
||||
if _, err := db.ExecContext(ctx, q); err != nil {
|
||||
t.Fatalf("seed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
room, hidden := int64(1), int64(2)
|
||||
// Channel messages are stored with to_agent NULL (InsertMessage).
|
||||
seed := []struct {
|
||||
from, to string
|
||||
ch *int64
|
||||
}{
|
||||
{"b", "a", nil}, // 1: DM to a
|
||||
{"b", "c", nil}, // 2: DM to c
|
||||
{"b", "", &room}, // 3: room, from b
|
||||
{"c", "", &hidden}, // 4: hidden, from c
|
||||
{"a", "b", nil}, // 5: DM from a to b
|
||||
{"a", "", &room}, // 6: room, from a
|
||||
{"b", "a", nil}, // 7: DM to a
|
||||
}
|
||||
for _, s := range seed {
|
||||
conv := &Conversation{Subject: "subj-" + s.from, CreatedBy: s.from}
|
||||
if err := store.InsertConversation(ctx, conv); err != nil {
|
||||
t.Fatalf("InsertConversation: %v", err)
|
||||
}
|
||||
m := &Message{ConversationID: conv.ID, FromAgent: s.from, ToAgent: s.to, ChannelID: s.ch,
|
||||
Body: "secret", Priority: 5, Status: StatusPending}
|
||||
if err := store.InsertMessage(ctx, m); err != nil {
|
||||
t.Fatalf("InsertMessage: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
agent string
|
||||
afterID int64
|
||||
limit int
|
||||
want []int64
|
||||
}{
|
||||
{"a sees own DMs and joined channel, not others or own", "a", 0, 100, []int64{1, 3, 7}},
|
||||
{"c sees DM and own channel only", "c", 0, 100, []int64{2}}, // 4 is c's own message
|
||||
{"b sees DM and room message from a", "b", 0, 100, []int64{5, 6}},
|
||||
{"unknown agent sees nothing", "zed", 0, 100, nil},
|
||||
{"afterID is exclusive", "a", 1, 100, []int64{3, 7}},
|
||||
{"afterID past end", "a", 7, 100, nil},
|
||||
{"limit applies, ascending", "a", 0, 2, []int64{1, 3}},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
got, err := store.ListEventMetaAfter(ctx, tt.agent, tt.afterID, tt.limit)
|
||||
if err != nil {
|
||||
t.Fatalf("ListEventMetaAfter: %v", err)
|
||||
}
|
||||
var ids []int64
|
||||
for _, e := range got {
|
||||
ids = append(ids, e.MessageID)
|
||||
}
|
||||
if !reflect.DeepEqual(ids, tt.want) {
|
||||
t.Fatalf("ids = %v, want %v", ids, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
t.Run("metadata fields", func(t *testing.T) {
|
||||
got, _ := store.ListEventMetaAfter(ctx, "a", 0, 100)
|
||||
if got[0].FromAgent != "b" || got[0].ToAgent != "a" || got[0].Channel != "" || got[0].Subject != "subj-b" {
|
||||
t.Errorf("DM meta = %+v", got[0])
|
||||
}
|
||||
if got[1].Channel != "room" || got[1].ToAgent != "" {
|
||||
t.Errorf("channel meta = %+v", got[1])
|
||||
}
|
||||
})
|
||||
}
|
||||
Reference in New Issue
Block a user