Un WoT fan out (le map) ; un gather ramène les branches (le
reduce). Tout se compose sur deux axes orthogonaux — une discipline de collecte × un kernel
de fusion — chaque enfant corrélé par un fan_id stable. Les scènes sont en 3D : glisse
pour tourner, molette pour zoomer.
Un WoT fan out (le map) ; un gather ramène les branches (le
reduce). Tout se compose sur deux axes : une discipline de collecte (§1) × un
kernel (§2). Chaque enfant porte un fan_id stable — c'est lui qui corrèle le map à
sa collecte. Les scènes sont en 3D : glisse pour tourner, molette pour zoomer.
Comment on ramène les branches : adjacent (câblé après le fan), ou des
collecteurs corrélés par fan_id — perfect (ferme à N) vs latent
(agglomère, tiré à l'étape n), et ambient (plusieurs collecteurs, sans câblage).
append rassemble (concat). Puis comment on traite la donnée en continu — DEUX stratégies opposées : summary (fenêtre glissante : on compresse) ⟷ transform (on traite / calcule). rlm = un contexte que des agents traitent en continu (façon summary, mais agent-driven). okf = par référence (§3). Flux à gauche → kernel → résultat à droite.
Les 4 designs ci-dessus sont des disciplines de collecte (le dispatch). L'OKF,
lui, est un kernel — il change ce que le reduce RENVOIE quand le résultat est gros. Inliner
1,24 M lignes de prêts dans le contexte de l'agent le fait exploser. OKF n'en renvoie que la poignée :
la ref + le schéma évolutif (nb de lignes, une description par champ). La donnée
reste dans le store ; l'agent en tire des parties à la demande.
# OKF — le reduce renvoie une REF + un schéma évolutif, JAMAIS la donnée handle = {"ref": "okf:claims#v1", "schema": {"rows": 1_240_000, "columns": {"branch": {"type": "str", "desc": "agence"}, "amount": {"type": "float", "desc": "encours €"}}}} set_okf_backend(ClickHouseOkfStore(host="ch")) # le store = une vraie table ClickHouse okf_read("okf:claims#v1", fields=["branch", "amount"], rows="0:50") # tranche -> SELECT ... LIMIT 50 okf_stats("okf:claims#v1", field="amount") # agregat server-side, pas les lignes
La balise schéma-only montre à l'agent le schéma, jamais la donnée :
<tckm domain="okf"/> → okf:claims#v1 — 1 240 000 rows; branch(agence), amount(encours €).
Où l'OKF gagne : le compliance qui doit vérifier des chiffres sur un résultat d'entrepôt énorme
(il fait okf_stats au lieu de charger le million de lignes) · passer un dataset entre agents
sans le re-sérialiser · un chart builder qui n'a besoin que du Top-20 d'un million de lignes.
| Discipline | Exemple concret | Le plus adapté quand |
|---|---|---|
| Adjacent join | un sinistre éclaté en fraude / couverture / montant → les 3 analyses jointes en bout de chaîne | le reduce est FIXE, juste après le fan ; le résultat transite tout de suite |
| Perfect · close at N | 8 agences analysées en parallèle → ferme au 8e, sort la liste des 8 | tu connais le nombre de branches et tu veux la barrière (liste complète, agrégation finale) |
| Latent · pull | un KPI de risque qui se met à jour pendant que les agences arrivent, lu par le nœud report à la fin | la taille est ouverte, ou tu veux une vue live/partielle consommée à une étape choisie |
| Ambient · multi | un même fan d'agences nourrit EN PARALLÈLE : la liste (perfect), un résumé glissant (window) et le dataset (okf) | plusieurs vues d'un même fan-out, ou un collecteur loin du fan, sans câblage |
| OKF · ref | la réduction = 1,24 M lignes de prêts → l'agent voit le schéma et fait okf_read(rows="0:50") / okf_stats(field="amount") |
le résultat est GROS et un agent doit le query, pas le lire en entier |
| Kernel | Exemple concret | Le plus adapté quand |
|---|---|---|
| append | les findings des 4 workers → ["f0","f1","f2","f3"] |
le défaut — tu veux juste toutes les parties ensemble |
| summary (window) | 12 events, fenêtre 3 → 4 lots repliés en UN résumé exécutif glissant | stratégie A — COMPRESSER un long flux au fil de l'eau (un exec-summary live), pas une liste |
| transform | somme courante des montants au fil des workers (latent) ; ou un calcul unique en fin de lot (end) | stratégie B (opposée à summary) — TRAITER / calculer la donnée, pas la compresser |
| rlm | les findings de chercheurs parallèles s'accumulent dans le ContextRLMRecord d'une session pour un tour suivant |
le reduce doit nourrir un contexte conversationnel/RLM géré ailleurs |
| okf | un dataset d'1,24 M lignes → une ref + schéma, lu par tranches | l'objet réduit est gros et se query (cf. ci-dessus) |
Un WoT fan out : un événement se disperse en N enfants (le map). Le
reduce les ramène. On ne code pas six logiques : on décompose sur deux axes orthogonaux
et on compose une discipline de collecte × un kernel de fusion. Chaque enfant d'un
fan-out porte un fan_id stable — c'est lui qui corrèle la dispersion à sa collecte.
adjacent = reduce entre deux nœuds voisins (le join actuel). gather =
des collecteurs placés plus loin dans le graphe, corrélés par fan_id — plusieurs
collecteurs peuvent pointer sur le même fan-out.
src.fan(into=[worker], fan_id="claims", expect="auto") # MAP : fan_id porté sur chaque enfant
worker >> Reduce("brief", kernel="append") # collecte ADJACENTE (le join actuel)
Gather("collect", fan_id="claims", n=8, mode="perfect") # COLLECTEUR distant, ferme à 8, transite
Gather("pool", fan_id="claims", mode="latent") # COLLECTEUR ouvert, tiré à l'étape n
Le kernel ne connaît PAS la discipline de collecte : on le branche sur n'importe
quelle collecte. Cinq kernels — append, fenêtre glissante, transform,
rlm, okf.
Un reduce concret = (mode de collecte) × (kernel). Ex. :
gather(latent) × window(3), gather(perfect, n=8) × append,
adjacent × okf.
Perfect : un nombre attendu n ; quand n events du même
fan_id sont arrivés, le collecteur ferme et transite en aval (barrière push).
Latent : aucune fermeture ; un StatefulRecord agglomère en continu et une
étape n le tire (pull).
mode="perfect", n=8 # ferme à 8, puis transite → généralise MergeNode.wait_n
mode="latent" # Record gather:claims agglomère ; "report" fait bag_read("gather:claims")
| Kernel | Ce qu'il fait | Syntaxe |
|---|---|---|
| append | concat / dict-merge (conflit détecté) | kernel="append" |
| fenêtre glissante | résumé progressif : new = f(résumé, N derniers) | kernel=Window(size=3, summarize=…) |
| transform | handler : latent (incrémental) | end (tout arrivé) | append | kernel=Transform(fn, when="latent") |
| rlm | ingéré en latent, géré par la logique du RLM context | kernel=Rlm(consume_at="report") |
| okf | ref vers un objet sérialisé + schéma + accès partiel | kernel=Okf(store=…, partial=True) |
Traitement de sinistres : on éclate par dossier, puis on ramasse de trois façons différentes en parallèle — un KPI qui se résume au fil de l'eau, la liste complète, et un gros dataset passé par référence.
hive = Hive(app, name="claims-wot", type="wot", tracer=tracer)
triage, assess, report = hive.stage("triage", triage_fn), assess_agent, report_agent
(hive.from_input()
>> triage
.fan(into=[assess], fan_id="claims", expect="auto") # MAP : 1 dossier → N
# (1) collecteur LATENT + fenêtre glissante — KPI résumé au fil de l'eau
>> Gather("kpi_pool", fan_id="claims", mode="latent",
kernel=Window(size=3, summarize=SUMMARY_PROMPT), consume_at="report")
# (2) collecteur PERFECT + append — liste complète, ferme à N, transite
>> Gather("all_claims", fan_id="claims", mode="perfect", kernel="append")
# (3) collecteur PERFECT + OKF — gros dataset → ref + schéma, pas d'inline
>> Gather("dataset", fan_id="claims", mode="perfect",
kernel=Okf(store="clickhouse", partial=True))
>> report.to_loop()) # report lit ce dont il a besoin, sans overflow
Quand le kernel est Okf, le reduce n'inline pas la donnée : il écrit
l'objet lourd dans un store (réutilise NodeRecord/resolve_node) et pose dans le
bag un item = la ref + le schéma évolutif (nb de lignes, une description par champ) + une
interface d'interaction. On récupère alors seulement certaines parties — pas d'overflow.
# item OKF dans le bag (compact — ce que l'agent VOIT, jamais la donnée) :
{ "ref": "okf:claims#v3",
"schema": { "rows": 1240,
"columns": { "branch": {"type":"str", "desc":"agence"},
"risk": {"type":"float", "desc":"score 0-1"} } } }
# rendu en contexte via un frame schéma-only :
<tckm domain="okf" ref="okf:claims#v3"/> → "okf:claims#v3 — 1240 lignes; branch, risk"
# l'interface d'interaction — on récupère QUE une tranche, jamais tout :
okf_read("okf:claims#v3", fields=["branch","risk"], rows="0:50") # tranche
okf_read("okf:claims#v3", where="risk>0.8", limit=20) # filtre
okf_stats("okf:claims#v3", field="amount") # agrégat, pas les lignes
Déjà là : le join adjacent
(_node_join, MergeNode), le digest progressif (dry_context), le
.transform par-event, l'agglomération durable (ContextRLMRecord), le
by-ref (NodeRecord/resolve_node) et le schéma compact (BAG_INFO +
render_bag). Neuf : un fan_id stable multi-hop,
le gather clé par fan_id, l'enum latent/perfect, plusieurs collecteurs par
fan-out, le record OKF unifié + okf_read partiel, la fenêtre à N messages, le gate
consume-at-step-n. Le reste est du câblage déclaratif ; le framework exécute la barrière.
fan_id + gather — spec d'implémentationLe pré-requis de tout l'axe « gather ». Deux moitiés : fan_id (un id de fan-out stable qui survit à tous les hops) et gather (le collecteur clé par ce fan_id, en latent ou perfect). Aucun des deux n'existe aujourd'hui — voici le plan, ancré sur l'existant.
Problème. Aujourd'hui seul WotRecord.parent_id existe
(wot_record.py:29) — lineage un seul hop : un collecteur 5 hops plus loin ne sait
pas de quel fan-out vient l'event. Solution : une pile de frames portée par le record,
minte au fan, recopiée à chaque hop.
# spec : une PILE (gère les fans imbriqués ET le multi-collecteur)
@dataclass
class FanFrame:
fan_id: str # id stable du fan-out
seq: int # index de cet enfant (0..n-1)
expected: int # n attendu
class WotRecord(StatefulRecord):
fan_ctx: list[FanFrame] = field(default_factory=list) # top = fan le plus interne
Minting + stamping — là où parent_id est déjà posé
(deployer.py:428) :
fid = f"{node}#{parent_event_id}" # déterministe → resume-safe
for i, child in enumerate(items):
child.fan_ctx = [*parent.fan_ctx, FanFrame(fid, i, len(items))] # hérite + pousse
Survie aux hops — LE point critique. Chaque couture qui produit un nouveau record recopie la pile :
def carry_fan(src, dst):
if src.fan_ctx and not dst.fan_ctx:
dst.fan_ctx = list(src.fan_ctx)
return dst
# appelé : advance()/child() (wot_record.py:46), .transform (core.py:2071),
# le voyageur _bag sur handoff, header x-fan-ctx (cast.py)
Un event appartient à tous les fan_id de sa pile ⇒ un gather
matche si son fan_id ∈ la pile (règle l'imbrication + le multi-collecteur). Le
fan_ctx étant un champ du StatefulRecord, il ride le snapshot/CRDT tout seul
(node_record.py:68) — rien à faire côté sérialisation Ray.
| wot_record.py | + FanFrame, + champ fan_ctx, préserver dans advance()/child() |
| deployer.py:428 | pousser une FanFrame en même temps que parent_id |
| loop_tools.py:621 | handoff_agent(list,…) mint le fan_id |
| core.py:2071 | .transform : carry_fan(event, result) |
| cast.py | header x-fan-ctx (events injectés par cast) |
Problème. Aucun gather clé par un fan_id runtime : Cast clé un
correlation_id 1:1 (cast.py:461), le join deployer clé par position de
niveau (deployer.py:441), MergeNode.wait_n clé par noms de sources statiques
(wot_merge.py:227). Solution : un GatherRecord par
(gather_node, fan_id), dans le bag sous gather:<fan_id>.
@sr_init(entity="gather", initial="open", states=["open","closed"])
class GatherRecord(StatefulRecord):
fan_id: str = ""; mode: str = "perfect"; expected: int = 0
seen: set = field(default_factory=set) # {seq} — dedup (Ray peut redélivrer)
kernel: str = "append"; kstate: dict = field(default_factory=dict)
result: Any = None; closed: bool = False
Le corps (framework-side). Latent & perfect partagent le chemin —
apply est toujours incrémental, finalize ne sert qu'à clore :
async def _node_gather(rec, ev, *, gather_node, fan_id, n, mode, kernel):
if fan_id not in [f.fan_id for f in ev.fan_ctx]:
return ev # pas notre fan-out → passe-through
g = ensure_gather(gather_node, fan_id, mode, n or fan_expected(ev, fan_id), kernel)
seq = fan_seq(ev, fan_id)
if seq in g.seen: return None # dedup / idempotence
g.seen.add(seq); apply_kernel(g, ev) # INCRÉMENTAL
if mode == "perfect" and len(g.seen) >= g.expected:
g.result = finalize_kernel(g); g.closed = True
return g.result # ferme + TRANSITE en aval
return None # latent, ou perfect pas plein → rien ne sort
Perfect : compte jusqu'à expected, finalize, transite
(généralise MergeNode.wait_n). Latent : ne ferme jamais ; le
GatherRecord agglomère, un nœud consume_at le tire
(channel_read(consume=True), loop_context.py:304 / bag_read).
Les kernels se branchent sur apply/finalize et
réutilisent l'existant :
KERNELS = { # apply = incrémental ; finalize = clore
"append": (push, lambda g: g.kstate["xs"]), # _node_join / _merge_dicts
"window": (window_fold, lambda g: g.kstate["digest"]), # étend dry_context
"transform": (user_fn, lambda g: user_fn(end=True)), # .transform
"rlm": (rlm.add_message, lambda g: rlm_record), # ContextRLMRecord
"okf": (okf_write, lambda g: g.kstate["ref"]), # la REF, pas la donnée
}
Placement. wired = arête directe (worker >> Gather(...),
comme un join). ambient = le collecteur est ailleurs et reçoit une copie de tout
event portant son fan_id (« des collecteurs à des endroits définis ») via un
fan_bus (un topic par fan_id, réutilise cast/channel). Timeout/dead-letter
(réutilise MergeNode) clôt un perfect si un enfant se perd.
| gather_record.py (neuf) | GatherRecord + KERNELS + apply/finalize_kernel |
| loop_wot_deploy.py | _node_gather, sélection de body, deploy ; fan_bus (ambient) |
| builder.py | .fan(...), Gather(...), Reduce(...) (fluent) |
| wot_compiler.py:374 | reconnaître un nœud gather (fan_id+mode+n) au lieu du MergeSpec statique |
| loop_tools.py | tools okf_read/okf_stats (kernel okf) |
Ordre : ① fan_ctx + stamping + carry_fan (test : fan_id intact à 3 hops) → ② gather perfect × append (ferme à N, sort la liste) → ③ brancher les kernels → ④ latent + consume_at → ⑤ ambient + multi-collecteurs → ⑥ timeout/dedup.