Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1,7 @@
*.exe
*.exe~
*.dll
*.so
*.dylib
*.test
*.out
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Binary file removed goqueue.exe
Binary file not shown.
13 changes: 1 addition & 12 deletions handlers.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
2 changes: 0 additions & 2 deletions main.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,6 @@ import (
)

func main() {
rand.Seed(time.Now().UnixNano())

logger := slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{
Level: slog.LevelInfo,
}))
Expand Down
38 changes: 31 additions & 7 deletions metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
//
Expand Down Expand Up @@ -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")
Expand All @@ -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")
Expand All @@ -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 {
Expand All @@ -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])
}
}

Expand Down
39 changes: 37 additions & 2 deletions store.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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...)
Expand Down Expand Up @@ -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
Expand Down