From patchwork Thu Jul 30 19:15:56 2026 Content-Type: text/plain; charset="utf-8" MIME-Version: 1.0 Content-Transfer-Encoding: 7bit X-Patchwork-Submitter: Robin Roevens X-Patchwork-Id: 10081 Return-Path: Received: from mail01.ipfire.org (mail01.haj.ipfire.org [172.28.1.202]) (using TLSv1.3 with cipher TLS_AES_256_GCM_SHA384 (256/256 bits) key-exchange x25519) (Client CN "mail01.haj.ipfire.org", Issuer "YR2" (not verified)) by web04.haj.ipfire.org (Postfix) with ESMTPS id 4hB0Lw4ggzz3wqJ for ; Thu, 30 Jul 2026 19:56:08 +0000 (UTC) Received: from mail02.haj.ipfire.org (mail02.haj.ipfire.org [172.28.1.201]) (using TLSv1.3 with cipher TLS_AES_256_GCM_SHA384 (256/256 bits) key-exchange x25519) (Client CN "mail02.haj.ipfire.org", Issuer "YE1" (not verified)) by mail01.ipfire.org (Postfix) with ESMTPS id 4hB0Lm357Yz5gj for ; Thu, 30 Jul 2026 19:56:00 +0000 (UTC) Received: from mail02.haj.ipfire.org (localhost [IPv6:::1]) by mail02.haj.ipfire.org (Postfix) with ESMTP id 4hB0HK6llGz36WB for ; Thu, 30 Jul 2026 19:53:01 +0000 (UTC) X-Original-To: development@lists.ipfire.org Received: from mail01.ipfire.org (mail01.haj.ipfire.org [IPv6:2001:678:b28::25]) (using TLSv1.3 with cipher TLS_AES_256_GCM_SHA384 (256/256 bits) key-exchange x25519) (Client CN "mail01.haj.ipfire.org", Issuer "YR2" (not verified)) by mail02.haj.ipfire.org (Postfix) with ESMTPS id 4hB0GS3XwVz37B0 for ; Thu, 30 Jul 2026 19:52:16 +0000 (UTC) Received: from layka.disroot.org (layka.disroot.org [178.21.23.139]) (using TLSv1.3 with cipher TLS_AES_256_GCM_SHA384 (256/256 bits) key-exchange x25519) (Client did not present a certificate) by mail01.ipfire.org (Postfix) with ESMTPS id 4hB0GH0T3bz2V for ; Thu, 30 Jul 2026 19:52:07 +0000 (UTC) Authentication-Results: mail01.ipfire.org; dkim=pass header.d=disroot.org header.s=mail header.b="C/BhnjPn"; spf=pass (mail01.ipfire.org: domain of robin.roevens@disroot.org designates 178.21.23.139 as permitted sender) smtp.mailfrom=robin.roevens@disroot.org; dmarc=pass (policy=reject) header.from=disroot.org ARC-Seal: i=1; a=rsa-sha256; d=lists.ipfire.org; s=202003rsa; cv=none; t=1785441127; b=dHIm5+hjNrPSBu+lczlrN1nQXY94hqb+VGBLODJydNMfl9wRUd+6Bbg7aagsTYUjxtu9Fr q8A06b7hRruThH3gspEES5FX7Qn/PLrvKgJD5Y/QlKWVE61aAR2i6pP+wQ+mE8SnhaMcmV 7fm14b+e7OY0m6jbGuwewZaENHT8aRQ1YzB1TrKs8Fs7xWL+oBKTVSWOha7Q8qB848USpP byhnJ+86gdBkAqiRHGs6ga3mCQ4Cm6kleIVv+NxZtNPrvZdwTA24NsMsqyaekgR0kGJoDV GQkrz3WaOAQGjrLT3E53PoDTsqhgQScALtQPRMc7R7BBQ1qhm4bqifLgu9Xxsw== ARC-Message-Signature: i=1; a=rsa-sha256; c=relaxed/relaxed; d=lists.ipfire.org; s=202003rsa; t=1785441127; h=from:from:reply-to:subject:subject:date:date:message-id:message-id: to:to:cc:cc:mime-version:mime-version: content-transfer-encoding:content-transfer-encoding: in-reply-to:in-reply-to:references:references:dkim-signature; bh=Lvv+RJk/W5YzW3xktjsJsIOwHRtJ3GhQum9nGBke6cQ=; b=IUmaM/SgYarFAsO11WD7CFsrngwAShaVuxz4h8CcTcIxRRS7HOz75m3Lhy2RhoO3+6fOBo cOlBtxjV/gEuf8H9gvA56R/VY3jN4hmHC6Ez9DAb86uwCuTZtn8iIq4wCBCerWuQBXc6pJ ekoIKsvk8wzh4lkBBqQVtosF5Qgyiq02QCufNlgLsAgipljPK7+lYNLyuKmYrd1Wm/AYUl I4ls54lF1060rJZ+M7icGcDsD1r//5JzwWXf5ltyEctT9228ryl1Lqr/LQeQzCAuPXe5bw lb6MIgigH5ihgA3T9tn8eEMaa6IAzBBecppeKif2nqWysUCbmD9FPrgQJZgxyg== ARC-Authentication-Results: i=1; mail01.ipfire.org; dkim=pass header.d=disroot.org header.s=mail header.b="C/BhnjPn"; spf=pass (mail01.ipfire.org: domain of robin.roevens@disroot.org designates 178.21.23.139 as permitted sender) smtp.mailfrom=robin.roevens@disroot.org; dmarc=pass (policy=reject) header.from=disroot.org Received: from mail01.layka.lan (localhost [127.0.0.1]) by disroot.org (Postfix) with ESMTP id 2414A84BE3 for ; Thu, 30 Jul 2026 21:52:01 +0200 (CEST) X-Virus-Scanned: SPAM Filter at disroot.org Received: from layka.disroot.org ([127.0.0.1]) by localhost (disroot.org [127.0.0.1]) (amavis, port 10024) with ESMTP id rMU8qh7r5RZc for ; Thu, 30 Jul 2026 21:51:59 +0200 (CEST) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/simple; d=disroot.org; s=mail; t=1785441119; bh=NURcr4XnAFuUDfdeapqS8IpXOdi9JFq/WhhmERhP0DY=; h=From:To:Cc:Subject:Date:In-Reply-To:References; b=C/BhnjPnX5VpjT8UddvYFbDAkEiM4XJ8dhksu4veadcz8E2pEdTsRmq6vFTPMHT+l 6yWOYyG6Gp7vecD6Qu0giuMjOtzz+uWXglNVxjMlQWlzw2/gxIYOlTePItkb5RjwU6 fwbVLa/Om1vDFYSZGNzOD1xZR7xWaCxIUUd5XmN/hgu0rycGWDXhkgg05Ao7SlTrrS nDjt83ySjsn3B6252LwuJv0EOUuUbgOawEtPZVKtE+1KAdxd266vtmPtEg+ihfr70Q 9Fo2nfIMq9bcnILoIDarmVhj4fmg6TNcfSWU6h+J2+6aG7ZhfMArdMqdxbzb1PQLQR XDw+mdluAW3Qg== Received: from chojin.roevenslambrechts.be (chojin.roevenslambrechts.be [192.168.0.50]) (using TLSv1.3 with cipher TLS_AES_256_GCM_SHA384 (256/256 bits)) (no client certificate requested) (Authenticated sender) by hachiman (MailScanner Milter) with SMTP id 50C23586005; Thu, 30 Jul 2026 21:51:56 +0200 (CEST) From: Robin Roevens To: development@lists.ipfire.org Cc: Robin Roevens Subject: [PATCH 5/5] Add background task to send all pending alerts to Zabbix every 1 second Date: Thu, 30 Jul 2026 21:15:56 +0200 Message-ID: <20260730195148.3278295-6-robin.roevens@disroot.org> In-Reply-To: <20260730195148.3278295-1-robin.roevens@disroot.org> References: <20260730195148.3278295-1-robin.roevens@disroot.org> Precedence: list List-Id: List-Subscribe: , List-Unsubscribe: , List-Post: List-Help: Sender: Mail-Followup-To: MIME-Version: 1.0 X-RoevensLambrechts-MailScanner-ID: 50C23586005.AD543 X-RoevensLambrechts-MailScanner: Found to be clean X-RoevensLambrechts-MailScanner-From: robin.roevens@disroot.org X-RoevensLambrechts-MailScanner-Watermark: 1786045917.0423@2NZt8INRh+ADTxeDe/V5PA X-Rspamd-Server: mail01.haj.ipfire.org X-Rspamd-Queue-Id: 4hB0GH0T3bz2V X-Rspamd-Action: no action X-Spamd-Result: default: False [-5.63 / 11.00]; BAYES_HAM(-2.99)[99.97%]; R_DKIM_ALLOW(-1.65)[disroot.org:s=mail]; MID_CONTAINS_FROM(1.00)[]; DKIM_REPUTATION(-0.92)[-0.92153661380942]; SPF_REPUTATION_HAM(-0.66)[-0.65685405590104]; DMARC_POLICY_ALLOW(-0.50)[disroot.org,reject]; R_MISSING_CHARSET(0.50)[]; R_SPF_ALLOW(-0.20)[+a:c]; MIME_GOOD(-0.10)[text/plain]; MX_GOOD(-0.10)[disroot.org]; ASN(0.00)[asn:50673, ipnet:178.21.23.0/24, country:NL]; ARC_SIGNED(0.00)[lists.ipfire.org:s=202003rsa:i=1]; TO_DN_SOME(0.00)[]; ARC_NA(0.00)[]; MIME_TRACE(0.00)[0:+]; MISSING_XM_UA(0.00)[]; RCPT_COUNT_TWO(0.00)[2]; RCVD_COUNT_THREE(0.00)[3]; FROM_EQ_ENVFROM(0.00)[]; FROM_HAS_DN(0.00)[]; IP_REPUTATION_HAM(0.00)[asn: 50673(0.00), country: NL(-0.01), ip: 178.21.23.139(0.00)]; TO_MATCH_ENVRCPT_SOME(0.00)[]; RCVD_TLS_LAST(0.00)[]; PREVIOUSLY_DELIVERED(0.00)[development@lists.ipfire.org]; DKIM_TRACE(0.00)[disroot.org:+] Background task that will send pending events every 1 second. On failure it will retry 3 times and then pauses until next alert comes in to prevent hammering a possibly already overloaded Zabbix server.. Signed-off-by: Robin Roevens --- src/suricata-reporter.in | 109 +++++++++++++++++++++++++++++++++------ 1 file changed, 94 insertions(+), 15 deletions(-) diff --git a/src/suricata-reporter.in b/src/suricata-reporter.in index 47482ab..a95f6b7 100644 --- a/src/suricata-reporter.in +++ b/src/suricata-reporter.in @@ -82,10 +82,6 @@ class Reporter(object): # Remember the last time the database was cleaned self.last_cleanup_at = None - # Initialize Zabbix sender - self.zabbix_sender = None - self.init_zabbix_sender() - # Register any signals for signo in (signal.SIGINT, signal.SIGTERM): self.loop.add_signal_handler(signo, self.terminate) @@ -99,6 +95,11 @@ class Reporter(object): # Create the socket self.sock = self._create_socket() + # Initialize Zabbix sender + self.zabbix_sender = None + self.init_zabbix_sender() + self.zabbix_sender_task = None + def read_config(self): """ Reads or re-reads the configuration. @@ -123,6 +124,12 @@ class Reporter(object): zabbix_server_host = self.config.get('zabbix', 'zabbix_server_host', fallback='') zabbix_server_port = self.config.getint('zabbix', 'zabbix_server_port', fallback=10051) + # Zabbix sender backoff state + self.zabbix_consecutive_failures = 0 + self.zabbix_pause_until_new_event = False + self.zabbix_sender_wakeup = asyncio.Event() + self.zabbix_pause_after_failures = 3 + if zabbix_config: if not os.path.isfile(zabbix_config): log.error(f"Zabbix agent config file {zabbix_config} does not exist.") @@ -242,6 +249,67 @@ class Reporter(object): # Return the socket return sock + def _start_zabbix_sender_task(self): + """ + Start a background Zabbix sender task + """ + if not self.config.getboolean("zabbix", "enabled", fallback=False): + return + + self.zabbix_sender_task = asyncio.create_task(self._periodic_zabbix_sender_flush()) + + async def _stop_zabbix_sender_task(self): + """ + Stops the background Zabbix sender task + """ + if not self.zabbix_sender_task: + return + + self.zabbix_sender_task.cancel() + try: + await self.zabbix_sender_task + except asyncio.CancelledError: + pass + + # Final flush of any remaining alerts + await self.flush_pending_to_zabbix() + + async def _periodic_zabbix_sender_flush(self): + """ + Background task that periodically sends all pending Zabbix alerts. + Runs every 1 second to batch-send alerts in bulk on heavy load, but + pauses after repeated failures until a new event arrives. + """ + log.debug("Starting periodic Zabbix sender task") + + try: + while not self.is_terminated.is_set(): + if self.zabbix_pause_until_new_event: + log.warning("Zabbix sending is paused after repeated failures; waiting for the next event") + self.zabbix_sender_wakeup.clear() + await self.zabbix_sender_wakeup.wait() + self.zabbix_sender_wakeup.clear() + self.zabbix_pause_until_new_event = False + self.zabbix_consecutive_failures = 0 + continue + + if self.is_terminated.is_set(): + break + + # Send pending alerts to Zabbix + if await self.flush_pending_to_zabbix(): + self.zabbix_consecutive_failures = 0 + else: + self.zabbix_consecutive_failures += 1 + if self.zabbix_consecutive_failures >= self.zabbix_pause_after_failures: + self.zabbix_pause_until_new_event = True + + await asyncio.sleep(1) + + except asyncio.CancelledError: + log.debug("Periodic Zabbix sender task cancelled") + raise + async def run(self): """ The main loop of the application. @@ -251,22 +319,29 @@ class Reporter(object): # Cleanup the database at startup self.cleanup() - # Wait until we have terminated - await self.is_terminated.wait() + # Start the periodic Zabbix sender task + self._start_zabbix_sender_task() - # Remove the socket so we won't receive any more data try: - os.unlink(self.socket_path) - except OSError as e: - log.error("Failed to remove %s: %s" % (self.socket_path, e)) + # Wait until we have terminated + await self.is_terminated.wait() + finally: + # Remove the socket so we won't receive any more data + try: + os.unlink(self.socket_path) + except OSError as e: + log.error("Failed to remove %s: %s" % (self.socket_path, e)) + + # Cancel the periodic Zabbix sender task + await self._stop_zabbix_sender_task() - # We will optimize the database before we exit - self.optimize() + # We will optimize the database before we exit + self.optimize() - # Close the database - self.db.close() + # Close the database + self.db.close() - log.debug("Reporter has exited") + log.debug("Reporter has exited") def terminate(self): """ @@ -485,6 +560,10 @@ class Reporter(object): # Store the alert self.store(event) + # Wake the periodic flush task so it can retry immediately on new input + if self.config.getboolean("zabbix", "enabled", fallback=False): + self.zabbix_sender_wakeup.set() + # Send to syslog if self.config.getboolean("syslog", "enabled", fallback=False): await self.send_to_syslog(event)