diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..39b3c48 --- /dev/null +++ b/.gitignore @@ -0,0 +1,7 @@ +*.exe +*.exe~ +*.dll +*.so +*.dylib +*.test +*.out diff --git a/README.md b/README.md index ded17ed..9308ee0 100644 --- a/README.md +++ b/README.md @@ -275,7 +275,7 @@ gofmt -l . # formatting go vet ./... # static analysis ``` -16 tests cover priority ordering, retry/backoff, timeout enforcement, +20 tests cover priority ordering, retry/backoff, timeout enforcement, delayed scheduling, rate limiting, deduplication (including window expiry), success/failure chaining, dead-letter archival and replay, pagination, and metrics output. diff --git a/goqueue.exe b/goqueue.exe deleted file mode 100644 index 78b5f21..0000000 Binary files a/goqueue.exe and /dev/null differ diff --git a/handlers.go b/handlers.go index 2be8304..6a9b963 100644 --- a/handlers.go +++ b/handlers.go @@ -280,19 +280,8 @@ func (h *Handlers) dlqReplayHandler(w http.ResponseWriter, r *http.Request) { return } - entries, err := h.store.ListDeadLetters() + entry, err := h.store.GetDeadLetter(id) if err != nil { - respondError(w, http.StatusInternalServerError, err.Error()) - return - } - var entry *DeadLetter - for _, e := range entries { - if e.ID == id { - entry = e - break - } - } - if entry == nil { respondError(w, http.StatusNotFound, fmt.Sprintf("dead letter %s not found", id)) return } diff --git a/main.go b/main.go index fc20705..8f986a4 100644 --- a/main.go +++ b/main.go @@ -13,8 +13,6 @@ import ( ) func main() { - rand.Seed(time.Now().UnixNano()) - logger := slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{ Level: slog.LevelInfo, })) diff --git a/metrics.go b/metrics.go index 3e23de1..13f1e8f 100644 --- a/metrics.go +++ b/metrics.go @@ -9,6 +9,30 @@ import ( "time" ) +// promQuote escapes a label value for the Prometheus text exposition +// format: double-quoted string with \, ", and newlines backslash-escaped. +// Unlike Go's %q, this avoids Go-specific escape sequences like \x or \u. +func promQuote(s string) string { + var b strings.Builder + b.WriteByte('"') + for _, c := range s { + switch c { + case '\\': + b.WriteString(`\\`) + case '"': + b.WriteString(`\"`) + case '\n': + b.WriteString(`\n`) + case '\r': + b.WriteString(`\r`) + default: + b.WriteRune(c) + } + } + b.WriteByte('"') + return b.String() +} + // Metrics is a small, dependency-free counter/histogram store that // renders in the standard Prometheus text exposition format. // @@ -107,7 +131,7 @@ func (m *Metrics) ServeHTTP(w http.ResponseWriter, r *http.Request) { return keys[i].status < keys[j].status }) for _, k := range keys { - fmt.Fprintf(&b, "goqueue_jobs_total{type=%q,status=%q} %d\n", k.jobType, k.status, m.jobsTotal[k]) + fmt.Fprintf(&b, "goqueue_jobs_total{type=%s,status=%s} %d\n", promQuote(k.jobType), promQuote(k.status), m.jobsTotal[k]) } b.WriteString("# HELP goqueue_retries_total Total retry attempts, by job type.\n") @@ -118,7 +142,7 @@ func (m *Metrics) ServeHTTP(w http.ResponseWriter, r *http.Request) { } sort.Strings(types) for _, t := range types { - fmt.Fprintf(&b, "goqueue_retries_total{type=%q} %d\n", t, m.retriesTotal[t]) + fmt.Fprintf(&b, "goqueue_retries_total{type=%s} %d\n", promQuote(t), m.retriesTotal[t]) } b.WriteString("# HELP goqueue_job_duration_seconds Job execution duration in seconds, by type.\n") @@ -131,11 +155,11 @@ func (m *Metrics) ServeHTTP(w http.ResponseWriter, r *http.Request) { for _, t := range durTypes { buckets := m.durationBuckets[t] for i, le := range histogramBuckets { - fmt.Fprintf(&b, "goqueue_job_duration_seconds_bucket{type=%q,le=\"%g\"} %d\n", t, le, buckets[i]) + fmt.Fprintf(&b, "goqueue_job_duration_seconds_bucket{type=%s,le=\"%g\"} %d\n", promQuote(t), le, buckets[i]) } - fmt.Fprintf(&b, "goqueue_job_duration_seconds_bucket{type=%q,le=\"+Inf\"} %d\n", t, m.durationCount[t]) - fmt.Fprintf(&b, "goqueue_job_duration_seconds_sum{type=%q} %g\n", t, m.durationSum[t]) - fmt.Fprintf(&b, "goqueue_job_duration_seconds_count{type=%q} %d\n", t, m.durationCount[t]) + fmt.Fprintf(&b, "goqueue_job_duration_seconds_bucket{type=%s,le=\"+Inf\"} %d\n", promQuote(t), m.durationCount[t]) + fmt.Fprintf(&b, "goqueue_job_duration_seconds_sum{type=%s} %g\n", promQuote(t), m.durationSum[t]) + fmt.Fprintf(&b, "goqueue_job_duration_seconds_count{type=%s} %d\n", promQuote(t), m.durationCount[t]) } if m.queueDepthFn != nil { @@ -148,7 +172,7 @@ func (m *Metrics) ServeHTTP(w http.ResponseWriter, r *http.Request) { } sort.Strings(lanes) for _, lane := range lanes { - fmt.Fprintf(&b, "goqueue_queue_depth{priority=%q} %d\n", lane, depths[lane]) + fmt.Fprintf(&b, "goqueue_queue_depth{priority=%s} %d\n", promQuote(lane), depths[lane]) } } diff --git a/store.go b/store.go index f449466..ab804d6 100644 --- a/store.go +++ b/store.go @@ -47,6 +47,7 @@ type Store interface { Stats() (map[string]int, error) SaveDeadLetter(entry *DeadLetter) error + GetDeadLetter(id string) (*DeadLetter, error) ListDeadLetters() ([]*DeadLetter, error) DeleteDeadLetter(id string) error } @@ -167,6 +168,16 @@ func (s *InMemoryStore) SaveDeadLetter(entry *DeadLetter) error { return nil } +func (s *InMemoryStore) GetDeadLetter(id string) (*DeadLetter, error) { + s.mu.RLock() + defer s.mu.RUnlock() + entry, exists := s.deadLetters[id] + if !exists { + return nil, fmt.Errorf("dead letter %s not found", id) + } + return entry, nil +} + func (s *InMemoryStore) ListDeadLetters() ([]*DeadLetter, error) { s.mu.RLock() defer s.mu.RUnlock() @@ -401,10 +412,12 @@ FROM jobs` query += " ORDER BY created_at DESC" if filter.Limit > 0 { - query += fmt.Sprintf(" LIMIT %d", filter.Limit) + query += fmt.Sprintf(" LIMIT $%d", len(args)+1) + args = append(args, filter.Limit) } if filter.Offset > 0 { - query += fmt.Sprintf(" OFFSET %d", filter.Offset) + query += fmt.Sprintf(" OFFSET $%d", len(args)+1) + args = append(args, filter.Offset) } rows, err := s.db.Query(query, args...) @@ -489,6 +502,28 @@ ON CONFLICT (id) DO NOTHING return err } +func (s *PostgresStore) GetDeadLetter(id string) (*DeadLetter, error) { + row := s.db.QueryRow(` +SELECT id, job_id, type, payload, priority, max_retries, error, failed_at, originally_created_at +FROM dead_letters +WHERE id = $1 +`, id) + var d DeadLetter + var payloadBytes []byte + var priority int + if err := row.Scan(&d.ID, &d.JobID, &d.Type, &payloadBytes, &priority, &d.MaxRetries, &d.Error, &d.FailedAt, &d.OriginallyCreatedAt); err != nil { + if err == sql.ErrNoRows { + return nil, fmt.Errorf("dead letter %s not found", id) + } + return nil, err + } + d.Priority = JobPriority(priority) + if len(payloadBytes) > 0 { + _ = json.Unmarshal(payloadBytes, &d.Payload) + } + return &d, nil +} + func (s *PostgresStore) ListDeadLetters() ([]*DeadLetter, error) { rows, err := s.db.Query(` SELECT id, job_id, type, payload, priority, max_retries, error, failed_at, originally_created_at