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
45 changes: 37 additions & 8 deletions cmd/xtcp2/xtcp2.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"flag"
"fmt"
"log"
"net"
"net/http"
"os"
"os/signal"
Expand All @@ -30,6 +31,7 @@ import (
"github.com/pkg/profile"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promhttp"
"github.com/randomizedcoder/xtcp2/pkg/ipsockopt"
"github.com/randomizedcoder/xtcp2/pkg/misc"
"github.com/randomizedcoder/xtcp2/pkg/xtcp"
"github.com/randomizedcoder/xtcp2/pkg/xtcp_config"
Expand Down Expand Up @@ -137,6 +139,9 @@ const (
hostnameCst = ""
resolveContainerIdCst = false

ipv4TtlCst uint = 0
ipv6HopLimitCst uint = 0

deserializersCst = "all"

grpcPortCst = 8889
Expand Down Expand Up @@ -207,6 +212,8 @@ type mainFlags struct {
location *string
hostname *string
resolveContainerId *bool
ipv4Ttl *uint
ipv6HopLimit *uint
grpcPort *uint
deserializers *string
promListen *string
Expand Down Expand Up @@ -264,6 +271,8 @@ func defineFlags() *mainFlags {
f.location = flag.String("location", locationCst, "deployment grouping/facility this daemon runs in (data center, PoP, region, site, …); stamped on every record's `location`. Falls back to LOCATION env.")
f.hostname = flag.String("hostname", hostnameCst, "hostname stamped on records; defaults to os.Hostname(). Set this in a container, where os.Hostname() returns the container id, not the host. Falls back to XTCP_HOSTNAME env (NOT HOSTNAME).")
f.resolveContainerId = flag.Bool("resolveContainerId", resolveContainerIdCst, "resolve each socket's owning container id from its cgroup into container_id/container_runtime; needs /sys/fs/cgroup readable (mount it + --cgroupns=host in a container). Falls back to CONTAINER_ID_RESOLVE env.")
f.ipv4Ttl = flag.Uint("ipv4Ttl", ipv4TtlCst, "outgoing IPv4 TTL for xtcp2's TCP listeners (Prometheus + gRPC); 0 = kernel default. A low value keeps replies from travelling far if the host is internet-exposed. Falls back to IPV4_TTL env.")
f.ipv6HopLimit = flag.Uint("ipv6HopLimit", ipv6HopLimitCst, "outgoing IPv6 unicast hop limit for xtcp2's TCP listeners; 0 = kernel default. Falls back to IPV6_HOP_LIMIT env.")
f.grpcPort = flag.Uint("grpcPort", grpcPortCst, "GRPC listening port")
f.deserializers = flag.String("deserializers", deserializersCst, fmt.Sprintf("Comma separated list of deserializers,%v", xtcp.GetAllDeserializers()))
f.promListen = flag.String("promListen", promListenCst, "Prometheus http listening socket")
Expand Down Expand Up @@ -378,6 +387,8 @@ func buildConfig(f *mainFlags, des *xtcp_config.EnabledDeserializers) *xtcp_conf
Location: *f.location,
Hostname: *f.hostname,
ResolveContainerId: *f.resolveContainerId,
Ipv4Ttl: uint32(*f.ipv4Ttl),
Ipv6HopLimit: uint32(*f.ipv6HopLimit),
GrpcPort: uint32(*f.grpcPort),
EnabledDeserializers: des,

Expand Down Expand Up @@ -529,8 +540,8 @@ var daemonRunner = runDaemonDefault
// promHandlerStarter is the indirection point for the prom-handler
// goroutine launch. Default starts the real handler; tests swap it for
// a no-op to skip the port-bind.
var promHandlerStarter = func(promPath, promListen string) {
go initPromHandler(promPath, promListen)
var promHandlerStarter = func(promPath, promListen string, ipv4TTL, ipv6HopLimit uint32) {
go initPromHandler(promPath, promListen, ipv4TTL, ipv6HopLimit)
}

func main() {
Expand Down Expand Up @@ -582,7 +593,7 @@ func runMain(parentCtx context.Context) int {
debugLevel)()

environmentOverrideProm(f.promListen, f.promPath, debugLevel)
promHandlerStarter(*f.promPath, *f.promListen)
promHandlerStarter(*f.promPath, *f.promListen, c.Ipv4Ttl, c.Ipv6HopLimit)
if debugLevel > 10 {
log.Println("Prometheus http listener started on:", *f.promListen, *f.promPath)
}
Expand Down Expand Up @@ -670,30 +681,38 @@ func awaitSignalAndShutdown(
// ListenAndServe error branch is exercisable without exiting.
var fatalf = log.Fatalf

func initPromHandler(promPath string, promListen string) {
func initPromHandler(promPath string, promListen string, ipv4TTL, ipv6HopLimit uint32) {
http.Handle(promPath, promhttp.HandlerFor(
prometheus.DefaultGatherer,
promhttp.HandlerOpts{
EnableOpenMetrics: promEnableOpenMetrics,
MaxRequestsInFlight: promMaxRequestsInFlight,
},
))
go servePromHandler(promListen)
go servePromHandler(promListen, ipv4TTL, ipv6HopLimit)
}

// servePromHandler runs the prom HTTP server on promListen. On
// ListenAndServe failure it invokes fatalf (default log.Fatalf in
// production, swapped to a capture by tests). Extracted from
// initPromHandler so tests can drive the error path in isolation.
func servePromHandler(promListen string) {
func servePromHandler(promListen string, ipv4TTL, ipv6HopLimit uint32) {
srv := &http.Server{
Addr: promListen,
ReadHeaderTimeout: 5 * time.Second,
ReadTimeout: 10 * time.Second,
WriteTimeout: 10 * time.Second,
IdleTimeout: 30 * time.Second,
}
if err := srv.ListenAndServe(); err != nil {
// net.Listen (not srv.ListenAndServe) so the IPv4 TTL / IPv6 hop limit can
// be clamped on the listening socket before bind (inherited by accepted
// connections). ipsockopt.Control is nil when both are 0 → kernel default.
lc := net.ListenConfig{Control: ipsockopt.Control(ipv4TTL, ipv6HopLimit)}
ln, err := lc.Listen(context.Background(), "tcp", promListen)
if err != nil {
fatalf("prometheus error, listen: %v", err)
return
}
if err := srv.Serve(ln); err != nil {
fatalf("prometheus error: %v", err)
}
}
Expand Down Expand Up @@ -1040,6 +1059,14 @@ func envOverrideLabeling(c *xtcp_config.XtcpConfig, debugLevel uint) {
c.ResolveContainerId = v
logEnv("CONTAINER_ID_RESOLVE", fmt.Sprintf("c.ResolveContainerId:%t", v), debugLevel)
}
if v, ok := envUint32("IPV4_TTL"); ok {
c.Ipv4Ttl = v
logEnv("IPV4_TTL", fmt.Sprintf("c.Ipv4Ttl:%d", v), debugLevel)
}
if v, ok := envUint32("IPV6_HOP_LIMIT"); ok {
c.Ipv6HopLimit = v
logEnv("IPV6_HOP_LIMIT", fmt.Sprintf("c.Ipv6HopLimit:%d", v), debugLevel)
}
if v, ok := envUint32("GRPC_PORT"); ok {
c.GrpcPort = v
logEnv("GRPC_PORT", fmt.Sprintf("c.GrpcPort:%d", v), debugLevel)
Expand Down Expand Up @@ -1076,6 +1103,8 @@ func printConfig(c *xtcp_config.XtcpConfig, comment string) {
fmt.Println("c.KafkaSchemaUrl:", c.KafkaSchemaUrl)
fmt.Println("c.KafkaProduceTimeout:", c.KafkaProduceTimeout.AsDuration())
fmt.Println("c.DebugLevel:", c.DebugLevel)
fmt.Println("c.Ipv4Ttl:", c.Ipv4Ttl)
fmt.Println("c.Ipv6HopLimit:", c.Ipv6HopLimit)
fmt.Println("c.Label:", c.Label)
fmt.Println("c.Tag:", c.Tag)
fmt.Println("c.Location:", c.Location)
Expand Down
20 changes: 15 additions & 5 deletions cmd/xtcp2/xtcp2_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -307,11 +307,16 @@ func TestEnvOverrideLabeling(t *testing.T) {
t.Setenv("LOCATION", "eu-ro-1")
t.Setenv("XTCP_HOSTNAME", "runpod435")
t.Setenv("CONTAINER_ID_RESOLVE", "true")
t.Setenv("IPV4_TTL", "3")
t.Setenv("IPV6_HOP_LIMIT", "9")
t.Setenv("GRPC_PORT", "9000")
envOverrideLabeling(c, 0)
if c.Label != "prod" || c.Tag != "host=foo" || c.GrpcPort != 9000 {
t.Errorf("envOverrideLabeling mismatch: %+v", c)
}
if c.Ipv4Ttl != 3 || c.Ipv6HopLimit != 9 {
t.Errorf("envOverrideLabeling ttl mismatch: Ipv4Ttl=%d Ipv6HopLimit=%d", c.Ipv4Ttl, c.Ipv6HopLimit)
}
if c.Location != "eu-ro-1" || c.Hostname != "runpod435" || !c.ResolveContainerId {
t.Errorf("envOverrideLabeling identity mismatch: Location=%q Hostname=%q Resolve=%v", c.Location, c.Hostname, c.ResolveContainerId)
}
Expand Down Expand Up @@ -415,7 +420,7 @@ func TestServePromHandler_bindError(t *testing.T) {
}
t.Cleanup(func() { fatalf = prev })

servePromHandler("invalid-host:-1")
servePromHandler("invalid-host:-1", 0, 0)
if !strings.Contains(captured, "prometheus error") {
t.Errorf("fatalf not invoked; got %q", captured)
}
Expand All @@ -432,7 +437,7 @@ func TestRunMain_version(t *testing.T) {

// Stub the prom handler starter so it doesn't bind a port.
prevProm := promHandlerStarter
promHandlerStarter = func(_, _ string) {}
promHandlerStarter = func(_, _ string, _, _ uint32) {}
t.Cleanup(func() { promHandlerStarter = prevProm })

// runMain spawns a signal-handler goroutine that blocks on signal.Notify.
Expand All @@ -452,7 +457,7 @@ func TestRunMain_conf(t *testing.T) {
t.Cleanup(func() { os.Args = prevArgs })

prevProm := promHandlerStarter
promHandlerStarter = func(_, _ string) {}
promHandlerStarter = func(_, _ string, _, _ uint32) {}
t.Cleanup(func() { promHandlerStarter = prevProm })

captureLog(t, func() {
Expand All @@ -471,7 +476,7 @@ func TestRunMain_stubbedDaemon(t *testing.T) {
t.Cleanup(func() { os.Args = prevArgs })

prevProm := promHandlerStarter
promHandlerStarter = func(_, _ string) {}
promHandlerStarter = func(_, _ string, _, _ uint32) {}
t.Cleanup(func() { promHandlerStarter = prevProm })

prevDaemon := daemonRunner
Expand Down Expand Up @@ -500,7 +505,7 @@ func TestInitPromHandler_smoke(t *testing.T) {
fatalf = func(string, ...any) {} // swallow
t.Cleanup(func() { fatalf = prevFatalf })

initPromHandler("/metrics", ":0")
initPromHandler("/metrics", ":0", 0, 0)
time.Sleep(10 * time.Millisecond)
}

Expand Down Expand Up @@ -773,6 +778,8 @@ func TestBuildConfig(t *testing.T) {
label := "lbl"
tag := "host=a"
gp := uint(8888)
ttl := uint(3)
hop := uint(9)
pl := ":9088"
pp := "/metrics"
gmp := uint(8)
Expand Down Expand Up @@ -809,6 +816,7 @@ func TestBuildConfig(t *testing.T) {
dest: &dst, destWriteFiles: &dwf,
topic: &topic, xtcpProtoFile: &xp, kafkaSchemaUrl: &ksu,
produceTimeout: &pto, label: &label, tag: &tag, grpcPort: &gp,
ipv4Ttl: &ttl, ipv6HopLimit: &hop,
deserializers: &ds, promListen: &pl, promPath: &pp, goMaxProcs: &gmp,
profileMode: &pm, v: &v, conf: &conf, d: &d,
ioUring: &iu, ioUringRecvBatch: &iurb, ioUringCqeBatch: &iucb,
Expand Down Expand Up @@ -843,6 +851,8 @@ func TestBuildConfig(t *testing.T) {
{"Hostname", c.Hostname, "protoText"},
{"ResolveContainerId", c.ResolveContainerId, true},
{"GrpcPort", c.GrpcPort, uint32(8888)},
{"Ipv4Ttl", c.Ipv4Ttl, uint32(3)},
{"Ipv6HopLimit", c.Ipv6HopLimit, uint32(9)},
{"S3SkipBucketProbe", c.S3SkipBucketProbe, true},
}
for _, ck := range checks {
Expand Down
71 changes: 53 additions & 18 deletions dart/xtcp_config/v1/xtcp_config.pb.dart

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading