diff --git a/.github/workflows/rss_source_pipeline.yml b/.github/workflows/rss_source_pipeline.yml
index 7bd252b..7d9dd66 100644
--- a/.github/workflows/rss_source_pipeline.yml
+++ b/.github/workflows/rss_source_pipeline.yml
@@ -42,11 +42,39 @@ concurrency:
cancel-in-progress: false
jobs:
+ validate-workflow-boundary:
+ if: >-
+ github.repository == 'QuantStrategyLab/PoliticalEventTrackingResearch' &&
+ github.ref == 'refs/heads/main' &&
+ github.workflow_ref == 'QuantStrategyLab/PoliticalEventTrackingResearch/.github/workflows/rss_source_pipeline.yml@refs/heads/main'
+ runs-on: ubuntu-latest
+ steps:
+ - name: Validate canonical workflow inputs
+ env:
+ FEEDS_PATH: ${{ github.event.inputs.feeds_path || 'config/free_rss_feeds.csv' }}
+ ALIASES_PATH: ${{ github.event.inputs.aliases_path || 'config/core_us_equity_aliases.csv' }}
+ WATCHLIST_PATH: ${{ github.event.inputs.watchlist_path || 'data/live/political_watchlist.csv' }}
+ run: |
+ set -euo pipefail
+ test "${FEEDS_PATH}" = "config/free_rss_feeds.csv"
+ test "${ALIASES_PATH}" = "config/core_us_equity_aliases.csv"
+ test "${WATCHLIST_PATH}" = "data/live/political_watchlist.csv"
+
build-rss-source-events:
+ needs: validate-workflow-boundary
+ if: needs.validate-workflow-boundary.result == 'success'
runs-on: ubuntu-latest
timeout-minutes: 30
steps:
- uses: actions/checkout@v6
+ with:
+ repository: QuantStrategyLab/PoliticalEventTrackingResearch
+ ref: refs/heads/main
+ persist-credentials: false
+ - name: Verify reviewed main checkout
+ env:
+ EXPECTED_SHA: ${{ github.sha }}
+ run: test "$(git rev-parse HEAD)" = "${EXPECTED_SHA}"
- uses: actions/setup-python@v6
with:
python-version: "3.11"
@@ -61,29 +89,47 @@ jobs:
run: |
set -euo pipefail
mkdir -p data/output/rss_source_pipeline
+ fetch_exit=0
python scripts/fetch_rss_sources.py \
--feeds "${FEEDS_PATH}" \
--output data/output/rss_source_pipeline/source_items.csv \
--max-items-per-feed "${MAX_ITEMS_PER_FEED}" \
--continue-on-feed-error \
- --status-output data/output/rss_source_pipeline/source_fetch_status.json
- python scripts/extract_source_mentions.py \
- --raw-items data/output/rss_source_pipeline/source_items.csv \
- --aliases "${ALIASES_PATH}" \
- --output data/output/rss_source_pipeline/source_events.csv
- python scripts/build_tracker.py \
- --watchlist "${WATCHLIST_PATH}" \
- --events data/output/rss_source_pipeline/source_events.csv \
- --output data/output/rss_source_pipeline/source_tracker.csv
+ --status-output data/output/rss_source_pipeline/source_fetch_status.json || fetch_exit=$?
+ if [ -f data/output/rss_source_pipeline/source_items.csv ]; then
+ python scripts/extract_source_mentions.py \
+ --raw-items data/output/rss_source_pipeline/source_items.csv \
+ --aliases "${ALIASES_PATH}" \
+ --output data/output/rss_source_pipeline/source_events.csv
+ python scripts/build_tracker.py \
+ --watchlist "${WATCHLIST_PATH}" \
+ --events data/output/rss_source_pipeline/source_events.csv \
+ --output data/output/rss_source_pipeline/source_tracker.csv
+ fi
+ printf '%s\n' "${fetch_exit}" > data/output/rss_source_pipeline/fetch_exit.txt
+
+ - name: Upload RSS source artifact
+ uses: actions/upload-artifact@v7
+ with:
+ name: rss-source-pipeline
+ path: data/output/rss_source_pipeline/
+ if-no-files-found: error
+
+ - name: Validate canonical feed status readback
+ run: python scripts/validate_fetch_status.py --status data/output/rss_source_pipeline/source_fetch_status.json --fetch-exit data/output/rss_source_pipeline/fetch_exit.txt
+
- name: Publish live CSV outputs to repository
env:
COMMIT_OUTPUTS: ${{ github.event_name == 'schedule' && 'true' || github.event.inputs.commit_outputs || 'false' }}
+ EXPECTED_SHA: ${{ github.sha }}
+ GITHUB_TOKEN: ${{ github.token }}
run: |
set -euo pipefail
if [ "${COMMIT_OUTPUTS}" != "true" ]; then
echo "Live output commit disabled."
exit 0
fi
+ test "$(git rev-parse HEAD)" = "${EXPECTED_SHA}"
mkdir -p data/live
cp data/output/rss_source_pipeline/source_items.csv data/live/source_items.csv
cp data/output/rss_source_pipeline/source_events.csv data/live/source_events.csv
@@ -105,11 +151,10 @@ jobs:
echo "No live RSS output changes to commit."
else
git commit -m "Update live RSS source events [skip ci]"
- git push
+ git config --local http.https://github.com/.extraheader "AUTHORIZATION: bearer ${GITHUB_TOKEN}"
+ cleanup_git_auth() { git config --local --unset-all http.https://github.com/.extraheader >/dev/null 2>&1 || true; }
+ trap cleanup_git_auth EXIT
+ git push origin HEAD:refs/heads/main
+ cleanup_git_auth
+ trap - EXIT
fi
- - name: Upload RSS source artifact
- uses: actions/upload-artifact@v7
- with:
- name: rss-source-pipeline
- path: data/output/rss_source_pipeline/
- if-no-files-found: error
diff --git a/scripts/validate_fetch_status.py b/scripts/validate_fetch_status.py
new file mode 100644
index 0000000..5de9af1
--- /dev/null
+++ b/scripts/validate_fetch_status.py
@@ -0,0 +1,43 @@
+#!/usr/bin/env python3
+"""Fail closed unless canonical fetch status permits live publication."""
+from __future__ import annotations
+
+import argparse
+import json
+import sys
+from pathlib import Path
+
+sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src"))
+
+from political_event_tracking_research.feed_status_canonical_h2c import DecisionContractError, read_status # noqa: E402
+from political_event_tracking_research.rss_source_fetch import FetchStatusError # noqa: E402
+
+
+def validate_status_file(status_path: Path, fetch_exit: int) -> bool:
+ try:
+ status_bytes = status_path.read_bytes()
+ payload = read_status(status_bytes)
+ if type(fetch_exit) is not int or fetch_exit < 0:
+ raise FetchStatusError("fetch_exit_invalid")
+ if fetch_exit != 0:
+ raise FetchStatusError("fetch_failed")
+ except (DecisionContractError, FetchStatusError, OSError, UnicodeError):
+ raise SystemExit("fetch_status_invalid") from None
+ return payload["eligible_for_live_publication"] is True
+
+
+def main() -> None:
+ parser = argparse.ArgumentParser()
+ parser.add_argument("--status", required=True, type=Path)
+ parser.add_argument("--fetch-exit", required=True, type=Path)
+ args = parser.parse_args()
+ try:
+ fetch_exit = int(args.fetch_exit.read_text(encoding="utf-8").strip())
+ except (OSError, UnicodeError, ValueError):
+ raise SystemExit("fetch_status_invalid") from None
+ if not validate_status_file(args.status, fetch_exit):
+ raise SystemExit("fetch_incomplete")
+
+
+if __name__ == "__main__":
+ main()
diff --git a/src/political_event_tracking_research/rss_source_fetch.py b/src/political_event_tracking_research/rss_source_fetch.py
index 0884243..795c7e8 100644
--- a/src/political_event_tracking_research/rss_source_fetch.py
+++ b/src/political_event_tracking_research/rss_source_fetch.py
@@ -1,10 +1,12 @@
from __future__ import annotations
import argparse
+import csv
import datetime as dt
import email.utils
import hashlib
import html
+import io
import json
import re
import urllib.request
@@ -15,7 +17,14 @@
import defusedxml.ElementTree as ET
from defusedxml.common import DefusedXmlException
-from .csv_utils import read_csv_rows, write_csv_rows
+from .csv_utils import read_csv_rows
+from .feed_status_canonical_h2c import (
+ CanonicalDecision,
+ DecisionContractError,
+ DecisionKind,
+ build_decision,
+ read_status,
+)
USER_AGENT = (
@@ -23,6 +32,7 @@
"+https://github.com/QuantStrategyLab/PoliticalEventTrackingResearch)"
)
MAX_XML_BYTES = 1024 * 1024
+SOURCE_ITEM_FIELDS = ("item_id", "published_at", "source_type", "source_url", "author", "text")
@dataclass(frozen=True)
@@ -37,17 +47,19 @@ class FeedConfig:
class FeedFetchStatus:
feed_id: str
feed_url: str
- ok: bool
- item_count: int
- error: str = ""
+ kind: str
+ state: str
+ rows: tuple[dict[str, str], ...] = ()
+ error_code: str | None = None
- def to_json(self) -> dict[str, object]:
+ def to_outcome(self) -> dict[str, object]:
return {
"feed_id": self.feed_id,
"feed_url": self.feed_url,
- "ok": self.ok,
- "item_count": self.item_count,
- "error": self.error,
+ "kind": self.kind,
+ "state": self.state,
+ "rows": list(self.rows),
+ "error_code": self.error_code,
}
@@ -55,6 +67,25 @@ class FeedXmlError(ValueError):
"""Sanitized producer-boundary XML failure."""
+class FetchStatusError(ValueError):
+ def __init__(self, code: str):
+ super().__init__(code)
+ self.code = code
+
+
+def validate_fetch_status(payload: object, *, fetch_exit: int = 0) -> bool:
+ if type(fetch_exit) is not int or fetch_exit < 0:
+ raise FetchStatusError("fetch_exit_invalid")
+ if fetch_exit != 0:
+ raise FetchStatusError("fetch_failed")
+ try:
+ status_bytes = json.dumps(payload, ensure_ascii=False, sort_keys=True, separators=(",", ":")).encode("utf-8")
+ parsed = read_status(status_bytes)
+ except (DecisionContractError, TypeError, UnicodeError, ValueError):
+ raise FetchStatusError("fetch_status_invalid") from None
+ return parsed["eligible_for_live_publication"] is True
+
+
def load_feed_config(path: str | Path) -> list[FeedConfig]:
feeds: list[FeedConfig] = []
for row in read_csv_rows(path):
@@ -131,7 +162,9 @@ def stable_item_id(feed_id: str, link: str, title: str) -> str:
return f"{feed_id}-{digest}"
-def parse_feed_items(feed_bytes: bytes, feed: FeedConfig, *, max_items: int = 25) -> list[dict[str, str]]:
+def parse_feed_snapshot(
+ feed_bytes: bytes, feed: FeedConfig, *, max_items: int = 25
+) -> tuple[str, list[dict[str, str]]]:
if type(feed_bytes) is not bytes:
raise FeedXmlError("feed_xml_invalid")
if len(feed_bytes) > MAX_XML_BYTES:
@@ -142,8 +175,11 @@ def parse_feed_items(feed_bytes: bytes, feed: FeedConfig, *, max_items: int = 25
raise FeedXmlError("feed_xml_invalid") from None
rows: list[dict[str, str]] = []
- rss_items = root.findall("./channel/item")
- if rss_items:
+ if root.tag == "rss":
+ channel = root.find("./channel")
+ if channel is None:
+ raise FeedXmlError("feed_xml_invalid")
+ rss_items = channel.findall("./item")
for item in rss_items[:max_items]:
title = child_text(item, ("title",))
link = rss_item_link(item)
@@ -160,8 +196,10 @@ def parse_feed_items(feed_bytes: bytes, feed: FeedConfig, *, max_items: int = 25
"text": text,
}
)
- return rows
+ return "rss2", rows
+ if root.tag != "{http://www.w3.org/2005/Atom}feed":
+ raise FeedXmlError("feed_xml_invalid")
atom_entries = root.findall("{http://www.w3.org/2005/Atom}entry")
for entry in atom_entries[:max_items]:
title = child_text(entry, ("{http://www.w3.org/2005/Atom}title", "title"))
@@ -182,25 +220,62 @@ def parse_feed_items(feed_bytes: bytes, feed: FeedConfig, *, max_items: int = 25
"text": text,
}
)
- return rows
+ return "atom", rows
+
+
+def parse_feed_items(feed_bytes: bytes, feed: FeedConfig, *, max_items: int = 25) -> list[dict[str, str]]:
+ return parse_feed_snapshot(feed_bytes, feed, max_items=max_items)[1]
def utc_now_iso() -> str:
return dt.datetime.now(dt.UTC).replace(microsecond=0).isoformat().replace("+00:00", "Z")
-def write_fetch_status(path: str | Path, statuses: list[FeedFetchStatus], *, item_count: int) -> None:
- payload = {
- "generated_at": utc_now_iso(),
- "feed_count": len(statuses),
- "successful_feed_count": sum(1 for item in statuses if item.ok),
- "failed_feed_count": sum(1 for item in statuses if not item.ok),
- "item_count": item_count,
- "feeds": [item.to_json() for item in statuses],
- }
+def build_fetch_status(statuses: list[FeedFetchStatus]) -> CanonicalDecision:
+ if not statuses:
+ raise FetchStatusError("feed_config_empty")
+ try:
+ decision = build_decision(item.to_outcome() for item in statuses)
+ read_status(decision.status_bytes)
+ except DecisionContractError as exc:
+ raise FetchStatusError(exc.code) from None
+ return decision
+
+
+def write_fetch_status(path: str | Path, statuses: list[FeedFetchStatus]) -> DecisionKind:
+ decision = build_fetch_status(statuses)
output_path = Path(path)
output_path.parent.mkdir(parents=True, exist_ok=True)
- output_path.write_text(json.dumps(payload, ensure_ascii=False, indent=2, sort_keys=True) + "\n", encoding="utf-8")
+ output_path.write_bytes(decision.status_bytes)
+ return decision.decision.kind
+
+
+def serialize_source_items(rows: list[dict[str, str]]) -> bytes:
+ buffer = io.StringIO(newline="")
+ writer = csv.DictWriter(buffer, fieldnames=SOURCE_ITEM_FIELDS, lineterminator="\n")
+ writer.writeheader()
+ for row in rows:
+ writer.writerow({key: row.get(key, "") for key in SOURCE_ITEM_FIELDS})
+ return buffer.getvalue().encode("utf-8")
+
+
+def readback_source_items(
+ path: str | Path, expected_bytes: bytes, expected_rows: list[dict[str, str]]
+) -> list[dict[str, str]]:
+ try:
+ actual_bytes = Path(path).read_bytes()
+ if actual_bytes != expected_bytes:
+ raise FetchStatusError("source_items_bytes_mismatch")
+ text = actual_bytes.decode("utf-8")
+ reader = csv.DictReader(io.StringIO(text, newline=""))
+ if tuple(reader.fieldnames or ()) != SOURCE_ITEM_FIELDS:
+ raise FetchStatusError("source_items_schema_invalid")
+ actual_rows = [dict(row) for row in reader]
+ except (OSError, UnicodeError, csv.Error):
+ raise FetchStatusError("source_items_readback_invalid") from None
+ if actual_rows != expected_rows or serialize_source_items(actual_rows) != actual_bytes:
+ raise FetchStatusError("source_items_rows_mismatch")
+ return actual_rows
def fetch_rss_sources(
@@ -214,39 +289,62 @@ def fetch_rss_sources(
) -> list[dict[str, str]]:
rows: list[dict[str, str]] = []
statuses: list[FeedFetchStatus] = []
- for feed in load_feed_config(feeds_path):
+ feeds = load_feed_config(feeds_path)
+ if not feeds:
+ raise FetchStatusError("feed_config_empty")
+ for index, feed in enumerate(feeds):
try:
- feed_rows = parse_feed_items(fetcher(feed.feed_url), feed, max_items=max_items_per_feed)
+ kind, feed_rows = parse_feed_snapshot(fetcher(feed.feed_url), feed, max_items=max_items_per_feed)
except Exception as exc:
statuses.append(
FeedFetchStatus(
feed_id=feed.feed_id,
feed_url=feed.feed_url,
- ok=False,
- item_count=0,
- error=f"{type(exc).__name__}: {exc}",
+ kind="unknown",
+ state="failed",
+ error_code="fetch_failed",
)
)
if not continue_on_feed_error:
- raise
+ statuses.extend(
+ FeedFetchStatus(
+ feed_id=unattempted.feed_id,
+ feed_url=unattempted.feed_url,
+ kind="unknown",
+ state="failed",
+ error_code="not_attempted",
+ )
+ for unattempted in feeds[index + 1 :]
+ )
+ break
continue
rows.extend(feed_rows)
statuses.append(
FeedFetchStatus(
feed_id=feed.feed_id,
feed_url=feed.feed_url,
- ok=True,
- item_count=len(feed_rows),
+ kind=kind,
+ state="accepted" if feed_rows else "quarantined",
+ rows=tuple(feed_rows),
+ error_code=None if feed_rows else "zero_entries",
)
)
- if statuses and not any(item.ok for item in statuses):
- if status_output:
- write_fetch_status(status_output, statuses, item_count=0)
- raise RuntimeError("all configured RSS/Atom feeds failed")
rows.sort(key=lambda row: (row["published_at"], row["item_id"]))
- write_csv_rows(output_path, ["item_id", "published_at", "source_type", "source_url", "author", "text"], rows)
+ source_bytes = serialize_source_items(rows)
+ output_file = Path(output_path)
+ output_file.parent.mkdir(parents=True, exist_ok=True)
+ output_file.write_bytes(source_bytes)
+ readback_rows = readback_source_items(output_file, source_bytes, rows)
+ if readback_rows != rows:
+ raise FetchStatusError("source_items_rows_mismatch")
+ decision = build_fetch_status(statuses)
if status_output:
- write_fetch_status(status_output, statuses, item_count=len(rows))
+ Path(status_output).parent.mkdir(parents=True, exist_ok=True)
+ Path(status_output).write_bytes(decision.status_bytes)
+ if read_status(Path(status_output).read_bytes()) != json.loads(decision.status_bytes):
+ raise FetchStatusError("fetch_status_readback_mismatch")
+ if decision.decision.kind is DecisionKind.HARD_FAIL:
+ raise RuntimeError("feed_fetch_failed")
return rows
diff --git a/tests/test_rss_source_fetch.py b/tests/test_rss_source_fetch.py
index d1f324e..ab7218f 100644
--- a/tests/test_rss_source_fetch.py
+++ b/tests/test_rss_source_fetch.py
@@ -7,7 +7,15 @@
import pytest
from political_event_tracking_research import rss_source_fetch
-from political_event_tracking_research.rss_source_fetch import FeedConfig, fetch_rss_sources, parse_feed_items
+from political_event_tracking_research.feed_status_canonical_h2c import read_status
+from political_event_tracking_research.rss_source_fetch import (
+ FeedConfig,
+ fetch_rss_sources,
+ parse_feed_items,
+ parse_feed_snapshot,
+ readback_source_items,
+ serialize_source_items,
+)
def test_parse_rss_feed_items_to_source_items() -> None:
@@ -70,6 +78,19 @@ def test_parse_atom_feed_items_to_source_items() -> None:
assert "EVT2" in rows[0]["text"]
+@pytest.mark.parametrize(
+ ("payload", "kind"),
+ [
+ (b"", "rss2"),
+ (b"", "atom"),
+ ],
+)
+def test_empty_feed_keeps_parser_kind(payload: bytes, kind: str) -> None:
+ parsed_kind, rows = parse_feed_snapshot(payload, FeedConfig("x", "https://example.test", "official", ""))
+ assert parsed_kind == kind
+ assert rows == []
+
+
@pytest.mark.parametrize(
"payload",
[
@@ -159,20 +180,22 @@ def fake_fetch(url: str) -> bytes:
output = tmp_path / "source_items.csv"
status = tmp_path / "status.json"
- rows = fetch_rss_sources(
- feeds_path,
- output,
- continue_on_feed_error=True,
- status_output=status,
- fetcher=fake_fetch,
- )
+ with pytest.raises(RuntimeError, match="feed_fetch_failed"):
+ fetch_rss_sources(
+ feeds_path,
+ output,
+ continue_on_feed_error=True,
+ status_output=status,
+ fetcher=fake_fetch,
+ )
- assert len(rows) == 1
- payload = json.loads(status.read_text(encoding="utf-8"))
+ assert output.exists()
+ payload = read_status(status.read_bytes())
assert payload["successful_feed_count"] == 1
assert payload["failed_feed_count"] == 1
- assert payload["feeds"][1]["feed_id"] == "bad"
- assert "RuntimeError" in payload["feeds"][1]["error"]
+ failed = next(item for item in payload["feeds"] if item["feed_id"] == "bad")
+ assert failed["state"] == "failed"
+ assert failed["error_code"] == "fetch_failed"
def test_fetch_rss_sources_fails_when_all_feeds_fail(tmp_path: Path) -> None:
@@ -186,7 +209,7 @@ def test_fetch_rss_sources_fails_when_all_feeds_fail(tmp_path: Path) -> None:
def fake_fetch(_url: str) -> bytes:
raise RuntimeError("blocked")
- with pytest.raises(RuntimeError, match="all configured"):
+ with pytest.raises(RuntimeError, match="feed_fetch_failed"):
fetch_rss_sources(
feeds_path,
tmp_path / "source_items.csv",
@@ -194,3 +217,87 @@ def fake_fetch(_url: str) -> bytes:
status_output=tmp_path / "status.json",
fetcher=fake_fetch,
)
+ payload = read_status((tmp_path / "status.json").read_bytes())
+ assert payload["failed_feed_count"] == 1
+ assert payload["eligible_for_live_publication"] is False
+
+
+def test_zero_entry_quarantine_writes_canonical_noneligible_status(tmp_path: Path) -> None:
+ feeds_path = tmp_path / "feeds.csv"
+ feeds_path.write_text(
+ "feed_id,feed_url,source_type,author\nempty,https://example.invalid/empty,official,\n",
+ encoding="utf-8",
+ )
+ output = tmp_path / "source_items.csv"
+ status = tmp_path / "status.json"
+
+ rows = fetch_rss_sources(
+ feeds_path,
+ output,
+ continue_on_feed_error=True,
+ status_output=status,
+ fetcher=lambda _url: b"",
+ )
+
+ assert rows == []
+ payload = read_status(status.read_bytes())
+ assert payload["quarantined_feed_count"] == 1
+ assert payload["eligible_for_live_publication"] is False
+ assert status.read_bytes() == json.dumps(
+ payload, ensure_ascii=False, sort_keys=True, separators=(",", ":")
+ ).encode()
+
+
+def test_source_item_readback_rejects_tamper_and_stale_bytes(tmp_path: Path) -> None:
+ rows = [
+ {
+ "item_id": "a-1",
+ "published_at": "2026-05-01T00:00:00Z",
+ "source_type": "official",
+ "source_url": "https://example.test/a/1",
+ "author": "",
+ "text": "event",
+ }
+ ]
+ path = tmp_path / "source_items.csv"
+ expected = serialize_source_items(rows)
+ path.write_bytes(expected)
+ assert readback_source_items(path, expected, rows) == rows
+ path.write_bytes(expected.replace(b"event", b"tampered"))
+ with pytest.raises(ValueError, match="source_items_bytes_mismatch"):
+ readback_source_items(path, expected, rows)
+
+
+def test_hard_fail_is_raised_without_status_output(tmp_path: Path) -> None:
+ feeds = tmp_path / "feeds.csv"
+ feeds.write_text(
+ "feed_id,feed_url,source_type,author\nbad,https://example.invalid/bad,official,\n",
+ encoding="utf-8",
+ )
+ with pytest.raises(RuntimeError, match="feed_fetch_failed"):
+ fetch_rss_sources(
+ feeds,
+ tmp_path / "items.csv",
+ fetcher=lambda _url: (_ for _ in ()).throw(RuntimeError("blocked")),
+ )
+
+
+def test_early_abort_status_covers_all_configured_feeds(tmp_path: Path) -> None:
+ feeds = tmp_path / "feeds.csv"
+ feeds.write_text(
+ "feed_id,feed_url,source_type,author\n"
+ "first,https://example.invalid/first,official,\n"
+ "second,https://example.invalid/second,official,\n",
+ encoding="utf-8",
+ )
+ status = tmp_path / "status.json"
+ with pytest.raises(RuntimeError, match="feed_fetch_failed"):
+ fetch_rss_sources(
+ feeds,
+ tmp_path / "items.csv",
+ status_output=status,
+ fetcher=lambda _url: (_ for _ in ()).throw(RuntimeError("blocked")),
+ )
+ payload = read_status(status.read_bytes())
+ assert payload["feed_count"] == 2
+ assert {item["error_code"] for item in payload["feeds"]} == {"fetch_failed", "not_attempted"}
diff --git a/tests/test_weekly_workflow_w0b.py b/tests/test_weekly_workflow_w0b.py
new file mode 100644
index 0000000..4afccb0
--- /dev/null
+++ b/tests/test_weekly_workflow_w0b.py
@@ -0,0 +1,85 @@
+from __future__ import annotations
+
+import json
+import sys
+from pathlib import Path
+
+import pytest
+
+sys.path.insert(0, str(Path(__file__).parents[1]))
+
+from political_event_tracking_research.feed_status_canonical_h2c import build_decision
+from political_event_tracking_research.rss_source_fetch import (
+ FetchStatusError,
+ fetch_rss_sources,
+ validate_fetch_status,
+)
+from scripts.validate_fetch_status import validate_status_file
+
+
+def accepted_status() -> dict[str, object]:
+ return json.loads(
+ build_decision(
+ [
+ {
+ "feed_id": "a",
+ "feed_url": "https://example.test/a",
+ "kind": "rss2",
+ "state": "accepted",
+ "rows": [
+ {
+ "item_id": "a-1",
+ "published_at": "2026-05-01T00:00:00Z",
+ "source_type": "official",
+ "source_url": "https://example.test/a/1",
+ "author": "",
+ "text": "event",
+ }
+ ],
+ "error_code": None,
+ }
+ ]
+ ).status_bytes
+ )
+
+
+def test_status_readback_requires_canonical_status_and_success_exit(tmp_path: Path) -> None:
+ status = tmp_path / "status.json"
+ status.write_text(json.dumps(accepted_status(), sort_keys=True, separators=(",", ":")), encoding="utf-8")
+ fetch_exit = tmp_path / "fetch_exit.txt"
+ fetch_exit.write_text("0\n", encoding="utf-8")
+ assert validate_status_file(status, 0) is True
+ assert validate_fetch_status(json.loads(status.read_text()), fetch_exit=0) is True
+
+ fetch_exit.write_text("7\n", encoding="utf-8")
+ with pytest.raises(SystemExit, match="fetch_status_invalid"):
+ validate_status_file(status, 7)
+
+
+def test_status_tamper_is_rejected_before_publication(tmp_path: Path) -> None:
+ status = tmp_path / "status.json"
+ payload = accepted_status()
+ payload["accepted_row_count"] = 0
+ status.write_text(json.dumps(payload, sort_keys=True, separators=(",", ":")), encoding="utf-8")
+ with pytest.raises(SystemExit, match="fetch_status_invalid"):
+ validate_status_file(status, 0)
+
+
+def test_empty_feed_configuration_fails_closed(tmp_path: Path) -> None:
+ feeds = tmp_path / "feeds.csv"
+ feeds.write_text("feed_id,feed_url,source_type,author\n", encoding="utf-8")
+ with pytest.raises(FetchStatusError, match="feed_config_empty"):
+ fetch_rss_sources(feeds, tmp_path / "items.csv", status_output=tmp_path / "status.json")
+
+
+def test_workflow_uploads_and_validates_before_live_publish() -> None:
+ workflow = Path(__file__).parents[1].joinpath(".github/workflows/rss_source_pipeline.yml").read_text(
+ encoding="utf-8"
+ )
+ assert workflow.index("Validate canonical workflow inputs") < workflow.index("actions/checkout@")
+ assert workflow.index("Upload RSS source artifact") < workflow.index("Validate canonical feed status readback")
+ assert workflow.index("Validate canonical feed status readback") < workflow.index("Publish live CSV outputs")
+ assert "git push origin HEAD:refs/heads/main" in workflow
+ assert "ref: refs/heads/main" in workflow
+ assert "needs: validate-workflow-boundary" in workflow
+ assert "weekly" not in workflow.lower()