fix: partition smoke проверяет fencing при свежем auth и отказ записи старого владельца после обрыва БД
This commit is contained in:
parent
b1797786d2
commit
5f8dd98e39
1 changed files with 95 additions and 17 deletions
|
|
@ -174,31 +174,62 @@ def ws_connect(ws_base: str, path: str, cookie: str):
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
async def wait_for_fence(socket, timeout: float = 10) -> dict:
|
# Текст `_close_auth_state_unavailable` в backend/app/api/auth.py: 1013 с ним —
|
||||||
|
# закрытие от устаревшего кэша поколений, а не произвольный 1013.
|
||||||
|
AUTH_UNAVAILABLE_REASON = "Состояние доступа временно недоступно"
|
||||||
|
|
||||||
|
|
||||||
|
async def wait_auth_fresh(base: str, cookie: str, timeout: float = 10) -> bool:
|
||||||
|
"""Узел принял cookie: кэш поколений свежий и логин ему известен."""
|
||||||
|
request = urllib.request.Request(f"{base}/api/auth/me", headers={"Cookie": cookie})
|
||||||
|
deadline = time.monotonic() + timeout
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
with urllib.request.urlopen(request, timeout=2) as response:
|
||||||
|
if response.status == 200:
|
||||||
|
return True
|
||||||
|
except OSError:
|
||||||
|
pass
|
||||||
|
if time.monotonic() > deadline:
|
||||||
|
return False
|
||||||
|
await asyncio.sleep(0.2)
|
||||||
|
|
||||||
|
|
||||||
|
async def wait_for_fence(socket, timeout: float = 10, ack_id: str | None = None) -> dict:
|
||||||
|
acked = False
|
||||||
|
|
||||||
|
def outcome(**closure) -> dict:
|
||||||
|
if ack_id is not None:
|
||||||
|
closure["command_acked"] = acked
|
||||||
|
return closure
|
||||||
|
|
||||||
deadline = time.monotonic() + timeout
|
deadline = time.monotonic() + timeout
|
||||||
while time.monotonic() < deadline:
|
while time.monotonic() < deadline:
|
||||||
try:
|
try:
|
||||||
message = await asyncio.wait_for(socket.recv(), timeout=deadline - time.monotonic())
|
message = await asyncio.wait_for(socket.recv(), timeout=deadline - time.monotonic())
|
||||||
except TimeoutError:
|
except TimeoutError:
|
||||||
return {"closed": False, "fence_notice": False, "reason": "timeout"}
|
return outcome(closed=False, fence_notice=False, reason="timeout")
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
close = getattr(exc, "rcvd", None) or getattr(exc, "sent", None)
|
close = getattr(exc, "rcvd", None) or getattr(exc, "sent", None)
|
||||||
close_code = getattr(close, "code", None) or getattr(exc, "code", None)
|
close_code = getattr(close, "code", None) or getattr(exc, "code", None)
|
||||||
return {
|
return outcome(
|
||||||
"closed": True,
|
closed=True,
|
||||||
"fence_notice": close_code == 1012,
|
fence_notice=close_code == 1012,
|
||||||
"close_code": close_code,
|
close_code=close_code,
|
||||||
"reason": getattr(close, "reason", "") if close else str(exc),
|
reason=getattr(close, "reason", "") if close else str(exc),
|
||||||
}
|
)
|
||||||
if isinstance(message, str):
|
if isinstance(message, str):
|
||||||
try:
|
try:
|
||||||
event = json.loads(message)
|
event = json.loads(message)
|
||||||
except json.JSONDecodeError:
|
except json.JSONDecodeError:
|
||||||
continue
|
continue
|
||||||
|
if ack_id is not None and event.get("type") == "command.ack" \
|
||||||
|
and event.get("command_id") == ack_id:
|
||||||
|
acked = True
|
||||||
if event.get("message") == "Занятие передано другому backend-узлу; переподключитесь.":
|
if event.get("message") == "Занятие передано другому backend-узлу; переподключитесь.":
|
||||||
return {"closed": True, "fence_notice": True, "close_code": 1012,
|
return outcome(closed=True, fence_notice=True, close_code=1012,
|
||||||
"reason": event.get("message")}
|
reason=event.get("message"))
|
||||||
return {"closed": False, "fence_notice": False, "reason": "timeout"}
|
return outcome(closed=False, fence_notice=False, reason="timeout")
|
||||||
|
|
||||||
|
|
||||||
async def main(args: argparse.Namespace) -> int:
|
async def main(args: argparse.Namespace) -> int:
|
||||||
|
|
@ -253,7 +284,10 @@ async def main(args: argparse.Namespace) -> int:
|
||||||
# Учётки созданы прямо в БД, вход — на узле A, а сокет занятия Nginx
|
# Учётки созданы прямо в БД, вход — на узле A, а сокет занятия Nginx
|
||||||
# может отдать узлу B. B узнаёт новый логин при очередной сверке
|
# может отдать узлу B. B узнаёт новый логин при очередной сверке
|
||||||
# поколений (раз в секунду); до неё middleware считает cookie отозванной.
|
# поколений (раз в секунду); до неё middleware считает cookie отозванной.
|
||||||
await asyncio.sleep(2.5)
|
for port in (args.backend_port, args.backend_b_port):
|
||||||
|
for cookie in (instructor_cookie, trainee_cookie):
|
||||||
|
if not await wait_auth_fresh(f"http://127.0.0.1:{port}", cookie):
|
||||||
|
raise RuntimeError(f"backend :{port} не узнал новый логин за 10 с")
|
||||||
await send_control(ws_base, session_id, instructor_cookie, {
|
await send_control(ws_base, session_id, instructor_cookie, {
|
||||||
"type": "scenario.start", "scenario_id": args.scenario,
|
"type": "scenario.start", "scenario_id": args.scenario,
|
||||||
"trainee": trainee.name, "trainee_id": str(trainee.id),
|
"trainee": trainee.name, "trainee_id": str(trainee.id),
|
||||||
|
|
@ -404,14 +438,37 @@ async def main(args: argparse.Namespace) -> int:
|
||||||
)
|
)
|
||||||
compose(args, "stop", proxy_name)
|
compose(args, "stop", proxy_name)
|
||||||
fault_stopped = True
|
fault_stopped = True
|
||||||
result["boundary_command_sent_to_old_owner"] = "arrived"
|
# Реальная мутация старому владельцу уже без БД: commit не пройдёт,
|
||||||
|
# ACK не придёт, а id не должен попасть в checkpoint нового
|
||||||
|
# владельца — прямое доказательство запрета записи после потери БД.
|
||||||
|
# Отправка сразу, до `compose ps`: через 2 с auth закроет сокет.
|
||||||
|
boundary_command_id = str(uuid4())
|
||||||
|
boundary = {"command_id": boundary_command_id, "status": "working"}
|
||||||
|
try:
|
||||||
|
await active_sockets["station"].send(json.dumps({
|
||||||
|
"type": "card.status", "service": service, "status": "working",
|
||||||
|
"comment": (
|
||||||
|
"Основание: доклад по карточке.\n"
|
||||||
|
f"Сведения: начало работ после обрыва БД; {service} notified."
|
||||||
|
),
|
||||||
|
"_command_id": boundary_command_id,
|
||||||
|
}, ensure_ascii=False))
|
||||||
|
boundary["sent"] = True
|
||||||
|
except websockets.exceptions.WebSocketException as exc:
|
||||||
|
boundary["sent"] = False
|
||||||
|
boundary["error"] = f"{type(exc).__name__}: {exc}"
|
||||||
|
result["boundary_command_sent_to_old_owner"] = boundary
|
||||||
|
result["checks"]["boundary_command_sent_after_partition"] = boundary["sent"]
|
||||||
running_services = compose(args, "ps", "--status", "running", "--services").splitlines()
|
running_services = compose(args, "ps", "--status", "running", "--services").splitlines()
|
||||||
result["checks"]["owner_process_running"] = owner_service in running_services
|
result["checks"]["owner_process_running"] = owner_service in running_services
|
||||||
if owner_service not in running_services:
|
if owner_service not in running_services:
|
||||||
raise RuntimeError("owner process did not remain running during DB partition")
|
raise RuntimeError("owner process did not remain running during DB partition")
|
||||||
closures = await asyncio.gather(*(
|
closures = await asyncio.gather(*(
|
||||||
wait_for_fence(socket, timeout=args.takeover_timeout)
|
wait_for_fence(
|
||||||
for socket in active_sockets.values()
|
socket, timeout=args.takeover_timeout,
|
||||||
|
ack_id=boundary_command_id if channel == "station" else None,
|
||||||
|
)
|
||||||
|
for channel, socket in active_sockets.items()
|
||||||
))
|
))
|
||||||
result["existing_channel_closures"] = dict(zip(active_sockets, closures))
|
result["existing_channel_closures"] = dict(zip(active_sockets, closures))
|
||||||
result["existing_channel_fence_notices"] = {
|
result["existing_channel_fence_notices"] = {
|
||||||
|
|
@ -433,10 +490,16 @@ async def main(args: argparse.Namespace) -> int:
|
||||||
}
|
}
|
||||||
result["checks"].update({
|
result["checks"].update({
|
||||||
f"existing_{channel}_closed_fail_closed": (
|
f"existing_{channel}_closed_fail_closed": (
|
||||||
closure["fence_notice"] or closure.get("close_code") == 1013
|
closure["fence_notice"] or (
|
||||||
|
closure.get("close_code") == 1013
|
||||||
|
and closure.get("reason") == AUTH_UNAVAILABLE_REASON
|
||||||
|
)
|
||||||
)
|
)
|
||||||
for channel, closure in result["existing_channel_closures"].items()
|
for channel, closure in result["existing_channel_closures"].items()
|
||||||
})
|
})
|
||||||
|
result["checks"]["boundary_command_not_acked"] = (
|
||||||
|
result["existing_channel_closures"]["station"].get("command_acked") is False
|
||||||
|
)
|
||||||
|
|
||||||
deadline = time.monotonic() + args.takeover_timeout
|
deadline = time.monotonic() + args.takeover_timeout
|
||||||
takeover = None
|
takeover = None
|
||||||
|
|
@ -472,12 +535,24 @@ async def main(args: argparse.Namespace) -> int:
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
result["checks"]["old_owner_process_healthy"] = False
|
result["checks"]["old_owner_process_healthy"] = False
|
||||||
result["old_owner_health_error"] = f"{type(exc).__name__}: {exc}"
|
result["old_owner_health_error"] = f"{type(exc).__name__}: {exc}"
|
||||||
|
# Auth на старом узле снова свежий: 1013 до accept исключён, и 403
|
||||||
|
# на рукопожатии ниже остаётся только за fencing занятия (1012).
|
||||||
|
result["checks"]["old_owner_auth_fresh"] = await wait_auth_fresh(
|
||||||
|
old_backend, instructor_cookie,
|
||||||
|
)
|
||||||
stale = await probe_ws(
|
stale = await probe_ws(
|
||||||
old_backend.replace("http://", "ws://"),
|
old_backend.replace("http://", "ws://"),
|
||||||
f"/ws/control/{session_id}", instructor_cookie,
|
f"/ws/control/{session_id}", instructor_cookie,
|
||||||
)
|
)
|
||||||
result["stale_owner_control"] = stale
|
result["stale_owner_control"] = stale
|
||||||
result["checks"]["stale_owner_control_rejected"] = not stale["connected"]
|
result["checks"]["stale_owner_control_rejected"] = (
|
||||||
|
not stale["connected"] and "HTTP 403" in stale.get("error", "")
|
||||||
|
)
|
||||||
|
# БД вернулась к старому владельцу, но его команда после обрыва так
|
||||||
|
# и не должна появиться в checkpoint: fencing epoch отверг запись.
|
||||||
|
result["checks"]["boundary_command_absent_from_checkpoint"] = (
|
||||||
|
not await wait_for_command_checkpoint(boundary_command_id, timeout=1)
|
||||||
|
)
|
||||||
|
|
||||||
request = urllib.request.Request(
|
request = urllib.request.Request(
|
||||||
f"{args.frontend_url}/api/sessions/{session_id}",
|
f"{args.frontend_url}/api/sessions/{session_id}",
|
||||||
|
|
@ -539,6 +614,9 @@ async def main(args: argparse.Namespace) -> int:
|
||||||
result["checks"]["boundary_command_reconciled"] = recovered_status in {
|
result["checks"]["boundary_command_reconciled"] = recovered_status in {
|
||||||
"responding", "arrived",
|
"responding", "arrived",
|
||||||
}
|
}
|
||||||
|
if args.failure_mode == "partition":
|
||||||
|
# `working` после обрыва не применён: статус остался на lost-ack `arrived`.
|
||||||
|
result["checks"]["boundary_command_not_applied"] = recovered_status == "arrived"
|
||||||
result["dds_progress_after_takeover"] = []
|
result["dds_progress_after_takeover"] = []
|
||||||
if recovered_status == "responding":
|
if recovered_status == "responding":
|
||||||
recovered_snapshot = await station_command(ws_base, session_id, trainee_cookie, {
|
recovered_snapshot = await station_command(ws_base, session_id, trainee_cookie, {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue