-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathuser_event.py
More file actions
55 lines (46 loc) · 1.97 KB
/
Copy pathuser_event.py
File metadata and controls
55 lines (46 loc) · 1.97 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
import json
import logging
from pydantic import ValidationError
import metrics
from db import DB
from validator import parse_event
logger = logging.getLogger(__name__)
class UserEvent:
@staticmethod
def save_user_event(ch, method, properties, body: bytes):
# --- Deserialise -------------------------------------------------
try:
payload = json.loads(body)
except json.JSONDecodeError as exc:
logger.error("Invalid JSON — rejecting without requeue: %s", exc)
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
metrics.increment_error("unknown", "validation")
return
event_type = payload.get("event", "unknown")
# --- Validate ----------------------------------------------------
try:
parsed = parse_event(payload)
except (ValidationError, ValueError) as exc:
logger.error("Validation error — rejecting without requeue: %s", exc)
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
metrics.increment_error(event_type, "validation")
return
user_id = str(getattr(parsed, "user_id", ""))
# --- Persist -----------------------------------------------------
db = DB()
try:
db.insert_event(event_type, user_id, payload)
ch.basic_ack(delivery_tag=method.delivery_tag)
logger.info("Saved event=%s user_id=%s", event_type, user_id)
metrics.increment_event(event_type)
except Exception as exc:
requeue = not method.redelivered
logger.error(
"MongoDB error for event=%s user_id=%s: %s — %s",
event_type, user_id, exc,
"requeueing" if requeue else "sending to DLQ",
)
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=requeue)
metrics.increment_error(event_type, "db")
finally:
db.close()