diff --git a/CLAUDE.md b/CLAUDE.md index d722b89..183abd6 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -165,7 +165,7 @@ This is the most intricate part of the codebase and where most current work happ **Chat IDs.** Length is the type discriminator: 32 bytes = DM (SHA-256 over a domain-separated, sorted, deduped member set, so creation is idempotent and order-independent), 16 bytes = group (server-derived: a truncated, domain-separated SHA-256 over the creator's user ID and the request's required `IdempotencyKey`, so a retried `StartChat` names the same group and is answered from the existing record; stamped as a version 8 UUID so every group ID is UUID-shaped, though the ID is opaque and nothing parses it). DM paths must reject 16-byte IDs and vice versa. `CONTACT_DM` uses the bare legacy hash domain; other DM types append their enum number. A chat type of `UNKNOWN` falls back to `CONTACT_DM` for legacy clients. -**Chat storage (`chat/dynamodb`).** Tables: `chats` (metadata), `dm_inbox` (per-user DM feed rows; GSI `by_type_activity` on a composite `feed` key, legacy `by_activity` GSI still maintained), `group_members` (pk chat, sk user; plus one `#meta` item per group holding `member_count`/`version`, CAS-updated in the same transaction as each transition, with bounded retries on contention). `Chat.Members` is populated only for DMs; groups return empty `Members` and a `RosterSummary{MemberCount, Version}`. Version is *state, not a delta*: each real transition bumps it by exactly one and no-ops leave it alone; clients keep the greater version. DM sends fan `last_activity` into each member's inbox row. A store built with users in `excludedFromFeed` (a required argument of `NewInMemory` / `NewInDynamoDB`, nil for none, nil entries ignored, duplicates collapsed; `chat.FeedExclusions` in `chat/feed.go`) creates every DM with those users **excluded from the feed**; the parent passes the team account, so no flow that creates a DM can leave it in. The exclusion is store-internal, not on `Chat`: decided at creation, recorded on the DM's canonical item as `excluded_from_feed` (a binary set of user IDs, absent when empty; a DM between two excluded users excludes both), and each excluded member's `dm_inbox` row carries neither `feed` nor `last_activity`, so it is in neither GSI and records membership alone (kept because DM `IsMember` reads it). `GetDmFeedPage` never lists the chat for them, and `AdvanceLastMessage` skips their row off the canonical item it already reads, so a process built excluding no one still advances such a DM safely (its condition could never hold on that row, so including it would cancel every advance); opening DMs never writes the newest end of one GSI key and no send moves the team's row; **group sends never fan out** — the group feed is assembled at read time with order computed once and a window of chat IDs carried in the paging token (`maxGroupFeedChats = 1000`), re-checking membership per page. A fourth table, `chat_user_state` (pk user, sk chat), holds `chat.ViewerState`: what a chat records about one user independent of membership — today a mute (`muted_until`, epoch seconds, present only while a mute is recorded; an indefinite mute is a store-internal far-future sentinel) and a per-row `version` with the same state-not-delta rule. Rows are sparse, never deleted, and survive leaving the chat. The viewer-state methods (`SetMute`, `ClearMute`, `GetViewerStates`, `GetMutedUsers`, `GetMutedUsersPage`, `GetMutedCount`) are part of `chat.Store` itself, not a separate interface. `GetViewerStates` is one strongly consistent `Query` on the user's partition bounded to the requested key range (a page of chats costs a few RCU, not one per key). The chat-scoped read has **two shapes, chosen by size**: `GetMutedUsers` is a key-range `Query` on the sparse `by_muted` GSI (chat, `muted_until`), billed by the active mutes it returns; `GetMutedUsersPage` ranges the inverted `by_user` GSI (chat, user; every record, full projection so future states need no new index) over an inclusive `[lo, hi]` user-ID key range in the same `user#` order as a group's roster, so a fan-out holding one roster page (`GetGroupMembersPage`, a cursor walk of the `group_members` partition in ascending user-ID order) asks for the mutes between that page's first and last user and never holds either whole. `GetMutedCount` reads a per-chat `#meta` item (pk `chat#`, sk `#meta`; carries neither `chat` nor `muted_until`, which keeps it out of both indexes) holding the number of records with a mute *recorded* — moved in the same transaction as a first mute or a clear, never by a replace or a lapse, so it bounds active mutes from above; the fan-out compares it to the roster size to pick a shape. Both GSIs are eventually consistent. A replaced mute is one conditional `UpdateItem`; a first mute or a clear is a two-item transaction on the `#meta` pattern (version compared, lost race retried from the returned item, `TransactionConflict` backed off). A no-op never creates an item. A fifth table, `chat_activity` (pk chat, sk user; `NewInDynamoDB` / `CreateTables` take its name), holds each group's **activity records** for mention suggestions (`chat.RecentSender` in `chat/model.go`), a fact about the message log, independent of membership: `RecordSend` is one blind conditional update of `last_sent_at` (epoch ms), throttled to one per `ActivityRecordInterval` (1 min) per user, never moving backwards, and expiring by TTL `ActivityRetention` (1 year) after it; `GetRecentSenders` is one eventually consistent query on the `by_last_sent_at` LSI, most recent first; `GetLastSentAt` is the eventually consistent point read of one user's record (both at half a strong read's cost; their readers only rank and show people). A second LSI, `by_activity_score`, is reserved for a future frequency-weighted ordering and empty today: nothing writes `activity_score`, and it exists only because an LSI cannot be added later. `messaging.Sender` records each sender's last message of every group send (`recordGroupSenders`, after the broadcast, best effort; DMs and system messages record nothing); `GetMentionSuggestions` reads `GetRecentSenders` (see "Mention suggestions" below). A sixth table, `chat_key_envelopes` (pk user, sk chat, no index; `NewInDynamoDB` / `CreateTables` take its name), holds each member's **key envelope** for a private group (`chat.KeyEnvelope`: scheme, nonce, ciphertext and `wrapped_by`, the user who stored it; see "Private groups"). `SetKeyEnvelope` is one conditional put, refused when the stored envelope is one the user wrapped themself, with the refusal returning the stored item; `GetKeyEnvelope` is a strongly consistent point read. Like viewer state, the methods are on `chat.Store`, write against the IDs alone and know nothing of the roster. The one tie to the roster is a departure: `RemoveGroupMember` takes a required `discardKeyEnvelope` flag and, when set, deletes the leaver's envelope **in the same transaction** as the membership transition and its `#meta` update (an unconditional third item, so it commits iff the departure happens and a no-op leave deletes nothing). A seventh table, `chat_lobbies` (pk user, sk chat; `NewInDynamoDB` / `CreateTables` take its name), holds each private group's **lobby entries** (`chat.LobbyEntry`: `entered_at` epoch nanos; the sparse `by_chat` **GSI** is hashed on the `sk` itself, the chat, and ranged on `entered_at`, so no attribute repeats the chat and the index pages a lobby earliest first, **eventually consistent** as the proto allows) and two `#meta` count items per entry, the lobby's size under `chat#` and the user's lobby count under `user#` (both omit `entered_at`, the index's range key, so they stay out of it). **Keyed by user, decided 2026-10-05**: a user's own entries, and the future listing of every lobby they wait in, are one strongly consistent query with no index; the price is the creator's page, which a chat-keyed table with an LSI could serve strongly consistent (an LSI cannot gather one chat's entries across user partitions), and the creator reconciles against `LobbyUpdate`s instead. `EnterLobby` is one transaction: a conditional put of the entry, both counts incremented under `attribute_not_exists OR count < cap`, and a `ConditionCheck` that the user's `group_members` row is absent or not joined, so the caps (`chat.LobbyLimits`, `DefaultLobbyLimits` 100 per lobby / 100 lobbies per user, user's choice 2026-10-05) hold under concurrent entries with no read and **a member never gains an entry** (closing the race where a retried `EnterLobby` lands after the creator's admission removed the entry); the cancellation reasons tell `ErrAlreadyMember` (judged first) from an existing entry (the no-op) from `ErrLobbyFull` (judged before) `ErrTooManyLobbies`. `LeaveLobby` is the inverse (conditional delete + two decrements; a missing entry is the no-op). `GetLobbyEntries` is a strongly consistent read of the user's partition (point read for one chat, bounded range query otherwise, like `GetViewerStates`). `AdmitFromLobby` is `transitionMembership`'s join carrying `alongside` the entry's **conditional** delete, both decrements and an unconditional put of the key envelope: `transitionMembership` now admits conditional alongside items and reports a failed one as `alongsideConditionFailed{Index}` (judged after the transition's own no-op and the summary CAS), which `AdmitFromLobby` translates to `ErrNotInLobby`; a user who is already a member is the transition's no-op (`changed` false, nothing written, a leftover entry is theirs to clear). An eighth table, `chat_featured_groups` (pk user, sk `pos#`, no index; `NewInDynamoDB` / `CreateTables` take its name), holds each user's **featured groups** (`chat/featured.go`: an ordered list of up to `MaxFeaturedGroups` (10) **public** group IDs, no membership required, nothing recorded against the group; the store reads no record, so existence and the public-only rule are the setter's (`DENIED` for a private group), and since `IsPrivate` is immutable the list needs no filtering on read; `ValidateFeaturedGroups` refuses a DM ID or a repeat rather than collapsing it). `SetFeaturedGroups` replaces the whole list in one transaction: the `#meta` item (`version`, `featured_count`) compare-and-set on version, a put of every position the new list fills (each stamped with the new version) and a delete of every position past it, retried from a fresh read on a lost CAS; an equal list is the no-op. `GetFeaturedGroups` is one **eventually consistent** query of the partition (it may trail a replace by a moment; the replace's own read of the list stays strongly consistent, since a stale one could skip a needed write as a no-op), verified (every position at `#meta`'s version, `featured_count` of them, contiguous) and re-read up to 3× when it straddles a replace, so a reader never sees a mix of two lists. Rows carry the raw `chat` so a `(chat, pk)` GSI for "who features this group" can be added later; `#meta` omits it. **RPCs (`chat/featured.go`) are on the Chat service**, not Profile (`profilepb` cannot import `chatpb`, which imports it). `SetFeaturedGroups` (auth required) validates (`InvalidArgument` for a DM ID, a repeat or too many), reads the groups' records with `Store.GetGroupChatsByID` (eventually consistent, so a group created a moment ago may briefly be `NOT_FOUND`; no membership check; a TODO marks revisiting its consistency if another reader comes to use it), answers `NOT_FOUND` for a missing group before `DENIED` for a private one, writes, and returns the list hydrated from the records it already read. `GetFeaturedGroups` takes a username alone (resolved with `ProfileReader.GetUserIDByUsername`; `NOT_FOUND` when nobody holds it), auth optional and changing nothing, and drops groups with no record. Both project through `featuredGroupsMetadata`: `hydrate` with no viewer, `ReadingDenied` and `listDetail`, so every caller gets the same list (no viewer fields, no messaging state, no cover); it also drops a private group, which should never be there, and logs a warning if it finds one. +**Chat storage (`chat/dynamodb`).** Tables: `chats` (metadata), `dm_inbox` (per-user DM feed rows; GSI `by_type_activity` on a composite `feed` key, legacy `by_activity` GSI still maintained), `group_members` (pk chat, sk user; plus one `#meta` item per group holding `member_count`/`version`, CAS-updated in the same transaction as each transition, with bounded retries on contention). `Chat.Members` is populated only for DMs; groups return empty `Members` and a `RosterSummary{MemberCount, Version}`. Version is *state, not a delta*: each real transition bumps it by exactly one and no-ops leave it alone; clients keep the greater version. DM sends fan `last_activity` into each member's inbox row. A store built with users in `excludedFromFeed` (a required argument of `NewInMemory` / `NewInDynamoDB`, nil for none, nil entries ignored, duplicates collapsed; `chat.FeedExclusions` in `chat/feed.go`) creates every DM with those users **excluded from the feed**; the parent passes the team account, so no flow that creates a DM can leave it in. The exclusion is store-internal, not on `Chat`: decided at creation, recorded on the DM's canonical item as `excluded_from_feed` (a binary set of user IDs, absent when empty; a DM between two excluded users excludes both), and each excluded member's `dm_inbox` row carries neither `feed` nor `last_activity`, so it is in neither GSI and records membership alone (kept because DM `IsMember` reads it). `GetDmFeedPage` never lists the chat for them, and `AdvanceLastMessage` skips their row off the canonical item it already reads, so a process built excluding no one still advances such a DM safely (its condition could never hold on that row, so including it would cancel every advance); opening DMs never writes the newest end of one GSI key and no send moves the team's row; **group sends never fan out** — the group feed is assembled at read time with order computed once and a window of chat IDs carried in the paging token (`maxGroupFeedChats = 1000`), re-checking membership per page. A fourth table, `chat_user_state` (pk user, sk chat), holds `chat.ViewerState`: what a chat records about one user independent of membership — today a mute (`muted_until`, epoch seconds, present only while a mute is recorded; an indefinite mute is a store-internal far-future sentinel) and a per-row `version` with the same state-not-delta rule. Rows are sparse, never deleted, and survive leaving the chat. The viewer-state methods (`SetMute`, `ClearMute`, `GetViewerStates`, `GetMutedUsers`, `GetMutedUsersPage`, `GetMutedCount`) are part of `chat.Store` itself, not a separate interface. `GetViewerStates` is one strongly consistent `Query` on the user's partition bounded to the requested key range (a page of chats costs a few RCU, not one per key). The chat-scoped read has **two shapes, chosen by size**: `GetMutedUsers` is a key-range `Query` on the sparse `by_muted` GSI (chat, `muted_until`), billed by the active mutes it returns; `GetMutedUsersPage` ranges the inverted `by_user` GSI (chat, user; every record, full projection so future states need no new index) over an inclusive `[lo, hi]` user-ID key range in the same `user#` order as a group's roster, so a fan-out holding one roster page (`GetGroupMembersPage`, a cursor walk of the `group_members` partition in ascending user-ID order) asks for the mutes between that page's first and last user and never holds either whole. `GetMutedCount` reads a per-chat `#meta` item (pk `chat#`, sk `#meta`; carries neither `chat` nor `muted_until`, which keeps it out of both indexes) holding the number of records with a mute *recorded* — moved in the same transaction as a first mute or a clear, never by a replace or a lapse, so it bounds active mutes from above; the fan-out compares it to the roster size to pick a shape. Both GSIs are eventually consistent. A replaced mute is one conditional `UpdateItem`; a first mute or a clear is a two-item transaction on the `#meta` pattern (version compared, lost race retried from the returned item, `TransactionConflict` backed off). A no-op never creates an item. A fifth table, `chat_activity` (pk chat, sk user; `NewInDynamoDB` / `CreateTables` take its name), holds each group's **activity records** for mention suggestions (`chat.RecentSender` in `chat/model.go`), a fact about the message log, independent of membership: `RecordSend` is an eventually consistent read then one update of `last_sent_at` (epoch ms) and `activity_score`, conditioned on the `last_sent_at` read (a lost race is checked once more against the item the failure returns, and written once more only when its send is the newest by a full interval; a second loss leaves it unrecorded, best effort), throttled to one per `ActivityRecordInterval` (1 min) per user (a throttled send is answered off the read, with no write), never moving backwards, and expiring by TTL `ActivityRetention` (1 year) after it; `GetRecentSenders` is one eventually consistent query on the `by_last_sent_at` LSI, most recent first; `GetLastSentAt` is the eventually consistent point read of one user's record (both at half a strong read's cost; their readers only rank and show people). A second LSI, `by_activity_score`, orders the partition by **activity score** (`chat/activity.go`), a frequency-weighted ordering of recorded sends that **nothing reads yet**: the log-sum-exp of the send times with a 3-day half-life (`ActivityScoreHalfLife`), each send weighted by the quiet before it over `ActivityScoreFullWeightGap` (1 day, at most one, so it counts **days active, not messages**: an hour-long burst barely counts past its first send), stored as epoch ms, so one send scores its time and the decay never moves a stored score (comparing decayed sums at any moment compares the scores); never below the last send, and capped at `ActivityScoreMaxLead` (7 days) past the send that set it, which no steady pattern reaches at these values (sending daily settles 6.8 days ahead, weekly about 1). `NextActivityScore` is the one formula, shared by every store. A record written before scores has no `activity_score`, is missing from the LSI, and reads and scores as its `last_sent_at` (`EffectiveActivityScore`, which also lifts a score trailing `last_sent_at`, as an old writer during rollout can leave) until a backfill sets `activity_score = last_sent_at`. `GetRecentSenders` returns each record's effective score. `messaging.Sender` records each sender's last message of every group send (`recordGroupSenders`, after the broadcast, best effort; DMs and system messages record nothing); `GetMentionSuggestions` reads `GetRecentSenders` (see "Mention suggestions" below). A sixth table, `chat_key_envelopes` (pk user, sk chat, no index; `NewInDynamoDB` / `CreateTables` take its name), holds each member's **key envelope** for a private group (`chat.KeyEnvelope`: scheme, nonce, ciphertext and `wrapped_by`, the user who stored it; see "Private groups"). `SetKeyEnvelope` is one conditional put, refused when the stored envelope is one the user wrapped themself, with the refusal returning the stored item; `GetKeyEnvelope` is a strongly consistent point read. Like viewer state, the methods are on `chat.Store`, write against the IDs alone and know nothing of the roster. The one tie to the roster is a departure: `RemoveGroupMember` takes a required `discardKeyEnvelope` flag and, when set, deletes the leaver's envelope **in the same transaction** as the membership transition and its `#meta` update (an unconditional third item, so it commits iff the departure happens and a no-op leave deletes nothing). A seventh table, `chat_lobbies` (pk user, sk chat; `NewInDynamoDB` / `CreateTables` take its name), holds each private group's **lobby entries** (`chat.LobbyEntry`: `entered_at` epoch nanos; the sparse `by_chat` **GSI** is hashed on the `sk` itself, the chat, and ranged on `entered_at`, so no attribute repeats the chat and the index pages a lobby earliest first, **eventually consistent** as the proto allows) and two `#meta` count items per entry, the lobby's size under `chat#` and the user's lobby count under `user#` (both omit `entered_at`, the index's range key, so they stay out of it). **Keyed by user, decided 2026-10-05**: a user's own entries, and the future listing of every lobby they wait in, are one strongly consistent query with no index; the price is the creator's page, which a chat-keyed table with an LSI could serve strongly consistent (an LSI cannot gather one chat's entries across user partitions), and the creator reconciles against `LobbyUpdate`s instead. `EnterLobby` is one transaction: a conditional put of the entry, both counts incremented under `attribute_not_exists OR count < cap`, and a `ConditionCheck` that the user's `group_members` row is absent or not joined, so the caps (`chat.LobbyLimits`, `DefaultLobbyLimits` 100 per lobby / 100 lobbies per user, user's choice 2026-10-05) hold under concurrent entries with no read and **a member never gains an entry** (closing the race where a retried `EnterLobby` lands after the creator's admission removed the entry); the cancellation reasons tell `ErrAlreadyMember` (judged first) from an existing entry (the no-op) from `ErrLobbyFull` (judged before) `ErrTooManyLobbies`. `LeaveLobby` is the inverse (conditional delete + two decrements; a missing entry is the no-op). `GetLobbyEntries` is a strongly consistent read of the user's partition (point read for one chat, bounded range query otherwise, like `GetViewerStates`). `AdmitFromLobby` is `transitionMembership`'s join carrying `alongside` the entry's **conditional** delete, both decrements and an unconditional put of the key envelope: `transitionMembership` now admits conditional alongside items and reports a failed one as `alongsideConditionFailed{Index}` (judged after the transition's own no-op and the summary CAS), which `AdmitFromLobby` translates to `ErrNotInLobby`; a user who is already a member is the transition's no-op (`changed` false, nothing written, a leftover entry is theirs to clear). An eighth table, `chat_featured_groups` (pk user, sk `pos#`, no index; `NewInDynamoDB` / `CreateTables` take its name), holds each user's **featured groups** (`chat/featured.go`: an ordered list of up to `MaxFeaturedGroups` (10) **public** group IDs, no membership required, nothing recorded against the group; the store reads no record, so existence and the public-only rule are the setter's (`DENIED` for a private group), and since `IsPrivate` is immutable the list needs no filtering on read; `ValidateFeaturedGroups` refuses a DM ID or a repeat rather than collapsing it). `SetFeaturedGroups` replaces the whole list in one transaction: the `#meta` item (`version`, `featured_count`) compare-and-set on version, a put of every position the new list fills (each stamped with the new version) and a delete of every position past it, retried from a fresh read on a lost CAS; an equal list is the no-op. `GetFeaturedGroups` is one **eventually consistent** query of the partition (it may trail a replace by a moment; the replace's own read of the list stays strongly consistent, since a stale one could skip a needed write as a no-op), verified (every position at `#meta`'s version, `featured_count` of them, contiguous) and re-read up to 3× when it straddles a replace, so a reader never sees a mix of two lists. Rows carry the raw `chat` so a `(chat, pk)` GSI for "who features this group" can be added later; `#meta` omits it. **RPCs (`chat/featured.go`) are on the Chat service**, not Profile (`profilepb` cannot import `chatpb`, which imports it). `SetFeaturedGroups` (auth required) validates (`InvalidArgument` for a DM ID, a repeat or too many), reads the groups' records with `Store.GetGroupChatsByID` (eventually consistent, so a group created a moment ago may briefly be `NOT_FOUND`; no membership check; a TODO marks revisiting its consistency if another reader comes to use it), answers `NOT_FOUND` for a missing group before `DENIED` for a private one, writes, and returns the list hydrated from the records it already read. `GetFeaturedGroups` takes a username alone (resolved with `ProfileReader.GetUserIDByUsername`; `NOT_FOUND` when nobody holds it), auth optional and changing nothing, and drops groups with no record. Both project through `featuredGroupsMetadata`: `hydrate` with no viewer, `ReadingDenied` and `listDetail`, so every caller gets the same list (no viewer fields, no messaging state, no cover); it also drops a private group, which should never be there, and logs a warning if it finds one. **Tombstones.** A departed group member's row stays with `state=2`, `left_at`, a 1h TTL and the departure's roster version, so re-joins are idempotent updates and delayed duplicates of an undone join can be distinguished from news. Readers trust `state`, never the clock. "Formerly a member" is not durable. diff --git a/chat/activity.go b/chat/activity.go new file mode 100644 index 0000000..ddb3a7b --- /dev/null +++ b/chat/activity.go @@ -0,0 +1,93 @@ +package chat + +import ( + "math" + "time" +) + +// A group's activity score is a frequency-weighted ordering of its senders +// kept on each activity record beside the last send (see RecentSender): a +// record ranks above another when its recorded sends, each weighted and each +// decaying with ActivityScoreHalfLife, sum to more as of any one moment. +// +// A send is weighted by the quiet before it: the time since the user's +// previous recorded send over ActivityScoreFullWeightGap, at most one, and +// one for their first. So the score counts the days a user shows up rather +// than the messages they send: a day of chatting weighs about what one +// message that day would, and an hour-long burst barely more than its first +// message, however many sends the throttle (ActivityRecordInterval) lets +// through. +// +// It is stored as a time, not a weight: the moment at which a single send +// would rank equal to the record. A user with one recorded send scores their +// send time, and further sends move the score ahead of the last one, +// possibly past now. In that form the decay never moves a score: comparing +// two decayed sums at any moment compares their scores, because the moment +// cancels out of the comparison, so a store can keep the score as a sort key +// and rewrite it only on a recorded send. A record written before scores +// existed reads as scoring its last send, which is exactly its score had it +// had one send (see EffectiveActivityScore). A score is never below its last +// send, so a record never ranks below someone who sent once, after it. +// +// A user who sends every day settles about 6.8 days ahead of their last +// send, which is the most any steady pattern reaches at a 3-day half-life +// and a 1-day full-weight gap; one who sends weekly, about a day. That lead +// is how long they rank above someone who sends once after them. +// ActivityScoreMaxLead bounds it whatever the history; at these values no +// steady pattern reaches it. +// +// Nothing ranks by it yet. +const ( + // ActivityScoreHalfLife is how long a recorded send takes to count half + // as much toward an activity score. + ActivityScoreHalfLife = 72 * time.Hour + + // ActivityScoreFullWeightGap is the quiet before a send at which it + // counts in full toward an activity score; a send after less counts in + // proportion. + ActivityScoreFullWeightGap = 24 * time.Hour + + // ActivityScoreMaxLead is the furthest an activity score may run ahead of + // the send that set it. + ActivityScoreMaxLead = 7 * 24 * time.Hour +) + +// NextActivityScore returns an activity record's score after a send at sentAt +// is recorded on it, given its score before (see EffectiveActivityScore) and +// its last recorded send before, or the zero time for both for a user with no +// record. The result is at millisecond precision, like the record's send +// time, never before sentAt, and never below prior, so a score only moves +// forward with the sends that set it. +func NextActivityScore(prior, lastSentAt, sentAt time.Time) time.Time { + t := float64(sentAt.UnixMilli()) + if prior.IsZero() { + return time.UnixMilli(int64(t)).UTC() + } + s := float64(prior.UnixMilli()) + + weight := min(1, float64(sentAt.Sub(lastSentAt))/float64(ActivityScoreFullWeightGap)) + next := math.Max(s, t) + if weight > 0 { + // τ·ln(e^(s/τ) + w·e^(t/τ)), arranged so neither exponential + // overflows. + tau := float64(ActivityScoreHalfLife.Milliseconds()) / math.Ln2 + u := t + tau*math.Log(weight) + next = math.Max(s, u) + tau*math.Log1p(math.Exp(-math.Abs(s-u)/tau)) + next = math.Min(next, t+float64(ActivityScoreMaxLead.Milliseconds())) + next = math.Max(next, math.Max(s, t)) + } + return time.UnixMilli(int64(math.Round(next))).UTC() +} + +// EffectiveActivityScore returns the score of an activity record whose last +// recorded send is lastSentAt and whose stored score is score, the zero time +// when it has none. A record with no score, written before scores existed, +// scores its last send. So does one whose score trails its last send, which +// only a writer that records sends without scoring them leaves behind, at the +// cost of the sends it did not score: a score is never below the last send. +func EffectiveActivityScore(score, lastSentAt time.Time) time.Time { + if score.Before(lastSentAt) { + return lastSentAt + } + return score +} diff --git a/chat/activity_test.go b/chat/activity_test.go new file mode 100644 index 0000000..18b5261 --- /dev/null +++ b/chat/activity_test.go @@ -0,0 +1,108 @@ +package chat + +import ( + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +func TestNextActivityScore(t *testing.T) { + start := time.Unix(1_700_000_000, 0).UTC() + day := 24 * time.Hour + + // A first send scores its own time, at millisecond precision. + require.Equal(t, start, NextActivityScore(time.Time{}, time.Time{}, start)) + require.Equal(t, start.Add(time.Millisecond), NextActivityScore(time.Time{}, time.Time{}, start.Add(1_500*time.Microsecond))) + + // A send after no quiet counts for nothing. + require.Equal(t, start, NextActivityScore(start, start, start)) + + // A send after part of the full-weight gap counts in proportion: weight w + // at t is one full send at t + half-life·log2(w), so a send after half + // the gap counts as a full send one half-life before it. + half := start.Add(ActivityScoreFullWeightGap / 2) + equivalent := half.Add(-ActivityScoreHalfLife) + requireWithinMilli(t, + NextActivityScore(start, equivalent.Add(-ActivityScoreFullWeightGap), equivalent), + NextActivityScore(start, start, half), + ) + + // A send long after the last counts almost alone. + requireWithinMilli(t, start.Add(365*day), NextActivityScore(start, start, start.Add(365*day))) + + // The lead a steady pattern settles at, in days, after 120 days of it. + settledLead := func(sendsPerDay int, spacing time.Duration) float64 { + var score, last time.Time + for d := range 120 { + for k := range sendsPerDay { + sentAt := start.Add(time.Duration(d)*day + time.Duration(k)*spacing) + score = NextActivityScore(score, last, sentAt) + require.False(t, score.Before(sentAt)) + require.False(t, score.After(sentAt.Add(ActivityScoreMaxLead))) + last = sentAt + } + } + return score.Sub(last).Hours() / 24 + } + + // One send a day settles about 6.8 days ahead of the last, and a daily + // half-hour session, or sending all day, about the same: the score counts + // days, not messages. + daily := settledLead(1, 0) + require.InDelta(t, 6.84, daily, 0.01) + require.InDelta(t, daily, settledLead(30, ActivityRecordInterval), 0.05) + require.InDelta(t, 6.4, settledLead(16, time.Hour), 0.05) + + // One send a week, about a day. + var score, last time.Time + for w := range 26 { + sentAt := start.Add(time.Duration(w) * 7 * day) + score, last = NextActivityScore(score, last, sentAt), sentAt + } + require.InDelta(t, 0.96, score.Sub(last).Hours()/24, 0.01) + + // A daily sender outranks a one-off sender who sent after them, until + // their lead runs out. + dailyLast := start.Add(119 * day) + dailyScore := dailyLast.Add(time.Duration(daily * float64(day))) + require.True(t, dailyScore.After(NextActivityScore(time.Time{}, time.Time{}, dailyLast.Add(6*day)))) + require.True(t, dailyScore.Before(NextActivityScore(time.Time{}, time.Time{}, dailyLast.Add(7*day)))) + + // An hour-long burst at the throttle's rate barely counts past its first + // send. + score, last = time.Time{}, time.Time{} + for k := range 60 { + sentAt := start.Add(time.Duration(k) * ActivityRecordInterval) + score, last = NextActivityScore(score, last, sentAt), sentAt + } + require.Less(t, score.Sub(last).Hours()/24, 0.15) + + // A score is held to ActivityScoreMaxLead ahead of the send that set it. + require.Equal(t, start.Add(ActivityScoreMaxLead), NextActivityScore(start.Add(6*day+21*time.Hour), start.Add(-day), start)) + + // But a score never moves backwards, even when it is already further + // ahead of a new send than that. + prior := start.Add(8 * day) + require.Equal(t, prior, NextActivityScore(prior, start, start.Add(day))) +} + +func TestEffectiveActivityScore(t *testing.T) { + lastSentAt := time.Unix(1_700_000_000, 0).UTC() + + // No score: the last send. + require.Equal(t, lastSentAt, EffectiveActivityScore(time.Time{}, lastSentAt)) + + // A score trailing the last send is lifted to it. + require.Equal(t, lastSentAt, EffectiveActivityScore(lastSentAt.Add(-time.Hour), lastSentAt)) + + // A score at or ahead of it stands. + require.Equal(t, lastSentAt, EffectiveActivityScore(lastSentAt, lastSentAt)) + require.Equal(t, lastSentAt.Add(time.Hour), EffectiveActivityScore(lastSentAt.Add(time.Hour), lastSentAt)) +} + +func requireWithinMilli(t *testing.T, want, got time.Time) { + t.Helper() + diff := got.Sub(want) + require.True(t, diff >= -time.Millisecond && diff <= time.Millisecond, "got %v, want %v", got, want) +} diff --git a/chat/cache/store.go b/chat/cache/store.go index a4ed4d4..89462a4 100644 --- a/chat/cache/store.go +++ b/chat/cache/store.go @@ -27,7 +27,8 @@ import ( // is held. // // One thing is held that is not fixed at creation: a lower bound on each -// activity record, so a throttled send costs no write (see RecordSend). It can +// activity record, so a throttled send costs no store request (see +// RecordSend). It can // be out of date but never wrong in the direction that matters, because the // record only moves forward. type Cache struct { @@ -239,18 +240,19 @@ func (c *Cache) GetMutedCount(ctx context.Context, chatID *commonpb.ChatId) (uin return c.db.GetMutedCount(ctx, chatID) } -// RecordSend answers a throttled send itself, with no write: the backing -// store is billed for a conditional write whether or not its condition holds, -// so leaving the throttle to it would cost a write per message. Each send this +// RecordSend answers a throttled send itself, with no store request: the +// backing store reads the record to find a send throttled (see +// chat.Store.RecordSend), so leaving the throttle to it would cost a read per +// message. Each send this // process saw recorded is held for ActivityRecordInterval, keyed by (group, // user), and a send the interval has not yet cleared since it is answered // false, exactly as the store would answer it: the record only moves forward, // so it holds at least what was seen recorded here, whatever other processes // have written since. A send this process has not seen within the interval, // or one the store refused, goes to the store, so a stale entry can cost a -// write but never skip one. Each process throttles on its own, so a user's -// sends spread across several processes cost up to one write per process per -// interval. A send the store would reject outright (a DM ID, a time before the +// store request but never skip a record. Each process throttles on its own, +// so a user's sends spread across several processes cost up to one store +// request per process per interval, of which the store records one. A send the store would reject outright (a DM ID, a time before the // epoch) is never answered here. func (c *Cache) RecordSend(ctx context.Context, chatID *commonpb.ChatId, userID *commonpb.UserId, sentAt time.Time) (bool, error) { sentAtMillis := sentAt.UnixMilli() diff --git a/chat/dynamodb/activity_test.go b/chat/dynamodb/activity_test.go new file mode 100644 index 0000000..f014bd0 --- /dev/null +++ b/chat/dynamodb/activity_test.go @@ -0,0 +1,99 @@ +//go:build integration + +package dynamodb + +import ( + "context" + "testing" + "time" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/service/dynamodb" + "github.com/aws/aws-sdk-go-v2/service/dynamodb/types" + "github.com/stretchr/testify/require" + + commonpb "github.com/code-payments/flipcash2-protobuf-api/generated/go/common/v1" + + "github.com/code-payments/flipcash2-server/chat" + "github.com/code-payments/flipcash2-server/model" +) + +// TestChat_ActivityScoreLegacyRecords pins how records no Store method can +// write are read and scored: one written before scores existed, with no +// activity_score, and one whose score trails its last send, as a writer that +// records sends without scoring them leaves behind during a rollout. Both +// read and score as if their score were their last send. +func TestChat_ActivityScoreLegacyRecords(t *testing.T) { + ctx := context.Background() + require.NoError(t, CreateTables(ctx, testEnv.Client, chatsTable, dmInboxTable, groupMembersTable, userStateTable, activityTable, keyEnvelopesTable, lobbiesTable, featuredGroupsTable)) + + testStore := NewInDynamoDB(testEnv.Client, chatsTable, dmInboxTable, groupMembersTable, userStateTable, activityTable, keyEnvelopesTable, lobbiesTable, featuredGroupsTable, nil) + defer testStore.(*store).reset() + + groupID := chat.MustGenerateGroupChatID() + lastSentAt := time.Unix(1_700_000_000, 0).UTC() + + unscored := model.MustGenerateUserID() + trailing := model.MustGenerateUserID() + put := func(item map[string]types.AttributeValue) { + _, err := testEnv.Client.PutItem(ctx, &dynamodb.PutItemInput{TableName: aws.String(activityTable), Item: item}) + require.NoError(t, err) + } + put(map[string]types.AttributeValue{ + attrPK: avS(chatPK(groupID)), + attrSK: avS(userPK(unscored)), + attrLastSentAt: avInt(lastSentAt.UnixMilli()), + }) + put(map[string]types.AttributeValue{ + attrPK: avS(chatPK(groupID)), + attrSK: avS(userPK(trailing)), + attrLastSentAt: avInt(lastSentAt.UnixMilli()), + attrActivityScore: avInt(lastSentAt.Add(-time.Hour).UnixMilli()), + }) + + scores := func() map[string]time.Time { + t.Helper() + senders, err := testStore.GetRecentSenders(ctx, groupID, 0) + require.NoError(t, err) + out := make(map[string]time.Time, len(senders)) + for _, sender := range senders { + out[string(sender.UserID.Value)] = sender.ActivityScore + } + return out + } + + // Read as scoring their last send. + got := scores() + require.Len(t, got, 2) + require.True(t, got[string(unscored.Value)].Equal(lastSentAt)) + require.True(t, got[string(trailing.Value)].Equal(lastSentAt)) + + // The throttle still holds against them. + recorded, err := testStore.RecordSend(ctx, groupID, unscored, lastSentAt.Add(chat.ActivityRecordInterval/2)) + require.NoError(t, err) + require.False(t, recorded) + + // And their next send is scored against their last. + sentAt := lastSentAt.Add(time.Hour) + want := chat.NextActivityScore(lastSentAt, lastSentAt, sentAt) + for _, user := range []*commonpb.UserId{unscored, trailing} { + recorded, err := testStore.RecordSend(ctx, groupID, user, sentAt) + require.NoError(t, err) + require.True(t, recorded) + } + got = scores() + require.True(t, got[string(unscored.Value)].Equal(want), "got %v, want %v", got[string(unscored.Value)], want) + require.True(t, got[string(trailing.Value)].Equal(want), "got %v, want %v", got[string(trailing.Value)], want) + + // Now in the score index too. + out, err := testEnv.Client.Query(ctx, &dynamodb.QueryInput{ + TableName: aws.String(activityTable), + IndexName: aws.String(lsiByActivityScore), + KeyConditionExpression: aws.String("#pk = :pk"), + ExpressionAttributeNames: map[string]string{"#pk": attrPK}, + ExpressionAttributeValues: map[string]types.AttributeValue{":pk": avS(chatPK(groupID))}, + ConsistentRead: aws.Bool(true), + }) + require.NoError(t, err) + require.Len(t, out.Items, 2) +} diff --git a/chat/dynamodb/store.go b/chat/dynamodb/store.go index 7bd8229..f0e04d1 100644 --- a/chat/dynamodb/store.go +++ b/chat/dynamodb/store.go @@ -113,26 +113,24 @@ import ( // chat_activity pk = "chat#", sk = "user#" (one item per (group, // user) the user has sent in; see chat.RecentSender). A group's // activity records, independent of membership: last_sent_at -// (epoch ms) and expires_at (DynamoDB TTL, epoch seconds, -// ActivityRetention after it). Keyed by (chat, user) so a send -// is a blind conditional update of one known item — no read -// first, no duplicate rows — while lsiByLastSentAt orders the -// partition by recency, which is the read: a group's recent -// senders are one eventually consistent query, billed by the -// page. +// (epoch ms), activity_score (epoch ms; see +// chat.NextActivityScore) and expires_at (DynamoDB TTL, epoch +// seconds, ActivityRetention after the last send). Keyed by +// (chat, user) so a send is a read and a conditional update of +// one known item, with no duplicate rows, while lsiByLastSentAt +// orders the partition by recency, which is the read: a group's +// recent senders are one eventually consistent query, billed by +// the page. // -// lsiByActivityScore orders the same partition by activity_score, -// a frequency-weighted ordering that nothing writes yet. It exists -// now only because an LSI cannot be added to a table later; while -// no item carries activity_score it is empty and costs nothing. -// The score is meant to be epoch-ms-denominated, equal to -// last_sent_at for a user with one recorded send and running ahead -// of it (possibly past now) as their sends accumulate, so rows -// written before it existed can be given activity_score = -// last_sent_at and compare correctly with the rest. The write that -// maintains it will need the item's prior score, so it will read -// first and condition its update on last_sent_at as read, which -// every recorded send moves forward. +// lsiByActivityScore orders the same partition by +// activity_score, a frequency-weighted ordering that every +// recorded send maintains (see RecordSend) and nothing reads yet. +// The score is epoch-ms-denominated, equal to last_sent_at for a +// user with one recorded send and running ahead of it (possibly +// past now) as their sends accumulate, so a row written before +// scores existed, which is missing from the index, can be given +// activity_score = last_sent_at and compare correctly with the +// rest; until it is, it reads and scores as if it had been. // // chat_key_envelopes pk = "user#", sk = "chat#" (one item per // (user, private group) the user holds a key envelope for; see @@ -250,8 +248,9 @@ const ( // group's activity records in recency order (see GetRecentSenders). lsiByLastSentAt = "by_last_sent_at" - // lsiByActivityScore is the (chat, activity_score) LSI on chat_activity, - // reserved for a frequency-weighted ordering; nothing writes its key yet. + // lsiByActivityScore is the (chat, activity_score) LSI on chat_activity: + // a group's activity records in activity-score order. Nothing reads it + // yet, and a record written before scores existed is missing from it. lsiByActivityScore = "by_activity_score" // gsiLobbyByChat is the (sk, entered_at) index on chat_lobbies: a group's @@ -298,7 +297,7 @@ const ( attrMutedUntil = "muted_until" // chat_user_state: epoch seconds, present only while a mute is recorded — see muteForeverUntil attrMutedCount = "muted_count" // chat_user_state #meta item: records with a mute recorded attrLastSentAt = "last_sent_at" // chat_activity: epoch ms of the latest recorded send - attrActivityScore = "activity_score" // chat_activity: reserved, see lsiByActivityScore + attrActivityScore = "activity_score" // chat_activity: epoch ms, see chat.NextActivityScore; absent on records written before scores attrScheme = "scheme" // chat_key_envelopes: chatpb.KeyEnvelope_Scheme, by number attrNonce = "nonce" // chat_key_envelopes: the envelope's nonce (B) attrCiphertext = "ciphertext" // chat_key_envelopes: the wrapped chat key (B) @@ -2727,11 +2726,29 @@ func viewerStateFromItem(item map[string]types.AttributeValue) (chat.ViewerState return state, nil } -// RecordSend is one conditional update of the user's chat_activity item, -// with no read first: the condition is the throttle, and its failure is the -// not-recorded answer, so the no-op path costs no second request and never -// creates an item. last_sent_at is epoch milliseconds, which keeps the -// throttle's arithmetic and the index's order exact. +// RecordSend reads the user's chat_activity item and, unless the read already +// shows the send throttled, writes the next last_sent_at and activity_score +// in one update conditioned on the last_sent_at it read (or on there being +// none). The score needs the prior one, which an update expression cannot +// compute, hence the read; the condition makes the read-compute-write atomic +// across writers, because every recorded send moves last_sent_at forward by +// at least ActivityRecordInterval and nothing moves it back, so a value that +// still matches proves nothing was recorded since the read. +// +// A writer that loses checks its send once more against the item the failure +// returned, with no second read. Usually the winner's send throttles it and +// it is done; when it is at least ActivityRecordInterval later than the +// winner's, it is the newest send and is written once more, scored against +// the winner's. Should that write lose too, another send has landed since, +// and the send goes unrecorded, as a throttled one does: a record is best +// effort, and its user's next send repairs it. +// +// The read is eventually consistent, at half a strong read's cost: a stale +// one only shows an older send, so a throttled answer from it is right, and a +// write from it fails its condition and is checked again against the current +// item. +// last_sent_at and activity_score are epoch milliseconds, which keeps the +// throttle's arithmetic and both indexes' orders exact. func (s *store) RecordSend(ctx context.Context, chatID *commonpb.ChatId, userID *commonpb.UserId, sentAt time.Time) (bool, error) { if !chat.IsGroupChatID(chatID) { return false, fmt.Errorf("not a group chat id") @@ -2740,38 +2757,97 @@ func (s *store) RecordSend(ctx context.Context, chatID *commonpb.ChatId, userID if sentAtMillis <= 0 { return false, fmt.Errorf("send time %v is not after the epoch", sentAt) } + sentAt = time.UnixMilli(sentAtMillis).UTC() + key := map[string]types.AttributeValue{ + attrPK: avS(chatPK(chatID)), + attrSK: avS(userPK(userID)), + } - _, err := s.client.UpdateItem(ctx, &dynamodb.UpdateItemInput{ - TableName: aws.String(s.activityTable), - Key: map[string]types.AttributeValue{ - attrPK: avS(chatPK(chatID)), - attrSK: avS(userPK(userID)), - }, - UpdateExpression: aws.String("SET #sent = :sent, #expires = :expires"), - ConditionExpression: aws.String("attribute_not_exists(#sent) OR #sent <= :stale"), - ExpressionAttributeNames: map[string]string{ - "#sent": attrLastSentAt, - "#expires": attrExpiresAt, - }, - ExpressionAttributeValues: map[string]types.AttributeValue{ + out, err := s.client.GetItem(ctx, &dynamodb.GetItemInput{ + TableName: aws.String(s.activityTable), + Key: key, + ProjectionExpression: aws.String("#sent, #score"), + ExpressionAttributeNames: map[string]string{"#sent": attrLastSentAt, "#score": attrActivityScore}, + }) + if err != nil { + return false, err + } + item := out.Item + + for range maxRecordSendAttempts { + var prior, lastSentAt time.Time + values := map[string]types.AttributeValue{ ":sent": avInt(sentAtMillis), - ":stale": avInt(sentAtMillis - chat.ActivityRecordInterval.Milliseconds()), ":expires": avInt(sentAt.Add(chat.ActivityRetention).Unix()), - }, - }) - if isConditionalCheckFailed(err) { - return false, nil + } + condition := "attribute_not_exists(#sent)" + if av, ok := item[attrLastSentAt]; ok { + lastSentMillis, err := parseInt(av) + if err != nil { + return false, fmt.Errorf("parsing %s: %w", attrLastSentAt, err) + } + if lastSentMillis > sentAtMillis-chat.ActivityRecordInterval.Milliseconds() { + return false, nil + } + score, err := activityScoreFromItem(item) + if err != nil { + return false, err + } + lastSentAt = time.UnixMilli(lastSentMillis).UTC() + prior = chat.EffectiveActivityScore(score, lastSentAt) + condition = "#sent = :read" + values[":read"] = av + } + values[":score"] = avInt(chat.NextActivityScore(prior, lastSentAt, sentAt).UnixMilli()) + + _, err := s.client.UpdateItem(ctx, &dynamodb.UpdateItemInput{ + TableName: aws.String(s.activityTable), + Key: key, + UpdateExpression: aws.String("SET #sent = :sent, #score = :score, #expires = :expires"), + ConditionExpression: aws.String(condition), + ExpressionAttributeNames: map[string]string{ + "#sent": attrLastSentAt, + "#score": attrActivityScore, + "#expires": attrExpiresAt, + }, + ExpressionAttributeValues: values, + ReturnValuesOnConditionCheckFailure: types.ReturnValuesOnConditionCheckFailureAllOld, + }) + if err == nil { + return true, nil + } + var ccf *types.ConditionalCheckFailedException + if !errors.As(err, &ccf) { + return false, err + } + item = ccf.Item } + return false, nil +} + +// maxRecordSendAttempts bounds RecordSend's writes: the first, and one more +// after a lost race (see RecordSend). +const maxRecordSendAttempts = 2 + +// activityScoreFromItem decodes a chat_activity item's activity_score, the +// zero time when it has none (a record written before scores existed). +func activityScoreFromItem(item map[string]types.AttributeValue) (time.Time, error) { + av, ok := item[attrActivityScore] + if !ok { + return time.Time{}, nil + } + scoreMillis, err := parseInt(av) if err != nil { - return false, err + return time.Time{}, fmt.Errorf("parsing %s: %w", attrActivityScore, err) } - return true, nil + return time.UnixMilli(scoreMillis).UTC(), nil } // GetRecentSenders queries lsiByLastSentAt descending, eventually // consistent at half the cost of a strong read (the index is local, so a // strong one is available if a reader ever needs it). It projects -// last_sent_at as its own key, so nothing is fetched from the table. Every item in the partition is +// last_sent_at as its own key and activity_score as an included attribute, +// so nothing is fetched from the table. Every item in the partition is // an activity record carrying last_sent_at, so the index holds them all. func (s *store) GetRecentSenders(ctx context.Context, chatID *commonpb.ChatId, limit int) ([]chat.RecentSender, error) { if !chat.IsGroupChatID(chatID) { @@ -2785,8 +2861,8 @@ func (s *store) GetRecentSenders(ctx context.Context, chatID *commonpb.ChatId, l TableName: aws.String(s.activityTable), IndexName: aws.String(lsiByLastSentAt), KeyConditionExpression: aws.String("#pk = :pk"), - ProjectionExpression: aws.String("#sk, #sent"), - ExpressionAttributeNames: map[string]string{"#pk": attrPK, "#sk": attrSK, "#sent": attrLastSentAt}, + ProjectionExpression: aws.String("#sk, #sent, #score"), + ExpressionAttributeNames: map[string]string{"#pk": attrPK, "#sk": attrSK, "#sent": attrLastSentAt, "#score": attrActivityScore}, ExpressionAttributeValues: map[string]types.AttributeValue{ ":pk": avS(chatPK(chatID)), }, @@ -2809,9 +2885,15 @@ func (s *store) GetRecentSenders(ctx context.Context, chatID *commonpb.ChatId, l if err != nil { return nil, fmt.Errorf("parsing %s: %w", attrLastSentAt, err) } + score, err := activityScoreFromItem(item) + if err != nil { + return nil, err + } + lastSentAt := time.UnixMilli(sentAtMillis).UTC() senders = append(senders, chat.RecentSender{ - UserID: userID, - LastSentAt: time.UnixMilli(sentAtMillis).UTC(), + UserID: userID, + LastSentAt: lastSentAt, + ActivityScore: chat.EffectiveActivityScore(score, lastSentAt), }) } if limit > 0 && len(senders) >= limit { diff --git a/chat/dynamodb/table.go b/chat/dynamodb/table.go index fd47148..e1c7ad4 100644 --- a/chat/dynamodb/table.go +++ b/chat/dynamodb/table.go @@ -24,7 +24,7 @@ import ( // when they end (see gsiByMuted) and an inverted GSI of a chat's records by // user (see gsiUserStateByUser); chat_activity is keyed by (pk, sk) = (chat, // user) with two LSIs, one by last_sent_at (lsiByLastSentAt) and one by -// activity_score (lsiByActivityScore, reserved and empty today), and TTL on +// activity_score (lsiByActivityScore, unread today), and TTL on // expires_at; chat_key_envelopes is keyed by (pk, sk) = (user, chat) with no // index; chat_lobbies is keyed by (pk, sk) = (user, chat) — plus one "#meta" // aggregates item per chat and per user — with a sparse GSI of a chat's diff --git a/chat/memory/store.go b/chat/memory/store.go index 0a0479f..dab156b 100644 --- a/chat/memory/store.go +++ b/chat/memory/store.go @@ -46,11 +46,10 @@ type memory struct { // one. excludedFromFeed map[string]map[string]struct{} - // lastSent is each group's activity records, keyed by chat ID then user - // ID: the latest recorded send, at the persistent store's millisecond - // precision (see chat.Store.RecordSend). Records never expire here; the + // activity is each group's activity records, keyed by chat ID then user + // ID (see chat.Store.RecordSend). Records never expire here; the // contract lets a reader see one past chat.ActivityRetention. - lastSent map[string]map[string]time.Time + activity map[string]map[string]activityRecord // keyEnvelopes holds each user's key envelope per chat, keyed by user ID // then chat ID, mirroring the persistent layout (see chat.Store). @@ -91,7 +90,7 @@ func NewInMemory(excludedFromFeed []*commonpb.UserId) chat.Store { groupVersions: make(map[string]uint64), viewerStates: make(map[string]map[string]*chat.ViewerState), mutedCounts: make(map[string]uint64), - lastSent: make(map[string]map[string]time.Time), + activity: make(map[string]map[string]activityRecord), keyEnvelopes: make(map[string]map[string]chat.KeyEnvelope), lobbies: make(map[string]map[string]time.Time), featured: make(map[string]chat.FeaturedGroups), @@ -108,7 +107,7 @@ func (m *memory) reset() { m.viewerStates = make(map[string]map[string]*chat.ViewerState) m.mutedCounts = make(map[string]uint64) m.excludedFromFeed = make(map[string]map[string]struct{}) - m.lastSent = make(map[string]map[string]time.Time) + m.activity = make(map[string]map[string]activityRecord) m.keyEnvelopes = make(map[string]map[string]chat.KeyEnvelope) m.lobbies = make(map[string]map[string]time.Time) m.featured = make(map[string]chat.FeaturedGroups) @@ -820,6 +819,14 @@ func (m *memory) GetMutedCount(_ context.Context, chatID *commonpb.ChatId) (uint return m.mutedCounts[string(chatID.Value)], nil } +// activityRecord is one user's activity record in a group, at the persistent +// store's millisecond precision. Every record here was written with its score, +// so score is the record's effective score as it stands. +type activityRecord struct { + lastSentAt time.Time + score time.Time +} + func (m *memory) RecordSend(_ context.Context, chatID *commonpb.ChatId, userID *commonpb.UserId, sentAt time.Time) (bool, error) { if !chat.IsGroupChatID(chatID) { return false, fmt.Errorf("not a group chat id") @@ -833,13 +840,17 @@ func (m *memory) RecordSend(_ context.Context, chatID *commonpb.ChatId, userID * defer m.Unlock() chatKey := string(chatID.Value) - if last, ok := m.lastSent[chatKey][string(userID.Value)]; ok && last.After(sentAt.Add(-chat.ActivityRecordInterval)) { + prior, ok := m.activity[chatKey][string(userID.Value)] + if ok && prior.lastSentAt.After(sentAt.Add(-chat.ActivityRecordInterval)) { return false, nil } - if m.lastSent[chatKey] == nil { - m.lastSent[chatKey] = make(map[string]time.Time) + if m.activity[chatKey] == nil { + m.activity[chatKey] = make(map[string]activityRecord) + } + m.activity[chatKey][string(userID.Value)] = activityRecord{ + lastSentAt: sentAt, + score: chat.NextActivityScore(prior.score, prior.lastSentAt, sentAt), } - m.lastSent[chatKey][string(userID.Value)] = sentAt return true, nil } @@ -851,11 +862,12 @@ func (m *memory) GetRecentSenders(_ context.Context, chatID *commonpb.ChatId, li m.Lock() defer m.Unlock() - senders := make([]chat.RecentSender, 0, len(m.lastSent[string(chatID.Value)])) - for user, last := range m.lastSent[string(chatID.Value)] { + senders := make([]chat.RecentSender, 0, len(m.activity[string(chatID.Value)])) + for user, record := range m.activity[string(chatID.Value)] { senders = append(senders, chat.RecentSender{ - UserID: &commonpb.UserId{Value: []byte(user)}, - LastSentAt: last, + UserID: &commonpb.UserId{Value: []byte(user)}, + LastSentAt: record.lastSentAt, + ActivityScore: record.score, }) } // Ties are in no particular order by contract; break them by user so @@ -880,8 +892,8 @@ func (m *memory) GetLastSentAt(_ context.Context, chatID *commonpb.ChatId, userI m.Lock() defer m.Unlock() - last, ok := m.lastSent[string(chatID.Value)][string(userID.Value)] - return last, ok, nil + record, ok := m.activity[string(chatID.Value)][string(userID.Value)] + return record.lastSentAt, ok, nil } func (m *memory) SetKeyEnvelope(_ context.Context, chatID *commonpb.ChatId, userID *commonpb.UserId, envelope chat.KeyEnvelope) (chat.KeyEnvelope, error) { diff --git a/chat/model.go b/chat/model.go index fe373b1..7e05e1e 100644 --- a/chat/model.go +++ b/chat/model.go @@ -840,10 +840,10 @@ func (v ViewerState) Clone() ViewerState { // send; expiry is garbage collection, not semantics, so a reader may still // see a record past it. // -// The record is meant to grow a frequency-weighted ordering beside recency -// (an activity score, see the DynamoDB store's chat_activity table), which is -// why it is an activity record and not a send time alone; nothing computes or -// reads one yet. +// Beside recency, each record carries an activity score, a frequency-weighted +// ordering of the same sends (see NextActivityScore), which is why it is an +// activity record and not a send time alone. The score is maintained with +// every recorded send; nothing ranks by it yet. const ( // ActivityRecordInterval is the least time between two recorded sends by // one user in one group. @@ -859,10 +859,12 @@ const ( ) // RecentSender is one user's activity record in a group, as recency reads -// it: who, and when their latest recorded send was, at millisecond precision. +// it: who, when their latest recorded send was, and the record's activity +// score (see EffectiveActivityScore), both at millisecond precision. type RecentSender struct { - UserID *commonpb.UserId - LastSentAt time.Time + UserID *commonpb.UserId + LastSentAt time.Time + ActivityScore time.Time } // A private group's lobby (see Chat.IsPrivate) is where a user waits to be diff --git a/chat/store.go b/chat/store.go index 6ecc0cb..14ad1ff 100644 --- a/chat/store.go +++ b/chat/store.go @@ -481,27 +481,30 @@ type Store interface { GetMutedCount(ctx context.Context, chatID *commonpb.ChatId) (uint64, error) // RecordSend records that userID sent in the group chatID at sentAt (see - // the activity record, RecentSender), and reports whether it did. It - // does not when the record already holds a send later than - // sentAt − ActivityRecordInterval, which is what throttles a burst and - // keeps a delayed or retried send from moving the record backwards, so - // any number of retries is safe. It records against the chat ID alone, - // reading neither the canonical record nor the roster: the caller is the - // send path, which has already gated both. A caching decorator may - // answer a throttled send itself, with no write (see cache.Cache). It - // returns an error if chatID is not a group chat ID, or if sentAt is not - // after the Unix epoch. + // the activity record, RecentSender), moving its activity score with it + // (see NextActivityScore), and reports whether it did. Each recorded send + // moves the score exactly once, however many writers record sends for the + // same user at once. It does not record when the record already holds a + // send later than sentAt − ActivityRecordInterval, which is what + // throttles a burst and keeps a delayed or retried send from moving the + // record backwards, so any number of retries is safe. It records against + // the chat ID alone, reading neither the canonical record nor the roster: + // the caller is the send path, which has already gated both. A caching + // decorator may answer a throttled send itself, with no store request + // (see cache.Cache). It returns an error if chatID is not a group chat + // ID, or if sentAt is not after the Unix epoch. RecordSend(ctx context.Context, chatID *commonpb.ChatId, userID *commonpb.UserId, sentAt time.Time) (recorded bool, err error) - // GetRecentSenders returns the group's activity records most recently - // sent first, at most limit of them (limit <= 0 means unbounded), from an - // eventually consistent read, so a send recorded a moment ago may be - // missing or at its previous time: its readers rank and show people, which - // a refetch corrects. Ties come back in no particular order. Records of - // users who have since left are included (see the activity record), and - // a record past ActivityRetention may be. A group with no records, or - // that does not exist, is an empty result. It returns an error if chatID - // is not a group chat ID. + // GetRecentSenders returns the group's activity records, with their + // activity scores, most recently sent first, at most limit of them + // (limit <= 0 means unbounded), from an eventually consistent read, so a + // send recorded a moment ago may be missing or at its previous time: its + // readers rank and show people, which a refetch corrects. Ties come back + // in no particular order. Records of users who have since left are + // included (see the activity record), and a record past + // ActivityRetention may be. A group with no records, or that does not + // exist, is an empty result. It returns an error if chatID is not a group + // chat ID. GetRecentSenders(ctx context.Context, chatID *commonpb.ChatId, limit int) ([]RecentSender, error) // GetLastSentAt returns userID's activity record in the group chatID: diff --git a/chat/tests/store.go b/chat/tests/store.go index 567e080..01cda06 100644 --- a/chat/tests/store.go +++ b/chat/tests/store.go @@ -88,6 +88,8 @@ func RunStoreTests(t *testing.T, s chat.Store, newStore func(excludedFromFeed [] testStore_Activity_Precision, testStore_Activity_Scope, testStore_Activity_Concurrent, + testStore_Activity_Score, + testStore_Activity_ScoreConcurrent, testStore_KeyEnvelope_SetAndGet, testStore_KeyEnvelope_OwnWrapStands, testStore_KeyEnvelope_DiscardedOnLeave, @@ -2791,6 +2793,133 @@ func testStore_Activity_Concurrent(t *testing.T, s chat.Store) { senders, err := s.GetRecentSenders(ctx, groupID, 0) require.NoError(t, err) requireRecentSenders(t, senders, user, at(70)) + require.True(t, senders[0].ActivityScore.Equal(at(70)), "scored once") +} + +func testStore_Activity_Score(t *testing.T, s chat.Store) { + ctx := context.Background() + + groupID := chat.MustGenerateGroupChatID() + user := model.MustGenerateUserID() + interval := chat.ActivityRecordInterval + + score := func() time.Time { + t.Helper() + senders, err := s.GetRecentSenders(ctx, groupID, 0) + require.NoError(t, err) + require.Len(t, senders, 1) + return senders[0].ActivityScore + } + + // A first send scores its own time. + recorded, err := s.RecordSend(ctx, groupID, user, at(100)) + require.NoError(t, err) + require.True(t, recorded) + require.True(t, score().Equal(at(100))) + + // Each recorded send moves the score once, by the formula, against the + // send before it. + want, last := at(100), at(100) + for _, sentAt := range []time.Time{ + at(100).Add(interval), + at(100).Add(3 * interval), + at(100).Add(24 * time.Hour), + } { + recorded, err = s.RecordSend(ctx, groupID, user, sentAt) + require.NoError(t, err) + require.True(t, recorded) + want, last = chat.NextActivityScore(want, last, sentAt), sentAt + require.True(t, score().Equal(want), "after send at %v: got %v, want %v", sentAt, score(), want) + } + require.True(t, want.After(at(100).Add(24*time.Hour)), "the score runs ahead of the last send") + + // A throttled send moves nothing. + recorded, err = s.RecordSend(ctx, groupID, user, at(100).Add(24*time.Hour).Add(interval/2)) + require.NoError(t, err) + require.False(t, recorded) + require.True(t, score().Equal(want)) + + // A frequent sender outranks a one-off sender who sent after them. + other := model.MustGenerateUserID() + recorded, err = s.RecordSend(ctx, groupID, other, at(100).Add(25*time.Hour)) + require.NoError(t, err) + require.True(t, recorded) + senders, err := s.GetRecentSenders(ctx, groupID, 0) + require.NoError(t, err) + requireRecentSenders(t, senders, other, at(100).Add(25*time.Hour), user, at(100).Add(24*time.Hour)) + require.True(t, senders[1].ActivityScore.After(senders[0].ActivityScore)) +} + +func testStore_Activity_ScoreConcurrent(t *testing.T, s chat.Store) { + ctx := context.Background() + + groupID := chat.MustGenerateGroupChatID() + user := model.MustGenerateUserID() + + // Sends far enough apart that none throttles another in order, recorded + // all at once by many writers, as several servers would. Whichever are + // recorded must have been recorded in time order, each scored once + // against the one before, so the score is exactly theirs: no lost update, + // no double count. A store may leave a send that loses its race twice + // unrecorded, so which ones are recorded is not fixed. + const writers = 4 + sendTimes := make([]time.Time, writers) + for i := range writers { + sendTimes[i] = at(1000).Add(time.Duration(i) * time.Hour) + } + var wg sync.WaitGroup + results := make([]bool, writers) + errs := make([]error, writers) + for i := range writers { + wg.Add(1) + go func() { + defer wg.Done() + results[i], errs[i] = s.RecordSend(ctx, groupID, user, sendTimes[i]) + }() + } + wg.Wait() + + var want, last time.Time + for i := range writers { + require.NoError(t, errs[i]) + if results[i] { + want = chat.NextActivityScore(want, last, sendTimes[i]) + last = sendTimes[i] + } + } + require.False(t, last.IsZero()) + senders, err := s.GetRecentSenders(ctx, groupID, 0) + require.NoError(t, err) + requireRecentSenders(t, senders, user, last) + require.True(t, senders[0].ActivityScore.Equal(want), "got %v, want %v", senders[0].ActivityScore, want) + + // Of two such writers, the later send is always recorded, whichever + // writes first: it loses at most once, to the earlier. (Against a store + // whose first read may be stale it could lose once to that as well; the + // stores under test read current data.) + for range 5 { + groupID := chat.MustGenerateGroupChatID() + earlier, later := at(2000), at(2000).Add(time.Hour) + var wg sync.WaitGroup + var errEarlier, errLater error + var recordedLater bool + wg.Add(2) + go func() { + defer wg.Done() + _, errEarlier = s.RecordSend(ctx, groupID, user, earlier) + }() + go func() { + defer wg.Done() + recordedLater, errLater = s.RecordSend(ctx, groupID, user, later) + }() + wg.Wait() + require.NoError(t, errEarlier) + require.NoError(t, errLater) + require.True(t, recordedLater) + senders, err := s.GetRecentSenders(ctx, groupID, 0) + require.NoError(t, err) + requireRecentSenders(t, senders, user, later) + } } // requireRecentSenders asserts senders is exactly the given (user, last sent)