11from __future__ import annotations
22
3- from uuid import UUID
3+ import json
4+ from datetime import datetime , timezone
5+ from uuid import UUID , uuid5
46
5- from fastapi import APIRouter , Depends
7+ from fastapi import APIRouter , Depends , Request
8+ from pydantic import ValidationError as PydanticValidationError
69from sqlalchemy import select
710from sqlalchemy .ext .asyncio import AsyncSession
811
912from app .core .database import get_db
13+ from app .core .exceptions import BadRequestError
1014from app .core .permissions import Permission , require_plan_permission
11- from app .models import InvalidPayloadError
12- from app .schemas .tracking_plan import InvalidPayloadErrorResponse
15+ from app .models import InvalidPayloadError , ValidationLog
16+ from app .schemas .tracking_plan import (
17+ GuardDLQImportRecord ,
18+ GuardDLQImportResponse ,
19+ InvalidPayloadErrorResponse ,
20+ )
21+ from app .services .snapshot_service import SnapshotService
1322
1423router = APIRouter (tags = ["dlq" ])
1524
25+ GUARD_DLQ_IMPORT_NAMESPACE = UUID ("f92c4282-7b14-4e2e-9ac1-0e6cf15f4a6c" )
26+ MAX_IMPORT_RECORDS = 1000
27+
1628
1729@router .get ("/plans/{plan_id}/dlq" , response_model = list [InvalidPayloadErrorResponse ])
1830async def get_dlq_errors (
@@ -27,3 +39,208 @@ async def get_dlq_errors(
2739 .limit (200 )
2840 )
2941 return list (result .scalars ().all ())
42+
43+
44+ @router .post ("/plans/{plan_id}/dlq/import" , response_model = GuardDLQImportResponse )
45+ async def import_guard_dlq_errors (
46+ plan_id : UUID ,
47+ request : Request ,
48+ access = Depends (require_plan_permission (Permission .EDIT )),
49+ db : AsyncSession = Depends (get_db ),
50+ ):
51+ del access
52+ records = await _parse_import_records (request )
53+ latest_version = await SnapshotService (db ).get_latest_version (plan_id )
54+ version_id = latest_version .id if latest_version else None
55+
56+ imported = 0
57+ skipped = 0
58+ upserted_errors = 0
59+ for record in records :
60+ request_id = uuid5 (GUARD_DLQ_IMPORT_NAMESPACE , f"{ plan_id } :{ record .guard_dlq_id } " )
61+ existing_log = await db .execute (
62+ select (ValidationLog .id ).where (
63+ ValidationLog .plan_id == plan_id ,
64+ ValidationLog .request_id == request_id ,
65+ )
66+ )
67+ if existing_log .scalar_one_or_none () is not None :
68+ skipped += 1
69+ continue
70+
71+ seen_at = _record_seen_at (record )
72+ validation_log = ValidationLog (
73+ plan_id = plan_id ,
74+ event_name = record .event_name ,
75+ payload = record .payload ,
76+ is_valid = False ,
77+ errors = _record_errors (record ),
78+ version_id = version_id ,
79+ api_key_id = None ,
80+ request_id = request_id ,
81+ source_ip = None ,
82+ source_label = "guard-dlq-import" ,
83+ validated_at = seen_at ,
84+ )
85+ db .add (validation_log )
86+ await db .flush ()
87+
88+ await _upsert_imported_invalid_payload (
89+ db ,
90+ plan_id = plan_id ,
91+ version_id = version_id ,
92+ validation_log_id = validation_log .id ,
93+ event_name = record .event_name ,
94+ payload = record .payload ,
95+ error_reason = _error_reason (record ),
96+ seen_at = seen_at ,
97+ )
98+ imported += 1
99+ upserted_errors += 1
100+
101+ return GuardDLQImportResponse (
102+ imported = imported ,
103+ skipped = skipped ,
104+ upserted_errors = upserted_errors ,
105+ )
106+
107+
108+ async def _parse_import_records (request : Request ) -> list [GuardDLQImportRecord ]:
109+ body = (await request .body ()).strip ()
110+ if not body :
111+ raise BadRequestError ("Import body is empty." , code = "dlq_import_empty" )
112+
113+ try :
114+ decoded = json .loads (body )
115+ except json .JSONDecodeError :
116+ decoded = _decode_ndjson (body )
117+
118+ if isinstance (decoded , dict ) and "records" in decoded :
119+ raw_records = decoded ["records" ]
120+ elif isinstance (decoded , dict ):
121+ raw_records = [decoded ]
122+ else :
123+ raw_records = decoded
124+
125+ if not isinstance (raw_records , list ):
126+ raise BadRequestError (
127+ "DLQ import body must be NDJSON, a JSON array, or an object with records." ,
128+ code = "dlq_import_invalid" ,
129+ )
130+ if not raw_records :
131+ raise BadRequestError ("DLQ import contains no records." , code = "dlq_import_empty" )
132+ if len (raw_records ) > MAX_IMPORT_RECORDS :
133+ raise BadRequestError (
134+ f"DLQ import accepts at most { MAX_IMPORT_RECORDS } records per request." ,
135+ code = "dlq_import_too_large" ,
136+ )
137+
138+ records : list [GuardDLQImportRecord ] = []
139+ for index , raw_record in enumerate (raw_records ):
140+ try :
141+ records .append (GuardDLQImportRecord .model_validate (raw_record ))
142+ except PydanticValidationError as exc :
143+ raise BadRequestError (
144+ "DLQ import record is invalid." ,
145+ code = "dlq_import_invalid_record" ,
146+ extra = {"record_index" : index , "errors" : exc .errors ()},
147+ ) from exc
148+ return records
149+
150+
151+ def _decode_ndjson (body : bytes ) -> list [dict ]:
152+ records : list [dict ] = []
153+ for line_number , raw_line in enumerate (body .splitlines (), start = 1 ):
154+ line = raw_line .strip ()
155+ if not line :
156+ continue
157+ try :
158+ decoded = json .loads (line )
159+ except json .JSONDecodeError as exc :
160+ raise BadRequestError (
161+ "DLQ import NDJSON line is invalid JSON." ,
162+ code = "dlq_import_invalid_ndjson" ,
163+ extra = {"line" : line_number },
164+ ) from exc
165+ if not isinstance (decoded , dict ):
166+ raise BadRequestError (
167+ "DLQ import NDJSON lines must be JSON objects." ,
168+ code = "dlq_import_invalid_ndjson" ,
169+ extra = {"line" : line_number },
170+ )
171+ records .append (decoded )
172+ return records
173+
174+
175+ def _record_seen_at (record : GuardDLQImportRecord ) -> datetime :
176+ seen_at = record .created_at or datetime .now (timezone .utc )
177+ if seen_at .tzinfo is None :
178+ return seen_at .replace (tzinfo = timezone .utc )
179+ return seen_at .astimezone (timezone .utc )
180+
181+
182+ def _record_errors (record : GuardDLQImportRecord ) -> list [dict ]:
183+ reason_codes = record .reason_codes or ["guard_dlq_reject" ]
184+ return [
185+ {
186+ "code" : code ,
187+ "path" : "payload" ,
188+ "property_name" : None ,
189+ "message" : f"Guard rejected payload: { code } " ,
190+ }
191+ for code in reason_codes
192+ ]
193+
194+
195+ def _error_reason (record : GuardDLQImportRecord ) -> str :
196+ reason_codes = ", " .join (record .reason_codes or ["guard_dlq_reject" ])
197+ return f"guard dlq rejected payload: { reason_codes } "
198+
199+
200+ async def _upsert_imported_invalid_payload (
201+ db : AsyncSession ,
202+ * ,
203+ plan_id : UUID ,
204+ version_id : UUID | None ,
205+ validation_log_id : UUID ,
206+ event_name : str ,
207+ payload : dict ,
208+ error_reason : str ,
209+ seen_at : datetime ,
210+ ) -> None :
211+ version_clause = (
212+ InvalidPayloadError .version_id == version_id
213+ if version_id is not None
214+ else InvalidPayloadError .version_id .is_ (None )
215+ )
216+ result = await db .execute (
217+ select (InvalidPayloadError ).where (
218+ InvalidPayloadError .plan_id == plan_id ,
219+ version_clause ,
220+ InvalidPayloadError .event_name == event_name ,
221+ InvalidPayloadError .error_reason == error_reason ,
222+ )
223+ )
224+ invalid_payload = result .scalar_one_or_none ()
225+ if invalid_payload is None :
226+ db .add (
227+ InvalidPayloadError (
228+ plan_id = plan_id ,
229+ version_id = version_id ,
230+ validation_log_id = validation_log_id ,
231+ event_name = event_name ,
232+ payload = payload ,
233+ error_reason = error_reason ,
234+ first_seen_at = seen_at ,
235+ last_seen_at = seen_at ,
236+ occurrence_count = 1 ,
237+ )
238+ )
239+ return
240+
241+ invalid_payload .payload = payload
242+ invalid_payload .validation_log_id = validation_log_id
243+ invalid_payload .last_seen_at = max (invalid_payload .last_seen_at or seen_at , seen_at )
244+ if invalid_payload .first_seen_at is None or seen_at < invalid_payload .first_seen_at :
245+ invalid_payload .first_seen_at = seen_at
246+ invalid_payload .occurrence_count += 1
0 commit comments