""" Tests for the FileProcessingStep model and step_manager utilities. This test module verifies the new explicit status tracking approach. """ from datetime import datetime, timedelta from unittest.mock import patch import pytest from sqlalchemy import create_engine from sqlalchemy.orm import Session, sessionmaker from app.config import settings from app.database import Base from app.models import FileProcessingStep, FileRecord from app.utils.step_manager import ( MAIN_PROCESSING_STEPS, TERMINAL_STEP, 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="create_file_record", 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 == "create_file_record" 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="create_file_record", 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="create_file_record", 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, "create_file_record", "in_progress", started_at=now) # Verify step was created step = ( db_session.query(FileProcessingStep) .filter(FileProcessingStep.file_id == file_record.id, FileProcessingStep.step_name == "create_file_record") .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, "create_file_record", "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 == "create_file_record") .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 - use actual step names from MAIN_PROCESSING_STEPS now = datetime.now() update_step_status(db_session, file_record.id, "create_file_record", "success", completed_at=now) update_step_status(db_session, file_record.id, "check_text", "in_progress", started_at=now) update_step_status( db_session, file_record.id, "extract_text", "failure", error_message="Failed to extract 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["create_file_record"]["status"] == "success" assert status_map["create_file_record"]["completed_at"] == now assert status_map["check_text"]["status"] == "in_progress" assert status_map["extract_text"]["status"] == "failure" assert status_map["extract_text"]["error_message"] == "Failed to extract 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, "create_file_record", "success") update_step_status(db_session, file_record.id, "check_text", "success") update_step_status(db_session, file_record.id, "extract_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_completed_with_pending_intermediate_steps(self, db_session: Session): """Test that a file is 'completed' when the terminal step succeeds even if some intermediate steps are pending. This covers the scenario where steps like check_for_duplicates or extract_text are left in 'pending' because the dynamic pipeline skipped them without explicitly marking them (e.g. dedup log written before file_id was available, or extract_text not marked for non-PDF). """ file_record = FileRecord( filehash="test_terminal_complete", original_filename="terminal.pdf", local_filename="/tmp/terminal.pdf", file_size=1024, ) db_session.add(file_record) db_session.commit() initialize_file_steps(db_session, file_record.id) # Simulate a pipeline where most steps succeed but check_for_duplicates # stays pending (the exact scenario from the bug report). for step_name in MAIN_PROCESSING_STEPS: if step_name == "check_for_duplicates": continue # Leave this as "pending" update_step_status(db_session, file_record.id, step_name, "success") # Also add a dynamically created OCR step as "skipped" update_step_status(db_session, file_record.id, "process_with_ocr", "skipped") status = get_file_overall_status(db_session, file_record.id) # The terminal step (send_to_all_destinations) succeeded, so the # file should be "completed" despite check_for_duplicates being pending. assert status["status"] == "completed" assert status["has_errors"] is False def test_get_file_overall_status_pending_without_terminal_step(self, db_session: Session): """Test that a file stays 'pending' when the terminal step hasn't been recorded.""" file_record = FileRecord( filehash="test_no_terminal", original_filename="no_terminal.pdf", local_filename="/tmp/no_terminal.pdf", file_size=1024, ) db_session.add(file_record) db_session.commit() initialize_file_steps(db_session, file_record.id) # Only first few steps completed - terminal step is still pending update_step_status(db_session, file_record.id, "create_file_record", "success") update_step_status(db_session, file_record.id, "check_text", "success") update_step_status(db_session, file_record.id, "extract_text", "success") status = get_file_overall_status(db_session, file_record.id) # Terminal step hasn't run yet, so file should remain pending assert status["status"] == "pending" assert status["has_errors"] is False assert status["completed_steps"] == 3 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, "create_file_record", "success") update_step_status(db_session, file_record.id, "check_text", "success") update_step_status(db_session, file_record.id, "extract_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, "create_file_record", "success") update_step_status(db_session, file_record.id, "check_text", "success") update_step_status(db_session, file_record.id, "extract_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, "upload_to_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 (only upload_to_* steps, not queue_* steps) 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"] == 3 # 3 upload_to_* steps counted def test_get_file_overall_status_duplicate_file(self, db_session: Session): """Test overall status returns 'duplicate' immediately for duplicate files.""" file_record = FileRecord( filehash="test_dup01", original_filename="dup.pdf", local_filename="/tmp/dup.pdf", file_size=1024, is_duplicate=True, ) db_session.add(file_record) db_session.commit() status = get_file_overall_status(db_session, file_record.id) assert status["status"] == "duplicate" assert status["has_errors"] is False assert status["total_steps"] == 0 assert status["completed_steps"] == 0 assert status["failed_steps"] == 0 assert status["in_progress_steps"] == 0 assert status["skipped_steps"] == 0 def test_get_file_overall_status_no_steps_initialized(self, db_session: Session): """Test overall status returns 'pending' with zero counts when no steps exist.""" file_record = FileRecord( filehash="test_nosteps", original_filename="nosteps.pdf", local_filename="/tmp/nosteps.pdf", file_size=1024, ) db_session.add(file_record) db_session.commit() # Do NOT initialize any steps status = get_file_overall_status(db_session, file_record.id) assert status["status"] == "pending" assert status["has_errors"] is False assert status["total_steps"] == 0 assert status["completed_steps"] == 0 def test_get_file_overall_status_all_success_without_terminal_step(self, db_session: Session): """Test overall status stays 'pending' when all steps succeed but terminal step is absent.""" file_record = FileRecord( filehash="test_no_terminal2", original_filename="no_terminal2.pdf", local_filename="/tmp/no_terminal2.pdf", file_size=1024, ) db_session.add(file_record) db_session.commit() # Add only non-terminal steps as success (no TERMINAL_STEP = send_to_all_destinations) for step_name in ["create_file_record", "check_text", "extract_text"]: step = FileProcessingStep(file_id=file_record.id, step_name=step_name, status="success") db_session.add(step) db_session.commit() status = get_file_overall_status(db_session, file_record.id) # All counted steps are success/skipped, but terminal step missing → pending assert status["status"] == "pending" assert status["has_errors"] is False assert status["completed_steps"] == 3 def test_get_file_overall_status_dedup_disabled(self, db_session: Session): """Test that check_for_duplicates is excluded from real steps when dedup is disabled.""" file_record = FileRecord( filehash="test_dedup_off", original_filename="dedup_off.pdf", local_filename="/tmp/dedup_off.pdf", file_size=1024, ) db_session.add(file_record) db_session.commit() # Add check_for_duplicates step as success and terminal step as success db_session.add(FileProcessingStep(file_id=file_record.id, step_name="check_for_duplicates", status="success")) db_session.add(FileProcessingStep(file_id=file_record.id, step_name=TERMINAL_STEP, status="success")) db_session.commit() with patch.object(settings, "enable_deduplication", False): status = get_file_overall_status(db_session, file_record.id) # check_for_duplicates should not be counted; only send_to_all_destinations counts assert status["status"] == "completed" assert status["total_steps"] == 1 # Only the terminal step assert status["completed_steps"] == 1 def test_get_step_summary_dedup_disabled(self, db_session: Session): """Test that get_step_summary excludes check_for_duplicates when dedup is disabled.""" file_record = FileRecord( filehash="test_sum_dedup_off", original_filename="sum_dedup_off.pdf", local_filename="/tmp/sum_dedup_off.pdf", file_size=1024, ) db_session.add(file_record) db_session.commit() db_session.add(FileProcessingStep(file_id=file_record.id, step_name="check_for_duplicates", status="success")) db_session.add(FileProcessingStep(file_id=file_record.id, step_name=TERMINAL_STEP, status="success")) db_session.commit() with patch.object(settings, "enable_deduplication", False): summary = get_step_summary(db_session, file_record.id) # check_for_duplicates should not be counted; only terminal step as a main step assert summary["main"]["success"] == 1 assert summary["total_main_steps"] == 1 def test_add_upload_steps_idempotent(self, db_session: Session): """Test that add_upload_steps does not create duplicate steps if called twice.""" file_record = FileRecord( filehash="test_idem01", original_filename="idem.pdf", local_filename="/tmp/idem.pdf", file_size=1024, ) db_session.add(file_record) db_session.commit() destinations = ["dropbox", "s3"] add_upload_steps(db_session, file_record.id, destinations) add_upload_steps(db_session, file_record.id, destinations) # Second call should be a no-op steps = db_session.query(FileProcessingStep).filter(FileProcessingStep.file_id == file_record.id).all() # Only 2 steps per destination (queue_ and upload_to_), NOT 4 assert len(steps) == len(destinations) * 2 def test_get_step_summary_internal_step_filtered(self, db_session: Session): """Test that non-real internal steps (e.g. poll_task) are ignored in get_step_summary.""" file_record = FileRecord( filehash="test_internal01", original_filename="internal.pdf", local_filename="/tmp/internal.pdf", file_size=1024, ) db_session.add(file_record) db_session.commit() # Add a real step and an internal/diagnostic step db_session.add(FileProcessingStep(file_id=file_record.id, step_name=TERMINAL_STEP, status="success")) db_session.add(FileProcessingStep(file_id=file_record.id, step_name="poll_task", status="success")) db_session.commit() summary = get_step_summary(db_session, file_record.id) # poll_task is not a real step and must be filtered out assert summary["total_main_steps"] == 1 # Only TERMINAL_STEP assert summary["main"]["success"] == 1 def test_get_step_summary_unusual_upload_status(self, db_session: Session): """Test get_step_summary with an upload step whose status is not in the counts dict.""" file_record = FileRecord( filehash="test_unk_upload", original_filename="unk_upload.pdf", local_filename="/tmp/unk_upload.pdf", file_size=1024, ) db_session.add(file_record) db_session.commit() db_session.add(FileProcessingStep(file_id=file_record.id, step_name=TERMINAL_STEP, status="success")) # "cancelled" is not in the upload_counts dict keys db_session.add(FileProcessingStep(file_id=file_record.id, step_name="upload_to_dropbox", status="cancelled")) db_session.commit() summary = get_step_summary(db_session, file_record.id) # The upload step is still counted in total_upload_tasks even with an unusual status assert summary["total_upload_tasks"] == 1 # "cancelled" not in counts, so no key incremented assert summary["uploads"]["success"] == 0 assert summary["uploads"]["failure"] == 0 def test_get_step_summary_unusual_main_status(self, db_session: Session): """Test get_step_summary with a main step whose status is not in the counts dict.""" file_record = FileRecord( filehash="test_unk_main", original_filename="unk_main.pdf", local_filename="/tmp/unk_main.pdf", file_size=1024, ) db_session.add(file_record) db_session.commit() db_session.add(FileProcessingStep(file_id=file_record.id, step_name=TERMINAL_STEP, status="success")) # "cancelled" is not in the main_counts dict keys db_session.add(FileProcessingStep(file_id=file_record.id, step_name="create_file_record", status="cancelled")) db_session.commit() summary = get_step_summary(db_session, file_record.id) # The main step is still counted in total_main_steps even with an unusual status assert summary["total_main_steps"] == 2 # terminal + create_file_record # "cancelled" not in counts, so no key incremented beyond terminal step's success assert summary["main"]["success"] == 1 assert summary["main"]["queued"] == 0 def test_get_step_summary_no_terminal_step_in_db(self, db_session: Session): """Test that get_step_summary adds terminal step as queued when it is absent from DB.""" file_record = FileRecord( filehash="test_no_term_sum", original_filename="no_term_sum.pdf", local_filename="/tmp/no_term_sum.pdf", file_size=1024, ) db_session.add(file_record) db_session.commit() # Add only non-terminal main steps db_session.add(FileProcessingStep(file_id=file_record.id, step_name="create_file_record", status="success")) db_session.commit() summary = get_step_summary(db_session, file_record.id) # Terminal step absent → automatically added as queued (exactly one queued, two total) assert summary["main"]["queued"] == 1 assert summary["total_main_steps"] == 2 # create_file_record + virtual terminal step