mirror of
https://github.com/vladkens/twscrape.git
synced 2026-10-10 15:17:19 -04:00
fix: cap NetworkError retries per account, backoff then rotate (#325)
Co-authored-by: vladkens <[email protected]>
This commit is contained in:
+152
-5
@@ -94,9 +94,16 @@ async def test_switch_acc_on_http_error(client_fixture: CF):
|
||||
assert locked1 == locked3
|
||||
|
||||
|
||||
async def test_retry_with_same_acc_on_network_error(client_fixture: CF):
|
||||
async def test_retry_with_same_acc_on_network_error(client_fixture: CF, monkeypatch):
|
||||
pool, client, mock = client_fixture
|
||||
|
||||
sleeps = []
|
||||
|
||||
async def fake_sleep(secs):
|
||||
sleeps.append(secs)
|
||||
|
||||
monkeypatch.setattr("twscrape.queue_client.asyncio.sleep", fake_sleep)
|
||||
|
||||
await client.__aenter__()
|
||||
locked1 = await get_locked(pool)
|
||||
assert len(locked1) == 1
|
||||
@@ -107,6 +114,7 @@ async def test_retry_with_same_acc_on_network_error(client_fixture: CF):
|
||||
rep = await client.get(URL)
|
||||
assert rep is not None
|
||||
assert rep.json() == {"foo": "2"}
|
||||
assert sleeps == [2]
|
||||
|
||||
assert await get_locked(pool) == locked1
|
||||
|
||||
@@ -114,6 +122,62 @@ async def test_retry_with_same_acc_on_network_error(client_fixture: CF):
|
||||
assert username is not None
|
||||
|
||||
|
||||
async def test_network_error_rotates_account_after_3_failures(client_fixture: CF, monkeypatch):
|
||||
pool, client, mock = client_fixture
|
||||
|
||||
sleeps = []
|
||||
|
||||
async def fake_sleep(secs):
|
||||
sleeps.append(secs)
|
||||
|
||||
monkeypatch.setattr("twscrape.queue_client.asyncio.sleep", fake_sleep)
|
||||
|
||||
await client.__aenter__()
|
||||
assert await get_locked(pool) == {"user1"}
|
||||
|
||||
for _ in range(3):
|
||||
mock.add_exception(NetworkError("timeout"))
|
||||
mock.add_response(json={"ok": True})
|
||||
|
||||
rep = await client.get(URL)
|
||||
assert rep is not None
|
||||
assert rep.json() == {"ok": True}
|
||||
assert getattr(rep, "__username", None) == "user2"
|
||||
assert sleeps == [2, 4]
|
||||
|
||||
# user1 is short-locked (~60s transport lock, not the 15-min unknown-error lock)
|
||||
user1 = next(x for x in await pool.get_all() if x.username == "user1")
|
||||
lock_secs = (user1.locks["SearchTimeline"] - utc.now()).total_seconds()
|
||||
assert 0 < lock_secs <= 61
|
||||
|
||||
await client.__aexit__(None, None, None)
|
||||
|
||||
|
||||
async def test_network_error_counter_resets_on_account_change(client_fixture: CF, monkeypatch):
|
||||
pool, client, mock = client_fixture
|
||||
|
||||
async def fake_sleep(secs):
|
||||
pass
|
||||
|
||||
monkeypatch.setattr("twscrape.queue_client.asyncio.sleep", fake_sleep)
|
||||
|
||||
await client.__aenter__()
|
||||
|
||||
# 3 failures rotate user1 -> user2, which must get its own 3 tries: if the
|
||||
# counter carried over, the first failure on user2 would rotate it too and
|
||||
# exhaust the pool (request would return None)
|
||||
for _ in range(5):
|
||||
mock.add_exception(NetworkError("timeout"))
|
||||
mock.add_response(json={"ok": True})
|
||||
|
||||
rep = await client.get(URL)
|
||||
assert rep is not None
|
||||
assert rep.json() == {"ok": True}
|
||||
assert getattr(rep, "__username", None) == "user2"
|
||||
|
||||
await client.__aexit__(None, None, None)
|
||||
|
||||
|
||||
async def test_ctx_closed_on_break(client_fixture: CF):
|
||||
pool, client, mock = client_fixture
|
||||
|
||||
@@ -199,22 +263,48 @@ async def test_queue_client_passes_effective_proxy_to_xclid(pool_mock: AccountsP
|
||||
# --- ConnectError ---
|
||||
|
||||
|
||||
async def test_connect_error_raises_after_3_retries(client_fixture: CF):
|
||||
async def test_connect_error_cools_account_and_rotates_after_3_retries(
|
||||
client_fixture: CF, monkeypatch
|
||||
):
|
||||
pool, client, mock = client_fixture
|
||||
|
||||
sleeps = []
|
||||
|
||||
async def fake_sleep(secs):
|
||||
sleeps.append(secs)
|
||||
|
||||
monkeypatch.setattr("twscrape.queue_client.asyncio.sleep", fake_sleep)
|
||||
|
||||
await client.__aenter__()
|
||||
assert await get_locked(pool) == {"user1"}
|
||||
|
||||
mock.add_exception(ConnectError("refused"))
|
||||
mock.add_exception(ConnectError("refused"))
|
||||
mock.add_exception(ConnectError("refused"))
|
||||
mock.add_response(json={"ok": True})
|
||||
|
||||
with pytest.raises(ConnectError):
|
||||
await client.get(URL)
|
||||
rep = await client.get(URL)
|
||||
assert rep is not None
|
||||
assert rep.json() == {"ok": True}
|
||||
assert getattr(rep, "__username", None) == "user2"
|
||||
assert sleeps == [2, 4]
|
||||
|
||||
# user1 got the short transport cooldown lock, not the 15min unknown-error lock
|
||||
user1 = next(x for x in await pool.get_all() if x.username == "user1")
|
||||
lock_secs = (user1.locks["SearchTimeline"] - utc.now()).total_seconds()
|
||||
assert 0 < lock_secs <= 61
|
||||
|
||||
await client.__aexit__(None, None, None)
|
||||
|
||||
|
||||
async def test_connect_error_recovers_before_3_retries(client_fixture: CF):
|
||||
async def test_connect_error_recovers_before_3_retries(client_fixture: CF, monkeypatch):
|
||||
pool, client, mock = client_fixture
|
||||
|
||||
async def fake_sleep(secs):
|
||||
pass
|
||||
|
||||
monkeypatch.setattr("twscrape.queue_client.asyncio.sleep", fake_sleep)
|
||||
|
||||
await client.__aenter__()
|
||||
|
||||
mock.add_exception(ConnectError("refused"))
|
||||
@@ -224,6 +314,63 @@ async def test_connect_error_recovers_before_3_retries(client_fixture: CF):
|
||||
rep = await client.get(URL)
|
||||
assert rep is not None
|
||||
assert rep.json() == {"ok": True}
|
||||
assert getattr(rep, "__username", None) == "user1"
|
||||
|
||||
await client.__aexit__(None, None, None)
|
||||
|
||||
|
||||
async def test_alternating_categories_trip_total_safety_net(client_fixture: CF, monkeypatch):
|
||||
pool, client, mock = client_fixture
|
||||
|
||||
async def fake_sleep(secs):
|
||||
pass
|
||||
|
||||
monkeypatch.setattr("twscrape.queue_client.asyncio.sleep", fake_sleep)
|
||||
|
||||
await client.__aenter__()
|
||||
assert await get_locked(pool) == {"user1"}
|
||||
|
||||
# neither category alone reaches its own limit (3), but alternating between
|
||||
# them should still trip the combined total-failure safety net (4)
|
||||
mock.add_exception(ConnectError("refused"))
|
||||
mock.add_exception(RuntimeError("boom"))
|
||||
mock.add_exception(ConnectError("refused"))
|
||||
mock.add_exception(RuntimeError("boom"))
|
||||
mock.add_response(json={"ok": True})
|
||||
|
||||
rep = await client.get(URL)
|
||||
assert rep is not None
|
||||
assert getattr(rep, "__username", None) == "user2"
|
||||
|
||||
user1 = next(x for x in await pool.get_all() if x.username == "user1")
|
||||
assert "SearchTimeline" in user1.locks
|
||||
|
||||
await client.__aexit__(None, None, None)
|
||||
|
||||
|
||||
async def test_transport_error_retry_budget_is_per_account(client_fixture: CF, monkeypatch):
|
||||
pool, client, mock = client_fixture
|
||||
|
||||
async def fake_sleep(secs):
|
||||
pass
|
||||
|
||||
monkeypatch.setattr("twscrape.queue_client.asyncio.sleep", fake_sleep)
|
||||
|
||||
await client.__aenter__()
|
||||
|
||||
# 3 failures rotate user1 -> user2, which must get its own fresh budget: if the
|
||||
# counter carried over, the first failure on user2 would immediately rotate it
|
||||
# too and exhaust the whole (2-account) pool.
|
||||
for _ in range(5):
|
||||
mock.add_exception(NetworkError("timeout"))
|
||||
mock.add_response(json={"ok": True})
|
||||
|
||||
rep = await client.get(URL)
|
||||
assert rep is not None
|
||||
assert rep.json() == {"ok": True}
|
||||
assert getattr(rep, "__username", None) == "user2"
|
||||
|
||||
await client.__aexit__(None, None, None)
|
||||
|
||||
await client.__aexit__(None, None, None)
|
||||
|
||||
|
||||
+37
-23
@@ -1,6 +1,7 @@
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
from enum import Enum, auto
|
||||
from typing import Any
|
||||
from urllib.parse import urlparse
|
||||
|
||||
@@ -34,6 +35,11 @@ class GqlFeaturesOutdatedError(AbortReqError):
|
||||
"""GQL_FEATURES in api.py no longer matches the X API. Retrying cannot help."""
|
||||
|
||||
|
||||
class FailKind(Enum):
|
||||
TRANSPORT = auto()
|
||||
UNKNOWN = auto()
|
||||
|
||||
|
||||
class XClIdGenStore:
|
||||
items: dict[str, XClIdGen] = {}
|
||||
|
||||
@@ -59,6 +65,13 @@ class Ctx:
|
||||
self.acc = acc
|
||||
self.clt = clt
|
||||
self.proxy = proxy
|
||||
self.fails = {FailKind.TRANSPORT: 0, FailKind.UNKNOWN: 0}
|
||||
|
||||
def fail(self, kind: FailKind) -> bool:
|
||||
"""Count a failed attempt of this kind, return whether it's still worth retrying."""
|
||||
fail_limit, total_fail_limit = 3, 4
|
||||
self.fails[kind] += 1
|
||||
return self.fails[kind] < fail_limit and sum(self.fails.values()) < total_fail_limit
|
||||
|
||||
async def aclose(self):
|
||||
await self.clt.aclose()
|
||||
@@ -276,10 +289,10 @@ class QueueClient:
|
||||
return await self.req("GET", url, params=params)
|
||||
|
||||
async def req(self, method: HttpMethod, url: str, params: ReqParams = None) -> Response | None:
|
||||
unknown_retry, connection_retry = 0, 0
|
||||
|
||||
while True:
|
||||
ctx = await self._get_ctx() # not need to close client, class implements __aexit__
|
||||
# 1. same ctx until _close_ctx() clears it — that's retry vs rotate
|
||||
# 2. no aclose() needed here, __aexit__ handles it
|
||||
ctx = await self._get_ctx()
|
||||
if ctx is None:
|
||||
return None
|
||||
|
||||
@@ -306,7 +319,6 @@ class QueueClient:
|
||||
await self._check_rep(rep)
|
||||
|
||||
ctx.req_count += 1 # count only successful
|
||||
unknown_retry, connection_retry = 0, 0
|
||||
return rep
|
||||
except GqlFeaturesOutdatedError:
|
||||
# structurally invalid request, retrying cannot help — let the caller see it
|
||||
@@ -328,23 +340,25 @@ class QueueClient:
|
||||
)
|
||||
await self._close_ctx()
|
||||
return None
|
||||
except NetworkError:
|
||||
# http transport failed, just retry with same account
|
||||
continue
|
||||
except ConnectError as e:
|
||||
# if proxy misconfigured or host unreachable
|
||||
connection_retry += 1
|
||||
if connection_retry >= 3:
|
||||
raise e
|
||||
except Exception as e:
|
||||
unknown_retry += 1
|
||||
if unknown_retry >= 3:
|
||||
msg = [
|
||||
"Unknown error. Account timeouted for 15 minutes.",
|
||||
"Create issue please: https://github.com/vladkens/twscrape/issues",
|
||||
"If it mistake, you can unlock accounts with `twscrape reset_locks`. "
|
||||
f"Err: {self._format_ctx_error(ctx, e)}",
|
||||
]
|
||||
except (NetworkError, ConnectError) as e:
|
||||
# transport failed, retry same account with backoff, then cool it down and rotate
|
||||
if ctx.fail(FailKind.TRANSPORT):
|
||||
await asyncio.sleep(2 ** ctx.fails[FailKind.TRANSPORT])
|
||||
continue
|
||||
|
||||
logger.warning(" ".join(msg))
|
||||
await self._close_ctx(utc.ts() + 60 * 15) # 15 minutes
|
||||
logger.warning(f"{self._format_ctx_error(ctx, e)}; cooling account for 60s")
|
||||
await self._close_ctx(utc.ts() + 60)
|
||||
continue
|
||||
except Exception as e:
|
||||
if ctx.fail(FailKind.UNKNOWN):
|
||||
continue
|
||||
|
||||
msg = [
|
||||
"Unknown error. Account timeouted for 15 minutes.",
|
||||
"Create issue please: https://github.com/vladkens/twscrape/issues",
|
||||
"If it mistake, you can unlock accounts with `twscrape reset_locks`. "
|
||||
f"Err: {self._format_ctx_error(ctx, e)}",
|
||||
]
|
||||
|
||||
logger.warning(" ".join(msg))
|
||||
await self._close_ctx(utc.ts() + 60 * 15) # 15 minutes
|
||||
|
||||
Reference in New Issue
Block a user