97 lines
3.6 KiB
Python
97 lines
3.6 KiB
Python
from __future__ import annotations
|
|
|
|
import tempfile
|
|
import threading
|
|
import unittest
|
|
from pathlib import Path
|
|
|
|
from backend.jobs import InProcessJobRunner, JobRegistry, SQLiteJobRunRepository
|
|
from database import ReviewDatabase
|
|
|
|
|
|
class JobRunnerTests(unittest.TestCase):
|
|
def setUp(self) -> None:
|
|
self.temporary = tempfile.TemporaryDirectory()
|
|
self.addCleanup(self.temporary.cleanup)
|
|
self.database = ReviewDatabase(Path(self.temporary.name) / "review.db")
|
|
self.repository = SQLiteJobRunRepository(self.database)
|
|
self.runner = InProcessJobRunner(JobRegistry.load(), self.repository)
|
|
|
|
def test_successful_idempotent_job_runs_once(self) -> None:
|
|
calls = []
|
|
self.assertTrue(
|
|
self.runner.run_inline(
|
|
"screener.automatic", "20260729:v8", lambda: calls.append("run")
|
|
)
|
|
)
|
|
self.assertFalse(
|
|
self.runner.run_inline(
|
|
"screener.automatic", "20260729:v8", lambda: calls.append("again")
|
|
)
|
|
)
|
|
self.assertEqual(calls, ["run"])
|
|
self.assertEqual(self.repository.recent(1)[0]["status"], "success")
|
|
|
|
def test_failure_is_persisted_and_can_be_retried_later(self) -> None:
|
|
def fail() -> None:
|
|
raise RuntimeError("provider unavailable")
|
|
|
|
self.assertTrue(self.runner.run_inline("market.refresh", "failed-key", fail))
|
|
failed = self.repository.recent(1)[0]
|
|
self.assertEqual(failed["status"], "failed")
|
|
self.assertEqual(failed["error_code"], "RuntimeError")
|
|
self.assertTrue(
|
|
self.runner.run_inline("market.refresh", "failed-key", lambda: None)
|
|
)
|
|
self.assertEqual(self.repository.recent(1)[0]["attempt"], 2)
|
|
|
|
def test_lock_key_rejects_concurrent_submission(self) -> None:
|
|
entered = threading.Event()
|
|
release = threading.Event()
|
|
|
|
def wait() -> None:
|
|
entered.set()
|
|
release.wait(2)
|
|
|
|
self.assertTrue(self.runner.submit("market.refresh", "first", wait))
|
|
self.assertTrue(entered.wait(1))
|
|
self.assertFalse(self.runner.submit("market.refresh", "second", lambda: None))
|
|
release.set()
|
|
self.assertTrue(self.runner.wait_for_idle())
|
|
|
|
def test_failed_status_payload_is_recorded_as_failure(self) -> None:
|
|
self.assertTrue(
|
|
self.runner.run_inline(
|
|
"screener.automatic", "reported-failure",
|
|
lambda: {"status": "failed", "error": "missing factors"},
|
|
)
|
|
)
|
|
self.assertEqual(self.repository.recent(1)[0]["status"], "failed")
|
|
|
|
def test_scheduler_start_is_idempotent_and_stop_waits_for_exit(self) -> None:
|
|
first = self.runner.start_scheduler(
|
|
lambda: None, interval_seconds=60, initial_delay_seconds=60
|
|
)
|
|
second = self.runner.start_scheduler(
|
|
lambda: None, interval_seconds=60, initial_delay_seconds=60
|
|
)
|
|
self.assertIs(first, second)
|
|
self.assertTrue(first.is_alive())
|
|
self.assertTrue(self.runner.stop_scheduler())
|
|
self.assertFalse(first.is_alive())
|
|
|
|
def test_scheduler_can_restart_after_an_orderly_stop(self) -> None:
|
|
first = self.runner.start_scheduler(
|
|
lambda: None, interval_seconds=60, initial_delay_seconds=60
|
|
)
|
|
self.assertTrue(self.runner.stop_scheduler())
|
|
second = self.runner.start_scheduler(
|
|
lambda: None, interval_seconds=60, initial_delay_seconds=60
|
|
)
|
|
self.assertIsNot(first, second)
|
|
self.assertTrue(self.runner.stop_scheduler())
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|