collector

package
v1.5.0 Latest Latest
Warning

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

Go to latest
Published: Sep 8, 2026 License: Apache-2.0 Imports: 28 Imported by: 0

Documentation

Index

Constants

View Source
const (
	// AWSNodeLabel is the label selector for aws-node pods.
	AWSNodeLabel = "k8s-app=aws-node"
	// AWSNodeNamespace is the namespace of the aws-node pods.
	AWSNodeNamespace = "kube-system"
	// AWSEksNodeagentContainer is the name of the container that has flow logs.
	AWSEksNodeagentContainer = "aws-eks-nodeagent"

	// MaxFlowAge bounds how far in the past an AWS VPC CNI flow may be and still be
	// worth sending. Both collectors re-read a pod/daemon's whole agent log on their
	// first scrape after a restart (the read position is not persisted), so the log
	// can contain flows far older than the backend's aggregation window. The backend
	// aggregates on a 1-minute window plus a 1-minute grace period and discards
	// anything older, so CacheFlowLine filters any flow whose timestamp is older than
	// MaxFlowAge before it is cached and sent.
	MaxFlowAge = 2 * time.Minute
)
View Source
const (
	// EKSComputeTypeLabel is the node label EKS applies to Auto Mode nodes.
	EKSComputeTypeLabel = "eks.amazonaws.com/compute-type"
	// EKSComputeTypeAuto is the value of EKSComputeTypeLabel on Auto Mode nodes.
	EKSComputeTypeAuto = "auto"
)
View Source
const (
	ICMP = "icmp"
	TCP  = "tcp"
	UDP  = "udp"
	SCTP = "sctp"

	IPv4 = "ipv4"
	IPv6 = "ipv6"
)

Protocol constants.

View Source
const DefaultNetworkPolicyAgentLogPath = "aws-routed-eni/network-policy-agent.log"

DefaultNetworkPolicyAgentLogPath is the node-local path (relative to the kubelet log root) of the AWS-managed Network Policy Agent log. It is the default for the configurable log path; AWS controls this component, so the path is overridable in case a future agent version relocates the file.

View Source
const (

	// HubbleRelayStatusCheckInterval is how often the peer-health watchdog polls
	// Hubble Relay's ServerStatus RPC while a flow stream is active.
	HubbleRelayStatusCheckInterval = 30 * time.Minute
)

Variables

View Source
var (
	ErrAWSVPCCNIInvalidLog       = errors.New("invalid AWS VPC CNI flow log format")
	ErrAWSVPCCNIInvalidIP        = errors.New("invalid IP address in AWS VPC CNI flow log")
	ErrAWSVPCCNINotFlowLog       = errors.New("log line is not a flow log")
	ErrAWSVPCCNIInvalidProtocol  = errors.New("unsupported protocol in AWS VPC CNI flow log")
	ErrAWSVPCCNIInvalidTimestamp = errors.New("invalid or missing timestamp in AWS VPC CNI flow log")
)

AWS VPC CNI flow log errors.

View Source
var (
	ErrFalcoEventIsNotFlow   = errors.New("ignoring falco event, not a network flow")
	ErrFalcoIncompleteL3Flow = errors.New("ignoring incomplete falco l3 network flow")
	ErrFalcoIncompleteL4Flow = errors.New("ignoring incomplete falco l4 network flow")
	ErrFalcoInvalidPort      = errors.New("ignoring incomplete falco flow due to bad ports")
	ErrFalcoTimestamp        = errors.New("incomplete or incorrectly formatted timestamp found in Falco flow")
)

Errors for Falco flow parsing.

Functions

func CacheFlowLine added in v1.5.0

func CacheFlowLine(ctx context.Context, sink FlowSink, line string, notBefore time.Time, logger *zap.Logger) (bool, error)

CacheFlowLine is the shared parse -> stale-filter -> cache path for the AWS VPC CNI Network Policy Agent flow log format. It is used by BOTH the standard aws-node pod-log collector and the EKS Auto Mode node-proxy log collector so the two share identical flow handling.

notBefore drops flows that fall outside the backend's stream time window. Neither collector persists its read position across restarts, so the first scrape of a pod/daemon re-reads the whole (up to ~200 MiB) active log; many of those records are far in the past. The backend discards flows outside its time window anyway, so sending them wastes the whole path from the operator onward. Dropping any flow whose log timestamp is before notBefore avoids that. Callers pass a ROLLING bound (typically time.Now().Add(-MaxFlowAge)) recomputed each poll, not a fixed startup time. A zero notBefore disables the filter (all flows pass).

Returns cached=true when a flow was cached. A parse failure is not an error: the agent log interleaves non-flow housekeeping lines with flow records, so only a context error is returned, letting the caller stop without advancing its checkpoint. A CacheFlow error is logged and skipped (best-effort delivery).

func ConvertCiliumFlow

func ConvertCiliumFlow(flowResp *observer.GetFlowsResponse) *pb.CiliumFlow

ConvertCiliumFlow converts a GetFlowsResponse object to a CiliumFlow object.

func CreateLayer3Message

func CreateLayer3Message(source string, destination string, ipVersion string) (*pb.IP, error)

CreateLayer3Message creates a Layer3 IP message from source/destination addresses.

func CreateLayer4Message

func CreateLayer4Message(proto string, srcPort, dstPort uint32, ipVersion string) (*pb.Layer4, error)

CreateLayer4Message converts event protocol and ports to a Layer4 proto message.

func FilterIllumioTraffic

func FilterIllumioTraffic(body string) bool

FilterIllumioTraffic filters out events related to Illumio network traffic.

func IsAWSVPCCNIAvailable added in v1.4.0

func IsAWSVPCCNIAvailable(ctx context.Context, logger *zap.Logger, k8sClient kubernetes.Interface) bool

IsAWSVPCCNIAvailable checks if AWS VPC CNI with flow logging is available in the cluster. It looks for aws-node pods with the aws-eks-nodeagent container. Checks multiple pods to handle rolling upgrades where some pods may not have nodeagent yet.

func IsCiliumAvailable

func IsCiliumAvailable(ctx context.Context, logger *zap.Logger, clientset kubernetes.Interface, ciliumNamespaces []string, tlsAuthProps tls.AuthProperties) bool

IsCiliumAvailable checks if Cilium Hubble Relay is available in the cluster.

func IsEKSAutoModeAvailable added in v1.5.0

func IsEKSAutoModeAvailable(ctx context.Context, logger *zap.Logger, k8sClient kubernetes.Interface) bool

IsEKSAutoModeAvailable reports whether the cluster has at least one EKS Auto Mode node. In Auto Mode the VPC CNI and Network Policy Agent are AWS-managed (no aws-node DaemonSet is present), so IsAWSVPCCNIAvailable returns false and standard pod-log collection cannot be used. Auto Mode nodes are identified by the node label eks.amazonaws.com/compute-type=auto.

This is a positive signal (the presence of the Auto Mode label) rather than inferring Auto Mode from the absence of aws-node, so a cluster with neither aws-node nor the Auto Mode label is not mistaken for Auto Mode.

func IsOVNKDeployed

func IsOVNKDeployed(ctx context.Context, logger *zap.Logger, ovnkNamespace string, clientset kubernetes.Interface) bool

IsOVNKDeployed checks for the presence of the OVN-Kubernetes namespace. https://ovn-kubernetes.io/installation/launching-ovn-kubernetes-on-kind/#run-the-kind-deployment-with-podman

func NetworkPolicyAgentLogPathSegments added in v1.5.0

func NetworkPolicyAgentLogPathSegments(logPath string) []string

NetworkPolicyAgentLogPathSegments returns the kubelet-proxy log path segments (prefixed with "logs") used to fetch the given node-local log path via the node proxy. An empty logPath falls back to DefaultNetworkPolicyAgentLogPath.

func NewFalcoEventHandler

func NewFalcoEventHandler(eventChan chan<- string) http.HandlerFunc

NewFalcoEventHandler creates a new HTTP handler function for processing Falco events.

func NewTemplateSystem

func NewTemplateSystem(logger *zap.Logger) (*netflows.BasicTemplateSystem, error)

NewTemplateSystem creates a template system for IPFIX message. It reads the template set from a binary file and adds it to the template system.

func ParseAWSVPCCNIFlowLog added in v1.4.0

func ParseAWSVPCCNIFlowLog(line string) (*pb.FiveTupleFlow, error)

ParseAWSVPCCNIFlowLog parses a VPC CNI flow log line into a FiveTupleFlow. Supports both old format (v1.0.x - v1.2.1) with separate JSON fields and new format (v1.2.2+) with embedded msg string.

func ParseIPVersion

func ParseIPVersion(decodedValue []byte) (string, error)

ParseIPVersion converts a byte slice into an IP version string (e.g., "ipv4" or "ipv6"). Returns an error if the slice is not the correct size or the IP version is unknown.

func ParseIPv4Address

func ParseIPv4Address(b []byte) (string, error)

ParseIPv4Address converts a byte slice into an IPv4 address string. Returns an error if the slice is not the correct size.

func ParseIPv6Address

func ParseIPv6Address(b []byte) (string, error)

ParseIPv6Address converts a byte slice into an IPv6 address string. Returns an error if the slice is not the correct size.

func ParsePodNetworkInfo

func ParsePodNetworkInfo(input string) (*pb.FiveTupleFlow, error)

ParsePodNetworkInfo parses the input string to extract network information into a FiveTupleFlow message.

func ParsePort

func ParsePort(decodedValue []byte) (uint16, error)

ParsePort converts a byte slice into a uint16 port number using BigEndian encoding. Returns an error if the slice is not the correct size.

func ParseProtocol

func ParseProtocol(decodedValue []byte) (string, error)

ParseProtocol converts a byte slice into a protocol string based on IANA protocol numbers. Returns an error if the slice is not the correct size or the protocol is unknown.

Types

type AWSVPCCNIFlowLog added in v1.4.0

type AWSVPCCNIFlowLog struct {
	Level     string `json:"level"`
	Timestamp string `json:"ts"`
	Logger    string `json:"logger"` // v1.0.x - v1.2.1
	Caller    string `json:"caller"` // v1.2.2+
	Message   string `json:"msg"`
	SrcIP     string `json:"Src IP"`

	SrcPort  uint32 `json:"Src Port"`
	DestIP   string `json:"Dest IP"`
	DestPort uint32 `json:"Dest Port"`
	Proto    string `json:"Proto"`   // TCP, UDP, ICMP, SCTP, UNKNOWN
	Verdict  string `json:"Verdict"` // ACCEPT, DENY, EXPIRED/DELETED
}

AWSVPCCNIFlowLog represents the flow log format from aws-eks-nodeagent.

Old format (v1.0.x - v1.2.1):

{"level":"info","ts":"2024-09-23T12:36:53.562Z","logger":"ebpf-client",
 "msg":"Flow Info: ","Src IP":"10.0.141.167","Src Port":39197,
 "Dest IP":"172.20.0.10","Dest Port":53,"Proto":"TCP","Verdict":"ACCEPT"}

New format (v1.2.2+):

{"level":"debug","ts":"2026-04-13T21:18:46.888Z","caller":"runtime/asm_amd64.s:1700",
 "msg":"Flow Info: Src IP: 10.0.1.28 Src Port: 55484 Dest IP: 10.0.1.132 Dest Port: 80 Proto TCP Verdict ACCEPT Direction egress"}

type CiliumFlowCollector

type CiliumFlowCollector struct {
	// contains filtered or unexported fields
}

CiliumFlowCollector collects flows from Cilium Hubble Relay running in this cluster.

func NewCiliumFlowCollector

func NewCiliumFlowCollector(ctx context.Context, logger *zap.Logger, clientset kubernetes.Interface, ciliumNamespaces []string, tlsAuthProperties tls.AuthProperties) (*CiliumFlowCollector, error)

NewCiliumFlowCollector connects to Cilium Hubble Relay, sets up an Observer client, and returns a new Collector using it. It tries namespaces until discovery succeeds.

func (*CiliumFlowCollector) ExportCiliumFlows

func (fm *CiliumFlowCollector) ExportCiliumFlows(ctx context.Context, flowSink FlowSink) error

ExportCiliumFlows makes one stream gRPC call to hubble-relay to collect, convert, and export flows into the given stream.

type FalcoEvent

type FalcoEvent struct {
	// Timestamp is the time the network event occurred. ISO 8601 format
	Timestamp *timestamppb.Timestamp `json:"time"`
	// SrcIP is the source IP address involved in the network event.
	SrcIP string `json:"srcip"`
	// DstIP is the destination IP address involved in the network event.
	DstIP string `json:"dstip"`
	// SrcPort is the source port number involved in the network event.
	SrcPort string `json:"srcport"`
	// DstPort is the destination port number involved in the network event.
	DstPort string `json:"dstport"`
	// Proto is the protocol used in the network event (e.g., TCP, UDP).
	Proto string `json:"proto"`
	// IpVersion is the version used in the network event (e.g. ipv4, ipv6).
	IpVersion string `json:"prototype"`
}

FalcoEvent represents the network information extracted from a Falco event.

type FlowSink

type FlowSink interface {
	CacheFlow(ctx context.Context, flow pb.Flow) error
	IncrementFlowsReceived()
}

FlowSink is the interface for caching network flows.

type K8sClientGetter

type K8sClientGetter interface {
	GetClientset() kubernetes.Interface
	GetDynamicClient() dynamic.Interface
	GetDiscoveryClient() discovery.DiscoveryInterface
}

K8sClientGetter provides access to Kubernetes clients.

type OVNKCollector

type OVNKCollector struct {
	// contains filtered or unexported fields
}

OVNKCollector collects IPFIX flows from OVN-Kubernetes.

func NewOVNKCollector

func NewOVNKCollector(logger *zap.Logger, ipfixCollectorPort string, flowSink FlowSink) *OVNKCollector

NewOVNKCollector creates a new OVN-K IPFIX collector.

func (*OVNKCollector) RunIPFIXCollector

func (c *OVNKCollector) RunIPFIXCollector(ctx context.Context) error

RunIPFIXCollector runs the UDP listener for OVN-K IPFIX flows. It blocks until the context is canceled.

type OVNKFlow

type OVNKFlow struct {
	SourceIP        string
	DestinationIP   string
	SourcePort      uint16
	DestinationPort uint16
	Protocol        string
	IPVersion       string
	StartTimestamp  *timestamppb.Timestamp
	EndTimestamp    *timestamppb.Timestamp
}

OVNKFlow represents a flow captured from OVN-Kubernetes.

func ProcessDataRecord

func ProcessDataRecord(dataRecord netflows.DataRecord, exportTime uint32) (OVNKFlow, error)

ProcessDataRecord processes a single data record and converts it into an OVNFlow. If any parsing step fails, it returns an error and skips the record.

Jump to

Keyboard shortcuts

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