@@ -187,14 +187,14 @@ def build_weekly_artifact(*, period_start: date, as_of: date, generated_at: date
187187 raise _invalid ("generated_at_invalid" )
188188 if source_provenance != "official_rss_source_pipeline_v1" or run_mode not in {"scheduled" , "manual" }:
189189 raise _invalid ("producer_contract_invalid" )
190- raw_digest = _snapshot_digest (source_events , watchlist )
191190 period_events = _filter_events (source_events , period_start , as_of )
191+ source_snapshot_digest = _snapshot_digest (period_events , watchlist )
192192 event_count , _ = _csv_snapshot (period_events , EVENT_HEADER , "events_csv_invalid" , allow_empty = True )
193193 watch_count , _ = _csv_snapshot (watchlist , WATCHLIST_HEADER , "watchlist_csv_invalid" )
194194 status = _status (feed_status )
195195 snapshot_id = f"rss_source_snapshot_{ as_of :%Y%m%d} _{ source_run_id } "
196196 try :
197- lock = PoliticalEventWeeklyPeriodLockV1 (period_start , period_end , as_of , workflow_ref , source_run_id , 1 , producer_ref , snapshot_id , raw_digest , source_provenance , (SourceSnapshotArtifact (EVENTS_NAME , _sha256 (period_events ), event_count ), SourceSnapshotArtifact (WATCHLIST_NAME , _sha256 (watchlist ), watch_count )))
197+ lock = PoliticalEventWeeklyPeriodLockV1 (period_start , period_end , as_of , workflow_ref , source_run_id , 1 , producer_ref , snapshot_id , source_snapshot_digest , source_provenance , (SourceSnapshotArtifact (EVENTS_NAME , _sha256 (period_events ), event_count ), SourceSnapshotArtifact (WATCHLIST_NAME , _sha256 (watchlist ), watch_count )))
198198 contract = WeeklySourceContract (as_of , period_start , period_end , generated_at , run_mode , producer_ref , source_provenance , (WeeklySourceArtifact (EVENTS_NAME , _sha256 (period_events ), event_count ), WeeklySourceArtifact (WATCHLIST_NAME , _sha256 (watchlist ), watch_count )), status )
199199 files : dict [str , bytes ] = {PERIOD_LOCK_NAME : serialize_period_lock (lock ), EVENTS_NAME : period_events , WATCHLIST_NAME : watchlist , WEEKLY_NAME : serialize_weekly_contract (contract )}
200200 except (PeriodLockError , WeeklyContractError , TypeError , ValueError , OverflowError ):
@@ -228,6 +228,8 @@ def parse_weekly_artifact(files: Mapping[str, bytes]) -> dict[str, bytes]:
228228 if tuple ((item .path , item .sha256 , item .row_count ) for item in lock .source_artifacts ) != expected or tuple ((item .path , item .sha256 , item .row_count ) for item in contract .source_artifacts ) != expected :
229229 raise _invalid ("source_artifact_mismatch" )
230230 expected_id = f"rss_source_snapshot_{ contract .as_of :%Y%m%d} _{ lock .source_run_id } "
231+ if lock .source_snapshot_digest != _snapshot_digest (files [EVENTS_NAME ], files [WATCHLIST_NAME ]):
232+ raise _invalid ("source_snapshot_digest_mismatch" )
231233 if lock .source_snapshot_id != expected_id or lock .source_attempt != 1 or lock .period_start != contract .period_start or lock .period_end_exclusive != contract .period_end_exclusive or lock .as_of != contract .as_of or lock .producer_ref != contract .producer_ref or lock .source_provenance != contract .source_provenance :
232234 raise _invalid ("period_contract_mismatch" )
233235 manifest_value = _parse_json (files [MANIFEST_NAME ], "manifest_wire_invalid" )
0 commit comments