From f10bacb29ee38690af523240f8fef58fd0d058d7 Mon Sep 17 00:00:00 2001 From: jsboige Date: Tue, 15 Sep 2026 01:05:18 +0200 Subject: [PATCH] fix(nb-tools,#16213): reprise EAGAIN partagee pour les deux derniers gardes Les deux gardes sans reprise sur la pression de fork -- check_slot_reservation et check_source_output_ratchet -- sont cables sur une primitive partagee, `fork_retry.run_with_fork_retry`, sur le modele de `naming_canon.py` (#15503) plutot qu'en quatrieme et cinquieme copie de la meme boucle. Chaque garde CONSERVE sa politique d'epuisement : fail-closed pour check_slot_reservation (un instrument muet ne vote pas), fail-open pour check_source_output_ratchet (politique preexistante, arbitrage #16164). Deux tests l'epinglent garde par garde -- c'est exactement ce qu'un refactor distrait harmoniserait. Le filtre reste etroit : seuls EAGAIN et EWOULDBLOCK sont retentes, controle positif inclus (ENOENT / EACCES / ENOMEM / EINVAL remontent au premier appel). See #16213, #16111, #16164. Co-Authored-By: Claude Opus 5 (1M context) --- .../notebook_tools/check_slot_reservation.py | 18 +- .../check_source_output_ratchet.py | 18 +- scripts/notebook_tools/fork_retry.py | 80 +++++ .../notebook_tools/tests/test_fork_retry.py | 287 ++++++++++++++++++ 4 files changed, 397 insertions(+), 6 deletions(-) create mode 100644 scripts/notebook_tools/fork_retry.py create mode 100644 scripts/notebook_tools/tests/test_fork_retry.py diff --git a/scripts/notebook_tools/check_slot_reservation.py b/scripts/notebook_tools/check_slot_reservation.py index c74d6ae73d..e8ebd23230 100644 --- a/scripts/notebook_tools/check_slot_reservation.py +++ b/scripts/notebook_tools/check_slot_reservation.py @@ -175,6 +175,12 @@ sys.path.insert(0, _here) from naming_canon import index_key, strip_lang # noqa: E402 +# Reprise bornee sur la pression de fork (#16213) : sous `pytest-xdist -n 4`, le +# spawn de `git` est refuse par intermittence (`EAGAIN`). Meme raison de +# mutualiser que la grammaire de nom ci-dessus -- cinq gardes reecrivaient la +# meme boucle de reprise, mot pour mot. +from fork_retry import run_with_fork_retry # noqa: E402 + DEFAULT_RESERVATIONS = Path(__file__).resolve().parent / "slot_reservations.json" SOURCE_BASE = "base" @@ -197,9 +203,17 @@ def _git(args): + """Appel git resilient a la pression de fork -- fail-closed a l'epuisement. + + `run_with_fork_retry` remonte la derniere `OSError` une fois ses tentatives + epuisees, et ce garde la laisse remonter DELIBEREMENT. Un instrument qui n'a + pas pu mesurer ne doit pas rendre de verdict : ici, une reservation de slot + silencieusement absente laisserait passer la collision que le garde existe + pour attraper. + """ env = dict(os.environ, MSYS_NO_PATHCONV="1") - r = subprocess.run(["git"] + args, capture_output=True, text=True, - encoding="utf-8", errors="replace", env=env) + r = run_with_fork_retry(["git"] + args, capture_output=True, text=True, + encoding="utf-8", errors="replace", env=env) if r.returncode != 0: raise RuntimeError("git %s -> %s" % (" ".join(args), (r.stderr or "").strip()[:200])) return r.stdout diff --git a/scripts/notebook_tools/check_source_output_ratchet.py b/scripts/notebook_tools/check_source_output_ratchet.py index 4578ccae54..3c54273ebd 100644 --- a/scripts/notebook_tools/check_source_output_ratchet.py +++ b/scripts/notebook_tools/check_source_output_ratchet.py @@ -77,7 +77,6 @@ import difflib import json import re -import subprocess import sys from pathlib import Path @@ -88,6 +87,10 @@ QC_CLOUD_PATHS, QUANTBOOK_PATTERN, ) +# Reprise bornee sur la pression de fork (#16213) : sous `pytest-xdist -n 4`, le +# spawn de `git` est refuse par intermittence (`EAGAIN`). Primitive partagee par +# les gardes, pas une cinquieme copie de la meme boucle. +from fork_retry import run_with_fork_retry # noqa: E402 # Same exclusions as check_papermill_ratchet.py / check_exec_ratchet.py: # archived, papermill-output and research copies are not deliverable @@ -117,10 +120,17 @@ def git(*args, cwd=None): - """Run a git command, returning stdout (utf-8) or None on failure.""" + """Run a git command, returning stdout (utf-8) or None on failure. + + Fail-open a l'epuisement de la reprise : l'`OSError` est avalee en `None`, + comme avant #16213. Cette politique est PREEXISTANTE et deliberement + conservee ici -- la mutualisation de la reprise ne doit pas trancher en + passant un arbitrage qui appartient a #16164 (un ratchet muet qui rend + « 0 changement » est un faux vert, pas une mesure). + """ try: - out = subprocess.run(["git", *args], cwd=cwd, capture_output=True, - encoding="utf-8", errors="replace", check=False) + out = run_with_fork_retry(["git", *args], cwd=cwd, capture_output=True, + encoding="utf-8", errors="replace", check=False) except OSError: return None return out.stdout if out.returncode == 0 else None diff --git a/scripts/notebook_tools/fork_retry.py b/scripts/notebook_tools/fork_retry.py new file mode 100644 index 0000000000..7802ae73b9 --- /dev/null +++ b/scripts/notebook_tools/fork_retry.py @@ -0,0 +1,80 @@ +"""Reprise bornee sur pression de fork (EAGAIN) -- primitive partagee par les gardes. + +Origine -- la vague de rouges du 2026-09-14. `f149f2fe93` (mergee sur `main` le +2026-09-13T23:08:17Z, #15833 sur #14598) a fait passer `Scripts Tests (CPU)` a +`pytest-xdist -n 4 --dist loadscope`. Quatre workers executant en parallele des +gardes qui forkent `git` saturent la table de processus du runner : le spawn +suivant est refuse avec `BlockingIOError: [Errno 11]` -- `EAGAIN`. + +`EAGAIN` est transitoire PAR DEFINITION : le noyau refuse *un processus de plus*, +il ne signale pas une panne. Echouer au premier refus convertit un pic de charge +en rouge de base. + +Trois gardes ont recu cette reprise separement -- `check_twin_parity` (#16125), +`check_kernel_suffix_canon` et `check_exec_ratchet` (#16157) : meme boucle, meme +backoff, meme filtre d'errno, trois ecritures. Ce module est l'ecriture unique, +sur le modele de `naming_canon.py` (#15503) : meme classe de probleme -- une +primitive redefinie par plusieurs organes -- et donc meme reponse. + +PORTEE -- ce que ce module fait, et ce qu'il ne fait pas +-------------------------------------------------------- +Il rend un echec de spawn TRANSITOIRE surmontable. Il ne reduit PAS la pression +de fork : le fan-out des gardes reste le meme, et le reduire (un +`git cat-file --batch` en flux plutot que N forks) est un axe distinct, laisse +ouvert par #16111 faute d'instrumentation du runner. + +Il ne decide pas non plus de ce qu'un garde fait quand la reprise est epuisee. +La derniere `OSError` REMONTE, et chaque garde garde sa politique par-dessus : +`check_slot_reservation` la laisse remonter (fail-closed), `check_exec_ratchet` +et `check_source_output_ratchet` l'enveloppent en `None` (fail-open -- arbitrage +en cours dans #16164). Mutualiser la reprise ne doit pas uniformiser la +semantique de sortie en passant : ces deux politiques repondent a des questions +differentes, et le choix appartient a chaque garde. + +Le filtre est ETROIT, et c'est le point qui compte : seuls `EAGAIN` et +`EWOULDBLOCK` sont retentes ; toute autre `OSError` remonte au premier appel. +Sans cela, une panne permanente -- git absent, `ENOENT` -- deviendrait trois +essais silencieux, et le remede serait pire que le mal. +""" +from __future__ import annotations + +import errno +import subprocess +import time + +# `EWOULDBLOCK` est un alias d'`EAGAIN` sous Linux ; il peut differer ailleurs. +FORK_PRESSURE_ERRNOS = (errno.EAGAIN, getattr(errno, "EWOULDBLOCK", errno.EAGAIN)) + +ATTEMPTS = 3 +BACKOFF = (0.05, 0.15) # avant les 2e et 3e tentatives + + +def is_fork_pressure(exc: BaseException) -> bool: + """Vrai pour un refus de spawn transitoire, et pour rien d'autre.""" + return isinstance(exc, OSError) and exc.errno in FORK_PRESSURE_ERRNOS + + +def run_with_fork_retry(args, **kwargs) -> subprocess.CompletedProcess: + """`subprocess.run(args, **kwargs)`, retente de facon bornee sur EAGAIN. + + Les `kwargs` sont transmis VERBATIM : ce module ne prend aucune position sur + le mode texte, l'encodage, `cwd`, `env` ou `check`. Un garde qui lit du + binaire (`text=False`, sans `encoding=`) et un garde qui lit de l'utf-8 + passent par la meme porte sans que l'un impose sa forme a l'autre -- c'est + la condition pour que la mutualisation n'introduise aucun changement de + comportement hors pression de fork. + + Rend le `CompletedProcess` du premier spawn qui aboutit. A l'epuisement des + tentatives, remonte la derniere `OSError` de pression de fork. + """ + derniere: OSError | None = None + for tentative in range(ATTEMPTS): + if tentative: + time.sleep(BACKOFF[tentative - 1]) + try: + return subprocess.run(args, **kwargs) + except OSError as exc: + if not is_fork_pressure(exc): + raise + derniere = exc + raise derniere # `ATTEMPTS >= 1` garantit que `derniere` est renseignee diff --git a/scripts/notebook_tools/tests/test_fork_retry.py b/scripts/notebook_tools/tests/test_fork_retry.py new file mode 100644 index 0000000000..2577f29861 --- /dev/null +++ b/scripts/notebook_tools/tests/test_fork_retry.py @@ -0,0 +1,287 @@ +"""Tests de `fork_retry` -- la reprise bornee sur EAGAIN, et sa portee exacte. + +Why this exists +--------------- +Un correctif de reprise est, vu de l'exterieur, indiscernable d'un avaleur +d'erreurs : les deux font disparaitre des exceptions. Ce qui les separe est le +CONTROLE POSITIF -- la preuve qu'une erreur qui n'est pas de la pression de fork +remonte AU PREMIER APPEL, sans reprise et sans attente. C'est la moitie de ce +fichier, et c'est la moitie qui compte. + +L'autre moitie tient a la mutualisation elle-meme. Extraire une primitive +partagee par cinq gardes fait courir un risque precis : uniformiser en passant ce +qui etait deliberement different. Deux choses etaient deliberement differentes et +doivent le rester -- + +1. la FORME de l'appel (un garde lit de l'utf-8, un autre du binaire, un autre + passe un `env` ou un `cwd`) : `run_with_fork_retry` transmet les `kwargs` + verbatim, et `test_kwargs_transmis_verbatim` l'epingle ; +2. la POLITIQUE D'EPUISEMENT (fail-closed pour `check_slot_reservation`, + fail-open pour `check_source_output_ratchet`) : les deux derniers tests + l'epinglent garde par garde, parce que c'est precisement ce qu'un refactor + distrait harmoniserait. + +Ce qui n'est PAS teste, volontairement : que la reprise reduise la pression de +fork. Elle ne la reduit pas -- elle rend l'echec transitoire surmontable. Le +fan-out reste le meme, et le reduire est un axe distinct (#16111). +""" +from __future__ import annotations + +import errno +import os +import subprocess +import sys + +import pytest + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + +import check_slot_reservation as csr # noqa: E402 +import check_source_output_ratchet as csor # noqa: E402 +import fork_retry # noqa: E402 + + +def _eagain() -> BlockingIOError: + return BlockingIOError(errno.EAGAIN, "Resource temporarily unavailable") + + +@pytest.fixture(autouse=True) +def _pas_d_attente_reelle(monkeypatch): + """Les tests mesurent le backoff, ils ne le subissent pas.""" + monkeypatch.setattr(fork_retry.time, "sleep", lambda _s: None) + + +# -------------------------------------------------------------------------- +# 1. La reprise +# -------------------------------------------------------------------------- + +def test_retente_sur_eagain(monkeypatch): + """EAGAIN est temporaire : les 2 premiers spawns echouent, le 3e passe.""" + appels = {"n": 0} + sentinelle = subprocess.CompletedProcess(["git"], 0, stdout="ok\n", stderr="") + + def faux_run(_args, **_kwargs): + appels["n"] += 1 + if appels["n"] < 3: + raise _eagain() + return sentinelle + + monkeypatch.setattr(subprocess, "run", faux_run) + + assert fork_retry.run_with_fork_retry(["git", "status"]) is sentinelle + assert appels["n"] == 3 + + +def test_retente_sur_ewouldblock(monkeypatch): + """`EWOULDBLOCK` alias `EAGAIN` sous Linux, mais pas partout : couvrir les deux.""" + appels = {"n": 0} + + def faux_run(args, **_kwargs): + appels["n"] += 1 + if appels["n"] < 2: + raise OSError(errno.EWOULDBLOCK, "Operation would block") + return subprocess.CompletedProcess(args, 0, stdout="", stderr="") + + monkeypatch.setattr(subprocess, "run", faux_run) + + fork_retry.run_with_fork_retry(["git", "status"]) + assert appels["n"] == 2 + + +def test_epuisement_remonte_la_derniere_erreur(monkeypatch): + """La borne est reelle : apres `ATTEMPTS` essais, l'erreur remonte.""" + appels = {"n": 0} + + def faux_run(_args, **_kwargs): + appels["n"] += 1 + raise _eagain() + + monkeypatch.setattr(subprocess, "run", faux_run) + + with pytest.raises(BlockingIOError): + fork_retry.run_with_fork_retry(["git", "status"]) + assert appels["n"] == fork_retry.ATTEMPTS + + +# -------------------------------------------------------------------------- +# 2. Controle positif -- le filtre est etroit +# -------------------------------------------------------------------------- + +def test_controle_positif_une_panne_reelle_n_est_pas_retentee(monkeypatch): + """Sans ce test, la reprise serait indiscernable d'un avaleur d'erreurs. + + Une panne permanente -- git absent, `ENOENT` -- doit remonter AU PREMIER + appel : la retenter trois fois en silence rendrait le remede pire que le mal. + """ + appels = {"n": 0} + + def faux_run(_args, **_kwargs): + appels["n"] += 1 + raise FileNotFoundError(errno.ENOENT, "No such file or directory: 'git'") + + monkeypatch.setattr(subprocess, "run", faux_run) + + with pytest.raises(FileNotFoundError): + fork_retry.run_with_fork_retry(["git", "status"]) + assert appels["n"] == 1, "une erreur non-EAGAIN ne doit pas etre retentee" + + +@pytest.mark.parametrize("errno_erreur", [errno.EACCES, errno.ENOMEM, errno.EINVAL]) +def test_les_autres_errno_ne_sont_pas_retentes(monkeypatch, errno_erreur): + """`ENOMEM` est le voisin dangereux : lui aussi vient de la charge, et lui + n'est pas transitoire. Le filtre ne doit pas glisser jusqu'a lui.""" + appels = {"n": 0} + + def faux_run(_args, **_kwargs): + appels["n"] += 1 + raise OSError(errno_erreur, os.strerror(errno_erreur)) + + monkeypatch.setattr(subprocess, "run", faux_run) + + with pytest.raises(OSError): + fork_retry.run_with_fork_retry(["git", "status"]) + assert appels["n"] == 1 + + +def test_is_fork_pressure_ne_repond_qu_aux_oserror(): + assert fork_retry.is_fork_pressure(_eagain()) + assert not fork_retry.is_fork_pressure(OSError(errno.ENOENT, "absent")) + assert not fork_retry.is_fork_pressure(ValueError("pas une OSError")) + + +# -------------------------------------------------------------------------- +# 3. Le repli reste borne en temps +# -------------------------------------------------------------------------- + +def test_backoff_borne(monkeypatch): + dors: list[float] = [] + appels = {"n": 0} + + def faux_run(args, **_kwargs): + appels["n"] += 1 + if appels["n"] < fork_retry.ATTEMPTS: + raise _eagain() + return subprocess.CompletedProcess(args, 0, stdout="", stderr="") + + monkeypatch.setattr(subprocess, "run", faux_run) + monkeypatch.setattr(fork_retry.time, "sleep", lambda s: dors.append(s)) + + fork_retry.run_with_fork_retry(["git", "status"]) + + assert dors == list(fork_retry.BACKOFF) + assert sum(dors) <= 1.0, "le repli doit rester borne en temps" + + +def test_un_backoff_par_reprise(): + """Invariant d'indexation : `BACKOFF[tentative - 1]` ne doit jamais deborder.""" + assert len(fork_retry.BACKOFF) >= fork_retry.ATTEMPTS - 1 + + +# -------------------------------------------------------------------------- +# 4. La mutualisation n'impose aucune forme d'appel +# -------------------------------------------------------------------------- + +def test_kwargs_transmis_verbatim(monkeypatch): + """Le module ne prend position ni sur le mode texte, ni sur `env`/`cwd`/`check`. + + C'est la condition pour qu'un garde binaire et un garde utf-8 partagent la + meme porte sans changement de comportement. + """ + vus: list[dict] = [] + + def faux_run(args, **kwargs): + vus.append(kwargs) + return subprocess.CompletedProcess(args, 0, stdout=b"", stderr=b"") + + monkeypatch.setattr(subprocess, "run", faux_run) + + fork_retry.run_with_fork_retry(["git", "show", "HEAD:x"], capture_output=True, text=False) + fork_retry.run_with_fork_retry( + ["git", "status"], cwd=".", capture_output=True, + encoding="utf-8", errors="replace", check=False, env={"A": "1"}, + ) + + assert vus[0] == {"capture_output": True, "text": False} + assert vus[1] == { + "cwd": ".", "capture_output": True, "encoding": "utf-8", + "errors": "replace", "check": False, "env": {"A": "1"}, + } + + +# -------------------------------------------------------------------------- +# 5. Chaque garde CONSERVE sa politique d'epuisement +# -------------------------------------------------------------------------- + +def test_slot_reservation_retente_puis_reste_fail_closed(monkeypatch): + """`check_slot_reservation` laisse remonter : un instrument muet ne vote pas.""" + appels = {"n": 0} + + def faux_run(args, **_kwargs): + appels["n"] += 1 + if appels["n"] < 3: + raise _eagain() + return subprocess.CompletedProcess(args, 0, stdout="a.ipynb\n", stderr="") + + monkeypatch.setattr(subprocess, "run", faux_run) + assert csr._git(["status"]) == "a.ipynb\n" + assert appels["n"] == 3 + + appels["n"] = 0 + + def toujours_eagain(_args, **_kwargs): + appels["n"] += 1 + raise _eagain() + + monkeypatch.setattr(subprocess, "run", toujours_eagain) + with pytest.raises(OSError): + csr._git(["status"]) + assert appels["n"] == fork_retry.ATTEMPTS + + +def test_source_output_ratchet_retente_puis_reste_fail_open(monkeypatch): + """`check_source_output_ratchet` rend `None` -- politique preexistante, #16164. + + Ce test ne l'approuve pas, il l'EPINGLE : mutualiser la reprise ne doit pas + la changer au passage. L'arbitrage fail-open/fail-closed appartient a #16164. + """ + appels = {"n": 0} + + def faux_run(args, **_kwargs): + appels["n"] += 1 + if appels["n"] < 3: + raise _eagain() + return subprocess.CompletedProcess(args, 0, stdout="ok\n", stderr="") + + monkeypatch.setattr(subprocess, "run", faux_run) + assert csor.git("status") == "ok\n" + assert appels["n"] == 3 + + appels["n"] = 0 + + def toujours_eagain(_args, **_kwargs): + appels["n"] += 1 + raise _eagain() + + monkeypatch.setattr(subprocess, "run", toujours_eagain) + assert csor.git("status") is None + assert appels["n"] == fork_retry.ATTEMPTS + + +def test_les_deux_gardes_ne_retentent_pas_une_panne_reelle(monkeypatch): + """Controle positif au niveau des gardes, pas seulement de la primitive.""" + appels = {"n": 0} + + def faux_run(_args, **_kwargs): + appels["n"] += 1 + raise FileNotFoundError(errno.ENOENT, "No such file or directory: 'git'") + + monkeypatch.setattr(subprocess, "run", faux_run) + + with pytest.raises(FileNotFoundError): + csr._git(["status"]) + assert appels["n"] == 1 + + appels["n"] = 0 + # Fail-open : le ratchet avale l'`OSError` comme avant -- mais UNE SEULE fois. + assert csor.git("status") is None + assert appels["n"] == 1