mirror of
https://github.com/vladkens/twscrape.git
synced 2026-10-10 15:17:19 -04:00
feat: self-heal missing GQL features instead of aborting (#339)
Co-authored-by: vladkens <[email protected]>
This commit is contained in:
@@ -1,3 +1,4 @@
|
||||
import json
|
||||
from collections import OrderedDict
|
||||
from contextlib import aclosing
|
||||
|
||||
@@ -6,8 +7,13 @@ import pytest
|
||||
import twscrape.queue_client as queue_client_module
|
||||
from twscrape.account import Account
|
||||
from twscrape.accounts_pool import AccountsPool
|
||||
from twscrape.http import ConnectError, NetworkError
|
||||
from twscrape.queue_client import GqlFeaturesOutdatedError, QueueClient, XClIdGenStore
|
||||
from twscrape.http import ConnectError, HttpMethod, NetworkError, Response
|
||||
from twscrape.queue_client import (
|
||||
GqlFeaturesOutdatedError,
|
||||
QueueClient,
|
||||
ReqParams,
|
||||
XClIdGenStore,
|
||||
)
|
||||
from twscrape.utils import utc
|
||||
from twscrape.xclid import XClIdAccountError, XClIdGen, XClIdParseError
|
||||
|
||||
@@ -552,6 +558,86 @@ async def test_gql_features_outdated_raises(client_fixture: CF):
|
||||
assert await get_locked(pool) == set()
|
||||
|
||||
|
||||
async def test_gql_features_self_healed_and_retried(client_fixture: CF, monkeypatch):
|
||||
pool, client, mock = client_fixture
|
||||
await client.__aenter__()
|
||||
|
||||
sent: list[dict] = []
|
||||
original = mock.request
|
||||
|
||||
async def spy(method: HttpMethod, url: str, **kwargs) -> Response:
|
||||
sent.append(dict(kwargs.get("params") or {}))
|
||||
return await original(method, url, **kwargs)
|
||||
|
||||
monkeypatch.setattr(mock, "request", spy)
|
||||
|
||||
mock.add_response(
|
||||
json={
|
||||
"errors": [{"code": 336, "message": "The following features cannot be null: foo, bar"}]
|
||||
}
|
||||
)
|
||||
mock.add_response(json={"data": {"user": {"id": 1}}})
|
||||
|
||||
params: ReqParams = {"variables": "{}", "features": json.dumps({"baz": True})}
|
||||
rep = await client.get(URL, params=params)
|
||||
assert rep is not None
|
||||
|
||||
# the flags X named were added and the request replayed
|
||||
assert len(sent) == 2
|
||||
assert json.loads(sent[0]["features"]) == {"baz": True}
|
||||
assert json.loads(sent[1]["features"]) == {"baz": True, "foo": True, "bar": True}
|
||||
|
||||
await client.__aexit__(None, None, None)
|
||||
|
||||
# the account must survive: not deactivated, not left locked
|
||||
assert await get_inactive(pool) == set()
|
||||
assert await get_locked(pool) == set()
|
||||
|
||||
|
||||
async def test_gql_features_self_heal_persists_across_requests(client_fixture: CF, monkeypatch):
|
||||
_pool, client, mock = client_fixture
|
||||
await client.__aenter__()
|
||||
|
||||
sent = []
|
||||
original = mock.request
|
||||
|
||||
async def spy(method: HttpMethod, url: str, **kwargs) -> Response:
|
||||
sent.append(json.loads(kwargs["params"]["features"]))
|
||||
return await original(method, url, **kwargs)
|
||||
|
||||
monkeypatch.setattr(mock, "request", spy)
|
||||
|
||||
mock.add_response(
|
||||
json={"errors": [{"code": 336, "message": "The following features cannot be null: foo"}]}
|
||||
)
|
||||
mock.add_response(json={"data": {"user": {"id": 1}}})
|
||||
mock.add_response(json={"data": {"user": {"id": 1}}})
|
||||
|
||||
assert await client.get(URL, params={"variables": "{}", "features": "{}"}) is not None
|
||||
assert await client.get(URL, params={"variables": "{}", "features": "{}"}) is not None
|
||||
assert sent == [{}, {"foo": True}, {"foo": True}]
|
||||
|
||||
await client.__aexit__(None, None, None)
|
||||
|
||||
|
||||
async def test_gql_features_self_heal_retries_only_once(client_fixture: CF):
|
||||
_pool, client, mock = client_fixture
|
||||
await client.__aenter__()
|
||||
|
||||
# X keeps rejecting even with the flags set — give up rather than loop
|
||||
for _ in range(2):
|
||||
mock.add_response(
|
||||
json={
|
||||
"errors": [{"code": 336, "message": "The following features cannot be null: foo"}]
|
||||
}
|
||||
)
|
||||
|
||||
with pytest.raises(GqlFeaturesOutdatedError):
|
||||
await client.get(URL, params={"variables": "{}", "features": "{}"})
|
||||
|
||||
await client.__aexit__(None, None, None)
|
||||
|
||||
|
||||
async def test_gql_features_outdated_raises_when_not_first(client_fixture: CF):
|
||||
_pool, client, mock = client_fixture
|
||||
await client.__aenter__()
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
from enum import Enum, auto
|
||||
from typing import Any
|
||||
from urllib.parse import urlparse
|
||||
@@ -22,6 +23,7 @@ from .utils import utc
|
||||
from .xclid import XClIdAccountError, XClIdGen, XClIdParseError
|
||||
|
||||
ReqParams = dict[str, str | int] | None
|
||||
MISSING_FEATURES_RE = re.compile(r"The following features cannot be null: ([^;]+)")
|
||||
TMP_TS = utc.now().isoformat().split(".")[0].replace("T", "_").replace(":", "-")[0:16]
|
||||
|
||||
|
||||
@@ -32,7 +34,7 @@ class AbortReqError(Exception): ...
|
||||
|
||||
|
||||
class GqlFeaturesOutdatedError(AbortReqError):
|
||||
"""GQL_FEATURES in api.py no longer matches the X API. Retrying cannot help."""
|
||||
"""GQL_FEATURES in api.py no longer matches the X API and self-healing did not resolve it."""
|
||||
|
||||
|
||||
class FailKind(Enum):
|
||||
@@ -134,6 +136,35 @@ def has_error(errors: list[str], prefix: str) -> bool:
|
||||
return any(error.startswith(prefix) for error in errors)
|
||||
|
||||
|
||||
def parse_missing_features(msg: str) -> list[str]:
|
||||
"""Feature flags X named in a (336) error: "...cannot be null: a, b"."""
|
||||
match = MISSING_FEATURES_RE.search(msg)
|
||||
return [x.strip() for x in match.group(1).split(",") if x.strip()] if match else []
|
||||
|
||||
|
||||
def add_missing_features(params: ReqParams, missing: list[str]) -> bool:
|
||||
"""Set `missing` to true in params["features"]. False if there is nothing to patch."""
|
||||
if not isinstance(params, dict) or "features" not in params:
|
||||
return False
|
||||
|
||||
features = params["features"]
|
||||
if isinstance(features, str): # encode_params() has already run
|
||||
try:
|
||||
features = json.loads(features)
|
||||
except json.JSONDecodeError:
|
||||
return False
|
||||
|
||||
if not isinstance(features, dict):
|
||||
return False
|
||||
|
||||
if all(features.get(x) is True for x in missing):
|
||||
return False # already true, so a retry would just loop
|
||||
|
||||
features.update(dict.fromkeys(missing, True))
|
||||
params["features"] = json.dumps(features, separators=(",", ":"))
|
||||
return True
|
||||
|
||||
|
||||
def dump_rep(rep: Response):
|
||||
count = getattr(dump_rep, "__count", -1) + 1
|
||||
setattr(dump_rep, "__count", count)
|
||||
@@ -168,6 +199,7 @@ class QueueClient:
|
||||
self.debug = debug
|
||||
self.ctx: Ctx | None = None
|
||||
self.proxy = proxy
|
||||
self._healed_features: set[str] = set()
|
||||
|
||||
async def __aenter__(self):
|
||||
await self._get_ctx()
|
||||
@@ -312,6 +344,9 @@ class QueueClient:
|
||||
return await self.req("GET", url, params=params)
|
||||
|
||||
async def req(self, method: HttpMethod, url: str, params: ReqParams = None) -> Response | None:
|
||||
if self._healed_features:
|
||||
add_missing_features(params, list(self._healed_features))
|
||||
features_retried = False
|
||||
while True:
|
||||
# 1. same ctx until _close_ctx() clears it — that's retry vs rotate
|
||||
# 2. no aclose() needed here, __aexit__ handles it
|
||||
@@ -343,9 +378,21 @@ class QueueClient:
|
||||
|
||||
ctx.req_count += 1 # count only successful
|
||||
return rep
|
||||
except GqlFeaturesOutdatedError:
|
||||
# structurally invalid request, retrying cannot help — let the caller see it
|
||||
raise
|
||||
except GqlFeaturesOutdatedError as e:
|
||||
# X names the flags it wants, so set them and retry once. GQL_FEATURES still
|
||||
# needs a permanent update — the warning says so — but a stale dict shouldn't
|
||||
# take every caller down until the next release.
|
||||
missing = parse_missing_features(str(e))
|
||||
if features_retried or not add_missing_features(params, missing):
|
||||
raise
|
||||
|
||||
self._healed_features.update(missing)
|
||||
features_retried = True
|
||||
logger.warning(
|
||||
f"{self.queue}: GQL_FEATURES missing {missing}, set true and retrying. "
|
||||
"Update GQL_FEATURES in api.py to fix this permanently."
|
||||
)
|
||||
continue
|
||||
except AbortReqError:
|
||||
# abort all queries
|
||||
return
|
||||
|
||||
Reference in New Issue
Block a user