refactor(detmir): retire python runtime paths

This commit is contained in:
igor04091968
2026-06-03 03:27:52 +03:00
parent 109c31f291
commit dbef90a09e
103 changed files with 914 additions and 19800 deletions
@@ -6,8 +6,8 @@ After=network-online.target
Type=oneshot
User=activitywatch
Group=activitywatch
WorkingDirectory=/opt/activitywatch/dlp-integrations
ExecStart=/opt/activitywatch/dlp-integrations/.venv/bin/python /opt/activitywatch/dlp-integrations/cef_exporter.py
WorkingDirectory=/var/lib/activitywatch
ExecStart=/usr/local/bin/dlp-cef-exporter-rust
[Install]
WantedBy=multi-user.target
-165
View File
@@ -1,165 +0,0 @@
#!/usr/bin/env python3
from __future__ import annotations
import json
import logging
import socket
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from urllib import error, request
import yaml
LOG = logging.getLogger("aw.dlp.cef_exporter")
def setup_logging() -> None:
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
def load_yaml(path: Path) -> dict[str, Any]:
if not path.exists():
return {}
data = yaml.safe_load(path.read_text(encoding="utf-8")) or {}
if not isinstance(data, dict):
return {}
return data
def load_json(path: Path) -> dict[str, Any]:
if not path.exists():
return {}
try:
data = json.loads(path.read_text(encoding="utf-8"))
if isinstance(data, dict):
return data
except Exception:
return {}
return {}
def save_json(path: Path, payload: dict[str, Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
def http_json(url: str, timeout: int = 15) -> Any:
req = request.Request(url, method="GET")
with request.urlopen(req, timeout=timeout) as resp:
return json.loads(resp.read().decode("utf-8", errors="ignore"))
def escape_cef(v: Any) -> str:
s = "" if v is None else str(v)
return s.replace("\\", "\\\\").replace("|", "\\|").replace("=", "\\=").replace("\n", "\\n").replace("\r", "")
def map_severity(name: str, mapping: dict[str, int]) -> int:
return int(mapping.get((name or "").lower(), 3))
def build_cef(event: dict[str, Any], mapping: dict[str, int]) -> str:
data = event.get("data") or {}
sev_name = str(data.get("severity") or "low").lower()
sev_num = map_severity(sev_name, mapping)
rt = event.get("timestamp") or datetime.now(timezone.utc).isoformat()
rule = data.get("ruleId") or "dlp-incident"
msg = data.get("message") or "AWatch DLP incident"
sig = data.get("signalType") or "unknown"
host = data.get("hostname") or "unknown"
user = data.get("username") or "unknown"
action = data.get("action") or "alert"
ext = (
f"rt={escape_cef(rt)} "
f"shost={escape_cef(host)} "
f"suser={escape_cef(user)} "
f"cs1Label=signalType cs1={escape_cef(sig)} "
f"cs2Label=action cs2={escape_cef(action)} "
f"cs3Label=ruleId cs3={escape_cef(rule)}"
)
return (
f"CEF:0|AWatch-rus|DLP|1.0|{escape_cef(rule)}|{escape_cef(msg)}|{sev_num}|{ext}"
)
def send_syslog_udp(line: str, host: str, port: int) -> None:
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
try:
sock.sendto(line.encode("utf-8", errors="ignore"), (host, port))
finally:
sock.close()
def send_syslog_tcp(line: str, host: str, port: int, timeout: int = 10) -> None:
sock = socket.create_connection((host, port), timeout=timeout)
try:
sock.sendall((line + "\n").encode("utf-8", errors="ignore"))
finally:
sock.close()
def iter_new_incidents(
aw_base: str,
state: dict[str, Any],
per_bucket_limit: int,
) -> tuple[list[dict[str, Any]], dict[str, int]]:
buckets = http_json(f"{aw_base}/buckets/")
bucket_ids = sorted([bid for bid in buckets.keys() if bid.startswith("aw-dlp-incidents_")])
last_ids = state.get("last_ids", {})
if not isinstance(last_ids, dict):
last_ids = {}
max_ids: dict[str, int] = {}
out: list[dict[str, Any]] = []
for bid in bucket_ids:
try:
events = http_json(f"{aw_base}/buckets/{bid}/events?limit={int(per_bucket_limit)}")
except error.HTTPError as exc:
LOG.warning("skip bucket %s: %s", bid, exc)
continue
prev = int(last_ids.get(bid, 0))
bucket_max = prev
for ev in events:
eid = int(ev.get("id") or 0)
if eid <= prev:
continue
out.append(ev)
if eid > bucket_max:
bucket_max = eid
max_ids[bid] = bucket_max
out.sort(key=lambda x: int(x.get("id") or 0))
return out, max_ids
def main() -> None:
setup_logging()
cfg_path = Path("/opt/activitywatch/dlp-integrations/cef-config.yaml")
cfg = load_yaml(cfg_path)
aw_base = str(cfg.get("aw_api_base", "http://127.0.0.1:5600/api/0")).rstrip("/")
syslog_host = str(cfg.get("syslog_host", "127.0.0.1"))
syslog_port = int(cfg.get("syslog_port", 514))
syslog_proto = str(cfg.get("syslog_proto", "udp")).lower()
per_bucket_limit = int(cfg.get("per_bucket_limit", 300))
state_path = Path(str(cfg.get("state_path", "/var/lib/activitywatch/dlp-integrations/cef-state.json")))
sev_mapping = cfg.get("severity_mapping", {"low": 3, "medium": 6, "high": 10})
if not isinstance(sev_mapping, dict):
sev_mapping = {"low": 3, "medium": 6, "high": 10}
state = load_json(state_path)
incidents, max_ids = iter_new_incidents(aw_base=aw_base, state=state, per_bucket_limit=per_bucket_limit)
sent = 0
for ev in incidents:
line = build_cef(ev, sev_mapping)
if syslog_proto == "tcp":
send_syslog_tcp(line, syslog_host, syslog_port)
else:
send_syslog_udp(line, syslog_host, syslog_port)
sent += 1
state["last_ids"] = max_ids
state["updated_at"] = datetime.now(timezone.utc).isoformat()
save_json(state_path, state)
LOG.info("CEF exporter done: sent=%d buckets=%d target=%s:%d/%s", sent, len(max_ids), syslog_host, syslog_port, syslog_proto)
if __name__ == "__main__":
main()
@@ -4,7 +4,7 @@ After=network-online.target
[Service]
Type=oneshot
WorkingDirectory=/opt/activitywatch/dlp-integrations
ExecStart=/opt/activitywatch/dlp-integrations/.venv/bin/python /opt/activitywatch/dlp-integrations/syslog_forwarder.py
WorkingDirectory=/var/lib/activitywatch
ExecStart=/usr/local/bin/dlp-syslog-forwarder-rust
User=activitywatch
Group=activitywatch
@@ -1,146 +0,0 @@
#!/usr/bin/env python3
from __future__ import annotations
import json
import logging
import socket
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from urllib import error, request
import yaml
LOG = logging.getLogger("aw.dlp.syslog_forwarder")
def setup_logging() -> None:
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
def load_yaml(path: Path) -> dict[str, Any]:
if not path.exists():
return {}
data = yaml.safe_load(path.read_text(encoding="utf-8")) or {}
return data if isinstance(data, dict) else {}
def load_json(path: Path) -> dict[str, Any]:
if not path.exists():
return {}
try:
data = json.loads(path.read_text(encoding="utf-8"))
except Exception:
return {}
return data if isinstance(data, dict) else {}
def save_json(path: Path, payload: dict[str, Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
def http_json(url: str, timeout: int = 15) -> Any:
req = request.Request(url, method="GET")
with request.urlopen(req, timeout=timeout) as resp:
return json.loads(resp.read().decode("utf-8", errors="ignore"))
def iter_new_incidents(aw_base: str, state: dict[str, Any], per_bucket_limit: int) -> tuple[list[dict[str, Any]], dict[str, int]]:
buckets = http_json(f"{aw_base}/buckets/")
bucket_ids = sorted([bid for bid in buckets.keys() if bid.startswith("aw-dlp-incidents_")])
last_ids = state.get("last_ids", {})
if not isinstance(last_ids, dict):
last_ids = {}
max_ids: dict[str, int] = {}
out: list[dict[str, Any]] = []
for bid in bucket_ids:
try:
events = http_json(f"{aw_base}/buckets/{bid}/events?limit={int(per_bucket_limit)}")
except (TimeoutError, OSError, error.URLError, error.HTTPError) as exc:
LOG.warning("skip bucket %s: %s", bid, exc)
continue
prev = int(last_ids.get(bid, 0))
bucket_max = prev
for ev in events:
eid = int(ev.get("id") or 0)
if eid <= prev:
continue
out.append(ev)
if eid > bucket_max:
bucket_max = eid
max_ids[bid] = bucket_max
out.sort(key=lambda x: int(x.get("id") or 0))
return out, max_ids
def build_message(event: dict[str, Any], app_name: str, facility: int) -> str:
pri = facility * 8 + 6
ts = datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
data = event.get("data") or {}
host = str(data.get("hostname") or "unknown")
payload = json.dumps(
{
"event_id": event.get("id"),
"timestamp": event.get("timestamp"),
"host": host,
"severity": data.get("severity"),
"signalType": data.get("signalType"),
"username": data.get("username"),
"action": data.get("action"),
"message": data.get("message"),
"data": data,
},
ensure_ascii=False,
separators=(",", ":"),
)
return f"<{pri}>1 {ts} {host} {app_name} - - - {payload}"
def send_syslog(line: str, host: str, port: int, proto: str, timeout: int = 10) -> None:
if proto.lower() == "tcp":
sock = socket.create_connection((host, port), timeout=timeout)
try:
sock.sendall((line + "\n").encode("utf-8", errors="ignore"))
finally:
sock.close()
return
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
try:
sock.sendto(line.encode("utf-8", errors="ignore"), (host, port))
finally:
sock.close()
def main() -> None:
setup_logging()
cfg_path = Path("/opt/activitywatch/dlp-integrations/syslog-forwarder-config.yaml")
cfg = load_yaml(cfg_path)
aw_base = str(cfg.get("aw_api_base", "http://127.0.0.1:5600/api/0")).rstrip("/")
state_path = Path(str(cfg.get("state_path", "/var/lib/activitywatch/dlp-integrations/syslog-forwarder-state.json")))
per_bucket_limit = int(cfg.get("per_bucket_limit", 300))
syslog_host = str(cfg.get("syslog_host", "127.0.0.1"))
syslog_port = int(cfg.get("syslog_port", 514))
syslog_proto = str(cfg.get("syslog_proto", "udp"))
facility = int(cfg.get("facility", 16))
app_name = str(cfg.get("app_name", "aw-dlp"))
state = load_json(state_path)
try:
incidents, max_ids = iter_new_incidents(aw_base=aw_base, state=state, per_bucket_limit=per_bucket_limit)
except (TimeoutError, OSError, error.URLError, error.HTTPError) as exc:
LOG.warning("skip syslog forwarder run: AW API unavailable: %s", exc)
return
sent = 0
for event in incidents:
line = build_message(event, app_name=app_name, facility=facility)
send_syslog(line=line, host=syslog_host, port=syslog_port, proto=syslog_proto)
sent += 1
save_json(state_path, {"last_ids": max_ids})
LOG.info("syslog forwarder sent=%d buckets=%d", sent, len(max_ids))
if __name__ == "__main__":
main()
@@ -1,59 +0,0 @@
#!/usr/bin/env python3
import importlib.util
import json
import sys
from pathlib import Path
MODULE_PATH = Path(__file__).with_name("syslog_forwarder.py")
SPEC = importlib.util.spec_from_file_location("syslog_forwarder", MODULE_PATH)
MODULE = importlib.util.module_from_spec(SPEC)
sys.modules[SPEC.name] = MODULE
SPEC.loader.exec_module(MODULE)
def test_iter_new_incidents_skips_timed_out_bucket(monkeypatch):
def fake_http_json(url, timeout=15):
if url.endswith("/buckets/"):
return {"aw-dlp-incidents_SHARKON2025": {}}
raise TimeoutError("timed out")
monkeypatch.setattr(MODULE, "http_json", fake_http_json)
incidents, max_ids = MODULE.iter_new_incidents(
aw_base="http://127.0.0.1:5600/api/0",
state={"last_ids": {"aw-dlp-incidents_SHARKON2025": 42}},
per_bucket_limit=300,
)
assert incidents == []
assert max_ids == {}
def test_main_skips_aw_api_timeout_without_overwriting_state(monkeypatch, tmp_path):
state_path = tmp_path / "syslog-forwarder-state.json"
original_state = {"last_ids": {"aw-dlp-incidents_SHARKON2025": 99}}
state_path.write_text(json.dumps(original_state), encoding="utf-8")
monkeypatch.setattr(
MODULE,
"load_yaml",
lambda path: {
"aw_api_base": "http://127.0.0.1:5600/api/0",
"state_path": str(state_path),
},
)
monkeypatch.setattr(
MODULE,
"iter_new_incidents",
lambda aw_base, state, per_bucket_limit: (_ for _ in ()).throw(TimeoutError("timed out")),
)
monkeypatch.setattr(
MODULE,
"save_json",
lambda path, payload: (_ for _ in ()).throw(AssertionError("state should not be saved on AW API timeout")),
)
MODULE.main()
assert json.loads(state_path.read_text(encoding="utf-8")) == original_state
@@ -6,8 +6,8 @@ After=network-online.target
Type=oneshot
User=activitywatch
Group=activitywatch
WorkingDirectory=/opt/activitywatch/dlp-integrations
ExecStart=/opt/activitywatch/dlp-integrations/.venv/bin/python /opt/activitywatch/dlp-integrations/webhook_sender.py
WorkingDirectory=/var/lib/activitywatch
ExecStart=/usr/local/bin/dlp-webhook-sender-rust
[Install]
WantedBy=multi-user.target
@@ -1,155 +0,0 @@
#!/usr/bin/env python3
from __future__ import annotations
import json
import logging
import time
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from urllib import error, request
import yaml
LOG = logging.getLogger("aw.dlp.webhook_sender")
def setup_logging() -> None:
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
def load_yaml(path: Path) -> dict[str, Any]:
if not path.exists():
return {}
data = yaml.safe_load(path.read_text(encoding="utf-8")) or {}
if not isinstance(data, dict):
return {}
return data
def load_json(path: Path) -> dict[str, Any]:
if not path.exists():
return {}
try:
data = json.loads(path.read_text(encoding="utf-8"))
if isinstance(data, dict):
return data
except Exception:
return {}
return {}
def save_json(path: Path, payload: dict[str, Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(json.dumps(payload, ensure_ascii=False, indent=2), encoding="utf-8")
def http_json(url: str, timeout: int = 15) -> Any:
req = request.Request(url, method="GET")
with request.urlopen(req, timeout=timeout) as resp:
return json.loads(resp.read().decode("utf-8", errors="ignore"))
def post_with_retry(url: str, payload: dict[str, Any], retries: int, timeout: int, backoff_base: float) -> bool:
body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
headers = {"Content-Type": "application/json; charset=utf-8"}
for attempt in range(1, retries + 1):
try:
req = request.Request(url, data=body, headers=headers, method="POST")
with request.urlopen(req, timeout=timeout) as resp:
code = getattr(resp, "status", 200)
if 200 <= code < 300:
return True
except error.HTTPError as exc:
LOG.warning("webhook http error url=%s code=%s attempt=%d/%d", url, exc.code, attempt, retries)
except Exception as exc:
LOG.warning("webhook transport error url=%s err=%s attempt=%d/%d", url, exc, attempt, retries)
if attempt < retries:
time.sleep(backoff_base ** (attempt - 1))
return False
def iter_new_incidents(aw_base: str, state: dict[str, Any], per_bucket_limit: int) -> tuple[list[dict[str, Any]], dict[str, int]]:
buckets = http_json(f"{aw_base}/buckets/")
bucket_ids = sorted([bid for bid in buckets.keys() if bid.startswith("aw-dlp-incidents_")])
last_ids = state.get("last_ids", {})
if not isinstance(last_ids, dict):
last_ids = {}
max_ids: dict[str, int] = {}
out: list[dict[str, Any]] = []
for bid in bucket_ids:
events = http_json(f"{aw_base}/buckets/{bid}/events?limit={int(per_bucket_limit)}")
prev = int(last_ids.get(bid, 0))
bucket_max = prev
for ev in events:
eid = int(ev.get("id") or 0)
if eid <= prev:
continue
out.append(ev)
if eid > bucket_max:
bucket_max = eid
max_ids[bid] = bucket_max
out.sort(key=lambda x: int(x.get("id") or 0))
return out, max_ids
def should_send(severity: str, allowed: list[str]) -> bool:
return severity.lower() in {s.lower() for s in allowed}
def main() -> None:
setup_logging()
cfg_path = Path("/opt/activitywatch/dlp-integrations/webhook-config.yaml")
cfg = load_yaml(cfg_path)
aw_base = str(cfg.get("aw_api_base", "http://127.0.0.1:5600/api/0")).rstrip("/")
state_path = Path(str(cfg.get("state_path", "/var/lib/activitywatch/dlp-integrations/webhook-state.json")))
retries = int(cfg.get("retries", 4))
timeout = int(cfg.get("timeout_sec", 15))
backoff_base = float(cfg.get("backoff_base", 2.0))
per_bucket_limit = int(cfg.get("per_bucket_limit", 300))
hooks = cfg.get("critical_webhooks", [])
if not isinstance(hooks, list):
hooks = []
state = load_json(state_path)
incidents, max_ids = iter_new_incidents(aw_base=aw_base, state=state, per_bucket_limit=per_bucket_limit)
sent = 0
for ev in incidents:
data = ev.get("data") or {}
severity = str(data.get("severity") or "low")
for hook in hooks:
if not isinstance(hook, dict):
continue
url = str(hook.get("url") or "").strip()
if not url:
continue
allowed = hook.get("severity", ["high"])
if isinstance(allowed, str):
allowed = [allowed]
if not should_send(severity, [str(x) for x in allowed]):
continue
payload = {
"source": "AWatch-rus DLP",
"timestamp": ev.get("timestamp"),
"event_id": ev.get("id"),
"severity": severity,
"message": data.get("message"),
"ruleId": data.get("ruleId"),
"signalType": data.get("signalType"),
"hostname": data.get("hostname"),
"username": data.get("username"),
"action": data.get("action"),
"raw": data,
}
if post_with_retry(url=url, payload=payload, retries=retries, timeout=timeout, backoff_base=backoff_base):
sent += 1
state["last_ids"] = max_ids
state["updated_at"] = datetime.now(timezone.utc).isoformat()
save_json(state_path, state)
LOG.info("Webhook sender done: delivered=%d incidents_seen=%d", sent, len(incidents))
if __name__ == "__main__":
main()