Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions CLAUDE.md

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion chat/activity.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ import (
// ActivityScoreMaxLead bounds it whatever the history; at these values no
// steady pattern reaches it.
//
// Nothing ranks by it yet.
// A group's chatter sample is ordered by it (see Store.GetActiveSenders).
const (
// ActivityScoreHalfLife is how long a recorded send takes to count half
// as much toward an activity score.
Expand Down
4 changes: 4 additions & 0 deletions chat/cache/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -282,6 +282,10 @@ func (c *Cache) GetRecentSenders(ctx context.Context, chatID *commonpb.ChatId, l
return c.db.GetRecentSenders(ctx, chatID, limit)
}

func (c *Cache) GetActiveSenders(ctx context.Context, chatID *commonpb.ChatId, limit int) ([]chat.RecentSender, error) {
return c.db.GetActiveSenders(ctx, chatID, limit)
}

func (c *Cache) GetLastSentAt(ctx context.Context, chatID *commonpb.ChatId, userID *commonpb.UserId) (time.Time, bool, error) {
return c.db.GetLastSentAt(ctx, chatID, userID)
}
Expand Down
46 changes: 32 additions & 14 deletions chat/dynamodb/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -124,13 +124,15 @@ import (
//
// 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.
// recorded send maintains (see RecordSend), read by
// GetActiveSenders. 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, was given activity_score = last_sent_at by a one-off
// backfill and compares correctly with the rest; one that
// somehow still lacks it reads and scores as if it had it, but
// GetActiveSenders does not find it.
//
// chat_key_envelopes pk = "user#<id>", sk = "chat#<id>" (one item per
// (user, private group) the user holds a key envelope for; see
Expand Down Expand Up @@ -249,8 +251,9 @@ const (
lsiByLastSentAt = "by_last_sent_at"

// 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.
// a group's activity records in activity-score order (see
// GetActiveSenders). It is sparse on activity_score, which records
// written before scores existed lacked until backfilled.
lsiByActivityScore = "by_activity_score"

// gsiLobbyByChat is the (sk, entered_at) index on chat_lobbies: a group's
Expand Down Expand Up @@ -2845,11 +2848,26 @@ func activityScoreFromItem(item map[string]types.AttributeValue) (time.Time, err

// 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 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.
// strong one is available if a reader ever needs it). 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) {
return s.querySenders(ctx, chatID, lsiByLastSentAt, limit)
}

// GetActiveSenders queries lsiByActivityScore descending, eventually
// consistent like GetRecentSenders. The index is sparse on activity_score, so
// a record written before scores existed and never backfilled is missing
// from it.
func (s *store) GetActiveSenders(ctx context.Context, chatID *commonpb.ChatId, limit int) ([]chat.RecentSender, error) {
return s.querySenders(ctx, chatID, lsiByActivityScore, limit)
}

// querySenders reads a group's activity records from one of the two LSIs on
// chat_activity, in descending order of its sort key. Each index projects the
// other's sort key, so a record's send time and score both come from the
// index and nothing is fetched from the table.
func (s *store) querySenders(ctx context.Context, chatID *commonpb.ChatId, index string, limit int) ([]chat.RecentSender, error) {
if !chat.IsGroupChatID(chatID) {
return nil, fmt.Errorf("not a group chat id")
}
Expand All @@ -2859,7 +2877,7 @@ func (s *store) GetRecentSenders(ctx context.Context, chatID *commonpb.ChatId, l
for {
input := &dynamodb.QueryInput{
TableName: aws.String(s.activityTable),
IndexName: aws.String(lsiByLastSentAt),
IndexName: aws.String(index),
KeyConditionExpression: aws.String("#pk = :pk"),
ProjectionExpression: aws.String("#sk, #sent, #score"),
ExpressionAttributeNames: map[string]string{"#pk": attrPK, "#sk": attrSK, "#sent": attrLastSentAt, "#score": attrActivityScore},
Expand Down
2 changes: 1 addition & 1 deletion chat/dynamodb/table.go
Original file line number Diff line number Diff line change
Expand Up @@ -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, unread today), and TTL on
// activity_score (lsiByActivityScore), 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
Expand Down
30 changes: 30 additions & 0 deletions chat/memory/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -884,6 +884,36 @@ func (m *memory) GetRecentSenders(_ context.Context, chatID *commonpb.ChatId, li
return senders, nil
}

func (m *memory) GetActiveSenders(_ context.Context, chatID *commonpb.ChatId, limit int) ([]chat.RecentSender, error) {
if !chat.IsGroupChatID(chatID) {
return nil, fmt.Errorf("not a group chat id")
}

m.Lock()
defer m.Unlock()

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: record.lastSentAt,
ActivityScore: record.score,
})
}
// Ties are in no particular order by contract; break them by user so
// this store is at least deterministic.
sort.Slice(senders, func(i, j int) bool {
if !senders[i].ActivityScore.Equal(senders[j].ActivityScore) {
return senders[i].ActivityScore.After(senders[j].ActivityScore)
}
return bytes.Compare(senders[i].UserID.Value, senders[j].UserID.Value) > 0
})
if limit > 0 && len(senders) > limit {
senders = senders[:limit]
}
return senders, nil
}

func (m *memory) GetLastSentAt(_ context.Context, chatID *commonpb.ChatId, userID *commonpb.UserId) (time.Time, bool, error) {
if !chat.IsGroupChatID(chatID) {
return time.Time{}, false, fmt.Errorf("not a group chat id")
Expand Down
11 changes: 7 additions & 4 deletions chat/model.go
Original file line number Diff line number Diff line change
Expand Up @@ -843,7 +843,8 @@ func (v ViewerState) Clone() ViewerState {
// 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.
// every recorded send, and orders a group's chatter sample (see
// Store.GetActiveSenders).
const (
// ActivityRecordInterval is the least time between two recorded sends by
// one user in one group.
Expand All @@ -858,9 +859,11 @@ const (
ActivityRetention = 365 * 24 * time.Hour
)

// RecentSender is one user's activity record in a group, as recency reads
// it: who, when their latest recorded send was, and the record's activity
// score (see EffectiveActivityScore), both at millisecond precision.
// RecentSender is one user's activity record in a group, as the reads of a
// group's senders return it (see Store.GetRecentSenders and
// Store.GetActiveSenders): 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
Expand Down
58 changes: 45 additions & 13 deletions chat/sample.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"bytes"
"context"
"errors"
"sort"
"time"

"go.uber.org/zap"
Expand All @@ -20,19 +21,23 @@ import (
)

// A sample of a public group's chatters: a short list of its members to show,
// the creator first while they are a member, then the members who have sent
// most recently, most recent first. It is how a group shows who is in it
// without listing members who only read: everyone in it is a member as of the
// read, but a member who has not sent recently is not in it.
// the creator first while they are a member, then the most active members who
// have sent, by activity score (see NextActivityScore) leaning toward
// recency (see sampleRank): a member who shows up every day ranks above one
// who sent once a little more recently, but of two with close scores the
// more recent sender comes first. It is how a group
// shows who is in it without listing members who only read: everyone in it
// is a member as of the read, but a member who has not sent recently is not
// in it.
//
// The candidates are the group's recent senders, read from its activity
// records (see RecentSender), which name people who have since left as well
// as members. So each candidate's membership is checked (see
// The candidates are the group's most active senders, read from its activity
// records (see Store.GetActiveSenders), which name people who have since left
// as well as members. So each candidate's membership is checked (see
// Store.GetGroupMembersByID), in chunks in candidate order, stopping as
// soon as one more member than the sample holds is found: that one proves
// has_more without being returned. Each chunk is as many candidates as
// members still to find, plus sampleChattersCheckBuffer, so a group whose
// recent senders are all still members is answered by one read, and a
// most active senders are all still members is answered by one read, and a
// later read checks only about as many as are still missing. The read of
// senders is bounded too; when it fills and too few of the senders are
// still members, has_more is set,
Expand All @@ -47,8 +52,8 @@ import (
// (see ActivityRetention) is not a candidate until they send again.
//
// The creator is a candidate whether or not they have sent, and is shown at
// their own last send whenever they have a record, however many others have
// sent since: one outside the senders read is looked up directly (see
// their own last send whenever they have a record, however many others rank
// above them: one outside the senders read is looked up directly (see
// Store.GetLastSentAt), only when they are shown and the read of senders
// filled, since otherwise it holds every record the group has.
//
Expand All @@ -69,16 +74,35 @@ const (
// the max_items on SampleChattersResponse.chatters.
sampleChattersSize = 20

// sampleChattersSenderWindow is how many of a group's most recent
// sampleChattersSenderWindow is how many of a group's most active
// senders are candidates for the sample.
sampleChattersSenderWindow = 100

// sampleChattersCheckBuffer is how many candidates beyond the members
// still to find are checked per read, to absorb a few who have left
// without another read.
sampleChattersCheckBuffer = 5

// sampleChattersRecencyWeight is how far a sender's rank leans from
// their activity score toward their last send (see sampleRank): 0 ranks
// by score alone, 1 by recency alone.
sampleChattersRecencyWeight = 0.2
)

// sampleRank is the key a sender is ranked by in a sample, highest first:
// their activity score less sampleChattersRecencyWeight of its lead over
// their last send. Comparing two senders, the more recent one ranks first
// exactly when the gap between their last sends is more than
// (1 − w) / w times the gap between their scores (4 times at w = 0.2): close
// scores go to the more recent sender, distant ones to the more active. It
// only reorders the senders read, which are the most active by score alone,
// so a sender just outside that read cannot be ranked in by recency; the
// read is five times the sample, so one that would is far down it.
func sampleRank(sender RecentSender) time.Time {
lead := sender.ActivityScore.Sub(sender.LastSentAt)
return sender.ActivityScore.Add(-time.Duration(sampleChattersRecencyWeight * float64(lead)))
}

func (s *Server) SampleChatters(ctx context.Context, req *chatpb.SampleChattersRequest) (*chatpb.SampleChattersResponse, error) {
log := s.log.With(zap.String("chat_id", model.ChatIDString(req.ChatId)))
if req.Auth != nil {
Expand Down Expand Up @@ -126,13 +150,21 @@ type sampleCandidate struct {

// sampleChatters builds the sample of the public group c, as described above.
func (s *Server) sampleChatters(ctx context.Context, c *Chat) ([]*chatpb.SampledChatter, bool, error) {
senders, err := s.chats.GetRecentSenders(ctx, c.ID, sampleChattersSenderWindow)
senders, err := s.chats.GetActiveSenders(ctx, c.ID, sampleChattersSenderWindow)
if err != nil {
return nil, false, err
}
// By rank, and of two with the same rank, the more recent sender first.
sort.SliceStable(senders, func(i, j int) bool {
ri, rj := sampleRank(senders[i]), sampleRank(senders[j])
if !ri.Equal(rj) {
return ri.After(rj)
}
return senders[i].LastSentAt.After(senders[j].LastSentAt)
})

// The creator first, at their send time if the senders read holds it
// (otherwise looked up below), then every other sender in recency order.
// (otherwise looked up below), then every other sender in rank order.
candidates := make([]sampleCandidate, 0, len(senders)+1)
if c.CreatorID != nil {
creator := sampleCandidate{userID: c.CreatorID, isCreator: true}
Expand Down
11 changes: 11 additions & 0 deletions chat/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -507,6 +507,17 @@ type Store interface {
// chat ID.
GetRecentSenders(ctx context.Context, chatID *commonpb.ChatId, limit int) ([]RecentSender, error)

// GetActiveSenders returns the group's activity records most active
// first, by activity score (see NextActivityScore), at most limit of them
// (limit <= 0 means unbounded), from an eventually consistent read, as
// GetRecentSenders. Ties come back in no particular order. A record
// written before scores existed and never since backfilled is not
// returned. Records of users who have since left are included, 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.
GetActiveSenders(ctx context.Context, chatID *commonpb.ChatId, limit int) ([]RecentSender, error)

// GetLastSentAt returns userID's activity record in the group chatID:
// when their latest recorded send was, and whether they have a record at
// all. It is the point read of what GetRecentSenders ranges over, for a
Expand Down
41 changes: 41 additions & 0 deletions chat/tests/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@ func RunServerTests(t *testing.T, s chat.Store, teardown func()) {
testServer_SampleChatters,
testServer_SampleChatters_Size,
testServer_SampleChatters_Window,
testServer_SampleChatters_Activity,
testServer_SampleChatters_Gates,
testServer_GetDmChatFeed_Empty,
testServer_GetDmChatFeed_OrderAndContent,
Expand Down Expand Up @@ -1250,6 +1251,46 @@ func testServer_SampleChatters_Window(t *testing.T, s chat.Store) {
require.True(t, resp.Chatters[0].LastSentAt.AsTime().Equal(at(0)))
}

func testServer_SampleChatters_Activity(t *testing.T, s chat.Store) {
e := newServerEnv(t, s)

creator := model.MustGenerateUserID()
regular := model.MustGenerateUserID() // sent daily for a week
oneOff := model.MustGenerateUserID() // sent once, a day after the regular
nearby := model.MustGenerateUserID() // sent once, later, at a score just below the regular's
groupID := e.putGroupWithCreator("Regulars", creator, regular, oneOff, nearby)

day := 24 * time.Hour
var score, last time.Time
for d := range 7 {
sentAt := at(0).Add(time.Duration(d) * day)
e.recordSend(groupID, regular, sentAt)
score, last = chat.NextActivityScore(score, last, sentAt), sentAt
}
e.recordSend(groupID, oneOff, last.Add(day))

// Distant scores go to the more active: the regular, about six days
// ahead of their last send, above someone who sent once a day after it.
// Each is shown at their own last send.
resp := e.sampleChatters(e.keys, groupID)
require.Equal(t, chatpb.SampleChattersResponse_OK, resp.Result)
require.Equal(t, [][]byte{creator.Value, regular.Value, oneOff.Value}, sampledUserIDs(resp.Chatters))
require.True(t, resp.Chatters[1].LastSentAt.AsTime().Equal(last))
require.True(t, resp.Chatters[2].LastSentAt.AsTime().Equal(last.Add(day)))

// Close scores go to the more recent: a single send at nine tenths of
// the regular's lead scores just below them, but its last send is so
// much later that it ranks first.
lead := score.Sub(last)
nearbyAt := last.Add(lead * 9 / 10)
e.recordSend(groupID, nearby, nearbyAt)
senders, err := s.GetActiveSenders(e.ctx, groupID, 0)
require.NoError(t, err)
require.Equal(t, regular.Value, senders[0].UserID.Value, "by score alone, the regular leads")
resp = e.sampleChatters(e.keys, groupID)
require.Equal(t, [][]byte{creator.Value, nearby.Value, regular.Value, oneOff.Value}, sampledUserIDs(resp.Chatters))
}

func testServer_SampleChatters_Gates(t *testing.T, s chat.Store) {
e := newServerEnv(t, s)
_, strangerKeys := e.addUser()
Expand Down
Loading
Loading