from __future__ import annotations import uuid from dataclasses import dataclass, field from datetime import UTC, datetime, timedelta from sqlalchemy import func, select, update from sqlalchemy.orm import Session from app.core.errors import AppError from app.models.booking import Booking from app.models.customer import Customer from app.models.data_quality import DataQualityIssue from app.models.vehicle import Vehicle from app.schemas import CurrentUser, ResolveOdometerRegressionRequest from app.services.audit import record_audit_event from app.services.data_quality_duplicate_scan import scan_duplicate_customers from app.services.vehicle_status import ( RECOMMENDATION_CODE_NO_CONFLICT, VehicleStatusRecommendation, compute_recommendation_token, evaluate_vehicle_status, gather_vehicle_status_facts, ) REQUIRED_CUSTOMER_FIELDS = ("first_name", "last_name") REQUIRED_VEHICLE_FIELDS = ("registration_number", "make", "model", "location") DATA_QUALITY_SCAN_LOCK_ID = 6_138_493_717_091_029_491 def issue_due_at(detected_at: datetime, severity: str) -> datetime: """Return the local operational SLA deadline for a newly detected issue.""" return detected_at + { "high": timedelta(hours=4), "medium": timedelta(days=1), "low": timedelta(days=3), }.get(severity, timedelta(days=1)) @dataclass class ScanResult: created: dict[str, int] = field(default_factory=dict) def bump(self, rule_type: str) -> None: self.created[rule_type] = self.created.get(rule_type, 0) + 1 def _has_open_issue(db: Session, rule_type: str, entity_type: str, entity_id: uuid.UUID) -> bool: return ( db.scalar( select(DataQualityIssue.id).where( DataQualityIssue.rule_type == rule_type, DataQualityIssue.entity_type == entity_type, DataQualityIssue.entity_id == entity_id, DataQualityIssue.status == "open", ) ) is not None ) def _new_scan_ref(prefix: str) -> str: """Generate a stable human-readable prefix with a concurrent-safe suffix.""" return f"{prefix}-{uuid.uuid4().hex[:10].upper()}" def _open_issue( db: Session, scan: ScanResult, *, rule_type: str, entity_type: str, entity_id: uuid.UUID, severity: str, summary: str, entity_ref: str, related_refs: list[str], signals: list[dict] | None = None, ) -> None: if _has_open_issue(db, rule_type, entity_type, entity_id): return now = datetime.now(UTC) # Reintroduced evidence creates a new issue rather than silently reopening the old # one, but it stays linked to whatever decision was made last time so an operator # doesn't re-litigate from a blank slate. previous = db.scalar( select(DataQualityIssue) .where( DataQualityIssue.rule_type == rule_type, DataQualityIssue.entity_type == entity_type, DataQualityIssue.entity_id == entity_id, DataQualityIssue.status != "open", ) .order_by(DataQualityIssue.detected_at.desc()) ) # `summary` is kept as a technical-fallback string (shown only under "Technical # details"); `signals` is the stable, localizable structure the frontend renders as # the primary evidence -- see docs/fleet-ops-correction/current-gap-audit.md §2/§6. evidence: dict = { "summary": summary, "entity_ref": entity_ref, "related_refs": related_refs, "signals": signals or [], } if previous is not None: evidence["reopened_from"] = previous.public_ref evidence["previous_decision"] = previous.status issue = DataQualityIssue( public_ref=_new_scan_ref("DQ-SCAN"), rule_type=rule_type, entity_type=entity_type, entity_id=entity_id, severity=severity, status="open", evidence_json=evidence, proposed_action_json={}, detected_at=now, due_at=issue_due_at(now, severity), ) db.add(issue) db.flush() scan.bump(rule_type) def _scan_missing_required_fields(db: Session, scan: ScanResult) -> None: # Anonymised customers have had their contact data removed on purpose; flagging # them as "missing required field" would only be resolvable by re-entering PII. for customer in db.scalars( select(Customer).where( Customer.merged_into_customer_id.is_(None), Customer.anonymized_at.is_(None), ) ).all(): missing = [f for f in REQUIRED_CUSTOMER_FIELDS if not getattr(customer, f)] if not customer.email and not customer.phone: missing.append("email_or_phone") if missing: _open_issue( db, scan, rule_type="missing_required_field", entity_type="customer", entity_id=customer.id, severity="low", summary=f"Missing: {', '.join(missing)}", entity_ref=customer.public_ref, related_refs=[], signals=[{"code": "missing_field", "params": {"field": f}} for f in missing], ) for vehicle in db.scalars(select(Vehicle).where(Vehicle.active.is_(True))).all(): missing = [f for f in REQUIRED_VEHICLE_FIELDS if not getattr(vehicle, f)] if missing: _open_issue( db, scan, rule_type="missing_required_field", entity_type="vehicle", entity_id=vehicle.id, severity="low", summary=f"Missing: {', '.join(missing)}", entity_ref=vehicle.public_ref, related_refs=[], signals=[{"code": "missing_field", "params": {"field": f}} for f in missing], ) def _scan_booking_overlaps(db: Session, scan: ScanResult) -> None: vehicles = db.scalars(select(Vehicle)).all() bookings_by_vehicle: dict[uuid.UUID, list[Booking]] = {} for booking in db.scalars( select(Booking).where(Booking.status.in_(["reserved", "active"])) ).all(): bookings_by_vehicle.setdefault(booking.vehicle_id, []).append(booking) vehicle_by_id = {v.id: v for v in vehicles} for vehicle_id, bookings in bookings_by_vehicle.items(): bookings.sort(key=lambda b: b.starts_at) for i, first in enumerate(bookings): for second in bookings[i + 1 :]: if second.starts_at < first.ends_at and first.starts_at < second.ends_at: vehicle = vehicle_by_id[vehicle_id] _open_issue( db, scan, rule_type="booking_overlap", entity_type="vehicle", entity_id=vehicle_id, severity="high", summary=f"Overlapping bookings {first.public_ref} and {second.public_ref}", entity_ref=vehicle.public_ref, related_refs=[first.public_ref, second.public_ref], signals=[ { "code": "overlap.reserved_bookings", "params": {"refs": [first.public_ref, second.public_ref]}, } ], ) def _scan_vehicle_status_conflicts(db: Session, scan: ScanResult) -> None: # Uses the same shared evaluator as the preview/apply flow (app.services.vehicle_status) # so detection and resolution can never structurally disagree -- see # docs/fleet-ops-correction/vehicle-status-decision-table.md. vehicles = db.scalars(select(Vehicle)).all() for vehicle in vehicles: facts = gather_vehicle_status_facts(db, vehicle) recommendation = evaluate_vehicle_status(vehicle, facts) if recommendation.recommendation_code == RECOMMENDATION_CODE_NO_CONFLICT: continue signals = [{"code": recommendation.recommendation_code, "params": facts.as_dict()}] summary = ( f"Recommended status: {recommendation.recommended_status}" if recommendation.recommended_status else "Manual review required: active rental conflicts with a blocking condition" ) _open_issue( db, scan, rule_type="vehicle_status_conflict", entity_type="vehicle", entity_id=vehicle.id, severity="high", summary=summary, entity_ref=vehicle.public_ref, related_refs=[ *facts.active_booking_refs, *(ref for pair in facts.overlapping_booking_pairs for ref in pair), ], signals=signals, ) def _scan_odometer_regressions(db: Session, scan: ScanResult) -> None: # The seed dataset's vehicle.odometer_km is generated independently of booking # history, so comparing every historical booking against it produces near-universal # false positives. Instead check the booking sequence's own internal consistency: # each vehicle's completed bookings should show a non-decreasing odometer reading. vehicles = {v.id: v for v in db.scalars(select(Vehicle)).all()} bookings_by_vehicle: dict[uuid.UUID, list[Booking]] = {} for booking in db.scalars( select(Booking).where(Booking.status == "returned", Booking.end_odometer_km.is_not(None)) ).all(): bookings_by_vehicle.setdefault(booking.vehicle_id, []).append(booking) for vehicle_id, bookings in bookings_by_vehicle.items(): bookings.sort(key=lambda b: b.ends_at) for earlier, later in zip(bookings, bookings[1:], strict=False): # The query above filters end_odometer_km IS NOT NULL, so both are ints here. assert earlier.end_odometer_km is not None assert later.end_odometer_km is not None if later.end_odometer_km < earlier.end_odometer_km: vehicle = vehicles[vehicle_id] _open_issue( db, scan, rule_type="odometer_regression", entity_type="vehicle", entity_id=vehicle_id, severity="medium", summary=( f"Booking {later.public_ref} recorded {later.end_odometer_km} km, " f"below the {earlier.end_odometer_km} km recorded by earlier " f"booking {earlier.public_ref}." ), entity_ref=vehicle.public_ref, related_refs=[earlier.public_ref, later.public_ref], signals=[ { "code": "odometer.regression", "params": { "later_ref": later.public_ref, "later_km": later.end_odometer_km, "earlier_ref": earlier.public_ref, "earlier_km": earlier.end_odometer_km, }, } ], ) break def run_scan( db: Session, *, actor_label: str | None = None, actor_type: str = "user" ) -> ScanResult: # The check-then-insert work below spans several rules. Serialise whole scans at the # database boundary so API and n8n triggers cannot both observe an empty condition. db.scalar(select(func.pg_advisory_xact_lock(DATA_QUALITY_SCAN_LOCK_ID))) scan = ScanResult() scan_duplicate_customers(db, scan, _open_issue) _scan_missing_required_fields(db, scan) _scan_odometer_regressions(db, scan) _scan_booking_overlaps(db, scan) _scan_vehicle_status_conflicts(db, scan) if actor_label is not None: record_audit_event( db, actor_type=actor_type, actor_label=actor_label, action="data_quality_scan_run", entity_type="system", metadata={"created": scan.created}, ) db.commit() return scan def _load_open_issue( db: Session, public_ref: str, *, lock: bool = True ) -> DataQualityIssue: statement = select(DataQualityIssue).where(DataQualityIssue.public_ref == public_ref) if lock: statement = statement.with_for_update() issue = db.scalar(statement) if issue is None: raise AppError("ISSUE_NOT_FOUND", "Data quality issue not found.", status_code=404) if issue.status != "open": raise AppError( "ISSUE_NOT_OPEN", f"Issue is '{issue.status}', not 'open'.", status_code=409, ) return issue def defer_issue(db: Session, public_ref: str, actor: CurrentUser) -> DataQualityIssue: issue = _load_open_issue(db, public_ref) issue.status = "deferred" issue.resolved_at = datetime.now(UTC) issue.resolved_by = actor.display_name record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="data_quality_issue_deferred", entity_type="data_quality_issue", entity_id=issue.id, before={"status": "open"}, after={"status": "deferred"}, ) db.commit() return issue def reject_issue(db: Session, public_ref: str, actor: CurrentUser) -> DataQualityIssue: issue = _load_open_issue(db, public_ref) issue.status = "rejected" issue.resolved_at = datetime.now(UTC) issue.resolved_by = actor.display_name record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="data_quality_issue_rejected", entity_type="data_quality_issue", entity_id=issue.id, before={"status": "open"}, after={"status": "rejected"}, ) db.commit() return issue def provide_missing_fields( db: Session, public_ref: str, fields: dict[str, str], actor: CurrentUser ) -> DataQualityIssue: issue = _load_open_issue(db, public_ref) if issue.rule_type != "missing_required_field": raise AppError( "NOT_A_MISSING_FIELD_ISSUE", "This issue is not a missing-required-field issue.", status_code=409, ) entity: Customer | Vehicle | None if issue.entity_type == "customer": entity = db.get(Customer, issue.entity_id) allowed = {*REQUIRED_CUSTOMER_FIELDS, "email", "phone"} elif issue.entity_type == "vehicle": entity = db.get(Vehicle, issue.entity_id) allowed = set(REQUIRED_VEHICLE_FIELDS) else: raise AppError( "UNSUPPORTED_ENTITY", f"Cannot provide fields for entity type '{issue.entity_type}'.", status_code=409, ) if entity is None: raise AppError( "ENTITY_NOT_FOUND", "The underlying record could not be found.", status_code=404 ) invalid = set(fields) - allowed if invalid: raise AppError( "INVALID_FIELD", f"Fields not permitted here: {', '.join(sorted(invalid))}.", status_code=422, ) if not fields: raise AppError( "NO_FIELDS_PROVIDED", "At least one field must be provided.", status_code=422 ) before = {f: getattr(entity, f) for f in allowed} for field_name, value in fields.items(): if not value.strip(): raise AppError("EMPTY_VALUE", f"Field '{field_name}' cannot be blank.", status_code=422) setattr(entity, field_name, value.strip()) after = {f: getattr(entity, f) for f in allowed} correlation_id = uuid.uuid4() record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="data_quality_fields_provided", entity_type=issue.entity_type, entity_id=entity.id, correlation_id=correlation_id, before=before, after=after, metadata={"issue_ref": issue.public_ref}, ) if isinstance(entity, Customer): missing = [f for f in REQUIRED_CUSTOMER_FIELDS if not getattr(entity, f)] if not entity.email and not entity.phone: missing.append("email_or_phone") else: missing = [f for f in REQUIRED_VEHICLE_FIELDS if not getattr(entity, f)] if not missing: issue.status = "resolved" issue.resolved_at = datetime.now(UTC) issue.resolved_by = actor.display_name record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="data_quality_issue_resolved", entity_type="data_quality_issue", entity_id=issue.id, correlation_id=correlation_id, before={"status": "open"}, after={"status": "resolved"}, ) else: issue.evidence_json = {**issue.evidence_json, "summary": f"Missing: {', '.join(missing)}"} db.commit() return issue def resolve_odometer_regression( db: Session, public_ref: str, body: ResolveOdometerRegressionRequest, actor: CurrentUser ) -> DataQualityIssue: issue = _load_open_issue(db, public_ref) if issue.rule_type != "odometer_regression": raise AppError( "NOT_AN_ODOMETER_ISSUE", "This issue is not an odometer_regression issue.", status_code=409, ) # Lock order is booking -> vehicle everywhere (checkout, return, reschedule); taking # the vehicle lock first here would be a deadlock waiting to happen under concurrency. booking: Booking | None = None if body.decision != "retain_canonical": related_refs = issue.evidence_json.get("related_refs", []) if body.booking_ref not in related_refs: raise AppError( "INVALID_BOOKING_REFERENCE", "booking_ref must be one of this issue's related bookings.", status_code=422, ) if body.corrected_odometer_km is None: raise AppError( "CORRECTED_VALUE_REQUIRED", "corrected_odometer_km is required when correcting a reading.", status_code=422, ) booking = db.scalar( select(Booking).where(Booking.public_ref == body.booking_ref).with_for_update() ) if booking is None: raise AppError( "BOOKING_NOT_FOUND", "The booking to correct was not found.", status_code=404 ) vehicle = db.scalar(select(Vehicle).where(Vehicle.id == issue.entity_id).with_for_update()) if vehicle is None: raise AppError( "VEHICLE_NOT_FOUND", "The vehicle for this issue was not found.", status_code=404 ) correlation_id = uuid.uuid4() if body.decision == "retain_canonical": record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="data_quality_odometer_retained", entity_type="vehicle", entity_id=vehicle.id, correlation_id=correlation_id, metadata={"issue_ref": issue.public_ref, "canonical_odometer_km": vehicle.odometer_km}, ) else: assert booking is not None and body.corrected_odometer_km is not None # Never silently lower the canonical odometer: a correction must be at or above # the current canonical value, otherwise it would just create a new regression. if body.corrected_odometer_km < vehicle.odometer_km: raise AppError( "CORRECTION_BELOW_CANONICAL", ( f"Corrected value {body.corrected_odometer_km} km is still below the " f"canonical {vehicle.odometer_km} km; it would not resolve the regression." ), status_code=422, ) before = { "booking_end_odometer_km": booking.end_odometer_km, "vehicle_odometer_km": vehicle.odometer_km, } booking.end_odometer_km = body.corrected_odometer_km vehicle.odometer_km = body.corrected_odometer_km vehicle.version += 1 record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="data_quality_odometer_corrected", entity_type="vehicle", entity_id=vehicle.id, correlation_id=correlation_id, before=before, after={ "booking_end_odometer_km": booking.end_odometer_km, "vehicle_odometer_km": vehicle.odometer_km, }, metadata={"issue_ref": issue.public_ref, "booking_ref": booking.public_ref}, ) issue.status = "resolved" issue.resolved_at = datetime.now(UTC) issue.resolved_by = actor.display_name record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="data_quality_issue_resolved", entity_type="data_quality_issue", entity_id=issue.id, correlation_id=correlation_id, before={"status": "open"}, after={"status": "resolved"}, metadata={"decision": body.decision, "note": body.note}, ) db.commit() return issue def resolve_booking_overlap( db: Session, public_ref: str, booking_ref: str, note: str | None, actor: CurrentUser ) -> DataQualityIssue: issue = _load_open_issue(db, public_ref) if issue.rule_type != "booking_overlap": raise AppError( "NOT_AN_OVERLAP_ISSUE", "This issue is not a booking_overlap issue.", status_code=409 ) related_refs = issue.evidence_json.get("related_refs", []) if booking_ref not in related_refs: raise AppError( "INVALID_BOOKING_REFERENCE", "booking_ref must be one of this issue's overlapping bookings.", status_code=422, ) booking = db.scalar(select(Booking).where(Booking.public_ref == booking_ref).with_for_update()) if booking is None: raise AppError("BOOKING_NOT_FOUND", "The booking to block was not found.", status_code=404) if booking.status not in ("reserved", "active"): raise AppError( "BOOKING_NOT_ACTIVE", f"Booking is '{booking.status}'; only a reserved or active booking can be blocked.", status_code=409, ) before = {"status": booking.status} booking.status = "blocked" # Verify the minimal safe resolution actually removed the conflict: no two # reserved/active bookings for this vehicle should still overlap. The session has # autoflush disabled, so exclude the just-blocked booking by id rather than relying # on the in-memory status change being visible to this query. remaining = db.scalars( select(Booking).where( Booking.vehicle_id == booking.vehicle_id, Booking.status.in_(["reserved", "active"]), Booking.public_ref.in_(related_refs), Booking.id != booking.id, ) ).all() for i, first in enumerate(remaining): for second in remaining[i + 1 :]: if second.starts_at < first.ends_at and first.starts_at < second.ends_at: raise AppError( "OVERLAP_STILL_PRESENT", "Blocking this booking did not remove the overlap; another commitment remains.", status_code=409, ) correlation_id = uuid.uuid4() record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="data_quality_booking_blocked", entity_type="booking", entity_id=booking.id, correlation_id=correlation_id, before=before, after={"status": booking.status}, metadata={"issue_ref": issue.public_ref, "note": note}, ) issue.status = "resolved" issue.resolved_at = datetime.now(UTC) issue.resolved_by = actor.display_name record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="data_quality_issue_resolved", entity_type="data_quality_issue", entity_id=issue.id, correlation_id=correlation_id, before={"status": "open"}, after={"status": "resolved"}, ) db.commit() return issue def _load_vehicle_status_conflict_issue( db: Session, public_ref: str, *, lock: bool = True ) -> DataQualityIssue: issue = _load_open_issue(db, public_ref, lock=lock) if issue.rule_type != "vehicle_status_conflict": raise AppError( "NOT_A_STATUS_CONFLICT_ISSUE", "This issue is not a vehicle_status_conflict issue.", status_code=409, ) return issue def preview_vehicle_status_recommendation( db: Session, public_ref: str ) -> tuple[DataQualityIssue, Vehicle, VehicleStatusRecommendation, str]: """Non-mutating: computes and returns the recommendation only. Never resolves the issue, never writes an audit event, never queues automation -- safe to call as often as the UI needs (e.g. every time the panel is opened) with zero side effects.""" issue = _load_vehicle_status_conflict_issue(db, public_ref, lock=False) vehicle = db.scalar(select(Vehicle).where(Vehicle.id == issue.entity_id)) if vehicle is None: raise AppError( "VEHICLE_NOT_FOUND", "The vehicle for this issue was not found.", status_code=404 ) facts = gather_vehicle_status_facts(db, vehicle, exclude_issue_id=issue.id) recommendation = evaluate_vehicle_status(vehicle, facts) token = compute_recommendation_token(vehicle, facts) return issue, vehicle, recommendation, token def apply_recommended_status( db: Session, public_ref: str, actor: CurrentUser, expected_token: str ) -> tuple[DataQualityIssue, str, str]: issue = _load_vehicle_status_conflict_issue(db, public_ref) # Lock the vehicle row for the remainder of this transaction so a concurrent apply # (or return/checkout) can't race between our fact-gathering and the write below. vehicle = db.scalar(select(Vehicle).where(Vehicle.id == issue.entity_id).with_for_update()) if vehicle is None: raise AppError( "VEHICLE_NOT_FOUND", "The vehicle for this issue was not found.", status_code=404 ) facts = gather_vehicle_status_facts(db, vehicle, exclude_issue_id=issue.id) recommendation = evaluate_vehicle_status(vehicle, facts) current_token = compute_recommendation_token(vehicle, facts) if current_token != expected_token: raise AppError( "RECOMMENDATION_STALE", "The underlying facts changed since this recommendation was shown; " "review the recommendation again before applying it.", status_code=409, ) if recommendation.manual_review_required or not recommendation.safe_to_apply: raise AppError( "MANUAL_REVIEW_REQUIRED", "This vehicle's state requires manual review; no automatic status change is safe.", status_code=409, ) if recommendation.recommended_status is None: raise AppError( "NO_CONFLICT_DETECTED", "The current vehicle state no longer conflicts; nothing to apply.", status_code=409, ) new_status = recommendation.recommended_status reason_code = recommendation.recommendation_code before = {"operational_status": vehicle.operational_status} vehicle.operational_status = new_status vehicle.version += 1 # Re-validate against the same shared evaluator, over freshly-gathered facts, that # applying this change actually leaves no conflict -- never trust the pre-computed # recommendation alone for the post-condition. post_facts = gather_vehicle_status_facts(db, vehicle, exclude_issue_id=issue.id) post_check = evaluate_vehicle_status(vehicle, post_facts) if post_check.recommendation_code not in (RECOMMENDATION_CODE_NO_CONFLICT,): raise AppError( "CONFLICT_STILL_PRESENT", "Applying the recommended status did not resolve the conflict.", status_code=409, ) correlation_id = uuid.uuid4() record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="data_quality_status_applied", entity_type="vehicle", entity_id=vehicle.id, correlation_id=correlation_id, before=before, after={"operational_status": vehicle.operational_status}, metadata={"issue_ref": issue.public_ref, "reason_code": reason_code}, ) issue.status = "resolved" issue.resolved_at = datetime.now(UTC) issue.resolved_by = actor.display_name record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="data_quality_issue_resolved", entity_type="data_quality_issue", entity_id=issue.id, correlation_id=correlation_id, before={"status": "open"}, after={"status": "resolved"}, ) db.commit() return issue, new_status, reason_code MERGEABLE_FIELDS = ("first_name", "last_name", "email", "phone", "postal_code", "city") # Mirrors the column lengths in app/models/customer.py so an override can never fail with # a database DataError (500) instead of a validation error. _MERGEABLE_FIELD_MAX_LENGTH = { "first_name": 80, "last_name": 80, "email": 200, "phone": 40, "postal_code": 20, "city": 120, } def merge_customers( db: Session, public_ref: str, survivor_ref: str, field_overrides: dict[str, str] | None, actor: CurrentUser, ) -> dict: issue = _load_open_issue(db, public_ref) if issue.rule_type != "possible_duplicate_customer": raise AppError( "NOT_A_DUPLICATE_ISSUE", "This issue is not a possible-duplicate-customer issue.", status_code=409, ) entity_ref = issue.evidence_json.get("entity_ref") related_refs = issue.evidence_json.get("related_refs", []) candidate_refs = {entity_ref, *related_refs} if survivor_ref not in candidate_refs: raise AppError( "INVALID_SURVIVOR", "The survivor reference must be one of the two customers in this issue.", status_code=422, details={"candidates": sorted(candidate_refs)}, ) loser_ref = next(ref for ref in candidate_refs if ref != survivor_ref) # Lock both rows in a deterministic order (by public_ref) so two concurrent merges # touching the same customers serialise instead of deadlocking or double-merging. survivor = None loser = None for ref in sorted((survivor_ref, loser_ref)): customer = db.scalar(select(Customer).where(Customer.public_ref == ref).with_for_update()) if ref == survivor_ref: survivor = customer else: loser = customer if survivor is None or loser is None: raise AppError( "CUSTOMER_NOT_FOUND", "One of the customers could not be found.", status_code=404 ) if survivor.merged_into_customer_id is not None or loser.merged_into_customer_id is not None: raise AppError( "CUSTOMER_ALREADY_MERGED", "One of the customers has already been merged into another record.", status_code=409, ) before = { "survivor": {f: getattr(survivor, f) for f in MERGEABLE_FIELDS}, "loser": {f: getattr(loser, f) for f in MERGEABLE_FIELDS}, } for field_name, value in (field_overrides or {}).items(): if field_name not in MERGEABLE_FIELDS: raise AppError( "INVALID_FIELD_OVERRIDE", f"Field '{field_name}' cannot be merged.", status_code=422 ) cleaned = value.strip() if isinstance(value, str) else value max_length = _MERGEABLE_FIELD_MAX_LENGTH[field_name] if not cleaned or len(cleaned) > max_length: raise AppError( "INVALID_FIELD_OVERRIDE", f"Field '{field_name}' must be 1 to {max_length} characters.", status_code=422, ) setattr(survivor, field_name, cleaned) rewired = db.execute( update(Booking).where(Booking.customer_id == loser.id).values(customer_id=survivor.id) ) rewired_count: int = rewired.rowcount # type: ignore[attr-defined] loser.merged_into_customer_id = survivor.id issue.status = "resolved" issue.resolved_at = datetime.now(UTC) issue.resolved_by = actor.display_name record_audit_event( db, actor_type="user", actor_label=actor.display_name, action="customer_merged", entity_type="customer", entity_id=survivor.id, before=before, after={"survivor": {f: getattr(survivor, f) for f in MERGEABLE_FIELDS}}, metadata={ "loser_ref": loser_ref, "survivor_ref": survivor_ref, "rewired_bookings": rewired_count, }, ) db.commit() return { "issue_ref": issue.public_ref, "survivor_ref": survivor_ref, "loser_ref": loser_ref, "rewired_bookings": rewired_count, }