feature: elastic integration

This commit is contained in:
Dobin Rutishauser
2026-01-01 13:30:33 +01:00
parent 91ae32105c
commit 80c6a9ecef
5 changed files with 2877 additions and 0 deletions
@@ -0,0 +1,92 @@
import os
import logging
from datetime import datetime, timedelta
from typing import Dict, List, Optional, Set, Tuple
import requests
class ElasticCloudClient:
def __init__(self, base_url: str, api_key: str):
self.base_url = base_url
self.api_key = api_key
def fetch_alerts(
self,
hostname: Optional[str],
start_time: datetime,
end_time: datetime,
) -> List[dict]:
# Build the search query for Elastic SIEM signals
# curl -k -H "Authorization: ApiKey bbbb==" \
# -X GET "https://10.10.20.20:9200/.siem-signals-*/_search" \
# -H "Content-Type: application/json" \
# -d '{
# "size": 10,
# "query": {
# "bool": {
# "must": [
# { "range": { "@timestamp": { "gte": "2025-12-31T23:00:00.000Z", "lte": "2026-01-01T22:59:59.999Z" } } },
# { "term": { "host.name": "desktop-h79u9ft" } }
# ]
# }
# }
# }'
query_body = {
"size": 32,
"query": {
"bool": {
"must": [ {
"range": {
"@timestamp": {
"gte": start_time.strftime("%Y-%m-%dT%H:%M:%S.%f")[:-3] + "Z",
"lte": end_time.strftime("%Y-%m-%dT%H:%M:%S.%f")[:-3] + "Z",
}
}
}, {
"term": {
"host.name": hostname,
}
}]
}
}
}
# Make request to .siem-signals-* index
headers = {
"Authorization": f"ApiKey {self.api_key}",
"Content-Type": "application/json"
}
url = f"{self.base_url.rstrip('/')}/.siem-signals-*/_search"
response = requests.get(url, headers=headers, json=query_body, verify=False, timeout=15)
if response.status_code >= 400:
raise RuntimeError(f"Elastic API GET {url} failed: {response.status_code} {response.text}")
return response.json().get("hits", {}).get("hits", [])
if __name__ == "__main__":
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
# Example usage
elastic_client = ElasticCloudClient(
base_url="https://your-elastic-cloud-instance:9200",
api_key="your_api_key_here"
)
start = datetime.utcnow() - timedelta(hours=1)
end = datetime.utcnow()
alerts = elastic_client.fetch_alerts(
hostname="example-hostname",
start_time=start,
end_time=end
)
for alert in alerts:
logger.info(alert)
@@ -0,0 +1,179 @@
import logging
import threading
import pprint
import time
from datetime import datetime, timedelta
from typing import Dict, Optional, Tuple, List
from sqlalchemy.orm import Session, joinedload
import json
from detonatorapi.database import get_db_direct, Submission, SubmissionAlert, Profile
from detonatorapi.db_interface import db_submission_add_log, db_submission_change_status_quick
from detonatorapi.edr_cloud.elastic_cloud_client import ElasticCloudClient
from .edr_cloud import EdrCloud
logger = logging.getLogger(__name__)
POLLING_TIME_MINUTES = 2 # post end monitoring duration
POLL_INTERVAL_SECONDS = 10 # polling interval
class CloudElasticPlugin(EdrCloud):
def __init__(self):
super().__init__()
self.elasticClient: ElasticCloudClient | None = None
self.thread: Optional[threading.Thread] = None
@staticmethod
def is_relevant(profile_data: dict) -> bool:
edr_info = profile_data.get("edr_elastic", None)
return edr_info is not None
def start_monitoring_thread(self, submission_id: int):
self.submission_id = submission_id
self.thread = threading.Thread(
target=self._monitor_loop,
name=f"elastic-monitor-{submission_id}",
daemon=True,
)
self.thread.start()
logger.info("Alert monitoring thread started")
def _monitor_loop(self):
while True:
db = None
try:
db = get_db_direct()
submission = db.query(Submission).filter(Submission.id == self.submission_id).first()
if not submission:
break
# Create elastic client
if self.elasticClient is None:
edr_info = submission.profile.data.get("edr_elastic", None)
if not edr_info:
raise RuntimeError(f"Profile {submission.profile.id} has no edr_elastic data")
required_keys = ("elastic_url", "elastic_apikey", "hostname")
for key in required_keys:
if key not in edr_info:
raise RuntimeError(f"Profile {submission.profile.id} edr_elastic missing key: {key}")
self.elasticClient = ElasticCloudClient(
base_url=edr_info.get("elastic_url"),
api_key=edr_info.get("elastic_apikey"),
)
# check if we are done
if submission.status in ("error", "finished"):
# check if we are > POLLING_TIME_MINUTES after completed_at
if submission.completed_at and \
submission.completed_at + timedelta(minutes=POLLING_TIME_MINUTES) < datetime.utcnow():
break
# poll
self._poll(db, submission)
db.commit()
# sleep a bit before next poll
time.sleep(POLL_INTERVAL_SECONDS)
except Exception as exc:
logger.error(f"Alert monitor loop error: {exc}")
finally:
if db:
db.close()
def _poll(self, db: Session, submission: Submission) -> bool:
# Submission Info
if not submission.profile:
return False
device_info = submission.profile.data.get("edr_elastic", None)
if not device_info:
return False
if not self.elasticClient:
return False
device_hostname = device_info.get("hostname", None)
# Determine polling window
time_from = submission.created_at
time_to = submission.completed_at or datetime.utcnow()
try:
poll_msg = f"ELASTIC poll for submission {submission.id}: from {time_from.isoformat()} to {time_to.isoformat()} "
#db_submission_add_log(db, submission, poll_msg)
logger.info(poll_msg)
elastic_alerts = self.elasticClient.fetch_alerts(
device_hostname, time_from, time_to
)
alerts = self.convert_elastic_alerts(elastic_alerts)
self._store_alerts(db, submission, alerts)
except Exception as exc:
db_submission_add_log(db, submission, f"ELASTIC poll: failed: {exc}")
return True
def convert_elastic_alerts(self, elastic_alerts: List[dict]) -> List[dict]:
converted_alerts = []
for alert in elastic_alerts:
source = alert.get("_source", {})
converted_alert = {
"alert_id": alert.get("_id"),
"title": source.get("kibana.alert.rule.name"),
"severity": source.get("kibana.alert.severity"),
"detection_source": source.get("message"),
"detected_at": source.get("@timestamp"),
"category": "", # TBD
"raw": json.dumps(alert),
"additional_data": {
"rule_id": source.get("kibana.alert.rule.rule_id"),
}
}
converted_alerts.append(converted_alert)
return converted_alerts
def _store_alerts(self, db: Session, submission: Submission, alerts: List[dict]) -> bool:
existing_ids = {alert.alert_id for alert in submission.alerts}
for alert in alerts:
alert_id = alert["alert_id"]
if alert_id in existing_ids:
continue
existing_ids.add(alert_id)
detected_at = alert.get("detected_at") # "2026-01-01T09:33:44.088Z"
detected_dt = None
if detected_at:
try:
# Handle Elastic's ISO 8601 format: remove 'Z' and parse
timestamp_str = detected_at.rstrip('Z')
detected_dt = datetime.fromisoformat(timestamp_str)
except ValueError:
logger.warning(f"submission {submission.id}: Invalid alert detected_at format: {detected_at}")
# Use current time as fallback to avoid NULL constraint violation
detected_dt = datetime.utcnow()
submission_alert = SubmissionAlert(
submission_id=submission.id,
source="Elastic Plugin",
raw=alert.get("raw", ""),
alert_id=alert_id,
title=alert.get("title"),
severity=alert.get("severity"),
category=alert.get("category"),
detection_source=alert.get("detection_source"),
detected_at=detected_dt,
additional_data=alert.get("additional_data", {}),
)
db.add(submission_alert)
submission.alerts.append(submission_alert)
logger.info(f"submission {submission.id}: New alert stored: {alert_id}")
return True
+90
View File
@@ -0,0 +1,90 @@
# Gather Elastic Defend Logs
## Install Elastic
Use this guide: https://github.com/peasead/elastic-container
## Setup Elastic for Detonator Integration
Create Role `alert_reader`:
```
curl -k -u elastic:ELASTIC_PASSWORD -X POST "https://10.10.20.20:9200/_security/role/alert_reader" -H 'Content-Type: application/json' -d '{
"indices": [
{
"names": [".siem-signals-*"],
"privileges": ["read"]
}
]
}'
```
Create user `alert_user` with role `alert_reader`:
```
curl -k -u elastic:ELASTIC_PASSWORD -X POST "https://10.10.20.20:9200/_security/user/alert_user" -H 'Content-Type: application/json' -d '{
"password" : "ALERT_USER_PASSWORD",
"roles" : ["alert_reader"],
"full_name" : "SIEM Alert Reader",
"email" : "alert_user@example.com"
}'
```
Create an API key for `alert_user`:
```
curl -k -u elastic:ELASTIC_PASSWORD -X POST "https://10.10.20.20:9200/_security/api_key" -H 'Content-Type: application/json' -d '{
"name": "alert_reader_key",
"role_descriptors": {
"alert_reader_role": {
"cluster": [],
"index": [
{
"names": [".siem-signals-*"],
"privileges": ["read"]
}
]
}
}
}'
```
Response:
```
{
"id": "ididididid",
"name": "alert_reader_key",
"api_key": "aaaa",
"encoded": "bbbb=="
}
```
Test the token:
```
curl -k -H "Authorization: ApiKey bbbb==" \
-X GET "https://10.10.20.20:9200/.siem-signals-*/_search" \
-H "Content-Type: application/json" \
-d '{
"size": 10,
"query": {
"range": { "@timestamp": { "gte": "2025-12-31T23:00:00.000Z", "lte": "2026-01-01T22:59:59.999Z" } }
}
}'
```
## Configure Detonator
```
myfirstvm:
connector: ...
...
data:
...
edr_elastic:
elastic_url: "https://10.10.20.20:9200"
elastic_apikey: "bbbb=="
hostname: "DESKTOP-12356"
```
File diff suppressed because one or more lines are too long
+26
View File
@@ -0,0 +1,26 @@
from unittest import TestCase
import json
from detonatorapi.edr_cloud.elastic_cloud_plugin import CloudElasticPlugin
class TestElastic(TestCase):
def test_parser(self):
elasticCloudPlugin = CloudElasticPlugin()
with open("tests/elastic_data.json", "r") as f:
elastic_data = json.load(f)
alerts = elasticCloudPlugin.convert_elastic_alerts(elastic_data)
self.assertIsInstance(alerts, list)
self.assertGreater(len(alerts), 0)
first_alert = alerts[0]
self.assertIsInstance(first_alert, dict)
self.assertEqual(first_alert.get("alert_id"), "dd3f01efa6c30e6e2c869414dd07b8764b09341f33d269f6c2a64c35786c216f")
self.assertEqual(first_alert.get("severity"), "medium")
self.assertEqual(first_alert.get("detection_source"), "Endpoint process event")
self.assertEqual(first_alert.get("detected_at"), "2026-01-01T09:33:44.088Z")
self.assertEqual(first_alert["additional_data"].get("rule_id"), "ebfe1448-7fac-4d59-acea-181bd89b1f7f")