Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 17 additions & 17 deletions THIRD-PARTY-NOTICES.md
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,7 @@ Total distributed third-party packages: 234
- [github.com/jbenet/go-temp-err-catcher](https://github.com/jbenet/go-temp-err-catcher/blob/v0.1.0/LICENSE)
- [github.com/jpillora/backoff](https://github.com/jpillora/backoff/blob/v1.0.0/LICENSE)
- [github.com/json-iterator/go](https://github.com/json-iterator/go/blob/v1.1.12/LICENSE)
- [github.com/klauspost/compress/zstd/internal/xxhash](https://github.com/klauspost/compress/blob/v1.19.0/zstd/internal/xxhash/LICENSE.txt)
- [github.com/klauspost/compress/zstd/internal/xxhash](https://github.com/klauspost/compress/blob/v1.19.1/zstd/internal/xxhash/LICENSE.txt)
- [github.com/klauspost/cpuid/v2](https://github.com/klauspost/cpuid/blob/v2.3.0/LICENSE)
- [github.com/knadh/koanf/maps](https://github.com/knadh/koanf/blob/maps/v0.1.2/maps/LICENSE)
- [github.com/knadh/koanf/providers/confmap](https://github.com/knadh/koanf/blob/providers/confmap/v1.0.0/providers/confmap/LICENSE)
Expand Down Expand Up @@ -138,8 +138,8 @@ Total distributed third-party packages: 234
- [cloud.google.com/go/auth](https://github.com/googleapis/google-cloud-go/blob/auth/v0.18.2/auth/LICENSE)
- [cloud.google.com/go/auth/oauth2adapt](https://github.com/googleapis/google-cloud-go/blob/auth/oauth2adapt/v0.2.8/auth/oauth2adapt/LICENSE)
- [cloud.google.com/go/compute/metadata](https://github.com/googleapis/google-cloud-go/blob/compute/metadata/v0.9.0/compute/metadata/LICENSE)
- [github.com/MicahParks/jwkset](https://github.com/MicahParks/jwkset/blob/v0.11.0/LICENSE)
- [github.com/MicahParks/keyfunc/v3](https://github.com/MicahParks/keyfunc/blob/v3.8.0/LICENSE)
- [github.com/MicahParks/jwkset](https://github.com/MicahParks/jwkset/blob/v0.11.1/LICENSE)
- [github.com/MicahParks/keyfunc/v3](https://github.com/MicahParks/keyfunc/blob/v3.8.1/LICENSE)
- [github.com/aws/aws-sdk-go-v2](https://github.com/aws/aws-sdk-go-v2/blob/v1.41.4/LICENSE.txt)
- [github.com/aws/aws-sdk-go-v2/config](https://github.com/aws/aws-sdk-go-v2/blob/config/v1.32.12/config/LICENSE.txt)
- [github.com/aws/aws-sdk-go-v2/credentials](https://github.com/aws/aws-sdk-go-v2/blob/credentials/v1.19.12/credentials/LICENSE.txt)
Expand All @@ -161,7 +161,7 @@ Total distributed third-party packages: 234
- [github.com/googleapis/enterprise-certificate-proxy/client](https://github.com/googleapis/enterprise-certificate-proxy/blob/v0.3.14/LICENSE)
- [github.com/ipfs/boxo](https://github.com/ipfs/boxo/blob/v0.30.0/LICENSE.md)
- [github.com/jackpal/go-nat-pmp](https://github.com/jackpal/go-nat-pmp/blob/v1.0.2/LICENSE)
- [github.com/klauspost/compress](https://github.com/klauspost/compress/blob/v1.19.0/LICENSE)
- [github.com/klauspost/compress](https://github.com/klauspost/compress/blob/v1.19.1/LICENSE)
- [github.com/kylelemons/godebug](https://github.com/kylelemons/godebug/blob/v1.1.0/LICENSE)
- [github.com/libp2p/go-libp2p-pubsub](https://github.com/libp2p/go-libp2p-pubsub/blob/41b11d5cb1a7/LICENSE-APACHE)
- [github.com/libp2p/go-libp2p/p2p/net/nat/internal/nat](https://github.com/libp2p/go-libp2p/blob/v0.48.0/p2p/net/nat/internal/nat/LICENSE)
Expand All @@ -177,11 +177,11 @@ Total distributed third-party packages: 234
- [github.com/open-telemetry/opentelemetry-collector-contrib/processor/deltatocumulativeprocessor](https://github.com/open-telemetry/opentelemetry-collector-contrib/blob/processor/deltatocumulativeprocessor/v0.148.0/processor/deltatocumulativeprocessor/LICENSE)
- [github.com/opentracing/opentracing-go](https://github.com/opentracing/opentracing-go/blob/v1.2.0/LICENSE)
- [github.com/prometheus/client_golang/exp](https://github.com/prometheus/client_golang/blob/d8591d0db856/exp/LICENSE)
- [github.com/prometheus/client_golang/prometheus](https://github.com/prometheus/client_golang/blob/v1.23.2/LICENSE)
- [github.com/prometheus/client_golang/prometheus](https://github.com/prometheus/client_golang/blob/v1.24.1/LICENSE)
- [github.com/prometheus/client_model/go](https://github.com/prometheus/client_model/blob/v0.6.2/LICENSE)
- [github.com/prometheus/common](https://github.com/prometheus/common/blob/v0.67.5/LICENSE)
- [github.com/prometheus/common](https://github.com/prometheus/common/blob/v0.70.1/LICENSE)
- [github.com/prometheus/otlptranslator](https://github.com/prometheus/otlptranslator/blob/v1.0.0/LICENSE)
- [github.com/prometheus/procfs](https://github.com/prometheus/procfs/blob/v0.19.2/LICENSE)
- [github.com/prometheus/procfs](https://github.com/prometheus/procfs/blob/v0.21.1/LICENSE)
- [github.com/prometheus/prometheus](https://github.com/prometheus/prometheus/blob/v0.311.3/LICENSE)
- [github.com/prometheus/sigv4](https://github.com/prometheus/sigv4/blob/v0.4.1/LICENSE)
- [github.com/puzpuzpuz/xsync/v4](https://github.com/puzpuzpuz/xsync/blob/v4.4.0/LICENSE)
Expand Down Expand Up @@ -231,26 +231,26 @@ Total distributed third-party packages: 234
- [github.com/googleapis/gax-go/v2/internallog](https://github.com/googleapis/gax-go/blob/v2.18.0/v2/LICENSE)
- [github.com/grafana/regexp](https://github.com/grafana/regexp/blob/f7b3be9d1853/LICENSE)
- [github.com/hashicorp/golang-lru/v2/simplelru](https://github.com/hashicorp/golang-lru/blob/v2.0.7/simplelru/LICENSE_list)
- [github.com/klauspost/compress/internal/snapref](https://github.com/klauspost/compress/blob/v1.19.0/internal/snapref/LICENSE)
- [github.com/klauspost/compress/s2](https://github.com/klauspost/compress/blob/v1.19.0/s2/LICENSE)
- [github.com/klauspost/compress/snappy](https://github.com/klauspost/compress/blob/v1.19.0/snappy/LICENSE)
- [github.com/klauspost/compress/internal/snapref](https://github.com/klauspost/compress/blob/v1.19.1/internal/snapref/LICENSE)
- [github.com/klauspost/compress/s2](https://github.com/klauspost/compress/blob/v1.19.1/s2/LICENSE)
- [github.com/klauspost/compress/snappy](https://github.com/klauspost/compress/blob/v1.19.1/snappy/LICENSE)
- [github.com/libp2p/go-netroute](https://github.com/libp2p/go-netroute/blob/v0.4.0/LICENSE)
- [github.com/miekg/dns](https://github.com/miekg/dns/blob/v1.1.72/LICENSE)
- [github.com/multiformats/go-base32](https://github.com/multiformats/go-base32/blob/v0.1.0/LICENSE)
- [github.com/munnerz/goautoneg](https://github.com/munnerz/goautoneg/blob/a7dc8b61c822/LICENSE)
- [github.com/pbnjay/memory](https://github.com/pbnjay/memory/blob/7b4eea64cf58/LICENSE)
- [github.com/pmezard/go-difflib/difflib](https://github.com/pmezard/go-difflib/blob/5d4384ee4fb2/LICENSE)
- [github.com/prometheus/client_golang/internal/github.com/golang/gddo/httputil](https://github.com/prometheus/client_golang/blob/v1.23.2/internal/github.com/golang/gddo/LICENSE)
- [github.com/prometheus/client_golang/internal/github.com/golang/gddo/httputil](https://github.com/prometheus/client_golang/blob/v1.24.1/internal/github.com/golang/gddo/LICENSE)
- [github.com/spaolacci/murmur3](https://github.com/spaolacci/murmur3/blob/v1.1.0/LICENSE)
- [github.com/wlynxg/anet](https://github.com/wlynxg/anet/blob/v0.0.5/LICENSE)
- [golang.org/x/crypto](https://cs.opensource.google/go/x/crypto/+/v0.53.0:LICENSE)
- [golang.org/x/crypto](https://cs.opensource.google/go/x/crypto/+/v0.54.0:LICENSE)
- [golang.org/x/exp](https://cs.opensource.google/go/x/exp/+/74f9aab9:LICENSE)
- [golang.org/x/net](https://cs.opensource.google/go/x/net/+/v0.56.0:LICENSE)
- [golang.org/x/net](https://cs.opensource.google/go/x/net/+/v0.57.0:LICENSE)
- [golang.org/x/oauth2](https://cs.opensource.google/go/x/oauth2/+/v0.36.0:LICENSE)
- [golang.org/x/sync/errgroup](https://cs.opensource.google/go/x/sync/+/v0.21.0:LICENSE)
- [golang.org/x/sys](https://cs.opensource.google/go/x/sys/+/v0.46.0:LICENSE)
- [golang.org/x/term](https://cs.opensource.google/go/x/term/+/v0.44.0:LICENSE)
- [golang.org/x/text](https://cs.opensource.google/go/x/text/+/v0.39.0:LICENSE)
- [golang.org/x/sync/errgroup](https://cs.opensource.google/go/x/sync/+/v0.22.0:LICENSE)
- [golang.org/x/sys](https://cs.opensource.google/go/x/sys/+/v0.47.0:LICENSE)
- [golang.org/x/term](https://cs.opensource.google/go/x/term/+/v0.45.0:LICENSE)
- [golang.org/x/text](https://cs.opensource.google/go/x/text/+/v0.40.0:LICENSE)
- [golang.org/x/time/rate](https://cs.opensource.google/go/x/time/+/v0.15.0:LICENSE)
- [gonum.org/v1/gonum/mathext](https://github.com/gonum/gonum/blob/v0.17.0/LICENSE)
- [google.golang.org/api](https://github.com/googleapis/google-api-go-client/blob/v0.272.0/LICENSE)
Expand Down
4 changes: 4 additions & 0 deletions buf.gen.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -4,3 +4,7 @@ plugins:
out: .
opt:
- module=github.com/getoptimum/optimum-gateway
- local: ["go", "tool", "protoc-gen-go-grpc"]
out: .
opt:
- module=github.com/getoptimum/optimum-gateway
61 changes: 60 additions & 1 deletion cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,17 @@ package main

import (
"context"
"errors"
"flag"
"fmt"
"io"
"net/http"
"os"
"os/signal"
"syscall"
"time"

"google.golang.org/grpc"

commonio "github.com/getoptimum/optimum-common/pkg/io"
"github.com/getoptimum/optimum-common/pkg/logger"
Expand All @@ -16,6 +21,8 @@ import (
"github.com/getoptimum/optimum-gateway/pkg/service/auth_token"
gateway "github.com/getoptimum/optimum-gateway/pkg/service/gossipsub-gateway"
"github.com/getoptimum/optimum-gateway/pkg/service/message_router"
"github.com/getoptimum/optimum-gateway/pkg/service/stream"
"github.com/getoptimum/optimum-gateway/pkg/service/streamhub"
"github.com/getoptimum/optimum-gateway/pkg/service/telemetry"
"github.com/getoptimum/optimum-gateway/pkg/utils"
)
Expand Down Expand Up @@ -155,7 +162,35 @@ func main() {
}
}

srvGateway, err := gateway.NewService(ctx, l, appConf, srvMessageRouter, authMgr)
// Consumer block-stream (ADR-0011), opt-in and off by default. The hub must
// exist before the gateway so it can be wired as an emit sink; WS and gRPC
// share the one hub.
var streamServer *stream.Server
var streamGRPCServer *stream.GRPCServer
var hub *streamhub.Service
if appConf.StreamEnable {
hub = streamhub.New()
authenticator := stream.NewConsumerAuthenticator(authMgr, appConf.StreamRequireAuth)
// One limiter across both transports keeps the caps global, not
// per-transport; config validation guarantees the caps are > 0.
limiter := stream.NewConnLimiter(appConf.StreamMaxConns, appConf.StreamMaxConnsPerSub)
streamServer = stream.NewServer(hub, authenticator, stream.Config{
Addr: appConf.StreamAddr,
MaxConns: appConf.StreamMaxConns,
MaxConnsPerSub: appConf.StreamMaxConnsPerSub,
BufferSize: appConf.StreamBufferSize,
Limiter: limiter,
}, l)
streamGRPCServer = stream.NewGRPCServer(hub, authenticator, stream.Config{
Addr: appConf.StreamGRPCAddr,
MaxConns: appConf.StreamMaxConns,
MaxConnsPerSub: appConf.StreamMaxConnsPerSub,
BufferSize: appConf.StreamBufferSize,
Limiter: limiter,
}, l)
Comment thread
swarna1101 marked this conversation as resolved.
}

srvGateway, err := gateway.NewService(ctx, l, appConf, srvMessageRouter, authMgr, gateway.WithStreamHub(hub))
if err != nil {
l.Fatal("unable to initialize gossipsub gateway", err)
}
Expand All @@ -171,9 +206,33 @@ func main() {
}
}()

if streamServer != nil {
go func() {
if runErr := streamServer.Run(); runErr != nil && !errors.Is(runErr, http.ErrServerClosed) {
l.Fatal("failed to run consumer stream server", runErr)
}
}()
}

if streamGRPCServer != nil {
go func() {
if runErr := streamGRPCServer.Run(); runErr != nil && !errors.Is(runErr, grpc.ErrServerStopped) {
l.Fatal("failed to run consumer stream grpc server", runErr)
}
}()
}

<-c // This blocks the main thread until an interrupt is received
cancel()
_ = appRouter.Stop()
if streamServer != nil {
shutdownCtx, cancelShutdown := context.WithTimeout(context.Background(), 5*time.Second)
_ = streamServer.Stop(shutdownCtx)
cancelShutdown()
}
if streamGRPCServer != nil {
streamGRPCServer.Stop()
}
srvGateway.Stop()
if lokiDone != nil {
<-lokiDone // wait for final Loki flush to complete
Expand Down
21 changes: 10 additions & 11 deletions docs/adr/0011-gateway-consumer-block-stream.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
# ADR-0011: Gateway consumer block-stream API (WebSocket + gRPC)

**Status:** Approved (implementation pending)
**Status:** Approved (implemented)
**Date:** 2026-08-05

## Context
Expand Down Expand Up @@ -47,7 +47,7 @@ The stream carries a small transport-neutral frame union, encoded per transport
connection's cumulative `dropped` count so the consumer knows it missed events.

Each subscriber has a bounded ring buffer (default 64). On overflow the hub drops
the oldest event, increments the per-connection `dropped` counter, and sends a
the incoming event, increments the per-connection `dropped` counter, and sends a
Comment thread
swarna1101 marked this conversation as resolved.
`lagged` frame. The emit from ingest is a non-blocking send — it never waits on a
consumer, so a slow/stalled subscriber cannot backpressure ingest. At-most-once,
no replay (see Non-goals).
Expand All @@ -59,7 +59,7 @@ Both read-only — consumers cannot publish into the mesh.
* **WebSocket** — `GET /api/v1/stream/blocks` on its own listener
(`OPT_STREAM_ADDR`) and Fiber app, so `/metrics` and `/health` stay off the
exposed port. Params: `mode=metadata|raw`, `topics=beacon_block`.
* **gRPC** — `BlockStream.Subscribe(SubscribeRequest) returns (stream
* **gRPC** — `BlockStreamService.Subscribe(SubscribeRequest) returns (stream
BlockEvent)` on `OPT_STREAM_GRPC_ADDR`, from a new
`proto/getoptimum/optimum_gateway/service/stream/v1/stream.proto`.

Expand All @@ -80,9 +80,9 @@ Consumers present a JWT minted by `auth.getoptimum.io` for a new audience
unauthenticated connections are rejected, never subscribed.
* No scope claim in v1 — `aud=stream` is the authorization. There is one topic
(`beacon_block`), so a valid stream token grants read of the whole stream. The
gateway still caps connections per `sub` and globally and rate-limits events
per connection, but those are config-driven, not per-token. (If topics beyond
`beacon_block` are added later, a scope claim can gate them then.)
gateway still caps connections per `sub` and globally, but those caps are
config-driven, not per-token. (If topics beyond `beacon_block` are added later,
a scope claim can gate them then.)
* The verifier sits behind a `ConsumerAuthenticator` interface so an
operator-local key mode can replace central auth later without touching the
transport or hub. Interface now, local impl later.
Expand All @@ -105,9 +105,8 @@ Consumers present a JWT minted by `auth.getoptimum.io` for a new audience
beyond `127.0.0.1`) requires TLS — native or a trusted TLS-terminating proxy.
Startup validation must **reject** a non-loopback listener when
`stream_require_auth=false`; disabling auth is allowed only on a loopback bind for
local dev. (Read/idle timeouts, max frame size, and the per-connection event-rate
cap are the other DoS mitigations — see Consequences — with concrete values fixed
in the implementation.)
local dev. (Read/idle timeouts, max frame size, and the connection caps are the
other DoS mitigations — see Consequences.)

### Data model — `BlockEvent`

Expand Down Expand Up @@ -153,8 +152,8 @@ flowchart LR
## Consequences

* New public surface means DoS exposure. Mitigations: auth-before-subscribe,
connection caps (global + per-`sub`), per-connection rate cap, read/idle
timeouts, max frame size, WS keepalive, TLS (proxy or native).
connection caps (global + per-`sub`), read/idle timeouts, max frame size, WS
keepalive, gRPC keepalive, TLS (proxy or native).
* Central-auth coupling, bounded by the `ConsumerAuthenticator` seam.
* Drop-on-lag means slow consumers miss events — surfaced via `lagged`/`dropped`
rather than silently.
Expand Down
Loading
Loading