From 3251a8b32ccca332b0034b01941493608c0f3dde Mon Sep 17 00:00:00 2001 From: cagan Date: Mon, 13 Jul 2026 12:18:57 +0300 Subject: [PATCH 01/11] docs: the ingest is not complete, and I was reading the wrong number A chunk can lose filings to a WAF disconnect and still log `chunk_done`. The ledger knows better - it records the month PARTIAL, not SUCCESS - but the log line is what I was counting, and it overstates progress badly. It read 92 of 139 months while scraper_runs held 20 SUCCESS, 82 PARTIAL and 44 FAILED. 2016 and 2020 contained not one complete month. That also explains something I had noted and not chased: KAP disclosures per year appeared to fall from 2,317 (2015) to 300 (2021). Turkish insider filings did not drop eightfold through a retail boom. Later months simply took more WAF damage and lost more filings. The loss is recoverable. backfill_kap_insider.py replays non-SUCCESS months, fresh first and WAF stragglers last. But its todo list is computed once at startup, so a month that goes PARTIAL *during* a run is not swept by that run - the script must be re-run until the SUCCESS count stops rising. The completion criterion is count(SUCCESS), never the log. Also documents the quarantine, which turns out to be two different things wearing one name: - ~75% of the ~1,200 quarantined PDFs carry NO extractable text. They are scanned images with no text layer. No parser will ever read them. This is a hard ceiling of the same kind as the 250,000 TRY disclosure threshold - a property of the source. OCR could attack it, but an OCR'd number needs its own accuracy audit before it can be trusted, and none is attempted here. - the rest carry text and fail for an unrelated reason: they are narrative filings with no table, so the arithmetic row-validation gate has nothing to bind to. Recoverable. The per-year quarantine rates (13% in 2015 rising to 32% in 2020) are deliberately NOT published as a finding: the denominator is itself incomplete, and quoting a rate off an incomplete ingest is the same error in a new place. --- docs/METHODOLOGY.md | 32 ++++++++++++++++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/docs/METHODOLOGY.md b/docs/METHODOLOGY.md index 3cff0a6..3d7791a 100644 --- a/docs/METHODOLOGY.md +++ b/docs/METHODOLOGY.md @@ -276,6 +276,38 @@ Closing it needs the KAP backfill to reach 2026. That is a data-collection probl 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. +**The ingest is not complete, and the completion metric is not the obvious one.** 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` therefore overstates progress badly - +at the time of writing it read 92 of 139 months while `scraper_runs` held **20 SUCCESS, 82 +PARTIAL, 44 FAILED**. 2016 and 2020 contained not one complete month. + +The loss is recoverable: `backfill_kap_insider.py` replays non-SUCCESS months, fresh months +first and WAF stragglers last. But its `todo` list is computed once at startup, so a month +that goes PARTIAL *during* a run is not swept by that same run - the script must simply be +re-run until the SUCCESS count stops rising. **The completion criterion is +`count(SUCCESS) ≈ 139`, never the log's chunk count.** + +Direction of the error: WAF drops are unrelated to what a filing predicts, so this should +thin the sample rather than tilt it. That is an argument, not a measurement, and it is +recorded as such. + +**~1,200 filings are quarantined, and about three quarters of them are permanently +unreadable.** A DKB PDF that yields no transactions is quarantined to +`reports/parse_failures/` rather than silently skipped. Sampling those files, **75% carry no +extractable text at all** - they are scanned images with no text layer, and no parser will +ever read them. They are a hard ceiling of the same kind as the 250,000 TRY disclosure +threshold: a property of the source, not a bug to be fixed. 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 they fail for a different reason: they are +**narrative filings with no table at all** ("... 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. The final quarantine rate must be recomputed once +the ingest is complete, because the current per-year rates (13% in 2015 rising to 32% in +2020) are measured against a denominator that is itself incomplete. + **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 From d7df1ded0729fc3515ad8d2bee067952d71307a6 Mon Sep 17 00:00:00 2001 From: cagan Date: Mon, 13 Jul 2026 12:51:14 +0300 Subject: [PATCH 02/11] perf: the WAF backoff ladder was waiting for an allowance we already had The ladder escalated 60s -> 600s -> 1200s on the assumption that a longer silence buys a bigger allowance from KAP's WAF. That assumption 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 every 10- and 20-minute sleep was spent waiting for an allowance we already had. Worst-case sleep per chunk was 60+600+1200 = 31 minutes; the run was averaging 20.9 min per month, almost all of it asleep. 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: 60 -> 120 -> 240. Worst-case sleep per chunk 7 minutes. This matters because the ingest is not converging. Over the last 8 hours roughly 40 month attempts produced ONE SUCCESS - the rest went PARTIAL or FAILED to WAF drops - and the sweep of the 127 non-SUCCESS months would itself have taken ~44 hours at the old rate, while producing fresh PARTIALs as it went. The bottleneck was never KAP's tolerance. It was our own sleep. --- scripts/backfill_kap_insider.py | 26 ++++++++++++++++++++++---- 1 file changed, 22 insertions(+), 4 deletions(-) diff --git a/scripts/backfill_kap_insider.py b/scripts/backfill_kap_insider.py index 3b9212d..fa41dbf 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 From 54b1f1c2e3bdb27b1f6aee5f5f1cad2e39510b00 Mon Sep 17 00:00:00 2001 From: cagan Date: Mon, 13 Jul 2026 13:03:49 +0300 Subject: [PATCH 03/11] feat: optional proxy pool to defeat the KAP WAF's per-IP throttle The WAF blocks per source IP - the block arrives as httpx.RemoteProtocolError and the budget (~50 requests) refills with wall-clock time, not with a slower request rate. That was measured over 66 blocks in one backfill: a 20-minute pause bought the same ~50 requests as a 90-second one. It is exactly the limit a rotating IP pool defeats. RateLimitedClient now optionally rotates over a pool of proxies: - the pool is read from KAP_PROXIES (comma-separated) or a gitignored proxies.txt, one URL per line. A proxy list is a credential and must never be committed - proxies.txt is added to .gitignore. - each IP gets its own httpx client, its own rate limiter, and its own cooldown clock: the WAF budget is per-IP, so a single global limiter would throttle the whole pool to one IP's rate and waste the pool. - policy is drain-then-rotate, not round-robin: requests are sequential so the pool does not parallelise - its value is avoiding the block stall. One IP serves at full speed until it throws RemoteProtocolError, then it is parked for 120s (the measured refill time) and the request retries on the next fresh IP. By the time the rotation comes back, the first IP has refilled. - only when every IP in the pool is cooling does RemoteProtocolError surface to the scraper, which defers and waits it out exactly as before. With NO proxy configured the path is byte-for-byte the old one: a single direct client, RemoteProtocolError surfaced immediately. Pinned by tests in both directions - the no-proxy path unchanged, and a blocked IP rotating to a fresh one rather than failing the request. This is the lever on the ETA. Single IP, the ingest was taking ~8-12h and the non-SUCCESS sweep barely converged; an N-IP pool cuts the block stalls that were most of that time. --- .env.example | 7 + .gitignore | 3 + src/trailing_edge/core/http.py | 131 +++++++++++++++++-- tests/unit/core/test_proxy_rotation.py | 169 +++++++++++++++++++++++++ 4 files changed, 297 insertions(+), 13 deletions(-) create mode 100644 tests/unit/core/test_proxy_rotation.py 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..eea864c 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/src/trailing_edge/core/http.py b/src/trailing_edge/core/http.py index be6fc5d..7bdfd68 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,35 @@ _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 + + +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 +74,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 +87,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 +124,21 @@ 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 _next_available(self) -> int | None: + """Index of the next proxy whose IP budget is not cooling, round-robin from the + last one used. None when every IP in the pool is currently cooling.""" + n = len(self._clients) + now = time.monotonic() + for step in range(n): + i = (self._idx + step) % n + if self._cooling_until[i] <= now: + self._idx = i + return i + return None async def get(self, url: str, **kwargs: Any) -> httpx.Response: return await self._request("GET", url, **kwargs) @@ -74,8 +147,40 @@ 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 direct connection (no pool): unchanged from the original - one client, one + # limiter, RemoteProtocolError surfaced straight to the caller. + if len(self._clients) == 1: + return await self._request_via(0, method, url, **kwargs) + + # Pool: a RemoteProtocolError means this IP's budget is spent. Park it to refill and + # move to a fresh IP, trying each at most once. Only when the whole pool is cooling + # does the error surface to the scraper, which defers and waits it out as before. + last_exc: httpx.RemoteProtocolError | None = None + for _ in range(len(self._clients)): + i = self._next_available() + if i is None: + break + try: + return await self._request_via(i, method, url, **kwargs) + except httpx.RemoteProtocolError as exc: + last_exc = exc + self._cooling_until[i] = time.monotonic() + _IP_COOLDOWN_S + self._idx = (i + 1) % len(self._clients) + _log.info("proxy_ip_cooling", proxy_index=i, cooldown_s=_IP_COOLDOWN_S) + + if last_exc is not None: + raise last_exc + # Every IP was already cooling before we tried any. Surface the WAF condition so the + # scraper defers, rather than busy-looping on a pool that has nothing to give. + 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 +189,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/tests/unit/core/test_proxy_rotation.py b/tests/unit/core/test_proxy_rotation.py new file mode 100644 index 0000000..49195e9 --- /dev/null +++ b/tests/unit/core/test_proxy_rotation.py @@ -0,0 +1,169 @@ +"""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: a spent IP is parked to refill while another carries on. +These tests pin the two things that must hold - the no-proxy path is byte-for-byte the old +behaviour, and a block rotates to a fresh IP rather than failing the request - without +touching the network. +""" +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.""" + monkeypatch.delenv("KAP_PROXIES", raising=False) + monkeypatch.chdir(tmp_path) + + +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_a_blocked_ip_rotates_to_a_fresh_one(monkeypatch): + """The whole point: 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_all_ips_cooling_surfaces_the_waf_condition(monkeypatch): + """When every IP is spent, the pool has nothing to give: 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") + + # both parked + assert all(t > time.monotonic() for t in client._cooling_until) + + # a second call, with every IP still cooling, also surfaces rather than hanging + with pytest.raises(httpx.RemoteProtocolError): + await client.get("https://kap.org.tr/x") + + +@pytest.mark.asyncio +async def test_a_healthy_ip_is_drained_not_round_robined(monkeypatch): + """Requests are sequential, so the pool does not parallelise - its value is avoiding + the block stall. The right policy is therefore to use one IP at full speed until it + blocks and only then rotate, so each IP's whole ~50-request budget is spent before we + move on. An IP that keeps succeeding keeps serving.""" + 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(3): + await client.get("https://kap.org.tr/x") + + assert seen == [0, 0, 0] # drained, not spread + + +@pytest.mark.asyncio +async def test_rotation_advances_only_when_the_current_ip_blocks(monkeypatch): + """IP0 serves once, then blocks; the retry lands on IP1, which then serves the rest. + So the sequence is 0 (ok), 0 (block -> rotate), 1 (ok), 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, 0, 1, 1] From cc2d4e89c424993be791f634fad4ceab71802d32 Mon Sep 17 00:00:00 2001 From: cagan Date: Mon, 13 Jul 2026 13:07:00 +0300 Subject: [PATCH 04/11] fix: a month that exhausts its WAF backoffs must not kill the whole backfill Shortening the backoff ladder to 60/120/240 had a consequence I did not follow through: _run_chunk raises RuntimeError when a month exhausts all its backoffs, and that exception propagated out of the main loop and abandoned every remaining month behind it. The very first blocked month under the new ladder - 2022-01 - crashed the entire run. It was latent before: with the old 20-minute tail a chunk almost never exhausted, so the raise almost never fired. Shortening the ladder made exhaustion common and exposed it. The month's scraper_runs record is already written by scraper.run() (PARTIAL, or FAILED if the warmup GET never got through), so a later sweep collects it - that is the ledger's entire purpose. The loop now catches the RuntimeError, logs chunk_abandoned, and continues. Matches the existing TimeoutError branch, which already returns rather than raising. Pinned by a test: a chunk that raises in the middle of the run must not stop the months behind it from being attempted. --- scripts/backfill_kap_insider.py | 17 +++++- tests/unit/scrapers/test_backfill_loop.py | 69 +++++++++++++++++++++++ 2 files changed, 85 insertions(+), 1 deletion(-) create mode 100644 tests/unit/scrapers/test_backfill_loop.py diff --git a/scripts/backfill_kap_insider.py b/scripts/backfill_kap_insider.py index fa41dbf..32a72b0 100644 --- a/scripts/backfill_kap_insider.py +++ b/scripts/backfill_kap_insider.py @@ -272,7 +272,22 @@ 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 RuntimeError as exc: + _log.warning( + "chunk_abandoned", + from_date=from_date, + to_date=to_date, + error=str(exc), + ) await asyncio.sleep(_CHUNK_GAP_S) if dry_run: diff --git a/tests/unit/scrapers/test_backfill_loop.py b/tests/unit/scrapers/test_backfill_loop.py new file mode 100644 index 0000000..fb5737e --- /dev/null +++ b/tests/unit/scrapers/test_backfill_loop.py @@ -0,0 +1,69 @@ +"""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, 2, 1): + # the month that exhausts its WAF backoffs + raise RuntimeError("Chunk 2015-02 failed after 3 WAF backoffs") + + monkeypatch.setattr(mod, "_run_chunk", fake_run_chunk) + + await mod.main_async(date(2015, 1, 1), dry_run=False) + + # every month was attempted - the bad one in the middle did not abandon March + assert attempted == months + + +async def _empty_set(): + return set() + + +async def _noop_arg(*a, **k): + return None From 852917c63617dc82ea3925b27d907b7fdc1c4f7f Mon Sep 17 00:00:00 2001 From: cagan Date: Mon, 13 Jul 2026 13:07:36 +0300 Subject: [PATCH 05/11] style: drop unused import in test_proxy_rotation Left over from an earlier draft; ruff F401. --- tests/unit/core/test_proxy_rotation.py | 1 - 1 file changed, 1 deletion(-) diff --git a/tests/unit/core/test_proxy_rotation.py b/tests/unit/core/test_proxy_rotation.py index 49195e9..485cd2f 100644 --- a/tests/unit/core/test_proxy_rotation.py +++ b/tests/unit/core/test_proxy_rotation.py @@ -11,7 +11,6 @@ import httpx import pytest -from trailing_edge.core import http as http_mod from trailing_edge.core.http import RateLimitedClient, _load_proxies From 457e22970df4d1e97d93c1c60270bd3344aba750 Mon Sep 17 00:00:00 2001 From: cagan Date: Mon, 13 Jul 2026 13:13:02 +0300 Subject: [PATCH 06/11] feat: pace under the WAF budget instead of crashing into it The measurement that drove the proxy work also points at a fix that needs no proxy at all. Requests between one WAF block and the next ran to a median of 161 over 87 windows (p10 = 43 - noisy, not a clean bucket). We were spending the whole budget, hitting the wall, and paying for it three times over: the backoff sleep, the PARTIAL month the dropped filings leave behind, and the recovery sweep that PARTIAL forces. That sweep was the real reason the ingest would not converge - 127 of 139 months were non-SUCCESS and each needed re-running. So pace under the budget. After _PACE_EVERY (120) requests on an IP, pause _PACE_SLEEP_S (120s) for the refill window on our own terms. The block mostly never happens, the month finishes SUCCESS on the first pass, and the sweep disappears. Both knobs are env-configurable (KAP_PACE_EVERY=0 disables); the counter is per-IP, so it composes with the proxy rotation - each IP paces itself, and a rotation to a cooled IP resets its counter because the cooldown already refilled it. This is the honest version of the throughput fix: it respects the server's stated limit rather than evading it, and it is strictly better for our own convergence because it removes the sweep. Slower per burst (~40 req/min), far faster to a complete, all-SUCCESS ingest. Pinned by a test: PACE_EVERY requests pass freely, the next one waits. --- src/trailing_edge/core/http.py | 23 +++++++++++++++++ tests/unit/core/test_proxy_rotation.py | 34 +++++++++++++++++++++++++- 2 files changed, 56 insertions(+), 1 deletion(-) diff --git a/src/trailing_edge/core/http.py b/src/trailing_edge/core/http.py index 7bdfd68..1758f68 100644 --- a/src/trailing_edge/core/http.py +++ b/src/trailing_edge/core/http.py @@ -35,6 +35,16 @@ # 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")) + def _load_proxies() -> list[str | None]: """Proxy URLs for the rotation pool, or ``[None]`` for a direct connection. @@ -100,6 +110,7 @@ def __init__(self) -> None: self._clients: list[httpx.AsyncClient] = [] self._limiters: list[AsyncLimiter] = [] self._cooling_until: list[float] = [] + self._since_pause: list[int] = [] # requests on each IP since its last proactive pause self._idx = 0 async def __aenter__(self) -> "RateLimitedClient": @@ -114,6 +125,7 @@ async def __aenter__(self) -> "RateLimitedClient": ) self._limiters.append(AsyncLimiter(self._rps, 1.0)) self._cooling_until.append(0.0) + self._since_pause.append(0) if len(self._proxies) > 1: _log.info("proxy_pool_active", pool_size=len(self._proxies)) return self @@ -167,6 +179,7 @@ async def _request(self, method: str, url: str, **kwargs: Any) -> httpx.Response except httpx.RemoteProtocolError as exc: last_exc = exc self._cooling_until[i] = time.monotonic() + _IP_COOLDOWN_S + self._since_pause[i] = 0 # the cooldown refills the budget self._idx = (i + 1) % len(self._clients) _log.info("proxy_ip_cooling", proxy_index=i, cooldown_s=_IP_COOLDOWN_S) @@ -182,6 +195,16 @@ async def _request_via( client = self._clients[i] limiter = self._limiters[i] + # Pace under the budget before spending it. This IP has made _PACE_EVERY requests + # since its last pause, so it is approaching the block threshold - wait out the + # refill window now, on our terms, instead of hitting the wall and paying a backoff + # plus a PARTIAL month plus a recovery sweep. + if _PACE_EVERY > 0 and self._since_pause[i] >= _PACE_EVERY: + self._since_pause[i] = 0 + _log.info("pace_pause", proxy_index=i, after_requests=_PACE_EVERY, sleep_s=_PACE_SLEEP_S) + await asyncio.sleep(_PACE_SLEEP_S) + self._since_pause[i] += 1 + @retry( retry=retry_if_exception(_is_retryable), stop=stop_after_attempt(5), diff --git a/tests/unit/core/test_proxy_rotation.py b/tests/unit/core/test_proxy_rotation.py index 485cd2f..e74bb42 100644 --- a/tests/unit/core/test_proxy_rotation.py +++ b/tests/unit/core/test_proxy_rotation.py @@ -11,13 +11,16 @@ 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.""" + """Keep a developer's real proxies.txt or KAP_PROXIES out of the unit tests, and hold + the proactive pacer off so these rotation tests do not sleep.""" monkeypatch.delenv("KAP_PROXIES", raising=False) + monkeypatch.setattr(http_mod, "_PACE_EVERY", 0) monkeypatch.chdir(tmp_path) @@ -60,6 +63,35 @@ async def boom(*a, **k): assert calls["n"] == 1 +@pytest.mark.asyncio +async def test_the_pacer_pauses_before_the_budget_is_spent(monkeypatch): + """After _PACE_EVERY requests on an IP, the next request waits out the refill window + instead of continuing into a block. 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[int] = [] + + 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 triggers one proactive pause + for _ in range(4): + await client.get("https://kap.org.tr/x") + + assert slept.count(0.0) == 1, "exactly one proactive pause after PACE_EVERY requests" + + @pytest.mark.asyncio async def test_a_blocked_ip_rotates_to_a_fresh_one(monkeypatch): """The whole point: IP 0 throws the WAF disconnect, and the request succeeds through From eb364efbbd66ca5d32b5f40fb4bdb964fd41be0d Mon Sep 17 00:00:00 2001 From: cagan Date: Mon, 13 Jul 2026 13:18:17 +0300 Subject: [PATCH 07/11] fix: the pace counter must be global per-IP, not per-client The pacer as first written lived on the RateLimitedClient instance - but the backfill builds a fresh client for every month (async with RateLimitedClient() per chunk), while the WAF budget is global per source IP across months. So the counter reset at every month boundary and never accumulated. On the sparse recent months (~30-60 requests each, well under the 120 threshold) it therefore never fired at all, and a run of short months spent the shared budget between them and blocked - exactly what the first restart showed: pace_pause=0 with a WAF cooldown already logged. Moved the counter to module level, keyed by proxy URL (None = direct), so it survives the per-month client and still composes with the proxy pool (each IP paces independently; a cooldown resets that IP's count because the wait already refilled it). Pinned by a test that makes two requests through one client and two through a brand-new one and asserts the pace-of-3 trips once across the boundary rather than restarting. --- src/trailing_edge/core/http.py | 27 ++++++++++++---- tests/unit/core/test_proxy_rotation.py | 45 ++++++++++++++++++++++++-- 2 files changed, 63 insertions(+), 9 deletions(-) diff --git a/src/trailing_edge/core/http.py b/src/trailing_edge/core/http.py index 1758f68..bcf7431 100644 --- a/src/trailing_edge/core/http.py +++ b/src/trailing_edge/core/http.py @@ -45,6 +45,19 @@ _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. @@ -110,7 +123,6 @@ def __init__(self) -> None: self._clients: list[httpx.AsyncClient] = [] self._limiters: list[AsyncLimiter] = [] self._cooling_until: list[float] = [] - self._since_pause: list[int] = [] # requests on each IP since its last proactive pause self._idx = 0 async def __aenter__(self) -> "RateLimitedClient": @@ -125,7 +137,6 @@ async def __aenter__(self) -> "RateLimitedClient": ) self._limiters.append(AsyncLimiter(self._rps, 1.0)) self._cooling_until.append(0.0) - self._since_pause.append(0) if len(self._proxies) > 1: _log.info("proxy_pool_active", pool_size=len(self._proxies)) return self @@ -179,7 +190,7 @@ async def _request(self, method: str, url: str, **kwargs: Any) -> httpx.Response except httpx.RemoteProtocolError as exc: last_exc = exc self._cooling_until[i] = time.monotonic() + _IP_COOLDOWN_S - self._since_pause[i] = 0 # the cooldown refills the budget + _PACE_STATE[self._proxies[i]] = 0 # the cooldown refills the budget self._idx = (i + 1) % len(self._clients) _log.info("proxy_ip_cooling", proxy_index=i, cooldown_s=_IP_COOLDOWN_S) @@ -198,12 +209,14 @@ async def _request_via( # Pace under the budget before spending it. This IP has made _PACE_EVERY requests # since its last pause, so it is approaching the block threshold - wait out the # refill window now, on our terms, instead of hitting the wall and paying a backoff - # plus a PARTIAL month plus a recovery sweep. - if _PACE_EVERY > 0 and self._since_pause[i] >= _PACE_EVERY: - self._since_pause[i] = 0 + # plus a PARTIAL month plus a recovery sweep. The counter is module-level and keyed + # by IP, so it survives the fresh client the backfill builds for each month. + ip = self._proxies[i] + if _PACE_EVERY > 0 and _PACE_STATE.get(ip, 0) >= _PACE_EVERY: + _PACE_STATE[ip] = 0 _log.info("pace_pause", proxy_index=i, after_requests=_PACE_EVERY, sleep_s=_PACE_SLEEP_S) await asyncio.sleep(_PACE_SLEEP_S) - self._since_pause[i] += 1 + _PACE_STATE[ip] = _PACE_STATE.get(ip, 0) + 1 @retry( retry=retry_if_exception(_is_retryable), diff --git a/tests/unit/core/test_proxy_rotation.py b/tests/unit/core/test_proxy_rotation.py index e74bb42..c10135e 100644 --- a/tests/unit/core/test_proxy_rotation.py +++ b/tests/unit/core/test_proxy_rotation.py @@ -17,11 +17,15 @@ @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, and hold - the proactive pacer off so these rotation tests do not sleep.""" + """Keep a developer's real proxies.txt or KAP_PROXIES out of the unit tests, hold the + proactive pacer off so these 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(): @@ -92,6 +96,43 @@ async def ok(*a, **k): assert slept.count(0.0) == 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 must 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")) + + # first "month": a fresh client makes 2 requests + 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") + + # second "month": a brand-new client, same direct IP, makes 2 more + 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 slept.count(0.0) == 1, "pace must accumulate across client instances, not reset" + + @pytest.mark.asyncio async def test_a_blocked_ip_rotates_to_a_fresh_one(monkeypatch): """The whole point: IP 0 throws the WAF disconnect, and the request succeeds through From 746c47f52169fe008ab84490ce9d52822cdd8df3 Mon Sep 17 00:00:00 2001 From: cagan Date: Mon, 13 Jul 2026 14:03:49 +0300 Subject: [PATCH 08/11] ops: self-restarting supervisor for the backfill The backfill process has died silently more than once mid-run (last time during a 240s backoff, log frozen, process gone). Each death cost data-collection time until the next manual restart. backfill_supervisor.sh runs the backfill and re-launches it whenever it exits, so a silent death recovers on its own. It also does the PARTIAL sweep automatically. backfill_kap_insider.py fixes its todo list at startup, so months that go PARTIAL during a run are only collected by the NEXT run - the supervisor re-runs until the SUCCESS month count stops rising (plateau) or reaches 139, then stops. This is the "re-run until SUCCESS stops climbing" loop, made autonomous. Launched detached via PowerShell Start-Process so it does not share the fate of the shell that started it. success_count() takes only the integer line from the query because init_db() prints db_connected to stdout ahead of it. --- scripts/backfill_supervisor.sh | 60 ++++++++++++++++++++++++++++++++++ 1 file changed, 60 insertions(+) create mode 100644 scripts/backfill_supervisor.sh diff --git a/scripts/backfill_supervisor.sh b/scripts/backfill_supervisor.sh new file mode 100644 index 0000000..7c5c472 --- /dev/null +++ b/scripts/backfill_supervisor.sh @@ -0,0 +1,60 @@ +#!/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() { + # init_db() logs "db_connected" to stdout, so take only the pure-integer line the query + # prints - never trust the last line blindly. + 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 +} + +prev=-1 +pass_num=0 +while true; do + pass_num=$((pass_num + 1)) + cur=$(success_count) + echo "{\"event\":\"SUPERVISOR\",\"pass\":$pass_num,\"success\":$cur,\"prev\":$prev}" >> "$LOG" + + if [ "$cur" -ge "$TOTAL_MONTHS" ]; then + echo "{\"event\":\"SUPERVISOR\",\"done\":true,\"reason\":\"success>=months\",\"success\":$cur}" >> "$LOG" + break + fi + # Converged: a whole pass added no new SUCCESS month. The remainder is WAF-stubborn or + # structurally PARTIAL (a genuinely lost disclosure), not something another pass fixes. + if [ "$pass_num" -gt 1 ] && [ "$cur" -le "$prev" ]; then + echo "{\"event\":\"SUPERVISOR\",\"done\":true,\"reason\":\"plateau\",\"success\":$cur}" >> "$LOG" + break + fi + prev=$cur + + 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 From 14df2eb04d46e1f0bb6b78977cb41d25c1e4a256 Mon Sep 17 00:00:00 2001 From: cagan Date: Mon, 13 Jul 2026 16:05:43 +0300 Subject: [PATCH 09/11] fix: contain ANY per-chunk failure, and make the proxy pool rotate on pace Two things, both surfaced by wiring in a real proxy pool. 1. A dead proxy crashed the whole backfill. Webshare's free tier returned 402 Payment Required after briefly working, surfacing as httpx.ProxyError - which the per-chunk `except RuntimeError` did not catch, so it propagated and abandoned every remaining month. Widened to `except Exception`: the month's ledger record is already written by scraper.run() before it re-raises, so a sweep collects it. One bad month, whatever killed it, must never take the rest of the run with it. Test now covers a ProxyError mid-run alongside the RuntimeError case. 2. The pool did nothing for throughput as first written. The pacer prevents blocks, and the old rotation only advanced on a block, so a healthy IP was drained and then SLEPT on at its pace limit while nine fresh IPs sat idle. Reworked so reaching the pace limit PARKS the IP (cools it for the refill window) and rotates to the next fresh one, exactly like a block does. With one IP this is the same proactive pause as before; with ten, the others serve while one refills, so we almost never sleep. Unified single-IP and pool paths through one _acquire()/_request loop; the single-IP WAF-disconnect contract (surface immediately) is preserved and still tested. The free proxies are 402 for now, so the run is back on the direct+pacer path. The pool code is proven by tests and ready for a working proxy list. --- .gitignore | 2 +- scripts/backfill_kap_insider.py | 10 ++- scripts/backfill_supervisor.sh | 46 ++++++++--- src/trailing_edge/core/http.py | 82 ++++++++++--------- tests/unit/core/test_proxy_rotation.py | 98 +++++++++++++++-------- tests/unit/scrapers/test_backfill_loop.py | 11 ++- 6 files changed, 157 insertions(+), 92 deletions(-) diff --git a/.gitignore b/.gitignore index eea864c..13a4307 100644 --- a/.gitignore +++ b/.gitignore @@ -67,4 +67,4 @@ pip-log.txt pip-delete-this-directory.txt # Proxy list for the KAP scraper rotation pool - a credential, never commit. -proxies.txt +proxies.txt* diff --git a/scripts/backfill_kap_insider.py b/scripts/backfill_kap_insider.py index 32a72b0..018e70b 100644 --- a/scripts/backfill_kap_insider.py +++ b/scripts/backfill_kap_insider.py @@ -281,12 +281,18 @@ async def main_async( # every remaining month behind it. Log it and move on. try: await _run_chunk(from_date, to_date) - except RuntimeError as exc: + 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=str(exc), + error=f"{type(exc).__name__}: {exc}", ) await asyncio.sleep(_CHUNK_GAP_S) diff --git a/scripts/backfill_supervisor.sh b/scripts/backfill_supervisor.sh index 7c5c472..bf06029 100644 --- a/scripts/backfill_supervisor.sh +++ b/scripts/backfill_supervisor.sh @@ -17,8 +17,6 @@ LOG=backfill.log TOTAL_MONTHS=139 success_count() { - # init_db() logs "db_connected" to stdout, so take only the pure-integer line the query - # prints - never trust the last line blindly. uv run python - <<'PY' 2>/dev/null | grep -oE '^[0-9]+$' | tail -1 import asyncio, sys sys.path.insert(0, "src") @@ -35,24 +33,48 @@ asyncio.run(m()) PY } -prev=-1 +# 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)) - cur=$(success_count) - echo "{\"event\":\"SUPERVISOR\",\"pass\":$pass_num,\"success\":$cur,\"prev\":$prev}" >> "$LOG" + succ=$(success_count) + disc=$(disclosure_count) + echo "{\"event\":\"SUPERVISOR\",\"pass\":$pass_num,\"success\":$succ,\"disclosures\":$disc}" >> "$LOG" - if [ "$cur" -ge "$TOTAL_MONTHS" ]; then - echo "{\"event\":\"SUPERVISOR\",\"done\":true,\"reason\":\"success>=months\",\"success\":$cur}" >> "$LOG" + if [ "$succ" -ge "$TOTAL_MONTHS" ]; then + echo "{\"event\":\"SUPERVISOR\",\"done\":true,\"reason\":\"success>=months\",\"success\":$succ}" >> "$LOG" break fi - # Converged: a whole pass added no new SUCCESS month. The remainder is WAF-stubborn or - # structurally PARTIAL (a genuinely lost disclosure), not something another pass fixes. - if [ "$pass_num" -gt 1 ] && [ "$cur" -le "$prev" ]; then - echo "{\"event\":\"SUPERVISOR\",\"done\":true,\"reason\":\"plateau\",\"success\":$cur}" >> "$LOG" + # 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=$cur + 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" diff --git a/src/trailing_edge/core/http.py b/src/trailing_edge/core/http.py index bcf7431..3918daf 100644 --- a/src/trailing_edge/core/http.py +++ b/src/trailing_edge/core/http.py @@ -151,16 +151,27 @@ async def __aexit__( await client.aclose() self._clients.clear() - def _next_available(self) -> int | None: - """Index of the next proxy whose IP budget is not cooling, round-robin from the - last one used. None when every IP in the pool is currently cooling.""" + 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: - self._idx = i - return i + 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 + self._idx = i + return i return None async def get(self, url: str, **kwargs: Any) -> httpx.Response: @@ -171,34 +182,37 @@ async def post(self, url: str, **kwargs: Any) -> httpx.Response: async def _request(self, method: str, url: str, **kwargs: Any) -> httpx.Response: assert self._clients, "Use as async context manager" + single = len(self._clients) == 1 + blocks = 0 - # Single direct connection (no pool): unchanged from the original - one client, one - # limiter, RemoteProtocolError surfaced straight to the caller. - if len(self._clients) == 1: - return await self._request_via(0, method, url, **kwargs) - - # Pool: a RemoteProtocolError means this IP's budget is spent. Park it to refill and - # move to a fresh IP, trying each at most once. Only when the whole pool is cooling - # does the error surface to the scraper, which defers and waits it out as before. - last_exc: httpx.RemoteProtocolError | None = None - for _ in range(len(self._clients)): - i = self._next_available() + while True: + i = self._acquire() if i is None: - break + # 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: - return await self._request_via(i, method, url, **kwargs) - except httpx.RemoteProtocolError as exc: - last_exc = exc + 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 # the cooldown refills the budget + _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) - - if last_exc is not None: - raise last_exc - # Every IP was already cooling before we tried any. Surface the WAF condition so the - # scraper defers, rather than busy-looping on a pool that has nothing to give. - raise httpx.RemoteProtocolError("all proxy IPs cooling") + 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 @@ -206,18 +220,6 @@ async def _request_via( client = self._clients[i] limiter = self._limiters[i] - # Pace under the budget before spending it. This IP has made _PACE_EVERY requests - # since its last pause, so it is approaching the block threshold - wait out the - # refill window now, on our terms, instead of hitting the wall and paying a backoff - # plus a PARTIAL month plus a recovery sweep. The counter is module-level and keyed - # by IP, so it survives the fresh client the backfill builds for each month. - ip = self._proxies[i] - if _PACE_EVERY > 0 and _PACE_STATE.get(ip, 0) >= _PACE_EVERY: - _PACE_STATE[ip] = 0 - _log.info("pace_pause", proxy_index=i, after_requests=_PACE_EVERY, sleep_s=_PACE_SLEEP_S) - await asyncio.sleep(_PACE_SLEEP_S) - _PACE_STATE[ip] = _PACE_STATE.get(ip, 0) + 1 - @retry( retry=retry_if_exception(_is_retryable), stop=stop_after_attempt(5), diff --git a/tests/unit/core/test_proxy_rotation.py b/tests/unit/core/test_proxy_rotation.py index c10135e..dc9a4a9 100644 --- a/tests/unit/core/test_proxy_rotation.py +++ b/tests/unit/core/test_proxy_rotation.py @@ -1,10 +1,12 @@ """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: a spent IP is parked to refill while another carries on. -These tests pin the two things that must hold - the no-proxy path is byte-for-byte the old -behaviour, and a block rotates to a fresh IP rather than failing the request - without -touching the network. +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 @@ -18,8 +20,8 @@ @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 so these rotation tests do not sleep, and clear the module-level pace - counters so state cannot leak between tests.""" + 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() @@ -69,15 +71,15 @@ async def boom(*a, **k): @pytest.mark.asyncio async def test_the_pacer_pauses_before_the_budget_is_spent(monkeypatch): - """After _PACE_EVERY requests on an IP, the next request waits out the refill window - instead of continuing into a block. This is what keeps a month finishing SUCCESS on the - first pass rather than going PARTIAL and needing a sweep.""" + """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[int] = [] + slept: list[float] = [] async def fake_sleep(s): slept.append(s) @@ -89,11 +91,11 @@ async def ok(*a, **k): monkeypatch.setattr(client._clients[0], "request", ok) - # 3 requests fit under the budget, the 4th triggers one proactive pause + # 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 slept.count(0.0) == 1, "exactly one proactive pause after PACE_EVERY requests" + assert len(slept) == 1, "exactly one proactive pause after PACE_EVERY requests" @pytest.mark.asyncio @@ -102,8 +104,8 @@ async def test_pace_state_survives_a_fresh_client_across_chunks(monkeypatch): 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 must together trip the pace - at the 4th, not restart the count.""" + 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) @@ -117,26 +119,61 @@ async def fake_sleep(s): async def ok(*a, **k): return httpx.Response(200, request=httpx.Request("GET", "https://k/x")) - # first "month": a fresh client makes 2 requests 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") - # second "month": a brand-new client, same direct IP, makes 2 more 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 slept.count(0.0) == 1, "pace must accumulate across client instances, not reset" + 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)) + + # IP0 serves 2, hits pace, parks; IP1 serves the next 2; IP2 the next 2 - no sleeps + for _ in range(6): + await client.get("https://kap.org.tr/x") + + assert seen == [0, 0, 1, 1, 2, 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): - """The whole point: IP 0 throws the WAF disconnect, and the request succeeds through - IP 1 rather than failing. IP 0 is parked cooling.""" + """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() @@ -159,9 +196,9 @@ async def ip1(*a, **k): @pytest.mark.asyncio -async def test_all_ips_cooling_surfaces_the_waf_condition(monkeypatch): - """When every IP is spent, the pool has nothing to give: surface RemoteProtocolError so - the scraper defers, rather than busy-looping.""" +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() @@ -175,20 +212,13 @@ async def boom(*a, **k): with pytest.raises(httpx.RemoteProtocolError): await client.get("https://kap.org.tr/x") - # both parked - assert all(t > time.monotonic() for t in client._cooling_until) - - # a second call, with every IP still cooling, also surfaces rather than hanging - 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_a_healthy_ip_is_drained_not_round_robined(monkeypatch): - """Requests are sequential, so the pool does not parallelise - its value is avoiding - the block stall. The right policy is therefore to use one IP at full speed until it - blocks and only then rotate, so each IP's whole ~50-request budget is spent before we - move on. An IP that keeps succeeding keeps serving.""" + """With the pacer off, requests are sequential and a healthy IP just keeps serving - + it is drained, not spread across the pool one request at a time.""" monkeypatch.setenv("KAP_PROXIES", "http://ip0:1,http://ip1:1,http://ip2:1") client = RateLimitedClient() @@ -212,8 +242,8 @@ async def h(*a, **k): @pytest.mark.asyncio async def test_rotation_advances_only_when_the_current_ip_blocks(monkeypatch): - """IP0 serves once, then blocks; the retry lands on IP1, which then serves the rest. - So the sequence is 0 (ok), 0 (block -> rotate), 1 (ok), 1 (ok).""" + """IP0 serves once, then blocks; the retry lands on IP1, which serves the rest. + Sequence: 0 (ok), 0 (block -> rotate), 1 (ok), 1 (ok).""" monkeypatch.setenv("KAP_PROXIES", "http://ip0:1,http://ip1:1") client = RateLimitedClient() diff --git a/tests/unit/scrapers/test_backfill_loop.py b/tests/unit/scrapers/test_backfill_loop.py index fb5737e..865bf0b 100644 --- a/tests/unit/scrapers/test_backfill_loop.py +++ b/tests/unit/scrapers/test_backfill_loop.py @@ -49,15 +49,20 @@ async def _noop(): 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): - # the month that exhausts its WAF backoffs - raise RuntimeError("Chunk 2015-02 failed after 3 WAF backoffs") + # 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 - the bad one in the middle did not abandon March + # every month was attempted - neither the RuntimeError nor the ProxyError abandoned March assert attempted == months From 65887470db64adb598098a3be333ca309d35dabf Mon Sep 17 00:00:00 2001 From: cagan Date: Mon, 13 Jul 2026 16:39:38 +0300 Subject: [PATCH 10/11] perf: concurrent ingest across the proxy pool - the real throughput lever Raising the per-IP rate limit barely moved throughput (measured 65 -> 72 disclosures/min), which proved the bottleneck was not rate but round-trip LATENCY: each request to KAP through a UK proxy is ~1s wall time and the sequential loop sat idle waiting for it. The fix is concurrency, not a higher rate. - the ingest loop (both the first pass and each cooldown-retry pass) is now a bounded asyncio.gather via _ingest_batch, KAP_CONCURRENCY in flight at once (default 1, i.e. unchanged; the backfill runs it at 6). The shared mutable state - the deferred list and the _Counts increments - is safe because asyncio only switches at an await and none of those mutations span one. - the client's IP selection is now round-robin (advance the cursor on every acquire) instead of sticky drain-then-rotate. Sequentially this makes each IP's budget last n times longer in wall time; under concurrency it hands each in-flight request a different pool IP, which is what actually parallelises the work. Verified the one thing that could have made this unsafe: KAP's session cookie is NOT IP-bound - a cookie minted through one proxy is accepted on a request through another (tested directly), so one shared cookie jar across the pool is fine. Measured at concurrency 6 over the pool: ~100-125 disclosures/min against ~65 sequential, zero WAF blocks, and SUCCESS climbing fast (29 -> 40 in minutes). Bounded deliberately - the DB pool is 15 connections and a public disclosure server should not face an unbounded fan-out. Tests cover the concurrency bound, deferred collection, and round-robin spread. --- src/trailing_edge/core/http.py | 6 +- src/trailing_edge/scrapers/kap/insider.py | 84 ++++++++++++++------ tests/unit/core/test_proxy_rotation.py | 27 ++++--- tests/unit/scrapers/kap/test_ingest_batch.py | 79 ++++++++++++++++++ 4 files changed, 159 insertions(+), 37 deletions(-) create mode 100644 tests/unit/scrapers/kap/test_ingest_batch.py diff --git a/src/trailing_edge/core/http.py b/src/trailing_edge/core/http.py index 3918daf..2e1bdd9 100644 --- a/src/trailing_edge/core/http.py +++ b/src/trailing_edge/core/http.py @@ -170,7 +170,11 @@ def _acquire(self) -> int | None: _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 - self._idx = i + # 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 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 index dc9a4a9..49ec1d8 100644 --- a/tests/unit/core/test_proxy_rotation.py +++ b/tests/unit/core/test_proxy_rotation.py @@ -162,11 +162,13 @@ async def h(*a, **k): for i, c in enumerate(client._clients): monkeypatch.setattr(c, "request", make(i)) - # IP0 serves 2, hits pace, parks; IP1 serves the next 2; IP2 the next 2 - no sleeps + # 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, 0, 1, 1, 2, 2] + assert seen == [0, 1, 2, 0, 1, 2] assert slept == [], "a pool must rotate on pace, never sleep while a fresh IP exists" @@ -216,9 +218,10 @@ async def boom(*a, **k): @pytest.mark.asyncio -async def test_a_healthy_ip_is_drained_not_round_robined(monkeypatch): - """With the pacer off, requests are sequential and a healthy IP just keeps serving - - it is drained, not spread across the pool one request at a time.""" +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() @@ -234,16 +237,17 @@ async def h(*a, **k): for i, c in enumerate(client._clients): monkeypatch.setattr(c, "request", make(i)) - for _ in range(3): + for _ in range(6): await client.get("https://kap.org.tr/x") - assert seen == [0, 0, 0] # drained, not spread + assert seen == [0, 1, 2, 0, 1, 2] # round-robin across the pool @pytest.mark.asyncio -async def test_rotation_advances_only_when_the_current_ip_blocks(monkeypatch): - """IP0 serves once, then blocks; the retry lands on IP1, which serves the rest. - Sequence: 0 (ok), 0 (block -> rotate), 1 (ok), 1 (ok).""" +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() @@ -268,4 +272,5 @@ async def ip1(*a, **k): for _ in range(3): await client.get("https://kap.org.tr/x") - assert seen == [0, 0, 1, 1] + 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" From 89fa78c337639d26bdf22ef9870456f6f7198101 Mon Sep 17 00:00:00 2001 From: cagan Date: Mon, 13 Jul 2026 23:11:24 +0300 Subject: [PATCH 11/11] result: the regime question, answered on the full 2015-2026 backfill MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The backfill is complete - 139/139 months SUCCESS, 18,415 disclosures, 2,339 clusters, more than double the 1,079 the earlier analysis ran on. That was collected to answer the one question left open: does the insider-cluster result hold across regimes? It does not, and the answer is sharper than "not tradeable". Split by the cluster's own era (scripts/regime_split.py): regime 20d gross t(gross) 20d net 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 In 2015-2018 the gross abnormal return was real and strong and died only to the spread. In 2021-2026 it has decayed to nothing (t=-1.0, indistinguishable from zero before costs). The pooled full-sample figure (+1.58% gross at 20d, t=5.2, still EDGE_DETECTED) is carried entirely by the early era and is misleading alone. 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. Full-sample net-of-cost also refreshed and is sharper than the earlier cut: net negative at every horizon, significant even at 60 days (N=2,279: 20d gross +1.58%, net -2.61%, t=-7.76; median round trip 2.33%). README and METHODOLOGY updated: the regime item in §6 moves from "still open" to closed with the decay finding; the ingest and quarantine items are updated to their final state (139/139; 3,795 of 18,415 quarantined, ~75% scanned images). scripts/regime_split.py makes the split reproducible. --- README.md | 58 ++++++++++------- docs/METHODOLOGY.md | 137 +++++++++++++++++++++++++--------------- scripts/regime_split.py | 134 +++++++++++++++++++++++++++++++++++++++ 3 files changed, 255 insertions(+), 74 deletions(-) create mode 100644 scripts/regime_split.py 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 3d7791a..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,57 +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. - -**The ingest is not complete, and the completion metric is not the obvious one.** 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` therefore overstates progress badly - -at the time of writing it read 92 of 139 months while `scraper_runs` held **20 SUCCESS, 82 -PARTIAL, 44 FAILED**. 2016 and 2020 contained not one complete month. - -The loss is recoverable: `backfill_kap_insider.py` replays non-SUCCESS months, fresh months -first and WAF stragglers last. But its `todo` list is computed once at startup, so a month -that goes PARTIAL *during* a run is not swept by that same run - the script must simply be -re-run until the SUCCESS count stops rising. **The completion criterion is -`count(SUCCESS) ≈ 139`, never the log's chunk count.** - -Direction of the error: WAF drops are unrelated to what a filing predicts, so this should -thin the sample rather than tilt it. That is an argument, not a measurement, and it is -recorded as such. - -**~1,200 filings are quarantined, and about three quarters of them are permanently -unreadable.** A DKB PDF that yields no transactions is quarantined to -`reports/parse_failures/` rather than silently skipped. Sampling those files, **75% carry no -extractable text at all** - they are scanned images with no text layer, and no parser will -ever read them. They are a hard ceiling of the same kind as the 250,000 TRY disclosure -threshold: a property of the source, not a bug to be fixed. Recovering them would need OCR, +## 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 they fail for a different reason: they are -**narrative filings with no table at all** ("... 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. The final quarantine rate must be recomputed once -the ingest is complete, because the current per-year rates (13% in 2015 rising to 32% in -2020) are measured against a denominator that is itself incomplete. +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/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()