from __future__ import annotations from datetime import UTC, datetime from fastapi import APIRouter, Depends, HTTPException, Query from sqlalchemy import case, func, select from sqlalchemy.orm import Session from app.api.deps import get_db, require_operations_manager from app.models.booking import Booking from app.models.customer import Customer from app.models.data_quality import DataQualityIssue from app.models.inspection import Inspection from app.models.user import User from app.models.vehicle import Vehicle from app.schemas import ( ApplyRecommendedStatusRequest, ApplyRecommendedStatusResult, BulkDataQualityWorkRequest, BulkDataQualityWorkResult, CurrentUser, DataQualityIssueDetailOut, DataQualityIssueOut, DataQualityIssuePageOut, MergeCustomersRequest, MergeCustomersResult, ProvideFieldsRequest, ResolveOdometerRegressionRequest, ResolveOverlapRequest, ScanResultOut, StatusRecommendationOut, VehicleStatusFactsOut, ) from app.services.audit import record_audit_event from app.services.data_quality import ( apply_recommended_status, defer_issue, merge_customers, preview_vehicle_status_recommendation, provide_missing_fields, reject_issue, resolve_booking_overlap, resolve_odometer_regression, run_scan, ) router = APIRouter(prefix="/api/v1/data-quality", tags=["data-quality"]) def _to_out(issue: DataQualityIssue) -> DataQualityIssueOut: assignee = issue.assigned_to_user return DataQualityIssueOut( public_ref=issue.public_ref, rule_type=issue.rule_type, entity_type=issue.entity_type, entity_ref=issue.evidence_json.get("entity_ref", ""), severity=issue.severity, status=issue.status, evidence=issue.evidence_json, detected_at=issue.detected_at, due_at=issue.due_at, assigned_to_ref=assignee.public_ref if assignee else None, assigned_to_name=assignee.display_name if assignee else None, overdue=( issue.status == "open" and issue.due_at is not None and issue.due_at < datetime.now(UTC) ), resolved_at=issue.resolved_at, ) @router.get("/issues", response_model=list[DataQualityIssueOut] | DataQualityIssuePageOut) def list_issues( status: str | None = Query(default=None), rule_type: str | None = Query(default=None), severity: str | None = Query(default=None), assigned_to_ref: str | None = Query(default=None), overdue: bool | None = Query(default=None), demo_only: bool | None = Query(default=None), page: int | None = Query(default=None, ge=1), page_size: int = Query(default=25, ge=1, le=25), db: Session = Depends(get_db), _user: CurrentUser = Depends(require_operations_manager), ) -> list[DataQualityIssueOut] | DataQualityIssuePageOut: severity_order = case( (DataQualityIssue.severity == "high", 0), (DataQualityIssue.severity == "medium", 1), else_=2, ) stmt = select(DataQualityIssue).order_by( DataQualityIssue.due_at.asc().nulls_last(), severity_order, DataQualityIssue.detected_at.desc(), ) if status: stmt = stmt.where(DataQualityIssue.status == status) if rule_type: stmt = stmt.where(DataQualityIssue.rule_type == rule_type) if severity: stmt = stmt.where(DataQualityIssue.severity == severity) if assigned_to_ref == "unassigned": stmt = stmt.where(DataQualityIssue.assigned_to_user_id.is_(None)) elif assigned_to_ref: stmt = stmt.join(DataQualityIssue.assigned_to_user).where( User.public_ref == assigned_to_ref ) if overdue is True: stmt = stmt.where( DataQualityIssue.status == "open", DataQualityIssue.due_at < datetime.now(UTC), ) if demo_only is True: # Server-side so the guided demo scenarios are found on any page, not only the # 25 rows currently loaded in the browser. stmt = stmt.where(DataQualityIssue.public_ref.like("DQ-DEMO-%")) total = db.scalar(select(func.count()).select_from(stmt.subquery())) or 0 page_number = page or 1 issues = db.scalars( stmt if page is None else stmt.offset((page_number - 1) * page_size).limit(page_size) ).all() items = [_to_out(i) for i in issues] if page is None: return items total_pages = max(1, (total + page_size - 1) // page_size) return DataQualityIssuePageOut( items=items, page=min(page_number, total_pages), page_size=page_size, total=total, total_pages=total_pages, ) @router.post("/issues/bulk-work", response_model=BulkDataQualityWorkResult) def update_issue_work_queue( body: BulkDataQualityWorkRequest, db: Session = Depends(get_db), user: CurrentUser = Depends(require_operations_manager), ) -> BulkDataQualityWorkResult: refs = list(dict.fromkeys(body.issue_refs)) if ( body.assigned_to_ref is None and not body.clear_assignment and body.due_at is None and not body.clear_due_at ): raise HTTPException(status_code=422, detail="No work queue change was requested") if body.assigned_to_ref is not None and body.clear_assignment: raise HTTPException(status_code=422, detail="Choose an assignee or clear assignment") if body.due_at is not None and body.clear_due_at: raise HTTPException(status_code=422, detail="Choose a due date or clear the due date") if body.due_at is not None and body.due_at.tzinfo is None: raise HTTPException(status_code=422, detail="Due date must include a timezone") assignee = None if body.assigned_to_ref is not None: assignee = db.scalar( select(User).where( User.public_ref == body.assigned_to_ref, User.active.is_(True), ) ) if assignee is None: raise HTTPException(status_code=422, detail="Active assignee not found") issues = list( db.scalars( select(DataQualityIssue).where(DataQualityIssue.public_ref.in_(refs)).with_for_update() ).all() ) if len(issues) != len(refs): found = {issue.public_ref for issue in issues} missing = next(ref for ref in refs if ref not in found) raise HTTPException(status_code=404, detail=f"Data quality issue {missing} not found") for issue in issues: if issue.status != "open": raise HTTPException( status_code=409, detail=f"Data quality issue {issue.public_ref} is not open", ) before = { "assigned_to_ref": issue.assigned_to_user.public_ref if issue.assigned_to_user else None, "due_at": issue.due_at.isoformat() if issue.due_at else None, } if body.assigned_to_ref is not None: issue.assigned_to_user = assignee elif body.clear_assignment: issue.assigned_to_user = None if body.due_at is not None: issue.due_at = body.due_at elif body.clear_due_at: issue.due_at = None after = { "assigned_to_ref": assignee.public_ref if body.assigned_to_ref is not None and assignee else (None if body.clear_assignment else before["assigned_to_ref"]), "due_at": issue.due_at.isoformat() if issue.due_at else None, } record_audit_event( db, actor_type="user", actor_label=user.display_name, action="data_quality_work_updated", entity_type="data_quality_issue", entity_id=issue.id, before=before, after=after, ) db.commit() for issue in issues: db.refresh(issue) return BulkDataQualityWorkResult(updated=[_to_out(issue) for issue in issues]) # Every public reference in this system carries its entity type in its own prefix # (CUS-/MO-/BK-/INSP-/DQ-). Related-entity typing is resolved from the reference itself, # not guessed from the issue's rule_type -- a booking_overlap issue's related refs are # bookings, not vehicles, and an inline odometer_regression issue's related refs mix a # booking and an inspection ref in the same list. _PREFIX_TO_TYPE = { "CUS-": "customer", "MO-": "vehicle", "BK-": "booking", "INSP-": "inspection", } def _entity_type_for_ref(ref: str) -> str | None: for prefix, entity_type in _PREFIX_TO_TYPE.items(): if ref.startswith(prefix): return entity_type return None def _snapshot(entity_type: str, ref: str, db: Session) -> dict | None: if entity_type == "customer": customer = db.scalar(select(Customer).where(Customer.public_ref == ref)) if customer is None: return None return { "entity_type": "customer", "public_ref": customer.public_ref, "first_name": customer.first_name, "last_name": customer.last_name, "email": customer.email, "phone": customer.phone, "postal_code": customer.postal_code, "city": customer.city, } if entity_type == "vehicle": vehicle = db.scalar(select(Vehicle).where(Vehicle.public_ref == ref)) if vehicle is None: return None return { "entity_type": "vehicle", "public_ref": vehicle.public_ref, "registration_number": vehicle.registration_number, "make": vehicle.make, "model": vehicle.model, "location": vehicle.location, "operational_status": vehicle.operational_status, "odometer_km": vehicle.odometer_km, } if entity_type == "booking": booking = db.scalar(select(Booking).where(Booking.public_ref == ref)) if booking is None: return None vehicle = db.get(Vehicle, booking.vehicle_id) customer = db.get(Customer, booking.customer_id) return { "entity_type": "booking", "public_ref": booking.public_ref, "status": booking.status, "starts_at": booking.starts_at.isoformat(), "ends_at": booking.ends_at.isoformat(), "vehicle_ref": vehicle.public_ref if vehicle else None, "customer_ref": customer.public_ref if customer else None, "end_odometer_km": booking.end_odometer_km, } if entity_type == "inspection": inspection = db.scalar(select(Inspection).where(Inspection.public_ref == ref)) if inspection is None: return None booking = db.get(Booking, inspection.booking_id) return { "entity_type": "inspection", "public_ref": inspection.public_ref, "type": inspection.type, "odometer_km": inspection.odometer_km, "completed_at": inspection.completed_at.isoformat(), "booking_ref": booking.public_ref if booking else None, } return None @router.get("/issues/{public_ref}", response_model=DataQualityIssueDetailOut) def get_issue( public_ref: str, db: Session = Depends(get_db), _user: CurrentUser = Depends(require_operations_manager), ) -> DataQualityIssueDetailOut: issue = db.scalar(select(DataQualityIssue).where(DataQualityIssue.public_ref == public_ref)) if issue is None: raise HTTPException(status_code=404, detail="Data quality issue not found") base = _to_out(issue) related_refs = issue.evidence_json.get("related_refs", []) related_snapshots = [] for ref in related_refs: entity_type = _entity_type_for_ref(ref) if entity_type is None: continue snap = _snapshot(entity_type, ref, db) if snap is not None: related_snapshots.append(snap) return DataQualityIssueDetailOut( **base.model_dump(), entity_snapshot=_snapshot(issue.entity_type, base.entity_ref, db), related_snapshots=related_snapshots, ) @router.post("/issues/{public_ref}/defer", response_model=DataQualityIssueOut) def defer( public_ref: str, db: Session = Depends(get_db), user: CurrentUser = Depends(require_operations_manager), ) -> DataQualityIssueOut: issue = defer_issue(db, public_ref, user) return _to_out(issue) @router.post("/issues/{public_ref}/reject", response_model=DataQualityIssueOut) def reject( public_ref: str, db: Session = Depends(get_db), user: CurrentUser = Depends(require_operations_manager), ) -> DataQualityIssueOut: issue = reject_issue(db, public_ref, user) return _to_out(issue) @router.post("/issues/{public_ref}/merge-customers", response_model=MergeCustomersResult) def merge( public_ref: str, body: MergeCustomersRequest, db: Session = Depends(get_db), user: CurrentUser = Depends(require_operations_manager), ) -> MergeCustomersResult: result = merge_customers(db, public_ref, body.survivor_ref, body.field_overrides, user) return MergeCustomersResult(**result) @router.post("/issues/{public_ref}/provide-fields", response_model=DataQualityIssueOut) def provide_fields( public_ref: str, body: ProvideFieldsRequest, db: Session = Depends(get_db), user: CurrentUser = Depends(require_operations_manager), ) -> DataQualityIssueOut: issue = provide_missing_fields(db, public_ref, body.fields, user) return _to_out(issue) @router.post("/issues/{public_ref}/resolve-odometer-regression", response_model=DataQualityIssueOut) def resolve_odometer( public_ref: str, body: ResolveOdometerRegressionRequest, db: Session = Depends(get_db), user: CurrentUser = Depends(require_operations_manager), ) -> DataQualityIssueOut: issue = resolve_odometer_regression(db, public_ref, body, user) return _to_out(issue) @router.post("/issues/{public_ref}/resolve-overlap", response_model=DataQualityIssueOut) def resolve_overlap( public_ref: str, body: ResolveOverlapRequest, db: Session = Depends(get_db), user: CurrentUser = Depends(require_operations_manager), ) -> DataQualityIssueOut: issue = resolve_booking_overlap(db, public_ref, body.booking_ref, body.note, user) return _to_out(issue) @router.post("/issues/{public_ref}/status-recommendation", response_model=StatusRecommendationOut) def status_recommendation( public_ref: str, db: Session = Depends(get_db), user: CurrentUser = Depends(require_operations_manager), ) -> StatusRecommendationOut: """Non-mutating preview: computes the recommendation without changing anything, resolving no issue and writing no audit event. Safe to call repeatedly.""" _issue, _vehicle, recommendation, token = preview_vehicle_status_recommendation(db, public_ref) return StatusRecommendationOut( current_status=recommendation.current_status, recommended_status=recommendation.recommended_status, recommendation_code=recommendation.recommendation_code, safe_to_apply=recommendation.safe_to_apply, manual_review_required=recommendation.manual_review_required, facts=VehicleStatusFactsOut(**recommendation.facts.as_dict()), blocking_reasons=recommendation.blocking_reasons, recommendation_token=token, ) @router.post( "/issues/{public_ref}/apply-recommended-status", response_model=ApplyRecommendedStatusResult ) def apply_status( public_ref: str, body: ApplyRecommendedStatusRequest, db: Session = Depends(get_db), user: CurrentUser = Depends(require_operations_manager), ) -> ApplyRecommendedStatusResult: issue, applied_status, reason_code = apply_recommended_status( db, public_ref, user, body.recommendation_token ) return ApplyRecommendedStatusResult( issue=_to_out(issue), applied_status=applied_status, reason_code=reason_code ) @router.post("/scan", response_model=ScanResultOut) def scan( db: Session = Depends(get_db), user: CurrentUser = Depends(require_operations_manager), ) -> ScanResultOut: result = run_scan(db, actor_label=user.display_name, actor_type="user") return ScanResultOut(created=result.created)