2026-08-17 10:14:44 -07:00
|
|
|
package cli
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"bytes"
|
|
|
|
|
"context"
|
|
|
|
|
"strings"
|
|
|
|
|
"testing"
|
|
|
|
|
|
2026-08-21 20:09:48 -07:00
|
|
|
kafkamgmtv1 "forgejo.riotpiao.com/rock/kmsvc-proto/gen/kafkamgmt/v1"
|
2026-08-17 10:14:44 -07:00
|
|
|
)
|
|
|
|
|
|
|
|
|
|
func TestMessageSendCmd(t *testing.T) {
|
|
|
|
|
fake := &fakeQueueService{
|
|
|
|
|
sendMessage: func(ctx context.Context, req *kafkamgmtv1.SendMessageRequest) (*kafkamgmtv1.SendMessageResponse, error) {
|
|
|
|
|
if req.QueueName != "orders" || string(req.MessageBody) != "hello" {
|
|
|
|
|
t.Errorf("unexpected request: %+v", req)
|
|
|
|
|
}
|
|
|
|
|
return &kafkamgmtv1.SendMessageResponse{MessageId: "m1"}, nil
|
|
|
|
|
},
|
|
|
|
|
}
|
|
|
|
|
addr := startTestServer(t, fake)
|
|
|
|
|
|
|
|
|
|
flags := &globalFlags{server: addr, output: "table", insecure: true}
|
|
|
|
|
cmd := newMessageSendCmd(flags)
|
|
|
|
|
cmd.SetArgs([]string{"--queue", "orders", "--body", "hello"})
|
|
|
|
|
|
|
|
|
|
var out bytes.Buffer
|
|
|
|
|
cmd.SetOut(&out)
|
|
|
|
|
if err := cmd.Execute(); err != nil {
|
|
|
|
|
t.Fatalf("Execute: %v", err)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if !strings.Contains(out.String(), "message_id=m1") {
|
|
|
|
|
t.Errorf("output = %q, want message_id=m1", out.String())
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func TestMessageReceiveCmdJSON(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("hi")}},
|
|
|
|
|
}, nil
|
|
|
|
|
},
|
|
|
|
|
}
|
|
|
|
|
addr := startTestServer(t, fake)
|
|
|
|
|
|
|
|
|
|
flags := &globalFlags{server: addr, output: "json", insecure: true}
|
|
|
|
|
cmd := newMessageReceiveCmd(flags)
|
|
|
|
|
cmd.SetArgs([]string{"--queue", "orders"})
|
|
|
|
|
|
|
|
|
|
var out bytes.Buffer
|
|
|
|
|
cmd.SetOut(&out)
|
|
|
|
|
if err := cmd.Execute(); err != nil {
|
|
|
|
|
t.Fatalf("Execute: %v", err)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if !strings.Contains(out.String(), `"message_id":"m1"`) && !strings.Contains(out.String(), `"MessageID":"m1"`) {
|
|
|
|
|
t.Errorf("output = %q, want JSON containing message id m1", out.String())
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func TestMessageDeleteCmd(t *testing.T) {
|
|
|
|
|
var gotHandle string
|
|
|
|
|
fake := &fakeQueueService{
|
|
|
|
|
deleteMessage: func(ctx context.Context, req *kafkamgmtv1.DeleteMessageRequest) (*kafkamgmtv1.DeleteMessageResponse, error) {
|
|
|
|
|
gotHandle = req.ReceiptHandle
|
|
|
|
|
return &kafkamgmtv1.DeleteMessageResponse{}, nil
|
|
|
|
|
},
|
|
|
|
|
}
|
|
|
|
|
addr := startTestServer(t, fake)
|
|
|
|
|
|
|
|
|
|
flags := &globalFlags{server: addr, output: "table", insecure: true}
|
|
|
|
|
cmd := newMessageDeleteCmd(flags)
|
|
|
|
|
cmd.SetArgs([]string{"--queue", "orders", "--receipt-handle", "rh-1"})
|
|
|
|
|
|
|
|
|
|
var out bytes.Buffer
|
|
|
|
|
cmd.SetOut(&out)
|
|
|
|
|
if err := cmd.Execute(); err != nil {
|
|
|
|
|
t.Fatalf("Execute: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if gotHandle != "rh-1" {
|
|
|
|
|
t.Errorf("ReceiptHandle = %q, want rh-1", gotHandle)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func TestMessageSendCmdRequiresQueueAndBody(t *testing.T) {
|
|
|
|
|
flags := &globalFlags{server: "unused:1", output: "table"}
|
|
|
|
|
cmd := newMessageSendCmd(flags)
|
|
|
|
|
cmd.SetArgs([]string{})
|
|
|
|
|
cmd.SilenceUsage = true
|
|
|
|
|
cmd.SilenceErrors = true
|
|
|
|
|
|
|
|
|
|
if err := cmd.Execute(); err == nil {
|
|
|
|
|
t.Fatal("expected error for missing required flags")
|
|
|
|
|
}
|
|
|
|
|
}
|