사용자별 credential 저장, lease, projection, runtime 전달과 OpenAI-compatible 계약 및 검증 근거를 함께 반영한다.
245 lines
8.4 KiB
Go
245 lines
8.4 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"crypto/tls"
|
|
"fmt"
|
|
"net"
|
|
"net/http"
|
|
"strings"
|
|
"time"
|
|
|
|
"go.uber.org/zap"
|
|
|
|
credentialleasesvc "iop/apps/control-plane/internal/credentiallease"
|
|
"iop/apps/control-plane/internal/credentialops"
|
|
"iop/apps/control-plane/internal/credentialseal"
|
|
"iop/apps/control-plane/internal/credentialstore"
|
|
"iop/apps/control-plane/internal/wire"
|
|
"iop/packages/go/auth"
|
|
"iop/packages/go/observability"
|
|
iop "iop/proto/gen/iop"
|
|
)
|
|
|
|
const principalProjectionTTL = 5 * time.Minute
|
|
|
|
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),
|
|
)
|
|
cred, err := composeCredentialRuntime(ctx, cfg, logger)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
store := cred.store
|
|
if store != nil {
|
|
defer store.Close()
|
|
logger.Info("control-plane credential store ready",
|
|
zap.String("dialect", dialectLabel(cfg.Database.URL)),
|
|
zap.Bool("secret_encryption", cred.service != nil),
|
|
)
|
|
}
|
|
if cfg.Redis.URL != "" {
|
|
logger.Info("control-plane redis configured", redisLogFields(cfg.Redis.URL, cfg.Redis.KeyPrefix)...)
|
|
}
|
|
if cfg.Metrics.Port > 0 {
|
|
go func() {
|
|
if err := observability.ServeMetrics(cfg.Metrics.Port); err != nil {
|
|
logger.Warn("control-plane metrics server exited", zap.Error(err))
|
|
}
|
|
}()
|
|
}
|
|
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() }()
|
|
|
|
var edgeTLS *tls.Config
|
|
if cfg.Server.EdgeWireTLS.Enabled {
|
|
edgeTLS, err = auth.LoadServerTLSWithIdentity(
|
|
cfg.Server.EdgeWireTLS.Cert,
|
|
cfg.Server.EdgeWireTLS.Key,
|
|
cfg.Server.EdgeWireTLS.CA,
|
|
cfg.Server.EdgeWireTLS.EffectivePeerRole("edge"),
|
|
cfg.Server.EdgeWireTLS.PeerName,
|
|
)
|
|
if err != nil {
|
|
return fmt.Errorf("load edge wire TLS: %w", err)
|
|
}
|
|
}
|
|
edgeServer, err := wire.NewEdgeServerTLS(cfg.Server.EdgeWireListen, edgeTLS, logger)
|
|
if err != nil {
|
|
return fmt.Errorf("edge wire server: %w", err)
|
|
}
|
|
if cfg.CredentialPlane.Enabled {
|
|
issuerKey, err := credentialleasesvc.LoadIssuerPrivateKey(cfg.CredentialPlane.IssuerPrivateKey)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
leaseService, err := credentialleasesvc.New(store, cred.keyring, cfg.CredentialPlane.IssuerKeyID, issuerKey, time.Duration(cfg.CredentialPlane.LeaseTTLSeconds)*time.Second, cfg.CredentialPlane.LeaseCacheSize, nil, nil)
|
|
if err != nil {
|
|
return fmt.Errorf("compose credential lease service: %w", err)
|
|
}
|
|
edgeServer.SetCredentialPlane(
|
|
func(callCtx context.Context) (*iop.PrincipalProjection, error) {
|
|
return store.BuildPrincipalProjection(callCtx, credentialstore.ProjectionBuildOptions{TTL: principalProjectionTTL})
|
|
},
|
|
leaseService.Acquire,
|
|
principalProjectionTTL/2,
|
|
)
|
|
}
|
|
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)
|
|
|
|
if cfg.CredentialPlane.Enabled {
|
|
registerCredentialHandlers(mux, cred.service, func(request *http.Request) error {
|
|
return edgeServer.BroadcastProjection(request.Context())
|
|
})
|
|
return startHTTPSServer(ctx, cfg.Server.Listen, mux, cfg.CredentialPlane.HTTPS.Cert, cfg.CredentialPlane.HTTPS.Key, logger)
|
|
}
|
|
return startHTTPServer(ctx, cfg.Server.Listen, mux, logger)
|
|
}
|
|
|
|
// credentialRuntime is the explicit composition of the credential store and
|
|
// the principal-scoped management service. The service is present only when
|
|
// at-rest encryption is configured and a store is open; it is not attached to
|
|
// any listener yet and is reserved for later secure handler composition.
|
|
type credentialRuntime struct {
|
|
store *credentialstore.Store
|
|
service *credentialops.Service
|
|
keyring *credentialseal.Keyring
|
|
}
|
|
|
|
// composeCredentialRuntime loads the external keyring, opens the credential
|
|
// store, and wires one keyring instance into both the store and the credential
|
|
// service. Failures are fail-closed: a partial or invalid encryption config, or
|
|
// any store error, returns before any network listener starts.
|
|
//
|
|
// When encryption is fully omitted the keyring is nil: the store keeps the
|
|
// legacy principal/metadata behavior and no sealer is injected, so every
|
|
// provider-secret mutation fails closed.
|
|
func composeCredentialRuntime(ctx context.Context, cfg controlPlaneConfig, logger *zap.Logger) (*credentialRuntime, error) {
|
|
keyring, err := credentialseal.LoadFile(cfg.CredentialEncryption)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("load credential encryption: %w", err)
|
|
}
|
|
|
|
// A complete encryption configuration without a database cannot host the
|
|
// credential store or service. Reject before any listener so an operator
|
|
// cannot start a credential-bearing process that silently has no store.
|
|
if keyring != nil && strings.TrimSpace(cfg.Database.URL) == "" {
|
|
return nil, fmt.Errorf("credential encryption requires database.url")
|
|
}
|
|
|
|
var opts []credentialstore.Option
|
|
if keyring != nil {
|
|
opts = append(opts, credentialstore.WithEnvelopeKeyRegistry(keyring))
|
|
}
|
|
store, err := credentialstore.Open(ctx, cfg.Database.URL, opts...)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("open credential store: %w", err)
|
|
}
|
|
|
|
rt := &credentialRuntime{store: store, keyring: keyring}
|
|
if store != nil && keyring != nil {
|
|
rt.service = credentialops.NewService(store, logger, keyring, nil)
|
|
}
|
|
return rt, nil
|
|
}
|
|
|
|
// dialectLabel returns a short label for the credential store dialect used
|
|
// by a database URL. It is used only for logging and never exposes credentials.
|
|
func dialectLabel(databaseURL string) string {
|
|
if databaseURL == "" {
|
|
return "unconfigured"
|
|
}
|
|
lower := strings.ToLower(databaseURL)
|
|
switch {
|
|
case strings.HasPrefix(lower, "postgres://") || strings.HasPrefix(lower, "postgresql://"):
|
|
return "postgres"
|
|
case strings.HasPrefix(lower, "file:") || strings.Contains(lower, ".db") || !strings.ContainsAny(lower, "://"):
|
|
return "sqlite"
|
|
default:
|
|
return "unknown"
|
|
}
|
|
}
|
|
|
|
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 {
|
|
return startHTTPServerWithTLS(ctx, listenAddr, handler, nil, logger)
|
|
}
|
|
|
|
func startHTTPSServer(ctx context.Context, listenAddr string, handler http.Handler, certFile, keyFile string, logger *zap.Logger) error {
|
|
tlsConfig, err := auth.LoadHTTPServerTLS(certFile, keyFile)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return startHTTPServerWithTLS(ctx, listenAddr, handler, tlsConfig, logger)
|
|
}
|
|
|
|
func startHTTPServerWithTLS(ctx context.Context, listenAddr string, handler http.Handler, tlsConfig *tls.Config, logger *zap.Logger) error {
|
|
server := &http.Server{
|
|
Addr: listenAddr,
|
|
Handler: handler,
|
|
ReadHeaderTimeout: 5 * time.Second,
|
|
TLSConfig: tlsConfig,
|
|
}
|
|
ln, err := net.Listen("tcp", listenAddr)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if tlsConfig != nil {
|
|
ln = tls.NewListener(ln, tlsConfig)
|
|
}
|
|
|
|
errCh := make(chan error, 1)
|
|
go func() {
|
|
logger.Info("control-plane http endpoint listening", zap.String("listen", listenAddr))
|
|
if err := server.Serve(ln); 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
|
|
}
|
|
}
|