diff --git a/pkg/atproto/backfill_walk_test.go b/pkg/atproto/backfill_walk_test.go index 3ff7168b..f9daff60 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 a659a57f..a0e73d27 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,