diff --git a/cmd/shithubd/worker.go b/cmd/shithubd/worker.go index 053ca552..d59d96cd 100644 --- a/cmd/shithubd/worker.go +++ b/cmd/shithubd/worker.go @@ -33,6 +33,7 @@ import ( "github.com/tenseleyFlow/shithub/internal/infra/storage" "github.com/tenseleyFlow/shithub/internal/notifications" "github.com/tenseleyFlow/shithub/internal/orgs" + repotraffic "github.com/tenseleyFlow/shithub/internal/repos/traffic" "github.com/tenseleyFlow/shithub/internal/secretscan" "github.com/tenseleyFlow/shithub/internal/webhook" "github.com/tenseleyFlow/shithub/internal/webhookrelay" @@ -161,6 +162,9 @@ var workerCmd = &cobra.Command{ p.Register(worker.KindJobsPurge, jobs.JobsPurge(jobs.JobsPurgeDeps{ Pool: pool, Logger: logger, })) + p.Register(repotraffic.KindTrafficPurge, jobs.TrafficPurge(jobs.TrafficPurgeDeps{ + Pool: pool, Logger: logger, + })) p.Register(worker.KindLifecycleSweep, jobs.LifecycleSweep(jobs.LifecycleSweepDeps{ Pool: pool, RepoFS: rfs, Audit: auditRecorder(), Logger: logger, })) diff --git a/deploy/systemd/shithubd-cron.service b/deploy/systemd/shithubd-cron.service index 00026d09..3ac57da8 100644 --- a/deploy/systemd/shithubd-cron.service +++ b/deploy/systemd/shithubd-cron.service @@ -14,3 +14,4 @@ ExecStart=/usr/local/bin/shithubd admin run-job lifecycle:sweep ExecStart=/usr/local/bin/shithubd admin run-job jobs:purge_completed ExecStart=/usr/local/bin/shithubd admin run-job webhook:purge_old ExecStart=/usr/local/bin/shithubd admin run-job workflow:cleanup +ExecStart=/usr/local/bin/shithubd admin run-job traffic:purge diff --git a/docs/internal/actions-runner-api.md b/docs/internal/actions-runner-api.md index e655e6eb..63566b3b 100644 --- a/docs/internal/actions-runner-api.md +++ b/docs/internal/actions-runner-api.md @@ -303,3 +303,9 @@ runner posts terminal job status `cancelled`. - `shithub_actions_step_timeouts_total` - `shithub_actions_storage_objects{kind="artifacts|step_logs|hot_log_chunks"}` - `shithub_actions_storage_bytes{kind="artifacts|step_logs|hot_log_chunks"}` + +The gauge observer refreshes queue, runner and object-count gauges every +15 s. `shithub_actions_storage_bytes{kind="hot_log_chunks"}` is refreshed +every 5 min instead: summing `octet_length(chunk)` detoasts every row of +`workflow_step_log_chunks`, so it is scraped on a slower cadence and can +lag the other gauges by up to five minutes. diff --git a/docs/internal/db-indexes.md b/docs/internal/db-indexes.md index 78a7ea21..e56abb02 100644 --- a/docs/internal/db-indexes.md +++ b/docs/internal/db-indexes.md @@ -94,6 +94,26 @@ rows for a 1M-row table; "medium" 1–10k; "low" most-of-the-table. | `signup_ip_throttle` PK `(cidr)` | per-/24 lookup | high | UPSERT | | `signup_ip_throttle_window_started_idx (window_started_at)` | periodic prune | low | scan-friendly | +## Repo traffic (S38 + retention, availability campaign) + +| Index | Covers query | Selectivity | Cost notes | +|---|---|---|---| +| `repo_traffic_daily` PK `(repo_id, day)` | Traffic chart for one repo | high | UPSERT hot path | +| `repo_traffic_daily_day_idx (day DESC)` | 400-day retention purge | low | scan-friendly | +| `repo_traffic_paths` PK `(repo_id, day, path)` | popular-content rollup | high | UPSERT hot path | +| `repo_traffic_paths_day_idx (day)` | 30-day retention purge | low | 0129; without it every purge batch is a seq scan | +| `repo_traffic_referrers` PK `(repo_id, day, referrer)` | referrer rollup | high | UPSERT hot path | +| `repo_traffic_referrers_day_idx (day)` | 30-day retention purge | low | 0129 | +| `repo_traffic_uniques` PK `(repo_id, day, metric, key, visitor_hash)` | dedupe a visitor within a day | high | INSERT ... ON CONFLICT DO NOTHING, one per pageview | +| `repo_traffic_uniques_created_idx (created_at)` | 30-day retention purge | low | 0112; the purge filters this table on `created_at` rather than `day` so it can reuse this index instead of adding a second one to a table that takes an insert per pageview | + +The purge (`traffic:purge`) deletes with +`WHERE ctid IN (SELECT ctid FROM WHERE < $1 LIMIT $2)`. +The `day`/`created_at` index drives the subselect; the outer delete is a +tid scan. Without the index the plan is a sequential scan **per batch**, +which on the production table sizes (1.26 M paths, 1.34 M uniques as of +2026-09-02) is a few hundred full scans per run. + ## Future considerations (deferred) - **`pg_stat_statements` extension.** S37's deploy doc owns the diff --git a/docs/internal/repository-insights.md b/docs/internal/repository-insights.md index a86b4653..fe6cd524 100644 --- a/docs/internal/repository-insights.md +++ b/docs/internal/repository-insights.md @@ -25,6 +25,36 @@ addresses, user agents, authenticated user IDs, and full referrer URLs are not persisted. Referrers are reduced to an external host and same-site referrers are dropped. +## Traffic retention + +The traffic tables are pruned by the `traffic:purge` worker job, enqueued +nightly from `deploy/systemd/shithubd-cron.service`. Retention windows live in +`internal/repos/traffic/purge.go`: + +| Table | Window | Cutoff column | Why | +|---|---|---|---| +| `repo_traffic_uniques` | 30 days | `created_at` | one row per visitor digest per repo/day/metric | +| `repo_traffic_paths` | 30 days | `day` | one row per distinct path per repo/day; crawlers inflate this hardest | +| `repo_traffic_referrers` | 30 days | `day` | one row per external host per repo/day | +| `repo_traffic_daily` | 400 days | `day` | one row per repo/day; the only long-term history, and small enough to keep | + +Thirty days is deliberately more than double the fourteen the Traffic UI reads +(`traffic.DefaultWindowDays`), so a purge can never eat a bar the chart would +draw. `repo_traffic_uniques` is filtered on `created_at` rather than `day` +because that is the column it already has an index on, and the request path +stamps both from the same instant. + +The job deletes in batches of 5,000 rows, each its own statement, up to 2,000 +batches per table per run; the payload (`retention_days`, `daily_retention_days`, +`batch_size`, `max_batches`) overrides any of that for an ad-hoc +`shithubd admin run-job traffic:purge`. Nothing is done in one big transaction: +the tables were left unpruned from 2026-05-18 until the 2026-09-02 availability +sitrep, by which point they held 881 MB of a 988 MB database, and a single +DELETE over that backlog would have locked the write path for minutes. A run +that stops on the batch cap re-enqueues itself so the backlog drains without +waiting for the next cron beat. Re-running is always safe — the cutoff is +recomputed from the clock and rows inside the window are never touched. + ## Refresh Flow `push:process` enqueues `repo:insights_recalc` whenever the repository default diff --git a/docs/internal/retro/2026-09-02-availability-sitrep.md b/docs/internal/retro/2026-09-02-availability-sitrep.md index e64fbfd2..008036a3 100644 --- a/docs/internal/retro/2026-09-02-availability-sitrep.md +++ b/docs/internal/retro/2026-09-02-availability-sitrep.md @@ -153,12 +153,21 @@ The verification items below are still the operator's. - [x] Key the anonymous HTML tier by `/24` for repo history/blob/raw routes (Meta rotates within `57.141.2.0/24`) — applied to the whole anonymous tier, not just those routes -- [ ] Retention job for `repo_traffic_paths` / `repo_traffic_uniques` - (14-day window, matches the Traffic UI) + one-off prune migration +- [x] Retention job for `repo_traffic_paths` / `repo_traffic_uniques` + (plus `_referrers`) — `traffic:purge`, nightly from + `shithubd-cron.service`, 30-day window rather than the UI's 14 so a + purge can never truncate the chart; `repo_traffic_daily` keeps 400 + days. Deletes run in 5k-row batches, capped per run, so the first + pass over the 2.5 M-row backlog never holds a long transaction — + which is also why the backfill is the job's first run and not a + bulk DELETE in a migration. 0129 adds the `day` indexes the purge + needs (`repo_traffic_uniques` reuses its existing `created_at` + index). See `docs/internal/repository-insights.md`. - [ ] Cache per-entry last-commit for the code tab (single `git log --name-only` walk, or an LRU keyed by tree OID) and cache `rev-list --count` / recursive `ls-tree` per head OID -- [ ] `actionsobserver`: drop the `octet_length` sum or run it every 5 min +- [x] `actionsobserver`: the `octet_length` sum now runs every 5 min on its + own cadence; the count and queue-depth gauges stay at 15 s ### Phase 4 — observability and docs diff --git a/docs/internal/runbooks/actions.md b/docs/internal/runbooks/actions.md index f1d7e6d8..dcae21ab 100644 --- a/docs/internal/runbooks/actions.md +++ b/docs/internal/runbooks/actions.md @@ -187,6 +187,12 @@ Important metrics: - `shithub_actions_log_chunk_bytes_total{location="server"}` - `shithub_actions_storage_objects{kind="artifacts|step_logs|hot_log_chunks"}` - `shithub_actions_storage_bytes{kind="artifacts|step_logs|hot_log_chunks"}` + +The gauge observer refreshes queue, runner and object-count gauges every +15 s. `shithub_actions_storage_bytes{kind="hot_log_chunks"}` is refreshed +every 5 min instead: summing `octet_length(chunk)` detoasts every row of +`workflow_step_log_chunks`, so it is scraped on a slower cadence and can +lag the other gauges by up to five minutes. - `shithub_actions_step_timeouts_total` The committed dashboard JSON lives at: diff --git a/docs/internal/worker.md b/docs/internal/worker.md index 504bc21e..e755a89b 100644 --- a/docs/internal/worker.md +++ b/docs/internal/worker.md @@ -41,6 +41,7 @@ backstop poll (every 5s by default) covers dropped notifications. | `workflow:cleanup` | cron / manual ad-hoc | retention cutoff + idempotent deletes | | `trending:compute` | recurring self-scheduled S42 job | append-only snapshots | | `org:scheduled_reminder_sweep` | cron / manual ad-hoc | reminder delivery rows | +| `traffic:purge` | cron / manual ad-hoc | retention cutoff recomputed per run; batched deletes | Adding a new kind: write the handler in `internal/worker/jobs/.go`, add the `Kind` constant to `internal/worker/types.go`, register it in diff --git a/internal/infra/metrics/actionsobserver.go b/internal/infra/metrics/actionsobserver.go index e74c0f13..d90ab987 100644 --- a/internal/infra/metrics/actionsobserver.go +++ b/internal/infra/metrics/actionsobserver.go @@ -11,31 +11,95 @@ import ( const actionsRunnerStaleAfter = 60 * time.Second +// defaultActionsInterval is the cadence for the cheap gauges when the caller +// does not pick one. +const defaultActionsInterval = 15 * time.Second + +// actionsStorageBytesInterval bounds how often the hot log-chunk byte sum +// runs. `sum(octet_length(chunk))` has to scan and detoast every row of +// workflow_step_log_chunks, which costs the same whether or not anything is +// running; at the 15s cadence of the other gauges it was a standing load on +// the database. Chunk volume moves slowly enough that a 5 minute gauge is +// still useful. +const actionsStorageBytesInterval = 5 * time.Minute + // ObserveActions starts a goroutine that periodically refreshes DB-backed // Actions gauges. The goroutine exits when ctx is canceled. +// +// interval drives the queue, runner and object-count gauges. The hot +// log-chunk byte sum is refreshed on the slower actionsStorageBytesInterval +// cadence; see refreshActionLogChunkBytes. func ObserveActions(ctx context.Context, pool *pgxpool.Pool, interval time.Duration) { if pool == nil { return } if interval <= 0 { - interval = 15 * time.Second + interval = defaultActionsInterval } + slowEvery := ticksBetween(interval, actionsStorageBytesInterval) + t := time.NewTicker(interval) go func() { - refreshActions(ctx, pool) - t := time.NewTicker(interval) defer t.Stop() - for { - select { - case <-ctx.Done(): + observeActionsLoop(ctx, t.C, slowEvery, + func(ctx context.Context) { refreshActionsFast(ctx, pool) }, + func(ctx context.Context) { refreshActionLogChunkBytes(ctx, pool) }, + ) + }() +} + +// ticksBetween returns how many ticks of length tick must elapse between two +// runs of a task that should run at most once per every. It rounds up, so the +// task never runs more often than requested, and never returns less than 1. +func ticksBetween(tick, every time.Duration) int { + if tick <= 0 || every <= tick { + return 1 + } + n := int((every + tick - 1) / tick) + if n < 1 { + return 1 + } + return n +} + +// observeActionsLoop runs fast on every tick and slow once every slowEvery +// ticks. Both run once up front so the gauges are populated before the first +// tick. It returns when ctx is canceled or ticks is closed. +func observeActionsLoop(ctx context.Context, ticks <-chan time.Time, slowEvery int, fast, slow func(context.Context)) { + if slowEvery < 1 { + slowEvery = 1 + } + fast(ctx) + slow(ctx) + sinceSlow := 0 + for { + select { + case <-ctx.Done(): + return + case _, ok := <-ticks: + if !ok { return - case <-t.C: - refreshActions(ctx, pool) + } + fast(ctx) + sinceSlow++ + if sinceSlow >= slowEvery { + sinceSlow = 0 + slow(ctx) } } - }() + } } +// refreshActions refreshes every Actions gauge, cheap and expensive alike. +// The observer loop splits the two cadences apart; this is the one-shot form. func refreshActions(ctx context.Context, pool *pgxpool.Pool) { + if pool == nil { + return + } + refreshActionsFast(ctx, pool) + refreshActionLogChunkBytes(ctx, pool) +} + +func refreshActionsFast(ctx context.Context, pool *pgxpool.Pool) { if pool == nil { return } @@ -155,13 +219,16 @@ GROUP BY r.id, r.name, r.status, r.capacity, r.last_heartbeat_at, r.draining_at, ActionsRunnerStaleTotal.Set(stale) } +// refreshActionStorageGauges publishes the object counts for all three storage +// kinds plus the two byte sums that read a plain integer column. The +// hot_log_chunks byte sum is deliberately absent: it is the only one that has +// to detoast, so refreshActionLogChunkBytes owns that gauge. func refreshActionStorageGauges(ctx context.Context, pool *pgxpool.Pool) { ActionsStorageObjects.WithLabelValues("artifacts").Set(0) ActionsStorageObjects.WithLabelValues("step_logs").Set(0) ActionsStorageObjects.WithLabelValues("hot_log_chunks").Set(0) ActionsStorageBytes.WithLabelValues("artifacts").Set(0) ActionsStorageBytes.WithLabelValues("step_logs").Set(0) - ActionsStorageBytes.WithLabelValues("hot_log_chunks").Set(0) rows, err := pool.Query(ctx, ` SELECT 'artifacts'::text AS kind, count(*)::double precision, COALESCE(sum(byte_count), 0)::double precision @@ -171,7 +238,7 @@ SELECT 'step_logs'::text AS kind, count(*)::double precision, COALESCE(sum(log_b FROM workflow_steps WHERE log_object_key IS NOT NULL UNION ALL -SELECT 'hot_log_chunks'::text AS kind, count(*)::double precision, COALESCE(sum(octet_length(chunk)), 0)::double precision +SELECT 'hot_log_chunks'::text AS kind, count(*)::double precision, 0::double precision FROM workflow_step_log_chunks`) if err != nil { return @@ -184,6 +251,27 @@ FROM workflow_step_log_chunks`) return } ActionsStorageObjects.WithLabelValues(kind).Set(objects) + if kind == "hot_log_chunks" { + continue + } ActionsStorageBytes.WithLabelValues(kind).Set(bytes) } } + +// refreshActionLogChunkBytes publishes shithub_actions_storage_bytes for the +// hot chunk table. Every row is a bytea that Postgres has to fetch out of the +// TOAST heap to measure, so this runs on actionsStorageBytesInterval rather +// than with the cheap gauges. +func refreshActionLogChunkBytes(ctx context.Context, pool *pgxpool.Pool) { + if pool == nil { + return + } + var bytes float64 + err := pool.QueryRow(ctx, ` +SELECT COALESCE(sum(octet_length(chunk)), 0)::double precision +FROM workflow_step_log_chunks`).Scan(&bytes) + if err != nil { + return + } + ActionsStorageBytes.WithLabelValues("hot_log_chunks").Set(bytes) +} diff --git a/internal/infra/metrics/actionsobserver_schedule_test.go b/internal/infra/metrics/actionsobserver_schedule_test.go new file mode 100644 index 00000000..ed3eb606 --- /dev/null +++ b/internal/infra/metrics/actionsobserver_schedule_test.go @@ -0,0 +1,117 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package metrics + +import ( + "context" + "testing" + "time" +) + +func TestTicksBetween(t *testing.T) { + tests := []struct { + name string + tick time.Duration + every time.Duration + want int + }{ + {"production cadence", 15 * time.Second, 5 * time.Minute, 20}, + {"rounds up", 7 * time.Second, 5 * time.Minute, 43}, + {"slow interval equals tick", 15 * time.Second, 15 * time.Second, 1}, + {"slow interval below tick", time.Minute, 5 * time.Second, 1}, + {"zero tick", 0, 5 * time.Minute, 1}, + {"negative tick", -time.Second, 5 * time.Minute, 1}, + {"zero slow interval", 15 * time.Second, 0, 1}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + if got := ticksBetween(tc.tick, tc.every); got != tc.want { + t.Fatalf("ticksBetween(%v, %v) = %d, want %d", tc.tick, tc.every, got, tc.want) + } + }) + } +} + +// The point of the split is that the expensive refresh must not keep pace +// with the cheap one, so assert the exact call counts over a run of ticks. +func TestObserveActionsLoopRunsSlowRefreshEveryNTicks(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + ticks := make(chan time.Time) + var fast, slow int + done := make(chan struct{}) + go func() { + defer close(done) + observeActionsLoop(ctx, ticks, 4, + func(context.Context) { fast++ }, + func(context.Context) { slow++ }, + ) + }() + + // Both refreshes run once before the first tick so the gauges are not + // empty for the first interval; 12 ticks at slowEvery=4 add 3 slow runs. + for i := 0; i < 12; i++ { + ticks <- time.Now() + } + // Closing the channel makes the loop return, which orders the counter + // writes before the reads below. + close(ticks) + <-done + + if fast != 13 { + t.Fatalf("fast refreshes = %d, want 13 (initial plus one per tick)", fast) + } + if slow != 4 { + t.Fatalf("slow refreshes = %d, want 4 (initial plus one per 4 ticks)", slow) + } +} + +func TestObserveActionsLoopStopsOnContextCancel(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + ticks := make(chan time.Time) + done := make(chan struct{}) + var fast, slow int + go func() { + defer close(done) + observeActionsLoop(ctx, ticks, 2, + func(context.Context) { fast++ }, + func(context.Context) { slow++ }, + ) + }() + + cancel() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("observeActionsLoop did not return after context cancel") + } + if fast != 1 || slow != 1 { + t.Fatalf("fast=%d slow=%d, want the single up-front refresh of each", fast, slow) + } +} + +func TestObserveActionsLoopClampsSlowEvery(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + ticks := make(chan time.Time) + var fast, slow int + done := make(chan struct{}) + go func() { + defer close(done) + observeActionsLoop(ctx, ticks, 0, + func(context.Context) { fast++ }, + func(context.Context) { slow++ }, + ) + }() + for i := 0; i < 3; i++ { + ticks <- time.Now() + } + close(ticks) + <-done + + if fast != 4 || slow != 4 { + t.Fatalf("fast=%d slow=%d, want 4 and 4 (slowEvery clamped to 1)", fast, slow) + } +} diff --git a/internal/migrationsfs/migrations/0129_repo_traffic_retention_indexes.sql b/internal/migrationsfs/migrations/0129_repo_traffic_retention_indexes.sql new file mode 100644 index 00000000..861e6ed3 --- /dev/null +++ b/internal/migrationsfs/migrations/0129_repo_traffic_retention_indexes.sql @@ -0,0 +1,44 @@ +-- SPDX-License-Identifier: AGPL-3.0-or-later +-- +-- Retention indexes for the repo traffic tables (availability campaign, +-- 2026-09-02 sitrep root cause #7). +-- +-- 0112 shipped these tables with no purge and no index that a purge can +-- use. In production they had grown to 881 MB of a 988 MB database: +-- 1.26 M rows in repo_traffic_paths and 1.34 M in repo_traffic_uniques, +-- ~30 k/day/table, most of it crawler paths. The `traffic:purge` worker +-- job now trims them to a 30-day window; without a day index every +-- batch of that purge would be a sequential scan of the whole table. +-- +-- repo_traffic_uniques is deliberately absent: it already carries +-- repo_traffic_uniques_created_idx (created_at), and created_at is set +-- on the same request that derives `day`, so the purge filters that +-- table on created_at and reuses the existing index rather than paying +-- for a second one on a table that takes an insert per pageview. +-- +-- repo_traffic_daily already has repo_traffic_daily_day_idx (day DESC), +-- which serves its (much longer) retention window. +-- +-- CONCURRENTLY, because repo_traffic_paths and repo_traffic_referrers +-- are written synchronously in the request path — a plain CREATE INDEX +-- would block pageviews for the length of the build. That forces +-- NO TRANSACTION. If a build is interrupted Postgres leaves an INVALID +-- index behind; drop it by name and re-run the migration. +-- +-- No bulk DELETE here on purpose. Pruning 2.5 M rows inside a migration +-- would hold one long transaction during a deploy, which is the failure +-- mode this campaign is trying to remove. The first run of +-- `traffic:purge` does the backfill in bounded batches instead. + +-- +goose NO TRANSACTION + +-- +goose Up +CREATE INDEX CONCURRENTLY IF NOT EXISTS repo_traffic_paths_day_idx + ON repo_traffic_paths (day); + +CREATE INDEX CONCURRENTLY IF NOT EXISTS repo_traffic_referrers_day_idx + ON repo_traffic_referrers (day); + +-- +goose Down +DROP INDEX CONCURRENTLY IF EXISTS repo_traffic_referrers_day_idx; +DROP INDEX CONCURRENTLY IF EXISTS repo_traffic_paths_day_idx; diff --git a/internal/repos/queries/repo_traffic.sql b/internal/repos/queries/repo_traffic.sql index 646e7c70..0053fe69 100644 --- a/internal/repos/queries/repo_traffic.sql +++ b/internal/repos/queries/repo_traffic.sql @@ -82,3 +82,52 @@ WHERE repo_id = $1 GROUP BY referrer ORDER BY views DESC, referrer ASC LIMIT $3; + +-- name: PurgeRepoTrafficUniquesBatch :execrows +-- Retention purge for the `traffic:purge` worker job: deletes at most +-- batch_size rows so a first run over a multi-million-row backlog never +-- holds a long transaction. The job loops until a batch comes back +-- short. Filtered on created_at rather than day because created_at is +-- stamped by the same request that derives day and is the column this +-- table already has an index on. +DELETE FROM repo_traffic_uniques +WHERE ctid IN ( + SELECT ctid + FROM repo_traffic_uniques + WHERE created_at < sqlc.arg(cutoff)::timestamptz + LIMIT sqlc.arg(batch_size)::bigint +); + +-- name: PurgeRepoTrafficPathsBatch :execrows +-- Batched retention purge; see PurgeRepoTrafficUniquesBatch. Uses +-- repo_traffic_paths_day_idx (0129). +DELETE FROM repo_traffic_paths +WHERE ctid IN ( + SELECT ctid + FROM repo_traffic_paths + WHERE day < sqlc.arg(cutoff)::date + LIMIT sqlc.arg(batch_size)::bigint +); + +-- name: PurgeRepoTrafficReferrersBatch :execrows +-- Batched retention purge; see PurgeRepoTrafficUniquesBatch. Uses +-- repo_traffic_referrers_day_idx (0129). +DELETE FROM repo_traffic_referrers +WHERE ctid IN ( + SELECT ctid + FROM repo_traffic_referrers + WHERE day < sqlc.arg(cutoff)::date + LIMIT sqlc.arg(batch_size)::bigint +); + +-- name: PurgeRepoTrafficDailyBatch :execrows +-- Batched retention purge; see PurgeRepoTrafficUniquesBatch. The daily +-- rollup is one row per repo per day, so it keeps a far longer window +-- than the per-path and per-visitor tables. +DELETE FROM repo_traffic_daily +WHERE ctid IN ( + SELECT ctid + FROM repo_traffic_daily + WHERE day < sqlc.arg(cutoff)::date + LIMIT sqlc.arg(batch_size)::bigint +); diff --git a/internal/repos/sqlc/querier.go b/internal/repos/sqlc/querier.go index 2dcfdbb9..e09c09ba 100644 --- a/internal/repos/sqlc/querier.go +++ b/internal/repos/sqlc/querier.go @@ -277,6 +277,23 @@ type Querier interface { // handler guarantees we don't pause an already-archived repo, but // the constraint defends in depth. PauseRepo(ctx context.Context, db DBTX, arg PauseRepoParams) error + // Batched retention purge; see PurgeRepoTrafficUniquesBatch. The daily + // rollup is one row per repo per day, so it keeps a far longer window + // than the per-path and per-visitor tables. + PurgeRepoTrafficDailyBatch(ctx context.Context, db DBTX, arg PurgeRepoTrafficDailyBatchParams) (int64, error) + // Batched retention purge; see PurgeRepoTrafficUniquesBatch. Uses + // repo_traffic_paths_day_idx (0129). + PurgeRepoTrafficPathsBatch(ctx context.Context, db DBTX, arg PurgeRepoTrafficPathsBatchParams) (int64, error) + // Batched retention purge; see PurgeRepoTrafficUniquesBatch. Uses + // repo_traffic_referrers_day_idx (0129). + PurgeRepoTrafficReferrersBatch(ctx context.Context, db DBTX, arg PurgeRepoTrafficReferrersBatchParams) (int64, error) + // Retention purge for the `traffic:purge` worker job: deletes at most + // batch_size rows so a first run over a multi-million-row backlog never + // holds a long transaction. The job loops until a batch comes back + // short. Filtered on created_at rather than day because created_at is + // stamped by the same request that derives day and is the column this + // table already has an index on. + PurgeRepoTrafficUniquesBatch(ctx context.Context, db DBTX, arg PurgeRepoTrafficUniquesBatchParams) (int64, error) RecordDependencyAutoTriageEvent(ctx context.Context, db DBTX, arg RecordDependencyAutoTriageEventParams) (DependencyAutoTriageEvent, error) RemoveIssueFromRepoProject(ctx context.Context, db DBTX, arg RemoveIssueFromRepoProjectParams) error RemoveRepoSecurityAdvisoryTeamCollaborator(ctx context.Context, db DBTX, arg RemoveRepoSecurityAdvisoryTeamCollaboratorParams) error diff --git a/internal/repos/sqlc/repo_traffic.sql.go b/internal/repos/sqlc/repo_traffic.sql.go index a5f8c53b..d28c8597 100644 --- a/internal/repos/sqlc/repo_traffic.sql.go +++ b/internal/repos/sqlc/repo_traffic.sql.go @@ -175,6 +175,111 @@ func (q *Queries) ListRepoTrafficReferrers(ctx context.Context, db DBTX, arg Lis return items, nil } +const purgeRepoTrafficDailyBatch = `-- name: PurgeRepoTrafficDailyBatch :execrows +DELETE FROM repo_traffic_daily +WHERE ctid IN ( + SELECT ctid + FROM repo_traffic_daily + WHERE day < $1::date + LIMIT $2::bigint +) +` + +type PurgeRepoTrafficDailyBatchParams struct { + Cutoff pgtype.Date + BatchSize int64 +} + +// Batched retention purge; see PurgeRepoTrafficUniquesBatch. The daily +// rollup is one row per repo per day, so it keeps a far longer window +// than the per-path and per-visitor tables. +func (q *Queries) PurgeRepoTrafficDailyBatch(ctx context.Context, db DBTX, arg PurgeRepoTrafficDailyBatchParams) (int64, error) { + result, err := db.Exec(ctx, purgeRepoTrafficDailyBatch, arg.Cutoff, arg.BatchSize) + if err != nil { + return 0, err + } + return result.RowsAffected(), nil +} + +const purgeRepoTrafficPathsBatch = `-- name: PurgeRepoTrafficPathsBatch :execrows +DELETE FROM repo_traffic_paths +WHERE ctid IN ( + SELECT ctid + FROM repo_traffic_paths + WHERE day < $1::date + LIMIT $2::bigint +) +` + +type PurgeRepoTrafficPathsBatchParams struct { + Cutoff pgtype.Date + BatchSize int64 +} + +// Batched retention purge; see PurgeRepoTrafficUniquesBatch. Uses +// repo_traffic_paths_day_idx (0129). +func (q *Queries) PurgeRepoTrafficPathsBatch(ctx context.Context, db DBTX, arg PurgeRepoTrafficPathsBatchParams) (int64, error) { + result, err := db.Exec(ctx, purgeRepoTrafficPathsBatch, arg.Cutoff, arg.BatchSize) + if err != nil { + return 0, err + } + return result.RowsAffected(), nil +} + +const purgeRepoTrafficReferrersBatch = `-- name: PurgeRepoTrafficReferrersBatch :execrows +DELETE FROM repo_traffic_referrers +WHERE ctid IN ( + SELECT ctid + FROM repo_traffic_referrers + WHERE day < $1::date + LIMIT $2::bigint +) +` + +type PurgeRepoTrafficReferrersBatchParams struct { + Cutoff pgtype.Date + BatchSize int64 +} + +// Batched retention purge; see PurgeRepoTrafficUniquesBatch. Uses +// repo_traffic_referrers_day_idx (0129). +func (q *Queries) PurgeRepoTrafficReferrersBatch(ctx context.Context, db DBTX, arg PurgeRepoTrafficReferrersBatchParams) (int64, error) { + result, err := db.Exec(ctx, purgeRepoTrafficReferrersBatch, arg.Cutoff, arg.BatchSize) + if err != nil { + return 0, err + } + return result.RowsAffected(), nil +} + +const purgeRepoTrafficUniquesBatch = `-- name: PurgeRepoTrafficUniquesBatch :execrows +DELETE FROM repo_traffic_uniques +WHERE ctid IN ( + SELECT ctid + FROM repo_traffic_uniques + WHERE created_at < $1::timestamptz + LIMIT $2::bigint +) +` + +type PurgeRepoTrafficUniquesBatchParams struct { + Cutoff pgtype.Timestamptz + BatchSize int64 +} + +// Retention purge for the `traffic:purge` worker job: deletes at most +// batch_size rows so a first run over a multi-million-row backlog never +// holds a long transaction. The job loops until a batch comes back +// short. Filtered on created_at rather than day because created_at is +// stamped by the same request that derives day and is the column this +// table already has an index on. +func (q *Queries) PurgeRepoTrafficUniquesBatch(ctx context.Context, db DBTX, arg PurgeRepoTrafficUniquesBatchParams) (int64, error) { + result, err := db.Exec(ctx, purgeRepoTrafficUniquesBatch, arg.Cutoff, arg.BatchSize) + if err != nil { + return 0, err + } + return result.RowsAffected(), nil +} + const upsertRepoTrafficDailyClone = `-- name: UpsertRepoTrafficDailyClone :exec INSERT INTO repo_traffic_daily ( repo_id, day, clones, unique_clones diff --git a/internal/repos/traffic/purge.go b/internal/repos/traffic/purge.go new file mode 100644 index 00000000..9cdb6bd3 --- /dev/null +++ b/internal/repos/traffic/purge.go @@ -0,0 +1,212 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package traffic + +import ( + "context" + "fmt" + "time" + + "github.com/jackc/pgx/v5/pgtype" + "github.com/jackc/pgx/v5/pgxpool" + + reposdb "github.com/tenseleyFlow/shithub/internal/repos/sqlc" + "github.com/tenseleyFlow/shithub/internal/worker" +) + +// KindTrafficPurge is the cron-driven retention sweep over the traffic +// tables. Registered by cmd/shithubd/worker.go, enqueued nightly by +// deploy/systemd/shithubd-cron.service. +const KindTrafficPurge worker.Kind = "traffic:purge" + +// Retention windows for the traffic tables. +// +// The per-path, per-referrer and per-visitor tables are the ones that +// grow without bound — one row per distinct path (or referrer, or +// visitor digest) per repo per day, which crawlers inflate hard. The +// Traffic UI only ever reads DefaultWindowDays (14) of them, so +// anything older is dead weight; the window here is deliberately more +// than double that so a purge can never eat a bar the UI would draw. +// +// repo_traffic_daily is different: one row per repo per day, no +// cardinality explosion, and it is the only long-term history the +// project has. It keeps a little over a year so a future +// year-over-year view has something to read, but it is still bounded. +const ( + DefaultRetentionDays = 30 + DefaultDailyRetentionDays = 400 +) + +// Batch sizing for the purge loop. +// +// DefaultPurgeBatch is small enough that one DELETE stays well inside a +// normal statement timeout and holds row locks briefly; DefaultMaxBatches +// stops a single run from grinding indefinitely if the backlog is larger +// than expected. The job is idempotent and cron-driven, so a run that +// stops early just leaves the rest for the next night. +const ( + DefaultPurgeBatch = 5000 + DefaultMaxBatches = 2000 +) + +// PurgeOptions tunes one run. Zero values select the defaults above. +type PurgeOptions struct { + // RetentionDays applies to repo_traffic_uniques, repo_traffic_paths + // and repo_traffic_referrers. + RetentionDays int + // DailyRetentionDays applies to repo_traffic_daily. + DailyRetentionDays int + // BatchSize is the row cap on a single DELETE statement. + BatchSize int + // MaxBatches caps the number of DELETEs per table per run. + MaxBatches int + // Now is injectable so tests can pin the cutoff. Defaults to time.Now. + Now func() time.Time +} + +func (o PurgeOptions) normalize() PurgeOptions { + if o.RetentionDays <= 0 { + o.RetentionDays = DefaultRetentionDays + } + if o.DailyRetentionDays <= 0 { + o.DailyRetentionDays = DefaultDailyRetentionDays + } + // A daily window shorter than the detail window would throw away the + // aggregate while its own source rows were still live, which is never + // what an operator means. + if o.DailyRetentionDays < o.RetentionDays { + o.DailyRetentionDays = o.RetentionDays + } + if o.BatchSize <= 0 { + o.BatchSize = DefaultPurgeBatch + } + if o.MaxBatches <= 0 { + o.MaxBatches = DefaultMaxBatches + } + if o.Now == nil { + o.Now = time.Now + } + return o +} + +// PurgeResult reports what one run removed. +type PurgeResult struct { + UniquesDeleted int64 + PathsDeleted int64 + ReferrersDeleted int64 + DailyDeleted int64 + // Remaining is true when at least one table hit MaxBatches, i.e. rows + // older than the cutoff are still there and the next run has work. + Remaining bool +} + +// Total is the row count deleted across all four tables. +func (r PurgeResult) Total() int64 { + return r.UniquesDeleted + r.PathsDeleted + r.ReferrersDeleted + r.DailyDeleted +} + +// Purge trims the traffic tables to their retention windows. +// +// Every DELETE runs on its own connection outside an explicit +// transaction, so no single statement holds locks for long and a run +// interrupted halfway through leaves the rows it already deleted +// deleted. Re-running is always safe: the cutoff is recomputed from the +// clock and rows inside the window are never touched. +func Purge(ctx context.Context, pool *pgxpool.Pool, opts PurgeOptions) (PurgeResult, error) { + opts = opts.normalize() + + now := opts.Now().UTC() + cutoffDay := startDate(now).AddDate(0, 0, -opts.RetentionDays) + dailyCutoffDay := startDate(now).AddDate(0, 0, -opts.DailyRetentionDays) + batch := int64(opts.BatchSize) + + q := reposdb.New() + var res PurgeResult + + // repo_traffic_uniques filters on created_at, which is the column it + // has an index on; midnight UTC of the cutoff day is the same instant + // the date comparison would pick. + deleted, more, err := purgeBatched(ctx, opts.MaxBatches, batch, + func(ctx context.Context, limit int64) (int64, error) { + return q.PurgeRepoTrafficUniquesBatch(ctx, pool, reposdb.PurgeRepoTrafficUniquesBatchParams{ + Cutoff: pgtype.Timestamptz{Time: cutoffDay, Valid: true}, + BatchSize: limit, + }) + }) + res.UniquesDeleted = deleted + res.Remaining = res.Remaining || more + if err != nil { + return res, fmt.Errorf("purge repo_traffic_uniques: %w", err) + } + + deleted, more, err = purgeBatched(ctx, opts.MaxBatches, batch, + func(ctx context.Context, limit int64) (int64, error) { + return q.PurgeRepoTrafficPathsBatch(ctx, pool, reposdb.PurgeRepoTrafficPathsBatchParams{ + Cutoff: pgtype.Date{Time: cutoffDay, Valid: true}, + BatchSize: limit, + }) + }) + res.PathsDeleted = deleted + res.Remaining = res.Remaining || more + if err != nil { + return res, fmt.Errorf("purge repo_traffic_paths: %w", err) + } + + deleted, more, err = purgeBatched(ctx, opts.MaxBatches, batch, + func(ctx context.Context, limit int64) (int64, error) { + return q.PurgeRepoTrafficReferrersBatch(ctx, pool, reposdb.PurgeRepoTrafficReferrersBatchParams{ + Cutoff: pgtype.Date{Time: cutoffDay, Valid: true}, + BatchSize: limit, + }) + }) + res.ReferrersDeleted = deleted + res.Remaining = res.Remaining || more + if err != nil { + return res, fmt.Errorf("purge repo_traffic_referrers: %w", err) + } + + deleted, more, err = purgeBatched(ctx, opts.MaxBatches, batch, + func(ctx context.Context, limit int64) (int64, error) { + return q.PurgeRepoTrafficDailyBatch(ctx, pool, reposdb.PurgeRepoTrafficDailyBatchParams{ + Cutoff: pgtype.Date{Time: dailyCutoffDay, Valid: true}, + BatchSize: limit, + }) + }) + res.DailyDeleted = deleted + res.Remaining = res.Remaining || more + if err != nil { + return res, fmt.Errorf("purge repo_traffic_daily: %w", err) + } + + return res, nil +} + +// purgeBatched calls del until it reports a short batch — meaning +// nothing older than the cutoff is left — or maxBatches statements have +// run. It returns the rows deleted, whether it stopped on the batch cap +// with work outstanding, and the first error. +// +// The short-batch check is what terminates the loop; without it a table +// whose row count is an exact multiple of the batch size would still +// exit, but only after one extra empty DELETE. That extra round trip is +// cheap and keeps the condition a single comparison. +func purgeBatched(ctx context.Context, maxBatches int, batch int64, del func(context.Context, int64) (int64, error)) (int64, bool, error) { + if batch <= 0 || maxBatches <= 0 { + return 0, false, nil + } + var total int64 + for i := 0; i < maxBatches; i++ { + if err := ctx.Err(); err != nil { + return total, true, err + } + n, err := del(ctx, batch) + total += n + if err != nil { + return total, true, err + } + if n < batch { + return total, false, nil + } + } + return total, true, nil +} diff --git a/internal/repos/traffic/purge_test.go b/internal/repos/traffic/purge_test.go new file mode 100644 index 00000000..223e679d --- /dev/null +++ b/internal/repos/traffic/purge_test.go @@ -0,0 +1,198 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package traffic + +import ( + "context" + "errors" + "testing" + "time" +) + +func TestPurgeBatchedStopsOnShortBatch(t *testing.T) { + // 12 rows at a batch size of 5: two full batches, then a short one. + remaining := int64(12) + var calls []int64 + total, more, err := purgeBatched(context.Background(), 100, 5, + func(_ context.Context, limit int64) (int64, error) { + calls = append(calls, limit) + n := min(limit, remaining) + remaining -= n + return n, nil + }) + if err != nil { + t.Fatalf("purgeBatched: %v", err) + } + if total != 12 { + t.Fatalf("total = %d, want 12", total) + } + if more { + t.Fatal("more = true, want false: the loop drained the table") + } + if len(calls) != 3 { + t.Fatalf("statements = %d (%v), want 3", len(calls), calls) + } + for i, limit := range calls { + if limit != 5 { + t.Fatalf("statement %d ran with limit %d, want the batch size 5", i, limit) + } + } +} + +// An exact multiple of the batch size costs one extra empty DELETE; the +// loop must still terminate rather than spin to MaxBatches. +func TestPurgeBatchedStopsOnExactMultiple(t *testing.T) { + remaining := int64(10) + calls := 0 + total, more, err := purgeBatched(context.Background(), 100, 5, + func(_ context.Context, limit int64) (int64, error) { + calls++ + n := min(limit, remaining) + remaining -= n + return n, nil + }) + if err != nil { + t.Fatalf("purgeBatched: %v", err) + } + if total != 10 || more { + t.Fatalf("total = %d, more = %v, want 10 and false", total, more) + } + if calls != 3 { + t.Fatalf("statements = %d, want 3 (two full batches plus the empty one)", calls) + } +} + +func TestPurgeBatchedHonoursMaxBatches(t *testing.T) { + calls := 0 + total, more, err := purgeBatched(context.Background(), 4, 5, + func(_ context.Context, limit int64) (int64, error) { + calls++ + return limit, nil // always a full batch: backlog never drains + }) + if err != nil { + t.Fatalf("purgeBatched: %v", err) + } + if calls != 4 { + t.Fatalf("statements = %d, want the 4 the cap allows", calls) + } + if total != 20 { + t.Fatalf("total = %d, want 20", total) + } + if !more { + t.Fatal("more = false, want true: the run stopped on the batch cap") + } +} + +func TestPurgeBatchedZeroBounds(t *testing.T) { + for _, tc := range []struct { + name string + maxBatches int + batch int64 + }{ + {"zero batch size", 100, 0}, + {"negative batch size", 100, -5}, + {"zero max batches", 0, 5}, + {"negative max batches", -1, 5}, + } { + t.Run(tc.name, func(t *testing.T) { + total, more, err := purgeBatched(context.Background(), tc.maxBatches, tc.batch, + func(context.Context, int64) (int64, error) { + t.Fatal("delete ran with a non-positive bound") + return 0, nil + }) + if err != nil || total != 0 || more { + t.Fatalf("total = %d, more = %v, err = %v; want 0, false, nil", total, more, err) + } + }) + } +} + +func TestPurgeBatchedReturnsPartialProgressOnError(t *testing.T) { + boom := errors.New("boom") + calls := 0 + total, more, err := purgeBatched(context.Background(), 100, 5, + func(_ context.Context, limit int64) (int64, error) { + calls++ + if calls == 2 { + return 2, boom + } + return limit, nil + }) + if !errors.Is(err, boom) { + t.Fatalf("err = %v, want boom", err) + } + // Rows the failing statement did delete still count: each DELETE is + // its own transaction, so a mid-loop failure does not roll them back. + if total != 7 { + t.Fatalf("total = %d, want 7 (5 from the first batch, 2 from the failed one)", total) + } + if !more { + t.Fatal("more = false, want true: the table was not drained") + } +} + +func TestPurgeBatchedStopsOnCanceledContext(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + total, more, err := purgeBatched(ctx, 100, 5, + func(context.Context, int64) (int64, error) { + t.Fatal("delete ran after the context was canceled") + return 0, nil + }) + if !errors.Is(err, context.Canceled) { + t.Fatalf("err = %v, want context.Canceled", err) + } + if total != 0 || !more { + t.Fatalf("total = %d, more = %v, want 0 and true", total, more) + } +} + +func TestPurgeOptionsNormalize(t *testing.T) { + got := PurgeOptions{}.normalize() + if got.RetentionDays != DefaultRetentionDays { + t.Fatalf("RetentionDays = %d, want %d", got.RetentionDays, DefaultRetentionDays) + } + if got.DailyRetentionDays != DefaultDailyRetentionDays { + t.Fatalf("DailyRetentionDays = %d, want %d", got.DailyRetentionDays, DefaultDailyRetentionDays) + } + if got.BatchSize != DefaultPurgeBatch { + t.Fatalf("BatchSize = %d, want %d", got.BatchSize, DefaultPurgeBatch) + } + if got.MaxBatches != DefaultMaxBatches { + t.Fatalf("MaxBatches = %d, want %d", got.MaxBatches, DefaultMaxBatches) + } + if got.Now == nil { + t.Fatal("Now = nil, want a default clock") + } + + // The detail window must never outlive the aggregate it rolls up to. + clamped := PurgeOptions{RetentionDays: 90, DailyRetentionDays: 10}.normalize() + if clamped.DailyRetentionDays != 90 { + t.Fatalf("DailyRetentionDays = %d, want it clamped up to 90", clamped.DailyRetentionDays) + } + + explicit := PurgeOptions{ + RetentionDays: 7, + DailyRetentionDays: 90, + BatchSize: 100, + MaxBatches: 3, + Now: func() time.Time { return time.Unix(0, 0) }, + }.normalize() + if explicit.RetentionDays != 7 || explicit.DailyRetentionDays != 90 || + explicit.BatchSize != 100 || explicit.MaxBatches != 3 { + t.Fatalf("normalize overwrote explicit options: %+v", explicit) + } +} + +// The retention window has to stay clear of the window the Traffic UI +// reads, or a purge silently truncates the leftmost bars of the chart. +func TestDefaultRetentionCoversTheReadWindow(t *testing.T) { + if DefaultRetentionDays <= DefaultWindowDays { + t.Fatalf("DefaultRetentionDays = %d, must exceed DefaultWindowDays = %d", + DefaultRetentionDays, DefaultWindowDays) + } + if DefaultDailyRetentionDays < DefaultRetentionDays { + t.Fatalf("DefaultDailyRetentionDays = %d, must be at least DefaultRetentionDays = %d", + DefaultDailyRetentionDays, DefaultRetentionDays) + } +} diff --git a/internal/worker/jobs/traffic_purge.go b/internal/worker/jobs/traffic_purge.go new file mode 100644 index 00000000..6b3dad96 --- /dev/null +++ b/internal/worker/jobs/traffic_purge.go @@ -0,0 +1,77 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package jobs + +import ( + "context" + "encoding/json" + "log/slog" + + "github.com/jackc/pgx/v5/pgxpool" + + repotraffic "github.com/tenseleyFlow/shithub/internal/repos/traffic" + "github.com/tenseleyFlow/shithub/internal/worker" +) + +// TrafficPurgeDeps wires the traffic retention sweep. +type TrafficPurgeDeps struct { + Pool *pgxpool.Pool + Logger *slog.Logger +} + +// TrafficPurgePayload overrides the defaults for one run. An empty +// object (what cron sends) selects the package defaults; operators can +// pass a shorter window through `shithubd admin run-job` to drain a +// backlog faster. +type TrafficPurgePayload struct { + RetentionDays int `json:"retention_days,omitempty"` + DailyRetentionDays int `json:"daily_retention_days,omitempty"` + BatchSize int `json:"batch_size,omitempty"` + MaxBatches int `json:"max_batches,omitempty"` +} + +// TrafficPurge trims repo_traffic_uniques / _paths / _referrers to a +// 30-day window and repo_traffic_daily to 400 days, in bounded batches. +// Cron-driven and idempotent: the cutoff is recomputed from the clock +// on every run and rows inside the window are never touched, so a +// re-run after a partial pass simply picks up where it stopped. +func TrafficPurge(deps TrafficPurgeDeps) worker.Handler { + return func(ctx context.Context, raw json.RawMessage) error { + var p TrafficPurgePayload + if len(raw) > 0 { + if err := json.Unmarshal(raw, &p); err != nil { + // A malformed payload will never parse on a retry. + return worker.PoisonError(err) + } + } + res, err := repotraffic.Purge(ctx, deps.Pool, repotraffic.PurgeOptions{ + RetentionDays: p.RetentionDays, + DailyRetentionDays: p.DailyRetentionDays, + BatchSize: p.BatchSize, + MaxBatches: p.MaxBatches, + }) + if deps.Logger != nil { + // Logged even on error: the counters describe what the run + // did manage to delete before it stopped. + deps.Logger.InfoContext(ctx, "traffic:purge", + "uniques_deleted", res.UniquesDeleted, + "paths_deleted", res.PathsDeleted, + "referrers_deleted", res.ReferrersDeleted, + "daily_deleted", res.DailyDeleted, + "remaining", res.Remaining) + } + if err != nil { + return err + } + // A run that stopped on the batch cap has rows left over. Kick + // another pass now rather than waiting for tomorrow's cron beat, + // so the first run after deploy drains the backlog in one night. + if res.Remaining { + if _, err := worker.Enqueue(ctx, deps.Pool, repotraffic.KindTrafficPurge, + p, worker.EnqueueOptions{}); err != nil && deps.Logger != nil { + deps.Logger.WarnContext(ctx, "traffic:purge self-enqueue failed", "error", err) + } + } + return nil + } +} diff --git a/internal/worker/jobs/traffic_purge_test.go b/internal/worker/jobs/traffic_purge_test.go new file mode 100644 index 00000000..ef6a7a1a --- /dev/null +++ b/internal/worker/jobs/traffic_purge_test.go @@ -0,0 +1,226 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package jobs_test + +import ( + "context" + "fmt" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgtype" + "github.com/jackc/pgx/v5/pgxpool" + + reposdb "github.com/tenseleyFlow/shithub/internal/repos/sqlc" + repotraffic "github.com/tenseleyFlow/shithub/internal/repos/traffic" + "github.com/tenseleyFlow/shithub/internal/testing/dbtest" + usersdb "github.com/tenseleyFlow/shithub/internal/users/sqlc" + "github.com/tenseleyFlow/shithub/internal/worker/jobs" +) + +const trafficPurgeFixtureHash = "$argon2id$v=19$m=16384,t=1,p=1$" + + "AAAAAAAAAAAAAAAA$" + + "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA" + +// The retention cutoff is `day < today - RetentionDays`, so a row landing +// exactly on the cutoff day survives. Seed one row on each side of it plus +// one well inside the window, and assert only the older-than-cutoff row goes. +func TestTrafficPurgeKeepsTheRetentionWindow(t *testing.T) { + ctx := context.Background() + pool := dbtest.NewTestDB(t) + repoID := insertTrafficPurgeRepo(t, ctx, pool, "traffic-purge-window") + + today := startOfUTCDay(time.Now()) + inWindow := today.AddDate(0, 0, -5) + onCutoff := today.AddDate(0, 0, -repotraffic.DefaultRetentionDays) + pastCutoff := today.AddDate(0, 0, -repotraffic.DefaultRetentionDays-1) + // repo_traffic_daily keeps a much longer window of its own. + dailyInWindow := today.AddDate(0, 0, -100) + dailyPastCutoff := today.AddDate(0, 0, -repotraffic.DefaultDailyRetentionDays-1) + + for _, day := range []time.Time{inWindow, onCutoff, pastCutoff} { + insertTrafficPath(t, ctx, pool, repoID, day, "/README.md") + insertTrafficReferrer(t, ctx, pool, repoID, day, "example.com") + insertTrafficUnique(t, ctx, pool, repoID, day, "view", "", 0x01) + } + for _, day := range []time.Time{inWindow, dailyInWindow, dailyPastCutoff} { + insertTrafficDaily(t, ctx, pool, repoID, day) + } + + handler := jobs.TrafficPurge(jobs.TrafficPurgeDeps{Pool: pool}) + if err := handler(ctx, nil); err != nil { + t.Fatalf("traffic:purge handler: %v", err) + } + + assertTrafficDays(t, ctx, pool, "repo_traffic_paths", []time.Time{onCutoff, inWindow}) + assertTrafficDays(t, ctx, pool, "repo_traffic_referrers", []time.Time{onCutoff, inWindow}) + assertTrafficDays(t, ctx, pool, "repo_traffic_uniques", []time.Time{onCutoff, inWindow}) + assertTrafficDays(t, ctx, pool, "repo_traffic_daily", []time.Time{dailyInWindow, inWindow}) + + // Idempotent: a second run over the already-trimmed tables is a no-op. + if err := handler(ctx, nil); err != nil { + t.Fatalf("traffic:purge second run: %v", err) + } + assertTrafficDays(t, ctx, pool, "repo_traffic_paths", []time.Time{onCutoff, inWindow}) + assertTrafficDays(t, ctx, pool, "repo_traffic_daily", []time.Time{dailyInWindow, inWindow}) +} + +// A single run must not delete more than BatchSize × MaxBatches rows per +// table, so a first pass over a multi-million-row backlog stays bounded. +func TestTrafficPurgeBoundsDeletesPerRun(t *testing.T) { + ctx := context.Background() + pool := dbtest.NewTestDB(t) + repoID := insertTrafficPurgeRepo(t, ctx, pool, "traffic-purge-bounded") + + today := startOfUTCDay(time.Now()) + old := today.AddDate(0, 0, -90) + const seeded = 25 + for i := range seeded { + insertTrafficPath(t, ctx, pool, repoID, old, fmt.Sprintf("/blob/%d", i)) + } + + res, err := repotraffic.Purge(ctx, pool, repotraffic.PurgeOptions{ + BatchSize: 4, + MaxBatches: 3, + }) + if err != nil { + t.Fatalf("Purge: %v", err) + } + if res.PathsDeleted != 12 { + t.Fatalf("PathsDeleted = %d, want 12 (3 batches of 4)", res.PathsDeleted) + } + if !res.Remaining { + t.Fatal("Remaining = false, want true: the run stopped on the batch cap") + } + if got := countTrafficRows(t, ctx, pool, "repo_traffic_paths"); got != seeded-12 { + t.Fatalf("rows left = %d, want %d", got, seeded-12) + } + + // The next pass picks up where this one stopped; the job is safe to + // re-run until the backlog drains. + res, err = repotraffic.Purge(ctx, pool, repotraffic.PurgeOptions{ + BatchSize: 4, + MaxBatches: 3, + }) + if err != nil { + t.Fatalf("Purge second run: %v", err) + } + if res.PathsDeleted != 12 { + t.Fatalf("second run PathsDeleted = %d, want 12", res.PathsDeleted) + } + if got := countTrafficRows(t, ctx, pool, "repo_traffic_paths"); got != 1 { + t.Fatalf("rows left after two runs = %d, want 1", got) + } +} + +func startOfUTCDay(t time.Time) time.Time { + y, m, d := t.UTC().Date() + return time.Date(y, m, d, 0, 0, 0, 0, time.UTC) +} + +func insertTrafficPurgeRepo(t *testing.T, ctx context.Context, pool *pgxpool.Pool, name string) int64 { + t.Helper() + user, err := usersdb.New().CreateUser(ctx, pool, usersdb.CreateUserParams{ + Username: name, + DisplayName: name, + PasswordHash: trafficPurgeFixtureHash, + }) + if err != nil { + t.Fatalf("CreateUser: %v", err) + } + repo, err := reposdb.New().CreateRepo(ctx, pool, reposdb.CreateRepoParams{ + OwnerUserID: pgtype.Int8{Int64: user.ID, Valid: true}, + Name: name, + DefaultBranch: "trunk", + Visibility: reposdb.RepoVisibilityPublic, + }) + if err != nil { + t.Fatalf("CreateRepo: %v", err) + } + return repo.ID +} + +func insertTrafficPath(t *testing.T, ctx context.Context, pool *pgxpool.Pool, repoID int64, day time.Time, path string) { + t.Helper() + if _, err := pool.Exec(ctx, ` +INSERT INTO repo_traffic_paths (repo_id, day, path, views, unique_views, created_at) +VALUES ($1, $2::date, $3, 1, 1, $4)`, repoID, day.UTC().Format(time.DateOnly), path, day.UTC()); err != nil { + t.Fatalf("insert repo_traffic_paths(%s, %s): %v", day.Format(time.DateOnly), path, err) + } +} + +func insertTrafficReferrer(t *testing.T, ctx context.Context, pool *pgxpool.Pool, repoID int64, day time.Time, referrer string) { + t.Helper() + if _, err := pool.Exec(ctx, ` +INSERT INTO repo_traffic_referrers (repo_id, day, referrer, views, unique_views, created_at) +VALUES ($1, $2::date, $3, 1, 1, $4)`, repoID, day.UTC().Format(time.DateOnly), referrer, day.UTC()); err != nil { + t.Fatalf("insert repo_traffic_referrers(%s): %v", day.Format(time.DateOnly), err) + } +} + +// created_at matters here: repo_traffic_uniques is purged on created_at, +// which the request path stamps on the same day it derives `day` from. +func insertTrafficUnique(t *testing.T, ctx context.Context, pool *pgxpool.Pool, repoID int64, day time.Time, metric, key string, seed byte) { + t.Helper() + hash := make([]byte, 32) + for i := range hash { + hash[i] = seed + } + if _, err := pool.Exec(ctx, ` +INSERT INTO repo_traffic_uniques (repo_id, day, metric, key, visitor_hash, created_at) +VALUES ($1, $2::date, $3, $4, $5, $6)`, repoID, day.UTC().Format(time.DateOnly), metric, key, hash, day.UTC()); err != nil { + t.Fatalf("insert repo_traffic_uniques(%s): %v", day.Format(time.DateOnly), err) + } +} + +func insertTrafficDaily(t *testing.T, ctx context.Context, pool *pgxpool.Pool, repoID int64, day time.Time) { + t.Helper() + if _, err := pool.Exec(ctx, ` +INSERT INTO repo_traffic_daily (repo_id, day, views, unique_views, clones, unique_clones, created_at) +VALUES ($1, $2::date, 1, 1, 0, 0, $3)`, repoID, day.UTC().Format(time.DateOnly), day.UTC()); err != nil { + t.Fatalf("insert repo_traffic_daily(%s): %v", day.Format(time.DateOnly), err) + } +} + +// assertTrafficDays checks the surviving rows of a table are exactly the +// supplied days, ordered ascending. +func assertTrafficDays(t *testing.T, ctx context.Context, pool *pgxpool.Pool, table string, want []time.Time) { + t.Helper() + rows, err := pool.Query(ctx, "SELECT day FROM "+table+" ORDER BY day ASC") + if err != nil { + t.Fatalf("select %s: %v", table, err) + } + defer rows.Close() + var got []string + for rows.Next() { + var day pgtype.Date + if err := rows.Scan(&day); err != nil { + t.Fatalf("scan %s: %v", table, err) + } + got = append(got, day.Time.UTC().Format(time.DateOnly)) + } + if err := rows.Err(); err != nil { + t.Fatalf("iterate %s: %v", table, err) + } + wantDays := make([]string, 0, len(want)) + for _, d := range want { + wantDays = append(wantDays, d.UTC().Format(time.DateOnly)) + } + if len(got) != len(wantDays) { + t.Fatalf("%s survivors = %v, want %v", table, got, wantDays) + } + for i := range got { + if got[i] != wantDays[i] { + t.Fatalf("%s survivors = %v, want %v", table, got, wantDays) + } + } +} + +func countTrafficRows(t *testing.T, ctx context.Context, pool *pgxpool.Pool, table string) int { + t.Helper() + var n int + if err := pool.QueryRow(ctx, "SELECT count(*) FROM "+table).Scan(&n); err != nil { + t.Fatalf("count %s: %v", table, err) + } + return n +}