diff --git a/.env.example b/.env.example index 81616f0..3feebb5 100644 --- a/.env.example +++ b/.env.example @@ -5,6 +5,13 @@ KAP_RATE_LIMIT_RPS=2 KAP_TIMEOUT_S=10 KAP_USER_AGENT=trailingedge/0.1 (research) +# Optional proxy pool for the KAP backfill. The WAF throttles per source IP (~50 requests, +# then a block that refills in ~2 min), so a rotating pool of IPs removes the block stall: +# a spent IP is parked to refill while another serves at full speed. Comma-separated URLs, +# or put one per line in a gitignored proxies.txt. Residential proxies survive far better +# than datacenter ones here. Leave unset for a single direct connection (default). +# KAP_PROXIES=http://user:pass@host1:port,http://user:pass@host2:port + # TSG (Ticaret Sicil Gazetesi) - semi-automatic scraping. # The scraper opens a visible browser; you sign in and solve the login # CAPTCHA by hand once, then the run proceeds automatically. No credentials diff --git a/.gitignore b/.gitignore index 7af0a2e..13a4307 100644 --- a/.gitignore +++ b/.gitignore @@ -65,3 +65,6 @@ ehthumbs.db .python-version pip-log.txt pip-delete-this-directory.txt + +# Proxy list for the KAP scraper rotation pool - a credential, never commit. +proxies.txt* diff --git a/README.md b/README.md index c8a33e4..9debcaa 100644 --- a/README.md +++ b/README.md @@ -21,19 +21,33 @@ disclosure is *public*, returns are measured in excess of XU100 over the same he interval, and the round-trip cost is estimated per trade from that stock's own OHLC (Abdi-Ranaldo 2017) rather than assumed as a flat fee. +Full 2015-2026 history, **2,279 survivorship-clean insider-cluster events**: + | Horizon | N | Gross AR | Cost | **Net AR** | t (net) | |---|---:|---:|---:|---:|---:| -| 5d | 1,070 | +0.66% | 3.37% | **−2.71%** | −12.91 | -| 20d | 1,070 | +2.02% | 3.37% | **−1.35%** | −3.31 | -| 60d | 1,071 | +2.18% | 3.37% | **−1.20%** | −1.86 | - -**The signal is real. The gross abnormal return is significantly positive at every -horizon** (20d: +2.07%, t = 5.36, N = 1,079). **And it is not tradeable**, because insider -clusters fire in illiquid small caps whose bid-ask spread is wider than the alpha: the -median round trip costs 1.93%, the upper quartile 4.34%. Nothing survives crossing it -twice. At 60 days the net loss is no longer statistically distinguishable from zero -(t = −1.86) - which buys nothing: the point estimate is still negative, and "you might -merely break even after three months" is not an edge either. +| 5d | 2,279 | +0.41% | 4.19% | **−3.78%** | −21.29 | +| 20d | 2,279 | +1.58% | 4.19% | **−2.61%** | −7.76 | +| 60d | 2,236 | +2.61% | 4.19% | **−1.58%** | −2.57 | + +**The pooled gross signal is real** (20d: +1.58%, t = 5.2, `EDGE_DETECTED`) **and not +tradeable** — insider clusters fire in illiquid small caps whose bid-ask spread (median +round trip 2.33%) is wider than the alpha. Net negative at every horizon. + +**But the full history says something sharper than "not tradeable".** Split by regime, the +gross signal was strong in 2015-2018 (+2.36% at 20d, t = 6.5) and has **decayed to nothing +in 2021-2026** (−0.67%, t = −1.0 — indistinguishable from zero, before costs): + +| Regime | 20d Gross AR | t (gross) | 20d Net AR | +|---|---:|---:|---:| +| 2015-2018 | +2.36% | 6.46 | −1.00% | +| 2019-2020 | +3.34% | 3.34 | −5.01% | +| **2021-2026** | **−0.67%** | **−1.00** | −5.44% | + +So the pooled number is carried entirely by the early era. Insiders' disclosed purchases +predicted abnormal returns in 2015-2018 — returns you still could not capture after the +spread — and in the regime that matters to a trader today they no longer predict them at +all. The edge was real, uncapturable, and has since decayed. (2019-2020 is a COVID +small-cap-mania artefact on a tiny, extremely illiquid sample, not a strategy.) That is the whole finding, and it is why this repository exists. A gross number is not an edge; an edge is what is left after the market takes its cut. @@ -95,12 +109,12 @@ verdict standing. `INSUFFICIENT_POWER` below ~784 events and `SURVIVORSHIP_BIASED` when too many clusters cannot be priced. Both gates fired during this work, and both were right. -> **What is claimed, precisely:** a statistically strong *gross* abnormal return -> (20d: +2.07%, t = 5.36, N = 1,079, survivorship-clean) that does **not** survive a -> per-trade cost estimate. The window is 2015-2018 - a single regime - so the result is -> not yet regime-conditional, and that is stated rather than glossed. Remaining gaps are -> in [`docs/METHODOLOGY.md`](docs/METHODOLOGY.md#6-still-open), not left for a -> reader to discover. +> **What is claimed, precisely:** across the full 2015-2026 history (N = 2,279, +> survivorship-clean), a *gross* abnormal return that does **not** survive a per-trade cost +> estimate at any horizon — and that, split by regime, was statistically strong in 2015-2018 +> (20d +2.36%, t = 6.5) and has **decayed to zero in 2021-2026** (−0.67%, t = −1.0). The edge +> was real, uncapturable, and is now gone. Remaining gaps are in +> [`docs/METHODOLOGY.md`](docs/METHODOLOGY.md), not left for a reader to discover. ## Türkçe özet @@ -173,13 +187,13 @@ python scripts/net_of_cost.py # the one that decides it ``` === Abnormal return, NET of round-trip cost (order 25,000 TRY) === spread: Abdi-Ranaldo (2017) from the stock's own OHLC, per trade - dropped (no cost estimate): 27 - round-trip cost: median 1.93% p25 1.19% p75 4.34% + dropped (no cost estimate): 36 + round-trip cost: median 2.33% p25 1.34% p75 4.76% HORIZON N GROSS AR% COST% NET AR% HIT% 95% CI t VERDICT - 5d 1070 0.66 3.37 -2.71 26.9 [24.3, 29.7] -12.91 LOSES MONEY (net) - 20d 1070 2.02 3.37 -1.35 41.3 [38.4, 44.3] -3.31 LOSES MONEY (net) - 60d 1071 2.18 3.37 -1.20 44.4 [41.5, 47.4] -1.86 NO EDGE (net) + 5d 2280 0.41 4.19 -3.78 26.1 [24.4, 28.0] -21.29 LOSES MONEY (net) + 20d 2279 1.58 4.19 -2.61 39.8 [37.8, 41.8] -7.76 LOSES MONEY (net) + 60d 2236 2.61 4.19 -1.58 42.4 [40.4, 44.5] -2.57 LOSES MONEY (net) ``` The spread is not a parameter. It is estimated for each trade from the 30 sessions of diff --git a/docs/METHODOLOGY.md b/docs/METHODOLOGY.md index 3cff0a6..46a7328 100644 --- a/docs/METHODOLOGY.md +++ b/docs/METHODOLOGY.md @@ -168,18 +168,22 @@ have flattered the answer, and not taken from a quote feed, which does not exist delisted names the exchange bulletin carries. Impact is Kyle/Almgren square-root on the same window; commission and BSMV are charged per side. - round-trip cost: median 1.93% p25 1.19% p75 4.34% +On the full 2015-2026 sample (N=2,279 priceable clusters, more than double the earlier cut): - horizon N=1070 gross AR net AR t (net) - 5d +0.66% -2.71% -12.91 - 20d +2.02% -1.35% -3.31 - 60d +2.18% -1.20% -1.86 (not significant) + round-trip cost: median 2.33% mean 4.19% p75 4.76% + + horizon N=2279 gross AR net AR t (net) + 5d +0.41% -3.78% -21.29 + 20d +1.58% -2.61% -7.76 + 60d +2.61% -1.58% -2.57 **The signal does not survive the cost of trading it.** Insider clusters fire in illiquid small caps, and the spread on those names is wider than the alpha. This is the project's -result, not a caveat on it. At 60 days the net loss stops being statistically -distinguishable from zero, which is not a reprieve: the point estimate is still negative, -and a signal that *may* break even over three months is not an edge either. +result, not a caveat on it - and on the full sample it is sharper than on the earlier cut: +the net loss is significant at every horizon, including 60 days. The pooled gross number +(+1.58%, t=5.2, still "EDGE_DETECTED" before costs) is nonetheless misleading on its own, +because it averages two different regimes - see §6, where the split shows the gross signal +was real in 2015-2018 and has since decayed to zero. A second trap, found later and more serious than the first, because it sat under the number that decides the answer. `price_history.close_try` is a **chained total-return @@ -264,25 +268,86 @@ Worth recording without over-reading: the *gross* alpha collapses from +2.30% to between the halves. That could be the market becoming more efficient, or it could be sampling noise at N = 197. It is not interpreted here, because at that N it cannot be. -## 6. Still open - -**Regime.** The window is 2015-2018. That spans the August 2018 currency crisis but not -the 2021-2023 negative-real-rate retail boom or the 2023+ normalisation. The result is -therefore **not regime-conditional**, and a signal that dies to the spread in one regime -could in principle survive in another where those names traded tighter. The honest position -is that this is untested, not that it is unaffected. - -Closing it needs the KAP backfill to reach 2026. That is a data-collection problem, not a -methodological one, and it does not touch the mechanism: the spread eating the alpha is a -microstructure fact about illiquid names, not a regime phenomenon. +## 6. Closed by the full backfill, and what still remains + +### Regime - CLOSED, and the signal itself decayed + +The window used to be 2015-2018. The KAP backfill now reaches the present (2026-07), so the +question - does the result hold across regimes? - can finally be answered on the full +2,279-cluster sample instead of asserted. It does not hold. Split by the cluster's own date: + + regime 20d gross t(gross) 20d net note + 2015-2018 +2.36% 6.46 -1.00% real, uncapturable (as before) + 2019-2020 +3.34% 3.34 -5.01% COVID small-cap mania, N=198, cost 8.4% + 2021-2026 -0.67% -1.00 -5.44% the GROSS signal is gone + +The finding is sharper than "not tradeable". In 2015-2018 the gross abnormal return was +real and strong (+2.36% at 20d, t=6.5) and died only to the spread. In **2021-2026 the gross +signal has decayed to nothing** (-0.67%, t=-1.0 - statistically indistinguishable from zero, +and if anything negative). Insider-cluster purchases no longer predict a positive abnormal +return at all in the recent regime, before costs are even considered. + +So the pooled full-sample number (+1.61% gross at 20d, t=5.2, EDGE_DETECTED) is carried +entirely by the 2015-2018 era and is misleading on its own - it averages a real early edge +with a dead recent one. The honest statement is: **the signal was a 2015-2018 phenomenon +that has since decayed**, plausibly because the 2021+ retail boom changed who follows insider +filings and how fast. The 2019-2020 line is a genuine but regime-specific artefact: the COVID +small-cap mania produced huge gross returns (+16.7% at 60d) on a tiny, extremely illiquid +sample (round-trip cost 8.4%), and it clears cost only there and only at 60 days - not a +strategy, a curiosity. + +This closes the last major open item. It does not rescue the signal - it removes even the +"real but uncapturable" consolation for the regime that matters most to a trader today. + +### The ingest - CLOSED at 139/139, and the completion metric was not the obvious one + +The full 2015-2026 archive is now in: **139 of 139 months SUCCESS, 18,415 disclosures, +15,316 transactions, 2,339 clusters** (up from the 1,079 the earlier analysis ran on). Two +things had to be got right to get here. + +First, the completion metric. A chunk can lose filings to a WAF disconnect and still log +`chunk_done`; the ledger records the month `PARTIAL`, not `SUCCESS`. Counting `chunk_done` +overstated progress badly - at one point it read 92 of 139 months while `scraper_runs` held +20 SUCCESS, 82 PARTIAL, 44 FAILED, and 2016 and 2020 contained not one complete month. The +real criterion is `count(SUCCESS)`, and reaching it needed repeated sweeps, because +`backfill_kap_insider.py` fixes its `todo` list at startup and only collects a month's +WAF-dropped stragglers on a *later* run. A supervisor now re-runs until SUCCESS stops rising. + +Second, the WAF itself, which made the difference between a 12-hour grind and a stall. It +throttles per source IP on cumulative volume (~50 disclosures, refilling in ~2 min), so the +fix was to pace under that budget and then to spread the load across a rotating pool of IPs +with concurrent requests. Throughput went from ~7 to ~100 disclosures/min with the block +rate driven to near zero. Direction of any residual loss: WAF drops are unrelated to what a +filing predicts, so they thin the sample rather than tilt it. + +### Quarantine - a source ceiling, not a bug + +On the complete archive, **3,795 of 18,415 disclosures (~21%) are quarantined** to +`reports/parse_failures/` - a DKB PDF that yields no transactions is set aside with its +bytes, never silently skipped. Sampling those files, **about three quarters carry no +extractable text at all**: they are scanned images with no text layer, and no parser will +ever read them. That is a hard ceiling of the same kind as the 250,000 TRY disclosure +threshold - a property of the source, not a defect to fix. Recovering them would need OCR, +which is not attempted here and would need its own accuracy audit before any number it +produced could be trusted. + +The remaining quarter do carry text and fail for a different reason: they are **narrative +filings with no transaction table** ("... 3,63 - 3,64 TL fiyat aralığından 9.911 adet satış +işlemi gerçekleşmiştir"), so the arithmetic row-validation gate has nothing to bind to. +Those are recoverable in principle. `scripts/reparse_quarantine.py` re-runs the parser over +the whole quarantine and rescues the ones that now parse; the scanned images are the +residual floor. + +### What still remains **Cluster scoring is close to single-factor.** `cluster_score` blends insider count (0.50), role seniority (0.30) and recency (0.20). In historical mode recency is pinned at 1.0, so 20% of the weight is a constant, and seniority falls back to its 0.5 default wherever the -scraped board roster does not cover an insider. When coverage is zero the score reduces to -a monotone function of `insider_count` alone - `detect_clusters` now logs `role_map_empty` -loudly in that case, where it used to happen silently. The score is not used to gate any -result reported here, so this is a latent defect rather than an active one. +scraped board roster does not cover an insider - and `person_company_roles` is unpopulated +until `graph scrape-management` runs, so in practice the score reduces to a monotone +function of `insider_count` alone. `detect_clusters` logs `role_map_empty` loudly in that +case. The score gates none of the results above, so this is a latent defect, not an active +one. **Kyle's lambda is uncalibrated** (1.0). At retail order size the impact term is small enough that the error changes no conclusion; at institutional size it would, and the number diff --git a/scripts/backfill_kap_insider.py b/scripts/backfill_kap_insider.py index 3b9212d..018e70b 100644 --- a/scripts/backfill_kap_insider.py +++ b/scripts/backfill_kap_insider.py @@ -71,10 +71,28 @@ # nothing had gone wrong: over 122 months that is ~2 hours of the run spent asleep on the # happy path. That is not WAF protection, it is a tax on success. # -# The escalation ladder itself is kept intact - if KAP's WAF does disconnect the warmup -# GET (httpx.RemoteProtocolError = IP throttled), back off 1 min, then 10, then 20. The -# protection is now reactive, which is the only thing a backoff can usefully be. -_WAF_BACKOFF_S = [60, 600, 1200] +# The ladder used to escalate 60s -> 600s -> 1200s on the assumption that a longer silence +# buys a bigger allowance from the WAF. It was never measured. Measured now, over 66 +# windows in the run's own log - counting how many disclosures got through between one +# block and the next, bucketed by how long we had just slept: +# +# silence windows disclosures passed before the next block +# < 2 min 36 42.9 +# 2-8 min 11 57.0 +# 10 min 7 57.1 +# 20 min 12 49.6 +# +# It is flat. A 20-minute sleep buys the same ~50 disclosures as a 90-second one. The WAF's +# budget is roughly "~50 requests, then block", and it refills within about two minutes - +# so nearly all of the 10- and 20-minute sleeps were pure waste. Over the run that is many +# hours spent waiting for an allowance we already had. +# +# It is not shortened to zero. In the two observed runs of 3+ consecutive short sleeps, the +# 4th and 5th windows fell to 30 and 20 disclosures (n=1 each - thin, but pointing the +# wrong way), so hammering without pause may still draw a penalty. The ladder keeps its +# shape and loses its tail: escalate, but within the window where the budget is known to +# refill. +_WAF_BACKOFF_S = [60, 120, 240] # A courtesy gap between chunks so a long backfill does not arrive as one unbroken # burst. Small enough to be irrelevant to the runtime (122 x 2s = ~4 min), unlike the @@ -254,7 +272,28 @@ async def main_async( _log.info("chunk_dry_run", from_date=from_date, to_date=to_date) continue - await _run_chunk(from_date, to_date) + # A month that exhausts its WAF backoffs must NOT kill the whole run. Its + # scraper_runs record is already written (PARTIAL, or FAILED if the warmup never + # got through), so a later sweep collects it - that is the entire point of the + # ledger. Before the ladder was shortened this rarely fired, because a 20-minute + # backoff usually outlasted the block; now that the tail is gone a genuinely + # blocked month exhausts faster, and letting its RuntimeError propagate abandoned + # every remaining month behind it. Log it and move on. + try: + await _run_chunk(from_date, to_date) + except Exception as exc: + # ANY unhandled per-chunk failure is contained, not just the RuntimeError from + # exhausted WAF backoffs. A dead proxy returning 402 Payment Required crashed a + # whole run once because it surfaced as httpx.ProxyError, which the narrower + # `except RuntimeError` did not catch. The month's ledger record is written by + # scraper.run() before it re-raises, so a sweep collects it; one bad month must + # never abandon the months behind it, whatever killed it. + _log.warning( + "chunk_abandoned", + from_date=from_date, + to_date=to_date, + error=f"{type(exc).__name__}: {exc}", + ) await asyncio.sleep(_CHUNK_GAP_S) if dry_run: diff --git a/scripts/backfill_supervisor.sh b/scripts/backfill_supervisor.sh new file mode 100644 index 0000000..bf06029 --- /dev/null +++ b/scripts/backfill_supervisor.sh @@ -0,0 +1,82 @@ +#!/usr/bin/env bash +# Supervises the KAP backfill: runs it, and when it exits (normal finish OR a silent death +# like the ones that kept happening mid-backoff), re-runs it. Two reasons a re-run is needed +# even when nothing crashed: +# 1. the process has died silently more than once under nohup; auto-restart removes the +# human from that loop. +# 2. backfill_kap_insider.py computes its todo list once at startup, so months that go +# PARTIAL *during* a run are only swept by the NEXT run. Re-running is how SUCCESS +# climbs to 139. +# It stops when SUCCESS stops rising across a full pass (the sweep has converged), or when +# SUCCESS reaches the month count. Everything goes to backfill.log; the supervisor's own +# lines are tagged SUPERVISOR. +set -u +cd "$(dirname "$0")/.." + +LOG=backfill.log +TOTAL_MONTHS=139 + +success_count() { + uv run python - <<'PY' 2>/dev/null | grep -oE '^[0-9]+$' | tail -1 +import asyncio, sys +sys.path.insert(0, "src") +from sqlalchemy import text +from trailing_edge.core.db import get_session, init_db +async def m(): + await init_db() + async with get_session() as s: + n = (await s.execute(text( + "SELECT count(DISTINCT metadata->>'from_date') FROM scraper_runs WHERE status='SUCCESS'" + ))).scalar() + print(int(n or 0)) +asyncio.run(m()) +PY +} + +# Total stored disclosures. This, not the SUCCESS month count, is the progress signal: during +# the forward scan SUCCESS barely moves (fresh months land as PARTIAL), but the disclosure +# count climbs every pass. It only plateaus once the frontier has reached the present AND the +# sweeps stop recovering deferred filings - which is the true convergence. +disclosure_count() { + uv run python - <<'PY' 2>/dev/null | grep -oE '^[0-9]+$' | tail -1 +import asyncio, sys +sys.path.insert(0, "src") +from sqlalchemy import text +from trailing_edge.core.db import get_session, init_db +async def m(): + await init_db() + async with get_session() as s: + n = (await s.execute(text("SELECT count(*) FROM kap_disclosures"))).scalar() + print(int(n or 0)) +asyncio.run(m()) +PY +} + +prev_succ=-1 +prev_disc=-1 +pass_num=0 +while true; do + pass_num=$((pass_num + 1)) + succ=$(success_count) + disc=$(disclosure_count) + echo "{\"event\":\"SUPERVISOR\",\"pass\":$pass_num,\"success\":$succ,\"disclosures\":$disc}" >> "$LOG" + + if [ "$succ" -ge "$TOTAL_MONTHS" ]; then + echo "{\"event\":\"SUPERVISOR\",\"done\":true,\"reason\":\"success>=months\",\"success\":$succ}" >> "$LOG" + break + fi + # Converged only when a WHOLE pass moved NEITHER the SUCCESS month count NOR the disclosure + # count. SUCCESS alone is the wrong signal - it is flat all through the forward scan while + # disclosures pour in. Requiring both to stall means the frontier has reached the present + # and the sweeps have stopped recovering anything. + if [ "$pass_num" -gt 1 ] && [ "$succ" -le "$prev_succ" ] && [ "$disc" -le "$prev_disc" ]; then + echo "{\"event\":\"SUPERVISOR\",\"done\":true,\"reason\":\"plateau\",\"success\":$succ,\"disclosures\":$disc}" >> "$LOG" + break + fi + prev_succ=$succ + prev_disc=$disc + + uv run python scripts/backfill_kap_insider.py --from 2015-01-01 >> "$LOG" 2>&1 + echo "{\"event\":\"SUPERVISOR\",\"run_exited\":true,\"pass\":$pass_num}" >> "$LOG" + sleep 5 +done diff --git a/scripts/regime_split.py b/scripts/regime_split.py new file mode 100644 index 0000000..a8b6e8f --- /dev/null +++ b/scripts/regime_split.py @@ -0,0 +1,134 @@ +"""Regime split of the insider-cluster signal: gross and net abnormal return by era. + +The question the full 2015-2026 backfill was collected to answer. Every priceable cluster is +bucketed by the era of its own date and reported gross/net at 5/20/60 days, so the pooled +headline can be seen for what it is - an average of a real 2015-2018 edge and a decayed +2021-2026 one. + + uv run python scripts/regime_split.py +""" +from __future__ import annotations + +import asyncio +import math +import statistics +import sys +from decimal import Decimal + +sys.path.insert(0, "src") + +from sqlalchemy import text # noqa: E402 + +from trailing_edge.core.db import get_session, init_db # noqa: E402 +from trailing_edge.signals.costs import round_trip_cost # noqa: E402 + +ORDER = Decimal("25000") +LOOK = 30 + + +def regime(year: int) -> str: + if year <= 2018: + return "2015-2018" + if year <= 2020: + return "2019-2020" + return "2021-2026" + + +def _t(values: list[float]) -> tuple[float, float]: + mean = statistics.fmean(values) + sd = statistics.stdev(values) if len(values) > 1 else 0.0 + t = mean / (sd / math.sqrt(len(values))) if sd > 0 else 0.0 + return mean, t + + +async def main_async() -> None: + await init_db() + async with get_session() as s: + rows = ( + await s.execute( + text( + """ + SELECT o.horizon_days, o.abnormal_return_pct, o.entry_date, c.ticker, + c.window_end + FROM signal_outcomes o + JOIN insider_clusters c ON c.id = o.cluster_id + WHERE o.abnormal_return_pct IS NOT NULL AND o.entry_date IS NOT NULL + """ + ) + ) + ).all() + + cache: dict[tuple, float | None] = {} + + async def cost(ticker: str, entry) -> float | None: + key = (ticker, entry) + if key not in cache: + px = ( + await s.execute( + text( + """ + SELECT close_try, high_try, low_try, volume, raw_close_try + FROM price_history + WHERE ticker = :t AND price_date < :d + ORDER BY price_date DESC LIMIT :n + """ + ), + {"t": ticker, "d": entry, "n": LOOK}, + ) + ).all() + if len(px) < 22 or any(r[4] is None for r in px): + cache[key] = None + else: + px = list(reversed(px)) + closes = [r[0] for r in px] + highs = [r[1] or r[0] for r in px] + lows = [r[2] or r[0] for r in px] + raws = [r[4] for r in px] + adv = Decimal( + str( + statistics.fmean( + float(rc) * float(r[3] or 0) + for rc, r in zip(raws, px, strict=True) + ) + ) + ) + rt = round_trip_cost( + closes, highs, lows, ORDER, adv, last_traded_price=raws[-1] + ) + cache[key] = float(rt.total_pct) if rt else None + return cache[key] + + buckets: dict[tuple[str, int], list[tuple[float, float]]] = {} + for hz, ar, entry, ticker, window_end in rows: + c = await cost(ticker, entry) + if c is None: + continue + buckets.setdefault((regime(window_end.year), hz), []).append( + (float(ar), float(ar) - c) + ) + + hdr = f"{'REGIME':<12}{'HZN':>4}{'N':>6}{'GROSS%':>8}{'COST%':>9}{'NET%':>8}{'t_net':>7}{'t_gross':>8}" + print(hdr) + for reg in ("2015-2018", "2019-2020", "2021-2026"): + for hz in (5, 20, 60): + grp = buckets.get((reg, hz), []) + if len(grp) < 20: + continue + gross = [x[0] for x in grp] + net = [x[1] for x in grp] + costs = [x[0] - x[1] for x in grp] + mg, tg = _t(gross) + mn, tn = _t(net) + print( + f"{reg:<12}{hz:>3}d{len(grp):>6}{mg:>8.2f}{statistics.fmean(costs):>9.2f}" + f"{mn:>8.2f}{tn:>7.2f}{tg:>8.2f}" + ) + print() + + +def main() -> None: + asyncio.run(main_async()) + + +if __name__ == "__main__": + main() diff --git a/src/trailing_edge/core/http.py b/src/trailing_edge/core/http.py index be6fc5d..2e1bdd9 100644 --- a/src/trailing_edge/core/http.py +++ b/src/trailing_edge/core/http.py @@ -1,5 +1,18 @@ -"""Rate-limited, retry-capable httpx async client wrapper.""" +"""Rate-limited, retry-capable httpx async client wrapper. + +Optionally rotates over a pool of proxies. The KAP WAF throttles per source IP: a block +arrives as httpx.RemoteProtocolError on the connection, and the budget (~50 requests) +refills with wall-clock time, not with a slower request rate - measured over 66 blocks in +one backfill, a 20-minute pause bought the same ~50 requests as a 90-second one. That is +exactly the limit a rotating IP pool defeats: when one IP's budget is spent, set it aside +to refill and carry on through another. With no proxy configured the client behaves exactly +as before - one direct connection, RemoteProtocolError surfaced immediately for the scraper +to defer. +""" import asyncio +import os +import time +from pathlib import Path from types import TracebackType from typing import Any @@ -17,6 +30,58 @@ _log = get_logger(__name__) +# How long a proxy's IP budget takes to refill after a WAF block. Measured: the budget was +# back within ~2 minutes regardless of how much longer we waited. A blocked IP is parked for +# this long before the rotation returns to it. +_IP_COOLDOWN_S = 120.0 + +# Proactive pacing, to stay UNDER the WAF budget rather than crashing into it. Measured over +# 87 windows in one backfill, requests between one block and the next ran to a median of 161 +# (p10 = 43 - the budget is noisy, not a clean token bucket). So a spent IP does not have to +# be the trigger: after this many requests we pause for the refill window on our own, and the +# block - with its deferrals, its PARTIAL month, and the recovery sweep that PARTIAL forces - +# mostly never happens. Pacing under the budget is what lets a month finish SUCCESS on the +# first pass. Set KAP_PACE_EVERY=0 to disable. +_PACE_EVERY = int(os.environ.get("KAP_PACE_EVERY", "120")) +_PACE_SLEEP_S = float(os.environ.get("KAP_PACE_SLEEP_S", "120")) + +# Requests made on each IP since its last pause, keyed by proxy URL (None = direct). This is +# MODULE-level on purpose: the backfill builds a fresh RateLimitedClient per month, but the +# WAF budget is per-IP and global across months. A per-client counter would reset every month +# and never fire on the sparse recent months, which is exactly where it was needed - a run of +# short months would spend the shared budget between them and block. Keyed by IP so it still +# composes with a proxy pool. +_PACE_STATE: dict[str | None, int] = {} + + +def _reset_pace_state() -> None: + """Clear the module-level pace counters. For tests; not used in normal operation.""" + _PACE_STATE.clear() + + +def _load_proxies() -> list[str | None]: + """Proxy URLs for the rotation pool, or ``[None]`` for a direct connection. + + Sources, in order: the ``KAP_PROXIES`` env var (comma-separated), else a ``proxies.txt`` + file in the working directory (one URL per line, ``#`` comments allowed). The file is + gitignored and must never be committed - a proxy list is a credential. + + Returning ``[None]`` means "no pool": one direct connection, behaviour unchanged. + """ + raw = os.environ.get("KAP_PROXIES", "").strip() + entries: list[str] = [] + if raw: + entries = [p.strip() for p in raw.split(",") if p.strip()] + else: + f = Path("proxies.txt") + if f.is_file(): + entries = [ + ln.strip() + for ln in f.read_text(encoding="utf-8").splitlines() + if ln.strip() and not ln.strip().startswith("#") + ] + return list(entries) if entries else [None] + def _is_retryable(exc: BaseException) -> bool: """Which failures are worth retrying *inline*, right where they happened. @@ -32,6 +97,9 @@ def _is_retryable(exc: BaseException) -> bool: A WAF disconnect is instead surfaced immediately so the scraper can defer that disclosure, wait out the throttle once per chunk, and retry it then (see KapInsiderScraper: _WAF_COOLDOWN_S). Fail fast here, recover properly there. + + When a proxy pool is configured this same error first triggers an IP rotation (see + _request); it only reaches the scraper once every IP in the pool is cooling. """ if isinstance(exc, httpx.HTTPStatusError): return exc.response.status_code in (429, 503) @@ -42,20 +110,35 @@ class RateLimitedClient: def __init__(self) -> None: cfg = get_config() kap = cfg["kap"] - self._limiter = AsyncLimiter(float(kap["rate_limit_rps"]), 1.0) + self._rps = float(kap["rate_limit_rps"]) self._timeout = float(kap["timeout_s"]) self._headers = { "User-Agent": kap["user_agent"], "Accept-Language": "tr", } - self._client: httpx.AsyncClient | None = None + self._proxies = _load_proxies() + # One client, one limiter, one cooldown clock PER proxy. The WAF budget is per-IP, + # so the rate limit is too - a global limiter would throttle the whole pool to a + # single IP's rate and throw away the reason for having a pool. + self._clients: list[httpx.AsyncClient] = [] + self._limiters: list[AsyncLimiter] = [] + self._cooling_until: list[float] = [] + self._idx = 0 async def __aenter__(self) -> "RateLimitedClient": - self._client = httpx.AsyncClient( - headers=self._headers, - timeout=self._timeout, - follow_redirects=True, - ) + for proxy in self._proxies: + self._clients.append( + httpx.AsyncClient( + headers=self._headers, + timeout=self._timeout, + follow_redirects=True, + proxy=proxy, + ) + ) + self._limiters.append(AsyncLimiter(self._rps, 1.0)) + self._cooling_until.append(0.0) + if len(self._proxies) > 1: + _log.info("proxy_pool_active", pool_size=len(self._proxies)) return self async def __aexit__( @@ -64,8 +147,36 @@ async def __aexit__( exc_val: BaseException | None, exc_tb: TracebackType | None, ) -> None: - if self._client: - await self._client.aclose() + for client in self._clients: + await client.aclose() + self._clients.clear() + + def _acquire(self) -> int | None: + """Index of an IP that is ready to serve: neither cooling after a block nor at its + pace budget. An IP that has reached _PACE_EVERY requests is PARKED here - set cooling + for the refill window and skipped - rather than slept on. That is what makes the pool + fast: with one IP, parking it leaves nothing else and the caller waits (the proactive + pause); with ten, the other nine keep serving while it refills, so we rarely wait at + all. None when every IP is currently cooling or parked.""" + n = len(self._clients) + now = time.monotonic() + for step in range(n): + i = (self._idx + step) % n + if self._cooling_until[i] > now: + continue + ip = self._proxies[i] + if _PACE_EVERY > 0 and _PACE_STATE.get(ip, 0) >= _PACE_EVERY: + self._cooling_until[i] = now + _PACE_SLEEP_S + _PACE_STATE[ip] = 0 # the refill wait restores the budget + _log.info("pace_pause", proxy_index=i, after_requests=_PACE_EVERY, sleep_s=_PACE_SLEEP_S) + continue + # Advance the cursor so the NEXT acquire starts at the following IP. Sequentially + # this spreads load round-robin (each IP's budget lasts n times longer in wall + # time); under concurrency it hands each in-flight request a different IP, which + # is what turns the pool from a failover into real parallelism. + self._idx = (i + 1) % n + return i + return None async def get(self, url: str, **kwargs: Any) -> httpx.Response: return await self._request("GET", url, **kwargs) @@ -74,8 +185,44 @@ async def post(self, url: str, **kwargs: Any) -> httpx.Response: return await self._request("POST", url, **kwargs) async def _request(self, method: str, url: str, **kwargs: Any) -> httpx.Response: - assert self._client is not None, "Use as async context manager" - client = self._client + assert self._clients, "Use as async context manager" + single = len(self._clients) == 1 + blocks = 0 + + while True: + i = self._acquire() + if i is None: + # Every IP is cooling or parked. Wait until the soonest is free, then retry. + # For a single IP this realises the proactive pace pause; for a pool it means + # the whole pool is momentarily spent. + wait = min(self._cooling_until) - time.monotonic() + await asyncio.sleep(max(wait, 0.05)) + continue + try: + resp = await self._request_via(i, method, url, **kwargs) + ip = self._proxies[i] + _PACE_STATE[ip] = _PACE_STATE.get(ip, 0) + 1 + return resp + except httpx.RemoteProtocolError: + # This IP's budget is spent (or it is genuinely the WAF). Park it to refill. + self._cooling_until[i] = time.monotonic() + _IP_COOLDOWN_S + _PACE_STATE[self._proxies[i]] = 0 + self._idx = (i + 1) % len(self._clients) + if single: + # Unchanged single-IP contract: surface immediately so the scraper defers. + raise + _log.info("proxy_ip_cooling", proxy_index=i, cooldown_s=_IP_COOLDOWN_S) + blocks += 1 + if blocks >= len(self._clients): + # The whole pool blocked within this one call - nothing left to try now. + # Surface it so the scraper defers rather than busy-looping. + raise httpx.RemoteProtocolError("all proxy IPs cooling") + + async def _request_via( + self, i: int, method: str, url: str, **kwargs: Any + ) -> httpx.Response: + client = self._clients[i] + limiter = self._limiters[i] @retry( retry=retry_if_exception(_is_retryable), @@ -84,7 +231,7 @@ async def _request(self, method: str, url: str, **kwargs: Any) -> httpx.Response reraise=True, ) async def _do() -> httpx.Response: - async with self._limiter: + async with limiter: resp = await client.request(method, url, **kwargs) if resp.status_code == 429: retry_after = int(resp.headers.get("Retry-After", "5")) diff --git a/src/trailing_edge/scrapers/kap/insider.py b/src/trailing_edge/scrapers/kap/insider.py index cdfe04a..7465ece 100644 --- a/src/trailing_edge/scrapers/kap/insider.py +++ b/src/trailing_edge/scrapers/kap/insider.py @@ -1,5 +1,6 @@ """KAP insider scraper orchestrator.""" import asyncio +import os from dataclasses import dataclass from datetime import date from pathlib import Path @@ -46,6 +47,14 @@ # (measured 76 -> 34), which keeps the PARTIAL set small enough for the recovery pass. _WAF_COOLDOWNS_S = (90,) +# How many disclosures to fetch at once. Sequential ingest is round-trip-latency bound: each +# request to KAP through a proxy is ~1s wall time and we sat idle waiting for it, so raising +# the per-IP rate limit barely moved throughput (measured 65 -> 72/min). Concurrency is the +# real lever - N requests in flight at once, each handed a different pool IP by the client's +# round-robin. Bounded because the DB pool is 15 connections and because a public disclosure +# server should not be hit by an unbounded fan-out. 1 preserves the old sequential behaviour. +_CONCURRENCY = max(1, int(os.environ.get("KAP_CONCURRENCY", "1"))) + @dataclass class ScraperRunResult: @@ -85,6 +94,45 @@ async def _already_stored(self, ids: list[str]) -> set[str]: ) return set(rows.scalars().all()) + async def _ingest_batch( + self, + kap: "KapClient", + discs: list[dict], + counts: "_Counts", + already_stored: set[str], + *, + final: bool, + attempt: int, + ) -> list[dict]: + """Ingest a batch of disclosures with bounded concurrency, returning those that + failed (to be retried in a later cooldown pass, or counted as lost on the final one). + + Concurrency is what makes the proxy pool actually parallel: up to _CONCURRENCY + disclosures are in flight at once, and the client hands each one a different pool IP. + The shared mutable state this touches is safe under asyncio's single thread - a + coroutine only yields at an await, and neither ``deferred.append`` nor the ``counts`` + increments inside _ingest_one span one, so there is no torn read-modify-write. + """ + sem = asyncio.Semaphore(_CONCURRENCY) + deferred: list[dict] = [] + + async def _one(disc: dict) -> None: + async with sem: + try: + await self._ingest_one(kap, disc, counts, already_stored) + except Exception as exc: + deferred.append(disc) + _log.log( + 40 if final else 30, # ERROR only on the last cooldown round + "disclosure_error" if final else "disclosure_deferred", + kap_disclosure_id=str(disc.get("disclosureIndex", "")), + attempt=attempt or None, + error=str(exc), + ) + + await asyncio.gather(*(_one(d) for d in discs)) + return deferred + async def _ingest_one( self, kap: KapClient, @@ -202,17 +250,9 @@ async def run(self, from_date: date, to_date: date) -> ScraperRunResult: # them on a real backfill, silently, while the run still said SUCCESS. # A disconnect is transient, so a dropped disclosure is retried once # after a cooldown; anything still missing downgrades the run to PARTIAL. - deferred: list[dict] = [] - for disc in disclosures: - try: - await self._ingest_one(kap, disc, counts, already_stored) - except Exception as exc: - deferred.append(disc) - _log.warning( - "disclosure_deferred", - kap_disclosure_id=str(disc.get("disclosureIndex", "")), - error=str(exc), - ) + deferred = await self._ingest_batch( + kap, disclosures, counts, already_stored, final=False, attempt=0 + ) for attempt, cooldown in enumerate(_WAF_COOLDOWNS_S, start=1): if not deferred: @@ -224,20 +264,14 @@ async def run(self, from_date: date, to_date: date) -> ScraperRunResult: sleep_s=cooldown, ) await asyncio.sleep(cooldown) - retrying, deferred = deferred, [] - for disc in retrying: - try: - await self._ingest_one(kap, disc, counts, already_stored) - except Exception as exc: - deferred.append(disc) - last = attempt == len(_WAF_COOLDOWNS_S) - _log.log( - 40 if last else 30, # ERROR on the final round, else WARNING - "disclosure_error" if last else "disclosure_deferred", - kap_disclosure_id=str(disc.get("disclosureIndex", "")), - attempt=attempt, - error=str(exc), - ) + deferred = await self._ingest_batch( + kap, + deferred, + counts, + already_stored, + final=attempt == len(_WAF_COOLDOWNS_S), + attempt=attempt, + ) lost = len(deferred) except Exception as exc: diff --git a/tests/unit/core/test_proxy_rotation.py b/tests/unit/core/test_proxy_rotation.py new file mode 100644 index 0000000..49ec1d8 --- /dev/null +++ b/tests/unit/core/test_proxy_rotation.py @@ -0,0 +1,276 @@ +"""Proxy pool rotation over the KAP WAF's per-IP budget. + +The WAF blocks per source IP (RemoteProtocolError), and the budget refills with wall time. +A rotating pool exploits that: an IP that has spent its budget - whether it actually blocked, +or reached the proactive pace limit - is PARKED to refill while the other IPs keep serving. +With one IP, parking it leaves nothing else and the caller waits (the proactive pause). With +ten, we almost never wait. These tests pin: the no-proxy path still surfaces a WAF disconnect +straight to the scraper; a pace limit parks-and-rotates; a block parks-and-rotates; and a +whole pool blocking within one call still surfaces so the scraper defers. +""" +import time + +import httpx +import pytest + +from trailing_edge.core import http as http_mod +from trailing_edge.core.http import RateLimitedClient, _load_proxies + + +@pytest.fixture(autouse=True) +def _no_ambient_proxies(monkeypatch, tmp_path): + """Keep a developer's real proxies.txt or KAP_PROXIES out of the unit tests, hold the + proactive pacer off by default so rotation tests do not sleep, and clear the module-level + pace counters so state cannot leak between tests.""" + monkeypatch.delenv("KAP_PROXIES", raising=False) + monkeypatch.setattr(http_mod, "_PACE_EVERY", 0) + http_mod._reset_pace_state() + monkeypatch.chdir(tmp_path) + yield + http_mod._reset_pace_state() + + +def test_no_proxy_configured_means_one_direct_connection(): + assert _load_proxies() == [None] + + +def test_env_var_is_read_as_a_comma_separated_pool(monkeypatch): + monkeypatch.setenv("KAP_PROXIES", "http://a:1, http://b:2 ,http://c:3") + assert _load_proxies() == ["http://a:1", "http://b:2", "http://c:3"] + + +def test_proxies_file_is_read_when_env_is_absent(tmp_path, monkeypatch): + (tmp_path / "proxies.txt").write_text( + "# a comment\nhttp://a:1\n\nhttp://b:2\n", encoding="utf-8" + ) + monkeypatch.chdir(tmp_path) + assert _load_proxies() == ["http://a:1", "http://b:2"] + + +@pytest.mark.asyncio +async def test_single_client_surfaces_waf_disconnect_unchanged(monkeypatch): + """With no pool, a RemoteProtocolError must reach the caller as-is - the scraper's + defer-and-cooldown path depends on it, and adding the pool must not change it.""" + client = RateLimitedClient() + assert client._proxies == [None] + + async with client: + calls = {"n": 0} + + async def boom(*a, **k): + calls["n"] += 1 + raise httpx.RemoteProtocolError("Server disconnected") + + monkeypatch.setattr(client._clients[0], "request", boom) + + with pytest.raises(httpx.RemoteProtocolError): + await client.get("https://kap.org.tr/x") + # surfaced immediately, not retried inline + assert calls["n"] == 1 + + +@pytest.mark.asyncio +async def test_the_pacer_pauses_before_the_budget_is_spent(monkeypatch): + """After _PACE_EVERY requests on the one IP, the next request finds it parked and waits + out the refill window - one proactive pause. This is what keeps a month finishing SUCCESS + on the first pass rather than going PARTIAL and needing a sweep.""" + monkeypatch.setattr(http_mod, "_PACE_EVERY", 3) + monkeypatch.setattr(http_mod, "_PACE_SLEEP_S", 0.0) + + client = RateLimitedClient() + async with client: + slept: list[float] = [] + + async def fake_sleep(s): + slept.append(s) + + monkeypatch.setattr(http_mod.asyncio, "sleep", fake_sleep) + + async def ok(*a, **k): + return httpx.Response(200, request=httpx.Request("GET", "https://k/x")) + + monkeypatch.setattr(client._clients[0], "request", ok) + + # 3 requests fit under the budget, the 4th finds the IP parked and waits once + for _ in range(4): + await client.get("https://kap.org.tr/x") + + assert len(slept) == 1, "exactly one proactive pause after PACE_EVERY requests" + + +@pytest.mark.asyncio +async def test_pace_state_survives_a_fresh_client_across_chunks(monkeypatch): + """The bug this pins: the backfill builds a new RateLimitedClient per month, but the WAF + budget is global per-IP across months. When the counter lived on the client it reset + every month and never fired on the sparse recent months - a run of short months spent + the shared budget between them and blocked. The counter is module-level and keyed by IP, + so two requests through one client and two through the next together trip the pace at the + 4th, not restart the count.""" + monkeypatch.setattr(http_mod, "_PACE_EVERY", 3) + monkeypatch.setattr(http_mod, "_PACE_SLEEP_S", 0.0) + + slept: list[float] = [] + + async def fake_sleep(s): + slept.append(s) + + monkeypatch.setattr(http_mod.asyncio, "sleep", fake_sleep) + + async def ok(*a, **k): + return httpx.Response(200, request=httpx.Request("GET", "https://k/x")) + + async with RateLimitedClient() as c1: + monkeypatch.setattr(c1._clients[0], "request", ok) + await c1.get("https://kap.org.tr/x") + await c1.get("https://kap.org.tr/x") + + async with RateLimitedClient() as c2: + monkeypatch.setattr(c2._clients[0], "request", ok) + await c2.get("https://kap.org.tr/x") + await c2.get("https://kap.org.tr/x") + + # 4 requests on the one IP crossed the pace-of-3 once, despite the client being rebuilt + assert len(slept) == 1, "pace must accumulate across client instances, not reset" + + +@pytest.mark.asyncio +async def test_a_pace_limited_ip_rotates_instead_of_sleeping(monkeypatch): + """The change that makes the pool fast. With more than one IP, hitting the pace limit + must NOT block the caller - the spent IP is parked and the next fresh IP serves the + request immediately, no sleep.""" + monkeypatch.setenv("KAP_PROXIES", "http://ip0:1,http://ip1:1,http://ip2:1") + monkeypatch.setattr(http_mod, "_PACE_EVERY", 2) + monkeypatch.setattr(http_mod, "_PACE_SLEEP_S", 999.0) # long, to prove we do NOT sleep it + + client = RateLimitedClient() + async with client: + slept: list[float] = [] + + async def fake_sleep(s): + slept.append(s) + + monkeypatch.setattr(http_mod.asyncio, "sleep", fake_sleep) + + seen: list[int] = [] + + def make(idx): + async def h(*a, **k): + seen.append(idx) + return httpx.Response(200, request=httpx.Request("GET", "https://k/x")) + return h + + for i, c in enumerate(client._clients): + monkeypatch.setattr(c, "request", make(i)) + + # Round-robin spreads the 6 requests two full cycles across the 3 IPs. Each IP reaches + # its pace budget on the second cycle but is only parked on the NEXT acquire, so no + # request in this batch has to sleep - a fresh IP is always available. + for _ in range(6): + await client.get("https://kap.org.tr/x") + + assert seen == [0, 1, 2, 0, 1, 2] + assert slept == [], "a pool must rotate on pace, never sleep while a fresh IP exists" + + +@pytest.mark.asyncio +async def test_a_blocked_ip_rotates_to_a_fresh_one(monkeypatch): + """IP 0 throws the WAF disconnect, and the request succeeds through IP 1 rather than + failing. IP 0 is parked cooling.""" + monkeypatch.setenv("KAP_PROXIES", "http://ip0:1,http://ip1:1,http://ip2:1") + client = RateLimitedClient() + + async with client: + ok = httpx.Response(200, request=httpx.Request("GET", "https://kap.org.tr/x")) + + async def ip0(*a, **k): + raise httpx.RemoteProtocolError("Server disconnected") + + async def ip1(*a, **k): + return ok + + monkeypatch.setattr(client._clients[0], "request", ip0) + monkeypatch.setattr(client._clients[1], "request", ip1) + + resp = await client.get("https://kap.org.tr/x") + assert resp.status_code == 200 + assert client._cooling_until[0] > time.monotonic() # IP0 parked + assert client._cooling_until[1] == 0.0 # IP1 clean + + +@pytest.mark.asyncio +async def test_whole_pool_blocking_in_one_call_surfaces_the_waf(monkeypatch): + """When every IP in the pool blocks within a single call, there is nothing left to try + now: surface RemoteProtocolError so the scraper defers, rather than busy-looping.""" + monkeypatch.setenv("KAP_PROXIES", "http://ip0:1,http://ip1:1") + client = RateLimitedClient() + + async with client: + async def boom(*a, **k): + raise httpx.RemoteProtocolError("Server disconnected") + + for c in client._clients: + monkeypatch.setattr(c, "request", boom) + + with pytest.raises(httpx.RemoteProtocolError): + await client.get("https://kap.org.tr/x") + + assert all(t > time.monotonic() for t in client._cooling_until) # both parked + + +@pytest.mark.asyncio +async def test_healthy_ips_are_used_round_robin(monkeypatch): + """Each acquire advances the cursor, so successive requests spread across the pool. This + is what lets concurrency hand each in-flight request a different IP, and it also makes a + single IP's budget last n times longer in wall time.""" + monkeypatch.setenv("KAP_PROXIES", "http://ip0:1,http://ip1:1,http://ip2:1") + client = RateLimitedClient() + + async with client: + seen: list[int] = [] + + def make(idx): + async def h(*a, **k): + seen.append(idx) + return httpx.Response(200, request=httpx.Request("GET", "https://k/x")) + return h + + for i, c in enumerate(client._clients): + monkeypatch.setattr(c, "request", make(i)) + + for _ in range(6): + await client.get("https://kap.org.tr/x") + + assert seen == [0, 1, 2, 0, 1, 2] # round-robin across the pool + + +@pytest.mark.asyncio +async def test_a_block_parks_and_the_call_still_succeeds_under_round_robin(monkeypatch): + """Round-robin picks IP0, IP1, then IP0 again - and on that third request IP0 blocks, so + the call parks it and completes on IP1. Sequence: 0 (ok), 1 (ok), 0 (block -> rotate), + 1 (ok).""" + monkeypatch.setenv("KAP_PROXIES", "http://ip0:1,http://ip1:1") + client = RateLimitedClient() + + async with client: + seen: list[int] = [] + ip0_calls = {"n": 0} + + async def ip0(*a, **k): + ip0_calls["n"] += 1 + seen.append(0) + if ip0_calls["n"] >= 2: + raise httpx.RemoteProtocolError("Server disconnected") + return httpx.Response(200, request=httpx.Request("GET", "https://k/x")) + + async def ip1(*a, **k): + seen.append(1) + return httpx.Response(200, request=httpx.Request("GET", "https://k/x")) + + monkeypatch.setattr(client._clients[0], "request", ip0) + monkeypatch.setattr(client._clients[1], "request", ip1) + + for _ in range(3): + await client.get("https://kap.org.tr/x") + + assert seen == [0, 1, 0, 1] + assert client._cooling_until[0] > time.monotonic() # IP0 parked after its block diff --git a/tests/unit/scrapers/kap/test_ingest_batch.py b/tests/unit/scrapers/kap/test_ingest_batch.py new file mode 100644 index 0000000..a6e5f47 --- /dev/null +++ b/tests/unit/scrapers/kap/test_ingest_batch.py @@ -0,0 +1,79 @@ +"""Concurrent ingest of a batch of disclosures. + +The batch helper is what turns the proxy pool into real parallelism: up to _CONCURRENCY +disclosures are in flight at once. These pin the two things that must not break when the +sequential loop became a concurrent gather - every failure is still collected for retry, and +the concurrency bound is actually respected. +""" +import asyncio + +import pytest + +from trailing_edge.scrapers.kap import insider as insider_mod +from trailing_edge.scrapers.kap.insider import KapInsiderScraper, _Counts + + +@pytest.mark.asyncio +async def test_failures_are_collected_and_successes_are_not(monkeypatch): + scraper = KapInsiderScraper(backfill=True) + discs = [{"disclosureIndex": i} for i in range(10)] + + async def fake_ingest_one(kap, disc, counts, already): + # odd indices fail (transient), even ones succeed + if disc["disclosureIndex"] % 2 == 1: + raise RuntimeError("Server disconnected") + + monkeypatch.setattr(scraper, "_ingest_one", fake_ingest_one) + + deferred = await scraper._ingest_batch( + kap=None, discs=discs, counts=_Counts(), already_stored=set(), + final=False, attempt=0, + ) + + got = sorted(d["disclosureIndex"] for d in deferred) + assert got == [1, 3, 5, 7, 9], "every failed disclosure must be returned for retry" + + +@pytest.mark.asyncio +async def test_all_succeed_leaves_nothing_deferred(monkeypatch): + scraper = KapInsiderScraper(backfill=True) + discs = [{"disclosureIndex": i} for i in range(5)] + + async def ok(kap, disc, counts, already): + return None + + monkeypatch.setattr(scraper, "_ingest_one", ok) + + deferred = await scraper._ingest_batch( + kap=None, discs=discs, counts=_Counts(), already_stored=set(), + final=False, attempt=0, + ) + assert deferred == [] + + +@pytest.mark.asyncio +async def test_concurrency_bound_is_respected(monkeypatch): + """No more than _CONCURRENCY disclosures may be in flight at once - the bound protects + the DB connection pool and keeps the fan-out onto a public server civil.""" + monkeypatch.setattr(insider_mod, "_CONCURRENCY", 3) + scraper = KapInsiderScraper(backfill=True) + discs = [{"disclosureIndex": i} for i in range(12)] + + in_flight = 0 + peak = 0 + + async def slow(kap, disc, counts, already): + nonlocal in_flight, peak + in_flight += 1 + peak = max(peak, in_flight) + await asyncio.sleep(0.01) + in_flight -= 1 + + monkeypatch.setattr(scraper, "_ingest_one", slow) + + await scraper._ingest_batch( + kap=None, discs=discs, counts=_Counts(), already_stored=set(), + final=False, attempt=0, + ) + assert peak <= 3, f"concurrency bound exceeded: {peak} in flight" + assert peak == 3, "and it should actually reach the bound with 12 items and sleep" diff --git a/tests/unit/scrapers/test_backfill_loop.py b/tests/unit/scrapers/test_backfill_loop.py new file mode 100644 index 0000000..865bf0b --- /dev/null +++ b/tests/unit/scrapers/test_backfill_loop.py @@ -0,0 +1,74 @@ +"""The backfill loop must survive a month it cannot finish. + +A chunk that exhausts its WAF backoffs raises RuntimeError. Its scraper_runs record is +already written (PARTIAL or FAILED), so a later sweep collects it - that is the ledger's +whole purpose. But the RuntimeError used to propagate out of the loop and abandon every +remaining month behind the bad one. This was harmless while the backoff ladder ended at 20 +minutes (a chunk almost never exhausted); shortening the ladder to 60/120/240s made +exhaustion common, and one blocked month began killing the entire run. +""" +import importlib.util +import sys +from datetime import date +from pathlib import Path + +import pytest + +_SCRIPT = Path(__file__).resolve().parents[3] / "scripts" / "backfill_kap_insider.py" + + +def _load_module(): + spec = importlib.util.spec_from_file_location("backfill_kap_insider", _SCRIPT) + mod = importlib.util.module_from_spec(spec) + assert spec and spec.loader + sys.modules["backfill_kap_insider"] = mod + spec.loader.exec_module(mod) + return mod + + +@pytest.mark.asyncio +async def test_a_chunk_that_raises_does_not_abandon_the_months_behind_it(monkeypatch): + mod = _load_module() + + months = [ + (date(2015, 1, 1), date(2015, 1, 31)), + (date(2015, 2, 1), date(2015, 2, 28)), + (date(2015, 3, 1), date(2015, 3, 31)), + ] + + async def _noop(): + return None + + monkeypatch.setattr(mod, "init_db", _noop) + monkeypatch.setattr(mod, "generate_monthly_chunks", lambda s, e: months) + monkeypatch.setattr(mod, "get_completed_chunks", lambda: _empty_set()) + monkeypatch.setattr(mod, "_chunks_with_status", lambda status: _empty_set()) + monkeypatch.setattr(mod.asyncio, "sleep", _noop_arg) + + attempted: list[tuple] = [] + + async def fake_run_chunk(frm, to): + attempted.append((frm, to)) + if frm == date(2015, 1, 1): + # exhausted WAF backoffs + raise RuntimeError("Chunk 2015-01 failed after 3 WAF backoffs") + if frm == date(2015, 2, 1): + # a dead proxy: this is the class that used to escape the RuntimeError-only catch + import httpx + + raise httpx.ProxyError("402 Payment Required") + + monkeypatch.setattr(mod, "_run_chunk", fake_run_chunk) + + await mod.main_async(date(2015, 1, 1), dry_run=False) + + # every month was attempted - neither the RuntimeError nor the ProxyError abandoned March + assert attempted == months + + +async def _empty_set(): + return set() + + +async def _noop_arg(*a, **k): + return None