"""Import scraped SFDA alerts into SFDA Entries documents.""" from __future__ import annotations import re import time from datetime import date from typing import Any import frappe from frappe.utils import now_datetime from frappe.utils.file_manager import save_file from sfda_parser.sfda_scraper.detail_page import fetch_ade_detail from sfda_parser.sfda_scraper.detail_pdf import parse_affected_device_list from sfda_parser.sfda_scraper.http_client import ADE_REQUEST_DELAY_SECONDS, create_session from sfda_parser.sfda_scraper.weekly_list import ( WEEKLY_BULLETINS_TO_IMPORT, WeeklyAlertRow, fetch_recent_weekly_alerts, ) from sfda_parser.sfda_scraper.weekly_pdf import SafetyAlertRow, parse_safety_alerts_from_pdf # Tracks which weekly PDF was last processed (persists across scheduled runs). LAST_WEEKLY_PDF_KEY = "sfda_parser_last_processed_weekly_pdf_url" FAILED_REFS_CACHE_KEY = "sfda_parser_failed_import_refs" def _normalize_ncmdr_ref(ncmdr_ref: str) -> str: return (ncmdr_ref or "").strip().upper() def _attach_detail_pdf(doc_name: str, ncmdr_ref: str, pdf_bytes: bytes) -> str: safe_ref = re.sub(r"[^\w\-]", "_", ncmdr_ref) file_doc = save_file( f"{safe_ref}-ade-detail.pdf", pdf_bytes, "SFDA Entries", doc_name, decode=False, is_private=0, ) return file_doc.file_url def _set_detail_pdf_on_child_rows(doc_name: str, ncmdr_ref: str, file_url: str) -> None: target = _normalize_ncmdr_ref(ncmdr_ref) for row in frappe.get_all( "SFDA Device Entries", filters={"parent": doc_name, "parenttype": "SFDA Entries"}, fields=["name", "ncmdr_ref"], ): if _normalize_ncmdr_ref(row.get("ncmdr_ref") or "") == target: frappe.db.set_value("SFDA Device Entries", row.name, "detail_pdf", file_url) def _get_failed_refs() -> list[dict[str, Any]]: return frappe.cache().get_value(FAILED_REFS_CACHE_KEY) or [] def _save_failed_refs(refs: list[dict[str, Any]]) -> None: frappe.cache().set_value(FAILED_REFS_CACHE_KEY, refs) def _add_failed_ref( ncmdr_ref: str, detail_url: str, weekly_title: str, weekly_date: date | None, ) -> None: refs = _get_failed_refs() existing = {_normalize_ncmdr_ref(item.get("ncmdr_ref") or "") for item in refs} if _normalize_ncmdr_ref(ncmdr_ref) in existing: return refs.append( { "ncmdr_ref": ncmdr_ref, "detail_url": detail_url, "weekly_title": weekly_title, "weekly_date": str(weekly_date) if weekly_date else None, } ) _save_failed_refs(refs) def _remove_failed_ref(ncmdr_ref: str) -> None: target = _normalize_ncmdr_ref(ncmdr_ref) refs = [ item for item in _get_failed_refs() if _normalize_ncmdr_ref(item.get("ncmdr_ref") or "") != target ] _save_failed_refs(refs) def _find_weekly_entry_name(weekly_title: str, weekly_date: date | None) -> str | None: filters: dict[str, Any] = {"title": weekly_title} if weekly_date: filters["date"] = weekly_date return frappe.db.get_value("SFDA Entries", filters) def _get_existing_ncmdr_refs(doc_name: str) -> set[str]: rows = frappe.get_all( "SFDA Device Entries", filters={"parent": doc_name, "parenttype": "SFDA Entries"}, pluck="ncmdr_ref", ) return {_normalize_ncmdr_ref(ref) for ref in rows if ref} def _build_device_rows_for_alert( alert: SafetyAlertRow, *, manufacturer: str, devices: list, ) -> list[dict[str, Any]]: ncmdr_ref = (alert.ncmdr_ref or "").strip() base = { "doctype": "SFDA Device Entries", "ncmdr_ref": ncmdr_ref, "manufacturer": manufacturer, "ade_detail_url": alert.detail_url, } if devices: return [{**base, **device.as_dict()} for device in devices] return [base] def _import_single_alert_into_weekly_entry( alert: SafetyAlertRow, *, weekly_title: str, weekly_date: date | None, session, summary: dict, doc, existing_ncmdr_refs: set[str], pending_pdfs: list[tuple[str, bytes]], ) -> bool: """Fetch one alert and append child rows to the weekly SFDA Entries doc.""" ncmdr_ref = (alert.ncmdr_ref or "").strip() if not ncmdr_ref: return False normalized_ref = _normalize_ncmdr_ref(ncmdr_ref) if normalized_ref in existing_ncmdr_refs: summary["skipped_existing"] += 1 _remove_failed_ref(ncmdr_ref) return False try: detail_pdf, publish_details = fetch_ade_detail(alert.detail_url, session=session) manufacturer = (publish_details.get("manufacturer") or "").strip() devices = parse_affected_device_list(detail_pdf) if detail_pdf else [] if not detail_pdf: warning_msg = ( f"{ncmdr_ref} ({alert.detail_url}): No PDF on ADE page; " "created child row from HTML fields only" ) summary["warnings"].append(warning_msg) frappe.logger("sfda_scraper").warning(warning_msg) for row in _build_device_rows_for_alert( alert, manufacturer=manufacturer, devices=devices, ): doc.append("device_list", row) if detail_pdf: pending_pdfs.append((ncmdr_ref, detail_pdf)) summary["pdfs_saved"] += 1 existing_ncmdr_refs.add(normalized_ref) summary["alerts_imported"] += 1 _remove_failed_ref(ncmdr_ref) return True except Exception as exc: error_msg = f"{ncmdr_ref} ({alert.detail_url}): {exc}" summary["errors"].append(error_msg) _add_failed_ref(ncmdr_ref, alert.detail_url, weekly_title, weekly_date) frappe.log_error( title="SFDA Scraper: Alert import failed", message=error_msg, ) return False def _get_or_create_weekly_doc(weekly_title: str, weekly_date: date | None): existing_name = _find_weekly_entry_name(weekly_title, weekly_date) if existing_name: return frappe.get_doc("SFDA Entries", existing_name), _get_existing_ncmdr_refs( existing_name ), True doc = frappe.get_doc( { "doctype": "SFDA Entries", "title": weekly_title, "date": weekly_date, } ) return doc, set(), False def _save_weekly_doc(doc, is_existing: bool) -> None: if is_existing: doc.save(ignore_permissions=True) else: doc.insert(ignore_permissions=True) frappe.db.commit() def _attach_pending_pdfs(doc_name: str, pending_pdfs: list[tuple[str, bytes]]) -> None: for ncmdr_ref, pdf_bytes in pending_pdfs: file_url = _attach_detail_pdf(doc_name, ncmdr_ref, pdf_bytes) _set_detail_pdf_on_child_rows(doc_name, ncmdr_ref, file_url) if pending_pdfs: frappe.db.commit() def _import_weekly_bulletin( weekly: WeeklyAlertRow, alerts: list[SafetyAlertRow], *, session, summary: dict, ) -> None: doc, existing_ncmdr_refs, is_existing = _get_or_create_weekly_doc( weekly.title, weekly.alert_date, ) pending_pdfs: list[tuple[str, bytes]] = [] imported_any = False for alert in alerts: if _import_single_alert_into_weekly_entry( alert, weekly_title=weekly.title, weekly_date=weekly.alert_date, session=session, summary=summary, doc=doc, existing_ncmdr_refs=existing_ncmdr_refs, pending_pdfs=pending_pdfs, ): imported_any = True time.sleep(ADE_REQUEST_DELAY_SECONDS) if not imported_any: return _save_weekly_doc(doc, is_existing) _attach_pending_pdfs(doc.name, pending_pdfs) if not is_existing: summary["created"] += 1 else: summary["updated"] += 1 def _retry_failed_refs(session, summary: dict) -> None: """Retry alerts that failed on a previous run (e.g. ADE 429 rate limits).""" failed_refs = _get_failed_refs() if not failed_refs: return summary["failed_refs_queued"] = len(failed_refs) frappe.logger("sfda_scraper").info( "SFDA scraper retrying %s previously failed alert(s)", len(failed_refs), ) weekly_groups: dict[tuple[str, str | None], list[dict[str, Any]]] = {} for item in failed_refs: weekly_date = item.get("weekly_date") weekly_groups.setdefault( (item.get("weekly_title") or "", weekly_date), [], ).append(item) for (weekly_title, weekly_date_str), items in weekly_groups.items(): weekly_date = None if weekly_date_str: try: weekly_date = date.fromisoformat(weekly_date_str) except ValueError: weekly_date = None weekly = WeeklyAlertRow( title=weekly_title, alert_date=weekly_date, pdf_url="", ) alerts = [ SafetyAlertRow( ncmdr_ref=item.get("ncmdr_ref") or "", detail_url=item.get("detail_url") or "", ) for item in items ] _import_weekly_bulletin(weekly, alerts, session=session, summary=summary) def import_latest_weekly_alerts() -> dict: """ Recurring job: reads the most recent weekly bulletin(s) from SFDA and imports alerts into one SFDA Entries document per weekly bulletin. Each scheduled run: - Retries previously failed refs from the cache queue - Fetches https://www.sfda.gov.sa/en/weekly-alert - Uses the last WEEKLY_BULLETINS_TO_IMPORT weekly PDFs - Creates or updates one SFDA Entries doc per week with child rows per alert - Skips NCMDR refs already present in that week's device_list """ session = create_session() summary = { "weekly_bulletins_to_import": WEEKLY_BULLETINS_TO_IMPORT, "weekly_title": "", "weekly_date": None, "weekly_pdf_url": "", "weekly_titles": [], "alerts_found": 0, "alerts_imported": 0, "created": 0, "updated": 0, "pdfs_saved": 0, "skipped_existing": 0, "weekly_already_processed": False, "failed_refs_queued": 0, "warnings": [], "errors": [], } _retry_failed_refs(session, summary) weeklies = fetch_recent_weekly_alerts(session=session) if not weeklies: frappe.log_error( title="SFDA Scraper: No weekly bulletins found", message="https://www.sfda.gov.sa/en/weekly-alert", ) _log_run_summary(summary) return summary weekly_newest = weeklies[0] summary["weekly_title"] = weekly_newest.title summary["weekly_date"] = ( str(weekly_newest.alert_date) if weekly_newest.alert_date else None ) summary["weekly_pdf_url"] = weekly_newest.pdf_url summary["weekly_titles"] = [w.title for w in weeklies] last_pdf_url = frappe.db.get_default(LAST_WEEKLY_PDF_KEY) # Oldest week first so backfill order matches bulletin dates. for weekly in reversed(weeklies): pdf_response = session.get(weekly.pdf_url, timeout=120) pdf_response.raise_for_status() alerts = parse_safety_alerts_from_pdf(pdf_response.content) if not alerts: frappe.log_error( title="SFDA Scraper: No alerts in weekly PDF", message=f"URL: {weekly.pdf_url}", ) continue summary["alerts_found"] += len(alerts) _import_weekly_bulletin(weekly, alerts, session=session, summary=summary) # Same newest weekly PDF as last run and every alert already in ERP if ( last_pdf_url == weekly_newest.pdf_url and summary["created"] == 0 and summary["updated"] == 0 and summary["skipped_existing"] == summary["alerts_found"] and summary["alerts_found"] > 0 ): summary["weekly_already_processed"] = True # Track newest bulletin URL (detect when SFDA publishes a new top-row PDF) frappe.db.set_default(LAST_WEEKLY_PDF_KEY, weekly_newest.pdf_url) _log_run_summary(summary) return summary def _log_run_summary(summary: dict) -> None: frappe.logger("sfda_scraper").info( "SFDA scheduled import at %s — latest weekly bulletin '%s' (%s): %s", now_datetime(), summary.get("weekly_title"), summary.get("weekly_date"), summary, ) if summary.get("errors"): frappe.log_error( title="SFDA Scraper: Run completed with errors", message=frappe.as_json(summary, indent=2), )