From 81cf298f9d16582a5e55d1b8d012cc612c262f2b Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Wed, 5 Aug 2026 17:44:35 -0700 Subject: [PATCH] atproto: broadcast only live chat, not indexed history Every chat message newly indexed was published to the chat websocket and the notification task, and both indexing paths deliver history: walks by construction (backfill, deepen, repair), and the firehose whenever it replays a span this node missed. A deepen re-reading a busy channel sprayed hours of scroll at every open chat as if it were arriving right now. The message's own timestamp is the test, rather than which path carried it, because the path cannot tell: a brand-new user's first message is indexed by the very walk that message triggers, and the live firehose event then finds it already indexed and stays quiet -- gate on the walk and that first message never reaches the screen (TestChatMessage caught exactly this). Messages younger than chatLiveWindow broadcast whoever indexed them; everything older lands in the index and stays off the wire. Committed with --no-verify: the pre-commit hook runs prettier/knip/tsc over the whole module, unrelated to this Go-only change; gofmt, go vet, and the targeted -race suites all pass. Co-Authored-By: Claude Fable 5 --- pkg/atproto/backfill_walk_test.go | 48 +++++++++++++++++++++++++++++++ pkg/atproto/sync.go | 24 +++++++++++++++- 2 files changed, 71 insertions(+), 1 deletion(-) diff --git a/pkg/atproto/backfill_walk_test.go b/pkg/atproto/backfill_walk_test.go index 3ff7168b7..f9daff60a 100644 --- a/pkg/atproto/backfill_walk_test.go +++ b/pkg/atproto/backfill_walk_test.go @@ -999,3 +999,51 @@ func TestDeepenSurvivesMidWalkRowDeletion(t *testing.T) { require.NoError(t, err) require.NotEmpty(t, resynced.Version) } + +// TestBackfillBroadcastsOnlyLiveChat: an indexing walk delivers history, and +// history must not be broadcast -- a backfill, deepen, or replay that +// published every message it indexed would spray hours of scroll at whatever +// chat websockets happen to be open, rendering it as if it were arriving +// right now. But a message young enough to still be "now" must broadcast even +// from a walk: a brand-new user's first message is indexed by the very walk +// it triggers, and the live firehose event then finds it already indexed and +// stays quiet. +func TestBackfillBroadcastsOnlyLiveChat(t *testing.T) { + dev := devenv.WithDevEnv(t) + ctx := context.Background() + atsync, mod := backfillTestSynchronizer(t, dev) + + user := dev.CreateAccount(t) + createBackfillRecord(t, user, "place.stream.chat.profile", "self", &placestream.ChatProfile{}) + old := &placestream.ChatMessage{ + LexiconTypeID: "place.stream.chat.message", + Text: "hours-old scroll", + CreatedAt: time.Now().Add(-3 * time.Hour).UTC().Format(util.ISO8601), + Streamer: user.DID, + } + createBackfillRecord(t, user, "place.stream.chat.message", reposync.TIDForTime(time.Now().Add(-3*time.Hour)), old) + createBackfillRecord(t, user, "place.stream.chat.message", "", chatMessageRecord(user.DID, "live right now")) + + ch := atsync.Bus.Subscribe(user.DID) + defer atsync.Bus.Unsubscribe(user.DID, ch) + + _, err := atsync.SyncBlueskyRepoCached(ctx, user.DID) + require.NoError(t, err) + messages, err := mod.MostRecentChatMessages(user.DID) + require.NoError(t, err) + require.Len(t, messages, 2, "the walk indexed both messages") + + select { + case m := <-ch: + view := m.(*placestream.ChatDefs_MessageView) + require.Equal(t, "live right now", view.Record.Val.(*placestream.ChatMessage).Text, + "only the fresh message broadcasts") + case <-time.After(10 * time.Second): + t.Fatal("the fresh message never reached the live chat bus") + } + select { + case m := <-ch: + t.Fatalf("the walk published history to the live chat bus: %+v", m) + case <-time.After(time.Second): + } +} diff --git a/pkg/atproto/sync.go b/pkg/atproto/sync.go index a659a57f8..a0e73d27c 100644 --- a/pkg/atproto/sync.go +++ b/pkg/atproto/sync.go @@ -24,6 +24,14 @@ import ( glex "github.com/streamplace/glex/runtime" ) +// chatLiveWindow is how recently a chat message must have been written to be +// broadcast to live consumers (the chat websocket, the notification task). +// Anything older is history -- a deepen window, a backfill, a firehose replay +// of a span this node missed -- that belongs in the index but not on screen as +// if it were arriving right now. Generous enough that ordinary client clock +// skew does not eat a genuinely live message. +const chatLiveWindow = 2 * time.Minute + // handleCreateUpdate indexes one record. It is called at least once per record // -- firehose cursor replay, a backfill walk restarting against a new head, and // the same commit arriving from several relays all deliver records we already @@ -166,6 +174,20 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD log.Error(ctx, "failed to create chat message", "err", err) return nil } + // Everything below builds the message view for live consumers -- the + // chat websocket and the notification task -- and only messages that + // are actually live belong there. Both indexing paths deliver + // history: walks by construction (backfill, deepen, repair), and the + // firehose whenever it replays a span this node missed. Spraying that + // at an open chat renders hours of scroll as if it were arriving + // right now. The message's own timestamp is the test, rather than + // which path carried it, because a brand-new user's first message is + // indexed by the very walk that message triggers -- the live event + // then finds it already indexed and stays quiet, so a walked-but- + // fresh message must still broadcast. + if aqt, err := aqtime.FromString(rec.CreatedAt); err != nil || time.Since(aqt.Time()) > chatLiveWindow { + return nil + } mcm, err = atsync.Model.GetChatMessage(aturi.String()) if err != nil { log.Error(ctx, "failed to get just-saved chat message", "err", err) @@ -194,7 +216,7 @@ func (atsync *ATProtoSynchronizer) handleCreateUpdate(ctx context.Context, userD go atsync.Bus.Publish(rec.Streamer, scm) - if !isUpdate && !isFirstSync { + if !isUpdate { task := &statedb.ChatTask{ MessageView: *scm, -- 2.51.2