-
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 1 commit
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 | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| @@ -0,0 +1,147 @@ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| package stream | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| import ( | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "context" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "net" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "strings" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "google.golang.org/grpc" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "google.golang.org/grpc/codes" | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| "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 | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // 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: newConnLimiter(cfg.MaxConns, cfg.MaxConnsPerSub), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| grpcSrv: grpc.NewServer(), | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
coderabbitai[bot] marked this conversation as resolved.
Outdated
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| 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 | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| if err := stream.Send(&streamv1.BlockEvent{Lagged: true, Dropped: d}); 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 { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| pe := &streamv1.BlockEvent{ | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| 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 { | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| pe.Raw = ev.Raw | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| return pe | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|
|
||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| // 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) | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| } | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
| Original file line number | Diff line number | Diff line change | ||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| @@ -0,0 +1,139 @@ | ||||||||||||||||||||||||
| package stream | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| import ( | ||||||||||||||||||||||||
| "context" | ||||||||||||||||||||||||
| "net" | ||||||||||||||||||||||||
| "testing" | ||||||||||||||||||||||||
| "time" | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| "github.com/stretchr/testify/require" | ||||||||||||||||||||||||
| "google.golang.org/grpc" | ||||||||||||||||||||||||
| "google.golang.org/grpc/codes" | ||||||||||||||||||||||||
| "google.golang.org/grpc/credentials/insecure" | ||||||||||||||||||||||||
| "google.golang.org/grpc/metadata" | ||||||||||||||||||||||||
| "google.golang.org/grpc/status" | ||||||||||||||||||||||||
| "google.golang.org/grpc/test/bufconn" | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| "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/test_utils" | ||||||||||||||||||||||||
| ) | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| // newGRPCTestServer starts a GRPCServer over an in-memory bufconn and returns a | ||||||||||||||||||||||||
| // connected client. Auth is always required; loopback no-auth is covered by WS. | ||||||||||||||||||||||||
| func newGRPCTestServer(t *testing.T, cfg Config) (client streamv1.BlockStreamServiceClient, hub *streamhub.Service, rig *test_utils.AuthTestRig) { | ||||||||||||||||||||||||
| t.Helper() | ||||||||||||||||||||||||
| var authenticator ConsumerAuthenticator | ||||||||||||||||||||||||
| authenticator, rig = testAuth(t, true) | ||||||||||||||||||||||||
| hub = streamhub.New() | ||||||||||||||||||||||||
| g := NewGRPCServer(hub, authenticator, cfg, logger.NewAppSLogger(logger.Debug)) | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| lis := bufconn.Listen(1 << 20) | ||||||||||||||||||||||||
| go func() { _ = g.grpcSrv.Serve(lis) }() | ||||||||||||||||||||||||
| t.Cleanup(g.grpcSrv.Stop) | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| conn, err := grpc.NewClient("passthrough:///bufnet", | ||||||||||||||||||||||||
| grpc.WithContextDialer(func(ctx context.Context, _ string) (net.Conn, error) { return lis.DialContext(ctx) }), | ||||||||||||||||||||||||
| grpc.WithTransportCredentials(insecure.NewCredentials())) | ||||||||||||||||||||||||
| require.NoError(t, err) | ||||||||||||||||||||||||
| t.Cleanup(func() { _ = conn.Close() }) | ||||||||||||||||||||||||
| return streamv1.NewBlockStreamServiceClient(conn), hub, rig | ||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| func authCtx(t *testing.T, rig *test_utils.AuthTestRig, subject string) context.Context { | ||||||||||||||||||||||||
| t.Helper() | ||||||||||||||||||||||||
| return metadata.AppendToOutgoingContext(context.Background(), "authorization", "Bearer "+streamToken(t, rig, subject)) | ||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| func TestGRPC_RejectsWithoutToken(t *testing.T) { | ||||||||||||||||||||||||
| client, hub, _ := newGRPCTestServer(t, Config{}) | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| sub, err := client.Subscribe(context.Background(), &streamv1.SubscribeRequest{}) | ||||||||||||||||||||||||
| require.NoError(t, err) | ||||||||||||||||||||||||
| _, err = sub.Recv() | ||||||||||||||||||||||||
| require.Equal(t, codes.Unauthenticated, status.Code(err)) | ||||||||||||||||||||||||
| require.Zero(t, hub.SubscriberCount(), "rejected consumer must not create a subscriber") | ||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| func TestGRPC_DeliversFraming(t *testing.T) { | ||||||||||||||||||||||||
| for _, tc := range []struct { | ||||||||||||||||||||||||
| name string | ||||||||||||||||||||||||
| mode string | ||||||||||||||||||||||||
| expectRaw bool | ||||||||||||||||||||||||
| }{ | ||||||||||||||||||||||||
| {"metadata omits raw", "metadata", false}, | ||||||||||||||||||||||||
| {"raw includes bytes", "raw", true}, | ||||||||||||||||||||||||
| } { | ||||||||||||||||||||||||
| t.Run(tc.name, func(t *testing.T) { | ||||||||||||||||||||||||
| client, hub, rig := newGRPCTestServer(t, Config{}) | ||||||||||||||||||||||||
| sub, err := client.Subscribe(authCtx(t, rig, "sub-1"), &streamv1.SubscribeRequest{Mode: tc.mode}) | ||||||||||||||||||||||||
| require.NoError(t, err) | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| waitSubscribed(t, hub, 1) | ||||||||||||||||||||||||
| hub.Emit(sampleEvent()) | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| ev, err := sub.Recv() | ||||||||||||||||||||||||
| require.NoError(t, err) | ||||||||||||||||||||||||
| require.EqualValues(t, 42, ev.GetSlot()) | ||||||||||||||||||||||||
| require.False(t, ev.GetLagged()) | ||||||||||||||||||||||||
| if tc.expectRaw { | ||||||||||||||||||||||||
| require.Equal(t, []byte("ssz-snappy-bytes"), ev.GetRaw()) | ||||||||||||||||||||||||
| } else { | ||||||||||||||||||||||||
| require.Empty(t, ev.GetRaw()) | ||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
| }) | ||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| func TestGRPC_LaggedOnOverflow(t *testing.T) { | ||||||||||||||||||||||||
| client, hub, rig := newGRPCTestServer(t, Config{BufferSize: 1}) | ||||||||||||||||||||||||
| sub, err := client.Subscribe(authCtx(t, rig, "sub-1"), &streamv1.SubscribeRequest{}) | ||||||||||||||||||||||||
| require.NoError(t, err) | ||||||||||||||||||||||||
|
Comment on lines
+89
to
+95
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. Since
Suggested change
Applied locally: build,
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 — Subscribe now uses a 5s context so a blocked Recv fails instead of hanging. |
||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| waitSubscribed(t, hub, 1) | ||||||||||||||||||||||||
| for range 3000 { | ||||||||||||||||||||||||
| hub.Emit(sampleEvent()) | ||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| var sawLagged bool | ||||||||||||||||||||||||
| deadline := time.Now().Add(3 * time.Second) | ||||||||||||||||||||||||
| for time.Now().Before(deadline) && !sawLagged { | ||||||||||||||||||||||||
| ev, rerr := sub.Recv() | ||||||||||||||||||||||||
| require.NoError(t, rerr) | ||||||||||||||||||||||||
| if ev.GetLagged() { | ||||||||||||||||||||||||
| require.Positive(t, ev.GetDropped()) | ||||||||||||||||||||||||
| sawLagged = true | ||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
| require.True(t, sawLagged, "a lagged frame must be sent after overflow") | ||||||||||||||||||||||||
|
swarna1101 marked this conversation as resolved.
|
||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| func TestGRPC_GlobalCapRejects(t *testing.T) { | ||||||||||||||||||||||||
| client, hub, rig := newGRPCTestServer(t, Config{MaxConns: 1}) | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| first, err := client.Subscribe(authCtx(t, rig, "sub-a"), &streamv1.SubscribeRequest{}) | ||||||||||||||||||||||||
| require.NoError(t, err) | ||||||||||||||||||||||||
| waitSubscribed(t, hub, 1) | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| second, err := client.Subscribe(authCtx(t, rig, "sub-b"), &streamv1.SubscribeRequest{}) | ||||||||||||||||||||||||
| require.NoError(t, err) | ||||||||||||||||||||||||
| _, err = second.Recv() | ||||||||||||||||||||||||
| require.Equal(t, codes.ResourceExhausted, status.Code(err)) | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| _ = first // keep the first stream open for the duration of the assertion | ||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| func TestGRPC_CleanupOnCancel(t *testing.T) { | ||||||||||||||||||||||||
| client, hub, rig := newGRPCTestServer(t, Config{}) | ||||||||||||||||||||||||
| ctx, cancel := context.WithCancel(authCtx(t, rig, "sub-1")) | ||||||||||||||||||||||||
| _, err := client.Subscribe(ctx, &streamv1.SubscribeRequest{}) | ||||||||||||||||||||||||
| require.NoError(t, err) | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| waitSubscribed(t, hub, 1) | ||||||||||||||||||||||||
| cancel() | ||||||||||||||||||||||||
|
|
||||||||||||||||||||||||
| // Cancel must unwind the handler, closing the subscriber (drop-counter entry | ||||||||||||||||||||||||
| // included) and releasing the cap slot, so nothing leaks. | ||||||||||||||||||||||||
| waitSubscribed(t, hub, 0) | ||||||||||||||||||||||||
|
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. Release of the cap slot is described in the comment, but only require.Eventually(t, func() bool {
srv.limiter.mu.Lock()
defer srv.limiter.mu.Unlock()
return srv.limiter.conns == 0 && len(srv.limiter.perSub) == 0
}, 2*time.Second, 10*time.Millisecond)Not a suggestion, because A black-box alternative (
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 — newGRPCTestServer returns *GRPCServer and TestGRPC_CleanupOnCancel asserts the limiter is empty. |
||||||||||||||||||||||||
| } | ||||||||||||||||||||||||
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.