ZaStoGram_desktop/Telegram/SourceFiles/tests/test_mtp_instance_request_registry.py
loop-uh fa9ef983dc Fix broker deadlock and over-broad cooldown reset from review
Self-review of the mtproxy scheduler rework surfaced three defects:

- Deadlock: claimFront popped dead queue fronts while holding _mutex, and
  ~RequestState can release an active admission lease, which now fires the
  broker's own admission-release listener synchronously -> wakeEndpoint
  -> relock of the non-recursive _mutex. Dead fronts are now moved out and
  destroyed after the lock is released (matching the cancel* idiom).

- Over-broad penalty reset: noteMtproxyEndpointSelected fired on every
  proxyChanges emission, including automatic rotation and blanket
  connection restarts (IPv6 toggle), so a dead proxy's cooldown ladder
  never escalated. A manual flag now threads Application::setCurrentProxy
  -> ProxyChange -> Instance::migrateProxy so only genuine user selection
  clears the penalty.

- Hardening: clear wakeScheduled under _mutex (it was read under the lock
  from other threads), and re-drain after cancelling a request so one
  queued behind a denied-and-cancelled front is not stranded until the
  slow poll.
2026-07-11 21:11:03 +03:00

147 lines
4.8 KiB
Python

from pathlib import Path
SOURCE_DIR = Path(__file__).resolve().parents[1]
ROOT = SOURCE_DIR.parents[1]
CMAKE = ROOT / "Telegram" / "CMakeLists.txt"
MTPROTO_DIR = SOURCE_DIR / "mtproto"
INSTANCE_DIR = MTPROTO_DIR / "instance"
INSTANCE_CPP = INSTANCE_DIR / "mtp_instance.cpp"
INSTANCE_H = INSTANCE_DIR / "mtp_instance.h"
REQUEST_REGISTRY_H = INSTANCE_DIR / "request_registry.h"
REQUEST_REGISTRY_CPP = INSTANCE_DIR / "request_registry.cpp"
SESSION_DELEGATE_H = MTPROTO_DIR / "session" / "session_delegate.h"
def read(path):
assert path.exists(), f"missing expected source file: {path}"
return path.read_text(encoding="utf-8")
def function_body(source, signature):
start = source.index(signature)
brace = source.index("{", start)
depth = 0
for index in range(brace, len(source)):
char = source[index]
if char == "{":
depth += 1
elif char == "}":
depth -= 1
if depth == 0:
return source[brace + 1:index]
raise AssertionError(f"function body not found: {signature}")
def test_request_registry_sources_are_registered():
cmake = read(CMAKE)
assert "mtproto/instance/request_registry.cpp" in cmake
assert "mtproto/instance/request_registry.h" in cmake
assert read(REQUEST_REGISTRY_H)
assert read(REQUEST_REGISTRY_CPP)
def test_instance_private_delegates_request_bookkeeping():
instance = read(INSTANCE_CPP)
private_body = function_body(instance, "class Instance::Private")
assert '#include "mtproto/instance/request_registry.h"' in instance
assert "RequestRegistry _requests;" in private_body
for moved in (
"_requestsByDc",
"_requestByDcLock",
"_parserMap",
"_parserMapLock",
"_requestMap",
"_requestMapLock",
"_delayedRequests",
"_dependentRequests",
"_dependentRequestsLock",
"_requestsDelays"):
assert moved not in private_body
for kept in (
"_authExportRequests",
"_authWaiters",
"_badGuestDcRequests"):
assert kept in private_body
def test_request_registry_owns_rpc_state_and_locks():
header = read(REQUEST_REGISTRY_H)
source = read(REQUEST_REGISTRY_CPP)
assert "class RequestRegistry final" in header
for field in (
"_requestsByDc",
"_requestByDcLock",
"_parserMap",
"_parserMapLock",
"_requestMap",
"_requestMapLock",
"_dependentRequests",
"_dependentRequestsLock",
"_delayedRequests",
"_requestsDelays"):
assert field in header
for method in (
"storeRequest(",
"registerRequest(",
"queryDc(",
"changeDc(",
"request(",
"hasCallback(",
"takeCallback(",
"restoreCallback(",
"unregisterRequest(",
"prepareDependency(",
"nextBackoffSeconds(",
"scheduleDelayed(",
"takeReadyDelayed(",
"nextDelayedAt("):
assert method in header
assert f"RequestRegistry::{method}" in source
def test_public_instance_and_session_delegate_surfaces_stay_stable():
instance_header = read(INSTANCE_H)
delegate = read(SESSION_DELEGATE_H)
thread_safe_block = instance_header.split(
"\t// Thread-safe.", 1)[1].split("\t// Main thread.", 1)[0]
main_thread_block = instance_header.split(
"\tvoid addKeysForDestroy(AuthKeysList &&keys);", 1)[1].split(
"\tvoid killSession", 1)[0]
assert "request_registry" not in instance_header
assert "class RequestRegistry" not in instance_header
assert "sendSerialized(" in instance_header
assert "sendProtocolMessage(" in instance_header
for method in (
"void restart();",
"void restart(ShiftedDcId shiftedDcId);",
"void migrateProxy(bool manual = true);",
"int32 dcstate(ShiftedDcId shiftedDcId = 0);",
"QString dctransport(ShiftedDcId shiftedDcId = 0);",
"ConnectionStatus &connectionStatus() const;",
"void ping();",
"void cancel(mtpRequestId requestId);",
"int32 state(mtpRequestId requestId);"):
assert method not in thread_safe_block
assert method in main_thread_block
for callback in (
"hasCallback(mtpRequestId requestId) const",
"processCallback(const Response &response)",
"processUpdate(const Response &message)"):
assert callback in delegate
if __name__ == "__main__":
test_request_registry_sources_are_registered()
test_instance_private_delegates_request_bookkeeping()
test_request_registry_owns_rpc_state_and_locks()
test_public_instance_and_session_delegate_surfaces_stay_stable()