Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/design.md
Original file line number Diff line number Diff line change
Expand Up @@ -110,7 +110,7 @@ SQLite 没有原生日期时间类型,状态存储需要直接对 TEXT 做范

RSS 是可选补充源,其失败不影响任务成功率;天气 API 是主要信息来源。

**获取失败** — 单个 RSS 源经可配置次数重试(默认 3 次,间隔 3–5 秒)后仍无法获取或解析。行为:记录警告日志,本次运行继续使用已成功获取的 RSS 内容和天气 API 数据。
**获取失败** — 单个 RSS 源只对传输错误和临时 HTTP 状态(408、425、429、500、502、503、504)重试,其他 HTTP 错误立即失败。重试默认最多 3 次、随机间隔 3–5 秒;有效 `Retry-After` 指定更长退避时优先遵守。请求仍无法获取或解析时记录警告日志,本次运行继续使用已成功获取的 RSS 内容和天气 API 数据。

**长期无更新** — RSS 源在配置的小时数(默认 24)内未曾见到任何新文章。判定基准为 `source_health` 表中记录的最后一次文章时间,不受当天本地日期筛选影响。

Expand Down
2 changes: 1 addition & 1 deletion docs/requirements.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@

## 可靠性

1. RSS 请求失败后随机等待 3–5 秒并重试;重试次数可配置。
1. RSS 只重试传输错误,以及 408、425、429、500、502、503、504 HTTP 状态;其他 HTTP 错误立即失败。重试默认随机等待 3–5 秒,服务端返回有效 `Retry-After` 时至少等待其指定时长;重试次数可配置。
2. 重试耗尽则记录警告日志并继续任务;RSS 是可选补充源,其获取失败不终止任务,但连续失败达到阈值时独立发出运维告警。连续失败和长期无更新告警均采用至少一次投递;投递或状态写入失败时记录日志并在后续任务中重试,不得使简报任务失败,接收端应按告警类型、源集合和失败周期去重。
3. 任务发生非 RSS 错误时立即向运维通知渠道报警;投递失败时在后续任务失败中重试,成功投递后同一轮失败不重复告警,任务成功后重新开放。RSS 获取连续失败达到可配置阈值(默认 3)时独立报警,与任务失败隔离。任一已配置 RSS 源超过配置的时长(默认 24 小时)没有新文章也触发报警。未配置 RSS 不属于故障。

Expand Down
97 changes: 96 additions & 1 deletion tests/test_sources.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,11 @@
from unittest.mock import AsyncMock

import httpx
import pendulum
import pytest

from weather_briefing.models import ContextSourceConfig, FeedConfig, SourceDocument
from weather_briefing.sources import HTTPContextSource, RSSSource, SourceFetchError
from weather_briefing.sources import HTTPContextSource, RSSSource, SourceFetchError, _retry_after_seconds


async def test_rss_source_marks_configured_verbatim_article(caplog) -> None:
Expand Down Expand Up @@ -126,6 +127,100 @@ def handler(_: httpx.Request) -> httpx.Response:
assert caught.value.__cause__ is None


async def test_rss_source_does_not_retry_permanent_http_status(monkeypatch) -> None:
handler = AsyncMock(return_value=httpx.Response(404))
sleep = AsyncMock()
monkeypatch.setattr("weather_briefing.sources.asyncio.sleep", sleep)

async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client:
with pytest.raises(SourceFetchError):
await RSSSource(client, max_attempts=3).fetch(
FeedConfig("source", "Source", "https://private.example.invalid")
)

assert handler.await_count == 1
sleep.assert_not_awaited()


async def test_rss_source_does_not_retry_other_http_errors(monkeypatch) -> None:
attempts = 0
sleep = AsyncMock()

def handler(_: httpx.Request) -> httpx.Response:
nonlocal attempts
attempts += 1
raise httpx.HTTPError("invalid request")

monkeypatch.setattr("weather_briefing.sources.asyncio.sleep", sleep)
async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client:
with pytest.raises(SourceFetchError):
await RSSSource(client, max_attempts=3).fetch(
FeedConfig("source", "Source", "https://private.example.invalid")
)

assert attempts == 1
sleep.assert_not_awaited()


async def test_rss_source_retries_transport_errors(monkeypatch) -> None:
attempts = 0
sleep = AsyncMock()

def handler(request: httpx.Request) -> httpx.Response:
nonlocal attempts
attempts += 1
if attempts == 1:
raise httpx.ConnectError("connection failed", request=request)
return httpx.Response(200, text="<?xml version='1.0'?><rss version='2.0'><channel /></rss>")

monkeypatch.setattr("weather_briefing.sources.asyncio.sleep", sleep)
async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client:
articles = await RSSSource(client, max_attempts=2).fetch(
FeedConfig("source", "Source", "https://private.example.invalid")
)

assert articles == ()
assert attempts == 2
assert sleep.await_count == 1


async def test_rss_source_respects_retry_after(monkeypatch) -> None:
sleep = AsyncMock()
monkeypatch.setattr("weather_briefing.sources.asyncio.sleep", sleep)
monkeypatch.setattr("weather_briefing.sources.random.uniform", lambda _low, _high: 4.0)

async with httpx.AsyncClient(
transport=httpx.MockTransport(lambda _: httpx.Response(429, headers={"Retry-After": "12"}))
) as client:
with pytest.raises(SourceFetchError):
await RSSSource(client, max_attempts=2).fetch(
FeedConfig("source", "Source", "https://private.example.invalid")
)

sleep.assert_awaited_once_with(12.0)


@pytest.mark.parametrize(
("value", "expected"),
[
(None, None),
("invalid", None),
("Thu, 16 Jul 2026 00:00:10", None),
("Thu, 16 Jul 2026 00:00:10 GMT", 10.0),
("Thu, 16 Jul 2025 00:00:00 GMT", 0.0),
("-1", 0.0),
],
)
def test_retry_after_parsing(monkeypatch, value: str | None, expected: float | None) -> None:
monkeypatch.setattr(
"weather_briefing.sources.pendulum.now",
lambda _timezone: pendulum.datetime(2026, 7, 16, tz="UTC"),
)
headers = {} if value is None else {"Retry-After": value}

assert _retry_after_seconds(httpx.Response(503, headers=headers)) == expected


async def test_rss_local_offset_is_normalized_then_restored() -> None:
xml = """<?xml version="1.0"?><rss version="2.0"><channel><title>x</title>
<item><guid>one</guid><title>Local publication</title>
Expand Down
32 changes: 30 additions & 2 deletions weather_briefing/sources.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
import hashlib
import logging
import random
from email.utils import parsedate_to_datetime
from time import struct_time
from typing import Protocol

Expand All @@ -18,6 +19,7 @@
from .models import Article, ContextSourceConfig, FeedConfig, SourceDocument

_LOGGER = logging.getLogger("weather_briefing.sources")
_RETRYABLE_STATUS_CODES = frozenset({408, 425, 429, 500, 502, 503, 504})


class SourceFetchError(RuntimeError):
Expand Down Expand Up @@ -54,6 +56,23 @@ def _entry_content(entry: feedparser.FeedParserDict) -> str:
return str(entry.get("summary", "")).strip()


def _retry_after_seconds(response: httpx.Response) -> float | None:
value = response.headers.get("Retry-After")
if value is None:
return None
try:
seconds = int(value)
except ValueError:
try:
retry_at = parsedate_to_datetime(value)
except (TypeError, ValueError, OverflowError):
return None
if retry_at.tzinfo is None:
return None
seconds = (pendulum.instance(retry_at) - pendulum.now("UTC")).total_seconds()
return max(0.0, float(seconds))


class RSSSource:
"""Fetch, retry, clean, and normalize RSS feed entries."""

Expand Down Expand Up @@ -121,16 +140,25 @@ async def fetch(self, config: FeedConfig) -> tuple[Article, ...]:

async def _fetch_with_retry(self, config: FeedConfig) -> str:
for attempt in range(1, self._max_attempts + 1):
retry_after: float | None = None
try:
response = await self._client.get(
config.url,
extensions=api_call_extensions("rss", "fetch"),
)
response.raise_for_status()
return response.text
except httpx.TransportError:
pass
except httpx.HTTPStatusError as exc:
if exc.response.status_code not in _RETRYABLE_STATUS_CODES:
break
retry_after = _retry_after_seconds(exc.response)
except httpx.HTTPError:
if attempt < self._max_attempts:
await asyncio.sleep(random.uniform(self._retry_min_seconds, self._retry_max_seconds))
break
if attempt < self._max_attempts:
delay = random.uniform(self._retry_min_seconds, self._retry_max_seconds)
await asyncio.sleep(max(delay, retry_after or 0.0))
raise SourceFetchError(f"RSS source {config.id} failed after {self._max_attempts} attempts") from None


Expand Down