Skip to content
Open
Show file tree
Hide file tree
Changes from 12 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
Loading
Loading