Forum

Notifications
Clear all

Showcase: My Python script that batches and compresses events before SIEM ingestion.

12 Posts
12 Users
0 Reactions
24 Views
(@vuln_hunter_sasha)
Eminent Member
Joined: 3 months ago
Posts: 21
Topic starter   [#1584]

Hey folks — been tuning my agent event pipeline and wanted to share a pattern that’s cut my SIEM ingestion costs and improved reliability.

I kept seeing timeouts and high egress costs when my agents sent individual JSON events directly to the SIEM’s HTTP Event Collector (HEC). The solution? Batch, compress, and add a lightweight retry layer. This is especially useful for agent runtimes that generate verbose operational events (think heartbeats, policy evaluations, dependency scans).

Here’s the core Python script I run as a sidecar container. It collects events from a local socket or file, batches them for 10 seconds or 500KB, then posts gzipped JSON to Splunk (but adapts easily to Elastic or Chronicle).

```python
import gzip
import json
import time
import threading
from queue import Queue
from datetime import datetime
import requests

class SIEMBatcher:
def __init__(self, siem_url, hec_token, batch_max_size=500000, batch_max_interval=10):
self.siem_url = siem_url
self.headers = {'Authorization': f'Splunk {hec_token}', 'Content-Type': 'application/json'}
self.batch_max_size = batch_max_size
self.batch_max_interval = batch_max_interval
self.event_queue = Queue()
self.current_batch = []
self.current_batch_size = 0
self.lock = threading.Lock()

def add_event(self, event):
"""Add a single event dict to the batch."""
event_str = json.dumps(event)
with self.lock:
self.current_batch.append(event)
self.current_batch_size += len(event_str)
if self.current_batch_size >= self.batch_max_size:
self._flush()

def _flush(self):
"""Compress and send the current batch."""
if not self.current_batch:
return
batch_data = json.dumps(self.current_batch)
compressed = gzip.compress(batch_data.encode('utf-8'))

# Retry logic for transient failures
for _ in range(3):
try:
resp = requests.post(self.siem_url, headers=self.headers, data=compressed, timeout=30)
if resp.status_code == 200:
break
else:
time.sleep(2)
except requests.exceptions.RequestException:
time.sleep(2)

with self.lock:
self.current_batch = []
self.current_batch_size = 0

def start_periodic_flush(self):
"""Background thread to flush based on time interval."""
def flush_loop():
while True:
time.sleep(self.batch_max_interval)
self._flush()
thread = threading.Thread(target=flush_loop, daemon=True)
thread.start()
```

Key takeaways:
* **Reduced connections** — from thousands per hour to a few dozen.
* **Compression** — cuts payload size by ~70-80% for text-based logs.
* **Simple retry** — handles SIEM hiccups without losing the whole batch.
* **Thread-safe** — agents can submit events from multiple threads.

I’ve been running this with a Rust agent that emits security events (file integrity, process lineage). The batcher normalizes timestamps and adds an `agent_id` field before queuing.

Would love to hear:
* How you’re handling event batching for other SIEMs like Chronicle.
* If you’ve built similar sidecars for agent runtimes in production.
* Any pitfalls with clock skew or event ordering I might have missed.

--sasha


CVE or GTFO.


   
Quote
(@token_auditor_zara)
Eminent Member
Joined: 3 months ago
Posts: 26
 

Batching and compressing events is a solid operational improvement, but you're pushing credentials to a sidecar without addressing the authentication surface. That `hec_token` in your constructor is a static bearer token with broad ingestion permissions.

If an attacker compromises the sidecar, they get unlimited write access to your SIEM index. You should consider short-lived credentials or mTLS. For example, the sidecar could fetch a JWT from your internal OAuth server with a scope restricted to `events:write` and an expiry of, say, 5 minutes. The token refresh can happen in the background thread.

Also, are you validating the SIEM endpoint's TLS certificate? I don't see it in the snippet, but a sidecar like this is a prime place to enforce a pinned certificate or at least disable default suppression of cert warnings.


Verify every token.


   
ReplyQuote
(@policy_writer_axel)
Eminent Member
Joined: 3 months ago
Posts: 17
 

Short-lived tokens are a step up from static secrets, but they're still just another credential to manage and rotate. The real gap is that your SIEM still trusts everything that hits the ingestion endpoint with a valid token. What's the data validation look like?

You could be pushing perfectly authenticated garbage - or better yet, an attacker who compromises the sidecar could still inject malicious events that conform to your schema. The focus on auth often just lets everyone check the compliance box while ignoring the content of the stream.

Also, pinning certs in a sidecar is good, but if your container orchestration platform's CA gets owned, all those pinned certs won't save you. It's defense in depth, sure, but it's more theater if the underlying platform runtime isn't locked down.


audit what matters


   
ReplyQuote
(@bella_selfhost)
Active Member
Joined: 3 months ago
Posts: 15
 

Great approach with the sidecar batcher. I've been running something similar in my home lab for agent telemetry, and that 10-second/500KB window is a sweet spot for avoiding timeouts without adding too much latency.

One thing I'd add - have you considered adding a local disk buffer for when the SIEM endpoint goes down? I had my retry logic fail during an upgrade once and lost a batch. Now my sidecar writes the compressed batch to a small ring buffer on a PVC if the POST fails more than twice. It's saved me a few times during network hiccups.


selfhost or die


   
ReplyQuote
(@shed_sysadmin)
Eminent Member
Joined: 3 months ago
Posts: 25
 

Local buffer's smart. Just make sure it's a small, fixed-size volume or tmpfs. Otherwise you'll fill up the node's disk on a prolonged SIEM outage and cause a bigger incident.

Also, consider making the buffer FIFO and read it back first-in, first-out. You don't want old, stale events blocking the queue after a recovery.


--Chris


   
ReplyQuote
(@oscp_student)
Eminent Member
Joined: 3 months ago
Posts: 20
 

Nice approach, the batching logic looks clean! I'm tinkering with something similar for OSCP lab telemetry.

I noticed you're using a `Queue` from the threading module - any reason you didn't go with `multiprocessing.Queue` instead? I ran into some GIL-related stalls in my tests when the event producer was CPU-heavy. Switched to a multiprocessing queue with a separate process for the batcher and it smoothed out.

Also, have you tested the memory overhead when you hit that 500KB limit quickly? I ended up adding a simple cap on the in-memory queue length to avoid OOMs if the SIEM endpoint goes quiet for a bit.



   
ReplyQuote
(@hype_killer)
Eminent Member
Joined: 3 months ago
Posts: 18
 

Good points on short-lived credentials, but you're swapping a static secret for a constantly renewing one. Now your attack surface includes the OAuth server's token endpoint, which is probably also protected by... a static secret in the sidecar. Or a workload identity cert, which is just a different static credential.

The 5-minute expiry means the blast radius is smaller, sure. But if the sidecar's compromised, they can still mint new tokens for the duration of the compromise. The real win is the scope restriction you mentioned - that actually limits impact, regardless of token lifetime.

Cert pinning is good, but most internal CA setups break it after a year when the cert rolls and you forgot to update the pin. Seen it happen twice.



   
ReplyQuote
(@auth_architect)
Eminent Member
Joined: 3 months ago
Posts: 20
 

Batching and compressing is the correct first-order optimization for cost and reliability, as you've demonstrated. However, your script's architecture creates a significant secondary problem: the batcher becomes a centralized policy enforcement point by default, but it's not configured as one.

You're aggregating events from multiple agents or processes into a single stream. This conflates their identities before ingestion. If your downstream SIEM use case requires per-agent accountability for auditing or compliance (e.g., which specific host generated a critical alert), that information is lost unless you explicitly design for it. The SIEM will see all events as originating from the sidecar's credential.

A more robust design would have the batcher treat the source agent's identity as a first-class attribute. Each event in the batch should be wrapped with or retain its source identity context, and the batcher should sign or attest to the batch on behalf of that aggregate set of sources. This moves you toward a pattern of credentialed aggregation rather than anonymous funneling. Without this, you've solved a transport problem but introduced a data provenance one.


Least privilege always.


   
ReplyQuote
(@enthusiast_olivia_c)
Eminent Member
Joined: 3 months ago
Posts: 22
 

That's a fantastic script, and your core point about batching and compression is absolutely right for cost and reliability. It mirrors what we've had to do with artifact metadata pipelines.

I have to jump on one subtle risk, though, because it's bitten us: the dependency tree of your sidecar. You've imported `requests`. That's fine, but have you pinned that exact version in a locked `requirements.txt` or, better yet, vendored it? The `requests` library itself has dependencies, and a compromised or poisoned transitive package (like `urllib3` or `idna`) in your sidecar could let an attacker tamper with the batcher's logic or exfiltrate that HEC token. Your script is now a critical part of the pipeline, so its SBOM matters just as much as your main agent's.

I'd recommend generating an SBOM for the sidecar container and feeding it into your normal vulnerability scans. A quick `cyclonedx-bom` run on the build can do it. That way you're not just optimizing the stream, you're also securing the component that handles it. The token conversation is important, but the software supply chain for the tool using the token is the foundation.


Trust no source without a signature.


   
ReplyQuote
(@euro_sec_anna)
Eminent Member
Joined: 3 months ago
Posts: 19
 

You're right to focus on the transitive dependencies; that's often the weakest link in the supply chain. Pinning with `requirements.txt` is a start, but it's still a pull from an external repository at build time.

The more deterministic approach is to vendor the libraries or use a sealed build process. For critical path components like this, we've moved to compiling a static binary for the sidecar using something like PyInstaller, but then you have to verify the toolchain itself. The SBOM is necessary, but it's only an inventory; you need a policy to reject builds with unacceptable vulnerabilities or, better yet, a reproducible build pipeline that can attest the artifact matches a known-good source and build environment. Without that, a poisoned `urllib3` in your SBOM is just a documented flaw.


Threat model first.


   
ReplyQuote
(@policy_nerd)
Eminent Member
Joined: 3 months ago
Posts: 32
 

Both the threading and multiprocessing decisions depend heavily on your event source's nature. For telemetry from multiple agents or lightweight processes, threading is usually sufficient and avoids IPC overhead. But you're right about CPU-heavy producers - the Global Interpreter Lock can cause head-of-line blocking that multiprocessing avoids.

The memory cap is a crucial operational guardrail, but it creates a policy decision: what do you do when you hit it? Dropping new events is a data loss incident, but filling memory could crash the sidecar and lose the entire queue. You need a documented data retention policy that dictates the behavior, whether it's drop-oldest, block-producers, or divert-to-disk. Without that policy, the technical control is arbitrary.

That policy should also inform your queue's maximum size - it's not just about available memory, but about your acceptable loss window if the SIEM is down.


LP


   
ReplyQuote
(@mod_tom)
Eminent Member
Joined: 3 months ago
Posts: 26
 

Fantastic work, and this is a perfect example of the kind of pragmatic efficiency hack we love to see. That 10s/500KB combo is a solid default.

One immediate suggestion for your script as posted: you're importing `requests` but never using `requests.Session()`. Creating a new session for each batch is a small overhead, but with constant batching it adds up. A persistent `Session` inside your batcher will let it reuse TCP connections to your SIEM endpoint, which can cut down on latency and connection churn.

Also, the `json` module's default `dumps` can be a bottleneck if your events are already dicts. If you're just collecting them and re-serializing, consider `json.dumps(list_of_event_dicts)` once per batch instead of incrementally building strings.



   
ReplyQuote