diff --git a/README.md b/README.md index 4eb2a37..d694ba3 100644 --- a/README.md +++ b/README.md @@ -25,6 +25,9 @@ correctness model, then the network/protocol layer, and only later transparent Redis proxying and persistence. The current repo is already useful as a compact Go reference for reliable queue internals. +Moxy is `v0.1.0-alpha` software. It is useful for learning, experiments, and +small reliability-model prototypes, but it is not production-hardened yet. + ## The Pain Plain Redis list consumption often starts with `LPOP`. It is fast, simple, and @@ -115,6 +118,8 @@ to engines. Queue backends own task storage; the core engine owns lease metadata expiration scheduling. See [ARCHITECTURE.md](ARCHITECTURE.md) for the deeper system notes. +See [docs/redis-production-caveats.md](docs/redis-production-caveats.md) before +using the Redis backend for anything beyond experimentation. ## Internal Commands @@ -167,6 +172,7 @@ It is deterministic, easy to test, and useful for validating lease behavior. - `moxy:{queue}:ready` - `moxy:{queue}:processing` +- `moxy:{queue}:dead` `Acquire` uses `LMOVE ready processing RIGHT LEFT`. `Complete` and `Requeue` use Lua scripts so finding a task by ID and removing or moving it happens atomically inside @@ -267,7 +273,7 @@ implemented yet: - crash recovery - Redis Streams - distributed coordination -- user-facing network protocol commands +- additional Redis-compatible commands beyond `PING` and `MOXY.*` ## Roadmap diff --git a/docs/adr/0001-redis-lists-vs-streams.md b/docs/adr/0001-redis-lists-vs-streams.md new file mode 100644 index 0000000..d9fd35a --- /dev/null +++ b/docs/adr/0001-redis-lists-vs-streams.md @@ -0,0 +1,46 @@ +# ADR 0001: Redis Lists vs Streams + +## Status + +Accepted for the current alpha backend. + +## Context + +Moxy's first Redis backend uses Redis lists plus `LMOVE` and small Lua scripts. +The backend models tasks moving through explicit lifecycle storage: + +```text +READY -> PROCESSING -> ACK/REQUEUE +``` + +Redis Streams are also a valid Redis-native way to model reliable work queues. + +## Decision + +Keep the current Redis backend on lists for now: + +- Ready tasks live in a Redis list. +- Processing tasks live in a Redis list. +- `LMOVE` moves one task from ready to processing. +- Lua scripts perform bounded atomic transitions for ACK, requeue, and + dead-letter moves. + +## Why Lists First + +Lists are the simplest explicit baseline. They make `READY -> PROCESSING -> +ACK/REQUEUE` easy to reason about, map cleanly to the current `queue.Backend` +abstraction, and are a good educational first backend for Moxy's lease model. + +## Streams Alternative + +Redis Streams provide consumer groups, pending entries, `XACK`, `XAUTOCLAIM`, +and built-in tools for inspecting and reclaiming stuck messages. + +Streams are more Redis-native for reliable work queues and may be a better fit +for production-oriented deployments. + +## Consequences + +The lists backend stays small and understandable, but Moxy must own more queue +semantics in code and Lua. A future `RedisStreamsQueue` backend can be added +without changing `core.Engine` if it satisfies the same backend interface. diff --git a/docs/redis-production-caveats.md b/docs/redis-production-caveats.md new file mode 100644 index 0000000..1ae535d --- /dev/null +++ b/docs/redis-production-caveats.md @@ -0,0 +1,72 @@ +# Redis Production Caveats + +Moxy is currently `v0.1.0-alpha`. It is useful for learning, experimentation, +and validating queue semantics, but it is not production-hardened yet. + +## Delivery Semantics + +Moxy currently provides at-least-once delivery, not exactly-once delivery. A task +can be delivered more than once after worker crashes, lease expiration, +replication failover, or client retry behavior. Workers should be idempotent. + +## Redis Persistence + +Moxy cannot provide stronger durability than the configured Redis persistence. + +- With AOF `appendfsync everysec`, Redis can lose a small recent write window + during a crash. +- With AOF `appendfsync always`, Redis fsyncs every write and is safer, but + slower. +- RDB-only or loosely configured persistence can lose more queue state than a + durable queue workload usually expects. + +Choose Redis persistence settings based on the durability required by the queue. + +## Redis Eviction + +Queue Redis instances should not use cache-style eviction policies such as +`allkeys-lru` for queue data. Evicting queue keys can silently drop ready, +processing, or dead-lettered work. + +Prefer `noeviction` and a persistence-oriented Redis configuration for queue +workloads. + +## Replication And Failover + +Redis replication and failover can still produce duplicates or replayed work. +Writes acknowledged by a primary may not be present on a promoted replica, and +clients may retry operations around failover boundaries. + +Moxy expects workers to treat task handling as idempotent. + +## Redis Cluster Key Slots + +Redis Cluster requires all keys touched by one atomic operation to hash to the +same slot. Moxy queue keys use a shared hash tag per logical queue: + +```text +moxy:{}:ready +moxy:{}:processing +moxy:{}:dead +``` + +The content inside braces must be identical for all keys belonging to the same +logical queue. + +## Lua Scripts + +Redis executes Lua scripts atomically and blocks other server activity while a +script runs. Moxy scripts should remain short and bounded to queue transition +work such as ACK, requeue, and dead-letter moves. + +## Dead-Letter Queues + +Dead-letter queue support is being added so tasks that expire too many times can +move out of the retry loop instead of requeueing forever. This is a baseline +safety feature, not full crash recovery. + +## Streams + +Redis Streams are a valid alternative for reliable work queues. Consumer groups, +pending entries, `XACK`, and reclaim operations may make Streams a better fit for +some production systems. Moxy may explore a separate Redis Streams backend later. diff --git a/internal/command/handler_test.go b/internal/command/handler_test.go index 75b47ff..ddba1b6 100644 --- a/internal/command/handler_test.go +++ b/internal/command/handler_test.go @@ -3,6 +3,7 @@ package command import ( "errors" "testing" + "time" "github.com/an8kk/moxy/internal/core" "github.com/an8kk/moxy/internal/queue" @@ -138,6 +139,37 @@ func TestStatsReportsQueueState(t *testing.T) { if stats.Stats.Ready != 1 || stats.Stats.Processing != 1 || stats.Stats.ActiveLeases != 1 { t.Fatalf("stats = %+v, want ready=1 processing=1 active=1", stats.Stats) } + if stats.Stats.Dead != 0 { + t.Fatalf("dead count = %d, want 0", stats.Stats.Dead) + } +} + +func TestStatsReportsDeadCount(t *testing.T) { + svc := service.NewWithConfig(func(queueName string) queue.Backend { + return queue.NewMemoryQueue() + }, service.ServiceConfig{ + Engine: core.EngineConfig{MaxAttempts: 1}, + }) + handler := NewHandler(svc) + + if _, err := handler.Handle(Command{Name: "MOXY.ENQUEUE", Args: []string{"jobs", "payload"}}); err != nil { + t.Fatalf("enqueue returned error: %v", err) + } + fetch, err := handler.Handle(Command{Name: "MOXY.FETCH", Args: []string{"jobs", "1"}}) + if err != nil { + t.Fatalf("fetch returned error: %v", err) + } + if _, err := svc.ReapExpired(fetch.ExpiresAt.Add(time.Nanosecond)); err != nil { + t.Fatalf("reap expired returned error: %v", err) + } + + stats, err := handler.Handle(Command{Name: "MOXY.STATS", Args: []string{"jobs"}}) + if err != nil { + t.Fatalf("stats returned error: %v", err) + } + if stats.Stats.Dead != 1 { + t.Fatalf("dead count = %d, want 1", stats.Stats.Dead) + } } func newTestHandler() *Handler { diff --git a/internal/core/engine.go b/internal/core/engine.go index 918508d..6846f14 100644 --- a/internal/core/engine.go +++ b/internal/core/engine.go @@ -22,16 +22,24 @@ const requeueRetryDelay = time.Second type Stats struct { Ready int Processing int + Dead int ActiveLeases int ExpirationHeap int } +// EngineConfig controls lease expiration behavior. +type EngineConfig struct { + MaxAttempts int + RequeueRetryDelay time.Duration +} + // Engine owns all mutable in-memory queue, lease, and expiration state. type Engine struct { mu sync.Mutex ready queue.Backend leases map[string]*Lease expirations expirationHeap + config EngineConfig } // NewEngine creates an empty in-memory lease engine. @@ -41,17 +49,31 @@ func NewEngine() *Engine { // NewEngineWithBackend creates a single-queue lease engine over the provided backend. func NewEngineWithBackend(backend queue.Backend) *Engine { - return newEngine(backend) + return newEngine(backend, EngineConfig{}) +} + +// NewEngineWithBackendAndConfig creates a lease engine with explicit expiration config. +func NewEngineWithBackendAndConfig(backend queue.Backend, config EngineConfig) *Engine { + return newEngine(backend, config) } -func newEngine(backend queue.Backend) *Engine { +func newEngine(backend queue.Backend, configs ...EngineConfig) *Engine { + var config EngineConfig + if len(configs) > 0 { + config = configs[0] + } + expirations := expirationHeap{} heap.Init(&expirations) + if config.RequeueRetryDelay <= 0 { + config.RequeueRetryDelay = requeueRetryDelay + } return &Engine{ ready: backend, leases: make(map[string]*Lease), expirations: expirations, + config: config, } } @@ -141,10 +163,11 @@ func (e *Engine) ReapExpired(now time.Time) (int, error) { continue } - if err := e.ready.Requeue(lease.Task.ID); err != nil { + err := e.expireLease(lease) + if err != nil { heap.Push(&e.expirations, expirationItem{ LeaseID: item.LeaseID, - ExpiresAt: now.Add(requeueRetryDelay), + ExpiresAt: now.Add(e.config.RequeueRetryDelay), }) return requeued, err } @@ -155,6 +178,18 @@ func (e *Engine) ReapExpired(now time.Time) (int, error) { return requeued, nil } +func (e *Engine) expireLease(lease *Lease) error { + if e.shouldDeadLetter(lease) { + return e.ready.DeadLetter(lease.Task.ID, "max attempts exceeded") + } + + return e.ready.Requeue(lease.Task.ID) +} + +func (e *Engine) shouldDeadLetter(lease *Lease) bool { + return e.config.MaxAttempts > 0 && lease.Task.Attempts+1 >= e.config.MaxAttempts +} + // Stats returns counts for tests and diagnostics. func (e *Engine) Stats() Stats { e.mu.Lock() @@ -164,6 +199,7 @@ func (e *Engine) Stats() Stats { return Stats{ Ready: queueStats.Ready, Processing: queueStats.Processing, + Dead: queueStats.Dead, ActiveLeases: len(e.leases), ExpirationHeap: e.expirations.Len(), } @@ -184,8 +220,9 @@ func cloneLease(lease *Lease) *Lease { func cloneTask(task Task) Task { return Task{ - ID: task.ID, - Payload: cloneBytes(task.Payload), + ID: task.ID, + Payload: cloneBytes(task.Payload), + Attempts: task.Attempts, } } diff --git a/internal/core/engine_test.go b/internal/core/engine_test.go index f00f308..f97842a 100644 --- a/internal/core/engine_test.go +++ b/internal/core/engine_test.go @@ -78,6 +78,81 @@ func TestAckRemovesLease(t *testing.T) { } } +func TestAckBeforeReapExpiredDoesNotRequeueTask(t *testing.T) { + engine := NewEngine() + engine.Enqueue([]byte("first")) + lease, err := engine.Fetch(time.Millisecond) + if err != nil { + t.Fatalf("fetch returned error: %v", err) + } + + if err := engine.Ack(lease.LeaseID); err != nil { + t.Fatalf("ack returned error: %v", err) + } + requeued, err := engine.ReapExpired(lease.ExpiresAt.Add(time.Hour)) + if err != nil { + t.Fatalf("reap expired returned error: %v", err) + } + if requeued != 0 { + t.Fatalf("requeued count = %d, want 0", requeued) + } + + stats := engine.Stats() + if stats.Ready != 0 || stats.Processing != 0 || stats.ActiveLeases != 0 || stats.Dead != 0 { + t.Fatalf("stats after ack-before-reap = %+v, want ready=0 processing=0 active=0 dead=0", stats) + } +} + +func TestReapExpiredBeforeAckRequeuesTaskAndAckFails(t *testing.T) { + engine := NewEngine() + engine.Enqueue([]byte("first")) + lease, err := engine.Fetch(time.Millisecond) + if err != nil { + t.Fatalf("fetch returned error: %v", err) + } + + requeued, err := engine.ReapExpired(lease.ExpiresAt.Add(time.Nanosecond)) + if err != nil { + t.Fatalf("reap expired returned error: %v", err) + } + if requeued != 1 { + t.Fatalf("requeued count = %d, want 1", requeued) + } + if err := engine.Ack(lease.LeaseID); !errors.Is(err, ErrLeaseNotFound) { + t.Fatalf("ack after reap returned %v, want ErrLeaseNotFound", err) + } + + stats := engine.Stats() + if stats.Ready != 1 || stats.Processing != 0 || stats.ActiveLeases != 0 || stats.Dead != 0 { + t.Fatalf("stats after reap-before-ack = %+v, want ready=1 processing=0 active=0 dead=0", stats) + } +} + +func TestReapExpiredBeforeAckDeadLettersWhenMaxAttemptsExceeded(t *testing.T) { + engine := NewEngineWithBackendAndConfig(queue.NewMemoryQueue(), EngineConfig{MaxAttempts: 1}) + engine.Enqueue([]byte("first")) + lease, err := engine.Fetch(time.Millisecond) + if err != nil { + t.Fatalf("fetch returned error: %v", err) + } + + requeued, err := engine.ReapExpired(lease.ExpiresAt.Add(time.Nanosecond)) + if err != nil { + t.Fatalf("reap expired returned error: %v", err) + } + if requeued != 1 { + t.Fatalf("reap count = %d, want 1", requeued) + } + if err := engine.Ack(lease.LeaseID); !errors.Is(err, ErrLeaseNotFound) { + t.Fatalf("ack after dead letter returned %v, want ErrLeaseNotFound", err) + } + + stats := engine.Stats() + if stats.Ready != 0 || stats.Processing != 0 || stats.ActiveLeases != 0 || stats.Dead != 1 { + t.Fatalf("stats after reap-before-ack dead letter = %+v, want ready=0 processing=0 active=0 dead=1", stats) + } +} + func TestFetchInvalidTimeoutFails(t *testing.T) { engine := NewEngine() enqueue(t, engine, []byte("first")) @@ -218,6 +293,91 @@ func TestExpiredLeaseRequeuesTask(t *testing.T) { } } +func TestDefaultMaxAttemptsAllowsRepeatedRequeue(t *testing.T) { + engine := NewEngine() + task := enqueue(t, engine, []byte("first")) + + for attempt := 1; attempt <= 2; attempt++ { + lease, err := engine.Fetch(time.Millisecond) + if err != nil { + t.Fatalf("fetch attempt %d returned error: %v", attempt, err) + } + requeued, err := engine.ReapExpired(lease.ExpiresAt.Add(time.Nanosecond)) + if err != nil { + t.Fatalf("reap attempt %d returned error: %v", attempt, err) + } + if requeued != 1 { + t.Fatalf("reap attempt %d count = %d, want 1", attempt, requeued) + } + stats := engine.Stats() + if stats.Ready != 1 || stats.Dead != 0 { + t.Fatalf("stats after attempt %d = %+v, want ready=1 dead=0", attempt, stats) + } + } + + lease, err := engine.Fetch(time.Minute) + if err != nil { + t.Fatalf("fetch final returned error: %v", err) + } + if lease.Task.ID != task.ID || lease.Task.Attempts != 2 { + t.Fatalf("final task = %+v, want id=%s attempts=2", lease.Task, task.ID) + } +} + +func TestMaxAttemptsOneMovesExpiredTaskToDeadLetter(t *testing.T) { + engine := NewEngineWithBackendAndConfig(queue.NewMemoryQueue(), EngineConfig{MaxAttempts: 1}) + engine.Enqueue([]byte("first")) + lease, err := engine.Fetch(time.Millisecond) + if err != nil { + t.Fatalf("fetch returned error: %v", err) + } + + requeued, err := engine.ReapExpired(lease.ExpiresAt.Add(time.Nanosecond)) + if err != nil { + t.Fatalf("reap expired returned error: %v", err) + } + if requeued != 1 { + t.Fatalf("reap count = %d, want 1", requeued) + } + + stats := engine.Stats() + if stats.Ready != 0 || stats.Processing != 0 || stats.ActiveLeases != 0 || stats.Dead != 1 { + t.Fatalf("stats after max attempts = %+v, want ready=0 processing=0 active=0 dead=1", stats) + } +} + +func TestMaxAttemptsTwoRequeuesOnceThenDeadLetters(t *testing.T) { + engine := NewEngineWithBackendAndConfig(queue.NewMemoryQueue(), EngineConfig{MaxAttempts: 2}) + engine.Enqueue([]byte("first")) + + first, err := engine.Fetch(time.Millisecond) + if err != nil { + t.Fatalf("fetch first returned error: %v", err) + } + if _, err := engine.ReapExpired(first.ExpiresAt.Add(time.Nanosecond)); err != nil { + t.Fatalf("first reap returned error: %v", err) + } + stats := engine.Stats() + if stats.Ready != 1 || stats.Dead != 0 { + t.Fatalf("stats after first reap = %+v, want ready=1 dead=0", stats) + } + + second, err := engine.Fetch(time.Millisecond) + if err != nil { + t.Fatalf("fetch second returned error: %v", err) + } + if second.Task.Attempts != 1 { + t.Fatalf("second lease attempts = %d, want 1", second.Task.Attempts) + } + if _, err := engine.ReapExpired(second.ExpiresAt.Add(time.Nanosecond)); err != nil { + t.Fatalf("second reap returned error: %v", err) + } + stats = engine.Stats() + if stats.Ready != 0 || stats.Processing != 0 || stats.ActiveLeases != 0 || stats.Dead != 1 { + t.Fatalf("stats after second reap = %+v, want ready=0 processing=0 active=0 dead=1", stats) + } +} + func TestDoubleAckFails(t *testing.T) { engine := NewEngine() enqueue(t, engine, []byte("first")) @@ -393,6 +553,80 @@ func TestReapExpiredRetriesFailedBackendRequeueLater(t *testing.T) { } } +func TestReapExpiredDoesNotDeleteLeaseIfBackendDeadLetterFails(t *testing.T) { + backend := queue.NewMemoryQueue() + engine := NewEngineWithBackendAndConfig(backend, EngineConfig{MaxAttempts: 1}) + engine.Enqueue([]byte("first")) + lease, err := engine.Fetch(time.Millisecond) + if err != nil { + t.Fatalf("fetch returned error: %v", err) + } + + engine.ready = &failingDeadLetterBackend{Backend: backend} + + requeued, err := engine.ReapExpired(lease.ExpiresAt.Add(time.Nanosecond)) + if !errors.Is(err, errDeadLetterFailed) { + t.Fatalf("reap expired error = %v, want errDeadLetterFailed", err) + } + if requeued != 0 { + t.Fatalf("requeued count = %d, want 0", requeued) + } + + stats := engine.Stats() + if stats.ActiveLeases != 1 || stats.Processing != 1 || stats.ExpirationHeap != 1 || stats.Dead != 0 { + t.Fatalf("stats after failed dead letter = %+v, want active=1 processing=1 heap=1 dead=0", stats) + } +} + +func TestReapExpiredRetriesFailedBackendDeadLetterLater(t *testing.T) { + backend := queue.NewMemoryQueue() + engine := NewEngineWithBackendAndConfig(backend, EngineConfig{ + MaxAttempts: 1, + RequeueRetryDelay: time.Second, + }) + engine.Enqueue([]byte("first")) + lease, err := engine.Fetch(time.Millisecond) + if err != nil { + t.Fatalf("fetch returned error: %v", err) + } + + failing := &failOnceDeadLetterBackend{Backend: backend} + engine.ready = failing + + firstReapAt := lease.ExpiresAt.Add(time.Nanosecond) + requeued, err := engine.ReapExpired(firstReapAt) + if !errors.Is(err, errDeadLetterFailed) { + t.Fatalf("first reap error = %v, want errDeadLetterFailed", err) + } + if requeued != 0 { + t.Fatalf("first reap count = %d, want 0", requeued) + } + + requeued, err = engine.ReapExpired(firstReapAt.Add(500 * time.Millisecond)) + if err != nil { + t.Fatalf("second reap before retry delay returned error: %v", err) + } + if requeued != 0 { + t.Fatalf("second reap count = %d, want 0", requeued) + } + stats := engine.Stats() + if stats.ActiveLeases != 1 || stats.Processing != 1 || stats.Dead != 0 { + t.Fatalf("stats before retry = %+v, want active=1 processing=1 dead=0", stats) + } + + requeued, err = engine.ReapExpired(firstReapAt.Add(2 * time.Second)) + if err != nil { + t.Fatalf("retry reap returned error: %v", err) + } + if requeued != 1 { + t.Fatalf("retry reap count = %d, want 1", requeued) + } + stats = engine.Stats() + if stats.ActiveLeases != 0 || stats.Processing != 0 || stats.Dead != 1 { + t.Fatalf("stats after retry = %+v, want active=0 processing=0 dead=1", stats) + } +} + func TestMultipleLeaseExpirationOrdering(t *testing.T) { engine := NewEngine() enqueue(t, engine, []byte("first")) @@ -509,3 +743,27 @@ func (b *failOnceRequeueBackend) Requeue(taskID string) error { return b.Backend.Requeue(taskID) } + +var errDeadLetterFailed = errors.New("dead letter failed") + +type failingDeadLetterBackend struct { + queue.Backend +} + +func (b *failingDeadLetterBackend) DeadLetter(taskID string, reason string) error { + return errDeadLetterFailed +} + +type failOnceDeadLetterBackend struct { + queue.Backend + failed bool +} + +func (b *failOnceDeadLetterBackend) DeadLetter(taskID string, reason string) error { + if !b.failed { + b.failed = true + return errDeadLetterFailed + } + + return b.Backend.DeadLetter(taskID, reason) +} diff --git a/internal/core/redis_integration_test.go b/internal/core/redis_integration_test.go index 896aefe..cb40d44 100644 --- a/internal/core/redis_integration_test.go +++ b/internal/core/redis_integration_test.go @@ -106,7 +106,8 @@ func cleanupRedisEngineQueue(t *testing.T, client *redis.Client, queueName strin defer cancel() readyKey := fmt.Sprintf("moxy:{%s}:ready", queueName) processingKey := fmt.Sprintf("moxy:{%s}:processing", queueName) - if err := client.Del(ctx, readyKey, processingKey).Err(); err != nil { + deadKey := fmt.Sprintf("moxy:{%s}:dead", queueName) + if err := client.Del(ctx, readyKey, processingKey, deadKey).Err(); err != nil { t.Fatalf("cleanup redis engine queue: %v", err) } } diff --git a/internal/protocol/adapter.go b/internal/protocol/adapter.go index a20bf87..9b6a19e 100644 --- a/internal/protocol/adapter.go +++ b/internal/protocol/adapter.go @@ -88,6 +88,8 @@ func responseToRESP(name string, response command.Response) resp.Value { resp.Integer(int64(response.Stats.Ready)), resp.BulkString("processing"), resp.Integer(int64(response.Stats.Processing)), + resp.BulkString("dead"), + resp.Integer(int64(response.Stats.Dead)), resp.BulkString("active_leases"), resp.Integer(int64(response.Stats.ActiveLeases)), resp.BulkString("heap"), diff --git a/internal/protocol/adapter_test.go b/internal/protocol/adapter_test.go index 4614dbc..ad6cf07 100644 --- a/internal/protocol/adapter_test.go +++ b/internal/protocol/adapter_test.go @@ -78,15 +78,17 @@ func TestRESPMoxyStatsReturnsQueueState(t *testing.T) { reply := adapter.Handle(resp.Array(resp.BulkString("MOXY.STATS"), resp.BulkString("jobs"))) - assertArrayLength(t, reply, 8) + assertArrayLength(t, reply, 10) assertBulkString(t, reply.Array[0], "ready") assertInteger(t, reply.Array[1], 1) assertBulkString(t, reply.Array[2], "processing") assertInteger(t, reply.Array[3], 1) - assertBulkString(t, reply.Array[4], "active_leases") - assertInteger(t, reply.Array[5], 1) - assertBulkString(t, reply.Array[6], "heap") + assertBulkString(t, reply.Array[4], "dead") + assertInteger(t, reply.Array[5], 0) + assertBulkString(t, reply.Array[6], "active_leases") assertInteger(t, reply.Array[7], 1) + assertBulkString(t, reply.Array[8], "heap") + assertInteger(t, reply.Array[9], 1) } func TestRESPInvalidCommandReturnsError(t *testing.T) { diff --git a/internal/queue/backend.go b/internal/queue/backend.go index d8c6d93..6af2772 100644 --- a/internal/queue/backend.go +++ b/internal/queue/backend.go @@ -2,6 +2,7 @@ package queue import ( "errors" + "time" "github.com/an8kk/moxy/internal/task" ) @@ -15,6 +16,14 @@ var ( type Stats struct { Ready int Processing int + Dead int +} + +// DeadTask records a task that has left the retry loop. +type DeadTask struct { + Task task.Task `json:"task"` + Reason string `json:"reason"` + DeadAt time.Time `json:"dead_at"` } // Backend is the minimal ready-queue storage boundary used by the core engine. @@ -23,5 +32,6 @@ type Backend interface { Acquire() (task.Task, error) Complete(taskID string) error Requeue(taskID string) error + DeadLetter(taskID string, reason string) error Stats() Stats } diff --git a/internal/queue/contract_test.go b/internal/queue/contract_test.go index 96c7267..d690293 100644 --- a/internal/queue/contract_test.go +++ b/internal/queue/contract_test.go @@ -31,9 +31,18 @@ func RunBackendContractTests(t *testing.T, factory backendFactory) { t.Run("RequeueMovesProcessingTaskBackToReady", func(t *testing.T) { testRequeueMovesProcessingTaskBackToReady(t, factory) }) + t.Run("RequeueIncrementsAttempts", func(t *testing.T) { + testRequeueIncrementsAttempts(t, factory) + }) t.Run("RequeueMissingProcessingTaskFails", func(t *testing.T) { testRequeueMissingProcessingTaskFails(t, factory) }) + t.Run("DeadLetterMovesProcessingTaskToDead", func(t *testing.T) { + testDeadLetterMovesProcessingTaskToDead(t, factory) + }) + t.Run("DeadLetterMissingProcessingTaskFails", func(t *testing.T) { + testDeadLetterMissingProcessingTaskFails(t, factory) + }) t.Run("StatsReportsReadyAndProcessing", func(t *testing.T) { testStatsReportsReadyAndProcessing(t, factory) }) @@ -172,11 +181,61 @@ func testRequeueMissingProcessingTaskFails(t *testing.T, factory backendFactory) } } +func testRequeueIncrementsAttempts(t *testing.T, factory backendFactory) { + backend := factory(t) + if err := backend.Enqueue(task.Task{ID: "task-1", Attempts: 2}); err != nil { + t.Fatalf("enqueue returned error: %v", err) + } + acquired, err := backend.Acquire() + if err != nil { + t.Fatalf("acquire returned error: %v", err) + } + + if err := backend.Requeue(acquired.ID); err != nil { + t.Fatalf("requeue returned error: %v", err) + } + again, err := backend.Acquire() + if err != nil { + t.Fatalf("acquire after requeue returned error: %v", err) + } + if again.Attempts != 3 { + t.Fatalf("attempts after requeue = %d, want 3", again.Attempts) + } +} + +func testDeadLetterMovesProcessingTaskToDead(t *testing.T, factory backendFactory) { + backend := factory(t) + if err := backend.Enqueue(task.Task{ID: "task-1", Attempts: 1}); err != nil { + t.Fatalf("enqueue returned error: %v", err) + } + acquired, err := backend.Acquire() + if err != nil { + t.Fatalf("acquire returned error: %v", err) + } + + if err := backend.DeadLetter(acquired.ID, "expired"); err != nil { + t.Fatalf("dead letter returned error: %v", err) + } + + stats := backend.Stats() + if stats.Ready != 0 || stats.Processing != 0 || stats.Dead != 1 { + t.Fatalf("stats after dead letter = %+v, want ready=0 processing=0 dead=1", stats) + } +} + +func testDeadLetterMissingProcessingTaskFails(t *testing.T, factory backendFactory) { + backend := factory(t) + + if err := backend.DeadLetter("missing", "expired"); !errors.Is(err, ErrTaskNotProcessing) { + t.Fatalf("dead letter returned %v, want ErrTaskNotProcessing", err) + } +} + func testStatsReportsReadyAndProcessing(t *testing.T, factory backendFactory) { backend := factory(t) - if got := backend.Stats(); got.Ready != 0 || got.Processing != 0 { - t.Fatalf("initial stats = %+v, want ready=0 processing=0", got) + if got := backend.Stats(); got.Ready != 0 || got.Processing != 0 || got.Dead != 0 { + t.Fatalf("initial stats = %+v, want ready=0 processing=0 dead=0", got) } if err := backend.Enqueue(task.Task{ID: "task-1"}); err != nil { t.Fatalf("enqueue first returned error: %v", err) @@ -184,22 +243,22 @@ func testStatsReportsReadyAndProcessing(t *testing.T, factory backendFactory) { if err := backend.Enqueue(task.Task{ID: "task-2"}); err != nil { t.Fatalf("enqueue second returned error: %v", err) } - if got := backend.Stats(); got.Ready != 2 || got.Processing != 0 { - t.Fatalf("stats after enqueue = %+v, want ready=2 processing=0", got) + if got := backend.Stats(); got.Ready != 2 || got.Processing != 0 || got.Dead != 0 { + t.Fatalf("stats after enqueue = %+v, want ready=2 processing=0 dead=0", got) } acquired, err := backend.Acquire() if err != nil { t.Fatalf("acquire returned error: %v", err) } - if got := backend.Stats(); got.Ready != 1 || got.Processing != 1 { - t.Fatalf("stats after acquire = %+v, want ready=1 processing=1", got) + if got := backend.Stats(); got.Ready != 1 || got.Processing != 1 || got.Dead != 0 { + t.Fatalf("stats after acquire = %+v, want ready=1 processing=1 dead=0", got) } if err := backend.Complete(acquired.ID); err != nil { t.Fatalf("complete returned error: %v", err) } - if got := backend.Stats(); got.Ready != 1 || got.Processing != 0 { - t.Fatalf("stats after complete = %+v, want ready=1 processing=0", got) + if got := backend.Stats(); got.Ready != 1 || got.Processing != 0 || got.Dead != 0 { + t.Fatalf("stats after complete = %+v, want ready=1 processing=0 dead=0", got) } } diff --git a/internal/queue/memory.go b/internal/queue/memory.go index a706f2b..6357b57 100644 --- a/internal/queue/memory.go +++ b/internal/queue/memory.go @@ -2,6 +2,7 @@ package queue import ( "sync" + "time" "github.com/an8kk/moxy/internal/task" ) @@ -11,6 +12,7 @@ type MemoryQueue struct { mu sync.Mutex ready []task.Task processing map[string]task.Task + dead []DeadTask } // NewMemoryQueue creates an empty in-memory queue backend. @@ -18,6 +20,7 @@ func NewMemoryQueue() *MemoryQueue { return &MemoryQueue{ ready: make([]task.Task, 0), processing: make(map[string]task.Task), + dead: make([]DeadTask, 0), } } @@ -73,11 +76,32 @@ func (q *MemoryQueue) Requeue(taskID string) error { } delete(q.processing, taskID) + task.Attempts++ q.ready = append(q.ready, cloneTask(task)) return nil } -// Stats reports ready and processing task counts. +// DeadLetter moves a processing task into dead-letter storage. +func (q *MemoryQueue) DeadLetter(taskID string, reason string) error { + q.mu.Lock() + defer q.mu.Unlock() + + task, ok := q.processing[taskID] + if !ok { + return ErrTaskNotProcessing + } + + delete(q.processing, taskID) + task.Attempts++ + q.dead = append(q.dead, DeadTask{ + Task: cloneTask(task), + Reason: reason, + DeadAt: time.Now().UTC(), + }) + return nil +} + +// Stats reports ready, processing, and dead task counts. func (q *MemoryQueue) Stats() Stats { q.mu.Lock() defer q.mu.Unlock() @@ -85,13 +109,15 @@ func (q *MemoryQueue) Stats() Stats { return Stats{ Ready: len(q.ready), Processing: len(q.processing), + Dead: len(q.dead), } } func cloneTask(item task.Task) task.Task { return task.Task{ - ID: item.ID, - Payload: cloneBytes(item.Payload), + ID: item.ID, + Payload: cloneBytes(item.Payload), + Attempts: item.Attempts, } } diff --git a/internal/queue/memory_test.go b/internal/queue/memory_test.go index 03860c6..fb9d046 100644 --- a/internal/queue/memory_test.go +++ b/internal/queue/memory_test.go @@ -1,6 +1,11 @@ package queue -import "testing" +import ( + "errors" + "testing" + + "github.com/an8kk/moxy/internal/task" +) func TestMemoryQueueContract(t *testing.T) { RunBackendContractTests(t, func(t *testing.T) Backend { @@ -21,6 +26,10 @@ func TestRequeueMovesProcessingTaskBackToReady(t *testing.T) { testRequeueMovesProcessingTaskBackToReady(t, newMemoryBackend) } +func TestMemoryQueueRequeueIncrementsAttempts(t *testing.T) { + testRequeueIncrementsAttempts(t, newMemoryBackend) +} + func TestCompleteMissingProcessingTaskFails(t *testing.T) { testCompleteMissingProcessingTaskFails(t, newMemoryBackend) } @@ -33,6 +42,63 @@ func TestMemoryQueueStatsReportsReadyAndProcessing(t *testing.T) { testStatsReportsReadyAndProcessing(t, newMemoryBackend) } +func TestMemoryQueueDeadLetterMovesTaskFromProcessingToDead(t *testing.T) { + queue := NewMemoryQueue() + if err := queue.Enqueue(task.Task{ID: "task-1", Payload: []byte("payload")}); err != nil { + t.Fatalf("enqueue returned error: %v", err) + } + acquired, err := queue.Acquire() + if err != nil { + t.Fatalf("acquire returned error: %v", err) + } + + if err := queue.DeadLetter(acquired.ID, "expired"); err != nil { + t.Fatalf("dead letter returned error: %v", err) + } + + if got := queue.Stats(); got.Dead != 1 || got.Processing != 0 { + t.Fatalf("stats after dead letter = %+v, want dead=1 processing=0", got) + } + if len(queue.dead) != 1 { + t.Fatalf("dead storage length = %d, want 1", len(queue.dead)) + } + if queue.dead[0].Task.Attempts != 1 { + t.Fatalf("dead task attempts = %d, want 1", queue.dead[0].Task.Attempts) + } + if queue.dead[0].Reason != "expired" { + t.Fatalf("dead reason = %q, want %q", queue.dead[0].Reason, "expired") + } +} + +func TestMemoryQueueDeadLetterMissingTaskReturnsErrTaskNotProcessing(t *testing.T) { + queue := NewMemoryQueue() + + if err := queue.DeadLetter("missing", "expired"); !errors.Is(err, ErrTaskNotProcessing) { + t.Fatalf("dead letter returned %v, want ErrTaskNotProcessing", err) + } +} + +func TestMemoryQueueDeadLetterClonesPayload(t *testing.T) { + queue := NewMemoryQueue() + payload := []byte("payload") + if err := queue.Enqueue(task.Task{ID: "task-1", Payload: payload}); err != nil { + t.Fatalf("enqueue returned error: %v", err) + } + acquired, err := queue.Acquire() + if err != nil { + t.Fatalf("acquire returned error: %v", err) + } + acquired.Payload[0] = 'P' + + if err := queue.DeadLetter(acquired.ID, "expired"); err != nil { + t.Fatalf("dead letter returned error: %v", err) + } + + if string(queue.dead[0].Task.Payload) != "payload" { + t.Fatalf("stored dead payload = %q, want payload", queue.dead[0].Task.Payload) + } +} + func TestMemoryQueuePayloadCloningStillWorks(t *testing.T) { testPayloadCloningPreventsExternalMutation(t, newMemoryBackend) } diff --git a/internal/queue/redis.go b/internal/queue/redis.go index 43437f2..0d09ba8 100644 --- a/internal/queue/redis.go +++ b/internal/queue/redis.go @@ -5,6 +5,8 @@ import ( "encoding/json" "errors" "fmt" + "strings" + "time" "github.com/an8kk/moxy/internal/task" "github.com/redis/go-redis/v9" @@ -15,16 +17,47 @@ type RedisQueue struct { client *redis.Client readyKey string processingKey string + deadKey string + scripts redisQueueScripts } func NewRedisQueue(client *redis.Client, queueName string) *RedisQueue { + keys := newRedisQueueKeys(queueName) return &RedisQueue{ client: client, - readyKey: fmt.Sprintf("moxy:{%s}:ready", queueName), - processingKey: fmt.Sprintf("moxy:{%s}:processing", queueName), + readyKey: keys.ready, + processingKey: keys.processing, + deadKey: keys.dead, + scripts: defaultRedisQueueScripts(), } } +type redisQueueKeys struct { + ready string + processing string + dead string +} + +func newRedisQueueKeys(queueName string) redisQueueKeys { + hashTag := fmt.Sprintf("{%s}", queueName) + prefix := "moxy:" + hashTag + return redisQueueKeys{ + ready: prefix + ":ready", + processing: prefix + ":processing", + dead: prefix + ":dead", + } +} + +func redisHashTag(key string) string { + start := strings.IndexByte(key, '{') + end := strings.IndexByte(key, '}') + if start < 0 || end <= start { + return "" + } + + return key[start : end+1] +} + // Enqueue serializes a task and appends it to Redis ready storage. func (q *RedisQueue) Enqueue(task task.Task) error { encoded, err := encodeTask(task) @@ -50,7 +83,7 @@ func (q *RedisQueue) Acquire() (task.Task, error) { // Complete atomically removes a processing task. func (q *RedisQueue) Complete(taskID string) error { - found, err := q.runTaskScript(completeTaskScript, taskID) + found, err := q.runTaskScript(q.scripts.complete, taskID) if err != nil { return err } @@ -63,7 +96,7 @@ func (q *RedisQueue) Complete(taskID string) error { // Requeue atomically moves a processing task back to ready storage. func (q *RedisQueue) Requeue(taskID string) error { - found, err := q.runTaskScript(requeueTaskScript, taskID) + found, err := q.runTaskScript(q.scripts.requeue, taskID) if err != nil { return err } @@ -74,15 +107,42 @@ func (q *RedisQueue) Requeue(taskID string) error { return nil } -// Stats reports Redis ready and processing list lengths. +// DeadLetter atomically moves a processing task to dead-letter storage. +func (q *RedisQueue) DeadLetter(taskID string, reason string) error { + needle, err := taskIDNeedle(taskID) + if err != nil { + return err + } + + result, err := q.scripts.dead.Run( + context.Background(), + q.client, + []string{q.processingKey, q.deadKey}, + needle, + reason, + time.Now().UTC().Format(time.RFC3339Nano), + ).Int() + if err != nil { + return err + } + if result == 0 { + return ErrTaskNotProcessing + } + + return nil +} + +// Stats reports Redis ready, processing, and dead list lengths. func (q *RedisQueue) Stats() Stats { ctx := context.Background() ready := q.client.LLen(ctx, q.readyKey).Val() processing := q.client.LLen(ctx, q.processingKey).Val() + dead := q.client.LLen(ctx, q.deadKey).Val() return Stats{ Ready: int(ready), Processing: int(processing), + Dead: int(dead), } } @@ -111,7 +171,7 @@ func taskIDNeedle(taskID string) (string, error) { return "", err } - return `"ID":` + string(encoded), nil + return `"id":` + string(encoded), nil } func encodeTask(item task.Task) (string, error) { @@ -132,40 +192,12 @@ func decodeTask(encoded string) (task.Task, error) { return cloneTask(decoded), nil } -var completeTaskScript = redis.NewScript(` -local processing = KEYS[1] -local needle = ARGV[1] -local tasks = redis.call("LRANGE", processing, 0, -1) - -for _, encoded in ipairs(tasks) do - if string.find(encoded, needle, 1, true) then - local removed = redis.call("LREM", processing, 1, encoded) - if removed == 0 then - return 0 - end - return 1 - end -end - -return 0 -`) - -var requeueTaskScript = redis.NewScript(` -local processing = KEYS[1] -local ready = KEYS[2] -local needle = ARGV[1] -local tasks = redis.call("LRANGE", processing, 0, -1) - -for _, encoded in ipairs(tasks) do - if string.find(encoded, needle, 1, true) then - local removed = redis.call("LREM", processing, 1, encoded) - if removed == 0 then - return 0 - end - redis.call("LPUSH", ready, encoded) - return 1 - end -end - -return 0 -`) +func decodeDeadTask(encoded string) (DeadTask, error) { + var decoded DeadTask + if err := json.Unmarshal([]byte(encoded), &decoded); err != nil { + return DeadTask{}, err + } + + decoded.Task = cloneTask(decoded.Task) + return decoded, nil +} diff --git a/internal/queue/redis_keys_test.go b/internal/queue/redis_keys_test.go new file mode 100644 index 0000000..1c931b3 --- /dev/null +++ b/internal/queue/redis_keys_test.go @@ -0,0 +1,28 @@ +package queue + +import "testing" + +func TestRedisQueueKeysUseClusterSafeHashTag(t *testing.T) { + keys := newRedisQueueKeys("jobs") + + if keys.ready != "moxy:{jobs}:ready" { + t.Fatalf("ready key = %q, want %q", keys.ready, "moxy:{jobs}:ready") + } + if keys.processing != "moxy:{jobs}:processing" { + t.Fatalf("processing key = %q, want %q", keys.processing, "moxy:{jobs}:processing") + } + if keys.dead != "moxy:{jobs}:dead" { + t.Fatalf("dead key = %q, want %q", keys.dead, "moxy:{jobs}:dead") + } + + wantTag := "{jobs}" + for name, key := range map[string]string{ + "ready": keys.ready, + "processing": keys.processing, + "dead": keys.dead, + } { + if got := redisHashTag(key); got != wantTag { + t.Fatalf("%s key hash tag = %q, want %q", name, got, wantTag) + } + } +} diff --git a/internal/queue/redis_scripts.go b/internal/queue/redis_scripts.go new file mode 100644 index 0000000..f231ea7 --- /dev/null +++ b/internal/queue/redis_scripts.go @@ -0,0 +1,86 @@ +package queue + +import "github.com/redis/go-redis/v9" + +type redisQueueScripts struct { + complete *redis.Script + requeue *redis.Script + dead *redis.Script +} + +func defaultRedisQueueScripts() redisQueueScripts { + return redisQueueScripts{ + complete: redis.NewScript(completeTaskScript), + requeue: redis.NewScript(requeueTaskScript), + dead: redis.NewScript(deadLetterTaskScript), + } +} + +const completeTaskScript = ` +local processing = KEYS[1] +local needle = ARGV[1] +local tasks = redis.call("LRANGE", processing, 0, -1) + +for _, encoded in ipairs(tasks) do + if string.find(encoded, needle, 1, true) then + local removed = redis.call("LREM", processing, 1, encoded) + if removed == 0 then + return 0 + end + return 1 + end +end + +return 0 +` + +const requeueTaskScript = ` +local processing = KEYS[1] +local ready = KEYS[2] +local needle = ARGV[1] +local tasks = redis.call("LRANGE", processing, 0, -1) + +for _, encoded in ipairs(tasks) do + if string.find(encoded, needle, 1, true) then + local removed = redis.call("LREM", processing, 1, encoded) + if removed == 0 then + return 0 + end + local task = cjson.decode(encoded) + task["attempts"] = (task["attempts"] or 0) + 1 + redis.call("LPUSH", ready, cjson.encode(task)) + return 1 + end +end + +return 0 +` + +const deadLetterTaskScript = ` +local processing = KEYS[1] +local dead = KEYS[2] +local needle = ARGV[1] +local reason = ARGV[2] +local dead_at = ARGV[3] +local tasks = redis.call("LRANGE", processing, 0, -1) + +for _, encoded in ipairs(tasks) do + if string.find(encoded, needle, 1, true) then + local removed = redis.call("LREM", processing, 1, encoded) + if removed == 0 then + return 0 + end + local task = cjson.decode(encoded) + task["attempts"] = (task["attempts"] or 0) + 1 + local dead_task = { + task = task, + reason = reason, + dead_at = dead_at + } + redis.call("LPUSH", dead, cjson.encode(dead_task)) + return 1 + end +end + +return 0 +` diff --git a/internal/queue/redis_test.go b/internal/queue/redis_test.go index b744384..fdc17fa 100644 --- a/internal/queue/redis_test.go +++ b/internal/queue/redis_test.go @@ -123,6 +123,123 @@ func TestRedisQueueRepeatedRequeueFailsCleanly(t *testing.T) { } } +func TestRedisQueueDeadLetterMovesTaskToDeadKey(t *testing.T) { + client := redisClientForTest(t) + queue := redisQueueForTest(t, client) + + if err := queue.Enqueue(task.Task{ID: "task-1", Payload: []byte("payload")}); err != nil { + t.Fatalf("enqueue returned error: %v", err) + } + acquired, err := queue.Acquire() + if err != nil { + t.Fatalf("acquire returned error: %v", err) + } + + if err := queue.DeadLetter(acquired.ID, "expired"); err != nil { + t.Fatalf("dead letter returned error: %v", err) + } + + stats := queue.Stats() + if stats.Ready != 0 || stats.Processing != 0 || stats.Dead != 1 { + t.Fatalf("stats after dead letter = %+v, want ready=0 processing=0 dead=1", stats) + } + dead := redisDeadTasks(t, client, queue) + if len(dead) != 1 { + t.Fatalf("dead task count = %d, want 1", len(dead)) + } + if dead[0].Task.ID != "task-1" || dead[0].Task.Attempts != 1 || dead[0].Reason != "expired" { + t.Fatalf("dead entry = %+v, want task-1 attempts=1 reason=expired", dead[0]) + } +} + +func TestRedisQueueDeadLetterMissingTaskReturnsErrTaskNotProcessing(t *testing.T) { + client := redisClientForTest(t) + queue := redisQueueForTest(t, client) + + if err := queue.DeadLetter("missing", "expired"); !errors.Is(err, ErrTaskNotProcessing) { + t.Fatalf("dead letter returned %v, want ErrTaskNotProcessing", err) + } +} + +func TestRedisQueueRepeatedDeadLetterFailsCleanly(t *testing.T) { + client := redisClientForTest(t) + queue := redisQueueForTest(t, client) + + if err := queue.Enqueue(task.Task{ID: "task-1"}); err != nil { + t.Fatalf("enqueue returned error: %v", err) + } + acquired, err := queue.Acquire() + if err != nil { + t.Fatalf("acquire returned error: %v", err) + } + + if err := queue.DeadLetter(acquired.ID, "expired"); err != nil { + t.Fatalf("first dead letter returned error: %v", err) + } + if err := queue.DeadLetter(acquired.ID, "expired"); !errors.Is(err, ErrTaskNotProcessing) { + t.Fatalf("second dead letter returned %v, want ErrTaskNotProcessing", err) + } +} + +func TestRedisQueueCachedScriptsRunRepeatedly(t *testing.T) { + client := redisClientForTest(t) + queue := redisQueueForTest(t, client) + + if queue.scripts.complete == nil { + t.Fatal("complete script is nil") + } + if queue.scripts.requeue == nil { + t.Fatal("requeue script is nil") + } + + for _, id := range []string{"complete-1", "complete-2"} { + if err := queue.Enqueue(task.Task{ID: id}); err != nil { + t.Fatalf("enqueue %s returned error: %v", id, err) + } + acquired, err := queue.Acquire() + if err != nil { + t.Fatalf("acquire %s returned error: %v", id, err) + } + if err := queue.Complete(acquired.ID); err != nil { + t.Fatalf("complete %s returned error: %v", id, err) + } + } + + for _, id := range []string{"requeue-1", "requeue-2"} { + if err := queue.Enqueue(task.Task{ID: id}); err != nil { + t.Fatalf("enqueue %s returned error: %v", id, err) + } + acquired, err := queue.Acquire() + if err != nil { + t.Fatalf("acquire %s returned error: %v", id, err) + } + if err := queue.Requeue(acquired.ID); err != nil { + t.Fatalf("requeue %s returned error: %v", id, err) + } + } +} + +func redisDeadTasks(t *testing.T, client *redis.Client, queue *RedisQueue) []DeadTask { + t.Helper() + + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + encoded, err := client.LRange(ctx, queue.deadKey, 0, -1).Result() + if err != nil { + t.Fatalf("read dead key: %v", err) + } + + dead := make([]DeadTask, 0, len(encoded)) + for _, item := range encoded { + decoded, err := decodeDeadTask(item) + if err != nil { + t.Fatalf("decode dead task %q: %v", item, err) + } + dead = append(dead, decoded) + } + return dead +} + func redisClientForTest(t *testing.T) *redis.Client { t.Helper() @@ -171,7 +288,7 @@ func cleanupRedisQueue(t *testing.T, client *redis.Client, queue *RedisQueue) { ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) defer cancel() - if err := client.Del(ctx, queue.readyKey, queue.processingKey).Err(); err != nil { + if err := client.Del(ctx, queue.readyKey, queue.processingKey, queue.deadKey).Err(); err != nil { t.Fatalf("cleanup redis queue: %v", err) } } diff --git a/internal/server/server_test.go b/internal/server/server_test.go index 33690c7..51c9461 100644 --- a/internal/server/server_test.go +++ b/internal/server/server_test.go @@ -95,12 +95,15 @@ func TestServerMoxyStatsReturnsStats(t *testing.T) { sendRESP(t, conn, resp.Array(resp.BulkString("MOXY.ENQUEUE"), resp.BulkString("jobs"), resp.BulkString("hello"))) stats := sendRESP(t, conn, resp.Array(resp.BulkString("MOXY.STATS"), resp.BulkString("jobs"))) - if stats.Type != resp.TypeArray || len(stats.Array) != 8 { - t.Fatalf("stats reply = %+v, want 8-element array", stats) + if stats.Type != resp.TypeArray || len(stats.Array) != 10 { + t.Fatalf("stats reply = %+v, want 10-element array", stats) } if stats.Array[0].String != "ready" || stats.Array[1].Integer != 1 { t.Fatalf("stats reply = %+v, want ready=1", stats) } + if stats.Array[4].String != "dead" || stats.Array[5].Integer != 0 { + t.Fatalf("stats reply = %+v, want dead=0", stats) + } } func TestServerMultipleCommandsOnSameConnection(t *testing.T) { diff --git a/internal/service/service.go b/internal/service/service.go index afaef2d..89ce949 100644 --- a/internal/service/service.go +++ b/internal/service/service.go @@ -19,18 +19,30 @@ type BackendFactory func(queueName string) queue.Backend // Stats reports the state of one service-managed queue. type Stats = core.Stats +// ServiceConfig controls engines created by the service. +type ServiceConfig struct { + Engine core.EngineConfig +} + // Service owns multiple named queues. type Service struct { mu sync.Mutex factory BackendFactory engines map[string]*core.Engine + config ServiceConfig } // New creates a service that lazily creates engines from the backend factory. func New(factory BackendFactory) *Service { + return NewWithConfig(factory, ServiceConfig{}) +} + +// NewWithConfig creates a service that applies config to lazily created engines. +func NewWithConfig(factory BackendFactory, config ServiceConfig) *Service { return &Service{ factory: factory, engines: make(map[string]*core.Engine), + config: config, } } @@ -115,7 +127,7 @@ func (s *Service) engineFor(queueName string) (*core.Engine, error) { return engine, nil } - engine = core.NewEngineWithBackend(s.factory(queueName)) + engine = core.NewEngineWithBackendAndConfig(s.factory(queueName), s.config.Engine) s.engines[queueName] = engine return engine, nil } diff --git a/internal/service/service_test.go b/internal/service/service_test.go index 0e92548..60d23e2 100644 --- a/internal/service/service_test.go +++ b/internal/service/service_test.go @@ -129,6 +129,35 @@ func TestReapExpiredRequeuesExpiredLeasesAcrossAllQueues(t *testing.T) { } } +func TestServiceConfigMaxAttemptsDeadLettersExpiredTasks(t *testing.T) { + service := NewWithConfig(memoryBackendFactory, ServiceConfig{ + Engine: core.EngineConfig{MaxAttempts: 1}, + }) + if _, err := service.Enqueue("jobs", []byte("payload")); err != nil { + t.Fatalf("enqueue returned error: %v", err) + } + lease, err := service.Fetch("jobs", time.Millisecond) + if err != nil { + t.Fatalf("fetch returned error: %v", err) + } + + requeued, err := service.ReapExpired(lease.ExpiresAt.Add(time.Nanosecond)) + if err != nil { + t.Fatalf("reap expired returned error: %v", err) + } + if requeued != 1 { + t.Fatalf("reap count = %d, want 1", requeued) + } + + stats, ok := service.Stats("jobs") + if !ok { + t.Fatal("stats did not report queue") + } + if stats.Ready != 0 || stats.Processing != 0 || stats.ActiveLeases != 0 || stats.Dead != 1 { + t.Fatalf("stats = %+v, want ready=0 processing=0 active=0 dead=1", stats) + } +} + func TestEmptyQueueNameFails(t *testing.T) { service := New(memoryBackendFactory) @@ -178,8 +207,8 @@ func TestStatsReportsPerQueueState(t *testing.T) { if !ok { t.Fatal("stats did not report queue") } - if stats.Ready != 1 || stats.Processing != 1 || stats.ActiveLeases != 1 || stats.ExpirationHeap != 1 { - t.Fatalf("stats = %+v, want ready=1 processing=1 active=1 heap=1", stats) + if stats.Ready != 1 || stats.Processing != 1 || stats.ActiveLeases != 1 || stats.ExpirationHeap != 1 || stats.Dead != 0 { + t.Fatalf("stats = %+v, want ready=1 processing=1 active=1 heap=1 dead=0", stats) } } diff --git a/internal/task/task.go b/internal/task/task.go index a122fee..fd6ef6a 100644 --- a/internal/task/task.go +++ b/internal/task/task.go @@ -2,6 +2,7 @@ package task // Task represents a unit of work managed by Moxy. type Task struct { - ID string - Payload []byte + ID string `json:"id"` + Payload []byte `json:"payload"` + Attempts int `json:"attempts"` }