From mboxrd@z Thu Jan 1 00:00:00 1970 Return-Path: Received: from mail02.haj.ipfire.org (localhost [IPv6:::1]) by mail02.haj.ipfire.org (Postfix) with ESMTP id 4hBMdb1y4lz37Ff for ; Fri, 31 Jul 2026 10:25:07 +0000 (UTC) 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 4hBMdW5d26z36V3 for ; Fri, 31 Jul 2026 10:25:03 +0000 (UTC) Received: from [127.0.0.1] (localhost [127.0.0.1]) (using TLSv1.2 with cipher ECDHE-RSA-AES256-GCM-SHA384 (256/256 bits)) (No client certificate requested) by mail01.ipfire.org (Postfix) with ESMTPSA id 4hBMdM3Vlkz40J; Fri, 31 Jul 2026 10:24:55 +0000 (UTC) DKIM-Signature: v=1; a=ed25519-sha256; c=relaxed/relaxed; d=ipfire.org; s=202003ed25519; t=1785493495; h=from:from:reply-to:subject:subject:date:date:message-id:message-id: to:to:cc:cc:mime-version:mime-version:content-type:content-type: content-transfer-encoding:content-transfer-encoding: in-reply-to:in-reply-to:references:references; bh=KMy3MVJDbxT1BcigApvEp2altszDBWweFavEfgV8eYs=; b=0snLVCbMfcOAUMclF5zHJXF+XTyd04s2ZzgN906IWWF52EfIhCqjAtaH1k/+ckLLhe9DFA 7YJ0voT+TSU12kCQ== DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=ipfire.org; s=202003rsa; t=1785493495; h=from:from:reply-to:subject:subject:date:date:message-id:message-id: to:to:cc:cc:mime-version:mime-version:content-type:content-type: content-transfer-encoding:content-transfer-encoding: in-reply-to:in-reply-to:references:references; bh=KMy3MVJDbxT1BcigApvEp2altszDBWweFavEfgV8eYs=; b=pzaRxv6488U72LivvapyaH0JXKaUGplqtHijNe77K+WXRTQM+FNzEuxAo7DxxK/m8bLENZ xPweSWgnmoaJZT7p/TrSjD0IwFiKBkx1s4qwqcjFBNqxYiMLu/ivX4iGU52GpShDLK0ur2 kij/zSMd5Nv92Iv8NpASnWFcJEPiSymyjGyZC2spO9UBCV8eKOxSslTPXCfMaWQha+D/VS 92FP2K+7fODb7AQXweLgFGkJy3pJKLfWOoyoF07KxcyjBfSJ+C4HosYtQOHrYMjo5h43CE 1zzvWFLdcgMWo8+sQ9hE/sB0I1MKf0XAZpKZsd568Mts7a0dHzGb1AgkczDiEw== Content-Type: text/plain; charset=us-ascii Precedence: list List-Id: List-Subscribe: , List-Unsubscribe: , List-Post: List-Help: Sender: Mail-Followup-To: Mime-Version: 1.0 Subject: Re: [PATCH 5/5] Add background task to send all pending alerts to Zabbix every 1 second From: Michael Tremer In-Reply-To: <20260730195148.3278295-6-robin.roevens@disroot.org> Date: Fri, 31 Jul 2026 11:24:32 +0100 Cc: development@lists.ipfire.org Content-Transfer-Encoding: quoted-printable Message-Id: <26C62116-293D-4646-98DE-B89FF7C24178@ipfire.org> References: <20260730195148.3278295-1-robin.roevens@disroot.org> <20260730195148.3278295-6-robin.roevens@disroot.org> To: Robin Roevens Hello, I like that Zabbix gets paused when we know for sure that there has not = been another event. But why the consecutive failure thing? -Michael > On 30 Jul 2026, at 20:15, Robin Roevens = wrote: >=20 > 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.. >=20 > Signed-off-by: Robin Roevens > --- > src/suricata-reporter.in | 109 +++++++++++++++++++++++++++++++++------ > 1 file changed, 94 insertions(+), 15 deletions(-) >=20 > 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 =3D None >=20 > - # Initialize Zabbix sender > - self.zabbix_sender =3D 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 =3D self._create_socket() >=20 > + # Initialize Zabbix sender > + self.zabbix_sender =3D None > + self.init_zabbix_sender() > + self.zabbix_sender_task =3D None > + > def read_config(self): > """ > Reads or re-reads the configuration. > @@ -123,6 +124,12 @@ class Reporter(object): > zabbix_server_host =3D self.config.get('zabbix', 'zabbix_server_host', = fallback=3D'') > zabbix_server_port =3D self.config.getint('zabbix', = 'zabbix_server_port', fallback=3D10051) >=20 > + # Zabbix sender backoff state > + self.zabbix_consecutive_failures =3D 0 > + self.zabbix_pause_until_new_event =3D False > + self.zabbix_sender_wakeup =3D asyncio.Event() > + self.zabbix_pause_after_failures =3D 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 >=20 > + def _start_zabbix_sender_task(self): > + """ > + Start a background Zabbix sender task > + """ > + if not self.config.getboolean("zabbix", "enabled", fallback=3DFalse): > + return > + > + self.zabbix_sender_task =3D = 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 > +=20 > + self.zabbix_sender_task.cancel() > + try: > + await self.zabbix_sender_task > + except asyncio.CancelledError: > + pass > +=20 > + # 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=20= > + pauses after repeated failures until a new event arrives. > + """ > + log.debug("Starting periodic Zabbix sender task") > +=20 > + 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 =3D False > + self.zabbix_consecutive_failures =3D 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 =3D 0 > + else: > + self.zabbix_consecutive_failures +=3D 1 > + if self.zabbix_consecutive_failures >=3D = self.zabbix_pause_after_failures: > + self.zabbix_pause_until_new_event =3D True > + > + await asyncio.sleep(1) > +=20 > + 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() >=20 > - # Wait until we have terminated > - await self.is_terminated.wait() > + # Start the periodic Zabbix sender task > + self._start_zabbix_sender_task() >=20 > - # 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() >=20 > - # We will optimize the database before we exit > - self.optimize() > + # We will optimize the database before we exit > + self.optimize() >=20 > - # Close the database > - self.db.close() > + # Close the database > + self.db.close() >=20 > - log.debug("Reporter has exited") > + log.debug("Reporter has exited") >=20 > def terminate(self): > """ > @@ -485,6 +560,10 @@ class Reporter(object): > # Store the alert > self.store(event) >=20 > + # Wake the periodic flush task so it can retry immediately on new = input > + if self.config.getboolean("zabbix", "enabled", fallback=3DFalse): > + self.zabbix_sender_wakeup.set() > + > # Send to syslog > if self.config.getboolean("syslog", "enabled", fallback=3DFalse): > await self.send_to_syslog(event) > --=20 > 2.54.0 >=20 >=20 > --=20 > Dit bericht is gescanned op virussen en andere gevaarlijke > inhoud door MailScanner en lijkt schoon te zijn. >=20 >=20