diff --git a/internal/queue/queue.go b/internal/queue/queue.go new file mode 100644 index 0000000..c24f82a --- /dev/null +++ b/internal/queue/queue.go @@ -0,0 +1,53 @@ +package queue + +import ( + "context" + "errors" + "fmt" + "time" + + goredis "github.com/redis/go-redis/v9" +) + +type Queue struct { + cli *goredis.Client +} + +func New(cli *goredis.Client) *Queue { + return &Queue{cli: cli} +} + +var ErrEmpty = errors.New("queue empty") + +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) 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 { + 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..e42c49a --- /dev/null +++ b/internal/queue/queue_integration_test.go @@ -0,0 +1,83 @@ +//go:build integration + +package queue_test + +import ( + "context" + "testing" + "time" + + "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) + } +} + +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) + } +}