Files
kmsvc-sdk/examples/sendreceive/main.go
T

60 lines
1.5 KiB
Go
Raw Normal View History

// Command sendreceive is a manual smoke test: send a message, receive it,
// then delete it. Requires a running kafaka_management_service instance.
package main
import (
"context"
"flag"
"log"
"time"
kmsvc "forgejo.riotpiao.com/rock/kmsvc-sdk"
)
func main() {
target := flag.String("target", "localhost:8443", "kmsvc gRPC address")
queue := flag.String("queue", "demo", "queue name")
token := flag.String("token", "", "bearer token")
flag.Parse()
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
var opts []kmsvc.Option
if *token != "" {
opts = append(opts, kmsvc.WithTokenSource(kmsvc.StaticToken(*token)))
}
client, err := kmsvc.New(ctx, *target, opts...)
if err != nil {
log.Fatalf("kmsvc.New: %v", err)
}
defer client.Close()
sendOut, err := client.SendMessage(ctx, kmsvc.SendMessageInput{
QueueName: *queue,
Body: []byte("hello from kmsvc-sdk"),
})
if err != nil {
log.Fatalf("SendMessage: %v", err)
}
log.Printf("sent message_id=%s", sendOut.MessageID)
msgs, err := client.ReceiveMessage(ctx, *queue, kmsvc.ReceiveOptions{
MaxNumberOfMessages: 1,
2026-08-30 09:46:58 -07:00
WaitTimeSeconds: 10,
})
if err != nil {
log.Fatalf("ReceiveMessage: %v", err)
}
if len(msgs) == 0 {
log.Fatal("no messages received within wait window")
}
log.Printf("received message_id=%s body=%q", msgs[0].MessageID, msgs[0].Body)
if err := client.DeleteMessage(ctx, *queue, msgs[0].ReceiptHandle); err != nil {
log.Fatalf("DeleteMessage: %v", err)
}
log.Println("deleted message")
}