server.go 3.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110
  1. package grpcserver
  2. import (
  3. "context"
  4. "fmt"
  5. "net"
  6. "time"
  7. pbv1 "git3.techno-world.net/lrosales/broad-announce/gen/go/broadannounce/v1"
  8. "git3.techno-world.net/lrosales/broad-announce/internal/observability"
  9. "git3.techno-world.net/lrosales/broad-announce/internal/pipeline"
  10. "google.golang.org/grpc"
  11. "google.golang.org/grpc/keepalive"
  12. )
  13. // Config for the gRPC server.
  14. type Config struct {
  15. Addr string // default ":9090"
  16. MaxInflight int // max in-flight per stream; default 256
  17. }
  18. // IngestdDeps is what ingestd's main.go passes to NewForIngestd.
  19. // Using an interface avoids an import cycle (internal/grpcserver cannot import cmd/ingestd).
  20. type IngestdDeps interface {
  21. GetLogger() interface{ Info(msg string, args ...any) }
  22. GetMetrics() *observability.IngestdMetrics
  23. GetPipeline() pipeline.Deps
  24. }
  25. // NewForIngestd is the constructor ingestd's main.go calls.
  26. // It registers the Ingest service on a grpc.Server but does not start listening.
  27. func NewForIngestd(cfg Config, deps IngestdDeps) (*grpc.Server, error) {
  28. if cfg.Addr == "" {
  29. cfg.Addr = ":9090"
  30. }
  31. if cfg.MaxInflight == 0 {
  32. cfg.MaxInflight = 256
  33. }
  34. logger := deps.GetLogger()
  35. m := deps.GetMetrics()
  36. // Prime the histogram buckets with a zero observation so Prometheus
  37. // registers them immediately (avoids missing bucket labels on first scrape).
  38. m.GRPCAckLatency.WithLabelValues("_init").Observe(0)
  39. ka := keepalive.ServerParameters{
  40. Time: 30 * time.Second,
  41. Timeout: 10 * time.Second,
  42. }
  43. enforcement := keepalive.EnforcementPolicy{
  44. MinTime: 5 * time.Second,
  45. PermitWithoutStream: false,
  46. }
  47. grpcOpts := []grpc.ServerOption{
  48. grpc.KeepaliveParams(ka),
  49. grpc.KeepaliveEnforcementPolicy(enforcement),
  50. grpc.WriteBufferSize(256 * 1024),
  51. grpc.ReadBufferSize(32 * 1024),
  52. // TODO(m11): grpc.Creds(tlsCredentials()) — enable TLS once certificates are available
  53. }
  54. srv := grpc.NewServer(grpcOpts...)
  55. pbv1.RegisterIngestServer(srv, &Server{
  56. Dep: deps.GetPipeline(),
  57. MaxInflight: cfg.MaxInflight,
  58. Logger: logger,
  59. })
  60. logger.Info("grpc server configured", "addr", cfg.Addr, "max_inflight", cfg.MaxInflight)
  61. return srv, nil
  62. }
  63. // Serve starts a TCP listener on addr and blocks serving.
  64. // On ctx cancellation, it initiates a graceful shutdown (drains for up to 15s
  65. // then hard-stops). The caller should run Serve in a goroutine.
  66. func Serve(ctx context.Context, cfg Config, deps IngestdDeps) error {
  67. srv, err := NewForIngestd(cfg, deps)
  68. if err != nil {
  69. return fmt.Errorf("grpcserver.NewForIngestd: %w", err)
  70. }
  71. ln, err := net.Listen("tcp", cfg.Addr)
  72. if err != nil {
  73. return fmt.Errorf("grpc listen %s: %w", cfg.Addr, err)
  74. }
  75. errCh := make(chan error, 1)
  76. go func() { errCh <- srv.Serve(ln) }()
  77. select {
  78. case <-ctx.Done():
  79. done := make(chan struct{})
  80. go func() {
  81. srv.GracefulStop()
  82. close(done)
  83. }()
  84. select {
  85. case <-done:
  86. return nil
  87. case <-time.After(15 * time.Second):
  88. srv.Stop()
  89. return ctx.Err()
  90. }
  91. case err := <-errCh:
  92. return err
  93. }
  94. }