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",
"clusternetworkpolicies",
}
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, obj *unstructured.Unstructured, 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 ownership filtering (via managedFields on the raw object) 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.