From 1f95e73bf57095bce6f32e97287c52dafc796a91 Mon Sep 17 00:00:00 2001 From: smallchungus Date: Fri, 17 Apr 2026 20:40:26 -0400 Subject: [PATCH 1/2] feat(queue): add Queue with Push and Depth --- internal/queue/queue.go | 35 ++++++++++++++++++++++ internal/queue/queue_integration_test.go | 37 ++++++++++++++++++++++++ 2 files changed, 72 insertions(+) create mode 100644 internal/queue/queue.go create mode 100644 internal/queue/queue_integration_test.go diff --git a/internal/queue/queue.go b/internal/queue/queue.go new file mode 100644 index 0000000..b265342 --- /dev/null +++ b/internal/queue/queue.go @@ -0,0 +1,35 @@ +package queue + +import ( + "context" + "fmt" + + goredis "github.com/redis/go-redis/v9" +) + +type Queue struct { + cli *goredis.Client +} + +func New(cli *goredis.Client) *Queue { + return &Queue{cli: cli} +} + +func key(stage string) string { + return "queue:" + stage +} + +func (q *Queue) Push(ctx context.Context, stage, jobID string) error { + if err := q.cli.LPush(ctx, key(stage), jobID).Err(); err != nil { + return fmt.Errorf("push: %w", err) + } + return nil +} + +func (q *Queue) Depth(ctx context.Context, stage string) (int64, error) { + n, err := q.cli.LLen(ctx, key(stage)).Result() + if err != nil { + return 0, fmt.Errorf("depth: %w", err) + } + return n, nil +} diff --git a/internal/queue/queue_integration_test.go b/internal/queue/queue_integration_test.go new file mode 100644 index 0000000..379ec86 --- /dev/null +++ b/internal/queue/queue_integration_test.go @@ -0,0 +1,37 @@ +//go:build integration + +package queue_test + +import ( + "context" + "testing" + + "github.com/smallchungus/disttaskqueue/internal/queue" + "github.com/smallchungus/disttaskqueue/internal/testutil" +) + +func newQueue(t *testing.T) *queue.Queue { + t.Helper() + cli := testutil.StartRedis(t) + return queue.New(cli) +} + +func TestPush_AppendsToStageList(t *testing.T) { + q := newQueue(t) + ctx := context.Background() + + if err := q.Push(ctx, "fetch", "job-1"); err != nil { + t.Fatalf("push: %v", err) + } + if err := q.Push(ctx, "fetch", "job-2"); err != nil { + t.Fatalf("push: %v", err) + } + + depth, err := q.Depth(ctx, "fetch") + if err != nil { + t.Fatalf("depth: %v", err) + } + if depth != 2 { + t.Fatalf("depth: got %d, want 2", depth) + } +} From ccf7018d3de0db3a6f46ebc615702bd2c60a78bb Mon Sep 17 00:00:00 2001 From: smallchungus Date: Fri, 17 Apr 2026 20:41:09 -0400 Subject: [PATCH 2/2] feat(queue): add BlockingPop (BRPop, FIFO with LPush) and ErrEmpty --- internal/queue/queue.go | 18 ++++++++++ internal/queue/queue_integration_test.go | 46 ++++++++++++++++++++++++ 2 files changed, 64 insertions(+) diff --git a/internal/queue/queue.go b/internal/queue/queue.go index b265342..c24f82a 100644 --- a/internal/queue/queue.go +++ b/internal/queue/queue.go @@ -2,7 +2,9 @@ package queue import ( "context" + "errors" "fmt" + "time" goredis "github.com/redis/go-redis/v9" ) @@ -15,6 +17,8 @@ func New(cli *goredis.Client) *Queue { return &Queue{cli: cli} } +var ErrEmpty = errors.New("queue empty") + func key(stage string) string { return "queue:" + stage } @@ -26,6 +30,20 @@ func (q *Queue) Push(ctx context.Context, stage, jobID string) error { return nil } +func (q *Queue) BlockingPop(ctx context.Context, stage string, timeout time.Duration) (string, error) { + res, err := q.cli.BRPop(ctx, timeout, key(stage)).Result() + if errors.Is(err, goredis.Nil) { + return "", ErrEmpty + } + if err != nil { + return "", fmt.Errorf("pop: %w", err) + } + if len(res) != 2 { + return "", fmt.Errorf("pop: unexpected response shape %v", res) + } + return res[1], nil +} + func (q *Queue) Depth(ctx context.Context, stage string) (int64, error) { n, err := q.cli.LLen(ctx, key(stage)).Result() if err != nil { diff --git a/internal/queue/queue_integration_test.go b/internal/queue/queue_integration_test.go index 379ec86..e42c49a 100644 --- a/internal/queue/queue_integration_test.go +++ b/internal/queue/queue_integration_test.go @@ -5,6 +5,7 @@ package queue_test import ( "context" "testing" + "time" "github.com/smallchungus/disttaskqueue/internal/queue" "github.com/smallchungus/disttaskqueue/internal/testutil" @@ -35,3 +36,48 @@ func TestPush_AppendsToStageList(t *testing.T) { t.Fatalf("depth: got %d, want 2", depth) } } + +func TestBlockingPop_ReturnsPushedJobID(t *testing.T) { + q := newQueue(t) + ctx := context.Background() + + if err := q.Push(ctx, "fetch", "job-abc"); err != nil { + t.Fatal(err) + } + + jobID, err := q.BlockingPop(ctx, "fetch", 2*time.Second) + if err != nil { + t.Fatalf("pop: %v", err) + } + if jobID != "job-abc" { + t.Fatalf("jobID: got %q, want %q", jobID, "job-abc") + } +} + +func TestBlockingPop_FIFOOrdering(t *testing.T) { + q := newQueue(t) + ctx := context.Background() + for _, id := range []string{"a", "b", "c"} { + if err := q.Push(ctx, "fetch", id); err != nil { + t.Fatal(err) + } + } + + for _, want := range []string{"a", "b", "c"} { + got, err := q.BlockingPop(ctx, "fetch", time.Second) + if err != nil { + t.Fatal(err) + } + if got != want { + t.Fatalf("got %q, want %q", got, want) + } + } +} + +func TestBlockingPop_ReturnsErrEmptyOnTimeout(t *testing.T) { + q := newQueue(t) + _, err := q.BlockingPop(context.Background(), "fetch", 200*time.Millisecond) + if err != queue.ErrEmpty { + t.Fatalf("got %v, want ErrEmpty", err) + } +}