|
6 | 6 | import logging |
7 | 7 | import os |
8 | 8 | import re |
9 | | -from datetime import datetime, timezone |
| 9 | +from datetime import UTC, datetime |
10 | 10 | from typing import Any |
11 | 11 |
|
| 12 | +from botocore.exceptions import BotoCoreError, ClientError |
| 13 | + |
12 | 14 | logger = logging.getLogger(__name__) |
13 | 15 |
|
14 | 16 | ALLOWED_EVENT_TYPES = {"impression", "click", "skip", "cart", "purchase", "dislike"} |
@@ -63,16 +65,14 @@ def parse_feedback_message(body: str) -> dict[str, Any]: |
63 | 65 | if not isinstance(payload["event_id"], str) or not payload["event_id"]: |
64 | 66 | raise TypeError("event_id must be a non-empty string") |
65 | 67 |
|
66 | | - occurred_at = str(payload["occurred_at"]) |
67 | | - parsed = datetime.fromisoformat(occurred_at.replace("Z", "+00:00")) |
| 68 | + parsed = datetime.fromisoformat(str(payload["occurred_at"])) |
68 | 69 | if parsed.tzinfo is None: |
69 | 70 | raise ValueError("occurred_at must include a timezone") |
70 | 71 | return payload |
71 | 72 |
|
72 | 73 |
|
73 | 74 | def _storage_key(payload: dict[str, Any], message_id: str) -> str: |
74 | | - occurred = datetime.fromisoformat(str(payload["occurred_at"]).replace("Z", "+00:00")) |
75 | | - occurred = occurred.astimezone(timezone.utc) |
| 75 | + occurred = datetime.fromisoformat(str(payload["occurred_at"])).astimezone(UTC) |
76 | 76 | safe_message_id = re.sub(r"[^A-Za-z0-9_.-]", "_", message_id) |
77 | 77 | return ( |
78 | 78 | f"feedback/event_date={occurred:%Y-%m-%d}/hour={occurred:%H}/" |
@@ -125,7 +125,7 @@ def lambda_handler(event: dict[str, Any], context: Any) -> dict[str, list[dict[s |
125 | 125 | raise TypeError("SQS record must be an object") |
126 | 126 | key = _process_record(record, bucket=bucket, s3_client=s3_client) |
127 | 127 | logger.info("feedback_landed message_id=%s key=%s", message_id, key) |
128 | | - except Exception as exc: # Lambda must isolate malformed records in a batch. |
| 128 | + except (TypeError, ValueError, BotoCoreError, ClientError) as exc: |
129 | 129 | logger.warning( |
130 | 130 | "feedback_ingestion_failed message_id=%s error_type=%s", |
131 | 131 | message_id, |
|
0 commit comments