408 lines
12 KiB
Python
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),
|
|
)
|