From 1b7f00e7a0ab3a8826787cd3865e995be7cbea17 Mon Sep 17 00:00:00 2001 From: Christian Krakau-Louis Date: Fri, 22 May 2026 19:49:43 +0200 Subject: [PATCH] feat: add per-source mail import trigger --- .../app/api/api_v1/endpoints/mail_sources.py | 112 ++++++++++ backend/app/templates/mail_sources.html | 37 ++++ backend/app/tests/test_mail_sources.py | 197 ++++++++++++++++++ docs/development/roadmap.md | 1 + docs/milestones.md | 1 + docs/todo.md | 1 + 6 files changed, 349 insertions(+) diff --git a/backend/app/api/api_v1/endpoints/mail_sources.py b/backend/app/api/api_v1/endpoints/mail_sources.py index d3b54f9..e1847de 100644 --- a/backend/app/api/api_v1/endpoints/mail_sources.py +++ b/backend/app/api/api_v1/endpoints/mail_sources.py @@ -202,6 +202,88 @@ def _import_to_response(row: MailSourceImport) -> MailSourceImportResponse: ) +def _fetch_response(source: MailSource, results: Dict[str, Any]) -> Dict[str, Any]: + """Build the common response payload for a manual source fetch.""" + return { + "source_id": source.id, + "name": source.name, + "success": bool(results.get("success", False)), + "processed": int(results.get("processed", 0)), + "reports_found": int(results.get("reports_found", 0)), + "duplicate_reports": int(results.get("duplicate_reports", 0)), + "new_domains": [str(d) for d in results.get("new_domains", [])], + "error_count": len(results.get("errors", [])), + "timestamp": datetime.now().isoformat(), + } + + +def _fetch_gmail_source(source: MailSource, db: Session) -> Dict[str, Any]: + """Run one Gmail API import and persist source/import metadata.""" + if not source.gmail_access_token: + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail="Gmail account not yet authorised. Complete OAuth2 flow first.", + ) + + already = GmailClient.load_ingested_ids(source.gmail_ingested_ids) + client = GmailClient( + client_id=source.gmail_client_id or "", + client_secret=source.gmail_client_secret or "", + access_token=source.gmail_access_token, + refresh_token=source.gmail_refresh_token or "", + already_ingested_ids=already, + db=db, + ) + + started_at = datetime.utcnow() + results = client.fetch_reports() + + if results.get("new_ingested_ids"): + all_ids = list(dict.fromkeys(already + results["new_ingested_ids"])) + source.gmail_ingested_ids = GmailClient.dump_ingested_ids(all_ids) + + refreshed = client.get_refreshed_tokens() + if refreshed: + source.gmail_access_token = refreshed["access_token"] + if "refresh_token" in refreshed: + source.gmail_refresh_token = refreshed["refresh_token"] + + source.last_checked = datetime.utcnow() + record_import_attempt(db, source, results, started_at=started_at, trigger="manual") + db.commit() + return results + + +def _fetch_imap_source(source: MailSource, db: Session, days: int) -> Dict[str, Any]: + """Run one IMAP import and persist source/import metadata.""" + client = IMAPClient( + server=source.server, + port=source.port or 993, + username=source.username, + password=source.password, + delete_emails=False, + db=db, + ) + started_at = datetime.utcnow() + results = client.fetch_reports(days=days) + source.last_checked = datetime.utcnow() + record_import_attempt(db, source, results, started_at=started_at, trigger="manual") + db.commit() + return results + + +def _fetch_source(source: MailSource, db: Session, days: int) -> Dict[str, Any]: + """Dispatch a manual fetch for one configured mail source.""" + if source.method == "GMAIL_API": + return _fetch_gmail_source(source, db) + if source.method == "IMAP": + return _fetch_imap_source(source, db, days) + raise HTTPException( + status_code=status.HTTP_400_BAD_REQUEST, + detail=f"Manual fetch is not available for method '{source.method}'.", + ) + + # --------------------------------------------------------------------------- # Routes # --------------------------------------------------------------------------- @@ -278,6 +360,36 @@ async def list_mail_source_imports( return [_import_to_response(row) for row in rows] +@router.post("/{source_id}/fetch", response_model=Dict[str, Any]) +async def fetch_mail_source( + source_id: int, + days: int = 7, + db: Session = Depends(get_db), + _auth: dict = Depends(require_admin_auth), +) -> Dict[str, Any]: + """Manually fetch DMARC reports for one configured mail source.""" + if days < 1 or days > 365: + raise HTTPException(status_code=400, detail="Days parameter must be between 1 and 365") + + source = _get_source_or_404(source_id, db) + results = _fetch_source(source, db, days) + logger.info( + "Manual fetch for source id=%d: processed=%d reports_found=%d duplicates=%d", + int(source_id), + int(results.get("processed", 0)), + int(results.get("reports_found", 0)), + int(results.get("duplicate_reports", 0)), + ) + for err in results.get("errors", []): + logger.warning( + "Manual fetch warning for source id=%d: %s", + int(source_id), + _sanitize_for_log(err), + ) + + return _fetch_response(source, results) + + @router.put("/{source_id}", response_model=MailSourceResponse) async def update_mail_source( source_id: int, diff --git a/backend/app/templates/mail_sources.html b/backend/app/templates/mail_sources.html index 2eed4bd..34cf85a 100644 --- a/backend/app/templates/mail_sources.html +++ b/backend/app/templates/mail_sources.html @@ -105,6 +105,15 @@ +