[5/5] Add background task to send all pending alerts to Zabbix every 1 second

Message ID 20260730195148.3278295-6-robin.roevens@disroot.org
State New
Headers
Series Add Zabbix functionality to suricata-reporter |

Commit Message

Robin Roevens 30 Jul 2026, 7:15 p.m. UTC
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

Michael Tremer 31 Jul 2026, 10:24 a.m. UTC | #1
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.
> 
>
  

Patch

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)