diff --git a/cmd/adapter/main.go b/cmd/adapter/main.go index 8735f038..721eb6d3 100644 --- a/cmd/adapter/main.go +++ b/cmd/adapter/main.go @@ -15,9 +15,8 @@ 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/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/logger" "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/metrics" @@ -304,88 +303,6 @@ func createAPIClient(apiConfig configloader.HyperfleetAPIConfig, log logger.Logg return hyperfleetapi.NewClient(log, opts...) } -// createTransportClient creates the appropriate transport client based on config. -func createTransportClient( - ctx context.Context, - config *configloader.Config, - log logger.Logger, -) (transportclient.TransportClient, error) { - if config.Clients.Maestro != nil { - log.Info(ctx, "Creating Maestro transport client...") - client, err := createMaestroClient(ctx, config.Clients.Maestro, log) - if err != nil { - return nil, fmt.Errorf("failed to create Maestro client: %w", err) - } - log.Info(ctx, "Maestro transport client created successfully") - return client, nil - } - - log.Info(ctx, "Creating Kubernetes transport client...") - client, err := createK8sClient(ctx, config.Clients.Kubernetes, log) - if err != nil { - return nil, fmt.Errorf("failed to create Kubernetes client: %w", err) - } - log.Info(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, - log logger.Logger, -) (*k8sclient.Client, error) { - clientConfig := k8sclient.ClientConfig{ - KubeConfigPath: k8sConfig.KubeConfigPath, - QPS: k8sConfig.QPS, - Burst: k8sConfig.Burst, - } - return k8sclient.NewClient(ctx, clientConfig, log) -} - -// createMaestroClient creates a Maestro client from the config -func createMaestroClient( - ctx context.Context, - maestroConfig *configloader.MaestroClientConfig, - log logger.Logger, -) (*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, log) -} - // buildExecutor creates the executor with the given clients. func buildExecutor( config *configloader.Config, @@ -545,12 +462,27 @@ func runServe(flags *pflag.FlagSet) error { return fmt.Errorf("failed to create HyperFleet API client: %w", err) } - tc, err := createTransportClient(ctx, config, log) + transportRuntime, err := transportregistry.Build(ctx, config, log) if err != nil { errCtx := logger.WithErrorField(ctx, err) - log.Errorf(errCtx, "Failed to create transport client") + log.Errorf(errCtx, "Failed to create transport registry") return err } + defer func() { + if closeErr := transportRuntime.Close(); closeErr != nil { + errCtx := logger.WithErrorField(ctx, closeErr) + log.Warnf(errCtx, "Failed to close transport registry") + } + }() + + 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 transport client: %w", err) + } // Build executor log.Info(ctx, "Creating event executor...") @@ -740,8 +672,28 @@ 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 { + errCtx := logger.WithErrorField(ctx, closeErr) + log.Warnf(errCtx, "Failed to close recording transport registry") + } + }() + + 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, log, nil) + exec, err := buildExecutor(config, dryrunAPI, tc, log, 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 b4bbc850..0fa08470 100644 --- a/go.mod +++ b/go.mod @@ -5,15 +5,18 @@ 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 github.com/mitchellh/copystructure v1.2.0 + github.com/openshift-hyperfleet/hyperfleet-applier v0.0.0-20260827133207-94b7a4d56697 github.com/openshift-hyperfleet/hyperfleet-broker v1.1.1 github.com/openshift-online/maestro v0.0.0-20260202062555-48b47506a254 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 @@ -128,6 +131,7 @@ require ( github.com/oklog/ulid v1.3.1 // indirect github.com/opencontainers/go-digest v1.0.0 // indirect github.com/opencontainers/image-spec v1.1.1 // indirect + github.com/openshift-hyperfleet/hyperfleet-logger v0.0.0-20260811173525-c9f9e282d029 // indirect github.com/pelletier/go-toml/v2 v2.4.2 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect @@ -146,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 @@ -158,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 @@ -180,7 +186,7 @@ require ( gopkg.in/yaml.v2 v2.4.0 // indirect k8s.io/api v0.37.0 // indirect k8s.io/klog/v2 v2.140.0 // indirect - k8s.io/kube-openapi v0.0.0-20260721132016-d427ff9ee9ad // indirect + k8s.io/kube-openapi v0.0.0-20260821135717-be32def86098 // indirect k8s.io/utils v0.0.0-20260707023825-cf1189d6abe3 // indirect sigs.k8s.io/json v0.0.0-20250730193827-2d320260d730 // indirect sigs.k8s.io/randfill v1.0.0 // indirect diff --git a/go.sum b/go.sum index 0df42505..fe933d38 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= @@ -273,8 +281,12 @@ github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8 github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM= github.com/opencontainers/image-spec v1.1.1 h1:y0fUlFfIZhPF1W537XOLg0/fcx6zcHCJwooC2xJA040= github.com/opencontainers/image-spec v1.1.1/go.mod h1:qpqAh3Dmcf36wStyyWU+kCeDgrGnAve2nCC8+7h8Q0M= +github.com/openshift-hyperfleet/hyperfleet-applier v0.0.0-20260827133207-94b7a4d56697 h1:9DFSlPuXlMY4zshztT4/MEnY6QAgu0WoVY1KFtVE9fg= +github.com/openshift-hyperfleet/hyperfleet-applier v0.0.0-20260827133207-94b7a4d56697/go.mod h1:tDZnkLwpVO5ksfPm3NH9dy5Yl0ZBSfbtOA3hj2wrmHc= github.com/openshift-hyperfleet/hyperfleet-broker v1.1.1 h1:3zbpNuFL+OEvKl6a/KJAlHFcYR4QqQBoGkf35higypU= github.com/openshift-hyperfleet/hyperfleet-broker v1.1.1/go.mod h1:E7Br4NnsaTTfWR2fEqHAtvFXUAgzFpksF+G5qTBMmy0= +github.com/openshift-hyperfleet/hyperfleet-logger v0.0.0-20260811173525-c9f9e282d029 h1:c3GdD3EUdR9lRjNot+aBIxVbAXbOSTtiQpWVTHnFHaU= +github.com/openshift-hyperfleet/hyperfleet-logger v0.0.0-20260811173525-c9f9e282d029/go.mod h1:5Nh2IMS2MouehZ4UiRXFMjhuwdys2DX26J4vArmXe6Y= github.com/openshift-online/maestro v0.0.0-20260202062555-48b47506a254 h1:v/jYqdzZpzB/bscVpajlbcKgCNeV4tx4fkm5m2JR8Ug= github.com/openshift-online/maestro v0.0.0-20260202062555-48b47506a254/go.mod h1:cyeif610uObNrbcyn5s1fZg7OWseVjaMAqgrEDA2Aec= github.com/openshift-online/ocm-sdk-go v0.1.510 h1:oPwgGHPi6LeY+RANXUhJwVvkY4yvddyRqHHThHZZRyM= @@ -303,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= @@ -350,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= @@ -392,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= @@ -508,16 +528,16 @@ honnef.co/go/tools v0.0.0-20190102054323-c2f93a96b099/go.mod h1:rf3lG4BRIbNafJWh honnef.co/go/tools v0.0.0-20190523083050-ea95bdfd59fc/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4= k8s.io/api v0.37.0 h1:Z//Vj9N7RA/yS2sDmxyeo7h+RR4zbUrd2vrd3Z0TbB4= k8s.io/api v0.37.0/go.mod h1:LKXgcJWMc+f4OLbP5SFR8rulEg07zZhpi/zMULiBImk= -k8s.io/apiextensions-apiserver v0.36.0 h1:Wt7E8J+VBCbj4FjiBfDTK/neXDDjyJVJc7xfuOHImZ0= -k8s.io/apiextensions-apiserver v0.36.0/go.mod h1:kGDjH0msuiIB3tgsYRV0kS9GqpMYMUsQ3GHv7TApyug= +k8s.io/apiextensions-apiserver v0.36.4 h1:SfvCVt+4CqKWvzuVytYDT5g9hyb9MztoiYELIkPVrFc= +k8s.io/apiextensions-apiserver v0.36.4/go.mod h1:JT9V2Ju7ys1FY4zbSpmX9XOvKB3/BwsODc4hFQEa+Xo= k8s.io/apimachinery v0.37.0 h1:Np2AbDtf8x6RDHiD8T9LbKJ9gaegeVNa8yNm5FuGKm0= k8s.io/apimachinery v0.37.0/go.mod h1:RN3nhprFSCxOi5Selxd7oMTXOe/c+ZbcE7Im+TS2zkE= k8s.io/client-go v0.37.0 h1:nsN31fy8wBySuZ+QRnKmrjRSQLOG2rvoGN0tKd12zhQ= k8s.io/client-go v0.37.0/go.mod h1:FcGqw+Ll/gNQiq+nPGY1Oyt9y7SgDh1d3MW3RFDEbn0= k8s.io/klog/v2 v2.140.0 h1:Tf+J3AH7xnUzZyVVXhTgGhEKnFqye14aadWv7bzXdzc= k8s.io/klog/v2 v2.140.0/go.mod h1:o+/RWfJ6PwpnFn7OyAG3QnO47BFsymfEfrz6XyYSSp0= -k8s.io/kube-openapi v0.0.0-20260721132016-d427ff9ee9ad h1:oXImqH8mQNk7PmvzKhmN3ddJoY6OnyM225MXwGHPm0A= -k8s.io/kube-openapi v0.0.0-20260721132016-d427ff9ee9ad/go.mod h1:0/mqHCVhlumdJ3BhCfnjSZQE037nAhNodh1/hK0T8/I= +k8s.io/kube-openapi v0.0.0-20260821135717-be32def86098 h1:z5+pcu1jTyKK5mNTe2/+x+U6Uuv9jRVOJQLaBJJMpeI= +k8s.io/kube-openapi v0.0.0-20260821135717-be32def86098/go.mod h1:0/mqHCVhlumdJ3BhCfnjSZQE037nAhNodh1/hK0T8/I= k8s.io/utils v0.0.0-20260707023825-cf1189d6abe3 h1:jVkFFVfXdXP74B/zbO3hM3hpSFD0xvhQ5U686DPurkE= k8s.io/utils v0.0.0-20260707023825-cf1189d6abe3/go.mod h1:M2s5JB1lIYP3jzZdorPLHXIPJzt9vv2muW5a6L9DtNM= open-cluster-management.io/api v1.3.0 h1:Q3miH38BE3N5+PesHQ0kcFi5nhX5350m7OJWapZcVqY= 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..4df81bb1 --- /dev/null +++ b/internal/configloader/registry_config_test.go @@ -0,0 +1,177 @@ +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) + }) + } +} 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 82cecc0d..29004c61 100644 --- a/internal/configloader/validator.go +++ b/internal/configloader/validator.go @@ -3,14 +3,17 @@ package configloader import ( "context" "fmt" + "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: %w", path, err) + } + 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/desireclient/apply.go b/internal/desireclient/apply.go new file mode 100644 index 00000000..b4970218 --- /dev/null +++ b/internal/desireclient/apply.go @@ -0,0 +1,166 @@ +package desireclient + +import ( + "context" + "encoding/json" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" +) + +// ApplyResource implements transportclient.TransportClient. It upserts an +// apply desire for the rendered manifest and auto-creates its paired read +// desire so the applied resource becomes visible to discovery. A read-desire +// pairing failure is returned as an error: an apply desire without its read +// desire is permanently invisible to discovery. +func (c *Client) ApplyResource( + ctx context.Context, + manifestBytes []byte, + opts *transportclient.ApplyOptions, + target transportclient.TransportContext, +) (*transportclient.ApplyResult, error) { + if len(manifestBytes) == 0 { + return nil, fmt.Errorf("desireclient: manifest bytes cannot be empty") + } + + tc, err := resolveTransportContext(target) + if err != nil { + return nil, err + } + + obj, err := parseToUnstructured(manifestBytes) + if err != nil { + return nil, fmt.Errorf("desireclient: failed to parse manifest: %w", err) + } + + // The store's ApplySpec.KubeContent must be valid JSON, but manifestBytes + // may have been YAML (parseToUnstructured accepts both). Re-marshal the + // parsed object rather than storing the original bytes verbatim. + kubeContent, err := json.Marshal(obj.Object) + if err != nil { + return nil, fmt.Errorf("desireclient: failed to marshal manifest to JSON: %w", err) + } + + gvk := obj.GroupVersionKind() + namespace, name := obj.GetNamespace(), obj.GetName() + + readID, err := buildIdentity(tc, desire.TypeRead, gvk, namespace, name) + if err != nil { + return nil, err + } + + if err = c.ensureReadDesire(ctx, readID, gvk.Version); err != nil { + return nil, fmt.Errorf( + "desireclient: failed to create paired read desire for %s/%s: %w", gvk.Kind, name, err) + } + + applyID, err := buildIdentity(tc, desire.TypeApply, gvk, namespace, name) + if err != nil { + return nil, err + } + + result, err := c.upsertApplyDesire(ctx, applyID, obj, kubeContent) + if err != nil { + return nil, err + } + + return result, nil +} + +// upsertApplyDesire creates or updates the apply desire, deciding the +// operation via the manifest's hyperfleet.io/generation annotation — the same +// signal k8sclient/maestroclient compare on apply. This is unrelated to the +// store's own CAS Version, used below purely as the optimistic-concurrency +// token for UpdateApplyDesireSpec. +func (c *Client) upsertApplyDesire( + ctx context.Context, + id desire.Identity, + obj *unstructured.Unstructured, + kubeContent []byte, +) (*transportclient.ApplyResult, error) { + existing, err := c.store.GetApplyDesire(ctx, id) + if err != nil && !errors.Is(err, desire.ErrNotFound) { + return nil, fmt.Errorf("desireclient: failed to get apply desire for %s/%s: %w", id.Namespace, id.Name, err) + } + exists := err == nil + + newGen := manifest.GetGenerationFromUnstructured(obj) + var existingGen int64 + if exists { + existingGen = generationFromKubeContent(existing.Spec.KubeContent) + } + + decision := manifest.CompareGenerations(newGen, existingGen, exists) + result := &transportclient.ApplyResult{Operation: decision.Operation, Reason: decision.Reason} + + switch decision.Operation { + case manifest.OperationCreate: + if _, err := c.store.CreateApplyDesire(ctx, desire.ApplyDesire{ + Identity: id, + Owner: c.owner, + Spec: desire.ApplySpec{KubeContent: kubeContent}, + }); err != nil { + return nil, fmt.Errorf("desireclient: failed to create apply desire for %s/%s: %w", id.Namespace, id.Name, err) + } + case manifest.OperationUpdate: + if _, err := c.store.UpdateApplyDesireSpec( + ctx, id, desire.ApplySpec{KubeContent: kubeContent}, c.owner, existing.Version, + ); err != nil { + return nil, fmt.Errorf("desireclient: failed to update apply desire for %s/%s: %w", id.Namespace, id.Name, err) + } + case manifest.OperationSkip: + // Nothing to do. + default: + return nil, fmt.Errorf("desireclient: unexpected apply decision operation %q", decision.Operation) + } + + c.log.Debugf(ctx, "ApplyResource %s/%s: operation=%s reason=%s", + id.Namespace, id.Name, result.Operation, result.Reason) + return result, nil +} + +// ensureReadDesire creates the paired read desire if it doesn't already +// exist. ErrAlreadyExists is a no-op success. +func (c *Client) ensureReadDesire(ctx context.Context, id desire.Identity, targetVersion string) error { + _, err := c.store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: id, + Owner: c.owner, + TargetVersion: targetVersion, + }) + if err == nil { + return nil + } + + if !errors.Is(err, desire.ErrAlreadyExists) { + return fmt.Errorf("desireclient: failed to create read desire for %s/%s: %w", id.Namespace, id.Name, err) + } + + des, err := c.store.GetReadDesire(ctx, id) + if err != nil { + return fmt.Errorf("desireclient: failed to get existing read desire for %s/%s: %w", id.Namespace, id.Name, err) + } + + if des.TargetVersion == targetVersion { + return nil + } + + err = c.store.DeleteReadDesire(ctx, id, c.owner, des.Version) + if err != nil { + return fmt.Errorf("desireclient: failed to delete existing read desire for %s/%s: %w", id.Namespace, id.Name, err) + } + + if _, err := c.store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: id, + Owner: c.owner, + TargetVersion: targetVersion, + }); err != nil { + return fmt.Errorf("desireclient: failed to create read desire for %s/%s: %w", id.Namespace, id.Name, err) + } + + return nil + +} diff --git a/internal/desireclient/apply_test.go b/internal/desireclient/apply_test.go new file mode 100644 index 00000000..b115b773 --- /dev/null +++ b/internal/desireclient/apply_test.go @@ -0,0 +1,232 @@ +package desireclient + +import ( + "context" + "encoding/json" + "errors" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestApplyResource_CreatesApplyAndReadDesire(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + result, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + assert.Equal(t, manifest.OperationCreate, result.Operation) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + applied, err := store.GetApplyDesire(ctx, applyID) + require.NoError(t, err) + assert.Equal(t, testOwner, applied.Owner) + + readID := applyID + readID.Type = desire.TypeRead + readDesire, err := store.GetReadDesire(ctx, readID) + require.NoError(t, err, "paired read desire must be auto-created") + assert.Equal(t, "v1", readDesire.TargetVersion) + assert.Equal(t, testOwner, readDesire.Owner) +} + +func TestApplyResource_UpdatesWithCASVersion(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + before, err := store.GetApplyDesire(ctx, applyID) + require.NoError(t, err) + require.Equal(t, int64(1), before.Version) + + result, err := c.ApplyResource(ctx, configMapManifest(2), nil, testTransportContext()) + require.NoError(t, err) + assert.Equal(t, manifest.OperationUpdate, result.Operation) + + after, err := store.GetApplyDesire(ctx, applyID) + require.NoError(t, err) + assert.Equal(t, int64(2), after.Version, "UpdateApplyDesireSpec must have used the fetched CAS version") + assert.Contains(t, string(after.Spec.KubeContent), `"key":"value"`) +} + +func TestApplyResource_AcceptsYAMLManifest(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + yamlManifest := []byte(` +apiVersion: v1 +kind: ConfigMap +metadata: + name: my-config + namespace: default + annotations: + hyperfleet.io/generation: "1" +data: + key: value +`) + + result, err := c.ApplyResource(ctx, yamlManifest, nil, testTransportContext()) + require.NoError(t, err) + assert.Equal(t, manifest.OperationCreate, result.Operation) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + applied, err := store.GetApplyDesire(ctx, applyID) + require.NoError(t, err) + // The store requires KubeContent to be valid JSON; a YAML input must be + // normalized before being persisted, not stored verbatim. + assert.True(t, json.Valid(applied.Spec.KubeContent), + "KubeContent must be valid JSON even when the input manifest was YAML") + assert.Contains(t, string(applied.Spec.KubeContent), `"key":"value"`) +} + +func TestApplyResource_SkipsWhenGenerationUnchanged(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + result, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + assert.Equal(t, manifest.OperationSkip, result.Operation) +} + +func TestApplyResource_ReadDesirePairingFailureIsError(t *testing.T) { + ctx := context.Background() + store := &failingReadDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.Error(t, err, "a read-desire pairing failure must surface as an error, not a warning") + assert.Contains(t, err.Error(), "paired read desire") +} + +func TestApplyResource_RequiresTransportContext(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, nil) + require.Error(t, err) +} + +func TestApplyResource_EmptyManifestIsError(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + _, err := c.ApplyResource(ctx, nil, nil, testTransportContext()) + require.Error(t, err) +} + +// spyDeleteReadDesireStore counts DeleteReadDesire calls so tests can assert +// whether ensureReadDesire actually attempted a recreate. +type spyDeleteReadDesireStore struct { + desire.SpecStore + deleteReadDesireCalls int +} + +func (s *spyDeleteReadDesireStore) DeleteReadDesire( + ctx context.Context, id desire.Identity, owner string, version int64, +) error { + s.deleteReadDesireCalls++ + return s.SpecStore.DeleteReadDesire(ctx, id, owner, version) +} + +func TestEnsureReadDesire_SameTargetVersionIsNoOp(t *testing.T) { + ctx := context.Background() + store := &spyDeleteReadDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + + require.NoError(t, c.ensureReadDesire(ctx, readID, "v1")) + require.NoError(t, c.ensureReadDesire(ctx, readID, "v1")) + + assert.Zero(t, store.deleteReadDesireCalls, "matching target version must not trigger a recreate") + + rd, err := store.GetReadDesire(ctx, readID) + require.NoError(t, err) + assert.Equal(t, "v1", rd.TargetVersion) +} + +func TestEnsureReadDesire_RecreatesOnTargetVersionChange(t *testing.T) { + ctx := context.Background() + store := &spyDeleteReadDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + + require.NoError(t, c.ensureReadDesire(ctx, readID, "v1")) + require.NoError(t, c.ensureReadDesire(ctx, readID, "v2"), + "a target version change must delete and recreate the read desire, not surface an error") + + assert.Equal(t, 1, store.deleteReadDesireCalls, "changed target version must delete the stale read desire") + + rd, err := store.GetReadDesire(ctx, readID) + require.NoError(t, err) + assert.Equal(t, "v2", rd.TargetVersion, "read desire must be recreated with the new target version") +} + +// staleApplyVersionStore always reports version 1 for GetApplyDesire, +// regardless of the real stored version, so the caller's subsequent +// UpdateApplyDesireSpec is exercised against a genuinely stale CAS token. +type staleApplyVersionStore struct { + desire.SpecStore +} + +func (s *staleApplyVersionStore) GetApplyDesire(ctx context.Context, id desire.Identity) (desire.ApplyDesire, error) { + d, err := s.SpecStore.GetApplyDesire(ctx, id) + if err != nil { + return d, err + } + d.Version = 1 + return d, nil +} + +func TestApplyResource_VersionConflictSurfacesAsError(t *testing.T) { + ctx := context.Background() + memStore := newMemoryStore() + c := newTestClient(memStore) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + // Bump the real version out from under the (stale-reporting) view the + // client will see next, so its CAS write genuinely conflicts. + _, err = memStore.UpdateApplyDesireSpec( + ctx, applyID, desire.ApplySpec{KubeContent: configMapManifest(2)}, testOwner, 1) + require.NoError(t, err) + + staleClient := newTestClient(&staleApplyVersionStore{SpecStore: memStore}) + _, err = staleClient.ApplyResource(ctx, configMapManifest(3), nil, testTransportContext()) + require.Error(t, err) + assert.True(t, errors.Is(err, desire.ErrVersionConflict)) +} diff --git a/internal/desireclient/client.go b/internal/desireclient/client.go new file mode 100644 index 00000000..37356326 --- /dev/null +++ b/internal/desireclient/client.go @@ -0,0 +1,31 @@ +// Package desireclient implements transportclient.TransportClient against the +// desire-store contract from github.com/openshift-hyperfleet/hyperfleet-applier. +// It is the producer half of desire-based delivery: the adapter writes intent +// (apply/delete desires) and reads mirrored status (read desires) through the +// store, while a separate applier reconciles that intent against the target +// cluster. See docs/adapter-authoring-guide.md's "Desire transport" section +// for the contract this client implements. +package desireclient + +import ( + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/logger" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" +) + +// Client implements transportclient.TransportClient against a desire store. +type Client struct { + store desire.SpecStore + log logger.Logger + owner string +} + +// NewClient builds a desire-backed transport client. store is the producer +// surface of the desire contract (never StatusStore, which is the applier's +// own status-writing surface). owner identifies this adapter as the writer +// for ownership/CAS checks on the store. +func NewClient(store desire.SpecStore, owner string, log logger.Logger) *Client { + return &Client{store: store, owner: owner, log: log} +} + +var _ transportclient.TransportClient = (*Client)(nil) diff --git a/internal/desireclient/delete.go b/internal/desireclient/delete.go new file mode 100644 index 00000000..b21effd6 --- /dev/null +++ b/internal/desireclient/delete.go @@ -0,0 +1,59 @@ +package desireclient + +import ( + "context" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +// DeleteResource implements transportclient.TransportClient. It ensures a +// paired read desire exists, then posts a delete desire — CreateDeleteDesire +// atomically removes any sibling apply desire for the same target (see +// desire.SpecStore), so nothing re-applies. The read desire is created first +// so a failure there leaves the apply desire (if any) untouched — a safe, +// retryable state — rather than risking a delete desire with no read desire +// to observe its disappearance through, which a CreateDeleteDesire-then-read +// ordering could leave behind if the read-desire write failed afterward. +// +// opts is unused: propagation policy has no equivalent in the desire model +// (the applier owns deletion semantics against the target cluster), the same +// way maestroclient ignores it for ManifestWork deletes. +func (c *Client) DeleteResource( + ctx context.Context, + gvk schema.GroupVersionKind, + namespace, name string, + opts *transportclient.DeleteOptions, + target transportclient.TransportContext, +) error { + tc, err := resolveTransportContext(target) + if err != nil { + return err + } + + deleteID, err := buildIdentity(tc, desire.TypeDelete, gvk, namespace, name) + if err != nil { + return err + } + readID, err := buildIdentity(tc, desire.TypeRead, gvk, namespace, name) + if err != nil { + return err + } + + if err = c.ensureReadDesire(ctx, readID, gvk.Version); err != nil { + return fmt.Errorf( + "desireclient: failed to create paired read desire for %s/%s: %w", gvk.Kind, name, err) + } + + if _, err = c.store.CreateDeleteDesire(ctx, desire.DeleteDesire{ + Identity: deleteID, + Owner: c.owner, + }); err != nil && !errors.Is(err, desire.ErrAlreadyExists) { + return fmt.Errorf("desireclient: failed to create delete desire for %s/%s: %w", namespace, name, err) + } + + return nil +} diff --git a/internal/desireclient/delete_test.go b/internal/desireclient/delete_test.go new file mode 100644 index 00000000..367be54e --- /dev/null +++ b/internal/desireclient/delete_test.go @@ -0,0 +1,170 @@ +package desireclient + +import ( + "context" + "errors" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestDeleteResource_RemovesApplyKeepsRead(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + readID := applyID + readID.Type = desire.TypeRead + deleteID := applyID + deleteID.Type = desire.TypeDelete + + err = c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.NoError(t, err) + + _, err = store.GetApplyDesire(ctx, applyID) + assert.True(t, errors.Is(err, desire.ErrNotFound), "apply desire must be removed so nothing re-applies") + + _, err = store.GetDeleteDesire(ctx, deleteID) + require.NoError(t, err, "delete desire must be posted") + + _, err = store.GetReadDesire(ctx, readID) + require.NoError(t, err, "read desire must be left in place so disappearance stays observable") +} + +func TestDeleteResource_WithoutApply_PostsDelete(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + err := c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.NoError(t, err) + + deleteID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeDelete, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + _, err = store.GetDeleteDesire(ctx, deleteID) + require.NoError(t, err) +} + +func TestDeleteResource_WithoutApply_PairsRead(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + // No ApplyResource call first: this identity was never applied through + // this client, so no read desire was ever paired in. + err := c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.NoError(t, err) + + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + readDesire, err := store.GetReadDesire(ctx, readID) + require.NoError(t, err, "paired read desire must be auto-created even without a prior apply") + assert.Equal(t, "v1", readDesire.TargetVersion) + assert.Equal(t, testOwner, readDesire.Owner) +} + +func TestDeleteResource_ReadFailure_ReturnsError(t *testing.T) { + ctx := context.Background() + store := &failingReadDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + err := c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.Error(t, err, "a read-desire pairing failure must surface as an error, not be swallowed") + assert.Contains(t, err.Error(), "paired read desire") +} + +// TestDeleteResource_ReadFailure_LeavesApplyIntact verifies +// the failure-injection contract for the read-before-delete ordering: an +// existing apply desire must survive a read-desire pairing failure untouched, +// and no delete desire must have been created, since the apply->delete +// transition (CreateDeleteDesire) never runs if the paired read desire can't +// be ensured first. +func TestDeleteResource_ReadFailure_LeavesApplyIntact(t *testing.T) { + ctx := context.Background() + inner := newMemoryStore() + + // Set up the pre-existing apply desire through a working client first — + // ApplyResource itself pairs a read desire via the same store call the + // failing wrapper below targets, so it must succeed here. + c := newTestClient(inner) + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + deleteID := applyID + deleteID.Type = desire.TypeDelete + + failingClient := newTestClient(&failingReadDesireStore{SpecStore: inner}) + err = failingClient.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.Error(t, err, "a read-desire pairing failure must surface as an error") + + _, err = inner.GetApplyDesire(ctx, applyID) + assert.NoError(t, err, + "apply desire must be untouched: the apply->delete transition must not run before the read desire is ensured") + + _, err = inner.GetDeleteDesire(ctx, deleteID) + assert.True(t, errors.Is(err, desire.ErrNotFound), + "delete desire must not exist: CreateDeleteDesire must not run when read-desire pairing fails first") +} + +// TestDeleteResource_DeleteFailure_LeavesApplyAndReadIntact +// verifies the second failure-injection path: if CreateDeleteDesire fails +// after the read desire has already been ensured, the apply desire must +// still be present (CreateDeleteDesire never got a chance to atomically +// remove it) and the read desire must remain — a safe, retryable state +// rather than an orphaned delete desire with no observable disappearance. +func TestDeleteResource_DeleteFailure_LeavesApplyAndReadIntact(t *testing.T) { + ctx := context.Background() + store := &failingCreateDeleteDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + readID := applyID + readID.Type = desire.TypeRead + deleteID := applyID + deleteID.Type = desire.TypeDelete + + err = c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.Error(t, err, "a delete-desire creation failure must surface as an error") + assert.Contains(t, err.Error(), "failed to create delete desire") + + _, err = store.GetApplyDesire(ctx, applyID) + assert.NoError(t, err, + "apply desire must still exist: CreateDeleteDesire failed before it could atomically remove it") + + _, err = store.GetReadDesire(ctx, readID) + assert.NoError(t, err, "read desire must still exist: it was ensured before the failed transition") + + _, err = store.GetDeleteDesire(ctx, deleteID) + assert.True(t, errors.Is(err, desire.ErrNotFound), "delete desire must not exist since its creation failed") +} + +func TestDeleteResource_RequiresTransportContext(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + err := c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, nil) + require.Error(t, err) +} diff --git a/internal/desireclient/desireclient_test.go b/internal/desireclient/desireclient_test.go new file mode 100644 index 00000000..9bf1f56c --- /dev/null +++ b/internal/desireclient/desireclient_test.go @@ -0,0 +1,78 @@ +package desireclient + +import ( + "context" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/constants" + "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/logger" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire/store/memory" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +const ( + testManagementCluster = "mgmt-cluster-01" + testResource = "configmaps" + testNamespace = "default" + testName = "my-config" + testOwner = "hyperfleet-adapter" +) + +func newTestClient(store desire.SpecStore) *Client { + return NewClient(store, testOwner, logger.NewTestLogger()) +} + +func testTransportContext() *TransportContext { + return &TransportContext{ManagementCluster: testManagementCluster, Resource: testResource} +} + +func testGVK() schema.GroupVersionKind { + return schema.GroupVersionKind{Version: "v1", Kind: "ConfigMap"} +} + +// configMapManifest builds a minimal ConfigMap manifest carrying the +// hyperfleet.io/generation annotation task authors are required to set +// (docs/adapter-authoring-guide.md:498). +func configMapManifest(generation int64) []byte { + return fmt.Appendf(nil, `{ + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": { + "name": %q, + "namespace": %q, + "annotations": {%q: %q} + }, + "data": {"key": "value"} + }`, testName, testNamespace, constants.AnnotationGeneration, fmt.Sprint(generation)) +} + +// failingReadDesireStore wraps a real SpecStore but forces CreateReadDesire +// to fail, simulating a read-desire pairing failure. +type failingReadDesireStore struct { + desire.SpecStore +} + +func (f *failingReadDesireStore) CreateReadDesire( + _ context.Context, _ desire.ReadDesire, +) (desire.ReadDesire, error) { + return desire.ReadDesire{}, errors.New("boom: read desire store unavailable") +} + +// failingCreateDeleteDesireStore wraps a real SpecStore but forces +// CreateDeleteDesire to fail, simulating a failure in the apply-to-delete +// transition after the paired read desire has already been ensured. +type failingCreateDeleteDesireStore struct { + desire.SpecStore +} + +func (f *failingCreateDeleteDesireStore) CreateDeleteDesire( + _ context.Context, _ desire.DeleteDesire, +) (desire.DeleteDesire, error) { + return desire.DeleteDesire{}, errors.New("boom: delete desire store unavailable") +} + +func newMemoryStore() *memory.Store { + return memory.New() +} diff --git a/internal/desireclient/discover.go b/internal/desireclient/discover.go new file mode 100644 index 00000000..8bcf2d1c --- /dev/null +++ b/internal/desireclient/discover.go @@ -0,0 +1,71 @@ +package desireclient + +import ( + "context" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/logger" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +// DiscoverResources implements transportclient.TransportClient. It lists +// every read desire in the target partition and filters by GVK and the +// discovery criteria. desire.Identity carries no labels, so label-selector +// discovery must scan client-side, the same shape as +// maestroclient.DiscoverResources scanning ManifestWorks. +func (c *Client) DiscoverResources( + ctx context.Context, + gvk schema.GroupVersionKind, + discovery manifest.Discovery, + target transportclient.TransportContext, +) (*unstructured.UnstructuredList, error) { + tc, err := resolveTransportContext(target) + if err != nil { + return nil, err + } + + reads, err := c.store.ListReadDesires(ctx, tc.ManagementCluster) + if err != nil { + return nil, fmt.Errorf("desireclient: failed to list read desires for partition %q: %w", tc.ManagementCluster, err) + } + + list := &unstructured.UnstructuredList{} + for _, rd := range reads { + if rd.Identity.Group != gvk.Group || rd.Identity.Resource != tc.Resource { + continue + } + + // Route through the same three-way interpretation GetResource uses, so a + // single item's outcome here can never drift from what a Get on that same + // identity would report. + obj, err := c.decodeReadDesire(gvk, rd.Identity.Namespace, rd.Identity.Name, rd) + switch { + case errors.Is(err, ErrNotSyncedYet): + // Not yet synced — nothing to match discovery criteria against, + // same as a live List not yet showing a slow-to-create resource. + continue + case apierrors.IsNotFound(err): + // Confirmed absent — same as a live List not showing a deleted resource. + continue + case err != nil: + // A True condition with undecodable content is a store invariant + // violation for that one record, not legitimate transience — but it + // must not fail discovery for every other resource in the partition. + errCtx := logger.WithErrorField(ctx, err) + c.log.Errorf(errCtx, "desireclient: discovery skipping %s/%s: %v", + rd.Identity.Namespace, rd.Identity.Name, err) + continue + } + + if manifest.MatchesDiscoveryCriteria(obj, discovery) { + list.Items = append(list.Items, *obj) + } + } + + return list, nil +} diff --git a/internal/desireclient/discover_test.go b/internal/desireclient/discover_test.go new file mode 100644 index 00000000..a71328c5 --- /dev/null +++ b/internal/desireclient/discover_test.go @@ -0,0 +1,258 @@ +package desireclient + +import ( + "context" + "errors" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire/store/memory" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +// putSyncedReadDesire creates a read desire and marks it Successful=True with +// the given content, the shape DiscoverResources expects for a converged mirror. +func putSyncedReadDesire(t *testing.T, ctx context.Context, store *memory.Store, name string, content []byte) { + t.Helper() + id := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: name, + } + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{Identity: id, Owner: testOwner, TargetVersion: "v1"}) + require.NoError(t, err) + + _, err = store.UpdateReadDesireStatus(ctx, id, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonSynced, + }}}, + KubeContent: content, + }) + require.NoError(t, err) +} + +// failingListReadDesiresStore wraps a real SpecStore but forces +// ListReadDesires to fail, simulating a store-level outage during discovery. +type failingListReadDesiresStore struct { + desire.SpecStore +} + +func (f *failingListReadDesiresStore) ListReadDesires( + _ context.Context, _ string, +) ([]desire.ReadDesire, error) { + return nil, errors.New("boom: read desire store unavailable") +} + +func TestDiscoverResources_EmptyPartitionReturnsEmptyList(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items) +} + +func TestDiscoverResources_ReturnsSyncedResourceByName(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, testName, configMapManifest(1)) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{ByName: testName}, testTransportContext()) + require.NoError(t, err) + require.Len(t, list.Items, 1) + assert.Equal(t, testName, list.Items[0].GetName()) +} + +func TestDiscoverResources_ByNameExcludesNonMatchingName(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, testName, configMapManifest(1)) + + discovery := &manifest.DiscoveryConfig{ByName: "other-name"} + list, err := c.DiscoverResources(ctx, testGVK(), discovery, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items, "discovery criteria must filter out non-matching names") +} + +func TestDiscoverResources_LabelSelectorMatchesSubset(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + labeledManifest := []byte(`{ + "apiVersion": "v1", "kind": "ConfigMap", + "metadata": {"name": "labeled", "namespace": "default", "labels": {"app": "myapp"}} + }`) + unlabeledManifest := []byte(`{ + "apiVersion": "v1", "kind": "ConfigMap", + "metadata": {"name": "unlabeled", "namespace": "default"} + }`) + putSyncedReadDesire(t, ctx, store, "labeled", labeledManifest) + putSyncedReadDesire(t, ctx, store, "unlabeled", unlabeledManifest) + + list, err := c.DiscoverResources(ctx, testGVK(), + &manifest.DiscoveryConfig{LabelSelector: "app=myapp"}, testTransportContext()) + require.NoError(t, err) + require.Len(t, list.Items, 1) + assert.Equal(t, "labeled", list.Items[0].GetName()) +} + +func TestDiscoverResources_SkipsNotYetSyncedResource(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + id := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{Identity: id, Owner: testOwner, TargetVersion: "v1"}) + require.NoError(t, err) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items, "a read desire with no Successful condition yet must not appear in discovery") +} + +func TestDiscoverResources_SkipsFailedRead(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + id := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{Identity: id, Owner: testOwner, TargetVersion: "v1"}) + require.NoError(t, err) + _, err = store.UpdateReadDesireStatus(ctx, id, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonKubeAPIError, + }}}, + KubeContent: configMapManifest(1), + }) + require.NoError(t, err) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items, "Successful=False must never surface content, even stale content") +} + +func TestDiscoverResources_SkipsConfirmedAbsent(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, testName, nil) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items, "empty KubeContent under Successful=True is confirmed-absent, not a match") +} + +func TestDiscoverResources_SkipsUndecodableContentButKeepsOthers(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, "bad", []byte("not-json")) + putSyncedReadDesire(t, ctx, store, "good", configMapManifest(1)) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err, "a single bad record must not fail discovery for the whole partition") + require.Len(t, list.Items, 1) + assert.Equal(t, testName, list.Items[0].GetName()) +} + +func TestDiscoverResources_FiltersOutOtherResourceType(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, testName, configMapManifest(1)) + + otherContext := &TransportContext{ManagementCluster: testManagementCluster, Resource: "secrets"} + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, otherContext) + require.NoError(t, err) + assert.Empty(t, list.Items, "a read desire for a different plural resource must not appear") +} + +func TestDiscoverResources_FiltersOutOtherGroup(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + id := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Group: "apps", Resource: testResource, Namespace: testNamespace, Name: testName, + } + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{Identity: id, Owner: testOwner, TargetVersion: "v1"}) + require.NoError(t, err) + _, err = store.UpdateReadDesireStatus(ctx, id, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonSynced, + }}}, + KubeContent: configMapManifest(1), + }) + require.NoError(t, err) + + // testGVK() has an empty Group ("core"), so a read desire recorded under + // group "apps" must not match even though Resource matches. + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items, "a read desire in a different API group must not appear") +} + +func TestDiscoverResources_ScopedToPartition(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, testName, configMapManifest(1)) + + otherPartition := &TransportContext{ManagementCluster: "other-cluster", Resource: testResource} + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, otherPartition) + require.NoError(t, err) + assert.Empty(t, list.Items, "discovery must not leak resources across management-cluster partitions") +} + +func TestDiscoverResources_MultipleMatches(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + first := []byte(`{"apiVersion": "v1", "kind": "ConfigMap", "metadata": {"name": "first", "namespace": "default"}}`) + second := []byte(`{"apiVersion": "v1", "kind": "ConfigMap", "metadata": {"name": "second", "namespace": "default"}}`) + putSyncedReadDesire(t, ctx, store, "first", first) + putSyncedReadDesire(t, ctx, store, "second", second) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + require.Len(t, list.Items, 2) + names := []string{list.Items[0].GetName(), list.Items[1].GetName()} + assert.ElementsMatch(t, []string{"first", "second"}, names) +} + +func TestDiscoverResources_ListFailure_ReturnsError(t *testing.T) { + ctx := context.Background() + store := &failingListReadDesiresStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + _, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.Error(t, err) + assert.Contains(t, err.Error(), "failed to list read desires") +} + +func TestDiscoverResources_RequiresTransportContext(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + _, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, nil) + require.Error(t, err) +} diff --git a/internal/desireclient/get.go b/internal/desireclient/get.go new file mode 100644 index 00000000..0bdf2d85 --- /dev/null +++ b/internal/desireclient/get.go @@ -0,0 +1,95 @@ +package desireclient + +import ( + "context" + "encoding/json" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + apierrors "k8s.io/apimachinery/pkg/api/errors" + apimeta "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +// GetResource implements transportclient.TransportClient. It reads the +// mirrored live object from the read desire's status, distinguishing three +// outcomes per the eventual-consistency contract +// (docs/adapter-authoring-guide.md): not-synced-yet (ErrNotSyncedYet), +// confirmed-absent (apierrors.NewNotFound), and present (the mirrored +// object). +func (c *Client) GetResource( + ctx context.Context, + gvk schema.GroupVersionKind, + namespace, name string, + target transportclient.TransportContext, +) (*unstructured.Unstructured, error) { + tc, err := resolveTransportContext(target) + if err != nil { + return nil, err + } + + id, err := buildIdentity(tc, desire.TypeRead, gvk, namespace, name) + if err != nil { + return nil, err + } + + rd, err := c.store.GetReadDesire(ctx, id) + if errors.Is(err, desire.ErrNotFound) { + return nil, ErrNotSyncedYet + } + if err != nil { + return nil, fmt.Errorf("desireclient: failed to get read desire for %s/%s: %w", namespace, name, err) + } + + return c.decodeReadDesire(gvk, namespace, name, rd) +} + +// decodeReadDesire translates a ReadDesire's status conditions into the +// three-way outcome the eventual-consistency contract defines. Successful=True +// covers both a synced mirror and a confirmed-absent resource (the applier's +// notFound() writes Successful=True with empty KubeContent) — decodeKubeContent +// tells those apart. Successful=False means the read itself didn't succeed +// (KubeAPIError/PreCheckFailed); content is never relayed in that case, even +// if a stale mirror is sitting there from a prior sync, since the caller has +// no way to know it's stale. +func (c *Client) decodeReadDesire( + gvk schema.GroupVersionKind, namespace, name string, rd desire.ReadDesire, +) (*unstructured.Unstructured, error) { + cond := apimeta.FindStatusCondition(rd.Status.Conditions, desire.TypeSuccessful) + + switch { + case cond == nil: + // Read desire exists but the applier hasn't observed it yet. + return nil, ErrNotSyncedYet + + case cond.Status == metav1.ConditionTrue: + return decodeKubeContent(rd.Status.KubeContent, gvk, rd.Identity.Resource, namespace, name) + + case cond.Reason == desire.ReasonNotFound: + return nil, apierrors.NewNotFound(schema.GroupResource{Group: gvk.Group, Resource: rd.Identity.Resource}, name) + default: + return nil, ErrNotSyncedYet + } +} + +// decodeKubeContent decodes a synced read desire's mirrored content, or +// reports confirmed-absence when content is empty. Only ever called under +// Successful=True, where the applier's status-writers guarantee emptiness +// means Reason=NotFound and non-emptiness means Reason=Synced — so content +// length alone is an unambiguous stand-in for Reason here. +func decodeKubeContent( + kubeContent []byte, gvk schema.GroupVersionKind, resource, namespace, name string, +) (*unstructured.Unstructured, error) { + if len(kubeContent) == 0 { + return nil, apierrors.NewNotFound(schema.GroupResource{Group: gvk.Group, Resource: resource}, name) + } + obj := &unstructured.Unstructured{} + if err := json.Unmarshal(kubeContent, obj); err != nil { + return nil, fmt.Errorf("desireclient: failed to decode mirrored content for %s/%s: %w", namespace, name, err) + } + return obj, nil +} diff --git a/internal/desireclient/get_test.go b/internal/desireclient/get_test.go new file mode 100644 index 00000000..2158d232 --- /dev/null +++ b/internal/desireclient/get_test.go @@ -0,0 +1,215 @@ +package desireclient + +import ( + "context" + "errors" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +func readIdentity() desire.Identity { + return desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } +} + +func TestGetResource_AbsentReadDesireIsNotSyncedYet(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.True(t, errors.Is(err, ErrNotSyncedYet)) + assert.False(t, apierrors.IsNotFound(err), "absent mirror must not collapse into NotFound") +} + +func TestGetResource_ExistsButNotYetObservedIsNotSyncedYet(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: readIdentity(), Owner: testOwner, TargetVersion: "v1", + }) + require.NoError(t, err) + + _, err = c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.True(t, errors.Is(err, ErrNotSyncedYet), "no Successful condition yet must read as not-synced-yet") +} + +func TestGetResource_SyncedReturnsMirroredObject(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + id := readIdentity() + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{Identity: id, Owner: testOwner, TargetVersion: "v1"}) + require.NoError(t, err) + + _, err = store.UpdateReadDesireStatus(ctx, id, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonSynced, + }}}, + KubeContent: configMapManifest(1), + }) + require.NoError(t, err) + + obj, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.NoError(t, err) + assert.Equal(t, testName, obj.GetName()) + assert.Equal(t, testNamespace, obj.GetNamespace()) +} + +func TestGetResource_ConfirmedNotFound(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + id := readIdentity() + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{Identity: id, Owner: testOwner, TargetVersion: "v1"}) + require.NoError(t, err) + + _, err = store.UpdateReadDesireStatus(ctx, id, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonNotFound, + }}}, + }) + require.NoError(t, err) + + _, err = c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.True(t, apierrors.IsNotFound(err), "confirmed-gone must be a real NotFound, not ErrNotSyncedYet") + assert.False(t, errors.Is(err, ErrNotSyncedYet)) +} + +func TestGetResource_TransientErrorIsNotSyncedYetEvenWithStaleContent(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + id := readIdentity() + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{Identity: id, Owner: testOwner, TargetVersion: "v1"}) + require.NoError(t, err) + + _, err = store.UpdateReadDesireStatus(ctx, id, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonKubeAPIError, + }}}, + KubeContent: configMapManifest(1), + }) + require.NoError(t, err) + + _, err = c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.True(t, errors.Is(err, ErrNotSyncedYet), + "a failed read must never relay content, stale or otherwise — the caller can't tell it's stale") + assert.False(t, apierrors.IsNotFound(err)) +} + +func TestDecodeKubeContent(t *testing.T) { + tests := []struct { + name string + content []byte + wantNotFound bool + wantErr bool + }{ + {name: "empty content is confirmed absent", content: nil, wantNotFound: true}, + {name: "valid content decodes", content: configMapManifest(1)}, + {name: "invalid json is a decode error", content: []byte("not-json"), wantErr: true}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + obj, err := decodeKubeContent(tt.content, testGVK(), testResource, testNamespace, testName) + switch { + case tt.wantNotFound: + require.Error(t, err) + assert.True(t, apierrors.IsNotFound(err)) + case tt.wantErr: + require.Error(t, err) + assert.False(t, apierrors.IsNotFound(err), "a decode failure is not the same outcome as confirmed-absence") + default: + require.NoError(t, err) + assert.Equal(t, testName, obj.GetName()) + assert.Equal(t, testNamespace, obj.GetNamespace()) + } + }) + } +} + +// successfulCondition builds the single summary condition every desire +// carries, shortening the table below to status+reason. +func successfulCondition(status metav1.ConditionStatus, reason string) *metav1.Condition { + return &metav1.Condition{Type: desire.TypeSuccessful, Status: status, Reason: reason} +} + +func TestDecodeReadDesire(t *testing.T) { + tests := []struct { + name string + condition *metav1.Condition + content []byte + wantNotSynced bool + wantNotFound bool + wantContent bool + }{ + { + name: "no condition yet is not synced", + condition: nil, + wantNotSynced: true, + }, + { + name: "successful true decodes content", + condition: successfulCondition(metav1.ConditionTrue, desire.ReasonSynced), + content: configMapManifest(1), + wantContent: true, + }, + { + name: "successful true with empty content is confirmed absent", + condition: successfulCondition(metav1.ConditionTrue, desire.ReasonNotFound), + wantNotFound: true, + }, + { + name: "false with notfound reason is confirmed absent", + condition: successfulCondition(metav1.ConditionFalse, desire.ReasonNotFound), + wantNotFound: true, + }, + { + name: "false with other reason is not synced yet, even with stale content", + condition: successfulCondition(metav1.ConditionFalse, desire.ReasonKubeAPIError), + content: configMapManifest(1), + wantNotSynced: true, + }, + } + + c := newTestClient(newMemoryStore()) + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + rd := desire.ReadDesire{Identity: readIdentity()} + if tt.condition != nil { + rd.Status = desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{*tt.condition}}, + KubeContent: tt.content, + } + } + + obj, err := c.decodeReadDesire(testGVK(), testNamespace, testName, rd) + switch { + case tt.wantNotSynced: + require.Error(t, err) + assert.True(t, errors.Is(err, ErrNotSyncedYet)) + case tt.wantNotFound: + require.Error(t, err) + assert.True(t, apierrors.IsNotFound(err)) + case tt.wantContent: + require.NoError(t, err) + assert.Equal(t, testName, obj.GetName()) + } + }) + } +} diff --git a/internal/desireclient/types.go b/internal/desireclient/types.go new file mode 100644 index 00000000..e8f9bfe1 --- /dev/null +++ b/internal/desireclient/types.go @@ -0,0 +1,104 @@ +package desireclient + +import ( + "encoding/json" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "sigs.k8s.io/yaml" +) + +// TransportContext carries per-request routing information for the desire +// transport backend. Pass this as the TransportContext (any) in ApplyResource, +// GetResource, DiscoverResources, and DeleteResource. +type TransportContext struct { + // ManagementCluster is the target managed cluster identifier — the desire + // store's partition key. Required for all operations. + ManagementCluster string + // Resource is the plural Kubernetes resource type (e.g. "networkpolicies"), + // declared in the task config alongside the manifest. The desire store's + // Identity is keyed by this plural form directly (see + // pkg/desire.Identity.Resource in hyperfleet-applier); the adapter never + // derives it from the manifest's Kind — no RESTMapper is used or needed. + // Required for all operations. + Resource string +} + +// ErrNotSyncedYet indicates the read desire's mirror has not yet been +// populated by the applier — a transient, non-terminal outcome distinct from +// a confirmed-absent resource (reported via apierrors.NewNotFound). Per the +// eventual-consistency contract (docs/adapter-authoring-guide.md), callers +// should treat this as "not converged yet" rather than a failure. Use +// errors.Is to check for it. +var ErrNotSyncedYet = errors.New("desireclient: resource not synced yet") + +// resolveTransportContext type-asserts the generic TransportContext and +// validates both fields are set. +func resolveTransportContext(target transportclient.TransportContext) (*TransportContext, error) { + tc, ok := target.(*TransportContext) + if !ok || tc == nil { + return nil, fmt.Errorf("desireclient: TransportContext with ManagementCluster and Resource is required") + } + if tc.ManagementCluster == "" { + return nil, fmt.Errorf("desireclient: TransportContext.ManagementCluster is required") + } + if tc.Resource == "" { + return nil, fmt.Errorf("desireclient: TransportContext.Resource is required") + } + return tc, nil +} + +// buildIdentity is the single shared helper for assembling a desire.Identity +// from a resolved TransportContext, GVK, namespace, and name. All four +// TransportClient methods route through this rather than duplicating identity +// construction. +func buildIdentity( + tc *TransportContext, dtype desire.DesireType, gvk schema.GroupVersionKind, namespace, name string, +) (desire.Identity, error) { + id := desire.Identity{ + ManagementCluster: tc.ManagementCluster, + Type: dtype, + Group: gvk.Group, + Resource: tc.Resource, + Namespace: namespace, + Name: name, + } + if err := id.Validate(); err != nil { + return desire.Identity{}, fmt.Errorf("desireclient: invalid identity: %w", err) + } + return id, nil +} + +// generationFromKubeContent extracts the hyperfleet.io/generation annotation +// from a previously-stored ApplyDesire's KubeContent. Returns 0 if the content +// can't be parsed or carries no annotation. +func generationFromKubeContent(kubeContent []byte) int64 { + obj, err := parseToUnstructured(kubeContent) + if err != nil { + return 0 + } + return manifest.GetGenerationFromUnstructured(obj) +} + +// parseToUnstructured parses JSON or YAML bytes into an unstructured resource, +// mirroring the pattern in internal/k8sclient/apply.go and +// internal/maestroclient/client.go. +func parseToUnstructured(data []byte) (*unstructured.Unstructured, error) { + obj := &unstructured.Unstructured{} + if err := json.Unmarshal(data, &obj.Object); err == nil && obj.Object != nil { + return obj, nil + } + jsonData, err := yaml.YAMLToJSON(data) + if err != nil { + return nil, fmt.Errorf("failed to convert YAML to JSON: %w", err) + } + if err := json.Unmarshal(jsonData, &obj.Object); err != nil { + return nil, fmt.Errorf("failed to parse manifest: %w", err) + } + return obj, nil +} diff --git a/internal/desireclient/types_test.go b/internal/desireclient/types_test.go new file mode 100644 index 00000000..c58732eb --- /dev/null +++ b/internal/desireclient/types_test.go @@ -0,0 +1,186 @@ +package desireclient + +import ( + "fmt" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/constants" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// otherTransportContext stands in for a different transport's context type +// (e.g. *maestroclient.TransportContext) to exercise the type-assertion branch +// of resolveTransportContext without importing another transport package. +type otherTransportContext struct{} + +func TestResolveTransportContext_Invalid(t *testing.T) { + tests := []struct { + name string + target any + wantErr string + }{ + {name: "nil", target: nil, wantErr: "TransportContext"}, + {name: "wrong type", target: &otherTransportContext{}, wantErr: "TransportContext"}, + { + name: "missing management cluster", + target: &TransportContext{Resource: testResource}, + wantErr: "ManagementCluster", + }, + { + name: "missing resource", + target: &TransportContext{ManagementCluster: testManagementCluster}, + wantErr: "Resource", + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + _, err := resolveTransportContext(tt.target) + require.Error(t, err) + assert.Contains(t, err.Error(), tt.wantErr) + }) + } +} + +func TestResolveTransportContext_Valid(t *testing.T) { + in := &TransportContext{ManagementCluster: testManagementCluster, Resource: testResource} + tc, err := resolveTransportContext(in) + require.NoError(t, err) + assert.Same(t, in, tc) +} + +func TestBuildIdentity(t *testing.T) { + tests := []struct { + name string + namespace string + resName string + wantErr string + }{ + {name: "valid", namespace: testNamespace, resName: testName}, + { + name: "invalid namespace is rejected", + namespace: "Invalid_Namespace", + resName: testName, + wantErr: "invalid identity", + }, + { + name: "invalid name is rejected", + namespace: testNamespace, + resName: "Invalid_Name", + wantErr: "invalid identity", + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + id, err := buildIdentity(testTransportContext(), desire.TypeRead, testGVK(), tt.namespace, tt.resName) + if tt.wantErr != "" { + require.Error(t, err) + assert.Contains(t, err.Error(), tt.wantErr) + return + } + require.NoError(t, err) + assert.Equal(t, desire.Identity{ + ManagementCluster: testManagementCluster, + Type: desire.TypeRead, + Group: testGVK().Group, + Resource: testResource, + Namespace: tt.namespace, + Name: tt.resName, + }, id) + }) + } +} + +func TestParseToUnstructured(t *testing.T) { + tests := []struct { + name string + data []byte + wantErr bool + wantNil bool + }{ + {name: "valid JSON", data: configMapManifest(1)}, + { + name: "valid YAML", + data: []byte("apiVersion: v1\nkind: ConfigMap\nmetadata:\n name: my-config\n namespace: default\n"), + }, + {name: "garbage is neither valid JSON nor YAML", data: []byte("{not valid: [json or yaml"), wantErr: true}, + { + // json.Unmarshal("null", &obj.Object) succeeds with a nil map instead + // of erroring, and the YAML fallback degrades the same way — so a + // literal "null" manifest parses successfully into an empty object + // rather than failing loudly. Pinning this down since it's the one + // input where the two-stage JSON-then-YAML parse doesn't behave like + // either "valid" or "invalid". + name: "literal null parses to an empty object rather than erroring", + data: []byte("null"), + wantNil: true, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + obj, err := parseToUnstructured(tt.data) + if tt.wantErr { + require.Error(t, err) + return + } + require.NoError(t, err) + if tt.wantNil { + assert.Nil(t, obj.Object) + return + } + assert.Equal(t, "ConfigMap", obj.GetKind()) + assert.Equal(t, testName, obj.GetName()) + }) + } +} + +// TestParseToUnstructured_JSONAndYAMLProduceEquivalentObject verifies the +// property the store's JSON-only contract depends on: a JSON manifest and its +// YAML equivalent must parse to the identical unstructured object, so +// ApplyResource's re-marshal-to-JSON step doesn't silently diverge in content +// depending on which format the caller happened to render. +func TestParseToUnstructured_JSONAndYAMLProduceEquivalentObject(t *testing.T) { + jsonManifest := configMapManifest(1) + yamlManifest := fmt.Appendf(nil, ` +apiVersion: v1 +kind: ConfigMap +metadata: + name: %s + namespace: %s + annotations: + %s: "1" +data: + key: value +`, testName, testNamespace, constants.AnnotationGeneration) + + jsonObj, err := parseToUnstructured(jsonManifest) + require.NoError(t, err) + yamlObj, err := parseToUnstructured(yamlManifest) + require.NoError(t, err) + + assert.Equal(t, jsonObj.Object, yamlObj.Object, + "JSON and YAML encodings of the same manifest must parse to the same object") +} + +func TestGenerationFromKubeContent(t *testing.T) { + tests := []struct { + name string + content []byte + want int64 + }{ + {name: "valid content with generation annotation", content: configMapManifest(5), want: 5}, + { + name: "valid content without generation annotation", + content: []byte(`{"apiVersion":"v1","kind":"ConfigMap","metadata":{"name":"x","namespace":"default"}}`), + want: 0, + }, + {name: "unparseable content returns zero rather than erroring", content: []byte("not-json"), want: 0}, + {name: "empty content returns zero rather than erroring", content: nil, want: 0}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, generationFromKubeContent(tt.content)) + }) + } +} 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..bc70e0a5 --- /dev/null +++ b/internal/transportregistry/registry.go @@ -0,0 +1,260 @@ +// Package transportregistry builds the configured transport clients and owns +// the resources whose lifetimes they require. +package transportregistry + +import ( + "context" + "fmt" + "io" + "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-adapter/pkg/logger" + "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, log logger.Logger) (*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, log); err != nil { + closeAfterBuildFailure(ctx, runtime, log) + return nil, err + } + return runtime, nil + } + + stores, err := runtime.buildStores(ctx, config.Stores, log) + if err != nil { + closeAfterBuildFailure(ctx, runtime, log) + return nil, err + } + + for _, name := range sortedTransportNames(config.Transports) { + definition := config.Transports[name] + client, err := buildTransport(ctx, config, definition, stores, log) + if err != nil { + closeAfterBuildFailure(ctx, runtime, log) + return nil, fmt.Errorf("build transport %q: %w", name, err) + } + runtime.Registry[name] = client + } + + return runtime, nil +} + +func closeAfterBuildFailure(ctx context.Context, runtime *Runtime, log logger.Logger) { + if err := runtime.Close(); err != nil { + log.Warnf( + logger.WithErrorField(ctx, err), + "Failed to close transport registry after build failure", + ) + } +} + +// 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 + } + // The executor still needs its legacy singleton key even when the deployment + // uses named entries. Mirror the production Maestro-first choice here, but + // without constructing any network clients. + if config.Clients.Maestro != nil { + runtime.Registry[configloader.TransportClientMaestro] = client + } else { + runtime.Registry[configloader.TransportClientKubernetes] = client + } + return runtime, nil +} + +// 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, + log logger.Logger, +) 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, log) + 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, log) + 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, + log logger.Logger, +) (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: parse Redis URL: %w", name, err) + } + 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 { + log.Warnf( + logger.WithErrorField(ctx, closeErr), + "Failed to close Redis client after ping failure", + ) + } + 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, + log logger.Logger, +) (transportclient.TransportClient, error) { + switch definition.Type { + case configloader.TransportTypeKubernetes: + return buildKubernetes(ctx, config.Clients.Kubernetes, log) + 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, log), nil + default: + return nil, fmt.Errorf("unsupported type %q", definition.Type) + } +} + +func buildKubernetes( + ctx context.Context, + config configloader.KubernetesConfig, + log logger.Logger, +) (*k8sclient.Client, error) { + return k8sclient.NewClient(ctx, k8sclient.ClientConfig{ + KubeConfigPath: config.KubeConfigPath, + QPS: config.QPS, + Burst: config.Burst, + }, log) +} + +func buildMaestro( + ctx context.Context, + config *configloader.MaestroClientConfig, + log logger.Logger, +) (*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, log) +} + +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..d98f94af --- /dev/null +++ b/internal/transportregistry/runtime_test.go @@ -0,0 +1,399 @@ +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/openshift-hyperfleet/hyperfleet-adapter/pkg/logger" + "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, logger.NewTestLogger()) + 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) + + 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 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, logger.NewTestLogger()) + require.Error(t, err) + assert.Nil(t, runtime) + assert.Contains(t, err.Error(), "unreachable") +} + +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, logger.NewTestLogger()) + 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, logger.NewTestLogger()) + + 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, + logger.NewTestLogger(), + ) + + 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, logger.NewTestLogger()) + + 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, logger.NewTestLogger()) + 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 +}