From 1abad432968fdbc8daf0cad6acde947c62541206 Mon Sep 17 00:00:00 2001 From: Dmitrii Andreev Date: Wed, 2 Sep 2026 13:09:53 -0500 Subject: [PATCH] HYPERFLEET-1437 - feat: Build the named transport registry and store construction --- cmd/adapter/main.go | 123 ++--- configs/adapter-config-template.yaml | 22 + go.mod | 4 + go.sum | 16 + internal/configloader/constants.go | 15 + internal/configloader/registry_config_test.go | 196 ++++++++ internal/configloader/types.go | 71 ++- internal/configloader/validator.go | 57 +++ internal/transportclient/registry.go | 17 + internal/transportclient/registry_test.go | 103 ++++ internal/transportregistry/registry.go | 263 ++++++++++ internal/transportregistry/runtime_test.go | 451 ++++++++++++++++++ 12 files changed, 1241 insertions(+), 97 deletions(-) create mode 100644 internal/configloader/registry_config_test.go create mode 100644 internal/transportclient/registry.go create mode 100644 internal/transportclient/registry_test.go create mode 100644 internal/transportregistry/registry.go create mode 100644 internal/transportregistry/runtime_test.go diff --git a/cmd/adapter/main.go b/cmd/adapter/main.go index c35319bf..63929f48 100644 --- a/cmd/adapter/main.go +++ b/cmd/adapter/main.go @@ -16,10 +16,9 @@ import ( "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/dryrun" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/executor" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/hyperfleetapi" - "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/k8sclient" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/logctx" - "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/maestroclient" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportregistry" "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/health" "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/metrics" "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/telemetry" @@ -335,85 +334,6 @@ func createAPIClient(apiConfig configloader.HyperfleetAPIConfig) (hyperfleetapi. return hyperfleetapi.NewClient(opts...) } -// createTransportClient creates the appropriate transport client based on config. -func createTransportClient( - ctx context.Context, - config *configloader.Config, -) (transportclient.TransportClient, error) { - if config.Clients.Maestro != nil { - slog.InfoContext(ctx, "creating maestro transport client...") - client, err := createMaestroClient(ctx, config.Clients.Maestro) - if err != nil { - return nil, err - } - slog.InfoContext(ctx, "maestro transport client created successfully") - return client, nil - } - - slog.InfoContext(ctx, "creating kubernetes transport client...") - client, err := createK8sClient(ctx, config.Clients.Kubernetes) - if err != nil { - return nil, err - } - slog.InfoContext(ctx, "kubernetes transport client created successfully") - return client, nil -} - -// createK8sClient creates a Kubernetes client from the config -func createK8sClient( - ctx context.Context, - k8sConfig configloader.KubernetesConfig, -) (*k8sclient.Client, error) { - clientConfig := k8sclient.ClientConfig{ - KubeConfigPath: k8sConfig.KubeConfigPath, - QPS: k8sConfig.QPS, - Burst: k8sConfig.Burst, - } - return k8sclient.NewClient(ctx, clientConfig) -} - -// createMaestroClient creates a Maestro client from the config -func createMaestroClient( - ctx context.Context, - maestroConfig *configloader.MaestroClientConfig, -) (*maestroclient.Client, error) { - config := &maestroclient.Config{ - MaestroServerAddr: maestroConfig.HTTPServerAddress, - GRPCServerAddr: maestroConfig.GRPCServerAddress, - SourceID: maestroConfig.SourceID, - Insecure: maestroConfig.Insecure, - } - - if maestroConfig.Timeout != "" { - d, err := time.ParseDuration(maestroConfig.Timeout) - if err != nil { - return nil, fmt.Errorf("invalid maestro timeout %q: %w", maestroConfig.Timeout, err) - } - config.HTTPTimeout = d - } - - if maestroConfig.ServerHealthinessTimeout != "" { - d, err := time.ParseDuration(maestroConfig.ServerHealthinessTimeout) - if err != nil { - return nil, fmt.Errorf( - "invalid maestro serverHealthinessTimeout %q: %w", - maestroConfig.ServerHealthinessTimeout, - err, - ) - } - config.ServerHealthinessTimeout = d - } - - if maestroConfig.Auth.TLSConfig != nil { - config.CAFile = maestroConfig.Auth.TLSConfig.CAFile - config.ClientCertFile = maestroConfig.Auth.TLSConfig.CertFile - config.ClientKeyFile = maestroConfig.Auth.TLSConfig.KeyFile - config.HTTPCAFile = maestroConfig.Auth.TLSConfig.HTTPCAFile - } - - return maestroclient.NewMaestroClient(ctx, config) -} - // buildExecutor creates the executor with the given clients. func buildExecutor( config *configloader.Config, @@ -563,10 +483,24 @@ func runServe(flags *pflag.FlagSet) error { return fmt.Errorf("failed to create HyperFleet API client: %w", err) } - tc, err := createTransportClient(ctx, config) + transportRuntime, err := transportregistry.Build(ctx, config) + if err != nil { + slog.ErrorContext(ctx, "failed to create transport registry", "error", err) + return fmt.Errorf("failed to create transport registry: %w", err) + } + defer func() { + if closeErr := transportRuntime.Close(); closeErr != nil { + slog.WarnContext(ctx, "failed to close transport registry", "error", closeErr) + } + }() + + compatibilityKey := configloader.TransportClientKubernetes + if config.Clients.Maestro != nil { + compatibilityKey = configloader.TransportClientMaestro + } + tc, err := transportRuntime.Registry.Get(compatibilityKey) if err != nil { - slog.ErrorContext(ctx, "failed to create transport client", "error", err) - return fmt.Errorf("failed to create transport client: %w", err) + return fmt.Errorf("failed to resolve transport client: %w", err) } // Build executor @@ -745,8 +679,27 @@ func runDryRun(flags *pflag.FlagSet) error { dryrunClient = dryrun.NewDryrunTransportClient() } + transportRuntime, err := transportregistry.BuildRecording(config, dryrunClient) + if err != nil { + return fmt.Errorf("failed to create recording transport registry: %w", err) + } + defer func() { + if closeErr := transportRuntime.Close(); closeErr != nil { + slog.WarnContext(ctx, "failed to close recording transport registry", "error", closeErr) + } + }() + + compatibilityKey := configloader.TransportClientKubernetes + if config.Clients.Maestro != nil { + compatibilityKey = configloader.TransportClientMaestro + } + tc, err := transportRuntime.Registry.Get(compatibilityKey) + if err != nil { + return fmt.Errorf("failed to resolve recording transport client: %w", err) + } + // Build executor with mock clients (same builder as serve, no metrics in dry-run) - exec, err := buildExecutor(config, dryrunAPI, dryrunClient, nil) + exec, err := buildExecutor(config, dryrunAPI, tc, nil) if err != nil { return fmt.Errorf("failed to create executor: %w", err) } diff --git a/configs/adapter-config-template.yaml b/configs/adapter-config-template.yaml index b84f5840..edb4811a 100644 --- a/configs/adapter-config-template.yaml +++ b/configs/adapter-config-template.yaml @@ -32,6 +32,28 @@ adapter: # Flag: --debug-config debug_config: false +# Named deployment transports and stores. These are infrastructure settings, +# loaded with the deployment config (including Viper overrides); they do not +# belong in the task configuration. +# +# A Kubernetes entry needs no store. A remote entry writes desires through a +# named memory or Redis store. Names are stable registry keys, so multiple +# entries can use the same transport or store type. +transports: + kubernetes: + type: kubernetes + remote: + type: remote + store: adapter-desires + +stores: + adapter-desires: + type: memory + # For Redis, replace the memory store above (or add another named store): + # redis-desires: + # type: redis + # url: "rediss://username:password@redis.example.com:6379/0" + # Logging configuration # Priority: CLI flag > LOG_LEVEL/LOG_FORMAT/LOG_OUTPUT env vars > this file > defaults log: diff --git a/go.mod b/go.mod index b04df972..938a76e5 100644 --- a/go.mod +++ b/go.mod @@ -5,6 +5,7 @@ go 1.26.0 require ( cel.dev/cel-go v0.32.0 github.com/Masterminds/semver/v3 v3.5.0 + github.com/alicebob/miniredis/v2 v2.39.0 github.com/cloudevents/sdk-go/v2 v2.16.2 github.com/go-playground/validator/v10 v10.30.3 github.com/go-viper/mapstructure/v2 v2.5.0 @@ -16,6 +17,7 @@ require ( github.com/openshift-online/ocm-sdk-go v0.1.510 github.com/prometheus/client_golang v1.24.1 github.com/prometheus/client_model v0.6.2 + github.com/redis/go-redis/v9 v9.22.0 github.com/spf13/cobra v1.10.2 github.com/spf13/pflag v1.0.10 github.com/spf13/viper v1.21.0 @@ -148,6 +150,7 @@ require ( github.com/tklauser/go-sysconf v0.4.0 // indirect github.com/tklauser/numcpus v0.12.0 // indirect github.com/x448/float16 v0.8.4 // indirect + github.com/yuin/gopher-lua v1.1.1 // indirect github.com/yusufpapurcu/wmi v1.2.4 // indirect go.opencensus.io v0.24.0 // indirect go.opentelemetry.io/auto/sdk v1.2.1 // indirect @@ -160,6 +163,7 @@ require ( go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.46.0 // indirect go.opentelemetry.io/otel/metric v1.46.0 // indirect go.opentelemetry.io/proto/otlp v1.11.0 // indirect + go.uber.org/atomic v1.11.0 // indirect go.uber.org/multierr v1.11.0 // indirect go.uber.org/zap v1.28.0 // indirect go.yaml.in/yaml/v2 v2.4.4 // indirect diff --git a/go.sum b/go.sum index 3f01c8b7..b897e559 100644 --- a/go.sum +++ b/go.sum @@ -32,10 +32,16 @@ github.com/ThreeDotsLabs/watermill-amqp/v3 v3.1.0 h1:2EhCSlRZyZZUpLMh7PvhaTKJus0 github.com/ThreeDotsLabs/watermill-amqp/v3 v3.1.0/go.mod h1:eYO5aoQNezSBHuuiW69vj8iyH90bJSld1+zUvhhfJsQ= github.com/ThreeDotsLabs/watermill-googlecloud/v2 v2.0.1 h1:UF8mC04XepJ5uP0EN6YOUFjfQRKEJvgPO/lVOe8ftms= github.com/ThreeDotsLabs/watermill-googlecloud/v2 v2.0.1/go.mod h1:dIL2o+0uhh3GUgk3yfayPorjO83oSvO4J9fKsxASnVw= +github.com/alicebob/miniredis/v2 v2.39.0 h1:M7WbmV5BmV56L8KTG0rw6vEQ+woTOghpDgin2xv4A0g= +github.com/alicebob/miniredis/v2 v2.39.0/go.mod h1:TcL7YfarKPGDAthEtl5NBeHZfeUQj6OXMm/+iu5cLMM= github.com/antlr4-go/antlr/v4 v4.13.1 h1:SqQKkuVZ+zWkMMNkjy5FZe5mr5WURWnlpmOuzYWrPrQ= github.com/antlr4-go/antlr/v4 v4.13.1/go.mod h1:GKmUxMtwp6ZgGwZSva4eWPC5mS6vUAmOABFgjdkM7Nw= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= +github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= +github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= +github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= +github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= github.com/bwmarrin/snowflake v0.3.0 h1:xm67bEhkKh6ij1790JB83OujPR5CzNe8QuQqAgISZN0= github.com/bwmarrin/snowflake v0.3.0/go.mod h1:NdZxfVWX+oR6y2K0o6qAYv6gIOP9rjG0/E9WsDpxqwE= github.com/cenkalti/backoff/v3 v3.2.2 h1:cfUAAO3yvKMYKPrvhDuHSwQnhZNk/RMHKdZqKTxfm6M= @@ -218,6 +224,8 @@ github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnr github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= github.com/klauspost/compress v1.19.1 h1:VsB4HPswih7mmZ8WleSFQ75c/Ui1M4trX5oAsJnhSlk= github.com/klauspost/compress v1.19.1/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/klauspost/cpuid/v2 v2.2.10 h1:tBs3QSyvjDyFTq3uoc/9xFpCuOsJQFNPiAhYdw2skhE= +github.com/klauspost/cpuid/v2 v2.2.10/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= @@ -307,6 +315,8 @@ github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+ github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY= github.com/rabbitmq/amqp091-go v1.12.0 h1:V0v14Iqfs+MwHWihJt/nGS5Ulu0vw572b2Co3mwunkI= github.com/rabbitmq/amqp091-go v1.12.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= +github.com/redis/go-redis/v9 v9.22.0 h1:laDvpYXTJtZLloinw1fA5Kqd6HAEH2XKxOkG/PDq2F0= +github.com/redis/go-redis/v9 v9.22.0/go.mod h1:y2g0Wj8rQvuK0ELM+oxSudcLtC09JScs98I/X9gRWY4= github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= @@ -354,8 +364,12 @@ github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6Kllzaw github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyCJ6HpOuEn7z0Csc= github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= +github.com/yuin/gopher-lua v1.1.1 h1:kYKnWBjvbNP4XLT3+bPEwAXJx262OhaHDWDVOPjL46M= +github.com/yuin/gopher-lua v1.1.1/go.mod h1:GBR0iDaNXjAgGg9zfCvksxSRnQx76gclCIb7kdAd1Pw= github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0= github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0= +github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= +github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s= go.einride.tech/aip v0.83.0 h1:TI21IdeOnLTwZEJ3BxtImIZk6bsN2Q+sd0x99SLiQ+M= go.einride.tech/aip v0.83.0/go.mod h1:E8+wdTApA70odnpFzJgsGogHozC2JCIhFJBKPr8bVig= go.opencensus.io v0.24.0 h1:y73uSU6J157QMP2kn2r30vwW1A2W2WFwSCGnAVxeaD0= @@ -396,6 +410,8 @@ go.opentelemetry.io/otel/trace v1.46.0 h1:OULy7ccdJnZtJ0UDYFOIGaCmiWzJ8Vi2G/Rsu6 go.opentelemetry.io/otel/trace v1.46.0/go.mod h1:J7GAXweO77XSFkB/rmAqk9D6ihszhFjLU+d9WuUxDLI= go.opentelemetry.io/proto/otlp v1.11.0 h1:5rrYs0Ykyj50sdU/JU0x8etU+LubXWb+gED6TbEdMIk= go.opentelemetry.io/proto/otlp v1.11.0/go.mod h1:SmVizdCOAm3XBtG1g1NnOdhW6jtddT72hLMhv8VwA8E= +go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= +go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= diff --git a/internal/configloader/constants.go b/internal/configloader/constants.go index d41f9651..b234ead1 100644 --- a/internal/configloader/constants.go +++ b/internal/configloader/constants.go @@ -15,11 +15,14 @@ const ( FieldPost = "post" FieldEnv = "env" FieldEvent = "event" + FieldTransports = "transports" + FieldStores = "stores" ) // Adapter field names const ( FieldVersion = "version" + FieldStore = "store" ) // Parameter field names @@ -83,6 +86,18 @@ const ( TransportClientMaestro = "maestro" ) +// Deployment transport types. +const ( + TransportTypeKubernetes = "kubernetes" + TransportTypeRemote = "remote" +) + +// Deployment store types. +const ( + StoreTypeMemory = "memory" + StoreTypeRedis = "redis" +) + // Resource field names const ( FieldManifest = "manifest" diff --git a/internal/configloader/registry_config_test.go b/internal/configloader/registry_config_test.go new file mode 100644 index 00000000..69853a1c --- /dev/null +++ b/internal/configloader/registry_config_test.go @@ -0,0 +1,196 @@ +package configloader + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gopkg.in/yaml.v3" +) + +func TestLoadConfigCopiesNamedTransportAndStoreDefinitions(t *testing.T) { + tmpDir := t.TempDir() + + adapterPath, taskPath := createTestConfigFiles(t, tmpDir, ` +adapter: + name: test-adapter + version: "1.0.0" +clients: + hyperfleet_api: + timeout: 5s + kubernetes: + api_version: v1 +stores: + desired-memory: + type: memory +transports: + remote-primary: + type: remote + store: desired-memory + remote-secondary: + type: remote + store: desired-memory +`, `{}`) + + config, err := LoadConfig( + WithAdapterConfigPath(adapterPath), + WithTaskConfigPath(taskPath), + WithSkipSemanticValidation(), + ) + require.NoError(t, err) + require.NotNil(t, config) + + assert.Len(t, config.Stores, 1) + assert.Equal(t, StoreDefinition{Type: StoreTypeMemory}, config.Stores["desired-memory"]) + assert.Len(t, config.Transports, 2) + assert.Equal( + t, + TransportDefinition{Type: TransportTypeRemote, Store: "desired-memory"}, + config.Transports["remote-primary"], + ) + assert.Equal( + t, + TransportDefinition{Type: TransportTypeRemote, Store: "desired-memory"}, + config.Transports["remote-secondary"], + ) +} + +func TestConfigRedactedRedactsRedisPasswordWithoutMutatingOriginal(t *testing.T) { + config := &Config{ + Stores: map[string]StoreDefinition{ + "authenticated": { + Type: StoreTypeRedis, + URL: "rediss://adapter:super-secret@redis.example.com:6379/0", + }, + "anonymous": {Type: StoreTypeRedis, URL: "redis://redis.example.com:6379/1"}, + "invalid": {Type: StoreTypeRedis, URL: "not a URL"}, + }, + } + + redacted := config.Redacted() + require.NotNil(t, redacted) + require.NotNil(t, redacted.Stores) + + assert.Equal( + t, + "rediss://adapter:%2A%2AREDACTED%2A%2A@redis.example.com:6379/0", + redacted.Stores["authenticated"].URL, + ) + assert.Equal(t, "redis://redis.example.com:6379/1", redacted.Stores["anonymous"].URL) + assert.Equal(t, "not a URL", redacted.Stores["invalid"].URL) + assert.Equal( + t, + "rediss://adapter:super-secret@redis.example.com:6379/0", + config.Stores["authenticated"].URL, + ) +} + +func TestAdapterConfigValidationRejectsInvalidNamedRegistryDefinitions(t *testing.T) { + const unsupportedError = "unsupported" + + tests := []struct { + name string + yaml string + errorMsg string + }{ + { + name: "unsupported transport type", + yaml: ` +adapter: + name: test-adapter +transports: + unknown: + type: unsupported +`, + errorMsg: unsupportedError, + }, + { + name: "remote transport without store", + yaml: ` +adapter: + name: test-adapter +transports: + remote: + type: remote +`, + errorMsg: "store is required", + }, + { + name: "remote transport references missing store", + yaml: ` +adapter: + name: test-adapter +transports: + remote: + type: remote + store: missing +`, + errorMsg: "missing", + }, + { + name: "unsupported store type", + yaml: ` +adapter: + name: test-adapter +stores: + desired: + type: unsupported +`, + errorMsg: unsupportedError, + }, + { + name: "redis store without URL", + yaml: ` +adapter: + name: test-adapter +stores: + desired: + type: redis +`, + errorMsg: "url is required", + }, + { + name: "redis store with malformed URL", + yaml: ` +adapter: + name: test-adapter +stores: + desired: + type: redis + url: not-a-redis-url +`, + errorMsg: "redis", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + var config AdapterConfig + err := yaml.Unmarshal([]byte(tt.yaml), &config) + require.NoError(t, err) + + err = NewAdapterConfigValidator(&config, "").ValidateStructure() + require.Error(t, err) + assert.Contains(t, err.Error(), tt.errorMsg) + }) + } +} + +func TestAdapterConfigValidationDoesNotExposeCredentialsFromInvalidRedisURL(t *testing.T) { + const password = "super-secret" + config := &AdapterConfig{ + Adapter: AdapterInfo{Name: "test-adapter"}, + Stores: map[string]StoreDefinition{ + "credentials": { + Type: StoreTypeRedis, + URL: "redis://adapter:" + password + "@redis.example.com:%", + }, + }, + } + + err := NewAdapterConfigValidator(config, "").ValidateStructure() + + require.Error(t, err) + assert.ErrorContains(t, err, "stores.credentials.url is invalid") + assert.NotContains(t, err.Error(), password) +} diff --git a/internal/configloader/types.go b/internal/configloader/types.go index 960e5607..d222d41d 100644 --- a/internal/configloader/types.go +++ b/internal/configloader/types.go @@ -2,6 +2,7 @@ package configloader import ( "fmt" + "net/url" "strings" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/hyperfleetapi" @@ -11,14 +12,16 @@ import ( // Config is the unified configuration passed throughout the application. // Created by merging AdapterConfig (deployment) and AdapterTaskConfig (task). type Config struct { - Post *PostConfig `yaml:"post,omitempty"` - Log LogConfig `yaml:"log,omitempty"` - Adapter AdapterInfo `yaml:"adapter"` - Params []Parameter `yaml:"params,omitempty"` - Preconditions []Precondition `yaml:"preconditions,omitempty"` - Resources []Resource `yaml:"resources,omitempty"` - Clients ClientsConfig `yaml:"clients"` - DebugConfig bool `yaml:"debug_config,omitempty"` + Transports map[string]TransportDefinition `yaml:"transports,omitempty"` + Stores map[string]StoreDefinition `yaml:"stores,omitempty"` + Post *PostConfig `yaml:"post,omitempty"` + Log LogConfig `yaml:"log,omitempty"` + Adapter AdapterInfo `yaml:"adapter"` + Params []Parameter `yaml:"params,omitempty"` + Preconditions []Precondition `yaml:"preconditions,omitempty"` + Resources []Resource `yaml:"resources,omitempty"` + Clients ClientsConfig `yaml:"clients"` + DebugConfig bool `yaml:"debug_config,omitempty"` } // Merge combines AdapterConfig (deployment) and AdapterTaskConfig (task) into a unified Config. @@ -32,6 +35,8 @@ func Merge(adapterCfg *AdapterConfig, taskCfg *AdapterTaskConfig) *Config { return &Config{ Adapter: adapterCfg.Adapter, Clients: adapterCfg.Clients, + Transports: adapterCfg.Transports, + Stores: adapterCfg.Stores, DebugConfig: adapterCfg.DebugConfig, Log: adapterCfg.Log, Params: taskCfg.Params, @@ -50,6 +55,7 @@ func (c *Config) Redacted() *Config { } copy := *c copy.Clients = redactedClients(c.Clients) + copy.Stores = redactedStores(c.Stores) return © } @@ -78,6 +84,32 @@ func redactedClients(clients ClientsConfig) ClientsConfig { return copy } +func redactedStores(stores map[string]StoreDefinition) map[string]StoreDefinition { + if stores == nil { + return nil + } + + copy := make(map[string]StoreDefinition, len(stores)) + for name, store := range stores { + store.URL = redactRedisURL(store.URL) + copy[name] = store + } + return copy +} + +func redactRedisURL(rawURL string) string { + parsedURL, err := url.Parse(rawURL) + if err != nil || parsedURL.User == nil { + return rawURL + } + if _, hasPassword := parsedURL.User.Password(); !hasPassword { + return rawURL + } + + parsedURL.User = url.UserPassword(parsedURL.User.Username(), redactedValue) + return parsedURL.String() +} + // FieldExpressionDef represents a common pattern for value extraction. // Used when a value should be computed via field extraction (JSONPath) or CEL expression. // Only one of Field or Expression should be set. @@ -644,10 +676,25 @@ func (ve *ValidationErrors) HasErrors() bool { // Contains infrastructure settings that can be overridden via environment variables // and CLI flags using Viper. type AdapterConfig struct { - Adapter AdapterInfo `yaml:"adapter" mapstructure:"adapter"` - Log LogConfig `yaml:"log,omitempty" mapstructure:"log"` - Clients ClientsConfig `yaml:"clients" mapstructure:"clients"` - DebugConfig bool `yaml:"debug_config,omitempty" mapstructure:"debug_config"` + Transports map[string]TransportDefinition `yaml:"transports,omitempty" mapstructure:"transports"` + Stores map[string]StoreDefinition `yaml:"stores,omitempty" mapstructure:"stores"` + Log LogConfig `yaml:"log,omitempty" mapstructure:"log"` + Adapter AdapterInfo `yaml:"adapter" mapstructure:"adapter"` + Clients ClientsConfig `yaml:"clients" mapstructure:"clients"` + DebugConfig bool `yaml:"debug_config,omitempty" mapstructure:"debug_config"` +} + +// TransportDefinition is a named deployment transport entry. It is distinct +// from TransportConfig, which belongs to the task resource DSL. +type TransportDefinition struct { + Type string `yaml:"type" mapstructure:"type"` + Store string `yaml:"store,omitempty" mapstructure:"store"` +} + +// StoreDefinition is a named deployment store entry used by remote transports. +type StoreDefinition struct { + Type string `yaml:"type" mapstructure:"type"` + URL string `yaml:"url,omitempty" mapstructure:"url"` } // ClientsConfig contains configuration for all external clients diff --git a/internal/configloader/validator.go b/internal/configloader/validator.go index 7692df70..cb817193 100644 --- a/internal/configloader/validator.go +++ b/internal/configloader/validator.go @@ -4,14 +4,17 @@ import ( "context" "fmt" "log/slog" + "maps" "os" "path/filepath" "reflect" "regexp" + "slices" "strings" "cel.dev/cel-go/cel" "github.com/Masterminds/semver/v3" + "github.com/redis/go-redis/v9" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/criteria" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" @@ -54,10 +57,64 @@ func (v *AdapterConfigValidator) ValidateStructure() error { if err := v.validateHyperfleetAuth(); err != nil { return err } + if err := v.validateTransportRegistry(); err != nil { + return err + } return nil } +func (v *AdapterConfigValidator) validateTransportRegistry() error { + // sort for deterministic validation order + storeNames := sortedStoreNames(v.config.Stores) + for _, name := range storeNames { + store := v.config.Stores[name] + path := fmt.Sprintf("%s.%s", FieldStores, name) + switch store.Type { + case StoreTypeMemory: + case StoreTypeRedis: + if strings.TrimSpace(store.URL) == "" { + return fmt.Errorf("%s.url is required for redis store", path) + } + if _, err := redis.ParseURL(store.URL); err != nil { + return fmt.Errorf("%s.url is invalid for redis store", path) + } + default: + return fmt.Errorf("%s.type %q is unsupported (supported: %s, %s)", + path, store.Type, StoreTypeMemory, StoreTypeRedis) + } + } + + // sort for deterministic validation order + for _, name := range sortedTransportNames(v.config.Transports) { + transport := v.config.Transports[name] + path := fmt.Sprintf("%s.%s", FieldTransports, name) + switch transport.Type { + case TransportTypeKubernetes: + case TransportTypeRemote: + if strings.TrimSpace(transport.Store) == "" { + return fmt.Errorf("%s.%s is required for remote transport", path, FieldStore) + } + if _, ok := v.config.Stores[transport.Store]; !ok { + return fmt.Errorf("%s.%s references unknown store %q", path, FieldStore, transport.Store) + } + default: + return fmt.Errorf("%s.type %q is unsupported (supported: %s, %s)", + path, transport.Type, TransportTypeKubernetes, TransportTypeRemote) + } + } + + return nil +} + +func sortedStoreNames(stores map[string]StoreDefinition) []string { + return slices.Sorted(maps.Keys(stores)) +} + +func sortedTransportNames(transports map[string]TransportDefinition) []string { + return slices.Sorted(maps.Keys(transports)) +} + func (v *AdapterConfigValidator) validateHyperfleetAuth() error { auth := v.config.Clients.HyperfleetAPI.Auth if auth == nil { diff --git a/internal/transportclient/registry.go b/internal/transportclient/registry.go new file mode 100644 index 00000000..d0302ee5 --- /dev/null +++ b/internal/transportclient/registry.go @@ -0,0 +1,17 @@ +package transportclient + +import "fmt" + +// Registry associates configured transport client names with their implementations. +// It deliberately performs no routing or fallback; callers choose the client name. +type Registry map[string]TransportClient + +// Get returns the transport client registered under name. +func (r Registry) Get(name string) (TransportClient, error) { + client, ok := r[name] + if !ok || client == nil { + return nil, fmt.Errorf("transport client %q not configured", name) + } + + return client, nil +} diff --git a/internal/transportclient/registry_test.go b/internal/transportclient/registry_test.go new file mode 100644 index 00000000..1f61c5b0 --- /dev/null +++ b/internal/transportclient/registry_test.go @@ -0,0 +1,103 @@ +package transportclient + +import ( + "context" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +func TestRegistryKeepsEntriesDistinctByConfiguredName(t *testing.T) { + registry := Registry{ + "remote-primary": nil, + "remote-secondary": nil, + } + + assert.Len(t, registry, 2) + assert.Contains(t, registry, "remote-primary") + assert.Contains(t, registry, "remote-secondary") +} + +func TestRegistryGetRejectsUnknownAndNilEntries(t *testing.T) { + registry := Registry{ + "configured-but-nil": nil, + } + + tests := []struct { + name string + key string + }{ + { + name: "unknown configured name", + key: "not-configured", + }, + { + name: "nil configured entry", + key: "configured-but-nil", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + client, err := registry.Get(tt.key) + require.Error(t, err) + assert.Nil(t, client) + assert.Contains(t, err.Error(), tt.key) + }) + } +} + +func TestRegistryGetReturnsRegisteredClient(t *testing.T) { + registered := &stubTransportClient{} + registry := Registry{"remote-primary": registered} + + client, err := registry.Get("remote-primary") + + require.NoError(t, err) + assert.Same(t, registered, client) +} + +type stubTransportClient struct{} + +func (*stubTransportClient) ApplyResource( + context.Context, + []byte, + *ApplyOptions, + TransportContext, +) (*ApplyResult, error) { + return nil, nil +} + +func (*stubTransportClient) GetResource( + context.Context, + schema.GroupVersionKind, + string, + string, + TransportContext, +) (*unstructured.Unstructured, error) { + return nil, nil +} + +func (*stubTransportClient) DiscoverResources( + context.Context, + schema.GroupVersionKind, + manifest.Discovery, + TransportContext, +) (*unstructured.UnstructuredList, error) { + return nil, nil +} + +func (*stubTransportClient) DeleteResource( + context.Context, + schema.GroupVersionKind, + string, + string, + *DeleteOptions, + TransportContext, +) error { + return nil +} diff --git a/internal/transportregistry/registry.go b/internal/transportregistry/registry.go new file mode 100644 index 00000000..976dde3b --- /dev/null +++ b/internal/transportregistry/registry.go @@ -0,0 +1,263 @@ +// Package transportregistry builds the configured transport clients and owns +// the resources whose lifetimes they require. +package transportregistry + +import ( + "context" + "fmt" + "io" + "log/slog" + "maps" + "slices" + "time" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/configloader" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/desireclient" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/k8sclient" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/maestroclient" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire/store/memory" + redisstore "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire/store/redis" + "github.com/redis/go-redis/v9" +) + +const redisPingTimeout = 5 * time.Second + +// Runtime is the configured client registry and the resources it owns. +type Runtime struct { + Registry transportclient.Registry + closers []io.Closer +} + +// Build constructs each declared transport and its backing stores. +func Build(ctx context.Context, config *configloader.Config) (*Runtime, error) { + if config == nil { + return nil, fmt.Errorf("transport registry config is required") + } + + runtime := &Runtime{Registry: make(transportclient.Registry)} + if len(config.Transports) == 0 { + if err := runtime.buildLegacy(ctx, config); err != nil { + closeAfterBuildFailure(ctx, runtime) + return nil, err + } + return runtime, nil + } + + stores, err := runtime.buildStores(ctx, config.Stores) + if err != nil { + closeAfterBuildFailure(ctx, runtime) + return nil, err + } + + for _, name := range sortedTransportNames(config.Transports) { + definition := config.Transports[name] + client, err := buildTransport(ctx, config, definition, stores) + if err != nil { + closeAfterBuildFailure(ctx, runtime) + return nil, fmt.Errorf("build transport %q: %w", name, err) + } + runtime.Registry[name] = client + } + runtime.registerCompatibilityAlias(config, nil) + + return runtime, nil +} + +func closeAfterBuildFailure(ctx context.Context, runtime *Runtime) { + if err := runtime.Close(); err != nil { + slog.WarnContext(ctx, "failed to close transport registry after build failure", "error", err) + } +} + +// BuildRecording builds a registry for dry-run execution without creating any +// network clients. Each configured name uses client directly. +func BuildRecording( + config *configloader.Config, + client transportclient.TransportClient, +) (*Runtime, error) { + if config == nil { + return nil, fmt.Errorf("transport registry config is required") + } + if client == nil { + return nil, fmt.Errorf("recording transport client is required") + } + + runtime := &Runtime{Registry: make(transportclient.Registry)} + for name := range config.Transports { + runtime.Registry[name] = client + } + runtime.registerCompatibilityAlias(config, client) + return runtime, nil +} + +// registerCompatibilityAlias keeps the singleton executor working until it +// resolves configured transport names itself. Prefer the historical key when +// it is declared; otherwise use the first configured name deterministically. +func (r *Runtime) registerCompatibilityAlias( + config *configloader.Config, + fallback transportclient.TransportClient, +) { + key := configloader.TransportClientKubernetes + if config.Clients.Maestro != nil { + key = configloader.TransportClientMaestro + } + if _, ok := r.Registry[key]; ok { + return + } + + names := sortedTransportNames(config.Transports) + if len(names) > 0 { + r.Registry[key] = r.Registry[names[0]] + } else if fallback != nil { + r.Registry[key] = fallback + } +} + +// Close releases all resources created by Build. It attempts every close and +// returns the first error encountered. +func (r *Runtime) Close() error { + if r == nil { + return nil + } + var firstErr error + for index := len(r.closers) - 1; index >= 0; index-- { + if err := r.closers[index].Close(); err != nil && firstErr == nil { + firstErr = err + } + } + r.closers = nil + return firstErr +} + +func (r *Runtime) buildLegacy(ctx context.Context, config *configloader.Config) error { + // Preserve the old selection behavior: Maestro is the sole default when it + // is configured; Kubernetes is otherwise the default. + if config.Clients.Maestro != nil { + client, err := buildMaestro(ctx, config.Clients.Maestro) + if err != nil { + return fmt.Errorf("build transport %q: %w", configloader.TransportClientMaestro, err) + } + r.Registry[configloader.TransportClientMaestro] = client + r.closers = append(r.closers, client) + return nil + } + + client, err := buildKubernetes(ctx, config.Clients.Kubernetes) + if err != nil { + return fmt.Errorf("build transport %q: %w", configloader.TransportClientKubernetes, err) + } + r.Registry[configloader.TransportClientKubernetes] = client + return nil +} + +func (r *Runtime) buildStores( + ctx context.Context, + definitions map[string]configloader.StoreDefinition, +) (map[string]desire.SpecStore, error) { + stores := make(map[string]desire.SpecStore, len(definitions)) + for _, name := range sortedStoreNames(definitions) { + definition := definitions[name] + switch definition.Type { + case configloader.StoreTypeMemory: + stores[name] = memory.New() + case configloader.StoreTypeRedis: + options, err := redis.ParseURL(definition.URL) + if err != nil { + return nil, fmt.Errorf("build store %q: redis URL is invalid", name) + } + client := redis.NewClient(options) + pingCtx, cancel := context.WithTimeout(ctx, redisPingTimeout) + err = client.Ping(pingCtx).Err() + cancel() + if err != nil { + if closeErr := client.Close(); closeErr != nil { + slog.WarnContext(ctx, "failed to close Redis client after ping failure", "error", closeErr) + } + return nil, fmt.Errorf("build store %q: ping Redis: %w", name, err) + } + r.closers = append(r.closers, client) + stores[name] = redisstore.New(client) + default: + return nil, fmt.Errorf("build store %q: unsupported type %q", name, definition.Type) + } + } + return stores, nil +} + +func buildTransport( + ctx context.Context, + config *configloader.Config, + definition configloader.TransportDefinition, + stores map[string]desire.SpecStore, +) (transportclient.TransportClient, error) { + switch definition.Type { + case configloader.TransportTypeKubernetes: + return buildKubernetes(ctx, config.Clients.Kubernetes) + case configloader.TransportTypeRemote: + store, ok := stores[definition.Store] + if !ok { + return nil, fmt.Errorf("store %q is not configured", definition.Store) + } + return desireclient.NewClient(store, config.Adapter.Name), nil + default: + return nil, fmt.Errorf("unsupported type %q", definition.Type) + } +} + +func buildKubernetes( + ctx context.Context, + config configloader.KubernetesConfig, +) (*k8sclient.Client, error) { + return k8sclient.NewClient(ctx, k8sclient.ClientConfig{ + KubeConfigPath: config.KubeConfigPath, + QPS: config.QPS, + Burst: config.Burst, + }) +} + +func buildMaestro( + ctx context.Context, + config *configloader.MaestroClientConfig, +) (*maestroclient.Client, error) { + maestroConfig := &maestroclient.Config{ + MaestroServerAddr: config.HTTPServerAddress, + GRPCServerAddr: config.GRPCServerAddress, + SourceID: config.SourceID, + Insecure: config.Insecure, + } + if config.Timeout != "" { + timeout, err := time.ParseDuration(config.Timeout) + if err != nil { + return nil, fmt.Errorf("invalid maestro timeout %q: %w", config.Timeout, err) + } + maestroConfig.HTTPTimeout = timeout + } + if config.ServerHealthinessTimeout != "" { + timeout, err := time.ParseDuration(config.ServerHealthinessTimeout) + if err != nil { + return nil, fmt.Errorf( + "invalid maestro serverHealthinessTimeout %q: %w", + config.ServerHealthinessTimeout, + err, + ) + } + maestroConfig.ServerHealthinessTimeout = timeout + } + if config.Auth.TLSConfig != nil { + maestroConfig.CAFile = config.Auth.TLSConfig.CAFile + maestroConfig.ClientCertFile = config.Auth.TLSConfig.CertFile + maestroConfig.ClientKeyFile = config.Auth.TLSConfig.KeyFile + maestroConfig.HTTPCAFile = config.Auth.TLSConfig.HTTPCAFile + } + return maestroclient.NewMaestroClient(ctx, maestroConfig) +} + +func sortedStoreNames(definitions map[string]configloader.StoreDefinition) []string { + return slices.Sorted(maps.Keys(definitions)) +} + +func sortedTransportNames(definitions map[string]configloader.TransportDefinition) []string { + return slices.Sorted(maps.Keys(definitions)) +} diff --git a/internal/transportregistry/runtime_test.go b/internal/transportregistry/runtime_test.go new file mode 100644 index 00000000..f7357eee --- /dev/null +++ b/internal/transportregistry/runtime_test.go @@ -0,0 +1,451 @@ +package transportregistry + +import ( + "errors" + "io" + "os" + "path/filepath" + "testing" + + "github.com/alicebob/miniredis/v2" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/configloader" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/desireclient" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/dryrun" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestBuildCreatesRemoteClientsForEachConfiguredName(t *testing.T) { + config := loadRuntimeConfig(t, ` +adapter: + name: test-adapter +stores: + desired-memory: + type: memory +transports: + remote-primary: + type: remote + store: desired-memory + remote-secondary: + type: remote + store: desired-memory +`) + + runtime, err := Build(t.Context(), config) + require.NoError(t, err) + require.NotNil(t, runtime) + t.Cleanup(func() { assert.NoError(t, runtime.Close()) }) + + primary, err := runtime.Registry.Get("remote-primary") + require.NoError(t, err) + assert.NotNil(t, primary) + + secondary, err := runtime.Registry.Get("remote-secondary") + require.NoError(t, err) + assert.NotNil(t, secondary) + + compatibilityClient, err := runtime.Registry.Get(configloader.TransportClientKubernetes) + require.NoError(t, err) + assert.Same(t, primary, compatibilityClient) + + primaryResult, err := primary.ApplyResource( + t.Context(), + testConfigMapManifest, + nil, + testTransportContext(), + ) + require.NoError(t, err) + assert.Equal(t, manifest.OperationCreate, primaryResult.Operation) + + secondaryResult, err := secondary.ApplyResource( + t.Context(), + testConfigMapManifest, + nil, + testTransportContext(), + ) + require.NoError(t, err) + assert.Equal(t, manifest.OperationSkip, secondaryResult.Operation, + "transports using the same named store must see the same desired state") +} + +func TestBuildRegistersMaestroCompatibilityAliasForNamedTransports(t *testing.T) { + config := loadRuntimeConfig(t, ` +adapter: + name: test-adapter +clients: + maestro: + source_id: test-adapter +stores: + desired-memory: + type: memory +transports: + kubernetes: + type: remote + store: desired-memory + remote: + type: remote + store: desired-memory +`) + + runtime, err := Build(t.Context(), config) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, runtime.Close()) }) + + namedClient, err := runtime.Registry.Get(configloader.TransportClientKubernetes) + require.NoError(t, err) + compatibilityClient, err := runtime.Registry.Get(configloader.TransportClientMaestro) + require.NoError(t, err) + assert.Same(t, namedClient, compatibilityClient) +} + +func TestBuildPingsRedisStoreBeforeReturning(t *testing.T) { + config := loadRuntimeConfig(t, ` +adapter: + name: test-adapter +stores: + unreachable: + type: redis + url: redis://127.0.0.1:1/0 +transports: + remote: + type: remote + store: unreachable +`) + + runtime, err := Build(t.Context(), config) + require.Error(t, err) + assert.Nil(t, runtime) + assert.Contains(t, err.Error(), "unreachable") +} + +func TestBuildDoesNotExposeCredentialsFromInvalidRedisURL(t *testing.T) { + const password = "super-secret" + config := &configloader.Config{ + Adapter: configloader.AdapterInfo{Name: "test-adapter"}, + Stores: map[string]configloader.StoreDefinition{ + "credentials": { + Type: configloader.StoreTypeRedis, + URL: "redis://adapter:" + password + "@redis.example.com:%", + }, + }, + Transports: map[string]configloader.TransportDefinition{ + "remote": {Type: configloader.TransportTypeRemote, Store: "credentials"}, + }, + } + + runtime, err := Build(t.Context(), config) + + require.Error(t, err) + assert.Nil(t, runtime) + assert.ErrorContains(t, err, "credentials") + assert.NotContains(t, err.Error(), password) +} + +func TestBuildCreatesAndClosesReachableRedisStore(t *testing.T) { + server := miniredis.RunT(t) + config := loadRuntimeConfig(t, ` +adapter: + name: test-adapter +stores: + redis-store: + type: redis + url: redis://`+server.Addr()+`/0 +transports: + remote: + type: remote + store: redis-store +`) + + runtime, err := Build(t.Context(), config) + require.NoError(t, err) + require.NotNil(t, runtime) + + client, err := runtime.Registry.Get("remote") + require.NoError(t, err) + assert.NotNil(t, client) + assert.NoError(t, runtime.Close()) +} + +func TestBuildRejectsMissingConfiguration(t *testing.T) { + runtime, err := Build(t.Context(), nil) + + require.Error(t, err) + assert.Nil(t, runtime) + assert.Contains(t, err.Error(), "config is required") +} + +func TestBuildRejectsInvalidNamedTransportBeforeStartingService(t *testing.T) { + tests := []struct { + name string + config *configloader.Config + want string + }{ + { + name: "remote transport references absent store", + config: &configloader.Config{ + Transports: map[string]configloader.TransportDefinition{ + "remote": { + Type: configloader.TransportTypeRemote, + Store: "missing", + }, + }, + }, + want: "store \"missing\" is not configured", + }, + { + name: "unsupported transport type", + config: &configloader.Config{ + Transports: map[string]configloader.TransportDefinition{ + "unknown": {Type: "unsupported"}, + }, + }, + want: "unsupported type \"unsupported\"", + }, + { + name: "Kubernetes transport has unusable kubeconfig", + config: &configloader.Config{ + Transports: map[string]configloader.TransportDefinition{ + "kubernetes": {Type: configloader.TransportTypeKubernetes}, + }, + Clients: configloader.ClientsConfig{ + Kubernetes: configloader.KubernetesConfig{ + KubeConfigPath: "/does/not/exist", + }, + }, + }, + want: "failed to load kubeconfig", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + runtime, err := Build(t.Context(), tt.config) + + require.Error(t, err) + assert.Nil(t, runtime) + assert.ErrorContains(t, err, tt.want) + }) + } +} + +func TestBuildLegacyRejectsInvalidMaestroConfiguration(t *testing.T) { + config := &configloader.Config{Clients: configloader.ClientsConfig{ + Maestro: &configloader.MaestroClientConfig{Timeout: "not-a-duration"}, + }} + + runtime, err := Build(t.Context(), config) + + require.Error(t, err) + assert.Nil(t, runtime) + assert.ErrorContains(t, err, "invalid maestro timeout") +} + +func TestBuildPreservesLegacyKubernetesDefault(t *testing.T) { + kubeconfigPath := filepath.Join(t.TempDir(), "kubeconfig") + require.NoError(t, os.WriteFile(kubeconfigPath, []byte(` +apiVersion: v1 +kind: Config +clusters: + - name: test + cluster: + server: https://127.0.0.1:65535 +contexts: + - name: test + context: + cluster: test + user: test +current-context: test +users: + - name: test + user: + token: test-token +`), 0644)) + config := &configloader.Config{Clients: configloader.ClientsConfig{ + Kubernetes: configloader.KubernetesConfig{ + KubeConfigPath: kubeconfigPath, + }, + }} + + runtime, err := Build(t.Context(), config) + require.NoError(t, err) + t.Cleanup(func() { assert.NoError(t, runtime.Close()) }) + + client, err := runtime.Registry.Get(configloader.TransportClientKubernetes) + require.NoError(t, err) + assert.NotNil(t, client) +} + +func TestRuntimeCloseAttemptsAllClosersAndIsIdempotent(t *testing.T) { + firstErr := errors.New("first close failed") + first := &testCloser{err: firstErr} + second := &testCloser{err: errors.New("second close failed")} + runtime := &Runtime{closers: []io.Closer{first, second}} + + err := runtime.Close() + + assert.ErrorIs( + t, + err, + second.err, + "closers are released in reverse construction order", + ) + assert.Equal(t, 1, first.calls) + assert.Equal(t, 1, second.calls) + assert.NoError(t, runtime.Close()) + assert.Equal(t, 1, first.calls) + assert.Equal(t, 1, second.calls) +} + +func TestBuildRecordingRejectsMissingInputs(t *testing.T) { + tests := []struct { + client transportclient.TransportClient + config *configloader.Config + name string + }{ + {name: "nil config", client: dryrun.NewDryrunTransportClient()}, + {name: "nil recording client", config: &configloader.Config{}}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + runtime, err := BuildRecording(tt.config, tt.client) + + require.Error(t, err) + assert.Nil(t, runtime) + }) + } +} + +func TestBuildRecordingMapsConfiguredNamesAndCompatibilityDefault( + t *testing.T, +) { + config := loadRuntimeConfig(t, ` +adapter: + name: test-adapter +stores: + desired-memory: + type: memory +transports: + remote-primary: + type: remote + store: desired-memory + remote-secondary: + type: remote + store: desired-memory +`) + recorder := dryrun.NewDryrunTransportClient() + + runtime, err := BuildRecording(config, recorder) + require.NoError(t, err) + require.NotNil(t, runtime) + + for _, name := range []string{"remote-primary", "remote-secondary", "kubernetes"} { + client, err := runtime.Registry.Get(name) + require.NoErrorf(t, err, "expected recording client for %q", name) + assert.Same(t, recorder, client) + } +} + +func TestBuildRecordingPreservesLegacyKubernetesAndMaestroDefaults( + t *testing.T, +) { + tests := []struct { + name string + config string + key string + }{ + { + name: "legacy Kubernetes configuration", + config: ` +adapter: + name: test-adapter +clients: + kubernetes: + api_version: v1 +`, + key: "kubernetes", + }, + { + name: "legacy Maestro configuration", + config: ` +adapter: + name: test-adapter +clients: + maestro: + source_id: test-adapter +`, + key: "maestro", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + runtime, err := BuildRecording( + loadRuntimeConfig(t, tt.config), + dryrun.NewDryrunTransportClient(), + ) + require.NoError(t, err) + require.NotNil(t, runtime) + + client, err := runtime.Registry.Get(tt.key) + require.NoError(t, err) + assert.NotNil(t, client) + }) + } +} + +func loadRuntimeConfig(t *testing.T, adapterYAML string) *configloader.Config { + t.Helper() + + adapterPath, taskPath := createRuntimeConfigFiles(t, adapterYAML, `{}`) + config, err := configloader.LoadConfig( + configloader.WithAdapterConfigPath(adapterPath), + configloader.WithTaskConfigPath(taskPath), + configloader.WithSkipSemanticValidation(), + ) + require.NoError(t, err) + return config +} + +func createRuntimeConfigFiles( + t *testing.T, + adapterYAML, taskYAML string, +) (string, string) { + t.Helper() + + tmpDir := t.TempDir() + adapterPath := filepath.Join(tmpDir, "adapter-config.yaml") + taskPath := filepath.Join(tmpDir, "task-config.yaml") + require.NoError(t, os.WriteFile(adapterPath, []byte(adapterYAML), 0644)) + require.NoError(t, os.WriteFile(taskPath, []byte(taskYAML), 0644)) + return adapterPath, taskPath +} + +func testTransportContext() *desireclient.TransportContext { + return &desireclient.TransportContext{ + ManagementCluster: "management-cluster", + Resource: "configmaps", + } +} + +var testConfigMapManifest = []byte(`{ + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": { + "name": "registry-test", + "namespace": "default", + "annotations": {"hyperfleet.io/generation": "1"} + } +}`) + +type testCloser struct { + err error + calls int +} + +func (c *testCloser) Close() error { + c.calls++ + return c.err +}