Commit d51efc8bb8
Verified · cmc ci/build: success ci/test: success ci/vuln: success
Layout: unified · split
cmd/gitbayd/main.go +31
| @@ -171,6 +171,7 @@ func serveCmd() *cobra.Command { | |||
| 171 | if d := cfg.Registration.PendingExpiryDuration(); d > 0 { | 171 | if d := cfg.Registration.PendingExpiryDuration(); d > 0 { |
| 172 | go reapPending(whCtx, st, d) | 172 | go reapPending(whCtx, st, d) |
| 173 | } | 173 | } |
| 174 | go sweep(whCtx, st, cfg) | ||
| 174 | go (&ci.Scheduler{St: st, SiteURL: cfg.Server.SiteURL, | 175 | go (&ci.Scheduler{St: st, SiteURL: cfg.Server.SiteURL, |
| 175 | RepoDir: func(owner, name string) string { | 176 | RepoDir: func(owner, name string) string { |
| 176 | return control.RepoDir(cfg.Server.Root, owner, name) | 177 | return control.RepoDir(cfg.Server.Root, owner, name) |
| @@ -494,6 +495,36 @@ func hostUserCreateCmd() *cobra.Command { | |||
| 494 | } | 495 | } |
| 495 | } | 496 | } |
| 496 | 497 | ||
| 498 | // sweep prunes expired sessions and tokens, and rows past their | ||
| 499 | // configured retention, hourly and once at start. GITBAY_SWEEP_TICK | ||
| 500 | // shortens the interval for tests. | ||
| 501 | func sweep(ctx context.Context, st *store.Store, cfg config.Config) { | ||
| 502 | tick := time.Hour | ||
| 503 | if v := os.Getenv("GITBAY_SWEEP_TICK"); v != "" { | ||
| 504 | if d, err := time.ParseDuration(v); err == nil { | ||
| 505 | tick = d | ||
| 506 | } | ||
| 507 | } | ||
| 508 | audit, events, deliveries, mail := cfg.Retention.Durations() | ||
| 509 | r := store.Retention{Audit: audit, Events: events, | ||
| 510 | WebhookDeliveries: deliveries, Mail: mail} | ||
| 511 | t := time.NewTicker(tick) | ||
| 512 | defer t.Stop() | ||
| 513 | for { | ||
| 514 | swept, err := st.Sweep(r, time.Now()) | ||
| 515 | if err != nil { | ||
| 516 | slog.Error("sweeping", "err", err, "removed", swept.Total()) | ||
| 517 | } else if n := swept.Total(); n > 0 { | ||
| 518 | slog.Info("swept expired rows", "removed", n, "tables", swept) | ||
| 519 | } | ||
| 520 | select { | ||
| 521 | case <-ctx.Done(): | ||
| 522 | return | ||
| 523 | case <-t.C: | ||
| 524 | } | ||
| 525 | } | ||
| 526 | } | ||
| 527 | |||
| 497 | // reapPending removes self-registered accounts still unverified after | 528 | // reapPending removes self-registered accounts still unverified after |
| 498 | // maxAge, hourly and once at start. GITBAY_REAP_TICK shortens the | 529 | // maxAge, hourly and once at start. GITBAY_REAP_TICK shortens the |
| 499 | // interval for tests. | 530 | // interval for tests. |
deploy/cloud-init.yaml +9
| @@ -190,6 +190,15 @@ write_files: | |||
| 190 | [registration] | 190 | [registration] |
| 191 | mode = "closed" | 191 | mode = "closed" |
| 192 | 192 | ||
| 193 | # How long the append-only tables keep a row. Unset means forever, | ||
| 194 | # which is the default: growing is a decision, but so is deleting an | ||
| 195 | # audit trail. Expired sessions and tokens are swept either way. | ||
| 196 | # [retention] | ||
| 197 | # audit = "8760h" # a year | ||
| 198 | # events = "4380h" # six months | ||
| 199 | # webhook_deliveries = "720h" # a month | ||
| 200 | # mail = "720h" | ||
| 201 | |||
| 193 | - path: /etc/systemd/system/gitbayd.service | 202 | - path: /etc/systemd/system/gitbayd.service |
| 194 | content: | | 203 | content: | |
| 195 | [Unit] | 204 | [Unit] |
e2e/audit_test.go +19
| @@ -45,6 +45,25 @@ func TestAuditAndHardening(t *testing.T) { | |||
| 45 | t.Fatalf("host audit: %s", out) | 45 | t.Fatalf("host audit: %s", out) |
| 46 | } | 46 | } |
| 47 | 47 | ||
| 48 | // Prose reaches argv through --title and --body. The entry records | ||
| 49 | // that the flags were given, not what was written: the issue itself is | ||
| 50 | // the record of its own text, and the audit log is not pruned by | ||
| 51 | // default (#122). | ||
| 52 | if _, errOut, code := inst.ssh(t, aliceKey, "", "issue", "create", "alice/app", | ||
| 53 | "--title", "'a short title'", "--body", "'prose that must not be copied'"); code != 0 { | ||
| 54 | t.Fatalf("issue create: %s", errOut) | ||
| 55 | } | ||
| 56 | out, _, code = inst.ssh(t, adminKey, "", "audit", "--json") | ||
| 57 | if code != 0 || !strings.Contains(out, "cmd issue create") { | ||
| 58 | t.Fatalf("issue create not audited: %s", out) | ||
| 59 | } | ||
| 60 | if strings.Contains(out, "prose that must not be copied") || strings.Contains(out, "a short title") { | ||
| 61 | t.Fatalf("audit log copied the issue text:\n%s", out) | ||
| 62 | } | ||
| 63 | if !strings.Contains(out, "--body") || !strings.Contains(out, "alice/app") { | ||
| 64 | t.Fatalf("audit log dropped the flag names or the target:\n%s", out) | ||
| 65 | } | ||
| 66 | |||
| 48 | // Disable: everything refused, sessions dropped, nothing deleted. | 67 | // Disable: everything refused, sessions dropped, nothing deleted. |
| 49 | inst.admin(t, "admin", "user", "disable", "bob") | 68 | inst.admin(t, "admin", "user", "disable", "bob") |
| 50 | if _, errOut, code := inst.ssh(t, bobKey, "", "whoami"); code != 4 || !strings.Contains(errOut, "disabled") { | 69 | if _, errOut, code := inst.ssh(t, bobKey, "", "whoami"); code != 4 || !strings.Contains(errOut, "disabled") { |
internal/config/config.go +26
| @@ -28,6 +28,7 @@ type Config struct { | |||
| 28 | Mail Mail `toml:"mail"` | 28 | Mail Mail `toml:"mail"` |
| 29 | Mirrors Mirrors `toml:"mirrors"` | 29 | Mirrors Mirrors `toml:"mirrors"` |
| 30 | Deps Deps `toml:"deps"` | 30 | Deps Deps `toml:"deps"` |
| 31 | Retention Retention `toml:"retention"` | ||
| 31 | // GoImport maps vanity Go module paths to repositories, e.g. | 32 | // GoImport maps vanity Go module paths to repositories, e.g. |
| 32 | // "gitbay.org/gitbay" = "krz/gitbay". Requests carrying ?go-get=1 | 33 | // "gitbay.org/gitbay" = "krz/gitbay". Requests carrying ?go-get=1 |
| 33 | // under a mapped path get a go-import meta tag. | 34 | // under a mapped path get a go-import meta tag. |
| @@ -101,6 +102,31 @@ func (r Registration) PendingExpiryDuration() time.Duration { | |||
| 101 | return d | 102 | return d |
| 102 | } | 103 | } |
| 103 | 104 | ||
| 105 | // Retention is how long the append-only tables keep a row. Each is a | ||
| 106 | // duration string ("2160h"); empty or zero keeps forever, which is what | ||
| 107 | // an instance that has never configured this gets. Expired sessions and | ||
| 108 | // tokens are swept regardless: they are dead weight the moment they | ||
| 109 | // expire and no setting makes them worth keeping. | ||
| 110 | type Retention struct { | ||
| 111 | Audit string `toml:"audit"` | ||
| 112 | Events string `toml:"events"` | ||
| 113 | WebhookDeliveries string `toml:"webhook_deliveries"` | ||
| 114 | // Mail is the outbound queue: rows already sent or given up on. | ||
| 115 | Mail string `toml:"mail"` | ||
| 116 | } | ||
| 117 | |||
| 118 | // Durations parses the four, mapping each to zero when unset or bad. | ||
| 119 | func (r Retention) Durations() (audit, events, deliveries, mail time.Duration) { | ||
| 120 | parse := func(s string) time.Duration { | ||
| 121 | d, err := time.ParseDuration(s) | ||
| 122 | if err != nil || d < 0 { | ||
| 123 | return 0 | ||
| 124 | } | ||
| 125 | return d | ||
| 126 | } | ||
| 127 | return parse(r.Audit), parse(r.Events), parse(r.WebhookDeliveries), parse(r.Mail) | ||
| 128 | } | ||
| 129 | |||
| 104 | // LFS stores large-file objects content-addressed under Root (default | 130 | // LFS stores large-file objects content-addressed under Root (default |
| 105 | // <server.root>/lfs). MaxObjectBytes caps a single object; 0 means the | 131 | // <server.root>/lfs). MaxObjectBytes caps a single object; 0 means the |
| 106 | // 512MB default. | 132 | // 512MB default. |
internal/control/audit_test.go +37
| @@ -26,3 +26,40 @@ func TestParseSince(t *testing.T) { | |||
| 26 | } | 26 | } |
| 27 | } | 27 | } |
| 28 | } | 28 | } |
| 29 | |||
| 30 | // TestAuditArgsDropFlagValues: prose reaches argv through --body and | ||
| 31 | // --message, and the audit log kept it verbatim for a repository that may | ||
| 32 | // be private. Identifiers are positional and survive; flag names survive | ||
| 33 | // so the entry still says what shape the command had (#122). | ||
| 34 | func TestAuditArgsDropFlagValues(t *testing.T) { | ||
| 35 | cases := []struct { | ||
| 36 | in []string | ||
| 37 | want []string | ||
| 38 | }{ | ||
| 39 | {[]string{"krz/gitbay", "--title", "a bug", "--body", "the whole issue text"}, | ||
| 40 | []string{"krz/gitbay", "--title", "--body"}}, | ||
| 41 | {[]string{"krz/gitbay", "7", "--add", "bug", "--add", "ui"}, | ||
| 42 | []string{"krz/gitbay", "7", "--add", "--add"}}, | ||
| 43 | // A switch takes no value, so the next argument is not swallowed. | ||
| 44 | {[]string{"krz/gitbay", "--private", "--name", "x"}, | ||
| 45 | []string{"krz/gitbay", "--private", "--name"}}, | ||
| 46 | // After "--" everything is positional, including a token that | ||
| 47 | // looks like a flag. | ||
| 48 | {[]string{"krz/gitbay", "--", "--not-a-flag"}, | ||
| 49 | []string{"krz/gitbay", "--", "--not-a-flag"}}, | ||
| 50 | {nil, []string{}}, | ||
| 51 | } | ||
| 52 | for _, tc := range cases { | ||
| 53 | got := auditArgs(tc.in) | ||
| 54 | if len(got) != len(tc.want) { | ||
| 55 | t.Errorf("auditArgs(%v) = %v, want %v", tc.in, got, tc.want) | ||
| 56 | continue | ||
| 57 | } | ||
| 58 | for i := range got { | ||
| 59 | if got[i] != tc.want[i] { | ||
| 60 | t.Errorf("auditArgs(%v) = %v, want %v", tc.in, got, tc.want) | ||
| 61 | break | ||
| 62 | } | ||
| 63 | } | ||
| 64 | } | ||
| 65 | } | ||
internal/control/control.go +32 −4
| @@ -124,18 +124,46 @@ func Dispatch(c *Ctx, argv []string) int { | |||
| 124 | c.Stdin = emptyReader{} | 124 | c.Stdin = emptyReader{} |
| 125 | } | 125 | } |
| 126 | code := cmd.Run(c, args) | 126 | code := cmd.Run(c, args) |
| 127 | // Every successful mutating command lands in the audit log. Argv is | 127 | // Every successful mutating command lands in the audit log. |
| 128 | // safe to record by construction: secrets travel on stdin, never as | ||
| 129 | // arguments. | ||
| 130 | if code == protocol.ExitOK && !cmd.ReadOnly { | 128 | if code == protocol.ExitOK && !cmd.ReadOnly { |
| 131 | c.Store.Audit(c.User.ID, "cmd "+joinPath(cmd.Path), map[string]any{ | 129 | c.Store.Audit(c.User.ID, "cmd "+joinPath(cmd.Path), map[string]any{ |
| 132 | "argv": args, | 130 | "argv": auditArgs(args), |
| 133 | "source": c.Source, | 131 | "source": c.Source, |
| 134 | }) | 132 | }) |
| 135 | } | 133 | } |
| 136 | return code | 134 | return code |
| 137 | } | 135 | } |
| 138 | 136 | ||
| 137 | // auditArgs is argv with flag values dropped. Secrets never reach argv — | ||
| 138 | // they travel on stdin — but prose does: `issue create a/b --title x | ||
| 139 | // --body <the whole issue>` used to store the body verbatim, in a table | ||
| 140 | // nothing pruned, for a repository that may be private. The identifiers | ||
| 141 | // are positional, so keeping those and the flag names says what was done | ||
| 142 | // without copying what was written (#122). | ||
| 143 | func auditArgs(args []string) []string { | ||
| 144 | out := make([]string, 0, len(args)) | ||
| 145 | for i := 0; i < len(args); i++ { | ||
| 146 | a := args[i] | ||
| 147 | if !strings.HasPrefix(a, "--") { | ||
| 148 | out = append(out, a) | ||
| 149 | continue | ||
| 150 | } | ||
| 151 | out = append(out, a) | ||
| 152 | // "--" ends flag parsing; everything after it is positional. | ||
| 153 | if a == "--" { | ||
| 154 | out = append(out, args[i+1:]...) | ||
| 155 | break | ||
| 156 | } | ||
| 157 | // A flag's value is the next argument unless that is itself a | ||
| 158 | // flag, which is how a switch is told from one that takes a value | ||
| 159 | // without consulting the command's spec. | ||
| 160 | if i+1 < len(args) && !strings.HasPrefix(args[i+1], "--") { | ||
| 161 | i++ | ||
| 162 | } | ||
| 163 | } | ||
| 164 | return out | ||
| 165 | } | ||
| 166 | |||
| 139 | // pendingAllowed lists what an unverified self-registered account may do. | 167 | // pendingAllowed lists what an unverified self-registered account may do. |
| 140 | func pendingAllowed(path []string) bool { | 168 | func pendingAllowed(path []string) bool { |
| 141 | key := joinPath(path) | 169 | key := joinPath(path) |
internal/store/retention.go added +103
| @@ -0,0 +1,103 @@ | |||
| 1 | package store | ||
| 2 | |||
| 3 | import ( | ||
| 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. | ||
| 15 | type Swept map[string]int64 | ||
| 16 | |||
| 17 | // Total is how many rows the sweep removed altogether. | ||
| 18 | func (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. | ||
| 30 | type Retention struct { | ||
| 31 | Audit time.Duration | ||
| 32 | Events time.Duration | ||
| 33 | WebhookDeliveries time.Duration | ||
| 34 | Mail time.Duration | ||
| 35 | } | ||
| 36 | |||
| 37 | // Sweep deletes expired sessions and tokens, then the rows older than | ||
| 38 | // each configured retention. Errors are returned with whatever was | ||
| 39 | // removed before them: a sweep that fails halfway has still done that | ||
| 40 | // much, and the next one picks up the rest. | ||
| 41 | func (s *Store) Sweep(r Retention, now time.Time) (Swept, error) { | ||
| 42 | out := Swept{} | ||
| 43 | // Dead the moment they expire, whatever retention says. A used login | ||
| 44 | // token cannot be replayed and an expired session cannot authenticate, | ||
| 45 | // so neither is evidence of anything. | ||
| 46 | expired := []struct { | ||
| 47 | table string | ||
| 48 | where string | ||
| 49 | }{ | ||
| 50 | {"web_sessions", "expires_at <= ?"}, | ||
| 51 | {"login_tokens", "expires_at <= ?"}, | ||
| 52 | {"email_tokens", "expires_at <= ?"}, | ||
| 53 | } | ||
| 54 | for _, e := range expired { | ||
| 55 | n, err := s.deleteBy("DELETE FROM "+e.table+" WHERE "+e.where, fmtTime(now)) | ||
| 56 | out[e.table] += n | ||
| 57 | if err != nil { | ||
| 58 | return out, fmt.Errorf("sweeping %s: %w", e.table, err) | ||
| 59 | } | ||
| 60 | } | ||
| 61 | |||
| 62 | // Order matters: deliveries before events. webhook_deliveries.event_id | ||
| 63 | // is ON DELETE CASCADE, so an event taken out from under a delivery | ||
| 64 | // takes the delivery with it — including one still queued for retry. | ||
| 65 | // Sweeping finished deliveries first, and skipping any event that | ||
| 66 | // still has an unfinished one, keeps that from happening. The effect | ||
| 67 | // is that a delivery is kept for the shorter of the two retentions, | ||
| 68 | // which is the honest reading of "keep deliveries for N". | ||
| 69 | aged := []struct { | ||
| 70 | table string | ||
| 71 | where string | ||
| 72 | keep time.Duration | ||
| 73 | }{ | ||
| 74 | {"audit_log", "created_at < ?", r.Audit}, | ||
| 75 | // Only deliveries that have finished: one still being retried is | ||
| 76 | // live state, however old its first attempt. | ||
| 77 | {"webhook_deliveries", "created_at < ? AND (delivered_at IS NOT NULL OR failed_at IS NOT NULL)", r.WebhookDeliveries}, | ||
| 78 | {"events", `created_at < ? AND NOT EXISTS ( | ||
| 79 | SELECT 1 FROM webhook_deliveries d | ||
| 80 | WHERE d.event_id = events.id AND d.delivered_at IS NULL AND d.failed_at IS NULL)`, r.Events}, | ||
| 81 | {"notifications", "created_at < ? AND (sent_at IS NOT NULL OR failed_at IS NOT NULL)", r.Mail}, | ||
| 82 | } | ||
| 83 | for _, a := range aged { | ||
| 84 | if a.keep <= 0 { | ||
| 85 | continue | ||
| 86 | } | ||
| 87 | n, err := s.deleteBy("DELETE FROM "+a.table+" WHERE "+a.where, fmtTime(now.Add(-a.keep))) | ||
| 88 | out[a.table] += n | ||
| 89 | if err != nil { | ||
| 90 | return out, fmt.Errorf("sweeping %s: %w", a.table, err) | ||
| 91 | } | ||
| 92 | } | ||
| 93 | return out, nil | ||
| 94 | } | ||
| 95 | |||
| 96 | func (s *Store) deleteBy(q string, args ...any) (int64, error) { | ||
| 97 | res, err := s.DB.Exec(q, args...) | ||
| 98 | if err != nil { | ||
| 99 | return 0, err | ||
| 100 | } | ||
| 101 | n, _ := res.RowsAffected() | ||
| 102 | return n, nil | ||
| 103 | } | ||
internal/store/retention_test.go added +151
| @@ -0,0 +1,151 @@ | |||
| 1 | package store | ||
| 2 | |||
| 3 | import ( | ||
| 4 | "testing" | ||
| 5 | "time" | ||
| 6 | ) | ||
| 7 | |||
| 8 | func retentionFixture(t *testing.T) (*Store, int64) { | ||
| 9 | t.Helper() | ||
| 10 | s := open(t) | ||
| 11 | if err := s.MigrateUp(); err != nil { | ||
| 12 | t.Fatal(err) | ||
| 13 | } | ||
| 14 | uid, err := s.CreateUser("cmc", true) | ||
| 15 | if err != nil { | ||
| 16 | t.Fatal(err) | ||
| 17 | } | ||
| 18 | return s, uid | ||
| 19 | } | ||
| 20 | |||
| 21 | func TestSweepRemovesExpiredSessionsAndTokens(t *testing.T) { | ||
| 22 | s, uid := retentionFixture(t) | ||
| 23 | if err := s.CreateWebSession("live", uid, time.Hour); err != nil { | ||
| 24 | t.Fatal(err) | ||
| 25 | } | ||
| 26 | if err := s.CreateWebSession("dead", uid, -time.Hour); err != nil { | ||
| 27 | t.Fatal(err) | ||
| 28 | } | ||
| 29 | if err := s.CreateLoginToken(uid, "stale", -time.Minute); err != nil { | ||
| 30 | t.Fatal(err) | ||
| 31 | } | ||
| 32 | |||
| 33 | // Retention is unset: expired rows still go, since nothing keeps them | ||
| 34 | // meaningful once they cannot authenticate. | ||
| 35 | got, err := s.Sweep(Retention{}, time.Now()) | ||
| 36 | if err != nil { | ||
| 37 | t.Fatal(err) | ||
| 38 | } | ||
| 39 | if got["web_sessions"] != 1 || got["login_tokens"] != 1 { | ||
| 40 | t.Fatalf("swept %v", got) | ||
| 41 | } | ||
| 42 | var n int | ||
| 43 | s.DB.QueryRow("SELECT COUNT(*) FROM web_sessions").Scan(&n) | ||
| 44 | if n != 1 { | ||
| 45 | t.Fatalf("%d sessions remain, want the live one", n) | ||
| 46 | } | ||
| 47 | if _, err := s.WebSessionUser("live"); err != nil { | ||
| 48 | t.Fatalf("live session swept: %v", err) | ||
| 49 | } | ||
| 50 | } | ||
| 51 | |||
| 52 | // Zero retention keeps forever: deleting an audit trail is not a default. | ||
| 53 | func TestSweepKeepsWhenUnconfigured(t *testing.T) { | ||
| 54 | s, uid := retentionFixture(t) | ||
| 55 | s.Audit(uid, "cmd repo create", map[string]any{"argv": []string{"a/b"}}) | ||
| 56 | s.DB.Exec("UPDATE audit_log SET created_at = '2020-01-01T00:00:00.000Z'") | ||
| 57 | |||
| 58 | if _, err := s.Sweep(Retention{}, time.Now()); err != nil { | ||
| 59 | t.Fatal(err) | ||
| 60 | } | ||
| 61 | var n int | ||
| 62 | s.DB.QueryRow("SELECT COUNT(*) FROM audit_log").Scan(&n) | ||
| 63 | if n != 1 { | ||
| 64 | t.Fatalf("audit row swept with retention unset (%d rows)", n) | ||
| 65 | } | ||
| 66 | |||
| 67 | if _, err := s.Sweep(Retention{Audit: 24 * time.Hour}, time.Now()); err != nil { | ||
| 68 | t.Fatal(err) | ||
| 69 | } | ||
| 70 | s.DB.QueryRow("SELECT COUNT(*) FROM audit_log").Scan(&n) | ||
| 71 | if n != 0 { | ||
| 72 | t.Fatalf("audit row survived its retention (%d rows)", n) | ||
| 73 | } | ||
| 74 | } | ||
| 75 | |||
| 76 | // A delivery still being retried is live state, however old its first | ||
| 77 | // attempt; only finished ones age out. | ||
| 78 | func TestSweepKeepsUnfinishedWork(t *testing.T) { | ||
| 79 | s, uid := retentionFixture(t) | ||
| 80 | repoID, err := s.CreateRepo("user", uid, "lib", "public") | ||
| 81 | if err != nil { | ||
| 82 | t.Fatal(err) | ||
| 83 | } | ||
| 84 | if _, err := s.AddWebhook(repoID, "https://example.test/h", "", "*"); err != nil { | ||
| 85 | t.Fatal(err) | ||
| 86 | } | ||
| 87 | for i := 0; i < 3; i++ { | ||
| 88 | if err := s.RecordEvent(repoID, uid, "push", "{}"); err != nil { | ||
| 89 | t.Fatal(err) | ||
| 90 | } | ||
| 91 | } | ||
| 92 | s.DB.Exec("UPDATE webhook_deliveries SET created_at = '2020-01-01T00:00:00.000Z'") | ||
| 93 | s.DB.Exec("UPDATE webhook_deliveries SET delivered_at = '2020-01-02T00:00:00.000Z' WHERE id = 1") | ||
| 94 | s.DB.Exec("UPDATE webhook_deliveries SET failed_at = '2020-01-02T00:00:00.000Z' WHERE id = 2") | ||
| 95 | |||
| 96 | got, err := s.Sweep(Retention{WebhookDeliveries: time.Hour}, time.Now()) | ||
| 97 | if err != nil { | ||
| 98 | t.Fatal(err) | ||
| 99 | } | ||
| 100 | if got["webhook_deliveries"] != 2 { | ||
| 101 | t.Fatalf("swept %v, want the delivered and the failed one", got) | ||
| 102 | } | ||
| 103 | var n int | ||
| 104 | s.DB.QueryRow("SELECT COUNT(*) FROM webhook_deliveries").Scan(&n) | ||
| 105 | if n != 1 { | ||
| 106 | t.Fatalf("%d deliveries remain, want the one still being retried", n) | ||
| 107 | } | ||
| 108 | } | ||
| 109 | |||
| 110 | // events.id is the parent of webhook_deliveries.event_id under ON DELETE | ||
| 111 | // CASCADE, so sweeping events would take a queued delivery with it. The | ||
| 112 | // sweep skips any event that still has one. | ||
| 113 | func TestSweepEventsSpareQueuedDeliveries(t *testing.T) { | ||
| 114 | s, uid := retentionFixture(t) | ||
| 115 | repoID, err := s.CreateRepo("user", uid, "lib", "public") | ||
| 116 | if err != nil { | ||
| 117 | t.Fatal(err) | ||
| 118 | } | ||
| 119 | if _, err := s.AddWebhook(repoID, "https://example.test/h", "", "*"); err != nil { | ||
| 120 | t.Fatal(err) | ||
| 121 | } | ||
| 122 | if err := s.RecordEvent(repoID, uid, "push", "{}"); err != nil { | ||
| 123 | t.Fatal(err) | ||
| 124 | } | ||
| 125 | s.DB.Exec("UPDATE events SET created_at = '2020-01-01T00:00:00.000Z'") | ||
| 126 | s.DB.Exec("UPDATE webhook_deliveries SET created_at = '2020-01-01T00:00:00.000Z'") | ||
| 127 | |||
| 128 | // The delivery has neither delivered_at nor failed_at: still queued. | ||
| 129 | got, err := s.Sweep(Retention{Events: time.Hour, WebhookDeliveries: time.Hour}, time.Now()) | ||
| 130 | if err != nil { | ||
| 131 | t.Fatal(err) | ||
| 132 | } | ||
| 133 | if got["events"] != 0 || got["webhook_deliveries"] != 0 { | ||
| 134 | t.Fatalf("swept %v, want nothing while the delivery is queued", got) | ||
| 135 | } | ||
| 136 | var n int | ||
| 137 | s.DB.QueryRow("SELECT COUNT(*) FROM webhook_deliveries").Scan(&n) | ||
| 138 | if n != 1 { | ||
| 139 | t.Fatal("queued delivery cascaded away with its event") | ||
| 140 | } | ||
| 141 | |||
| 142 | // Once it finishes, both go. | ||
| 143 | s.DB.Exec("UPDATE webhook_deliveries SET delivered_at = '2020-01-02T00:00:00.000Z'") | ||
| 144 | if _, err := s.Sweep(Retention{Events: time.Hour, WebhookDeliveries: time.Hour}, time.Now()); err != nil { | ||
| 145 | t.Fatal(err) | ||
| 146 | } | ||
| 147 | s.DB.QueryRow("SELECT COUNT(*) FROM events").Scan(&n) | ||
| 148 | if n != 0 { | ||
| 149 | t.Fatalf("%d events remain after the delivery finished", n) | ||
| 150 | } | ||
| 151 | } | ||