🪄 L'illusion de la fibre : ce que RxJS nous a appris sur les pipelines asynchrones réels en PHP
La plupart des développeurs découvrent RxJS grâce à sa syntaxe fluide :
source.pipe(map(...), mergeMap(...), retry(...));
Mémorable — et la partie la moins intéressante lorsqu'on conçoit un moteur de pipeline PHP.
Nous avons utilisé RxJS comme modèle architectural pour Darkwood Flow. Les équipes qui développent des scrapers, des importateurs et des pipelines multimédias se demandent sans cesse si elles ont besoin d'API « réactives ». Notre question était plus directe :
Que peut apprendre Flow de RxJS sans devenir RxJS ?
RxJS décrit les valeurs au fil du temps selon un cycle de vie explicite : composition avec pipe, démarrage avec subscribe, suivi des tâches en cours, attente de la fin, propagation des erreurs, et suppression lors de la désinscription. Les planificateurs sont importants pour les minuteurs ; ils ne sont pas ce qui permet le fonctionnement des requêtes d'URL simultanées. Le concept fondamental est le cycle de vie d'exécution.
Sous le capot, cette machine ressemble à ceci :
subscription starts execution → values propagate through a chain → asynchronous inner work is tracked → completion waits for active work → errors terminate or restart the run → unsubscription triggers teardown
Flow se heurte déjà aux mêmes questions à mesure qu'il dépasse le simple cadre de la composition de scènes :
-
Quand un pipeline démarre-t-il ?
-
Quand sera-t-il terminé ?
-
Qu'est-ce qui est encore actif ?
-
Que se passe-t-il en cas d'erreur ?
-
Qui ferme les ressources ?
-
Comment arrêter une course ?
Dans Flow, chaque valeur est transportée dans un paquet d'instructions nommé Ip. Les paquets sont admis par des stratégies, planifiés par un pilote et collectés après l'appel à await(). Nous prenions déjà en charge Fibers, Amp, React, Swoole et d'autres technologies similaires, et nous traitions déjà de nombreux paquets via le processus fetch → hash → save. Ce qui nous manquait, c'était une prise en charge transparente des E/S non bloquantes natives : la possibilité de gérer les attentes réelles des sockets sans imposer Amp ou React comme dépendance stricte, et sans supposer qu'une Fiber autour d'une lecture bloquante soit suffisante.
JavaScript fournit une boucle d'événements à RxJS. PHP CLI n'en fournit pas. Un Fiber peut suspendre l'exécution PHP coopérative ; il ne peut pas transformer à lui seul un appel socket bloquant en E/S chevauchantes. Flow s'adapte avec des primitives explicites :
PHP streams → non-blocking mode → stream_select() → Fibers → Driver → Flow
À retenir : Adoptez la rigueur du cycle de vie de RxJS. Conservez le vocabulaire de Flow (tâches, paquets, pilotes). Cette étude n’a jamais été une digression sur RxPHP ; ce qui est intéressant, ce sont les principes, et non les types.
Pour rester intègres, nous avons utilisé un pipeline en béton – et non pas seulement des minuteries artificielles :
URLs → fetch concurrently → checksum → save → per-URL result → summary
Les sockets réelles révèlent un comportement bloquant. Les limites de concurrence sont respectées ou non. Les délais d'attente et le nettoyage deviennent visibles. Une URL ayant échoué ne doit pas annuler cinq succès lorsque c'est la règle du produit. Les démonstrations artificielles de delay() prouvent la planification Fiber ; elles ne peuvent pas prouver que les sockets ouvertes partagent une attente stream_select.
Thèse : Flow devrait emprunter la discipline de RxJS en matière d’exécution, de concurrence, d’achèvement et de nettoyage, tout en restant un moteur PHP natif construit sur des tâches, des paquets, des stratégies et des pilotes.
Autrement dit : nous pensions qu’un pilote basé sur Fiber permettait déjà l’exécution asynchrone de Flow. L’étude de RxJS nous a contraints à définir ce que signifie « asynchrone ». La suite de cet article relate cette découverte.
Sous les opérateurs : comment RxJS s’exécute réellement
Nous avons examiné RxJS 8.0.0-alpha.14 (types principaux dans @rxjs/observable). La structure des packages est spécifique à cette version ; les principes du cycle de vie sont plus généraux que dans cette version alpha.
from(urls).pipe( mergeMap(fetchImage, 4), map(hashBinary), retry(2), finalize(cleanup), );
La composition n'est pas l'exécution
La méthode Observable.pipe transforme les opérateurs en nouveaux Observables. Aucune donnée n'est encore récupérée. Les opérateurs de cette arborescence renvoient généralement new Observable(subscriber => { ... }) ; l'ancienne méthode lift est obsolète à partir de la version 8. Flow sépare déjà la composition (FlowFactory / fn) de l'exécution (await()) : même séparation, noms différents.
L'abonnement lance l'exécution
La fonction subscribe encapsule le consommateur dans un Subscriber, connecte le producteur et enregistre la destruction. Les opérateurs s'abonnent en amont ; next, error et complete envoient les données vers le récepteur. Des fonctions auxiliaires telles que operate({ destination, ... }) lient les processus enfants à l'arbre de destruction de la destination, ce qui permet aux désabonnements de se propager.
Un Subscription stocke les finaliseurs et les exécute de manière idempotente lors de la désinscription. Les signaux terminaux arrêtent les notifications et désinscrivent. Le nettoyage fait partie de la fin d'une exécution. Les sources froides (y compris from(array)) isolent les exécutions : chaque abonnement possède son propre producteur.
mergeMap suit les éléments internes actifs
mergeMap(project, concurrent) délègue la logique de fusion partagée : un compteur d’entrées internes actives, une mémoire tampon pour les débordements et la fin de la fusion uniquement lorsque la source externe est terminée, que la mémoire tampon est vide et qu’il ne reste plus d’entrées internes. Les entrées internes passent par from(project(...)).
Nous n'avions pas besoin de mergeMap en PHP. Nous avions besoin de sa garantie : un travail en cours limité et une achèvement qui attend la fin de chaque opération active.
réessayer et finaliser
En cas d'erreur, retry désinscrit l'abonnement interne courant et se réinscrit à la même source. Pour les requêtes from(urls) effectuées à froid, cela redémarre la liste — une méthode puissante, mais souvent inadaptée aux E/S par lots. finalize(callback) enregistre le nettoyage de l'abonnement externe (terminé, erreur ou désinscription), et non entre les tentatives de nouvelle connexion.
Ce qu'il faut voler, et ce qu'il ne faut pas voler
C'est le cœur de l'étude.
VolerLaisser derrière soiLazy compose → démarrage expliciteObservable / Observer / Subscriber comme types d'applicationTravail en vol délimitéCatalogues d'opérateurs et noms de méthodes JSAchèvement en attente de travail actifZoos du planificateur publicNettoyage sur chaque chemin terminalSujets / multidiffusion « car Rx les prend en charge »
Les analogies sont utiles, mais il ne faut pas s'arrêter là :
-
Subscription ≈ un gestionnaire d'exécution de flux futur
-
Planificateur ≈ Pilote
-
mergeMap(n) ≈ étape concurrente bornée (MaxIpStrategy + tâches générantes)
À retenir : Copiez le contrat de cycle de vie. Ne copiez pas le système de types.
Ce que Flow possédait déjà
Avant l'étude, Flow n'était pas une page blanche. C'est ce point de départ qui a rendu l'illusion de la fibre convaincante.
Nous pourrions déjà :
-
Composez des étapes avec FlowFactory (notation do du générateur) ou fn()
-
Déplacer les valeurs dans les paquets Ip
-
Exécuter des tâches (JobInterface ou fermetures)
-
Limiter l'admission avec MaxIpStrategy($n)
-
Planification via les pilotes (FiberDriver par défaut, ainsi que Amp, React, Swoole, …)
-
Barrière sur await() (void—les résultats nécessitent une tâche de collecte)
-
Gérer les erreurs de routage via errorJob (optionnel)
build pipeline → push Ip(s) → await() → Driver admits packets, runs jobs, forwards results → when nothing remains queued or active, await returns
L'envoi d'un paquet le place dans la file d'attente de la première étape. Dans la fonction await(), le pilote interroge la stratégie de chaque étape pour déterminer quels paquets peuvent être exécutés, lance les fibres de tâches (ou coroutines), puis transmet le succès en aval ou appelle errorJob. Les paquets terminés libèrent des emplacements de concurrence.
Modèle mental : les paquets entrent en étapes, les stratégies limitent leur nombre d'exécutions, le pilote gère l'attente.
Avec FiberDriver, la gestion de plusieurs paquets en transit, delay() et MaxIpStrategy(4), l'exécution ressemblait à de l'asynchronisme. Les démonstrations de minuterie d'amplification se chevauchaient parfaitement. Le problème venait des E/S bloquantes natives : nous avions une interprétation erronée du terme « asynchrone » dans ce contexte.
L'illusion de la fibre
Ce que nous croyions
Envoi de quatre URL → admission de quatre processus → démarrage de quatre fibres → await(). Si les tâches ont effectué des E/S, les attentes doivent se chevaucher.
RxJS a rendu cette hypothèse gênante : si mergeMap suit les internes actifs, que suivait Flow — des paquets, des fibres ou des sockets ?
Ce qui a révélé l'affaire
file_get_contents($url); // still blocks the process
Dans une tâche FiberDriver, cet appel n'est jamais suspendu en cas de disponibilité. Le pilote attend sur le bloc. MaxIpStrategy(4) peut accepter quatre paquets pendant que le processus effectue encore la sérialisation au niveau du socket, ou bloquer le planificateur pendant la recherche de concurrence.
Même piège : blocage de fread. Les fibres ne sont utiles que lorsque la tâche cède la main. Sans API de disponibilité, il n'y a rien sur quoi céder la main.
Une fibre optique peut suspendre l'exécution de PHP. Elle ne peut pas rendre non bloquant un appel socket bloquant.
ÉtiquetteSignificationFibres seulement ?PHP simultanéPlusieurs fibres planifiéesOuiChevauchement des E/SPlusieurs sockets dans un seul stream_selectNon — nécessite une sélection non bloquanteFaux asynchroneBlocage des appels dans FibersDangereusement facile
RxJS a imposé ce vocabulaire : les limites d’admission ne correspondent pas à un chevauchement d’E/S. Flow a dû adopter cette transparence.
La pile dont nous avions besoin
non-blocking streams → stream_select() → Fiber suspension on wait → readiness-based resume → Driver await loop → Flow stages
Cette pile est devenue le centre de l'étude, et c'est pourquoi StreamSelectDriver a été intégré à darkwood/flow.
À retenir : « Pilote asynchrone » sans mécanisme de gestion de la disponibilité reste un slogan. Avec stream_select, cela devient un engagement.
Résultats concrets : le pipeline d’images
La démo se trouve dans content/rxjs-flow et utilise Flow local via un lien symbolique Composer. Un moteur local éphémère a permis de tester stream_select, puis a été volontairement supprimé. Un seul environnement d'exécution réutilisable.
php bin/console app:fetch-images --concurrency=4 --timeout=5
ScèneRôleFetchImageJobstream_socket_client, E/S non bloquantes, waitWritable / waitReadable, analyse HTTP minimale, fermeture finallyHashBinaryJobSHA-256EnregistrerFichierJobPersister le corpsCollecteurAgrégat ImageResult (await() renvoie void)
La gestion de la concurrence s'effectue lors de la récupération via MaxIpStrategy($concurrency). L'erreur HTTP 404 est une erreur limitée à chaque élément sur les DTO (et non sur errorJob), ce qui signifie que cinq succès peuvent coexister avec un échec. Les traces génèrent flow.started, item.started / completed / failed, flow.completed.
Exécution représentative (concurrence=2) :
[flow.started] items=6 concurrency=2 [item.started] url=http://127.0.0.1:8765/image-a.png [item.started] url=http://127.0.0.1:8765/image-b.png [item.started] url=http://127.0.0.1:8765/image-c.png [item.completed] url=…/image-a.png checksum=… path=… … [item.failed] url=…/fail.png error=HTTP 404 [flow.completed] ok=5 failed=1 elapsed_ms=711.23
Les fixtures utilisent php -S (monothread) avec de courts délais de routage pour que le transfert de pool soit visible. Les temps d'attente limités côté client sont bien réels ; l'accélération réelle nécessite un serveur multi-connexions. La tâche de récupération (fetch) illustre le fonctionnement du protocole HTTP : les redirections, les cas particuliers de TLS, les spécificités du traitement par blocs, les proxys, la compression et les E/S partielles ne sont pas abordés. Les clients de production restent gérés par des tâches ; les pilotes restent génériques.
Prochain test d'acceptation (prévu, non revendiqué tel que mesuré ici) :
5 × ~1s delayed endpoints on a multi-connection backend concurrency=1 → ~5s concurrency=4 → ~2s
StreamSelectDriver : ce que nous avons livré
StreamSelectDriver multiplexe la disponibilité des flux natifs sur les fibres :
$ready = $driver->waitReadable($stream, $timeoutSeconds); $ready = $driver->waitWritable($stream, $timeoutSeconds);
Les deux renvoient un booléen et doivent être exécutés dans un job Fiber. En interne, le pilote :
-
Enregistrements Fibre, flux, mode, date limite facultative
-
Suspend
-
Construit des ensembles stream_select à partir des attentes en attente
-
Reprend les fibres prêtes avec true, les attentes expirées avec false
-
Reprend toujours les suspensions delay() simples à l'itération normale de la boucle.
Même DriverInterface que les autres Drivers (async, defer, await, delay, tick), plus ces API d'attente auxquelles les Jobs s'inscrivent par type.
Correction FIFO pour MaxIpStrategy : à pleine capacité, ne pas extraire et remettre en file d'attente (dans cet ordre de rotation). Les traces simultanées ont permis de mettre en évidence le bogue.
Pourquoi ne pas améliorer FiberDriver ? La reprise par défaut de Fiber diffère de la reprise conditionnée par l'état de disponibilité. Un pilote dédié permet de définir explicitement la règle de planification.
Driver → schedule, wait, resume, deadlines Job → protocol + application work IpStrategy → packet admission / concurrency Collector → terminal aggregate
$driver = new StreamSelectDriver(); $flow = (new FlowFactory($driver))->create(static function () use (...) { yield [$fetchJob, null, new MaxIpStrategy($concurrency)]; yield $hashJob; yield $saveJob; yield $collectorFn; });
À retenir : Les tâches sont exécutées. Les pilotes attendent. Le protocole HTTP ne réside jamais dans le pilote — à moins que chaque protocole n’invente un framework.
Ce que « concurrence=4 » limite réellement
Une limite de concurrence n'a de sens que si l'environnement d'exécution peut indiquer précisément ce qu'elle limite.
MaxIpStrategy(4) limite le traitement des paquets à une étape donnée. À elle seule, elle ne crée pas d'attentes réseau qui se chevauchent.
Couche« 4 » pourrait signifierPaquets en attenteNombreux sont ceux qui attendentOffres d'emploi acceptées≤ 4 (traitement < max)Fibres actives≤ 4 une fois commencéesSockets ouverts≤ 4 si chaque récupération donneRequêtes en coursIdem, mais avec des temps d'attente non bloquants
Avec les E/S bloquantes, la limitation est purement cosmétique. Avec StreamSelectDriver et les tâches de récupération avec catch :
MaxIpStrategy(4) + non-blocking fetch + StreamSelectDriver → ≤ 4 fetch packets admitted → ≤ 4 fetch Fibers active → ≤ 4 sockets waiting on select
Il s'agit de la traduction de Flow de mergeMap(..., 4). La portée est d'une seule étape : le hachage/l'enregistrement peuvent encore chevaucher la récupération suivante — il s'agit d'un chevauchement de pipeline, et non d'une rupture de la limite.
Erreurs, achèvement, nettoyage
Collecte vs gestion des échecs rapides. Les importations par lots nécessitent une gestion des erreurs par élément ; les publications critiques peuvent nécessiter errorJob ou le traitement des échecs. Flow prend en charge les deux. Les délais d'attente existent déjà ; les opérateurs de nouvelle tentative/délai d'attente complets peuvent attendre que les schémas se répètent entre les applications.
Terminé signifie : source épuisée, files d'attente vides, aucune tâche active, aucune attente de flux en cours, ressources fermées. Actuellement, cette étape est franchie par await() et un collecteur. L'ajout futur de $flow->start() / Execution (annulation, résultat, erreurs) est une amélioration ergonomique optionnelle, non requise pour cette section.
Le nettoyage, c'est enfin respecter les délais – le nettoyage RxJS implémenté en PHP, et non l'espoir d'un destructeur. Les traces structurées sont plus efficaces qu'un bus d'événements interne pour une observabilité dès le premier jour.
Correspondance entre RxJS et Flow
RxJSFluxAppelObservableSource des adresses IPAdaptProtocole observateurCollecteur + errorJobSimplifierAbonnementGestionnaire d'exécution futureReporterOpérateurPoste / étapeAdopterPlanificateurPiloteConserver le nom du fluxmergeMap(n)MaxIpStrategy(n) + tâches générantesCapacité d'adoptionfinaliserNettoyage expliciteAdopterSujet / partage—Refuser pour le moment
À conserver : composition paresseuse / exécution explicite, pipelines multi-éléments, concurrence limitée, achèvement déterministe, nettoyage, isolation par exécution, observabilité.
Simplification : tâches + paquets ; pilotes ; collecteurs — pas un protocole public next/error/complete ; limites — pas des noms d'opérateurs.
Postpone: public Execution, cancel tokens, first-class retry/timeout ops, switchMap / exhaustMap, hot streams.
Rejet : API publique de type RxJS, Subjects, zoos de promesses/opérateurs, HTTP dans le pilote, moteurs doubles.
Flow reste un moteur de pipeline mature et pragmatique avec un véritable exécuteur asynchrone, et non un framework réactif utilisant la syntaxe PHP.
Où Flow va-t-il ensuite ?
Keep: lifecycle discipline (start → bound concurrency → wait → tear down) Shape: Job + Ip + MaxIpStrategy + Driver + await + collector Next: multi-connection overlap benchmarks cancel tokens on StreamSelectDriver waits richer await ergonomics if collectors become painful
Darkwood Flow ne devient pas RxJS pour PHP. Il mise tout sur un moteur de pipeline natif PHP : concurrence limitée, E/S basées sur la disponibilité, achèvement déterministe et nettoyage explicite, avec des pilotes capables de justifier leurs temps d’attente.
Lorsqu'une tâche est en attente d'E/S, l'environnement d'exécution attend-il réellement avec elle, ou fait-il seulement semblant ?
💫 Bonzai Creator - 🌿🩸VITALITÉ<>OREXIS 🌿🩸
🤖 Veille Darkwood - 2026-07-25
🤖 Veille Darkwood - 2026-07-24
💫 Créateur de Hacker News - Les laboratoires d'IA abusent-ils du système PelicanMaxx ?
💫 Créateur Reddit - r/opensource : Twigg : Gestion de versions et forge logicielle sous licence AGPL
🤖 Veille Darkwood - 2026-07-23
💫 Créateur de Hacker News - OverpAId – Virez votre PDG. Embauchez l'avenir.
💫 GitHub Creator - Laravel: laravel/framework: v13.21.1
🤖 Veille Darkwood - 2026-07-22
💫 Hacker News Creator - La stratégie chinoise d'IA à pondération ouverte est gagnante
💫 GitHub Creator - php: frankenphp v1.12.5
💫 Bonzai Créateur - Arnaud "Tugan" Labossière
🤖 Veille Darkwood - 2026-07-21
💫 Bonzaï Créateur - Mini - Jules
🤖 Veille Darkwood - 2026-07-20
💫 Créateur de Bluesky - @alexdaubois : Mise à jour de sécurité Vulcain
💫 Créateur Reddit - r/opensource : canid : carte de contact/lien dans la bio accessible par URL
💫 Créateur arXiv - cs.AI : RoboTTT : Mise à l'échelle du contexte pour les politiques robotiques
💫 Bonzaï Créateur - Intiméa Studio
🤖 Veille Darkwood - 2026-07-19