From 26b6922ad416ab54d67b86a304d4de0592593869 Mon Sep 17 00:00:00 2001 From: riotpiaole <19826264+Riotpiaole@users.noreply.github.com> Date: Mon, 22 Jun 2026 06:30:25 -0700 Subject: [PATCH] feat: implement gRPC server + grpc-gateway wiring (task 9) Assembles tasks 1/6/7/8 into a runnable cmd/server binary: QueueServiceServer handlers translating kafkamgmt.v1 proto to internal/core/queue's plain Go types, a lazy per-queue Kafka consumer registry for ReceiveMessage, and a Redis-scan-based queue discovery loop that starts a reaper goroutine per queue (queue lifecycle isn't exposed over gRPC, so this is the server's only signal). Promotes kmsvc-proto to a direct go.mod dependency. Handler-level integration tests run against kfake+miniredis (same documented tradeoff as tasks 5-7's envtest/testcontainers substitution), exercising send->receive->delete through the real QueueServiceServer implementation. --- cmd/server/main.go | 185 ++++++++++++++++++++ go.mod | 3 +- go.sum | 2 + internal/api/handlers/consumer_registry.go | 92 ++++++++++ internal/api/handlers/queue_service.go | 160 +++++++++++++++++ internal/api/handlers/queue_service_test.go | 149 ++++++++++++++++ internal/redis/queuemeta.go | 17 ++ 7 files changed, 607 insertions(+), 1 deletion(-) create mode 100644 cmd/server/main.go create mode 100644 internal/api/handlers/consumer_registry.go create mode 100644 internal/api/handlers/queue_service.go create mode 100644 internal/api/handlers/queue_service_test.go diff --git a/cmd/server/main.go b/cmd/server/main.go new file mode 100644 index 0000000..574239a --- /dev/null +++ b/cmd/server/main.go @@ -0,0 +1,185 @@ +// Command server runs the kafkamgmt.v1 message-plane gRPC+REST API +// (design.md §1, §9): config load, Kafka/Redis client init, gRPC server + +// grpc-gateway mux sharing one auth interceptor, and per-queue reaper +// goroutines for redelivery/DLQ sweep (design.md §5). +package main + +import ( + "context" + "errors" + "log" + "net" + "net/http" + "os" + "os/signal" + "sync" + "syscall" + "time" + + kafkamgmtv1 "forgejo.riotpiao.homelab.com/rock/kmsvc-proto/gen/kafkamgmt/v1" + "github.com/grpc-ecosystem/grpc-gateway/v2/runtime" + goredis "github.com/redis/go-redis/v9" + "google.golang.org/grpc" + + "github.com/rockliang/kafka-management-service/internal/api/handlers" + "github.com/rockliang/kafka-management-service/internal/api/interceptors" + "github.com/rockliang/kafka-management-service/internal/auth" + "github.com/rockliang/kafka-management-service/internal/config" + "github.com/rockliang/kafka-management-service/internal/core/queue" + "github.com/rockliang/kafka-management-service/internal/core/reaper" + "github.com/rockliang/kafka-management-service/internal/kafka" + kmsvcredis "github.com/rockliang/kafka-management-service/internal/redis" +) + +// queueDiscoveryInterval is how often the server rescans Redis for newly +// reconciled queues to start a reaper goroutine for (design.md §5) — queue +// lifecycle isn't exposed over gRPC, so this is the server's only signal. +const queueDiscoveryInterval = 10 * time.Second + +func main() { + if err := run(); err != nil { + log.Fatalf("kmsvc-server: %v", err) + } +} + +func run() error { + ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer stop() + + cfg, err := config.Load() + if err != nil { + return err + } + + rdb := kmsvcredis.NewClient(cfg.RedisAddr, cfg.RedisPassword, cfg.RedisDB) + defer rdb.Close() + + admin, err := kafka.NewAdmin(cfg.KafkaBrokers) + if err != nil { + return err + } + defer admin.Close() + + producer, err := kafka.NewProducer(cfg.KafkaBrokers) + if err != nil { + return err + } + defer producer.Close() + + validator, err := auth.NewValidator(ctx, cfg.AuthentikIssuerURL, cfg.AuthentikAudience) + if err != nil { + return err + } + + router := &queue.ShardRouter{Redis: rdb} + + svc := &handlers.QueueService{ + Redis: rdb, + Router: router, + Send: &queue.SendMessageService{ + Redis: rdb, + Producer: producer, + Router: router, + DedupWindow: time.Duration(cfg.DedupWindowSeconds) * time.Second, + }, + Delete: &queue.DeleteMessageService{Redis: rdb, Committer: admin}, + Visibility: &queue.ChangeVisibilityService{Redis: rdb}, + Consumers: &handlers.ConsumerRegistry{Brokers: cfg.KafkaBrokers, Router: router}, + } + defer svc.Consumers.Close() + + grpcSrv := grpc.NewServer( + grpc.ChainUnaryInterceptor(interceptors.UnaryServerInterceptor(validator)), + grpc.ChainStreamInterceptor(interceptors.StreamServerInterceptor(validator)), + ) + kafkamgmtv1.RegisterQueueServiceServer(grpcSrv, svc) + + mux := runtime.NewServeMux() + if err := kafkamgmtv1.RegisterQueueServiceHandlerServer(ctx, mux, svc); err != nil { + return err + } + httpSrv := &http.Server{Addr: cfg.HTTPListenAddr, Handler: mux} + + r := &reaper.Reaper{ + Redis: rdb, + DLQSender: &queue.SendMessageService{ + Redis: rdb, + Producer: producer, + Router: router, + DedupWindow: time.Duration(cfg.DedupWindowSeconds) * time.Second, + }, + } + go runReaperDiscovery(ctx, rdb, r) + + var wg sync.WaitGroup + wg.Add(2) + + go func() { + defer wg.Done() + lis, err := net.Listen("tcp", cfg.GRPCListenAddr) + if err != nil { + log.Printf("grpc listen %s: %v", cfg.GRPCListenAddr, err) + return + } + log.Printf("grpc listening on %s", cfg.GRPCListenAddr) + if err := grpcSrv.Serve(lis); err != nil { + log.Printf("grpc serve: %v", err) + } + }() + + go func() { + defer wg.Done() + log.Printf("http listening on %s", cfg.HTTPListenAddr) + if err := httpSrv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { + log.Printf("http serve: %v", err) + } + }() + + <-ctx.Done() + log.Println("shutting down") + + shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + _ = httpSrv.Shutdown(shutdownCtx) + grpcSrv.GracefulStop() + + wg.Wait() + return nil +} + +// runReaperDiscovery periodically scans for known queues and ensures each +// has a running reaper sweep goroutine. A queue is only ever added, never +// removed mid-run: if it's deleted, the next sweep finds no queue meta and +// is a harmless no-op (design.md §5's Sweep already handles this). +func runReaperDiscovery(ctx context.Context, rdb *goredis.Client, r *reaper.Reaper) { + started := make(map[string]bool) + ticker := time.NewTicker(queueDiscoveryInterval) + defer ticker.Stop() + + discover := func() { + names, err := kmsvcredis.ListQueueNames(ctx, rdb) + if err != nil { + log.Printf("reaper discovery: %v", err) + return + } + for _, name := range names { + if started[name] { + continue + } + started[name] = true + go r.Run(ctx, name, func(err error) { + log.Printf("reaper sweep %s: %v", name, err) + }) + } + } + + discover() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + discover() + } + } +} diff --git a/go.mod b/go.mod index 3b5f0f0..3551727 100644 --- a/go.mod +++ b/go.mod @@ -3,6 +3,7 @@ module github.com/rockliang/kafka-management-service go 1.26.0 require ( + forgejo.riotpiao.homelab.com/rock/kmsvc-proto v1.1.0 github.com/alicebob/miniredis/v2 v2.38.0 github.com/google/uuid v1.6.0 github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 @@ -11,7 +12,6 @@ require ( github.com/twmb/franz-go v1.21.3 github.com/twmb/franz-go/pkg/kadm v1.18.0 github.com/twmb/franz-go/pkg/kfake v0.0.0-20260615024848-f17c00130060 - google.golang.org/genproto/googleapis/api v0.0.0-20260618152121-87f3d3e198d3 google.golang.org/grpc v1.81.1 google.golang.org/protobuf v1.36.12-0.20260120151049-f2248ac996af k8s.io/apimachinery v0.36.2 @@ -72,6 +72,7 @@ require ( golang.org/x/text v0.37.0 // indirect golang.org/x/time v0.14.0 // indirect gomodules.xyz/jsonpatch/v2 v2.4.0 // indirect + google.golang.org/genproto/googleapis/api v0.0.0-20260618152121-87f3d3e198d3 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260610212136-7ab31c22f7ad // indirect gopkg.in/evanphx/json-patch.v4 v4.13.0 // indirect gopkg.in/inf.v0 v0.9.1 // indirect diff --git a/go.sum b/go.sum index 8a7bd71..c0532c6 100644 --- a/go.sum +++ b/go.sum @@ -1,3 +1,5 @@ +forgejo.riotpiao.homelab.com/rock/kmsvc-proto v1.1.0 h1:T7Y0aWFucwtwoOKy0XTCb+MAtay8U5j4SHCgi6PCrQo= +forgejo.riotpiao.homelab.com/rock/kmsvc-proto v1.1.0/go.mod h1:pKNhLE2KPpNUeu2BL5BFDgXmPlbyB9LHCfi0G4aE38s= github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1Xbatp0= github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM= github.com/alicebob/miniredis/v2 v2.38.0 h1:nZAzCR+Lj+Vxk4ZXzm2NuKq2O33RXj1XxJ2e2uP9jiw= diff --git a/internal/api/handlers/consumer_registry.go b/internal/api/handlers/consumer_registry.go new file mode 100644 index 0000000..5a0d341 --- /dev/null +++ b/internal/api/handlers/consumer_registry.go @@ -0,0 +1,92 @@ +package handlers + +import ( + "context" + "fmt" + "sync" + + "github.com/rockliang/kafka-management-service/internal/core/queue" + "github.com/rockliang/kafka-management-service/internal/kafka" +) + +// ConsumerRegistry lazily creates and caches one Kafka consumer per queue, +// subscribed to that queue's currently-consumable shard topics (design.md +// §6 — every Active+Closing shard). It is the only place ReceiveMessage's +// gRPC handler needs to know about per-queue Kafka clients. +// +// Known v1 limitation (tracked in design.md §11 item 5): if the operator +// splits a shard after a queue's consumer was created, the new child topics +// are picked up the next time Get is called for that queue (each call +// re-syncs against the live shard map), but there is a short window between +// a split and the next ReceiveMessage call where the consumer hasn't yet +// subscribed to the new topics. +type ConsumerRegistry struct { + Brokers []string + Router *queue.ShardRouter + + mu sync.Mutex + consumers map[string]*registeredConsumer +} + +type registeredConsumer struct { + consumer *kafka.Consumer + topics map[string]bool +} + +// Get returns a Fetcher subscribed to queueName's current consumable shard +// topics, creating the underlying Kafka consumer-group client on first use. +func (r *ConsumerRegistry) Get(ctx context.Context, queueName string) (queue.Fetcher, error) { + shards, err := r.Router.ConsumableShards(ctx, queueName) + if err != nil { + return nil, err + } + if len(shards) == 0 { + return nil, fmt.Errorf("queue %s has no consumable shards", queueName) + } + + r.mu.Lock() + defer r.mu.Unlock() + + if r.consumers == nil { + r.consumers = make(map[string]*registeredConsumer) + } + + rc, ok := r.consumers[queueName] + if !ok { + topics := make([]string, 0, len(shards)) + topicSet := make(map[string]bool, len(shards)) + for _, s := range shards { + topics = append(topics, s.Topic) + topicSet[s.Topic] = true + } + cl, err := kafka.NewConsumer(r.Brokers, kafka.ConsumerGroup(queueName), topics...) + if err != nil { + return nil, fmt.Errorf("create consumer for queue %s: %w", queueName, err) + } + rc = ®isteredConsumer{consumer: cl, topics: topicSet} + r.consumers[queueName] = rc + return rc.consumer, nil + } + + var newTopics []string + for _, s := range shards { + if !rc.topics[s.Topic] { + newTopics = append(newTopics, s.Topic) + rc.topics[s.Topic] = true + } + } + if len(newTopics) > 0 { + rc.consumer.AddTopics(newTopics...) + } + return rc.consumer, nil +} + +// Close closes every consumer this registry has created, used on server +// shutdown. +func (r *ConsumerRegistry) Close() { + r.mu.Lock() + defer r.mu.Unlock() + for _, rc := range r.consumers { + rc.consumer.Close() + } +} diff --git a/internal/api/handlers/queue_service.go b/internal/api/handlers/queue_service.go new file mode 100644 index 0000000..2683e06 --- /dev/null +++ b/internal/api/handlers/queue_service.go @@ -0,0 +1,160 @@ +// Package handlers implements the kafkamgmt.v1 QueueServiceServer interface +// (design.md §2b) by delegating to internal/core/queue's business logic — +// this layer only translates between proto messages and that package's +// plain Go types, and maps domain errors to gRPC status codes. +package handlers + +import ( + "context" + "errors" + "strconv" + "strings" + "time" + + kafkamgmtv1 "forgejo.riotpiao.homelab.com/rock/kmsvc-proto/gen/kafkamgmt/v1" + goredis "github.com/redis/go-redis/v9" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + "google.golang.org/protobuf/types/known/timestamppb" + + "github.com/rockliang/kafka-management-service/internal/core/queue" +) + +// QueueService implements kafkamgmtv1.QueueServiceServer. +type QueueService struct { + kafkamgmtv1.UnimplementedQueueServiceServer + + Redis *goredis.Client + Router *queue.ShardRouter + + Send *queue.SendMessageService + Delete *queue.DeleteMessageService + Visibility *queue.ChangeVisibilityService + Consumers *ConsumerRegistry + + // ReceivePollInterval overrides the default poll interval used while + // long-polling; primarily for tests. Zero uses the package default. + ReceivePollInterval time.Duration +} + +func (s *QueueService) SendMessage(ctx context.Context, req *kafkamgmtv1.SendMessageRequest) (*kafkamgmtv1.SendMessageResponse, error) { + out, err := s.Send.SendMessage(ctx, queue.SendMessageInput{ + QueueName: req.GetQueueName(), + Body: string(req.GetMessageBody()), + MessageGroupID: req.GetMessageGroupId(), + MessageDeduplicationID: req.GetMessageDeduplicationId(), + }) + if err != nil { + return nil, mapError(err) + } + return &kafkamgmtv1.SendMessageResponse{ + MessageId: out.MessageID, + SequenceNumber: strconv.FormatInt(out.SequenceNumber, 10), + }, nil +} + +func (s *QueueService) SendMessageBatch(ctx context.Context, req *kafkamgmtv1.SendMessageBatchRequest) (*kafkamgmtv1.SendMessageBatchResponse, error) { + resp := &kafkamgmtv1.SendMessageBatchResponse{} + for _, entry := range req.GetEntries() { + out, err := s.Send.SendMessage(ctx, queue.SendMessageInput{ + QueueName: req.GetQueueName(), + Body: string(entry.GetMessageBody()), + MessageGroupID: entry.GetMessageGroupId(), + MessageDeduplicationID: entry.GetMessageDeduplicationId(), + }) + if err != nil { + resp.Failed = append(resp.Failed, &kafkamgmtv1.BatchResultEntry{Id: entry.GetId(), Error: err.Error()}) + continue + } + resp.Successful = append(resp.Successful, &kafkamgmtv1.BatchResultEntry{Id: entry.GetId(), MessageId: out.MessageID}) + } + return resp, nil +} + +func (s *QueueService) ReceiveMessage(ctx context.Context, req *kafkamgmtv1.ReceiveMessageRequest) (*kafkamgmtv1.ReceiveMessageResponse, error) { + fetcher, err := s.Consumers.Get(ctx, req.GetQueueName()) + if err != nil { + return nil, mapError(err) + } + // A fresh ReceiveMessageService per call: Redis/Router are shared and + // safe for concurrent use, but Fetcher is request-scoped (each call may + // target a different queue's consumer), so the service struct itself + // must not be shared across concurrent requests. + recv := &queue.ReceiveMessageService{ + Redis: s.Redis, + Fetcher: fetcher, + Router: s.Router, + PollInterval: s.ReceivePollInterval, + } + + msgs, err := recv.ReceiveMessage(ctx, queue.ReceiveMessageInput{ + QueueName: req.GetQueueName(), + MaxNumberOfMessages: req.GetMaxNumberOfMessages(), + WaitTime: time.Duration(req.GetWaitTimeSeconds()) * time.Second, + VisibilityTimeoutOverride: time.Duration(req.GetVisibilityTimeoutSeconds()) * time.Second, + }) + if err != nil { + return nil, mapError(err) + } + + out := make([]*kafkamgmtv1.Message, 0, len(msgs)) + for _, m := range msgs { + out = append(out, &kafkamgmtv1.Message{ + MessageId: m.ReceiptHandle, + ReceiptHandle: m.ReceiptHandle, + Body: []byte(m.Body), + ReceiveCount: m.ReceiveCount, + EnqueuedAt: timestamppb.Now(), + }) + } + return &kafkamgmtv1.ReceiveMessageResponse{Messages: out}, nil +} + +func (s *QueueService) DeleteMessage(ctx context.Context, req *kafkamgmtv1.DeleteMessageRequest) (*kafkamgmtv1.DeleteMessageResponse, error) { + if err := s.Delete.DeleteMessage(ctx, req.GetQueueName(), req.GetReceiptHandle()); err != nil { + return nil, mapError(err) + } + return &kafkamgmtv1.DeleteMessageResponse{}, nil +} + +func (s *QueueService) DeleteMessageBatch(ctx context.Context, req *kafkamgmtv1.DeleteMessageBatchRequest) (*kafkamgmtv1.DeleteMessageBatchResponse, error) { + resp := &kafkamgmtv1.DeleteMessageBatchResponse{} + for _, entry := range req.GetEntries() { + if err := s.Delete.DeleteMessage(ctx, req.GetQueueName(), entry.GetReceiptHandle()); err != nil { + resp.Failed = append(resp.Failed, &kafkamgmtv1.BatchResultEntry{Id: entry.GetId(), Error: err.Error()}) + continue + } + resp.Successful = append(resp.Successful, &kafkamgmtv1.BatchResultEntry{Id: entry.GetId()}) + } + return resp, nil +} + +func (s *QueueService) ChangeMessageVisibility(ctx context.Context, req *kafkamgmtv1.ChangeMessageVisibilityRequest) (*kafkamgmtv1.ChangeMessageVisibilityResponse, error) { + timeout := time.Duration(req.GetVisibilityTimeoutSeconds()) * time.Second + if err := s.Visibility.ChangeMessageVisibility(ctx, req.GetQueueName(), req.GetReceiptHandle(), timeout); err != nil { + return nil, mapError(err) + } + return &kafkamgmtv1.ChangeMessageVisibilityResponse{}, nil +} + +// mapError maps internal/core/queue's plain errors to gRPC status codes. +// These services return fmt.Errorf-wrapped strings rather than typed +// sentinels, so this matches on substring rather than errors.Is — a +// pragmatic v1 choice documented here rather than introducing typed errors +// purely for this translation layer. +func mapError(err error) error { + if err == nil { + return nil + } + msg := err.Error() + switch { + case strings.Contains(msg, "not found"): + return status.Error(codes.NotFound, msg) + case strings.Contains(msg, "exceeds") || strings.Contains(msg, "required for FIFO") || strings.Contains(msg, "no consumable shards") || strings.Contains(msg, "not reconciled"): + return status.Error(codes.InvalidArgument, msg) + case errors.Is(err, context.DeadlineExceeded): + return status.Error(codes.DeadlineExceeded, msg) + default: + return status.Error(codes.Internal, msg) + } +} diff --git a/internal/api/handlers/queue_service_test.go b/internal/api/handlers/queue_service_test.go new file mode 100644 index 0000000..1ecd24f --- /dev/null +++ b/internal/api/handlers/queue_service_test.go @@ -0,0 +1,149 @@ +package handlers + +import ( + "context" + "testing" + "time" + + kafkamgmtv1 "forgejo.riotpiao.homelab.com/rock/kmsvc-proto/gen/kafkamgmt/v1" + "github.com/alicebob/miniredis/v2" + goredis "github.com/redis/go-redis/v9" + "github.com/twmb/franz-go/pkg/kfake" + + "github.com/rockliang/kafka-management-service/internal/core/queue" + "github.com/rockliang/kafka-management-service/internal/kafka" + kmsvcredis "github.com/rockliang/kafka-management-service/internal/redis" +) + +// newTestKafka starts an in-memory, wire-protocol-compatible fake Kafka +// cluster (kfake), the same documented tradeoff used by internal/core/queue's +// own tests — no Docker/testcontainers needed in this environment. +func newTestKafka(t *testing.T) []string { + t.Helper() + cluster, err := kfake.NewCluster(kfake.NumBrokers(1)) + if err != nil { + t.Fatalf("starting kfake cluster: %v", err) + } + t.Cleanup(cluster.Close) + return cluster.ListenAddrs() +} + +func newTestRedis(t *testing.T) *goredis.Client { + t.Helper() + mr, err := miniredis.Run() + if err != nil { + t.Fatalf("starting miniredis: %v", err) + } + t.Cleanup(mr.Close) + return goredis.NewClient(&goredis.Options{Addr: mr.Addr()}) +} + +func newTestQueueService(t *testing.T, brokers []string, rdb *goredis.Client) *QueueService { + t.Helper() + + producer, err := kafka.NewProducer(brokers) + if err != nil { + t.Fatalf("new producer: %v", err) + } + t.Cleanup(producer.Close) + + router := &queue.ShardRouter{Redis: rdb} + svc := &QueueService{ + Redis: rdb, + Router: router, + Send: &queue.SendMessageService{Redis: rdb, Producer: producer, Router: router}, + Delete: &queue.DeleteMessageService{Redis: rdb}, + Visibility: &queue.ChangeVisibilityService{Redis: rdb}, + Consumers: &ConsumerRegistry{Brokers: brokers, Router: router}, + ReceivePollInterval: 20 * time.Millisecond, + } + t.Cleanup(svc.Consumers.Close) + return svc +} + +func seedTestQueue(t *testing.T, ctx context.Context, brokers []string, rdb *goredis.Client, name string) { + t.Helper() + admin, err := kafka.NewAdmin(brokers) + if err != nil { + t.Fatalf("new admin: %v", err) + } + t.Cleanup(admin.Close) + + shard := kafka.Shard{ID: "0", Topic: kafka.ShardTopicName(name, false, "0"), HashRangeStart: 0, HashRangeEnd: kafka.FullHashRangeEnd, Phase: "Active"} + if err := admin.CreateTopic(ctx, shard.Topic, kafka.TopicConfig{PartitionCount: 1, ReplicationFactor: 1, RetentionSeconds: 345600, MinInsyncReplicas: 1}); err != nil { + t.Fatalf("create shard topic: %v", err) + } + if err := kmsvcredis.PutQueueMeta(ctx, rdb, name, kmsvcredis.QueueMeta{VisibilityTimeoutSeconds: 5, MaxReceiveCount: 3, PartitionsPerShard: 1, RetentionSeconds: 345600}); err != nil { + t.Fatalf("put queue meta: %v", err) + } + if err := kmsvcredis.PutShardMap(ctx, rdb, name, []kafka.Shard{shard}); err != nil { + t.Fatalf("put shard map: %v", err) + } +} + +// TestSendReceiveDeleteThroughGRPCHandlers exercises the proto<->core +// translation layer end-to-end: SendMessage -> ReceiveMessage -> +// DeleteMessage via the same QueueService a real grpc.Server would dispatch +// to, against a fake Kafka+Redis backend. +func TestSendReceiveDeleteThroughGRPCHandlers(t *testing.T) { + ctx := context.Background() + brokers := newTestKafka(t) + rdb := newTestRedis(t) + seedTestQueue(t, ctx, brokers, rdb, "orders") + svc := newTestQueueService(t, brokers, rdb) + + sendResp, err := svc.SendMessage(ctx, &kafkamgmtv1.SendMessageRequest{QueueName: "orders", MessageBody: []byte("hello")}) + if err != nil { + t.Fatalf("SendMessage: %v", err) + } + if sendResp.GetMessageId() == "" { + t.Fatal("SendMessage: expected a non-empty message id") + } + + deadline := time.Now().Add(3 * time.Second) + var msgs []*kafkamgmtv1.Message + for time.Now().Before(deadline) { + recvResp, err := svc.ReceiveMessage(ctx, &kafkamgmtv1.ReceiveMessageRequest{QueueName: "orders", MaxNumberOfMessages: 1, WaitTimeSeconds: 1}) + if err != nil { + t.Fatalf("ReceiveMessage: %v", err) + } + if len(recvResp.GetMessages()) > 0 { + msgs = recvResp.GetMessages() + break + } + } + if len(msgs) != 1 || string(msgs[0].GetBody()) != "hello" { + t.Fatalf("messages = %+v, want one body=hello", msgs) + } + + if _, err := svc.DeleteMessage(ctx, &kafkamgmtv1.DeleteMessageRequest{QueueName: "orders", ReceiptHandle: msgs[0].GetReceiptHandle()}); err != nil { + t.Fatalf("DeleteMessage: %v", err) + } +} + +func TestSendMessageUnknownQueueMapsToNotFound(t *testing.T) { + ctx := context.Background() + brokers := newTestKafka(t) + rdb := newTestRedis(t) + svc := newTestQueueService(t, brokers, rdb) + + _, err := svc.SendMessage(ctx, &kafkamgmtv1.SendMessageRequest{QueueName: "missing", MessageBody: []byte("hi")}) + if err == nil { + t.Fatal("expected an error for an unknown queue") + } +} + +func TestChangeMessageVisibilityUnknownHandleErrors(t *testing.T) { + ctx := context.Background() + brokers := newTestKafka(t) + rdb := newTestRedis(t) + seedTestQueue(t, ctx, brokers, rdb, "orders") + svc := newTestQueueService(t, brokers, rdb) + + _, err := svc.ChangeMessageVisibility(ctx, &kafkamgmtv1.ChangeMessageVisibilityRequest{ + QueueName: "orders", ReceiptHandle: "does-not-exist", VisibilityTimeoutSeconds: 30, + }) + if err == nil { + t.Fatal("expected an error for an unknown receipt handle") + } +} diff --git a/internal/redis/queuemeta.go b/internal/redis/queuemeta.go index 8c26c3e..fe80bdb 100644 --- a/internal/redis/queuemeta.go +++ b/internal/redis/queuemeta.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "strconv" + "strings" "time" "github.com/redis/go-redis/v9" @@ -79,3 +80,19 @@ func DeleteQueueMeta(ctx context.Context, rdb *redis.Client, queue string) error } return nil } + +// ListQueueNames scans for every queue with a published `kmsvc:queue:` row, +// used by the message-plane server to discover which queues need a reaper +// goroutine (design.md §5) since queue lifecycle isn't exposed over gRPC. +func ListQueueNames(ctx context.Context, rdb *redis.Client) ([]string, error) { + const prefix = "kmsvc:queue:" + var names []string + iter := rdb.Scan(ctx, 0, prefix+"*", 0).Iterator() + for iter.Next(ctx) { + names = append(names, strings.TrimPrefix(iter.Val(), prefix)) + } + if err := iter.Err(); err != nil { + return nil, fmt.Errorf("list queue names: %w", err) + } + return names, nil +}