| 1 | """Pipeline Reddit dispatch: free-first, SC thinness-floor backfill (U7), and |
| 2 | that the keyless path never calls search.json (U2).""" |
| 3 | |
| 4 | from unittest import mock |
| 5 | |
| 6 | from pathlib import Path |
| 7 | |
| 8 | from lib import env, health, http, pipeline, reddit_keyless, schema |
| 9 | |
| 10 | |
| 11 | def _subquery(): |
| 12 | return schema.SubQuery(label="t", search_query="kanye", ranking_query="kanye", |
| 13 | sources=["reddit"]) |
| 14 | |
| 15 | |
| 16 | def _runtime(): |
| 17 | return schema.ProviderRuntime(reasoning_provider="mock", planner_model="mock", |
| 18 | rerank_model="mock") |
| 19 | |
| 20 | |
| 21 | def _item(rid): |
| 22 | return {"url": f"https://www.reddit.com/r/x/comments/{rid}/t/", "title": rid} |
| 23 | |
| 24 | |
| 25 | def _ids(items): |
| 26 | return [pipeline._reddit_post_key(i) for i in items] |
| 27 | |
| 28 | |
| 29 | class TestThinnessFloor: |
| 30 | KEY = {"SCRAPECREATORS_API_KEY": "k"} |
| 31 | FLOOR = env.REDDIT_SC_MIN_ITEMS_VAR |
| 32 | |
| 33 | def _run(self, config, public, sc_parsed): |
| 34 | with mock.patch("lib.reddit_public.search_reddit_public", return_value=public), \ |
| 35 | mock.patch("lib.reddit.search_and_enrich", return_value={"raw": 1}) as sc, \ |
| 36 | mock.patch("lib.reddit.parse_reddit_response", return_value=sc_parsed): |
| 37 | items, _ = _stream(config) |
| 38 | return items, sc |
| 39 | |
| 40 | def test_unset_floor_backfills_thin_free_run_free_first(self): |
| 41 | # Default floor is 5: 3 free items < 5 -> one SC call, free first. |
| 42 | free = [_item("a"), _item("b"), _item("c")] |
| 43 | items, sc = self._run(self.KEY, free, [_item("b"), _item("z")]) |
| 44 | sc.assert_called_once() |
| 45 | assert _ids(items) == ["a", "b", "c", "z"] |
| 46 | |
| 47 | def test_unset_floor_six_free_items_spends_nothing(self): |
| 48 | free = [_item(c) for c in "abcdef"] |
| 49 | items, sc = self._run(self.KEY, free, [_item("z")]) |
| 50 | sc.assert_not_called() |
| 51 | assert len(items) == 6 |
| 52 | |
| 53 | def test_exactly_default_floor_is_acceptable(self): |
| 54 | free = [_item(c) for c in "abcde"] |
| 55 | items, sc = self._run(self.KEY, free, [_item("z")]) |
| 56 | sc.assert_not_called() |
| 57 | assert len(items) == 5 |
| 58 | |
| 59 | def test_empty_string_floor_behaves_as_unset(self): |
| 60 | cfg = {**self.KEY, self.FLOOR: ""} |
| 61 | items, sc = self._run(cfg, [_item("a"), _item("b"), _item("c")], [_item("z")]) |
| 62 | sc.assert_called_once() |
| 63 | assert _ids(items) == ["a", "b", "c", "z"] |
| 64 | |
| 65 | def test_explicit_zero_is_empty_only(self): |
| 66 | cfg = {**self.KEY, self.FLOOR: "0"} |
| 67 | items, sc = self._run(cfg, [_item("a"), _item("b"), _item("c")], [_item("z")]) |
| 68 | sc.assert_not_called() |
| 69 | assert len(items) == 3 |
| 70 | |
| 71 | def test_explicit_zero_still_backfills_empty_free_run(self): |
| 72 | cfg = {**self.KEY, self.FLOOR: "0"} |
| 73 | items, sc = self._run(cfg, [], [_item("z")]) |
| 74 | sc.assert_called_once() |
| 75 | assert _ids(items) == ["z"] |
| 76 | |
| 77 | def test_threshold_not_fired_when_free_above_floor(self): |
| 78 | cfg = {**self.KEY, self.FLOOR: "2"} |
| 79 | items, sc = self._run(cfg, [_item("a"), _item("b"), _item("c")], [_item("z")]) |
| 80 | sc.assert_not_called() |
| 81 | assert len(items) == 3 |
| 82 | |
| 83 | def test_no_key_never_calls_sc(self): |
| 84 | for floor in (None, "", "0", "5", "junk"): |
| 85 | cfg = {} if floor is None else {self.FLOOR: floor} |
| 86 | items, sc = self._run(cfg, [_item("a")], [_item("z")]) |
| 87 | sc.assert_not_called() |
| 88 | assert len(items) == 1 |
| 89 | |
| 90 | def test_malformed_floor_spends_nothing(self): |
| 91 | # Malformed means 0 (empty-only): 3 free items -> no SC call. |
| 92 | cfg = {**self.KEY, self.FLOOR: "not-an-int"} |
| 93 | items, sc = self._run(cfg, [_item("a"), _item("b"), _item("c")], [_item("z")]) |
| 94 | sc.assert_not_called() |
| 95 | assert len(items) == 3 |
| 96 | |
| 97 | |
| 98 | def _stream(config): |
| 99 | return pipeline._retrieve_stream( |
| 100 | topic="kanye", subquery=_subquery(), source="reddit", config=config, |
| 101 | depth="quick", date_range=("2026-05-26", "2026-06-25"), |
| 102 | runtime=_runtime(), mock=False, |
| 103 | ) |
| 104 | |
| 105 | |
| 106 | def _sc_402(*_args, **_kwargs): |
| 107 | # Mirrors the real transport: http records the failure in the run's |
| 108 | # capture sink, then raises. |
| 109 | http._raise(http.HTTPError("HTTP 402: Payment Required", status_code=402)) |
| 110 | |
| 111 | |
| 112 | class TestBackfillOutcome: |
| 113 | KEY = {"SCRAPECREATORS_API_KEY": "k"} |
| 114 | |
| 115 | def test_backfill_402_after_free_items_keeps_source_working(self): |
| 116 | free = [_item("a"), _item("b"), _item("c")] |
| 117 | with mock.patch("lib.reddit_public.search_reddit_public", return_value=free), \ |
| 118 | mock.patch("lib.reddit.search_and_enrich", side_effect=_sc_402): |
| 119 | items, artifact = _stream(self.KEY) |
| 120 | assert _ids(items) == ["a", "b", "c"] |
| 121 | assert not artifact.get("_source_outcome") # not branded failed |
| 122 | detail = artifact.get("_source_outcome_detail") or "" |
| 123 | assert "402" in detail |
| 124 | assert "ScrapeCreators backfill failed" in detail |
| 125 | assert artifact.get("_source_outcome_detail_state") == health.PAYMENT_REQUIRED |
| 126 | |
| 127 | def test_backfill_generic_failure_after_free_items_is_detail(self): |
| 128 | free = [_item("a")] |
| 129 | with mock.patch("lib.reddit_public.search_reddit_public", return_value=free), \ |
| 130 | mock.patch("lib.reddit.search_and_enrich", side_effect=Exception("down")): |
| 131 | items, artifact = _stream(self.KEY) |
| 132 | assert _ids(items) == ["a"] |
| 133 | assert not artifact.get("_source_outcome") |
| 134 | assert "down" in (artifact.get("_source_outcome_detail") or "") |
| 135 | |
| 136 | def test_backfill_failure_with_no_free_items_keeps_explicit_failure(self): |
| 137 | with mock.patch("lib.reddit_public.search_reddit_public", return_value=[]), \ |
| 138 | mock.patch("lib.reddit.search_and_enrich", side_effect=_sc_402): |
| 139 | items, artifact = _stream(self.KEY) |
| 140 | assert items == [] |
| 141 | outcome = artifact.get("_source_outcome") or {} |
| 142 | assert outcome.get("state") == health.PAYMENT_REQUIRED |
| 143 | |
| 144 | def test_backfill_note_and_swallowed_keyless_403_both_survive(self): |
| 145 | free = [_item("a"), _item("b"), _item("c")] |
| 146 | |
| 147 | def _public(*_args, **_kwargs): |
| 148 | # A keyless lane swallowed a 403 but the source still delivered. |
| 149 | http._record_failure(http.HTTPError("HTTP 403: Blocked", status_code=403)) |
| 150 | return free |
| 151 | |
| 152 | with mock.patch("lib.reddit_public.search_reddit_public", side_effect=_public), \ |
| 153 | mock.patch("lib.reddit.search_and_enrich", return_value={"raw": 1}), \ |
| 154 | mock.patch("lib.reddit.parse_reddit_response", return_value=[_item("z")]): |
| 155 | items, artifact = _stream(self.KEY) |
| 156 | assert _ids(items) == ["a", "b", "c", "z"] |
| 157 | assert not artifact.get("_source_outcome") |
| 158 | detail = artifact.get("_source_outcome_detail") or "" |
| 159 | assert "ScrapeCreators backfill ran" in detail |
| 160 | assert "1 sub-request blocked (HTTP 403)" in detail |
| 161 | |
| 162 | def test_backfill_that_ran_is_noted(self): |
| 163 | with mock.patch("lib.reddit_public.search_reddit_public", return_value=[_item("a")]), \ |
| 164 | mock.patch("lib.reddit.search_and_enrich", return_value={"raw": 1}), \ |
| 165 | mock.patch("lib.reddit.parse_reddit_response", |
| 166 | return_value=[_item("a"), _item("z")]): |
| 167 | items, artifact = _stream(self.KEY) |
| 168 | assert _ids(items) == ["a", "z"] |
| 169 | assert not artifact.get("_source_outcome") |
| 170 | detail = artifact.get("_source_outcome_detail") or "" |
| 171 | assert "ScrapeCreators backfill ran" in detail |
| 172 | assert "added 1" in detail |
| 173 | |
| 174 | def test_swallowed_sc_429_after_free_items_does_not_rate_limit_reddit(self): |
| 175 | # ScrapeCreators 429s inside the backfill are swallowed by lib.reddit |
| 176 | # (recorded in the capture sink, then []). reddit.com was never |
| 177 | # rate-limited, so the stream must not carry RATE_LIMITED: the fan-out |
| 178 | # would add 'reddit' to rate_limited_sources and skip later Reddit |
| 179 | # streams and the thin-source retry. |
| 180 | free = [_item("a"), _item("b"), _item("c")] |
| 181 | |
| 182 | def _sc_http_429(*_args, **_kwargs): |
| 183 | http._raise(http.HTTPError("HTTP 429: Too Many Requests", status_code=429)) |
| 184 | |
| 185 | with mock.patch("lib.reddit_public.search_reddit_public", return_value=free), \ |
| 186 | mock.patch("lib.http.get", side_effect=_sc_http_429) as sc_get: |
| 187 | items, artifact = _stream(self.KEY) |
| 188 | assert sc_get.called |
| 189 | assert _ids(items) == ["a", "b", "c"] |
| 190 | assert not artifact.get("_source_outcome") |
| 191 | assert artifact.get("_source_outcome_detail_state") != health.RATE_LIMITED |
| 192 | detail = artifact.get("_source_outcome_detail") or "" |
| 193 | assert "ScrapeCreators backfill" in detail |
| 194 | assert "HTTP 429" in detail |
| 195 | |
| 196 | def test_raised_sc_429_after_free_items_does_not_rate_limit_reddit(self): |
| 197 | free = [_item("a"), _item("b"), _item("c")] |
| 198 | |
| 199 | def _sc_raise_429(*_args, **_kwargs): |
| 200 | http._raise(http.HTTPError("HTTP 429: Too Many Requests", status_code=429)) |
| 201 | |
| 202 | with mock.patch("lib.reddit_public.search_reddit_public", return_value=free), \ |
| 203 | mock.patch("lib.reddit.search_and_enrich", side_effect=_sc_raise_429): |
| 204 | items, artifact = _stream(self.KEY) |
| 205 | assert _ids(items) == ["a", "b", "c"] |
| 206 | assert not artifact.get("_source_outcome") |
| 207 | assert artifact.get("_source_outcome_detail_state") != health.RATE_LIMITED |
| 208 | detail = artifact.get("_source_outcome_detail") or "" |
| 209 | assert "ScrapeCreators backfill failed" in detail |
| 210 | assert "429" in detail |
| 211 | |
| 212 | def test_sc_429_with_no_free_items_keeps_explicit_rate_limited_outcome(self): |
| 213 | def _sc_raise_429(*_args, **_kwargs): |
| 214 | http._raise(http.HTTPError("HTTP 429: Too Many Requests", status_code=429)) |
| 215 | |
| 216 | with mock.patch("lib.reddit_public.search_reddit_public", return_value=[]), \ |
| 217 | mock.patch("lib.reddit.search_and_enrich", side_effect=_sc_raise_429): |
| 218 | items, artifact = _stream(self.KEY) |
| 219 | assert items == [] |
| 220 | outcome = artifact.get("_source_outcome") or {} |
| 221 | assert outcome.get("state") == health.RATE_LIMITED |
| 222 | |
| 223 | def test_no_backfill_no_note(self): |
| 224 | free = [_item(c) for c in "abcdef"] |
| 225 | with mock.patch("lib.reddit_public.search_reddit_public", return_value=free), \ |
| 226 | mock.patch("lib.reddit.search_and_enrich") as sc: |
| 227 | _items, artifact = _stream(self.KEY) |
| 228 | sc.assert_not_called() |
| 229 | assert not artifact.get("_source_outcome_detail") |
| 230 | |
| 231 | |
| 232 | class TestMergeHelper: |
| 233 | def test_dedup_by_post_id_free_first(self): |
| 234 | out = pipeline._merge_reddit_items([_item("a"), _item("b")], [_item("b"), _item("c")]) |
| 235 | assert _ids(out) == ["a", "b", "c"] |
| 236 | |
| 237 | |
| 238 | class TestNoSearchJson: |
| 239 | def test_keyless_discovery_never_calls_searchjson(self): |
| 240 | # reddit_public.search (the .json caller) must never run in the keyless flow. |
| 241 | with mock.patch("lib.reddit_public.search") as json_search, \ |
| 242 | mock.patch("lib.reddit_keyless.reddit_search.search", return_value=[]), \ |
| 243 | mock.patch("lib.reddit_keyless.reddit_listing.fetch_listings", return_value=[]): |
| 244 | reddit_keyless._discover("topic", "default", ["test"]) |
| 245 | json_search.assert_not_called() |
| 246 | |
| 247 | |
| 248 | class TestEnvConstantParity: |
| 249 | """F2 regression (restate-as-mirror drift): pipeline's Reddit gating must |
| 250 | key off env's declared constants (env.REDDIT_BACKEND_PIN_VAR / |
| 251 | env.REDDIT_SC_MIN_ITEMS_VAR) — never restated raw strings that can drift |
| 252 | from the single source of truth in lib/env.py.""" |
| 253 | |
| 254 | def _run(self, config, public, sc_parsed): |
| 255 | with mock.patch("lib.reddit_public.search_reddit_public", |
| 256 | return_value=public) as pub, \ |
| 257 | mock.patch("lib.reddit.search_and_enrich", return_value={"raw": 1}) as sc, \ |
| 258 | mock.patch("lib.reddit.parse_reddit_response", return_value=sc_parsed): |
| 259 | items, _ = pipeline._retrieve_stream( |
| 260 | topic="kanye", subquery=_subquery(), source="reddit", config=config, |
| 261 | depth="quick", date_range=("2026-05-26", "2026-06-25"), |
| 262 | runtime=_runtime(), mock=False, |
| 263 | ) |
| 264 | return items, pub, sc |
| 265 | |
| 266 | def test_pipeline_source_has_no_raw_reddit_env_literals(self): |
| 267 | # The declared constants live in env.py; pipeline.py must not restate |
| 268 | # the raw LAST30DAYS_REDDIT_* strings (comments included — they drift too). |
| 269 | source = Path(pipeline.__file__).read_text() |
| 270 | assert "LAST30DAYS_REDDIT_" not in source |
| 271 | |
| 272 | def test_backend_pin_constant_flips_gating_to_sc_primary(self): |
| 273 | # Keyed via the env constant, not a raw string: pin=scrapecreators |
| 274 | # makes SC primary and skips the free path entirely. |
| 275 | cfg = {"SCRAPECREATORS_API_KEY": "k", env.REDDIT_BACKEND_PIN_VAR: "scrapecreators"} |
| 276 | items, pub, sc = self._run(cfg, [_item("a")], [_item("z")]) |
| 277 | sc.assert_called_once() |
| 278 | pub.assert_not_called() |
| 279 | assert _ids(items) == ["z"] |
| 280 | |
| 281 | def test_min_items_constant_drives_thinness_backfill(self): |
| 282 | # Keyed via the env constant: floor of 5 vs 1 free item -> SC backfill. |
| 283 | cfg = {"SCRAPECREATORS_API_KEY": "k", env.REDDIT_SC_MIN_ITEMS_VAR: "5"} |
| 284 | items, pub, sc = self._run(cfg, [_item("a")], [_item("z")]) |
| 285 | sc.assert_called_once() |
| 286 | assert _ids(items) == ["a", "z"] |
| 287 |