diff --git a/config/sample_webconfig.conf b/config/sample_webconfig.conf index c8b4417..075dc47 100644 --- a/config/sample_webconfig.conf +++ b/config/sample_webconfig.conf @@ -3,7 +3,9 @@ webconfig { encryption_key_env_name = "WEBCONFIG_KEY" } + app_name = "webconfig" panic_exit_enabled = false + factory_reset_enabled = false tracing { moracide_tag_prefix = "X-Cl-Experiment" otel { @@ -36,6 +38,8 @@ webconfig { level = "info" file = "/tmp/webconfig.log" sarama_logger_enabled = false + // List of request header names to suppress from logs (case-insensitive) + ignored_headers = [] } metrics { @@ -104,6 +108,8 @@ webconfig { // Defaults to false. Setting this to true disables server certificate // verification for all outbound HTTPS connectors and is insecure. tls_insecure_skip_verify = false + // Optional User-Agent header sent on all outbound HTTP requests. + user_agent = "" } jwt { @@ -112,7 +118,6 @@ webconfig { // north bound APIs from orc/customers, no mac embedded in tokens api_token { - enabled = true kids = [ "sat-prod-k1-1024", ] @@ -126,7 +131,6 @@ webconfig { // south bound APIs from CPEs, macs embedded in tokens cpe_token { - enabled = true kids = [ "webconfig_key", ] @@ -244,20 +248,58 @@ webconfig { } // TLS configuration for secure Kafka connections - tls { - enabled = false - // Path to client certificate file (PEM format) - // NOTE: Certificate files (cert, key, ca) are REQUIRED when insecure_skip_verify=false (default) - // Certificate files are OPTIONAL when insecure_skip_verify=true (insecure mode) - cert_file = "/etc/webconfig/kafka/client-cert.pem" - // Path to client private key file (PEM format) - key_file = "/etc/webconfig/kafka/client-key.pem" - // Path to CA certificate file for broker verification (PEM format) - ca_cert_file = "/etc/webconfig/kafka/ca-cert.pem" - // Skip certificate verification (WARNING: use ONLY in dev/test environments, NEVER in production) - // When true, certificate files become optional (TLS without client authentication) - insecure_skip_verify = false - } + // NOTE: These are flat keys read by LoadKafkaTLSConfig at prefix "webconfig.kafka" + // Certificate files (tls_cert_file, tls_key_file, tls_ca_cert_file) are REQUIRED when + // tls_insecure_skip_verify=false (default). They are OPTIONAL when tls_insecure_skip_verify=true. + tls_enabled = false + tls_cert_file = "/etc/webconfig/kafka/client-cert.pem" + tls_key_file = "/etc/webconfig/kafka/client-key.pem" + tls_ca_cert_file = "/etc/webconfig/kafka/ca-cert.pem" + // WARNING: set true only in dev/test environments, NEVER in production + tls_insecure_skip_verify = false + + // Reconnect resilience — tuned for broker restart recovery. + + // metadata.retry_max: max metadata fetch retries before giving up. + // lib default: 3 | configured: 10 + // Extends the retry window to 10×2s = 20s, covering the cluster + // stabilisation window during leader election after a broker restart. + metadata.retry_max = 10 + + // metadata.retry_backoff_sec: seconds between metadata fetch retries. + // lib default: 0.25s | configured: 2s + // Prevents rapid-fire retries against a restarting broker; gives the + // cluster time to stabilise between attempts. + metadata.retry_backoff_sec = 2 + + // metadata.refresh_frequency_sec: background metadata refresh interval in seconds. + // lib default: 600s (10min) | configured: 300s (5min) + // Proactive background refresh catches stale partition leaders up to 5min + // sooner than the library default after a broker restart. + metadata.refresh_frequency_sec = 300 + + // consumer.session_timeout_sec: consumer group session timeout in seconds. + // lib default: 10s | configured: 30s + // Must be within broker group.min.session.timeout.ms (6s) / + // group.max.session.timeout.ms (300s). 30s gives the consumer group + // time to survive a restarting group coordinator without triggering + // an unnecessary rebalance. + // Heartbeat interval is derived as session_timeout_sec / 6 (= 5s for + // the default 30s). The Kafka protocol ceiling is session_timeout / 3 + // (10s); using /6 gives double the margin against transient jitter. + // lib heartbeat default: 3s | configured: 5s + consumer.session_timeout_sec = 30 + + // consumer.retry_backoff_sec: seconds a partition reader waits before retrying a failed fetch. + // lib default: 2s | configured: 2s + // Set explicitly to make the value visible and tunable. + consumer.retry_backoff_sec = 2 + + // net.timeout_sec: dial/read/write timeout for broker connections in seconds. + // lib default: 30s | configured: 10s + // Tighter timeout lets the retry loop engage within 10s per attempt + // instead of blocking for 30s (fail-fast on broken broker connections). + net.timeout_sec = 10 // if we want to use more than 1 cluster clusters { @@ -274,16 +316,12 @@ webconfig { messages_per_second = 10 } - // TLS configuration for mesh cluster - tls { - enabled = false - // NOTE: Certificate files (cert, key, ca) are REQUIRED when insecure_skip_verify=false (default) - // Certificate files are OPTIONAL when insecure_skip_verify=true (insecure mode) - cert_file = "/etc/webconfig/kafka/mesh-client-cert.pem" - key_file = "/etc/webconfig/kafka/mesh-client-key.pem" - ca_cert_file = "/etc/webconfig/kafka/mesh-ca-cert.pem" - insecure_skip_verify = false - } + // TLS configuration for mesh cluster (flat keys at prefix "webconfig.kafka.clusters.mesh") + tls_enabled = false + tls_cert_file = "/etc/webconfig/kafka/mesh-client-cert.pem" + tls_key_file = "/etc/webconfig/kafka/mesh-client-key.pem" + tls_ca_cert_file = "/etc/webconfig/kafka/mesh-ca-cert.pem" + tls_insecure_skip_verify = false } east { enabled = false @@ -298,16 +336,12 @@ webconfig { messages_per_second = 10 } - // TLS configuration for east cluster - tls { - enabled = false - // NOTE: Certificate files (cert, key, ca) are REQUIRED when insecure_skip_verify=false (default) - // Certificate files are OPTIONAL when insecure_skip_verify=true (insecure mode) - cert_file = "/etc/webconfig/kafka/east-client-cert.pem" - key_file = "/etc/webconfig/kafka/east-client-key.pem" - ca_cert_file = "/etc/webconfig/kafka/east-ca-cert.pem" - insecure_skip_verify = false - } + // TLS configuration for east cluster (flat keys at prefix "webconfig.kafka.clusters.east") + tls_enabled = false + tls_cert_file = "/etc/webconfig/kafka/east-client-cert.pem" + tls_key_file = "/etc/webconfig/kafka/east-client-key.pem" + tls_ca_cert_file = "/etc/webconfig/kafka/east-ca-cert.pem" + tls_insecure_skip_verify = false } } } @@ -327,22 +361,59 @@ webconfig { // correct subdoc states if versions match but not "deployed" state_correction_enabled= false + // pre-cook supplementary subdocs into the database + supplementary_precook_enabled = false + // number of days to retain supplementary precook state (default 7) + supplementary_precook_state_ttl_days = 7 + // forward kafka messages if needed kafka_producer { enabled = false brokers = "localhost:9092" topic = "webconfig_downstream" - // TLS configuration for Kafka producer - tls { - enabled = false - // NOTE: Certificate files (cert, key, ca) are REQUIRED when insecure_skip_verify=false (default) - // Certificate files are OPTIONAL when insecure_skip_verify=true (insecure mode) - cert_file = "/etc/webconfig/kafka/producer-client-cert.pem" - key_file = "/etc/webconfig/kafka/producer-client-key.pem" - ca_cert_file = "/etc/webconfig/kafka/producer-ca-cert.pem" - insecure_skip_verify = false - } + // TLS configuration for Kafka producer (flat keys at prefix "webconfig.kafka_producer") + tls_enabled = false + tls_cert_file = "/etc/webconfig/kafka/producer-client-cert.pem" + tls_key_file = "/etc/webconfig/kafka/producer-client-key.pem" + tls_ca_cert_file = "/etc/webconfig/kafka/producer-ca-cert.pem" + tls_insecure_skip_verify = false + + // Reconnect resilience — tuned for broker restart recovery. + + // metadata.retry_max: max metadata fetch retries before giving up. + // lib default: 3 | configured: 10 + // Extends the retry window to 10×2s = 20s, covering the cluster + // stabilisation window during leader election after a broker restart. + metadata.retry_max = 10 + + // metadata.retry_backoff_sec: seconds between metadata fetch retries. + // lib default: 0.25s | configured: 2s + // Prevents rapid-fire retries against a restarting broker. + metadata.retry_backoff_sec = 2 + + // metadata.refresh_frequency_sec: background metadata refresh interval in seconds. + // lib default: 600s (10min) | configured: 300s (5min) + // Proactive background refresh catches stale partition leaders up to 5min + // sooner than the library default after a broker restart. + metadata.refresh_frequency_sec = 300 + + // producer.retry_max: max internal send retries before a message is put on the error channel. + // lib default: 3 | configured: 10 + // Mirrors metadata.retry_max — gives the producer the same tolerance for + // transient broker unavailability as the consumer has for metadata fetches. + producer.retry_max = 10 + + // producer.retry_backoff_sec: seconds between internal send retries. + // lib default: 0.1s | configured: 2s + // Matches metadata.retry_backoff_sec cadence; avoids hammering a restarting broker. + producer.retry_backoff_sec = 2 + + // net.timeout_sec: dial/read/write timeout for broker connections in seconds. + // lib default: 30s | configured: 10s + // Tighter timeout lets the retry loop engage within 10s per attempt + // instead of blocking for 30s (fail-fast on broken broker connections). + net.timeout_sec = 10 } // this allows the root document locked if needed @@ -363,4 +434,7 @@ webconfig { filter_output_by_bitmap_enabled = false bitmap_filter_exempt_subdoc_ids = [] + + // return an empty profile document when no upstream data is available + default_empty_profile_enabled = false } diff --git a/db/sqlite/sqlite_client.go b/db/sqlite/sqlite_client.go index 01b1e1e..7b1b5d9 100644 --- a/db/sqlite/sqlite_client.go +++ b/db/sqlite/sqlite_client.go @@ -21,6 +21,7 @@ import ( "database/sql" "errors" "fmt" + "os" "github.com/go-akka/configuration" "github.com/rdkcentral/webconfig/common" @@ -56,6 +57,9 @@ func NewSqliteClient(conf *configuration.Config, testOnly bool) (*SqliteClient, var dbfile string if testOnly { dbfile = conf.GetString("webconfig.database.sqlite.unittest_db_file", defaultSqliteTestDbFile) + if x := os.Getenv("WEBCONFIG_TESTDB_SQLITE_FILE"); len(x) > 0 { + dbfile = x + } } else { dbfile = conf.GetString("webconfig.database.sqlite.db_file", defaultSqliteDbFile) } diff --git a/go.mod b/go.mod index abe8783..1a7ca88 100644 --- a/go.mod +++ b/go.mod @@ -66,14 +66,14 @@ require ( go.opentelemetry.io/otel/metric v1.43.0 // indirect go.opentelemetry.io/proto/otlp v1.10.0 // indirect go.yaml.in/yaml/v2 v2.4.4 // indirect - golang.org/x/crypto v0.50.0 // indirect - golang.org/x/net v0.53.0 // indirect - golang.org/x/sys v0.43.0 // indirect - golang.org/x/text v0.36.0 // indirect + golang.org/x/crypto v0.52.0 // indirect + golang.org/x/net v0.55.0 // indirect + golang.org/x/sys v0.45.0 // indirect + golang.org/x/text v0.37.0 // indirect google.golang.org/appengine v1.6.8 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260504160031-60b97b32f348 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260504160031-60b97b32f348 // indirect - google.golang.org/grpc v1.81.0 // indirect + google.golang.org/grpc v1.82.1 // indirect google.golang.org/protobuf v1.36.11 // indirect gopkg.in/inf.v0 v0.9.1 // indirect modernc.org/libc v1.72.2 // indirect diff --git a/go.sum b/go.sum index 1adcd3b..6af8102 100644 --- a/go.sum +++ b/go.sum @@ -169,11 +169,11 @@ go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= golang.org/x/crypto v0.6.0/go.mod h1:OFC/31mSvZgRz0V1QTNCzfAI1aIRzbiufJtkMIlEp58= -golang.org/x/crypto v0.50.0 h1:zO47/JPrL6vsNkINmLoo/PH1gcxpls50DNogFvB5ZGI= -golang.org/x/crypto v0.50.0/go.mod h1:3muZ7vA7PBCE6xgPX7nkzzjiUq87kRItoJQM1Yo8S+Q= +golang.org/x/crypto v0.52.0 h1:RMs7fP2rXdep0CftQlK8Uf+kibLm7qkCcradZWYz988= +golang.org/x/crypto v0.52.0/go.mod h1:1QgfPxDqh0T2M/elOJtp9RvuR95kVjir0e6/BvEmGbc= golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= -golang.org/x/mod v0.34.0 h1:xIHgNUUnW6sYkcM5Jleh05DvLOtwc6RitGHbDk4akRI= -golang.org/x/mod v0.34.0/go.mod h1:ykgH52iCZe79kzLLMhyCUzhMci+nQj+0XkbXpNYtVjY= +golang.org/x/mod v0.35.0 h1:Ww1D637e6Pg+Zb2KrWfHQUnH2dQRLBQyAtpr/haaJeM= +golang.org/x/mod v0.35.0/go.mod h1:+GwiRhIInF8wPm+4AoT6L0FA1QWAad3OMdTRx4tFYlU= golang.org/x/net v0.0.0-20190603091049-60506f45cf65/go.mod h1:HSz+uSET+XFnRR8LxR5pz3Of3rY3CfYBVs4xY44aLks= golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20200114155413-6afb5195e5aa/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= @@ -182,8 +182,8 @@ golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c= golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs= golang.org/x/net v0.7.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs= -golang.org/x/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA= -golang.org/x/net v0.53.0/go.mod h1:JvMuJH7rrdiCfbeHoo3fCQU24Lf5JJwT9W3sJFulfgs= +golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8= +golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww= golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= @@ -194,8 +194,8 @@ golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBc golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.43.0 h1:Rlag2XtaFTxp19wS8MXlJwTvoh8ArU6ezoyFsMyCTNI= -golang.org/x/sys v0.43.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY= +golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= golang.org/x/term v0.5.0/go.mod h1:jMB1sMXY+tzblOD4FWmEbocvup2/aLOaQEp7JmGp78k= @@ -205,13 +205,13 @@ golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= golang.org/x/text v0.3.8/go.mod h1:E6s5w1FMmriuDzIBO73fBruAKo1PCIq6d2Q6DHfQ8WQ= golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= -golang.org/x/text v0.36.0 h1:JfKh3XmcRPqZPKevfXVpI1wXPTqbkE5f7JA92a55Yxg= -golang.org/x/text v0.36.0/go.mod h1:NIdBknypM8iqVmPiuco0Dh6P5Jcdk8lJL0CUebqK164= +golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc= +golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= -golang.org/x/tools v0.43.0 h1:12BdW9CeB3Z+J/I/wj34VMl8X+fEXBxVR90JeMX5E7s= -golang.org/x/tools v0.43.0/go.mod h1:uHkMso649BX2cZK6+RpuIPXS3ho2hZo4FVwfoy1vIk0= +golang.org/x/tools v0.44.0 h1:UP4ajHPIcuMjT1GqzDWRlalUEoY+uzoZKnhOjbIPD2c= +golang.org/x/tools v0.44.0/go.mod h1:KA0AfVErSdxRZIsOVipbv3rQhVXTnlU6UhKxHd1seDI= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= @@ -223,8 +223,8 @@ google.golang.org/genproto/googleapis/api v0.0.0-20260504160031-60b97b32f348 h1: google.golang.org/genproto/googleapis/api v0.0.0-20260504160031-60b97b32f348/go.mod h1:Yzdzr5OOZFgSsEV2D/Xi9NL3bszpXFAg0hFJiRohcD8= google.golang.org/genproto/googleapis/rpc v0.0.0-20260504160031-60b97b32f348 h1:pfIbyB44sWzHiCpRqIen67ZQnVXSfIxWrqUMk1qwODE= google.golang.org/genproto/googleapis/rpc v0.0.0-20260504160031-60b97b32f348/go.mod h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8= -google.golang.org/grpc v1.81.0 h1:W3G9N3KQf3BU+YuCtGKJk0CmxQNbAISICD/9AORxLIw= -google.golang.org/grpc v1.81.0/go.mod h1:xGH9GfzOyMTGIOXBJmXt+BX/V0kcdQbdcuwQ/zNw42I= +google.golang.org/grpc v1.82.1 h1:NnAxzGRA0677vCa4BUkOAnO5+FfQqVl9iUXeD0IqcGE= +google.golang.org/grpc v1.82.1/go.mod h1:yzTZ1TB1Z3SG+LIYaI+WiE8D5+PZ3ArnrSp8zF3+/ZA= google.golang.org/protobuf v1.26.0-rc.1/go.mod h1:jlhhOSvTdKEhbULTjvd4ARK9grFBp09yW+WbY/TyQbw= google.golang.org/protobuf v1.26.0/go.mod h1:9q0QmTI4eRPtz6boOQmLYwt+qCgq0jsYwAQnmE0givc= google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= diff --git a/http/webconfig_server.go b/http/webconfig_server.go index 633201b..dbbbe0b 100644 --- a/http/webconfig_server.go +++ b/http/webconfig_server.go @@ -287,7 +287,7 @@ func NewWebconfigServer(sc *common.ServerConfig, testOnly bool) *WebconfigServer kafkaEnabled := conf.GetBoolean("webconfig.kafka.enabled") upstreamEnabled := conf.GetBoolean("webconfig.upstream.enabled") - appName := conf.GetString("webconfig.app_name") + appName := conf.GetString("webconfig.app_name", "webconfig") validateMacEnabled := conf.GetBoolean("webconfig.validate_device_id_as_mac_address", tokenApiEnabledDefault) configValidPartners := conf.GetStringList("webconfig.valid_partners") validPartners := []string{} @@ -314,6 +314,41 @@ func NewWebconfigServer(sc *common.ServerConfig, testOnly bool) *WebconfigServer saramaConfig := sarama.NewConfig() saramaConfig.Producer.Return.Errors = true + // Resilience tuning — survive Kafka broker restarts without restarting the app. + // + // Metadata.Retry.Max: lib default 3 → now 10 (configurable) + // Extends retry window to 10×2s = 20s, covering the cluster stabilisation + // window during leader election after a broker restart. + saramaConfig.Metadata.Retry.Max = int(conf.GetInt32("webconfig.kafka_producer.metadata.retry_max", 10)) + + // Metadata.Retry.Backoff: lib default 250ms → now 2s (configurable) + // Prevents rapid-fire retries against a restarting broker; paced to give + // the cluster time to stabilise between attempts. + saramaConfig.Metadata.Retry.Backoff = time.Duration(conf.GetInt32("webconfig.kafka_producer.metadata.retry_backoff_sec", 2)) * time.Second + + // Metadata.RefreshFrequency: lib default 10min → now 5min (configurable) + // Proactive background refresh. Catches stale partition leaders (caused by + // broker restarts) up to 5min sooner than the library default. + saramaConfig.Metadata.RefreshFrequency = time.Duration(conf.GetInt32("webconfig.kafka_producer.metadata.refresh_frequency_sec", 300)) * time.Second + + // Producer.Retry.Max: lib default 3 → now 10 (configurable) + // Internal send retries before a message is put on the Errors() channel. + // Mirrors Metadata.Retry.Max — gives the producer the same tolerance for + // transient broker unavailability as the consumer has for metadata fetches. + saramaConfig.Producer.Retry.Max = int(conf.GetInt32("webconfig.kafka_producer.producer.retry_max", 10)) + + // Producer.Retry.Backoff: lib default 100ms → now 2s (configurable) + // Matches Metadata.Retry.Backoff cadence; avoids hammering a restarting broker. + saramaConfig.Producer.Retry.Backoff = time.Duration(conf.GetInt32("webconfig.kafka_producer.producer.retry_backoff_sec", 2)) * time.Second + + // Net.DialTimeout / ReadTimeout / WriteTimeout: lib default 30s → now 10s (configurable) + // Tighter timeouts let the retry loop engage within 10s per attempt instead + // of blocking for 30s. Fail-fast on broken broker connections. + producerNetTimeout := time.Duration(conf.GetInt32("webconfig.kafka_producer.net.timeout_sec", 10)) * time.Second + saramaConfig.Net.DialTimeout = producerNetTimeout + saramaConfig.Net.ReadTimeout = producerNetTimeout + saramaConfig.Net.WriteTimeout = producerNetTimeout + // Load TLS configuration for producer tlsConfig, err := common.LoadKafkaTLSConfig(conf, "webconfig.kafka_producer") if err != nil { @@ -1085,6 +1120,21 @@ func (s *WebconfigServer) ForwardKafkaMessage(kbytes []byte, m *common.EventMess Key: sarama.ByteEncoder(kbytes), Value: sarama.ByteEncoder(bbytes), } + + defer func() { + if r := recover(); r != nil { + if fmt.Sprint(r) != "send on closed channel" { + panic(r) + } + if m := s.Metrics(); m != nil { + m.ObserveKafkaProducerErr(s.KafkaProducerTopic(), -1) + } + tfields["logger"] = "kafkaproducer" + tfields["error"] = r + log.WithFields(tfields).Warn("dropped: producer closed during shutdown") + } + }() + s.Input() <- outMessage tfields["logger"] = "kafkaproducer" @@ -1119,7 +1169,24 @@ func (s *WebconfigServer) ForwardSuccessKafkaMessages(messages []common.EventMes Key: sarama.ByteEncoder(strings.ToLower(mac)), Value: sarama.ByteEncoder(bbytes), } - s.Input() <- outMessage + + sent := func() (ok bool) { + defer func() { + if r := recover(); r != nil { + if fmt.Sprint(r) != "send on closed channel" { + panic(r) + } + tfields["error"] = r + log.WithFields(tfields).Warn("dropped: producer closed during shutdown") + ok = false + } + }() + s.Input() <- outMessage + return true + }() + if !sent { + return + } tfields["output_key"] = mac tfields["output_body"] = m @@ -1165,14 +1232,19 @@ func (s *WebconfigServer) LogToken(xw *XResponseWriter, authorization, token str log.WithFields(tfields).Debug(tokenErr) } -func (s *WebconfigServer) HandleKafkaProducerResults() { +func (s *WebconfigServer) HandleKafkaProducerResults(ctx context.Context) { if s.AsyncProducer == nil { return } for { select { - case success := <-s.Successes(): + case <-ctx.Done(): + return + case success, ok := <-s.Successes(): + if !ok { + return + } if success == nil { continue } @@ -1182,7 +1254,10 @@ func (s *WebconfigServer) HandleKafkaProducerResults() { fields["output_partition"] = success.Partition fields["output_offset"] = success.Offset log.WithFields(fields).Debug("sent") - case pErr := <-s.Errors(): + case pErr, ok := <-s.Errors(): + if !ok { + return + } if pErr == nil || pErr.Msg == nil { continue } diff --git a/kafka/kafka_consumer_group.go b/kafka/kafka_consumer_group.go index 9cc24e9..cb915fe 100644 --- a/kafka/kafka_consumer_group.go +++ b/kafka/kafka_consumer_group.go @@ -29,6 +29,14 @@ import ( wchttp "github.com/rdkcentral/webconfig/http" ) +const ( + // Resilience defaults — configurable via webconfig.kafka.* keys. + defaultMetadataRetryMax = 10 + defaultMetadataRetryBackoffSec = 2 + defaultNetTimeoutSec = 10 + defaultSessionTimeoutSec = 30 +) + type KafkaConsumerGroup struct { sarama.ConsumerGroup db.DatabaseClient @@ -103,6 +111,48 @@ func NewKafkaConsumerGroup(conf *configuration.Config, s *wchttp.WebconfigServer } } + // Resilience tuning — survive Kafka broker restarts without restarting the app. + // + // Metadata.Retry.Max: lib default 3 → now 10 (configurable) + // Extends retry window to 10×2s = 20s, covering the cluster stabilisation + // window during leader election after a broker restart. + sconfig.Metadata.Retry.Max = int(conf.GetInt32(prefix+".metadata.retry_max", defaultMetadataRetryMax)) + + // Metadata.Retry.Backoff: lib default 250ms → now 2s (configurable) + // Prevents rapid-fire retries against a restarting broker; paced to give + // the cluster time to stabilise between attempts. + sconfig.Metadata.Retry.Backoff = time.Duration(conf.GetInt32(prefix+".metadata.retry_backoff_sec", defaultMetadataRetryBackoffSec)) * time.Second + + // Metadata.RefreshFrequency: lib default 10min → now 5min (configurable) + // Proactive background refresh. Catches stale partition leaders (caused by + // broker restarts) up to 5min sooner than the library default. + sconfig.Metadata.RefreshFrequency = time.Duration(conf.GetInt32(prefix+".metadata.refresh_frequency_sec", 300)) * time.Second + + // Consumer.Group.Session.Timeout: lib default 10s → now 30s (configurable) + // Time before the broker evicts a consumer whose heartbeats have stopped. + // 30s gives the consumer group time to survive a restarting group coordinator + // without triggering an unnecessary rebalance. + // + // Consumer.Group.Heartbeat.Interval: lib default 3s → now Session.Timeout/6 = 5s + // The Kafka protocol ceiling is Session.Timeout/3 (10s). Using /6 keeps the + // interval at half the ceiling, providing headroom for transient network jitter + // during broker restarts without risking false session expiry. + sessionTimeout := time.Duration(conf.GetInt32(prefix+".consumer.session_timeout_sec", defaultSessionTimeoutSec)) * time.Second + sconfig.Consumer.Group.Session.Timeout = sessionTimeout + sconfig.Consumer.Group.Heartbeat.Interval = sessionTimeout / 6 + + // Consumer.Retry.Backoff: lib default 2s (configurable) + // How long a partition reader waits before retrying a failed fetch. + sconfig.Consumer.Retry.Backoff = time.Duration(conf.GetInt32(prefix+".consumer.retry_backoff_sec", 2)) * time.Second + + // Net.DialTimeout / ReadTimeout / WriteTimeout: lib default 30s → now 10s (configurable) + // Tighter timeouts let the retry loop engage within 10s per attempt instead + // of blocking for 30s. Fail-fast on broken broker connections. + netTimeout := time.Duration(conf.GetInt32(prefix+".net.timeout_sec", defaultNetTimeoutSec)) * time.Second + sconfig.Net.DialTimeout = netTimeout + sconfig.Net.ReadTimeout = netTimeout + sconfig.Net.WriteTimeout = netTimeout + // Load TLS configuration tlsConfig, err := common.LoadKafkaTLSConfig(conf, prefix) if err != nil { diff --git a/main.go b/main.go index a6e5f80..ced6179 100644 --- a/main.go +++ b/main.go @@ -23,6 +23,7 @@ import ( "fmt" "os" "os/signal" + "sync" "syscall" "time" @@ -107,9 +108,6 @@ func main() { router.Handle("/metrics", promhttp.Handler()) metrics = common.NewMetrics(sc.Config) server.SetMetrics(metrics) - if server.KafkaProducerEnabled() { - go server.HandleKafkaProducerResults() - } handler := metrics.WebMetrics(router) server.Handler = handler } else { @@ -119,6 +117,10 @@ func main() { // setup contexts groups g, gCtx := errgroup.WithContext(mainCtx) + if server.KafkaProducerEnabled() { + go server.HandleKafkaProducerResults(gCtx) + } + // setup http server g.Go( func() error { @@ -126,54 +128,79 @@ func main() { }, ) - g.Go( - func() error { - <-gCtx.Done() - fmt.Printf("HTTP server shutdown NOW !!\n") - if server.KafkaProducerEnabled() { - if err := server.AsyncProducer.Close(); err != nil { - fmt.Fprintf(os.Stderr, "%v AsyncProducer.Close() err=%v\n", time.Now().Format(common.LoggingTimeFormat), err) - } - } - return server.Shutdown(context.Background()) - }, - ) - // setup kafka consumer, if config kafka.enabled=false, then kcgroup=nil, err=nil kcgroups, err := kafka.NewKafkaConsumerGroups(sc, server, metrics) if err != nil { panic(err) } + // consumeWg tracks active consume loops so shutdown can wait for them to + // finish before closing the AsyncProducer (prevents "send on closed channel"). + var consumeWg sync.WaitGroup + for _, kcgroup := range kcgroups { - consumer := *(kcgroup.Consumer()) + consumer := kcgroup.Consumer() topics := kcgroup.Topics() + consumeWg.Add(1) g.Go( func() error { + defer consumeWg.Done() for { - if err := kcgroup.Consume(gCtx, topics, &consumer); err != nil { - fmt.Printf("kcgroup.Consumer: err=%v\n", err) - return err + if err := kcgroup.Consume(gCtx, topics, consumer); err != nil { + select { + case <-gCtx.Done(): + fmt.Printf("kcgroup.Consumer: topics=|%v|, shutdown gracefully\n", topics) + return nil + default: + fmt.Printf("kcgroup.Consumer: topics=|%v|, err=%v, retrying in 2s\n", topics, err) + select { + case <-time.After(2 * time.Second): + case <-gCtx.Done(): + return nil + } + } } consumer.Ready = make(chan bool) } }, ) - // This is to setup notify AFTER the sarama is running - // it is more or less optional, without this reading from the chan, - // the consumer runs anyway. - <-consumer.Ready - - g.Go( - func() error { - <-gCtx.Done() - fmt.Printf("SARAMA shutdown NOW !!\n") - return kcgroup.Close() - }, - ) + // Wait for the first session to be established before starting the next + // consumer group. Context-aware: if Kafka is unreachable on startup and + // the context is cancelled (e.g. SIGTERM), we stop waiting rather than + // deadlocking until Setup() is eventually called. + select { + case <-consumer.Ready: + case <-gCtx.Done(): + } } + // Single shutdown goroutine with ordered cleanup: + // 1. close consumer groups (signals ConsumeClaim goroutines to stop) + // 2. wait for all consume loops to drain (ensures no goroutine is mid-send) + // 3. close AsyncProducer (safe now that no one sends to its Input channel) + // 4. shut down HTTP server + g.Go( + func() error { + <-gCtx.Done() + fmt.Printf("shutdown: initiating consumer group shutdown\n") + for _, kcgroup := range kcgroups { + if err := kcgroup.Close(); err != nil { + fmt.Fprintf(os.Stderr, "kcgroup.Close() err=%v\n", err) + } + } + consumeWg.Wait() + fmt.Printf("shutdown: all consumer loops drained\n") + if server.KafkaProducerEnabled() { + if err := server.AsyncProducer.Close(); err != nil { + fmt.Fprintf(os.Stderr, "%v AsyncProducer.Close() err=%v\n", time.Now().Format(common.LoggingTimeFormat), err) + } + } + fmt.Printf("shutdown: stopping HTTP server\n") + return server.Shutdown(context.Background()) + }, + ) + if err := g.Wait(); err != nil { fmt.Printf("exit reason: %s \n", err) }