diff --git a/scripts/migrate_csv.py b/scripts/migrate_csv.py new file mode 100755 index 0000000..573a7ce --- /dev/null +++ b/scripts/migrate_csv.py @@ -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="", + 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 "", + 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)