Skip to content
Open
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
244 changes: 244 additions & 0 deletions src/political_event_tracking_research/fetch_acceptance.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,244 @@
"""Pure RSS/Atom content acceptance and canonical fetch-status contract."""
from __future__ import annotations

import email.utils
import json
import re
import xml.etree.ElementTree as ET
from collections.abc import Iterable, Mapping
from datetime import datetime, timezone

_STATUS_KEYS = frozenset(
{
"status_version",
"configured_feed_count",
"feed_count",
"successful_feed_count",
"failed_feed_count",
"stale_feed_count",
"missing_feed_count",
"quarantined_feed_count",
"accepted_row_count",
"rejected_row_count",
"publication_complete",
"eligible_for_live_publication",
"zero_entry_policy",
"feeds",
}
)
_FEED_KEYS = frozenset({"feed_id", "feed_url", "kind", "accepted_row_count", "rejected_row_count", "state", "error_code"})
_STATES = frozenset({"accepted", "failed", "stale", "missing", "quarantined"})
_KINDS = frozenset({"rss", "atom", "unknown"})
_SAFE_ERROR = re.compile(r"^[a-z][a-z0-9_]*$")
_ATOM = "http://www.w3.org/2005/Atom"
MAX_XML_BYTES = 1024 * 1024
MAX_XML_DEPTH = 32
MAX_XML_NODES = 10000
MAX_XML_TEXT_BYTES = 256 * 1024
MAX_XML_ATTRIBUTES = 128
_FORBIDDEN_DECLARATION = re.compile(rb"<!\s*(?:doctype|entity|element|attlist|notation)\b|\b(?:system|public)\b", re.IGNORECASE)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Restrict forbidden XML keywords to declarations

This pattern matches standalone public or system anywhere in the raw payload, not just inside <!DOCTYPE ... PUBLIC/SYSTEM ...> declarations. A valid political RSS item such as <title>Public hearing</title> is therefore classified as xml_forbidden_declaration, causing the feed to fail and preventing live publication even though the XML is well-formed RSS/Atom; scope those alternatives to declaration markup instead of scanning all content.

Useful? React with 👍 / 👎.



class FetchAcceptanceError(ValueError):
def __init__(self, code: str) -> None:
self.code = code
super().__init__(code)


def _fail(code: str) -> FetchAcceptanceError:
return FetchAcceptanceError(code)


def _text(element: ET.Element | None, names: tuple[str, ...]) -> str:
if element is None:
return ""
for name in names:
child = element.find(name)
if child is not None and child.text:
return child.text.strip()
return ""


def _atom_link(entry: ET.Element) -> str:
for child in entry:
if child.tag == f"{{{_ATOM}}}link" and child.attrib.get("href"):
return child.attrib["href"]
Comment on lines +64 to +65

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Strip Atom hrefs before accepting rows

For Atom entries whose only link is href=" ", this condition treats the whitespace as a usable link, so row_valid counts the entry as accepted and build_acceptance_status can mark the feed publishable even though the downstream source_url would be blank/whitespace. Strip the attribute and require a non-empty value before returning it, matching how RSS text fields are handled.

Useful? React with 👍 / 👎.

return _text(entry, (f"{{{_ATOM}}}id", "id"))


def _valid_timestamp(value: str) -> bool:
if not value:
return False
try:
parsed = email.utils.parsedate_to_datetime(value)
except (TypeError, ValueError):
try:
normalized = value[:-1] + "+00:00" if value.endswith("Z") else value
parsed = datetime.fromisoformat(normalized)
except (TypeError, ValueError):
return False
return parsed.tzinfo is not None and parsed.astimezone(timezone.utc) is not None

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Treat overflowed timestamps as invalid rows

For feeds containing a syntactically valid timestamp that overflows when converted to UTC (for example 9999-12-31T23:59:59-23:59), fromisoformat/parsedate_to_datetime succeeds but this astimezone call raises OverflowError. That exception escapes classify_feed_payload, so a single malformed entry aborts classification instead of being counted as rejected and quarantining the feed.

Useful? React with 👍 / 👎.



def _feed_result(feed_id: str, feed_url: str, kind: str, accepted: int, rejected: int, state: str, error_code: str | None) -> dict[str, object]:
if type(feed_id) is not str or not feed_id or type(feed_url) is not str or not feed_url or kind not in _KINDS or state not in _STATES or type(accepted) is not int or accepted < 0 or type(rejected) is not int or rejected < 0 or type(error_code) not in {str, type(None)} or (error_code is not None and not _SAFE_ERROR.fullmatch(error_code)):
raise _fail("feed_result_invalid")
return {"feed_id": feed_id, "feed_url": feed_url, "kind": kind, "accepted_row_count": accepted, "rejected_row_count": rejected, "state": state, "error_code": error_code}


def _parse_bounded_xml(payload: bytes) -> ET.Element:
if type(payload) is not bytes:
raise _fail("feed_payload_invalid")
if len(payload) > MAX_XML_BYTES:
raise _fail("xml_oversize")
if _FORBIDDEN_DECLARATION.search(payload):

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Normalize encoded XML before scanning declarations

When a feed declares UTF-16/UTF-16LE, this raw byte scan never sees contiguous <!DOCTYPE or <!ENTITY tokens, but ET.fromstring below honors the encoding and expands the internal entity; a UTF-16 RSS payload with a forbidden DTD is currently returned as an accepted feed instead of xml_forbidden_declaration. Since the contract is fail-closed for forbidden declarations, normalize or restrict encodings before this check.

Useful? React with 👍 / 👎.

raise _fail("xml_forbidden_declaration")
try:
root = ET.fromstring(payload)
except (ET.ParseError, UnicodeError, ValueError, RecursionError):
raise _fail("xml_invalid") from None
nodes = 0
text_bytes = 0
stack: list[tuple[ET.Element, int]] = [(root, 1)]
while stack:
element, depth = stack.pop()
nodes += 1
if nodes > MAX_XML_NODES or depth > MAX_XML_DEPTH or len(element.attrib) > MAX_XML_ATTRIBUTES:
raise _fail("xml_structure_over_limit")
for text in (element.text, element.tail):
if text is not None:
text_bytes += len(text.encode("utf-8", errors="strict"))
if text_bytes > MAX_XML_TEXT_BYTES:
raise _fail("xml_structure_over_limit")
stack.extend((child, depth + 1) for child in reversed(list(element)))
return root


def classify_feed_payload(feed_id: str, feed_url: str, payload: bytes) -> dict[str, object]:
if type(payload) is not bytes:
raise _fail("feed_payload_invalid")
try:
root = _parse_bounded_xml(payload)
except FetchAcceptanceError as error:
return _feed_result(feed_id, feed_url, "unknown", 0, 0, "failed", error.code)
if root.tag == "rss":
if root.attrib.get("version") not in {"2.0"}:
return _feed_result(feed_id, feed_url, "unknown", 0, 0, "failed", "rss_schema_unsupported")
channel = root.find("./channel")
if channel is None:
return _feed_result(feed_id, feed_url, "rss", 0, 0, "failed", "rss_schema_invalid")
entries = channel.findall("./item")
kind = "rss"
def row_valid(entry: ET.Element) -> bool:
return bool(_text(entry, ("title",)) and (_text(entry, ("link",)) or _text(entry, ("guid",))) and _valid_timestamp(_text(entry, ("pubDate", "{http://purl.org/dc/elements/1.1/}date"))))
elif root.tag == f"{{{_ATOM}}}feed":
entries = root.findall(f"{{{_ATOM}}}entry")
kind = "atom"
def row_valid(entry: ET.Element) -> bool:
published = _text(entry, (f"{{{_ATOM}}}published", f"{{{_ATOM}}}updated"))
return bool(_text(entry, (f"{{{_ATOM}}}title",)) and _atom_link(entry) and _valid_timestamp(published))
else:
return _feed_result(feed_id, feed_url, "unknown", 0, 0, "failed", "payload_not_rss_atom")
accepted = sum(row_valid(entry) for entry in entries)
rejected = len(entries) - accepted
if not entries:
return _feed_result(feed_id, feed_url, kind, 0, 0, "quarantined", "zero_entries")
if rejected:
return _feed_result(feed_id, feed_url, kind, accepted, rejected, "quarantined", "entry_invalid")
return _feed_result(feed_id, feed_url, kind, accepted, 0, "accepted", None)


def _validate_feed_result(value: object) -> dict[str, object]:
if not isinstance(value, Mapping) or set(value) != _FEED_KEYS:
raise _fail("feed_result_shape_invalid")
result = _feed_result(value["feed_id"], value["feed_url"], value["kind"], value["accepted_row_count"], value["rejected_row_count"], value["state"], value["error_code"])
if result != dict(value):
raise _fail("feed_result_noncanonical")
state = result["state"]
accepted_rows = result["accepted_row_count"]
rejected_rows = result["rejected_row_count"]
error_code = result["error_code"]
if state == "accepted" and (accepted_rows <= 0 or rejected_rows != 0 or error_code is not None):
raise _fail("feed_result_invalid")
if state in {"failed", "stale", "missing"} and (accepted_rows != 0 or rejected_rows != 0 or not error_code):
raise _fail("feed_result_invalid")
if state == "quarantined" and not error_code:
raise _fail("feed_result_invalid")
Comment on lines +165 to +166

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Validate quarantined counters against the reason

When external feed states are passed directly, this check only requires any error code for quarantined feeds, so a result like state="quarantined", error_code="zero_entries", and a positive accepted_row_count is accepted and serialized as canonical. That makes the strict status report accepted rows for a feed that claims to have zero entries, undermining the exact counter contract; reject quarantine reasons whose row counts are inconsistent with the reason.

Useful? React with 👍 / 👎.

return result


def build_acceptance_status(feed_results: Iterable[Mapping[str, object]]) -> dict[str, object]:
if not isinstance(feed_results, list) or not feed_results:

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Accept all iterables promised by the API

The public signature accepts Iterable[Mapping[str, object]], but this guard rejects any non-list iterable before reading it. A natural integration such as passing a tuple or generator of classify_feed_payload(...) results will raise configured_feed_empty even when results are present; convert the iterable to a list before checking emptiness, or narrow the signature to list.

Useful? React with 👍 / 👎.

raise _fail("configured_feed_empty")
feeds = [_validate_feed_result(item) for item in feed_results]
if len({item["feed_id"] for item in feeds}) != len(feeds):
raise _fail("feed_duplicate")
feeds.sort(key=lambda item: (item["feed_id"], item["feed_url"]))
accepted = sum(item["state"] == "accepted" for item in feeds)
failed = sum(item["state"] == "failed" for item in feeds)
stale = sum(item["state"] == "stale" for item in feeds)
missing = sum(item["state"] == "missing" for item in feeds)
quarantined = sum(item["state"] == "quarantined" for item in feeds)
accepted_rows = sum(item["accepted_row_count"] for item in feeds)
rejected_rows = sum(item["rejected_row_count"] for item in feeds)
complete = len(feeds) > 0 and accepted == len(feeds) and not any((failed, stale, missing, quarantined)) and accepted_rows > 0
return {
"status_version": "pert.fetch_acceptance.v1",
"configured_feed_count": len(feeds),
"feed_count": len(feeds),
"successful_feed_count": accepted,
"failed_feed_count": failed,
"stale_feed_count": stale,
"missing_feed_count": missing,
"quarantined_feed_count": quarantined,
"accepted_row_count": accepted_rows,
"rejected_row_count": rejected_rows,
"publication_complete": complete,
"eligible_for_live_publication": complete,
"zero_entry_policy": "quarantine",
"feeds": feeds,
}


def _validate_status(value: object) -> dict[str, object]:
if not isinstance(value, Mapping) or set(value) != _STATUS_KEYS or value.get("status_version") != "pert.fetch_acceptance.v1" or value.get("zero_entry_policy") != "quarantine":
raise _fail("status_shape_invalid")
integer_keys = ("configured_feed_count", "feed_count", "successful_feed_count", "failed_feed_count", "stale_feed_count", "missing_feed_count", "quarantined_feed_count", "accepted_row_count", "rejected_row_count")
if any(type(value[key]) is not int or value[key] < 0 for key in integer_keys) or type(value["publication_complete"]) is not bool or type(value["eligible_for_live_publication"]) is not bool or value["publication_complete"] != value["eligible_for_live_publication"] or not isinstance(value["feeds"], list):
raise _fail("status_counter_invalid")
expected = build_acceptance_status(value["feeds"])
if dict(value) != expected:
raise _fail("status_counter_mismatch")
return expected


def serialize_status(value: Mapping[str, object]) -> bytes:
try:
payload = _validate_status(value)
return json.dumps(payload, ensure_ascii=False, sort_keys=True, separators=(",", ":"), allow_nan=False).encode("utf-8")
except FetchAcceptanceError:
raise
except (TypeError, ValueError, UnicodeError, OverflowError, RecursionError):
raise _fail("status_serialization_invalid") from None


def parse_status_bytes(wire: bytes) -> dict[str, object]:
if type(wire) is not bytes:
raise _fail("status_wire_invalid")
def pairs(items: list[tuple[str, object]]) -> dict[str, object]:
result: dict[str, object] = {}
for key, item in items:
if key in result:
raise _fail("status_duplicate_key")
result[key] = item
return result
try:
value = json.loads(wire.decode("utf-8"), object_pairs_hook=pairs)
except FetchAcceptanceError:
raise
except (UnicodeError, json.JSONDecodeError, TypeError, ValueError, RecursionError):
raise _fail("status_wire_invalid") from None
parsed = _validate_status(value)
if serialize_status(parsed) != wire:
raise _fail("status_noncanonical")
return parsed
120 changes: 120 additions & 0 deletions tests/test_fetch_acceptance.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,120 @@
from __future__ import annotations

import json

import pytest

from political_event_tracking_research.fetch_acceptance import (
FetchAcceptanceError,
build_acceptance_status,
classify_feed_payload,
parse_status_bytes,
serialize_status,
)


RSS = b"""<?xml version=\"1.0\"?><rss version=\"2.0\"><channel><title>Feed</title><item><title>Event</title><link>https://example.test/event</link><pubDate>Fri, 01 May 2026 12:30:00 GMT</pubDate><description>Text</description></item></channel></rss>"""
ATOM = b"""<?xml version=\"1.0\"?><feed xmlns=\"http://www.w3.org/2005/Atom\"><title>Feed</title><entry><title>Event</title><link href=\"https://example.test/event\"/><updated>2026-05-01T12:30:00Z</updated><summary>Text</summary></entry></feed>"""
NON_FEED = b"<root><item><title>Not RSS</title></item></root>"


def test_rss_and_atom_are_recognized_with_exact_rows() -> None:
rss = classify_feed_payload("rss", "https://example.test/rss", RSS)
atom = classify_feed_payload("atom", "https://example.test/atom", ATOM)
assert (rss["kind"], rss["accepted_row_count"], rss["rejected_row_count"]) == ("rss", 1, 0)
assert (atom["kind"], atom["accepted_row_count"], atom["rejected_row_count"]) == ("atom", 1, 0)


@pytest.mark.parametrize("payload", [NON_FEED, b"<rss>", b"not xml"])
def test_non_rss_or_malformed_payload_fails_closed(payload: bytes) -> None:
result = classify_feed_payload("feed", "https://example.test/feed", payload)
assert result["state"] == "failed"
assert result["accepted_row_count"] == 0


def test_unsupported_rss_schema_is_not_successful() -> None:
result = classify_feed_payload("feed", "https://example.test/feed", RSS.replace(b'version="2.0"', b'version="1.0"'))
assert result["state"] == "failed"
assert result["error_code"] == "rss_schema_unsupported"


def test_zero_entry_is_quarantined_not_complete() -> None:
result = classify_feed_payload("empty", "https://example.test/empty", b"<rss version='2.0'><channel><title>Empty</title></channel></rss>")
status = build_acceptance_status([result])
assert result["state"] == "quarantined"
assert status["publication_complete"] is False
assert status["eligible_for_live_publication"] is False
assert status["zero_entry_policy"] == "quarantine"


def test_mixed_partial_counts_are_exact_and_incomplete() -> None:
accepted = classify_feed_payload("ok", "https://example.test/ok", RSS)
failed = classify_feed_payload("bad", "https://example.test/bad", NON_FEED)
status = build_acceptance_status([accepted, failed])
assert status["configured_feed_count"] == 2
assert status["successful_feed_count"] == 1
assert status["failed_feed_count"] == 1
assert status["accepted_row_count"] == 1
assert status["publication_complete"] is False


def test_empty_config_and_counter_mismatch_fail_closed() -> None:
with pytest.raises(FetchAcceptanceError):
build_acceptance_status([])
status = build_acceptance_status([classify_feed_payload("ok", "https://example.test/ok", RSS)])
tampered = dict(status)
tampered["accepted_row_count"] = 999
with pytest.raises(FetchAcceptanceError):
serialize_status(tampered)


def test_status_canonical_roundtrip_and_tamper_rejection() -> None:
status = build_acceptance_status([classify_feed_payload("ok", "https://example.test/ok", RSS)])
wire = serialize_status(status)
assert parse_status_bytes(wire) == status
assert serialize_status(parse_status_bytes(wire)) == wire
with pytest.raises(FetchAcceptanceError):
parse_status_bytes(b" " + wire)
payload = json.loads(wire)
payload["unknown"] = True
with pytest.raises(FetchAcceptanceError):
serialize_status(payload)


def test_invalid_entry_is_counted_not_accepted() -> None:
payload = RSS.replace(b"<link>https://example.test/event</link>", b"")
result = classify_feed_payload("feed", "https://example.test/feed", payload)
assert result["accepted_row_count"] == 0
assert result["rejected_row_count"] == 1
assert result["state"] == "quarantined"


@pytest.mark.parametrize("declaration", ["<!DOCTYPE rss>", "<!doctype rss>", "<! DOCTYPE rss>", "<!\nENTITY x \"x\">", "<!SYSTEM \"file:///tmp/x\">"])
def test_forbidden_xml_declarations_are_sanitized(declaration: str) -> None:
result = classify_feed_payload("feed", "https://example.test/feed", (declaration + "<rss version='2.0'><channel/></rss>").encode())
assert result["state"] == "failed"
assert result["error_code"] == "xml_forbidden_declaration"


def test_xml_size_depth_nodes_text_and_attributes_are_bounded() -> None:
assert classify_feed_payload("feed", "https://example.test/feed", RSS + b"x" * (1024 * 1024))["error_code"] == "xml_oversize"
deep = "<rss version='2.0'><channel>" + "<x>" * 40 + "</x>" * 40 + "</channel></rss>"
assert classify_feed_payload("feed", "https://example.test/feed", deep.encode())["error_code"] == "xml_structure_over_limit"
attrs = "<rss version='2.0' " + " ".join(f"a{i}='x'" for i in range(130)) + "><channel/></rss>"
assert classify_feed_payload("feed", "https://example.test/feed", attrs.encode())["error_code"] == "xml_structure_over_limit"


@pytest.mark.parametrize(
"feed",
[
{"feed_id": "x", "feed_url": "u", "kind": "rss", "accepted_row_count": 0, "rejected_row_count": 0, "state": "failed", "error_code": None},
{"feed_id": "x", "feed_url": "u", "kind": "rss", "accepted_row_count": 1, "rejected_row_count": 0, "state": "failed", "error_code": "bad"},
{"feed_id": "x", "feed_url": "u", "kind": "rss", "accepted_row_count": 0, "rejected_row_count": 0, "state": "stale", "error_code": ""},
{"feed_id": "x", "feed_url": "u", "kind": "rss", "accepted_row_count": 0, "rejected_row_count": 0, "state": "missing", "error_code": None},
{"feed_id": "x", "feed_url": "u", "kind": "rss", "accepted_row_count": 1, "rejected_row_count": 0, "state": "accepted", "error_code": "bad"},
{"feed_id": "x", "feed_url": "u", "kind": "rss", "accepted_row_count": 0, "rejected_row_count": 0, "state": "quarantined", "error_code": None},
],
)
def test_feed_state_row_and_error_invariants_are_strict(feed: dict[str, object]) -> None:
with pytest.raises(FetchAcceptanceError):
build_acceptance_status([feed])
Loading