multicluster

package
v0.0.0-...-6a36b9b Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 9, 2026 License: Apache-2.0 Imports: 29 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var DefaultProbeOptions = ProbeOptions{
	Interval:         10 * time.Second,
	Timeout:          5 * time.Second,
	FailureThreshold: 3,
}

DefaultProbeOptions probes every 10s with a 5s per-probe timeout and treats a remote as lost after 3 consecutive failures (~30s of sustained unreachability).

View Source
var DefaultResourceRouters = map[schema.GroupVersionKind]ResourceRouter{
	{Group: "kvm.cloud.sap", Version: "v1", Kind: "Hypervisor"}:               HypervisorResourceRouter{},
	{Group: "cortex.cloud", Version: "v1alpha1", Kind: "Reservation"}:         ReservationsResourceRouter{},
	{Group: "cortex.cloud", Version: "v1alpha1", Kind: "History"}:             HistoryResourceRouter{},
	{Group: "cortex.cloud", Version: "v1alpha1", Kind: "CommittedResource"}:   CommittedResourceRouter{},
	{Group: "cortex.cloud", Version: "v1alpha1", Kind: "ProjectQuota"}:        ProjectQuotaResourceRouter{},
	{Group: "cortex.cloud", Version: "v1alpha1", Kind: "FlavorGroupCapacity"}: FlavorGroupCapacityResourceRouter{},
}

DefaultResourceRouters defines all mappings of GroupVersionKinds to RRs for the multicluster client that cortex supports by default. This is used to route resources to the correct cluster in a multicluster setup.

Functions

func IsDuplicateError

func IsDuplicateError(err error) bool

IsDuplicateError returns true if the error indicates that a resource was found in multiple clusters. This can be used by callers of the Get and List methods to keep using the result even if a duplicate exists, as long as they don't mind that the result is potentially inconsistent.

func IsNoClusterMatchedError

func IsNoClusterMatchedError(err error) bool

IsNoClusterMatchedError returns true if the error indicates that no configured cluster matched the resource for a write operation. Callers can use this to skip resources targeting unavailable AZs gracefully.

Types

type APIServersConfig

type APIServersConfig struct {
	// Resources managed in the cluster where cortex is deployed.
	Home HomeConfig `json:"home"`
	// Resources managed in remote clusters.
	Remotes []RemoteConfig `json:"remotes,omitempty"`
}

APIServersConfig separates resources into home and remote clusters.

type Client

type Client struct {
	// ResourceRouters determine which cluster a resource should be written to
	// when multiple clusters serve the same GVK.
	ResourceRouters map[schema.GroupVersionKind]ResourceRouter

	// Wrappers are applied to every cluster (home and remotes) during InitFromConf.
	// Each wrapper transforms the cluster before it is stored for routing.
	// Applied in slice order; the raw inner cluster is always added to the manager
	// so its informers start independently of wrapping.
	Wrappers []ClusterWrapper

	// The cluster in which cortex is deployed.
	HomeCluster cluster.Cluster
	// The REST config for the home cluster in which cortex is deployed.
	HomeRestConfig *rest.Config
	// The scheme for the home cluster in which cortex is deployed.
	// This scheme should include all types used in the remote clusters.
	HomeScheme *runtime.Scheme

	// Optional monitor for Prometheus metrics. A nil Monitor causes recording
	// to be skipped, so the client can be used without wiring metrics.
	Monitor Monitor
	// contains filtered or unexported fields
}

func (*Client) AddRemote

func (c *Client) AddRemote(ctx context.Context, host, caCert string, insecureSkipTLSVerify bool, labels map[string]string, gvks ...schema.GroupVersionKind) (cluster.Cluster, error)

Add a remote cluster which uses the same REST config as the home cluster, but a different host, for the given resource gvks.

If insecureSkipTLSVerify is true, the remote apiserver's TLS certificate is not verified and caCert is ignored. This is useful for apiservers whose CA certificate rotates frequently and does not chain to a stable root.

This can be used when the remote cluster accepts the home cluster's service account tokens. See the kubernetes documentation on structured auth to learn more about jwt-based authentication across clusters. AddRemote returns the raw inner cluster.Cluster (which the caller must add to the manager so its informers/caches start). Each registered Wrapper is responsible for adding its own lifecycle Runnables to mgr directly. The wrapped cluster is stored in remoteClusters so all routing goes through any per-cluster wrapper.

func (*Client) Apply

Apply is not supported in the multicluster client as the group version kind cannot be inferred from the ApplyConfiguration.

func (*Client) ClustersForGVK

func (c *Client) ClustersForGVK(gvk schema.GroupVersionKind) ([]cluster.Cluster, error)

ClustersForGVK returns all clusters that serve the given GVK. The GVK must be explicitly configured in either homeGVKs or remoteClusters. Returns an error if the GVK is unknown.

func (*Client) ConfiguredRouteLabels

func (c *Client) ConfiguredRouteLabels(gvk schema.GroupVersionKind) []map[string]string

ConfiguredRouteLabels returns the routing label sets of all configured remote clusters for the given GVK. This can be used to determine which availability zones (or other routing dimensions) are served. Returns nil if the GVK is only configured for the home cluster.

func (*Client) Create

func (c *Client) Create(ctx context.Context, obj client.Object, opts ...client.CreateOption) error

Create routes the object to the matching cluster using the ResourceRouter and performs a Create operation.

Before writing, it performs a best-effort Get against the other clusters serving the same GVK to detect a cross-cluster name collision. If the object name already exists on another cluster, a duplicateError is returned (checkable with IsDuplicateError) and no create is performed. Non-NotFound errors from the probe clusters are logged and ignored so that a single unavailable cluster does not block writes.

func (*Client) Delete

func (c *Client) Delete(ctx context.Context, obj client.Object, opts ...client.DeleteOption) error

Delete routes the object to the matching cluster using the ResourceRouter and performs a Delete operation.

func (*Client) DeleteAllOf

func (c *Client) DeleteAllOf(ctx context.Context, obj client.Object, opts ...client.DeleteAllOfOption) error

DeleteAllOf iterates over all clusters with the GVK and performs DeleteAllOf on each.

func (*Client) GVKFromHomeScheme

func (c *Client) GVKFromHomeScheme(obj runtime.Object) (gvk schema.GroupVersionKind, err error)

Get the gvk registered for the given resource in the home cluster's scheme.

func (*Client) Get

func (c *Client) Get(ctx context.Context, key client.ObjectKey, obj client.Object, opts ...client.GetOption) error

Get iterates over all clusters with the GVK and returns the result.

If the requested resource is encountered in multiple clusters, this function will return the first one, but will set an error message that can be checked with IsDuplicateError. In that way the result can be used if the caller just cares about the resource existing in at least one cluster, and doesn't mind which one is returned.

If no cluster has the resource, a NotFound error is returned.

Non-NotFound errors from individual clusters are logged and silently skipped so that a single unavailable cluster does not block the entire read path.

func (*Client) GetEventRecorder

func (c *Client) GetEventRecorder(name string) recorder.EventRecorder

GetEventRecorder creates a multi-cluster-aware EventRecorder. It pre-creates a per-cluster recorder for the home cluster and every remote cluster currently registered in the client. The name parameter is passed through to each cluster's GetEventRecorder method (it becomes the reportingController in the Kubernetes Event).

func (*Client) GroupVersionKindFor

func (c *Client) GroupVersionKindFor(obj runtime.Object) (schema.GroupVersionKind, error)

Return the GroupVersionKind for the given object using the home cluster's RESTMapper.

func (*Client) IndexField

func (c *Client) IndexField(ctx context.Context, obj client.Object, list client.ObjectList, field string, extractValue client.IndexerFunc) error

Index a field for a resource in all matching cluster caches. Usually, you want to index the same field in both the object and list type, as both would be mapped to individual clients based on their GVK.

Per-cluster errors are logged and skipped — a temporarily-unreachable cluster does not prevent index setup for the other clusters. This mirrors the List error-handling pattern: the affected cluster's objects will be silently absent from MatchingFields queries until the index is established (e.g. after restart).

func (*Client) InitFromConf

func (c *Client) InitFromConf(ctx context.Context, mgr ctrl.Manager, conf ClientConfig) error

Helper function to initialize a new multicluster client during service startup, using the conf module provided by cortex.

func (*Client) IsObjectNamespaced

func (c *Client) IsObjectNamespaced(obj runtime.Object) (bool, error)

Return true if the GroupVersionKind of the object is namespaced using the home cluster's RESTMapper.

func (*Client) List

func (c *Client) List(ctx context.Context, list client.ObjectList, opts ...client.ListOption) error

List iterates over all clusters with the GVK and returns a combined list containing all resources found in any cluster.

If resources are encountered in multiple clusters with the same namespace/name, this function will still return a combined list of all resources, but will set an error message that can be checked with IsDuplicateError. In that way the result can be used if duplicates are ok and disambiguated by the caller.

Errors from individual clusters are logged and silently skipped so that a single unavailable cluster does not block the entire read path.

func (*Client) ListMetadataPerCluster

func (c *Client) ListMetadataPerCluster(ctx context.Context, gvk schema.GroupVersionKind, opts ...client.ListOption) ([]ClusterObjectMetadata, error)

ListMetadataPerCluster returns the object metadata of the given GVK for each configured cluster. It uses PartialObjectMetadataList so only object metadata crosses the wire — no spec or status — making it efficient even for large object counts. Callers that only need counts can use len(Items). Clusters that return an error are logged and skipped (same policy as List). The home cluster is included with IsHome set to true.

func (*Client) Patch

func (c *Client) Patch(ctx context.Context, obj client.Object, patch client.Patch, opts ...client.PatchOption) error

Patch routes the object to the matching cluster using the ResourceRouter and performs a Patch operation.

func (*Client) ProbeRemotes

func (c *Client) ProbeRemotes(ctx context.Context, opts ProbeOptions, onLost func(host string))

ProbeRemotes runs one reachability probe goroutine per unique remote apiserver until ctx is cancelled. Each goroutine periodically probes its remote, updates the reachability gauge on the Monitor (if any), and tracks consecutive failures; when a remote has been unreachable FailureThreshold times in a row it calls onLost(host) once and stops probing that remote (the caller is expected to tear down the manager cycle in response).

A single transient blip does not trigger onLost: the failure counter resets on the first reachable probe. "Unreachable" means a transport-level failure (connection refused, no route, TLS/handshake failure, timeout) — an apiserver that responds with any HTTP status (including 401/403/5xx) is considered reachable, since a manager rebuild cannot fix an authz/server error and would only cause a restart storm.

ProbeRemotes returns immediately when no remotes are configured.

func (*Client) RESTMapper

func (c *Client) RESTMapper() meta.RESTMapper

Return the RESTMapper of the home cluster.

func (*Client) Scheme

func (c *Client) Scheme() *runtime.Scheme

Return the scheme of the home cluster.

func (*Client) Status

func (c *Client) Status() client.StatusWriter

Provide a wrapper around the status subresource client which picks the right cluster based on the resource type.

func (*Client) SubResource

func (c *Client) SubResource(subResource string) client.SubResourceClient

Provide a wrapper around the given subresource client which picks the right cluster based on the resource type.

func (*Client) UniqueRemotes

func (c *Client) UniqueRemotes() []RemoteEndpoint

UniqueRemotes returns each distinct remote cluster exactly once, across all configured GVKs. The same remote apiserver is stored once per GVK it serves, so this dedupes by cluster identity (the same policy IndexField uses to dedupe caches). The home cluster is not included.

func (*Client) Update

func (c *Client) Update(ctx context.Context, obj client.Object, opts ...client.UpdateOption) error

Update routes the object to the matching cluster using the ResourceRouter and performs an Update operation.

type ClientConfig

type ClientConfig struct {
	// Apiserver configuration mapping GVKs to home or remote clusters.
	// Every GVK used through the multicluster client must be listed
	// in either Home or Remotes. Unknown GVKs will cause an error.
	APIServers APIServersConfig `json:"apiservers"`
}

type ClusterObjectMetadata

type ClusterObjectMetadata struct {
	Labels map[string]string
	Items  []metav1.PartialObjectMetadata
	IsHome bool
}

ClusterObjectMetadata is one entry in the result of ListMetadataPerCluster. Labels holds the routing labels for the cluster. Items holds the object metadata returned by the cluster (no spec or status). IsHome is true for the home cluster, which has no routing labels.

type ClusterWrapper

type ClusterWrapper interface {
	WrapCluster(cl cluster.Cluster) (cluster.Cluster, error)
}

ClusterWrapper can be registered on the Client to transparently transform each cluster before it is stored for routing. WrapCluster receives the current cluster: raw for the first wrapper in the chain, or the previous wrapper's result for subsequent ones. It must return the (possibly wrapped) cluster to use for routing and is responsible for registering any lifecycle Runnables with the manager itself. Because the manager is not passed to WrapCluster, a wrapper that needs it must capture it at construction time (e.g. by storing it on the wrapper). Wrappers are applied in order for every home and remote cluster during InitFromConf. For remote clusters the original unwrapped cluster is always added to the manager separately so its informers start independently of wrapping.

type CommittedResourceRouter

type CommittedResourceRouter struct{}

CommittedResourceRouter routes committed resources to clusters based on availability zone.

func (CommittedResourceRouter) ExtractClusterSelector

func (c CommittedResourceRouter) ExtractClusterSelector(obj any) (string, error)

func (CommittedResourceRouter) Match

func (c CommittedResourceRouter) Match(obj any, labels map[string]string) (bool, error)

type FlavorGroupCapacityResourceRouter

type FlavorGroupCapacityResourceRouter struct{}

FlavorGroupCapacityResourceRouter routes flavor group capacity CRDs to clusters based on availability zone.

func (FlavorGroupCapacityResourceRouter) ExtractClusterSelector

func (f FlavorGroupCapacityResourceRouter) ExtractClusterSelector(obj any) (string, error)

func (FlavorGroupCapacityResourceRouter) Match

func (f FlavorGroupCapacityResourceRouter) Match(obj any, labels map[string]string) (bool, error)

type HistoryResourceRouter

type HistoryResourceRouter struct{}

HistoryResourceRouter routes histories to clusters based on availability zone.

func (HistoryResourceRouter) ExtractClusterSelector

func (h HistoryResourceRouter) ExtractClusterSelector(obj any) (string, error)

func (HistoryResourceRouter) Match

func (h HistoryResourceRouter) Match(obj any, labels map[string]string) (bool, error)

type HomeConfig

type HomeConfig struct {
	// The resource GVKs formatted as "<group>/<version>/<Kind>".
	GVKs []string `json:"gvks"`
}

HomeConfig lists GVKs that are managed in the home cluster.

type HypervisorResourceRouter

type HypervisorResourceRouter struct{}

HypervisorResourceRouter routes hypervisors to clusters based on availability zone.

func (HypervisorResourceRouter) ExtractClusterSelector

func (h HypervisorResourceRouter) ExtractClusterSelector(obj any) (string, error)

func (HypervisorResourceRouter) Match

func (h HypervisorResourceRouter) Match(obj any, labels map[string]string) (bool, error)

type Monitor

type Monitor interface {
	prometheus.Collector
	// contains filtered or unexported methods
}

Monitor is the metrics sink for the multicluster client. It is optional on the Client: a nil Monitor causes recording to be skipped entirely. It embeds prometheus.Collector so a concrete implementation can be registered with a Prometheus registry.

func NewMonitor

func NewMonitor(prefix string) Monitor

NewMonitor creates a new Prometheus-backed multicluster client monitor. The prefix is prepended to every metric name (e.g. pass "cortex_" to produce "cortex_multicluster_cross_cluster_name_conflicts_total").

type MultiClusterRecorder

type MultiClusterRecorder struct {
	// contains filtered or unexported fields
}

MultiClusterRecorder implements recorder.EventRecorder and routes events to the correct cluster based on the GVK of the "regarding" object. It uses the same routing logic as the multicluster Client's write path.

func (*MultiClusterRecorder) AnnotatedEventf

func (r *MultiClusterRecorder) AnnotatedEventf(regarding, related runtime.Object, annotations map[string]string, eventtype, reason, action, note string, args ...any)

AnnotatedEventf routes the annotated event to the cluster that owns the "regarding" object. Falls back to the home cluster recorder if routing fails.

func (*MultiClusterRecorder) Eventf

func (r *MultiClusterRecorder) Eventf(regarding, related runtime.Object, eventtype, reason, action, note string, args ...any)

Eventf routes the event to the cluster that owns the "regarding" object. Falls back to the home cluster recorder if routing fails.

type MulticlusterBuilder

type MulticlusterBuilder struct {
	// Wrapped builder provided by controller-runtime.
	*builder.Builder
	// contains filtered or unexported fields
}

Builder which provides special methods to watch resources across multiple clusters.

func BuildController

func BuildController(c *Client, mgr manager.Manager) MulticlusterBuilder

Build a multicluster controller using the multicluster client. Use this builder to watch resources across multiple clusters.

func (MulticlusterBuilder) WatchesMulticluster

func (b MulticlusterBuilder) WatchesMulticluster(object client.Object, eventHandler handler.TypedEventHandler[client.Object, reconcile.Request], predicates ...predicate.Predicate) (MulticlusterBuilder, error)

WatchesMulticluster watches a resource across all clusters that serve its GVK. If the GVK is served by multiple remote clusters, a watch is set up on each. Returns an error if the GVK is not configured in any cluster.

type NoClusterMatchedError

type NoClusterMatchedError struct {
	GVK             schema.GroupVersionKind
	ClusterSelector string
}

NoClusterMatchedError is returned when a write operation cannot find a matching remote cluster for the given GVK and resource. This typically happens in multi-AZ setups where the resource targets an AZ that has no configured cluster.

func (*NoClusterMatchedError) Error

func (e *NoClusterMatchedError) Error() string

type ProbeOptions

type ProbeOptions struct {
	// Interval is the time between reachability probes for each remote.
	Interval time.Duration
	// Timeout bounds a single probe request so it never blocks on a dead socket
	// longer than intended.
	Timeout time.Duration
	// FailureThreshold is the number of consecutive unreachable probes that must
	// occur before a remote is considered lost and onLost is called. It is chosen
	// above the supervisor's backoff floor so a doomed manager cycle outlives the
	// growing backoff, keeping the rebuild loop period bounded rather than hot.
	FailureThreshold int
}

ProbeOptions configures the per-remote reachability probe.

type ProjectQuotaResourceRouter

type ProjectQuotaResourceRouter struct{}

ProjectQuotaResourceRouter routes project quotas to clusters based on availability zone.

func (ProjectQuotaResourceRouter) ExtractClusterSelector

func (p ProjectQuotaResourceRouter) ExtractClusterSelector(obj any) (string, error)

func (ProjectQuotaResourceRouter) Match

func (p ProjectQuotaResourceRouter) Match(obj any, labels map[string]string) (bool, error)

type RemoteConfig

type RemoteConfig struct {
	// The remote kubernetes apiserver url, e.g. "https://my-apiserver:6443".
	Host string `json:"host"`
	// The root CA certificate to verify the remote apiserver.
	// Ignored if InsecureSkipTLSVerify is true.
	CACert string `json:"caCert,omitempty"`
	// InsecureSkipTLSVerify disables verification of the remote apiserver's
	// TLS certificate. Use this for apiservers whose CA certificate rotates
	// frequently and does not chain to a stable root. Mutually exclusive
	// with CACert: when true, CACert is ignored.
	InsecureSkipTLSVerify bool `json:"insecureSkipTLSVerify,omitempty"`
	// The resource GVKs this apiserver serves, formatted as "<group>/<version>/<Kind>".
	GVKs []string `json:"gvks"`
	// Labels used by ResourceRouters to match resources to this cluster
	// for write operations (Create/Update/Delete/Patch).
	Labels map[string]string `json:"labels,omitempty"`
}

RemoteConfig maps multiple GVKs to a remote kubernetes apiserver with routing labels. It is assumed that the remote apiserver accepts the serviceaccount tokens issued by the local cluster.

type RemoteEndpoint

type RemoteEndpoint struct {
	// Host is the remote apiserver URL (from the cluster's rest.Config).
	Host string
	// Cluster is the controller-runtime cluster whose apiserver is probed.
	Cluster cluster.Cluster
}

RemoteEndpoint identifies a unique remote cluster to probe for reachability.

type ReservationsResourceRouter

type ReservationsResourceRouter struct{}

ReservationsResourceRouter routes reservations to clusters based on availability zone.

func (ReservationsResourceRouter) ExtractClusterSelector

func (r ReservationsResourceRouter) ExtractClusterSelector(obj any) (string, error)

func (ReservationsResourceRouter) Match

func (r ReservationsResourceRouter) Match(obj any, labels map[string]string) (bool, error)

type ResourceRouter

type ResourceRouter interface {
	Match(obj any, labels map[string]string) (bool, error)
	// ExtractClusterSelector extracts the routing key from the object (e.g. availability zone).
	// Used to enrich error messages when no cluster matches.
	ExtractClusterSelector(obj any) (string, error)
}

ResourceRouter determines which remote cluster a resource should be written to by matching the resource content against the cluster's labels.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL