From 056b1a512c10149058f8ab5590ca2d0f9e4954e1 Mon Sep 17 00:00:00 2001 From: Gabiar <189609572+gabiar-maker@users.noreply.github.com> Date: Fri, 12 Jun 2026 15:38:36 +0200 Subject: [PATCH 1/3] feat(jb_outbound): attribution des drafts au departement du job cron MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Quand une interception a lieu pendant un job cron (skill/casquette), le DraftRequest porte desormais le departement de la tache : champs ADDITIFS de premier niveau `department`, `skill_id`, `job_id` (omis hors contexte job — chat libre ; le daemon Go actuel ignore les champs inconnus). - job_context.py : ContextVar du job courant, posee par le scheduler au lancement du job. ContextVar et PAS env : les jobs cron tournent dans des THREADS du gateway (pool parallele de tick()) — os.environ s'ecraserait entre jobs concurrents. Hermes propage deja le contexte a chaque saut de thread (copy_context dans _run_job_impl, propagate_context_to_thread pour les outils) → la valeur est visible du middleware pendant tout le job. - Casquette lue dans le front-matter du skill du job : `casquette:` (gold) puis `department:` (custom). Resolution best-effort, parser minimal autonome (pas de dependance yaml ni du coeur) ; toute erreur → champs absents, jamais d'exception. - cron/scheduler.py : pont OPTIONNEL _jb_job_hooks() dans run_job — resout le module DEJA charge par le PluginManager (hermes_plugins.jb_outbound), jamais d'import a froid : sans plugin, comportement stock inchange. Double garde : un hook ne bloque ni ne fait echouer un job. Tests : 9 nouveaux (stamp present/absent/efface, gold/custom/sans casquette, categorie, passif, pont scheduler reel) ; suite 21/21. Co-Authored-By: Claude Fable 5 --- cron/scheduler.py | 42 +++- plugins/jb_outbound/job_context.py | 181 +++++++++++++++++ plugins/jb_outbound/middleware.py | 11 +- .../jb_outbound/test_attribution_activity.py | 184 ++++++++++++++++++ 4 files changed, 416 insertions(+), 2 deletions(-) create mode 100644 plugins/jb_outbound/job_context.py create mode 100644 plugins/jb_outbound/test_attribution_activity.py diff --git a/cron/scheduler.py b/cron/scheduler.py index f5c71ceed4f0..a615b18d0315 100644 --- a/cron/scheduler.py +++ b/cron/scheduler.py @@ -1341,11 +1341,51 @@ def _scan_assembled_cron_prompt(assembled: str, job: dict, *, has_skills: bool = return assembled +def _jb_job_hooks(): + """Pont OPTIONNEL vers le plugin jb_outbound (greffe Jean-Billie). + + Expose au scheduler le contexte d'attribution (casquette/skill/job, lu ensuite par le + middleware du plugin) et les signaux début/fin de job. On résout le module DÉJÀ chargé — + le PluginManager importe les plugins sous ``hermes_plugins.`` ; jamais d'import à + froid : sans plugin chargé, renvoie ``None`` et le scheduler garde son comportement + d'origine (stock Hermes inchangé). + """ + import importlib + for pkg in ("hermes_plugins.jb_outbound", "jb_outbound", "plugins.jb_outbound"): + if pkg in sys.modules: + try: + return importlib.import_module(pkg + ".job_context") + except Exception: + logger.debug("jb_outbound job_context unavailable", exc_info=True) + return None + return None + + def run_job(job: dict) -> tuple[bool, str, str, Optional[str]]: """Execute a single cron job, applying any per-job profile override.""" job_id = job["id"] with _job_profile_context(job_id, job.get("profile")): - return _run_job_impl(job) + jb_hooks = _jb_job_hooks() + if jb_hooks is None: + return _run_job_impl(job) + # Greffe Jean-Billie : attribution + signaux d'activité. Les deux hooks n'échouent + # jamais côté plugin (best-effort) ; double garde ici — un hook ne doit en AUCUN cas + # bloquer ni faire échouer le job. + jb_token = None + try: + jb_token = jb_hooks.job_started(job) + except Exception: + logger.debug("jb_outbound job_started hook failed", exc_info=True) + success = False + try: + result = _run_job_impl(job) + success = bool(result and result[0]) + return result + finally: + try: + jb_hooks.job_finished(job, success=success, token=jb_token) + except Exception: + logger.debug("jb_outbound job_finished hook failed", exc_info=True) def _run_job_impl(job: dict) -> tuple[bool, str, str, Optional[str]]: diff --git a/plugins/jb_outbound/job_context.py b/plugins/jb_outbound/job_context.py new file mode 100644 index 000000000000..bd4c0dccc125 --- /dev/null +++ b/plugins/jb_outbound/job_context.py @@ -0,0 +1,181 @@ +"""Contexte d'attribution du job cron courant (casquette / skill / job). + +Posé par le scheduler cron (``cron/scheduler.py::run_job``, via le pont ``_jb_job_hooks``) au +lancement de chaque job, lu par le middleware (``middleware.py``) pour estampiller les +DraftRequest avec le « département » de la tâche : ``department`` / ``skill_id`` / ``job_id``. + +Pourquoi une ``ContextVar`` et pas une variable d'environnement : les jobs cron tournent dans des +THREADS du processus gateway (pool parallèle de ``tick()``) — ``os.environ`` est global au +processus, des jobs concurrents s'écraseraient mutuellement. Hermes a déjà migré l'état de +session vers des ContextVars pour cette raison (``gateway/session_context.py``) et propage le +contexte à chaque saut de thread (``copy_context`` dans ``_run_job_impl``, +``propagate_context_to_thread`` pour les outils) : une ContextVar posée dans ``run_job`` est donc +visible du middleware pendant TOUTE l'exécution du job, sans fuite entre jobs. + +La casquette vient du front-matter du skill attaché au job : champ ``casquette:`` (skills gold +Jean-Billie) ou ``department:`` (skills custom). Résolution best-effort : toute erreur → champs +absents, jamais d'exception — l'attribution ne doit JAMAIS faire échouer un job. +""" + +from __future__ import annotations + +import logging +import os +import re +from contextvars import ContextVar +from pathlib import Path +from typing import Any, Dict, List, Optional, Tuple + +logger = logging.getLogger(__name__) + +# Contexte du job cron courant ({job_id, skill_id, department, label}) ou None hors job. +_JOB_CTX: ContextVar[Optional[Dict[str, Any]]] = ContextVar("jb_job_ctx", default=None) + + +def current() -> Optional[Dict[str, Any]]: + """Contexte d'attribution du job courant, ou ``None`` hors job cron (chat libre).""" + return _JOB_CTX.get() + + +def job_started(job: Optional[Dict[str, Any]]) -> Optional[object]: + """Pose le contexte d'attribution au lancement d'un job cron. N'échoue jamais. + + Appelé par le scheduler. Renvoie le token de reset à repasser à ``job_finished`` + (ou ``None`` si le plugin est passif — ``JB_DECISION_PUSH_URL`` absent). + """ + try: + from . import config + + if not config.enabled(): + return None + token = _JOB_CTX.set(_build_ctx(job)) + return token + except Exception: + logger.debug("jb_outbound: job_started en échec (ignoré)", exc_info=True) + return None + + +def job_finished(job: Optional[Dict[str, Any]], success: bool = True, token: Optional[object] = None) -> None: + """Nettoie le contexte à la fin du job. N'échoue jamais. + + Le nettoyage est indispensable : les threads du pool cron sont RÉUTILISÉS — sans reset, un + job suivant sans skill hériterait de l'attribution du précédent. + """ + try: + if token is not None: + _JOB_CTX.reset(token) + else: + _JOB_CTX.set(None) + except Exception: + _JOB_CTX.set(None) + + +# --------------------------------------------------------------------------- +# Construction du contexte (job → {job_id, skill_id, department, label}) +# --------------------------------------------------------------------------- + +def _build_ctx(job: Optional[Dict[str, Any]]) -> Dict[str, Any]: + job = job or {} + skills = _skill_names(job) + department, skill_id = _resolve_department(skills) + return { + "job_id": str(job.get("id") or "").strip() or None, + "skill_id": skill_id, + "department": department, + "label": str(job.get("name") or "").strip() or None, + } + + +def _skill_names(job: Dict[str, Any]) -> List[str]: + """Skills du job, dans l'ordre (champ canonique ``skills``, repli legacy ``skill``).""" + raw = job.get("skills") + if raw is None: + raw = [job.get("skill")] if job.get("skill") else [] + elif isinstance(raw, str): + raw = [raw] + out: List[str] = [] + for item in raw if isinstance(raw, list) else []: + text = str(item or "").strip() + if text and text not in out: + out.append(text) + return out + + +def _resolve_department(skills: List[str]) -> Tuple[Optional[str], Optional[str]]: + """(department, skill_id) du job : premier skill qui déclare une casquette. + + ``casquette:`` (gold) est lu avant ``department:`` (custom). Si aucun skill ne déclare de + département, ``skill_id`` retombe sur le premier skill du job (attribution partielle). + """ + fallback = skills[0] if skills else None + for name in skills: + try: + fm = _skill_frontmatter(name) + except Exception: + continue + for key in ("casquette", "department"): + value = fm.get(key) + if isinstance(value, str) and value.strip(): + return value.strip(), name + return None, fallback + + +def _home() -> Path: + return Path(os.getenv("HERMES_HOME", str(Path.home() / ".hermes"))) + + +def _find_skill_md(name: str) -> Optional[Path]: + """Localise le fichier d'un skill sous ``/skills`` (miroir allégé de skills_tool). + + Stratégies : chemin direct (``name/SKILL.md``, couvre aussi ``catégorie/name``), fichier plat + ``name.md``, puis recherche récursive par nom de dossier. Refuse toute forme de traversée. + """ + if not name or ".." in name.replace("\\", "/").split("/") or Path(name).is_absolute() or Path(name).drive: + return None + skills_dir = _home() / "skills" + if not skills_dir.is_dir(): + return None + direct = skills_dir / name + if (direct / "SKILL.md").is_file(): + return direct / "SKILL.md" + if direct.with_suffix(".md").is_file(): + return direct.with_suffix(".md") + leaf = name.replace("\\", "/").split("/")[-1] + for cand in skills_dir.rglob("SKILL.md"): + if cand.parent.name == leaf: + return cand + for cand in skills_dir.rglob(f"{leaf}.md"): + if cand.name != "SKILL.md": + return cand + return None + + +def _skill_frontmatter(name: str) -> Dict[str, str]: + path = _find_skill_md(name) + if path is None: + return {} + try: + return _parse_simple_frontmatter(path.read_text(encoding="utf-8")) + except Exception: + return {} + + +def _parse_simple_frontmatter(text: str) -> Dict[str, str]: + """Extraction minimale du front-matter YAML : clés scalaires de premier niveau. + + Suffisant pour ``casquette:`` / ``department:`` — pas de dépendance yaml ni du cœur Hermes + (le plugin reste autonome, même esprit que ``http_client.py``). Les lignes indentées (blocs, + listes) sont ignorées. + """ + if not text.startswith("---"): + return {} + end = re.search(r"\n---\s*(\n|$)", text[3:]) + if not end: + return {} + out: Dict[str, str] = {} + for line in text[3 : end.start() + 3].splitlines(): + if not line.strip() or line[:1] in (" ", "\t") or ":" not in line: + continue + key, _, value = line.partition(":") + out[key.strip().lower()] = value.strip().strip("'\"") + return out diff --git a/plugins/jb_outbound/middleware.py b/plugins/jb_outbound/middleware.py index 78f8131f3400..e01306b2af6f 100644 --- a/plugins/jb_outbound/middleware.py +++ b/plugins/jb_outbound/middleware.py @@ -34,7 +34,7 @@ def jb_outbound_tool_execution( next_call: Callable[[Any], Any], **_: Any, ) -> Any: - from . import classify, config, http_client, mapping, store + from . import classify, config, http_client, job_context, mapping, store # Plugin passif hors box Jean-Billie (JB_DECISION_PUSH_URL non posé) : ne rien changer. if not config.enabled(): @@ -67,6 +67,15 @@ def jb_outbound_tool_execution( body = dict(draft) body["payload"] = {**draft.get("payload", {}), "jb_id": jb_id} + # Attribution : si l'interception a lieu pendant un job cron (skill → casquette), le draft + # porte le département. Champs ADDITIFS, omis hors contexte job (chat libre) — le daemon Go + # actuel ignore les champs inconnus (contrat répliqué côté Go en vague 2). + ctx = job_context.current() or {} + for key in ("department", "skill_id", "job_id"): + value = ctx.get(key) + if value: + body[key] = value + try: http_client.post_json(config.draft_url(), body) except Exception as exc: # dépôt impossible → on n'a rien envoyé, on le dit franchement. diff --git a/plugins/jb_outbound/test_attribution_activity.py b/plugins/jb_outbound/test_attribution_activity.py new file mode 100644 index 000000000000..1ef5ac93d96e --- /dev/null +++ b/plugins/jb_outbound/test_attribution_activity.py @@ -0,0 +1,184 @@ +"""Tests de l'attribution des drafts (stamp department/skill_id/job_id). + +Autonomes comme test_jb_outbound.py : HTTP loopback mocké, pas d'environnement Hermes complet. +Couvre : stamp présent en contexte job / absent hors contexte, résolution de la casquette depuis +le front-matter des skills (``casquette:`` gold, ``department:`` custom), et le pont +scheduler → plugin (cron/scheduler.py::run_job). +""" + +from __future__ import annotations + +import json +import sys +from pathlib import Path + +import pytest + +# Rendre le paquet `jb_outbound` importable comme paquet de premier niveau (plugins/ sur le path). +_PLUGINS_DIR = Path(__file__).resolve().parents[1] +if str(_PLUGINS_DIR) not in sys.path: + sys.path.insert(0, str(_PLUGINS_DIR)) + +import jb_outbound.http_client as http_client # noqa: E402 +import jb_outbound.job_context as job_context # noqa: E402 +import jb_outbound.middleware as middleware # noqa: E402 + + +@pytest.fixture(autouse=True) +def _isolate_env(tmp_path, monkeypatch): + """Isole HERMES_HOME et neutralise les gates JB pour CHAQUE test (baseline passive).""" + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + monkeypatch.delenv("JB_DECISION_PUSH_URL", raising=False) + # Filet de sécurité : jamais de contexte résiduel d'un test précédent. + job_context._JOB_CTX.set(None) + + +@pytest.fixture +def posts(monkeypatch): + """Active la boucle de proposition et capture tous les POST loopback.""" + monkeypatch.setenv("JB_DECISION_PUSH_URL", "http://127.0.0.1:8444/jb/decision") + monkeypatch.setenv("JB_DRAFT_ADDR", "127.0.0.1:8442") + captured: list = [] + monkeypatch.setattr( + http_client, "post_json", + lambda url, payload, timeout=10.0: captured.append((url, payload)) or 200, + ) + return captured + + +def _write_skill(tmp_path, name: str, body: str) -> None: + d = tmp_path / "skills" / name + d.mkdir(parents=True, exist_ok=True) + (d / "SKILL.md").write_text(body, encoding="utf-8") + + +_GOLD = "---\nname: relance-devis\ncasquette: Le Commercial\n---\n\n# Relance devis\n" +_CUSTOM = "---\nname: veille-presse\ndepartment: Le Marketing\n---\n\n# Veille presse\n" +_SANS = "---\nname: tri-boite\ndescription: Trie la boîte mail.\n---\n\n# Tri\n" + + +def _job(**over) -> dict: + job = {"id": "a1b2c3d4e5f6", "name": "Relances du matin", "skills": ["relance-devis"]} + job.update(over) + return job + + +def _propose(tool: str = "send_message", args: dict | None = None) -> dict: + """Passe un appel d'envoi dans le middleware (court-circuit attendu) et rend le résultat.""" + def next_call(_a): + raise AssertionError("l'outil d'envoi NE doit PAS s'exécuter avant validation") + + return json.loads( + middleware.make_middleware()( + tool_name=tool, args=args or {"chat_id": "1", "content": "Hi"}, next_call=next_call + ) + ) + + +def _drafts(posts) -> list: + return [p[1] for p in posts if p[0].endswith("/v1/draft")] + + +# --------------------------------------------------------------------------- +# Stamp d'attribution sur les drafts +# --------------------------------------------------------------------------- + +def test_stamp_draft_en_contexte_job(posts, tmp_path): + _write_skill(tmp_path, "relance-devis", _GOLD) + token = job_context.job_started(_job()) + + _propose() + draft = _drafts(posts)[-1] + assert draft["department"] == "Le Commercial" + assert draft["skill_id"] == "relance-devis" + assert draft["job_id"] == "a1b2c3d4e5f6" + # L'attribution est portée au premier niveau du DraftRequest, pas dans le payload round-trip. + assert "department" not in draft["payload"] + + job_context.job_finished(_job(), success=True, token=token) + + +def test_stamp_absent_hors_contexte(posts): + _propose() + draft = _drafts(posts)[-1] + for key in ("department", "skill_id", "job_id"): + assert key not in draft # chat libre → champs OMIS, pas de null + + +def test_stamp_efface_apres_job(posts, tmp_path): + _write_skill(tmp_path, "relance-devis", _GOLD) + token = job_context.job_started(_job()) + job_context.job_finished(_job(), success=True, token=token) + + _propose() + assert "department" not in _drafts(posts)[-1] + assert job_context.current() is None + + +def test_casquette_gold_prioritaire_sur_department(posts, tmp_path): + # Un skill qui porte les DEUX champs : `casquette:` (gold) gagne. + _write_skill(tmp_path, "relance-devis", "---\ncasquette: Le Commercial\ndepartment: Autre\n---\n") + token = job_context.job_started(_job()) + assert job_context.current()["department"] == "Le Commercial" + job_context.job_finished(_job(), token=token) + + +def test_department_custom(posts, tmp_path): + _write_skill(tmp_path, "veille-presse", _CUSTOM) + token = job_context.job_started(_job(skills=["veille-presse"])) + ctx = job_context.current() + assert ctx["department"] == "Le Marketing" + assert ctx["skill_id"] == "veille-presse" + job_context.job_finished(_job(), token=token) + + +def test_skill_sans_casquette_stamp_partiel(posts, tmp_path): + _write_skill(tmp_path, "tri-boite", _SANS) + token = job_context.job_started(_job(skills=["tri-boite"])) + + _propose() + draft = _drafts(posts)[-1] + assert "department" not in draft # pas de casquette déclarée → champ omis + assert draft["skill_id"] == "tri-boite" # attribution partielle conservée + assert draft["job_id"] == "a1b2c3d4e5f6" + + job_context.job_finished(_job(), token=token) + + +def test_resolution_skill_en_categorie(posts, tmp_path): + # Skill rangé sous une catégorie (ex. casquettes/relance-devis), référencé par nom nu. + _write_skill(tmp_path, "casquettes/relance-devis", _GOLD) + token = job_context.job_started(_job()) + assert job_context.current()["department"] == "Le Commercial" + job_context.job_finished(_job(), token=token) + + +def test_passif_sans_box(): + # JB_DECISION_PUSH_URL absent → job_started est un no-op total (plugin passif). + assert job_context.job_started(_job()) is None + assert job_context.current() is None + + +# --------------------------------------------------------------------------- +# Pont scheduler → plugin (cron/scheduler.py::run_job) +# --------------------------------------------------------------------------- + +def test_scheduler_run_job_pose_et_nettoie_le_contexte(posts, tmp_path, monkeypatch): + """run_job (scheduler réel) pose le contexte pendant le job et le nettoie après.""" + _write_skill(tmp_path, "relance-devis", _GOLD) + + import cron.scheduler as scheduler + + seen = {} + + def _fake_impl(job): + seen["ctx"] = job_context.current() # visible PENDANT l'exécution du job + return True, "doc", "réponse", None + + monkeypatch.setattr(scheduler, "_run_job_impl", _fake_impl) + result = scheduler.run_job(_job()) + + assert result[0] is True + assert seen["ctx"]["department"] == "Le Commercial" + assert seen["ctx"]["job_id"] == "a1b2c3d4e5f6" + assert job_context.current() is None # nettoyé après le job (threads de pool réutilisés) From 241d110a18600266fe20f9b669bfabdfb2122733 Mon Sep 17 00:00:00 2001 From: Gabiar <189609572+gabiar-maker@users.noreply.github.com> Date: Fri, 12 Jun 2026 15:39:19 +0200 Subject: [PATCH 2/3] =?UTF-8?q?feat(jb=5Foutbound):=20fil=20d=20activite?= =?UTF-8?q?=20=E2=80=94=20signaux=20debut/fin=20de=20job=20cron?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit POST fire-and-forget vers http://{JB_DRAFT_ADDR}/v1/activity au debut et a la fin de chaque job cron (et uniquement cron) : {phase: started|finished, status: ok|error, department?, skill_id?, job_id?, label?} — label = nom lisible du job (jobs.json), champs d attribution OMIS quand absents (meme convention que le stamp des drafts). - GATED par JB_ACTIVITY_EVENTS=1, defaut OFF : la route daemon n existe pas encore (vague 2). Timeout 2 s, toute exception avalee (log debug au plus) — un signal d activite ne bloque ni ne fait echouer JAMAIS un job. - activity.py : gate + emission ; config.py : activity_url() ; job_context.job_started/job_finished : emission cablée sur le cycle de vie du job (status ok au started, ok/error au finished selon le succes du run). Tests : +6 (gate OFF par defaut, started→finished avec payload complet, status error, echec reseau avale sans casser le contexte, emit sans contexte, pont scheduler reel ok/error) ; suite 27/27. Co-Authored-By: Claude Fable 5 --- plugins/jb_outbound/activity.py | 45 ++++++++ plugins/jb_outbound/config.py | 5 + plugins/jb_outbound/job_context.py | 36 ++++--- .../jb_outbound/test_attribution_activity.py | 102 ++++++++++++++++-- 4 files changed, 167 insertions(+), 21 deletions(-) create mode 100644 plugins/jb_outbound/activity.py diff --git a/plugins/jb_outbound/activity.py b/plugins/jb_outbound/activity.py new file mode 100644 index 000000000000..17a57d06fdcf --- /dev/null +++ b/plugins/jb_outbound/activity.py @@ -0,0 +1,45 @@ +"""Fil d'activité Jean-Billie : signaux début/fin de job cron (fire-and-forget). + +POST ``http://{JB_DRAFT_ADDR}/v1/activity`` avec +``{phase: "started"|"finished", status: "ok"|"error", department?, skill_id?, job_id?, label?}``. +``label`` = nom lisible du job (``jobs.json``). Les champs d'attribution absents sont OMIS du +JSON (pas de ``null``) — même convention que le stamp des DraftRequest. + +La route daemon n'existe PAS ENCORE (vague 2), d'où le gate ``JB_ACTIVITY_EVENTS=1`` +(défaut OFF). Timeout court, toute erreur avalée (log debug au plus) : un signal d'activité ne +doit JAMAIS bloquer ni faire échouer un job. +""" + +from __future__ import annotations + +import logging +import os +from typing import Any, Dict, Optional + +logger = logging.getLogger(__name__) + +# Timeout volontairement court : le daemon est sur le loopback, et un signal d'activité ne vaut +# pas la peine de retenir un job plus de 2 s. +_TIMEOUT = 2.0 + + +def enabled() -> bool: + """Vrai uniquement quand l'opérateur a activé le fil d'activité (``JB_ACTIVITY_EVENTS=1``).""" + return os.getenv("JB_ACTIVITY_EVENTS", "").strip() == "1" + + +def emit(phase: str, status: str, ctx: Optional[Dict[str, Any]]) -> None: + """Émet un évènement d'activité (best-effort). Silencieux si le gate est fermé ou en échec.""" + if not enabled(): + return + try: + from . import config, http_client + + event: Dict[str, Any] = {"phase": phase, "status": status} + for key in ("department", "skill_id", "job_id", "label"): + value = (ctx or {}).get(key) + if value: + event[key] = value + http_client.post_json(config.activity_url(), event, timeout=_TIMEOUT) + except Exception as exc: + logger.debug("jb_outbound: signal d'activité non délivré (%s/%s) : %s", phase, status, exc) diff --git a/plugins/jb_outbound/config.py b/plugins/jb_outbound/config.py index f4a76a3545db..e0cdcb5ce0a1 100644 --- a/plugins/jb_outbound/config.py +++ b/plugins/jb_outbound/config.py @@ -30,6 +30,11 @@ def result_url() -> str: return f"http://{_draft_addr()}/v1/result" +def activity_url() -> str: + """URL où POSTer un évènement d'activité (début/fin de job — cf. activity.py).""" + return f"http://{_draft_addr()}/v1/activity" + + def _push_url() -> str: return os.getenv("JB_DECISION_PUSH_URL", _DEFAULT_PUSH_URL).strip() or _DEFAULT_PUSH_URL diff --git a/plugins/jb_outbound/job_context.py b/plugins/jb_outbound/job_context.py index bd4c0dccc125..ea070d34d45e 100644 --- a/plugins/jb_outbound/job_context.py +++ b/plugins/jb_outbound/job_context.py @@ -38,17 +38,20 @@ def current() -> Optional[Dict[str, Any]]: def job_started(job: Optional[Dict[str, Any]]) -> Optional[object]: - """Pose le contexte d'attribution au lancement d'un job cron. N'échoue jamais. + """Pose le contexte d'attribution et signale le début du job. N'échoue jamais. - Appelé par le scheduler. Renvoie le token de reset à repasser à ``job_finished`` - (ou ``None`` si le plugin est passif — ``JB_DECISION_PUSH_URL`` absent). + Appelé par le scheduler au lancement d'un job cron. Renvoie le token de reset à repasser + à ``job_finished`` (ou ``None`` si le plugin est passif : ni boucle de proposition + ``JB_DECISION_PUSH_URL``, ni fil d'activité ``JB_ACTIVITY_EVENTS``). """ try: - from . import config + from . import activity, config - if not config.enabled(): + if not (config.enabled() or activity.enabled()): return None - token = _JOB_CTX.set(_build_ctx(job)) + ctx = _build_ctx(job) + token = _JOB_CTX.set(ctx) + activity.emit("started", "ok", ctx) return token except Exception: logger.debug("jb_outbound: job_started en échec (ignoré)", exc_info=True) @@ -56,18 +59,27 @@ def job_started(job: Optional[Dict[str, Any]]) -> Optional[object]: def job_finished(job: Optional[Dict[str, Any]], success: bool = True, token: Optional[object] = None) -> None: - """Nettoie le contexte à la fin du job. N'échoue jamais. + """Signale la fin du job et nettoie le contexte. N'échoue jamais. Le nettoyage est indispensable : les threads du pool cron sont RÉUTILISÉS — sans reset, un job suivant sans skill hériterait de l'attribution du précédent. """ try: - if token is not None: - _JOB_CTX.reset(token) - else: - _JOB_CTX.set(None) + from . import activity + + if activity.enabled(): + ctx = _JOB_CTX.get() or _build_ctx(job) + activity.emit("finished", "ok" if success else "error", ctx) except Exception: - _JOB_CTX.set(None) + logger.debug("jb_outbound: job_finished en échec (ignoré)", exc_info=True) + finally: + try: + if token is not None: + _JOB_CTX.reset(token) + else: + _JOB_CTX.set(None) + except Exception: + _JOB_CTX.set(None) # --------------------------------------------------------------------------- diff --git a/plugins/jb_outbound/test_attribution_activity.py b/plugins/jb_outbound/test_attribution_activity.py index 1ef5ac93d96e..d4e7e0c15f0d 100644 --- a/plugins/jb_outbound/test_attribution_activity.py +++ b/plugins/jb_outbound/test_attribution_activity.py @@ -1,9 +1,9 @@ -"""Tests de l'attribution des drafts (stamp department/skill_id/job_id). +"""Tests de l'attribution (stamp department/skill_id/job_id) et du fil d'activité. Autonomes comme test_jb_outbound.py : HTTP loopback mocké, pas d'environnement Hermes complet. Couvre : stamp présent en contexte job / absent hors contexte, résolution de la casquette depuis -le front-matter des skills (``casquette:`` gold, ``department:`` custom), et le pont -scheduler → plugin (cron/scheduler.py::run_job). +le front-matter des skills (``casquette:`` gold, ``department:`` custom), gate ``JB_ACTIVITY_EVENTS``, +innocuité des échecs réseau, et le pont scheduler → plugin (cron/scheduler.py::run_job). """ from __future__ import annotations @@ -19,6 +19,7 @@ if str(_PLUGINS_DIR) not in sys.path: sys.path.insert(0, str(_PLUGINS_DIR)) +import jb_outbound.activity as activity # noqa: E402 import jb_outbound.http_client as http_client # noqa: E402 import jb_outbound.job_context as job_context # noqa: E402 import jb_outbound.middleware as middleware # noqa: E402 @@ -29,13 +30,14 @@ def _isolate_env(tmp_path, monkeypatch): """Isole HERMES_HOME et neutralise les gates JB pour CHAQUE test (baseline passive).""" monkeypatch.setenv("HERMES_HOME", str(tmp_path)) monkeypatch.delenv("JB_DECISION_PUSH_URL", raising=False) + monkeypatch.delenv("JB_ACTIVITY_EVENTS", raising=False) # Filet de sécurité : jamais de contexte résiduel d'un test précédent. job_context._JOB_CTX.set(None) @pytest.fixture def posts(monkeypatch): - """Active la boucle de proposition et capture tous les POST loopback.""" + """Active la boucle de proposition et capture tous les POST loopback (drafts + activité).""" monkeypatch.setenv("JB_DECISION_PUSH_URL", "http://127.0.0.1:8444/jb/decision") monkeypatch.setenv("JB_DRAFT_ADDR", "127.0.0.1:8442") captured: list = [] @@ -79,8 +81,12 @@ def _drafts(posts) -> list: return [p[1] for p in posts if p[0].endswith("/v1/draft")] +def _activities(posts) -> list: + return [p[1] for p in posts if p[0].endswith("/v1/activity")] + + # --------------------------------------------------------------------------- -# Stamp d'attribution sur les drafts +# Tâche 1 — stamp d'attribution sur les drafts # --------------------------------------------------------------------------- def test_stamp_draft_en_contexte_job(posts, tmp_path): @@ -153,18 +159,79 @@ def test_resolution_skill_en_categorie(posts, tmp_path): job_context.job_finished(_job(), token=token) -def test_passif_sans_box(): - # JB_DECISION_PUSH_URL absent → job_started est un no-op total (plugin passif). +def test_passif_sans_box_ni_activite(): + # Ni JB_DECISION_PUSH_URL ni JB_ACTIVITY_EVENTS → job_started est un no-op total. assert job_context.job_started(_job()) is None assert job_context.current() is None +# --------------------------------------------------------------------------- +# Tâche 2 — fil d'activité (début/fin de job), gated par JB_ACTIVITY_EVENTS +# --------------------------------------------------------------------------- + +def test_activity_off_par_defaut(posts, tmp_path): + _write_skill(tmp_path, "relance-devis", _GOLD) + token = job_context.job_started(_job()) + job_context.job_finished(_job(), success=True, token=token) + assert _activities(posts) == [] # gate fermé → aucun évènement + + +def test_activity_on_emet_started_puis_finished(posts, tmp_path, monkeypatch): + monkeypatch.setenv("JB_ACTIVITY_EVENTS", "1") + _write_skill(tmp_path, "relance-devis", _GOLD) + + token = job_context.job_started(_job()) + job_context.job_finished(_job(), success=True, token=token) + + events = _activities(posts) + assert [e["phase"] for e in events] == ["started", "finished"] + started, finished = events + assert started["status"] == "ok" and finished["status"] == "ok" + for e in events: + assert e["department"] == "Le Commercial" + assert e["skill_id"] == "relance-devis" + assert e["job_id"] == "a1b2c3d4e5f6" + assert e["label"] == "Relances du matin" # nom lisible du job (jobs.json) + + +def test_activity_status_error_si_echec(posts, monkeypatch): + monkeypatch.setenv("JB_ACTIVITY_EVENTS", "1") + token = job_context.job_started(_job(skills=[])) + job_context.job_finished(_job(skills=[]), success=False, token=token) + + finished = _activities(posts)[-1] + assert finished["phase"] == "finished" and finished["status"] == "error" + assert "department" not in finished and "skill_id" not in finished # job sans skill → omis + + +def test_activity_echec_reseau_avale(monkeypatch, tmp_path): + # Daemon injoignable (route /v1/activity inexistante, conteneur down…) : le job continue. + monkeypatch.setenv("JB_ACTIVITY_EVENTS", "1") + _write_skill(tmp_path, "relance-devis", _GOLD) + + def _boom(url, payload, timeout=10.0): + raise ConnectionError("connexion refusée") + + monkeypatch.setattr(http_client, "post_json", _boom) + token = job_context.job_started(_job()) # ne lève pas + assert job_context.current()["department"] == "Le Commercial" # le contexte reste posé + job_context.job_finished(_job(), success=True, token=token) # ne lève pas + assert job_context.current() is None + + +def test_activity_emit_sans_contexte_ne_leve_pas(posts, monkeypatch): + monkeypatch.setenv("JB_ACTIVITY_EVENTS", "1") + activity.emit("started", "ok", None) + assert _activities(posts) == [{"phase": "started", "status": "ok"}] + + # --------------------------------------------------------------------------- # Pont scheduler → plugin (cron/scheduler.py::run_job) # --------------------------------------------------------------------------- -def test_scheduler_run_job_pose_et_nettoie_le_contexte(posts, tmp_path, monkeypatch): - """run_job (scheduler réel) pose le contexte pendant le job et le nettoie après.""" +def test_scheduler_run_job_pose_contexte_et_signale(posts, tmp_path, monkeypatch): + """run_job (scheduler réel) pose le contexte pendant le job et émet started/finished.""" + monkeypatch.setenv("JB_ACTIVITY_EVENTS", "1") _write_skill(tmp_path, "relance-devis", _GOLD) import cron.scheduler as scheduler @@ -182,3 +249,20 @@ def _fake_impl(job): assert seen["ctx"]["department"] == "Le Commercial" assert seen["ctx"]["job_id"] == "a1b2c3d4e5f6" assert job_context.current() is None # nettoyé après le job (threads de pool réutilisés) + assert [e["phase"] for e in _activities(posts)] == ["started", "finished"] + assert _activities(posts)[-1]["status"] == "ok" + + +def test_scheduler_run_job_echec_signale_error(posts, tmp_path, monkeypatch): + monkeypatch.setenv("JB_ACTIVITY_EVENTS", "1") + _write_skill(tmp_path, "relance-devis", _GOLD) + + import cron.scheduler as scheduler + + monkeypatch.setattr(scheduler, "_run_job_impl", lambda job: (False, "doc", "", "boom")) + result = scheduler.run_job(_job()) + + assert result[0] is False + finished = _activities(posts)[-1] + assert finished["phase"] == "finished" + assert finished["status"] == "error" From 2c9067f0a436e3fd8dc2fb52963e83cb05fbd8eb Mon Sep 17 00:00:00 2001 From: Gabiar <189609572+gabiar-maker@users.noreply.github.com> Date: Fri, 12 Jun 2026 15:40:09 +0200 Subject: [PATCH 3/3] =?UTF-8?q?docs(jb=5Foutbound):=20README=20=E2=80=94?= =?UTF-8?q?=20attribution=20des=20departements=20+=20fil=20d=20activite?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Documente le stamp des DraftRequest (department/skill_id/job_id), les signaux /v1/activity et le gate JB_ACTIVITY_EVENTS ; precise que la greffe compte desormais un point optionnel dans cron/scheduler.py (no-op sans plugin) et la commande de test Windows. Co-Authored-By: Claude Fable 5 --- plugins/jb_outbound/README.md | 29 ++++++++++++++++++++++++----- 1 file changed, 24 insertions(+), 5 deletions(-) diff --git a/plugins/jb_outbound/README.md b/plugins/jb_outbound/README.md index f2e27e9e8a80..3836cdcce4ca 100644 --- a/plugins/jb_outbound/README.md +++ b/plugins/jb_outbound/README.md @@ -1,7 +1,8 @@ # jb_outbound — « rien ne part sans accord » (Jean-Billie) -Greffe Jean-Billie sur **Hermes Agent** (Nous Research, MIT). **Un seul plugin, zéro modification du -cœur** → suivi de l'upstream trivial. +Greffe Jean-Billie sur **Hermes Agent** (Nous Research, MIT). Un seul plugin + **un point de +greffe optionnel dans le scheduler cron** (`cron/scheduler.py::run_job`, no-op sans le plugin) +→ suivi de l'upstream trivial. ## Ce que ça fait @@ -28,6 +29,22 @@ outil d'envoi appelé → POST ResultRequest {id, executed|failed} → http://127.0.0.1:8442/v1/result ``` +## Attribution (départements) & fil d'activité + +Au lancement d'un job cron, le scheduler pose le **contexte d'attribution** du job +(`job_context.py`, ContextVar — les jobs tournent dans des threads du gateway) : casquette lue +dans le front-matter du skill du job (`casquette:` pour les skills gold, `department:` pour les +customs), id du skill, id du job. + +- **Stamp des drafts** : tout DraftRequest émis pendant un job porte les champs additifs de + premier niveau `department`, `skill_id`, `job_id` (omis hors contexte job — chat libre). Le + daemon ignore les champs inconnus tant que le contrat Go n'est pas étendu (vague 2). +- **Fil d'activité** (`activity.py`) : au début et à la fin de chaque job cron, POST + fire-and-forget `http://{JB_DRAFT_ADDR}/v1/activity` avec + `{phase: "started"|"finished", status: "ok"|"error", department?, skill_id?, job_id?, label?}` + (`label` = nom lisible du job). **Gated par `JB_ACTIVITY_EVENTS=1`** (défaut OFF — la route + daemon n'existe pas encore). Timeout 2 s, échecs avalés : ne bloque jamais un job. + ## Règles - **Fail-closed** : un outil d'envoi composio non répertorié est **bloqué** (jamais auto-envoyé). On @@ -49,9 +66,11 @@ plugins: Endpoints lus dans l'environnement (posés par le bundle Jean-Billie / le `daemon.env`) : `JB_DRAFT_ADDR` (défaut `127.0.0.1:8442`), `JB_DECISION_PUSH_URL` (défaut -`http://127.0.0.1:8444/jb/decision`). Sans `JB_DECISION_PUSH_URL`, le plugin reste **passif**. +`http://127.0.0.1:8444/jb/decision`), `JB_ACTIVITY_EVENTS` (`1` pour activer le fil d'activité, +défaut OFF). Sans `JB_DECISION_PUSH_URL`, le plugin reste **passif**. ## Tests -`python -m pytest plugins/jb_outbound/test_jb_outbound.py` — autonome (mocke le HTTP loopback et le -registre d'outils, n'a pas besoin d'un environnement Hermes complet). +`python -m pytest plugins/jb_outbound/` — autonome (mocke le HTTP loopback et le registre +d'outils, n'a pas besoin d'un environnement Hermes complet). Sous Windows : +`pytest -o addopts=""` (pytest-timeout/SIGALRM indisponible).