diff --git a/UPGRADE.md b/UPGRADE.md index 074f7922..4b897650 100644 --- a/UPGRADE.md +++ b/UPGRADE.md @@ -24,6 +24,40 @@ 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 + +`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 +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 les trois méthodes ; +rendre `null` reproduit exactement le comportement d'avant, sans garde sur la charge. + +**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 `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..0b0074f4 100644 --- a/src/Bridge/Temporal/Worker/TemporalExecutionHistory.php +++ b/src/Bridge/Temporal/Worker/TemporalExecutionHistory.php @@ -53,6 +53,15 @@ 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> 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 = []; @@ -178,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; @@ -241,6 +262,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; @@ -433,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; @@ -525,6 +574,26 @@ 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; + 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..b71b8484 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,13 @@ public function activity(string $name, array $payload = [], ?ActivityOptions $op { $slotIndex = $this->activitySlotIndex++; $this->refuseActivityDivergence($slotIndex, $name); + $this->refusePayloadDivergence( + 'activity', + $slotIndex, + $name, + $this->historySource->activityPayloadForSlot($slotIndex), + $payload, + ); $replay = $this->historySource->findActivitySlotResult($slotIndex); if (null !== $replay) { $deferred = new \Gplanchat\Durable\Awaitable\Deferred(); @@ -159,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(); @@ -274,6 +291,153 @@ private function refuseActivityDivergence(int $slotIndex, string $requested): vo $this->refuseDivergence('activity', $slotIndex, $this->historySource->activityNameForSlot($slotIndex), $requested); } + /** + * 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 + * 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. + * + * 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 le même appel avec une autre charge + */ + 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. + 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 %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 identity or the slot, not the payload alone.', + $slotKind, + $slotIndex, + $this->executionId, + $identity, + $slotKind, + $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é. * @@ -364,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 5f1d292d..78f78c78 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. @@ -105,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. * @@ -114,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 b13c566c..d4d362d0 100644 --- a/src/Durable/Store/EventStoreHistorySource.php +++ b/src/Durable/Store/EventStoreHistorySource.php @@ -127,6 +127,57 @@ 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 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/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), + ); + } +} 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();