internal/store/retention.go

8dcfa45a8ac03a5ff9c36828274d05acadcf846c
gitbay/internal/store/retention.go history · blame · raw

109 lines · 3911 bytes

  1package store
  2
  3import (
  4	"fmt"
  5	"time"
  6)
  7
  8// Nothing was ever deleted from this database. Expired sessions and
  9// tokens were filtered on read and left in place; the audit log, the
 10// activity feed, webhook deliveries and the mail queue grew without
 11// bound. Sweep removes them (#122).
 12
 13// Swept counts what one sweep removed, per table. Zero-valued entries are
 14// left in so a caller logging the result sees every table it asked about.
 15type Swept map[string]int64
 16
 17// Total is how many rows the sweep removed altogether.
 18func (s Swept) Total() int64 {
 19	var n int64
 20	for _, v := range s {
 21		n += v
 22	}
 23	return n
 24}
 25
 26// Retention says how long each capped table keeps a row. A zero duration
 27// means keep forever, which is what an instance that has not configured
 28// retention gets: growing is a decision, but so is deleting an audit
 29// trail, and the second one is not made on an operator's behalf.
 30type Retention struct {
 31	Audit             time.Duration
 32	Events            time.Duration
 33	WebhookDeliveries time.Duration
 34	Mail              time.Duration
 35	Push              time.Duration
 36}
 37
 38// Sweep deletes expired sessions and tokens, then the rows older than
 39// each configured retention. Errors are returned with whatever was
 40// removed before them: a sweep that fails halfway has still done that
 41// much, and the next one picks up the rest.
 42func (s *Store) Sweep(r Retention, now time.Time) (Swept, error) {
 43	out := Swept{}
 44	// Dead the moment they expire, whatever retention says. A used login
 45	// token cannot be replayed and an expired session cannot authenticate,
 46	// so neither is evidence of anything.
 47	expired := []struct {
 48		table string
 49		where string
 50	}{
 51		{"web_sessions", "expires_at <= ?"},
 52		{"login_tokens", "expires_at <= ?"},
 53		{"email_tokens", "expires_at <= ?"},
 54		{"push_tokens", "expires_at <= ?"},
 55	}
 56	for _, e := range expired {
 57		n, err := s.deleteBy("DELETE FROM "+e.table+" WHERE "+e.where, fmtTime(now))
 58		out[e.table] += n
 59		if err != nil {
 60			return out, fmt.Errorf("sweeping %s: %w", e.table, err)
 61		}
 62	}
 63
 64	// Order matters: deliveries before events. webhook_deliveries.event_id
 65	// is ON DELETE CASCADE, so an event taken out from under a delivery
 66	// takes the delivery with it — including one still queued for retry.
 67	// Sweeping finished deliveries first, and skipping any event that
 68	// still has an unfinished one, keeps that from happening. The effect
 69	// is that a delivery is kept for the shorter of the two retentions,
 70	// which is the honest reading of "keep deliveries for N".
 71	aged := []struct {
 72		table string
 73		where string
 74		keep  time.Duration
 75	}{
 76		// By id, so a backwards clock step cannot leave a newer row
 77		// removed and an older one kept: the audit chain would read
 78		// that as tampering.
 79		{"audit_log", "id <= (SELECT MAX(id) FROM audit_log WHERE created_at < ?)", r.Audit},
 80		// Only deliveries that have finished: one still being retried is
 81		// live state, however old its first attempt.
 82		{"webhook_deliveries", "created_at < ? AND (delivered_at IS NOT NULL OR failed_at IS NOT NULL)", r.WebhookDeliveries},
 83		{"events", `created_at < ? AND NOT EXISTS (
 84			SELECT 1 FROM webhook_deliveries d
 85			WHERE d.event_id = events.id AND d.delivered_at IS NULL AND d.failed_at IS NULL)`, r.Events},
 86		{"notifications", "created_at < ? AND (sent_at IS NOT NULL OR failed_at IS NOT NULL)", r.Mail},
 87		{"push_queue", "created_at < ? AND (sent_at IS NOT NULL OR failed_at IS NOT NULL)", r.Push},
 88	}
 89	for _, a := range aged {
 90		if a.keep <= 0 {
 91			continue
 92		}
 93		n, err := s.deleteBy("DELETE FROM "+a.table+" WHERE "+a.where, fmtTime(now.Add(-a.keep)))
 94		out[a.table] += n
 95		if err != nil {
 96			return out, fmt.Errorf("sweeping %s: %w", a.table, err)
 97		}
 98	}
 99	return out, nil
100}
101
102func (s *Store) deleteBy(q string, args ...any) (int64, error) {
103	res, err := s.DB.Exec(q, args...)
104	if err != nil {
105		return 0, err
106	}
107	n, _ := res.RowsAffected()
108	return n, nil
109}