diff --git a/js/components/src/components/chat/chat.tsx b/js/components/src/components/chat/chat.tsx index 8e2f1389f..3f3dd11e8 100644 --- a/js/components/src/components/chat/chat.tsx +++ b/js/components/src/components/chat/chat.tsx @@ -1,3 +1,4 @@ +import { chatMessageOpacity, isChatMessageGone } from "@streamplace/core"; import { ChevronDown, Ellipsis, Reply } from "lucide-react-native"; import { ComponentProps, @@ -73,6 +74,20 @@ function LeftAction(prog: SharedValue, drag: SharedValue) { const SHOWN_MSGS = Platform.OS === "ios" || Platform.OS === "android" ? 25 : 100; +// Chat ages out on the hour, so a slow tick is plenty: it only exists so the +// fade and the disappearance happen while the viewer watches, rather than +// waiting for the next message to arrive and re-render the list. +const CHAT_EXPIRY_TICK_MS = 30_000; + +function useChatExpiryTick(): number { + const [now, setNow] = useState(() => Date.now()); + useEffect(() => { + const timer = setInterval(() => setNow(Date.now()), CHAT_EXPIRY_TICK_MS); + return () => clearInterval(timer); + }, []); + return now; +} + const keyExtractor = (item: ChatMessageViewHydrated, index: number) => { return `${item.uri}`; }; @@ -165,7 +180,14 @@ const ActionsBar = memo( }, ); -const ChatLine = memo(({ item }: { item: ChatMessageViewHydrated }) => { +const ChatLine = memo(function ChatLine({ + item, + opacity, +}: { + item: ChatMessageViewHydrated; + /** 0..1 from chatMessageOpacity; 1 when the viewer is reading history. */ + opacity: number; +}) { const { theme } = useTheme(); const setReply = useSetReplyToMessage(); const setModMsg = usePlayerStore((state) => state.setModMessage); @@ -198,12 +220,14 @@ const ChatLine = memo(({ item }: { item: ChatMessageViewHydrated }) => { if (item.author.did === "did:sys:system") { return ( - + + + ); } @@ -218,6 +242,7 @@ const ChatLine = memo(({ item }: { item: ChatMessageViewHydrated }) => { borderRadius: borderRadius.md, minWidth: 0, maxWidth: "100%", + opacity, }, isHovered ? { backgroundColor: theme.colors.surfaceHover } : {}, ]} @@ -238,7 +263,7 @@ const ChatLine = memo(({ item }: { item: ChatMessageViewHydrated }) => { } return ( - <> + { > - + ); }); @@ -311,16 +336,24 @@ export function Chat({ }, []); const [isVisible, setIsVisible] = useState(true); const flatListRef = useRef(null); + const now = useChatExpiryTick(); // The store keeps chat oldest-first. An inverted FlatList renders index 0 at // the bottom, so feed it newest-first to keep the latest message at the // bottom (or at the top when reverse is set, where inverted is off). + // + // Old messages fade out of the live view and then leave it, but they are + // never dropped from the store: a viewer who scrolls back up is reading + // history, so there everything is shown at full strength. const displayMessages = useMemo(() => { if (!chat) return []; const visible = hideSystemMessages ? chat.filter((m) => m.author.did !== "did:sys:system") : chat; - return visible.slice(-shownMessages).reverse(); - }, [chat, shownMessages, hideSystemMessages]); + const live = isScrolledUp + ? visible + : visible.filter((m) => !isChatMessageGone(m.record.createdAt, now)); + return live.slice(-shownMessages).reverse(); + }, [chat, shownMessages, hideSystemMessages, isScrolledUp, now]); const latestMessageTime = displayMessages[0] ? new Date(displayMessages[0].record.createdAt).getTime() : null; @@ -424,7 +457,14 @@ export function Chat({ keyExtractor={keyExtractor} renderItem={({ item, index }) => ( - + )} removeClippedSubviews={true} diff --git a/js/core/src/livestream-store/chat-expiry.test.ts b/js/core/src/livestream-store/chat-expiry.test.ts new file mode 100644 index 000000000..6db163e1d --- /dev/null +++ b/js/core/src/livestream-store/chat-expiry.test.ts @@ -0,0 +1,76 @@ +import { describe, expect, it } from "vitest"; +import { + CHAT_MESSAGE_FADE_DURATION_MS, + CHAT_MESSAGE_FADE_START_MS, + CHAT_MESSAGE_GONE_MS, + chatMessageOpacity, + isChatMessageGone, +} from "./chat-expiry"; + +// Fixed clock so age arithmetic is exact rather than wall-clock dependent. +const NOW = Date.UTC(2024, 0, 1, 12, 0, 0); +const createdAtMsAgo = (ms: number) => new Date(NOW - ms).toISOString(); + +describe("chatMessageOpacity", () => { + it("leaves a fresh message at full strength", () => { + expect(chatMessageOpacity(createdAtMsAgo(0), NOW)).toBe(1); + expect( + chatMessageOpacity(createdAtMsAgo(CHAT_MESSAGE_FADE_START_MS - 1), NOW), + ).toBe(1); + }); + + it("starts fading exactly at the fade start", () => { + expect( + chatMessageOpacity(createdAtMsAgo(CHAT_MESSAGE_FADE_START_MS), NOW), + ).toBe(1); + }); + + it("is half faded at the midpoint of the fade window", () => { + const half = CHAT_MESSAGE_FADE_START_MS + CHAT_MESSAGE_FADE_DURATION_MS / 2; + expect(chatMessageOpacity(createdAtMsAgo(half), NOW)).toBeCloseTo(0.5); + }); + + it("is gone at the end of the fade window and stays gone", () => { + expect(chatMessageOpacity(createdAtMsAgo(CHAT_MESSAGE_GONE_MS), NOW)).toBe( + 0, + ); + expect( + chatMessageOpacity( + createdAtMsAgo(CHAT_MESSAGE_GONE_MS + 86_400_000), + NOW, + ), + ).toBe(0); + }); + + it("does not fade a message timestamped in the future", () => { + expect(chatMessageOpacity(new Date(NOW + 60_000).toISOString(), NOW)).toBe( + 1, + ); + }); + + it("treats an unreadable timestamp as fresh rather than hiding it", () => { + expect(chatMessageOpacity("not a date", NOW)).toBe(1); + }); +}); + +describe("isChatMessageGone", () => { + it("keeps messages until the fade finishes", () => { + expect(isChatMessageGone(createdAtMsAgo(0), NOW)).toBe(false); + expect( + isChatMessageGone(createdAtMsAgo(CHAT_MESSAGE_GONE_MS - 1), NOW), + ).toBe(false); + }); + + it("drops a message once the fade finishes", () => { + expect(isChatMessageGone(createdAtMsAgo(CHAT_MESSAGE_GONE_MS), NOW)).toBe( + true, + ); + expect( + isChatMessageGone(createdAtMsAgo(CHAT_MESSAGE_GONE_MS + 1), NOW), + ).toBe(true); + }); + + it("keeps a message with an unreadable timestamp", () => { + expect(isChatMessageGone("not a date", NOW)).toBe(false); + }); +}); diff --git a/js/core/src/livestream-store/chat-expiry.ts b/js/core/src/livestream-store/chat-expiry.ts new file mode 100644 index 000000000..e3272f769 --- /dev/null +++ b/js/core/src/livestream-store/chat-expiry.ts @@ -0,0 +1,61 @@ +// Age-based expiry for live chat. +// +// Chat is a conversation, not an archive. A message holds full strength for +// the first hour, fades over the next hour, and is gone from the live view +// once it has faded out. Nothing is discarded client-side: a viewer who +// scrolls back up to read history sees expired messages again, and it is the +// node's own --chat-message-retention window that eventually stops handing +// them out at all. +// +// These helpers are pure so the React Native chat and the web chat share one +// definition of "old", and so the arithmetic can be tested without a renderer. + +/** Age at which a message starts fading out of the live view. */ +export const CHAT_MESSAGE_FADE_START_MS = 60 * 60 * 1000; // 1 hour + +/** How long the fade lasts; a message is gone at START + DURATION. */ +export const CHAT_MESSAGE_FADE_DURATION_MS = 60 * 60 * 1000; // 1 hour + +/** Age at which a message has fully faded and leaves the live view. */ +export const CHAT_MESSAGE_GONE_MS = + CHAT_MESSAGE_FADE_START_MS + CHAT_MESSAGE_FADE_DURATION_MS; + +export type ChatMessageTimestamp = string | number | Date; + +// Epoch milliseconds, or NaN for a timestamp we cannot read. Callers treat a +// bad timestamp as "no age", so a malformed record never hides real chat. +const messageTime = (createdAt: ChatMessageTimestamp): number => + createdAt instanceof Date + ? createdAt.getTime() + : typeof createdAt === "number" + ? createdAt + : new Date(createdAt).getTime(); + +/** + * Opacity for a message in the live view: 1 while fresh, ramping linearly to + * 0 across the fade window, and 0 once it is gone. + */ +export function chatMessageOpacity( + createdAt: ChatMessageTimestamp, + now: number = Date.now(), +): number { + const created = messageTime(createdAt); + if (!Number.isFinite(created)) return 1; + const age = now - created; + if (age <= CHAT_MESSAGE_FADE_START_MS) return 1; + if (age >= CHAT_MESSAGE_GONE_MS) return 0; + return 1 - (age - CHAT_MESSAGE_FADE_START_MS) / CHAT_MESSAGE_FADE_DURATION_MS; +} + +/** + * True once a message has faded out entirely, meaning the live view drops it. + * Scrolled back into history the message is still rendered, at full opacity. + */ +export function isChatMessageGone( + createdAt: ChatMessageTimestamp, + now: number = Date.now(), +): boolean { + const created = messageTime(createdAt); + if (!Number.isFinite(created)) return false; + return now - created >= CHAT_MESSAGE_GONE_MS; +} diff --git a/js/core/src/livestream-store/index.tsx b/js/core/src/livestream-store/index.tsx index be6c0f8ca..c164cfeb8 100644 --- a/js/core/src/livestream-store/index.tsx +++ b/js/core/src/livestream-store/index.tsx @@ -2,6 +2,7 @@ // Contains only platform-agnostic, React-free code: state, factory, // reducers, and pure utilities. The React hooks and context live in // @streamplace/components. +export * from "./chat-expiry"; export * from "./chat-reducer"; export * from "./connect"; export * from "./problems"; diff --git a/js/web/src/components/stream/chat-panel.tsx b/js/web/src/components/stream/chat-panel.tsx index aaab622f5..cf173b27e 100644 --- a/js/web/src/components/stream/chat-panel.tsx +++ b/js/web/src/components/stream/chat-panel.tsx @@ -5,8 +5,10 @@ import { } from "@atproto/api/dist/client/types/app/bsky/richtext/facet"; import type { LivestreamStore } from "@streamplace/core"; import { + chatMessageOpacity, formatBadgeIssuer, formatBadgeLabel, + isChatMessageGone, segmentize, type Facet, type FacetFeature, @@ -43,6 +45,20 @@ import { import { getAdjacentBadgeIndex } from "./badge-navigation"; import { initializeChatScroll } from "./chat-scroll"; +// Chat ages out on the hour, so a slow tick is plenty: it only exists so the +// fade and the disappearance happen while the viewer watches, rather than +// waiting for the next message to arrive and re-render the list. +const CHAT_EXPIRY_TICK_MS = 30_000; + +function useChatExpiryTick(): number { + const [now, setNow] = useState(() => Date.now()); + useEffect(() => { + const timer = setInterval(() => setNow(Date.now()), CHAT_EXPIRY_TICK_MS); + return () => clearInterval(timer); + }, []); + return now; +} + export function ChatPanel({ store, reversed = false, @@ -64,6 +80,7 @@ export function ChatPanel({ const scrollRef = useRef(null); const anchorRef = useRef(null); const [isAtAnchor, setIsAtAnchor] = useState(true); + const now = useChatExpiryTick(); const [newMessageCount, setNewMessageCount] = useState(0); const prevChatLenRef = useRef(chat.length); const initialScrollDoneRef = useRef(false); @@ -152,8 +169,14 @@ export function ChatPanel({ const displayMessages = useMemo(() => { const sliced = chat.slice(-1500); - return reversed ? [...sliced].reverse() : sliced; - }, [chat, reversed]); + // Old messages fade out of the live view and then leave it, but they stay + // in the store: a viewer scrolled back into history is reading, so there + // everything is shown at full strength. + const live = isAtAnchor + ? sliced.filter((m) => !isChatMessageGone(m.record.createdAt, now)) + : sliced; + return reversed ? [...live].reverse() : live; + }, [chat, reversed, isAtAnchor, now]); const badgeIssuerDids = useMemo(() => { const issuers = new Set(); for (const message of displayMessages) { @@ -211,6 +234,9 @@ export function ChatPanel({ store={store} isGrouped={isGrouped} issuerProfiles={issuerProfiles} + opacity={ + isAtAnchor ? chatMessageOpacity(msg.record.createdAt, now) : 1 + } /> ); }) @@ -249,6 +275,7 @@ function ChatMessage({ store, isGrouped = false, issuerProfiles, + opacity = 1, }: { message: ChatMessageViewHydrated; profile: ChatMessageViewHydrated["chatProfile"]; @@ -256,6 +283,8 @@ function ChatMessage({ store: LivestreamStore; isGrouped?: boolean; issuerProfiles: ReturnType; + /** 0..1 from chatMessageOpacity; 1 when the viewer is reading history. */ + opacity?: number; }) { const { t } = useTranslation("common"); const { state, pdsAgent, did } = useSession(); @@ -288,7 +317,10 @@ function ChatMessage({ if (isSystem) { return ( -
+

{message.record.text}

); @@ -296,6 +328,7 @@ function ChatMessage({ return (
{/* Hover actions; visible on group hover */} diff --git a/pkg/api/websocket.go b/pkg/api/websocket.go index ac91a28bc..fea801cce 100644 --- a/pkg/api/websocket.go +++ b/pkg/api/websocket.go @@ -265,7 +265,14 @@ func (a *StreamplaceAPI) HandleWebsocket(ctx context.Context) httprouter.Handle }() go func() { - messages, err := a.Model.MostRecentChatMessages(repoDID) + // Withhold chat older than the node's retention window: a viewer + // arriving the next day should land on an empty chat, not + // yesterday's conversation. A zero retention serves everything. + var since time.Time + if a.CLI.ChatMessageRetention > 0 { + since = time.Now().Add(-a.CLI.ChatMessageRetention) + } + messages, err := a.Model.MostRecentChatMessages(repoDID, since) if err != nil { log.Error(ctx, "could not get chat messages", "error", err) return diff --git a/pkg/atproto/backfill_walk_test.go b/pkg/atproto/backfill_walk_test.go index aa4e13bae..bda910f65 100644 --- a/pkg/atproto/backfill_walk_test.go +++ b/pkg/atproto/backfill_walk_test.go @@ -109,7 +109,7 @@ func TestBackfillWalk(t *testing.T) { require.Equal(t, head.Rev, stored.Version) require.Equal(t, head.Root.String(), stored.RootCID) - messages, err := mod.MostRecentChatMessages(user.DID) + messages, err := mod.MostRecentChatMessages(user.DID, time.Time{}) require.NoError(t, err) require.Len(t, messages, 2) texts := []string{ @@ -156,7 +156,7 @@ func TestBackfillWedgeHeals(t *testing.T) { if repo.Version == "" { return fmt.Errorf("repo still has no version") } - messages, err := mod.MostRecentChatMessages(user.DID) + messages, err := mod.MostRecentChatMessages(user.DID, time.Time{}) if err != nil { return err } @@ -244,7 +244,7 @@ func TestBackfillFallsBackToGetRepo(t *testing.T) { if err != nil { return err } - messages, err := mod.MostRecentChatMessages(user.DID) + messages, err := mod.MostRecentChatMessages(user.DID, time.Time{}) if err != nil { return err } @@ -1029,7 +1029,7 @@ func TestBackfillBroadcastsOnlyLiveChat(t *testing.T) { _, err := atsync.SyncBlueskyRepoCached(ctx, user.DID) require.NoError(t, err) - messages, err := mod.MostRecentChatMessages(user.DID) + messages, err := mod.MostRecentChatMessages(user.DID, time.Time{}) require.NoError(t, err) require.Len(t, messages, 2, "the walk indexed both messages") diff --git a/pkg/atproto/chat_message_test.go b/pkg/atproto/chat_message_test.go index 9dd2b6940..2935f935c 100644 --- a/pkg/atproto/chat_message_test.go +++ b/pkg/atproto/chat_message_test.go @@ -109,7 +109,7 @@ func TestChatMessage(t *testing.T) { messages := []placestream.ChatDefs_MessageView{} err = untilNoErrors(t, func() error { - messages, err = mod.MostRecentChatMessages(user.DID) + messages, err = mod.MostRecentChatMessages(user.DID, time.Time{}) if err != nil { return err } @@ -161,7 +161,7 @@ func TestChatMessage(t *testing.T) { require.NoError(t, err) err = untilNoErrors(t, func() error { - messages, err = mod.MostRecentChatMessages(user.DID) + messages, err = mod.MostRecentChatMessages(user.DID, time.Time{}) if err != nil { return err } diff --git a/pkg/atproto/deepen_test.go b/pkg/atproto/deepen_test.go index cbb1b125f..402d6396f 100644 --- a/pkg/atproto/deepen_test.go +++ b/pkg/atproto/deepen_test.go @@ -60,7 +60,7 @@ func TestDeepenForeverTrickles(t *testing.T) { return nil }), "the deepener should find the repo by scanning and walk it to the end") - messages, err := mod.MostRecentChatMessages(user.DID) + messages, err := mod.MostRecentChatMessages(user.DID, time.Time{}) require.NoError(t, err) require.Len(t, messages, 1, "the message in the deepest window was indexed") diff --git a/pkg/atproto/firehose_multirelay_test.go b/pkg/atproto/firehose_multirelay_test.go index 027b9143f..95b8db39f 100644 --- a/pkg/atproto/firehose_multirelay_test.go +++ b/pkg/atproto/firehose_multirelay_test.go @@ -84,7 +84,7 @@ func TestMultiRelayDedup(t *testing.T) { // The message is indexed exactly once, even though two relays delivered it. err = untilNoErrors(t, func() error { - messages, err := mod.MostRecentChatMessages(user.DID) + messages, err := mod.MostRecentChatMessages(user.DID, time.Time{}) if err != nil { return err } diff --git a/pkg/atproto/handle_test.go b/pkg/atproto/handle_test.go index 9e8faf1e2..be1124db1 100644 --- a/pkg/atproto/handle_test.go +++ b/pkg/atproto/handle_test.go @@ -74,7 +74,7 @@ func TestHandleChange(t *testing.T) { var message placestream.ChatDefs_MessageView err = untilNoErrors(t, func() error { - messages, err := mod.MostRecentChatMessages(user.DID) + messages, err := mod.MostRecentChatMessages(user.DID, time.Time{}) if err != nil { return err } @@ -94,7 +94,7 @@ func TestHandleChange(t *testing.T) { require.NoError(t, err) err = untilNoErrors(t, func() error { - messages, err := mod.MostRecentChatMessages(user.DID) + messages, err := mod.MostRecentChatMessages(user.DID, time.Time{}) if err != nil { return err } diff --git a/pkg/atproto/headcheck_test.go b/pkg/atproto/headcheck_test.go index 0eb17608f..a4715a893 100644 --- a/pkg/atproto/headcheck_test.go +++ b/pkg/atproto/headcheck_test.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "testing" + "time" "github.com/bluesky-social/indigo/atproto/identity" "github.com/bluesky-social/indigo/atproto/syntax" @@ -51,7 +52,7 @@ func TestHeadCheckHealsSilentGap(t *testing.T) { require.NoError(t, err) require.NotEmpty(t, indexed.Version) require.True(t, indexed.BackfillDone, "the sweep read the whole repo") - messages, err := mod.MostRecentChatMessages(user.DID) + messages, err := mod.MostRecentChatMessages(user.DID, time.Time{}) require.NoError(t, err) require.Len(t, messages, 1) @@ -85,7 +86,7 @@ func TestHeadCheckHealsSilentGap(t *testing.T) { require.Equal(t, hostRev, healed.Version, "the repair caught the index up to the host") require.True(t, healed.BackfillDone, "repairing a day of history does not un-index the rest") require.Empty(t, healed.RepairFrom, "the repair it asked for is the one that ran") - messages, err = mod.MostRecentChatMessages(user.DID) + messages, err = mod.MostRecentChatMessages(user.DID, time.Time{}) require.NoError(t, err) require.Len(t, messages, 2, "the record written during the gap is indexed") @@ -95,7 +96,7 @@ func TestHeadCheckHealsSilentGap(t *testing.T) { current, err := mod.GetRepo(user.DID) require.NoError(t, err) require.Equal(t, hostRev, current.Version) - messages, err = mod.MostRecentChatMessages(user.DID) + messages, err = mod.MostRecentChatMessages(user.DID, time.Time{}) require.NoError(t, err) require.Len(t, messages, 2) } diff --git a/pkg/atproto/redelivery_test.go b/pkg/atproto/redelivery_test.go index c80460ccb..88f5f81fb 100644 --- a/pkg/atproto/redelivery_test.go +++ b/pkg/atproto/redelivery_test.go @@ -111,7 +111,7 @@ func TestHandleCreateUpdateRedelivery(t *testing.T) { time.Sleep(250 * time.Millisecond) require.Equal(t, 1, countPublished(), "a redelivered chat message must not be published again") - messages, err := mod.MostRecentChatMessages(did) + messages, err := mod.MostRecentChatMessages(did, time.Time{}) require.NoError(t, err) require.Len(t, messages, 1, "a redelivered chat message must not be indexed again") require.Equal(t, "hello twice", messages[0].Record.Val.(*placestream.ChatMessage).Text) diff --git a/pkg/atproto/sweep_test.go b/pkg/atproto/sweep_test.go index 77dd900ba..7ee7d2f83 100644 --- a/pkg/atproto/sweep_test.go +++ b/pkg/atproto/sweep_test.go @@ -52,7 +52,7 @@ func TestBackfillWindowedHistory(t *testing.T) { plant(200*24*time.Hour, "two hundred days ago") countMessages := func() int { - messages, err := mod.MostRecentChatMessages(user.DID) + messages, err := mod.MostRecentChatMessages(user.DID, time.Time{}) require.NoError(t, err) return len(messages) } @@ -191,7 +191,7 @@ func TestSweepShallowThenDeepens(t *testing.T) { require.NoError(t, err) require.NotEmpty(t, stored.Version, "%s should have been synced", did) require.True(t, stored.BackfillDone, "%s should have been deepened to the end", did) - messages, err := mod.MostRecentChatMessages(did) + messages, err := mod.MostRecentChatMessages(did, time.Time{}) require.NoError(t, err, did) require.Len(t, messages, 2, "both messages of %s should be indexed", did) } @@ -205,7 +205,7 @@ func TestSweepShallowThenDeepens(t *testing.T) { // Idempotent: a second sweep is a few head fetches and nothing else. require.NoError(t, atsync.Sweep(ctx)) for _, did := range []string{fresh.DID, legacy.DID} { - messages, err := mod.MostRecentChatMessages(did) + messages, err := mod.MostRecentChatMessages(did, time.Time{}) require.NoError(t, err) require.Len(t, messages, 2, "a second sweep must not duplicate anything") } diff --git a/pkg/config/config.go b/pkg/config/config.go index 50948c0ef..6c69c0cf6 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -153,6 +153,7 @@ type CLI struct { StreamSessionTimeout time.Duration LegacySegmentCleaner bool SegmentArchiveRetention time.Duration + ChatMessageRetention time.Duration Replicators []string WebsocketURL string BehindHTTPSProxy bool @@ -210,6 +211,16 @@ type CLI struct { // increase in traffic against every PDS on the network. const DefaultSweepInterval = 6 * time.Hour +// DefaultChatMessageRetention is how much chat history a node serves when +// --chat-message-retention is unset. +// +// Chat is a live conversation, not an archive: viewers who arrive the next day +// should land on an empty chat rather than a wall of stale messages. A day is +// long enough that a viewer rejoining the same day still sees what was said, +// and short enough that the backlog stays a conversation. The rows themselves +// stay in the index; this only bounds what the websocket backlog hands out. +const DefaultChatMessageRetention = 24 * time.Hour + // DefaultSweepConcurrency is how many PDS hosts the atproto backfill sweep // works on at once when --sweep-concurrency is unset or zero. // @@ -1030,6 +1041,13 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { Destination: &cli.SegmentArchiveRetention, Sources: urfavecli.EnvVars("SP_SEGMENT_ARCHIVE_RETENTION"), }, + &urfavecli.DurationFlag{ + Name: "chat-message-retention", + Usage: "how much chat history a node serves to newly connected viewers. Messages older than this are withheld from the chat backlog, so a viewer arriving later lands on an empty chat. 0 serves the full backlog.", + Value: DefaultChatMessageRetention, + Destination: &cli.ChatMessageRetention, + Sources: urfavecli.EnvVars("SP_CHAT_MESSAGE_RETENTION"), + }, &urfavecli.StringFlag{ Name: "replicators", Usage: "comma-separated list of replication protocols to use (websocket, iroh)", diff --git a/pkg/model/chat_message.go b/pkg/model/chat_message.go index a93f9537a..63a8b77f8 100644 --- a/pkg/model/chat_message.go +++ b/pkg/model/chat_message.go @@ -131,15 +131,23 @@ func (m *DBModel) GetChatMessage(uri string) (*ChatMessage, error) { return &message, nil } -func (m *DBModel) MostRecentChatMessages(repoDID string) ([]placestream.ChatDefs_MessageView, error) { +func (m *DBModel) MostRecentChatMessages(repoDID string, since time.Time) ([]placestream.ChatDefs_MessageView, error) { dbmessages := []ChatMessage{} - err := m.DB. + query := m.DB. Preload("Repo"). Preload("ChatProfile"). Preload("ReplyTo"). Preload("ReplyTo.Repo"). Preload("ReplyTo.ChatProfile"). - Where("streamer_repo_did = ?", repoDID). + Where("streamer_repo_did = ?", repoDID) + if !since.IsZero() { + // Chat is a live conversation, not an archive: withhold history older + // than the node's retention window so a viewer arriving later lands on + // an empty chat instead of a wall of stale messages. A zero `since` + // serves everything, which is what a 0 retention asks for. + query = query.Where("chat_messages.created_at >= ?", since) + } + err := query. // Exclude messages from users blocked by the streamer Joins("LEFT JOIN blocks ON blocks.repo_did = chat_messages.streamer_repo_did AND blocks.subject_did = chat_messages.repo_did"). Where("blocks.rkey IS NULL"). // Only include messages where no block exists diff --git a/pkg/model/chat_message_test.go b/pkg/model/chat_message_test.go new file mode 100644 index 000000000..6ba5cb366 --- /dev/null +++ b/pkg/model/chat_message_test.go @@ -0,0 +1,63 @@ +package model + +import ( + "bytes" + "context" + "testing" + "time" + + "github.com/stretchr/testify/require" + "stream.place/streamplace/pkg/placestream" +) + +// TestMostRecentChatMessagesRetention: a node serves chat only inside its +// retention window, so a viewer arriving the next day lands on an empty chat +// instead of yesterday's conversation. The filter is a read, not a delete -- +// the rows stay indexed -- and a zero cutoff still serves the full backlog. +func TestMostRecentChatMessagesRetention(t *testing.T) { + ctx := context.Background() + mod := indexedTestDB(t) + + now := time.Now().UTC() + rec := &placestream.ChatMessage{ + LexiconTypeID: "place.stream.chat.message", + Text: "hello", + CreatedAt: now.Format(time.RFC3339), + Streamer: "did:plc:streamer", + } + var buf bytes.Buffer + require.NoError(t, rec.MarshalCBOR(&buf)) + body := buf.Bytes() + + require.NoError(t, mod.CreateChatMessage(ctx, &ChatMessage{ + CID: "bafyold", + URI: "at://did:plc:chatter/place.stream.chat.message/3lold", + CreatedAt: now.Add(-25 * time.Hour), + ChatMessage: &body, + RepoDID: "did:plc:chatter", + StreamerRepoDID: "did:plc:streamer", + IndexedAt: &now, + })) + require.NoError(t, mod.CreateChatMessage(ctx, &ChatMessage{ + CID: "bafyrecent", + URI: "at://did:plc:chatter/place.stream.chat.message/3lrecent", + CreatedAt: now.Add(-time.Hour), + ChatMessage: &body, + RepoDID: "did:plc:chatter", + StreamerRepoDID: "did:plc:streamer", + IndexedAt: &now, + })) + + all, err := mod.MostRecentChatMessages("did:plc:streamer", time.Time{}) + require.NoError(t, err) + require.Len(t, all, 2, "a zero cutoff serves the whole backlog") + + recent, err := mod.MostRecentChatMessages("did:plc:streamer", now.Add(-24*time.Hour)) + require.NoError(t, err) + require.Len(t, recent, 1, "a message older than the window is withheld") + require.Equal(t, + "at://did:plc:chatter/place.stream.chat.message/3lrecent", recent[0].Uri) + + require.Equal(t, int64(2), countRows(t, mod, &ChatMessage{}), + "withholding history must not delete it") +} diff --git a/pkg/model/model.go b/pkg/model/model.go index caf9383a8..cbe0c0d5c 100644 --- a/pkg/model/model.go +++ b/pkg/model/model.go @@ -92,7 +92,7 @@ type Model interface { DeleteBlock(ctx context.Context, rkey string) error CreateChatMessage(ctx context.Context, message *ChatMessage) error - MostRecentChatMessages(repoDID string) ([]placestream.ChatDefs_MessageView, error) + MostRecentChatMessages(repoDID string, since time.Time) ([]placestream.ChatDefs_MessageView, error) GetChatMessage(uri string) (*ChatMessage, error) DeleteChatMessage(ctx context.Context, uri string, deletedAt *time.Time) error