Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
2929d98
fix(cpm): re-resolve in-cluster peer IP at profile flush (SUB-8865)
matthyx Oct 5, 2026
37fb999
fix(cpm): address Copilot review 5413180027
matthyx Oct 5, 2026
1564663
fix(cpm): address Copilot review 5413292601
matthyx Oct 5, 2026
a0c188b
fix(cpm): address Copilot review 5413479461
matthyx Oct 5, 2026
8baa483
perf(cpm): address Copilot review 5413610359
matthyx Oct 5, 2026
3d5860e
fix(cpm): address Copilot review 5413716935
matthyx Oct 5, 2026
50bc912
fix(cpm): address CodeRabbit review 5413700702
matthyx Oct 5, 2026
5973490
fix(cpm): address Copilot review 5413961676
matthyx Oct 5, 2026
1eedb87
fix(queue): preserve optional split data on enqueue errors
matthyx Oct 5, 2026
88c1c21
fix(queue): split dominant peer ports before depth exhaustion
matthyx Oct 5, 2026
5b4fb86
fix(queue): balance peer splits by payload bytes
matthyx Oct 5, 2026
23c1a95
fix(cpm): budget all materialized profile payloads
matthyx Oct 5, 2026
bb96600
fix(cpm): retry Service ports after failed ingestion lookup
matthyx Oct 5, 2026
87499c4
fix(queue): reserve final split level for storage rejection
matthyx Oct 5, 2026
3800bc9
fix(queue): isolate recursive split slice capacities
matthyx Oct 5, 2026
0254c97
fix(queue): reject optional splits that grow protobuf payloads
matthyx Oct 5, 2026
39730cf
fix(cpm): preserve retry windows and pending status reports
matthyx Oct 5, 2026
76e1380
fix(profiles): bound pressure flush work to fresh events
matthyx Oct 5, 2026
449c7e6
fix(profiles): cap deferred network retention
matthyx Oct 5, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
354 changes: 305 additions & 49 deletions pkg/containerprofilemanager/v1/container_data.go

Large diffs are not rendered by default.

11 changes: 7 additions & 4 deletions pkg/containerprofilemanager/v1/container_data_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,7 @@ func serviceNetworkEvent(port uint16, protocol string) NetworkEvent {
}
}

// TestCreateNetworkNeighbor_ServiceTargetPortMatrix checks Service port remapping and observed-port fallbacks across protocols and endpoint sources.
func TestCreateNetworkNeighbor_ServiceTargetPortMatrix(t *testing.T) {
tests := []struct {
name string
Expand Down Expand Up @@ -270,7 +271,7 @@ func TestCreateNetworkNeighbor_ServiceTargetPortMatrix(t *testing.T) {
}

cd := &containerData{}
neighbor := cd.createNetworkNeighbor("", tc.event, "default", client, nil)
neighbor := cd.createNetworkNeighbor("", tc.event, "default", client, nil, nil, nil, false)
require.NotNil(t, neighbor)
require.Equal(t, map[string]string{"app": "api"}, neighbor.PodSelector.MatchLabels)
require.Equal(t, tc.wantPorts, networkPortValues(neighbor.Ports))
Expand All @@ -284,6 +285,7 @@ func TestCreateNetworkNeighbor_ServiceTargetPortMatrix(t *testing.T) {
}
}

// TestCreateNetworkNeighbor_NonServiceDestinationsUnchanged checks that pod and raw-IP peers retain their observed ports.
func TestCreateNetworkNeighbor_NonServiceDestinationsUnchanged(t *testing.T) {
cd := &containerData{}

Expand All @@ -298,7 +300,7 @@ func TestCreateNetworkNeighbor_NonServiceDestinationsUnchanged(t *testing.T) {
},
}
podEvent.SetDestinationPodLabels(map[string]string{"app": "web"})
podNeighbor := cd.createNetworkNeighbor("", podEvent, "default", nil, nil)
podNeighbor := cd.createNetworkNeighbor("", podEvent, "default", nil, nil, nil, nil, false)
require.NotNil(t, podNeighbor)
require.Equal(t, []int32{8080}, networkPortValues(podNeighbor.Ports))

Expand All @@ -310,11 +312,12 @@ func TestCreateNetworkNeighbor_NonServiceDestinationsUnchanged(t *testing.T) {
IPAddress: "93.184.216.34",
},
}
rawNeighbor := cd.createNetworkNeighbor("", rawEvent, "default", nil, nil)
rawNeighbor := cd.createNetworkNeighbor("", rawEvent, "default", nil, nil, nil, nil, false)
require.NotNil(t, rawNeighbor)
require.Equal(t, []int32{443}, networkPortValues(rawNeighbor.Ports))
}

// TestGenerateNetworkPolicy_ServiceTargetPortRoundTrip checks that generated policies use the backend target port instead of the Service port.
func TestGenerateNetworkPolicy_ServiceTargetPortRoundTrip(t *testing.T) {
service := newServiceWorkload("api", map[string]any{"app.kubernetes.io/name": "api"}, map[string]any{
"port": 80, "targetPort": 8080, "protocol": "TCP",
Expand All @@ -326,7 +329,7 @@ func TestGenerateNetworkPolicy_ServiceTargetPortRoundTrip(t *testing.T) {

cd := &containerData{}
event := serviceNetworkEvent(80, "tcp")
neighbor := cd.createNetworkNeighbor("", event, "default", client, nil)
neighbor := cd.createNetworkNeighbor("", event, "default", client, nil, nil, nil, false)
require.NotNil(t, neighbor)

egressPorts := make([]softwarecomposition.NetworkPort, 0, len(neighbor.Ports))
Expand Down
55 changes: 44 additions & 11 deletions pkg/containerprofilemanager/v1/containerprofile_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
mapset "github.com/deckarep/golang-set/v2"
"github.com/goradd/maps"
containercollection "github.com/inspektor-gadget/inspektor-gadget/pkg/container-collection"
"github.com/inspektor-gadget/inspektor-gadget/pkg/operators/common"
"github.com/kubescape/go-logger"
"github.com/kubescape/go-logger/helpers"
"github.com/kubescape/node-agent/pkg/config"
Expand Down Expand Up @@ -47,22 +48,35 @@ type containerData struct {
monitorDone chan struct{}
monitorDoneOnce sync.Once

// Apparent size
// Apparent size of observations collected since the last flush; deferred peers
// remain in networks without being charged to each new active batch.
size atomic.Int64

// Cleanup resources
timer *time.Timer // For max sniffing time

// Events reported for this container that need to be saved to the profile
capabilites mapset.Set[string]
syscalls mapset.Set[string]
endpoints *maps.SafeMap[string, *v1beta1.HTTPEndpoint]
execs *maps.SafeMap[string, []string] // Map of execs, key is SHA256 hash
opens *maps.SafeMap[string, mapset.Set[string]] // Map of opens, key is file path
rulePolicies *maps.SafeMap[string, *v1beta1.RulePolicy] // Map of rule policies, key is rule ID
callStacks *maps.SafeMap[string, *v1beta1.IdentifiedCallStack] // Map of callstacks, key is SHA256 hash
networks mapset.Set[NetworkEvent]
droppedEvents bool // Indicates if any events were dropped during monitoring
capabilites mapset.Set[string]
syscalls mapset.Set[string]
endpoints *maps.SafeMap[string, *v1beta1.HTTPEndpoint]
execs *maps.SafeMap[string, []string] // Map of execs, key is SHA256 hash
opens *maps.SafeMap[string, mapset.Set[string]] // Map of opens, key is file path
rulePolicies *maps.SafeMap[string, *v1beta1.RulePolicy] // Map of rule policies, key is rule ID
callStacks *maps.SafeMap[string, *v1beta1.IdentifiedCallStack] // Map of callstacks, key is SHA256 hash
networks mapset.Set[NetworkEvent] // Union used for deduplication and interval/final retries.
activeNetworks mapset.Set[NetworkEvent] // Newly collected observations since the last flush.
networkFlushForSize bool // Size-triggered saves visit only activeNetworks.
deferredNetworks mapset.Set[NetworkEvent]
prevDeferredNetworks mapset.Set[NetworkEvent]
droppedEvents bool // Indicates if any events were dropped during monitoring

// Positive durations give unresolved peers an informer catch-up window across rapid saves.
networkDeferralDuration time.Duration
networkDeferredUntil map[NetworkEvent]time.Time
// Deferred admission is tracked independently from the active flush budget.
networkDeferredSizeLimit int64
networkDeferredSize int64
networkDeferredSizes map[NetworkEvent]int64

// Service port snapshots keep report-time accounting and serialization consistent.
servicePorts map[NetworkEvent][]uint16
Expand All @@ -86,6 +100,7 @@ type ContainerProfileManager struct {
cfg config.Config
k8sClient k8sclient.K8sClientInterface
k8sObjectCache objectcache.K8sObjectCache
k8sInventory common.K8sInventoryCache
storageClient storage.ProfileCreator
dnsResolverClient dnsmanager.DNSResolver
seccompManager seccompmanager.SeccompManagerClient
Expand Down Expand Up @@ -122,6 +137,11 @@ func (cpm *ContainerProfileManager) SetCompletionNotifier(n objectcache.Completi
cpm.completionNotifier = n
}

// SetK8sInventory sets the k8s inventory cache (primarily used in tests)
func (cpm *ContainerProfileManager) SetK8sInventory(k8sInventory common.K8sInventoryCache) {
cpm.k8sInventory = k8sInventory
}

// SetSyscallFlusher implements containerprofilemanager.ContainerProfileManagerClient.
func (cpm *ContainerProfileManager) SetSyscallFlusher(flush func()) {
cpm.syscallFlusher.Store(&flush)
Expand Down Expand Up @@ -162,6 +182,15 @@ func NewContainerProfileManager(
lifecycleTracker: otelsetup.NewProfileLifecycleTracker(),
}

if cfg.KubernetesMode {
if k8sInventory, err := common.GetK8sInventoryCache(); err == nil && k8sInventory != nil {
containerProfileManager.k8sInventory = k8sInventory
k8sInventory.Start()
} else if err != nil {
logger.L().Debug("failed to initialize k8s inventory cache in container profile manager", helpers.Error(err))
}
}

// Initialize queue
queueDir := os.Getenv("QUEUE_DIR")
if queueDir == "" {
Expand Down Expand Up @@ -203,7 +232,7 @@ func NewContainerProfileManager(
return containerProfileManager, nil
}

// Stop stops the container profile manager
// Close stops container timers, the persistent queue, and the Kubernetes inventory.
func (cpm *ContainerProfileManager) Close() {
// Stop all container timers and clear container map
cpm.containersMu.Lock()
Expand All @@ -221,6 +250,10 @@ func (cpm *ContainerProfileManager) Close() {
if cpm.queueData != nil {
_ = cpm.queueData.Close()
}

if cpm.k8sInventory != nil {
cpm.k8sInventory.Stop()
}
}

var _ containerprofilemanager.ContainerProfileManagerClient = (*ContainerProfileManager)(nil)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -545,6 +545,7 @@ func TestContainerProfileManagerCreation(t *testing.T) {
assert.NotNil(t, cpm.maxSniffTimeNotificationChan)
}

// TestContainerDataMethods checks that an empty container produces no profile events or network neighbors.
func TestContainerDataMethods(t *testing.T) {
cd := &containerData{}

Expand Down Expand Up @@ -573,11 +574,11 @@ func TestContainerDataMethods(t *testing.T) {
assert.Empty(t, callStacks)

// Test getIngressNetworkNeighbors with nil networks
ingress := cd.getIngressNetworkNeighbors("", "default", nil, nil)
ingress := cd.getIngressNetworkNeighbors("", "default", nil, nil, nil, nil, false)
assert.Empty(t, ingress)

// Test getEgressNetworkNeighbors with nil networks
egress := cd.getEgressNetworkNeighbors("", "default", nil, nil)
egress := cd.getEgressNetworkNeighbors("", "default", nil, nil, nil, nil, false)
assert.Empty(t, egress)
}

Expand Down
26 changes: 15 additions & 11 deletions pkg/containerprofilemanager/v1/event_reporting.go
Original file line number Diff line number Diff line change
Expand Up @@ -353,27 +353,31 @@ func (cpm *ContainerProfileManager) ReportNetworkEvent(containerID string, event
networkEvent.SetPodLabels(event.GetPodLabels())
networkEvent.SetDestinationPodLabels(dstEndpoint.PodLabels)

resolveEndpoint(&networkEvent, cpm.k8sInventory, cpm.k8sObjectCache)
Comment thread
matthyx marked this conversation as resolved.

// Skip if we already saved this event
if data.networks.Contains(networkEvent) {
return 0, nil
}

if networkEvent.Destination.Kind == EndpointKindService {
ports := []uint16{networkEvent.Port}
if cpm.k8sClient != nil {
svc, err := cpm.k8sClient.GetWorkload(networkEvent.Destination.Namespace, "Service", networkEvent.Destination.Name)
if err == nil {
ports = resolveServiceEnforcementPorts(cpm.k8sClient, networkEvent.Destination.Namespace,
networkEvent.Destination.Name, svc, networkEvent.Port, networkEvent.Protocol)
if networkEvent.Destination.Kind == EndpointKindService && cpm.k8sClient != nil {
svc, err := cpm.k8sClient.GetWorkload(networkEvent.Destination.Namespace, "Service", networkEvent.Destination.Name)
// Failed lookups leave no snapshot so flush-time resolution can retry.
if err == nil && svc != nil {
ports := resolveServiceEnforcementPorts(cpm.k8sClient, networkEvent.Destination.Namespace,
networkEvent.Destination.Name, svc, networkEvent.Port, networkEvent.Protocol)
if data.servicePorts == nil {
data.servicePorts = make(map[NetworkEvent][]uint16)
}
data.servicePorts[networkEvent] = ports
}
if data.servicePorts == nil {
data.servicePorts = make(map[NetworkEvent][]uint16)
}
data.servicePorts[networkEvent] = ports
}

data.networks.Add(networkEvent)
if data.activeNetworks == nil {
data.activeNetworks = mapset.NewSet[NetworkEvent]()
}
data.activeNetworks.Add(networkEvent)
return size.Of(networkEvent) + networkNeighborIncrement(data, networkEvent), nil
})

Expand Down
16 changes: 9 additions & 7 deletions pkg/containerprofilemanager/v1/event_reporting_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -117,7 +117,7 @@ func TestNetworkNeighborIncrementCoversMaxDNSName(t *testing.T) {
}

cd := &containerData{}
neighbor := cd.createNetworkNeighbor("", networkEvent, "default", nil, fakeDNSResolver{domain: maxDNSName})
neighbor := cd.createNetworkNeighbor("", networkEvent, "default", nil, fakeDNSResolver{domain: maxDNSName}, nil, nil, false)
if !assert.NotNil(t, neighbor) {
return
}
Expand Down Expand Up @@ -157,7 +157,7 @@ func TestNetworkNeighborIncrementCoversSelectorPayload(t *testing.T) {
// neighbor. watchedContainerData.Namespace is what networkNeighborIncrement reads to make
// the same "different namespace" call createNetworkNeighbor's own namespace arg does below.
cd := &containerData{watchedContainerData: &objectcache.WatchedContainerData{Namespace: "default"}}
neighbor := cd.createNetworkNeighbor("", networkEvent, "default", nil, nil)
neighbor := cd.createNetworkNeighbor("", networkEvent, "default", nil, nil, nil, nil, false)
if !assert.NotNil(t, neighbor) {
return
}
Expand Down Expand Up @@ -236,6 +236,7 @@ func (r *trackingDNSResolver) ResolveContainerProcessToCloudServices(string, uin
return nil
}

// TestCreateNetworkNeighbor_EmptyContainerIDWithWatchedContainerData checks that DNS lookup preserves an explicitly empty container ID.
func TestCreateNetworkNeighbor_EmptyContainerIDWithWatchedContainerData(t *testing.T) {
cd := &containerData{
watchedContainerData: &objectcache.WatchedContainerData{
Expand All @@ -252,13 +253,14 @@ func TestCreateNetworkNeighbor_EmptyContainerIDWithWatchedContainerData(t *testi
}

resolver := &trackingDNSResolver{}
neighbor := cd.createNetworkNeighbor("", networkEvent, "default", nil, resolver)
neighbor := cd.createNetworkNeighbor("", networkEvent, "default", nil, resolver, nil, nil, false)
assert.NotNil(t, neighbor)
assert.Equal(t, "", resolver.lastContainerID, "empty containerID must be preserved without falling back to watchedContainerData")
assert.Equal(t, "93.184.216.34", resolver.lastIPAddress)
assert.Equal(t, "resolved.domain", neighbor.DNS)
}

// TestReportNetworkEventServicePortMultiplicity checks that all backend ports count toward the size budget and stay fixed within a batch.
func TestReportNetworkEventServicePortMultiplicity(t *testing.T) {
cpm, entry := newTestManager(t, "container1")
client := &servicePortTestClient{
Expand All @@ -279,7 +281,7 @@ func TestReportNetworkEventServicePortMultiplicity(t *testing.T) {
DstPort: 80, Proto: "tcp", PktType: utils.OutgoingPktType,
}
cpm.ReportNetworkEvent("container1", event)
neighbor := entry.data.createNetworkNeighbor("", serviceNetworkEvent(80, "tcp"), "default", client, nil)
neighbor := entry.data.createNetworkNeighbor("", serviceNetworkEvent(80, "tcp"), "default", client, nil, nil, nil, false)
require.NotNil(t, neighbor)
require.Equal(t, []int32{8080, 9090, 10000}, networkPortValues(neighbor.Ports))
// Isolate the port budget so unused selector headroom cannot hide an undercount.
Expand All @@ -305,13 +307,13 @@ func TestReportNetworkEventServicePortMultiplicity(t *testing.T) {

// Endpoint changes after reporting must not change the budgeted port list.
require.NoError(t, client.kubeClient.DiscoveryV1().EndpointSlices("default").Delete(context.Background(), "c", metav1.DeleteOptions{}))
neighbor = entry.data.createNetworkNeighbor("", serviceNetworkEvent(80, "tcp"), "default", client, nil)
neighbor = entry.data.createNetworkNeighbor("", serviceNetworkEvent(80, "tcp"), "default", client, nil, nil, nil, false)
require.Equal(t, []int32{8080, 9090, 10000}, networkPortValues(neighbor.Ports))

// A new profile batch resolves fresh ports instead of keeping the old snapshot.
entry.data.emptyEvents()
cpm.ReportNetworkEvent("container1", event)
neighbor = entry.data.createNetworkNeighbor("", serviceNetworkEvent(80, "tcp"), "default", client, nil)
neighbor = entry.data.createNetworkNeighbor("", serviceNetworkEvent(80, "tcp"), "default", client, nil, nil, nil, false)
require.Equal(t, []int32{8080, 9090}, networkPortValues(neighbor.Ports))
require.Less(t, entry.data.size.Load(), recordedSize)
}
Expand Down Expand Up @@ -349,7 +351,7 @@ func TestCreateNetworkNeighbor_StatefulSetPeerStripsPodIdentityLabels(t *testing
})

cd := &containerData{}
neighbor := cd.createNetworkNeighbor("", networkEvent, "default", nil, nil)
neighbor := cd.createNetworkNeighbor("", networkEvent, "default", nil, nil, nil, nil, false)
if !assert.NotNil(t, neighbor) {
return
}
Expand Down
Loading
Loading