diff --git a/README.md b/README.md index e1a77bce..2ad450c6 100644 --- a/README.md +++ b/README.md @@ -39,7 +39,15 @@ uv run --frozen weather-briefing run hourly `LLM_PROVIDER=deepseek` 使用 `DEEPSEEK_API_KEY`、`DEEPSEEK_MODEL` 和可选的 `DEEPSEEK_BASE_URL`;DeepSeek provider 已预置官方 Base URL。`LLM_PROVIDER=openai-compatible` 使用 `LLM_API_KEY`、`LLM_MODEL` 和 `LLM_BASE_URL`。两套配置互不回退。 -应用将带时间、级别和 logger 名称的运行日志写入标准错误;设置 `DEBUG=true` 可输出 RSS 获取和 LLM 重试等诊断信息。 +应用将带时间、级别和 logger 名称的运行日志写入标准错误;设置 `DEBUG=true` 可输出 RSS 获取、LLM 重试,以及从 RSS 清洗、权威预报转发和平台渲染到 Telegram 分片接受状态的非敏感诊断信息。该链路只记录来源、发布时间、字符数和分片状态,不记录标题、正文、URL、token、chat ID 或请求 endpoint。若仍需排查平台渲染或分片内容,可在不重启 daemon 的情况下临时记录完整渲染正文: + +```bash +weather-briefing diagnostics rendered-text enable --for 15m +weather-briefing diagnostics rendered-text status +weather-briefing diagnostics rendered-text disable +``` + +容器部署通过同一运行实例执行,例如 `docker exec weather-briefing weather-briefing diagnostics rendered-text enable --for 15m`。该开关最长启用 24 小时并自动过期,状态保存在 `BRIEFING_STATE_PATH`。只有同时启用 `DEBUG` 和临时开关时才记录正文;日志包含简报、告警、权威预报以及 Telegram 分片的完整文本,可能暴露来源内容和位置上下文,排障后应立即关闭并妥善保护日志。token、chat ID 和请求 endpoint 不会写入这些诊断日志。 定位层从地名解析国家或行政区代码。Open-Meteo 负责城市/邮编查询,空结果时由 OpenStreetMap Nominatim 解析详细地名;结果会持久缓存。只有坐标时使用中国大陆服务范围四至宽松包围盒作快速可能性判断。省略 `WEATHER_PROVIDERS` 时,中国大陆地点使用 QWeather、Open-Meteo,其他地点只使用 Open-Meteo;显式配置时首项是主要来源,后续项依次作为备用。 diff --git a/docs/design.md b/docs/design.md index a61d949e..458fa9d5 100644 --- a/docs/design.md +++ b/docs/design.md @@ -103,7 +103,9 @@ RSS 是可选补充源,其失败不影响任务成功率;天气 API 是主 小时 LLM 结果包含布尔字段 `should_publish`。模型比较当前及历史 API 快照,仅在降雨、显著天气变化、预警或灾害动态值得打扰时设为真;活动预警不允许与 false 同时出现。false 结果不投递消息,但当前快照、文章去重和预警状态仍持久化。 -CLI 在读取运行配置前以 INFO 幂等配置单个标准错误 handler,配置成功后再按 `DEBUG` 更新级别,避免配置错误绕过统一格式、daemon 每轮任务重复追加 handler 或向 root logger 重复传播。默认记录生命周期、文章数量、陈旧来源和失败信息;`DEBUG` 启用 RSS 获取及 LLM 重试诊断。业务层只给异常追加失败计数等上下文,完整堆栈由 CLI 入口或 APScheduler 单点记录。 +CLI 在读取运行配置前以 INFO 幂等配置单个标准错误 handler,配置成功后再按 `DEBUG` 更新级别,避免配置错误绕过统一格式、daemon 每轮任务重复追加 handler 或向 root logger 重复传播。默认记录生命周期、文章数量、陈旧来源和失败信息;`DEBUG` 启用 RSS 获取及 LLM 重试诊断,并以非敏感元数据串联 RSS 清洗后的来源、发布时间、正文长度和权威预报标记,权威预报投递前后的来源、发布时间和正文长度,renderer 的可见与 payload 长度,以及 Telegram 分片总数、逐片长度和平台接受状态。这条链路不记录标题、正文、URL、token、chat ID 或请求 endpoint。业务层只给异常追加失败计数等上下文,完整堆栈由 CLI 入口或 APScheduler 单点记录。 + +完整渲染文本诊断采用第二道运行时开关,避免仅凭长期 `DEBUG` 配置泄露正文。`diagnostics rendered-text enable --for `、`status` 和 `disable` 直接读写 `BRIEFING_STATE_PATH` 中带 UTC 过期时间的 SQLite 状态;每次投递在记录前重新读取,因此另一个 CLI 进程的修改无需重启 daemon 即可生效。开关最长 24 小时,首次观察到过期状态时删除记录并发出警告。`DeliveryProvider` 记录 renderer 的完整输出,Telegram publisher 另记实际发送的每个分片,以区分渲染与传输边界;两处都只在 DEBUG 与开关同时有效时输出,并且不记录 token、chat ID 或请求 endpoint。诊断状态后端是可选能力:初始化或状态检查失败只记录警告并视为关闭,不改变消息投递。 ## 依赖边界 diff --git a/docs/requirements.md b/docs/requirements.md index 572cdc56..c5fc545c 100644 --- a/docs/requirements.md +++ b/docs/requirements.md @@ -48,6 +48,7 @@ 7. 配置项应经过类型、范围和必填校验;provider、source 和状态存储保持明确扩展边界,避免厂商逻辑进入核心编排。 8. DeepSeek 使用 `DEEPSEEK_*` 配置,通用 OpenAI-compatible provider 使用 `LLM_*` 配置;RSS 只从命名 JSON 文件读取。SQLite 没有原生日期时间类型,应用仅在持久化边界将时区感知时间转换为固定宽度 UTC 文本,使文本字典序等同绝对时间顺序。 9. 地理范围、空气质量分级与健康提示、正文清洗默认规则、provider 默认顺序及厂商指数代码等纯领域数据应存放在独立数据文件中,由实现代码加载并校验。 +10. DEBUG 日志应以非敏感元数据覆盖 RSS 清洗、权威预报转发、平台渲染和分片投递边界,使运维可以根据来源、发布时间、字符数、分片数和平台接受状态定位内容丢失阶段,而无需记录标题、正文、URL、投递凭据、接收方标识或私有 endpoint。完整渲染正文属于敏感诊断数据,默认不得记录;运行时可以通过 CLI 临时启用、查询或关闭该行为,无需重启 daemon。启用时长必须为正且不超过 24 小时,到期自动失效。只有 DEBUG 日志级别与临时开关同时有效时才输出正文,并覆盖投递 provider 生成的完整消息及平台分片。诊断状态初始化或读取失败不得阻断正常投递。 ## 投递假设 diff --git a/tests/test_cli.py b/tests/test_cli.py index 85d24dab..cc544ab5 100644 --- a/tests/test_cli.py +++ b/tests/test_cli.py @@ -1,5 +1,6 @@ import base64 import logging +import sqlite3 from collections.abc import AsyncIterator from dataclasses import replace from pathlib import Path @@ -22,6 +23,7 @@ _in_schedule, _llm_provider, _location_state_path, + _manage_rendered_text_diagnostics, _parse_run_time, _precision_reduction_notice, _qweather_is_configured, @@ -197,6 +199,69 @@ def test_version_flag() -> None: parser.parse_args(["--version"]) +def test_rendered_text_diagnostics_parser_accepts_bounded_duration() -> None: + args = build_parser().parse_args(["diagnostics", "rendered-text", "enable", "--for", "15m"]) + + assert args.diagnostics_action == "enable" + assert args.duration_seconds == 900 + + +@pytest.mark.parametrize("duration", ("0m", "15", "25h")) +def test_rendered_text_diagnostics_parser_rejects_invalid_duration(duration: str) -> None: + with pytest.raises(SystemExit): + build_parser().parse_args(["diagnostics", "rendered-text", "enable", "--for", duration]) + + +def test_main_manages_rendered_text_diagnostics_without_loading_service_settings( + monkeypatch, + tmp_path: Path, + capsys, +) -> None: + state_path = tmp_path / "state.sqlite3" + monkeypatch.setenv("BRIEFING_STATE_PATH", str(state_path)) + monkeypatch.setattr("weather_briefing.cli.load_dotenv", lambda *, override: True) + monkeypatch.setattr("weather_briefing.cli._configure_logging", lambda *, debug: None) + + monkeypatch.setattr( + "sys.argv", + ["weather-briefing", "diagnostics", "rendered-text", "enable", "--for", "15m"], + ) + main() + assert "enabled until" in capsys.readouterr().out + + monkeypatch.setattr("sys.argv", ["weather-briefing", "diagnostics", "rendered-text", "status"]) + main() + assert "is enabled until" in capsys.readouterr().out + + monkeypatch.setattr("sys.argv", ["weather-briefing", "diagnostics", "rendered-text", "disable"]) + main() + assert capsys.readouterr().out.strip() == "Rendered text diagnostic logging disabled" + + monkeypatch.setattr("sys.argv", ["weather-briefing", "diagnostics", "rendered-text", "status"]) + main() + assert capsys.readouterr().out.strip() == "Rendered text diagnostic logging is disabled" + + +@pytest.mark.parametrize( + ("action", "duration", "message"), + ( + ("enable", None, "require a duration"), + ("unsupported", None, "Unsupported rendered text diagnostics action"), + ), +) +def test_rendered_text_diagnostics_reject_invalid_internal_requests( + action: str, + duration: int | None, + message: str, + monkeypatch, + tmp_path: Path, +) -> None: + monkeypatch.setenv("BRIEFING_STATE_PATH", str(tmp_path / "state.sqlite3")) + + with pytest.raises(ValueError, match=message): + _manage_rendered_text_diagnostics(action, duration) + + def test_main_loads_dotenv_with_supported_arguments(monkeypatch) -> None: calls: list[bool] = [] @@ -417,7 +482,7 @@ async def test_run_skips_and_logs_when_enforce_window_outside_schedule(monkeypat logging.root.setLevel(original_root_level) -async def test_run_logs_start_resolve_and_publish(monkeypatch, capsys) -> None: +async def test_run_continues_when_runtime_diagnostics_are_unavailable(monkeypatch, capsys) -> None: from types import SimpleNamespace from unittest.mock import patch @@ -434,7 +499,11 @@ async def test_run_logs_start_resolve_and_publish(monkeypatch, capsys) -> None: monkeypatch.setattr("weather_briefing.cli._parse_run_time", lambda v, t: now) monkeypatch.setattr("weather_briefing.cli._in_schedule", lambda k, n, s: True) - monkeypatch.setattr("weather_briefing.cli._delivery_provider", lambda s, c: None) + + def delivery_without_diagnostics(s: object, c: object, diagnostics: object) -> None: + assert diagnostics is None + + monkeypatch.setattr("weather_briefing.cli._delivery_provider", delivery_without_diagnostics) monkeypatch.setattr("weather_briefing.cli._llm_provider", lambda s, c: None) monkeypatch.setattr("weather_briefing.cli._weather_context_provider", lambda s, c, loc: None) @@ -453,6 +522,11 @@ def __exit__(self, *args: object) -> None: monkeypatch.setattr("weather_briefing.cli.SQLiteStateStore", lambda p: FakeState()) + def unavailable_diagnostics(path: Path) -> None: + raise sqlite3.OperationalError("database is locked") + + monkeypatch.setattr("weather_briefing.cli.SQLiteRuntimeDiagnostics", unavailable_diagnostics) + async def fake_service_run(kind: str, n: object) -> str: return "published body" @@ -472,6 +546,7 @@ async def fake_service_run(kind: str, n: object) -> str: stderr = capsys.readouterr().err assert "Starting hourly briefing run" in stderr + assert "Runtime diagnostics unavailable; continuing without sensitive rendered text logging" in stderr assert "Resolving 1 location(s)" in stderr assert "Processing location test (Test City)" in stderr assert "briefing published (14 characters)" in stderr @@ -517,7 +592,7 @@ async def publish_alert(self, title: str, body: str) -> None: monkeypatch.setattr("weather_briefing.cli._parse_run_time", lambda v, t: now) monkeypatch.setattr("weather_briefing.cli._in_schedule", lambda k, n, s: True) - monkeypatch.setattr("weather_briefing.cli._delivery_provider", lambda s, c: AlertDelivery()) + monkeypatch.setattr("weather_briefing.cli._delivery_provider", lambda s, c, d: AlertDelivery()) monkeypatch.setattr("weather_briefing.cli._llm_provider", lambda s, c: None) monkeypatch.setattr("weather_briefing.cli._weather_context_provider", lambda s, c, loc: None) @@ -535,6 +610,7 @@ def __exit__(self, *args: object) -> None: pass monkeypatch.setattr("weather_briefing.cli.SQLiteStateStore", lambda p: FakeState()) + monkeypatch.setattr("weather_briefing.cli.SQLiteRuntimeDiagnostics", lambda p: FakeState()) async def fake_service_run(kind: str, n: object) -> str: return "published body" @@ -581,7 +657,7 @@ async def test_run_logs_skipped_when_no_content(monkeypatch, capsys) -> None: monkeypatch.setattr("weather_briefing.cli._parse_run_time", lambda v, t: now) monkeypatch.setattr("weather_briefing.cli._in_schedule", lambda k, n, s: True) - monkeypatch.setattr("weather_briefing.cli._delivery_provider", lambda s, c: None) + monkeypatch.setattr("weather_briefing.cli._delivery_provider", lambda s, c, d: None) monkeypatch.setattr("weather_briefing.cli._llm_provider", lambda s, c: None) monkeypatch.setattr("weather_briefing.cli._weather_context_provider", lambda s, c, loc: None) @@ -599,6 +675,7 @@ def __exit__(self, *args: object) -> None: pass monkeypatch.setattr("weather_briefing.cli.SQLiteStateStore", lambda p: FakeState()) + monkeypatch.setattr("weather_briefing.cli.SQLiteRuntimeDiagnostics", lambda p: FakeState()) async def fake_service_run(kind: str, n: object) -> str | None: return None diff --git a/tests/test_publishers.py b/tests/test_publishers.py index c1456968..996a470f 100644 --- a/tests/test_publishers.py +++ b/tests/test_publishers.py @@ -20,6 +20,33 @@ async def publish(self, message: RenderedMessage, *, single_message: bool = Fals pass +class EnabledDiagnostics: + def rendered_text_logging_enabled(self) -> bool: + return True + + +class FailingDiagnostics: + def rendered_text_logging_enabled(self) -> bool: + raise RuntimeError("diagnostic state unavailable") + + +class CountingDiagnostics: + def __init__(self) -> None: + self.checks = 0 + + def rendered_text_logging_enabled(self) -> bool: + self.checks += 1 + return True + + +class RecordingPublisher: + def __init__(self) -> None: + self.messages: list[RenderedMessage] = [] + + async def publish(self, message: RenderedMessage, *, single_message: bool = False) -> None: + self.messages.append(message) + + def test_delivery_provider_applies_platform_limit_without_leaking_it_into_config() -> None: unrestricted = DeliveryProvider(PlainTextRenderer(), NoopPublisher()) telegram_like = DeliveryProvider(PlainTextRenderer(), NoopPublisher(), 4096) @@ -33,21 +60,82 @@ def test_split_message_prefers_line_boundary() -> None: assert _split_message("first line\nsecond line", 12) == ("first line", "second line") -async def test_telegram_publisher_uses_runtime_values() -> None: +async def test_telegram_publisher_uses_runtime_values(caplog) -> None: requests: list[httpx.Request] = [] def handler(request: httpx.Request) -> httpx.Response: requests.append(request) return httpx.Response(200, json={"ok": True}) - async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client: - publisher = TelegramPublisher(client, "runtime-token", "runtime-chat") - await publisher.publish(RenderedMessage("Title\n\nBody", 11)) + with caplog.at_level("DEBUG", logger="weather_briefing.publishers"): + async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client: + publisher = TelegramPublisher(client, "runtime-token", "runtime-chat", EnabledDiagnostics()) + await publisher.publish(RenderedMessage("Title\n\nBody", 11)) assert requests[0].url.path == "/botruntime-token/sendMessage" payload = json.loads(requests[0].content) assert payload["chat_id"] == "runtime-chat" assert payload["parse_mode"] == "HTML" + assert "Telegram delivery prepared: visible_characters=11 payload_characters=18 chunks=1" in caplog.text + assert "Telegram chunk accepted: index=1/1 payload_characters=18" in caplog.text + assert ( + "Sensitive rendered text diagnostic: stage=telegram-chunk-1-of-1 body='Title\\n\\nBody'" + ) in caplog.text + assert "runtime-token" not in caplog.text + assert "runtime-chat" not in caplog.text + + +async def test_telegram_checks_runtime_diagnostics_once_for_multiple_chunks(caplog) -> None: + diagnostics = CountingDiagnostics() + body = "x" * (TelegramPublisher.MAX_MESSAGE_LENGTH + 1) + + with caplog.at_level("DEBUG", logger="weather_briefing.publishers"): + async with httpx.AsyncClient(transport=httpx.MockTransport(lambda _: httpx.Response(200))) as client: + publisher = TelegramPublisher(client, "runtime-token", "runtime-chat", diagnostics) + await publisher.publish(RenderedMessage(body, len(body))) + + assert diagnostics.checks == 1 + assert caplog.text.count("Sensitive rendered text diagnostic: stage=telegram-chunk-") == 2 + + +async def test_rendered_text_is_not_logged_without_runtime_diagnostics(caplog) -> None: + delivery = DeliveryProvider(PlainTextRenderer(), NoopPublisher()) + + with caplog.at_level("DEBUG", logger="weather_briefing.publishers"): + await delivery.publish_alert("Private diagnostic title", "Private diagnostic body") + + assert "Private diagnostic" not in caplog.text + + +async def test_delivery_logs_rendered_text_when_runtime_diagnostics_are_enabled(caplog) -> None: + delivery = DeliveryProvider(PlainTextRenderer(), NoopPublisher(), diagnostics=EnabledDiagnostics()) + + with caplog.at_level("DEBUG", logger="weather_briefing.publishers"): + await delivery.publish_alert("Diagnostic title", "Diagnostic body") + + assert "Sensitive rendered text diagnostic: stage=alert body='Diagnostic title\\n\\nDiagnostic body'" in caplog.text + + +async def test_runtime_diagnostics_are_checked_without_debug_logging(caplog) -> None: + diagnostics = CountingDiagnostics() + delivery = DeliveryProvider(PlainTextRenderer(), NoopPublisher(), diagnostics=diagnostics) + + with caplog.at_level("INFO", logger="weather_briefing.publishers"): + await delivery.publish_alert("Private diagnostic title", "Private diagnostic body") + + assert diagnostics.checks == 1 + assert "Private diagnostic" not in caplog.text + + +async def test_runtime_diagnostic_failure_does_not_block_delivery(caplog) -> None: + publisher = RecordingPublisher() + delivery = DeliveryProvider(PlainTextRenderer(), publisher, diagnostics=FailingDiagnostics()) + + with caplog.at_level("DEBUG", logger="weather_briefing.publishers"): + await delivery.publish_alert("Diagnostic title", "Diagnostic body") + + assert publisher.messages == [RenderedMessage("Diagnostic title\n\nDiagnostic body", 33)] + assert "Rendered text diagnostic state check failed" in caplog.text async def test_telegram_error_does_not_expose_token() -> None: diff --git a/tests/test_service.py b/tests/test_service.py index 518c71c6..995d36fc 100644 --- a/tests/test_service.py +++ b/tests/test_service.py @@ -443,7 +443,7 @@ async def test_task_failure_alert_delivery_failure_is_retried( assert "Failed to publish or record briefing failure alert" in caplog.text -async def test_daily_briefing_publishes_verbatim_articles(tmp_path: Path) -> None: +async def test_daily_briefing_publishes_verbatim_articles(tmp_path: Path, caplog) -> None: timezone = pendulum.timezone("Asia/Shanghai") now = pendulum.datetime(2026, 7, 13, 8, tz=timezone) verbatim = Article( @@ -470,7 +470,7 @@ async def test_daily_briefing_publishes_verbatim_articles(tmp_path: Path) -> Non publisher = RecordingPublisher() delivery = DeliveryProvider(PlainTextRenderer(), publisher) - with SQLiteStateStore(tmp_path / "v.sqlite3") as state: + with caplog.at_level("DEBUG"), SQLiteStateStore(tmp_path / "v.sqlite3") as state: service = BriefingService( cast(Any, settings), _location(), @@ -483,9 +483,14 @@ async def test_daily_briefing_publishes_verbatim_articles(tmp_path: Path) -> Non ) await service.run("daily", now) - assert len(publisher.messages) == 2 - assert publisher.messages[1][0].body == "Forecast bulletin\n\nRaw forecast" - assert publisher.messages[1][1] is False + assert len(publisher.messages) == 2 + assert publisher.messages[1][0].body == "Forecast bulletin\n\nRaw forecast" + assert publisher.messages[1][1] is False + assert ( + "Publishing verbatim article: source=feed published_at=2026-07-13T08:00:00+08:00 content_characters=12" + ) in caplog.text + assert "Rendered verbatim message: visible_characters=31 payload_characters=31" in caplog.text + assert "Verbatim article published: source=feed published_at=2026-07-13T08:00:00+08:00" in caplog.text async def test_run_returns_none_when_no_content_and_no_warnings(tmp_path: Path) -> None: diff --git a/tests/test_sources.py b/tests/test_sources.py index 0527524d..95a1f9d2 100644 --- a/tests/test_sources.py +++ b/tests/test_sources.py @@ -7,7 +7,7 @@ from weather_briefing.sources import HTTPContextSource, RSSSource, SourceFetchError -async def test_rss_source_marks_configured_verbatim_article() -> None: +async def test_rss_source_marks_configured_verbatim_article(caplog) -> None: xml = """x oneOfficial forecast bulletin https://example.invalid/oneSun, 12 Jul 2026 23:30:00 GMT @@ -16,22 +16,27 @@ async def test_rss_source_marks_configured_verbatim_article() -> None: def handler(_: httpx.Request) -> httpx.Response: return httpx.Response(200, text=xml) - async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client: - source = RSSSource(client, max_attempts=1) - articles = await source.fetch( - FeedConfig( - id="authority", - name="Authority", - url="https://example.invalid/feed", - verbatim_title_patterns=("forecast bulletin",), + with caplog.at_level("DEBUG", logger="weather_briefing.sources"): + async with httpx.AsyncClient(transport=httpx.MockTransport(handler)) as client: + source = RSSSource(client, max_attempts=1) + articles = await source.fetch( + FeedConfig( + id="authority", + name="Authority", + url="https://example.invalid/feed", + verbatim_title_patterns=("forecast bulletin",), + ) ) - ) assert len(articles) == 1 assert articles[0].published_at.to_iso8601_string() == "2026-07-12T23:30:00Z" assert articles[0].published_at.in_timezone("Asia/Shanghai").to_iso8601_string() == "2026-07-13T07:30:00+08:00" assert articles[0].is_verbatim is True assert articles[0].content == "full body" + assert ( + "Parsed RSS article: source=authority published_at=2026-07-12T23:30:00+00:00 content_characters=9 verbatim=True" + ) in caplog.text + assert "example.invalid" not in caplog.text async def test_rss_source_skips_verbatim_article_without_cleaned_content() -> None: diff --git a/tests/test_state.py b/tests/test_state.py index 68ece135..c48c5dd6 100644 --- a/tests/test_state.py +++ b/tests/test_state.py @@ -1,9 +1,80 @@ +import sqlite3 +from contextlib import closing from pathlib import Path +from time import monotonic import pendulum +import pytest from weather_briefing.models import Article, SourceDocument, Warning -from weather_briefing.state import SQLiteStateStore +from weather_briefing.state import SQLiteRuntimeDiagnostics, SQLiteStateStore + + +def test_rendered_text_diagnostics_can_be_enabled_and_disabled(tmp_path: Path) -> None: + now = pendulum.datetime(2026, 7, 14, 7, tz="UTC") + expires_at = now.add(minutes=15) + with SQLiteRuntimeDiagnostics(tmp_path / "state.db") as diagnostics: + assert diagnostics.rendered_text_logging_until(now) is None + + diagnostics.enable_rendered_text_logging(expires_at) + assert diagnostics.rendered_text_logging_until(now) == expires_at + + diagnostics.disable_rendered_text_logging() + assert diagnostics.rendered_text_logging_until(now) is None + + +def test_rendered_text_diagnostics_reports_current_enabled_state(tmp_path: Path) -> None: + with SQLiteRuntimeDiagnostics(tmp_path / "state.db") as diagnostics: + diagnostics.enable_rendered_text_logging(pendulum.now("UTC").add(minutes=15)) + assert diagnostics.rendered_text_logging_enabled() + + diagnostics.disable_rendered_text_logging() + assert not diagnostics.rendered_text_logging_enabled() + + +def test_rendered_text_diagnostics_changes_are_visible_across_connections(tmp_path: Path) -> None: + now = pendulum.datetime(2026, 7, 14, 7, tz="UTC") + state_path = tmp_path / "state.db" + with ( + SQLiteRuntimeDiagnostics(state_path) as daemon_diagnostics, + SQLiteRuntimeDiagnostics(state_path) as cli_diagnostics, + ): + assert daemon_diagnostics.rendered_text_logging_until(now) is None + + expires_at = now.add(minutes=15) + cli_diagnostics.enable_rendered_text_logging(expires_at) + assert daemon_diagnostics.rendered_text_logging_until(now) == expires_at + + cli_diagnostics.disable_rendered_text_logging() + assert daemon_diagnostics.rendered_text_logging_until(now) is None + + +def test_rendered_text_diagnostics_fail_fast_when_database_is_locked(tmp_path: Path) -> None: + state_path = tmp_path / "state.db" + with ( + SQLiteRuntimeDiagnostics(state_path) as diagnostics, + closing(sqlite3.connect(state_path)) as lock_connection, + ): + lock_connection.execute("BEGIN EXCLUSIVE") + started_at = monotonic() + + with pytest.raises(sqlite3.OperationalError, match="locked"): + diagnostics.rendered_text_logging_until() + + assert monotonic() - started_at < 1 + + +def test_rendered_text_diagnostics_expire_automatically(tmp_path: Path, caplog) -> None: + now = pendulum.datetime(2026, 7, 14, 7, tz="UTC") + with SQLiteRuntimeDiagnostics(tmp_path / "state.db") as diagnostics: + diagnostics.enable_rendered_text_logging(now.add(minutes=15)) + caplog.clear() + + with caplog.at_level("WARNING", logger="weather_briefing.state"): + assert diagnostics.rendered_text_logging_until(now.add(minutes=16)) is None + assert diagnostics.rendered_text_logging_until(now.add(minutes=17)) is None + + assert caplog.text.count("Sensitive rendered text diagnostic logging expired") == 1 def test_articles_are_deduplicated(tmp_path: Path) -> None: diff --git a/weather_briefing/cli.py b/weather_briefing/cli.py index 2f151c6a..255a6275 100644 --- a/weather_briefing/cli.py +++ b/weather_briefing/cli.py @@ -3,8 +3,11 @@ import argparse import asyncio import logging +import re +import sqlite3 import sys -from collections.abc import Callable +from collections.abc import Callable, Iterator +from contextlib import AsyncExitStack, contextmanager from datetime import UTC, datetime from pathlib import Path @@ -19,7 +22,7 @@ AirQualityProvider, AQICNProvider, ) -from .config import Settings, weather_providers_for +from .config import Settings, state_path_from_env, weather_providers_for from .geocoding import ( CachedLocationResolver, FallbackGeocodingProvider, @@ -33,11 +36,11 @@ OpenAICompatibleChatCompletionsProvider, ) from .models import ResolvedLocation -from .publishers import DeliveryProvider, StdoutPublisher, TelegramPublisher +from .publishers import DeliveryProvider, RenderedTextDiagnostics, StdoutPublisher, TelegramPublisher from .render import PlainTextRenderer, TelegramHTMLRenderer from .service import BriefingService from .sources import HTTPContextSource, RSSSource -from .state import SQLiteStateStore +from .state import SQLiteRuntimeDiagnostics, SQLiteStateStore from .time_utils import parse_aware_datetime from .weather_context import ( AirQualitySupplementingWeatherProvider, @@ -59,9 +62,38 @@ def build_parser() -> argparse.ArgumentParser: run_parser.add_argument("--at", help="Override run time with an ISO-8601 timestamp including UTC offset") daemon_parser = subparsers.add_parser("daemon") daemon_parser.add_argument("--run-now", action="store_true", help="Run a briefing immediately before scheduling") + diagnostics_parser = subparsers.add_parser("diagnostics") + diagnostics_topics = diagnostics_parser.add_subparsers(dest="diagnostics_topic", required=True) + rendered_text_parser = diagnostics_topics.add_parser("rendered-text") + rendered_text_actions = rendered_text_parser.add_subparsers(dest="diagnostics_action", required=True) + enable_parser = rendered_text_actions.add_parser("enable") + enable_parser.add_argument( + "--for", + dest="duration_seconds", + required=True, + type=_diagnostic_duration_seconds, + metavar="DURATION", + help="Enable sensitive rendered-text logging temporarily, for example 15m or 1h (maximum 24h)", + ) + rendered_text_actions.add_parser("status") + rendered_text_actions.add_parser("disable") return parser +_DIAGNOSTIC_DURATION_PATTERN = re.compile(r"^(?P[1-9][0-9]*)(?P[smh])$") + + +def _diagnostic_duration_seconds(value: str) -> int: + match = _DIAGNOSTIC_DURATION_PATTERN.fullmatch(value) + if match is None: + raise argparse.ArgumentTypeError("duration must use a positive value followed by s, m, or h") + multipliers = {"s": 1, "m": 60, "h": 3600} + seconds = int(match.group("value")) * multipliers[match.group("unit")] + if seconds > 24 * 60 * 60: + raise argparse.ArgumentTypeError("duration cannot exceed 24h") + return seconds + + def _in_schedule(kind: str, now: pendulum.DateTime, settings: Settings) -> bool: if kind == "daily": return now.hour == settings.greeting_hour @@ -79,6 +111,21 @@ def _hour_in_cron(hour: int, cron_hour: str) -> bool: _LOGGER = logging.getLogger("weather_briefing") +@contextmanager +def _runtime_diagnostics(path: Path) -> Iterator[RenderedTextDiagnostics | None]: + try: + diagnostics = SQLiteRuntimeDiagnostics(path) + except (OSError, sqlite3.Error): + _LOGGER.warning( + "Runtime diagnostics unavailable; continuing without sensitive rendered text logging", + exc_info=True, + ) + yield None + return + with diagnostics: + yield diagnostics + + def _configure_logging(*, debug: bool) -> None: level = logging.DEBUG if debug else logging.INFO _fmt = logging.Formatter( @@ -108,8 +155,12 @@ async def run(kind: str, enforce_window: bool, at: str | None = None) -> None: _LOGGER.info("Skipping delayed %s run outside configured local-time window", kind) return _LOGGER.info("Starting %s briefing run at %s", kind, now.to_iso8601_string()) - async with httpx.AsyncClient(timeout=settings.http_timeout_seconds, follow_redirects=True) as client: - delivery = _delivery_provider(settings, client) + async with AsyncExitStack() as stack: + diagnostics = stack.enter_context(_runtime_diagnostics(settings.state_path)) + client = await stack.enter_async_context( + httpx.AsyncClient(timeout=settings.http_timeout_seconds, follow_redirects=True) + ) + delivery = _delivery_provider(settings, client, diagnostics) llm_provider = _llm_provider(settings, client) resolver = CachedLocationResolver( PrecisionReducingGeocodingProvider( @@ -193,16 +244,26 @@ def _llm_provider(settings: Settings, client: httpx.AsyncClient) -> LLMProvider: raise ValueError(f"Unsupported LLM provider: {settings.llm_provider}") -def _delivery_provider(settings: Settings, client: httpx.AsyncClient) -> DeliveryProvider: +def _delivery_provider( + settings: Settings, + client: httpx.AsyncClient, + diagnostics: RenderedTextDiagnostics | None = None, +) -> DeliveryProvider: if settings.publisher == "stdout": - return DeliveryProvider(PlainTextRenderer(), StdoutPublisher()) + return DeliveryProvider(PlainTextRenderer(), StdoutPublisher(), diagnostics=diagnostics) if settings.publisher == "telegram": if not settings.telegram_bot_token or not settings.telegram_chat_id: raise ValueError("Telegram publisher requires TELEGRAM_BOT_TOKEN and TELEGRAM_CHAT_ID") return DeliveryProvider( TelegramHTMLRenderer(), - TelegramPublisher(client, settings.telegram_bot_token, settings.telegram_chat_id), - TelegramPublisher.MAX_MESSAGE_LENGTH, + TelegramPublisher( + client, + settings.telegram_bot_token, + settings.telegram_chat_id, + diagnostics, + ), + single_message_limit=TelegramPublisher.MAX_MESSAGE_LENGTH, + diagnostics=diagnostics, ) raise ValueError(f"Unsupported publisher: {settings.publisher}") @@ -348,6 +409,35 @@ async def daemon(run_now: bool = False) -> None: await asyncio.Event().wait() +def _manage_rendered_text_diagnostics(action: str, duration_seconds: int | None = None) -> None: + with SQLiteRuntimeDiagnostics(state_path_from_env()) as diagnostics: + if action == "enable": + if duration_seconds is None: + raise ValueError("Rendered text diagnostics require a duration") + expires_at = pendulum.now("UTC").add(seconds=duration_seconds) + diagnostics.enable_rendered_text_logging(expires_at) + print( + "Rendered text diagnostic logging enabled until " + f"{expires_at.to_iso8601_string()}; rendered bodies require DEBUG logging" + ) + return + if action == "disable": + diagnostics.disable_rendered_text_logging() + print("Rendered text diagnostic logging disabled") + return + if action == "status": + expires_at = diagnostics.rendered_text_logging_until() + if expires_at is None: + print("Rendered text diagnostic logging is disabled") + else: + print( + "Rendered text diagnostic logging is enabled until " + f"{expires_at.to_iso8601_string()}; rendered bodies require DEBUG logging" + ) + return + raise ValueError(f"Unsupported rendered text diagnostics action: {action}") + + def main() -> None: load_dotenv(override=False) args = build_parser().parse_args() @@ -355,6 +445,11 @@ def main() -> None: try: if args.command == "daemon": asyncio.run(daemon(args.run_now)) + elif args.command == "diagnostics": + _manage_rendered_text_diagnostics( + args.diagnostics_action, + getattr(args, "duration_seconds", None), + ) else: asyncio.run(run(args.kind, args.enforce_window, args.at)) except Exception: diff --git a/weather_briefing/config.py b/weather_briefing/config.py index f11a43eb..7c80fdc6 100644 --- a/weather_briefing/config.py +++ b/weather_briefing/config.py @@ -116,6 +116,10 @@ def _configured_weather_providers() -> tuple[str, ...] | None: return providers +def state_path_from_env() -> Path: + return Path(_clean_env(os.getenv("BRIEFING_STATE_PATH", "state/weather.sqlite3"))) + + def weather_providers_for(location: ResolvedLocation, configured: tuple[str, ...] | None) -> tuple[str, ...]: if configured is not None: return configured @@ -319,7 +323,7 @@ def from_env(cls) -> Settings: open_meteo_api_key=_clean_env(os.getenv("OPEN_METEO_API_KEY")) or None, aqicn_api_token=_clean_env(os.getenv("AQICN_API_TOKEN")) or None, aqicn_base_url=_clean_env(os.getenv("AQICN_BASE_URL", "https://api.waqi.info")).rstrip("/"), - state_path=Path(_clean_env(os.getenv("BRIEFING_STATE_PATH", "state/weather.sqlite3"))), + state_path=state_path_from_env(), publisher=_clean_env(os.getenv("PUBLISHER", "telegram")), telegram_bot_token=_clean_env(os.getenv("TELEGRAM_BOT_TOKEN")) or None, telegram_chat_id=_clean_env(os.getenv("TELEGRAM_CHAT_ID")) or None, diff --git a/weather_briefing/publishers.py b/weather_briefing/publishers.py index 40ffdb8a..bb69967e 100644 --- a/weather_briefing/publishers.py +++ b/weather_briefing/publishers.py @@ -1,5 +1,6 @@ from __future__ import annotations +import logging from dataclasses import dataclass from typing import Protocol @@ -8,11 +9,17 @@ from .models import Article, BriefingResult, RenderedMessage, SourceDocument from .render import MessageRenderer +_LOGGER = logging.getLogger("weather_briefing.publishers") + class Publisher(Protocol): async def publish(self, message: RenderedMessage, *, single_message: bool = False) -> None: ... +class RenderedTextDiagnostics(Protocol): + def rendered_text_logging_enabled(self) -> bool: ... + + @dataclass(frozen=True, slots=True) class DeliveryProvider: """Bind a platform renderer to its message transport.""" @@ -20,6 +27,7 @@ class DeliveryProvider: renderer: MessageRenderer publisher: Publisher single_message_limit: int | None = None + diagnostics: RenderedTextDiagnostics | None = None def briefing_limit(self, configured_limit: int) -> int: if self.single_message_limit is None: @@ -40,13 +48,23 @@ async def publish_rendered( *, single_message: bool = False, ) -> None: + _log_rendered_text(self.diagnostics, "briefing", message.body) await self.publisher.publish(message, single_message=single_message) async def publish_verbatim(self, article: Article) -> None: - await self.publisher.publish(self.renderer.render_verbatim(article)) + message = self.renderer.render_verbatim(article) + _LOGGER.debug( + "Rendered verbatim message: visible_characters=%d payload_characters=%d", + message.visible_length, + len(message.body), + ) + _log_rendered_text(self.diagnostics, "verbatim", message.body) + await self.publisher.publish(message) async def publish_alert(self, title: str, body: str) -> None: - await self.publisher.publish(self.renderer.render_alert(title, body)) + message = self.renderer.render_alert(title, body) + _log_rendered_text(self.diagnostics, "alert", message.body) + await self.publisher.publish(message) class StdoutPublisher: @@ -61,16 +79,38 @@ class DeliveryError(RuntimeError): class TelegramPublisher: MAX_MESSAGE_LENGTH = 4096 - def __init__(self, client: httpx.AsyncClient, token: str, chat_id: str) -> None: + def __init__( + self, + client: httpx.AsyncClient, + token: str, + chat_id: str, + diagnostics: RenderedTextDiagnostics | None = None, + ) -> None: self._client = client self._url = f"https://api.telegram.org/bot{token}/sendMessage" self._chat_id = chat_id + self._diagnostics = diagnostics async def publish(self, message: RenderedMessage, *, single_message: bool = False) -> None: if single_message and message.visible_length > self.MAX_MESSAGE_LENGTH: raise DeliveryError("Telegram single message exceeds the platform limit") chunks = (message.body,) if single_message else _split_message(message.body, self.MAX_MESSAGE_LENGTH) - for chunk in chunks: + _LOGGER.debug( + "Telegram delivery prepared: visible_characters=%d payload_characters=%d chunks=%d single_message=%s", + message.visible_length, + len(message.body), + len(chunks), + single_message, + ) + log_rendered_text = _rendered_text_logging_enabled(self._diagnostics) + for index, chunk in enumerate(chunks, start=1): + if log_rendered_text: + _LOGGER.debug( + "Sensitive rendered text diagnostic: stage=telegram-chunk-%d-of-%d body=%r", + index, + len(chunks), + chunk, + ) try: response = await self._client.post( self._url, @@ -84,6 +124,32 @@ async def publish(self, message: RenderedMessage, *, single_message: bool = Fals response.raise_for_status() except httpx.HTTPError: raise DeliveryError("Telegram delivery failed") from None + _LOGGER.debug( + "Telegram chunk accepted: index=%d/%d payload_characters=%d", + index, + len(chunks), + len(chunk), + ) + + +def _log_rendered_text( + diagnostics: RenderedTextDiagnostics | None, + stage: str, + body: str, +) -> None: + if _rendered_text_logging_enabled(diagnostics): + _LOGGER.debug("Sensitive rendered text diagnostic: stage=%s body=%r", stage, body) + + +def _rendered_text_logging_enabled(diagnostics: RenderedTextDiagnostics | None) -> bool: + if diagnostics is None: + return False + try: + enabled = diagnostics.rendered_text_logging_enabled() + except Exception: + _LOGGER.warning("Rendered text diagnostic state check failed", exc_info=True) + return False + return enabled and _LOGGER.isEnabledFor(logging.DEBUG) def _split_message(body: str, limit: int) -> tuple[str, ...]: diff --git a/weather_briefing/service.py b/weather_briefing/service.py index 54532ccc..4fb31e30 100644 --- a/weather_briefing/service.py +++ b/weather_briefing/service.py @@ -205,7 +205,18 @@ def validate_length(candidate: BriefingResult) -> None: await self._delivery.publish_rendered(message, single_message=True) for article in new_articles: if article.is_verbatim: + _LOGGER.debug( + "Publishing verbatim article: source=%s published_at=%s content_characters=%d", + article.source_id, + article.published_at.isoformat(), + len(article.content), + ) await self._delivery.publish_verbatim(article) + _LOGGER.info( + "Verbatim article published: source=%s published_at=%s", + article.source_id, + article.published_at.isoformat(), + ) self._save_result_state( kind, now, diff --git a/weather_briefing/sources.py b/weather_briefing/sources.py index 3c3dc71f..66bd84e8 100644 --- a/weather_briefing/sources.py +++ b/weather_briefing/sources.py @@ -2,6 +2,7 @@ import asyncio import hashlib +import logging import random from time import struct_time @@ -12,6 +13,8 @@ from .content_cleaners import ContentCleaner, ContentCleaningRules, HTMLContentCleaner from .models import Article, ContextSourceConfig, FeedConfig, SourceDocument +_LOGGER = logging.getLogger("weather_briefing.sources") + class SourceFetchError(RuntimeError): """Raised after a source exhausts all retry attempts.""" @@ -69,6 +72,13 @@ async def fetch(self, config: FeedConfig) -> tuple[Article, ...]: ), ) is_verbatim = any(pattern in title for pattern in config.verbatim_title_patterns) + _LOGGER.debug( + "Parsed RSS article: source=%s published_at=%s content_characters=%d verbatim=%s", + config.id, + published_at.isoformat(), + len(content), + is_verbatim, + ) if is_verbatim and not content: continue articles.append( diff --git a/weather_briefing/state.py b/weather_briefing/state.py index 54989d55..f5359847 100644 --- a/weather_briefing/state.py +++ b/weather_briefing/state.py @@ -1,6 +1,7 @@ from __future__ import annotations import json +import logging import re import sqlite3 from pathlib import Path @@ -11,6 +12,81 @@ from .time_utils import require_aware_datetime _STORAGE_TIME_PATTERN = re.compile(r"^[0-9]{4}-[0-9]{2}-[0-9]{2}T[0-9]{2}:[0-9]{2}:[0-9]{2}\.[0-9]{6}Z$") +_RENDERED_TEXT_DIAGNOSTIC = "rendered_text" +_RUNTIME_DIAGNOSTIC_BUSY_TIMEOUT_SECONDS = 0.1 +_LOGGER = logging.getLogger("weather_briefing.state") + + +class SQLiteRuntimeDiagnostics: + def __init__(self, path: Path) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + self._connection = sqlite3.connect(path, timeout=_RUNTIME_DIAGNOSTIC_BUSY_TIMEOUT_SECONDS) + self._connection.row_factory = sqlite3.Row + self._connection.execute( + "CREATE TABLE IF NOT EXISTS runtime_diagnostics (name TEXT PRIMARY KEY, expires_at TEXT NOT NULL)" + ) + self._connection.commit() + + def close(self) -> None: + self._connection.close() + + def __enter__(self) -> SQLiteRuntimeDiagnostics: + return self + + def __exit__(self, *_: object) -> None: + self.close() + + def enable_rendered_text_logging(self, expires_at: pendulum.DateTime) -> None: + expires_at = require_aware_datetime(expires_at, context="Rendered text diagnostic expiration") + self._connection.execute( + """INSERT INTO runtime_diagnostics(name, expires_at) VALUES (?, ?) + ON CONFLICT(name) DO UPDATE SET expires_at = excluded.expires_at""", + (_RENDERED_TEXT_DIAGNOSTIC, _storage_time(expires_at)), + ) + self._connection.commit() + _LOGGER.warning( + "Sensitive rendered text diagnostic logging enabled until %s", + expires_at.to_iso8601_string(), + ) + + def disable_rendered_text_logging(self) -> None: + self._connection.execute( + "DELETE FROM runtime_diagnostics WHERE name = ?", + (_RENDERED_TEXT_DIAGNOSTIC,), + ) + self._connection.commit() + _LOGGER.warning("Sensitive rendered text diagnostic logging disabled") + + def rendered_text_logging_until( + self, + now: pendulum.DateTime | None = None, + ) -> pendulum.DateTime | None: + current_time = require_aware_datetime( + now or pendulum.now("UTC"), + context="Rendered text diagnostic check time", + ) + row = self._connection.execute( + "SELECT expires_at FROM runtime_diagnostics WHERE name = ?", + (_RENDERED_TEXT_DIAGNOSTIC,), + ).fetchone() + if row is None: + return None + expires_at = _parse_time(str(row["expires_at"])) + if expires_at > current_time: + return expires_at + self._connection.execute( + "DELETE FROM runtime_diagnostics WHERE name = ?", + (_RENDERED_TEXT_DIAGNOSTIC,), + ) + self._connection.commit() + _LOGGER.warning( + "Sensitive rendered text diagnostic logging expired at %s", + expires_at.to_iso8601_string(), + ) + return None + + def rendered_text_logging_enabled(self) -> bool: + return self.rendered_text_logging_until() is not None class SQLiteStateStore: