408 lines
12 KiB
Python

"""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),
)