package cli import ( "bytes" "context" "strings" "sync/atomic" "testing" kafkamgmtv1 "forgejo.riotpiao.com/rock/kmsvc-proto/gen/kafkamgmt/v1" ) func TestDLQRedriveHappyPath(t *testing.T) { var sendCalled, deleteCalled atomic.Bool var sendQueue string fake := &fakeQueueService{ receiveMessage: func(ctx context.Context, req *kafkamgmtv1.ReceiveMessageRequest) (*kafkamgmtv1.ReceiveMessageResponse, error) { return &kafkamgmtv1.ReceiveMessageResponse{ Messages: []*kafkamgmtv1.Message{{MessageId: "m1", ReceiptHandle: "rh1", Body: []byte("payload")}}, }, nil }, sendMessage: func(ctx context.Context, req *kafkamgmtv1.SendMessageRequest) (*kafkamgmtv1.SendMessageResponse, error) { sendCalled.Store(true) sendQueue = req.QueueName if deleteCalled.Load() { t.Error("delete called before send completed") } return &kafkamgmtv1.SendMessageResponse{MessageId: "m1-redriven"}, nil }, deleteMessage: func(ctx context.Context, req *kafkamgmtv1.DeleteMessageRequest) (*kafkamgmtv1.DeleteMessageResponse, error) { deleteCalled.Store(true) if !sendCalled.Load() { t.Error("delete called before send") } return &kafkamgmtv1.DeleteMessageResponse{}, nil }, } addr := startTestServer(t, fake) flags := &globalFlags{server: addr, output: "table", insecure: true} cmd := newDLQRedriveCmd(flags) cmd.SetArgs([]string{"--queue", "orders.dlq", "--to", "orders"}) var out bytes.Buffer cmd.SetOut(&out) if err := cmd.Execute(); err != nil { t.Fatalf("Execute: %v", err) } if !sendCalled.Load() || !deleteCalled.Load() { t.Fatal("expected both send and delete to be called") } if sendQueue != "orders" { t.Errorf("send queue = %q, want orders", sendQueue) } } func TestDLQRedriveSurfacesDeleteFailure(t *testing.T) { fake := &fakeQueueService{ receiveMessage: func(ctx context.Context, req *kafkamgmtv1.ReceiveMessageRequest) (*kafkamgmtv1.ReceiveMessageResponse, error) { return &kafkamgmtv1.ReceiveMessageResponse{ Messages: []*kafkamgmtv1.Message{{MessageId: "m1", ReceiptHandle: "rh1", Body: []byte("payload")}}, }, nil }, sendMessage: func(ctx context.Context, req *kafkamgmtv1.SendMessageRequest) (*kafkamgmtv1.SendMessageResponse, error) { return &kafkamgmtv1.SendMessageResponse{MessageId: "m1-redriven"}, nil }, deleteMessage: func(ctx context.Context, req *kafkamgmtv1.DeleteMessageRequest) (*kafkamgmtv1.DeleteMessageResponse, error) { return nil, errBoom }, } addr := startTestServer(t, fake) flags := &globalFlags{server: addr, output: "table", insecure: true} cmd := newDLQRedriveCmd(flags) cmd.SetArgs([]string{"--queue", "orders.dlq", "--to", "orders"}) cmd.SilenceUsage = true cmd.SilenceErrors = true var out bytes.Buffer cmd.SetOut(&out) err := cmd.Execute() if err == nil { t.Fatal("expected redrive to report failure when delete fails") } output := out.String() if !strings.Contains(output, "may be duplicated") { t.Errorf("output = %q, want a duplicate-risk warning", output) } } var errBoom = &boomError{} type boomError struct{} func (e *boomError) Error() string { return "boom" }