356 lines
12 KiB
Python
356 lines
12 KiB
Python
from __future__ import annotations
|
|
|
|
import sys
|
|
import unittest
|
|
import csv
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from tempfile import TemporaryDirectory
|
|
from unittest.mock import call, patch
|
|
|
|
|
|
SCRIPTS_DIR = Path(__file__).resolve().parents[1] / "scripts"
|
|
sys.path.insert(0, str(SCRIPTS_DIR))
|
|
|
|
from request_nsrdb_county_polygon_archives import ( # noqa: E402
|
|
ArchiveQueueMonitor,
|
|
CountyRequestEvent,
|
|
CountyRequestState,
|
|
CountyRequestStateMachine,
|
|
EXCEPTIONAL_POLYGON_WAIT,
|
|
InvalidCountyRequestTransition,
|
|
LARGE_POLYGON_WAIT,
|
|
LocalQueueCapacityError,
|
|
MEDIUM_POLYGON_WAIT,
|
|
SMALL_POLYGON_WAIT,
|
|
)
|
|
|
|
|
|
class CountyRequestStateMachineTests(unittest.TestCase):
|
|
def test_successful_request_lifecycle(self) -> None:
|
|
machine = CountyRequestStateMachine("01001 Autauga County, AL")
|
|
|
|
machine.transition(CountyRequestEvent.START)
|
|
machine.transition(CountyRequestEvent.REQUEST_SUCCEEDED)
|
|
|
|
self.assertEqual(machine.state, CountyRequestState.REQUESTED)
|
|
self.assertEqual(
|
|
machine.history,
|
|
[
|
|
CountyRequestState.PENDING,
|
|
CountyRequestState.SUBMITTING,
|
|
CountyRequestState.REQUESTED,
|
|
],
|
|
)
|
|
|
|
def test_retry_returns_to_submitting(self) -> None:
|
|
machine = CountyRequestStateMachine("01001 Autauga County, AL")
|
|
|
|
machine.transition(CountyRequestEvent.START)
|
|
machine.transition(CountyRequestEvent.RETRY_REQUIRED)
|
|
machine.transition(CountyRequestEvent.RETRY_STARTED)
|
|
|
|
self.assertEqual(machine.state, CountyRequestState.SUBMITTING)
|
|
|
|
def test_existing_request_is_skipped(self) -> None:
|
|
machine = CountyRequestStateMachine("01001 Autauga County, AL")
|
|
|
|
machine.transition(CountyRequestEvent.EXISTING_REQUEST_FOUND)
|
|
|
|
self.assertEqual(machine.state, CountyRequestState.SKIPPED)
|
|
|
|
def test_invalid_transition_is_rejected(self) -> None:
|
|
machine = CountyRequestStateMachine("01001 Autauga County, AL")
|
|
|
|
with self.assertRaisesRegex(
|
|
InvalidCountyRequestTransition,
|
|
"request_succeeded while pending",
|
|
):
|
|
machine.transition(CountyRequestEvent.REQUEST_SUCCEEDED)
|
|
|
|
|
|
class ArchiveQueueMonitorTests(unittest.TestCase):
|
|
def make_monitor(self) -> ArchiveQueueMonitor:
|
|
temp_dir = TemporaryDirectory()
|
|
self.addCleanup(temp_dir.cleanup)
|
|
return ArchiveQueueMonitor(Path(temp_dir.name), timeout=1, workers=1)
|
|
|
|
def test_dynamic_wait_uses_largest_tracked_polygon(self) -> None:
|
|
monitor = self.make_monitor()
|
|
|
|
for site_count, expected_wait in (
|
|
(110, SMALL_POLYGON_WAIT),
|
|
(111, MEDIUM_POLYGON_WAIT),
|
|
(250, MEDIUM_POLYGON_WAIT),
|
|
(251, LARGE_POLYGON_WAIT),
|
|
(550, LARGE_POLYGON_WAIT),
|
|
(551, EXCEPTIONAL_POLYGON_WAIT),
|
|
):
|
|
monitor.tracked_site_counts = {"https://example.com/job.zip": site_count}
|
|
self.assertEqual(monitor.dynamic_wait_seconds(), expected_wait)
|
|
|
|
def test_downloaded_archive_is_not_reloaded_as_pending(self) -> None:
|
|
temp_dir = TemporaryDirectory()
|
|
self.addCleanup(temp_dir.cleanup)
|
|
manifest_dir = Path(temp_dir.name)
|
|
response_dir = manifest_dir / "polygon_cloud_request_responses"
|
|
archive_dir = manifest_dir / "polygon_cloud_archives"
|
|
response_dir.mkdir()
|
|
archive_dir.mkdir()
|
|
response_path = response_dir / "01001_goes-tmy_tmy-2024_cloud_response.json"
|
|
archive_path = archive_dir / "01001_goes-tmy_tmy-2024_cloud.zip"
|
|
response_path.write_text("{}", encoding="utf-8")
|
|
archive_path.write_bytes(b"downloaded")
|
|
|
|
manifest_path = manifest_dir / "county_polygon_cloud_request_manifest.csv"
|
|
with manifest_path.open("w", newline="", encoding="utf-8") as handle:
|
|
writer = csv.DictWriter(
|
|
handle,
|
|
fieldnames=[
|
|
"county_fips",
|
|
"site_count",
|
|
"download_url",
|
|
"response_json",
|
|
"submitted_at_utc",
|
|
],
|
|
)
|
|
writer.writeheader()
|
|
writer.writerow(
|
|
{
|
|
"county_fips": "01001",
|
|
"site_count": "100",
|
|
"download_url": "https://example.amazonaws.com/job.zip",
|
|
"response_json": str(response_path),
|
|
"submitted_at_utc": datetime.now(timezone.utc).isoformat(),
|
|
}
|
|
)
|
|
|
|
monitor = ArchiveQueueMonitor(manifest_dir, timeout=1, workers=1)
|
|
|
|
with patch(
|
|
"request_nsrdb_county_polygon_archives.archive_download_status",
|
|
return_value="pending",
|
|
) as status_mock:
|
|
counts = monitor.summarize()
|
|
|
|
self.assertEqual(counts["pending"], 0)
|
|
self.assertEqual(counts["total"], 0)
|
|
status_mock.assert_not_called()
|
|
|
|
def test_historical_s3_403_does_not_consume_a_queue_slot(self) -> None:
|
|
monitor = self.make_monitor()
|
|
url = "https://example.amazonaws.com/historical.zip"
|
|
monitor.jobs[url] = "01001"
|
|
|
|
with patch(
|
|
"request_nsrdb_county_polygon_archives.archive_download_status",
|
|
return_value="pending",
|
|
):
|
|
counts = monitor.summarize()
|
|
|
|
self.assertEqual(counts["pending"], 0)
|
|
self.assertEqual(counts["ambiguous_403"], 1)
|
|
|
|
def test_current_process_s3_403_remains_pending(self) -> None:
|
|
monitor = self.make_monitor()
|
|
url = "https://example.amazonaws.com/current.zip"
|
|
monitor.record_submission(
|
|
{
|
|
"county_fips": "01001",
|
|
"tile_id": "",
|
|
"site_count": "100",
|
|
"download_url": url,
|
|
}
|
|
)
|
|
|
|
with patch(
|
|
"request_nsrdb_county_polygon_archives.archive_download_status",
|
|
return_value="pending",
|
|
):
|
|
counts = monitor.summarize()
|
|
|
|
self.assertEqual(counts["pending"], 1)
|
|
self.assertEqual(counts["ambiguous_403"], 0)
|
|
|
|
def test_site_count_expires_after_recent_download_window(self) -> None:
|
|
monitor = self.make_monitor()
|
|
url = "https://example.com/job.zip"
|
|
monitor.jobs[url] = "01001"
|
|
monitor.pending_urls.add(url)
|
|
monitor.tracked_site_counts[url] = 551
|
|
|
|
with (
|
|
patch(
|
|
"request_nsrdb_county_polygon_archives.archive_download_status",
|
|
return_value="ready",
|
|
),
|
|
patch(
|
|
"request_nsrdb_county_polygon_archives.time.monotonic",
|
|
side_effect=[100.0, 100.0, 401.0],
|
|
),
|
|
):
|
|
first_counts = monitor.summarize()
|
|
second_counts = monitor.summarize()
|
|
|
|
self.assertEqual(first_counts["recently_downloaded"], 1)
|
|
self.assertEqual(second_counts["recently_downloaded"], 0)
|
|
self.assertNotIn(url, monitor.tracked_site_counts)
|
|
|
|
def test_fifteen_visible_slots_bypass_pacing_until_queue_is_full(self) -> None:
|
|
monitor = self.make_monitor()
|
|
monitor.paced_submissions = 1
|
|
counts_with_fifteen_slots = {
|
|
"pending": 5,
|
|
"recently_downloaded": 2,
|
|
"unobserved": 0,
|
|
"unknown": 0,
|
|
"ready": 0,
|
|
"total": 7,
|
|
}
|
|
counts_with_fourteen_slots = {
|
|
**counts_with_fifteen_slots,
|
|
"pending": 6,
|
|
}
|
|
full_queue_counts = {
|
|
**counts_with_fifteen_slots,
|
|
"pending": 20,
|
|
}
|
|
|
|
with (
|
|
patch.object(
|
|
monitor,
|
|
"summarize",
|
|
side_effect=[
|
|
counts_with_fifteen_slots,
|
|
counts_with_fourteen_slots,
|
|
full_queue_counts,
|
|
],
|
|
),
|
|
patch(
|
|
"request_nsrdb_county_polygon_archives.wait_for_queue",
|
|
) as wait_for_queue_mock,
|
|
):
|
|
monitor.wait_for_capacity("first", retries=0)
|
|
monitor.wait_for_capacity("second", retries=0)
|
|
with self.assertRaisesRegex(LocalQueueCapacityError, "20 possible"):
|
|
monitor.wait_for_capacity("third", retries=0)
|
|
|
|
wait_for_queue_mock.assert_not_called()
|
|
self.assertFalse(monitor.filling_open_queue)
|
|
|
|
def test_fourteen_visible_slots_still_use_recent_download_pacing(self) -> None:
|
|
monitor = self.make_monitor()
|
|
monitor.paced_submissions = 1
|
|
counts = {
|
|
"pending": 6,
|
|
"recently_downloaded": 2,
|
|
"unobserved": 0,
|
|
"unknown": 0,
|
|
"ready": 0,
|
|
"total": 8,
|
|
}
|
|
|
|
with (
|
|
patch.object(monitor, "summarize", return_value=counts),
|
|
patch(
|
|
"request_nsrdb_county_polygon_archives.wait_for_queue",
|
|
) as wait_for_queue_mock,
|
|
):
|
|
monitor.wait_for_capacity("next", retries=0)
|
|
|
|
wait_for_queue_mock.assert_called_once_with(
|
|
SMALL_POLYGON_WAIT,
|
|
"Pacing wait complete; refreshing archive status.",
|
|
)
|
|
|
|
def test_six_minute_wait_checks_queue_after_three_minutes(self) -> None:
|
|
monitor = self.make_monitor()
|
|
counts = {
|
|
"pending": 6,
|
|
"recently_downloaded": 0,
|
|
"unobserved": 0,
|
|
"unknown": 0,
|
|
"ready": 0,
|
|
"total": 6,
|
|
}
|
|
|
|
with (
|
|
patch.object(monitor, "summarize", return_value=counts) as summarize_mock,
|
|
patch(
|
|
"request_nsrdb_county_polygon_archives.time.sleep",
|
|
) as sleep_mock,
|
|
):
|
|
entered_fill_mode = monitor.wait_for_dynamic_queue(
|
|
LARGE_POLYGON_WAIT,
|
|
"Queue wait complete.",
|
|
)
|
|
|
|
self.assertFalse(entered_fill_mode)
|
|
self.assertEqual(sleep_mock.call_args_list, [call(180.0)] * 2)
|
|
summarize_mock.assert_called_once_with()
|
|
|
|
def test_six_minute_wait_enters_fill_mode_at_checkpoint(self) -> None:
|
|
monitor = self.make_monitor()
|
|
counts = {
|
|
"pending": 5,
|
|
"recently_downloaded": 0,
|
|
"unobserved": 0,
|
|
"unknown": 0,
|
|
"ready": 0,
|
|
"total": 5,
|
|
}
|
|
|
|
with (
|
|
patch.object(monitor, "summarize", return_value=counts),
|
|
patch(
|
|
"request_nsrdb_county_polygon_archives.time.sleep",
|
|
) as sleep_mock,
|
|
):
|
|
entered_fill_mode = monitor.wait_for_dynamic_queue(
|
|
LARGE_POLYGON_WAIT,
|
|
"Queue wait complete.",
|
|
)
|
|
|
|
self.assertTrue(entered_fill_mode)
|
|
self.assertTrue(monitor.filling_open_queue)
|
|
sleep_mock.assert_called_once_with(180.0)
|
|
|
|
def test_fifteen_minute_wait_checks_at_five_and_ten_minutes(self) -> None:
|
|
monitor = self.make_monitor()
|
|
counts_with_fourteen_slots = {
|
|
"pending": 6,
|
|
"recently_downloaded": 0,
|
|
"unobserved": 0,
|
|
"unknown": 0,
|
|
"ready": 0,
|
|
"total": 6,
|
|
}
|
|
counts_with_fifteen_slots = {
|
|
**counts_with_fourteen_slots,
|
|
"pending": 5,
|
|
}
|
|
|
|
with (
|
|
patch.object(
|
|
monitor,
|
|
"summarize",
|
|
side_effect=[counts_with_fourteen_slots, counts_with_fifteen_slots],
|
|
) as summarize_mock,
|
|
patch(
|
|
"request_nsrdb_county_polygon_archives.time.sleep",
|
|
) as sleep_mock,
|
|
):
|
|
entered_fill_mode = monitor.wait_for_dynamic_queue(
|
|
EXCEPTIONAL_POLYGON_WAIT,
|
|
"Queue wait complete.",
|
|
)
|
|
|
|
self.assertTrue(entered_fill_mode)
|
|
self.assertTrue(monitor.filling_open_queue)
|
|
self.assertEqual(sleep_mock.call_args_list, [call(300.0)] * 2)
|
|
self.assertEqual(summarize_mock.call_count, 2)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|