Single-node · In-memory · gRPC

Heads up: I ended up switching to NATS + PostgreSQL. If you need consumer groups, bare NATS won't help — you need JetStream, which requires persistence. tinybroker fills the gap when you want in-memory consumer groups with no persistence overhead. The code works and is free to use; just know why you're reaching for it.

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)
        }
    }
}