resources

package
v1.3.15-beta2 Latest Latest
Warning

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

Go to latest
Published: Jul 17, 2026 License: Apache-2.0 Imports: 25 Imported by: 0

Documentation

Index

Constants

View Source
const MutationCheckpointInterval = 60 * time.Second

MutationCheckpointInterval is the interval for logging mutation checkpoint messages.

Variables

View Source
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) Name

func (f *Factory) Name() string

Name returns the stream name for logging.

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

type ResourceInfo struct {
	Group   string
	Version string
}

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

func (r *Watcher) DynamicListResources(ctx context.Context, logger *zap.Logger) (string, error)

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.

Jump to

Keyboard shortcuts

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