import pytest from datetime import datetime from dataclasses import dataclass import logging from bt_df_data_retention_manager.jobs.engines.interface import GenericJob, GenericJobRuntimeParams from bt_df_data_retention_manager.state.record.retention import JobStatusEnum from bt_df_data_retention_manager.jobs.job_processor import JobProcessor from bt_df_data_retention_manager.state.record.retention import TableRetentionState logging.basicConfig(level=logging.INFO) @dataclass(slots=True) class TestParams(GenericJobRuntimeParams): pass @dataclass(slots=True) class TestJob(GenericJob): """Class for testing state machine with predefined behavior""" job_params: GenericJobRuntimeParams timed_out: bool in_progress: bool failed: bool completed: bool def set_job_params(self, params: TestParams) -> None: self.job_params = params def is_timed_out(self) -> bool: return self.timed_out def is_running(self) -> bool: return self.in_progress def is_failed(self) -> bool: return self.failed def is_completed(self) -> bool: return self.completed def create(self) -> None: print("Job has been created") self.in_progress = True self.failed = False self.completed = False def cancel(self) -> None: print("Job has been canceled") self.in_progress = False self.failed = True self.completed = False def set_job_completed(self): self.in_progress = False self.failed = False self.completed = True def get_job_id(self) -> str: return 'test_job_id' def get_end_retention_date(self) -> datetime: return datetime(year=2022, month=11, day=30) @pytest.fixture def mock_processor() -> JobProcessor: params = TestParams() job = TestJob(job_params=params, completed=False, failed=False, in_progress=False, timed_out=False) state = TableRetentionState(table_id="charts", job_status=JobStatusEnum.ready, latest_job_id='', current_end_date=job.get_end_retention_date(), retries=0) return JobProcessor( job=job, table_state=state, runnable=True, max_retries=1, ) def test_job_processor_timeout_and_retries(mock_processor: JobProcessor): """Test start, timeout and retry behavior in job processor""" assert mock_processor.is_ready() mock_processor.to_next_state() assert mock_processor.is_in_progress() # Test cancellation via timeout mock_processor.job.timed_out = True mock_processor.to_next_state() assert mock_processor.is_failed() # Moving from FAILED to READY mock_processor.to_next_state() assert mock_processor.is_ready() # Start new job mock_processor.to_next_state() assert mock_processor.is_in_progress() # And cancel it by timeout once again mock_processor.to_next_state() assert mock_processor.is_failed() # Try to reset the job mock_processor.to_next_state() # Should still be failed since there was one retry assert mock_processor.is_failed() def test_job_processor_completion_update_params(mock_processor: JobProcessor): """Test start, timeout and retry behavior in job processor""" mock_processor.to_next_state() mock_processor.to_next_state() assert mock_processor.is_in_progress() mock_processor.job.set_job_completed() mock_processor.to_next_state() assert mock_processor.is_completed() # Should stay completed mock_processor.to_next_state() assert mock_processor.is_completed() mock_processor.table_state.current_end_date = datetime.now() mock_processor.to_next_state() assert mock_processor.is_ready() assert mock_processor.table_state.current_end_date == mock_processor.job.get_end_retention_date( )