61 lines
4.9 KiB
Python
61 lines
4.9 KiB
Python
"""Repair verified mileage provenance and offline baselines. Dry-run by default."""
|
|
import argparse,datetime,json,pathlib
|
|
from decimal import Decimal
|
|
from db_access import connect
|
|
p=argparse.ArgumentParser();p.add_argument('--apply',action='store_true');args=p.parse_args()
|
|
root=pathlib.Path('/opt/lingniu-go-native/backups/mileage-reconciliation-20260916')
|
|
legacy="('legacy-mysql.lingniu-prod','manual-lingniu-prod-day-mileage')"
|
|
c=connect();report={}
|
|
with c.cursor() as q:
|
|
q.execute('START TRANSACTION READ ONLY')
|
|
q.execute("SELECT COUNT(*) AS n,COUNT(DISTINCT vin) AS vehicles FROM vehicle_daily_mileage_source WHERE protocol='GB32960' AND source_ip IN "+legacy);report['legacy']=q.fetchone()
|
|
q.execute("SELECT COUNT(*) AS n FROM vehicle_daily_mileage m JOIN vehicle_data_source d ON d.id=m.source_id WHERE m.protocol='GB32960' AND d.source_ip IN "+legacy);report['legacy_projections']=q.fetchone()
|
|
q.execute("""SELECT s.vin,s.stat_date,s.protocol,s.source_key,s.latest_total_mileage_km,s.first_total_mileage_km,s.daily_mileage_km,s.first_event_time,s.latest_event_time,
|
|
p.latest_total_mileage_km AS previous_total,p.latest_event_time AS previous_time
|
|
FROM vehicle_daily_mileage_source s
|
|
JOIN vehicle_daily_mileage_source p ON p.vin=s.vin AND p.protocol=s.protocol AND p.source_key=s.source_key
|
|
AND p.stat_date=(SELECT MAX(b.stat_date) FROM vehicle_daily_mileage_source b
|
|
WHERE b.vin=s.vin AND b.protocol=s.protocol AND b.source_key=s.source_key AND b.stat_date<s.stat_date
|
|
AND b.quality_status='OK' AND b.latest_total_mileage_km>0
|
|
AND b.latest_event_time>=TIMESTAMP(b.stat_date) AND b.latest_event_time<TIMESTAMP(b.stat_date)+INTERVAL 1 DAY)
|
|
WHERE s.stat_date<='2026-09-16' AND s.quality_status='OK'
|
|
AND s.quality_reason='current_day_first_baseline'
|
|
AND s.latest_total_mileage_km>0 AND s.first_event_time>=TIMESTAMP(s.stat_date)
|
|
AND s.source_ip NOT IN """+legacy+" ORDER BY s.vin,s.protocol,s.source_key,s.stat_date")
|
|
changes=q.fetchall();report['baseline_candidates']=len(changes)
|
|
q.execute("SELECT id FROM vehicle_data_source WHERE source_ip IN "+legacy);legacy_ids=[r['id'] for r in q.fetchall()]
|
|
c.rollback()
|
|
if args.apply:
|
|
assert (root/'manifest.json').is_file() and (root/'database.sql.gz').is_file(),'Complete backup required'
|
|
(root/'baseline-before.json').write_text(json.dumps(changes,default=str,ensure_ascii=False))
|
|
q.execute('START TRANSACTION')
|
|
q.execute("UPDATE vehicle_daily_mileage_source SET quality_status='INVALID_DELTA',quality_reason='UNVERIFIED_LEGACY_PROTOCOL',is_selected=0 WHERE protocol='GB32960' AND source_ip IN "+legacy)
|
|
report['quarantined_source_rows']=q.rowcount
|
|
q.execute("UPDATE vehicle_data_source SET enabled=0,remark='Protocol provenance unverified; quarantined 2026-09-16; original records retained' WHERE protocol='GB32960' AND source_ip IN "+legacy)
|
|
q.execute("DELETE m FROM vehicle_daily_mileage m JOIN vehicle_data_source d ON d.id=m.source_id WHERE m.protocol='GB32960' AND d.source_ip IN "+legacy)
|
|
report['removed_mislabelled_projections']=q.rowcount
|
|
affected=set()
|
|
for s in changes:
|
|
delta=s['latest_total_mileage_km']-s['previous_total']
|
|
days=max(1,(s['latest_event_time'].date()-s['previous_time'].date()).days)
|
|
status='OK';reason='historical_source_baseline'
|
|
if Decimal('-1')<=delta<0:delta=Decimal(0);reason='negative_jitter_clamped'
|
|
elif delta<0:status='INVALID_DELTA';reason='TOTAL_MILEAGE_ROLLBACK'
|
|
elif delta>Decimal(2500*days):status='INVALID_DELTA';reason='outside_daily_range'
|
|
q.execute("""UPDATE vehicle_daily_mileage_source SET first_total_mileage_km=%s,first_event_time=%s,daily_mileage_km=%s,
|
|
quality_status=%s,quality_reason=%s WHERE vin=%s AND stat_date=%s AND protocol=%s AND source_key=%s
|
|
AND first_total_mileage_km <=> %s AND daily_mileage_km=%s""",(s['previous_total'],s['previous_time'],delta,status,reason,s['vin'],s['stat_date'],s['protocol'],s['source_key'],s['first_total_mileage_km'],s['daily_mileage_km']))
|
|
if q.rowcount:affected.add((s['vin'],s['stat_date'],s['protocol']))
|
|
# Restore valid non-import GB32960 projections that were hidden by migrated rows.
|
|
q.execute("""SELECT DISTINCT s.vin,s.stat_date,s.protocol FROM vehicle_daily_mileage_source s
|
|
JOIN vehicle_daily_mileage_source bad ON bad.vin=s.vin AND bad.stat_date=s.stat_date AND bad.protocol=s.protocol
|
|
WHERE bad.source_ip IN """+legacy+" AND s.source_ip NOT IN "+legacy+" AND s.quality_status='OK' AND s.latest_total_mileage_km IS NOT NULL")
|
|
affected.update((r['vin'],r['stat_date'],r['protocol']) for r in q.fetchall())
|
|
report['reproject_days']=len(affected)
|
|
# The Go projection utility is run before resuming writers, using canonical SQL.
|
|
(root/'reproject.json').write_text(json.dumps([{'vin':v,'date':str(d),'protocol':p} for v,d,p in sorted(affected)]))
|
|
c.commit()
|
|
(root if args.apply else pathlib.Path('/tmp')).joinpath('repair-result.json' if args.apply else 'mileage-repair-dry-run.json').write_text(json.dumps(report,default=str,indent=2))
|
|
print(json.dumps(report,default=str),flush=True)
|
|
c.close()
|