From bcfc0af0777d8b9b7ed48d05f251d7a9e8d3ce8e Mon Sep 17 00:00:00 2001 From: Rico Date: Sun, 3 Aug 2025 19:12:26 +0200 Subject: [PATCH 01/30] docs: fix typo in docker-compose config --- docker/docker-compose.metrics.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docker/docker-compose.metrics.yml b/docker/docker-compose.metrics.yml index 316a8ec..301a165 100644 --- a/docker/docker-compose.metrics.yml +++ b/docker/docker-compose.metrics.yml @@ -1,6 +1,6 @@ version: '2' -# Make sure to create the sub directories "prometheus", "prometheus_data", "grafana", "grafana_data" and "certstream" +# Make sure to create the subdirectories "prometheus", "prometheus_data", "grafana", "grafana_data" and "certstream" # and create the config files for all three services. For further details please refer to https://github.com/d-Rickyy-b/certstream-server-go/wiki/Collecting-and-Visualizing-Metrics networks: From 6865b6b20445a0ab48f7923abbe5d0d658678182 Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 4 Aug 2025 01:57:24 +0200 Subject: [PATCH 02/30] feat: first implementation for streaming proessors such as kafka Relates to #23 and #35 --- config.sample.yaml | 13 +- go.mod | 3 + go.sum | 10 ++ internal/broadcast/baseclient.go | 48 ++++++ internal/broadcast/broadcastmanager.go | 150 ++++++++++++++++++ internal/broadcast/kafkaclient.go | 122 ++++++++++++++ .../websocketclient.go} | 49 +++--- .../certificatetransparency/ct-watcher.go | 5 +- internal/certstream/certstream.go | 15 +- internal/config/config.go | 18 ++- internal/metrics/logmetrics.go | 2 + internal/metrics/prometheus.go | 2 + internal/web/broadcastmanager.go | 133 ---------------- internal/web/examplecert.go | 2 +- internal/web/server.go | 21 +-- 15 files changed, 410 insertions(+), 183 deletions(-) create mode 100644 internal/broadcast/baseclient.go create mode 100644 internal/broadcast/broadcastmanager.go create mode 100644 internal/broadcast/kafkaclient.go rename internal/{web/client.go => broadcast/websocketclient.go} (84%) delete mode 100644 internal/web/broadcastmanager.go diff --git a/config.sample.yaml b/config.sample.yaml index 74fe9a6..d22dbeb 100644 --- a/config.sample.yaml +++ b/config.sample.yaml @@ -37,6 +37,14 @@ prometheus: whitelist: - "127.0.0.1/8" +# Configuration related to external stream processing tools go here. +streamprocessing: + kafka: + enabled: true + server_address: "127.0.0.1" + server_port: 9092 + topic: "certstream" + general: # DisableDefaultLogs indicates whether the default logs used in Google Chrome and provided by Google should be disabled. disable_default_logs: false @@ -66,9 +74,8 @@ general: websocket: 300 # Buffer for each CT log connection ctlog: 1000 - # Combined buffer for the broadcast manager - broadcastmanager: 10000 - + # Combined buffer for the cert dispatcher + dispatcher: 10000 # Google regularly updates the log list. If this option is set to true, the server will remove all logs no longer listed in the Google log list. # This option defaults to true. See https://github.com/d-Rickyy-b/certstream-server-go/issues/51 drop_old_logs: true diff --git a/go.mod b/go.mod index 162d32c..8aadea7 100644 --- a/go.mod +++ b/go.mod @@ -8,6 +8,7 @@ require ( github.com/google/certificate-transparency-go v1.3.3 github.com/google/trillian v1.7.3 github.com/gorilla/websocket v1.5.4-0.20250319132907-e064f32e3674 + github.com/segmentio/kafka-go v0.4.51 github.com/spf13/cobra v1.10.2 github.com/spf13/viper v1.21.0 golang.org/x/crypto v0.49.0 @@ -18,7 +19,9 @@ require ( github.com/go-logr/logr v1.4.3 // indirect github.com/go-viper/mapstructure/v2 v2.5.0 // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect + github.com/klauspost/compress v1.15.9 // indirect github.com/pelletier/go-toml/v2 v2.3.0 // indirect + github.com/pierrec/lz4/v4 v4.1.15 // indirect github.com/sagikazarmark/locafero v0.12.0 // indirect github.com/spf13/afero v1.15.0 // indirect github.com/spf13/cast v1.10.0 // indirect diff --git a/go.sum b/go.sum index 6fb7886..105cf71 100644 --- a/go.sum +++ b/go.sum @@ -25,6 +25,8 @@ github.com/gorilla/websocket v1.5.4-0.20250319132907-e064f32e3674 h1:JeSE6pjso5T github.com/gorilla/websocket v1.5.4-0.20250319132907-e064f32e3674/go.mod h1:r4w70xmWCQKmi1ONH4KIaBptdivuRPyosB9RmPlGEwA= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= +github.com/klauspost/compress v1.15.9 h1:wKRjX6JRtDdrE9qwa4b/Cip7ACOshUI4smpCQanqjSY= +github.com/klauspost/compress v1.15.9/go.mod h1:PhcZ0MbTNciWF3rruxRgKxI5NkcHHrHUDtV4Yw2GlzU= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= @@ -33,6 +35,8 @@ github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0 github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= github.com/pelletier/go-toml/v2 v2.3.0 h1:k59bC/lIZREW0/iVaQR8nDHxVq8OVlIzYCOJf421CaM= github.com/pelletier/go-toml/v2 v2.3.0/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= +github.com/pierrec/lz4/v4 v4.1.15 h1:MO0/ucJhngq7299dKLwIMtgTfbkoSPF6AoMYDd8Q4q0= +github.com/pierrec/lz4/v4 v4.1.15/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/rogpeppe/go-internal v1.9.0 h1:73kH8U+JUqXU8lRuOHeVHaa/SZPifC7BkcraZVejAe8= @@ -61,6 +65,12 @@ github.com/valyala/fastrand v1.1.0 h1:f+5HkLW4rsgzdNoleUOB69hyT9IlD2ZQh9GyDMfb5G github.com/valyala/fastrand v1.1.0/go.mod h1:HWqCzkrkg6QXT8V2EXWvXCoow7vLwOFN002oeRzjapQ= github.com/valyala/histogram v1.2.0 h1:wyYGAZZt3CpwUiIb9AU/Zbllg1llXyrtApRS815OLoQ= github.com/valyala/histogram v1.2.0/go.mod h1:Hb4kBwb4UxsaNbbbh+RRz8ZR6pdodR57tzWUS3BUzXY= +github.com/xdg-go/pbkdf2 v1.0.0 h1:Su7DPu48wXMwC3bs7MCNG+z4FhcyEuz5dlvchbq0B0c= +github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI= +github.com/xdg-go/scram v1.1.2 h1:FHX5I5B4i4hKRVRBCFRxq1iQRej7WO3hhBuJf+UUySY= +github.com/xdg-go/scram v1.1.2/go.mod h1:RT/sEzTbU5y00aCK8UOx6R7YryM0iF1N2MOmC3kKLN4= +github.com/xdg-go/stringprep v1.0.4 h1:XLI/Ng3O1Atzq0oBs3TWm+5ZVgkq2aqdlvP9JtoZ6c8= +github.com/xdg-go/stringprep v1.0.4/go.mod h1:mPGuuIYwz7CmR2bT9j4GbQqutWS1zV24gijq1dTyGkM= go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= golang.org/x/crypto v0.49.0 h1:+Ng2ULVvLHnJ/ZFEq4KdcDd/cfjrrjjNSXNzxg0Y4U4= diff --git a/internal/broadcast/baseclient.go b/internal/broadcast/baseclient.go new file mode 100644 index 0000000..9b69e3e --- /dev/null +++ b/internal/broadcast/baseclient.go @@ -0,0 +1,48 @@ +package broadcast + +import "log" + +// BaseClient defines the basic structure for a client that can receive broadcast messages. +// Other client types can embed this struct to inherit its functionality. +type BaseClient struct { + broadcastChan chan []byte + stopChan chan struct{} + name string + subType SubscriptionType + skippedCerts uint64 +} + +// Close cleans up the client's resources by closing the stop and broadcast channels. +func (c *BaseClient) Close() { + close(c.stopChan) + close(c.broadcastChan) +} + +// Name returns the name of the client. +func (c *BaseClient) Name() string { + return c.name +} + +// SubType returns the subscription type of the client. +func (c *BaseClient) SubType() SubscriptionType { + return c.subType +} + +// SkippedCerts returns the number of certificates that were skipped due to the client's broadcast channel being full. +func (c *BaseClient) SkippedCerts() uint64 { + return c.skippedCerts +} + +// Write sends a message to the client's broadcast channel. +// If the channel is full, it increments the skippedCerts counter and logs a message. +func (c *BaseClient) Write(data []byte) { + select { + case c.broadcastChan <- data: + default: + // Default case is executed if the client's broadcast channel is full. + c.skippedCerts++ + if c.skippedCerts%1000 == 1 { + log.Printf("Not providing client '%s' with cert because client's buffer is full. The client can't keep up. Skipped certs: %d\n", c.name, c.skippedCerts) + } + } +} diff --git a/internal/broadcast/broadcastmanager.go b/internal/broadcast/broadcastmanager.go new file mode 100644 index 0000000..2803231 --- /dev/null +++ b/internal/broadcast/broadcastmanager.go @@ -0,0 +1,150 @@ +package broadcast + +import ( + "log" + "sync" + + "github.com/d-Rickyy-b/certstream-server-go/internal/config" + "github.com/d-Rickyy-b/certstream-server-go/internal/metrics" + "github.com/d-Rickyy-b/certstream-server-go/internal/models" +) + +// ClientHandler dispatches certificate entries to registered clients. +var ClientHandler *Dispatcher + +type Dispatcher struct { + MessageQueue chan models.Entry + clients []CertProcessor + clientLock sync.RWMutex +} + +// NewDispatcher creates a new Dispatcher instance and assigns it to the ClientHandler variable. +func NewDispatcher() *Dispatcher { + d := &Dispatcher{} + d.MessageQueue = make(chan models.Entry, config.AppConfig.General.BufferSizes.Dispatcher) + ClientHandler = d + + metrics.Prometheus.RegisterGaugeMetricInt("certstreamservergo_clients_total{type=\"full\"}", d.ClientFullCount) + metrics.Prometheus.RegisterGaugeMetricInt("certstreamservergo_clients_total{type=\"lite\"}", d.ClientLiteCount) + metrics.Prometheus.RegisterGaugeMetricInt("certstreamservergo_clients_total{type=\"domain\"}", d.ClientDomainsCount) + return d +} + +// RegisterClient adds a client to the list of clients of the Dispatcher. +// The client will receive certificate broadcasts right after registration. +func (bm *Dispatcher) RegisterClient(c CertProcessor) { + // TODO: check if the client is already registered + bm.clientLock.Lock() + bm.clients = append(bm.clients, c) + log.Printf("Added new client. Clients: %d, Capacity: %d\n", len(bm.clients), cap(bm.clients)) + metrics.Prometheus.RegisterClient(c.Name(), func() float64 { return float64(c.SkippedCerts()) }) + bm.clientLock.Unlock() +} + +// UnregisterClient removes a client from the list of clients of the Dispatcher. +// The client will no longer receive certificate broadcasts right after unregistering. +func (bm *Dispatcher) UnregisterClient(clientName string) { + bm.clientLock.Lock() + + for i, client := range bm.clients { + if clientName == client.Name() { + // Copy the last element of the slice to the position of the removed element + // Then remove the last element by re-slicing + bm.clients[i] = bm.clients[len(bm.clients)-1] + bm.clients[len(bm.clients)-1] = nil + bm.clients = bm.clients[:len(bm.clients)-1] + + // Close the broadcast channel of the client, otherwise this leads to a memory leak + client.Close() + + metrics.Prometheus.UnregisterClient(c.name) + + break + } + } + + bm.clientLock.Unlock() +} + +// ClientFullCount returns the current number of clients connected to the service on the `full` endpoint. +func (bm *Dispatcher) ClientFullCount() (count int64) { + return bm.clientCountByType(SubTypeFull) +} + +// ClientLiteCount returns the current number of clients connected to the service on the `lite` endpoint. +func (bm *Dispatcher) ClientLiteCount() (count int64) { + return bm.clientCountByType(SubTypeLite) +} + +// ClientDomainsCount returns the current number of clients connected to the service on the `domains-only` endpoint. +func (bm *Dispatcher) ClientDomainsCount() (count int64) { + return bm.clientCountByType(SubTypeDomain) +} + +// clientCountByType returns the current number of clients connected to the service on the endpoint matching +// the specified SubscriptionType. +func (bm *Dispatcher) clientCountByType(subType SubscriptionType) (count int64) { + bm.clientLock.RLock() + defer bm.clientLock.RUnlock() + + for _, c := range bm.clients { + if c.SubType() == subType { + count++ + } + } + + return count +} + +// GetSkippedCerts returns a map of client names to the number of skipped certificates for each client. +func (bm *Dispatcher) GetSkippedCerts() map[string]uint64 { + bm.clientLock.RLock() + defer bm.clientLock.RUnlock() + + skippedCerts := make(map[string]uint64, len(bm.clients)) + for _, c := range bm.clients { + skippedCerts[c.Name()] = c.SkippedCerts() + } + + return skippedCerts +} + +// broadcaster is run in a goroutine and handles the dispatching of certs to clients. +func (bm *Dispatcher) broadcaster() { + for { + var data []byte + + // Take entry out of broadcast channel and generate JSON representations for the entry. + entry := <-bm.Broadcast + dataLite := entry.JSONLite() + dataFull := entry.JSON() + dataDomain := entry.JSONDomains() + + bm.clientLock.RLock() + + for _, c := range bm.clients { + switch c.SubType() { + case SubTypeLite: + data = dataLite + case SubTypeFull: + data = dataFull + case SubTypeDomain: + data = dataDomain + default: + // This should never happen, but if it does, we log it and skip the client. + log.Printf("Unknown subscription type '%d' for client '%s'. Skipping this client!\n", c.SubType(), c.Name()) + continue + } + + c.Write(data) + } + + bm.clientLock.RUnlock() + } +} + +// Start starts the broadcaster goroutine. +func (bm *Dispatcher) Start() { + go bm.broadcaster() + log.Println("Dispatcher started. Listening for certificate entries...") +} diff --git a/internal/broadcast/kafkaclient.go b/internal/broadcast/kafkaclient.go new file mode 100644 index 0000000..b89b845 --- /dev/null +++ b/internal/broadcast/kafkaclient.go @@ -0,0 +1,122 @@ +package broadcast + +import ( + "context" + "log" + "net" + "strconv" + "time" + + "github.com/d-Rickyy-b/certstream-server-go/internal/config" + + "github.com/segmentio/kafka-go" +) + +const ( + TOPIC = "certstream" +) + +// KafkaClient connects to a Kafka server in order to provide it with certificates. +type KafkaClient struct { + conn *kafka.Conn // Kafka connection + addr string + isConnected bool + BaseClient +} + +// NewKafkaClient creates a new Kafka client that immediately connects to the configured Kafka server. +func NewKafkaClient(subType SubscriptionType, name string, certBufferSize int) *KafkaClient { + addr := net.JoinHostPort(config.AppConfig.StreamProcessing.Kafka.ServerAddr, strconv.Itoa(config.AppConfig.StreamProcessing.Kafka.ServerPort)) + + // Connect to the Kafka server + // TODO make topic configurable + conn, err := kafka.DialLeader(context.Background(), "tcp", addr, TOPIC, 0) + if err != nil { + log.Println("failed to connect to kafka:", err) + } + + kc := &KafkaClient{ + conn: conn, + addr: addr, + BaseClient: BaseClient{ + broadcastChan: make(chan []byte, certBufferSize), + stopChan: make(chan struct{}), + name: name, + subType: subType, + }, + } + + go kc.broadcastHandler() + go kc.reconnectHandler() + + return kc +} + +// reconnectHandler is a background job that attempts to reconnect to the Kafka server if the connection is lost. +func (c *KafkaClient) reconnectHandler() { + for { + select { + case <-c.stopChan: + log.Println("Stopping reconnectHandler for kafka producer:", c.addr) + c.conn.Close() + + return + default: + if c.isConnected { + // If already connected or no connection exists, skip reconnection + time.Sleep(5 * time.Second) + continue + } + conn, err := kafka.DialLeader(context.Background(), "tcp", c.addr, TOPIC, 0) + if err != nil { + log.Printf("Reconnect failed: %v. Retrying in 5s...", err) + time.Sleep(5 * time.Second) + + continue + } + // Close old connection if exists + if c.conn != nil { + _ = c.conn.Close() + } + c.conn = conn + c.isConnected = true + log.Println("Reconnected to Kafka at", c.addr) + } + } +} + +// Each client has a broadcastHandler that runs in the background and sends out the broadcast messages to the client. +func (c *KafkaClient) broadcastHandler() { + writeWait := 60 * time.Second + + defer func() { + log.Println("Closing broadcast handler for kafka producer:", c.addr) + if err := c.conn.Close(); err != nil { + log.Println("failed to close writer:", err) + } + + ClientHandler.UnregisterClient(c.name) + }() + + for { + select { + case <-c.stopChan: + return + case message := <-c.broadcastChan: + if !c.isConnected { + continue + } + + _ = c.conn.SetWriteDeadline(time.Now().Add(writeWait)) + + c.conn.Broker() + _, err := c.conn.WriteMessages( + kafka.Message{Value: message}, + ) + if err != nil { + c.isConnected = false + log.Println("Failed to write messages to kafka:", err) + } + } + } +} diff --git a/internal/web/client.go b/internal/broadcast/websocketclient.go similarity index 84% rename from internal/web/client.go rename to internal/broadcast/websocketclient.go index 6626c2a..642600a 100644 --- a/internal/web/client.go +++ b/internal/broadcast/websocketclient.go @@ -1,4 +1,4 @@ -package web +package broadcast import ( "fmt" @@ -21,33 +21,26 @@ const ( type SubscriptionType int -// client represents a single client's connection to the server. -type client struct { - clientData - - id string - conn *websocket.Conn - broadcastChan chan []byte - subType SubscriptionType - skippedCerts uint64 -} - -type clientData struct { - userAgent string - connectionIP string - connectionPort string - realIPFromHeader string +// WebsocketClient represents a single WebSocket client's connection to the server. +type WebsocketClient struct { + conn *websocket.Conn + *BaseClient } -// newClient creates a new client struct that holds information about a connected client. -func newClient(conn *websocket.Conn, subType SubscriptionType, data clientData, certBufferSize int) *client { - return &client{ - clientData: data, - id: generateClientID(), - conn: conn, - broadcastChan: make(chan []byte, certBufferSize), - subType: subType, +// NewWebsocketClient creates a new WebSocket client from the given connection. +func NewWebsocketClient(conn *websocket.Conn, subType SubscriptionType, name string, certBufferSize int) *WebsocketClient { + c := &WebsocketClient{ + conn: conn, + BaseClient: &BaseClient{ + broadcastChan: make(chan []byte, certBufferSize), + name: name, + subType: subType, + }, } + go c.broadcastHandler() + go c.listenWebsocket() + + return c } // generateClientID generates a random 8-char identifier for the client. @@ -62,7 +55,7 @@ func generateClientID() string { } // Each client has a broadcastHandler that runs in the background and sends out the broadcast messages to the client. -func (c *client) broadcastHandler() { +func (c *WebsocketClient) broadcastHandler() { writeWait := 60 * time.Second pingTicker := time.NewTicker(30 * time.Second) @@ -109,10 +102,10 @@ func (c *client) broadcastHandler() { // listenWebsocket is running in the background on a goroutine and listens for messages from the client. // It responds to ping messages with a pong message. It closes the connection if the client sends // a close message or no ping is received within 65 seconds. -func (c *client) listenWebsocket() { +func (c *WebsocketClient) listenWebsocket() { defer func() { _ = c.conn.Close() - ClientHandler.unregisterClient(c) + ClientHandler.UnregisterClient(c.name) }() readWait := 65 * time.Second diff --git a/internal/certificatetransparency/ct-watcher.go b/internal/certificatetransparency/ct-watcher.go index b4eb458..6db1174 100644 --- a/internal/certificatetransparency/ct-watcher.go +++ b/internal/certificatetransparency/ct-watcher.go @@ -13,6 +13,9 @@ import ( "sync/atomic" "time" + "github.com/google/trillian/client/backoff" + + "github.com/d-Rickyy-b/certstream-server-go/internal/broadcast" "github.com/d-Rickyy-b/certstream-server-go/internal/config" "github.com/d-Rickyy-b/certstream-server-go/internal/metrics" "github.com/d-Rickyy-b/certstream-server-go/internal/models" @@ -553,7 +556,7 @@ func certHandler(entryChan chan models.Entry) { } // Run JSON encoding in the background and send the result to the clients. - web.ClientHandler.Broadcast <- entry + broadcast.ClientHandler.MessageQueue <- entry // Update metrics url := entry.Data.Source.NormalizedURL diff --git a/internal/certstream/certstream.go b/internal/certstream/certstream.go index 3bf4065..af00b44 100644 --- a/internal/certstream/certstream.go +++ b/internal/certstream/certstream.go @@ -11,6 +11,7 @@ import ( "os/signal" "syscall" + "github.com/d-Rickyy-b/certstream-server-go/internal/broadcast" "github.com/d-Rickyy-b/certstream-server-go/internal/certificatetransparency" "github.com/d-Rickyy-b/certstream-server-go/internal/config" "github.com/d-Rickyy-b/certstream-server-go/internal/metrics" @@ -35,6 +36,11 @@ func NewRawCertstream(config config.Config) *Certstream { func NewCertstreamServer(config config.Config) (*Certstream, error) { cs := NewRawCertstream(config) + // Start the broadcast dispatcher + broadcast.NewDispatcher() + broadcast.ClientHandler.Start() + + // TODO: add support do disable websocket Server // Initialize the webserver used for the websocket server webserver := web.NewWebsocketServer( config.Webserver.ListenAddr, @@ -48,7 +54,14 @@ func NewCertstreamServer(config config.Config) (*Certstream, error) { // Setup metrics server cs.setupMetrics(webserver) - return cs, nil + if config.StreamProcessing.Kafka.Enabled { + log.Println("Initializing Kafka client...") + + kc := broadcast.NewKafkaClient(broadcast.SubTypeFull, "kafka-producer", config.General.BufferSizes.Websocket) + broadcast.ClientHandler.RegisterClient(kc) + } + + return &cs, nil } // NewCertstreamFromConfigFile creates a new Certstream server from a config file. diff --git a/internal/config/config.go b/internal/config/config.go index 02eefea..7f3a167 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -38,6 +38,7 @@ type BufferSizes struct { Websocket int `mapstructure:"websocket"` CTLog int `mapstructure:"ctlog"` BroadcastManager int `mapstructure:"broadcastmanager"` + Dispatcher int `mapstructure:"dispatcher"` } type Config struct { @@ -56,6 +57,14 @@ type Config struct { MetricsURL string `mapstructure:"metrics_url"` ExposeSystemMetrics bool `mapstructure:"expose_system_metrics"` } + StreamProcessing struct { + Kafka struct { + Enabled bool `yaml:"enabled"` + ServerAddr string `yaml:"server_addr"` + ServerPort int `yaml:"server_port"` + Topic string `yaml:"topic"` + } + } General struct { // DisableDefaultLogs indicates whether the default logs used in Google Chrome and provided by Google should be disabled. DisableDefaultLogs bool `mapstructure:"disable_default_logs"` @@ -327,8 +336,13 @@ func validateConfig(config *Config) bool { config.General.BufferSizes.CTLog = 1000 } - if config.General.BufferSizes.BroadcastManager <= 0 { - config.General.BufferSizes.BroadcastManager = 10000 + // For backward compatibility, copy value from deprecated BroadcastManager field + if config.General.BufferSizes.BroadcastManager != 0 { + config.General.BufferSizes.Dispatcher = config.General.BufferSizes.BroadcastManager + } + + if config.General.BufferSizes.Dispatcher <= 0 { + config.General.BufferSizes.Dispatcher = 10000 } // If the cleanup flag is not set, default to true diff --git a/internal/metrics/logmetrics.go b/internal/metrics/logmetrics.go index 4f45102..fb10415 100644 --- a/internal/metrics/logmetrics.go +++ b/internal/metrics/logmetrics.go @@ -327,10 +327,12 @@ func GetProcessedPrecerts() int64 { return ProcessedPrecerts } +// GetCertMetrics returns a copy of the internal metrics map. func GetCertMetrics() CTMetrics { return Metrics.GetCTMetrics() } +// GetLogOperators returns a map of operator names to a list of CT logs. func GetLogOperators() map[string][]string { return Metrics.OperatorLogMapping() } diff --git a/internal/metrics/prometheus.go b/internal/metrics/prometheus.go index caca426..2f72bb7 100644 --- a/internal/metrics/prometheus.go +++ b/internal/metrics/prometheus.go @@ -10,6 +10,8 @@ import ( "sync" "time" + "github.com/d-Rickyy-b/certstream-server-go/internal/certificatetransparency" + "github.com/VictoriaMetrics/metrics" ) diff --git a/internal/web/broadcastmanager.go b/internal/web/broadcastmanager.go deleted file mode 100644 index e2e4350..0000000 --- a/internal/web/broadcastmanager.go +++ /dev/null @@ -1,133 +0,0 @@ -package web - -import ( - "log" - "sync" - - "github.com/d-Rickyy-b/certstream-server-go/internal/metrics" - "github.com/d-Rickyy-b/certstream-server-go/internal/models" -) - -type BroadcastManager struct { - Broadcast chan models.Entry - clients []*client - clientLock sync.RWMutex -} - -func NewBroadcastManager() *BroadcastManager { - bm := &BroadcastManager{} - metrics.Prometheus.RegisterGaugeMetricInt("certstreamservergo_clients_total{type=\"full\"}", bm.ClientFullCount) - metrics.Prometheus.RegisterGaugeMetricInt("certstreamservergo_clients_total{type=\"lite\"}", bm.ClientLiteCount) - metrics.Prometheus.RegisterGaugeMetricInt("certstreamservergo_clients_total{type=\"domain\"}", bm.ClientDomainsCount) - - return bm -} - -// registerClient adds a client to the list of clients of the BroadcastManager. -// The client will receive certificate broadcasts right after registration. -func (bm *BroadcastManager) registerClient(c *client) { - bm.clientLock.Lock() - bm.clients = append(bm.clients, c) - log.Printf("Clients: %d, Capacity: %d\n", len(bm.clients), cap(bm.clients)) - metrics.Prometheus.RegisterClient(c.id, c.connectionIP, c.connectionPort, c.realIPFromHeader, c.userAgent, func() float64 { return float64(c.skippedCerts) }) - bm.clientLock.Unlock() -} - -// unregisterClient removes a client from the list of clients of the BroadcastManager. -// The client will no longer receive certificate broadcasts right after unregistering. -func (bm *BroadcastManager) unregisterClient(c *client) { - bm.clientLock.Lock() - log.Println("Unregistering client:", c.conn.RemoteAddr()) - - // Close the broadcast channel of the client, otherwise this leads to a memory leak - close(c.broadcastChan) - metrics.Prometheus.UnregisterClient(c.id, c.connectionIP, c.connectionPort, c.realIPFromHeader, c.userAgent) - - // Remove client from internal client list - for i, storedClient := range bm.clients { - if c != storedClient { - continue - } - - // Copy the last element of the slice to the position of the removed element - // Then remove the last element by re-slicing - bm.clients[i] = bm.clients[len(bm.clients)-1] - bm.clients[len(bm.clients)-1] = nil - bm.clients = bm.clients[:len(bm.clients)-1] - - break - } - - bm.clientLock.Unlock() -} - -// ClientFullCount returns the current number of clients connected to the service on the `full` endpoint. -func (bm *BroadcastManager) ClientFullCount() (count int64) { - return bm.clientCountByType(SubTypeFull) -} - -// ClientLiteCount returns the current number of clients connected to the service on the `lite` endpoint. -func (bm *BroadcastManager) ClientLiteCount() (count int64) { - return bm.clientCountByType(SubTypeLite) -} - -// ClientDomainsCount returns the current number of clients connected to the service on the `domains-only` endpoint. -func (bm *BroadcastManager) ClientDomainsCount() (count int64) { - return bm.clientCountByType(SubTypeDomain) -} - -// clientCountByType returns the current number of clients connected to the service on the endpoint matching -// the specified SubscriptionType. -func (bm *BroadcastManager) clientCountByType(subType SubscriptionType) (count int64) { - bm.clientLock.RLock() - defer bm.clientLock.RUnlock() - - for _, c := range bm.clients { - if c.subType == subType { - count++ - } - } - - return count -} - -// broadcaster is run in a goroutine and handles the dispatching of entries to clients. -func (bm *BroadcastManager) broadcaster() { - for { - var data []byte - - // Take entry out of broadcast channel and generate JSON representations for the entry. - entry := <-bm.Broadcast - dataLite := entry.JSONLite() - dataFull := entry.JSON() - dataDomain := entry.JSONDomains() - - bm.clientLock.RLock() - - for _, c := range bm.clients { - switch c.subType { - case SubTypeLite: - data = dataLite - case SubTypeFull: - data = dataFull - case SubTypeDomain: - data = dataDomain - default: - log.Printf("Unknown subscription type '%d'. Skipping client %s\n", c.subType, c.Name()) - continue - } - - select { - case c.broadcastChan <- data: - default: - // Default case is executed if the client's broadcast channel is full. - c.skippedCerts++ - if c.skippedCerts%1000 == 1 { - log.Printf("Client can't keep up. Skipped certs: %d | Client: %s\n", c.skippedCerts, c.Name()) - } - } - } - - bm.clientLock.RUnlock() - } -} diff --git a/internal/web/examplecert.go b/internal/web/examplecert.go index e841323..05576f3 100644 --- a/internal/web/examplecert.go +++ b/internal/web/examplecert.go @@ -29,7 +29,7 @@ func exampleDomains(w http.ResponseWriter, _ *http.Request) { w.Write(exampleCert.JSONDomains()) //nolint:errcheck } -// SetExampleCert sets one certificate as the example Cert that is returned by the example endpoints. +// SetExampleCert sets the example certificate to be used in the example endpoints. func SetExampleCert(cert models.Entry) { exampleCert = cert } diff --git a/internal/web/server.go b/internal/web/server.go index f3e1d80..4d391fd 100644 --- a/internal/web/server.go +++ b/internal/web/server.go @@ -228,7 +228,7 @@ func initFullWebsocket(w http.ResponseWriter, r *http.Request) { return } - setupClient(connection, SubTypeFull, r) + setupClient(connection, broadcast.SubTypeFull, r.RemoteAddr) } // initLiteWebsocket is called when a client connects to the / endpoint. @@ -240,7 +240,7 @@ func initLiteWebsocket(w http.ResponseWriter, r *http.Request) { return } - setupClient(connection, SubTypeLite, r) + setupClient(connection, broadcast.SubTypeLite, r.RemoteAddr) } // initDomainWebsocket is called when a client connects to the /domains-only endpoint. @@ -252,7 +252,7 @@ func initDomainWebsocket(w http.ResponseWriter, r *http.Request) { return } - setupClient(connection, SubTypeDomain, r) + setupClient(connection, broadcast.SubTypeDomain, r.RemoteAddr) } // upgradeConnection upgrades the connection to a websocket and returns the connection. @@ -283,7 +283,7 @@ func upgradeConnection(w http.ResponseWriter, r *http.Request) (*websocket.Conn, } // setupClient initializes a client struct and starts the broadcastHandler and websocket listener. -func setupClient(connection *websocket.Conn, subType SubscriptionType, r *http.Request) { +func setupClient(connection *websocket.Conn, subscriptionType broadcast.SubscriptionType, name string) { // Extract data from request origConnAddr, _ := r.Context().Value(origConnAddrKey).(string) @@ -305,11 +305,8 @@ func setupClient(connection *websocket.Conn, subType SubscriptionType, r *http.R realIPFromHeader: realIPFromHeader, } - c := newClient(connection, subType, data, config.AppConfig.General.BufferSizes.Websocket) - go c.broadcastHandler() - go c.listenWebsocket() - - ClientHandler.registerClient(c) + c := broadcast.NewWebsocketClient(connection, subscriptionType, name, config.AppConfig.General.BufferSizes.Websocket) + broadcast.ClientHandler.RegisterClient(c) } // setupWebsocketRoutes configures all the routes necessary for the websocket webserver. @@ -385,8 +382,7 @@ func NewMetricsServer(networkIf string, port int, certPath, keyPath string) *Ser } // NewWebsocketServer starts a new webserver and initialized it with the necessary routes. -// It also starts the broadcaster in ClientHandler as a background job and takes care of -// setting up websocket.Upgrader. +// It also takes care of setting up websocket.Upgrader. func NewWebsocketServer(networkIf string, port int, certPath, keyPath string) *Server { websocketServer := &Server{ networkIf: networkIf, @@ -417,9 +413,6 @@ func NewWebsocketServer(networkIf string, port int, certPath, keyPath string) *S setupWebsocketRoutes(websocketServer.routes) websocketServer.initServer() - ClientHandler.Broadcast = make(chan models.Entry, config.AppConfig.General.BufferSizes.BroadcastManager) - go ClientHandler.broadcaster() - return websocketServer } From 8e7e43091bafb1a86f62d610b035ae7b8c7c50f2 Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 4 Aug 2025 01:57:53 +0200 Subject: [PATCH 03/30] chore: surround yaml strings by quotes --- docker/docker-compose.metrics.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/docker/docker-compose.metrics.yml b/docker/docker-compose.metrics.yml index 301a165..9c17a7f 100644 --- a/docker/docker-compose.metrics.yml +++ b/docker/docker-compose.metrics.yml @@ -30,7 +30,7 @@ services: ports: # Exposing Prometheus is NOT required, if you don't want to access it from outside the Docker network. # Using localhost enables you to use a reverse proxy (e.g. with basic auth) to access Prometheus in a more secure way. - - 127.0.0.1:9090:9090 + - "127.0.0.1:9090:9090" networks: - monitoring extra_hosts: @@ -44,7 +44,7 @@ services: depends_on: - prometheus ports: - - 127.0.0.1:8082:3000 + - "127.0.0.1:8082:3000" volumes: - ./grafana_data:/var/lib/grafana - ./grafana/provisioning/:/etc/grafana/provisioning/ @@ -60,7 +60,7 @@ services: # Configure the service to run as specific user. # user: "1000:1000" ports: - - 127.0.0.1:8080:80 + - "127.0.0.1:8080:80" # Don't forget to open the other port in case you run the Prometheus endpoint on another port than the websocket server. # - 127.0.0.1:8081:81 volumes: From 0fdcb45423b5e99f83c299fe330987031b84e0e3 Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 4 Aug 2025 01:58:13 +0200 Subject: [PATCH 04/30] docs: add missing options to sample config --- config.sample.yaml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/config.sample.yaml b/config.sample.yaml index d22dbeb..cc8fb60 100644 --- a/config.sample.yaml +++ b/config.sample.yaml @@ -23,6 +23,11 @@ webserver: cert_key_path: "" # specify if the server should attempt to negotiate per message compression (RFC 7692) compression_enabled: false + # Use True-Client-IP, X-Real-IP or the X-Forwarded-For headers (in that order) to determine the real IP address of the client. + # If you are using a reverse proxy, you should set this to true. + real_ip: false + whitelist: + - "127.0.0.1/8" prometheus: enabled: true From 6703a32633352f44dde58727a7ee06a63f4603c5 Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 4 Aug 2025 01:59:30 +0200 Subject: [PATCH 05/30] fix: add missing imports and sort them --- internal/web/server.go | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/internal/web/server.go b/internal/web/server.go index 4d391fd..0e05542 100644 --- a/internal/web/server.go +++ b/internal/web/server.go @@ -12,12 +12,11 @@ import ( "strings" "time" - "github.com/go-chi/chi/v5" - "github.com/go-chi/chi/v5/middleware" - + "github.com/d-Rickyy-b/certstream-server-go/internal/broadcast" "github.com/d-Rickyy-b/certstream-server-go/internal/config" - "github.com/d-Rickyy-b/certstream-server-go/internal/models" + "github.com/go-chi/chi/v5" + "github.com/go-chi/chi/v5/middleware" "github.com/gorilla/websocket" ) From 28bfc2ae1a4b6564187a518e01c261d3e879fa98 Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 4 Aug 2025 02:04:35 +0200 Subject: [PATCH 06/30] docs: add comments for config --- internal/config/config.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/internal/config/config.go b/internal/config/config.go index 7f3a167..025b7fe 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -72,7 +72,9 @@ type Config struct { AdditionalLogs []LogConfig `mapstructure:"additional_logs"` AdditionalTiledLogs []LogConfig `mapstructure:"additional_tiled_logs"` ExcludedLogs []LogConfig `mapstructure:"excluded_logs"` + // BufferSizes contains the buffer sizes for the different components of the server. They usually don't need any adjustments. BufferSizes BufferSizes `mapstructure:"buffer_sizes"` + // DropOldLogs indicates whether downloading CT-Logs should start at the latest index (true) or should from the beginning (false). DropOldLogs *bool `mapstructure:"drop_old_logs"` Recovery struct { Enabled bool `mapstructure:"enabled"` From 92d983a9b7e8a0a3bd62a3e42585716a1ae07847 Mon Sep 17 00:00:00 2001 From: Rico Date: Tue, 5 Aug 2025 17:37:09 +0200 Subject: [PATCH 07/30] feat: add certprocessor interface --- internal/broadcast/certprocessor.go | 26 ++++++++++++++++++++++++++ 1 file changed, 26 insertions(+) create mode 100644 internal/broadcast/certprocessor.go diff --git a/internal/broadcast/certprocessor.go b/internal/broadcast/certprocessor.go new file mode 100644 index 0000000..d10d47a --- /dev/null +++ b/internal/broadcast/certprocessor.go @@ -0,0 +1,26 @@ +package broadcast + +// CertProcessor defines the interface for processing certificate messages. +// Different implementations can be used for different types of clients, such as WebSocket clients, kafka, or other message queues for example. +type CertProcessor interface { + // Write sends a message to the client's broadcast channel. + Write(message []byte) + // Close closes the client's connection and cleans up resources. + Close() + + // To obtain client details, there are getter methods for several values + Name() string + SubType() SubscriptionType + SkippedCerts() uint64 +} + +const ( + // SubTypeFull represents full certificate updates. + SubTypeFull SubscriptionType = iota + // SubTypeLite represents certificate updates with less details. + SubTypeLite + // SubTypeDomain represents updates that only include domain information. + SubTypeDomain +) + +type SubscriptionType int From e422401892792c539f42c77f57693aebb0382f54 Mon Sep 17 00:00:00 2001 From: Rico Date: Tue, 5 Aug 2025 18:08:46 +0200 Subject: [PATCH 08/30] feat: implement NSQ client refers to #23 and #35 I needed to change the config to an arrary in order to have proper support for multiple kinds of queues and stream processors. Currently NSQ and kafka are supported. Future releases could support amazon SNS/SQS by creating a new client in the broadcast package and by extending the config. --- config.sample.yaml | 12 ++- go.mod | 2 + go.sum | 4 + internal/broadcast/kafkaclient.go | 21 +++--- internal/broadcast/nsqclient.go | 119 ++++++++++++++++++++++++++++++ internal/certstream/certstream.go | 39 +++++++++- internal/config/config.go | 16 ++-- 7 files changed, 186 insertions(+), 27 deletions(-) create mode 100644 internal/broadcast/nsqclient.go diff --git a/config.sample.yaml b/config.sample.yaml index cc8fb60..2a52eee 100644 --- a/config.sample.yaml +++ b/config.sample.yaml @@ -43,10 +43,16 @@ prometheus: - "127.0.0.1/8" # Configuration related to external stream processing tools go here. -streamprocessing: - kafka: +stream_processing: + - name: "kafka" + enabled: false + server_addr: "127.0.0.1" + server_port: 9092 + topic: "certstream" + + - name: "nqs" enabled: true - server_address: "127.0.0.1" + server_addr: "127.0.0.1" server_port: 9092 topic: "certstream" diff --git a/go.mod b/go.mod index 8aadea7..1757d6f 100644 --- a/go.mod +++ b/go.mod @@ -8,6 +8,7 @@ require ( github.com/google/certificate-transparency-go v1.3.3 github.com/google/trillian v1.7.3 github.com/gorilla/websocket v1.5.4-0.20250319132907-e064f32e3674 + github.com/nsqio/go-nsq v1.1.0 github.com/segmentio/kafka-go v0.4.51 github.com/spf13/cobra v1.10.2 github.com/spf13/viper v1.21.0 @@ -18,6 +19,7 @@ require ( github.com/fsnotify/fsnotify v1.9.0 // indirect github.com/go-logr/logr v1.4.3 // indirect github.com/go-viper/mapstructure/v2 v2.5.0 // indirect + github.com/golang/snappy v0.0.1 // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect github.com/klauspost/compress v1.15.9 // indirect github.com/pelletier/go-toml/v2 v2.3.0 // indirect diff --git a/go.sum b/go.sum index 105cf71..609977f 100644 --- a/go.sum +++ b/go.sum @@ -15,6 +15,8 @@ github.com/go-viper/mapstructure/v2 v2.5.0 h1:vM5IJoUAy3d7zRSVtIwQgBj7BiWtMPfmPE github.com/go-viper/mapstructure/v2 v2.5.0/go.mod h1:oJDH3BJKyqBA2TXFhDsKDGDTlndYOZ6rGS0BRZIxGhM= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= +github.com/golang/snappy v0.0.1 h1:Qgr9rKW7uDUkrbSmQeiDsGa8SjGyCOGtuasMWwvp2P4= +github.com/golang/snappy v0.0.1/go.mod h1:/XxbfmMg8lxefKM7IXC3fBNl/7bRcc72aCRzEWrmP2Q= github.com/google/certificate-transparency-go v1.3.3 h1:hq/rSxztSkXN2tx/3jQqF6Xc0O565UQPdHrOWvZwybo= github.com/google/certificate-transparency-go v1.3.3/go.mod h1:iR17ZgSaXRzSa5qvjFl8TnVD5h8ky2JMVio+dzoKMgA= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= @@ -33,6 +35,8 @@ github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= +github.com/nsqio/go-nsq v1.1.0 h1:PQg+xxiUjA7V+TLdXw7nVrJ5Jbl3sN86EhGCQj4+FYE= +github.com/nsqio/go-nsq v1.1.0/go.mod h1:vKq36oyeVXgsS5Q8YEO7WghqidAVXQlcFxzQbQTuDEY= github.com/pelletier/go-toml/v2 v2.3.0 h1:k59bC/lIZREW0/iVaQR8nDHxVq8OVlIzYCOJf421CaM= github.com/pelletier/go-toml/v2 v2.3.0/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= github.com/pierrec/lz4/v4 v4.1.15 h1:MO0/ucJhngq7299dKLwIMtgTfbkoSPF6AoMYDd8Q4q0= diff --git a/internal/broadcast/kafkaclient.go b/internal/broadcast/kafkaclient.go index b89b845..a10a938 100644 --- a/internal/broadcast/kafkaclient.go +++ b/internal/broadcast/kafkaclient.go @@ -3,12 +3,8 @@ package broadcast import ( "context" "log" - "net" - "strconv" "time" - "github.com/d-Rickyy-b/certstream-server-go/internal/config" - "github.com/segmentio/kafka-go" ) @@ -20,24 +16,23 @@ const ( type KafkaClient struct { conn *kafka.Conn // Kafka connection addr string + topic string isConnected bool BaseClient } // NewKafkaClient creates a new Kafka client that immediately connects to the configured Kafka server. -func NewKafkaClient(subType SubscriptionType, name string, certBufferSize int) *KafkaClient { - addr := net.JoinHostPort(config.AppConfig.StreamProcessing.Kafka.ServerAddr, strconv.Itoa(config.AppConfig.StreamProcessing.Kafka.ServerPort)) - +func NewKafkaClient(subType SubscriptionType, addr, name, topic string, certBufferSize int) *KafkaClient { // Connect to the Kafka server - // TODO make topic configurable - conn, err := kafka.DialLeader(context.Background(), "tcp", addr, TOPIC, 0) + conn, err := kafka.DialLeader(context.Background(), "tcp", addr, topic, 0) if err != nil { log.Println("failed to connect to kafka:", err) } kc := &KafkaClient{ - conn: conn, - addr: addr, + conn: conn, + addr: addr, + topic: topic, BaseClient: BaseClient{ broadcastChan: make(chan []byte, certBufferSize), stopChan: make(chan struct{}), @@ -67,7 +62,9 @@ func (c *KafkaClient) reconnectHandler() { time.Sleep(5 * time.Second) continue } - conn, err := kafka.DialLeader(context.Background(), "tcp", c.addr, TOPIC, 0) + + // Attempt to connect to the Kafka server + conn, err := kafka.DialLeader(context.Background(), "tcp", c.addr, c.topic, 0) if err != nil { log.Printf("Reconnect failed: %v. Retrying in 5s...", err) time.Sleep(5 * time.Second) diff --git a/internal/broadcast/nsqclient.go b/internal/broadcast/nsqclient.go new file mode 100644 index 0000000..56091ec --- /dev/null +++ b/internal/broadcast/nsqclient.go @@ -0,0 +1,119 @@ +package broadcast + +import ( + "github.com/nsqio/go-nsq" + + "log" + "time" +) + +// NSQClient connects to a NSQ server in order to provide it with certificates. +type NSQClient struct { + conn *nsq.Producer // Kafka connection + addr string + topic string + isConnected bool + BaseClient +} + +// nullLogger is a logger that does not output anything. +// It is used to silence the NSQ logger. +type nullLogger struct{} + +func (l nullLogger) Output(calldepth int, s string) error { + // Do nothing, effectively silencing the logger + return nil +} + +// NewNSQClient creates a new NSQ client that immediately connects to the configured NSQ server +func NewNSQClient(subType SubscriptionType, addr, name, topic string, certBufferSize int) *NSQClient { + log.Println("Initializing NSQ client...") + + // Instantiate a producer. + conf := nsq.NewConfig() + conn, err := nsq.NewProducer(addr, conf) + if err != nil { + log.Println(err) + } + + log.Println("Connected to NSQ server at", addr) + + // Silence log output from NSQ + conn.SetLogger(nullLogger{}, nsq.LogLevelError) + + nsqc := &NSQClient{ + conn: conn, + addr: addr, + topic: topic, + BaseClient: BaseClient{ + broadcastChan: make(chan []byte, certBufferSize), + stopChan: make(chan struct{}), + name: name, + subType: subType, + }, + } + nsqc.isConnected = true + + go nsqc.broadcastHandler() + go nsqc.reconnectHandler() + + return nsqc +} + +// reconnectHandler is a background job that attempts to reconnect to the Kafka server if the connection is lost. +func (c *NSQClient) reconnectHandler() { + for { + select { + case <-c.stopChan: + log.Println("Stopping reconnectHandler for kafka producer:", c.addr) + return + default: + if c.isConnected { + // If already connected or no connection exists, skip reconnection + time.Sleep(5 * time.Second) + continue + } + // Attempt to connect to the NSQ server + err := c.conn.Ping() + if err != nil { + log.Printf("Reconnect to NSQ server failed: '%v'. Retrying in 5s...", err) + time.Sleep(5 * time.Second) + + continue + } + + c.isConnected = true + log.Println("Reconnected to NSQ server at", c.addr) + } + } +} + +// Each client has a broadcastHandler that runs in the background and sends out the broadcast messages to the client. +func (c *NSQClient) broadcastHandler() { + // writeWait := 60 * time.Second + + defer func() { + log.Println("Closing broadcast handler for nsq producer:", c.addr) + // Gracefully stop the producer when appropriate (e.g. before shutting down the service) + c.conn.Stop() + }() + + for { + select { + case <-c.stopChan: + return + case message := <-c.broadcastChan: + if !c.isConnected { + continue + } + + // Synchronously publish a single message to the specified topic. + // Messages can also be sent asynchronously and/or in batches. + err := c.conn.Publish(c.topic, message) + if err != nil { + log.Println("Error writing to NSQ topic:", err) + c.isConnected = false + } + } + } +} diff --git a/internal/certstream/certstream.go b/internal/certstream/certstream.go index af00b44..49c7edf 100644 --- a/internal/certstream/certstream.go +++ b/internal/certstream/certstream.go @@ -7,8 +7,10 @@ package certstream import ( "fmt" "log" + "net" "os" "os/signal" + "strconv" "syscall" "github.com/d-Rickyy-b/certstream-server-go/internal/broadcast" @@ -54,11 +56,40 @@ func NewCertstreamServer(config config.Config) (*Certstream, error) { // Setup metrics server cs.setupMetrics(webserver) - if config.StreamProcessing.Kafka.Enabled { - log.Println("Initializing Kafka client...") + log.Println(config.StreamProcessing) + // Initialize the stream processors if configured and enabled. + for _, streamProcessor := range config.StreamProcessing { + if !streamProcessor.Enabled { + continue + } - kc := broadcast.NewKafkaClient(broadcast.SubTypeFull, "kafka-producer", config.General.BufferSizes.Websocket) - broadcast.ClientHandler.RegisterClient(kc) + addr := net.JoinHostPort(streamProcessor.ServerAddr, strconv.Itoa(streamProcessor.ServerPort)) + log.Printf("Initializing stream processor: %s at %s\n", streamProcessor.Name, addr) + + switch streamProcessor.Type { + case "nsq": + log.Println("Initializing NSQ client...") + nc := broadcast.NewNSQClient( + broadcast.SubTypeFull, + addr, + streamProcessor.Name, + streamProcessor.Topic, + config.General.BufferSizes.Websocket, + ) + broadcast.ClientHandler.RegisterClient(nc) + case "kafka": + log.Println("Initializing Kafka client...") + kc := broadcast.NewKafkaClient( + broadcast.SubTypeFull, + addr, + streamProcessor.Name, + streamProcessor.Topic, + config.General.BufferSizes.Websocket, + ) + broadcast.ClientHandler.RegisterClient(kc) + default: + log.Printf("Unknown stream processor type '%s' for %s. Skipping...\n", streamProcessor.Type, streamProcessor.Name) + } } return &cs, nil diff --git a/internal/config/config.go b/internal/config/config.go index 025b7fe..f3a938f 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -57,14 +57,14 @@ type Config struct { MetricsURL string `mapstructure:"metrics_url"` ExposeSystemMetrics bool `mapstructure:"expose_system_metrics"` } - StreamProcessing struct { - Kafka struct { - Enabled bool `yaml:"enabled"` - ServerAddr string `yaml:"server_addr"` - ServerPort int `yaml:"server_port"` - Topic string `yaml:"topic"` - } - } + StreamProcessing []struct { + Name string `yaml:"name"` + Type string `yaml:"type"` + Enabled bool `yaml:"enabled"` + ServerAddr string `yaml:"server_addr"` + ServerPort int `yaml:"server_port"` + Topic string `yaml:"topic"` + } `yaml:"stream_processing"` General struct { // DisableDefaultLogs indicates whether the default logs used in Google Chrome and provided by Google should be disabled. DisableDefaultLogs bool `mapstructure:"disable_default_logs"` From 66ba556ccbacc8a4990b6d57cedd841586db66d7 Mon Sep 17 00:00:00 2001 From: Rico Date: Tue, 5 Aug 2025 22:29:01 +0200 Subject: [PATCH 09/30] refactor: remove unnecessary log --- internal/certstream/certstream.go | 1 - 1 file changed, 1 deletion(-) diff --git a/internal/certstream/certstream.go b/internal/certstream/certstream.go index 49c7edf..d51ba71 100644 --- a/internal/certstream/certstream.go +++ b/internal/certstream/certstream.go @@ -56,7 +56,6 @@ func NewCertstreamServer(config config.Config) (*Certstream, error) { // Setup metrics server cs.setupMetrics(webserver) - log.Println(config.StreamProcessing) // Initialize the stream processors if configured and enabled. for _, streamProcessor := range config.StreamProcessing { if !streamProcessor.Enabled { From 38c53a182768f61b15424b559157488f68779acf Mon Sep 17 00:00:00 2001 From: Rico Date: Wed, 6 Aug 2025 01:01:37 +0200 Subject: [PATCH 10/30] refactor: change "kafka" in comments to "nsq" --- internal/broadcast/nsqclient.go | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/internal/broadcast/nsqclient.go b/internal/broadcast/nsqclient.go index 56091ec..9aa6b4e 100644 --- a/internal/broadcast/nsqclient.go +++ b/internal/broadcast/nsqclient.go @@ -9,7 +9,7 @@ import ( // NSQClient connects to a NSQ server in order to provide it with certificates. type NSQClient struct { - conn *nsq.Producer // Kafka connection + conn *nsq.Producer // nsq connection addr string topic string isConnected bool @@ -60,12 +60,12 @@ func NewNSQClient(subType SubscriptionType, addr, name, topic string, certBuffer return nsqc } -// reconnectHandler is a background job that attempts to reconnect to the Kafka server if the connection is lost. +// reconnectHandler is a background job that attempts to reconnect to the NSQ server if the connection is lost. func (c *NSQClient) reconnectHandler() { for { select { case <-c.stopChan: - log.Println("Stopping reconnectHandler for kafka producer:", c.addr) + log.Println("Stopping reconnectHandler for nsq producer:", c.addr) return default: if c.isConnected { From 543b98b4f5e62eb1657a7bc07a83c1a13da967d6 Mon Sep 17 00:00:00 2001 From: Rico Date: Sun, 23 Nov 2025 18:51:31 +0100 Subject: [PATCH 11/30] fix: broken code due to rebase --- config.sample.yaml | 3 --- go.sum | 2 ++ internal/broadcast/broadcastmanager.go | 8 ++++---- internal/certificatetransparency/ct-watcher.go | 2 -- internal/certstream/certstream.go | 2 +- internal/metrics/prometheus.go | 2 -- internal/web/server.go | 5 +---- 7 files changed, 8 insertions(+), 16 deletions(-) diff --git a/config.sample.yaml b/config.sample.yaml index 2a52eee..aeb3bcd 100644 --- a/config.sample.yaml +++ b/config.sample.yaml @@ -23,9 +23,6 @@ webserver: cert_key_path: "" # specify if the server should attempt to negotiate per message compression (RFC 7692) compression_enabled: false - # Use True-Client-IP, X-Real-IP or the X-Forwarded-For headers (in that order) to determine the real IP address of the client. - # If you are using a reverse proxy, you should set this to true. - real_ip: false whitelist: - "127.0.0.1/8" diff --git a/go.sum b/go.sum index 609977f..7b2c8cd 100644 --- a/go.sum +++ b/go.sum @@ -48,6 +48,8 @@ github.com/rogpeppe/go-internal v1.9.0/go.mod h1:WtVeX8xhTBvf0smdhujwtBcq4Qrzq/f github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= github.com/sagikazarmark/locafero v0.12.0 h1:/NQhBAkUb4+fH1jivKHWusDYFjMOOKU88eegjfxfHb4= github.com/sagikazarmark/locafero v0.12.0/go.mod h1:sZh36u/YSZ918v0Io+U9ogLYQJ9tLLBmM4eneO6WwsI= +github.com/segmentio/kafka-go v0.4.51 h1:JgDPPG75tC1rWIS2Me6MwcvXJ6f49UQ4HjAOef71Hno= +github.com/segmentio/kafka-go v0.4.51/go.mod h1:Y1gn60kzLEEaW28YshXyk2+VCUKbJ3Qr6DrnT3i4+9E= github.com/sergi/go-diff v1.4.0 h1:n/SP9D5ad1fORl+llWyN+D6qoUETXNZARKjyY2/KVCw= github.com/sergi/go-diff v1.4.0/go.mod h1:A0bzQcvG0E7Rwjx0REVgAGH58e96+X0MeOfepqsbeW4= github.com/spf13/afero v1.15.0 h1:b/YBCLWAJdFWJTN9cLhiXXcD7mzKn9Dm86dNnfyQw1I= diff --git a/internal/broadcast/broadcastmanager.go b/internal/broadcast/broadcastmanager.go index 2803231..b1a7fb6 100644 --- a/internal/broadcast/broadcastmanager.go +++ b/internal/broadcast/broadcastmanager.go @@ -57,10 +57,10 @@ func (bm *Dispatcher) UnregisterClient(clientName string) { // Close the broadcast channel of the client, otherwise this leads to a memory leak client.Close() - metrics.Prometheus.UnregisterClient(c.name) + metrics.Prometheus.UnregisterClient(client.Name()) - break - } + break + } } bm.clientLock.Unlock() @@ -115,7 +115,7 @@ func (bm *Dispatcher) broadcaster() { var data []byte // Take entry out of broadcast channel and generate JSON representations for the entry. - entry := <-bm.Broadcast + entry := <-bm.MessageQueue dataLite := entry.JSONLite() dataFull := entry.JSON() dataDomain := entry.JSONDomains() diff --git a/internal/certificatetransparency/ct-watcher.go b/internal/certificatetransparency/ct-watcher.go index 6db1174..6a271eb 100644 --- a/internal/certificatetransparency/ct-watcher.go +++ b/internal/certificatetransparency/ct-watcher.go @@ -13,8 +13,6 @@ import ( "sync/atomic" "time" - "github.com/google/trillian/client/backoff" - "github.com/d-Rickyy-b/certstream-server-go/internal/broadcast" "github.com/d-Rickyy-b/certstream-server-go/internal/config" "github.com/d-Rickyy-b/certstream-server-go/internal/metrics" diff --git a/internal/certstream/certstream.go b/internal/certstream/certstream.go index d51ba71..b2bd095 100644 --- a/internal/certstream/certstream.go +++ b/internal/certstream/certstream.go @@ -91,7 +91,7 @@ func NewCertstreamServer(config config.Config) (*Certstream, error) { } } - return &cs, nil + return cs, nil } // NewCertstreamFromConfigFile creates a new Certstream server from a config file. diff --git a/internal/metrics/prometheus.go b/internal/metrics/prometheus.go index 2f72bb7..caca426 100644 --- a/internal/metrics/prometheus.go +++ b/internal/metrics/prometheus.go @@ -10,8 +10,6 @@ import ( "sync" "time" - "github.com/d-Rickyy-b/certstream-server-go/internal/certificatetransparency" - "github.com/VictoriaMetrics/metrics" ) diff --git a/internal/web/server.go b/internal/web/server.go index 0e05542..bc6809f 100644 --- a/internal/web/server.go +++ b/internal/web/server.go @@ -20,10 +20,7 @@ import ( "github.com/gorilla/websocket" ) -var ( - ClientHandler = NewBroadcastManager() - upgrader websocket.Upgrader -) +var upgrader websocket.Upgrader type contextKey int From dea5349df3b46d222a396169f1ce74bb16fe26c7 Mon Sep 17 00:00:00 2001 From: Rico Date: Tue, 7 Jul 2026 01:50:46 +0200 Subject: [PATCH 12/30] refactor: multiple code fixes after rebase --- config.sample.yaml | 2 -- internal/broadcast/nsqclient.go | 4 +-- internal/broadcast/websocketclient.go | 36 +++++++++---------- .../websocketclient_test.go} | 18 +++++----- internal/metrics/prometheus.go | 17 ++------- internal/web/server.go | 17 +++------ 6 files changed, 33 insertions(+), 61 deletions(-) rename internal/{web/client_test.go => broadcast/websocketclient_test.go} (92%) diff --git a/config.sample.yaml b/config.sample.yaml index aeb3bcd..224a555 100644 --- a/config.sample.yaml +++ b/config.sample.yaml @@ -23,8 +23,6 @@ webserver: cert_key_path: "" # specify if the server should attempt to negotiate per message compression (RFC 7692) compression_enabled: false - whitelist: - - "127.0.0.1/8" prometheus: enabled: true diff --git a/internal/broadcast/nsqclient.go b/internal/broadcast/nsqclient.go index 9aa6b4e..ff9c029 100644 --- a/internal/broadcast/nsqclient.go +++ b/internal/broadcast/nsqclient.go @@ -1,10 +1,10 @@ package broadcast import ( - "github.com/nsqio/go-nsq" - "log" "time" + + "github.com/nsqio/go-nsq" ) // NSQClient connects to a NSQ server in order to provide it with certificates. diff --git a/internal/broadcast/websocketclient.go b/internal/broadcast/websocketclient.go index 642600a..29b913e 100644 --- a/internal/broadcast/websocketclient.go +++ b/internal/broadcast/websocketclient.go @@ -13,24 +13,25 @@ import ( const idChars = "abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789" -const ( - SubTypeFull SubscriptionType = iota - SubTypeLite - SubTypeDomain -) - -type SubscriptionType int - // WebsocketClient represents a single WebSocket client's connection to the server. type WebsocketClient struct { - conn *websocket.Conn + conn *websocket.Conn + userAgent string + hostIP string + hostPort string + realIPFromHeader string + *BaseClient } // NewWebsocketClient creates a new WebSocket client from the given connection. -func NewWebsocketClient(conn *websocket.Conn, subType SubscriptionType, name string, certBufferSize int) *WebsocketClient { +func NewWebsocketClient(conn *websocket.Conn, subType SubscriptionType, name, userAgent, hostIP, hostPort, realIPFromHeader string, certBufferSize int) *WebsocketClient { c := &WebsocketClient{ - conn: conn, + conn: conn, + userAgent: userAgent, + hostIP: hostIP, + hostPort: hostPort, + realIPFromHeader: realIPFromHeader, BaseClient: &BaseClient{ broadcastChan: make(chan []byte, certBufferSize), name: name, @@ -177,19 +178,14 @@ func sanitizeInput(s string) string { } // Name returns the name/identifier for this client. -func (c *client) Name() string { - var clientName string - - clientName = fmt.Sprintf("[%s] - ", c.id) +func (c *WebsocketClient) Name() string { + clientName := fmt.Sprintf("[%s] - ", c.name) - connIP := sanitizeInput(c.connectionIP) - connPort := sanitizeInput(c.connectionPort) realIP := sanitizeInput(c.realIPFromHeader) - - socket := net.JoinHostPort(connIP, connPort) + socket := net.JoinHostPort(c.hostIP, c.hostPort) // If the realIP is set and if it differs from the connection IP, return both the connection IP and the real IP. - if realIP != "" && realIP != connIP { + if realIP != "" && realIP != c.hostIP { clientName += fmt.Sprintf("%s (via %s)", socket, realIP) } else { clientName += socket diff --git a/internal/web/client_test.go b/internal/broadcast/websocketclient_test.go similarity index 92% rename from internal/web/client_test.go rename to internal/broadcast/websocketclient_test.go index b1e2b30..cc39dd3 100644 --- a/internal/web/client_test.go +++ b/internal/broadcast/websocketclient_test.go @@ -1,19 +1,17 @@ -package web +package broadcast import ( "fmt" "testing" ) -func newTestClient(id, connectionIP, connectionPort, realIP, userAgent string) *client { - c := &client{ - clientData: clientData{ - connectionIP: connectionIP, - connectionPort: connectionPort, - realIPFromHeader: realIP, - userAgent: userAgent, - }, - id: id, +func newTestClient(name, connectionIP, connectionPort, realIP, userAgent string) *WebsocketClient { + c := &WebsocketClient{ + hostIP: connectionIP, + hostPort: connectionPort, + realIPFromHeader: realIP, + userAgent: userAgent, + BaseClient: &BaseClient{name: name}, } return c } diff --git a/internal/metrics/prometheus.go b/internal/metrics/prometheus.go index caca426..a6a9ac4 100644 --- a/internal/metrics/prometheus.go +++ b/internal/metrics/prometheus.go @@ -64,35 +64,22 @@ func (pm *PrometheusExporter) RegisterGaugeMetricInt(label string, callback func } // RegisterClient registers a new gauge metric for the client with the given name. -func (pm *PrometheusExporter) RegisterClient(id, connIP, connPort, realIP, useragent string, skippedCertsCallback func() float64) { +func (pm *PrometheusExporter) RegisterClient(id string, skippedCertsCallback func() float64) { argMap := make(map[string]string) argMap["id"] = id - argMap["conn_ip"] = connIP - argMap["conn_port"] = connPort - argMap["real_ip"] = realIP - argMap["useragent"] = useragent label := createMetric("certstreamservergo_skipped_certs", argMap) - // label := fmt.Sprintf("certstreamservergo_skipped_certs{id=\"%s\",conn_ip=\"%s\",conn_port=\"%s\",real_ip=\"%s\",useragent=\"%s\"}", - // id, connIP, connPort, realIP, useragent) metrics.GetOrCreateGauge(label, skippedCertsCallback) } // UnregisterClient unregisters the metric for the client with the given name. -func (pm *PrometheusExporter) UnregisterClient(id, connIP, connPort, realIP, useragent string) { +func (pm *PrometheusExporter) UnregisterClient(id string) { argMap := make(map[string]string) argMap["id"] = id - argMap["conn_ip"] = connIP - argMap["conn_port"] = connPort - argMap["real_ip"] = realIP - argMap["useragent"] = useragent label := createMetric("certstreamservergo_skipped_certs", argMap) - // label := fmt.Sprintf("certstreamservergo_skipped_certs{id=\"%s\",conn_ip=\"%s\",conn_port=\"%s\",real_ip=\"%s\",useragent=\"%s\"}", - // id, connIP, connPort, realIP, useragent) - ok := metrics.UnregisterMetric(label) if !ok { log.Printf("failed to unregister metric '%s'", label) diff --git a/internal/web/server.go b/internal/web/server.go index bc6809f..72d0c8e 100644 --- a/internal/web/server.go +++ b/internal/web/server.go @@ -224,7 +224,7 @@ func initFullWebsocket(w http.ResponseWriter, r *http.Request) { return } - setupClient(connection, broadcast.SubTypeFull, r.RemoteAddr) + setupClient(connection, broadcast.SubTypeFull, r.RemoteAddr, r) } // initLiteWebsocket is called when a client connects to the / endpoint. @@ -236,7 +236,7 @@ func initLiteWebsocket(w http.ResponseWriter, r *http.Request) { return } - setupClient(connection, broadcast.SubTypeLite, r.RemoteAddr) + setupClient(connection, broadcast.SubTypeLite, r.RemoteAddr, r) } // initDomainWebsocket is called when a client connects to the /domains-only endpoint. @@ -248,7 +248,7 @@ func initDomainWebsocket(w http.ResponseWriter, r *http.Request) { return } - setupClient(connection, broadcast.SubTypeDomain, r.RemoteAddr) + setupClient(connection, broadcast.SubTypeDomain, r.RemoteAddr, r) } // upgradeConnection upgrades the connection to a websocket and returns the connection. @@ -279,7 +279,7 @@ func upgradeConnection(w http.ResponseWriter, r *http.Request) (*websocket.Conn, } // setupClient initializes a client struct and starts the broadcastHandler and websocket listener. -func setupClient(connection *websocket.Conn, subscriptionType broadcast.SubscriptionType, name string) { +func setupClient(connection *websocket.Conn, subscriptionType broadcast.SubscriptionType, name string, r *http.Request) { // Extract data from request origConnAddr, _ := r.Context().Value(origConnAddrKey).(string) @@ -294,14 +294,7 @@ func setupClient(connection *websocket.Conn, subscriptionType broadcast.Subscrip realIPFromHeader = r.RemoteAddr } - data := clientData{ - userAgent: r.Header.Get("User-Agent"), - connectionIP: hostIP, - connectionPort: hostPort, - realIPFromHeader: realIPFromHeader, - } - - c := broadcast.NewWebsocketClient(connection, subscriptionType, name, config.AppConfig.General.BufferSizes.Websocket) + c := broadcast.NewWebsocketClient(connection, subscriptionType, name, r.Header.Get("User-Agent"), hostIP, hostPort, realIPFromHeader, config.AppConfig.General.BufferSizes.Websocket) broadcast.ClientHandler.RegisterClient(c) } From d214a0c87a4268a7d556cb39581cfe8cec8380c6 Mon Sep 17 00:00:00 2001 From: Rico Date: Tue, 7 Jul 2026 01:50:59 +0200 Subject: [PATCH 13/30] refactor: fix typo --- cmd/certpicker/main.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/cmd/certpicker/main.go b/cmd/certpicker/main.go index cc82bf5..7323fad 100644 --- a/cmd/certpicker/main.go +++ b/cmd/certpicker/main.go @@ -65,7 +65,7 @@ func main() { log.Fatalln("Error getting entry from CT log: ", getEntryErr) } - // Loop over entries and pars each one. + // Loop over entries and parse each one. for _, leafEntry := range entries.Entries { rawLogEntry, err := ct.RawLogEntryFromLeaf(certID, &leafEntry) if err != nil { From 52a5c91da5000eb641929d0287207bf2087c5d3d Mon Sep 17 00:00:00 2001 From: Rico Date: Tue, 7 Jul 2026 01:59:18 +0200 Subject: [PATCH 14/30] refactor: fix config to use mapstructure --- internal/config/config.go | 22 +++++++++++----------- 1 file changed, 11 insertions(+), 11 deletions(-) diff --git a/internal/config/config.go b/internal/config/config.go index f3a938f..2a930fd 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -38,7 +38,7 @@ type BufferSizes struct { Websocket int `mapstructure:"websocket"` CTLog int `mapstructure:"ctlog"` BroadcastManager int `mapstructure:"broadcastmanager"` - Dispatcher int `mapstructure:"dispatcher"` + Dispatcher int `mapstructure:"dispatcher"` } type Config struct { @@ -58,13 +58,13 @@ type Config struct { ExposeSystemMetrics bool `mapstructure:"expose_system_metrics"` } StreamProcessing []struct { - Name string `yaml:"name"` - Type string `yaml:"type"` - Enabled bool `yaml:"enabled"` - ServerAddr string `yaml:"server_addr"` - ServerPort int `yaml:"server_port"` - Topic string `yaml:"topic"` - } `yaml:"stream_processing"` + Name string `mapstructure:"name"` + Type string `mapstructure:"type"` + Enabled bool `mapstructure:"enabled"` + ServerAddr string `mapstructure:"server_addr"` + ServerPort int `mapstructure:"server_port"` + Topic string `mapstructure:"topic"` + } `mapstructure:"stream_processing"` General struct { // DisableDefaultLogs indicates whether the default logs used in Google Chrome and provided by Google should be disabled. DisableDefaultLogs bool `mapstructure:"disable_default_logs"` @@ -73,10 +73,10 @@ type Config struct { AdditionalTiledLogs []LogConfig `mapstructure:"additional_tiled_logs"` ExcludedLogs []LogConfig `mapstructure:"excluded_logs"` // BufferSizes contains the buffer sizes for the different components of the server. They usually don't need any adjustments. - BufferSizes BufferSizes `mapstructure:"buffer_sizes"` + BufferSizes BufferSizes `mapstructure:"buffer_sizes"` // DropOldLogs indicates whether downloading CT-Logs should start at the latest index (true) or should from the beginning (false). - DropOldLogs *bool `mapstructure:"drop_old_logs"` - Recovery struct { + DropOldLogs *bool `mapstructure:"drop_old_logs"` + Recovery struct { Enabled bool `mapstructure:"enabled"` CTIndexFile string `mapstructure:"ct_index_file"` } `mapstructure:"recovery"` From c1781281be78b47aa65f3abe73bf1de441553e62 Mon Sep 17 00:00:00 2001 From: Rico Date: Tue, 7 Jul 2026 02:50:31 +0200 Subject: [PATCH 15/30] feat: add compression options for kafka --- internal/broadcast/kafkaclient.go | 15 ++++++++++++++- internal/certstream/certstream.go | 1 + internal/config/config.go | 20 ++++++++++++++------ 3 files changed, 29 insertions(+), 7 deletions(-) diff --git a/internal/broadcast/kafkaclient.go b/internal/broadcast/kafkaclient.go index a10a938..bebbb35 100644 --- a/internal/broadcast/kafkaclient.go +++ b/internal/broadcast/kafkaclient.go @@ -17,12 +17,13 @@ type KafkaClient struct { conn *kafka.Conn // Kafka connection addr string topic string + compression kafka.Compression isConnected bool BaseClient } // NewKafkaClient creates a new Kafka client that immediately connects to the configured Kafka server. -func NewKafkaClient(subType SubscriptionType, addr, name, topic string, certBufferSize int) *KafkaClient { +func NewKafkaClient(subType SubscriptionType, addr, name, topic, compression string, certBufferSize int) *KafkaClient { // Connect to the Kafka server conn, err := kafka.DialLeader(context.Background(), "tcp", addr, topic, 0) if err != nil { @@ -41,6 +42,18 @@ func NewKafkaClient(subType SubscriptionType, addr, name, topic string, certBuff }, } + switch compression { + case "gzip": + kc.compression = kafka.Gzip + case "snappy": + kc.compression = kafka.Snappy + case "lz4": + kc.compression = kafka.Lz4 + case "none": + default: + log.Println("invalid compression type:", compression) + } + go kc.broadcastHandler() go kc.reconnectHandler() diff --git a/internal/certstream/certstream.go b/internal/certstream/certstream.go index b2bd095..438e541 100644 --- a/internal/certstream/certstream.go +++ b/internal/certstream/certstream.go @@ -83,6 +83,7 @@ func NewCertstreamServer(config config.Config) (*Certstream, error) { addr, streamProcessor.Name, streamProcessor.Topic, + streamProcessor.Compression, config.General.BufferSizes.Websocket, ) broadcast.ClientHandler.RegisterClient(kc) diff --git a/internal/config/config.go b/internal/config/config.go index 2a930fd..a2f24d0 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -58,12 +58,13 @@ type Config struct { ExposeSystemMetrics bool `mapstructure:"expose_system_metrics"` } StreamProcessing []struct { - Name string `mapstructure:"name"` - Type string `mapstructure:"type"` - Enabled bool `mapstructure:"enabled"` - ServerAddr string `mapstructure:"server_addr"` - ServerPort int `mapstructure:"server_port"` - Topic string `mapstructure:"topic"` + Name string `mapstructure:"name"` + Type string `mapstructure:"type"` + Enabled bool `mapstructure:"enabled"` + ServerAddr string `mapstructure:"server_addr"` + ServerPort int `mapstructure:"server_port"` + Topic string `mapstructure:"topic"` + Compression string `mapstructure:"compression"` } `mapstructure:"stream_processing"` General struct { // DisableDefaultLogs indicates whether the default logs used in Google Chrome and provided by Google should be disabled. @@ -303,6 +304,13 @@ func validateConfig(config *Config) bool { } } + if len(config.StreamProcessing) > 0 { + for _, streamProcessing := range config.StreamProcessing { + streamProcessing.Type = strings.ToLower(streamProcessing.Type) + streamProcessing.Compression = strings.ToLower(streamProcessing.Compression) + } + } + if len(config.General.ExcludedLogs) > 0 { for _, excludedLog := range config.General.ExcludedLogs { excludedLog.Operator = strings.TrimSpace(excludedLog.Operator) From 0f96e5995c4d9d1665ba923e3bf5db26c5409609 Mon Sep 17 00:00:00 2001 From: Rico Date: Tue, 7 Jul 2026 02:51:52 +0200 Subject: [PATCH 16/30] feat: batch kafka messages for efficiency Sends messages to the broker in batches to improve throughput. Messages are collected and flushed when a maximum batch size is reached or after a configured wait time. This significantly reduces individual write calls. Messages are also dropped if the client isn't connected. --- internal/broadcast/kafkaclient.go | 46 +++++++++++++++++++++++-------- 1 file changed, 35 insertions(+), 11 deletions(-) diff --git a/internal/broadcast/kafkaclient.go b/internal/broadcast/kafkaclient.go index bebbb35..763462e 100644 --- a/internal/broadcast/kafkaclient.go +++ b/internal/broadcast/kafkaclient.go @@ -9,7 +9,9 @@ import ( ) const ( - TOPIC = "certstream" + maxBatchSize = 50 + maxBatchWait = 60 * time.Second + writeWait = 60 * time.Second ) // KafkaClient connects to a Kafka server in order to provide it with certificates. @@ -97,8 +99,6 @@ func (c *KafkaClient) reconnectHandler() { // Each client has a broadcastHandler that runs in the background and sends out the broadcast messages to the client. func (c *KafkaClient) broadcastHandler() { - writeWait := 60 * time.Second - defer func() { log.Println("Closing broadcast handler for kafka producer:", c.addr) if err := c.conn.Close(); err != nil { @@ -108,25 +108,49 @@ func (c *KafkaClient) broadcastHandler() { ClientHandler.UnregisterClient(c.name) }() + batch := make([]kafka.Message, 0, maxBatchSize) + t := time.NewTimer(maxBatchWait) + for { select { case <-c.stopChan: return case message := <-c.broadcastChan: + // Drop messages if not connected if !c.isConnected { continue } - _ = c.conn.SetWriteDeadline(time.Now().Add(writeWait)) + msg := kafka.Message{Value: message} + batch = append(batch, msg) - c.conn.Broker() - _, err := c.conn.WriteMessages( - kafka.Message{Value: message}, - ) - if err != nil { - c.isConnected = false - log.Println("Failed to write messages to kafka:", err) + // Write batch if it reaches max size + if len(batch) >= maxBatchSize { + c.writeBatch(batch) + batch = batch[:0] + t.Reset(maxBatchWait) } + case <-t.C: + if len(batch) == 0 { + continue + } + + // Write any remaining batch after maxBatchWait + c.writeBatch(batch) + batch = batch[:0] } } } + +func (c *KafkaClient) writeBatch(batch []kafka.Message) { + if len(batch) == 0 { + return + } + + _ = c.conn.SetWriteDeadline(time.Now().Add(writeWait)) + _, err := c.conn.WriteMessages(batch...) + if err != nil { + c.isConnected = false + log.Println("Failed to write messages to kafka:", err) + } +} From d66529c1c2316b3bde4f54f834abf775257f84a9 Mon Sep 17 00:00:00 2001 From: Rico Date: Fri, 10 Jul 2026 21:47:16 +0200 Subject: [PATCH 17/30] refactor: split config into multiple individual files The code for reading and validating the config grew quite big, so it made sense to split it up into multiple files and have a method called Valid() that checks for config validity on each individual struct. In a further iteration we should split the config_test.go into multiple test files each testing only the respective individual files. --- internal/config/buffersizes.go | 16 ++ internal/config/config.go | 247 ++-------------------------- internal/config/general.go | 115 +++++++++++++ internal/config/prometheus.go | 62 +++++++ internal/config/serverconfig.go | 11 ++ internal/config/streamprocessing.go | 72 ++++++++ internal/config/webserver.go | 69 ++++++++ 7 files changed, 358 insertions(+), 234 deletions(-) create mode 100644 internal/config/buffersizes.go create mode 100644 internal/config/general.go create mode 100644 internal/config/prometheus.go create mode 100644 internal/config/serverconfig.go create mode 100644 internal/config/streamprocessing.go create mode 100644 internal/config/webserver.go diff --git a/internal/config/buffersizes.go b/internal/config/buffersizes.go new file mode 100644 index 0000000..d826394 --- /dev/null +++ b/internal/config/buffersizes.go @@ -0,0 +1,16 @@ +package config + +type BufferSizes struct { + Websocket int `mapstructure:"websocket"` + CTLog int `mapstructure:"ctlog"` + BroadcastManager int `mapstructure:"broadcastmanager"` + Dispatcher int `mapstructure:"dispatcher"` +} + +func (b *BufferSizes) Valid() bool { + if b.Websocket <= 0 || b.CTLog <= 0 || b.BroadcastManager <= 0 || b.Dispatcher <= 0 { + return false + } + + return true +} diff --git a/internal/config/config.go b/internal/config/config.go index a2f24d0..30a3505 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -4,7 +4,6 @@ import ( "errors" "fmt" "log" - "net" "regexp" "strings" @@ -12,76 +11,20 @@ import ( ) var ( + // AppConfig holds the parsed configuration. AppConfig Config Version = "1.9.0" ErrInvalidConfig = errors.New("invalid configuration") + URLPathRegex = regexp.MustCompile(`^(/[a-zA-Z0-9\-._]+)+$`) + URLRegex = regexp.MustCompile(`^https?://[a-zA-Z0-9\-._]+(:[0-9]+)?(/[a-zA-Z0-9\-._]+)*/?$`) ) -type ServerConfig struct { - ListenAddr string `mapstructure:"listen_addr"` - ListenPort int `mapstructure:"listen_port"` - CertPath string `mapstructure:"cert_path"` - CertKeyPath string `mapstructure:"cert_key_path"` - RealIP bool `mapstructure:"real_ip"` - TrustedProxies []string `mapstructure:"trusted_proxies"` - Whitelist []string `mapstructure:"whitelist"` -} - -type LogConfig struct { - Operator string `mapstructure:"operator"` - URL string `mapstructure:"url"` - Description string `mapstructure:"description"` -} - -type BufferSizes struct { - Websocket int `mapstructure:"websocket"` - CTLog int `mapstructure:"ctlog"` - BroadcastManager int `mapstructure:"broadcastmanager"` - Dispatcher int `mapstructure:"dispatcher"` -} - type Config struct { - Webserver struct { - ServerConfig `mapstructure:",squash"` - - FullURL string `mapstructure:"full_url"` - LiteURL string `mapstructure:"lite_url"` - DomainsOnlyURL string `mapstructure:"domains_only_url"` - CompressionEnabled bool `mapstructure:"compression_enabled"` - } - Prometheus struct { - ServerConfig `mapstructure:",squash"` - - Enabled bool `mapstructure:"enabled"` - MetricsURL string `mapstructure:"metrics_url"` - ExposeSystemMetrics bool `mapstructure:"expose_system_metrics"` - } - StreamProcessing []struct { - Name string `mapstructure:"name"` - Type string `mapstructure:"type"` - Enabled bool `mapstructure:"enabled"` - ServerAddr string `mapstructure:"server_addr"` - ServerPort int `mapstructure:"server_port"` - Topic string `mapstructure:"topic"` - Compression string `mapstructure:"compression"` - } `mapstructure:"stream_processing"` - General struct { - // DisableDefaultLogs indicates whether the default logs used in Google Chrome and provided by Google should be disabled. - DisableDefaultLogs bool `mapstructure:"disable_default_logs"` - // AdditionalLogs contains additional logs provided by the user that can be used in addition to the default logs. - AdditionalLogs []LogConfig `mapstructure:"additional_logs"` - AdditionalTiledLogs []LogConfig `mapstructure:"additional_tiled_logs"` - ExcludedLogs []LogConfig `mapstructure:"excluded_logs"` - // BufferSizes contains the buffer sizes for the different components of the server. They usually don't need any adjustments. - BufferSizes BufferSizes `mapstructure:"buffer_sizes"` - // DropOldLogs indicates whether downloading CT-Logs should start at the latest index (true) or should from the beginning (false). - DropOldLogs *bool `mapstructure:"drop_old_logs"` - Recovery struct { - Enabled bool `mapstructure:"enabled"` - CTIndexFile string `mapstructure:"ct_index_file"` - } `mapstructure:"recovery"` - } + Webserver Webserver + Prometheus Prometheus + StreamProcessing []StreamProcessor `mapstructure:"stream_processing"` + General General } // ReadConfig reads the configuration using Viper and returns a filled Config struct. @@ -185,188 +128,24 @@ func loadConfigFromViper(v *viper.Viper) (Config, error) { // validateConfig validates the config values and sets defaults for missing values. func validateConfig(config *Config) bool { // Still matches invalid IP addresses but good enough for detecting completely wrong formats - URLPathRegex := regexp.MustCompile(`^(/[a-zA-Z0-9\-._]+)+$`) - URLRegex := regexp.MustCompile(`^https?://[a-zA-Z0-9\-._]+(:[0-9]+)?(/[a-zA-Z0-9\-._]+)*/?$`) // Check webserver config - if config.Webserver.ListenAddr == "" || net.ParseIP(config.Webserver.ListenAddr) == nil { - log.Fatalln("Webhook listen IP is not a valid IP: ", config.Webserver.ListenAddr) + if !config.Webserver.Valid() { return false } - if config.Webserver.ListenPort == 0 { - log.Fatalln("Webhook listen port is not set") + if !config.Prometheus.Valid() { return false } - if config.Webserver.FullURL == "" || !URLPathRegex.MatchString(config.Webserver.FullURL) { - log.Println("Webhook full URL is not set or does not match pattern '/...'") - - config.Webserver.FullURL = "/full-stream" - } - - if config.Webserver.LiteURL == "" || !URLPathRegex.MatchString(config.Webserver.FullURL) { - log.Println("Webhook lite URL is not set or does not match pattern '/...'") - - config.Webserver.LiteURL = "/" - } - - if config.Webserver.DomainsOnlyURL == "" || !URLPathRegex.MatchString(config.Webserver.DomainsOnlyURL) { - log.Println("Webhook domains only URL is not set or does not match pattern '/...'") - - config.Webserver.FullURL = "/domains-only" - } - - if config.Webserver.FullURL == config.Webserver.LiteURL { - log.Fatalln("Webhook full URL is the same as lite URL - please fix the config!") - } - - if config.Webserver.DomainsOnlyURL == "" { - config.Webserver.FullURL = "/domains-only" - } - - for _, ip := range config.Webserver.TrustedProxies { - if net.ParseIP(ip) != nil { - continue - } - - _, _, err := net.ParseCIDR(ip) - if err != nil { - log.Fatalln("Invalid IP/CIDR in webserver trusted_proxies: ", ip) + for _, processor := range config.StreamProcessing { + if !processor.Valid() { return false } } - //nolint:nestif - if config.Prometheus.Enabled { - if config.Prometheus.ListenAddr == "" || net.ParseIP(config.Prometheus.ListenAddr) == nil { - log.Fatalln("Metrics export IP is not a valid IP") - return false - } - - if config.Prometheus.ListenPort == 0 { - log.Fatalln("Metrics export port is not set") - return false - } - - if config.Prometheus.Whitelist == nil { - config.Prometheus.Whitelist = []string{} - } - - // Check if IPs in whitelist match pattern - for _, ip := range config.Prometheus.Whitelist { - if net.ParseIP(ip) != nil { - continue - } - - // Provided entry is not an IP, check if it's a CIDR range - _, _, err := net.ParseCIDR(ip) - if err != nil { - log.Fatalln("Invalid IP in metrics whitelist: ", ip) - return false - } - } - - for _, ip := range config.Prometheus.TrustedProxies { - if net.ParseIP(ip) != nil { - continue - } - - _, _, err := net.ParseCIDR(ip) - if err != nil { - log.Fatalln("Invalid IP/CIDR in prometheus trusted_proxies: ", ip) - return false - } - } - } - - var validLogs, validTiledLogs, validExcludedLogs []LogConfig - - if len(config.General.AdditionalLogs) > 0 { - for _, ctLog := range config.General.AdditionalLogs { - if !URLRegex.MatchString(ctLog.URL) { - log.Println("Ignoring invalid additional log URL: ", ctLog.URL) - continue - } - - validLogs = append(validLogs, ctLog) - } - } - - if len(config.General.AdditionalTiledLogs) > 0 { - for _, ctLog := range config.General.AdditionalTiledLogs { - if !URLRegex.MatchString(ctLog.URL) { - log.Println("Ignoring invalid additional log URL: ", ctLog.URL) - continue - } - - validTiledLogs = append(validTiledLogs, ctLog) - } - } - - if len(config.StreamProcessing) > 0 { - for _, streamProcessing := range config.StreamProcessing { - streamProcessing.Type = strings.ToLower(streamProcessing.Type) - streamProcessing.Compression = strings.ToLower(streamProcessing.Compression) - } - } - - if len(config.General.ExcludedLogs) > 0 { - for _, excludedLog := range config.General.ExcludedLogs { - excludedLog.Operator = strings.TrimSpace(excludedLog.Operator) - excludedLog.URL = strings.TrimSpace(excludedLog.URL) - - if excludedLog.Operator == "" && excludedLog.URL == "" { - log.Println("Ignoring empty excluded_logs entry. Set operator and/or url.") - continue - } - - if excludedLog.URL != "" && !URLRegex.MatchString(excludedLog.URL) { - log.Println("Ignoring invalid excluded log URL: ", excludedLog.URL) - continue - } - - validExcludedLogs = append(validExcludedLogs, excludedLog) - } - } - - config.General.AdditionalLogs = validLogs - config.General.AdditionalTiledLogs = validTiledLogs - config.General.ExcludedLogs = validExcludedLogs - - if len(config.General.AdditionalLogs) == 0 && len(config.General.AdditionalTiledLogs) == 0 && config.General.DisableDefaultLogs { - log.Fatalln("Default logs are disabled, but no additional logs are configured. Please add at least one log to the config or enable default logs.") - } - - if config.General.BufferSizes.Websocket <= 0 { - config.General.BufferSizes.Websocket = 300 - } - - if config.General.BufferSizes.CTLog <= 0 { - config.General.BufferSizes.CTLog = 1000 - } - - // For backward compatibility, copy value from deprecated BroadcastManager field - if config.General.BufferSizes.BroadcastManager != 0 { - config.General.BufferSizes.Dispatcher = config.General.BufferSizes.BroadcastManager - } - - if config.General.BufferSizes.Dispatcher <= 0 { - config.General.BufferSizes.Dispatcher = 10000 - } - - // If the cleanup flag is not set, default to true - if config.General.DropOldLogs == nil { - log.Println("drop_old_logs is not set, defaulting to true") - - defaultCleanup := true - config.General.DropOldLogs = &defaultCleanup - } - - if config.General.Recovery.Enabled && config.General.Recovery.CTIndexFile == "" { - log.Println("Recovery enabled but no index file specified. Defaulting to ./ct_index.json") - - config.General.Recovery.CTIndexFile = "./ct_index.json" + if !config.General.Valid() { + return false } return true diff --git a/internal/config/general.go b/internal/config/general.go new file mode 100644 index 0000000..4a244e9 --- /dev/null +++ b/internal/config/general.go @@ -0,0 +1,115 @@ +package config + +import ( + "log" + "strings" +) + +type LogConfig struct { + Operator string `mapstructure:"operator"` + URL string `mapstructure:"url"` + Description string `mapstructure:"description"` +} + +type General struct { + // DisableDefaultLogs indicates whether the default logs used in Google Chrome and provided by Google should be disabled. + DisableDefaultLogs bool `mapstructure:"disable_default_logs"` + // AdditionalLogs contains additional logs provided by the user that can be used in addition to the default logs. + AdditionalLogs []LogConfig `mapstructure:"additional_logs"` + AdditionalTiledLogs []LogConfig `mapstructure:"additional_tiled_logs"` + ExcludedLogs []LogConfig `mapstructure:"excluded_logs"` + // BufferSizes contains the buffer sizes for the different components of the server. They usually don't need any adjustments. + BufferSizes BufferSizes `mapstructure:"buffer_sizes"` + // DropOldLogs indicates whether downloading CT-Logs should start at the latest index (true) or should from the beginning (false). + DropOldLogs *bool `mapstructure:"drop_old_logs"` + Recovery struct { + Enabled bool `mapstructure:"enabled"` + CTIndexFile string `mapstructure:"ct_index_file"` + } `mapstructure:"recovery"` +} + +func (g *General) Valid() bool { + var validLogs, validTiledLogs, validExcludedLogs []LogConfig + + if len(g.AdditionalLogs) > 0 { + for _, ctLog := range g.AdditionalLogs { + if !URLRegex.MatchString(ctLog.URL) { + log.Println("Ignoring invalid additional log URL: ", ctLog.URL) + continue + } + + validLogs = append(validLogs, ctLog) + } + } + + if len(g.AdditionalTiledLogs) > 0 { + for _, ctLog := range g.AdditionalTiledLogs { + if !URLRegex.MatchString(ctLog.URL) { + log.Println("Ignoring invalid additional log URL: ", ctLog.URL) + continue + } + + validTiledLogs = append(validTiledLogs, ctLog) + } + } + + if len(g.ExcludedLogs) > 0 { + for _, excludedLog := range g.ExcludedLogs { + excludedLog.Operator = strings.TrimSpace(excludedLog.Operator) + excludedLog.URL = strings.TrimSpace(excludedLog.URL) + + if excludedLog.Operator == "" && excludedLog.URL == "" { + log.Println("Ignoring empty excluded_logs entry. Set operator and/or url.") + continue + } + + if excludedLog.URL != "" && !URLRegex.MatchString(excludedLog.URL) { + log.Println("Ignoring invalid excluded log URL: ", excludedLog.URL) + continue + } + + validExcludedLogs = append(validExcludedLogs, excludedLog) + } + } + + g.AdditionalLogs = validLogs + g.AdditionalTiledLogs = validTiledLogs + g.ExcludedLogs = validExcludedLogs + + if len(g.AdditionalLogs) == 0 && len(g.AdditionalTiledLogs) == 0 && g.DisableDefaultLogs { + log.Fatalln("Default logs are disabled, but no additional logs are configured. Please add at least one log to the config or enable default logs.") + } + + if g.BufferSizes.Websocket <= 0 { + g.BufferSizes.Websocket = 300 + } + + if g.BufferSizes.CTLog <= 0 { + g.BufferSizes.CTLog = 1000 + } + + // For backward compatibility, copy value from deprecated BroadcastManager field + if g.BufferSizes.BroadcastManager != 0 { + g.BufferSizes.Dispatcher = g.BufferSizes.BroadcastManager + } + + if g.BufferSizes.Dispatcher <= 0 { + g.BufferSizes.Dispatcher = 10000 + } + + // If the cleanup flag is not set, default to true + if g.DropOldLogs == nil { + log.Println("drop_old_logs is not set, defaulting to true") + + defaultCleanup := true + g.DropOldLogs = &defaultCleanup + } + + if g.Recovery.Enabled && g.Recovery.CTIndexFile == "" { + log.Println("Recovery enabled but no index file specified. Defaulting to ./ct_index.json") + + g.Recovery.CTIndexFile = "./ct_index.json" + } + + return true +} diff --git a/internal/config/prometheus.go b/internal/config/prometheus.go new file mode 100644 index 0000000..cfbcf58 --- /dev/null +++ b/internal/config/prometheus.go @@ -0,0 +1,62 @@ +package config + +import ( + "log" + "net" +) + +type Prometheus struct { + ServerConfig `mapstructure:",squash"` + + Enabled bool `mapstructure:"enabled"` + MetricsURL string `mapstructure:"metrics_url"` + ExposeSystemMetrics bool `mapstructure:"expose_system_metrics"` +} + +func (p *Prometheus) Valid() bool { + if !p.Enabled { + return true + } + + if p.ListenAddr == "" || net.ParseIP(p.ListenAddr) == nil { + log.Fatalln("Metrics export IP is not a valid IP") + return false + } + + if p.ListenPort == 0 { + log.Fatalln("Metrics export port is not set") + return false + } + + if p.Whitelist == nil { + p.Whitelist = []string{} + } + + // Check if IPs in whitelist match pattern + for _, ip := range p.Whitelist { + if net.ParseIP(ip) != nil { + continue + } + + // Provided entry is not an IP, check if it's a CIDR range + _, _, err := net.ParseCIDR(ip) + if err != nil { + log.Fatalln("Invalid IP in metrics whitelist: ", ip) + return false + } + } + + for _, ip := range p.TrustedProxies { + if net.ParseIP(ip) != nil { + continue + } + + _, _, err := net.ParseCIDR(ip) + if err != nil { + log.Fatalln("Invalid IP/CIDR in prometheus trusted_proxies: ", ip) + return false + } + } + + return true +} diff --git a/internal/config/serverconfig.go b/internal/config/serverconfig.go new file mode 100644 index 0000000..a246d4a --- /dev/null +++ b/internal/config/serverconfig.go @@ -0,0 +1,11 @@ +package config + +type ServerConfig struct { + ListenAddr string `mapstructure:"listen_addr"` + ListenPort int `mapstructure:"listen_port"` + CertPath string `mapstructure:"cert_path"` + CertKeyPath string `mapstructure:"cert_key_path"` + RealIP bool `mapstructure:"real_ip"` + TrustedProxies []string `mapstructure:"trusted_proxies"` + Whitelist []string `mapstructure:"whitelist"` +} diff --git a/internal/config/streamprocessing.go b/internal/config/streamprocessing.go new file mode 100644 index 0000000..4c47ab1 --- /dev/null +++ b/internal/config/streamprocessing.go @@ -0,0 +1,72 @@ +package config + +import ( + "log" + "net" +) + +type StreamProcessorType string + +const ( + StreamProcessorTypeKafka StreamProcessorType = "kafka" + StreamProcessorTypeNQS StreamProcessorType = "nqs" +) + +type StreamProcessorCompression string + +const ( + StreamProcessorCompressionNone StreamProcessorCompression = "none" + StreamProcessorCompressionGzip StreamProcessorCompression = "gzip" + StreamProcessorCompressionSnappy StreamProcessorCompression = "snappy" +) + +func (c StreamProcessorCompression) Valid() bool { + switch c { + case StreamProcessorCompressionNone, StreamProcessorCompressionGzip, StreamProcessorCompressionSnappy, "": + return true + default: + return false + } +} + +type StreamProcessor struct { + Name string `mapstructure:"name"` + Type StreamProcessorType `mapstructure:"type"` + Enabled bool `mapstructure:"enabled"` + ServerAddr string `mapstructure:"server_addr"` + ServerPort int `mapstructure:"server_port"` + Topic string `mapstructure:"topic"` + Compression StreamProcessorCompression `mapstructure:"compression"` +} + +func (s *StreamProcessor) Valid() bool { + ip := net.ParseIP(s.ServerAddr) + if ip == nil { + log.Fatalln("Invalid IP for stream processor:", s.ServerAddr) + return false + } + + if s.ServerPort <= 0 { + log.Fatalln("Invalid server port for stream processor:", s.ServerPort) + return false + } + + switch s.Type { + case StreamProcessorTypeKafka, StreamProcessorTypeNQS: + default: + log.Fatalf("Invalid stream processor type '%s' for name '%s'\n", s.Type, s.Name) + return false + } + + if s.Topic == "" { + log.Println("Found stream processing config with empty topic - using \"certstream\" as topic for name", s.Name) + s.Topic = "certstream" + } + + if !s.Compression.Valid() { + log.Fatalf("Invalid compression '%s' for stream processor '%s'\n", s.Compression, s.Name) + return false + } + + return true +} diff --git a/internal/config/webserver.go b/internal/config/webserver.go new file mode 100644 index 0000000..fb846d9 --- /dev/null +++ b/internal/config/webserver.go @@ -0,0 +1,69 @@ +package config + +import ( + "log" + "net" +) + +type Webserver struct { + ServerConfig `mapstructure:",squash"` + + FullURL string `mapstructure:"full_url"` + LiteURL string `mapstructure:"lite_url"` + DomainsOnlyURL string `mapstructure:"domains_only_url"` + CompressionEnabled bool `mapstructure:"compression_enabled"` +} + +func (w *Webserver) Valid() bool { + // Still matches invalid IP addresses but good enough for detecting completely wrong formats + + if w.ListenAddr == "" || net.ParseIP(w.ListenAddr) == nil { + log.Fatalln("Webhook listen IP is not a valid IP: ", w.ListenAddr) + return false + } + + if w.ListenPort == 0 { + log.Fatalln("Webhook listen port is not set") + return false + } + + if w.FullURL == "" || !URLPathRegex.MatchString(w.FullURL) { + log.Println("Webhook full URL is not set or does not match pattern '/...'") + + w.FullURL = "/full-stream" + } + + if w.LiteURL == "" || !URLPathRegex.MatchString(w.FullURL) { + log.Println("Webhook lite URL is not set or does not match pattern '/...'") + + w.LiteURL = "/" + } + + if w.DomainsOnlyURL == "" || !URLPathRegex.MatchString(w.DomainsOnlyURL) { + log.Println("Webhook domains only URL is not set or does not match pattern '/...'") + + w.FullURL = "/domains-only" + } + + if w.FullURL == w.LiteURL { + log.Fatalln("Webhook full URL is the same as lite URL - please fix the config!") + } + + if w.DomainsOnlyURL == "" { + w.FullURL = "/domains-only" + } + + for _, ip := range w.TrustedProxies { + if net.ParseIP(ip) != nil { + continue + } + + _, _, err := net.ParseCIDR(ip) + if err != nil { + log.Fatalln("Invalid IP/CIDR in webserver trusted_proxies: ", ip) + return false + } + } + + return true +} From 8fc04dec542bcd68621a9863b5d63f2b3ebabfec Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 01:32:18 +0200 Subject: [PATCH 18/30] refactor: rename config variable to cfg --- internal/certstream/certstream.go | 21 +++++++++++---------- 1 file changed, 11 insertions(+), 10 deletions(-) diff --git a/internal/certstream/certstream.go b/internal/certstream/certstream.go index 438e541..d03c759 100644 --- a/internal/certstream/certstream.go +++ b/internal/certstream/certstream.go @@ -35,8 +35,8 @@ func NewRawCertstream(config config.Config) *Certstream { } // NewCertstreamServer creates a new Certstream server from a config struct. -func NewCertstreamServer(config config.Config) (*Certstream, error) { - cs := NewRawCertstream(config) +func NewCertstreamServer(cfg config.Config) (*Certstream, error) { + cs := NewRawCertstream(cfg) // Start the broadcast dispatcher broadcast.NewDispatcher() @@ -45,10 +45,10 @@ func NewCertstreamServer(config config.Config) (*Certstream, error) { // TODO: add support do disable websocket Server // Initialize the webserver used for the websocket server webserver := web.NewWebsocketServer( - config.Webserver.ListenAddr, - config.Webserver.ListenPort, - config.Webserver.CertPath, - config.Webserver.CertKeyPath, + cfg.Webserver.ListenAddr, + cfg.Webserver.ListenPort, + cfg.Webserver.CertPath, + cfg.Webserver.CertKeyPath, ) cs.webserver = webserver cs.watcher = certificatetransparency.NewWatcher() @@ -57,7 +57,7 @@ func NewCertstreamServer(config config.Config) (*Certstream, error) { cs.setupMetrics(webserver) // Initialize the stream processors if configured and enabled. - for _, streamProcessor := range config.StreamProcessing { + for _, streamProcessor := range cfg.StreamProcessing { if !streamProcessor.Enabled { continue } @@ -73,7 +73,7 @@ func NewCertstreamServer(config config.Config) (*Certstream, error) { addr, streamProcessor.Name, streamProcessor.Topic, - config.General.BufferSizes.Websocket, + cfg.General.BufferSizes.Websocket, ) broadcast.ClientHandler.RegisterClient(nc) case "kafka": @@ -83,9 +83,10 @@ func NewCertstreamServer(config config.Config) (*Certstream, error) { addr, streamProcessor.Name, streamProcessor.Topic, - streamProcessor.Compression, - config.General.BufferSizes.Websocket, + string(streamProcessor.Compression), + cfg.General.BufferSizes.Websocket, ) + broadcast.ClientHandler.RegisterClient(kc) default: log.Printf("Unknown stream processor type '%s' for %s. Skipping...\n", streamProcessor.Type, streamProcessor.Name) From 4d005db3f6d8029bbef3d592b1f26ed061edf86f Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 01:32:57 +0200 Subject: [PATCH 19/30] refactor: make stream type configurable for stream processing --- internal/certstream/certstream.go | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/internal/certstream/certstream.go b/internal/certstream/certstream.go index d03c759..b3b3815 100644 --- a/internal/certstream/certstream.go +++ b/internal/certstream/certstream.go @@ -65,11 +65,21 @@ func NewCertstreamServer(cfg config.Config) (*Certstream, error) { addr := net.JoinHostPort(streamProcessor.ServerAddr, strconv.Itoa(streamProcessor.ServerPort)) log.Printf("Initializing stream processor: %s at %s\n", streamProcessor.Name, addr) + var subscriptionType broadcast.SubscriptionType + switch streamProcessor.Stream { + case config.StreamTypeFull: + subscriptionType = broadcast.SubTypeFull + case config.StreamTypeLite: + subscriptionType = broadcast.SubTypeLite + case config.StreamTypeDomainsOnly: + subscriptionType = broadcast.SubTypeDomain + } + switch streamProcessor.Type { case "nsq": log.Println("Initializing NSQ client...") nc := broadcast.NewNSQClient( - broadcast.SubTypeFull, + subscriptionType, addr, streamProcessor.Name, streamProcessor.Topic, @@ -79,7 +89,7 @@ func NewCertstreamServer(cfg config.Config) (*Certstream, error) { case "kafka": log.Println("Initializing Kafka client...") kc := broadcast.NewKafkaClient( - broadcast.SubTypeFull, + subscriptionType, addr, streamProcessor.Name, streamProcessor.Topic, From 258a6c5fef1f871578da65d1a02ce87ff45c6bbf Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 01:34:07 +0200 Subject: [PATCH 20/30] refactor(nsqclient): improve formatting and error handling --- internal/broadcast/nsqclient.go | 21 +++++++++++++++------ 1 file changed, 15 insertions(+), 6 deletions(-) diff --git a/internal/broadcast/nsqclient.go b/internal/broadcast/nsqclient.go index ff9c029..9b99bad 100644 --- a/internal/broadcast/nsqclient.go +++ b/internal/broadcast/nsqclient.go @@ -31,6 +31,7 @@ func NewNSQClient(subType SubscriptionType, addr, name, topic string, certBuffer // Instantiate a producer. conf := nsq.NewConfig() + conn, err := nsq.NewProducer(addr, conf) if err != nil { log.Println(err) @@ -66,6 +67,7 @@ func (c *NSQClient) reconnectHandler() { select { case <-c.stopChan: log.Println("Stopping reconnectHandler for nsq producer:", c.addr) + return default: if c.isConnected { @@ -73,6 +75,7 @@ func (c *NSQClient) reconnectHandler() { time.Sleep(5 * time.Second) continue } + // Attempt to connect to the NSQ server err := c.conn.Ping() if err != nil { @@ -90,19 +93,25 @@ func (c *NSQClient) reconnectHandler() { // Each client has a broadcastHandler that runs in the background and sends out the broadcast messages to the client. func (c *NSQClient) broadcastHandler() { - // writeWait := 60 * time.Second - defer func() { log.Println("Closing broadcast handler for nsq producer:", c.addr) - // Gracefully stop the producer when appropriate (e.g. before shutting down the service) - c.conn.Stop() + if c.conn != nil { + // Gracefully stop the producer when appropriate (e.g. before shutting down the service) + c.conn.Stop() + } + }() for { select { case <-c.stopChan: return - case message := <-c.broadcastChan: + case message, ok := <-c.broadcastChan: + if !ok { + log.Println("broadcastChan closed for nsqClient:", c.addr) + return + } + if !c.isConnected { continue } @@ -111,7 +120,7 @@ func (c *NSQClient) broadcastHandler() { // Messages can also be sent asynchronously and/or in batches. err := c.conn.Publish(c.topic, message) if err != nil { - log.Println("Error writing to NSQ topic:", err) + log.Println("Failed to write messages to NSQ:", err) c.isConnected = false } } From 4cd6c16c76543db987555023fecb8d287f7a2dbe Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 01:36:06 +0200 Subject: [PATCH 21/30] refactor(nsqclient): minor logic fixes --- internal/broadcast/nsqclient.go | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/internal/broadcast/nsqclient.go b/internal/broadcast/nsqclient.go index 9b99bad..de4d4c0 100644 --- a/internal/broadcast/nsqclient.go +++ b/internal/broadcast/nsqclient.go @@ -67,6 +67,7 @@ func (c *NSQClient) reconnectHandler() { select { case <-c.stopChan: log.Println("Stopping reconnectHandler for nsq producer:", c.addr) + c.conn.Stop() return default: @@ -100,6 +101,7 @@ func (c *NSQClient) broadcastHandler() { c.conn.Stop() } + ClientHandler.UnregisterClient(c.name) }() for { @@ -112,7 +114,9 @@ func (c *NSQClient) broadcastHandler() { return } + // Drop messages if not connected if !c.isConnected { + time.Sleep(5 * time.Second) continue } From fec2d29240f8584acb14e160e0fe0e5a11808ef5 Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 01:36:29 +0200 Subject: [PATCH 22/30] feat(nsqclient): add batch processing feature --- internal/broadcast/nsqclient.go | 47 ++++++++++++++++++++++++++++++--- 1 file changed, 44 insertions(+), 3 deletions(-) diff --git a/internal/broadcast/nsqclient.go b/internal/broadcast/nsqclient.go index de4d4c0..f8e042e 100644 --- a/internal/broadcast/nsqclient.go +++ b/internal/broadcast/nsqclient.go @@ -7,6 +7,11 @@ import ( "github.com/nsqio/go-nsq" ) +const ( + nsqMaxBatchSize = 50 + nsqMaxBatchWait = 1 * time.Second +) + // NSQClient connects to a NSQ server in order to provide it with certificates. type NSQClient struct { conn *nsq.Producer // nsq connection @@ -104,6 +109,9 @@ func (c *NSQClient) broadcastHandler() { ClientHandler.UnregisterClient(c.name) }() + batch := make([][]byte, 0, nsqMaxBatchSize) + t := time.NewTicker(nsqMaxBatchWait) + for { select { case <-c.stopChan: @@ -120,13 +128,46 @@ func (c *NSQClient) broadcastHandler() { continue } - // Synchronously publish a single message to the specified topic. - // Messages can also be sent asynchronously and/or in batches. - err := c.conn.Publish(c.topic, message) + batch = append(batch, message) + + // Write batch if it reaches max size + if len(batch) >= nsqMaxBatchSize { + err := c.writeBatch(batch) + if err != nil { + log.Println("Failed to write messages to NSQ:", err) + c.isConnected = false + } + batch = batch[:0] + t.Reset(kafkaMaxBatchWait) + } + case <-t.C: + // If batch size has not reached nsqMaxBatchSize, write the batch after nsqMaxBatchWait + if len(batch) == 0 { + continue + } + + err := c.writeBatch(batch) if err != nil { log.Println("Failed to write messages to NSQ:", err) c.isConnected = false } + batch = batch[:0] } } } + +func (c *NSQClient) writeBatch(batch [][]byte) error { + if len(batch) == 0 { + return nil + } + + // Synchronously publish a batch of messages to the specified topic. + err := c.conn.MultiPublish(c.topic, batch) + if err != nil { + log.Println("Failed to write messages to NSQ:", err) + c.isConnected = false + return err + } + + return nil +} From f8ca66e9715f420e20a0a512f08e71328ef3b85f Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 01:37:30 +0200 Subject: [PATCH 23/30] refactor(config.sample.yml): add new options to sample config --- config.sample.yaml | 23 +++++++++++++++++++---- 1 file changed, 19 insertions(+), 4 deletions(-) diff --git a/config.sample.yaml b/config.sample.yaml index 224a555..bed04c4 100644 --- a/config.sample.yaml +++ b/config.sample.yaml @@ -39,14 +39,29 @@ prometheus: # Configuration related to external stream processing tools go here. stream_processing: - - name: "kafka" - enabled: false + # List of stream processing configurations. + # Each configuration specifies a stream processing tool to use and its settings. + - name: "kafka-prod" + # Type describes the type of stream processing tool to use. + # Supported types: "kafka", "nqs" + type: "kafka" + # The stream type can be any of "full", "lite", "domains-only" - same as the certstream websocket urls + # Default is "full" + stream: "full" + # Enable or disable this stream processing configuration. + enabled: true + # The server address and port to connect to the stream processing tool server_addr: "127.0.0.1" server_port: 9092 topic: "certstream" + # Available compression types depend on the stream processing tool. + # E.g. for Kafka, supported types are "none" / "", "gzip", "snappy", "lz4", "zstd" + compression: "" - - name: "nqs" - enabled: true + - name: "nqs-prod" + type: "nqs" + stream: "lite" + enabled: false server_addr: "127.0.0.1" server_port: 9092 topic: "certstream" From 0c52b420f486dc44dcec0d45cd65ccd76ba94218 Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 01:54:39 +0200 Subject: [PATCH 24/30] refactor: improve streamprocessing struct and methods Made stream type configurable so that users can decide which stream they want to feed into kafka. --- internal/certstream/certstream.go | 2 +- internal/config/streamprocessing.go | 87 ++++++++++++++++++++++++----- 2 files changed, 75 insertions(+), 14 deletions(-) diff --git a/internal/certstream/certstream.go b/internal/certstream/certstream.go index b3b3815..01cda11 100644 --- a/internal/certstream/certstream.go +++ b/internal/certstream/certstream.go @@ -58,7 +58,7 @@ func NewCertstreamServer(cfg config.Config) (*Certstream, error) { // Initialize the stream processors if configured and enabled. for _, streamProcessor := range cfg.StreamProcessing { - if !streamProcessor.Enabled { + if !*streamProcessor.Enabled { continue } diff --git a/internal/config/streamprocessing.go b/internal/config/streamprocessing.go index 4c47ab1..7dd8478 100644 --- a/internal/config/streamprocessing.go +++ b/internal/config/streamprocessing.go @@ -5,6 +5,8 @@ import ( "net" ) +// StreamProcessorType represents the type of stream processing tool to use. +// Supported types are "kafka" and "nqs". type StreamProcessorType string const ( @@ -12,34 +14,86 @@ const ( StreamProcessorTypeNQS StreamProcessorType = "nqs" ) -type StreamProcessorCompression string +// StreamType represents the type of stream to process. +// Supported types are "full", "lite", and "domains-only". +type StreamType string const ( - StreamProcessorCompressionNone StreamProcessorCompression = "none" - StreamProcessorCompressionGzip StreamProcessorCompression = "gzip" - StreamProcessorCompressionSnappy StreamProcessorCompression = "snappy" + StreamTypeFull StreamType = "full" + StreamTypeLite StreamType = "lite" + StreamTypeDomainsOnly StreamType = "domains-only" ) -func (c StreamProcessorCompression) Valid() bool { +// Compression represents the compression type for stream processing. +type Compression string + +const ( + CompressionNone Compression = "none" + CompressionGzip Compression = "gzip" + CompressionSnappy Compression = "snappy" + CompressionZstd Compression = "zstd" + CompressionLz4 Compression = "lz4" + CompressionDeflate Compression = "deflate" +) + +func (c Compression) setDefaults() { + if c == "" { + c = CompressionNone + } +} + +// Valid returns true if the compression type is valid. +func (c Compression) Valid() bool { + c.setDefaults() + switch c { - case StreamProcessorCompressionNone, StreamProcessorCompressionGzip, StreamProcessorCompressionSnappy, "": + case CompressionNone, CompressionGzip, CompressionSnappy, CompressionZstd, CompressionLz4: return true default: return false } } +// SupportedBy returns true if the compression type is supported by the stream processing tool. +func (c Compression) SupportedBy(t StreamProcessorType) bool { + if c == CompressionNone { + return true + } + + switch t { + case StreamProcessorTypeKafka: + return c == CompressionNone || c == CompressionGzip || c == CompressionSnappy || c == CompressionZstd || c == CompressionLz4 + case StreamProcessorTypeNQS: + return c == CompressionNone || c == CompressionDeflate || c == CompressionSnappy + default: + return false + } +} + type StreamProcessor struct { - Name string `mapstructure:"name"` - Type StreamProcessorType `mapstructure:"type"` - Enabled bool `mapstructure:"enabled"` - ServerAddr string `mapstructure:"server_addr"` - ServerPort int `mapstructure:"server_port"` - Topic string `mapstructure:"topic"` - Compression StreamProcessorCompression `mapstructure:"compression"` + Name string `mapstructure:"name"` + Type StreamProcessorType `mapstructure:"type"` + Stream StreamType `mapstructure:"stream"` + Enabled *bool `mapstructure:"enabled"` + ServerAddr string `mapstructure:"server_addr"` + ServerPort int `mapstructure:"server_port"` + Topic string `mapstructure:"topic"` + Compression Compression `mapstructure:"compression"` +} + +func (s *StreamProcessor) setDefaults() { + if s.Enabled == nil { + enabled := true + s.Enabled = &enabled + } + if s.Compression == "" { + s.Compression = CompressionNone + } } func (s *StreamProcessor) Valid() bool { + s.setDefaults() + ip := net.ParseIP(s.ServerAddr) if ip == nil { log.Fatalln("Invalid IP for stream processor:", s.ServerAddr) @@ -63,10 +117,17 @@ func (s *StreamProcessor) Valid() bool { s.Topic = "certstream" } + // Check that the compression type is generally valid if !s.Compression.Valid() { log.Fatalf("Invalid compression '%s' for stream processor '%s'\n", s.Compression, s.Name) return false } + // Check that the compression type is supported by the stream processing tool + if !s.Compression.SupportedBy(s.Type) { + log.Fatalf("Compression '%s' is not supported by stream processor type '%s' for name '%s'\n", s.Compression, s.Type, s.Name) + return false + } + return true } From 4cf6d25756f3c2382d19990686082f3c16151b60 Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 02:27:12 +0200 Subject: [PATCH 25/30] feat: implement custom exponential backoff handler --- internal/backoff/backoff.go | 32 ++++++++++++++++++++++++++++++++ internal/backoff/backoff_test.go | 18 ++++++++++++++++++ 2 files changed, 50 insertions(+) create mode 100644 internal/backoff/backoff.go create mode 100644 internal/backoff/backoff_test.go diff --git a/internal/backoff/backoff.go b/internal/backoff/backoff.go new file mode 100644 index 0000000..afbcb26 --- /dev/null +++ b/internal/backoff/backoff.go @@ -0,0 +1,32 @@ +package backoff + +import ( + "sync" + "time" +) + +// NewBackoff returns a function that implements exponential backoff for the given reset window. +func NewBackoff(resetWindow time.Duration) func(target func()) { + mu := sync.Mutex{} + var errorCount int + var errorLastTime time.Time + errorLogResetWindow := resetWindow + + return func(target func()) { + mu.Lock() + defer mu.Unlock() + now := time.Now() + + // Reset counter if enough time has passed since the last error + if now.Sub(errorLastTime) > errorLogResetWindow { + errorCount = 0 + } + + errorCount++ + errorLastTime = now + + if errorCount == 1 || (errorCount&(errorCount-1)) == 0 { + target() + } + } +} diff --git a/internal/backoff/backoff_test.go b/internal/backoff/backoff_test.go new file mode 100644 index 0000000..1db85e7 --- /dev/null +++ b/internal/backoff/backoff_test.go @@ -0,0 +1,18 @@ +package backoff + +import ( + "testing" + "time" +) + +func TestErrorBackoff(t *testing.T) { + backoffHandler := NewBackoff(1 * time.Second) + + for i := range 20 { + backoffHandler(func() { + test := i + 1 + t.Log("Hello World", test) + time.Sleep(1100 * time.Millisecond) + }) + } +} From 4cbf903f4c7159efa655228f36d789731ca16e5c Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 02:28:07 +0200 Subject: [PATCH 26/30] feat(kafkaclient): improve error handling and logging --- CHANGELOG.md | 1 + internal/broadcast/kafkaclient.go | 117 +++++++++++++++++++++++------- 2 files changed, 92 insertions(+), 26 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5232c3e..db51105 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] ### Added +- Support for stream processing tools like kafka or nsq - Option to configure trusted_proxies for which the X-Forwarded-For and X-Real-IP headers are used as client IP (b062edca) - Added multiple unit tests for ensuring correct functionality - Add `excluded_logs` config option to skip specific logs and operators. Works both for classic and tiled logs. diff --git a/internal/broadcast/kafkaclient.go b/internal/broadcast/kafkaclient.go index 763462e..2c965a5 100644 --- a/internal/broadcast/kafkaclient.go +++ b/internal/broadcast/kafkaclient.go @@ -2,16 +2,19 @@ package broadcast import ( "context" + "errors" "log" + "net" "time" + "github.com/d-Rickyy-b/certstream-server-go/internal/backoff" "github.com/segmentio/kafka-go" ) -const ( - maxBatchSize = 50 - maxBatchWait = 60 * time.Second - writeWait = 60 * time.Second +var ( + kafkaMaxBatchSize = 100 + kafkaMaxBatchWait = 1 * time.Second + kafkaConnTimeout = 5 * time.Second ) // KafkaClient connects to a Kafka server in order to provide it with certificates. @@ -27,10 +30,13 @@ type KafkaClient struct { // NewKafkaClient creates a new Kafka client that immediately connects to the configured Kafka server. func NewKafkaClient(subType SubscriptionType, addr, name, topic, compression string, certBufferSize int) *KafkaClient { // Connect to the Kafka server - conn, err := kafka.DialLeader(context.Background(), "tcp", addr, topic, 0) + ctx, cancel := context.WithTimeout(context.Background(), kafkaConnTimeout) + defer cancel() + conn, err := kafka.DialLeader(ctx, "tcp", addr, topic, 0) if err != nil { log.Println("failed to connect to kafka:", err) } + // TODO implement explicit topic creation kc := &KafkaClient{ conn: conn, @@ -44,18 +50,21 @@ func NewKafkaClient(subType SubscriptionType, addr, name, topic, compression str }, } + var kafkaCompression kafka.Compression switch compression { case "gzip": - kc.compression = kafka.Gzip + kafkaCompression = kafka.Gzip case "snappy": - kc.compression = kafka.Snappy + kafkaCompression = kafka.Snappy case "lz4": - kc.compression = kafka.Lz4 - case "none": + kafkaCompression = kafka.Lz4 + case "none", "": default: log.Println("invalid compression type:", compression) } + kc.compression = kafkaCompression + go kc.broadcastHandler() go kc.reconnectHandler() @@ -79,7 +88,9 @@ func (c *KafkaClient) reconnectHandler() { } // Attempt to connect to the Kafka server - conn, err := kafka.DialLeader(context.Background(), "tcp", c.addr, c.topic, 0) + ctx, cancel := context.WithTimeout(context.Background(), kafkaConnTimeout) + defer cancel() + conn, err := kafka.DialLeader(ctx, "tcp", c.addr, c.topic, 0) if err != nil { log.Printf("Reconnect failed: %v. Retrying in 5s...", err) time.Sleep(5 * time.Second) @@ -90,6 +101,7 @@ func (c *KafkaClient) reconnectHandler() { if c.conn != nil { _ = c.conn.Close() } + c.conn = conn c.isConnected = true log.Println("Reconnected to Kafka at", c.addr) @@ -101,23 +113,33 @@ func (c *KafkaClient) reconnectHandler() { func (c *KafkaClient) broadcastHandler() { defer func() { log.Println("Closing broadcast handler for kafka producer:", c.addr) - if err := c.conn.Close(); err != nil { - log.Println("failed to close writer:", err) + if c.conn != nil { + if err := c.conn.Close(); err != nil { + log.Println("failed to close conn:", err) + return + } } ClientHandler.UnregisterClient(c.name) }() - batch := make([]kafka.Message, 0, maxBatchSize) - t := time.NewTimer(maxBatchWait) + backoffHandler := backoff.NewBackoff(60 * time.Second) + batch := make([]kafka.Message, 0, kafkaMaxBatchSize) + t := time.NewTimer(kafkaMaxBatchWait) for { select { case <-c.stopChan: return - case message := <-c.broadcastChan: + case message, ok := <-c.broadcastChan: + if !ok { + log.Println("broadcastChan closed for kafkaClient:", c.addr) + return + } + // Drop messages if not connected if !c.isConnected { + time.Sleep(5 * time.Second) continue } @@ -125,32 +147,75 @@ func (c *KafkaClient) broadcastHandler() { batch = append(batch, msg) // Write batch if it reaches max size - if len(batch) >= maxBatchSize { - c.writeBatch(batch) + if len(batch) >= kafkaMaxBatchSize { + err := c.writeBatch(batch) + if err != nil { + // Without using a backoff strategy, the errors would massively spam the log + backoffHandler(func() { + var netErr *net.OpError + if errors.As(err, &netErr) { + c.isConnected = false + } + + log.Printf("Error writing messages to kafka: %v", err) + }) + } + batch = batch[:0] - t.Reset(maxBatchWait) + t.Reset(kafkaMaxBatchWait) } case <-t.C: + // If batch size has not reached kafkaMaxBatchSize, write the batch after kafkaMaxBatchWait if len(batch) == 0 { continue } - // Write any remaining batch after maxBatchWait - c.writeBatch(batch) + err := c.writeBatch(batch) + if err != nil { + // Without using a backoff strategy, the errors would massively spam the log + backoffHandler(func() { + var netErr *net.OpError + if errors.As(err, &netErr) { + c.isConnected = false + } + + log.Printf("Error writing messages to kafka: %v", err) + }) + } + batch = batch[:0] } } } -func (c *KafkaClient) writeBatch(batch []kafka.Message) { +// writeBatch writes a batch of messages to Kafka, handling exponential backoff and connection state +func (c *KafkaClient) writeBatch(batch []kafka.Message) error { if len(batch) == 0 { - return + return nil + } + + if c.conn == nil || !c.isConnected { + return errors.New("no connection to kafka") } - _ = c.conn.SetWriteDeadline(time.Now().Add(writeWait)) - _, err := c.conn.WriteMessages(batch...) + _ = c.conn.SetWriteDeadline(time.Now().Add(kafkaConnTimeout)) + + _, err := c.conn.WriteCompressedMessages( + c.compression.Codec(), + batch..., + ) if err != nil { - c.isConnected = false - log.Println("Failed to write messages to kafka:", err) + // Treat kafka errors specially + var kafkaErr kafka.Error + if errors.As(err, &kafkaErr) { + if errors.Is(kafkaErr, kafka.MessageSizeTooLarge) { + log.Printf("Message size is too large for kafka broker '%s' - reducing batch size to %d", c.addr, kafkaMaxBatchSize/2) + kafkaMaxBatchSize = kafkaMaxBatchSize / 2 + // TODO: currently there is no retry mechanism implemented. We should try to resend the current batch with the reduced batch size. + } + } + + return err } + return nil } From 748cf7e554edc094f45461060f0e31b3cd69a721 Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 02:28:31 +0200 Subject: [PATCH 27/30] chore: add task to run the server --- taskfile.yml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/taskfile.yml b/taskfile.yml index 3fb3fc3..9140c8f 100644 --- a/taskfile.yml +++ b/taskfile.yml @@ -15,6 +15,11 @@ tasks: cmds: - go test -v ./... + run: + desc: Run the server binary. + cmds: + - go run ./cmd/certstream-server-go + build: desc: Build the server binary for your current OS and architecture. cmds: From 8120312925ebb11cc4736f4d3ddc39267d61c523 Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 02:33:08 +0200 Subject: [PATCH 28/30] refactor: set default kafka batch size to 50 This should work well with a default kafka configuration. In the future this could become a config option. --- internal/broadcast/kafkaclient.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/internal/broadcast/kafkaclient.go b/internal/broadcast/kafkaclient.go index 2c965a5..e0e7ada 100644 --- a/internal/broadcast/kafkaclient.go +++ b/internal/broadcast/kafkaclient.go @@ -12,7 +12,7 @@ import ( ) var ( - kafkaMaxBatchSize = 100 + kafkaMaxBatchSize = 50 kafkaMaxBatchWait = 1 * time.Second kafkaConnTimeout = 5 * time.Second ) From 0ddf935ce8bc40366187d627708dfea1f5fb0afa Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 13:17:30 +0200 Subject: [PATCH 29/30] docs: reformat and update changelog for 1.10.0 --- CHANGELOG.md | 99 ++++++++++++++++++++++++++++++++++++++++++---------- 1 file changed, 81 insertions(+), 18 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index db51105..7c7ffd9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,48 +6,63 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). ## [Unreleased] + ### Added + +### Changed + +### Removed + +### Fixed + +## [1.10.0] - 2026-07-xx + +### Added + - Support for stream processing tools like kafka or nsq -- Option to configure trusted_proxies for which the X-Forwarded-For and X-Real-IP headers are used as client IP (b062edca) +- Option to configure trusted_proxies for which the X-Forwarded-For and X-Real-IP headers are used as client IP (b062edc) +- Add `excluded_logs` config option to skip specific logs and operators. Works both for classic and tiled logs. (170fbeb) - Added multiple unit tests for ensuring correct functionality -- Add `excluded_logs` config option to skip specific logs and operators. Works both for classic and tiled logs. ### Changed -- Updated weak cipher suites to stronger ones (939517cd) + +- Updated weak cipher suites to stronger ones (939517c) - Updated http client settings to prevent timeouts and other connectivity issues - Updated http server settings to allow for higher delays -- Minor code improvements and refactoring, mostly style related - Updated batch size from 100 to 256 (#97) -### Removed - ### Fixed -- Calculate correct starting tile index for static ct logs -- Use proper websocket close code (1008) instead of 1005, which wasn't sent to the client -- Respect ct_index path in config file for the `create-index` command -- Ensure all additionalLog entries are processed correctly -### Docs +- Prevent excessive tile fetching for static CT logs (#104) +- Calculate correct starting tile index for static ct logs (a2b4a9f) +- Use proper websocket close code (1008) instead of 1005, which wasn't sent to the client (76ff240) +- Respect ct_index path in config file for the `create-index` command (67838a9) +- Ensure all additionalLog entries are processed correctly (c8bd735) ## [1.9.0] - 2026-04-03 + ### Added + - Ability to store and resume processing of certs from where it left off after a restart - see sample config "recovery" (#49) - New CLI switch for creating an index file from a CT log (#49) -- Support for [Static CT](https://github.com/C2SP/C2SP/blob/main/static-ct-api.md) logs +- Support for [Static CT](https://github.com/C2SP/C2SP/blob/main/static-ct-api.md) logs - Check for retired CT logs and prevent them from being watched / stop watching them (#77) - Accept websocket connections from all origins - Option to disable the default logs provided by Google - see sample config "disable_default_logs" - Use of cobra for CLI argument parsing. New commands for displaying version and creating an index file -- Override confiuration options via environment variables (e.g. `CERTSTREAM_WEBSERVER_LISTEN_PORT=1234` to change the listen port) +- Override configuration options via environment variables (e.g. `CERTSTREAM_WEBSERVER_LISTEN_PORT=1234` to change the listen port) - New field in json `source.timestamp` representing the timestamp of when the cert was processed by the CT log (#98) ### Changed + - **Breaking:** The configuration file for the docker container is now read from the /app/config/ directory (b9e5e6) ### Removed + - Deleted non-functional Dodo log from sample config (#78) ### Fixed + - Properly remove stopped ct log workers (#74) - Added missing fields certificatePolicies and ctlPoisonByte (#85) - Prevent race condition caused by simultaneous rw access to logmetrics (#91) @@ -55,140 +70,188 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Properly add new metrics for all newly found logs (#96) ## [1.8.2] - 2025-11-22 + ### Fixed + - Added missing fields certificatePolicies and ctlPoisonByte (#85) ## [1.8.1] - 2025-05-04 + ### Fixed + - No longer reject URLs with trailing slashes defined in the `additional_logs` config (#62) - When using `drop_old_logs` in the config, the server won't remove logs defined in `additional_logs` anymore (#64) ## [1.8.0] - 2025-05-03 + ### Security + - Close several CVEs in x/crypto and x/net dependencies (#59) ### Added + - New CLI tool for fetching certificates from a CT log (#47) - Ability to add custom CT logs to the config (#56) - Remove old CT logs as soon as they are removed from the Google CT Loglist (#60) - New configuration for buffer sizes (#58) ### Fixed + - Properly handle IPv6 addresses in config (#61) ## [1.7.1] - 2025-05-03 + ### Fixed + - Properly handle IPv6 addresses in config (#61) ## [1.7.0] - 2024-08-20 + ### Added + - Support for websocket compression - disabled by default (#40) - Support for non-browsers by implementing server initiated heartbeats (#39) - Start new ct-watchers as new ct logs become available (#42) - More logging to document currently watched logs (03d878e) - + ### Changed + - Changed log output to be better grepable (5c055cc) - Update ct log update interval to once per hour instead of once per 6 hours as previously (9b6e77d) ### Fixed + - Fixed a possible race condition when accessing metrics ## [1.6.0] - 2024-03-05 + ### Added + - New metric for skipped certs per client (#34) ## [1.5.2] - 2024-02-17 + ### Fixed + - Fixed an issue with ip whitelists for the websocket server (#33) ## [1.5.1] - 2024-01-18 + ### Fixed -- Fixed a rare issue where it was possible for the all_domains json property (or data property in case of the domains-only endpoint) to be null + +- Fixed a rare issue where it was possible for the all_domains json property (or data property in case of the domains-only endpoint) to be null ## [1.5.0] - 2023-12-21 + ### Added + - New `-version` switch to print version and exit afterwards - Print version on every run of the tool - Count and log number of skipped certificates per client ### Changed + - Update to chi/v5 - Update ct-watcher timeout from 5 to 30 seconds ### Fixed + - Prevent invalid subscription types to be used - Kill connection after broadcasthandler was stopped ## [1.4.0] - 2023-11-29 + ### Added + - Config option to use X-Forwarded-For or X-Real-IP header as client IP - Config option to whitelist client IPs for both websocket and metrics endpoints - Config option to enable system metrics (cpu, memory, etc.) ## [1.3.2] - 2023-11-28 + ### Fixed + - Memory leak related to clients disconnecting from the websocket not being handled properly ## [1.3.1] - 2023-09-18 + ### Changed + - Updated config.sample.yaml to run both certstream and prometheus metrics on same socket ### Docs + - Fixed wrong docker command in readme ## [1.3.0] - 2023-04-11 + ### Added + - Calculate and display Sha256 sum of certificate ### Changed + - Update dependencies - Better logging for CT log errors ### Fixed + - End execution after all workers stopped - Implement timeout for the http client - Keep ct watcher from crashing upon a connection reset from server ## [1.2.2] - 2023-01-10 + ### Added + - Two docker-compose files - Check for presence of .yml or .yaml files in the current directory ### Fixed + - Handle sudden disconnects of CT logs ### Docs -- Added [wiki entry for docker-compose](https://github.com/d-Rickyy-b/certstream-server-go/wiki/Collecting-and-Visualizing-Metrics) + +- Added [wiki entry for docker-compose](https://github.com/d-Rickyy-b/certstream-server-go/wiki/Collecting-and-Visualizing-Metrics) ## [1.2.1] - 2022-12-16 + ### Changed + - Updated ci pipeline to use new setup-go and checkout actions - Use correct package name `github.com/d-Rickyy-b/certstream-server-go` ## [1.2.0] - 2022-12-15 + ### Added + - Log x-Forwarded-For header for requests - More logging for certain error situations - Add operator to ct log cert count metrics ### Changed + - Updated certificate-transparency-go dependency to v1.1.4 - Code improvements, adhering to styleguide - Rename module to certstream-server-go - Use log_list.json instead of all_logs_list.json ## [1.1.0] - 2022-10-19 + Fix for missing loglist urls. ### Fixed + Fixed the connection issue due to the offline Google loglist urls. ## [1.0.0] - 2022-08-08 + Initial release! First stable version of certstream-server-go is published as v1.0.0 -[unreleased]: https://github.com/d-Rickyy-b/certstream-server-go/compare/v1.9.0...HEAD -[1.8.2]: https://github.com/d-Rickyy-b/certstream-server-go/compare/v1.8.2...v1.9.0 +[unreleased]: https://github.com/d-Rickyy-b/certstream-server-go/compare/v1.10.0...HEAD +[1.10.0]: https://github.com/d-Rickyy-b/certstream-server-go/compare/v1.9.0...v1.10.0 +[1.9.0]: https://github.com/d-Rickyy-b/certstream-server-go/compare/v1.8.2...v1.9.0 [1.8.2]: https://github.com/d-Rickyy-b/certstream-server-go/compare/v1.8.1...v1.8.2 [1.8.1]: https://github.com/d-Rickyy-b/certstream-server-go/compare/v1.8.0...v1.8.1 [1.8.0]: https://github.com/d-Rickyy-b/certstream-server-go/compare/v1.7.1...v1.8.0 From b5b5d11fd439badf54defa229673de811a93e8a3 Mon Sep 17 00:00:00 2001 From: Rico Date: Mon, 13 Jul 2026 13:40:39 +0200 Subject: [PATCH 30/30] refactor: rename broadcastmanager to dispatcher --- internal/broadcast/{broadcastmanager.go => dispatcher.go} | 0 1 file changed, 0 insertions(+), 0 deletions(-) rename internal/broadcast/{broadcastmanager.go => dispatcher.go} (100%) diff --git a/internal/broadcast/broadcastmanager.go b/internal/broadcast/dispatcher.go similarity index 100% rename from internal/broadcast/broadcastmanager.go rename to internal/broadcast/dispatcher.go