Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
1 change: 1 addition & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ require (
github.com/hashicorp/golang-lru/v2 v2.0.7
github.com/ianlancetaylor/demangle v0.0.0-20251118225945-96ee0021ea0f
github.com/mariomac/guara v0.0.0-20250408105519-1e4dbdfb7136
github.com/oschwald/maxminddb-golang v1.13.1
github.com/prometheus/client_golang v1.23.2
github.com/prometheus/client_model v0.6.2
github.com/prometheus/common v0.66.1
Expand Down
2 changes: 2 additions & 0 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -221,6 +221,8 @@ github.com/onsi/ginkgo/v2 v2.23.4 h1:ktYTpKJAVZnDT4VjxSbiBenUjmlL/5QkBEocaWXiQus
github.com/onsi/ginkgo/v2 v2.23.4/go.mod h1:Bt66ApGPBFzHyR+JO10Zbt0Gsp4uWxu5mIOTusL46e8=
github.com/onsi/gomega v1.37.0 h1:CdEG8g0S133B4OswTDC/5XPSzE1OeP29QOioj2PID2Y=
github.com/onsi/gomega v1.37.0/go.mod h1:8D9+Txp43QWKhM24yyOBEdpkzN8FvJyAwecBgsU4KU0=
github.com/oschwald/maxminddb-golang v1.13.1 h1:G3wwjdN9JmIK2o/ermkHM+98oX5fS+k5MbwsmL4MRQE=
github.com/oschwald/maxminddb-golang v1.13.1/go.mod h1:K4pgV9N/GcK694KSTmVSDTODk4IsCNThNdTmnaBZ/F8=
github.com/pelletier/go-toml/v2 v2.2.2 h1:aYUidT7k73Pcl9nb2gScu7NSrKCSHIDE89b3+6Wq+LM=
github.com/pelletier/go-toml/v2 v2.2.2/go.mod h1:1t835xjRzz80PqgE6HHgN2JOsmgYu/h4qDAS4n929Rs=
github.com/pierrec/lz4/v4 v4.1.22 h1:cKFw6uJDK+/gfw5BcDL0JL5aBsAFdsIT18eRtLj7VIU=
Expand Down
Binary file added internal/test/geoip/GeoIP2-Country-Test.mmdb
Binary file not shown.
Binary file added internal/test/geoip/GeoLite2-ASN-Test.mmdb
Binary file not shown.
Binary file added internal/test/geoip/ipinfo_lite_sample.mmdb
Binary file not shown.
18 changes: 16 additions & 2 deletions pkg/export/attributes/attr_defs.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ const (
GroupHTTPCommon
GroupHost
GroupMessaging
GroupNetGeoIP
)

func (e *AttrGroups) Has(groups AttrGroups) bool {
Expand All @@ -52,6 +53,7 @@ func getDefinitions(
promEnabled := groups.Has(GroupPrometheus)
ifaceDirEnabled := groups.Has(GroupNetIfaceDirection)
cidrEnabled := groups.Has(GroupNetCIDR)
geoipEnabled := groups.Has(GroupNetGeoIP)

// attributes to be reported exclusively for prometheus exporters
prometheusAttributes := NewAttrReportGroup(
Expand Down Expand Up @@ -139,6 +141,18 @@ func getDefinitions(
extraGroupAttributes[GroupNetCIDR],
)

networkGeoIP := NewAttrReportGroup(
!geoipEnabled,
nil,
map[attr.Name]Default{
attr.SrcCountry: true,
attr.DstCountry: true,
attr.SrcASN: true,
attr.DstASN: true,
},
extraGroupAttributes[GroupNetGeoIP],
)

// networkInterZone* supports the same attributes as
// network* counterpart, but all of them disabled by default, to keep cardinality low
networkInterZone := copyDisabled(networkAttributes)
Expand Down Expand Up @@ -237,10 +251,10 @@ func getDefinitions(

return map[Section]AttrReportGroup{
NetworkFlow.Section: {
SubGroups: []*AttrReportGroup{&networkAttributes, &networkCIDR, &networkKubeAttributes},
SubGroups: []*AttrReportGroup{&networkAttributes, &networkCIDR, &networkGeoIP, &networkKubeAttributes},
},
NetworkInterZone.Section: {
SubGroups: []*AttrReportGroup{&networkInterZone, &networkInterZoneCIDR, &networkInterZoneKube},
SubGroups: []*AttrReportGroup{&networkInterZone, &networkInterZoneCIDR, &networkGeoIP, &networkInterZoneKube},
},
HTTPServerDuration.Section: {
SubGroups: []*AttrReportGroup{&appAttributes, &appKubeAttributes, &httpCommon, &serverInfo},
Expand Down
5 changes: 5 additions & 0 deletions pkg/export/attributes/names/attrs.go
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,11 @@ const (
K8sDstOwnerType = Name("k8s.dst.owner.type")
K8sDstNodeIP = Name("k8s.dst.node.ip")
K8sDstNodeName = Name("k8s.dst.node.name")

SrcCountry = Name("src.country")
DstCountry = Name("dst.country")
SrcASN = Name("src.asn")
DstASN = Name("dst.asn")
)

// other OBI-specific attributes
Expand Down
3 changes: 3 additions & 0 deletions pkg/instrumenter/instrumenter.go
Original file line number Diff line number Diff line change
Expand Up @@ -232,4 +232,7 @@ func attributeGroups(config *obi.Config, ctxInfo *global.ContextInfo) {
if config.NetworkFlows.CIDRs.Enabled() {
ctxInfo.MetricAttributeGroups.Add(attributes.GroupNetCIDR)
}
if config.NetworkFlows.GeoIP.Enabled() {
ctxInfo.MetricAttributeGroups.Add(attributes.GroupNetGeoIP)
}
}
180 changes: 180 additions & 0 deletions pkg/internal/netolly/flow/geoip.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,180 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0

package flow

import (
"context"
"errors"
"fmt"
"log/slog"
"net"
"time"

"github.com/hashicorp/golang-lru/v2/expirable"
"github.com/oschwald/maxminddb-golang"

attr "go.opentelemetry.io/obi/pkg/export/attributes/names"
"go.opentelemetry.io/obi/pkg/internal/netolly/ebpf"
"go.opentelemetry.io/obi/pkg/pipe/msg"
"go.opentelemetry.io/obi/pkg/pipe/swarm"
)

// GeoIP is currently experimental. It is kept disabled by default and will be hidden
// from the documentation. This means that it does not impact in the overall Beyla performance.
type GeoIP struct {
IPInfo IPInfoConfig `yaml:"ipinfo"`
MaxMindInfo MaxMindConfig `yaml:"maxmind"`
CacheLen int `yaml:"cache_len" env:"OTEL_EBPF_NETWORK_GEOIP_CACHE_LEN" validate:"gte=0"`
CacheTTL time.Duration `yaml:"cache_expiry" env:"OTEL_EBPF_NETWORK_GEOIP_CACHE_TTL" validate:"gte=0"`
}

type IPInfoConfig struct {
Path string `yaml:"path" env:"OTEL_EBPF_NETWORK_GEOIP_IPINFO_PATH"`
}
type MaxMindConfig struct {
CountryPath string `yaml:"country_path" env:"OTEL_EBPF_NETWORK_GEOIP_MAXMIND_COUNTRY_PATH"`
ASNPath string `yaml:"asn_path" env:"OTEL_EBPF_NETWORK_GEOIP_MAXMIND_ASN_PATH"`
}

type ipInfo struct {
Country string
ASN string
}

func (g GeoIP) Enabled() bool {
return g.IPInfo.Path != "" || (g.MaxMindInfo.ASNPath != "" && g.MaxMindInfo.CountryPath != "")
}

func geoiplog() *slog.Logger {
return slog.With("component", "flow.GeoIP")
}

func GeoIPProvider(cfg *GeoIP, input, output *msg.Queue[[]*ebpf.Record]) swarm.InstanceFunc {
return func(_ context.Context) (swarm.RunFunc, error) {
if !cfg.Enabled() {
return swarm.Bypass(input, output)
}
lookupFn, err := getLookupFn(cfg)
if err != nil {
return nil, err
}

log := geoiplog()
in := input.Subscribe(msg.SubscriberName("flow.GeoIP"))
cache := expirable.NewLRU[ebpf.IPAddr, ipInfo](cfg.CacheLen, nil, cfg.CacheTTL)
cachedLookup := func(addr *ebpf.IPAddr) (ipInfo, error) {
info, ok := cache.Get(*addr)
if ok {
return info, nil
}
info, err := lookupFn(addr.IP())
if err != nil {
return info, err
}
cache.Add(*addr, info)
return info, nil
}

// only warn the first time to prevent log flooding
var failureLogFn func(string, ...any)
failureLogFn = func(msg string, args ...any) {
log.Warn(msg, args...)
failureLogFn = log.Debug
}

return func(_ context.Context) {
defer output.Close()
log.Debug("starting GeoIP node")
for flows := range in {
for _, flow := range flows {
srcInfo, err := cachedLookup(flow.Id.SrcIP())
if err != nil {
failureLogFn("failed to perform geoip lookup for source", "err", err)
}
dstInfo, err := cachedLookup(flow.Id.DstIP())
if err != nil {
failureLogFn("failed to perform geoip lookup for destination", "err", err)
}
if flow.Attrs.Metadata == nil {
flow.Attrs.Metadata = map[attr.Name]string{}
}
flow.Attrs.Metadata[attr.SrcCountry] = srcInfo.Country
flow.Attrs.Metadata[attr.DstCountry] = dstInfo.Country
flow.Attrs.Metadata[attr.SrcASN] = srcInfo.ASN
flow.Attrs.Metadata[attr.DstASN] = dstInfo.ASN
}
output.Send(flows)
}
}, nil
}
}

type IPLookupFn func(addr net.IP) (ipInfo, error)

func getLookupFn(cfg *GeoIP) (IPLookupFn, error) {
if cfg.IPInfo.Path != "" {
return ipinfoLookup(cfg.IPInfo.Path)
}
if cfg.MaxMindInfo.ASNPath != "" && cfg.MaxMindInfo.CountryPath != "" {
return maxmindlookup(cfg.MaxMindInfo.CountryPath, cfg.MaxMindInfo.ASNPath)
}
return nil, errors.New("no provider configured")
}

type ipInfoLiteRecord struct {
Country string `maxminddb:"country_code"`
ASN string `maxminddb:"asn"`
}

func ipinfoLookup(path string) (IPLookupFn, error) {
db, err := maxminddb.Open(path)
if err != nil {
return nil, fmt.Errorf("opening ipinfo database: %w", err)
}
return func(addr net.IP) (ipInfo, error) {
record := ipInfoLiteRecord{}
err := db.Lookup(addr, &record)
if err != nil {
return ipInfo{}, fmt.Errorf("looking up address: %w", err)
}
return ipInfo(record), nil
}, nil
}

type maxmindCountryRecord struct {
Country struct {
Code string `maxminddb:"iso_code"`
} `maxminddb:"country"`
}
type maxmindASNRecord struct {
ASN uint64 `maxminddb:"autonomous_system_number"`
}

func maxmindlookup(countryPath, asnPath string) (IPLookupFn, error) {
countryDB, err := maxminddb.Open(countryPath)
if err != nil {
return nil, fmt.Errorf("opening maxmind country database: %w", err)
}
asnDB, err := maxminddb.Open(asnPath)
if err != nil {
return nil, fmt.Errorf("opening maxmind country database: %w", err)
}
return func(addr net.IP) (ipInfo, error) {
countryRecord := maxmindCountryRecord{}
if err := countryDB.Lookup(addr, &countryRecord); err != nil {
return ipInfo{}, fmt.Errorf("looking up country for address: %w", err)
}
asnRecord := maxmindASNRecord{}
if err := asnDB.Lookup(addr, &asnRecord); err != nil {
return ipInfo{}, fmt.Errorf("looking up country for address: %w", err)
}
out := ipInfo{
Country: countryRecord.Country.Code,
}
if asnRecord.ASN != 0 {
out.ASN = fmt.Sprintf("AS%d", asnRecord.ASN)
}
return out, nil
}, nil
}
116 changes: 116 additions & 0 deletions pkg/internal/netolly/flow/geoip_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,116 @@
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0

package flow

import (
"encoding/binary"
"flag"
"fmt"
"math"
"math/rand/v2"
"net"
"testing"
"time"

"github.com/hashicorp/golang-lru/v2/expirable"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"

"go.opentelemetry.io/obi/pkg/internal/netolly/ebpf"
)

func TestMaxMindLookup(t *testing.T) {
lookupFn, err := getLookupFn(&GeoIP{
MaxMindInfo: MaxMindConfig{
ASNPath: "../../../../internal/test/geoip/GeoLite2-ASN-Test.mmdb",
CountryPath: "../../../../internal/test/geoip/GeoIP2-Country-Test.mmdb",
},
})
require.NoError(t, err)
info, err := lookupFn(net.IPv4(216, 160, 83, 57))
require.NoError(t, err)
assert.Equal(t, "AS209", info.ASN)
assert.Equal(t, "US", info.Country)
}

func TestIPInfoLookup(t *testing.T) {
lookupFn, err := getLookupFn(&GeoIP{
IPInfo: IPInfoConfig{
Path: "../../../../internal/test/geoip/ipinfo_lite_sample.mmdb",
},
})
require.NoError(t, err)
info, err := lookupFn(net.IPv4(1, 7, 0, 17))
require.NoError(t, err)
assert.Equal(t, "AS9583", info.ASN)
assert.Equal(t, "IN", info.Country)
}

var fileFlag = flag.String("db", "../../../../internal/test/geoip/ipinfo_lite_sample.mmdb", "db to use for geoip benchmarks")

func BenchmarkDBLookup(b *testing.B) {
flag.Parse()
lookupFn, err := getLookupFn(&GeoIP{
IPInfo: IPInfoConfig{
Path: *fileFlag,
},
})
if err != nil {
b.Fatalf("failed to load database: %s", err.Error())
}
for _, addrSpace := range []uint32{256, 512, 1024, 2048, math.MaxUint32} {
b.Run(fmt.Sprintf("addr=%d", addrSpace), func(b *testing.B) {
for b.Loop() {
ipnum := rand.Uint32N(addrSpace)
bytes := make([]byte, 16)
binary.LittleEndian.PutUint32(bytes[12:], ipnum)
ip := ebpf.IPAddr(bytes)
_, err := lookupFn(ip.IP())
if err != nil {
b.Fatal(err.Error())
}
}
})
}
}

func BenchmarkDBLookupCached(b *testing.B) {
runBench := func(b *testing.B, cacheSize int, addrSpace uint32) {
cache := expirable.NewLRU[ebpf.IPAddr, ipInfo](cacheSize, nil, time.Hour)
lookupFn, err := getLookupFn(&GeoIP{
IPInfo: IPInfoConfig{
Path: *fileFlag,
},
})
if err != nil {
b.Fatalf("failed to load database: %s", err.Error())
}
lookups := 0
hits := 0
for b.Loop() {
lookups++
ipnum := rand.Uint32N(addrSpace)
bytes := make([]byte, 16)
binary.LittleEndian.PutUint32(bytes[12:], ipnum)
ip := ebpf.IPAddr(bytes)
_, ok := cache.Get(ip)
if !ok {
i, err := lookupFn(ip.IP())
if err != nil {
b.Fatal(err.Error())
}
cache.Add(ip, i)
} else {
hits++
}
}
}
for _, cacheSize := range []int{256, 512, 1024} {
for _, addrSpace := range []uint32{256, 512, 1024, math.MaxUint32} {
b.Run(fmt.Sprintf("cache=%d;addr=%d", cacheSize, addrSpace), func(b *testing.B) {
runBench(b, cacheSize, addrSpace)
})
}
}
}
6 changes: 5 additions & 1 deletion pkg/netolly/agent/pipeline.go
Original file line number Diff line number Diff line change
Expand Up @@ -76,8 +76,12 @@ func (f *Flows) buildPipeline(ctx context.Context) (*swarm.Runner, error) {
swi.Add(flow.ReverseDNSProvider(&f.cfg.NetworkFlows.ReverseDNS, kubeDecoratedFlows, dnsDecoratedFlows),
swarm.WithID("ReverseDNS"))

geoIPDecoratedFlows := newQueue("geoIPDecoratedFlows")
swi.Add(flow.GeoIPProvider(&f.cfg.NetworkFlows.GeoIP,
dnsDecoratedFlows, geoIPDecoratedFlows), swarm.WithID("GeoIPDecorator"))

cidrDecoratedFlows := newQueue("cidrDecoratedFlows")
swi.Add(cidr.DecoratorProvider(f.cfg.NetworkFlows.CIDRs, dnsDecoratedFlows, cidrDecoratedFlows),
swi.Add(cidr.DecoratorProvider(f.cfg.NetworkFlows.CIDRs, geoIPDecoratedFlows, cidrDecoratedFlows),
swarm.WithID("CIDRDecorator"))

decoratedFlows := newQueue("decoratedFlows")
Expand Down
Loading