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
42 changes: 41 additions & 1 deletion cron/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.<slug>`` ; 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]]:
Expand Down
29 changes: 24 additions & 5 deletions plugins/jb_outbound/README.md
Original file line number Diff line number Diff line change
@@ -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

Expand All @@ -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
Expand All @@ -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).
45 changes: 45 additions & 0 deletions plugins/jb_outbound/activity.py
Original file line number Diff line number Diff line change
@@ -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)
5 changes: 5 additions & 0 deletions plugins/jb_outbound/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
193 changes: 193 additions & 0 deletions plugins/jb_outbound/job_context.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,193 @@
"""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 et signale le début du job. N'échoue jamais.

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 activity, config

if not (config.enabled() or activity.enabled()):
return None
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)
return None


def job_finished(job: Optional[Dict[str, Any]], success: bool = True, token: Optional[object] = None) -> None:
"""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:
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:
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)


# ---------------------------------------------------------------------------
# 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 ``<HERMES_HOME>/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
11 changes: 10 additions & 1 deletion plugins/jb_outbound/middleware.py
Original file line number Diff line number Diff line change
Expand Up @@ -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():
Expand Down Expand Up @@ -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.
Expand Down
Loading
Loading