Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 85 additions & 10 deletions scripts/ci/xdist_watchdog.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,27 @@
d'acceptance de l'issue : un run bloque doit echouer vite, en disant
pourquoi.

Mode 2 -- le garde lui-meme tuait un run sain (run 35276661841, PR #16240,
attempt 2, 2026-09-17 23:27Z) : en fin de parcours sous ``-q``, pytest
ecrit ses points de test SANS saut de ligne tant que la ligne de ~72
caracteres n'est pas pleine. Le fil de lecture, base sur ``readline``,
restait bloque sur le fragment sans ``\\n`` pendant que des octets vivants
traversaient le tube : verdict « silence 480 s ... ``[99%]`` ... workers
morts : gw2 », kill d'un run a 99 % sans echec -- et le flush d'EOF du
kill a laisse la preuve sur le log : une ligne partielle de 43 resultats
emis PENDANT la fenetre dite muette (23:19:06 -> 23:27:07). Le master
n'etait pas bloque, il finissait la queue (tests de queue lents, puis
re-execution du lot du worker remplace). Correctif : la fraicheur se
mesure desormais par OCTET lu (``os.read`` par chunks, lignes
reconstituees en interne), pas par ligne complete. Un run qui emet ne
peut plus etre tue ; un blocage reel (zero octet, signature originelle)
l'est toujours.

Une tolerance au pourcentage de fin de parcours (« [95%+] => grace ») a
ete ecarte a dessin : le blocage originel #16288 s'est AUSSI produit a
``[99%]`` -- le pourcentage ne discrimine rien, seul le flux d'octets le
fait.

Pistes rejetees, pour memoire (detail dans #16288) :

- ``pytest-timeout`` + ``--timeout`` : borne un TEST qui bloque ; ici
Expand Down Expand Up @@ -92,6 +113,26 @@ def __init__(self) -> None:
self.last_progress_line: str | None = None
self.dead_workers: list[str] = []
self.line_count = 0
# Forensique mode 2 : des octets sans aucune ligne complete (points
# de fin de parcours sans \n) sont le signe d'un run VIVANT -- le
# verdict doit pouvoir les citer apres coup.
self.byte_count = 0

def record_bytes(self, nbytes: int) -> None:
"""Fraicheur par OCTET, pas par ligne complete (mode 2).

Sous ``-q``, les points de test n'emportent pas de ``\\n`` avant
que la ligne de ~72 caracteres soit pleine : un fil base sur
``readline`` ne verrait rien pendant que le run finit sa queue
(run 35276661841 : 43 resultats emis pendant une fenetre que le
garde croyait muette). La boucle de surveillance appelle donc
ceci sur CHAQUE chunk lu du tube.
"""
if nbytes <= 0:
return
with self.lock:
self.last_output = time.monotonic()
self.byte_count += nbytes

def record(self, line: str) -> None:
now = time.monotonic()
Expand All @@ -112,13 +153,45 @@ def idle_for(self) -> float:
return time.monotonic() - self.last_output


# Taille d'un chunk de lecture du tube. Un seul os.read = un seul appel
# systeme : il retourne des le PREMIER octet disponible (jamais apres un
# remplissage complet ni un \n) -- c'est la propriete qui repare le mode 2.
PIPE_CHUNK_BYTES = 65536


def _pump(stream, state: _StreamState, echo) -> None:
"""Lit les lignes d'un pipe du fils, les horodate et les recopie."""
# Sentinelle en BYTES : le pipe est binaire (Popen sans text=True),
# et b"" == "" est faux -- un sentinelle str ferait boucler le fil
# a l'infini apres l'EOF, inondant l'echo de lignes vides.
for raw in iter(stream.readline, b""):
line = raw.decode("utf-8", errors="replace") if isinstance(raw, bytes) else raw
"""Recopie la sortie du fils en mesurant l'activite par OCTET.

On lit le tube par chunks bruts (``os.read`` sur le fd) et on
reconstitue les lignes en interne : la fraicheur (``record_bytes``)
est mise a jour pour CHAQUE chunk, le decoupage en lignes ne sert
qu'au bookkeeping (progression, workers morts) et a l'echo.

Pourquoi pas ``readline`` : en fin de parcours ``-q``, pytest ecrit
ses points SANS saut de ligne -- un readline resterait bloque sur le
fragment pendant que le run finit sa queue, et le garde tuerait un
vivant (mode 2, run 35276661841). Un fragment final sans ``\\n``
n'est recopie qu'a l'EOF (ou au kill), exactement comme la ligne
partielle de 43 resultats revelee par le flush d'EOF du run cite.
"""
fd = stream.fileno()
pending = b""
while True:
try:
chunk = os.read(fd, PIPE_CHUNK_BYTES)
except OSError:
break # tube ferme sous nos pieds : equivaut a l'EOF
if not chunk:
break
state.record_bytes(len(chunk))
pending += chunk
while b"\n" in pending:
raw, pending = pending.split(b"\n", 1)
line = raw.decode("utf-8", errors="replace") + "\n"
state.record(line)
echo(line)
if pending:
line = pending.decode("utf-8", errors="replace")
state.record(line)
echo(line)
try:
Expand Down Expand Up @@ -210,11 +283,13 @@ def _verdict_blocked(state: _StreamState, idle: float, idle_limit: float,
emit(f"##[error]{VERDICT_PREFIX}: BLOQUE -- silence de sortie depuis "
f"{idle:.0f} s (limite {idle_limit:.0f} s), mur du job non atteint")
emit(f"##[error]{VERDICT_PREFIX}: derniere progression pytest : "
f"\"{progress}\" ; {state.line_count} lignes emises au total ; "
f"wall du wrapper {wall:.0f} s")
f"\"{progress}\" ; {state.line_count} lignes ({state.byte_count} "
f"octets) emises au total ; wall du wrapper {wall:.0f} s")
emit(f"##[error]{VERDICT_PREFIX}: workers morts : {workers}")
emit(f"##[error]{VERDICT_PREFIX}: le master etait vivant mais n'attendait "
f"pas du travail -- signature #16288 ; kill du groupe de processus")
emit(f"##[error]{VERDICT_PREFIX}: zero octet emis pendant la fenetre "
f"(ni ligne ni fragment) -- le master etait vivant mais "
f"n'attendait pas du travail, signature #16288 ; kill du groupe "
f"de processus")


def main(argv: list[str] | None = None) -> int:
Expand Down
75 changes: 75 additions & 0 deletions scripts/tests/test_xdist_watchdog.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,15 @@
mur, et c'est ce qui fait lire un blocage comme un depassement.
4. **Arme des le demarrage** -- un enfant muet depuis sa naissance (hang de
collection) est aussi tue : l'armement ne depend pas d'une premiere ligne.
5. **La fraicheur se mesure par OCTET, pas par ligne** (mode 2, run
35276661841 / PR #16240) : en fin de parcours ``-q``, pytest emet ses
points SANS saut de ligne tant que la ligne de ~72 caracteres n'est pas
pleine. Un garde qui n'ecouterait que les lignes completes croit a un
silence de 8 min et tue un run a ``[99%]`` en train de finir -- le
flush d'EOF du kill avait laisse sur le log une ligne partielle de 43
resultats emis PENDANT la fenetre dite muette. Des octets vivants sans
``\\n`` doivent donc maintenir la fraicheur ; et un vrai blocage (zero
octet) doit toujours mourir.

Les enfants sont des ``python -c`` mono-processus : aucun xdist requis
(l'issue note le defaut propre a la classe de runner ; le garde doit etre
Expand Down Expand Up @@ -113,6 +122,72 @@ def test_bloque_apres_progression_tue_et_nomme_le_worker():
assert "XDIST-WATCHDOG" not in out # le verdict ne pollue pas la sortie pilote


def test_fin_de_parcours_points_partiels_non_tue():
# Mode 2, run 35276661841 : des octets vivants SANS \n (points de fin
# de parcours -q) doivent maintenir la fraicheur. Un garde qui
# n'ecouterait que les lignes completes croirait a un silence et
# tuerait ce run a 1,0 s de limite -- c'est exactement le faux
# positif qui a bloque la PR #16240 a [99%].
code, out, verdict = _run_watchdog(
_child("""
import sys, time
for i in range(16):
sys.stdout.write(".")
sys.stdout.flush()
time.sleep(0.15)
raise SystemExit(0)
"""),
idle_limit=1.0,
)
assert code == 0
assert "XDIST-WATCHDOG" not in verdict
assert "." in out


def test_fragment_final_sans_saut_de_ligne_recopie():
# Pass-through du fragment final : le dernier chunk sans \n doit etre
# recopie a l'EOF. La ligne partielle de 43 resultats du run
# 35276661841 n'aurait jamais du etre invisible jusqu'au kill.
code, out, verdict = _run_watchdog(
_child("""
import sys
print(".... [ 99%]", flush=True)
sys.stdout.write("...............s...........................")
sys.stdout.flush()
raise SystemExit(0)
"""),
idle_limit=10.0,
)
assert code == 0
assert "[ 99%]" in out
assert "...............s" in out
assert verdict == ""


def test_bloque_apres_fragment_partiel_tue_quand_meme():
# Garde-fou anti-regression : la fraicheur par octet ne doit pas
# epargner les vrais blocages. Signature [99%] + worker mort, un
# dernier fragment partiel, puis ZERO octet : le kill doit partir
# (fenetre comptee depuis le DERNIER OCTET, pas la derniere ligne)
# et le verdict doit citer les octets pour la forensique.
code, out, verdict = _run_watchdog(
_child("""
import sys, time
print("............................ [ 99%]", flush=True)
print("[gw2] node down: Not properly terminated", flush=True)
sys.stdout.write("..")
sys.stdout.flush()
time.sleep(300)
"""),
idle_limit=1.0,
)
assert code == wd.EXIT_BLOCKED
assert "BLOQUE" in verdict
assert "gw2" in verdict
assert "99%" in verdict
assert "octets" in verdict


def test_bloque_muet_des_la_naissance_tue_aussi():
# Hang de collection : aucune ligne jamais emise. Le garde doit etre
# arme des le demarrage, pas apres une premiere ligne.
Expand Down
Loading