collector

package
v1.3.15-beta Latest Latest
Warning

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

Go to latest
Published: Jun 15, 2026 License: Apache-2.0 Imports: 27 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"
)
View Source
const (
	ICMP = "icmp"
	TCP  = "tcp"
	UDP  = "udp"
	SCTP = "sctp"

	IPv4 = "ipv4"
	IPv6 = "ipv6"
)

Protocol constants.

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 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 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 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