From def4bb55a87ce31f6f9f8245d0a0e13641d7e4b3 Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Mon, 5 Oct 2026 16:15:13 +0800 Subject: [PATCH] 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 --- internal/api/agent_events.go | 19 ++++-- internal/api/agent_events_test.go | 53 +++++++++++++++ internal/api/broadcaster.go | 3 + internal/messaging/event_meta.go | 4 ++ internal/messaging/event_meta_test.go | 97 +++++++++++++++++++++++++++ 5 files changed, 170 insertions(+), 6 deletions(-) create mode 100644 internal/messaging/event_meta_test.go diff --git a/internal/api/agent_events.go b/internal/api/agent_events.go index 159ac62..e7e0874 100644 --- a/internal/api/agent_events.go +++ b/internal/api/agent_events.go @@ -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), diff --git a/internal/api/agent_events_test.go b/internal/api/agent_events_test.go index ca71685..1ac67b8 100644 --- a/internal/api/agent_events_test.go +++ b/internal/api/agent_events_test.go @@ -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") + } +} diff --git a/internal/api/broadcaster.go b/internal/api/broadcaster.go index b897351..3bff5c1 100644 --- a/internal/api/broadcaster.go +++ b/internal/api/broadcaster.go @@ -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 { diff --git a/internal/messaging/event_meta.go b/internal/messaging/event_meta.go index 98e1cf4..60eea28 100644 --- a/internal/messaging/event_meta.go +++ b/internal/messaging/event_meta.go @@ -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) { diff --git a/internal/messaging/event_meta_test.go b/internal/messaging/event_meta_test.go new file mode 100644 index 0000000..2b7c0ec --- /dev/null +++ b/internal/messaging/event_meta_test.go @@ -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]) + } + }) +}