Control Plane HTTP API가 fleet polling과 client 상태 확장 전에 더 커지지 않도록 server lifecycle, DTO conversion, Edge/fleet handler 책임을 파일 단위로 나눈다.
97 lines
3 KiB
Go
97 lines
3 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"net/http"
|
|
"time"
|
|
|
|
"go.uber.org/zap"
|
|
|
|
"iop/apps/control-plane/internal/wire"
|
|
)
|
|
|
|
func run(ctx context.Context, cfg controlPlaneConfig, logger *zap.Logger) error {
|
|
mux := newHTTPMux()
|
|
|
|
wireEndpoint := wire.Endpoint{Listen: cfg.Server.WireListen}
|
|
logger.Info("control-plane client wire endpoint reserved",
|
|
zap.String("protocol", wire.Protocol),
|
|
zap.String("transport", wire.ClientTransport),
|
|
zap.String("listen", wireEndpoint.Listen),
|
|
)
|
|
logger.Info("control-plane edge wire endpoint reserved",
|
|
zap.String("protocol", wire.Protocol),
|
|
zap.String("transport", wire.EdgeTransport),
|
|
zap.String("listen", cfg.Server.EdgeWireListen),
|
|
)
|
|
if cfg.Database.URL != "" {
|
|
logger.Info("control-plane database configured", databaseLogFields(cfg.Database.URL)...)
|
|
}
|
|
if cfg.Redis.URL != "" {
|
|
logger.Info("control-plane redis configured", redisLogFields(cfg.Redis.URL, cfg.Redis.KeyPrefix)...)
|
|
}
|
|
clientServer, err := wire.NewClientServer(cfg.Server.WireListen, logger)
|
|
if err != nil {
|
|
return fmt.Errorf("wire server: %w", err)
|
|
}
|
|
if err := clientServer.Start(ctx); err != nil {
|
|
return fmt.Errorf("start wire server: %w", err)
|
|
}
|
|
defer func() { _ = clientServer.Stop() }()
|
|
|
|
edgeServer, err := wire.NewEdgeServer(cfg.Server.EdgeWireListen, logger)
|
|
if err != nil {
|
|
return fmt.Errorf("edge wire server: %w", err)
|
|
}
|
|
if err := edgeServer.Start(ctx); err != nil {
|
|
return fmt.Errorf("start edge wire server: %w", err)
|
|
}
|
|
defer func() { _ = edgeServer.Stop() }()
|
|
registerEdgeRegistryHandlers(mux, edgeServer.Registry(), edgeServer.RequestStatus, edgeServer.SendCommand)
|
|
registerFleetHandlers(mux, edgeServer.Registry(), edgeServer.RequestStatus, edgeServer.SendCommand)
|
|
|
|
return startHTTPServer(ctx, cfg.Server.Listen, mux, logger)
|
|
}
|
|
|
|
func newHTTPMux() *http.ServeMux {
|
|
mux := http.NewServeMux()
|
|
mux.HandleFunc("/healthz", func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusOK)
|
|
_, _ = w.Write([]byte("ok\n"))
|
|
})
|
|
mux.HandleFunc("/readyz", func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusOK)
|
|
_, _ = w.Write([]byte("ready\n"))
|
|
})
|
|
return mux
|
|
}
|
|
|
|
func startHTTPServer(ctx context.Context, listenAddr string, handler http.Handler, logger *zap.Logger) error {
|
|
server := &http.Server{
|
|
Addr: listenAddr,
|
|
Handler: handler,
|
|
ReadHeaderTimeout: 5 * time.Second,
|
|
}
|
|
|
|
errCh := make(chan error, 1)
|
|
go func() {
|
|
logger.Info("control-plane http endpoint listening", zap.String("listen", listenAddr))
|
|
if err := server.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
|
errCh <- err
|
|
return
|
|
}
|
|
errCh <- nil
|
|
}()
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
// Use the original context for shutdown (signal.NotifyContext may close it).
|
|
// Re-create a shutdown context with the same deadline or a small timeout.
|
|
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
return server.Shutdown(shutdownCtx)
|
|
case err := <-errCh:
|
|
return err
|
|
}
|
|
}
|