Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 34 additions & 0 deletions UPGRADE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
69 changes: 69 additions & 0 deletions src/Bridge/Temporal/Worker/TemporalExecutionHistory.php
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,15 @@ final class TemporalExecutionHistory implements WorkflowHistorySourceInterface
/** @var array<string, string> activityId → nom d'activité (pour typer les échecs) */
private array $activityNames = [];

/** @var array<string, array<string, mixed>> activityId → charge planifiée (garde DUR042) */
private array $activityPayloads = [];

/** @var array<int, array<string, mixed>> slot → charge Nexus planifiée (garde DUR042) */
private array $nexusOperationPayloads = [];

/** @var array<int, array<string, mixed>> slot → input du workflow enfant (garde DUR042) */
private array $childWorkflowInputs = [];

/** @var array<int, string> scheduled event ID → activity ID */
private array $scheduledEventIdToActivityId = [];

Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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;

Expand Down Expand Up @@ -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;

Expand Down Expand Up @@ -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;
Expand Down
171 changes: 171 additions & 0 deletions src/Durable/ExecutionContext.php
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, \Gplanchat\Durable\Awaitable\Deferred> */
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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();

Expand Down Expand Up @@ -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<string, mixed>|null $recorded
* @param array<string, mixed> $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<string, mixed> $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é.
*
Expand Down Expand Up @@ -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) {
Expand Down
Loading
Loading