"""Unit tests for app.scheduler module.""" import asyncio from datetime import datetime, timezone from types import SimpleNamespace import pytest from app import job_catalog from app.scheduler import ( _DAILY_PIPELINE_STEPS, _NEAR_CLOSE_PIPELINE_STEPS, _consume_backtest_options, _consume_backtest_target_model, _parse_frequency, _repause_after_manual_run, _resume_tickers, _last_successful, _run_source_import, run_sec_fundamentals_import, configure_scheduler, get_job_runtime_snapshot, queue_backtest_options, queue_backtest_target_model, scheduler, ) from app.services.data_import import STATUS_DEFERRED def test_manual_backtest_target_model_is_one_shot(): assert queue_backtest_target_model("structural_sr") == "structural_sr" assert _consume_backtest_target_model() == "structural_sr" assert _consume_backtest_target_model() == "production_gtl" def test_manual_backtest_target_model_rejects_removed_research_arms(): with pytest.raises(ValueError, match="Unknown backtest target model"): queue_backtest_target_model("production_control") def test_manual_backtest_options_are_one_shot_and_default_back_to_weekly(): assert queue_backtest_options("structural_sr", "daily") == ( "structural_sr", "daily", ) assert _consume_backtest_options() == ("structural_sr", "daily") assert _consume_backtest_options() == ("production_gtl", "weekly") def test_only_near_close_fetch_skips_redundant_sr_refresh(): assert dict(_DAILY_PIPELINE_STEPS)["data_collector"] == "collect_ohlcv" assert ( dict(_NEAR_CLOSE_PIPELINE_STEPS)["data_collector"] == "collect_ohlcv_for_scan" ) class TestParseFrequency: def test_hourly(self): assert _parse_frequency("hourly") == {"hours": 1} def test_daily(self): assert _parse_frequency("daily") == {"hours": 24} def test_case_insensitive(self): assert _parse_frequency("Hourly") == {"hours": 1} assert _parse_frequency("DAILY") == {"hours": 24} def test_weekly_maps_to_one_week(self): assert _parse_frequency("weekly") == {"weeks": 1} def test_unknown_defaults_to_daily(self): assert _parse_frequency("monthly") == {"hours": 24} assert _parse_frequency("") == {"hours": 24} class TestResumeTickers: def test_no_previous_returns_full_list(self): symbols = ["AAPL", "GOOG", "MSFT"] _last_successful["test_job"] = None result = _resume_tickers(symbols, "test_job") assert result == ["AAPL", "GOOG", "MSFT"] def test_resume_after_first(self): symbols = ["AAPL", "GOOG", "MSFT"] _last_successful["test_job"] = "AAPL" result = _resume_tickers(symbols, "test_job") # Should start from GOOG, then wrap around assert result == ["GOOG", "MSFT", "AAPL"] def test_resume_after_middle(self): symbols = ["AAPL", "GOOG", "MSFT", "TSLA"] _last_successful["test_job"] = "GOOG" result = _resume_tickers(symbols, "test_job") assert result == ["MSFT", "TSLA", "AAPL", "GOOG"] def test_resume_after_last(self): symbols = ["AAPL", "GOOG", "MSFT"] _last_successful["test_job"] = "MSFT" result = _resume_tickers(symbols, "test_job") # All already processed, wraps to full list assert result == ["AAPL", "GOOG", "MSFT"] def test_unknown_last_returns_full_list(self): symbols = ["AAPL", "GOOG", "MSFT"] _last_successful["test_job"] = "NVDA" result = _resume_tickers(symbols, "test_job") assert result == ["AAPL", "GOOG", "MSFT"] def test_empty_list(self): _last_successful["test_job"] = "AAPL" result = _resume_tickers([], "test_job") assert result == [] class TestConfigureScheduler: def test_configure_adds_all_jobs(self): # Derived from the catalog, not a fourth hand-maintained copy of the # job list: a job added to the catalog but never registered now fails # here instead of silently rendering "Not registered" in the admin UI. scheduler.remove_all_jobs() configure_scheduler() assert {j.id for j in scheduler.get_jobs()} == set(job_catalog.VALID_JOB_NAMES) def test_configure_is_idempotent(self): scheduler.remove_all_jobs() configure_scheduler() configure_scheduler() # Should replace, not duplicate job_ids = [j.id for j in scheduler.get_jobs()] assert sorted(job_ids) == sorted(job_catalog.VALID_JOB_NAMES) def test_independent_jobs_use_cron_not_interval(self): """Interval countdowns restart on every deploy, so a weekly interval on a frequently-redeployed box can defer forever. Both standalone jobs were migrated to cron; this pins them there.""" scheduler.remove_all_jobs() configure_scheduler() for job_id in ("backtest", "ticker_universe_sync"): trigger = type(scheduler.get_job(job_id).trigger).__name__ assert trigger == "CronTrigger", f"{job_id} regressed to {trigger}" class TestJobCatalog: def test_pipeline_members_are_derived_from_step_lists(self): derived = { step for steps in job_catalog.PIPELINE_STEPS.values() for step, _ in steps } assert job_catalog.PIPELINE_MEMBERS == derived # ...and reproduces the set that used to be maintained by hand, so the # derivation is behaviour-preserving rather than merely self-consistent. assert job_catalog.PIPELINE_MEMBERS == { "data_collector", "benchmark_collector", "sentiment_collector", "rr_scanner", "shadow_book", "outcome_evaluator", "alerts", "market_regime", "regime_monitor", } def test_categories_partition_every_job_exactly_once(self): buckets = [ job_catalog.PIPELINE_JOBS, job_catalog.PIPELINE_STEP_JOBS, job_catalog.SCHEDULED_JOBS, job_catalog.MANUAL_JOBS, ] flat = [name for bucket in buckets for name in bucket] assert len(flat) == len(set(flat)), "a job is in two categories" assert set(flat) == set(job_catalog.VALID_JOB_NAMES) assert all(name in job_catalog.JOB_CATEGORY for name in flat) def test_every_job_has_a_label_and_a_unique_sort_order(self): names = job_catalog.VALID_JOB_NAMES assert set(job_catalog.JOB_LABELS) == set(names) assert len({job_catalog.sort_order(n) for n in names}) == len(names) def test_multi_pipeline_members_report_every_parent(self): """Membership is many-to-many — the reason the UI groups into sections rather than nesting steps under one parent.""" by_member = job_catalog.PIPELINES_BY_MEMBER assert set(by_member["data_collector"]) == set(job_catalog.PIPELINE_JOBS) assert set(by_member["alerts"]) == {"daily_pipeline", "near_close_pipeline"} assert set(by_member["outcome_evaluator"]) == { "intraday_pipeline", "after_close_pipeline", } assert "backtest" not in by_member def test_every_job_has_a_runtime_row_before_it_first_runs(self): """The old private _JOB_NAMES list held 16 of 19, so three jobs showed no last-run line until their first run in a given process.""" assert set(get_job_runtime_snapshot()) == set(job_catalog.VALID_JOB_NAMES) class TestRepauseListener: def _configured(self): scheduler.remove_all_jobs() configure_scheduler() def test_manual_job_is_repaused_after_running(self): """Triggering a paused job re-arms its 520-week backstop, which used to surface as a "next run in ~87600h".""" self._configured() scheduler.modify_job("event_study", next_run_time=datetime.now(timezone.utc)) _repause_after_manual_run(SimpleNamespace(job_id="event_study")) assert scheduler.get_job("event_study").next_run_time is None def test_pipeline_step_is_repaused_after_running(self): self._configured() scheduler.modify_job("rr_scanner", next_run_time=datetime.now(timezone.utc)) _repause_after_manual_run(SimpleNamespace(job_id="rr_scanner")) assert scheduler.get_job("rr_scanner").next_run_time is None def test_cron_jobs_are_left_alone(self): # Set an explicit next run first: an unstarted scheduler leaves the # attribute unset, so comparing None to None would prove nothing. self._configured() due = datetime.now(timezone.utc) scheduler.modify_job("daily_pipeline", next_run_time=due) _repause_after_manual_run(SimpleNamespace(job_id="daily_pipeline")) assert scheduler.get_job("daily_pipeline").next_run_time == due def test_unknown_job_is_ignored(self): self._configured() _repause_after_manual_run(SimpleNamespace(job_id="not_a_job")) class _SessionContext: async def __aenter__(self): return object() async def __aexit__(self, *exc): return None class TestSourceImportJobs: @staticmethod def _session_factory(): return _SessionContext() async def test_promoted_run_surfaces_completion(self, monkeypatch): async def enabled(db, job_name): return True async def imported(importer): return SimpleNamespace( status="promoted", revision="abcdef1234567890", error_details=None ) monkeypatch.setattr("app.scheduler.async_session_factory", self._session_factory) monkeypatch.setattr("app.scheduler._is_job_enabled", enabled) monkeypatch.setattr("app.scheduler.run_import", imported) await _run_source_import("dolt_earnings_import", object()) runtime = get_job_runtime_snapshot("dolt_earnings_import") assert runtime["status"] == "completed" assert runtime["processed"] == 1 assert runtime["message"] == "promoted · abcdef123456" async def test_failed_run_surfaces_error(self, monkeypatch): async def enabled(db, job_name): return True async def imported(importer): return SimpleNamespace( status="failed", revision=None, error_details="validation failed" ) monkeypatch.setattr("app.scheduler.async_session_factory", self._session_factory) monkeypatch.setattr("app.scheduler._is_job_enabled", enabled) monkeypatch.setattr("app.scheduler.run_import", imported) await _run_source_import("sec_fundamentals_import", object()) runtime = get_job_runtime_snapshot("sec_fundamentals_import") assert runtime["status"] == "error" assert runtime["processed"] == 0 assert runtime["message"] == "validation failed" async def test_deferred_run_is_visible_without_error_status(self, monkeypatch): async def enabled(db, job_name): return True async def imported(importer): return SimpleNamespace( status=STATUS_DEFERRED, revision="abcdef1234567890", error_details="Company Facts publication lag; retrying", ) monkeypatch.setattr("app.scheduler.async_session_factory", self._session_factory) monkeypatch.setattr("app.scheduler._is_job_enabled", enabled) monkeypatch.setattr("app.scheduler.run_import", imported) await _run_source_import("sec_fundamentals_import", object()) runtime = get_job_runtime_snapshot("sec_fundamentals_import") assert runtime["status"] == STATUS_DEFERRED assert runtime["processed"] == 0 assert runtime["message"] == "Company Facts publication lag; retrying" async def test_source_lock_surfaces_skipped(self, monkeypatch): async def enabled(db, job_name): return True async def imported(importer): return None monkeypatch.setattr("app.scheduler.async_session_factory", self._session_factory) monkeypatch.setattr("app.scheduler._is_job_enabled", enabled) monkeypatch.setattr("app.scheduler.run_import", imported) await _run_source_import("dolt_earnings_import", object()) runtime = get_job_runtime_snapshot("dolt_earnings_import") assert runtime["status"] == "skipped" assert "already running" in runtime["message"] async def test_disabled_job_never_runs_importer(self, monkeypatch): async def disabled(db, job_name): return False async def should_not_run(importer): raise AssertionError("disabled job ran importer") monkeypatch.setattr("app.scheduler.async_session_factory", self._session_factory) monkeypatch.setattr("app.scheduler._is_job_enabled", disabled) monkeypatch.setattr("app.scheduler.run_import", should_not_run) await _run_source_import("sec_fundamentals_import", object()) runtime = get_job_runtime_snapshot("sec_fundamentals_import") assert runtime["status"] == "skipped" assert runtime["message"] == "Disabled" async def test_sec_failure_still_runs_local_cache_refresh(self, monkeypatch): calls = [] events = [] async def enabled(db, job_name): return True async def unavailable(importer): raise RuntimeError("SEC unavailable") async def refreshed(db): calls.append(db) return { "refreshed": 511, "score_inputs_changed": 2, "dimension_scores_staled": 2, "composite_scores_staled": 2, } async def record(**kwargs): events.append(kwargs) monkeypatch.setattr("app.scheduler.async_session_factory", self._session_factory) monkeypatch.setattr("app.scheduler._is_job_enabled", enabled) monkeypatch.setattr("app.scheduler.run_import", unavailable) monkeypatch.setattr("app.scheduler._record_system_event", record) monkeypatch.setattr( "app.scheduler.fundamental_data_refresh_service.refresh", refreshed, ) await run_sec_fundamentals_import() await asyncio.sleep(0) # let the fire-and-forget event task run assert len(calls) == 1 runtime = get_job_runtime_snapshot("sec_fundamentals_import") assert runtime["status"] == "error" # the failure stays the headline, but the cache result is still visible assert runtime["message"] == ( "SEC unavailable · cache 511 · 2 score inputs changed" ) # Rewording the outcome must not duplicate the durable event: the dedup # key includes the message, so a second finish would show up twice in # Admin → System Events. assert len(events) == 1, events async def test_sec_success_surfaces_cache_refresh_summary(self, monkeypatch): async def enabled(db, job_name): return True async def imported(importer): return SimpleNamespace( status="no_op", revision="abcdef1234567890", error_details=None ) async def refreshed(db): return { "refreshed": 511, "score_inputs_changed": 2, "dimension_scores_staled": 2, "composite_scores_staled": 2, } monkeypatch.setattr("app.scheduler.async_session_factory", self._session_factory) monkeypatch.setattr("app.scheduler._is_job_enabled", enabled) monkeypatch.setattr("app.scheduler.run_import", imported) monkeypatch.setattr( "app.scheduler.fundamental_data_refresh_service.refresh", refreshed, ) await run_sec_fundamentals_import() runtime = get_job_runtime_snapshot("sec_fundamentals_import") assert runtime["status"] == "completed" assert runtime["message"] == ( "no_op · abcdef123456 · cache 511 · 2 score inputs changed" ) async def test_source_locked_sec_run_still_reports_the_cache_refresh( self, monkeypatch ): """A skipped import keeps its skip status but shows the cache advanced.""" async def enabled(db, job_name): return True async def locked(importer): return None # another import owns the source lock async def refreshed(db): return { "refreshed": 511, "score_inputs_changed": 0, "dimension_scores_staled": 0, "composite_scores_staled": 0, } monkeypatch.setattr("app.scheduler.async_session_factory", self._session_factory) monkeypatch.setattr("app.scheduler._is_job_enabled", enabled) monkeypatch.setattr("app.scheduler.run_import", locked) monkeypatch.setattr( "app.scheduler.fundamental_data_refresh_service.refresh", refreshed, ) await run_sec_fundamentals_import() runtime = get_job_runtime_snapshot("sec_fundamentals_import") assert runtime["status"] == "skipped" assert runtime["message"] == ( "Another import for this source is already running · " "cache 511 · 0 score inputs changed" ) async def test_disabled_sec_job_still_refreshes_local_cache(self, monkeypatch): """Disabling the job stops the SEC fetch, not the local cache. The cache is derived from stored snapshots, earnings events and closes. Prices and earnings move daily even when no filing does, and there is no provider fallback since A6 — freezing it would silently stale scoring. """ calls = [] async def disabled(db, job_name): return False async def should_not_run(*args, **kwargs): raise AssertionError("disabled SEC job hit the network") async def refreshed(db): calls.append(db) return { "refreshed": 511, "score_inputs_changed": 2, "dimension_scores_staled": 2, "composite_scores_staled": 2, } monkeypatch.setattr("app.scheduler.async_session_factory", self._session_factory) monkeypatch.setattr("app.scheduler._is_job_enabled", disabled) monkeypatch.setattr("app.scheduler.run_import", should_not_run) monkeypatch.setattr( "app.scheduler.fundamental_data_refresh_service.refresh", refreshed, ) await run_sec_fundamentals_import() assert len(calls) == 1 runtime = get_job_runtime_snapshot("sec_fundamentals_import") assert runtime["status"] == "completed" assert runtime["message"] == ( "Import disabled · cache 511 · 2 score inputs changed" )