diff --git a/app/utils/migrate_logs_to_steps.py b/app/utils/migrate_logs_to_steps.py new file mode 100644 index 00000000..8d148804 --- /dev/null +++ b/app/utils/migrate_logs_to_steps.py @@ -0,0 +1,326 @@ +""" +Migration utility to backfill FileProcessingStep table from existing ProcessingLog entries. + +This script analyzes historical logs and populates the FileProcessingStep table +for files that were processed before the status tracking table was created. +""" + +import logging +from datetime import datetime +from typing import Dict, List + +from sqlalchemy import func +from sqlalchemy.orm import Session + +from app.models import FileProcessingStep, FileRecord, ProcessingLog +from app.utils.step_manager import MAIN_PROCESSING_STEPS + +logger = logging.getLogger(__name__) + + +def migrate_logs_to_steps(db: Session, file_id: int, dry_run: bool = False) -> Dict: + """ + Migrate ProcessingLog entries to FileProcessingStep for a single file. + + Analyzes all logs for the file and creates/updates FileProcessingStep entries + based on the latest status per step found in the logs. + + Args: + db: Database session + file_id: ID of the file to migrate + dry_run: If True, don't commit changes, just report what would be done + + Returns: + Dictionary with migration results: + { + "file_id": 123, + "steps_created": 5, + "steps_updated": 3, + "steps_skipped": 2, + "errors": [] + } + """ + results = {"file_id": file_id, "steps_created": 0, "steps_updated": 0, "steps_skipped": 0, "errors": []} + + try: + # Get all logs for this file, ordered by timestamp + logs = ( + db.query(ProcessingLog) + .filter(ProcessingLog.file_id == file_id) + .order_by(ProcessingLog.timestamp.asc()) + .all() + ) + + if not logs: + logger.info(f"No logs found for file {file_id}") + return results + + # Parse logs to extract latest status per step + step_states = _parse_logs_to_step_states(logs) + + # Create or update FileProcessingStep entries + for step_name, state in step_states.items(): + existing_step = ( + db.query(FileProcessingStep) + .filter(FileProcessingStep.file_id == file_id, FileProcessingStep.step_name == step_name) + .first() + ) + + if existing_step: + # Update existing step only if it differs + if ( + existing_step.status != state["status"] + or existing_step.started_at != state["started_at"] + or existing_step.completed_at != state["completed_at"] + ): + logger.info( + f"Updating step {step_name} for file {file_id}: {existing_step.status} -> {state['status']}" + ) + existing_step.status = state["status"] + existing_step.started_at = state["started_at"] + existing_step.completed_at = state["completed_at"] + existing_step.error_message = state["error_message"] + results["steps_updated"] += 1 + else: + logger.debug(f"Step {step_name} for file {file_id} already up to date") + results["steps_skipped"] += 1 + else: + # Create new step + logger.info(f"Creating step {step_name} for file {file_id} with status {state['status']}") + new_step = FileProcessingStep( + file_id=file_id, + step_name=step_name, + status=state["status"], + started_at=state["started_at"], + completed_at=state["completed_at"], + error_message=state["error_message"], + ) + db.add(new_step) + results["steps_created"] += 1 + + if not dry_run: + db.commit() + logger.info( + f"Migration complete for file {file_id}: " + f"{results['steps_created']} created, " + f"{results['steps_updated']} updated, " + f"{results['steps_skipped']} skipped" + ) + else: + db.rollback() + logger.info(f"Dry run complete for file {file_id} (no changes committed)") + + except Exception as e: + logger.error(f"Error migrating file {file_id}: {str(e)}") + results["errors"].append(str(e)) + db.rollback() + + return results + + +def _parse_logs_to_step_states(logs: List[ProcessingLog]) -> Dict[str, Dict]: + """ + Parse logs to determine the final state of each step. + + Processes logs in chronological order to build a state machine + tracking the progression of each step. + + Args: + logs: List of ProcessingLog entries ordered by timestamp + + Returns: + Dictionary mapping step_name to state dict: + { + "step_name": { + "status": "success", + "started_at": datetime, + "completed_at": datetime, + "error_message": None + } + } + """ + step_states = {} + + for log in logs: + step_name = log.step_name + status = log.status.lower() + + # Initialize step state if first time seeing this step + if step_name not in step_states: + step_states[step_name] = { + "status": "pending", + "started_at": None, + "completed_at": None, + "error_message": None, + } + + # Update state based on log entry + if status == "in_progress": + # Step started + if step_states[step_name]["started_at"] is None: + step_states[step_name]["started_at"] = log.timestamp + step_states[step_name]["status"] = "in_progress" + + elif status == "success": + # Step completed successfully + if step_states[step_name]["started_at"] is None: + # If we never saw in_progress, assume it started around this time + step_states[step_name]["started_at"] = log.timestamp + step_states[step_name]["status"] = "success" + step_states[step_name]["completed_at"] = log.timestamp + step_states[step_name]["error_message"] = None + + elif status == "failure": + # Step failed + if step_states[step_name]["started_at"] is None: + step_states[step_name]["started_at"] = log.timestamp + step_states[step_name]["status"] = "failure" + step_states[step_name]["completed_at"] = log.timestamp + step_states[step_name]["error_message"] = log.message + + elif status in ["pending", "queued"]: + # Only update if step hasn't started yet + if step_states[step_name]["status"] == "pending": + step_states[step_name]["status"] = "pending" + + # Note: If we see success/failure after a previous success/failure, + # the later one wins (represents retry/reprocessing) + + return step_states + + +def migrate_all_files(db: Session, batch_size: int = 100, dry_run: bool = False) -> Dict: + """ + Migrate all files that have logs but no FileProcessingStep entries. + + Args: + db: Database session + batch_size: Number of files to process in each batch + dry_run: If True, don't commit changes + + Returns: + Dictionary with overall migration statistics: + { + "total_files": 150, + "files_migrated": 145, + "files_failed": 5, + "total_steps_created": 1200, + "total_steps_updated": 350, + "errors": [...] + } + """ + summary = { + "total_files": 0, + "files_migrated": 0, + "files_failed": 0, + "total_steps_created": 0, + "total_steps_updated": 0, + "total_steps_skipped": 0, + "errors": [], + } + + # Find all files that have logs but no steps + files_with_logs = db.query(ProcessingLog.file_id).distinct().all() + file_ids_with_logs = {file_id for (file_id,) in files_with_logs if file_id is not None} + + files_with_steps = db.query(FileProcessingStep.file_id).distinct().all() + file_ids_with_steps = {file_id for (file_id,) in files_with_steps} + + files_to_migrate = list(file_ids_with_logs - file_ids_with_steps) + + summary["total_files"] = len(files_to_migrate) + logger.info(f"Found {len(files_to_migrate)} files to migrate") + + # Process in batches + for i in range(0, len(files_to_migrate), batch_size): + batch = files_to_migrate[i : i + batch_size] + logger.info( + f"Processing batch {i // batch_size + 1}: files {i + 1} to {min(i + batch_size, len(files_to_migrate))}" + ) + + for file_id in batch: + result = migrate_logs_to_steps(db, file_id, dry_run=dry_run) + + if result["errors"]: + summary["files_failed"] += 1 + summary["errors"].extend(result["errors"]) + else: + summary["files_migrated"] += 1 + summary["total_steps_created"] += result["steps_created"] + summary["total_steps_updated"] += result["steps_updated"] + summary["total_steps_skipped"] += result["steps_skipped"] + + logger.info( + f"Migration summary: {summary['files_migrated']} files migrated, " + f"{summary['files_failed']} failed, " + f"{summary['total_steps_created']} steps created, " + f"{summary['total_steps_updated']} steps updated" + ) + + return summary + + +def verify_migration(db: Session, file_id: int) -> Dict: + """ + Verify that migration for a file is correct by comparing logs to steps. + + Args: + db: Database session + file_id: ID of the file to verify + + Returns: + Dictionary with verification results: + { + "file_id": 123, + "is_valid": True, + "discrepancies": [], + "log_steps": ["step1", "step2"], + "table_steps": ["step1", "step2"] + } + """ + result = {"file_id": file_id, "is_valid": True, "discrepancies": [], "log_steps": [], "table_steps": []} + + # Get logs and parse them + logs = ( + db.query(ProcessingLog).filter(ProcessingLog.file_id == file_id).order_by(ProcessingLog.timestamp.asc()).all() + ) + + if not logs: + result["discrepancies"].append("No logs found for file") + result["is_valid"] = False + return result + + expected_states = _parse_logs_to_step_states(logs) + result["log_steps"] = sorted(expected_states.keys()) + + # Get actual steps from table + actual_steps = db.query(FileProcessingStep).filter(FileProcessingStep.file_id == file_id).all() + result["table_steps"] = sorted([s.step_name for s in actual_steps]) + + # Compare + actual_states = {s.step_name: s for s in actual_steps} + + # Check for missing steps + for step_name in expected_states: + if step_name not in actual_states: + result["discrepancies"].append(f"Step '{step_name}' missing from table") + result["is_valid"] = False + continue + + # Compare status + expected = expected_states[step_name] + actual = actual_states[step_name] + + if expected["status"] != actual.status: + result["discrepancies"].append( + f"Step '{step_name}' status mismatch: " f"expected '{expected['status']}', got '{actual.status}'" + ) + result["is_valid"] = False + + # Check for extra steps in table + for step_name in actual_states: + if step_name not in expected_states: + result["discrepancies"].append(f"Extra step '{step_name}' in table not found in logs") + # Not marking as invalid since this might be intentional + + return result diff --git a/tests/test_step_manager.py b/tests/test_step_manager.py new file mode 100644 index 00000000..37c13305 --- /dev/null +++ b/tests/test_step_manager.py @@ -0,0 +1,331 @@ +""" +Tests for the FileProcessingStep model and step_manager utilities. + +This test module verifies the new explicit status tracking approach. +""" + +from datetime import datetime, timedelta + +import pytest +from sqlalchemy import create_engine +from sqlalchemy.orm import Session, sessionmaker + +from app.database import Base +from app.models import FileProcessingStep, FileRecord +from app.utils.step_manager import ( + MAIN_PROCESSING_STEPS, + add_upload_steps, + get_file_overall_status, + get_file_step_status, + get_step_summary, + initialize_file_steps, + update_step_status, +) + + +@pytest.fixture +def db_session(): + """Create an in-memory SQLite database for testing.""" + engine = create_engine("sqlite:///:memory:") + Base.metadata.create_all(engine) + SessionLocal = sessionmaker(bind=engine) + session = SessionLocal() + + yield session + + session.close() + Base.metadata.drop_all(engine) + + +@pytest.mark.unit +class TestFileProcessingStepModel: + """Test the FileProcessingStep model.""" + + def test_create_step(self, db_session: Session): + """Test creating a processing step.""" + # Create a file record first + file_record = FileRecord( + filehash="test123", original_filename="test.pdf", local_filename="/tmp/test.pdf", file_size=1024 + ) + db_session.add(file_record) + db_session.commit() + + # Create a processing step + step = FileProcessingStep(file_id=file_record.id, step_name="hash_file", status="success") + db_session.add(step) + db_session.commit() + + assert step.id is not None + assert step.file_id == file_record.id + assert step.step_name == "hash_file" + assert step.status == "success" + + def test_unique_constraint(self, db_session: Session): + """Test that the unique constraint on (file_id, step_name) works.""" + # Create a file record + file_record = FileRecord( + filehash="test456", original_filename="test2.pdf", local_filename="/tmp/test2.pdf", file_size=2048 + ) + db_session.add(file_record) + db_session.commit() + + # Create first step + step1 = FileProcessingStep(file_id=file_record.id, step_name="hash_file", status="in_progress") + db_session.add(step1) + db_session.commit() + + # Try to create duplicate step - should fail + step2 = FileProcessingStep(file_id=file_record.id, step_name="hash_file", status="success") + db_session.add(step2) + + with pytest.raises(Exception): # SQLAlchemy will raise an integrity error + db_session.commit() + + +@pytest.mark.unit +class TestStepManager: + """Test step_manager utility functions.""" + + def test_initialize_file_steps(self, db_session: Session): + """Test initializing processing steps for a file.""" + # Create a file record + file_record = FileRecord( + filehash="test789", original_filename="test3.pdf", local_filename="/tmp/test3.pdf", file_size=4096 + ) + db_session.add(file_record) + db_session.commit() + + # Initialize steps + initialize_file_steps(db_session, file_record.id) + + # Verify steps were created + steps = db_session.query(FileProcessingStep).filter(FileProcessingStep.file_id == file_record.id).all() + + assert len(steps) == len(MAIN_PROCESSING_STEPS) + for step in steps: + assert step.status == "pending" + assert step.step_name in MAIN_PROCESSING_STEPS + + def test_add_upload_steps(self, db_session: Session): + """Test adding upload destination steps.""" + # Create a file record + file_record = FileRecord( + filehash="test101112", original_filename="test4.pdf", local_filename="/tmp/test4.pdf", file_size=8192 + ) + db_session.add(file_record) + db_session.commit() + + # Add upload steps + destinations = ["dropbox", "s3", "nextcloud"] + add_upload_steps(db_session, file_record.id, destinations) + + # Verify upload steps were created + steps = db_session.query(FileProcessingStep).filter(FileProcessingStep.file_id == file_record.id).all() + + # Should have 2 steps per destination (queue_ and upload_to_) + assert len(steps) == len(destinations) * 2 + + expected_steps = [] + for dest in destinations: + expected_steps.append(f"queue_{dest}") + expected_steps.append(f"upload_to_{dest}") + + for step in steps: + assert step.step_name in expected_steps + assert step.status == "pending" + + def test_update_step_status_new_step(self, db_session: Session): + """Test updating status creates step if it doesn't exist.""" + # Create a file record + file_record = FileRecord( + filehash="test131415", original_filename="test5.pdf", local_filename="/tmp/test5.pdf", file_size=16384 + ) + db_session.add(file_record) + db_session.commit() + + # Update a step that doesn't exist yet + now = datetime.now() + update_step_status(db_session, file_record.id, "hash_file", "in_progress", started_at=now) + + # Verify step was created + step = ( + db_session.query(FileProcessingStep) + .filter(FileProcessingStep.file_id == file_record.id, FileProcessingStep.step_name == "hash_file") + .first() + ) + + assert step is not None + assert step.status == "in_progress" + assert step.started_at == now + + def test_update_step_status_existing_step(self, db_session: Session): + """Test updating status of existing step.""" + # Create a file record and initialize steps + file_record = FileRecord( + filehash="test161718", original_filename="test6.pdf", local_filename="/tmp/test6.pdf", file_size=32768 + ) + db_session.add(file_record) + db_session.commit() + + initialize_file_steps(db_session, file_record.id) + + # Update an existing step + now = datetime.now() + update_step_status( + db_session, file_record.id, "hash_file", "success", started_at=now - timedelta(seconds=5), completed_at=now + ) + + # Verify step was updated + step = ( + db_session.query(FileProcessingStep) + .filter(FileProcessingStep.file_id == file_record.id, FileProcessingStep.step_name == "hash_file") + .first() + ) + + assert step.status == "success" + assert step.started_at == now - timedelta(seconds=5) + assert step.completed_at == now + + def test_get_file_step_status(self, db_session: Session): + """Test retrieving all step statuses for a file.""" + # Create a file record and initialize steps + file_record = FileRecord( + filehash="test192021", original_filename="test7.pdf", local_filename="/tmp/test7.pdf", file_size=65536 + ) + db_session.add(file_record) + db_session.commit() + + initialize_file_steps(db_session, file_record.id) + + # Update some steps + now = datetime.now() + update_step_status(db_session, file_record.id, "hash_file", "success", completed_at=now) + update_step_status(db_session, file_record.id, "create_file_record", "in_progress", started_at=now) + update_step_status(db_session, file_record.id, "check_text", "failure", error_message="Failed to check text") + + # Get all step statuses + status_map = get_file_step_status(db_session, file_record.id) + + assert len(status_map) == len(MAIN_PROCESSING_STEPS) + assert status_map["hash_file"]["status"] == "success" + assert status_map["hash_file"]["completed_at"] == now + assert status_map["create_file_record"]["status"] == "in_progress" + assert status_map["check_text"]["status"] == "failure" + assert status_map["check_text"]["error_message"] == "Failed to check text" + + def test_get_file_overall_status_pending(self, db_session: Session): + """Test overall status for a file with pending steps.""" + file_record = FileRecord( + filehash="test222324", original_filename="test8.pdf", local_filename="/tmp/test8.pdf", file_size=131072 + ) + db_session.add(file_record) + db_session.commit() + + initialize_file_steps(db_session, file_record.id) + + status = get_file_overall_status(db_session, file_record.id) + + assert status["status"] == "pending" + assert status["has_errors"] is False + assert status["total_steps"] == len(MAIN_PROCESSING_STEPS) + assert status["completed_steps"] == 0 + assert status["in_progress_steps"] == 0 + + def test_get_file_overall_status_processing(self, db_session: Session): + """Test overall status for a file with in-progress steps.""" + file_record = FileRecord( + filehash="test252627", original_filename="test9.pdf", local_filename="/tmp/test9.pdf", file_size=262144 + ) + db_session.add(file_record) + db_session.commit() + + initialize_file_steps(db_session, file_record.id) + + # Mark some steps as complete and one as in progress + update_step_status(db_session, file_record.id, "hash_file", "success") + update_step_status(db_session, file_record.id, "create_file_record", "success") + update_step_status(db_session, file_record.id, "check_text", "in_progress") + + status = get_file_overall_status(db_session, file_record.id) + + assert status["status"] == "processing" + assert status["has_errors"] is False + assert status["completed_steps"] == 2 + assert status["in_progress_steps"] == 1 + + def test_get_file_overall_status_completed(self, db_session: Session): + """Test overall status for a completed file.""" + file_record = FileRecord( + filehash="test282930", original_filename="test10.pdf", local_filename="/tmp/test10.pdf", file_size=524288 + ) + db_session.add(file_record) + db_session.commit() + + initialize_file_steps(db_session, file_record.id) + + # Mark all steps as success + for step_name in MAIN_PROCESSING_STEPS: + update_step_status(db_session, file_record.id, step_name, "success") + + status = get_file_overall_status(db_session, file_record.id) + + assert status["status"] == "completed" + assert status["has_errors"] is False + assert status["completed_steps"] == len(MAIN_PROCESSING_STEPS) + assert status["in_progress_steps"] == 0 + + def test_get_file_overall_status_failed(self, db_session: Session): + """Test overall status for a file with failed steps.""" + file_record = FileRecord( + filehash="test313233", original_filename="test11.pdf", local_filename="/tmp/test11.pdf", file_size=1048576 + ) + db_session.add(file_record) + db_session.commit() + + initialize_file_steps(db_session, file_record.id) + + # Mark some steps as success and one as failure + update_step_status(db_session, file_record.id, "hash_file", "success") + update_step_status(db_session, file_record.id, "create_file_record", "success") + update_step_status(db_session, file_record.id, "check_text", "failure", error_message="OCR failed") + + status = get_file_overall_status(db_session, file_record.id) + + assert status["status"] == "failed" + assert status["has_errors"] is True + assert status["failed_steps"] == 1 + + def test_get_step_summary(self, db_session: Session): + """Test getting step summary with counts.""" + file_record = FileRecord( + filehash="test343536", original_filename="test12.pdf", local_filename="/tmp/test12.pdf", file_size=2097152 + ) + db_session.add(file_record) + db_session.commit() + + # Initialize main steps and add upload steps + initialize_file_steps(db_session, file_record.id) + add_upload_steps(db_session, file_record.id, ["dropbox", "s3", "nextcloud"]) + + # Update statuses + update_step_status(db_session, file_record.id, "hash_file", "success") + update_step_status(db_session, file_record.id, "create_file_record", "success") + update_step_status(db_session, file_record.id, "check_text", "in_progress") + update_step_status(db_session, file_record.id, "upload_to_dropbox", "success") + update_step_status(db_session, file_record.id, "upload_to_s3", "failure") + update_step_status(db_session, file_record.id, "queue_nextcloud", "in_progress") + + summary = get_step_summary(db_session, file_record.id) + + # Check main step counts + assert summary["main"]["success"] == 2 + assert summary["main"]["in_progress"] == 1 + assert summary["main"]["queued"] >= 5 # Remaining pending steps + + # Check upload counts + assert summary["uploads"]["success"] == 1 + assert summary["uploads"]["failure"] == 1 + assert summary["uploads"]["in_progress"] == 1 + + assert summary["total_main_steps"] == len(MAIN_PROCESSING_STEPS) + assert summary["total_upload_tasks"] == 6 # 3 destinations x 2 steps each