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
193 changes: 180 additions & 13 deletions scripts/pick_idle_grain.py
Original file line number Diff line number Diff line change
Expand Up @@ -1257,11 +1257,110 @@ def delivered_probe_inert(issue_number: int, lane: str | None = None) -> bool:
return False


# --- Troisieme surface de livraison : PR MERGEE citant l'issue (#19907) ------
#
# Mesure du 2026-10-08 (cycles c.1450..c.1453, Tell c.1392 picker-delivered
# gap confirme 9x) : 24-28 cycles successifs sans grain actionnable. Le label
# `candidate-delivered` est pose par un workflow quotidien 05:49Z et le
# marqueur `[INFO] candidate-delivered` est poste par les lanes worker quand
# elles refutent un candidat. Mais l'urne `grain` reapparait regulierement des
# LIVREURS dont le label a ete retracte par le sweep post-merge, et dont
# aucune lane n'a encore refute (donc pas de marqueur). Cas fondateur : #16031
# (perf life_compose, PR #17347 MERGED 09/22, label retracte 09/24 par
# l'advisory post-merge, tapis narrow-cache l'a servi 28 cycles de suite).
#
# Le sweep post-merge ne s'arrete pas a 7 jours : il retracte le label sur
# TOUTE activite de commentaire post-merge, sans limite de temps (#15744). Et
# l'urne `delivered` (#15069) ne s'applique qu'aux LIVREURS portes par un
# label ou un marqueur -- les LIVREURS anciens dont le label a ete retracte
# retombent dans `grain` sans aucun signal. La PR mergee est la SEULE
# surface de verification qui survit au sweep.
#
# Cout : 1 requete `gh pr list --state merged --search "N in:title,body"`
# par candidat tire, partagee avec les autres sondes du meme plafond
# (`DELIVERED_SIGNAL_MAX_PROBES`). Borne de 90 jours : au-dela, un merge
# de cette envergure a deja ete documente ailleurs (release notes, ledger,
# ou ferme par `git log --grep=#N`). 90 = compromis entre couverture des
# livraisons recentes et exclusion des PRs historiques (avant #19907, la
# mesure c.1450 a releve 4 LIVREURS dans les 90 derniers jours, 0 au-dela
# sur le meme echantillon).
MERGED_PR_WINDOW_DAYS = 90


def merged_pr_probe_inert(issue_number: int, lane: str | None = None):
"""Sonde inerte : aucun signal, aucun appel reseau.

Meme doctrine que ``delivered_probe_inert`` : un test unitaire ne doit
pas emettre de requete par candidat tire. La signature retourne
``None`` (TRI-ETAT, voie fail-OPEN : pas de signal mesurable) plutot
que ``False`` parce qu'un verdict True/False signifierait « il y a un
signal, il vaut X » -- ce qu'un test inerte ne peut pas affirmer sans
avoir sonde. Le retour ``None`` laisse la voie du tapis intacte.
"""
return None


def merged_pr_signal(issue_number: int, lane: str | None = None):
"""Une PR MERGEE citant `#N` est-elle dans la fenetre de 90 jours ?

TRI-ETAT, meme doctrine que ``has_delivered_signal`` et
``open_cover_signal`` :
- ``True`` : au moins une PR MERGED cite `#N` et `mergedAt` est dans
les 90 derniers jours (jour UTC) ;
- ``False`` : aucune PR MERGED ne cite `#N`, OU toutes les PRs
MERGEes qui la citent datent de plus de 90 jours ;
- ``None`` : la requete a echoue (reseau, 403, payload illisible) --
l'appelant doit tirer quand meme EN LE DISANT.

Le `--search "N in:title,body"` est un filtre serveur (GitHub ne peut
pas tenir la borne de 90 jours cote serveur) ; le post-filtre temporel
se fait ici, sur le payload limite a 20 resultats. Meme doctrine
d'ancre `#N\\b` que ``open_cover_signal`` (post-filtre, #17760) pour
eviter les collisions de sous-chaine (ex. #11703 vs #1170391).
"""
try:
out = subprocess.run(
["gh", "pr", "list", "--repo", REPO, "--state", "merged",
"--limit", "20", "--search", f"{issue_number} in:title,body",
"--json", "number,state,mergedAt,title,body"],
capture_output=True, text=True, encoding="utf-8", check=True,
timeout=30,
).stdout
prs = json.loads(out)
except Exception: # noqa: BLE001 - sonde best-effort ; l'echec est DIT
return None
if not isinstance(prs, list):
return None
anchor = re.compile(rf"#{issue_number}\b")
now = dt.datetime.now(dt.timezone.utc)
window = dt.timedelta(days=MERGED_PR_WINDOW_DAYS)
for pr in prs:
if pr.get("state") != "MERGED":
continue
if not anchor.search((pr.get("title") or "") + "\n" +
(pr.get("body") or "")):
continue
merged_at = pr.get("mergedAt")
if not merged_at:
continue
try:
# GitHub renvoie des timestamps ISO-8601 en Z ; fromisoformat
# les accepte depuis Python 3.11 (la toolchain workers est 3.11+
# par regle F, et 3.13 sur po-2026 ML-Training-Pipeline).
merged_dt = dt.datetime.fromisoformat(merged_at.replace("Z", "+00:00"))
except ValueError:
continue
if (now - merged_dt) <= window:
return True
return False


def delivered_signal_reason(
item: dict,
lane: str | None = None,
probe=None,
failures: list[int] | None = None,
merged_pr_probe=None,
) -> str | None:
"""Pourquoi ecarter ce candidat de l'urne `grain`, ou ``None``.

Expand All @@ -1280,12 +1379,17 @@ def delivered_signal_reason(
retirerait une source de grains de CONTENU, l'inverse du but.

Le label est teste EN PREMIER parce qu'il ne coute rien ; la sonde de
commentaire n'est atteinte que s'il est absent. Ce n'est pas une
micro-optimisation : c'est ce qui fait que 59 des 113 issues signalees
du pool du 2026-09-12 sont ecartees sans un seul appel reseau.
commentaire n'est atteinte que s'il est absent ; la sonde PR-mergee
(#19907) n'est atteinte qu'en dernier, parce qu'elle est la plus
couteuse (meme `gh pr list` que ``open_cover_signal`` mais avec un
post-filtre temporel). Ce n'est pas une micro-optimisation : c'est ce
qui fait que 59 des 113 issues signalees du pool du 2026-09-12 sont
ecartees sans un seul appel reseau.
"""
if probe is None:
probe = delivered_probe_inert
if merged_pr_probe is None:
merged_pr_probe = merged_pr_probe_inert
labels = {str(label).casefold() for label in item.get("labels") or []}
if DELIVERED_LABEL in labels:
return (
Expand All @@ -1311,6 +1415,29 @@ def delivered_signal_reason(
"refutant. La re-servir comme grain de production fait bruler un "
"cycle a la lane qui la recoit."
)
# Troisieme surface : PR MERGEE <90j citant `#N` (#19907). Couvre les
# LIVREURS dont le label a ete retracte par le sweep post-merge et dont
# aucune lane n'a encore refute (donc pas de marqueur). Cas fondateur
# c.1450 : #16031 (perf life_compose), label retracte 09/24, servi 28
# cycles de suite par le tapis narrow-cache. La sonde est partagee
# avec `open_cover_reason` via le meme `_counted_probe` parent (meme
# plafond `DELIVERED_SIGNAL_MAX_PROBES`).
merged = merged_pr_probe(item["number"], lane)
if merged == DELIVERED_SIGNAL_UNPROBED:
return None
if merged is None:
if failures is not None:
failures.append(item["number"])
return None
if merged:
return (
"SIGNAL LIVRAISON (PR MERGEE <" + str(MERGED_PR_WINDOW_DAYS) +
"j) : une PR MERGEE recente cite cette issue, la livraison est "
"faite par cette PR. Le sweep post-merge a retracte le label "
"`candidate-delivered` (#15744) et aucune lane n'a encore "
"refute (donc pas de marqueur en commentaire), mais la "
"livraison est reelle. Cf #19907, cas fondateur c.1450 #16031."
)
return None


Expand Down Expand Up @@ -2018,7 +2145,8 @@ def check_claims(numbers: list[int], lane: str) -> dict[int, tuple[str, str]]:
def draw_unclaimed(by_class, args, rng, visits, series, issue_to_family,
delivery=None, delivered_probe=None, delivered_state=None,
fallback_by_class=None, continuity_state=None,
cover_probe=None, long_visits=None):
cover_probe=None, long_visits=None,
merged_pr_probe=None):
"""Tire, puis REMPLACE tout candidat qu une autre lane tient deja.

Trois raisons de remplacer plutot que d annoter :
Expand Down Expand Up @@ -2081,6 +2209,22 @@ def _counted_cover_probe(number):
budget[0] -= 1
return (cover_probe or open_cover_inert)(number)

def _counted_merged_pr_probe(number, lane_name):
# Troisieme surface de livraison (#19907) : PR MERGEE <90j citant
# l'issue. Partage le MEME budget que les deux autres sondes
# (label/marqueur + PR couvrante) : un candidat peut consommer
# jusqu'a 3 unites de budget (commentaire + PR couvrante + PR
# mergée recente). C'est le prix de la couverture des LIVREURS
# dont le label a ete retracte par le sweep post-merge (#15744)
# et dont aucune lane n'a refute (donc pas de marqueur en
# commentaire) -- le cas fondateur c.1450 #16031 servi 28
# cycles de suite par le tapis narrow-cache faute de cette sonde.
if budget[0] <= 0:
state["budget_hit"] = True
return DELIVERED_SIGNAL_UNPROBED
budget[0] -= 1
return (merged_pr_probe or merged_pr_probe_inert)(number, lane_name)

for cls, want, prev in urnes:
primary = list(by_class[cls])
primary_numbers = {item["number"] for item in primary}
Expand Down Expand Up @@ -2163,7 +2307,8 @@ def _counted_cover_probe(number):
# commentaire est plafonnee, et son epuisement est
# fail-OPEN (le candidat est conserve).
reason = delivered_signal_reason(
c, args.lane, _counted_probe, failures)
c, args.lane, _counted_probe, failures,
_counted_merged_pr_probe)
if reason is not None:
conflicts.append((c, "LIVRAISON : " + reason + (
" Candidat remplace dans la meme urne.")))
Expand Down Expand Up @@ -5358,7 +5503,8 @@ def belt_probe_budget(belt_check_window: int) -> int:


def belt_pick_with_replacements(belt_pool, belt_claims, args, probe_budget,
delivered_probe=None, claims_probe=None):
delivered_probe=None, claims_probe=None,
merged_pr_probe=None):
"""La boucle de service du tapis : claims, livraison, remplacement.

Extraite de ``main`` (#19390) pour etre testable sans harnais complet :
Expand All @@ -5373,9 +5519,16 @@ def belt_pick_with_replacements(belt_pool, belt_claims, args, probe_budget,
``belt_probe_budget(belt_check_window)`` -- la valeur doit couvrir la tete
que la boucle peut servir, sinon les candidats au-dela sont servis sans
lecture (#19969). Epuisement = fail-OPEN (``DELIVERED_SIGNAL_UNPROBED``),
rapporte dans l'etat de retour. ``delivered_probe`` (tests) remplace la
sonde reseau ; ``claims_probe`` (tests) remplace la verification au fil
de l'eau des items hors fenetre.
rapporte dans l'etat de retour. ``delivered_probe`` et ``merged_pr_probe``
(tests) remplacent les sondes reseau ; ``claims_probe`` (tests) remplace
la verification au fil de l'eau des items hors fenetre.

Defauts INERTES (doctrine ``delivered_probe_inert`` / #19913) : un test
qui appelle cette boucle sans injecter ses sondes n'emet AUCUNE requete
reseau. Qui veut le vrai signal l'injecte -- ``main`` le fait
explicitement. Le defaut precedent (sonde reseau) etait un piege : un
test neuf qui oubliait l'injection pendait sur un vrai ``gh`` (mesure
2026-10-09 : 99 s pour un seul test, cf. coordinateur #19913).

Rend ``(picks, withheld, etat)`` avec ``etat = {"failures": [numeros
illisibles], "budget_hit": bool}``.
Expand All @@ -5386,7 +5539,8 @@ def belt_pick_with_replacements(belt_pool, belt_claims, args, probe_budget,
# Le budget enrobe TOUTE sonde, injectee ou reelle : un test qui fournit
# sa sonde doit voir le plafond s'appliquer aussi -- sinon la borne de
# cout ne serait testable qu'avec le reseau reel.
user_probe = delivered_probe or has_delivered_signal
user_probe = delivered_probe or delivered_probe_inert
user_merged_pr_probe = merged_pr_probe or merged_pr_probe_inert

def counted_probe(number, lane_name):
if budget[0] <= 0:
Expand All @@ -5395,6 +5549,13 @@ def counted_probe(number, lane_name):
budget[0] -= 1
return user_probe(number, lane_name)

def counted_merged_pr_probe(number, lane_name):
if budget[0] <= 0:
state["budget_hit"] = True
return DELIVERED_SIGNAL_UNPROBED
budget[0] -= 1
return user_merged_pr_probe(number, lane_name)

if claims_probe is None:
def claims_probe(numbers):
return check_claims(numbers, args.lane)
Expand Down Expand Up @@ -5438,7 +5599,8 @@ def claims_probe(numbers):
# mesuree du canal label).
if it.get("klass") == "grain" and not args.include_delivered:
reason = delivered_signal_reason(
it, args.lane, counted_probe, state["failures"])
it, args.lane, counted_probe, state["failures"],
counted_merged_pr_probe)
if reason is not None:
belt_withheld.append(
(it, "LIVRAISON : " + reason
Expand Down Expand Up @@ -6038,10 +6200,15 @@ def main(argv: list[str] | None = None) -> int:
# Boucle de service extraite (#19390) : claims + sonde de livraison
# bornee a la fenetre de tete, remplacement dans la meme urne.
# Epuisement du budget de sondes = fail-OPEN, rapporte en banniere
# et en JSON (cf belt_pick_with_replacements).
# et en JSON (cf belt_pick_with_replacements). Les sondes reelles
# sont injectees ICI explicitement : les defauts de la boucle sont
# inertes (#19913) pour que les tests ne sortent jamais sur le
# reseau par oubli.
belt_picks, belt_withheld, belt_pick_state = (
belt_pick_with_replacements(belt_pool, belt_claims, args,
belt_probe_budget(belt_check_window)))
belt_probe_budget(belt_check_window),
delivered_probe=has_delivered_signal,
merged_pr_probe=merged_pr_signal))
belt_delivered_failures = belt_pick_state["failures"]
belt_probe_budget_hit = [belt_pick_state["budget_hit"]]
# Banniere legere : le tapis ne refuse jamais, mais rappelle
Expand Down
28 changes: 19 additions & 9 deletions scripts/tests/test_pick_delivered_gate.py
Original file line number Diff line number Diff line change
Expand Up @@ -542,20 +542,26 @@ def failing(number):


def test_plafond_partage_entre_sondes_livraison_et_couverture(monkeypatch):
"""« Au meme rang » (#16589) = AUSSI au meme cout : les deux sondes
"""« Au meme rang » (#16589) = AUSSI au meme cout : les 3 sondes
parent le MEME budget de DELIVERED_SIGNAL_MAX_PROBES, pas un plafond
double. Avec un plafond de 2 et deux candidats, exactement deux sondes
reelles partent (livraison du 1er, couverture du 1er) -- le 2e candidat
n'en a plus aucune."""
monkeypatch.setattr(pig, "DELIVERED_SIGNAL_MAX_PROBES", 2)
double. Avec un plafond de 4 et deux candidats, exactement quatre sondes
reelles partent (livraison 1er, merged_pr 1er, couverture 1er,
livraison 2e) -- le 2e candidat n'a plus que la livraison, et la
couverture n'est pas appelee une 2e fois."""
monkeypatch.setattr(pig, "DELIVERED_SIGNAL_MAX_PROBES", 4)
_patch_draw(monkeypatch, [1, 2])
delivered_calls = []
merged_pr_calls = []
cover_calls = []

def counting_delivered(number, lane=None):
delivered_calls.append(number)
return False

def counting_merged_pr(number, lane=None):
merged_pr_calls.append(number)
return False

def counting_cover(number):
cover_calls.append(number)
return ""
Expand All @@ -565,11 +571,15 @@ def counting_cover(number):
picks, _, _ = draw_unclaimed(
by, _args(grains=2), random.Random(7), None, None, None,
delivered_probe=counting_delivered,
cover_probe=counting_cover, delivered_state=state)
assert delivered_calls == [1], (
f"plafond partage : une seule sonde livraison, or {delivered_calls}")
cover_probe=counting_cover,
merged_pr_probe=counting_merged_pr,
delivered_state=state)
assert delivered_calls == [1, 2], (
f"plafond partage : une sonde livraison par candidat survivant, or {delivered_calls}")
assert merged_pr_calls == [1], (
f"plafond partage : une seule sonde merged_pr (1er candidat), or {merged_pr_calls}")
assert cover_calls == [1], (
f"plafond partage : une seule sonde couverture, or {cover_calls}")
f"plafond partage : une seule sonde couverture (1er candidat), or {cover_calls}")
assert state["budget_hit"] is True
assert len(picks) == 2, "les non sondes sont CONSERVES (fail-open)"

Expand Down
Loading
Loading