resources

package
v1.4.1-beta Latest Latest
Warning

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

Go to latest
Published: Aug 18, 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",
	"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) 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, 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

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