-
Notifications
You must be signed in to change notification settings - Fork 1
feat: add consumer block-stream gRPC transport #68
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -263,6 +263,9 @@ func (c *AppConfig) Validate() error { | |
| if err := validateStreamListener("stream_grpc_addr", c.StreamGRPCAddr, c.StreamRequireAuth); err != nil { | ||
| return err | ||
| } | ||
| if strings.TrimSpace(c.StreamAddr) == strings.TrimSpace(c.StreamGRPCAddr) { | ||
| return fmt.Errorf("stream_addr and stream_grpc_addr must differ, got %q", c.StreamAddr) | ||
| } | ||
|
Comment on lines
+266
to
+268
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🩺 Stability & Availability | 🟠 Major | 🏗️ Heavy lift 🧩 Analysis chain🏁 Script executed: #!/bin/bash
set -euo pipefail
ast-grep outline pkg/config/config.go --lang go
rg -n -C 4 'StreamAddr|StreamGRPCAddr|validateStreamListener|net\.Listen' \
pkg/config/config.go pkg/config/config_test.go cmd/main.goRepository: getoptimum/optimum-gateway Length of output: 7348 🏁 Script executed: #!/bin/bash
set -euo pipefail
printf '%s\n' '--- validation and listener startup ---'
sed -n '220,320p' pkg/config/config.go
sed -n '150,230p' cmd/main.go
printf '%s\n' '--- focused tests and listener implementations ---'
rg -n -C 5 'Run\(|Fatal|Listen\(|NewServer|NewGRPCServer' cmd pkg/service/stream pkg/config --glob '*.go'
printf '%s\n' '--- bind-equivalence probe ---'
python3 - <<'PY'
import socket
cases = [
(":0", "0.0.0.0:0"),
("127.0.0.1:0", "localhost:0"),
(":9600", "0.0.0.0:9600"),
]
for first, second in cases:
def split(addr):
host, port = addr.rsplit(":", 1)
return host or "0.0.0.0", int(port)
h1, p1 = split(first)
h2, p2 = split(second)
if p1 == 0:
# Use one concrete port so the second bind tests the same endpoint.
probe = socket.socket()
probe.bind((h1, 0))
p1 = probe.getsockname()[1]
probe.close()
p2 = p1
s1 = socket.socket()
s2 = socket.socket()
try:
s1.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
s2.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
s1.bind((h1, p1))
try:
s2.bind((h2, p2))
result = "both bind"
except OSError as exc:
result = f"second bind fails: {exc}"
finally:
s1.close()
s2.close()
print(f"{first} vs {second}: {result}")
PYRepository: getoptimum/optimum-gateway Length of output: 28245 🏁 Script executed: #!/bin/bash
set -euo pipefail
python3 - <<'PY'
import socket
port = 39641
for reuse in (False, True):
first = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
second = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
try:
if reuse:
first.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
second.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
first.bind(("0.0.0.0", port))
try:
second.bind(("0.0.0.0", port))
result = "both bind"
except OSError as exc:
result = f"second bind fails: {exc}"
print(f"Python SO_REUSEADDR={reuse}: {result}")
finally:
first.close()
second.close()
PY
if command -v go >/dev/null 2>&1; then
tmpdir="$(mktemp -d)"
trap 'rm -rf "$tmpdir"' EXIT
cat >"$tmpdir/main.go" <<'GO'
package main
import (
"fmt"
"net"
)
func main() {
first, err := net.Listen("tcp", ":39642")
if err != nil {
panic(err)
}
defer first.Close()
second, err := net.Listen("tcp", "0.0.0.0:39642")
if err != nil {
fmt.Printf("Go net.Listen: second bind fails: %v\n", err)
return
}
defer second.Close()
fmt.Println("Go net.Listen: both bind")
}
GO
(cd "$tmpdir" && GO111MODULE=off go run main.go)
else
echo "Go compiler unavailable"
fi
sed -n '286,335p' pkg/config/config_test.goRepository: getoptimum/optimum-gateway Length of output: 392 🏁 Script executed: #!/bin/bash
set -euo pipefail
python3 - <<'PY'
import errno
import socket
port = 39643
for reuse in (False, True):
first = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
second = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
try:
if reuse:
first.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
second.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
first.bind(("0.0.0.0", port))
first.listen(16)
try:
second.bind(("0.0.0.0", port))
second.listen(16)
result = "both listen"
except OSError as exc:
result = f"second bind/listen fails: {exc}"
print(f"Python SO_REUSEADDR={reuse}: {result}")
finally:
first.close()
second.close()
PY
printf '%s\n' '--- Go listener socket-option implementation, if present ---'
if command -v go >/dev/null 2>&1; then
goroot="$(go env GOROOT 2>/dev/null || true)"
if [ -n "$goroot" ]; then
rg -n -C 4 'setDefaultListenerSockopts|SO_REUSEADDR' \
"$goroot/src/net" 2>/dev/null || true
fi
fiRepository: getoptimum/optimum-gateway Length of output: 9361 Reject bind-equivalent listener addresses.
Reject listener address pairs that can overlap, including wildcard forms. Add focused configuration tests for these cases. 🤖 Prompt for AI AgentsSource: Coding guidelines |
||
| if c.StreamMaxConns <= 0 { | ||
| return fmt.Errorf("stream_max_conns must be > 0") | ||
| } | ||
|
|
||
| Original file line number | Diff line number | Diff line change | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| @@ -0,0 +1,157 @@ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| package stream | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| import ( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "context" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "net" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "strings" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "google.golang.org/grpc" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "google.golang.org/grpc/codes" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "google.golang.org/grpc/keepalive" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "google.golang.org/grpc/metadata" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "google.golang.org/grpc/status" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
Comment on lines
+8
to
+12
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Import for the suggestion on the constructor below. Ordering matches
Suggested change
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Added with the keepalive / MaxConcurrentStreams constructor change. |
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "github.com/getoptimum/optimum-common/pkg/logger" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| streamv1 "github.com/getoptimum/optimum-gateway/pkg/service/stream/v1" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "github.com/getoptimum/optimum-gateway/pkg/service/streamhub" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "github.com/getoptimum/optimum-gateway/pkg/service/telemetry" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // GRPCServer serves the consumer block-stream over gRPC on its own listener, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // reusing the hub, authenticator, and connection caps of the WS transport. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| type GRPCServer struct { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| streamv1.UnimplementedBlockStreamServiceServer | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| hub *streamhub.Service | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| auth ConsumerAuthenticator | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| cfg Config | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| log logger.AppLogger | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| limiter *ConnLimiter | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| grpcSrv *grpc.Server | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // maxConcurrentStreams bounds Subscribe streams per connection. The caps run | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // after auth, so one unauthenticated socket would otherwise be unbounded. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| const maxConcurrentStreams = 256 | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // NewGRPCServer builds the consumer gRPC server. It does not start listening; | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // call Run. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| func NewGRPCServer(hub *streamhub.Service, auth ConsumerAuthenticator, cfg Config, log logger.AppLogger) *GRPCServer { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| cfg = withDefaults(cfg) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| g := &GRPCServer{ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| hub: hub, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| auth: auth, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| cfg: cfg, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| log: log.With(logger.WithService("stream-grpc")), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| limiter: cfg.Limiter, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| grpcSrv: grpc.NewServer( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // Reap dead peers on the WS clock; the gRPC default is a 2h ping. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| grpc.KeepaliveParams(keepalive.ServerParameters{Time: pingPeriod, Timeout: writeWait}), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| grpc.MaxConcurrentStreams(maxConcurrentStreams), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| streamv1.RegisterBlockStreamServiceServer(g.grpcSrv, g) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return g | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
Comment on lines
+36
to
+54
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Calling Dead-peer reaping. In grpc-go v1.82.1 the server keepalive defaults are Worth noting that Concurrent streams. A constant rather than
Suggested change
The
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Done , keepalive uses the WS ping clock (Time=pingPeriod, Timeout=writeWait) and MaxConcurrentStreams is 256. |
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // Run serves until Stop is called; grpc.ErrServerStopped on a clean stop is | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // treated as normal by the caller. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| func (g *GRPCServer) Run() error { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| lis, err := net.Listen("tcp", g.cfg.Addr) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err != nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return err | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| g.log.Info("starting consumer stream grpc server", logger.WithString("addr", g.cfg.Addr)) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return g.grpcSrv.Serve(lis) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // Stop hard-stops the server, canceling active Subscribe streams so shutdown | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // does not block on long-lived consumers. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| func (g *GRPCServer) Stop() { g.grpcSrv.Stop() } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // Subscribe authenticates and enforces caps before opening the stream, then | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // drains the buffer as proto frames (metadata omits Raw); lagged on overflow. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| func (g *GRPCServer) Subscribe(req *streamv1.SubscribeRequest, stream grpc.ServerStreamingServer[streamv1.BlockEvent]) error { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| mode, ok := normalizeMode(req.GetMode()) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if !ok { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return status.Error(codes.InvalidArgument, "invalid mode") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if !topicsOK(req.GetTopics()...) { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return status.Error(codes.InvalidArgument, "unsupported topics") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ctx := stream.Context() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| subject, err := g.auth.Authenticate(metadataToken(ctx)) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err != nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| telemetry.RecordStreamAuthFailure() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return status.Error(codes.Unauthenticated, "unauthorized") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if !g.limiter.acquire(subject) { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return status.Error(codes.ResourceExhausted, "too many connections") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| defer g.limiter.release(subject) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| sub := g.hub.Subscribe(g.cfg.BufferSize) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| defer sub.Close() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| raw := mode == modeRaw | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| var lastDropped uint64 | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| for { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| select { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| case <-ctx.Done(): | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return ctx.Err() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| case ev, ok := <-sub.Events(): | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if !ok { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return nil | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if d := sub.Dropped(); d != lastDropped { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| lastDropped = d | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| lag := &streamv1.BlockEvent{Frame: &streamv1.BlockEvent_Lagged{Lagged: &streamv1.Lagged{Dropped: d}}} | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err := stream.Send(lag); err != nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return err | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err := stream.Send(toProto(ev, raw)); err != nil { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return err | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| telemetry.RecordStreamEventSent() | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| func toProto(ev *streamhub.BlockEvent, raw bool) *streamv1.BlockEvent { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| b := &streamv1.Block{ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Slot: ev.Slot, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ProposerIndex: ev.ProposerIndex, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ParentRoot: ev.ParentRoot, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| StateRoot: ev.StateRoot, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| BlockSizeBytes: ev.BlockSizeBytes, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Topic: ev.Topic, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Source: string(ev.Source), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ReceivedAtMs: ev.ReceivedAtMs, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| GatewayId: ev.GatewayID, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| ForkDigest: ev.ForkDigest, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Stale: ev.Stale, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if raw { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| b.Raw = ev.Raw | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return &streamv1.BlockEvent{Frame: &streamv1.BlockEvent_Block{Block: b}} | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // metadataToken reads the consumer JWT from the "authorization" gRPC metadata, | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // accepting either a bare token or a "Bearer <jwt>" value. | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| func metadataToken(ctx context.Context) string { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Possibly worth noting as a deliberate divergence: a bare token is accepted here, whereas empty is returned by WS's
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Left as-is (as you noted, harmless): gRPC still accepts a bare token or Bearer prefix; WS still requires Bearer. |
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| md, ok := metadata.FromIncomingContext(ctx) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if !ok { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return "" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| vals := md.Get("authorization") | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if len(vals) == 0 { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return "" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| tok := vals[0] | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if after, ok := strings.CutPrefix(tok, "Bearer "); ok { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| tok = after | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return strings.TrimSpace(tok) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
stream_addr == stream_grpc_addrdoes not appear to be rejected anywhere. If both are set to the same value, each field passesvalidateStreamListenerindependently, and thenl.Fatalis called by whicheverRun()goroutine loses the bind race, so the failure would be nondeterministic and would read as unrelated to the config. An equality check next tovalidateStreamListener(pkg/config/config.go:295) could fail fast with a clearer message.There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Done — config validation now rejects stream_addr == stream_grpc_addr.