Add legacy CSV migration support
This commit is contained in:
745
scripts/migrate_csv.py
Executable file
745
scripts/migrate_csv.py
Executable file
@ -0,0 +1,745 @@
|
||||
#!/usr/bin/env python3
|
||||
|
||||
"""
|
||||
Validate and migrate the legacy WCX CSV index into SQLite.
|
||||
|
||||
The script reads the historical CSV index, validates and normalizes its
|
||||
contents, and can later insert or update matching rows in the SQLite movie
|
||||
table.
|
||||
|
||||
By default, no database changes are made. The --apply option must be supplied
|
||||
explicitly to perform the migration.
|
||||
|
||||
Validation includes:
|
||||
- Required CSV headers
|
||||
- Duplicate IDs
|
||||
- ID consistency with the filename portion of WEBURL
|
||||
- ISO date validation
|
||||
- Duration conversion from mm:ss or hh:mm:ss to seconds
|
||||
- Integer conversion for AGE
|
||||
- Decimal conversion for RATING
|
||||
- Detection of rows missing required ID or NAME fields
|
||||
|
||||
The CSV file is opened using UTF-8 with BOM support.
|
||||
|
||||
Usage:
|
||||
migrate_csv.py --dry-run
|
||||
migrate_csv.py --dry-run --limit 20
|
||||
migrate_csv.py --database /path/to/wcx-test.db --apply
|
||||
|
||||
No rows are written unless --apply is explicitly specified.
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import csv
|
||||
import re
|
||||
import sqlite3
|
||||
import sys
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import date
|
||||
from pathlib import Path
|
||||
from urllib.parse import urlparse
|
||||
|
||||
|
||||
DEFAULT_CSV_FILE = Path(
|
||||
"/storage/disk1/WCX/migration/persons.utf8bom.csv"
|
||||
)
|
||||
DEFAULT_DATABASE_FILE = Path(
|
||||
"/storage/disk1/WCX/database/wcx.db"
|
||||
)
|
||||
|
||||
EXPECTED_HEADERS = [
|
||||
"ID",
|
||||
"NAME",
|
||||
"AKA",
|
||||
"NATIONALITY",
|
||||
"AGE",
|
||||
"SHOOT_LOCATION",
|
||||
"SHOOT_DATE",
|
||||
"DURATION",
|
||||
"DESCRIPTION",
|
||||
"RATING",
|
||||
"WEBURL",
|
||||
"THUMBNAIL",
|
||||
"PUBLISHED",
|
||||
"UPDATED",
|
||||
]
|
||||
|
||||
|
||||
@dataclass
|
||||
class ValidationIssue:
|
||||
row_number: int
|
||||
movie_id: str
|
||||
message: str
|
||||
|
||||
@dataclass
|
||||
class MigrationStats:
|
||||
rows_read: int = 0
|
||||
valid_rows: int = 0
|
||||
invalid_rows: int = 0
|
||||
duplicate_ids: int = 0
|
||||
existing_in_database: int = 0
|
||||
missing_from_database: int = 0
|
||||
id_url_mismatches: int = 0
|
||||
inserted: int = 0
|
||||
updated: int = 0
|
||||
warnings: list[ValidationIssue] = field(default_factory=list)
|
||||
errors: list[ValidationIssue] = field(default_factory=list)
|
||||
|
||||
@dataclass
|
||||
class MovieRecord:
|
||||
movie_id: str
|
||||
name: str
|
||||
aka: str | None
|
||||
nationality: str | None
|
||||
age: int | None
|
||||
shoot_location: str | None
|
||||
shoot_date: str | None
|
||||
duration_seconds: int | None
|
||||
description: str | None
|
||||
rating: float | None
|
||||
web_url: str | None
|
||||
thumbnail: str | None
|
||||
published: str | None
|
||||
updated: str | None
|
||||
ocr_status: str
|
||||
|
||||
|
||||
def clean_text(value: str | None) -> str | None:
|
||||
if value is None:
|
||||
return None
|
||||
|
||||
cleaned = value.strip()
|
||||
return cleaned if cleaned else None
|
||||
|
||||
|
||||
def parse_duration(value: str | None) -> int | None:
|
||||
value = clean_text(value)
|
||||
|
||||
if value is None:
|
||||
return None
|
||||
|
||||
parts = value.split(":")
|
||||
|
||||
try:
|
||||
numbers = [int(part) for part in parts]
|
||||
except ValueError as error:
|
||||
raise ValueError(f"Invalid duration: {value!r}") from error
|
||||
|
||||
if len(numbers) == 2:
|
||||
hours = 0
|
||||
minutes, seconds = numbers
|
||||
elif len(numbers) == 3:
|
||||
hours, minutes, seconds = numbers
|
||||
else:
|
||||
raise ValueError(
|
||||
f"Expected mm:ss or hh:mm:ss, got {value!r}"
|
||||
)
|
||||
|
||||
if hours < 0 or minutes < 0 or seconds < 0:
|
||||
raise ValueError(f"Duration cannot be negative: {value!r}")
|
||||
|
||||
if seconds >= 60:
|
||||
raise ValueError(
|
||||
f"Seconds must be less than 60: {value!r}"
|
||||
)
|
||||
|
||||
if len(numbers) == 3 and minutes >= 60:
|
||||
raise ValueError(
|
||||
f"Minutes must be less than 60 in hh:mm:ss: {value!r}"
|
||||
)
|
||||
|
||||
return hours * 3600 + minutes * 60 + seconds
|
||||
|
||||
|
||||
def parse_iso_date(
|
||||
value: str | None,
|
||||
field_name: str,
|
||||
) -> str | None:
|
||||
value = clean_text(value)
|
||||
|
||||
if value is None:
|
||||
return None
|
||||
|
||||
try:
|
||||
return date.fromisoformat(value).isoformat()
|
||||
except ValueError as error:
|
||||
raise ValueError(
|
||||
f"Invalid {field_name} date: {value!r}; "
|
||||
"expected YYYY-MM-DD"
|
||||
) from error
|
||||
|
||||
def parse_age(value: str | None) -> int | None:
|
||||
value = clean_text(value)
|
||||
|
||||
if value is None or value == "?":
|
||||
return None
|
||||
|
||||
match = re.match(r"^(\d+)", value)
|
||||
|
||||
if match is None:
|
||||
return None
|
||||
|
||||
age = int(match.group(1))
|
||||
|
||||
if age < 0:
|
||||
raise ValueError(f"Age cannot be negative: {value!r}")
|
||||
|
||||
return age
|
||||
|
||||
def parse_rating(value: str | None) -> float | None:
|
||||
value = clean_text(value)
|
||||
|
||||
if value is None:
|
||||
return None
|
||||
|
||||
normalized = value.replace(",", ".")
|
||||
|
||||
try:
|
||||
rating = float(normalized)
|
||||
except ValueError as error:
|
||||
raise ValueError(f"Invalid rating: {value!r}") from error
|
||||
|
||||
if not 0 <= rating <= 10:
|
||||
raise ValueError(
|
||||
f"Rating must be between 0 and 10: {value!r}"
|
||||
)
|
||||
|
||||
return rating
|
||||
|
||||
|
||||
def id_from_web_url(web_url: str | None) -> str | None:
|
||||
web_url = clean_text(web_url)
|
||||
|
||||
if web_url is None:
|
||||
return None
|
||||
|
||||
path = urlparse(web_url).path
|
||||
filename = Path(path).name
|
||||
|
||||
if not filename:
|
||||
return None
|
||||
|
||||
if filename.lower().endswith(".html"):
|
||||
return filename[:-5]
|
||||
|
||||
return Path(filename).stem
|
||||
|
||||
|
||||
def determine_ocr_status(
|
||||
nationality: str | None,
|
||||
shoot_location: str | None,
|
||||
shoot_date: str | None,
|
||||
) -> str:
|
||||
if nationality and shoot_location and shoot_date:
|
||||
return "completed"
|
||||
|
||||
return "manual_review"
|
||||
|
||||
|
||||
def validate_headers(fieldnames: list[str] | None) -> None:
|
||||
if fieldnames is None:
|
||||
raise ValueError("The CSV file has no header row.")
|
||||
|
||||
missing = [
|
||||
header
|
||||
for header in EXPECTED_HEADERS
|
||||
if header not in fieldnames
|
||||
]
|
||||
|
||||
extra = [
|
||||
header
|
||||
for header in fieldnames
|
||||
if header not in EXPECTED_HEADERS
|
||||
]
|
||||
|
||||
if missing:
|
||||
raise ValueError(
|
||||
"Missing CSV headers: " + ", ".join(missing)
|
||||
)
|
||||
|
||||
if extra:
|
||||
print(
|
||||
"Warning: unexpected CSV headers: "
|
||||
+ ", ".join(extra),
|
||||
file=sys.stderr,
|
||||
)
|
||||
|
||||
|
||||
def parse_row(
|
||||
row: dict[str, str],
|
||||
row_number: int,
|
||||
stats: MigrationStats,
|
||||
) -> MovieRecord | None:
|
||||
movie_id = clean_text(row.get("ID")) or ""
|
||||
name = clean_text(row.get("NAME")) or ""
|
||||
|
||||
row_errors: list[str] = []
|
||||
|
||||
if not movie_id:
|
||||
stats.warnings.append(
|
||||
ValidationIssue(
|
||||
row_number=row_number,
|
||||
movie_id="<missing>",
|
||||
message="Missing ID; row skipped",
|
||||
)
|
||||
)
|
||||
stats.invalid_rows += 1
|
||||
return None
|
||||
|
||||
if not name:
|
||||
row_errors.append("Missing NAME")
|
||||
|
||||
try:
|
||||
age = parse_age(row.get("AGE"))
|
||||
except ValueError as error:
|
||||
age = None
|
||||
row_errors.append(str(error))
|
||||
|
||||
try:
|
||||
shoot_date = parse_iso_date(
|
||||
row.get("SHOOT_DATE"),
|
||||
"SHOOT_DATE",
|
||||
)
|
||||
except ValueError as error:
|
||||
shoot_date = None
|
||||
row_errors.append(str(error))
|
||||
|
||||
try:
|
||||
duration_seconds = parse_duration(row.get("DURATION"))
|
||||
except ValueError as error:
|
||||
duration_seconds = None
|
||||
row_errors.append(str(error))
|
||||
|
||||
try:
|
||||
rating = parse_rating(row.get("RATING"))
|
||||
except ValueError as error:
|
||||
rating = None
|
||||
row_errors.append(str(error))
|
||||
|
||||
try:
|
||||
published = parse_iso_date(
|
||||
row.get("PUBLISHED"),
|
||||
"PUBLISHED",
|
||||
)
|
||||
except ValueError as error:
|
||||
published = None
|
||||
row_errors.append(str(error))
|
||||
|
||||
try:
|
||||
updated = parse_iso_date(
|
||||
row.get("UPDATED"),
|
||||
"UPDATED",
|
||||
)
|
||||
except ValueError as error:
|
||||
updated = None
|
||||
row_errors.append(str(error))
|
||||
|
||||
web_url = clean_text(row.get("WEBURL"))
|
||||
url_id = id_from_web_url(web_url)
|
||||
|
||||
if movie_id and url_id and movie_id != url_id:
|
||||
stats.id_url_mismatches += 1
|
||||
stats.warnings.append(
|
||||
ValidationIssue(
|
||||
row_number=row_number,
|
||||
movie_id=movie_id,
|
||||
message=(
|
||||
f"ID does not match WEBURL ID: "
|
||||
f"{movie_id!r} != {url_id!r}"
|
||||
),
|
||||
)
|
||||
)
|
||||
|
||||
if row_errors:
|
||||
for message in row_errors:
|
||||
stats.errors.append(
|
||||
ValidationIssue(
|
||||
row_number=row_number,
|
||||
movie_id=movie_id or "<missing>",
|
||||
message=message,
|
||||
)
|
||||
)
|
||||
|
||||
stats.invalid_rows += 1
|
||||
return None
|
||||
|
||||
nationality = clean_text(row.get("NATIONALITY"))
|
||||
shoot_location = clean_text(row.get("SHOOT_LOCATION"))
|
||||
|
||||
return MovieRecord(
|
||||
movie_id=movie_id,
|
||||
name=name,
|
||||
aka=clean_text(row.get("AKA")),
|
||||
nationality=nationality,
|
||||
age=age,
|
||||
shoot_location=shoot_location,
|
||||
shoot_date=shoot_date,
|
||||
duration_seconds=duration_seconds,
|
||||
description=clean_text(row.get("DESCRIPTION")),
|
||||
rating=rating,
|
||||
web_url=web_url,
|
||||
thumbnail=clean_text(row.get("THUMBNAIL")),
|
||||
published=published,
|
||||
updated=updated,
|
||||
ocr_status=determine_ocr_status(
|
||||
nationality,
|
||||
shoot_location,
|
||||
shoot_date,
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def read_csv_records(
|
||||
csv_file: Path,
|
||||
limit: int | None,
|
||||
stats: MigrationStats,
|
||||
) -> list[MovieRecord]:
|
||||
records: list[MovieRecord] = []
|
||||
seen_ids: set[str] = set()
|
||||
|
||||
with csv_file.open(
|
||||
"r",
|
||||
encoding="utf-8-sig",
|
||||
newline="",
|
||||
) as file:
|
||||
reader = csv.DictReader(file)
|
||||
validate_headers(reader.fieldnames)
|
||||
|
||||
for row_number, row in enumerate(reader, start=2):
|
||||
if limit is not None and stats.rows_read >= limit:
|
||||
break
|
||||
|
||||
stats.rows_read += 1
|
||||
|
||||
record = parse_row(row, row_number, stats)
|
||||
|
||||
if record is None:
|
||||
continue
|
||||
|
||||
if record.movie_id in seen_ids:
|
||||
stats.duplicate_ids += 1
|
||||
stats.invalid_rows += 1
|
||||
stats.errors.append(
|
||||
ValidationIssue(
|
||||
row_number=row_number,
|
||||
movie_id=record.movie_id,
|
||||
message="Duplicate ID in CSV",
|
||||
)
|
||||
)
|
||||
continue
|
||||
|
||||
seen_ids.add(record.movie_id)
|
||||
records.append(record)
|
||||
stats.valid_rows += 1
|
||||
|
||||
return records
|
||||
|
||||
|
||||
def inspect_database_matches(
|
||||
database_file: Path,
|
||||
records: list[MovieRecord],
|
||||
stats: MigrationStats,
|
||||
) -> None:
|
||||
if not database_file.is_file():
|
||||
raise FileNotFoundError(
|
||||
f"Database file does not exist: {database_file}"
|
||||
)
|
||||
|
||||
with sqlite3.connect(database_file) as connection:
|
||||
for record in records:
|
||||
exists = connection.execute(
|
||||
"""
|
||||
SELECT 1
|
||||
FROM movie
|
||||
WHERE id = ?
|
||||
""",
|
||||
(record.movie_id,),
|
||||
).fetchone()
|
||||
|
||||
if exists:
|
||||
stats.existing_in_database += 1
|
||||
else:
|
||||
stats.missing_from_database += 1
|
||||
|
||||
|
||||
def apply_migration(
|
||||
database_file: Path,
|
||||
records: list[MovieRecord],
|
||||
stats: MigrationStats,
|
||||
) -> None:
|
||||
if stats.errors:
|
||||
raise ValueError(
|
||||
"Migration cannot be applied while validation errors exist."
|
||||
)
|
||||
|
||||
inserted = 0
|
||||
updated = 0
|
||||
|
||||
with sqlite3.connect(database_file) as connection:
|
||||
for record in records:
|
||||
existing = connection.execute(
|
||||
"""
|
||||
SELECT 1
|
||||
FROM movie
|
||||
WHERE id = ?
|
||||
""",
|
||||
(record.movie_id,),
|
||||
).fetchone()
|
||||
|
||||
if existing is None:
|
||||
connection.execute(
|
||||
"""
|
||||
INSERT INTO movie (
|
||||
id,
|
||||
name,
|
||||
aka,
|
||||
nationality,
|
||||
age,
|
||||
shoot_location,
|
||||
shoot_date,
|
||||
duration_seconds,
|
||||
description,
|
||||
rating,
|
||||
web_url,
|
||||
thumbnail,
|
||||
published,
|
||||
updated,
|
||||
ocr_status
|
||||
)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""",
|
||||
(
|
||||
record.movie_id,
|
||||
record.name,
|
||||
record.aka,
|
||||
record.nationality,
|
||||
record.age,
|
||||
record.shoot_location,
|
||||
record.shoot_date,
|
||||
record.duration_seconds,
|
||||
record.description,
|
||||
record.rating,
|
||||
record.web_url,
|
||||
record.thumbnail,
|
||||
record.published,
|
||||
record.updated,
|
||||
record.ocr_status,
|
||||
),
|
||||
)
|
||||
inserted += 1
|
||||
else:
|
||||
connection.execute(
|
||||
"""
|
||||
UPDATE movie
|
||||
SET
|
||||
name = ?,
|
||||
aka = ?,
|
||||
nationality = ?,
|
||||
age = ?,
|
||||
shoot_location = ?,
|
||||
shoot_date = ?,
|
||||
duration_seconds = ?,
|
||||
description = ?,
|
||||
rating = ?,
|
||||
web_url = ?,
|
||||
thumbnail = ?,
|
||||
published = ?,
|
||||
updated = ?,
|
||||
ocr_status = ?,
|
||||
modified_at = CURRENT_TIMESTAMP
|
||||
WHERE id = ?
|
||||
""",
|
||||
(
|
||||
record.name,
|
||||
record.aka,
|
||||
record.nationality,
|
||||
record.age,
|
||||
record.shoot_location,
|
||||
record.shoot_date,
|
||||
record.duration_seconds,
|
||||
record.description,
|
||||
record.rating,
|
||||
record.web_url,
|
||||
record.thumbnail,
|
||||
record.published,
|
||||
record.updated,
|
||||
record.ocr_status,
|
||||
record.movie_id,
|
||||
),
|
||||
)
|
||||
updated += 1
|
||||
|
||||
stats.inserted = inserted
|
||||
stats.updated = updated
|
||||
|
||||
def print_issues(
|
||||
heading: str,
|
||||
issues: list[ValidationIssue],
|
||||
maximum: int = 25,
|
||||
) -> None:
|
||||
if not issues:
|
||||
return
|
||||
|
||||
print()
|
||||
print(heading)
|
||||
|
||||
for issue in issues[:maximum]:
|
||||
print(
|
||||
f" Row {issue.row_number}, "
|
||||
f"ID {issue.movie_id}: {issue.message}"
|
||||
)
|
||||
|
||||
remaining = len(issues) - maximum
|
||||
|
||||
if remaining > 0:
|
||||
print(f" ... and {remaining} more")
|
||||
|
||||
|
||||
def print_summary(
|
||||
csv_file: Path,
|
||||
database_file: Path,
|
||||
limit: int | None,
|
||||
stats: MigrationStats,
|
||||
) -> None:
|
||||
print()
|
||||
print("CSV migration analysis")
|
||||
print("----------------------")
|
||||
print(f"CSV file: {csv_file}")
|
||||
print(f"Database: {database_file}")
|
||||
print(
|
||||
f"Limit: "
|
||||
f"{limit if limit is not None else 'none'}"
|
||||
)
|
||||
print(f"Rows read: {stats.rows_read}")
|
||||
print(f"Valid rows: {stats.valid_rows}")
|
||||
print(f"Invalid rows: {stats.invalid_rows}")
|
||||
print(f"Duplicate IDs: {stats.duplicate_ids}")
|
||||
print(
|
||||
f"Existing in database: {stats.existing_in_database}"
|
||||
)
|
||||
print(
|
||||
f"Missing from database:{stats.missing_from_database:>5}"
|
||||
)
|
||||
print(
|
||||
f"ID/WEBURL mismatches: {stats.id_url_mismatches}"
|
||||
)
|
||||
print(f"Warnings: {len(stats.warnings)}")
|
||||
print(f"Errors: {len(stats.errors)}")
|
||||
|
||||
print_issues("Warnings:", stats.warnings)
|
||||
print_issues("Errors:", stats.errors)
|
||||
print(f"Inserted: {stats.inserted}")
|
||||
print(f"Updated: {stats.updated}")
|
||||
|
||||
|
||||
def parse_arguments() -> argparse.Namespace:
|
||||
parser = argparse.ArgumentParser(
|
||||
description=(
|
||||
"Validate the legacy WCX CSV migration source."
|
||||
)
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--input",
|
||||
type=Path,
|
||||
default=DEFAULT_CSV_FILE,
|
||||
help=f"CSV input file; default: {DEFAULT_CSV_FILE}",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--database",
|
||||
type=Path,
|
||||
default=DEFAULT_DATABASE_FILE,
|
||||
help=(
|
||||
"SQLite database used for match analysis; "
|
||||
f"default: {DEFAULT_DATABASE_FILE}"
|
||||
),
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--limit",
|
||||
type=int,
|
||||
help="Only analyze the first N data rows.",
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--dry-run",
|
||||
action="store_true",
|
||||
help=(
|
||||
"Validate and analyze without changing the database. "
|
||||
"This is currently the only supported mode."
|
||||
),
|
||||
)
|
||||
|
||||
parser.add_argument(
|
||||
"--apply",
|
||||
action="store_true",
|
||||
help=(
|
||||
"Reserved for the next implementation step. "
|
||||
"Database writes are not implemented yet."
|
||||
),
|
||||
)
|
||||
|
||||
return parser.parse_args()
|
||||
|
||||
|
||||
def main() -> None:
|
||||
args = parse_arguments()
|
||||
|
||||
if args.limit is not None and args.limit <= 0:
|
||||
raise ValueError("--limit must be greater than zero.")
|
||||
|
||||
if args.apply and args.dry_run:
|
||||
raise ValueError(
|
||||
"Use either --dry-run or --apply, not both."
|
||||
)
|
||||
|
||||
if not args.input.is_file():
|
||||
raise FileNotFoundError(
|
||||
f"CSV input file does not exist: {args.input}"
|
||||
)
|
||||
|
||||
stats = MigrationStats()
|
||||
|
||||
records = read_csv_records(
|
||||
csv_file=args.input,
|
||||
limit=args.limit,
|
||||
stats=stats,
|
||||
)
|
||||
|
||||
inspect_database_matches(
|
||||
database_file=args.database,
|
||||
records=records,
|
||||
stats=stats,
|
||||
)
|
||||
|
||||
if args.apply:
|
||||
apply_migration(
|
||||
database_file=args.database,
|
||||
records=records,
|
||||
stats=stats,
|
||||
)
|
||||
|
||||
print_summary(
|
||||
csv_file=args.input,
|
||||
database_file=args.database,
|
||||
limit=args.limit,
|
||||
stats=stats,
|
||||
)
|
||||
|
||||
if stats.errors:
|
||||
sys.exit(2)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
try:
|
||||
main()
|
||||
except (
|
||||
FileNotFoundError,
|
||||
ValueError,
|
||||
csv.Error,
|
||||
sqlite3.Error,
|
||||
) as error:
|
||||
print(f"Error: {error}", file=sys.stderr)
|
||||
sys.exit(1)
|
||||
Reference in New Issue
Block a user