Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
Prev Previous commit
Next Next commit
feat: Add serverTSL() method.
Signed-off-by: Shuchu Han <shuchu.han@gmail.com>
  • Loading branch information
shuchu authored and SHUCHU HAN committed Apr 14, 2026
commit 064f6128e79d034737eb41d125ad38ece5e45097
14 changes: 14 additions & 0 deletions go/internal/feast/server/http_server.go
Original file line number Diff line number Diff line change
Expand Up @@ -396,6 +396,20 @@ func (s *httpServer) Serve(host string, port int) error {
return err
}

func (s *httpServer) ServeTLS(host string, port int, certFile string, keyFile string) error {
mux := http.NewServeMux()
mux.Handle("/get-online-features", metricsMiddleware(recoverMiddleware(http.HandlerFunc(s.getOnlineFeatures))))
mux.Handle("/health", metricsMiddleware(http.HandlerFunc(healthCheckHandler)))
s.server = &http.Server{Addr: fmt.Sprintf("%s:%d", host, port), Handler: mux, ReadTimeout: 5 * time.Second, WriteTimeout: 10 * time.Second, IdleTimeout: 15 * time.Second}
err := s.server.ListenAndServeTLS(certFile, keyFile)
// Don't return the error if it's caused by graceful shutdown using Stop()
if err == http.ErrServerClosed {
return nil
}
log.Fatal().Stack().Err(err).Msg("Failed to start HTTPS server")
return err
}

func healthCheckHandler(w http.ResponseWriter, r *http.Request) {
w.WriteHeader(http.StatusOK)
fmt.Fprintf(w, "Healthy")
Expand Down
161 changes: 75 additions & 86 deletions go/main.go
Comment thread
shuchu marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (
"strings"
"sync"
"syscall"
"time"

"github.com/feast-dev/feast/go/internal/feast"
"github.com/feast-dev/feast/go/internal/feast/registry"
Expand Down Expand Up @@ -70,14 +71,14 @@ func main() {
log.Error().Stack().Err(err).Msg("Failed to get current directory")
}

flag.StringVar(&serverType, "type", serverType, "Specify the server type (http or grpc)")
flag.StringVar(&serverType, "type", serverType, "Specify the server type (http, https or grpc)")
flag.StringVar(&repoPath, "chdir", repoPath, "Repository path where feature store yaml file is stored")

flag.StringVar(&host, "host", host, "Specify a host for the server")
flag.IntVar(&port, "port", port, "Specify a port for the server")
flag.IntVar(&metricsPort, "metrics-port", metricsPort, "Specify a port for the metrics server")
flag.StringVar(&certFile, "tls-cert-file", "", "Path to the TLS certificate file")
flag.StringVar(&keyFile, "tls-key-file", "", "Path to the TLS key file" )
flag.StringVar(&keyFile, "tls-key-file", "", "Path to the TLS key file")
flag.Parse()

// Initialize tracer
Expand Down Expand Up @@ -244,20 +245,24 @@ func StartHttpServer(fs *feast.FeatureStore, host string, port int, metricsPort
}
ser := server.NewHttpServer(fs, loggingService)
log.Info().Msgf("Starting a HTTP server on host %s, port %d", host, port)
Comment thread
shuchu marked this conversation as resolved.

// Start metrics server
metricsServer := &http.Server{Addr: fmt.Sprintf(":%d", metricsPort)}
mux := http.NewServeMux()
mux.Handle("/metrics", promhttp.Handler())
metricsServer := &http.Server{
Addr: fmt.Sprintf(":%d", metricsPort),
Handler: mux,
}
go func() {
log.Info().Msgf("Starting metrics server on port %d", metricsPort)
mux := http.NewServeMux()
mux.Handle("/metrics", promhttp.Handler())
metricsServer.Handler = mux
if err := metricsServer.ListenAndServe(); err != nil && err != http.ErrServerClosed {
log.Error().Err(err).Msg("Failed to start metrics server")
}
}()

stop := make(chan os.Signal, 1)
signal.Notify(stop, syscall.SIGINT, syscall.SIGTERM)
defer signal.Stop(stop)

var wg sync.WaitGroup
wg.Add(1)
Expand All @@ -273,7 +278,9 @@ func StartHttpServer(fs *feast.FeatureStore, host string, port int, metricsPort
log.Error().Err(err).Msg("Error when stopping the HTTP server")
}
log.Info().Msg("Stopping metrics server...")
if err := metricsServer.Shutdown(context.Background()); err != nil {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := metricsServer.Shutdown(ctx); err != nil {
log.Error().Err(err).Msg("Error stopping metrics server")
}
if loggingService != nil {
Expand All @@ -295,6 +302,44 @@ func StartHttpServer(fs *feast.FeatureStore, host string, port int, metricsPort
return err
}

func OTELTracingEnabled() bool {
return strings.ToLower(os.Getenv("ENABLE_OTEL_TRACING")) == "true"
}

func newExporter(ctx context.Context) (*otlptrace.Exporter, error) {
exp, err := otlptracehttp.New(ctx,
otlptracehttp.WithInsecure())
if err != nil {
return nil, err
}
return exp, nil
}

func newTracerProvider(exp sdktrace.SpanExporter) (*sdktrace.TracerProvider, error) {
serviceName := os.Getenv("OTEL_SERVICE_NAME")
if serviceName == "" {
serviceName = "FeastGoFeatureServer"
}
r, err := resource.Merge(
resource.Default(),
resource.NewWithAttributes(
semconv.SchemaURL,
semconv.ServiceName(serviceName),
),
)

if err != nil {
return nil, err
}

return sdktrace.NewTracerProvider(
sdktrace.WithBatcher(exp),
sdktrace.WithResource(r),
), nil
}



// StartHttpsServer starts HTTP server with TLS. Requires TLS_CERT_FILE and TLS_KEY_FILE env vars.
func StartHttpsServer(fs *feast.FeatureStore, host string, port int, metricsPort int, certFile string, keyFile string, writeLoggedFeaturesCallback logging.OfflineStoreWriteCallback, loggingOpts *logging.LoggingOptions) error {
if certFile == "" || keyFile == "" {
Expand All @@ -305,24 +350,26 @@ func StartHttpsServer(fs *feast.FeatureStore, host string, port int, metricsPort
if err != nil {
return err
}

// Try to obtain an http.Handler from the concrete server if possible.
ser := server.NewHttpServer(fs, loggingService)
log.Info().Msgf("Starting a HTTPS server on host %s, port %d", host, port)

// Start metrics server (same as HTTP)
metricsServer := &http.Server{Addr: fmt.Sprintf(":%d", metricsPort)}
// Start metrics server
mux := http.NewServeMux()
mux.Handle("/metrics", promhttp.Handler())
metricsServer := &http.Server{
Addr: fmt.Sprintf(":%d", metricsPort),
Handler: mux,
}
go func() {
log.Info().Msgf("Starting metrics server on port %d", metricsPort)
mux := http.NewServeMux()
mux.Handle("/metrics", promhttp.Handler())
metricsServer.Handler = mux
if err := metricsServer.ListenAndServe(); err != nil && err != http.ErrServerClosed {
log.Error().Err(err).Msg("Failed to start metrics server")
}
}()

stop := make(chan os.Signal, 1)
signal.Notify(stop, syscall.SIGINT, syscall.SIGTERM)
defer signal.Stop(stop)

var wg sync.WaitGroup
wg.Add(1)
Expand All @@ -331,91 +378,33 @@ func StartHttpsServer(fs *feast.FeatureStore, host string, port int, metricsPort
defer wg.Done()
select {
case <-stop:
log.Info().Msg("Stopping the HTTPS server...")
// Try to stop underlying server if it exposes Stop()
if stopper, ok := interface{}(ser).(interface{ Stop() error }); ok {
if err := stopper.Stop(); err != nil {
log.Error().Err(err).Msg("Error when stopping the HTTPS server")
}
// Received SIGINT/SIGTERM. Perform graceful shutdown.
log.Info().Msg("Stopping the HTTP server...")
err := ser.Stop()
if err != nil {
log.Error().Err(err).Msg("Error when stopping the HTTP server")
}
if err := metricsServer.Shutdown(context.Background()); err != nil {
log.Info().Msg("Stopping metrics server...")
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := metricsServer.Shutdown(ctx); err != nil {
log.Error().Err(err).Msg("Error stopping metrics server")
}
if loggingService != nil {
loggingService.Stop()
}
log.Info().Msg("HTTPS server terminated")
log.Info().Msg("HTTP server terminated")
case <-serverExited:
// Server exited (e.g. startup error), ensure metrics server is stopped
metricsServer.Shutdown(context.Background())
if loggingService != nil {
loggingService.Stop()
}
}
}()

// If the concrete server exposes a Handler, use it with a tls-enabled http.Server.
if hProvider, ok := interface{}(ser).(interface{ Handler() http.Handler }); ok {
handler := hProvider.Handler()
srv := &http.Server{Addr: fmt.Sprintf("%s:%d", host, port), Handler: handler}
go func() {
log.Info().Msgf("Starting HTTPS server on host %s, port %d", host, port)
if err := srv.ListenAndServeTLS(certFile, keyFile); err != nil && err != http.ErrServerClosed {
log.Error().Err(err).Msg("HTTPS server failed")
}
}()
close(serverExited)
wg.Wait()
return nil
}

// If concrete server supports ServeTLS(host,port,cert,key), call it.
if tlsServ, ok := interface{}(ser).(interface {
ServeTLS(string, int, string, string) error
}); ok {
err := tlsServ.ServeTLS(host, port, certFile, keyFile)
close(serverExited)
wg.Wait()
return err
}

// Fallback: cannot enable TLS for this server implementation.
err = ser.ServeTLS(host, port, certFile, keyFile)
close(serverExited)
wg.Wait()
return fmt.Errorf("HTTPS not supported by underlying HTTP server implementation")
}

func OTELTracingEnabled() bool {
return strings.ToLower(os.Getenv("ENABLE_OTEL_TRACING")) == "true"
}

func newExporter(ctx context.Context) (*otlptrace.Exporter, error) {
exp, err := otlptracehttp.New(ctx,
otlptracehttp.WithInsecure())
if err != nil {
return nil, err
}
return exp, nil
}

func newTracerProvider(exp sdktrace.SpanExporter) (*sdktrace.TracerProvider, error) {
serviceName := os.Getenv("OTEL_SERVICE_NAME")
if serviceName == "" {
serviceName = "FeastGoFeatureServer"
}
r, err := resource.Merge(
resource.Default(),
resource.NewWithAttributes(
semconv.SchemaURL,
semconv.ServiceName(serviceName),
),
)

if err != nil {
return nil, err
}

return sdktrace.NewTracerProvider(
sdktrace.WithBatcher(exp),
sdktrace.WithResource(r),
), nil
}
return err
}