Documentation
¶
Index ¶
- Variables
- func IsDuplicateError(err error) bool
- func IsNoClusterMatchedError(err error) bool
- type APIServersConfig
- type Client
- func (c *Client) AddRemote(ctx context.Context, host, caCert string, insecureSkipTLSVerify bool, ...) (cluster.Cluster, error)
- func (c *Client) Apply(ctx context.Context, obj runtime.ApplyConfiguration, ...) error
- func (c *Client) ClustersForGVK(gvk schema.GroupVersionKind) ([]cluster.Cluster, error)
- func (c *Client) ConfiguredRouteLabels(gvk schema.GroupVersionKind) []map[string]string
- func (c *Client) Create(ctx context.Context, obj client.Object, opts ...client.CreateOption) error
- func (c *Client) Delete(ctx context.Context, obj client.Object, opts ...client.DeleteOption) error
- func (c *Client) DeleteAllOf(ctx context.Context, obj client.Object, opts ...client.DeleteAllOfOption) error
- func (c *Client) GVKFromHomeScheme(obj runtime.Object) (gvk schema.GroupVersionKind, err error)
- func (c *Client) Get(ctx context.Context, key client.ObjectKey, obj client.Object, ...) error
- func (c *Client) GetEventRecorder(name string) recorder.EventRecorder
- func (c *Client) GroupVersionKindFor(obj runtime.Object) (schema.GroupVersionKind, error)
- func (c *Client) IndexField(ctx context.Context, obj client.Object, list client.ObjectList, field string, ...) error
- func (c *Client) InitFromConf(ctx context.Context, mgr ctrl.Manager, conf ClientConfig) error
- func (c *Client) IsObjectNamespaced(obj runtime.Object) (bool, error)
- func (c *Client) List(ctx context.Context, list client.ObjectList, opts ...client.ListOption) error
- func (c *Client) ListMetadataPerCluster(ctx context.Context, gvk schema.GroupVersionKind, opts ...client.ListOption) ([]ClusterObjectMetadata, error)
- func (c *Client) Patch(ctx context.Context, obj client.Object, patch client.Patch, ...) error
- func (c *Client) ProbeRemotes(ctx context.Context, opts ProbeOptions, onLost func(host string))
- func (c *Client) RESTMapper() meta.RESTMapper
- func (c *Client) Scheme() *runtime.Scheme
- func (c *Client) Status() client.StatusWriter
- func (c *Client) SubResource(subResource string) client.SubResourceClient
- func (c *Client) UniqueRemotes() []RemoteEndpoint
- func (c *Client) Update(ctx context.Context, obj client.Object, opts ...client.UpdateOption) error
- type ClientConfig
- type ClusterObjectMetadata
- type ClusterWrapper
- type CommittedResourceRouter
- type FlavorGroupCapacityResourceRouter
- type HistoryResourceRouter
- type HomeConfig
- type HypervisorResourceRouter
- type Monitor
- type MultiClusterRecorder
- type MulticlusterBuilder
- type NoClusterMatchedError
- type ProbeOptions
- type ProjectQuotaResourceRouter
- type RemoteConfig
- type RemoteEndpoint
- type ReservationsResourceRouter
- type ResourceRouter
Constants ¶
This section is empty.
Variables ¶
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).
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 ¶
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 ¶
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 ¶
func (c *Client) Apply(ctx context.Context, obj runtime.ApplyConfiguration, opts ...client.ApplyOption) error
Apply is not supported in the multicluster client as the group version kind cannot be inferred from the ApplyConfiguration.
func (*Client) ClustersForGVK ¶
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 ¶
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 ¶
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 ¶
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 ¶
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 ¶
Helper function to initialize a new multicluster client during service startup, using the conf module provided by cortex.
func (*Client) IsObjectNamespaced ¶
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) 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.
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 ¶
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)
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)
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)
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)
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 ¶
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.
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)
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)
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.