Respect SpaceDevs request limits and cooldowns
This commit is contained in:
parent
69992bebb7
commit
3f6ee4f2ec
3 changed files with 124 additions and 13 deletions
|
|
@ -1,6 +1,9 @@
|
||||||
"""Get details of upcoming space launches."""
|
"""Get details of upcoming space launches."""
|
||||||
|
|
||||||
|
import fcntl
|
||||||
import json
|
import json
|
||||||
|
import re
|
||||||
|
import time
|
||||||
import os
|
import os
|
||||||
import typing
|
import typing
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
|
|
@ -19,6 +22,47 @@ LIMIT = 500
|
||||||
ACTIVE_CREWED_FLIGHTS_CACHE_FILE = "active_crewed_flights.json"
|
ACTIVE_CREWED_FLIGHTS_CACHE_FILE = "active_crewed_flights.json"
|
||||||
|
|
||||||
|
|
||||||
|
class RateLimitDeferred(Exception):
|
||||||
|
"""The shared API budget or server cooldown requires waiting."""
|
||||||
|
|
||||||
|
|
||||||
|
def api_get(
|
||||||
|
rocket_dir: str, url: str, params: dict[str, str | int] | None = None
|
||||||
|
) -> requests.Response:
|
||||||
|
"""Share a conservative rolling request budget across updater processes."""
|
||||||
|
# Keep three calls spare for other clients sharing the public IP.
|
||||||
|
filename = os.path.join(rocket_dir, "api_budget.json")
|
||||||
|
with open(filename + ".lock", "a") as lock:
|
||||||
|
fcntl.flock(lock, fcntl.LOCK_EX)
|
||||||
|
now = time.time()
|
||||||
|
try:
|
||||||
|
with open(filename) as cache:
|
||||||
|
state = json.load(cache)
|
||||||
|
except (OSError, ValueError):
|
||||||
|
state = {}
|
||||||
|
calls = [float(t) for t in state.get("calls", []) if float(t) > now - 3600]
|
||||||
|
retry_at = float(state.get("retry_at", 0))
|
||||||
|
if now < retry_at or len(calls) >= 12:
|
||||||
|
raise RateLimitDeferred()
|
||||||
|
calls.append(now)
|
||||||
|
# Count attempts even if the connection fails.
|
||||||
|
with open(filename, "w") as cache:
|
||||||
|
json.dump({"calls": calls, "retry_at": retry_at}, cache)
|
||||||
|
response = requests.get(url, params=params, timeout=30)
|
||||||
|
if response.status_code == 429:
|
||||||
|
delay = 3600.0
|
||||||
|
try:
|
||||||
|
delay = float(response.headers["Retry-After"])
|
||||||
|
except (KeyError, ValueError):
|
||||||
|
match = re.search(r"available in (\d+) seconds", response.text)
|
||||||
|
if match:
|
||||||
|
delay = float(match[1])
|
||||||
|
with open(filename, "w") as cache:
|
||||||
|
json.dump({"calls": calls, "retry_at": now + delay + 1}, cache)
|
||||||
|
raise RateLimitDeferred()
|
||||||
|
return response
|
||||||
|
|
||||||
|
|
||||||
def next_launch_api_data(rocket_dir: str, limit: int = LIMIT) -> StrDict | None:
|
def next_launch_api_data(rocket_dir: str, limit: int = LIMIT) -> StrDict | None:
|
||||||
"""Get the next upcoming launches from the API."""
|
"""Get the next upcoming launches from the API."""
|
||||||
now = datetime.now()
|
now = datetime.now()
|
||||||
|
|
@ -26,7 +70,7 @@ def next_launch_api_data(rocket_dir: str, limit: int = LIMIT) -> StrDict | None:
|
||||||
url = "https://ll.thespacedevs.com/2.2.0/launch/upcoming/"
|
url = "https://ll.thespacedevs.com/2.2.0/launch/upcoming/"
|
||||||
|
|
||||||
params: dict[str, str | int] = {"limit": limit}
|
params: dict[str, str | int] = {"limit": limit}
|
||||||
r = requests.get(url, params=params)
|
r = api_get(rocket_dir, url, params)
|
||||||
try:
|
try:
|
||||||
data: StrDict = r.json()
|
data: StrDict = r.json()
|
||||||
except requests.exceptions.JSONDecodeError:
|
except requests.exceptions.JSONDecodeError:
|
||||||
|
|
@ -99,7 +143,9 @@ def is_active_crewed_spaceflight(flight: Launch, now: datetime) -> bool:
|
||||||
return True
|
return True
|
||||||
|
|
||||||
|
|
||||||
def active_crewed_flights_api(limit: int = LIMIT) -> list[Summary] | None:
|
def active_crewed_flights_api(
|
||||||
|
rocket_dir: str, limit: int = LIMIT
|
||||||
|
) -> list[Summary] | None:
|
||||||
"""
|
"""
|
||||||
Get active crewed spaceflights from the SpaceDevs API.
|
Get active crewed spaceflights from the SpaceDevs API.
|
||||||
|
|
||||||
|
|
@ -115,7 +161,7 @@ def active_crewed_flights_api(limit: int = LIMIT) -> list[Summary] | None:
|
||||||
max_pages = 20
|
max_pages = 20
|
||||||
|
|
||||||
while url and page < max_pages:
|
while url and page < max_pages:
|
||||||
r = requests.get(url, params=params if page == 0 else None, timeout=30)
|
r = api_get(rocket_dir, url, params if page == 0 else None)
|
||||||
if not r.ok:
|
if not r.ok:
|
||||||
return None
|
return None
|
||||||
try:
|
try:
|
||||||
|
|
@ -125,7 +171,7 @@ def active_crewed_flights_api(limit: int = LIMIT) -> list[Summary] | None:
|
||||||
|
|
||||||
results = data.get("results")
|
results = data.get("results")
|
||||||
if not isinstance(results, list):
|
if not isinstance(results, list):
|
||||||
break
|
return None
|
||||||
|
|
||||||
for flight in results:
|
for flight in results:
|
||||||
if not isinstance(flight, dict):
|
if not isinstance(flight, dict):
|
||||||
|
|
@ -151,7 +197,8 @@ def active_crewed_flights_api(limit: int = LIMIT) -> list[Summary] | None:
|
||||||
params = {}
|
params = {}
|
||||||
page += 1
|
page += 1
|
||||||
|
|
||||||
return launches
|
# Never replace the cache with an incomplete paginated result.
|
||||||
|
return None if url else launches
|
||||||
|
|
||||||
|
|
||||||
def get_active_crewed_flights_cache_filename(rocket_dir: str) -> str:
|
def get_active_crewed_flights_cache_filename(rocket_dir: str) -> str:
|
||||||
|
|
@ -199,11 +246,11 @@ def get_active_crewed_flights(
|
||||||
"""Get active crewed flights with cache and API fallback."""
|
"""Get active crewed flights with cache and API fallback."""
|
||||||
now = datetime.now()
|
now = datetime.now()
|
||||||
cached = load_active_crewed_flights_cache(rocket_dir)
|
cached = load_active_crewed_flights_cache(rocket_dir)
|
||||||
if cached and not refresh and (now - cached[0]).seconds <= ttl:
|
if cached and not refresh and (now - cached[0]).total_seconds() <= ttl:
|
||||||
return cached[1]
|
return cached[1]
|
||||||
|
|
||||||
try:
|
try:
|
||||||
active_flights = active_crewed_flights_api()
|
active_flights = active_crewed_flights_api(rocket_dir)
|
||||||
if active_flights is not None:
|
if active_flights is not None:
|
||||||
write_active_crewed_flights_cache(rocket_dir, active_flights)
|
write_active_crewed_flights_cache(rocket_dir, active_flights)
|
||||||
return active_flights
|
return active_flights
|
||||||
|
|
@ -359,7 +406,7 @@ def get_launches(
|
||||||
|
|
||||||
existing.sort(reverse=True)
|
existing.sort(reverse=True)
|
||||||
|
|
||||||
if refresh or not existing or (now - existing[0][0]).seconds > ttl:
|
if refresh or not existing or (now - existing[0][0]).total_seconds() > ttl:
|
||||||
try:
|
try:
|
||||||
upcoming = next_launch_api(rocket_dir, limit=limit)
|
upcoming = next_launch_api(rocket_dir, limit=limit)
|
||||||
if upcoming is None:
|
if upcoming is None:
|
||||||
|
|
|
||||||
62
tests/test_spacedevs_budget.py
Normal file
62
tests/test_spacedevs_budget.py
Normal file
|
|
@ -0,0 +1,62 @@
|
||||||
|
"""Regression tests for the shared SpaceDevs request budget."""
|
||||||
|
|
||||||
|
import json
|
||||||
|
from pathlib import Path
|
||||||
|
from unittest.mock import Mock
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from agenda import thespacedevs
|
||||||
|
|
||||||
|
|
||||||
|
def test_budget_survives_separate_calls(
|
||||||
|
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
||||||
|
) -> None:
|
||||||
|
request = Mock(return_value=Mock(status_code=200))
|
||||||
|
monkeypatch.setattr("agenda.thespacedevs.requests.get", request)
|
||||||
|
monkeypatch.setattr("agenda.thespacedevs.time.time", lambda: 10000.0)
|
||||||
|
for _ in range(12):
|
||||||
|
thespacedevs.api_get(str(tmp_path), "https://example.com")
|
||||||
|
with pytest.raises(thespacedevs.RateLimitDeferred):
|
||||||
|
thespacedevs.api_get(str(tmp_path), "https://example.com")
|
||||||
|
assert request.call_count == 12
|
||||||
|
monkeypatch.setattr("agenda.thespacedevs.time.time", lambda: 13601.0)
|
||||||
|
thespacedevs.api_get(str(tmp_path), "https://example.com")
|
||||||
|
assert request.call_count == 13
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("headers", [{}, {"Retry-After": "1832"}])
|
||||||
|
def test_throttle_persists_cooldown(
|
||||||
|
tmp_path: Path, monkeypatch: pytest.MonkeyPatch, headers: dict[str, str]
|
||||||
|
) -> None:
|
||||||
|
request = Mock(
|
||||||
|
return_value=Mock(
|
||||||
|
status_code=429,
|
||||||
|
headers=headers,
|
||||||
|
text='{"detail": "Request was throttled. Expected available in 1832 seconds."}',
|
||||||
|
)
|
||||||
|
)
|
||||||
|
monkeypatch.setattr("agenda.thespacedevs.requests.get", request)
|
||||||
|
monkeypatch.setattr("agenda.thespacedevs.time.time", lambda: 10000.0)
|
||||||
|
with pytest.raises(thespacedevs.RateLimitDeferred):
|
||||||
|
thespacedevs.api_get(str(tmp_path), "https://example.com")
|
||||||
|
state = json.loads((tmp_path / "api_budget.json").read_text())
|
||||||
|
assert state["retry_at"] == 11833.0
|
||||||
|
monkeypatch.setattr("agenda.thespacedevs.time.time", lambda: 11000.0)
|
||||||
|
with pytest.raises(thespacedevs.RateLimitDeferred):
|
||||||
|
thespacedevs.api_get(str(tmp_path), "https://example.com")
|
||||||
|
assert request.call_count == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_incomplete_pagination_keeps_cached_flights(
|
||||||
|
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
|
||||||
|
) -> None:
|
||||||
|
cached = [{"slug": "still-in-orbit"}]
|
||||||
|
thespacedevs.write_active_crewed_flights_cache(str(tmp_path), cached)
|
||||||
|
response = Mock(status_code=200, ok=True)
|
||||||
|
response.json.return_value = {"results": [], "next": "https://example.com/page2"}
|
||||||
|
request = Mock(side_effect=[response, thespacedevs.RateLimitDeferred()])
|
||||||
|
monkeypatch.setattr(thespacedevs, "api_get", request)
|
||||||
|
assert thespacedevs.get_active_crewed_flights(str(tmp_path), refresh=True) == cached
|
||||||
|
loaded = thespacedevs.load_active_crewed_flights_cache(str(tmp_path))
|
||||||
|
assert loaded is not None and loaded[1] == cached
|
||||||
12
update.py
12
update.py
|
|
@ -342,10 +342,6 @@ def update_thespacedevs(config: flask.config.Config) -> None:
|
||||||
if agenda.thespacedevs.is_launches_cache_fresh(rocket_dir):
|
if agenda.thespacedevs.is_launches_cache_fresh(rocket_dir):
|
||||||
return
|
return
|
||||||
|
|
||||||
# Update active crewed mission cache used by the launches page.
|
|
||||||
# Uses the 2-hour TTL; failures are handled internally with cache fallback.
|
|
||||||
active_crewed = agenda.thespacedevs.get_active_crewed_flights(rocket_dir)
|
|
||||||
|
|
||||||
# Always follow configured slugs
|
# Always follow configured slugs
|
||||||
follow_slugs: set[str] = set(config["FOLLOW_LAUNCHES"])
|
follow_slugs: set[str] = set(config["FOLLOW_LAUNCHES"])
|
||||||
|
|
||||||
|
|
@ -357,7 +353,10 @@ def update_thespacedevs(config: flask.config.Config) -> None:
|
||||||
}
|
}
|
||||||
|
|
||||||
t0 = time()
|
t0 = time()
|
||||||
data = agenda.thespacedevs.next_launch_api_data(rocket_dir)
|
try:
|
||||||
|
data = agenda.thespacedevs.next_launch_api_data(rocket_dir)
|
||||||
|
except agenda.thespacedevs.RateLimitDeferred:
|
||||||
|
return
|
||||||
if not data:
|
if not data:
|
||||||
send_thespacedevs_payload_alert(
|
send_thespacedevs_payload_alert(
|
||||||
config,
|
config,
|
||||||
|
|
@ -375,6 +374,9 @@ def update_thespacedevs(config: flask.config.Config) -> None:
|
||||||
)
|
)
|
||||||
return
|
return
|
||||||
|
|
||||||
|
# Fetch upcoming launches first so historical pagination cannot starve them.
|
||||||
|
active_crewed = agenda.thespacedevs.get_active_crewed_flights(rocket_dir)
|
||||||
|
|
||||||
# Identify test-flight slugs present in the current data
|
# Identify test-flight slugs present in the current data
|
||||||
cur_test_slugs: set[str] = {
|
cur_test_slugs: set[str] = {
|
||||||
typing.cast(str, item["slug"])
|
typing.cast(str, item["slug"])
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue