From a8a7f04a79b7d5214d81b106383a9c55e4e7eeab Mon Sep 17 00:00:00 2001 From: Dylan Jeffers Date: Wed, 23 Sep 2026 16:40:22 -0700 Subject: [PATCH] fix(api): pick Weekly Rotation push recipients from plays challenge_listen_streak is only written by ListenStreakProcessor, which is a no-op while challenge `e` is inactive (it is in prod), so the table is stale and the weekly_rotation job was selecting almost no one. Use a recent play instead, via ix_plays_user_hour. Also trims the job's doc comment and drops the claim that the send window picks up listeners who become eligible after the cursor passes them. Co-Authored-By: Claude Opus 5.5 --- jobs/create_weekly_rotation_notifications.go | 63 ++++++------------- ...eate_weekly_rotation_notifications_test.go | 16 ++--- 2 files changed, 28 insertions(+), 51 deletions(-) diff --git a/jobs/create_weekly_rotation_notifications.go b/jobs/create_weekly_rotation_notifications.go index f18b7ec9..54b48343 100644 --- a/jobs/create_weekly_rotation_notifications.go +++ b/jobs/create_weekly_rotation_notifications.go @@ -14,44 +14,16 @@ import ( "go.uber.org/zap" ) -// WeeklyRotationNotificationsJob tells listeners their Weekly Rotation has -// rolled over. +// WeeklyRotationNotificationsJob inserts one weekly_rotation notification per +// recent listener per period (weeklyrotation.Period); pedalboard turns the row +// into a push, gated by the push_weekly_rotation remote-config variable. // -// The mix itself is computed on demand by GET /v1/users/:id/weekly-rotation -// and rolls over every Wednesday 00:00 UTC (weeklyrotation.Period), so there -// is nothing to precompute here. This job only fans out one -// `weekly_rotation` notification row per eligible listener per period; the -// pedalboard notifications app turns the row into a push (gated there by -// the push_weekly_rotation remote-config variable), and the clients render -// it in the notification feed. -// -// WHO. Anyone who has listened in the last weeklyRotationActiveWindow and -// whose account is live. Recency comes from challenge_listen_streak, which -// is one row per listener with a last_listen_date and cheap to scan, unlike -// plays. Recent listening is a proxy for "has a non-empty mix": the -// endpoint's affinity and history terms need plays to work with, and -// computing every user's mix here just to check would cost far more than -// the pushes are worth. -// -// WHEN. From weeklyRotationSendHourUTC on the Wednesday the period opens, -// for weeklyRotationSendWindow. The mix is ready at 00:00 UTC, but a push at -// 5pm Pacific on a Tuesday reads as noise; 16:00 UTC is 9am Pacific / noon -// Eastern. Someone who only becomes eligible after the window closes waits -// for the next Wednesday rather than getting "your rotation is ready" on a -// Saturday. Stage sends throughout the period so the flow can be exercised -// without waiting for a Wednesday, mirroring ListenStreakReminderJob. -// -// PACING. Each run inserts at most batchSize rows, walking user_id upward -// from a per-period cursor. Every row costs the notifications app an -// identity lookup and an SNS publish, and a single INSERT of the whole -// listener base would land there as one burst. At the scheduled interval -// this drains a few hundred thousand listeners inside the window without -// spiking anything. The cursor is in-memory: after a restart the job -// rescans from the start and the NOT EXISTS skips what was already sent. -// -// IDEMPOTENCY. group_id = weekly_rotation:: and specifier -// = user_id, so uq_notification (group_id, specifier) makes reruns, -// restarts and concurrent replicas safe. +// Recipients are live users with a play in the last weeklyRotationActiveWindow. +// Sends start at weeklyRotationSendHourUTC on rollover Wednesday and last for +// weeklyRotationSendWindow; stage sends all period. Each run inserts up to +// batchSize rows, walking user_id upward from an in-memory cursor, so one pass +// covers the listeners eligible when it reaches them. uq_notification +// (group_id, specifier) makes reruns, restarts and multiple replicas safe. type WeeklyRotationNotificationsJob struct { pool database.DbPool logger *zap.Logger @@ -71,8 +43,7 @@ const ( // weeklyRotationSendHourUTC is the hour (UTC) on rollover day the fan-out // begins. weeklyRotationSendHourUTC = 16 - // weeklyRotationSendWindow is how long after that the job keeps picking - // up newly eligible listeners. + // weeklyRotationSendWindow is how long after that the job keeps sending. weeklyRotationSendWindow = 24 * time.Hour // weeklyRotationActiveWindow is how recently someone must have listened // to be told about their mix. @@ -168,14 +139,18 @@ func (j *WeeklyRotationNotificationsJob) run(ctx context.Context) error { 'weekly_rotation', jsonb_build_object('year', @year::int, 'week', @week::int), @now - FROM challenge_listen_streak cls - JOIN users u ON u.user_id = cls.user_id - WHERE cls.last_listen_date >= @active_since - AND u.user_id > @after_user_id + FROM users u + WHERE u.user_id > @after_user_id AND u.is_current AND NOT u.is_deactivated AND u.is_available AND u.handle IS NOT NULL + -- Matches ix_plays_user_hour's expression so this is an index range scan. + AND EXISTS ( + SELECT 1 FROM plays p + WHERE p.user_id = u.user_id + AND date_trunc('hour', p.created_at) >= date_trunc('hour', @active_since::timestamp) + ) AND NOT EXISTS ( SELECT 1 FROM notification n WHERE n.group_id = @group_prefix || u.user_id::text @@ -189,7 +164,7 @@ func (j *WeeklyRotationNotificationsJob) run(ctx context.Context) error { "year": year, "week": week, "now": now, - "active_since": now.Add(-weeklyRotationActiveWindow), + "active_since": now.Add(-weeklyRotationActiveWindow).UTC(), "after_user_id": j.cursorUserId, "batch_size": j.batchSize, }) diff --git a/jobs/create_weekly_rotation_notifications_test.go b/jobs/create_weekly_rotation_notifications_test.go index fee1619e..8c710f77 100644 --- a/jobs/create_weekly_rotation_notifications_test.go +++ b/jobs/create_weekly_rotation_notifications_test.go @@ -24,16 +24,18 @@ func seedWeeklyRotationListeners(pool *pgxpool.Pool, now time.Time) { {"user_id": 4, "handle": "four", "wallet": "0x04"}, {"user_id": 5, "handle": "five", "wallet": "0x05"}, }, - "challenge_listen_streak": { + "plays": { // Listened this week -> notified. - {"user_id": 1, "listen_streak": 5, "last_listen_date": now.Add(-2 * 24 * time.Hour)}, + {"id": 1, "user_id": 1, "play_item_id": 100, "created_at": now.Add(-2 * 24 * time.Hour)}, // Last listen too long ago -> not notified. - {"user_id": 2, "listen_streak": 1, "last_listen_date": now.Add(-40 * 24 * time.Hour)}, + {"id": 2, "user_id": 2, "play_item_id": 100, "created_at": now.Add(-40 * 24 * time.Hour)}, // Deactivated -> not notified even though active. - {"user_id": 3, "listen_streak": 9, "last_listen_date": now.Add(-24 * time.Hour)}, - // Listened right at the edge of the window -> notified. - {"user_id": 4, "listen_streak": 1, "last_listen_date": now.Add(-29 * 24 * time.Hour)}, - // user 5 has never listened: no streak row -> not notified. + {"id": 3, "user_id": 3, "play_item_id": 100, "created_at": now.Add(-24 * time.Hour)}, + // Listened near the edge of the window -> notified. + {"id": 4, "user_id": 4, "play_item_id": 100, "created_at": now.Add(-29 * 24 * time.Hour)}, + // Anonymous play -> no one to notify. + {"id": 5, "user_id": nil, "play_item_id": 100, "created_at": now.Add(-time.Hour)}, + // user 5 has never listened -> not notified. }, }) }