Skip to content

Commit e28094c

Browse files
Pigbibicodex
andcommitted
Bind publication evidence to source CSV
Co-Authored-By: Codex <noreply@openai.com>
1 parent ddac894 commit e28094c

6 files changed

Lines changed: 113 additions & 16 deletions

File tree

.github/workflows/rss_source_pipeline.yml

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -52,7 +52,7 @@ jobs:
5252
- uses: actions/checkout@v6
5353
with:
5454
repository: QuantStrategyLab/PoliticalEventTrackingResearch
55-
ref: refs/heads/main
55+
ref: ${{ github.sha }}
5656
persist-credentials: false
5757
- name: Verify reviewed main checkout and canonical paths
5858
env:
@@ -76,7 +76,7 @@ jobs:
7676
- uses: actions/checkout@v6
7777
with:
7878
repository: QuantStrategyLab/PoliticalEventTrackingResearch
79-
ref: refs/heads/main
79+
ref: ${{ github.sha }}
8080
persist-credentials: false
8181
- name: Verify reviewed main checkout
8282
env:

scripts/validate_publish_input_evidence.py

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -2,15 +2,14 @@
22
from __future__ import annotations
33

44
import argparse
5-
import csv
65
from pathlib import Path
76
import sys
87

98
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src"))
109

11-
from political_event_tracking_research.feed_status_canonical_h2c import read_status # noqa: E402
1210
from political_event_tracking_research.publish_input_policy import ( # noqa: E402
1311
build_publication_evidence,
12+
recompute_source_items_binding,
1413
read_input_policy_evidence,
1514
read_publication_evidence,
1615
)
@@ -30,19 +29,17 @@ def main() -> None:
3029
evidence = read_publication_evidence(args.evidence.read_bytes())
3130
source_bytes = args.source_items.read_bytes()
3231
status_bytes = args.status.read_bytes()
33-
status = read_status(status_bytes)
34-
with args.source_items.open(newline="", encoding="utf-8") as handle:
35-
actual_row_count = sum(1 for _ in csv.DictReader(handle))
32+
actual_row_count, aggregate_digest, status_eligible = recompute_source_items_binding(source_bytes, status_bytes)
3633
if actual_row_count != evidence["source_items_row_count"]:
3734
raise ValueError("source_items_row_count_mismatch")
3835
actual = build_publication_evidence(
3936
policy,
4037
root=args.root,
4138
source_items_bytes=source_bytes,
4239
status_bytes=status_bytes,
43-
status_eligible=status["eligible_for_live_publication"],
40+
status_eligible=status_eligible,
4441
source_items_row_count=actual_row_count,
45-
aggregate_row_digest=status["aggregate_row_digest"],
42+
aggregate_row_digest=aggregate_digest,
4643
)
4744
if actual != evidence:
4845
raise ValueError("publication_evidence_mismatch")

scripts/write_publish_input_evidence.py

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -2,15 +2,14 @@
22
from __future__ import annotations
33

44
import argparse
5-
import csv
65
from pathlib import Path
76
import sys
87

98
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "src"))
109

11-
from political_event_tracking_research.feed_status_canonical_h2c import read_status # noqa: E402
1210
from political_event_tracking_research.publish_input_policy import ( # noqa: E402
1311
build_publication_evidence,
12+
recompute_source_items_binding,
1413
read_input_policy_evidence,
1514
read_publication_evidence,
1615
serialize_publication_evidence,
@@ -29,17 +28,15 @@ def main() -> None:
2928
policy = read_input_policy_evidence(args.policy.read_bytes())
3029
source_bytes = args.source_items.read_bytes()
3130
status_bytes = args.status.read_bytes()
32-
status = read_status(status_bytes)
33-
with args.source_items.open(newline="", encoding="utf-8") as handle:
34-
row_count = sum(1 for _ in csv.DictReader(handle))
31+
row_count, aggregate_digest, status_eligible = recompute_source_items_binding(source_bytes, status_bytes)
3532
evidence = build_publication_evidence(
3633
policy,
3734
root=args.root,
3835
source_items_bytes=source_bytes,
3936
status_bytes=status_bytes,
40-
status_eligible=status["eligible_for_live_publication"],
37+
status_eligible=status_eligible,
4138
source_items_row_count=row_count,
42-
aggregate_row_digest=status["aggregate_row_digest"],
39+
aggregate_row_digest=aggregate_digest,
4340
)
4441
if read_publication_evidence(evidence) != evidence:
4542
raise ValueError("publication_evidence_readback_invalid")

src/political_event_tracking_research/feed_status_canonical_h2c.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -102,6 +102,11 @@ def _digest(rows: list[dict[str, str]]) -> str:
102102
return hashlib.sha256(payload).hexdigest()
103103

104104

105+
def digest_rows(rows: list[dict[str, str]]) -> str:
106+
"""Return the H2C aggregate digest for an already validated row snapshot."""
107+
return _digest(rows)
108+
109+
105110
def _parse_outcome(value: object) -> tuple[dict[str, object], list[dict[str, str]]]:
106111
data = _mapping(value, _OUTCOME_KEYS, "outcome_invalid")
107112
kind = data["kind"]

src/political_event_tracking_research/publish_input_policy.py

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,16 @@
11
from __future__ import annotations
22

3+
import csv
34
import hashlib
5+
import io
46
import json
57
import re
68
import stat
79
from collections.abc import Mapping
810
from pathlib import Path
911

12+
from .feed_status_canonical_h2c import EMPTY_DIGEST, digest_rows, read_status
13+
1014

1115
PUBLISH_MAX_ITEMS_PER_FEED = 50
1216
MAX_DEBUG_ITEMS_PER_FEED = 10_000
@@ -19,6 +23,7 @@
1923
WORKFLOW_PATH = ".github/workflows/rss_source_pipeline.yml"
2024
WORKFLOW_REF = f"{REPOSITORY}/{WORKFLOW_PATH}@refs/heads/main"
2125
POLICY_VERSION = "pert.rss_publish_input.v1"
26+
SOURCE_ITEM_FIELDS = ("item_id", "published_at", "source_type", "source_url", "author", "text")
2227
_SHA256_RE = re.compile(r"^[0-9a-f]{64}$")
2328
_SHA1_RE = re.compile(r"^[0-9a-f]{40}$")
2429
_POLICY_KEYS = frozenset(
@@ -58,6 +63,40 @@ def _fail(code: str) -> None:
5863
raise PublishInputPolicyError(code)
5964

6065

66+
def recompute_source_items_binding(source_items_bytes: bytes, status_bytes: bytes) -> tuple[int, str, bool]:
67+
if type(source_items_bytes) is not bytes or type(status_bytes) is not bytes:
68+
_fail("publication_evidence_invalid")
69+
try:
70+
reader = csv.DictReader(io.StringIO(source_items_bytes.decode("utf-8"), newline=""))
71+
if tuple(reader.fieldnames or ()) != SOURCE_ITEM_FIELDS:
72+
_fail("source_items_schema_invalid")
73+
rows: list[dict[str, str]] = []
74+
for row in reader:
75+
if set(row) != set(SOURCE_ITEM_FIELDS) or any(type(value) is not str for value in row.values()):
76+
_fail("source_items_schema_invalid")
77+
rows.append({key: row[key] for key in SOURCE_ITEM_FIELDS})
78+
canonical = io.StringIO(newline="")
79+
writer = csv.DictWriter(canonical, fieldnames=SOURCE_ITEM_FIELDS, lineterminator="\n")
80+
writer.writeheader()
81+
writer.writerows(rows)
82+
if canonical.getvalue().encode("utf-8") != source_items_bytes:
83+
_fail("source_items_noncanonical")
84+
status = read_status(status_bytes)
85+
except (csv.Error, UnicodeError, ValueError):
86+
_fail("source_items_invalid")
87+
row_count = len(rows)
88+
aggregate_digest = digest_rows(rows) if rows else EMPTY_DIGEST
89+
derived_eligible = status["failed_feed_count"] == 0 and status["quarantined_feed_count"] == 0
90+
if (
91+
row_count != status["accepted_row_count"]
92+
or aggregate_digest != status["aggregate_row_digest"]
93+
or status["publication_complete"] is not derived_eligible
94+
or status["eligible_for_live_publication"] is not derived_eligible
95+
):
96+
_fail("source_items_status_mismatch")
97+
return row_count, aggregate_digest, derived_eligible
98+
99+
61100
def _mapping(value: object, keys: frozenset[str], code: str) -> dict[str, object]:
62101
if not isinstance(value, Mapping):
63102
_fail(code)

tests/test_publish_input_policy.py

Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,9 @@
1313
build_publication_evidence,
1414
read_input_policy_evidence,
1515
read_publication_evidence,
16+
recompute_source_items_binding,
1617
)
18+
from political_event_tracking_research.feed_status_canonical_h2c import build_decision
1719

1820

1921
def write_inputs(root: Path) -> None:
@@ -171,6 +173,62 @@ def test_debug_evidence_can_readback_without_becoming_publishable(tmp_path: Path
171173
assert evidence["status_eligible_for_live_publication"] is True
172174

173175

176+
def test_source_items_binding_recomputes_count_digest_and_publishability() -> None:
177+
row = {
178+
"item_id": "a-1",
179+
"published_at": "2026-05-01T12:30:00Z",
180+
"source_type": "official",
181+
"source_url": "https://example.test/a",
182+
"author": "",
183+
"text": "event",
184+
}
185+
status = build_decision(
186+
[
187+
{
188+
"feed_id": "feed-a",
189+
"feed_url": "https://example.test/feed-a",
190+
"kind": "rss2",
191+
"state": "accepted",
192+
"rows": [row],
193+
"error_code": None,
194+
}
195+
]
196+
)
197+
source = (
198+
"item_id,published_at,source_type,source_url,author,text\n"
199+
"a-1,2026-05-01T12:30:00Z,official,https://example.test/a,,event\n"
200+
).encode()
201+
count, digest, eligible = recompute_source_items_binding(source, status.status_bytes)
202+
assert (count, digest, eligible) == (1, json.loads(status.status_bytes)["aggregate_row_digest"], True)
203+
204+
205+
def test_source_items_binding_rejects_count_or_digest_drift() -> None:
206+
status = build_decision(
207+
[
208+
{
209+
"feed_id": "feed-a",
210+
"feed_url": "https://example.test/feed-a",
211+
"kind": "rss2",
212+
"state": "accepted",
213+
"rows": [
214+
{
215+
"item_id": "a-1",
216+
"published_at": "2026-05-01T12:30:00Z",
217+
"source_type": "official",
218+
"source_url": "https://example.test/a",
219+
"author": "",
220+
"text": "event",
221+
}
222+
],
223+
"error_code": None,
224+
}
225+
]
226+
)
227+
tampered = b"item_id,published_at,source_type,source_url,author,text\na-1,2026-05-01T12:30:00Z,official,https://example.test/a,,changed\n"
228+
with pytest.raises(PublishInputPolicyError, match="source_items_status_mismatch"):
229+
recompute_source_items_binding(tampered, status.status_bytes)
230+
231+
174232
def test_workflow_guard_and_evidence_precede_fetch_and_publish() -> None:
175233
workflow = Path(__file__).parents[1].joinpath(".github/workflows/rss_source_pipeline.yml").read_text(
176234
encoding="utf-8"
@@ -182,4 +240,5 @@ def test_workflow_guard_and_evidence_precede_fetch_and_publish() -> None:
182240
assert workflow.index("Upload RSS source artifact") < workflow.index("Build and validate publication evidence")
183241
assert workflow.index("Build and validate publication evidence") < workflow.index("Publish live CSV outputs")
184242
assert '--max-items-per-feed "${MAX_ITEMS_PER_FEED}"' in workflow
243+
assert workflow.count("ref: ${{ github.sha }}") == 2
185244
assert "git push origin HEAD:refs/heads/main" in workflow

0 commit comments

Comments
 (0)