From 823d40eada30f71336cf9fdaf15499c9c58b31db Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gr=C3=A9gory=20Planchat?= Date: Fri, 4 Sep 2026 20:15:22 +0200 Subject: [PATCH 1/2] =?UTF-8?q?fix(coeur):=20la=20garde=20de=20divergence?= =?UTF-8?q?=20compare=20aussi=20la=20charge=20de=20l'activit=C3=A9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit DUR042 ne comparait que le nom de l'activité au slot. La charge ne traversait jamais la comparaison : un replay qui recalculait un payload différent voyait le journal servir l'ancien résultat, la charge fraîche partir à la poubelle, et l'exécution se terminer en succès en ayant menti sur ce qu'elle avait demandé. Mesuré par mutation sur une maquette d'agent, avant correctif : neuf charges calculées, trois journalisées, six divergences avalées sans un mot, suite verte. `activityPayloadForSlot(): ?array` rejoint le port, et la comparaison passe par une empreinte canonique calculée **des deux côtés au même endroit** — c'est ce qui la rend symétrique entre un côté enregistré qui a fait l'aller-retour JSON du magasin et un côté frais qui ne l'a pas fait. Ce que la garde ne voit volontairement pas, parce qu'un faux positif arrête une exécution saine et coûte plus cher que le trou qu'il bouche : - un objet dont le journal ne retient rien — DTO à propriétés privées, le style de la maison — rend la même empreinte des deux côtés ; - l'ordre des clés d'un objet JSON, qui n'est pas de l'information ; l'ordre d'une liste, lui, en est ; - une charge inencodable (ressource, NAN) désarme la garde ; - une histoire écrite avant ce changement n'a rien à comparer. Le message nomme l'octet où les deux empreintes divergent et montre une fenêtre autour : sur une charge d'agent de 4 ko, deux préfixes tronqués se ressemblaient trait pour trait et n'apprenaient rien. Hors périmètre : les slots Nexus et workflow enfant comparent toujours leur seule identité. Le trou y est le même, et l'enjeu plus grand côté Nexus, où le doublon part chez un tiers. Co-Authored-By: Claude Opus 5 (1M context) --- UPGRADE.md | 31 +++ .../Worker/TemporalExecutionHistory.php | 27 +++ src/Durable/ExecutionContext.php | 139 +++++++++++++ .../Port/WorkflowHistorySourceInterface.php | 16 ++ src/Durable/Store/EventStoreHistorySource.php | 19 ++ .../Replay/ActivityPayloadDivergenceTest.php | 189 ++++++++++++++++++ 6 files changed, 421 insertions(+) create mode 100644 tests/unit/Durable/Replay/ActivityPayloadDivergenceTest.php diff --git a/UPGRADE.md b/UPGRADE.md index 074f7922..ee8d5416 100644 --- a/UPGRADE.md +++ b/UPGRADE.md @@ -24,6 +24,37 @@ ne contient que ce que Rector sait faire sans deviner ; tout le reste est écrit ## 0.1.0-alpha8 +### La garde de divergence compare aussi la charge de l'activité + +`WorkflowHistorySourceInterface` gagne `activityPayloadForSlot(int $slot): ?array`. La garde de +divergence (DUR042) ne comparait que le **nom** de l'activité au slot ; elle compare désormais aussi +ses arguments. Un replay qui redemande la même activité avec une autre charge lève un +`WorkflowTaskFailure` au lieu de continuer en silence. + +**Pourquoi** — le nom seul laissait passer la moitié du problème. Le journal servait l'ancien +résultat, la charge fraîchement calculée partait à la poubelle, et l'exécution se terminait **en +succès** en ayant menti sur ce qu'elle avait demandé. Mesuré sur une maquette d'agent : neuf charges +calculées, trois journalisées, six divergences avalées sans un mot, suite de tests verte. + +**Ce que ça change pour du code existant** — un workflow déjà déterministe ne voit rien. Un workflow +qui construisait sa charge avec une horloge, un aléa ou une lecture hors journal échoue désormais sa +tâche de replay, en nommant l'octet où les deux empreintes divergent. C'est le défaut qu'il fallait +voir : ces exécutions-là rendaient déjà un résultat faux. + +**Ce qui reste hors de portée de la garde, volontairement** — la comparaison passe par l'empreinte +que le journal sait tenir (aller-retour JSON, clés triées). Un objet dont le journal ne retient rien +— un DTO à propriétés privées, le style de la maison — ne fait donc pas diverger un replay fidèle. +Une charge inencodable (ressource, `NAN`) désarme la garde plutôt que d'accuser ce qu'elle ne sait +pas lire. Les histoires écrites avant ce changement n'ont rien à comparer et passent inchangées. + +**Ce que Rector ne peut pas faire** — rien à réécrire dans le code appelant. Seules les +implémentations tierces de `WorkflowHistorySourceInterface` doivent ajouter la méthode ; rendre +`null` reproduit exactement le comportement d'avant, sans garde sur la charge. + +**Hors périmètre** — les slots Nexus et workflow enfant comparent toujours leur seule identité. Le +trou y est le même, et l'enjeu plus grand côté Nexus, où un doublon part chez un tiers. + + ### Laravel refuse au démarrage un workflow dont les noms de paramètres divergent du contrat `gplanchat/durable-laravel` enregistrait sans vérifier. Un workflow portant diff --git a/src/Bridge/Temporal/Worker/TemporalExecutionHistory.php b/src/Bridge/Temporal/Worker/TemporalExecutionHistory.php index 44fc1390..e5fd0c82 100644 --- a/src/Bridge/Temporal/Worker/TemporalExecutionHistory.php +++ b/src/Bridge/Temporal/Worker/TemporalExecutionHistory.php @@ -53,6 +53,9 @@ final class TemporalExecutionHistory implements WorkflowHistorySourceInterface /** @var array activityId → nom d'activité (pour typer les échecs) */ private array $activityNames = []; + /** @var array> activityId → charge planifiée (garde DUR042) */ + private array $activityPayloads = []; + /** @var array scheduled event ID → activity ID */ private array $scheduledEventIdToActivityId = []; @@ -241,6 +244,20 @@ private function consumeEvent(HistoryEvent $event): void $this->activityIdToScheduledEventId[$activityId] = $eventId; $this->activityNames[$activityId] = (string) ($attr->getActivityType()?->getName() ?? ''); $this->scheduledEventIdToActivityId[$eventId] = $activityId; + + // L'entrée porte l'enveloppe écrite par TemporalActivityScheduleInput ; sa case + // `payload` tient les arguments. Absente ou illisible, on n'enregistre rien : + // la garde n'a alors rien à comparer, ce qui est son cas de repos. + $input = $attr->getInput(); + if (null !== $input) { + $payloads = $input->getPayloads(); + if ($payloads->count() > 0) { + $envelope = JsonPlainPayload::decode($payloads[0]); + if (\is_array($envelope) && \is_array($envelope['payload'] ?? null)) { + $this->activityPayloads[$activityId] = $envelope['payload']; + } + } + } } break; @@ -525,6 +542,16 @@ public function activityNameForSlot(int $slot): ?string return '' === $name ? null : $name; } + public function activityPayloadForSlot(int $slot): ?array + { + $activityId = $this->scheduledActivityIds[$slot] ?? null; + if (null === $activityId) { + return null; + } + + return $this->activityPayloads[$activityId] ?? null; + } + public function findTimerSlotResult(int $slot): ?array { $timerId = $this->scheduledTimerIds[$slot] ?? null; diff --git a/src/Durable/ExecutionContext.php b/src/Durable/ExecutionContext.php index 25cfadd0..42f12909 100644 --- a/src/Durable/ExecutionContext.php +++ b/src/Durable/ExecutionContext.php @@ -31,6 +31,9 @@ final class ExecutionContext { + /** Ce qu'un message de divergence montre d'une empreinte de charge avant de la couper. */ + private const DIVERGENCE_PRINT_LIMIT = 256; + private ?QueryHandlerRegistry $queryHandlers = null; /** @var array */ @@ -100,6 +103,7 @@ public function activity(string $name, array $payload = [], ?ActivityOptions $op { $slotIndex = $this->activitySlotIndex++; $this->refuseActivityDivergence($slotIndex, $name); + $this->refuseActivityPayloadDivergence($slotIndex, $name, $payload); $replay = $this->historySource->findActivitySlotResult($slotIndex); if (null !== $replay) { $deferred = new \Gplanchat\Durable\Awaitable\Deferred(); @@ -274,6 +278,141 @@ private function refuseActivityDivergence(int $slotIndex, string $requested): vo $this->refuseDivergence('activity', $slotIndex, $this->historySource->activityNameForSlot($slotIndex), $requested); } + /** + * Refuse un slot dont le **nom** concorde mais dont les **arguments** ont changé au replay. + * + * Le nom seul laissait passer la moitié du problème : le journal servait l'ancien résultat, la + * charge fraîchement calculée partait à la poubelle, et l'exécution se terminait en succès en + * ayant menti sur ce qu'elle avait demandé. Mesuré : neuf charges calculées, trois + * journalisées, six divergences avalées sans un mot. + * + * La comparaison passe par {@see canonicalPayload()}, qui fait voir aux deux côtés **ce que le + * journal sait tenir** et rien de plus. C'est la règle qui évite les faux positifs : un objet + * dont le journal ne retient rien ne doit pas faire diverger un replay fidèle, sous peine + * d'arrêter en production des exécutions parfaitement saines. + * + * @param array $requested + * + * @throws WorkflowTaskFailure si le code redemande la même activité avec une autre charge + */ + private function refuseActivityPayloadDivergence(int $slotIndex, string $name, array $requested): void + { + $recorded = $this->historySource->activityPayloadForSlot($slotIndex); + if (null === $recorded) { + // Rien d'enregistré ici : slot neuf, ou histoire écrite avant que la charge soit + // lisible. Refuser laisserait sans garde exactement les exécutions qu'elle protège. + return; + } + + $recordedPrint = $this->canonicalPayload($recorded); + $requestedPrint = $this->canonicalPayload($requested); + if (null === $recordedPrint || null === $requestedPrint || $recordedPrint === $requestedPrint) { + // Une charge que le journal ne sait pas rendre comparable ne prouve rien : on se tait. + return; + } + + $at = self::firstDifference($recordedPrint, $requestedPrint); + + throw new WorkflowTaskFailure(\sprintf( + 'Replay divergence at activity slot %d of execution "%s": the activity name "%s" matches, ' + . 'but its payload changed at byte %d. History recorded %s, code scheduled %s ' + . '(%d and %d bytes). ' + . 'This is non-deterministic workflow code — the payload is rebuilt on every replay pass, ' + . 'so something in it reads the clock, draws a random value, or is resolved from outside ' + . 'the journal. This is not a version skew: a declared change point would have changed ' + . 'the name or the slot, not the arguments alone.', + $slotIndex, + $this->executionId, + $name, + $at, + self::windowAround($recordedPrint, $at), + self::windowAround($requestedPrint, $at), + \strlen($recordedPrint), + \strlen($requestedPrint), + )); + } + + /** + * Rend d'une charge l'empreinte que le journal en retiendrait, ou null si elle n'en a pas. + * + * Les deux côtés passent par ici, et c'est tout l'intérêt : le côté enregistré a fait + * l'aller-retour JSON du magasin, le côté frais non. Sans cette normalisation, un objet à + * propriétés privées — le style de la maison — rend `{}` à gauche et `[]` à droite, et toute + * exécution qui en porte un divergerait à chaque reprise. Mesuré avant d'être écrit. + * + * Les clés sont triées parce qu'un objet JSON n'a pas d'ordre : le voir changer ne prouve rien. + * Les listes gardent le leur, où l'ordre est l'information. + * + * Null veut dire « incomparable » — une ressource, un NAN, une récursion. La garde se tait + * alors, plutôt que d'accuser une charge qu'elle ne sait pas lire. + * + * @param array $payload + */ + private function canonicalPayload(array $payload): ?string + { + try { + $throughTheJournal = json_decode( + json_encode($payload, \JSON_THROW_ON_ERROR), + true, + 512, + \JSON_THROW_ON_ERROR, + ); + + self::sortKeysDeeply($throughTheJournal); + + return json_encode($throughTheJournal, \JSON_THROW_ON_ERROR); + } catch (\JsonException) { + return null; + } + } + + private static function sortKeysDeeply(mixed &$value): void + { + if (!\is_array($value)) { + return; + } + + ksort($value); + foreach ($value as &$nested) { + self::sortKeysDeeply($nested); + } + } + + /** + * Le premier octet où les deux empreintes cessent de coïncider. + * + * C'est ce qui rend le message utilisable : une charge d'agent pèse des kilo-octets, et deux + * préfixes identiques n'apprennent rien à qui lit. Mesuré avant d'être écrit — le premier + * message montrait 256 octets de tête et les deux côtés se ressemblaient trait pour trait. + */ + private static function firstDifference(string $recorded, string $requested): int + { + $shortest = min(\strlen($recorded), \strlen($requested)); + $at = 0; + while ($at < $shortest && $recorded[$at] === $requested[$at]) { + ++$at; + } + + return $at; + } + + /** + * Rend la fenêtre d'empreinte autour de la divergence, avec de quoi la situer de chaque côté. + */ + private static function windowAround(string $print, int $at): string + { + $margin = intdiv(self::DIVERGENCE_PRINT_LIMIT, 4); + $from = max(0, $at - $margin); + $window = substr($print, $from, self::DIVERGENCE_PRINT_LIMIT); + + return \sprintf( + '%s%s%s', + $from > 0 ? '…' : '', + $window, + $from + \strlen($window) < \strlen($print) ? '…' : '', + ); + } + /** * La règle, une fois, pour les trois types de slot qui portent une identité. * diff --git a/src/Durable/Port/WorkflowHistorySourceInterface.php b/src/Durable/Port/WorkflowHistorySourceInterface.php index 5f1d292d..ef36fb30 100644 --- a/src/Durable/Port/WorkflowHistorySourceInterface.php +++ b/src/Durable/Port/WorkflowHistorySourceInterface.php @@ -42,6 +42,22 @@ public function findScheduledActivityId(int $slot): ?string; */ public function activityNameForSlot(int $slot): ?string; + /** + * Returns the payload recorded for activity slot N, or null if there is nothing to compare. + * + * Null and `[]` are different answers, and callers rely on it: `[]` is an activity that was + * genuinely scheduled with no arguments, null is "this journal says nothing here" — an empty + * slot, or a history written before the payload was readable. Only null waives the guard. + * Conflating the two is the {@see findSideEffectForSlot()} trap one file over. + * + * The name answers "what was this"; this one answers "with what". A replay that keeps the + * name and changes the arguments is non-deterministic workflow code, and without this the + * journal serves the old result while the freshly computed payload is dropped in silence. + * + * @return array|null + */ + public function activityPayloadForSlot(int $slot): ?array; + /** * Returns the version this execution recorded for a declared change point, or null if it has * not reached that point yet. diff --git a/src/Durable/Store/EventStoreHistorySource.php b/src/Durable/Store/EventStoreHistorySource.php index b13c566c..3f44190e 100644 --- a/src/Durable/Store/EventStoreHistorySource.php +++ b/src/Durable/Store/EventStoreHistorySource.php @@ -127,6 +127,25 @@ public function activityNameForSlot(int $slot): ?string return null; } + public function activityPayloadForSlot(int $slot): ?array + { + $index = 0; + foreach ($this->eventStore->readStream($this->executionId) as $event) { + if ($event instanceof ActivityScheduled) { + if ($index === $slot) { + // `payload()` rend l'enveloppe de l'événement ; les arguments de l'activité en + // sont une case. Un non-tableau vaut « rien à comparer », pas « tableau vide ». + $arguments = $event->payload()['payload'] ?? null; + + return \is_array($arguments) ? $arguments : null; + } + ++$index; + } + } + + return null; + } + public function childWorkflowTypeForSlot(int $slot): ?string { $index = 0; diff --git a/tests/unit/Durable/Replay/ActivityPayloadDivergenceTest.php b/tests/unit/Durable/Replay/ActivityPayloadDivergenceTest.php new file mode 100644 index 00000000..7d277659 --- /dev/null +++ b/tests/unit/Durable/Replay/ActivityPayloadDivergenceTest.php @@ -0,0 +1,189 @@ +journalWith(['city' => 'Paris']); + $context = $this->context($store); + + $this->expectException(WorkflowTaskFailure::class); + $this->expectExceptionMessageMatches('/payload changed/'); + + $context->activity('weather', ['city' => 'Lyon']); + } + + public function testTheMessageNamesTheSlotTheActivityAndTheCause(): void + { + $store = $this->journalWith(['nonce' => 1]); + $context = $this->context($store); + + try { + $context->activity('weather', ['nonce' => 2]); + self::fail('La garde aurait dû refuser cette charge.'); + } catch (WorkflowTaskFailure $refusal) { + $message = $refusal->getMessage(); + self::assertStringContainsString('activity slot 0', $message); + self::assertStringContainsString('"weather"', $message, 'le nom concorde : le dire évite de chercher un slot décalé'); + self::assertStringContainsString('non-deterministic workflow code', $message); + self::assertStringNotContainsString('different version of the workflow', $message, 'ce message-là enverrait vers ChangePoint, qui n\'y peut rien'); + } + } + + public function testAFaithfulReplayPasses(): void + { + $store = $this->journalWith(['city' => 'Paris']); + $context = $this->context($store); + + $awaitable = $context->activity('weather', ['city' => 'Paris']); + + self::assertTrue($awaitable->isSettled(), 'le slot se résout : la charge est celle du journal'); + } + + public function testKeyOrderIsNotADivergence(): void + { + // Un objet JSON n'a pas d'ordre. Le voir changer ne prouve rien, et l'ériger en divergence + // arrêterait des exécutions dont la charge est identique. + $store = $this->journalWith(['b' => 2, 'a' => 1]); + $context = $this->context($store); + + $awaitable = $context->activity('weather', ['a' => 1, 'b' => 2]); + + self::assertTrue($awaitable->isSettled()); + } + + public function testAListThatChangedOrderIsADivergence(): void + { + // Dans une liste, l'ordre *est* l'information. + $store = $this->journalWith(['cities' => ['Paris', 'Lyon']]); + $context = $this->context($store); + + $this->expectException(WorkflowTaskFailure::class); + + $context->activity('weather', ['cities' => ['Lyon', 'Paris']]); + } + + public function testAnEmptyRecordedPayloadIsStillCompared(): void + { + // `[]` est une activité planifiée sans argument, pas « rien d'enregistré ». La confondre + // avec null est le piège de findSideEffectForSlot(), et il ne se reproduit pas ici. + $store = $this->journalWith([]); + $context = $this->context($store); + + $this->expectException(WorkflowTaskFailure::class); + + $context->activity('weather', ['city' => 'Paris']); + } + + public function testAnObjectTheJournalCannotSeeIntoDoesNotDiverge(): void + { + // Le style de la maison : des DTO en lecture seule à propriétés privées. Le journal n'en + // retient rien — `{}` à l'aller, `[]` au retour. Sans normalisation des deux côtés, toute + // exécution qui en porte un divergerait à chaque reprise. C'est le faux positif qui aurait + // arrêté des workflows sains, et il est mesuré ici pour qu'il ne revienne pas. + $freshPayload = [ + 'amount' => new class (90, 'EUR') { + public function __construct(private int $cents, private string $currency) {} + }, + 'ref' => 'A-1', + ]; + + // L'aller-retour n'est pas supposé, il est fait : ce que `DbalEventStore` écrit, puis ce + // que `EventDataMapper` relit. Écrire son résultat en dur ferait passer ce test le jour où + // l'aplatissement changerait — c'est-à-dire le jour où le faux positif reviendrait. + $throughTheDatabase = json_decode( + json_encode( + (new ActivityScheduled(self::EXECUTION, 'act-1', 'charge', $freshPayload))->payload(), + \JSON_THROW_ON_ERROR, + ), + true, + 512, + \JSON_THROW_ON_ERROR, + )['payload']; + + self::assertNotSame($freshPayload, $throughTheDatabase, 'sans aplatissement, ce test ne prouve rien'); + + $store = new InMemoryEventStore(); + $store->append(new ActivityScheduled(self::EXECUTION, 'act-1', 'charge', $throughTheDatabase)); + $store->append(new ActivityCompleted(self::EXECUTION, 'act-1', 'ok')); + + $context = $this->context($store); + + // Ce que le code recalcule au replay : l'objet, toujours vivant. + $awaitable = $context->activity('charge', $freshPayload); + + self::assertTrue($awaitable->isSettled(), 'un replay fidèle ne doit pas mourir sur un objet opaque'); + } + + public function testAnIncomparablePayloadWaivesTheGuard(): void + { + // Une ressource ne s'encode pas. La garde n'a alors rien à comparer : elle se tait plutôt + // que d'accuser une charge qu'elle ne sait pas lire. + $store = $this->journalWith(['handle' => null]); + $context = $this->context($store); + + $handle = fopen('php://memory', 'r'); + self::assertIsResource($handle); + + $awaitable = $context->activity('weather', ['handle' => $handle]); + + self::assertTrue($awaitable->isSettled()); + fclose($handle); + } + + public function testAHistoryWithoutARecordedSlotIsNotADivergence(): void + { + // Un slot que personne n'a enregistré, c'est un workflow qui grandit — pas une divergence. + $context = $this->context(new InMemoryEventStore()); + + $awaitable = $context->activity('weather', ['city' => 'Paris']); + + self::assertFalse($awaitable->isSettled(), 'le slot est neuf : il part en planification'); + } + + /** + * @param array $payload + */ + private function journalWith(array $payload): InMemoryEventStore + { + $store = new InMemoryEventStore(); + $store->append(new ActivityScheduled(self::EXECUTION, 'act-1', 'weather', $payload)); + $store->append(new ActivityCompleted(self::EXECUTION, 'act-1', '22°C')); + + return $store; + } + + private function context(InMemoryEventStore $store): ExecutionContext + { + return new ExecutionContext( + self::EXECUTION, + new EventStoreHistorySource($store, self::EXECUTION), + new EventStoreCommandBuffer($store, new NoopActivityTransport(), self::EXECUTION), + ); + } +} From e983bdf0c9e2e26f558126c60121bb81abc65fcf Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gr=C3=A9gory=20Planchat?= Date: Fri, 4 Sep 2026 22:49:48 +0200 Subject: [PATCH 2/2] =?UTF-8?q?fix(coeur):=20la=20m=C3=AAme=20garde=20sur?= =?UTF-8?q?=20les=20slots=20Nexus=20et=20workflow=20enfant?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit La garde de charge ne couvrait que les activités. Les deux autres types de slot avaient le même trou : l'identité comparée, la charge jamais. Une seule méthode pour les trois — `refusePayloadDivergence()` — comme `refuseDivergence()` tient déjà l'identité des trois. Ce qui identifie un slot n'étant pas de même nature partout, le type et l'identité viennent de l'appelant : un nom pour une activité, un type pour un enfant, un triplet pour Nexus. Aucun champ ajouté à aucun événement : les trois charges étaient déjà sur le fil. Le désenveloppage, lui, diffère aux trois endroits, et c'est là qu'une erreur se serait cachée : - activité — `Payloads` d'un élément portant l'enveloppe TemporalActivityScheduleInput ; - Nexus — un `Payload` **nu**, l'enveloppe ayant été retirée du tampon (tâche 1.1) ; - enfant — `Payloads` d'un élément portant l'input nu. Les fixtures sont écrites d'après ce que le tampon produit, pas d'après le voisin : celle de NexusSlotDivergenceTest est antérieure au retrait de l'enveloppe et porte encore `{operationId, payload}`. Nexus n'est gardé que côté pont Temporal, et c'est structurel : le backend journal refuse ces opérations par construction (DUR036) et son événement ne porte que le site d'appel. `fromEvents()` permet de fabriquer un historique Temporal synthétique sans serveur — ce qui lève au passage la réserve de la PR sur le chemin Temporal des activités, désormais mesuré lui aussi. Vérifié par mutation : neutraliser la garde partagée fait tomber exactement les sept tests qui l'éprouvent, sur les trois types de slot. Co-Authored-By: Claude Opus 5 (1M context) --- UPGRADE.md | 21 +- .../Worker/TemporalExecutionHistory.php | 42 ++++ src/Durable/ExecutionContext.php | 52 ++++- .../Port/WorkflowHistorySourceInterface.php | 26 +++ src/Durable/Store/EventStoreHistorySource.php | 32 +++ .../Temporal/Worker/PayloadForSlotTest.php | 214 ++++++++++++++++++ .../ChildWorkflowSlotDivergenceTest.php | 38 ++++ 7 files changed, 406 insertions(+), 19 deletions(-) create mode 100644 tests/unit/Bridge/Temporal/Worker/PayloadForSlotTest.php diff --git a/UPGRADE.md b/UPGRADE.md index ee8d5416..4b897650 100644 --- a/UPGRADE.md +++ b/UPGRADE.md @@ -24,12 +24,13 @@ ne contient que ce que Rector sait faire sans deviner ; tout le reste est écrit ## 0.1.0-alpha8 -### La garde de divergence compare aussi la charge de l'activité +### La garde de divergence compare aussi la charge -`WorkflowHistorySourceInterface` gagne `activityPayloadForSlot(int $slot): ?array`. La garde de -divergence (DUR042) ne comparait que le **nom** de l'activité au slot ; elle compare désormais aussi -ses arguments. Un replay qui redemande la même activité avec une autre charge lève un -`WorkflowTaskFailure` au lieu de continuer en silence. +`WorkflowHistorySourceInterface` gagne trois méthodes — `activityPayloadForSlot()`, +`nexusOperationPayloadForSlot()` et `childWorkflowInputForSlot()`, toutes `?array`. La garde de +divergence (DUR042) ne comparait que l'**identité** du slot — nom d'activité, type d'enfant, +triplet Nexus ; elle compare désormais aussi la charge, sur les trois. Un replay qui redemande le +même appel avec une autre charge lève un `WorkflowTaskFailure` au lieu de continuer en silence. **Pourquoi** — le nom seul laissait passer la moitié du problème. Le journal servait l'ancien résultat, la charge fraîchement calculée partait à la poubelle, et l'exécution se terminait **en @@ -48,11 +49,13 @@ Une charge inencodable (ressource, `NAN`) désarme la garde plutôt que d'accuse pas lire. Les histoires écrites avant ce changement n'ont rien à comparer et passent inchangées. **Ce que Rector ne peut pas faire** — rien à réécrire dans le code appelant. Seules les -implémentations tierces de `WorkflowHistorySourceInterface` doivent ajouter la méthode ; rendre -`null` reproduit exactement le comportement d'avant, sans garde sur la charge. +implémentations tierces de `WorkflowHistorySourceInterface` doivent ajouter les trois méthodes ; +rendre `null` reproduit exactement le comportement d'avant, sans garde sur la charge. -**Hors périmètre** — les slots Nexus et workflow enfant comparent toujours leur seule identité. Le -trou y est le même, et l'enjeu plus grand côté Nexus, où un doublon part chez un tiers. +**Nexus** — la garde s'y exerce côté pont Temporal uniquement, et c'est structurel : le backend +journal refuse les opérations Nexus par construction (DUR036), et son événement +`NexusOperationScheduled` ne porte que le site d'appel. Aucun champ n'a été ajouté à aucun +événement : les trois charges étaient déjà sur le fil. ### Laravel refuse au démarrage un workflow dont les noms de paramètres divergent du contrat diff --git a/src/Bridge/Temporal/Worker/TemporalExecutionHistory.php b/src/Bridge/Temporal/Worker/TemporalExecutionHistory.php index e5fd0c82..0b0074f4 100644 --- a/src/Bridge/Temporal/Worker/TemporalExecutionHistory.php +++ b/src/Bridge/Temporal/Worker/TemporalExecutionHistory.php @@ -56,6 +56,12 @@ final class TemporalExecutionHistory implements WorkflowHistorySourceInterface /** @var array> activityId → charge planifiée (garde DUR042) */ private array $activityPayloads = []; + /** @var array> slot → charge Nexus planifiée (garde DUR042) */ + private array $nexusOperationPayloads = []; + + /** @var array> slot → input du workflow enfant (garde DUR042) */ + private array $childWorkflowInputs = []; + /** @var array scheduled event ID → activity ID */ private array $scheduledEventIdToActivityId = []; @@ -181,6 +187,18 @@ private function consumeEvent(HistoryEvent $event): void 'service' => (string) $attr->getService(), 'operation' => (string) $attr->getOperation(), ]; + + // La charge de l'appelant, **nue** : une opération Nexus porte un + // `Payload` et non des `Payloads`, et l'enveloppe `{operationId, payload}` + // a été retirée du tampon (tâche 1.1). Décoder autre chose ici comparerait + // une forme que le fil ne porte plus. + $nexusInput = $attr->getInput(); + if (null !== $nexusInput) { + $decodedInput = JsonPlainPayload::decode($nexusInput); + if (\is_array($decodedInput)) { + $this->nexusOperationPayloads[\count($this->scheduledNexusOperationIds) - 1] = $decodedInput; + } + } } } break; @@ -450,6 +468,20 @@ private function consumeEvent(HistoryEvent $event): void // Le type, en parallèle et au même index : c'est lui l'identité du slot, // l'identifiant d'exécution étant engendré. $this->childWorkflowTypes[] = (string) ($attr->getWorkflowType()?->getName() ?? ''); + + // `singlePayloads(encode($input))` côté tampon : une liste d'un élément, dont + // le premier est l'input nu. Une forme de plus que Nexus (Payload nu) et que + // l'activité (enveloppe) — les trois se ressemblent et ne se valent pas. + $childInput = $attr->getInput(); + if (null !== $childInput) { + $childPayloads = $childInput->getPayloads(); + if ($childPayloads->count() > 0) { + $decodedChild = JsonPlainPayload::decode($childPayloads[0]); + if (\is_array($decodedChild)) { + $this->childWorkflowInputs[\count($this->childExecutionIds) - 1] = $decodedChild; + } + } + } } break; @@ -542,6 +574,16 @@ public function activityNameForSlot(int $slot): ?string return '' === $name ? null : $name; } + public function nexusOperationPayloadForSlot(int $slot): ?array + { + return $this->nexusOperationPayloads[$slot] ?? null; + } + + public function childWorkflowInputForSlot(int $slot): ?array + { + return $this->childWorkflowInputs[$slot] ?? null; + } + public function activityPayloadForSlot(int $slot): ?array { $activityId = $this->scheduledActivityIds[$slot] ?? null; diff --git a/src/Durable/ExecutionContext.php b/src/Durable/ExecutionContext.php index 42f12909..b71b8484 100644 --- a/src/Durable/ExecutionContext.php +++ b/src/Durable/ExecutionContext.php @@ -103,7 +103,13 @@ public function activity(string $name, array $payload = [], ?ActivityOptions $op { $slotIndex = $this->activitySlotIndex++; $this->refuseActivityDivergence($slotIndex, $name); - $this->refuseActivityPayloadDivergence($slotIndex, $name, $payload); + $this->refusePayloadDivergence( + 'activity', + $slotIndex, + $name, + $this->historySource->activityPayloadForSlot($slotIndex), + $payload, + ); $replay = $this->historySource->findActivitySlotResult($slotIndex); if (null !== $replay) { $deferred = new \Gplanchat\Durable\Awaitable\Deferred(); @@ -163,6 +169,13 @@ public function nexusOperation( $this->historySource->nexusOperationSignatureForSlot($slotIndex), \sprintf('%s/%s/%s', $endpoint->name(), $service->name(), $operation->name()), ); + $this->refusePayloadDivergence( + 'Nexus operation', + $slotIndex, + \sprintf('%s/%s/%s', $endpoint->name(), $service->name(), $operation->name()), + $this->historySource->nexusOperationPayloadForSlot($slotIndex), + $payload, + ); $scheduled = $this->historySource->findScheduledNexusOperation($slotIndex); $operationId = $scheduled ?? $this->uuid(); @@ -279,7 +292,7 @@ private function refuseActivityDivergence(int $slotIndex, string $requested): vo } /** - * Refuse un slot dont le **nom** concorde mais dont les **arguments** ont changé au replay. + * Refuse un slot dont l'**identité** concorde mais dont la **charge** a changé au replay. * * Le nom seul laissait passer la moitié du problème : le journal servait l'ancien résultat, la * charge fraîchement calculée partait à la poubelle, et l'exécution se terminait en succès en @@ -291,13 +304,23 @@ private function refuseActivityDivergence(int $slotIndex, string $requested): vo * dont le journal ne retient rien ne doit pas faire diverger un replay fidèle, sous peine * d'arrêter en production des exécutions parfaitement saines. * - * @param array $requested + * La règle, une fois, pour les trois types de slot qui portent une charge — comme + * {@see refuseDivergence()} le fait pour les trois qui portent une identité. Le `$slotKind` et + * l'`$identity` viennent de l'appelant parce que ce qui identifie un slot n'est pas de même + * nature partout : un nom pour une activité, un type pour un enfant, un triplet pour Nexus. + * + * @param array|null $recorded + * @param array $requested * - * @throws WorkflowTaskFailure si le code redemande la même activité avec une autre charge + * @throws WorkflowTaskFailure si le code redemande le même appel avec une autre charge */ - private function refuseActivityPayloadDivergence(int $slotIndex, string $name, array $requested): void - { - $recorded = $this->historySource->activityPayloadForSlot($slotIndex); + private function refusePayloadDivergence( + string $slotKind, + int $slotIndex, + string $identity, + ?array $recorded, + array $requested, + ): void { if (null === $recorded) { // Rien d'enregistré ici : slot neuf, ou histoire écrite avant que la charge soit // lisible. Refuser laisserait sans garde exactement les exécutions qu'elle protège. @@ -314,16 +337,18 @@ private function refuseActivityPayloadDivergence(int $slotIndex, string $name, a $at = self::firstDifference($recordedPrint, $requestedPrint); throw new WorkflowTaskFailure(\sprintf( - 'Replay divergence at activity slot %d of execution "%s": the activity name "%s" matches, ' + 'Replay divergence at %s slot %d of execution "%s": "%s" is still the same %s, ' . 'but its payload changed at byte %d. History recorded %s, code scheduled %s ' . '(%d and %d bytes). ' . 'This is non-deterministic workflow code — the payload is rebuilt on every replay pass, ' . 'so something in it reads the clock, draws a random value, or is resolved from outside ' . 'the journal. This is not a version skew: a declared change point would have changed ' - . 'the name or the slot, not the arguments alone.', + . 'the identity or the slot, not the payload alone.', + $slotKind, $slotIndex, $this->executionId, - $name, + $identity, + $slotKind, $at, self::windowAround($recordedPrint, $at), self::windowAround($requestedPrint, $at), @@ -503,6 +528,13 @@ public function executeChildWorkflow(string $childWorkflowType, array $input = [ $this->historySource->childWorkflowTypeForSlot($slotIndex), $childWorkflowType, ); + $this->refusePayloadDivergence( + 'child workflow', + $slotIndex, + $childWorkflowType, + $this->historySource->childWorkflowInputForSlot($slotIndex), + $input, + ); $replay = $this->historySource->findChildWorkflowForSlot($slotIndex); $deferred = new \Gplanchat\Durable\Awaitable\Deferred(); if (null !== $replay) { diff --git a/src/Durable/Port/WorkflowHistorySourceInterface.php b/src/Durable/Port/WorkflowHistorySourceInterface.php index ef36fb30..78f78c78 100644 --- a/src/Durable/Port/WorkflowHistorySourceInterface.php +++ b/src/Durable/Port/WorkflowHistorySourceInterface.php @@ -121,6 +121,19 @@ public function findScheduledChildExecutionId(int $slot): ?string; */ public function childWorkflowTypeForSlot(int $slot): ?string; + /** + * Returns the input recorded for child workflow slot N, or null if there is nothing to compare. + * + * Same null/`[]` rule as {@see activityPayloadForSlot()}: only null waives the guard. + * + * The type identifies the child; this is what it was started with. A replay that keeps the + * type and changes the input starts nothing new — the journal already holds the child's + * outcome — so the divergence would otherwise never surface. + * + * @return array|null + */ + public function childWorkflowInputForSlot(int $slot): ?array; + /** * Returns the identity of the Nexus operation recorded at slot N, or null if none was. * @@ -130,6 +143,19 @@ public function childWorkflowTypeForSlot(int $slot): ?string; */ public function nexusOperationSignatureForSlot(int $slot): ?string; + /** + * Returns the payload recorded for Nexus operation slot N, or null if there is nothing to compare. + * + * Same null/`[]` rule as {@see activityPayloadForSlot()}: only null waives the guard. + * + * The identity above answers "which service"; this one answers "with what". Getting it wrong + * costs more here than anywhere else — a duplicated activity lands on a worker of one's own, a + * Nexus operation lands on a third party, where the duplicate is theirs. + * + * @return array|null + */ + public function nexusOperationPayloadForSlot(int $slot): ?array; + /** * Returns the Nth recorded message, in recorded order, or null past the end. * diff --git a/src/Durable/Store/EventStoreHistorySource.php b/src/Durable/Store/EventStoreHistorySource.php index 3f44190e..d4d362d0 100644 --- a/src/Durable/Store/EventStoreHistorySource.php +++ b/src/Durable/Store/EventStoreHistorySource.php @@ -146,6 +146,38 @@ public function activityPayloadForSlot(int $slot): ?array return null; } + public function childWorkflowInputForSlot(int $slot): ?array + { + $index = 0; + foreach ($this->eventStore->readStream($this->executionId) as $event) { + if ($event instanceof ChildWorkflowScheduled) { + if ($index === $slot) { + $input = $event->payload()['input'] ?? null; + + return \is_array($input) ? $input : null; + } + ++$index; + } + } + + return null; + } + + /** + * Toujours null, et ce n'est pas un oubli. + * + * Ce backend refuse les opérations Nexus par construction (DUR036) : aucun de ses historiques + * n'en porte une que le workflow aurait planifiée. Le seul `NexusOperationScheduled` qui puisse + * traverser un flux vient du convertisseur du profileur, qui l'écrit pour l'affichage — et cet + * événement ne porte que le site d'appel, jamais la charge. Il n'y a donc rien à comparer. + * + * La garde s'exerce là où Nexus existe : {@see \Gplanchat\Bridge\Temporal\Worker\TemporalExecutionHistory}. + */ + public function nexusOperationPayloadForSlot(int $slot): ?array + { + return null; + } + public function childWorkflowTypeForSlot(int $slot): ?string { $index = 0; diff --git a/tests/unit/Bridge/Temporal/Worker/PayloadForSlotTest.php b/tests/unit/Bridge/Temporal/Worker/PayloadForSlotTest.php new file mode 100644 index 00000000..9a1d1063 --- /dev/null +++ b/tests/unit/Bridge/Temporal/Worker/PayloadForSlotTest.php @@ -0,0 +1,214 @@ +activityScheduled(5, 'act-1', 'weather', ['city' => 'Paris']), + ]); + + // Les arguments, pas l'enveloppe : rendre l'enveloppe comparerait `activityName` deux fois + // et laisserait passer un changement d'argument. + self::assertSame(['city' => 'Paris'], $history->activityPayloadForSlot(0)); + } + + public function testTheNexusPayloadIsReadNaked(): void + { + $history = TemporalExecutionHistory::fromEvents([ + $this->nexusScheduled(5, 'paiements', 'facturation', 'encaisser', ['amount' => 90]), + ]); + + self::assertSame(['amount' => 90], $history->nexusOperationPayloadForSlot(0)); + } + + public function testTheChildInputIsReadFromItsSinglePayload(): void + { + $history = TemporalExecutionHistory::fromEvents([ + $this->childScheduled('child-1', 'ChargeCardWorkflow', ['sku' => 'ABC']), + ]); + + self::assertSame(['sku' => 'ABC'], $history->childWorkflowInputForSlot(0)); + } + + public function testEachSlotKeepsItsOwnPayload(): void + { + $history = TemporalExecutionHistory::fromEvents([ + $this->nexusScheduled(5, 'paiements', 'facturation', 'encaisser', ['amount' => 90]), + $this->nexusScheduled(9, 'stocks', 'entrepot', 'reserver', ['sku' => 'ABC']), + $this->childScheduled('child-1', 'A', ['n' => 1]), + $this->childScheduled('child-2', 'B', ['n' => 2]), + ]); + + self::assertSame(['amount' => 90], $history->nexusOperationPayloadForSlot(0)); + self::assertSame(['sku' => 'ABC'], $history->nexusOperationPayloadForSlot(1)); + self::assertSame(['n' => 1], $history->childWorkflowInputForSlot(0)); + self::assertSame(['n' => 2], $history->childWorkflowInputForSlot(1)); + } + + public function testASlotNobodyScheduledHasNoPayload(): void + { + // Null, et non `[]` : « rien enregistré » désarme la garde, « planifié sans argument » non. + $history = TemporalExecutionHistory::fromEvents([]); + + self::assertNull($history->activityPayloadForSlot(0)); + self::assertNull($history->nexusOperationPayloadForSlot(0)); + self::assertNull($history->childWorkflowInputForSlot(0)); + } + + public function testAnEmptyPayloadIsRecordedAsEmptyNotAsAbsent(): void + { + $history = TemporalExecutionHistory::fromEvents([ + $this->nexusScheduled(5, 'paiements', 'facturation', 'ping', []), + ]); + + self::assertSame([], $history->nexusOperationPayloadForSlot(0)); + } + + public function testTheGuardRefusesANexusOperationWhosePayloadChanged(): void + { + // Le câblage, pas seulement la lecture : c'est ici que Nexus mérite le plus la garde — + // une activité replanifiée retombe sur un worker à soi, une opération Nexus part chez un + // tiers, où le doublon est le sien. + $context = new ExecutionContext( + 'exec-nexus', + TemporalExecutionHistory::fromEvents([ + $this->nexusScheduled(5, 'paiements', 'facturation', 'encaisser', ['amount' => 90]), + ]), + $this->createStub(WorkflowCommandBufferInterface::class), + ); + + try { + $context->nexusOperation( + NexusEndpoint::named('paiements'), + NexusService::named('facturation'), + NexusOperationName::named('encaisser'), + ['amount' => 120], + ); + self::fail('La divergence de charge aurait dû être refusée.'); + } catch (WorkflowTaskFailure $refusal) { + $message = $refusal->getMessage(); + } + + self::assertStringContainsString('Nexus operation slot 0', $message); + self::assertStringContainsString('"paiements/facturation/encaisser" is still the same Nexus operation', $message); + } + + public function testAFaithfulNexusReplayIsNotRefused(): void + { + $context = new ExecutionContext( + 'exec-nexus', + TemporalExecutionHistory::fromEvents([ + $this->nexusScheduled(5, 'paiements', 'facturation', 'encaisser', ['amount' => 90]), + ]), + $this->createStub(WorkflowCommandBufferInterface::class), + ); + + $awaitable = $context->nexusOperation( + NexusEndpoint::named('paiements'), + NexusService::named('facturation'), + NexusOperationName::named('encaisser'), + ['amount' => 90], + ); + + self::assertNotNull($awaitable); + } + + /** + * @param array $payload + */ + private function activityScheduled(int $eventId, string $activityId, string $name, array $payload): HistoryEvent + { + $attrs = new ActivityTaskScheduledEventAttributes(); + $attrs->setActivityId($activityId); + $attrs->setActivityType(new \Temporal\Api\Common\V1\ActivityType(['name' => $name])); + $attrs->setInput(JsonPlainPayload::singlePayloads(JsonPlainPayload::encode([ + 'executionId' => 'exec-1', + 'activityId' => $activityId, + 'activityName' => $name, + 'payload' => $payload, + 'metadata' => [], + ]))); + + $event = new HistoryEvent(); + $event->setEventType(EventType::EVENT_TYPE_ACTIVITY_TASK_SCHEDULED); + $event->setEventId($eventId); + $event->setActivityTaskScheduledEventAttributes($attrs); + + return $event; + } + + /** + * @param array $payload + */ + private function nexusScheduled(int $eventId, string $endpoint, string $service, string $operation, array $payload): HistoryEvent + { + $attrs = new NexusOperationScheduledEventAttributes(); + $attrs->setEndpoint($endpoint); + $attrs->setService($service); + $attrs->setOperation($operation); + // Nu, comme le tampon l'écrit. + $attrs->setInput(JsonPlainPayload::encode($payload)); + + $event = new HistoryEvent(); + $event->setEventType(EventType::EVENT_TYPE_NEXUS_OPERATION_SCHEDULED); + $event->setEventId($eventId); + $event->setNexusOperationScheduledEventAttributes($attrs); + + return $event; + } + + /** + * @param array $input + */ + private function childScheduled(string $workflowId, string $type, array $input): HistoryEvent + { + $attrs = new StartChildWorkflowExecutionInitiatedEventAttributes(); + $attrs->setWorkflowId($workflowId); + $attrs->setWorkflowType(new WorkflowType(['name' => $type])); + $attrs->setInput(JsonPlainPayload::singlePayloads(JsonPlainPayload::encode($input))); + + $event = new HistoryEvent(); + $event->setEventType(EventType::EVENT_TYPE_START_CHILD_WORKFLOW_EXECUTION_INITIATED); + $event->setStartChildWorkflowExecutionInitiatedEventAttributes($attrs); + + return $event; + } +} diff --git a/tests/unit/Durable/Replay/ChildWorkflowSlotDivergenceTest.php b/tests/unit/Durable/Replay/ChildWorkflowSlotDivergenceTest.php index 7aa30048..aa43a7dc 100644 --- a/tests/unit/Durable/Replay/ChildWorkflowSlotDivergenceTest.php +++ b/tests/unit/Durable/Replay/ChildWorkflowSlotDivergenceTest.php @@ -78,6 +78,44 @@ public function testAnUnchangedChildTypeStillReplays(): void self::assertNotNull($awaitable, "Le type inchangé ne doit pas diverger : l'identifiant d'exécution engendré n'entre pas dans la comparaison."); } + public function testTheSameChildStartedWithAnotherInputIsRefused(): void + { + // Le type concorde, l'input non. Sans cette garde la divergence ne remonterait jamais : + // le journal tient déjà l'issue de l'enfant, donc rien de neuf n'est démarré et le nouvel + // input part à la poubelle en silence. + $context = $this->contextWithChild('ChargeCardWorkflow'); + + $this->expectException(WorkflowTaskFailure::class); + $this->expectExceptionMessageMatches('/payload changed/'); + + $context->executeChildWorkflow('ChargeCardWorkflow', ['sku' => 'XYZ']); + } + + public function testTheInputRefusalNamesTheChildAndTheKind(): void + { + $context = $this->contextWithChild('ChargeCardWorkflow'); + + try { + $context->executeChildWorkflow('ChargeCardWorkflow', ['sku' => 'XYZ']); + self::fail('La divergence de charge aurait dû être refusée.'); + } catch (WorkflowTaskFailure $e) { + $message = $e->getMessage(); + } + + self::assertStringContainsString('child workflow slot 0', $message); + self::assertStringContainsString('"ChargeCardWorkflow" is still the same child workflow', $message); + self::assertStringContainsString('non-deterministic workflow code', $message); + } + + public function testAnUnchangedInputStillReplays(): void + { + $context = $this->contextWithChild('ChargeCardWorkflow'); + + $awaitable = $context->executeChildWorkflow('ChargeCardWorkflow', ['sku' => 'ABC']); + + self::assertNotNull($awaitable, 'Un replay fidèle ne doit pas diverger sur son propre input.'); + } + private function contextWithChild(string $childType): ExecutionContext { $store = new InMemoryEventStore();