Go
Dependencies
go get google.golang.org/grpc
go get google.golang.org/protobuf
Generate stubs from the proto files:
buf generate https://codeberg.org/phughk/tinybroker-rs.git#branch=main,subdir=proto
Or manually with protoc:
protoc -I proto \
--go_out=gen --go_opt=paths=source_relative \
--go-grpc_out=gen --go-grpc_opt=paths=source_relative \
tinybroker/v1/service.proto \
tinybroker/v1/client.proto \
tinybroker/v1/server.proto \
tinybroker/v1/shared.proto
Connect
import (
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
pb "your/module/gen/tinybroker/v1"
)
conn, err := grpc.NewClient("tinybroker:50051",
grpc.WithTransportCredentials(insecure.NewCredentials()))
if err != nil {
log.Fatal(err)
}
defer conn.Close()
client := pb.NewBrokerClient(conn)
Publish
ctx := context.Background()
_, err = client.Publish(ctx, &pb.PublishRequest{
Topic: "events.user.signup",
Payload: []byte(`{"user_id":"abc123","email":"alice@example.com"}`),
})
Subscribe (fan-out)
stream, err := client.Subscribe(ctx)
if err != nil {
log.Fatal(err)
}
// Open subscription
err = stream.Send(&pb.SubscribeCommand{
Command: &pb.SubscribeCommand_Open{
Open: &pb.OpenSubscription{
TopicPatterns: []string{"events.user.*", "events.order.*"},
Mode: pb.DeliveryMode_DELIVERY_MODE_MULTI,
},
},
})
// Read the subscription_id confirmation
firstFrame, _ := stream.Recv()
subscriptionID := firstFrame.GetSubscriptionId()
// Receive messages
for {
ev, err := stream.Recv()
if err != nil {
break // reconnect
}
if msg := ev.GetMessage(); msg != nil {
fmt.Printf("topic=%s payload=%s\n", msg.Topic, msg.Payload)
}
}
Subscribe (consumer group)
stream.Send(&pb.SubscribeCommand{
Command: &pb.SubscribeCommand_Open{
Open: &pb.OpenSubscription{
SubscriptionId: "job-workers", // shared across all workers
TopicPatterns: []string{"jobs.*"},
Mode: pb.DeliveryMode_DELIVERY_MODE_BALANCED,
ConsumerGroup: "job-workers",
},
},
})
for {
ev, err := stream.Recv()
if msg := ev.GetMessage(); msg != nil {
process(msg.Payload)
client.Ack(ctx, &pb.AckRequest{
SubscriptionId: "job-workers",
SystemId: msg.SystemId,
})
}
}
Dynamic pattern management
// Add patterns to a running subscription
stream.Send(&pb.SubscribeCommand{
Command: &pb.SubscribeCommand_AddPatterns{
AddPatterns: &pb.AddTopicPatterns{
TopicPatterns: []string{"alerts.*"},
},
},
})
// Remove patterns
stream.Send(&pb.SubscribeCommand{
Command: &pb.SubscribeCommand_RemovePatterns{
RemovePatterns: &pb.RemoveTopicPatterns{
TopicPatterns: []string{"events.user.*"},
},
},
})
Get topics
resp, err := client.GetTopics(ctx, &pb.GetTopicsRequest{
Pattern: "events.*",
})
fmt.Println(resp.Topics)
Reconnect pattern
func subscribeWithReconnect(ctx context.Context, client pb.BrokerClient) {
for {
if err := subscribe(ctx, client); err != nil {
log.Printf("subscription error: %v — reconnecting in 2s", err)
time.Sleep(2 * time.Second)
}
}
}