A gRPC event server that bridges the event/v3 pub-sub library to remote clients over gRPC, REST/JSON, WebSocket, and Server-Sent Events (SSE). Clients connect without needing direct access to the underlying transport (Redis, Kafka, NATS, etc.).
- gRPC EventService: Register/unregister events, publish, subscribe (server-streaming), batch ack, health check
- HTTP Gateway: REST/JSON via gRPC-Gateway with full OpenAPI mapping
- WebSocket Streaming: Bidirectional JSON protocol with ack support and heartbeats
- SSE Streaming: One-way server-sent events for browser-friendly subscriptions
- Pluggable Authorization: Interface-based auth with per-operation granularity
- Ack Tracking: Timeout-based acknowledgment with automatic nack on expiration
- RemoteTransport: Implements
transport.Transport- drop-in replacement for any local transport - Circuit Breaker: Configurable failure threshold and recovery timeout
- Retry with Backoff: Exponential backoff with jitter for transient failures
- Subscribe Reconnect: Automatic stream reconnection with configurable max errors
- Connection Lifecycle: State tracking with callbacks (disconnected, connecting, connected, closed)
- Structured logging via
slog - Health endpoint with transport status and latency
- Panic recovery interceptors
- Request/error logging interceptors
- OpenTelemetry metrics emitted via the global
MeterProvider(no-op until you register an exporter):event_server.publishes_total,event_server.publish_duration_seconds,event_server.subscribe_streams_active,event_server.messages_sent_total,event_server.acks_total
go get github.com/rbaliyan/event-serverA complete gRPC + HTTP server with an in-memory transport:
package main
import (
"context"
"log"
"log/slog"
"net"
"net/http"
"os"
"os/signal"
"syscall"
"github.com/rbaliyan/event-server/gateway"
eventpb "github.com/rbaliyan/event-server/proto/event/v1"
"github.com/rbaliyan/event-server/service"
"github.com/rbaliyan/event/v3/transport/channel"
"google.golang.org/grpc"
)
func main() {
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer cancel()
logger := slog.New(slog.NewTextHandler(os.Stdout, nil))
// Create transport (channel, redis, nats, kafka - any event/v3 transport)
ch := channel.New()
defer func() { _ = ch.Close(ctx) }()
// Create event service
svc, _ := service.NewService(ch,
service.WithSecurityGuard(service.AllowAll()),
service.WithLogger(logger),
)
defer svc.Stop()
// Start gRPC server — wire auth interceptors first
grpcServer := grpc.NewServer(
grpc.ChainUnaryInterceptor(
svc.UnaryInterceptor(),
service.LoggingInterceptor(logger),
service.RecoveryInterceptor(logger),
),
grpc.ChainStreamInterceptor(
svc.StreamInterceptor(),
service.StreamLoggingInterceptor(logger),
service.StreamRecoveryInterceptor(logger),
),
)
eventpb.RegisterEventServiceServer(grpcServer, svc)
lis, _ := net.Listen("tcp", ":9090")
go grpcServer.Serve(lis)
// Start HTTP gateway (REST + WebSocket + SSE)
handler, _ := gateway.NewHandler(ctx, "localhost:9090", gateway.WithInsecure())
httpServer := &http.Server{Addr: ":8080", Handler: handler}
go httpServer.ListenAndServe()
<-ctx.Done()
grpcServer.GracefulStop()
httpServer.Shutdown(context.Background())
}Connect to the server using RemoteTransport with the standard event.Bus API:
package main
import (
"context"
"fmt"
"log"
"time"
"github.com/rbaliyan/event-server/client"
"github.com/rbaliyan/event/v3"
)
type Order struct {
ID string `json:"id"`
Total float64 `json:"total"`
}
func main() {
ctx := context.Background()
// Connect to event server - no Redis/Kafka/NATS needed on the client!
remote, _ := client.New("localhost:9090",
client.WithInsecure(),
client.WithRetry(3, 100*time.Millisecond, 5*time.Second),
client.WithCircuitBreaker(5, 30*time.Second),
)
remote.Connect(ctx)
defer remote.Close(ctx)
// Use with event.Bus - same API as any local transport
bus, _ := event.NewBus("my-app", event.WithTransport(remote))
defer bus.Close(ctx)
orderEvent := event.New[Order]("order.created")
event.Register(ctx, bus, orderEvent)
orderEvent.Subscribe(ctx, func(ctx context.Context, ev event.Event[Order], order Order) error {
fmt.Printf("order: %s ($%.2f)\n", order.ID, order.Total)
return nil // Automatically acks via gRPC
})
orderEvent.Publish(ctx, Order{ID: "ORD-1", Total: 99.99})
}┌──────────────────────────────────────────────────────────┐
│ Event Server │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ service.Service │ │
│ │ (gRPC EventService + SecurityGuard + AckTracker) │ │
│ └────────────────────────┬────────────────────────────┘ │
│ │ │
│ transport.Transport │
│ (channel/redis/nats/kafka) │
│ │
│ ┌────────────┐ ┌──────────────────────────────────────┐ │
│ │ gRPC │ │ gateway.Handler │ │
│ │ :9090 │ │ REST (gRPC-Gateway) + WS + SSE │ │
│ │ │ │ :8080 │ │
│ └────────────┘ └──────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────┘
▲ ▲ ▲
│ │ │
gRPC Client HTTP/REST WebSocket/SSE
(RemoteTransport) (curl, etc.) (browser, etc.)
| RPC | HTTP Mapping | Description |
|---|---|---|
RegisterEvent |
POST /v1/events/{name} |
Create transport resources for a named event |
UnregisterEvent |
DELETE /v1/events/{name} |
Remove transport resources |
ListEvents |
GET /v1/events |
List all registered event names |
Publish |
POST /v1/events/{event}/messages |
Send a message to an event |
Subscribe |
gRPC only | Server-streaming subscription |
Ack |
POST /v1/ack |
Batch acknowledge messages |
Health |
GET /v1/health |
Server and transport health status |
Connect to GET /v1/events/{name}/subscribe for bidirectional JSON streaming.
Query parameters:
| Parameter | Values | Description |
|---|---|---|
delivery_mode |
broadcast, worker_pool |
Message distribution mode |
worker_group |
string | Worker group name (worker_pool mode) |
start_from |
beginning, latest |
Where to start reading |
latest_only |
true |
Only deliver most recent message |
consumer_id |
string | Stable consumer ID for checkpoint resume |
Server sends:
{"type": "message", "id": "msg-1", "source": "my-app", "payload": "...", "ack_id": "ack-123"}
{"type": "heartbeat"}
{"type": "error", "error": "subscribe failed: ..."}Client sends:
{"type": "ack", "ack_id": "ack-123"}
{"type": "ack", "ack_id": "ack-456", "error": "processing failed"}Connect to GET /v1/events/{name}/stream for one-way server-sent events. Supports the same query parameters as WebSocket.
event: message
data: {"id":"msg-1","source":"my-app","payload":"...","ack_id":"ack-123"}
event: heartbeat
data: {}
# Register an event
curl -X POST http://localhost:8080/v1/events/order.created
# List events
curl http://localhost:8080/v1/events
# Publish a message
curl -X POST http://localhost:8080/v1/events/order.created/messages \
-H 'Content-Type: application/json' \
-d '{"payload":"eyJpZCI6Im9yZGVyLTEifQ==","metadata":{"source":"api"}}'
# Acknowledge a message
curl -X POST http://localhost:8080/v1/ack \
-H 'Content-Type: application/json' \
-d '{"entries":[{"ack_id":"ack-123"}]}'
# Health check
curl http://localhost:8080/v1/healthsvc, _ := service.NewService(transport,
service.WithSecurityGuard(myGuard), // Default: DenyAll()
service.WithLogger(logger), // Default: slog.Default()
service.WithAckTimeout(30*time.Second), // Default: 30s
)
defer svc.Stop()SecurityGuard interface:
type SecurityGuard interface {
// Authenticate extracts identity from the incoming context (e.g. gRPC metadata).
Authenticate(ctx context.Context) (Identity, error)
// Authorize checks whether identity may perform action (the gRPC full method name).
Authorize(ctx context.Context, id Identity, action string) (Decision, error)
}
type Identity interface {
UserID() string
Claims() map[string]any
}
type Decision struct {
Allowed bool
Scope string // e.g. "all", "owned", "tenant"
Reason string // included in PermissionDenied error
}Wire auth into gRPC with the service's interceptors (must be registered on the server):
grpcServer := grpc.NewServer(
grpc.ChainUnaryInterceptor(svc.UnaryInterceptor(), ...),
grpc.ChainStreamInterceptor(svc.StreamInterceptor(), ...),
)Built-in: AllowAll() (dev only), DenyAll() (default).
Interceptors:
service.LoggingInterceptor(logger) // Log method calls and errors
service.RecoveryInterceptor(logger) // Catch panics, return Internal
service.StreamLoggingInterceptor(logger) // Log stream lifecycle
service.StreamRecoveryInterceptor(logger) // Catch panics in streamsremote, err := client.New("localhost:9090",
client.WithInsecure(),
client.WithTLS(tlsConfig),
client.WithRetry(3, 100*time.Millisecond, 5*time.Second),
client.WithCallTimeout(10*time.Second),
client.WithCircuitBreaker(5, 30*time.Second),
client.WithSubscribeReconnect(true, time.Second),
client.WithSubscribeMaxErrors(10),
client.WithSubscribeBufferSize(100),
client.WithKeepalive(30*time.Second, 10*time.Second),
client.WithStateCallback(func(s client.ConnState) { /* ... */ }),
client.WithStreamErrorCallback(func(err error) { /* ... */ }),
client.WithDialOptions(grpc.WithBlock()),
)
remote.Connect(ctx) // Establish connection
remote.Ready() // Check if connected
remote.State() // ConnStateDisconnected|Connecting|Connected|Closed
remote.Close(ctx) // Graceful shutdownThe transport maps gRPC status codes back to the standard transport sentinels
(ErrEventNotRegistered, ErrTransportClosed, ErrPublishTimeout, …) so callers
can use errors.Is. Permission failures surface as *client.PermissionDeniedError
(matches errors.Is(err, client.ErrPermissionDenied)); other non-standard codes
surface as *client.RemoteError carrying the original codes.Code and message.
// Remote mode: connects to a gRPC server
handler, err := gateway.NewHandler(ctx, "localhost:9090",
gateway.WithInsecure(),
gateway.WithTLS(tlsConfig),
gateway.WithHeartbeatInterval(30*time.Second),
gateway.WithWSOriginPatterns("example.com", "*.example.com"),
gateway.WithMuxOptions(runtime.WithMarshalerOption(...)),
gateway.WithDialOptions(grpc.WithBlock()),
)
// In-process mode: calls service directly (no network hop)
handler, err := gateway.NewInProcessHandler(ctx, svc,
gateway.WithHeartbeatInterval(30*time.Second),
gateway.WithWSOriginPatterns("example.com"),
)eventctl is an operational CLI for event-server, event-scheduler, and the schema
registry over their HTTP/REST gateways. It uses the standard library only. Install
with just install (or go install ./cmd/eventctl).
eventctl [--server http://host:port] [--scheduler http://host:port] [--schema http://host:port] <command>
events list # list registered events
events pub <event> [payload] # publish a message ("-" reads payload from stdin)
events sub <event> # subscribe and stream messages over SSE (ctrl+c to stop)
events health # server health check
scheduler list # list scheduled messages (-event/-limit/-before/-after)
scheduler get <id> # get a scheduled message by ID
scheduler health # scheduler health check
schema list # list event schemas
schema get <event> # get a schema
schema set <event> [flags] # create/update a schema (-monitor/-idempotency/-poison/-timeout/...)
schema delete <event> # delete a schemaThe --scheduler and --schema flags default to --server, so the extra flags
are only needed when the services listen on different ports.
gRPC server + HTTP gateway as a standalone service:
[Event Server] ── gRPC :9090 + HTTP :8080
│
transport (Redis/Kafka/NATS)
Add the EventService to an existing gRPC server with custom auth:
eventpb.RegisterEventServiceServer(yourGRPCServer, eventSvc)See examples/embedded for a complete example with role-based authorization.
Use RemoteTransport with event.Bus for transparent remote event handling:
bus, _ := event.NewBus("my-app", event.WithTransport(remote))
// Use Publish/Subscribe exactly like a local transportPart of the event/v3 ecosystem:
| Package | Description |
|---|---|
| event | Core event bus with transports (channel, Redis, NATS, Kafka) |
| event-server | gRPC server + HTTP gateway + remote client (this package) |
| event-mongodb | MongoDB Change Stream transport (CDC) |
| event-dlq | Dead Letter Queue management |
| event-scheduler | Delayed/scheduled message delivery |
| event-extras | Rate limiting and saga orchestration |
Install tools via mise:
mise installjust build # go build ./...
just smoke # fast pre-merge gate: go test -tags smoke -race ./smoke/...
just test # go test -v ./...
just test-race # go test -race ./...
just test-cover # go test -cover ./...
just test-integration # go test -tags integration -race ./... (needs docker or podman)
just lint # golangci-lint run ./...
just fmt # go fmt ./...
just tidy # go mod tidy
just vulncheck # govulncheck ./...
just proto # Regenerate protobuf code
just release # Tag and push a new patch releaseEnd-to-end tests that run the full stack against a real backend live behind the
integration build tag and read the backend address from EVENT_REDIS_ADDR:
just test-integrationThe recipe provisions throwaway Redis and NATS JetStream containers
(auto-detecting docker, falling back to podman), exports
EVENT_REDIS_ADDR / EVENT_NATS_URL, and runs the tagged tests. The shared
scenario suite (round-trip, ordering, worker-pool, broadcast) runs against every
configured backend, proving the server is backend-agnostic. In CI the backends
are provided as service/step containers. Tests skip cleanly when the env vars
are unset (e.g. no container runtime installed).
just proto # Generate to root (for development)
just proto-local # Generate to proto/ directoryRequires protoc, protoc-gen-go, protoc-gen-go-grpc, and protoc-gen-grpc-gateway (all managed by mise).
MIT License - see LICENSE for details.