feat: Add logging and reporting capabilities with GDPR compliance

Agent-Logs-Url: https://github.com/christianlouis/pop_puller_to_gmail/sessions/d6e65ca9-f07f-43a2-bac2-09eb44588ded

Co-authored-by: christianlouis <361235+christianlouis@users.noreply.github.com>
This commit is contained in:
copilot-swe-agent[bot]
2026-03-26 18:12:22 +00:00
parent 276ef1f8b2
commit fd45c838fc
13 changed files with 977 additions and 20 deletions
+4
View File
@@ -13,6 +13,7 @@ from app.api.v1.endpoints import (
admin,
providers,
app_settings,
logs,
)
api_router = APIRouter()
@@ -34,3 +35,6 @@ api_router.include_router(
)
api_router.include_router(admin.router, prefix="/admin", tags=["Admin"])
api_router.include_router(app_settings.router, prefix="/settings", tags=["Settings"])
api_router.include_router(
logs.router, prefix="/processing-runs", tags=["Processing Logs"]
)
+171 -2
View File
@@ -1,15 +1,18 @@
"""Admin endpoints"""
from typing import List
from fastapi import APIRouter, Depends, HTTPException, status
import math
from typing import List, Optional
from fastapi import APIRouter, Depends, HTTPException, Query, status
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, func
from app.core.database import get_db
from app.core.deps import get_current_superuser
from app.core.gdpr import mask_email, mask_from_header
from app.models.database_models import (
User,
MailAccount,
ProcessingLog,
ProcessingRun,
SubscriptionPlan,
SubscriptionTier,
@@ -21,6 +24,10 @@ from app.models.schemas import (
SubscriptionPlanResponse,
SubscriptionPlanCreate,
SubscriptionPlanUpdate,
AdminProcessingRunResponse,
AdminProcessingLogResponse,
PaginatedAdminRunsResponse,
PaginatedAdminLogsResponse,
)
router = APIRouter()
@@ -282,3 +289,165 @@ async def delete_plan(
)
await db.delete(plan)
await db.commit()
# ── Admin Logs ─────────────────────────────────────────────────────────────────
def _admin_paginate(total: int, page: int, page_size: int) -> dict:
pages = max(1, math.ceil(total / page_size)) if total else 1
return {"total": total, "page": page, "page_size": page_size, "pages": pages}
@router.get(
"/processing-runs",
response_model=PaginatedAdminRunsResponse,
summary="List all processing runs across all users (admin only)",
)
async def admin_list_processing_runs(
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
user_id: Optional[int] = Query(None, description="Filter by user ID"),
account_id: Optional[int] = Query(None, description="Filter by mail account ID"),
status_filter: Optional[str] = Query(
None,
alias="status",
description="Filter by run status",
),
current_user: User = Depends(get_current_superuser),
db: AsyncSession = Depends(get_db),
):
"""
Return a paginated list of all processing runs in the system with
GDPR-masked user / account email addresses.
"""
base = (
select(
ProcessingRun,
MailAccount.name.label("account_name"),
MailAccount.email_address.label("account_email"),
MailAccount.user_id.label("uid"),
User.email.label("user_email"),
)
.join(MailAccount, ProcessingRun.mail_account_id == MailAccount.id)
.join(User, MailAccount.user_id == User.id)
)
if user_id is not None:
base = base.where(MailAccount.user_id == user_id)
if account_id is not None:
base = base.where(ProcessingRun.mail_account_id == account_id)
if status_filter:
base = base.where(ProcessingRun.status == status_filter)
total = (
await db.execute(select(func.count()).select_from(base.subquery()))
).scalar_one()
offset = (page - 1) * page_size
rows = (
await db.execute(
base.order_by(ProcessingRun.started_at.desc())
.offset(offset)
.limit(page_size)
)
).all()
items = [
AdminProcessingRunResponse(
id=row.ProcessingRun.id, # type: ignore[arg-type]
mail_account_id=row.ProcessingRun.mail_account_id, # type: ignore[arg-type]
started_at=row.ProcessingRun.started_at, # type: ignore[arg-type]
completed_at=row.ProcessingRun.completed_at, # type: ignore[arg-type]
duration_seconds=row.ProcessingRun.duration_seconds, # type: ignore[arg-type]
emails_fetched=row.ProcessingRun.emails_fetched, # type: ignore[arg-type]
emails_forwarded=row.ProcessingRun.emails_forwarded, # type: ignore[arg-type]
emails_failed=row.ProcessingRun.emails_failed, # type: ignore[arg-type]
status=row.ProcessingRun.status, # type: ignore[arg-type]
error_message=row.ProcessingRun.error_message, # type: ignore[arg-type]
account_name=row.account_name,
account_email=mask_email(row.account_email) if row.account_email else None,
user_id=row.uid,
user_email=mask_email(row.user_email) if row.user_email else None,
)
for row in rows
]
return PaginatedAdminRunsResponse(
items=items, **_admin_paginate(total, page, page_size) # type: ignore[arg-type]
)
@router.get(
"/processing-logs",
response_model=PaginatedAdminLogsResponse,
summary="List all per-email processing logs across all users (admin only)",
)
async def admin_list_processing_logs(
page: int = Query(1, ge=1),
page_size: int = Query(50, ge=1, le=200),
user_id: Optional[int] = Query(None, description="Filter by user ID"),
account_id: Optional[int] = Query(None, description="Filter by mail account ID"),
run_id: Optional[int] = Query(None, description="Filter by processing run ID"),
level: Optional[str] = Query(
None, description="Filter by level (INFO, WARNING, ERROR)"
),
current_user: User = Depends(get_current_superuser),
db: AsyncSession = Depends(get_db),
):
"""
Return paginated per-email log entries with GDPR-masked sender addresses.
Subject lines are shown as-is (the user owns their own mail content);
sender addresses are pseudonymised for operator privacy.
"""
base = select(
ProcessingLog,
User.email.label("user_email"),
).join(User, ProcessingLog.user_id == User.id)
if user_id is not None:
base = base.where(ProcessingLog.user_id == user_id)
if account_id is not None:
base = base.where(ProcessingLog.mail_account_id == account_id)
if run_id is not None:
base = base.where(ProcessingLog.processing_run_id == run_id)
if level:
base = base.where(ProcessingLog.level == level.upper())
total = (
await db.execute(select(func.count()).select_from(base.subquery()))
).scalar_one()
offset = (page - 1) * page_size
rows = (
await db.execute(
base.order_by(ProcessingLog.timestamp.desc())
.offset(offset)
.limit(page_size)
)
).all()
items = [
AdminProcessingLogResponse(
id=row.ProcessingLog.id, # type: ignore[arg-type]
timestamp=row.ProcessingLog.timestamp, # type: ignore[arg-type]
level=row.ProcessingLog.level, # type: ignore[arg-type]
message=row.ProcessingLog.message, # type: ignore[arg-type]
email_subject=row.ProcessingLog.email_subject, # type: ignore[arg-type]
email_from=(
mask_from_header(row.ProcessingLog.email_from)
if row.ProcessingLog.email_from
else None
),
success=row.ProcessingLog.success, # type: ignore[arg-type]
mail_account_id=row.ProcessingLog.mail_account_id, # type: ignore[arg-type]
processing_run_id=row.ProcessingLog.processing_run_id, # type: ignore[arg-type]
email_size_bytes=row.ProcessingLog.email_size_bytes, # type: ignore[arg-type]
error_details=row.ProcessingLog.error_details, # type: ignore[arg-type]
user_id=row.ProcessingLog.user_id, # type: ignore[arg-type]
user_email=mask_email(row.user_email) if row.user_email else None,
)
for row in rows
]
return PaginatedAdminLogsResponse(
items=items, **_admin_paginate(total, page, page_size) # type: ignore[arg-type]
)
+234
View File
@@ -0,0 +1,234 @@
"""
Processing logs and run history endpoints for users.
Users can view the full history of processing runs and per-email logs
for their own mailboxes. Admin equivalents live in admin.py.
"""
from __future__ import annotations
import math
from typing import Optional
from fastapi import APIRouter, Depends, HTTPException, Query, status
from sqlalchemy import select, func
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.database import get_db
from app.core.deps import get_current_active_user
from app.models.database_models import (
MailAccount,
ProcessingLog,
ProcessingRun,
User,
)
from app.models.schemas import (
PaginatedProcessingLogsResponse,
PaginatedProcessingRunsResponse,
ProcessingLogDetailResponse,
ProcessingRunDetailResponse,
)
router = APIRouter()
# ---------------------------------------------------------------------------
# Helper
# ---------------------------------------------------------------------------
def _paginate(total: int, page: int, page_size: int) -> dict:
pages = max(1, math.ceil(total / page_size)) if total else 1
return {"total": total, "page": page, "page_size": page_size, "pages": pages}
# ---------------------------------------------------------------------------
# Processing Runs (all accounts belonging to the current user)
# ---------------------------------------------------------------------------
@router.get(
"/processing-runs",
response_model=PaginatedProcessingRunsResponse,
summary="List processing runs for the current user",
)
async def list_processing_runs(
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
account_id: Optional[int] = Query(None, description="Filter by mail account ID"),
status_filter: Optional[str] = Query(
None,
alias="status",
description="Filter by run status (completed, failed, partial_failure, running)",
),
current_user: User = Depends(get_current_active_user),
db: AsyncSession = Depends(get_db),
):
"""
Return a paginated list of processing runs for all mail accounts owned by
the authenticated user, optionally filtered by account or status.
"""
# Base query: join with MailAccount to enforce ownership
base = (
select(ProcessingRun, MailAccount.name, MailAccount.email_address)
.join(MailAccount, ProcessingRun.mail_account_id == MailAccount.id)
.where(MailAccount.user_id == current_user.id) # type: ignore[arg-type]
)
if account_id is not None:
base = base.where(ProcessingRun.mail_account_id == account_id)
if status_filter:
base = base.where(ProcessingRun.status == status_filter)
# Total count
count_q = select(func.count()).select_from(base.subquery())
total = (await db.execute(count_q)).scalar_one()
offset = (page - 1) * page_size
rows = (
await db.execute(
base.order_by(ProcessingRun.started_at.desc())
.offset(offset)
.limit(page_size)
)
).all()
items = [
ProcessingRunDetailResponse(
id=row.ProcessingRun.id, # type: ignore[arg-type]
mail_account_id=row.ProcessingRun.mail_account_id, # type: ignore[arg-type]
started_at=row.ProcessingRun.started_at, # type: ignore[arg-type]
completed_at=row.ProcessingRun.completed_at, # type: ignore[arg-type]
duration_seconds=row.ProcessingRun.duration_seconds, # type: ignore[arg-type]
emails_fetched=row.ProcessingRun.emails_fetched, # type: ignore[arg-type]
emails_forwarded=row.ProcessingRun.emails_forwarded, # type: ignore[arg-type]
emails_failed=row.ProcessingRun.emails_failed, # type: ignore[arg-type]
status=row.ProcessingRun.status, # type: ignore[arg-type]
error_message=row.ProcessingRun.error_message, # type: ignore[arg-type]
account_name=row.name,
account_email=row.email_address,
)
for row in rows
]
return PaginatedProcessingRunsResponse(
items=items, **_paginate(total, page, page_size) # type: ignore[arg-type]
)
@router.get(
"/processing-runs/{run_id}",
response_model=ProcessingRunDetailResponse,
summary="Get a single processing run",
)
async def get_processing_run(
run_id: int,
current_user: User = Depends(get_current_active_user),
db: AsyncSession = Depends(get_db),
):
"""Return details for a single processing run owned by the current user."""
row = (
await db.execute(
select(ProcessingRun, MailAccount.name, MailAccount.email_address)
.join(MailAccount, ProcessingRun.mail_account_id == MailAccount.id)
.where(
ProcessingRun.id == run_id,
MailAccount.user_id == current_user.id, # type: ignore[arg-type]
)
)
).one_or_none()
if not row:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND, detail="Processing run not found"
)
return ProcessingRunDetailResponse(
id=row.ProcessingRun.id, # type: ignore[arg-type]
mail_account_id=row.ProcessingRun.mail_account_id, # type: ignore[arg-type]
started_at=row.ProcessingRun.started_at, # type: ignore[arg-type]
completed_at=row.ProcessingRun.completed_at, # type: ignore[arg-type]
duration_seconds=row.ProcessingRun.duration_seconds, # type: ignore[arg-type]
emails_fetched=row.ProcessingRun.emails_fetched, # type: ignore[arg-type]
emails_forwarded=row.ProcessingRun.emails_forwarded, # type: ignore[arg-type]
emails_failed=row.ProcessingRun.emails_failed, # type: ignore[arg-type]
status=row.ProcessingRun.status, # type: ignore[arg-type]
error_message=row.ProcessingRun.error_message, # type: ignore[arg-type]
account_name=row.name,
account_email=row.email_address,
)
@router.get(
"/processing-runs/{run_id}/logs",
response_model=PaginatedProcessingLogsResponse,
summary="Get per-email logs for a processing run",
)
async def get_run_logs(
run_id: int,
page: int = Query(1, ge=1),
page_size: int = Query(50, ge=1, le=200),
current_user: User = Depends(get_current_active_user),
db: AsyncSession = Depends(get_db),
):
"""
Return the detailed per-email log entries recorded during a specific
processing run. Ownership is verified by joining with MailAccount.
"""
# Verify the run belongs to this user
run_row = (
await db.execute(
select(ProcessingRun)
.join(MailAccount, ProcessingRun.mail_account_id == MailAccount.id)
.where(
ProcessingRun.id == run_id,
MailAccount.user_id == current_user.id, # type: ignore[arg-type]
)
)
).scalar_one_or_none()
if not run_row:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND, detail="Processing run not found"
)
count_q = select(func.count(ProcessingLog.id)).where(
ProcessingLog.processing_run_id == run_id
)
total = (await db.execute(count_q)).scalar_one()
offset = (page - 1) * page_size
logs = (
(
await db.execute(
select(ProcessingLog)
.where(ProcessingLog.processing_run_id == run_id)
.order_by(ProcessingLog.timestamp.asc())
.offset(offset)
.limit(page_size)
)
)
.scalars()
.all()
)
items = [
ProcessingLogDetailResponse(
id=log.id, # type: ignore[arg-type]
timestamp=log.timestamp, # type: ignore[arg-type]
level=log.level, # type: ignore[arg-type]
message=log.message, # type: ignore[arg-type]
email_subject=log.email_subject, # type: ignore[arg-type]
email_from=log.email_from, # type: ignore[arg-type]
success=log.success, # type: ignore[arg-type]
mail_account_id=log.mail_account_id, # type: ignore[arg-type]
processing_run_id=log.processing_run_id, # type: ignore[arg-type]
email_size_bytes=log.email_size_bytes, # type: ignore[arg-type]
error_details=log.error_details, # type: ignore[arg-type]
)
for log in logs
]
return PaginatedProcessingLogsResponse(
items=items, **_paginate(total, page, page_size) # type: ignore[arg-type]
)
+155 -3
View File
@@ -1,9 +1,10 @@
"""Mail account management endpoints"""
from typing import List
from fastapi import APIRouter, Depends, HTTPException, status
import math
from typing import List, Optional
from fastapi import APIRouter, Depends, HTTPException, Query, status
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, desc
from sqlalchemy import select, desc, func
from app.core.database import get_db
from app.core.deps import get_current_active_user
@@ -11,6 +12,8 @@ from app.core.security import encrypt_credential
from app.models.database_models import (
User,
MailAccount,
ProcessingLog,
ProcessingRun,
AccountStatus,
SubscriptionPlan,
)
@@ -22,6 +25,10 @@ from app.models.schemas import (
MailAccountTestResponse,
MailAccountAutoDetectRequest,
MailAccountAutoDetectResponse,
PaginatedProcessingRunsResponse,
PaginatedProcessingLogsResponse,
ProcessingRunDetailResponse,
ProcessingLogDetailResponse,
)
from app.services.mail_processor import MailProcessor, MailServerAutoDetect
from app.core.config import settings
@@ -274,3 +281,148 @@ async def auto_detect_mail_settings(
return MailAccountAutoDetectResponse(
success=len(suggestions) > 0, suggestions=suggestions
)
# ---------------------------------------------------------------------------
# Per-account processing runs & logs
# ---------------------------------------------------------------------------
@router.get(
"/{account_id}/processing-runs",
response_model=PaginatedProcessingRunsResponse,
summary="List processing runs for a specific mail account",
)
async def list_account_runs(
account_id: int,
page: int = Query(1, ge=1),
page_size: int = Query(20, ge=1, le=100),
current_user: User = Depends(get_current_active_user),
db: AsyncSession = Depends(get_db),
):
"""Return paginated processing runs for a mail account owned by the user."""
result = await db.execute(
select(MailAccount).where(
MailAccount.id == account_id,
MailAccount.user_id == current_user.id,
)
)
account = result.scalar_one_or_none()
if not account:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND, detail="Mail account not found"
)
base = select(ProcessingRun).where(ProcessingRun.mail_account_id == account_id)
total = (
await db.execute(select(func.count()).select_from(base.subquery()))
).scalar_one()
offset = (page - 1) * page_size
runs = (
(
await db.execute(
base.order_by(desc(ProcessingRun.started_at))
.offset(offset)
.limit(page_size)
)
)
.scalars()
.all()
)
pages = max(1, math.ceil(total / page_size)) if total else 1
items = [
ProcessingRunDetailResponse(
id=r.id, # type: ignore[arg-type]
mail_account_id=r.mail_account_id, # type: ignore[arg-type]
started_at=r.started_at, # type: ignore[arg-type]
completed_at=r.completed_at, # type: ignore[arg-type]
duration_seconds=r.duration_seconds, # type: ignore[arg-type]
emails_fetched=r.emails_fetched, # type: ignore[arg-type]
emails_forwarded=r.emails_forwarded, # type: ignore[arg-type]
emails_failed=r.emails_failed, # type: ignore[arg-type]
status=r.status, # type: ignore[arg-type]
error_message=r.error_message, # type: ignore[arg-type]
account_name=account.name, # type: ignore[arg-type]
account_email=account.email_address, # type: ignore[arg-type]
)
for r in runs
]
return PaginatedProcessingRunsResponse(
items=items, total=total, page=page, page_size=page_size, pages=pages
)
@router.get(
"/{account_id}/logs",
response_model=PaginatedProcessingLogsResponse,
summary="List processing logs for a specific mail account",
)
async def list_account_logs(
account_id: int,
page: int = Query(1, ge=1),
page_size: int = Query(50, ge=1, le=200),
level: Optional[str] = Query(
None, description="Filter by log level (INFO, WARNING, ERROR)"
),
current_user: User = Depends(get_current_active_user),
db: AsyncSession = Depends(get_db),
):
"""Return paginated per-email log entries for a mail account owned by the user."""
result = await db.execute(
select(MailAccount).where(
MailAccount.id == account_id,
MailAccount.user_id == current_user.id,
)
)
account = result.scalar_one_or_none()
if not account:
raise HTTPException(
status_code=status.HTTP_404_NOT_FOUND, detail="Mail account not found"
)
base = select(ProcessingLog).where(
ProcessingLog.mail_account_id == account_id,
ProcessingLog.user_id == current_user.id, # type: ignore[arg-type]
)
if level:
base = base.where(ProcessingLog.level == level.upper())
total = (
await db.execute(select(func.count()).select_from(base.subquery()))
).scalar_one()
offset = (page - 1) * page_size
logs = (
(
await db.execute(
base.order_by(ProcessingLog.timestamp.desc())
.offset(offset)
.limit(page_size)
)
)
.scalars()
.all()
)
pages = max(1, math.ceil(total / page_size)) if total else 1
items = [
ProcessingLogDetailResponse(
id=log.id, # type: ignore[arg-type]
timestamp=log.timestamp, # type: ignore[arg-type]
level=log.level, # type: ignore[arg-type]
message=log.message, # type: ignore[arg-type]
email_subject=log.email_subject, # type: ignore[arg-type]
email_from=log.email_from, # type: ignore[arg-type]
success=log.success, # type: ignore[arg-type]
mail_account_id=log.mail_account_id, # type: ignore[arg-type]
processing_run_id=log.processing_run_id, # type: ignore[arg-type]
email_size_bytes=log.email_size_bytes, # type: ignore[arg-type]
error_details=log.error_details, # type: ignore[arg-type]
)
for log in logs
]
return PaginatedProcessingLogsResponse(
items=items, total=total, page=page, page_size=page_size, pages=pages
)