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
2 changes: 1 addition & 1 deletion CLAUDE.md

Large diffs are not rendered by default.

93 changes: 93 additions & 0 deletions chat/activity.go
Original file line number Diff line number Diff line change
@@ -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
}
108 changes: 108 additions & 0 deletions chat/activity_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
16 changes: 9 additions & 7 deletions chat/cache/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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()
Expand Down
99 changes: 99 additions & 0 deletions chat/dynamodb/activity_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
Loading
Loading