Documentation
¶
Index ¶
- Constants
- Variables
- func BuildResourceAPIGroupMap(resources []string, clientset kubernetes.Interface, logger *zap.Logger) (map[string]ResourceInfo, error)
- type Factory
- type KubernetesResourcesStream
- type ResourceConverter
- type ResourceInfo
- type ResourceStreamSender
- type RuntimeCacheHandler
- type Watcher
- func (r *Watcher) DynamicListResources(ctx context.Context, logger *zap.Logger) (string, error)
- func (r *Watcher) FetchResources(ctx context.Context, namespace string) (*unstructured.UnstructuredList, error)
- func (r *Watcher) WatchK8sResources(ctx context.Context, cancel context.CancelFunc, resourceVersion string, ...)
- type WatcherConfig
Constants ¶
const MutationCheckpointInterval = 60 * time.Second
MutationCheckpointInterval is the interval for logging mutation checkpoint messages.
Variables ¶
var ManagedResourceNames = []string{
"ciliumcidrgroups",
"ciliumclusterwidenetworkpolicies",
"ciliumnetworkpolicies",
}
ManagedResourceNames lists the plural resource names managed by the reconciler.
Functions ¶
func BuildResourceAPIGroupMap ¶ added in v1.4.0
func BuildResourceAPIGroupMap(resources []string, clientset kubernetes.Interface, logger *zap.Logger) (map[string]ResourceInfo, error)
BuildResourceAPIGroupMap creates a mapping between Kubernetes resources and their API groups with preferred versions. Exported for use by the reconciler.
Types ¶
type Factory ¶
type Factory struct {
Logger *zap.Logger
K8sClient k8sclient.Client
Stats *stream.Stats
FlowCollectorType pb.FlowCollector
ClusterName string // Optional: cluster name for self-managed clusters
Cache *cache.ConfiguredObjectCache
}
Factory creates resources stream clients.
func (*Factory) NewStreamClient ¶
func (f *Factory) NewStreamClient(ctx context.Context, grpcConn grpc.ClientConnInterface) (stream.StreamClient, error)
NewStreamClient creates a new resources stream client.
type KubernetesResourcesStream ¶
type KubernetesResourcesStream interface {
Send(req *pb.SendKubernetesResourcesRequest) error
Recv() (*pb.SendKubernetesResourcesResponse, error)
}
KubernetesResourcesStream abstracts the SendKubernetesResources gRPC stream.
type ResourceConverter ¶ added in v1.4.0
type ResourceConverter func(ctx context.Context, obj *unstructured.Unstructured) (*pb.KubernetesObjectData, error)
ResourceConverter converts an unstructured Kubernetes object into a KubernetesObjectData proto. Core and Cilium resources use separate implementations.
type ResourceInfo ¶ added in v1.3.14
ResourceInfo holds the API group and preferred version for a resource.
type ResourceStreamSender ¶
type ResourceStreamSender interface {
SendObjectData(logger *zap.Logger, metadata *pb.KubernetesObjectData) error
CreateMutationObject(metadata *pb.KubernetesObjectData, eventType watch.EventType) *pb.KubernetesResourceMutation
}
ResourceStreamSender abstracts the operations for sending resources to CloudSecure. Implemented by resourcesClient.
type RuntimeCacheHandler ¶ added in v1.4.0
type RuntimeCacheHandler func(ctx context.Context, eventType watch.EventType, metadata *pb.KubernetesObjectData) error
RuntimeCacheHandler processes a K8s object for runtime cache population. Called by the watcher after converting the K8s object to proto. The handler is responsible for filtering (e.g. checking operator-managed labels) and updating the runtime cache.
For watch events, eventType is Added/Modified/Deleted and the handler inserts or deletes from the runtime cache directly.
For list operations, eventType is empty — the handler accumulates the object for a later bulk ReplaceAll without modifying the cache directly.
type Watcher ¶
type Watcher struct {
// contains filtered or unexported fields
}
Watcher encapsulates components for listing and managing Kubernetes resources.
func NewWatcher ¶
func NewWatcher(config WatcherConfig) *Watcher
NewWatcher creates a new Watcher for a specific resource type.
func (*Watcher) DynamicListResources ¶
DynamicListResources lists a specified resource dynamically and sends each item down the current gRPC stream. If a RuntimeCacheHandler is set, it is called for each item to perform any cache-related processing.
func (*Watcher) FetchResources ¶
func (r *Watcher) FetchResources(ctx context.Context, namespace string) (*unstructured.UnstructuredList, error)
func (*Watcher) WatchK8sResources ¶
func (r *Watcher) WatchK8sResources(ctx context.Context, cancel context.CancelFunc, resourceVersion string, mutationChan chan *pb.KubernetesResourceMutation)
WatchK8sResources initiates a watch stream for the specified Kubernetes resource.
type WatcherConfig ¶
type WatcherConfig struct {
ResourceName string // plural-lowercase (e.g., "pods")
ApiGroup string
ApiVersion string
BaseLogger *zap.Logger
DynamicClient dynamic.Interface
ResourcesClient ResourceStreamSender // for CloudSecure streaming
RuntimeCacheHandler RuntimeCacheHandler // optional: processes events for runtime cache population
Limiter *rate.Limiter
Converter ResourceConverter
}
WatcherConfig holds the configuration for creating a new Watcher.