Tachikoma · Docs · Web-of-Trust

Fan-out — map & reduce dans un Web-of-Trust

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.

Pattern map · reduce gouverné Collecte adjacent · perfect · latent · ambient Kernels append · window · transform · rlm · okf

Le fan-out — map & reduce dans un WoT

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.

1 · Le dispatch — les disciplines de collecte

Comment on ramène les branches : adjacent (câblé après le fan), ou des collecteurs corrélés par fan_idperfect (ferme à N) vs latent (agglomère, tiré à l'étape n), et ambient (plusieurs collecteurs, sans câblage).

glisser pour tourner · molette pour zoomer

2 · Les kernels — ce que le reduce CALCULE

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.

payloads → kernel → résultat

3 · OKF — le reduce par référence (le gros dataset)

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.

objet réduit 1 240 000 lignes table ClickHouse reduce okf handle — dans le bag { ref, schema } rows: 1 240 000 branch(agence), amount(encours €) compact — jamais la donnée agent contexte compact okf_read(rows="0:50") · okf_stats(field) — tranche / agrégat, jamais tout
# 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.

4 · Use cases — quel pattern quand

DisciplineExemple concret Le plus adapté quand
Adjacent joinun 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 N8 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 · pullun 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 · multiun 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 · refla 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
KernelExemple concret Le plus adapté quand
appendles 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
transformsomme 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
rlmles 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
okfun 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)
08

Les stratégies de reduce — comment ramener un fan-out

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.

Axe A — collecte & fermeture

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

Axe B — le kernel (fonction de fusion)

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.

Latent vs Perfect — la discipline de fermeture

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")

Les cinq kernels — comment on les écrit

KernelCe qu'il fait Syntaxe
appendconcat / dict-merge (conflit détecté) kernel="append"
fenêtre glissanterésumé progressif : new = f(résumé, N derniers) kernel=Window(size=3, summarize=…)
transformhandler : latent (incrémental) | end (tout arrivé) | append kernel=Transform(fn, when="latent")
rlmingéré en latent, géré par la logique du RLM context kernel=Rlm(consume_at="report")
okfref vers un objet sérialisé + schéma + accès partiel kernel=Okf(store=…, partial=True)

Un WoT à la main — trois collecteurs sur un même fan-out

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

OKF — la ref est dans le bag, avec son schéma et son interface d'accès partiel

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.

09

Le seam fan_id + gather — spec d'implémentation

Le 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.

Diff des fichiers

wot_record.py+ FanFrame, + champ fan_ctx, préserver dans advance()/child()
deployer.py:428pousser une FanFrame en même temps que parent_id
loop_tools.py:621handoff_agent(list,…) mint le fan_id
core.py:2071.transform : carry_fan(event, result)
cast.pyheader 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 cheminapply 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.

Diff des fichiers

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:374reconnaître un nœud gather (fan_id+mode+n) au lieu du MergeSpec statique
loop_tools.pytools 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.