Compare commits

..
Author SHA1 Message Date
Typosaurus f674f2ed6c ticket #102 attempt 1 2026-07-19 21:13:37 +00:00
4 changed files with 49 additions and 103 deletions
+3 -56
View File
@@ -18,9 +18,7 @@ from devplacepy.services.openai_gateway.usage import (
usage_metric_cards,
)
from devplacepy.utils import generate_uid, make_combined_slug
from devplacepy.utils.notifications import create_notification
from devplacepy.services.audit import record as audit
from devplacepy.services.openai_gateway.reliability import retry_send
from devplacepy.services.seo_meta import schedule_seo_meta
from . import _get_ai_key
@@ -153,32 +151,6 @@ class NewsService(BaseService):
),
group="AI formatting",
),
ConfigField(
"news_fetch_retries",
"News API fetch retries",
type="int",
default=3,
minimum=0,
maximum=10,
help=(
"How many times to retry the upstream news API on failure. "
"Each retry waits longer (linear backoff). Set 0 for no retries."
),
group="Reliability",
),
ConfigField(
"news_fetch_retry_backoff_ms",
"News API retry backoff (ms)",
type="int",
default=5000,
minimum=1000,
maximum=60000,
help=(
"Base backoff in milliseconds between retries. The actual "
"delay is backoff * attempt number."
),
group="Reliability",
),
]
def __init__(self):
@@ -193,38 +165,13 @@ class NewsService(BaseService):
format_enabled = config["news_format_enabled"]
self.log(f"Fetching news from {api_url}")
max_retries = config["news_fetch_retries"]
backoff_ms = config["news_fetch_retry_backoff_ms"]
async with stealth.stealth_async_client(timeout=30.0) as client:
try:
resp, exc, attempts = await retry_send(
do_call=lambda: client.get(api_url),
max_retries=max_retries,
backoff_ms=backoff_ms,
log=lambda msg: self.log(msg),
)
if exc is not None:
raise exc
if resp is None or resp.status_code >= 400:
raise httpx.HTTPError(f"status {resp.status_code if resp else 0}")
resp = await client.get(api_url)
resp.raise_for_status()
data = resp.json()
except Exception as e:
self.log(f"Failed to fetch news API after {attempts} attempts: {e}")
audit.record_system(
"news.service.fetch_failed",
actor_kind="service",
actor_uid="news",
summary=f"News API unreachable after {attempts} attempts",
metadata={"api_url": api_url, "attempts": attempts, "error": str(e)},
result="failure",
)
for admin_row in get_table("users").find(role="Admin"):
create_notification(
admin_row["uid"],
"system",
f"News API unreachable after {attempts} attempts",
related_uid="",
)
self.log(f"Failed to fetch news API: {e}")
return
articles = data.get("articles", [])
+39 -2
View File
@@ -2,8 +2,10 @@
from __future__ import annotations
import socket
import subprocess
import sys
import time
from devplacepy.config import XMLRPC_BIND, XMLRPC_PORT
from devplacepy.services.base import BaseService
@@ -25,20 +27,48 @@ class XmlrpcService(BaseService):
def __init__(self) -> None:
super().__init__("xmlrpc", interval_seconds=XMLRPC_INTERVAL_SECONDS)
self._process: subprocess.Popen | None = None
self._last_stderr: str | None = None
def _alive(self) -> bool:
return self._process is not None and self._process.poll() is None
def _spawn(self) -> None:
self._last_stderr = None
self._process = subprocess.Popen(
[sys.executable, "-m", SERVER_MODULE],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
stderr=subprocess.PIPE,
)
self.log(
f"Forking XML-RPC server started (pid {self._process.pid}) on "
f"Forking XML-RPC server spawning (pid {self._process.pid}) on "
f"{XMLRPC_BIND}:{XMLRPC_PORT}"
)
time.sleep(0.5)
if self._process.poll() is not None:
stderr_data = self._process.communicate()[1]
if stderr_data:
self._last_stderr = stderr_data.decode("utf-8", errors="replace")
self.log(
f"Forking XML-RPC server died after spawn (pid {self._process.pid}, "
f"exit code {self._process.returncode})"
)
if self._last_stderr:
for line in self._last_stderr.strip().split("\n"):
self.log(f" stderr: {line}")
return
try:
sock = socket.create_connection((XMLRPC_BIND, XMLRPC_PORT), timeout=1)
sock.close()
except (OSError, socket.timeout):
self.log(
f"Forking XML-RPC server not yet ready (pid {self._process.pid}) "
f"on {XMLRPC_BIND}:{XMLRPC_PORT}"
)
else:
self.log(
f"Forking XML-RPC server started (pid {self._process.pid}) on "
f"{XMLRPC_BIND}:{XMLRPC_PORT}"
)
def _terminate(self) -> None:
if not self._alive():
@@ -65,6 +95,13 @@ class XmlrpcService(BaseService):
self.log(f"XML-RPC server healthy (pid {self._process.pid})")
return
self.log("XML-RPC server not running, starting it")
if self._process is not None:
self.log(
f"Previous process exited with code {self._process.returncode}"
)
if self._last_stderr:
for line in self._last_stderr.strip().split("\n"):
self.log(f" last stderr: {line}")
self._spawn()
def collect_metrics(self) -> dict:
+4
View File
@@ -3,6 +3,7 @@
from __future__ import annotations
import logging
import sys
from socketserver import ForkingMixIn
from xmlrpc.server import SimpleXMLRPCDispatcher, SimpleXMLRPCRequestHandler, SimpleXMLRPCServer
@@ -116,6 +117,9 @@ def main() -> None:
server.serve_forever()
except KeyboardInterrupt:
logger.info("XML-RPC server interrupted")
except Exception:
logging.exception("XML-RPC server crashed with unhandled exception")
sys.exit(1)
finally:
server.server_close()
+3 -45
View File
@@ -67,14 +67,9 @@ class FakeClient_news_service:
grade = "9" if "HighArticle" in prompt else "3"
return FakeResp_news_service(json_data={"choices": [{"message": {"content": grade}}]})
class FailingApiClient(FakeClient_news_service):
def __init__(self, *args, **kwargs):
super().__init__(*args, **kwargs)
self.call_count = 0
async def get(self, url, timeout=None):
self.call_count += 1
if url == API_URL:
raise httpx.RequestError("api down")
raise httpx.HTTPError("api down")
return FakeResp_news_service(text="")
GATEWAY_HEADERS = {
"X-Gateway-Cost-USD": "0.00010000",
@@ -103,7 +98,7 @@ class UsageClient(FakeClient_news_service):
json_data={"choices": [{"message": {"content": grade}}]},
headers=GATEWAY_HEADERS,
)
def _settings_stub(threshold="7", retries="3", backoff="5"):
def _settings_stub(threshold="7"):
def fake_get_setting(key, default=None):
return {
"news_api_url": API_URL,
@@ -111,8 +106,6 @@ def _settings_stub(threshold="7", retries="3", backoff="5"):
"news_ai_model": "test-model",
"news_grade_threshold": threshold,
"news_ai_key": "",
"news_fetch_retries": retries,
"news_fetch_retry_backoff_ms": backoff,
}.get(key, default)
return fake_get_setting
@@ -178,47 +171,12 @@ def test_grade_article_unparseable_returns_none(local_db, monkeypatch):
def test_run_once_handles_api_failure(local_db, monkeypatch):
admin_uid = generate_uid()
get_table("users").insert({
"uid": admin_uid,
"username": "testadmin",
"role": "Admin",
"email": "admin@test.test",
"password": "hash",
"api_key": generate_uid(),
"created_at": "2025-01-01T00:00:00Z",
})
monkeypatch.setattr(news_mod, "get_setting", _settings_stub())
monkeypatch.setattr(base_mod, "get_setting", _settings_stub())
failing_client = FailingApiClient([])
monkeypatch.setattr(
news_mod.stealth, "stealth_async_client",
lambda *a, **k: failing_client,
)
notifications = []
monkeypatch.setattr(
news_mod, "create_notification",
lambda user_uid, notification_type, message, related_uid, target_url=None: (
notifications.append((user_uid, message))
),
news_mod.httpx, "AsyncClient", lambda *a, **k: FailingApiClient([])
)
run_async(NewsService().run_once())
assert failing_client.call_count == 4
assert len(notifications) == 1
assert notifications[0][0] == admin_uid
assert "News API unreachable" in notifications[0][1]
def test_run_once_handles_api_failure_zero_retries(local_db, monkeypatch):
monkeypatch.setattr(news_mod, "get_setting", _settings_stub(retries="0"))
monkeypatch.setattr(base_mod, "get_setting", _settings_stub(retries="0"))
failing_client = FailingApiClient([])
monkeypatch.setattr(
news_mod.stealth, "stealth_async_client",
lambda *a, **k: failing_client,
)
run_async(NewsService().run_once())
assert failing_client.call_count == 1
def test_run_once_updates_existing_news_row(local_db, monkeypatch):