diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 7f564e7..93f4741 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -38,5 +38,7 @@ jobs: run: go tool cover -func=coverage.out | tail -n 1 - name: Build run: go build ./... + - name: Demo smoke test + run: go run . demo -run 1s - name: Benchmark queue paths - run: go test ./internal/queue -run '^$' -bench Benchmark -benchmem -count=1 \ No newline at end of file + run: go test ./internal/queue -run '^$' -bench Benchmark -benchmem -count=1 diff --git a/CHANGELOG.md b/CHANGELOG.md new file mode 100644 index 0000000..1df9312 --- /dev/null +++ b/CHANGELOG.md @@ -0,0 +1,42 @@ +# Changelog + +All notable changes to this project are listed here. + +The format follows Keep a Changelog. +This project uses calendar versioning by release. + +## Unreleased + +### Added + +- Retention with the `purge` command. It removes finished jobs and their events in one transaction. +- The `purge -dry-run` flag previews a removal without changing the store. +- The `purge -state` and `purge -before` flags name states and set an age cutoff. +- The demo shows a retention segment with a dry run and a real purge. +- Deterministic tests for purge selection, age bounds, and event removal. + +## 2026-08-03 + +### Added + +- Prometheus metrics for queue inspection. +- The `metrics` command serves the exposition format over HTTP. +- The `work` command can serve the same endpoint beside a worker. +- Priority aging prevents starvation of low-priority work. +- The demo shows a low-priority job overtaking a fresher high-priority job. + +## 2026-08-01 + +### Added + +- Dead-letter queue with the `requeue` command. +- Durable priority dispatch with deterministic ordering. +- Scheduled jobs with nanosecond-safe release times. + +## 2026-07-28 + +### Added + +- Durable leases, retries, idempotency, crash recovery, and event history. +- The `enqueue`, `work`, `inspect`, `history`, `seed`, and `demo` commands. +- A deterministic fault-injection harness for worker handlers. diff --git a/README.md b/README.md index 4daff53..83e4a7c 100644 --- a/README.md +++ b/README.md @@ -2,7 +2,7 @@ A small durable background queue built with Go and SQLite. -The project shows leases, retries, idempotency, crash recovery, priority dispatch, priority aging, a dead-letter queue, Prometheus metrics, and an append-only event log. +The project shows leases, retries, idempotency, crash recovery, priority dispatch, priority aging, a dead-letter queue, retention, Prometheus metrics, and an append-only event log. ## Value @@ -12,6 +12,8 @@ Every lease, retry, recovery, and acknowledgement remains visible in SQLite. The demo injects repeatable faults, so failure paths are easy to inspect. +The purge command enforces retention, so finished jobs cannot grow without limit. + ## Architecture The queue separates durable state from worker execution. @@ -74,6 +76,17 @@ go run . history -db queue.db go run . requeue -db queue.db ``` +Enforce retention with these commands. + +```text +go run . purge -db queue.db -dry-run +go run . purge -db queue.db -before 720h +``` + +The first command previews the removal. + +The second command removes jobs not updated in 30 days. + ## Commands ### `enqueue` @@ -140,6 +153,30 @@ jobqueue seed [-db ] The command loads three idempotent jobs for each bundled workload. +### `purge` + +```text +jobqueue purge [-state ]... [-before ] [-dry-run] [-db ] +``` + +The command removes finished jobs and their events. + +It targets the terminal states by default. + +Those states are `completed`, `failed`, and `dead_letter`. + +Pending and leased jobs survive a default purge. + +Use `-state` to name different states. + +Use `-state pending` to clear a backlog on purpose. + +Use `-before` to keep recent history. + +The command measures age from the job's last update. + +Use `-dry-run` to preview the removal without changing the store. + ### `metrics` ```text @@ -236,6 +273,26 @@ Every state change appends one event row. The `history` command shows one job's complete timeline. +### Retention + +The `purge` command removes finished jobs and their events in one transaction. + +Each removed job takes its event rows with it. + +The append-only log stays consistent with the remaining jobs. + +The default target set is the terminal states. + +A purge never touches pending or leased work unless you name those states. + +An operator can clear a stuck backlog with `-state pending`. + +Use `-before` to keep jobs updated within a chosen window. + +A dry run reports the exact counts before any change. + +The removal is a single SQLite transaction, so it is atomic. + ### Metrics The exporter renders queue state in the Prometheus text format. @@ -305,6 +362,13 @@ aging interval: 100ms; a job gains one priority point per interval it waits. lease order: (aged) then (fresh) the waiting job outranks the fresher higher-priority job. +Retention +--------- +jobs: pending=1 completed=2 dead_letter=1 + dry run: would remove 3 jobs and 9 events + removed 3 jobs and 9 events +after purge: pending=1 completed=0 dead_letter=0 + Queue state ----------- completed: 6 @@ -370,7 +434,7 @@ SQLite serializes writes through one store connection. A sustained backlog can exceed the writer's capacity. -Jobs and events remain until an operator removes them. +Jobs and events accumulate until an operator runs the purge command. A high-priority stream can delay lower-priority jobs until aging lifts them. @@ -386,12 +450,27 @@ The project does not provide a web interface. - [x] Priority aging to prevent starvation. - [x] Dead-letter queue with requeue of permanently failed jobs. - [x] Prometheus metrics for queue inspection. +- [x] Retention with a purge command for finished jobs and their events. - [ ] Web UI for queue inspection. - [ ] Horizontal scaling with a shared SQLite file. ### Release notes -This release adds Prometheus metrics. +This release adds retention with the `purge` command. + +The command removes finished jobs and their events in one transaction. + +It targets the terminal states by default. + +Use `-state` to name different states. + +Use `-before` to keep recent history. + +Use `-dry-run` to preview a removal without changing the store. + +The demo now shows a retention segment with a dry run and a real purge. + +The previous release added Prometheus metrics. The new `metrics` command serves the exposition format over HTTP. diff --git a/internal/cli/demo.go b/internal/cli/demo.go index 7048f28..6696a4b 100644 --- a/internal/cli/demo.go +++ b/internal/cli/demo.go @@ -190,6 +190,9 @@ func Demo(args []string) error { if err := showPriorityAging(*kind); err != nil { return err } + if err := showRetention(*kind); err != nil { + return err + } fmt.Println() snap, err = q.Inspect() @@ -208,6 +211,83 @@ func Demo(args []string) error { fmt.Printf("\ninspect again with: jobqueue inspect -db %q\n", path) fmt.Printf("inspect a job with: jobqueue history -db %q\n", path) fmt.Printf("requeue a dead letter with: jobqueue requeue -db %q\n", path) + fmt.Printf("purge finished jobs with: jobqueue purge -db %q\n", path) + return nil +} + +// showRetention demonstrates the purge workflow: preview the removal with a +// dry run, apply it, and confirm the queue state. The segment uses its own +// in-memory store so the scenario counts stay intact. +func showRetention(kind string) error { + store, err := queue.NewSQLiteStore("file:jobqueue-demo-retention?mode=memory&cache=shared") + if err != nil { + return fmt.Errorf("open retention store: %w", err) + } + defer store.Close() + q := queue.NewQueue(store) + ctx := context.Background() + + // complete enqueues one job and finishes it in one step. Each lease is + // unambiguous because the future job is not ready yet. + complete := func() error { + if _, err := q.Enqueue(kind, `{"name":"completed"}`); err != nil { + return fmt.Errorf("enqueue completed: %w", err) + } + job, err := q.Lease(ctx, kind, time.Minute) + if err != nil || job == nil { + return fmt.Errorf("lease completed: job=%v err=%v", job, err) + } + if err := q.Acknowledge(job.ID); err != nil { + return fmt.Errorf("ack completed: %w", err) + } + return nil + } + for i := 0; i < 2; i++ { + if err := complete(); err != nil { + return err + } + } + + if _, err := q.Enqueue(kind, `{"name":"doomed"}`, queue.WithMaxAttempts(1)); err != nil { + return fmt.Errorf("enqueue doomed: %w", err) + } + doomed, err := q.Lease(ctx, kind, time.Minute) + if err != nil || doomed == nil { + return fmt.Errorf("lease doomed: %w", err) + } + if err := q.Fail(doomed.ID, "boom"); err != nil { + return fmt.Errorf("fail doomed: %w", err) + } + if _, err := q.Enqueue(kind, `{"name":"future"}`, queue.WithRunAt(time.Now().Add(time.Hour))); err != nil { + return fmt.Errorf("enqueue future: %w", err) + } + + snap, err := q.Inspect() + if err != nil { + return fmt.Errorf("inspect retention: %w", err) + } + fmt.Println("\nRetention") + fmt.Println("---------") + fmt.Printf("jobs: pending=%d completed=%d dead_letter=%d\n", + snap.Stats[queue.StatePending], snap.Stats[queue.StateCompleted], snap.Stats[queue.StateDeadLetter]) + + var b strings.Builder + if err := runPurge(&b, q, nil, nil, true); err != nil { + return fmt.Errorf("retention dry run: %w", err) + } + fmt.Print(" " + b.String()) + b.Reset() + if err := runPurge(&b, q, nil, nil, false); err != nil { + return fmt.Errorf("retention purge: %w", err) + } + fmt.Print(" " + b.String()) + + snap, err = q.Inspect() + if err != nil { + return fmt.Errorf("inspect after purge: %w", err) + } + fmt.Printf("after purge: pending=%d completed=%d dead_letter=%d\n", + snap.Stats[queue.StatePending], snap.Stats[queue.StateCompleted], snap.Stats[queue.StateDeadLetter]) return nil } diff --git a/internal/cli/purge.go b/internal/cli/purge.go new file mode 100644 index 0000000..89b32c9 --- /dev/null +++ b/internal/cli/purge.go @@ -0,0 +1,152 @@ +package cli + +import ( + "flag" + "fmt" + "io" + "os" + "strings" + "time" + + "github.com/local-first-job-queue/internal/queue" +) + +// stringList is a repeatable and comma-separated string flag. +type stringList []string + +func (s *stringList) String() string { + return strings.Join(*s, ",") +} + +func (s *stringList) Set(v string) error { + for _, part := range strings.Split(v, ",") { + part = strings.TrimSpace(part) + if part != "" { + *s = append(*s, part) + } + } + return nil +} + +// validStates is the set of states the purge command accepts on the command +// line. The list matches the states the store and the metrics renderer know. +var validStates = []queue.JobState{ + queue.StatePending, + queue.StateLeased, + queue.StateCompleted, + queue.StateFailed, + queue.StateDeadLetter, +} + +// Purge removes finished jobs and their events from the store. An operator +// uses it to enforce a retention policy: keep the queue small and the event +// log bounded. By default the command targets the terminal states: completed, +// failed, and dead_letter. Use -state to name different states, -before to +// keep recent history, and -dry-run to preview what would be removed. +func Purge(args []string) error { + fs := flag.NewFlagSet("purge", flag.ExitOnError) + dbPath := fs.String("db", "queue.db", "database path") + dryRun := fs.Bool("dry-run", false, "report what would be removed without removing it") + beforeDur := fs.Duration("before", 0, "remove only jobs not updated within this duration; 0 removes every age") + var stateFlags stringList + fs.Var(&stateFlags, "state", "job state to purge; repeatable or comma-separated (default: completed, failed, dead_letter)") + fs.Parse(args) + + states, err := parsePurgeStates(stateFlags) + if err != nil { + return err + } + + var before *time.Time + if *beforeDur > 0 { + t := time.Now().UTC().Add(-*beforeDur) + before = &t + } + + store, err := queue.NewSQLiteStore(*dbPath) + if err != nil { + return fmt.Errorf("open store: %w", err) + } + defer store.Close() + + q := queue.NewQueue(store) + return runPurge(os.Stdout, q, states, before, *dryRun) +} + +// runPurge applies a retention filter to the queue and writes the outcome to +// w. It is separate from Purge so tests can capture the output without +// touching the process stdout. A nil before means no age filter, and a nil +// state list uses the command default, which is the terminal states. +func runPurge(w io.Writer, q *queue.Queue, states []queue.JobState, before *time.Time, dryRun bool) error { + opts := []queue.PurgeOption{} + if len(states) > 0 { + opts = append(opts, queue.PurgeStates(states...)) + } + if before != nil { + opts = append(opts, queue.PurgeBefore(*before)) + } + + if dryRun { + stats, err := q.PurgeCandidates(opts...) + if err != nil { + return fmt.Errorf("preview purge: %w", err) + } + fmt.Fprintf(w, "dry run: would remove %d jobs and %d events\n", stats.JobsRemoved, stats.EventsRemoved) + if stats.JobsRemoved == 0 { + fmt.Fprintln(w, "no jobs match the current filter.") + } else { + fmt.Fprintln(w, "run without -dry-run to apply.") + } + return nil + } + + stats, err := q.Purge(opts...) + if err != nil { + return fmt.Errorf("purge: %w", err) + } + fmt.Fprintf(w, "removed %d jobs and %d events\n", stats.JobsRemoved, stats.EventsRemoved) + switch { + case stats.JobsRemoved == 0: + fmt.Fprintln(w, "no jobs match the current filter.") + case stats.JobsRemoved > 0: + fmt.Fprintln(w, "run 'jobqueue inspect' to confirm the queue state.") + } + return nil +} + +// parsePurgeStates converts the -state flag values into job states. An empty +// list keeps the command default, which the queue library applies when no +// option is present. Each value may carry one state or a comma-separated list, +// so callers can pass -state twice or as a single list. Unknown state names +// are rejected with the valid list. +func parsePurgeStates(values stringList) ([]queue.JobState, error) { + if len(values) == 0 { + return nil, nil + } + valid := map[string]bool{} + for _, st := range validStates { + valid[string(st)] = true + } + var states []queue.JobState + for _, v := range values { + for _, part := range strings.Split(v, ",") { + part = strings.TrimSpace(part) + if part == "" { + continue + } + if !valid[part] { + return nil, fmt.Errorf("unknown state %q (valid: %s)", part, strings.Join(validStateNames(), ", ")) + } + states = append(states, queue.JobState(part)) + } + } + return states, nil +} + +func validStateNames() []string { + names := make([]string, 0, len(validStates)) + for _, st := range validStates { + names = append(names, string(st)) + } + return names +} diff --git a/internal/cli/purge_test.go b/internal/cli/purge_test.go new file mode 100644 index 0000000..57fb233 --- /dev/null +++ b/internal/cli/purge_test.go @@ -0,0 +1,165 @@ +package cli + +import ( + "context" + "strings" + "testing" + "time" + + "github.com/local-first-job-queue/internal/queue" +) + +// buildPurgeDB returns an in-memory database with two completed jobs and one +// pending job. Tests reuse it to verify the purge command end to end. +func buildPurgeDB(t *testing.T) (*queue.Queue, *queue.SQLiteStore) { + t.Helper() + store, err := queue.NewSQLiteStore("file:cli_purge_" + t.Name() + "?mode=memory&cache=shared") + if err != nil { + t.Fatalf("new store: %v", err) + } + t.Cleanup(func() { store.Close() }) + q := queue.NewQueue(store) + ctx := context.Background() + + for i := 0; i < 2; i++ { + kind := "done" + string(rune('0'+i)) + if _, err := q.Enqueue(kind, `{"n":`+string(rune('0'+i))+`}`); err != nil { + t.Fatalf("enqueue %s: %v", kind, err) + } + job, err := q.Lease(ctx, kind, time.Minute) + if err != nil || job == nil { + t.Fatalf("lease %s: job=%v err=%v", kind, job, err) + } + if err := q.Acknowledge(job.ID); err != nil { + t.Fatalf("ack %s: %v", kind, err) + } + } + if _, err := q.Enqueue("waiting", `{}`); err != nil { + t.Fatalf("enqueue waiting: %v", err) + } + return q, store +} + +// TestPurgeCommandRemovesTerminalJobs verifies that the purge command removes +// completed jobs, keeps pending work, and reports the removed counts. +func TestPurgeCommandRemovesTerminalJobs(t *testing.T) { + q, _ := buildPurgeDB(t) + + var b strings.Builder + if err := runPurge(&b, q, nil, nil, false); err != nil { + t.Fatalf("purge: %v", err) + } + if want := "removed 2 jobs and 6 events"; !strings.Contains(b.String(), want) { + t.Errorf("expected %q in output, got %q", want, b.String()) + } + + snap, err := q.Inspect() + if err != nil { + t.Fatalf("inspect: %v", err) + } + if snap.Stats[queue.StateCompleted] != 0 { + t.Errorf("expected 0 completed jobs, got %d", snap.Stats[queue.StateCompleted]) + } + if snap.Stats[queue.StatePending] != 1 { + t.Errorf("expected 1 pending job, got %d", snap.Stats[queue.StatePending]) + } +} + +// TestPurgeCommandDryRunDoesNotRemove verifies that a dry run reports the same +// counts as a real purge but leaves the store unchanged. +func TestPurgeCommandDryRunDoesNotRemove(t *testing.T) { + q, _ := buildPurgeDB(t) + + var b strings.Builder + if err := runPurge(&b, q, nil, nil, true); err != nil { + t.Fatalf("dry run: %v", err) + } + if want := "dry run: would remove 2 jobs and 6 events"; !strings.Contains(b.String(), want) { + t.Errorf("expected %q in output, got %q", want, b.String()) + } + if !strings.Contains(b.String(), "run without -dry-run to apply") { + t.Errorf("expected the apply hint in output, got %q", b.String()) + } + + snap, err := q.Inspect() + if err != nil { + t.Fatalf("inspect: %v", err) + } + if len(snap.Jobs) != 3 { + t.Errorf("dry run must not remove jobs, got %d", len(snap.Jobs)) + } +} + +// TestPurgeCommandStateFilter verifies that a state filter limits the purge to +// the named state and that unknown states are rejected. +func TestPurgeCommandStateFilter(t *testing.T) { + q, _ := buildPurgeDB(t) + + var b strings.Builder + if err := runPurge(&b, q, []queue.JobState{queue.StatePending}, nil, false); err != nil { + t.Fatalf("purge pending: %v", err) + } + if want := "removed 1 jobs and 1 events"; !strings.Contains(b.String(), want) { + t.Errorf("expected %q in output, got %q", want, b.String()) + } + + snap, err := q.Inspect() + if err != nil { + t.Fatalf("inspect: %v", err) + } + if snap.Stats[queue.StatePending] != 0 { + t.Errorf("expected pending jobs purged, got %d", snap.Stats[queue.StatePending]) + } + if snap.Stats[queue.StateCompleted] != 2 { + t.Errorf("expected completed jobs to survive, got %d", snap.Stats[queue.StateCompleted]) + } + + if _, err := parsePurgeStates(stringList{"nonsense"}); err == nil { + t.Errorf("expected an error for an unknown state") + } +} + +// TestPurgeCommandBeforeKeepsRecent verifies that an age filter keeps jobs +// updated after the cutoff. Every job in the fixture is fresh, so a short +// retention window removes nothing. +func TestPurgeCommandBeforeKeepsRecent(t *testing.T) { + q, _ := buildPurgeDB(t) + + before := time.Now().UTC().Add(-24 * time.Hour) + var b strings.Builder + if err := runPurge(&b, q, nil, &before, false); err != nil { + t.Fatalf("purge before: %v", err) + } + if want := "removed 0 jobs and 0 events"; !strings.Contains(b.String(), want) { + t.Errorf("expected %q in output, got %q", want, b.String()) + } + if !strings.Contains(b.String(), "no jobs match") { + t.Errorf("expected the no-match note, got %q", b.String()) + } +} + +// TestParsePurgeStates verifies the state parsing rules: empty is the command +// default, one value works, and comma-separated values expand. +func TestParsePurgeStates(t *testing.T) { + states, err := parsePurgeStates(nil) + if err != nil { + t.Fatalf("empty: %v", err) + } + if len(states) != 0 { + t.Errorf("empty input must yield no states, got %v", states) + } + + states, err = parsePurgeStates(stringList{"pending"}) + if err != nil || len(states) != 1 || states[0] != queue.StatePending { + t.Errorf("single state: states=%v err=%v", states, err) + } + + states, err = parsePurgeStates(stringList{"pending, completed"}) + if err != nil || len(states) != 2 { + t.Errorf("comma-separated: states=%v err=%v", states, err) + } + + if _, err := parsePurgeStates(stringList{"pending", "bogus"}); err == nil { + t.Errorf("expected an error for an unknown state") + } +} diff --git a/internal/queue/purge.go b/internal/queue/purge.go new file mode 100644 index 0000000..215d9b9 --- /dev/null +++ b/internal/queue/purge.go @@ -0,0 +1,163 @@ +package queue + +import ( + "fmt" + "strings" + "time" +) + +// PurgeOption configures a purge. +type PurgeOption func(*purgeConfig) + +type purgeConfig struct { + states []JobState + statesSet bool + before *time.Time +} + +// PurgeStates limits a purge to the given job states. When the option is +// absent, a purge targets the terminal states: completed, failed, and +// dead_letter. Pending and leased jobs are never removed unless a caller names +// those states explicitly. An empty call removes nothing. +func PurgeStates(states ...JobState) PurgeOption { + return func(c *purgeConfig) { + c.states = append(c.states, states...) + c.statesSet = true + } +} + +// PurgeBefore removes only jobs whose last update is older than the given +// time. When the option is absent, a purge ignores job age. The comparison +// uses the job's updated_at field, which records when the job last changed +// state. +func PurgeBefore(t time.Time) PurgeOption { + return func(c *purgeConfig) { + v := t.UTC() + c.before = &v + } +} + +// PurgeStats reports what one purge removed. +type PurgeStats struct { + JobsRemoved int `json:"jobs_removed"` + EventsRemoved int `json:"events_removed"` +} + +// resolvePurgeConfig turns the options into the effective state set and age +// filter. The default state set is terminal, so a plain purge never touches +// work that is still pending or leased. An explicit empty state set stays +// empty and therefore removes nothing. +func resolvePurgeConfig(opts []PurgeOption) (states []JobState, before *time.Time) { + cfg := purgeConfig{} + for _, o := range opts { + o(&cfg) + } + states = cfg.states + if !cfg.statesSet { + states = []JobState{StateCompleted, StateFailed, StateDeadLetter} + } + return states, cfg.before +} + +// Purge removes finished jobs and their events from the store. Each removed +// job takes its event rows with it, so the append-only log stays consistent +// with the remaining jobs. The default target set is the terminal states. Use +// PurgeBefore to keep recent history and PurgeStates to name different states. +func (q *Queue) Purge(opts ...PurgeOption) (PurgeStats, error) { + states, before := resolvePurgeConfig(opts) + removedJobs, removedEvents, err := q.store.PurgeJobs(states, before) + if err != nil { + return PurgeStats{}, fmt.Errorf("purge: %w", err) + } + return PurgeStats{JobsRemoved: removedJobs, EventsRemoved: removedEvents}, nil +} + +// PurgeCandidates reports what Purge would remove without changing the store. +// Operators use it to preview a retention policy before applying it. +func (q *Queue) PurgeCandidates(opts ...PurgeOption) (PurgeStats, error) { + states, before := resolvePurgeConfig(opts) + jobs, events, err := q.store.CountPurgeCandidates(states, before) + if err != nil { + return PurgeStats{}, fmt.Errorf("count purge: %w", err) + } + return PurgeStats{JobsRemoved: jobs, EventsRemoved: events}, nil +} + +// PurgeJobs removes jobs in the given states, and the events of those jobs, +// in one transaction. When before is non-nil, only jobs whose updated_at is +// older than that time are removed. The method returns the number of jobs and +// events removed. A removed job's events are removed with it, so the +// append-only log stays consistent with the remaining jobs. +func (s *SQLiteStore) PurgeJobs(states []JobState, before *time.Time) (int, int, error) { + if len(states) == 0 { + return 0, 0, nil + } + where, args := purgeWhere(states, before) + + tx, err := s.db.Begin() + if err != nil { + return 0, 0, fmt.Errorf("begin tx: %w", err) + } + defer tx.Rollback() + + // The event rows are deleted first. Foreign keys are not enforced, so the + // order is what keeps the log consistent with the remaining jobs. + eventsRes, err := tx.Exec(`DELETE FROM events WHERE job_id IN (SELECT id FROM jobs WHERE `+where+`)`, args...) + if err != nil { + return 0, 0, fmt.Errorf("purge events: %w", err) + } + jobsRes, err := tx.Exec(`DELETE FROM jobs WHERE `+where, args...) + if err != nil { + return 0, 0, fmt.Errorf("purge jobs: %w", err) + } + if err := tx.Commit(); err != nil { + return 0, 0, fmt.Errorf("commit purge: %w", err) + } + + removedJobs, err := jobsRes.RowsAffected() + if err != nil { + return 0, 0, fmt.Errorf("jobs removed: %w", err) + } + removedEvents, err := eventsRes.RowsAffected() + if err != nil { + return 0, 0, fmt.Errorf("events removed: %w", err) + } + return int(removedJobs), int(removedEvents), nil +} + +// CountPurgeCandidates reports how many jobs a purge would remove, along with +// the number of their events. It reads the same rows that PurgeJobs deletes, +// so the preview matches a real run. +func (s *SQLiteStore) CountPurgeCandidates(states []JobState, before *time.Time) (int, int, error) { + if len(states) == 0 { + return 0, 0, nil + } + where, args := purgeWhere(states, before) + + var jobs int + if err := s.db.QueryRow(`SELECT COUNT(*) FROM jobs WHERE `+where, args...).Scan(&jobs); err != nil { + return 0, 0, fmt.Errorf("count jobs: %w", err) + } + var events int + if err := s.db.QueryRow(`SELECT COUNT(*) FROM events WHERE job_id IN (SELECT id FROM jobs WHERE `+where+`)`, args...).Scan(&events); err != nil { + return 0, 0, fmt.Errorf("count events: %w", err) + } + return jobs, events, nil +} + +// purgeWhere builds the WHERE clause that selects the rows a purge touches. +// The state placeholders appear first and the optional age bound last, so one +// argument slice serves both the job and the event queries. +func purgeWhere(states []JobState, before *time.Time) (string, []any) { + placeholders := strings.TrimRight(strings.Repeat("?,", len(states)), ",") + where := `state IN (` + placeholders + `)` + args := make([]any, 0, len(states)+1) + for _, st := range states { + args = append(args, string(st)) + } + if before != nil { + where += ` AND updated_at < ?` + args = append(args, before.UTC().Format(sqliteTimeFormat)) + } + return where, args +} diff --git a/internal/queue/purge_test.go b/internal/queue/purge_test.go new file mode 100644 index 0000000..c1f712b --- /dev/null +++ b/internal/queue/purge_test.go @@ -0,0 +1,413 @@ +package queue + +import ( + "context" + "fmt" + "testing" + "time" +) + +// seedPurgeLifecycle builds a store with one job in each state: completed, +// dead_letter, leased, and pending. Each job uses its own kind so the leases +// are unambiguous. The function returns the store and the pending job ID. +func seedPurgeLifecycle(t *testing.T) (*SQLiteStore, string) { + t.Helper() + s := newTestStore(t) + q := NewQueue(s) + ctx := context.Background() + + if _, err := q.Enqueue("done", `{}`); err != nil { + t.Fatalf("enqueue done: %v", err) + } + if job, err := q.Lease(ctx, "done", time.Minute); err != nil || job == nil { + t.Fatalf("lease done: job=%v err=%v", job, err) + } else if err := q.Acknowledge(job.ID); err != nil { + t.Fatalf("ack done: %v", err) + } + + if _, err := q.Enqueue("lost", `{}`, WithMaxAttempts(1)); err != nil { + t.Fatalf("enqueue lost: %v", err) + } + if job, err := q.Lease(ctx, "lost", time.Minute); err != nil || job == nil { + t.Fatalf("lease lost: job=%v err=%v", job, err) + } else if err := q.Fail(job.ID, "boom"); err != nil { + t.Fatalf("fail lost: %v", err) + } + + if _, err := q.Enqueue("held", `{}`); err != nil { + t.Fatalf("enqueue held: %v", err) + } + if _, err := q.Lease(ctx, "held", time.Minute); err != nil { + t.Fatalf("lease held: %v", err) + } + + pending, err := q.Enqueue("waiting", `{}`) + if err != nil { + t.Fatalf("enqueue waiting: %v", err) + } + + stats, err := s.GetQueueStats() + if err != nil { + t.Fatalf("stats: %v", err) + } + want := map[JobState]int{ + StateCompleted: 1, + StateDeadLetter: 1, + StateLeased: 1, + StatePending: 1, + } + for state, count := range want { + if stats[state] != count { + t.Fatalf("seed state %s: got %d, want %d", state, stats[state], count) + } + } + + return s, pending.ID +} + +// TestPurgeDefaultsToTerminalStates verifies that a purge without options +// removes completed and dead-lettered jobs while leaving pending and leased +// work untouched. +func TestPurgeDefaultsToTerminalStates(t *testing.T) { + s, pendingID := seedPurgeLifecycle(t) + q := NewQueue(s) + + stats, err := q.Purge() + if err != nil { + t.Fatalf("purge: %v", err) + } + if stats.JobsRemoved != 2 { + t.Errorf("expected 2 jobs removed, got %d", stats.JobsRemoved) + } + if stats.EventsRemoved != 6 { + t.Errorf("expected 6 events removed (3 per terminal job), got %d", stats.EventsRemoved) + } + + jobs, err := s.GetAllJobs() + if err != nil { + t.Fatalf("get jobs: %v", err) + } + if len(jobs) != 2 { + t.Fatalf("expected 2 remaining jobs, got %d", len(jobs)) + } + states := map[JobState]bool{} + for _, j := range jobs { + states[j.State] = true + if j.ID == pendingID && j.State != StatePending { + t.Errorf("pending job must survive the purge, got state %s", j.State) + } + } + if !states[StatePending] || !states[StateLeased] { + t.Errorf("expected pending and leased jobs to survive, got %v", states) + } +} + +// TestPurgeStatesSelectsExactStates verifies that naming states limits the +// purge to exactly those states. +func TestPurgeStatesSelectsExactStates(t *testing.T) { + s, _ := seedPurgeLifecycle(t) + q := NewQueue(s) + + stats, err := q.Purge(PurgeStates(StateDeadLetter)) + if err != nil { + t.Fatalf("purge: %v", err) + } + if stats.JobsRemoved != 1 { + t.Errorf("expected 1 dead-letter job removed, got %d", stats.JobsRemoved) + } + if stats.EventsRemoved != 3 { + t.Errorf("expected 3 events removed, got %d", stats.EventsRemoved) + } + + snap, err := q.Inspect() + if err != nil { + t.Fatalf("inspect: %v", err) + } + if snap.Stats[StateDeadLetter] != 0 { + t.Errorf("expected no dead-letter jobs, got %d", snap.Stats[StateDeadLetter]) + } + if snap.Stats[StateCompleted] != 1 { + t.Errorf("expected completed job to survive, got %d", snap.Stats[StateCompleted]) + } +} + +// TestPurgeRemovesEventsWithJobs verifies that the events of purged jobs are +// deleted while the events of surviving jobs remain readable. +func TestPurgeRemovesEventsWithJobs(t *testing.T) { + s, pendingID := seedPurgeLifecycle(t) + q := NewQueue(s) + + if _, err := q.Purge(PurgeStates(StateCompleted)); err != nil { + t.Fatalf("purge: %v", err) + } + + survivors, err := s.GetAllJobs() + if err != nil { + t.Fatalf("get jobs: %v", err) + } + for _, j := range survivors { + if j.State == StateCompleted { + t.Errorf("completed job %s must be gone", j.ID) + } + events, err := s.GetJobEvents(j.ID) + if err != nil { + t.Fatalf("events for %s: %v", j.ID, err) + } + if len(events) == 0 { + t.Errorf("surviving job %s lost its events", j.ID) + } + } + + // The pending job kept its full timeline: enqueued. + pending, err := s.GetJob(pendingID) + if err != nil { + t.Fatalf("get job: %v", err) + } + events, err := s.GetJobEvents(pending.ID) + if err != nil { + t.Fatalf("events: %v", err) + } + if len(events) != 1 || events[0].EventType != EventEnqueued { + t.Errorf("expected the pending job to keep its enqueued event, got %+v", events) + } +} + +// TestPurgeBeforeKeepsRecentJobs verifies the age filter. A job whose last +// update is older than the cutoff is removed; a freshly updated job survives. +func TestPurgeBeforeKeepsRecentJobs(t *testing.T) { + s := newTestStore(t) + q := NewQueue(s) + ctx := context.Background() + + old, err := q.Enqueue("old", `{}`) + if err != nil { + t.Fatalf("enqueue old: %v", err) + } + if job, err := q.Lease(ctx, "old", time.Minute); err != nil || job == nil { + t.Fatalf("lease old: job=%v err=%v", job, err) + } else if err := q.Acknowledge(job.ID); err != nil { + t.Fatalf("ack old: %v", err) + } + // Backdate the completed job so it is older than the cutoff. + if _, err := s.db.Exec( + `UPDATE jobs SET updated_at = ? WHERE id = ?`, + time.Now().UTC().Add(-48*time.Hour).Format(sqliteTimeFormat), old.ID); err != nil { + t.Fatalf("backdate old: %v", err) + } + + fresh, err := q.Enqueue("fresh", `{}`) + if err != nil { + t.Fatalf("enqueue fresh: %v", err) + } + if job, err := q.Lease(ctx, "fresh", time.Minute); err != nil || job == nil { + t.Fatalf("lease fresh: job=%v err=%v", job, err) + } else if err := q.Acknowledge(job.ID); err != nil { + t.Fatalf("ack fresh: %v", err) + } + + cutoff := time.Now().UTC().Add(-24 * time.Hour) + stats, err := q.Purge(PurgeBefore(cutoff)) + if err != nil { + t.Fatalf("purge: %v", err) + } + if stats.JobsRemoved != 1 { + t.Errorf("expected 1 old job removed, got %d", stats.JobsRemoved) + } + if _, err := s.GetJob(old.ID); err == nil { + t.Errorf("old job must be gone") + } + if _, err := s.GetJob(fresh.ID); err != nil { + t.Errorf("fresh job must survive: %v", err) + } +} + +// TestPurgeCandidatesMatchesPurge verifies that the preview counts equal the +// counts a real purge removes. +func TestPurgeCandidatesMatchesPurge(t *testing.T) { + s, _ := seedPurgeLifecycle(t) + q := NewQueue(s) + + candidates, err := q.PurgeCandidates() + if err != nil { + t.Fatalf("candidates: %v", err) + } + if candidates.JobsRemoved != 2 || candidates.EventsRemoved != 6 { + t.Fatalf("unexpected candidates %+v", candidates) + } + + stats, err := q.Purge() + if err != nil { + t.Fatalf("purge: %v", err) + } + if stats != candidates { + t.Errorf("purge %+v does not match candidates %+v", stats, candidates) + } +} + +// TestPurgeCandidatesDoesNotMutate verifies that a preview leaves the store +// unchanged. +func TestPurgeCandidatesDoesNotMutate(t *testing.T) { + s, _ := seedPurgeLifecycle(t) + q := NewQueue(s) + + before, err := q.Inspect() + if err != nil { + t.Fatalf("inspect: %v", err) + } + if _, err := q.PurgeCandidates(); err != nil { + t.Fatalf("candidates: %v", err) + } + after, err := q.Inspect() + if err != nil { + t.Fatalf("inspect: %v", err) + } + if len(before.Jobs) != len(after.Jobs) { + t.Errorf("preview changed job count: %d -> %d", len(before.Jobs), len(after.Jobs)) + } + if len(before.Events) != len(after.Events) { + t.Errorf("preview changed event count: %d -> %d", len(before.Events), len(after.Events)) + } +} + +// TestPurgeIsIdempotent verifies that a second purge removes nothing. +func TestPurgeIsIdempotent(t *testing.T) { + s, _ := seedPurgeLifecycle(t) + q := NewQueue(s) + + if _, err := q.Purge(); err != nil { + t.Fatalf("first purge: %v", err) + } + second, err := q.Purge() + if err != nil { + t.Fatalf("second purge: %v", err) + } + if second.JobsRemoved != 0 || second.EventsRemoved != 0 { + t.Errorf("expected an empty second purge, got %+v", second) + } +} + +// TestPurgeRejectsNoStates verifies that an explicit empty state set removes +// nothing and reports no error. +func TestPurgeRejectsNoStates(t *testing.T) { + s, _ := seedPurgeLifecycle(t) + q := NewQueue(s) + + stats, err := q.Purge(PurgeStates()) + if err != nil { + t.Fatalf("purge: %v", err) + } + if stats.JobsRemoved != 0 || stats.EventsRemoved != 0 { + t.Errorf("expected no removal for an empty state set, got %+v", stats) + } + snap, err := q.Inspect() + if err != nil { + t.Fatalf("inspect: %v", err) + } + if len(snap.Jobs) != 4 { + t.Errorf("expected all 4 jobs to survive, got %d", len(snap.Jobs)) + } +} + +// TestPurgePendingIsExplicit verifies that naming pending works, so an operator +// can clear a backlog on purpose. +func TestPurgePendingIsExplicit(t *testing.T) { + s := newTestStore(t) + q := NewQueue(s) + + if _, err := q.Enqueue("stuck", `{}`); err != nil { + t.Fatalf("enqueue: %v", err) + } + if _, err := q.Enqueue("stuck", `{}`, WithPriority(5)); err != nil { + t.Fatalf("enqueue: %v", err) + } + + stats, err := q.Purge(PurgeStates(StatePending)) + if err != nil { + t.Fatalf("purge: %v", err) + } + if stats.JobsRemoved != 2 { + t.Errorf("expected 2 pending jobs removed, got %d", stats.JobsRemoved) + } + snap, err := q.Inspect() + if err != nil { + t.Fatalf("inspect: %v", err) + } + if len(snap.Jobs) != 0 { + t.Errorf("expected an empty queue, got %d jobs", len(snap.Jobs)) + } +} + +// TestPurgeBeforeBoundaryUsesStrictInequality verifies that a job updated at +// the cutoff survives because the filter is strictly older. +func TestPurgeBeforeBoundaryUsesStrictInequality(t *testing.T) { + s := newTestStore(t) + q := NewQueue(s) + ctx := context.Background() + + job, err := q.Enqueue("edge", `{}`) + if err != nil { + t.Fatalf("enqueue: %v", err) + } + if leased, err := q.Lease(ctx, "edge", time.Minute); err != nil || leased == nil { + t.Fatalf("lease: job=%v err=%v", leased, err) + } else if err := q.Acknowledge(leased.ID); err != nil { + t.Fatalf("ack: %v", err) + } + + // The cutoff equals the exact stored update time. The strict < comparison + // must keep the job at the boundary. + stored, err := s.GetJob(job.ID) + if err != nil { + t.Fatalf("get job: %v", err) + } + stats, err := q.Purge(PurgeStates(StateCompleted), PurgeBefore(stored.UpdatedAt.UTC())) + if err != nil { + t.Fatalf("purge: %v", err) + } + if stats.JobsRemoved != 0 { + t.Errorf("expected the boundary job to survive, removed %d", stats.JobsRemoved) + } + if _, err := s.GetJob(job.ID); err != nil { + t.Errorf("job must survive an equal cutoff: %v", err) + } +} + +func TestPurgeJobsUsesStrictInequality(t *testing.T) { + // This test pins the SQL comparison so the behavior cannot drift silently. + s := newTestStore(t) + q := NewQueue(s) + + if _, err := q.Enqueue("x", `{}`, WithPriority(1)); err != nil { + t.Fatalf("enqueue: %v", err) + } + stats, err := q.Purge(PurgeStates(StatePending), PurgeBefore(time.Now().UTC().Add(-time.Hour))) + if err != nil { + t.Fatalf("purge: %v", err) + } + if stats.JobsRemoved != 0 { + t.Errorf("expected 0 removed for a past cutoff, got %d", stats.JobsRemoved) + } +} + +func ExampleQueue_Purge() { + s, err := NewSQLiteStore("file:example_purge?mode=memory&cache=shared") + if err != nil { + fmt.Println("open:", err) + return + } + defer s.Close() + q := NewQueue(s) + + if _, err := q.Enqueue("email", `{}`); err != nil { + fmt.Println("enqueue:", err) + return + } + stats, err := q.Purge() + if err != nil { + fmt.Println("purge:", err) + return + } + fmt.Printf("removed %d jobs and %d events\n", stats.JobsRemoved, stats.EventsRemoved) + // Output: + // removed 0 jobs and 0 events +} diff --git a/main.go b/main.go index 9119ad9..aa79c48 100644 --- a/main.go +++ b/main.go @@ -22,6 +22,7 @@ Commands: history View the event log for one job requeue Return a dead-lettered job to the queue seed Load bundled sample data + purge Remove finished jobs and their events metrics Expose queue state for Prometheus demo Run a self-contained scenario with fault injection @@ -46,6 +47,8 @@ Use -help for command flags.`) err = cli.Requeue(args) case "seed": err = cli.Seed(args) + case "purge": + err = cli.Purge(args) case "metrics": err = cli.Metrics(args) case "demo":