[5/5] Add background task to send all pending alerts to Zabbix every 1 second
Commit Message
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 <robin.roevens@disroot.org>
---
src/suricata-reporter.in | 109 +++++++++++++++++++++++++++++++++------
1 file changed, 94 insertions(+), 15 deletions(-)
Comments
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 <robin.roevens@disroot.org> wrote:
>
> 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 <robin.roevens@disroot.org>
> ---
> 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)
> --
> 2.54.0
>
>
> --
> Dit bericht is gescanned op virussen en andere gevaarlijke
> inhoud door MailScanner en lijkt schoon te zijn.
>
>
@@ -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)