diff --git a/.changeset/ordered-barrier-safety-net.md b/.changeset/ordered-barrier-safety-net.md new file mode 100644 index 0000000000..cb272809fb --- /dev/null +++ b/.changeset/ordered-barrier-safety-net.md @@ -0,0 +1,6 @@ +--- +'@workflow/core': patch +'workflow': patch +--- + +Fix a replay-determinism gap where branch wake order — and therefore step correlation ids — could depend on how much of the event log an invocation had loaded, corrupting runs under concurrent replays (CORRUPTED_EVENT_LOG). diff --git a/packages/core/src/__fixtures__/wrun-41KZYJ92TP-storm-log.json b/packages/core/src/__fixtures__/wrun-41KZYJ92TP-storm-log.json new file mode 100644 index 0000000000..2aa9b51350 --- /dev/null +++ b/packages/core/src/__fixtures__/wrun-41KZYJ92TP-storm-log.json @@ -0,0 +1,1326 @@ +[ + { "s": 3, "t": "hook_created", "k": "hook", "r": 0, "n": "" }, + { "s": 4, "t": "wait_created", "k": "wait", "r": 12, "n": "" }, + { "s": 5, "t": "step_created", "k": "step", "r": 7, "n": "settleStep" }, + { "s": 6, "t": "step_created", "k": "step", "r": 15, "n": "settleStep" }, + { "s": 7, "t": "wait_created", "k": "wait", "r": 6, "n": "" }, + { "s": 8, "t": "wait_created", "k": "wait", "r": 2, "n": "" }, + { "s": 9, "t": "step_created", "k": "step", "r": 9, "n": "settleStep" }, + { "s": 10, "t": "wait_created", "k": "wait", "r": 8, "n": "" }, + { "s": 11, "t": "wait_created", "k": "wait", "r": 4, "n": "" }, + { "s": 12, "t": "step_created", "k": "step", "r": 11, "n": "settleStep" }, + { "s": 13, "t": "step_created", "k": "step", "r": 13, "n": "settleStep" }, + { "s": 14, "t": "wait_created", "k": "wait", "r": 16, "n": "" }, + { "s": 15, "t": "wait_created", "k": "wait", "r": 14, "n": "" }, + { "s": 16, "t": "wait_created", "k": "wait", "r": 10, "n": "" }, + { "s": 17, "t": "step_created", "k": "step", "r": 5, "n": "settleStep" }, + { "s": 18, "t": "step_started", "k": "step", "r": 5, "n": "settleStep" }, + { "s": 19, "t": "step_started", "k": "step", "r": 9, "n": "settleStep" }, + { "s": 20, "t": "step_started", "k": "step", "r": 15, "n": "settleStep" }, + { "s": 21, "t": "step_created", "k": "step", "r": 1, "n": "settleStep" }, + { "s": 22, "t": "step_started", "k": "step", "r": 1, "n": "settleStep" }, + { "s": 23, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 24, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 25, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 26, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 27, "t": "step_started", "k": "step", "r": 13, "n": "settleStep" }, + { "s": 28, "t": "step_started", "k": "step", "r": 7, "n": "settleStep" }, + { "s": 29, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 30, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 31, "t": "step_started", "k": "step", "r": 11, "n": "settleStep" }, + { "s": 32, "t": "step_created", "k": "step", "r": 3, "n": "settleStep" }, + { "s": 33, "t": "step_started", "k": "step", "r": 3, "n": "settleStep" }, + { "s": 34, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 35, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 36, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 37, "t": "wait_completed", "k": "wait", "r": 12, "n": "" }, + { "s": 38, "t": "step_completed", "k": "step", "r": 5, "n": "settleStep" }, + { "s": 39, "t": "wait_completed", "k": "wait", "r": 6, "n": "" }, + { "s": 40, "t": "step_completed", "k": "step", "r": 9, "n": "settleStep" }, + { "s": 41, "t": "step_completed", "k": "step", "r": 1, "n": "settleStep" }, + { "s": 42, "t": "step_completed", "k": "step", "r": 15, "n": "settleStep" }, + { "s": 43, "t": "step_completed", "k": "step", "r": 13, "n": "settleStep" }, + { "s": 44, "t": "step_completed", "k": "step", "r": 7, "n": "settleStep" }, + { "s": 45, "t": "step_completed", "k": "step", "r": 3, "n": "settleStep" }, + { "s": 46, "t": "wait_completed", "k": "wait", "r": 2, "n": "" }, + { "s": 47, "t": "step_completed", "k": "step", "r": 11, "n": "settleStep" }, + { "s": 48, "t": "wait_completed", "k": "wait", "r": 8, "n": "" }, + { "s": 49, "t": "wait_completed", "k": "wait", "r": 4, "n": "" }, + { "s": 50, "t": "wait_completed", "k": "wait", "r": 16, "n": "" }, + { "s": 51, "t": "wait_completed", "k": "wait", "r": 14, "n": "" }, + { "s": 52, "t": "wait_completed", "k": "wait", "r": 10, "n": "" }, + { "s": 53, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 54, "t": "step_created", "k": "step", "r": 20, "n": "finalizeStep" }, + { "s": 55, "t": "step_created", "k": "step", "r": 22, "n": "finalizeStep" }, + { "s": 56, "t": "step_created", "k": "step", "r": 23, "n": "finalizeStep" }, + { "s": 57, "t": "step_created", "k": "step", "r": 24, "n": "finalizeStep" }, + { "s": 58, "t": "step_created", "k": "step", "r": 21, "n": "finalizeStep" }, + { "s": 59, "t": "step_started", "k": "step", "r": 24, "n": "finalizeStep" }, + { "s": 60, "t": "step_started", "k": "step", "r": 22, "n": "finalizeStep" }, + { "s": 61, "t": "step_completed", "k": "step", "r": 24, "n": "finalizeStep" }, + { "s": 62, "t": "step_completed", "k": "step", "r": 22, "n": "finalizeStep" }, + { "s": 63, "t": "step_created", "k": "step", "r": 19, "n": "finalizeStep" }, + { "s": 64, "t": "step_started", "k": "step", "r": 19, "n": "finalizeStep" }, + { "s": 65, "t": "step_created", "k": "step", "r": 18, "n": "finalizeStep" }, + { "s": 66, "t": "step_started", "k": "step", "r": 18, "n": "finalizeStep" }, + { "s": 67, "t": "step_created", "k": "step", "r": 17, "n": "recoverStep" }, + { "s": 68, "t": "step_started", "k": "step", "r": 17, "n": "recoverStep" }, + { "s": 69, "t": "step_started", "k": "step", "r": 21, "n": "finalizeStep" }, + { "s": 70, "t": "step_started", "k": "step", "r": 23, "n": "finalizeStep" }, + { "s": 71, "t": "step_completed", "k": "step", "r": 19, "n": "finalizeStep" }, + { "s": 72, "t": "step_completed", "k": "step", "r": 21, "n": "finalizeStep" }, + { "s": 73, "t": "step_started", "k": "step", "r": 20, "n": "finalizeStep" }, + { "s": 74, "t": "step_completed", "k": "step", "r": 17, "n": "recoverStep" }, + { "s": 75, "t": "step_completed", "k": "step", "r": 18, "n": "finalizeStep" }, + { "s": 76, "t": "step_completed", "k": "step", "r": 23, "n": "finalizeStep" }, + { "s": 77, "t": "step_completed", "k": "step", "r": 20, "n": "finalizeStep" }, + { "s": 78, "t": "step_created", "k": "step", "r": 25, "n": "releaseStep" }, + { "s": 79, "t": "step_started", "k": "step", "r": 25, "n": "releaseStep" }, + { "s": 80, "t": "step_created", "k": "step", "r": 26, "n": "releaseStep" }, + { "s": 81, "t": "step_started", "k": "step", "r": 26, "n": "releaseStep" }, + { "s": 82, "t": "step_created", "k": "step", "r": 27, "n": "releaseStep" }, + { "s": 83, "t": "step_started", "k": "step", "r": 27, "n": "releaseStep" }, + { "s": 84, "t": "step_completed", "k": "step", "r": 25, "n": "releaseStep" }, + { "s": 85, "t": "step_completed", "k": "step", "r": 26, "n": "releaseStep" }, + { "s": 86, "t": "step_created", "k": "step", "r": 30, "n": "releaseStep" }, + { "s": 87, "t": "step_completed", "k": "step", "r": 27, "n": "releaseStep" }, + { "s": 88, "t": "step_created", "k": "step", "r": 29, "n": "finalizeStep" }, + { "s": 89, "t": "step_created", "k": "step", "r": 28, "n": "releaseStep" }, + { "s": 90, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 91, "t": "step_started", "k": "step", "r": 28, "n": "releaseStep" }, + { "s": 92, "t": "step_started", "k": "step", "r": 29, "n": "finalizeStep" }, + { "s": 93, "t": "step_created", "k": "step", "r": 31, "n": "releaseStep" }, + { "s": 94, "t": "step_started", "k": "step", "r": 31, "n": "releaseStep" }, + { "s": 95, "t": "step_created", "k": "step", "r": 32, "n": "releaseStep" }, + { "s": 96, "t": "step_started", "k": "step", "r": 32, "n": "releaseStep" }, + { "s": 97, "t": "step_completed", "k": "step", "r": 28, "n": "releaseStep" }, + { "s": 98, "t": "step_completed", "k": "step", "r": 32, "n": "releaseStep" }, + { "s": 99, "t": "step_completed", "k": "step", "r": 29, "n": "finalizeStep" }, + { "s": 100, "t": "step_completed", "k": "step", "r": 31, "n": "releaseStep" }, + { "s": 101, "t": "step_started", "k": "step", "r": 30, "n": "releaseStep" }, + { "s": 102, "t": "step_created", "k": "step", "r": 33, "n": "releaseStep" }, + { "s": 103, "t": "step_started", "k": "step", "r": 33, "n": "releaseStep" }, + { "s": 104, "t": "step_completed", "k": "step", "r": 30, "n": "releaseStep" }, + { "s": 105, "t": "step_completed", "k": "step", "r": 33, "n": "releaseStep" }, + { "s": 106, "t": "step_created", "k": "step", "r": 36, "n": "reconcileStep" }, + { "s": 107, "t": "step_started", "k": "step", "r": 36, "n": "reconcileStep" }, + { "s": 108, "t": "step_created", "k": "step", "r": 34, "n": "reconcileStep" }, + { "s": 109, "t": "step_started", "k": "step", "r": 34, "n": "reconcileStep" }, + { "s": 110, "t": "step_created", "k": "step", "r": 35, "n": "reconcileStep" }, + { "s": 111, "t": "step_started", "k": "step", "r": 35, "n": "reconcileStep" }, + { + "s": 112, + "t": "step_completed", + "k": "step", + "r": 36, + "n": "reconcileStep" + }, + { + "s": 113, + "t": "step_completed", + "k": "step", + "r": 34, + "n": "reconcileStep" + }, + { + "s": 114, + "t": "step_completed", + "k": "step", + "r": 35, + "n": "reconcileStep" + }, + { "s": 115, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 116, "t": "wait_created", "k": "wait", "r": 37, "n": "" }, + { "s": 117, "t": "wait_completed", "k": "wait", "r": 37, "n": "" }, + { "s": 118, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 119, "t": "wait_created", "k": "wait", "r": 53, "n": "" }, + { "s": 120, "t": "wait_created", "k": "wait", "r": 49, "n": "" }, + { "s": 121, "t": "step_created", "k": "step", "r": 52, "n": "settleStep" }, + { "s": 122, "t": "wait_created", "k": "wait", "r": 47, "n": "" }, + { "s": 123, "t": "step_created", "k": "step", "r": 48, "n": "settleStep" }, + { "s": 124, "t": "wait_created", "k": "wait", "r": 51, "n": "" }, + { "s": 125, "t": "wait_created", "k": "wait", "r": 45, "n": "" }, + { "s": 126, "t": "wait_created", "k": "wait", "r": 43, "n": "" }, + { "s": 127, "t": "step_created", "k": "step", "r": 50, "n": "settleStep" }, + { "s": 128, "t": "step_created", "k": "step", "r": 46, "n": "settleStep" }, + { "s": 129, "t": "wait_created", "k": "wait", "r": 41, "n": "" }, + { "s": 130, "t": "step_created", "k": "step", "r": 44, "n": "settleStep" }, + { "s": 131, "t": "wait_created", "k": "wait", "r": 39, "n": "" }, + { "s": 132, "t": "step_created", "k": "step", "r": 38, "n": "settleStep" }, + { "s": 133, "t": "step_started", "k": "step", "r": 38, "n": "settleStep" }, + { "s": 134, "t": "step_created", "k": "step", "r": 42, "n": "settleStep" }, + { "s": 135, "t": "step_started", "k": "step", "r": 42, "n": "settleStep" }, + { "s": 136, "t": "step_started", "k": "step", "r": 52, "n": "settleStep" }, + { "s": 137, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 138, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 139, "t": "step_started", "k": "step", "r": 46, "n": "settleStep" }, + { "s": 140, "t": "step_created", "k": "step", "r": 40, "n": "settleStep" }, + { "s": 141, "t": "step_started", "k": "step", "r": 40, "n": "settleStep" }, + { "s": 142, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 143, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 144, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 145, "t": "step_started", "k": "step", "r": 50, "n": "settleStep" }, + { "s": 146, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 147, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 148, "t": "step_started", "k": "step", "r": 48, "n": "settleStep" }, + { "s": 149, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 150, "t": "step_started", "k": "step", "r": 44, "n": "settleStep" }, + { "s": 151, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 152, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 153, "t": "wait_completed", "k": "wait", "r": 53, "n": "" }, + { "s": 154, "t": "wait_completed", "k": "wait", "r": 49, "n": "" }, + { "s": 155, "t": "wait_completed", "k": "wait", "r": 47, "n": "" }, + { "s": 156, "t": "wait_completed", "k": "wait", "r": 51, "n": "" }, + { "s": 157, "t": "step_completed", "k": "step", "r": 38, "n": "settleStep" }, + { "s": 158, "t": "wait_completed", "k": "wait", "r": 45, "n": "" }, + { "s": 159, "t": "step_completed", "k": "step", "r": 42, "n": "settleStep" }, + { "s": 160, "t": "wait_completed", "k": "wait", "r": 43, "n": "" }, + { "s": 161, "t": "step_completed", "k": "step", "r": 40, "n": "settleStep" }, + { "s": 162, "t": "wait_completed", "k": "wait", "r": 41, "n": "" }, + { "s": 163, "t": "step_completed", "k": "step", "r": 46, "n": "settleStep" }, + { "s": 164, "t": "wait_completed", "k": "wait", "r": 39, "n": "" }, + { "s": 165, "t": "step_completed", "k": "step", "r": 52, "n": "settleStep" }, + { "s": 166, "t": "step_completed", "k": "step", "r": 50, "n": "settleStep" }, + { "s": 167, "t": "step_created", "k": "step", "r": 60, "n": "finalizeStep" }, + { "s": 168, "t": "step_created", "k": "step", "r": 57, "n": "recoverStep" }, + { "s": 169, "t": "step_created", "k": "step", "r": 59, "n": "recoverStep" }, + { "s": 170, "t": "step_created", "k": "step", "r": 61, "n": "finalizeStep" }, + { "s": 171, "t": "step_created", "k": "step", "r": 58, "n": "finalizeStep" }, + { "s": 172, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 173, "t": "step_created", "k": "step", "r": 54, "n": "recoverStep" }, + { "s": 174, "t": "step_started", "k": "step", "r": 54, "n": "recoverStep" }, + { "s": 175, "t": "step_created", "k": "step", "r": 56, "n": "recoverStep" }, + { "s": 176, "t": "step_started", "k": "step", "r": 56, "n": "recoverStep" }, + { "s": 177, "t": "step_created", "k": "step", "r": 55, "n": "recoverStep" }, + { "s": 178, "t": "step_started", "k": "step", "r": 55, "n": "recoverStep" }, + { "s": 179, "t": "step_started", "k": "step", "r": 61, "n": "finalizeStep" }, + { "s": 180, "t": "step_completed", "k": "step", "r": 56, "n": "recoverStep" }, + { "s": 181, "t": "step_completed", "k": "step", "r": 55, "n": "recoverStep" }, + { "s": 182, "t": "step_completed", "k": "step", "r": 48, "n": "settleStep" }, + { + "s": 183, + "t": "step_completed", + "k": "step", + "r": 61, + "n": "finalizeStep" + }, + { "s": 184, "t": "step_started", "k": "step", "r": 59, "n": "recoverStep" }, + { "s": 185, "t": "step_started", "k": "step", "r": 57, "n": "recoverStep" }, + { "s": 186, "t": "step_started", "k": "step", "r": 60, "n": "finalizeStep" }, + { "s": 187, "t": "step_completed", "k": "step", "r": 59, "n": "recoverStep" }, + { "s": 188, "t": "step_completed", "k": "step", "r": 54, "n": "recoverStep" }, + { "s": 189, "t": "step_completed", "k": "step", "r": 44, "n": "settleStep" }, + { "s": 190, "t": "step_completed", "k": "step", "r": 57, "n": "recoverStep" }, + { "s": 191, "t": "step_started", "k": "step", "r": 58, "n": "finalizeStep" }, + { + "s": 192, + "t": "step_completed", + "k": "step", + "r": 60, + "n": "finalizeStep" + }, + { "s": 193, "t": "step_created", "k": "step", "r": 63, "n": "finalizeStep" }, + { "s": 194, "t": "step_started", "k": "step", "r": 63, "n": "finalizeStep" }, + { + "s": 195, + "t": "step_completed", + "k": "step", + "r": 58, + "n": "finalizeStep" + }, + { "s": 196, "t": "step_created", "k": "step", "r": 64, "n": "releaseStep" }, + { "s": 197, "t": "step_started", "k": "step", "r": 64, "n": "releaseStep" }, + { + "s": 198, + "t": "step_completed", + "k": "step", + "r": 63, + "n": "finalizeStep" + }, + { "s": 199, "t": "step_created", "k": "step", "r": 62, "n": "finalizeStep" }, + { "s": 200, "t": "step_started", "k": "step", "r": 62, "n": "finalizeStep" }, + { "s": 201, "t": "step_completed", "k": "step", "r": 64, "n": "releaseStep" }, + { "s": 202, "t": "step_created", "k": "step", "r": 66, "n": "finalizeStep" }, + { "s": 203, "t": "step_created", "k": "step", "r": 65, "n": "finalizeStep" }, + { "s": 204, "t": "step_created", "k": "step", "r": 68, "n": "releaseStep" }, + { + "s": 205, + "t": "step_completed", + "k": "step", + "r": 62, + "n": "finalizeStep" + }, + { "s": 206, "t": "step_created", "k": "step", "r": 67, "n": "finalizeStep" }, + { "s": 207, "t": "step_started", "k": "step", "r": 66, "n": "finalizeStep" }, + { "s": 208, "t": "step_started", "k": "step", "r": 65, "n": "finalizeStep" }, + { + "s": 209, + "t": "step_completed", + "k": "step", + "r": 65, + "n": "finalizeStep" + }, + { "s": 210, "t": "step_started", "k": "step", "r": 67, "n": "finalizeStep" }, + { "s": 211, "t": "step_created", "k": "step", "r": 70, "n": "releaseStep" }, + { "s": 212, "t": "step_started", "k": "step", "r": 70, "n": "releaseStep" }, + { "s": 213, "t": "step_created", "k": "step", "r": 69, "n": "releaseStep" }, + { "s": 214, "t": "step_started", "k": "step", "r": 69, "n": "releaseStep" }, + { "s": 215, "t": "step_completed", "k": "step", "r": 70, "n": "releaseStep" }, + { + "s": 216, + "t": "step_completed", + "k": "step", + "r": 66, + "n": "finalizeStep" + }, + { "s": 217, "t": "step_completed", "k": "step", "r": 69, "n": "releaseStep" }, + { "s": 218, "t": "step_started", "k": "step", "r": 68, "n": "releaseStep" }, + { + "s": 219, + "t": "step_completed", + "k": "step", + "r": 67, + "n": "finalizeStep" + }, + { "s": 220, "t": "step_created", "k": "step", "r": 71, "n": "releaseStep" }, + { "s": 221, "t": "step_created", "k": "step", "r": 72, "n": "releaseStep" }, + { "s": 222, "t": "step_completed", "k": "step", "r": 68, "n": "releaseStep" }, + { "s": 223, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 224, "t": "step_started", "k": "step", "r": 71, "n": "releaseStep" }, + { "s": 225, "t": "step_completed", "k": "step", "r": 71, "n": "releaseStep" }, + { "s": 226, "t": "step_created", "k": "step", "r": 73, "n": "releaseStep" }, + { "s": 227, "t": "step_started", "k": "step", "r": 73, "n": "releaseStep" }, + { "s": 228, "t": "step_created", "k": "step", "r": 74, "n": "releaseStep" }, + { "s": 229, "t": "step_started", "k": "step", "r": 74, "n": "releaseStep" }, + { "s": 230, "t": "step_completed", "k": "step", "r": 74, "n": "releaseStep" }, + { "s": 231, "t": "step_completed", "k": "step", "r": 73, "n": "releaseStep" }, + { "s": 232, "t": "step_started", "k": "step", "r": 72, "n": "releaseStep" }, + { "s": 233, "t": "step_completed", "k": "step", "r": 72, "n": "releaseStep" }, + { "s": 234, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 235, "t": "step_created", "k": "step", "r": 79, "n": "reconcileStep" }, + { "s": 236, "t": "step_created", "k": "step", "r": 78, "n": "reconcileStep" }, + { "s": 237, "t": "step_created", "k": "step", "r": 81, "n": "reconcileStep" }, + { "s": 238, "t": "step_created", "k": "step", "r": 80, "n": "reconcileStep" }, + { "s": 239, "t": "step_created", "k": "step", "r": 75, "n": "reconcileStep" }, + { "s": 240, "t": "step_started", "k": "step", "r": 75, "n": "reconcileStep" }, + { "s": 241, "t": "step_started", "k": "step", "r": 79, "n": "reconcileStep" }, + { "s": 242, "t": "step_created", "k": "step", "r": 76, "n": "reconcileStep" }, + { "s": 243, "t": "step_started", "k": "step", "r": 76, "n": "reconcileStep" }, + { "s": 244, "t": "step_started", "k": "step", "r": 78, "n": "reconcileStep" }, + { + "s": 245, + "t": "step_completed", + "k": "step", + "r": 79, + "n": "reconcileStep" + }, + { + "s": 246, + "t": "step_completed", + "k": "step", + "r": 76, + "n": "reconcileStep" + }, + { "s": 247, "t": "step_started", "k": "step", "r": 80, "n": "reconcileStep" }, + { + "s": 248, + "t": "step_completed", + "k": "step", + "r": 78, + "n": "reconcileStep" + }, + { + "s": 249, + "t": "step_completed", + "k": "step", + "r": 75, + "n": "reconcileStep" + }, + { + "s": 250, + "t": "step_completed", + "k": "step", + "r": 80, + "n": "reconcileStep" + }, + { "s": 251, "t": "step_started", "k": "step", "r": 81, "n": "reconcileStep" }, + { "s": 252, "t": "step_created", "k": "step", "r": 77, "n": "reconcileStep" }, + { "s": 253, "t": "step_started", "k": "step", "r": 77, "n": "reconcileStep" }, + { + "s": 254, + "t": "step_completed", + "k": "step", + "r": 81, + "n": "reconcileStep" + }, + { "s": 255, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { + "s": 256, + "t": "step_completed", + "k": "step", + "r": 77, + "n": "reconcileStep" + }, + { "s": 257, "t": "wait_created", "k": "wait", "r": 82, "n": "" }, + { "s": 258, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 259, "t": "wait_completed", "k": "wait", "r": 82, "n": "" }, + { "s": 260, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 261, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 262, "t": "wait_created", "k": "wait", "r": 90, "n": "" }, + { "s": 263, "t": "wait_created", "k": "wait", "r": 84, "n": "" }, + { "s": 264, "t": "step_created", "k": "step", "r": 97, "n": "settleStep" }, + { "s": 265, "t": "wait_created", "k": "wait", "r": 96, "n": "" }, + { "s": 266, "t": "step_created", "k": "step", "r": 95, "n": "settleStep" }, + { "s": 267, "t": "step_created", "k": "step", "r": 89, "n": "settleStep" }, + { "s": 268, "t": "wait_created", "k": "wait", "r": 86, "n": "" }, + { "s": 269, "t": "wait_created", "k": "wait", "r": 88, "n": "" }, + { "s": 270, "t": "wait_created", "k": "wait", "r": 98, "n": "" }, + { "s": 271, "t": "step_created", "k": "step", "r": 93, "n": "settleStep" }, + { "s": 272, "t": "wait_created", "k": "wait", "r": 94, "n": "" }, + { "s": 273, "t": "wait_created", "k": "wait", "r": 92, "n": "" }, + { "s": 274, "t": "step_created", "k": "step", "r": 91, "n": "settleStep" }, + { "s": 275, "t": "step_created", "k": "step", "r": 85, "n": "settleStep" }, + { "s": 276, "t": "step_started", "k": "step", "r": 85, "n": "settleStep" }, + { "s": 277, "t": "step_started", "k": "step", "r": 89, "n": "settleStep" }, + { "s": 278, "t": "step_started", "k": "step", "r": 97, "n": "settleStep" }, + { "s": 279, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 280, "t": "step_started", "k": "step", "r": 95, "n": "settleStep" }, + { "s": 281, "t": "step_created", "k": "step", "r": 83, "n": "settleStep" }, + { "s": 282, "t": "step_started", "k": "step", "r": 83, "n": "settleStep" }, + { "s": 283, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 284, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 285, "t": "step_created", "k": "step", "r": 87, "n": "settleStep" }, + { "s": 286, "t": "step_started", "k": "step", "r": 87, "n": "settleStep" }, + { "s": 287, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 288, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 289, "t": "step_started", "k": "step", "r": 93, "n": "settleStep" }, + { "s": 290, "t": "step_started", "k": "step", "r": 91, "n": "settleStep" }, + { "s": 291, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 292, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 293, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 294, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 295, "t": "step_completed", "k": "step", "r": 85, "n": "settleStep" }, + { "s": 296, "t": "step_completed", "k": "step", "r": 97, "n": "settleStep" }, + { "s": 297, "t": "step_completed", "k": "step", "r": 83, "n": "settleStep" }, + { "s": 298, "t": "step_completed", "k": "step", "r": 89, "n": "settleStep" }, + { "s": 299, "t": "wait_completed", "k": "wait", "r": 90, "n": "" }, + { "s": 300, "t": "step_completed", "k": "step", "r": 93, "n": "settleStep" }, + { "s": 301, "t": "step_completed", "k": "step", "r": 87, "n": "settleStep" }, + { "s": 302, "t": "wait_completed", "k": "wait", "r": 84, "n": "" }, + { "s": 303, "t": "wait_completed", "k": "wait", "r": 96, "n": "" }, + { "s": 304, "t": "step_completed", "k": "step", "r": 91, "n": "settleStep" }, + { "s": 305, "t": "wait_completed", "k": "wait", "r": 86, "n": "" }, + { "s": 306, "t": "wait_completed", "k": "wait", "r": 88, "n": "" }, + { "s": 307, "t": "wait_completed", "k": "wait", "r": 98, "n": "" }, + { "s": 308, "t": "wait_completed", "k": "wait", "r": 94, "n": "" }, + { "s": 309, "t": "step_completed", "k": "step", "r": 95, "n": "settleStep" }, + { "s": 310, "t": "wait_completed", "k": "wait", "r": 92, "n": "" }, + { "s": 311, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 312, "t": "step_created", "k": "step", "r": 102, "n": "finalizeStep" }, + { "s": 313, "t": "step_created", "k": "step", "r": 101, "n": "finalizeStep" }, + { "s": 314, "t": "step_started", "k": "step", "r": 101, "n": "finalizeStep" }, + { "s": 315, "t": "step_created", "k": "step", "r": 100, "n": "finalizeStep" }, + { "s": 316, "t": "step_started", "k": "step", "r": 100, "n": "finalizeStep" }, + { "s": 317, "t": "step_created", "k": "step", "r": 99, "n": "finalizeStep" }, + { "s": 318, "t": "step_started", "k": "step", "r": 99, "n": "finalizeStep" }, + { "s": 319, "t": "step_started", "k": "step", "r": 102, "n": "finalizeStep" }, + { + "s": 320, + "t": "step_completed", + "k": "step", + "r": 100, + "n": "finalizeStep" + }, + { + "s": 321, + "t": "step_completed", + "k": "step", + "r": 99, + "n": "finalizeStep" + }, + { + "s": 322, + "t": "step_completed", + "k": "step", + "r": 102, + "n": "finalizeStep" + }, + { + "s": 323, + "t": "step_completed", + "k": "step", + "r": 101, + "n": "finalizeStep" + }, + { "s": 324, "t": "step_created", "k": "step", "r": 107, "n": "releaseStep" }, + { "s": 325, "t": "step_created", "k": "step", "r": 109, "n": "releaseStep" }, + { "s": 326, "t": "step_created", "k": "step", "r": 108, "n": "releaseStep" }, + { "s": 327, "t": "step_created", "k": "step", "r": 106, "n": "finalizeStep" }, + { "s": 328, "t": "step_created", "k": "step", "r": 110, "n": "releaseStep" }, + { "s": 329, "t": "step_created", "k": "step", "r": 104, "n": "finalizeStep" }, + { "s": 330, "t": "step_started", "k": "step", "r": 104, "n": "finalizeStep" }, + { "s": 331, "t": "step_created", "k": "step", "r": 103, "n": "finalizeStep" }, + { "s": 332, "t": "step_started", "k": "step", "r": 103, "n": "finalizeStep" }, + { + "s": 333, + "t": "step_completed", + "k": "step", + "r": 103, + "n": "finalizeStep" + }, + { "s": 334, "t": "step_created", "k": "step", "r": 105, "n": "recoverStep" }, + { "s": 335, "t": "step_started", "k": "step", "r": 105, "n": "recoverStep" }, + { "s": 336, "t": "step_started", "k": "step", "r": 109, "n": "releaseStep" }, + { + "s": 337, + "t": "step_completed", + "k": "step", + "r": 105, + "n": "recoverStep" + }, + { + "s": 338, + "t": "step_completed", + "k": "step", + "r": 109, + "n": "releaseStep" + }, + { + "s": 339, + "t": "step_completed", + "k": "step", + "r": 104, + "n": "finalizeStep" + }, + { "s": 340, "t": "step_started", "k": "step", "r": 106, "n": "finalizeStep" }, + { "s": 341, "t": "step_started", "k": "step", "r": 110, "n": "releaseStep" }, + { + "s": 342, + "t": "step_completed", + "k": "step", + "r": 106, + "n": "finalizeStep" + }, + { "s": 343, "t": "step_started", "k": "step", "r": 107, "n": "releaseStep" }, + { "s": 344, "t": "step_started", "k": "step", "r": 108, "n": "releaseStep" }, + { + "s": 345, + "t": "step_completed", + "k": "step", + "r": 107, + "n": "releaseStep" + }, + { + "s": 346, + "t": "step_completed", + "k": "step", + "r": 110, + "n": "releaseStep" + }, + { + "s": 347, + "t": "step_completed", + "k": "step", + "r": 108, + "n": "releaseStep" + }, + { "s": 348, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 349, "t": "step_created", "k": "step", "r": 111, "n": "releaseStep" }, + { "s": 350, "t": "step_started", "k": "step", "r": 111, "n": "releaseStep" }, + { "s": 351, "t": "step_created", "k": "step", "r": 112, "n": "finalizeStep" }, + { "s": 352, "t": "step_started", "k": "step", "r": 112, "n": "finalizeStep" }, + { "s": 353, "t": "step_created", "k": "step", "r": 113, "n": "releaseStep" }, + { "s": 354, "t": "step_started", "k": "step", "r": 113, "n": "releaseStep" }, + { + "s": 355, + "t": "step_completed", + "k": "step", + "r": 111, + "n": "releaseStep" + }, + { + "s": 356, + "t": "step_completed", + "k": "step", + "r": 112, + "n": "finalizeStep" + }, + { + "s": 357, + "t": "step_completed", + "k": "step", + "r": 113, + "n": "releaseStep" + }, + { "s": 358, "t": "step_created", "k": "step", "r": 114, "n": "releaseStep" }, + { "s": 359, "t": "step_created", "k": "step", "r": 115, "n": "releaseStep" }, + { "s": 360, "t": "step_started", "k": "step", "r": 115, "n": "releaseStep" }, + { + "s": 361, + "t": "step_completed", + "k": "step", + "r": 115, + "n": "releaseStep" + }, + { "s": 362, "t": "step_started", "k": "step", "r": 114, "n": "releaseStep" }, + { "s": 363, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { + "s": 364, + "t": "step_completed", + "k": "step", + "r": 114, + "n": "releaseStep" + }, + { "s": 365, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { + "s": 366, + "t": "step_created", + "k": "step", + "r": 118, + "n": "reconcileStep" + }, + { + "s": 367, + "t": "step_started", + "k": "step", + "r": 118, + "n": "reconcileStep" + }, + { + "s": 368, + "t": "step_created", + "k": "step", + "r": 117, + "n": "reconcileStep" + }, + { + "s": 369, + "t": "step_started", + "k": "step", + "r": 117, + "n": "reconcileStep" + }, + { + "s": 370, + "t": "step_created", + "k": "step", + "r": 116, + "n": "reconcileStep" + }, + { + "s": 371, + "t": "step_started", + "k": "step", + "r": 116, + "n": "reconcileStep" + }, + { + "s": 372, + "t": "step_completed", + "k": "step", + "r": 118, + "n": "reconcileStep" + }, + { + "s": 373, + "t": "step_completed", + "k": "step", + "r": 116, + "n": "reconcileStep" + }, + { + "s": 374, + "t": "step_completed", + "k": "step", + "r": 117, + "n": "reconcileStep" + }, + { "s": 375, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 376, "t": "wait_created", "k": "wait", "r": 119, "n": "" }, + { "s": 377, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 378, "t": "wait_completed", "k": "wait", "r": 119, "n": "" }, + { "s": 379, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 380, "t": "wait_created", "k": "wait", "r": 133, "n": "" }, + { "s": 381, "t": "step_created", "k": "step", "r": 128, "n": "settleStep" }, + { "s": 382, "t": "wait_created", "k": "wait", "r": 123, "n": "" }, + { "s": 383, "t": "wait_created", "k": "wait", "r": 131, "n": "" }, + { "s": 384, "t": "wait_created", "k": "wait", "r": 125, "n": "" }, + { "s": 385, "t": "step_created", "k": "step", "r": 130, "n": "settleStep" }, + { "s": 386, "t": "step_created", "k": "step", "r": 134, "n": "settleStep" }, + { "s": 387, "t": "step_created", "k": "step", "r": 132, "n": "settleStep" }, + { "s": 388, "t": "wait_created", "k": "wait", "r": 127, "n": "" }, + { "s": 389, "t": "wait_created", "k": "wait", "r": 129, "n": "" }, + { "s": 390, "t": "wait_created", "k": "wait", "r": 121, "n": "" }, + { "s": 391, "t": "step_created", "k": "step", "r": 126, "n": "settleStep" }, + { "s": 392, "t": "wait_created", "k": "wait", "r": 135, "n": "" }, + { "s": 393, "t": "step_created", "k": "step", "r": 124, "n": "settleStep" }, + { "s": 394, "t": "step_started", "k": "step", "r": 124, "n": "settleStep" }, + { "s": 395, "t": "step_created", "k": "step", "r": 120, "n": "settleStep" }, + { "s": 396, "t": "step_started", "k": "step", "r": 120, "n": "settleStep" }, + { "s": 397, "t": "step_started", "k": "step", "r": 130, "n": "settleStep" }, + { "s": 398, "t": "step_created", "k": "step", "r": 122, "n": "settleStep" }, + { "s": 399, "t": "step_started", "k": "step", "r": 122, "n": "settleStep" }, + { "s": 400, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 401, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 402, "t": "step_started", "k": "step", "r": 132, "n": "settleStep" }, + { "s": 403, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 404, "t": "step_started", "k": "step", "r": 134, "n": "settleStep" }, + { "s": 405, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 406, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 407, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 408, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 409, "t": "step_started", "k": "step", "r": 126, "n": "settleStep" }, + { "s": 410, "t": "step_started", "k": "step", "r": 128, "n": "settleStep" }, + { "s": 411, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 412, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 413, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 414, "t": "wait_completed", "k": "wait", "r": 133, "n": "" }, + { "s": 415, "t": "wait_completed", "k": "wait", "r": 123, "n": "" }, + { "s": 416, "t": "wait_completed", "k": "wait", "r": 131, "n": "" }, + { "s": 417, "t": "wait_completed", "k": "wait", "r": 125, "n": "" }, + { "s": 418, "t": "step_completed", "k": "step", "r": 120, "n": "settleStep" }, + { "s": 419, "t": "step_completed", "k": "step", "r": 124, "n": "settleStep" }, + { "s": 420, "t": "step_completed", "k": "step", "r": 130, "n": "settleStep" }, + { "s": 421, "t": "wait_completed", "k": "wait", "r": 127, "n": "" }, + { "s": 422, "t": "step_completed", "k": "step", "r": 122, "n": "settleStep" }, + { "s": 423, "t": "step_completed", "k": "step", "r": 134, "n": "settleStep" }, + { "s": 424, "t": "wait_completed", "k": "wait", "r": 129, "n": "" }, + { "s": 425, "t": "step_completed", "k": "step", "r": 132, "n": "settleStep" }, + { "s": 426, "t": "wait_completed", "k": "wait", "r": 121, "n": "" }, + { "s": 427, "t": "wait_completed", "k": "wait", "r": 135, "n": "" }, + { "s": 428, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 429, "t": "step_created", "k": "step", "r": 139, "n": "recoverStep" }, + { "s": 430, "t": "step_completed", "k": "step", "r": 126, "n": "settleStep" }, + { "s": 431, "t": "step_created", "k": "step", "r": 141, "n": "recoverStep" }, + { "s": 432, "t": "step_created", "k": "step", "r": 143, "n": "recoverStep" }, + { "s": 433, "t": "step_completed", "k": "step", "r": 128, "n": "settleStep" }, + { "s": 434, "t": "step_created", "k": "step", "r": 142, "n": "finalizeStep" }, + { "s": 435, "t": "step_created", "k": "step", "r": 140, "n": "finalizeStep" }, + { "s": 436, "t": "step_started", "k": "step", "r": 141, "n": "recoverStep" }, + { "s": 437, "t": "step_started", "k": "step", "r": 142, "n": "finalizeStep" }, + { "s": 438, "t": "step_created", "k": "step", "r": 138, "n": "recoverStep" }, + { "s": 439, "t": "step_started", "k": "step", "r": 138, "n": "recoverStep" }, + { "s": 440, "t": "step_started", "k": "step", "r": 143, "n": "recoverStep" }, + { "s": 441, "t": "step_created", "k": "step", "r": 137, "n": "recoverStep" }, + { "s": 442, "t": "step_started", "k": "step", "r": 137, "n": "recoverStep" }, + { + "s": 443, + "t": "step_completed", + "k": "step", + "r": 141, + "n": "recoverStep" + }, + { + "s": 444, + "t": "step_completed", + "k": "step", + "r": 143, + "n": "recoverStep" + }, + { + "s": 445, + "t": "step_completed", + "k": "step", + "r": 138, + "n": "recoverStep" + }, + { "s": 446, "t": "step_started", "k": "step", "r": 140, "n": "finalizeStep" }, + { + "s": 447, + "t": "step_completed", + "k": "step", + "r": 137, + "n": "recoverStep" + }, + { "s": 448, "t": "step_started", "k": "step", "r": 139, "n": "recoverStep" }, + { + "s": 449, + "t": "step_completed", + "k": "step", + "r": 140, + "n": "finalizeStep" + }, + { + "s": 450, + "t": "step_completed", + "k": "step", + "r": 142, + "n": "finalizeStep" + }, + { + "s": 451, + "t": "step_completed", + "k": "step", + "r": 139, + "n": "recoverStep" + }, + { "s": 452, "t": "step_created", "k": "step", "r": 136, "n": "recoverStep" }, + { "s": 453, "t": "step_started", "k": "step", "r": 136, "n": "recoverStep" }, + { "s": 454, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { + "s": 455, + "t": "step_completed", + "k": "step", + "r": 136, + "n": "recoverStep" + }, + { "s": 456, "t": "step_created", "k": "step", "r": 147, "n": "finalizeStep" }, + { "s": 457, "t": "step_created", "k": "step", "r": 150, "n": "finalizeStep" }, + { "s": 458, "t": "step_created", "k": "step", "r": 149, "n": "releaseStep" }, + { "s": 459, "t": "step_created", "k": "step", "r": 148, "n": "releaseStep" }, + { "s": 460, "t": "step_created", "k": "step", "r": 146, "n": "finalizeStep" }, + { "s": 461, "t": "step_created", "k": "step", "r": 145, "n": "finalizeStep" }, + { "s": 462, "t": "step_started", "k": "step", "r": 145, "n": "finalizeStep" }, + { "s": 463, "t": "step_created", "k": "step", "r": 144, "n": "finalizeStep" }, + { "s": 464, "t": "step_started", "k": "step", "r": 144, "n": "finalizeStep" }, + { "s": 465, "t": "step_started", "k": "step", "r": 148, "n": "releaseStep" }, + { + "s": 466, + "t": "step_completed", + "k": "step", + "r": 144, + "n": "finalizeStep" + }, + { + "s": 467, + "t": "step_completed", + "k": "step", + "r": 145, + "n": "finalizeStep" + }, + { "s": 468, "t": "step_started", "k": "step", "r": 146, "n": "finalizeStep" }, + { "s": 469, "t": "step_started", "k": "step", "r": 147, "n": "finalizeStep" }, + { + "s": 470, + "t": "step_completed", + "k": "step", + "r": 146, + "n": "finalizeStep" + }, + { + "s": 471, + "t": "step_completed", + "k": "step", + "r": 148, + "n": "releaseStep" + }, + { + "s": 472, + "t": "step_completed", + "k": "step", + "r": 147, + "n": "finalizeStep" + }, + { "s": 473, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 474, "t": "step_created", "k": "step", "r": 151, "n": "finalizeStep" }, + { "s": 475, "t": "step_started", "k": "step", "r": 149, "n": "releaseStep" }, + { "s": 476, "t": "step_created", "k": "step", "r": 152, "n": "releaseStep" }, + { "s": 477, "t": "step_started", "k": "step", "r": 152, "n": "releaseStep" }, + { "s": 478, "t": "step_created", "k": "step", "r": 153, "n": "releaseStep" }, + { "s": 479, "t": "step_started", "k": "step", "r": 153, "n": "releaseStep" }, + { + "s": 480, + "t": "step_completed", + "k": "step", + "r": 152, + "n": "releaseStep" + }, + { + "s": 481, + "t": "step_completed", + "k": "step", + "r": 153, + "n": "releaseStep" + }, + { + "s": 482, + "t": "step_completed", + "k": "step", + "r": 149, + "n": "releaseStep" + }, + { "s": 483, "t": "step_started", "k": "step", "r": 151, "n": "finalizeStep" }, + { "s": 484, "t": "step_started", "k": "step", "r": 150, "n": "finalizeStep" }, + { + "s": 485, + "t": "step_completed", + "k": "step", + "r": 150, + "n": "finalizeStep" + }, + { + "s": 486, + "t": "step_completed", + "k": "step", + "r": 151, + "n": "finalizeStep" + }, + { "s": 487, "t": "step_created", "k": "step", "r": 154, "n": "releaseStep" }, + { "s": 488, "t": "step_started", "k": "step", "r": 154, "n": "releaseStep" }, + { "s": 489, "t": "step_created", "k": "step", "r": 155, "n": "releaseStep" }, + { "s": 490, "t": "step_started", "k": "step", "r": 155, "n": "releaseStep" }, + { + "s": 491, + "t": "step_completed", + "k": "step", + "r": 154, + "n": "releaseStep" + }, + { + "s": 492, + "t": "step_completed", + "k": "step", + "r": 155, + "n": "releaseStep" + }, + { "s": 493, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 494, "t": "step_created", "k": "step", "r": 157, "n": "releaseStep" }, + { "s": 495, "t": "step_started", "k": "step", "r": 157, "n": "releaseStep" }, + { "s": 496, "t": "step_created", "k": "step", "r": 156, "n": "releaseStep" }, + { "s": 497, "t": "step_started", "k": "step", "r": 156, "n": "releaseStep" }, + { + "s": 498, + "t": "step_completed", + "k": "step", + "r": 157, + "n": "releaseStep" + }, + { + "s": 499, + "t": "step_completed", + "k": "step", + "r": 156, + "n": "releaseStep" + }, + { + "s": 500, + "t": "step_created", + "k": "step", + "r": 162, + "n": "reconcileStep" + }, + { + "s": 501, + "t": "step_created", + "k": "step", + "r": 163, + "n": "reconcileStep" + }, + { + "s": 502, + "t": "step_created", + "k": "step", + "r": 161, + "n": "reconcileStep" + }, + { + "s": 503, + "t": "step_created", + "k": "step", + "r": 165, + "n": "reconcileStep" + }, + { + "s": 504, + "t": "step_created", + "k": "step", + "r": 164, + "n": "reconcileStep" + }, + { + "s": 505, + "t": "step_created", + "k": "step", + "r": 158, + "n": "reconcileStep" + }, + { + "s": 506, + "t": "step_started", + "k": "step", + "r": 158, + "n": "reconcileStep" + }, + { + "s": 507, + "t": "step_created", + "k": "step", + "r": 160, + "n": "reconcileStep" + }, + { + "s": 508, + "t": "step_started", + "k": "step", + "r": 160, + "n": "reconcileStep" + }, + { + "s": 509, + "t": "step_started", + "k": "step", + "r": 161, + "n": "reconcileStep" + }, + { + "s": 510, + "t": "step_created", + "k": "step", + "r": 159, + "n": "reconcileStep" + }, + { + "s": 511, + "t": "step_started", + "k": "step", + "r": 159, + "n": "reconcileStep" + }, + { + "s": 512, + "t": "step_started", + "k": "step", + "r": 162, + "n": "reconcileStep" + }, + { + "s": 513, + "t": "step_completed", + "k": "step", + "r": 158, + "n": "reconcileStep" + }, + { + "s": 514, + "t": "step_started", + "k": "step", + "r": 165, + "n": "reconcileStep" + }, + { + "s": 515, + "t": "step_completed", + "k": "step", + "r": 159, + "n": "reconcileStep" + }, + { + "s": 516, + "t": "step_completed", + "k": "step", + "r": 165, + "n": "reconcileStep" + }, + { + "s": 517, + "t": "step_completed", + "k": "step", + "r": 160, + "n": "reconcileStep" + }, + { + "s": 518, + "t": "step_completed", + "k": "step", + "r": 162, + "n": "reconcileStep" + }, + { + "s": 519, + "t": "step_started", + "k": "step", + "r": 163, + "n": "reconcileStep" + }, + { + "s": 520, + "t": "step_completed", + "k": "step", + "r": 161, + "n": "reconcileStep" + }, + { + "s": 521, + "t": "step_completed", + "k": "step", + "r": 163, + "n": "reconcileStep" + }, + { + "s": 522, + "t": "step_started", + "k": "step", + "r": 164, + "n": "reconcileStep" + }, + { "s": 523, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { + "s": 524, + "t": "step_completed", + "k": "step", + "r": 164, + "n": "reconcileStep" + }, + { "s": 525, "t": "wait_created", "k": "wait", "r": 166, "n": "" }, + { "s": 526, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 527, "t": "wait_completed", "k": "wait", "r": 166, "n": "" }, + { "s": 528, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 529, "t": "wait_created", "k": "wait", "r": 172, "n": "" }, + { "s": 530, "t": "step_created", "k": "step", "r": 173, "n": "settleStep" }, + { "s": 531, "t": "step_created", "k": "step", "r": 181, "n": "settleStep" }, + { "s": 532, "t": "wait_created", "k": "wait", "r": 168, "n": "" }, + { "s": 533, "t": "wait_created", "k": "wait", "r": 180, "n": "" }, + { "s": 534, "t": "step_created", "k": "step", "r": 175, "n": "settleStep" }, + { "s": 535, "t": "wait_created", "k": "wait", "r": 178, "n": "" }, + { "s": 536, "t": "step_created", "k": "step", "r": 177, "n": "settleStep" }, + { "s": 537, "t": "wait_created", "k": "wait", "r": 170, "n": "" }, + { "s": 538, "t": "wait_created", "k": "wait", "r": 182, "n": "" }, + { "s": 539, "t": "wait_created", "k": "wait", "r": 174, "n": "" }, + { "s": 540, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 541, "t": "wait_created", "k": "wait", "r": 176, "n": "" }, + { "s": 542, "t": "step_created", "k": "step", "r": 179, "n": "settleStep" }, + { "s": 543, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 544, "t": "step_created", "k": "step", "r": 171, "n": "settleStep" }, + { "s": 545, "t": "step_started", "k": "step", "r": 171, "n": "settleStep" }, + { "s": 546, "t": "step_created", "k": "step", "r": 169, "n": "settleStep" }, + { "s": 547, "t": "step_started", "k": "step", "r": 169, "n": "settleStep" }, + { "s": 548, "t": "step_started", "k": "step", "r": 177, "n": "settleStep" }, + { "s": 549, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 550, "t": "step_started", "k": "step", "r": 179, "n": "settleStep" }, + { "s": 551, "t": "step_started", "k": "step", "r": 173, "n": "settleStep" }, + { "s": 552, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 553, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 554, "t": "step_created", "k": "step", "r": 167, "n": "settleStep" }, + { "s": 555, "t": "step_started", "k": "step", "r": 167, "n": "settleStep" }, + { "s": 556, "t": "step_started", "k": "step", "r": 175, "n": "settleStep" }, + { "s": 557, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 558, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 559, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 560, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 561, "t": "wait_completed", "k": "wait", "r": 172, "n": "" }, + { "s": 562, "t": "step_started", "k": "step", "r": 181, "n": "settleStep" }, + { "s": 563, "t": "wait_completed", "k": "wait", "r": 168, "n": "" }, + { "s": 564, "t": "attr_set", "k": "run", "r": -1, "n": "" }, + { "s": 565, "t": "wait_completed", "k": "wait", "r": 180, "n": "" }, + { "s": 566, "t": "wait_completed", "k": "wait", "r": 178, "n": "" }, + { "s": 567, "t": "wait_completed", "k": "wait", "r": 170, "n": "" }, + { "s": 568, "t": "wait_completed", "k": "wait", "r": 182, "n": "" }, + { "s": 569, "t": "wait_completed", "k": "wait", "r": 174, "n": "" }, + { "s": 570, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 571, "t": "step_completed", "k": "step", "r": 171, "n": "settleStep" }, + { "s": 572, "t": "step_created", "k": "step", "r": 189, "n": "recoverStep" }, + { "s": 573, "t": "step_completed", "k": "step", "r": 179, "n": "settleStep" }, + { "s": 574, "t": "step_completed", "k": "step", "r": 169, "n": "settleStep" }, + { "s": 575, "t": "step_completed", "k": "step", "r": 173, "n": "settleStep" }, + { "s": 576, "t": "step_created", "k": "step", "r": 186, "n": "recoverStep" }, + { "s": 577, "t": "step_completed", "k": "step", "r": 175, "n": "settleStep" }, + { "s": 578, "t": "step_created", "k": "step", "r": 188, "n": "recoverStep" }, + { "s": 579, "t": "step_completed", "k": "step", "r": 177, "n": "settleStep" }, + { "s": 580, "t": "step_created", "k": "step", "r": 187, "n": "recoverStep" }, + { "s": 581, "t": "step_completed", "k": "step", "r": 181, "n": "settleStep" }, + { "s": 582, "t": "step_completed", "k": "step", "r": 167, "n": "settleStep" }, + { "s": 583, "t": "wait_completed", "k": "wait", "r": 176, "n": "" }, + { "s": 584, "t": "step_created", "k": "step", "r": 184, "n": "recoverStep" }, + { "s": 585, "t": "step_started", "k": "step", "r": 184, "n": "recoverStep" }, + { "s": 586, "t": "step_created", "k": "step", "r": 185, "n": "recoverStep" }, + { "s": 587, "t": "step_started", "k": "step", "r": 185, "n": "recoverStep" }, + { + "s": 588, + "t": "step_completed", + "k": "step", + "r": 184, + "n": "recoverStep" + }, + { "s": 589, "t": "step_started", "k": "step", "r": 186, "n": "recoverStep" }, + { "s": 590, "t": "step_started", "k": "step", "r": 187, "n": "recoverStep" }, + { + "s": 591, + "t": "step_completed", + "k": "step", + "r": 185, + "n": "recoverStep" + }, + { "s": 592, "t": "step_created", "k": "step", "r": 183, "n": "recoverStep" }, + { "s": 593, "t": "step_started", "k": "step", "r": 183, "n": "recoverStep" }, + { + "s": 594, + "t": "step_completed", + "k": "step", + "r": 186, + "n": "recoverStep" + }, + { + "s": 595, + "t": "step_completed", + "k": "step", + "r": 187, + "n": "recoverStep" + }, + { + "s": 596, + "t": "step_completed", + "k": "step", + "r": 183, + "n": "recoverStep" + }, + { "s": 597, "t": "step_started", "k": "step", "r": 188, "n": "recoverStep" }, + { + "s": 598, + "t": "step_completed", + "k": "step", + "r": 188, + "n": "recoverStep" + }, + { "s": 599, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 600, "t": "step_created", "k": "step", "r": 190, "n": "finalizeStep" }, + { "s": 601, "t": "step_started", "k": "step", "r": 190, "n": "finalizeStep" }, + { + "s": 602, + "t": "step_completed", + "k": "step", + "r": 190, + "n": "finalizeStep" + }, + { "s": 603, "t": "step_started", "k": "step", "r": 189, "n": "recoverStep" }, + { + "s": 604, + "t": "step_completed", + "k": "step", + "r": 189, + "n": "recoverStep" + }, + { "s": 605, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 606, "t": "step_created", "k": "step", "r": 195, "n": "finalizeStep" }, + { "s": 607, "t": "step_created", "k": "step", "r": 196, "n": "finalizeStep" }, + { "s": 608, "t": "step_created", "k": "step", "r": 194, "n": "finalizeStep" }, + { "s": 609, "t": "step_created", "k": "step", "r": 192, "n": "finalizeStep" }, + { "s": 610, "t": "step_started", "k": "step", "r": 192, "n": "finalizeStep" }, + { "s": 611, "t": "step_started", "k": "step", "r": 194, "n": "finalizeStep" }, + { + "s": 612, + "t": "step_completed", + "k": "step", + "r": 192, + "n": "finalizeStep" + }, + { "s": 613, "t": "step_started", "k": "step", "r": 195, "n": "finalizeStep" }, + { "s": 614, "t": "step_started", "k": "step", "r": 196, "n": "finalizeStep" }, + { "s": 615, "t": "step_created", "k": "step", "r": 191, "n": "finalizeStep" }, + { "s": 616, "t": "step_started", "k": "step", "r": 191, "n": "finalizeStep" }, + { "s": 617, "t": "step_created", "k": "step", "r": 193, "n": "finalizeStep" }, + { "s": 618, "t": "step_started", "k": "step", "r": 193, "n": "finalizeStep" }, + { + "s": 619, + "t": "step_completed", + "k": "step", + "r": 191, + "n": "finalizeStep" + }, + { + "s": 620, + "t": "step_completed", + "k": "step", + "r": 196, + "n": "finalizeStep" + }, + { + "s": 621, + "t": "step_completed", + "k": "step", + "r": 194, + "n": "finalizeStep" + }, + { + "s": 622, + "t": "step_completed", + "k": "step", + "r": 193, + "n": "finalizeStep" + }, + { + "s": 623, + "t": "step_completed", + "k": "step", + "r": 195, + "n": "finalizeStep" + }, + { "s": 624, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 625, "t": "step_created", "k": "step", "r": 200, "n": "releaseStep" }, + { "s": 626, "t": "step_created", "k": "step", "r": 202, "n": "releaseStep" }, + { "s": 627, "t": "step_created", "k": "step", "r": 204, "n": "releaseStep" }, + { "s": 628, "t": "step_created", "k": "step", "r": 201, "n": "releaseStep" }, + { "s": 629, "t": "step_created", "k": "step", "r": 203, "n": "releaseStep" }, + { "s": 630, "t": "step_created", "k": "step", "r": 198, "n": "finalizeStep" }, + { "s": 631, "t": "step_created", "k": "step", "r": 199, "n": "finalizeStep" }, + { "s": 632, "t": "step_started", "k": "step", "r": 199, "n": "finalizeStep" }, + { "s": 633, "t": "step_created", "k": "step", "r": 197, "n": "releaseStep" }, + { "s": 634, "t": "step_started", "k": "step", "r": 197, "n": "releaseStep" }, + { "s": 635, "t": "step_started", "k": "step", "r": 204, "n": "releaseStep" }, + { + "s": 636, + "t": "step_completed", + "k": "step", + "r": 199, + "n": "finalizeStep" + }, + { + "s": 637, + "t": "step_completed", + "k": "step", + "r": 197, + "n": "releaseStep" + }, + { "s": 638, "t": "step_started", "k": "step", "r": 198, "n": "finalizeStep" }, + { + "s": 639, + "t": "step_completed", + "k": "step", + "r": 204, + "n": "releaseStep" + }, + { + "s": 640, + "t": "step_completed", + "k": "step", + "r": 198, + "n": "finalizeStep" + }, + { "s": 641, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 642, "t": "step_started", "k": "step", "r": 201, "n": "releaseStep" }, + { "s": 643, "t": "step_started", "k": "step", "r": 200, "n": "releaseStep" }, + { "s": 644, "t": "step_started", "k": "step", "r": 203, "n": "releaseStep" }, + { + "s": 645, + "t": "step_completed", + "k": "step", + "r": 201, + "n": "releaseStep" + }, + { + "s": 646, + "t": "step_completed", + "k": "step", + "r": 203, + "n": "releaseStep" + }, + { + "s": 647, + "t": "step_completed", + "k": "step", + "r": 200, + "n": "releaseStep" + }, + { "s": 648, "t": "step_started", "k": "step", "r": 202, "n": "releaseStep" }, + { + "s": 649, + "t": "step_completed", + "k": "step", + "r": 202, + "n": "releaseStep" + }, + { "s": 650, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 651, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 652, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 653, "t": "hook_received", "k": "hook", "r": 0, "n": "" }, + { "s": 654, "t": "hook_received", "k": "hook", "r": 0, "n": "" } +] diff --git a/packages/core/src/delivery-barrier-dispenser.test.ts b/packages/core/src/delivery-barrier-dispenser.test.ts new file mode 100644 index 0000000000..bb714f1d73 --- /dev/null +++ b/packages/core/src/delivery-barrier-dispenser.test.ts @@ -0,0 +1,323 @@ +/** + * Unit coverage for the barrier safety-net dispenser (`ensureBarrierSafetyNet` + * in private.ts) — the pieces of vercel/workflow#3554 that previously only + * end-to-end storm lanes exercised: + * + * 1. End-of-log suspension must not preempt deliveries parked behind an + * unclaimed buffered hook payload. The pass suspends only after the + * parked chain delivered and the woken branch made its follow-up draws + * (the regression surfaced as whole storm runs going dormant/`stuck`). + * 2. With SEVERAL parked segments, chains must wake in log order — the + * dispenser retires heads lowest-first and re-blocks while a woken chain + * drains, so the ULIDs the branches draw next are position-determined, + * not net-timing-determined. + * 3. The dispenser must survive a rejected `promiseQueue`: the registry now + * gates `isDeliveryIdle`, so a silently-dead dispenser would wedge the + * run rather than merely skip a cleanup. + */ +import { WorkflowRuntimeError } from '@workflow/errors'; +import { withResolvers } from '@workflow/utils'; +import type { Event } from '@workflow/world'; +import * as nanoid from 'nanoid'; +import { monotonicFactory } from 'ulid'; +import { describe, expect, it } from 'vitest'; +import { EventsConsumer } from './events-consumer.js'; +import { WorkflowSuspension } from './global.js'; +import { + isDeliveryIdle, + registerDeliveryBarrier, + type WorkflowOrchestratorContext, +} from './private.js'; +import { ReplayPayloadCache } from './replay-payload-cache.js'; +import { dehydrateStepReturnValue } from './serialization.js'; +import { createUseStep } from './step.js'; +import { createContext } from './vm/index.js'; +import { createCreateHook } from './workflow/hook.js'; +import { createSleep } from './workflow/sleep.js'; + +const CORR_IDS = [ + '01K11TFZ62YS0YYFDQ3E8B9YCV', + '01K11TFZ62YS0YYFDQ3E8B9YCW', + '01K11TFZ62YS0YYFDQ3E8B9YCX', + '01K11TFZ62YS0YYFDQ3E8B9YCY', + '01K11TFZ62YS0YYFDQ3E8B9YCZ', +]; + +function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext { + const context = createContext({ + seed: 'test', + fixedTimestamp: 1753481739458, + }); + const ulid = monotonicFactory(() => context.globalThis.Math.random()); + const workflowStartedAt = context.globalThis.Date.now(); + const promiseQueueHolder = { current: Promise.resolve() }; + const ctxRef: { current?: WorkflowOrchestratorContext } = {}; + const ctx: WorkflowOrchestratorContext = { + suspensionGeneration: 0, + runId: 'wrun_test', + encryptionKey: undefined, + replayPayloadCache: new ReplayPayloadCache(undefined), + globalThis: context.globalThis, + eventsConsumer: new EventsConsumer(events, { + isDeliveryIdle: () => + ctxRef.current ? isDeliveryIdle(ctxRef.current) : true, + onUnconsumedEvent: (event) => { + ctxRef.current?.onWorkflowError( + new WorkflowRuntimeError( + `Unconsumed event: eventType=${event.eventType}, correlationId=${event.correlationId}, eventId=${event.eventId}.` + ) + ); + }, + getPromiseQueue: () => promiseQueueHolder.current, + }), + invocationsQueue: new Map(), + generateUlid: () => ulid(workflowStartedAt), + generateNanoid: nanoid.customRandom(nanoid.urlAlphabet, 21, (size) => + new Uint8Array(size).map(() => 256 * context.globalThis.Math.random()) + ), + onWorkflowError: () => {}, + get promiseQueue() { + return promiseQueueHolder.current; + }, + set promiseQueue(value: Promise) { + promiseQueueHolder.current = value; + }, + pendingDeliveries: 0, + pendingDeliveryBarriers: new Map(), + }; + ctxRef.current = ctx; + return ctx; +} + +async function runWithDiscontinuation( + ctx: WorkflowOrchestratorContext, + workflowFn: () => Promise +): Promise<{ result?: any; error?: any }> { + const workflowDiscontinuation = withResolvers(); + ctx.onWorkflowError = workflowDiscontinuation.reject; + let result: any; + let error: any; + try { + result = await Promise.race([ + workflowFn(), + workflowDiscontinuation.promise, + ]); + } catch (err) { + error = err; + } + return { result, error }; +} + +function pendingStepNames(ctx: WorkflowOrchestratorContext): string[] { + return [...ctx.invocationsQueue.values()] + .filter((item) => item.type === 'step') + .map((item) => (item.type === 'step' ? item.stepName : '')); +} + +const resumeAtA = new Date('2026-07-27T12:00:05.000Z'); +const resumeAtB = new Date('2026-07-27T12:00:06.000Z'); + +describe('barrier safety-net dispenser', () => { + it('suspends only after deliveries parked behind an unclaimed payload have run', async () => { + const ops: Promise[] = []; + const payload = await dehydrateStepReturnValue( + { poke: 1 }, + 'wrun_test', + undefined, + ops + ); + // Draw order in the body: c0 = the never-read hook, c1 = the sleep. The + // poke payload is consumed before the wait completion, so the wait's + // delivery gates on the unclaimed payload's barrier, which only the + // dispenser retires. afterSleep must still be drawn (and become the + // suspension's pending step) BEFORE the pass is allowed to suspend. + const events: Event[] = [ + { + eventId: 'evnt_0', + runId: 'wrun_test', + eventType: 'hook_created', + correlationId: `hook_${CORR_IDS[0]}`, + eventData: { token: 'dispenser-token', isWebhook: false }, + createdAt: new Date(), + }, + { + eventId: 'evnt_1', + runId: 'wrun_test', + eventType: 'wait_created', + correlationId: `wait_${CORR_IDS[1]}`, + eventData: { resumeAt: resumeAtA }, + createdAt: new Date(), + }, + { + eventId: 'evnt_2', + runId: 'wrun_test', + eventType: 'hook_received', + correlationId: `hook_${CORR_IDS[0]}`, + eventData: { payload }, + createdAt: new Date(), + }, + { + eventId: 'evnt_3', + runId: 'wrun_test', + eventType: 'wait_completed', + correlationId: `wait_${CORR_IDS[1]}`, + eventData: { resumeAt: resumeAtA }, + createdAt: new Date(), + }, + ] as unknown as Event[]; + + const ctx = setupWorkflowContext(events); + const useStep = createUseStep(ctx); + const sleep = createSleep(ctx); + const createHook = createCreateHook(ctx); + const body = async () => { + const afterSleep = useStep('afterSleep'); + createHook({ token: 'dispenser-token' }); + await sleep(resumeAtA); + await afterSleep(); + }; + + const { error } = await runWithDiscontinuation(ctx, body); + expect(error).toBeDefined(); + if (!WorkflowSuspension.is(error)) { + throw error; + } + // The wait delivered (despite parking behind the unclaimed payload) and + // the branch ran to its next draw before the suspension was raised. + expect(pendingStepNames(ctx)).toEqual(['afterSleep']); + }); + + it('wakes chains parked behind SEVERAL unclaimed payloads in log order', async () => { + const ops: Promise[] = []; + const payload = await dehydrateStepReturnValue( + { poke: 1 }, + 'wrun_test', + undefined, + ops + ); + // Two parked segments: waitA parks behind the first payload, waitB + // behind both. The branches draw afterA / afterB when woken, and the + // ULIDs they draw are position-determined only if A wakes before B — + // lowest-first retirement with re-blocking. Under the per-barrier polls + // this order was scheduling noise. + const events: Event[] = [ + { + eventId: 'evnt_0', + runId: 'wrun_test', + eventType: 'hook_created', + correlationId: `hook_${CORR_IDS[0]}`, + eventData: { token: 'dispenser-token', isWebhook: false }, + createdAt: new Date(), + }, + { + eventId: 'evnt_1', + runId: 'wrun_test', + eventType: 'wait_created', + correlationId: `wait_${CORR_IDS[1]}`, + eventData: { resumeAt: resumeAtA }, + createdAt: new Date(), + }, + { + eventId: 'evnt_2', + runId: 'wrun_test', + eventType: 'wait_created', + correlationId: `wait_${CORR_IDS[2]}`, + eventData: { resumeAt: resumeAtB }, + createdAt: new Date(), + }, + { + eventId: 'evnt_3', + runId: 'wrun_test', + eventType: 'hook_received', + correlationId: `hook_${CORR_IDS[0]}`, + eventData: { payload }, + createdAt: new Date(), + }, + { + eventId: 'evnt_4', + runId: 'wrun_test', + eventType: 'wait_completed', + correlationId: `wait_${CORR_IDS[1]}`, + eventData: { resumeAt: resumeAtA }, + createdAt: new Date(), + }, + { + eventId: 'evnt_5', + runId: 'wrun_test', + eventType: 'hook_received', + correlationId: `hook_${CORR_IDS[0]}`, + eventData: { payload }, + createdAt: new Date(), + }, + { + eventId: 'evnt_6', + runId: 'wrun_test', + eventType: 'wait_completed', + correlationId: `wait_${CORR_IDS[2]}`, + eventData: { resumeAt: resumeAtB }, + createdAt: new Date(), + }, + ] as unknown as Event[]; + + const ctx = setupWorkflowContext(events); + const useStep = createUseStep(ctx); + const sleep = createSleep(ctx); + const createHook = createCreateHook(ctx); + const body = async () => { + const afterA = useStep('afterA'); + const afterB = useStep('afterB'); + createHook({ token: 'dispenser-token' }); + await Promise.all([ + (async () => { + await sleep(resumeAtA); + await afterA(); + })(), + (async () => { + await sleep(resumeAtB); + await afterB(); + })(), + ]); + }; + + const { error } = await runWithDiscontinuation(ctx, body); + expect(error).toBeDefined(); + if (!WorkflowSuspension.is(error)) { + throw error; + } + const pending = [...ctx.invocationsQueue.values()].filter( + (item) => item.type === 'step' + ); + expect(pending.map((item) => item.stepName).sort()).toEqual([ + 'afterA', + 'afterB', + ]); + // Log order: waitA completed at evnt_4, waitB at evnt_6, so branch A + // draws first and afterA's correlation id sorts below afterB's. + const idOf = (name: string) => + pending.find((item) => item.stepName === name)?.correlationId ?? ''; + expect(idOf('afterA') < idOf('afterB')).toBe(true); + }); + + it('drains the registry even when the promiseQueue is rejected', async () => { + const ctx = setupWorkflowContext([]); + // One unclaimed-payload barrier that only the dispenser can retire, with + // retirement initially blocked so the dispenser has to go through its + // promiseQueue re-arm path — against a queue that is already rejected. + ctx.pendingDeliveries = 1; + const rejected = Promise.reject(new Error('poisoned queue')); + rejected.catch(() => {}); + ctx.promiseQueue = rejected as Promise; + const barrier = registerDeliveryBarrier(ctx, 0, 'hook', { armed: false }); + void barrier; // retired by the dispenser, never marked delivered + expect(isDeliveryIdle(ctx)).toBe(false); + await new Promise((resolve) => setTimeout(resolve, 25)); + // Still parked: retirement is blocked by the in-flight delivery. + expect(ctx.pendingDeliveryBarriers?.size).toBe(1); + ctx.pendingDeliveries = 0; + // The dispenser must come back from the rejected-queue backoff, retire + // the entry, and restore delivery idle so a suspension could fire. + await new Promise((resolve) => setTimeout(resolve, 200)); + expect(ctx.pendingDeliveryBarriers?.size).toBe(0); + expect(isDeliveryIdle(ctx)).toBe(true); + }); +}); diff --git a/packages/core/src/private.ts b/packages/core/src/private.ts index fccc2b3996..1e13b195b3 100644 --- a/packages/core/src/private.ts +++ b/packages/core/src/private.ts @@ -239,6 +239,13 @@ interface DeliveryBarrierEntry { * once a consumer takes the payload. */ armed: boolean; + /** + * Retire this entry: resolve `delivered` and remove it from the registry, + * exactly as `markDelivered` would. Called only by the context's safety-net + * dispenser ({@link ensureBarrierSafetyNet}), and only on the lowest-index + * entry at delivery idle. Idempotent. + */ + retire: () => void; } /** @@ -411,18 +418,18 @@ function computeResolvesOnItsOwn( * in log order; {@link hasParkedCommittedDelivery} deliberately reports such a * step as not self-resolving so that idle stays reachable. * - * "The whole chain then delivers in log order" rests on the PAYLOAD's safety - * net observing idle before the net of the wait parked behind it. If the - * wait's net fired first, the step's gate would open while the wait was still - * parked on the payload barrier and the inversion above would reappear. Within - * one drain window that order is structural, and carried by FIFO of the - * safety-net polls: nets arm via `setTimeout` in log order during synchronous - * consumption, each polling round re-arms through `promiseQueue.then(...)` in - * the order the checks ran, and each net that fires flips - * {@link hasParkedCommittedDelivery} back to true, re-blocking the rest until - * the released delivery completes. Replay — where divergence manifests — - * always consumes the log in one window. Do not "optimize" the net scheduling - * in a way that breaks that per-window FIFO. + * "The whole chain then delivers in log order" rests on the PAYLOAD's barrier + * being retired before that of anything parked behind it. That order is + * structural: safety-net retirements go through one per-context dispenser that + * only ever retires the lowest-index entry at delivery idle, and every + * retirement that wakes a chain flips {@link hasParkedCommittedDelivery} back + * to true, re-blocking the dispenser until the chain has drained — see + * {@link ensureBarrierSafetyNet}. (This used to rest on the FIFO of one idle + * poll per barrier, which held for a single parked segment but decayed to + * timing noise with several — the release order, and therefore the ULIDs + * drawn by the woken branches, then depended on how much log the replay had + * loaded. storm-log-replay.test.ts replays a production log corrupted exactly + * that way.) */ export async function awaitEarlierDeliveries( ctx: WorkflowOrchestratorContext, @@ -517,12 +524,6 @@ export function registerDeliveryBarrier( let done = false; const { promise, resolve } = withResolvers(); - const entry: DeliveryBarrierEntry = { - kind, - delivered: promise, - armed: options.armed ?? true, - }; - barriers.set(eventIndex, entry); const finish = () => { if (done) { @@ -535,12 +536,23 @@ export function registerDeliveryBarrier( resolve(); }; + const entry: DeliveryBarrierEntry = { + kind, + delivered: promise, + armed: options.armed ?? true, + retire: finish, + }; + barriers.set(eventIndex, entry); + // Safety net: if this delivery is never delivered to the workflow (its // branch was not taken / the run is suspending, or a buffered hook payload // is only claimed after a later delivery the workflow is still waiting on), - // resolve at idle so a later delivery gated on it cannot deadlock and the - // registry cannot leak an entry per abandoned delivery. - scheduleWhenIdle(ctx, finish); + // it is retired at idle so a later delivery gated on it cannot deadlock and + // the registry cannot leak an entry per abandoned delivery. Retirement goes + // through the context's single ordered dispenser rather than a per-barrier + // idle poll — see {@link ensureBarrierSafetyNet} for why the ORDER of these + // retirements is load-bearing. + ensureBarrierSafetyNet(ctx); return { markDelivered: finish, @@ -550,6 +562,116 @@ export function registerDeliveryBarrier( }; } +/** + * Contexts whose barrier safety-net dispenser is currently armed. Module-level + * so the context interface (constructed literally by many test harnesses) + * needs no new field; entries drop with the context. + */ +const activeBarrierSafetyNets = new WeakSet(); + +/** + * The barrier registry's safety net: ONE idle-gated dispenser per context that + * retires, at each observation of delivery idle, only the LOWEST-index entry + * still registered, then yields so the chain it released can run before the + * next retirement is considered. + * + * Why one ordered dispenser and not a poll per barrier (which is what this + * replaced): the order of safety-net retirements decides the delivery order of + * every chain parked behind an unclaimed buffered hook payload — a hook the + * workflow never reads (a fire-and-forget `createHook`) parks every later + * armed wait/hook behind a barrier that only this net can retire. Per-barrier + * polls fire in whatever order their re-arm cycles land, and each re-arm + * attaches to a `promiseQueue` that grows between checks, so with several + * parked segments the release order decays to timing noise. Draws (`useStep` + * correlation ids) then depend on which segment happened to release first — + * concretely, on how MUCH log the replay loaded, since that decides what is in + * the registry. That is the mechanism behind slot-mode CORRUPTED_EVENT_LOG on + * storm-shaped runs (see storm-log-replay.test.ts, built from a production + * log): two replays of the same run holding different-length prefixes bound + * the same correlation ordinal to different steps. + * + * Retiring lowest-first is not merely tidy, it is the only order that cannot + * invert the log: every gate points from a higher index to a strictly lower + * one, so at delivery idle the lowest undelivered entry gates on nothing + * still registered — it is the head of every parked chain (in practice, the + * unclaimed payload itself). Releasing it lets the chain above deliver + * through the ordinary barrier order; anything the release wakes flips + * {@link hasParkedCommittedDelivery} back to true, which re-blocks this + * dispenser until the chain has fully drained. A higher entry must never be + * retired while a lower one is registered — that is exactly the inversion + * described on {@link awaitEarlierDeliveries}. + * + * The dispenser goes dormant when the registry empties and is re-armed by the + * next registration, so an idle context holds no live timer. + */ +function ensureBarrierSafetyNet(ctx: WorkflowOrchestratorContext): void { + const barriers = ctx.pendingDeliveryBarriers; + if (!barriers || activeBarrierSafetyNets.has(ctx)) { + return; + } + activeBarrierSafetyNets.add(ctx); + const rearm = () => { + setTimeout(check, 0); + }; + // A rejected promiseQueue settles immediately and forever, so re-arming + // through it at the normal cadence would degenerate into a busy loop on an + // abandoned context. Back off instead: correctness only needs the dispenser + // to still exist, since the registry gates suspension via isDeliveryIdle + // and a dead dispenser would wedge the run. + const rearmAfterRejection = () => { + setTimeout(check, 50); + }; + const check = () => { + if (barriers.size === 0) { + // Dormant. The next registerDeliveryBarrier re-arms. + activeBarrierSafetyNets.delete(ctx); + return; + } + if (!canRetireAbandonedBarriers(ctx)) { + // A delivery is hydrating or committed-but-parked on its deferral; let + // the queue drain and re-check a tick later (same cadence as + // scheduleWhenIdle). + ctx.promiseQueue.then(rearm, rearmAfterRejection); + return; + } + // Idle with entries left: nothing remaining delivers on its own, so + // release parked chains from the head — lowest index first, one at a + // time, re-reading idle between retirements. A retirement that wakes a + // chain flips {@link hasParkedCommittedDelivery} synchronously (it is + // computed from the registry this loop just mutated), which stops the + // sweep so the chain delivers before anything above it is released. A + // retirement that wakes nothing (a stale payload no delivery gates on) + // keeps the sweep going, so a backlog of those drains in ONE idle + // observation — pacing them one per timer tick would hold consumed-but- + // undelivered events hostage long enough to trip the events consumer's + // unconsumed-event deadline and fail healthy replays. + while (barriers.size > 0 && canRetireAbandonedBarriers(ctx)) { + let lowestIndex: number | undefined; + let lowestEntry: DeliveryBarrierEntry | undefined; + for (const [index, entry] of barriers) { + if (lowestIndex === undefined || index < lowestIndex) { + lowestIndex = index; + lowestEntry = entry; + } + } + lowestEntry?.retire(); + } + setTimeout(check, 0); + }; + setTimeout(check, 0); +} + +/** + * Whether the safety-net dispenser may retire abandoned barriers right now: + * no hydration in flight and no committed delivery still working through its + * detached deferral. This is deliberately WEAKER than {@link isDeliveryIdle}: + * the dispenser is what empties the registry, so gating it on registry + * emptiness would gate its own work. + */ +function canRetireAbandonedBarriers(ctx: WorkflowOrchestratorContext): boolean { + return ctx.pendingDeliveries === 0 && !hasParkedCommittedDelivery(ctx); +} + /** * Whether some registered branch-deciding delivery is going to reach the * workflow without any further help (it is armed and not transitively parked @@ -612,9 +734,24 @@ export function hasParkedCommittedDelivery( * what it has and has not done yet says nothing about the run. Two callers * read it, for the two such decisions: {@link scheduleWhenIdle} for the * suspension, and the events consumer's unconsumed-event check for divergence. + * + * A non-empty barrier registry counts as in flight, even when every remaining + * entry is parked behind an unclaimed buffered payload. Those entries only + * move when the safety-net dispenser retires them (lowest-first, see + * {@link ensureBarrierSafetyNet}), and the deliveries they release are real + * workflow reactions — a suspension raised before they run would be computed + * from a VM that has not seen them, scheduling none of their follow-up work + * and leaving the run dormant (the vercel/workflow#3183 shape). The dispenser + * itself is gated on {@link canRetireAbandonedBarriers}, the weaker predicate + * without the registry term, precisely so it can do the draining that this + * predicate waits for; registry size strictly decreases at each retirement, + * so idle is always reached. */ export function isDeliveryIdle(ctx: WorkflowOrchestratorContext): boolean { - return ctx.pendingDeliveries === 0 && !hasParkedCommittedDelivery(ctx); + return ( + ctx.pendingDeliveries === 0 && + (!ctx.pendingDeliveryBarriers || ctx.pendingDeliveryBarriers.size === 0) + ); } /** diff --git a/packages/core/src/race-padded-draw-ordering.test.ts b/packages/core/src/race-padded-draw-ordering.test.ts new file mode 100644 index 0000000000..42c81aa64d --- /dev/null +++ b/packages/core/src/race-padded-draw-ordering.test.ts @@ -0,0 +1,377 @@ +/** + * Reproduction attempt for the residual CORRUPTED_EVENT_LOG shape observed on + * spec-6 (slot-identity) runs, most recently + * `wrun_41KZYJ92TP0GYBNDKW3FJBWQ3Y` (step-storm repro on preview, + * 2026-08-13): concurrent writers bound the same drawn correlation id to + * different steps (`step_…QMWZ` = finalizeStep at slot 630 vs the canonical + * replay's releaseStep), and one logical finalize step was created and + * executed under multiple ids. + * + * The workflow shape in that run is `Promise.race([settleStep(), sleep(t)])` + * per branch: the watchdog path draws two follow-up ids (recover, finalize) + * where the settled path draws one (finalize). The race adds microtask hops + * BETWEEN a step result's barrier-ordered resolution and the branch's next + * draw — exactly the "padded consumer" residual that + * `step-delivery-ordering.test.ts` calls out of scope and + * `step-delivery-hop-count.test.ts` pins for plain awaits. + * + * Two orderings are asserted, cold and warm (shared ReplayPayloadCache): + * + * 1. A wait-woken branch's draw (recover) vs a step-woken branch's draw + * (finalize) with the wait earlier in the log — the covered class, with + * race padding on both consumers. + * 2. The settled branch's finalize draw (step_completed at slot i) vs the + * watchdog branch's finalize draw (recover step_completed at slot j > i) + * — the inversion the failed run's writer actually committed (its + * settled-branch finalize minted AFTER all recovery finalizes). + */ +import { WorkflowRuntimeError } from '@workflow/errors'; +import { withResolvers } from '@workflow/utils'; +import type { Event } from '@workflow/world'; +import * as nanoid from 'nanoid'; +import { monotonicFactory } from 'ulid'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { EventsConsumer } from './events-consumer.js'; +import { WorkflowSuspension } from './global.js'; +import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js'; +import { ReplayPayloadCache } from './replay-payload-cache.js'; +import { dehydrateStepReturnValue } from './serialization.js'; +import { createUseStep } from './step.js'; +import { createContext } from './vm/index.js'; +import { createSleep } from './workflow/sleep.js'; + +function setupWorkflowContext( + events: Event[], + replayPayloadCache: ReplayPayloadCache = new ReplayPayloadCache(undefined) +): WorkflowOrchestratorContext { + const context = createContext({ + seed: 'test', + fixedTimestamp: 1753481739458, + }); + const ulid = monotonicFactory(() => context.globalThis.Math.random()); + const workflowStartedAt = context.globalThis.Date.now(); + const promiseQueueHolder = { current: Promise.resolve() }; + const ctxRef: { current?: WorkflowOrchestratorContext } = {}; + const ctx: WorkflowOrchestratorContext = { + suspensionGeneration: 0, + runId: 'wrun_test', + encryptionKey: undefined, + replayPayloadCache, + globalThis: context.globalThis, + eventsConsumer: new EventsConsumer(events, { + isDeliveryIdle: () => + ctxRef.current ? isDeliveryIdle(ctxRef.current) : true, + onUnconsumedEvent: (event) => { + ctxRef.current?.onWorkflowError( + new WorkflowRuntimeError( + `Unconsumed event in event log: eventType=${event.eventType}, correlationId=${event.correlationId}, eventId=${event.eventId}.` + ) + ); + }, + getPromiseQueue: () => promiseQueueHolder.current, + }), + invocationsQueue: new Map(), + generateUlid: () => ulid(workflowStartedAt), + generateNanoid: nanoid.customRandom(nanoid.urlAlphabet, 21, (size) => + new Uint8Array(size).map(() => 256 * context.globalThis.Math.random()) + ), + onWorkflowError: vi.fn(), + get promiseQueue() { + return promiseQueueHolder.current; + }, + set promiseQueue(value: Promise) { + promiseQueueHolder.current = value; + }, + pendingDeliveries: 0, + pendingDeliveryBarriers: new Map(), + }; + ctxRef.current = ctx; + return ctx; +} + +// Deterministic correlation IDs from the ULID generator with seed 'test'. +// Draw order in the body below: +// c0 = settleW (branch W's raced step) +// c1 = sleepW (branch W's watchdog) +// c2 = settleS (branch S's raced step) +// c3 = sleepS (branch S's watchdog) +// c4..c6 = the follow-up draws whose order is under test +const CORR_IDS = [ + '01K11TFZ62YS0YYFDQ3E8B9YCV', + '01K11TFZ62YS0YYFDQ3E8B9YCW', + '01K11TFZ62YS0YYFDQ3E8B9YCX', + '01K11TFZ62YS0YYFDQ3E8B9YCY', + '01K11TFZ62YS0YYFDQ3E8B9YCZ', + '01K11TFZ62YS0YYFDQ3E8B9YD0', + '01K11TFZ62YS0YYFDQ3E8B9YD1', +]; + +const WATCHDOG = Symbol.for('race-padded-draw-ordering:watchdog'); + +function pendingStepNames(ctx: WorkflowOrchestratorContext): string[] { + return [...ctx.invocationsQueue.values()] + .filter((item) => item.type === 'step') + .map((item) => (item.type === 'step' ? item.stepName : '')); +} + +async function runWithDiscontinuation( + ctx: WorkflowOrchestratorContext, + workflowFn: () => Promise +): Promise<{ result?: any; error?: any }> { + const workflowDiscontinuation = withResolvers(); + ctx.onWorkflowError = workflowDiscontinuation.reject; + + let result: any; + let error: any; + try { + result = await Promise.race([ + workflowFn(), + workflowDiscontinuation.promise, + ]); + } catch (err) { + error = err; + } + return { result, error }; +} + +function delayHydration() { + const hydrateSpy = vi.fn(); + return { + hydrateSpy, + install: async () => { + const serialization = await import('./serialization.js'); + const originalHydrate = serialization.hydrateStepReturnValue; + return vi + .spyOn(serialization, 'hydrateStepReturnValue') + .mockImplementation(async (...args) => { + hydrateSpy(); + await new Promise((r) => setTimeout(r, 10)); + return originalHydrate(...args); + }); + }, + }; +} + +describe('race-padded consumers draw in event-log order', () => { + let spy: ReturnType | undefined; + + afterEach(() => { + spy?.mockRestore(); + spy = undefined; + }); + + const resumeAtW = new Date('2026-07-27T12:00:05.000Z'); + const resumeAtS = new Date('2026-07-27T12:00:06.000Z'); + + /** + * The live invocation's history, exactly as the failed run's final round + * recorded it (two branches instead of eight): + * + * - branch W's watchdog fired first (wait_completed lowest), + * so W drew c4 = recoverStep; + * - branch S's raced step then completed, so S drew c5 = finalizeS; + * - W's recover completed last, so W drew c6 = finalizeW. + */ + async function buildEventLog(): Promise { + const ops: Promise[] = []; + const [settleSResult, recoverResult] = await Promise.all([ + dehydrateStepReturnValue('settled', 'wrun_test', undefined, ops), + dehydrateStepReturnValue('recovered', 'wrun_test', undefined, ops), + ]); + + const at = () => new Date(); + return [ + // Round setup: both branches suspend together. + { + eventId: 'evnt_00', + runId: 'wrun_test', + eventType: 'step_created', + correlationId: `step_${CORR_IDS[0]}`, + eventData: { stepName: 'settleW' }, + createdAt: at(), + }, + { + eventId: 'evnt_01', + runId: 'wrun_test', + eventType: 'wait_created', + correlationId: `wait_${CORR_IDS[1]}`, + eventData: { resumeAt: resumeAtW }, + createdAt: at(), + }, + { + eventId: 'evnt_02', + runId: 'wrun_test', + eventType: 'step_created', + correlationId: `step_${CORR_IDS[2]}`, + eventData: { stepName: 'settleS' }, + createdAt: at(), + }, + { + eventId: 'evnt_03', + runId: 'wrun_test', + eventType: 'wait_created', + correlationId: `wait_${CORR_IDS[3]}`, + eventData: { resumeAt: resumeAtS }, + createdAt: at(), + }, + { + eventId: 'evnt_04', + runId: 'wrun_test', + eventType: 'step_started', + correlationId: `step_${CORR_IDS[0]}`, + eventData: { stepName: 'settleW' }, + createdAt: at(), + }, + { + eventId: 'evnt_05', + runId: 'wrun_test', + eventType: 'step_started', + correlationId: `step_${CORR_IDS[2]}`, + eventData: { stepName: 'settleS' }, + createdAt: at(), + }, + // Branch W's watchdog fires: W's race resolves 'watchdog', W draws c4. + { + eventId: 'evnt_06', + runId: 'wrun_test', + eventType: 'wait_completed', + correlationId: `wait_${CORR_IDS[1]}`, + eventData: { resumeAt: resumeAtW }, + createdAt: at(), + }, + // Branch S's raced step completes: S's race resolves 'settled', + // S draws c5. Hydration-sensitive delivery. + { + eventId: 'evnt_07', + runId: 'wrun_test', + eventType: 'step_completed', + correlationId: `step_${CORR_IDS[2]}`, + eventData: { stepName: 'settleS', result: settleSResult }, + createdAt: at(), + }, + // W's recovery step, drawn at evnt_06. + { + eventId: 'evnt_08', + runId: 'wrun_test', + eventType: 'step_created', + correlationId: `step_${CORR_IDS[4]}`, + eventData: { stepName: 'recoverW' }, + createdAt: at(), + }, + { + eventId: 'evnt_09', + runId: 'wrun_test', + eventType: 'step_started', + correlationId: `step_${CORR_IDS[4]}`, + eventData: { stepName: 'recoverW' }, + createdAt: at(), + }, + { + eventId: 'evnt_10', + runId: 'wrun_test', + eventType: 'step_completed', + correlationId: `step_${CORR_IDS[4]}`, + eventData: { stepName: 'recoverW', result: recoverResult }, + createdAt: at(), + }, + // S's finalize, drawn at evnt_07 — BEFORE W's finalize in ULID order. + { + eventId: 'evnt_11', + runId: 'wrun_test', + eventType: 'step_created', + correlationId: `step_${CORR_IDS[5]}`, + eventData: { stepName: 'finalizeS' }, + createdAt: at(), + }, + // W's finalize, drawn at evnt_10. + { + eventId: 'evnt_12', + runId: 'wrun_test', + eventType: 'step_created', + correlationId: `step_${CORR_IDS[6]}`, + eventData: { stepName: 'finalizeW' }, + createdAt: at(), + }, + ]; + } + + function workflowBody(ctx: WorkflowOrchestratorContext) { + const useStep = createUseStep(ctx); + const sleep = createSleep(ctx); + + return async () => { + const settleW = useStep('settleW'); + const settleS = useStep('settleS'); + const recoverW = useStep('recoverW'); + const finalizeW = useStep('finalizeW'); + const finalizeS = useStep('finalizeS'); + + const branchW = (async () => { + const winner = await Promise.race([ + settleW(), + sleep(resumeAtW).then(() => WATCHDOG), + ]); + if (winner === WATCHDOG) { + await recoverW(); + } + await finalizeW(); + })(); + + const branchS = (async () => { + const winner = await Promise.race([ + settleS(), + sleep(resumeAtS).then(() => WATCHDOG), + ]); + if (winner === WATCHDOG) { + throw new Error('branch S must settle in this log'); + } + await finalizeS(); + })(); + + await Promise.all([branchW, branchS]); + }; + } + + async function assertLogOrderReproduced( + events: Event[], + cache: ReplayPayloadCache + ) { + const ctx = setupWorkflowContext(events, cache); + const { error } = await runWithDiscontinuation(ctx, workflowBody(ctx)); + expect(error).toBeDefined(); + if (!WorkflowSuspension.is(error)) { + throw error; + } + // Correct behavior: the replay agrees with every committed binding. + // `settleW` stays pending — it lost its race and never completed, so its + // consumer legitimately outlives the round. + expect(pendingStepNames(ctx).sort()).toEqual([ + 'finalizeS', + 'finalizeW', + 'settleW', + ]); + expect(ctx.eventsConsumer.eventIndex).toBe(events.length); + } + + it('reproduces the recorded draw order on a cold replay', async () => { + const hydration = delayHydration(); + spy = await hydration.install(); + const events = await buildEventLog(); + await assertLogOrderReproduced(events, new ReplayPayloadCache(undefined)); + }); + + it('reproduces the recorded draw order on a warm replay sharing the payload cache', async () => { + const hydration = delayHydration(); + spy = await hydration.install(); + const events = await buildEventLog(); + const sharedCache = new ReplayPayloadCache(undefined); + // Cold pass primes the cache the way the first replay of a queue + // delivery does. + await assertLogOrderReproduced(events, sharedCache); + expect(hydration.hydrateSpy).toHaveBeenCalled(); + // Warm pass: the memoized primitive result now resolves in fewer hops + // than the wait, which is the asymmetry that reordered draws in + // production. + await assertLogOrderReproduced(events, sharedCache); + }); +}); diff --git a/packages/core/src/storm-log-replay.test.ts b/packages/core/src/storm-log-replay.test.ts new file mode 100644 index 0000000000..4e7a14c7d6 --- /dev/null +++ b/packages/core/src/storm-log-replay.test.ts @@ -0,0 +1,357 @@ +/** + * Offline replay of the ACTUAL corrupted production log from + * `wrun_41KZYJ92TP0GYBNDKW3FJBWQ3Y` (step-storm repro, preview, 2026-08-13, + * CORRUPTED_EVENT_LOG after 4 deterministic divergences at slot 630). + * + * The fixture (/tmp/rr_fixture.json, built from ClickHouse + * workflow_observability_staging) preserves every committed event's slot + * order, type, entity kind, step name, and the ULID *rank* of its correlation + * id. Ranks are remapped onto this harness's deterministic ULID sequence, so + * the workflow body below (a faithful port of `stepStormReproWorkflow`) + * mints ids that line up rank-for-rank with the committed ones. + * + * Two questions, one replay each: + * + * 1. Full log: does a faithful replay diverge at the slot-630 equivalent + * with "belongs to finalizeStep, but the current step consumer is + * releaseStep"? (Validates the canonical-order derivation and that the + * committed bindings really are mutually inconsistent.) + * + * 2. Writer B's prefix (slots 1..610): what does a faithful replay of + * exactly what 45cb4c904d25 loaded (`eventCount: 610`) draw for the + * pending creates? If it binds rank 198 (WZ) to finalizeStep, B was + * prefix-determined; if not, B's committed binding deviated from its own + * prefix and the bug is in the live loop, not in replay determinism. + */ +import { readFileSync } from 'node:fs'; +import { dirname, join } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { WorkflowRuntimeError } from '@workflow/errors'; +import { withResolvers } from '@workflow/utils'; +import type { Event } from '@workflow/world'; +import * as nanoid from 'nanoid'; +import { monotonicFactory } from 'ulid'; +import { describe, expect, it } from 'vitest'; +import { EventsConsumer } from './events-consumer.js'; +import { WorkflowSuspension } from './global.js'; +import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js'; +import { ReplayPayloadCache } from './replay-payload-cache.js'; +import { dehydrateStepReturnValue } from './serialization.js'; +import { createUseStep } from './step.js'; +import { createContext } from './vm/index.js'; +import { createCreateHook } from './workflow/hook.js'; +import { createSleep } from './workflow/sleep.js'; + +const SEED = 'test'; +const FIXED_TS = 1753481739458; + +/** The same deterministic ULID sequence the replay context will draw from. */ +function generateUlidSequence(count: number): string[] { + const context = createContext({ seed: SEED, fixedTimestamp: FIXED_TS }); + const ulid = monotonicFactory(() => context.globalThis.Math.random()); + const at = context.globalThis.Date.now(); + return Array.from({ length: count }, () => ulid(at)); +} + +function setupWorkflowContext( + events: Event[], + replayPayloadCache: ReplayPayloadCache = new ReplayPayloadCache(undefined) +): WorkflowOrchestratorContext { + const context = createContext({ seed: SEED, fixedTimestamp: FIXED_TS }); + const ulid = monotonicFactory(() => context.globalThis.Math.random()); + const workflowStartedAt = context.globalThis.Date.now(); + const promiseQueueHolder = { current: Promise.resolve() }; + const ctxRef: { current?: WorkflowOrchestratorContext } = {}; + const ctx: WorkflowOrchestratorContext = { + suspensionGeneration: 0, + runId: 'wrun_test', + encryptionKey: undefined, + replayPayloadCache, + globalThis: context.globalThis, + eventsConsumer: new EventsConsumer(events, { + isDeliveryIdle: () => + ctxRef.current ? isDeliveryIdle(ctxRef.current) : true, + onUnconsumedEvent: (event) => { + ctxRef.current?.onWorkflowError( + new WorkflowRuntimeError( + `Unconsumed event: eventType=${event.eventType}, correlationId=${event.correlationId}, eventId=${event.eventId}.` + ) + ); + }, + getPromiseQueue: () => promiseQueueHolder.current, + }), + invocationsQueue: new Map(), + generateUlid: () => ulid(workflowStartedAt), + generateNanoid: nanoid.customRandom(nanoid.urlAlphabet, 21, (size) => + new Uint8Array(size).map(() => 256 * context.globalThis.Math.random()) + ), + onWorkflowError: () => {}, + get promiseQueue() { + return promiseQueueHolder.current; + }, + set promiseQueue(value: Promise) { + promiseQueueHolder.current = value; + }, + pendingDeliveries: 0, + pendingDeliveryBarriers: new Map(), + }; + ctxRef.current = ctx; + return ctx; +} + +interface FixtureEvent { + s: number; // slot + t: string; // eventType + k: string; // step | wait | hook | run (attr_set) + r: number; // ULID rank of the correlation id, -1 for run-scoped + n: string; // step name (short) +} + +const WATCHDOG = Symbol.for('storm-log-replay:watchdog'); +const TOKEN = 'storm-log-replay-token'; + +async function buildEvents( + fixture: FixtureEvent[], + ulids: string[] +): Promise { + const ops: Promise[] = []; + const stepResult = await dehydrateStepReturnValue( + { ok: true }, + 'wrun_test', + undefined, + ops + ); + const hookPayload = await dehydrateStepReturnValue( + { round: 0, index: 0, sentAt: 1 }, + 'wrun_test', + undefined, + ops + ); + const events: Event[] = []; + for (const f of fixture) { + const eventId = `evnt_${String(f.s).padStart(26, '0')}`; + const createdAt = new Date(FIXED_TS + f.s); + if (f.t === 'attr_set') { + events.push({ + eventId, + runId: 'wrun_test', + eventType: 'attr_set', + correlationId: 'wrun_test', + eventData: { attributes: { settle: 'x' } }, + createdAt, + } as unknown as Event); + continue; + } + const correlationId = `${f.k}_${ulids[f.r]}`; + let eventData: Record; + switch (f.t) { + case 'hook_created': + eventData = { token: `${TOKEN}:poke`, isWebhook: false }; + break; + case 'hook_received': + eventData = { payload: hookPayload }; + break; + case 'wait_created': + case 'wait_completed': + // Same value on created/completed per entity is all the consumer + // requires (it re-reads resumeAt off wait_created). + eventData = { resumeAt: new Date(FIXED_TS + 1_000_000 + f.r) }; + break; + case 'step_created': + case 'step_started': + eventData = { stepName: f.n }; + break; + case 'step_completed': + eventData = { stepName: f.n, result: stepResult }; + break; + default: + throw new Error(`unhandled event type ${f.t}`); + } + events.push({ + eventId, + runId: 'wrun_test', + eventType: f.t, + correlationId, + eventData, + createdAt, + } as unknown as Event); + } + return events; +} + +function workflowBody(ctx: WorkflowOrchestratorContext) { + const useStep = createUseStep(ctx); + const sleep = createSleep(ctx); + const createHook = createCreateHook(ctx); + const settleStep = useStep('settleStep'); + const recoverStep = useStep('recoverStep'); + const finalizeStep = useStep('finalizeStep'); + const releaseStep = useStep('releaseStep'); + const reconcileStep = useStep('reconcileStep'); + + const cfg = { + rounds: 6, + width: 8, + watchdogMs: 2500, + betweenRoundSleepMs: 1000, + reconcileBase: 2, + }; + + return async () => { + const pokeHook = createHook({ token: `${TOKEN}:poke` }); + try { + for (let round = 0; round < cfg.rounds; round += 1) { + const branches = await Promise.all( + Array.from({ length: cfg.width }, (_, index) => + (async () => { + try { + const winner = await Promise.race([ + settleStep({ round, index }), + sleep(cfg.watchdogMs).then(() => WATCHDOG), + ]); + if (winner === WATCHDOG) { + await recoverStep({ round, index }); + await finalizeStep({ round, index, winner: 'watchdog' }); + return { winner: 'watchdog' as const }; + } + await finalizeStep({ round, index, winner: 'settled' }); + return { winner: 'settled' as const }; + } finally { + await releaseStep({ round, index }); + } + })() + ) + ); + const stragglers = branches.filter( + (b) => b.winner === 'watchdog' + ).length; + await Promise.all( + Array.from({ length: cfg.reconcileBase + stragglers }, (_, index) => + reconcileStep({ round, index, stragglers }) + ) + ); + if (cfg.betweenRoundSleepMs > 0) { + await sleep(cfg.betweenRoundSleepMs); + } + } + } finally { + pokeHook.dispose(); + } + }; +} + +async function replay(events: Event[]): Promise<{ + error?: any; + ctx: WorkflowOrchestratorContext; +}> { + const ctx = setupWorkflowContext(events); + const discontinuation = withResolvers(); + ctx.onWorkflowError = discontinuation.reject; + let error: any; + try { + await Promise.race([workflowBody(ctx)(), discontinuation.promise]); + } catch (err) { + error = err; + } + return { error, ctx }; +} + +describe('replaying the corrupted production storm log', () => { + const fixture: FixtureEvent[] = JSON.parse( + readFileSync( + join( + dirname(fileURLToPath(import.meta.url)), + '__fixtures__', + 'wrun-41KZYJ92TP-storm-log.json' + ), + 'utf8' + ) + ); + const maxRank = Math.max(...fixture.map((f) => f.r)); + const ulids = generateUlidSequence(maxRank + 8); + + /** Pending-step bindings by ULID rank after replaying `len` slots. */ + async function bindingsAtLength( + len: number + ): Promise<{ error: any; byRank: Map }> { + const prefix = fixture.filter((f) => f.s <= len); + const events = await buildEvents(prefix, ulids); + const { error, ctx } = await replay(events); + const byRank = new Map(); + for (const item of ctx.invocationsQueue.values()) { + if (item.type !== 'step') continue; + const rank = ulids.indexOf(item.correlationId.split('_', 2)[1]); + byRank.set(rank, item.stepName); + } + return { error, byRank }; + } + + it('full log: still diverges — the committed log holds bindings from two incompatible trajectories', async () => { + // The production writers created the SAME logical finalize step under two + // correlation ids (ranks 198 and 199, slots 630 and 631) from + // different-length prefixes under the pre-fix scheduler. No single + // deterministic trajectory can satisfy both creates, so a faithful replay + // of the full log must reject one of them. What the fix guarantees is not + // that this log becomes readable, but that new logs cannot acquire this + // shape: writers holding different-length prefixes now draw identical + // bindings (the tests below). + const events = await buildEvents(fixture, ulids); + const { error } = await replay(events); + expect(error).toBeDefined(); + expect(WorkflowSuspension.is(error)).toBe(false); + expect(String(error?.message)).toContain('Replay divergence'); + }); + + it("writer B's exact 610-event prefix reproduces writer B's committed binding", async () => { + // Rank 198 is `step_…QMWZ`, which the writer holding this exact prefix + // (eventCount: 610 in its runtime logs) committed as finalizeStep at slot + // 630. A faithful replay of its prefix must derive the same pending + // create, or the writer was never prefix-determined and replay itself is + // nondeterministic. + const { error, byRank } = await bindingsAtLength(610); + expect(WorkflowSuspension.is(error)).toBe(true); + expect(byRank.get(197)).toContain('releaseStep'); + expect(byRank.get(198)).toContain('finalizeStep'); + }); + + it('draw bindings are stable under log extension', async () => { + // THE regression assertion for the ordered safety-net dispenser + // (ensureBarrierSafetyNet): pending-step bindings derived from a prefix + // must never change when the same replay code is handed MORE of the same + // log. Before the fix, extending this log from 611 to 612 slots moved + // rank 198 from finalizeStep to releaseStep — two honest replayers with + // different-length snapshots then committed conflicting creates, which is + // the residual slot-mode CORRUPTED_EVENT_LOG mechanism + // (wrun_41KZYJ92TP0GYBNDKW3FJBWQ3Y). + // + // 630 is the longest clean prefix: 631 holds the second of the two + // incompatible committed creates, past which replay rightly diverges. + const lengths = [610, 611, 612, 619, 630]; + const results = new Map>(); + for (const len of lengths) { + const { error, byRank } = await bindingsAtLength(len); + if (!WorkflowSuspension.is(error)) { + throw new Error( + `prefix len ${len} did not suspend: ${error?.constructor?.name}: ${error?.message}`, + { cause: error } + ); + } + results.set(len, byRank); + } + for (let i = 1; i < lengths.length; i++) { + const shorter = results.get(lengths[i - 1])!; + const longer = results.get(lengths[i])!; + for (const [rank, name] of shorter) { + const extended = longer.get(rank); + // A rank absent from the longer replay was consumed by its (matching) + // created event arriving in the extension — only disagreement fails. + if (extended !== undefined) { + expect( + `${lengths[i]}:r${rank}=${extended}`, + `rank ${rank} rebound between len ${lengths[i - 1]} and ${lengths[i]}` + ).toBe(`${lengths[i]}:r${rank}=${name}`); + } + } + } + }); +}); diff --git a/packages/core/src/storm-log-sweep.test.ts b/packages/core/src/storm-log-sweep.test.ts new file mode 100644 index 0000000000..822621e516 --- /dev/null +++ b/packages/core/src/storm-log-sweep.test.ts @@ -0,0 +1,269 @@ +/** + * Companion to storm-log-replay.test.ts: sweep prefix lengths of the real + * corrupted log and report, per length, which step name ranks 197 (WY) and + * 198 (WZ) get bound to. Finds the exact log-extension point where a + * byte-identical shared prefix changes its own draw bindings. + */ +import { readFileSync } from 'node:fs'; +import { dirname, join } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { WorkflowRuntimeError } from '@workflow/errors'; +import { withResolvers } from '@workflow/utils'; +import type { Event } from '@workflow/world'; +import * as nanoid from 'nanoid'; +import { monotonicFactory } from 'ulid'; +import { describe, it } from 'vitest'; +import { EventsConsumer } from './events-consumer.js'; +import { isDeliveryIdle, type WorkflowOrchestratorContext } from './private.js'; +import { ReplayPayloadCache } from './replay-payload-cache.js'; +import { dehydrateStepReturnValue } from './serialization.js'; +import { createUseStep } from './step.js'; +import { createContext } from './vm/index.js'; +import { createCreateHook } from './workflow/hook.js'; +import { createSleep } from './workflow/sleep.js'; + +const SEED = 'test'; +const FIXED_TS = 1753481739458; + +function generateUlidSequence(count: number): string[] { + const context = createContext({ seed: SEED, fixedTimestamp: FIXED_TS }); + const ulid = monotonicFactory(() => context.globalThis.Math.random()); + const at = context.globalThis.Date.now(); + return Array.from({ length: count }, () => ulid(at)); +} + +function setupWorkflowContext(events: Event[]): WorkflowOrchestratorContext { + const context = createContext({ seed: SEED, fixedTimestamp: FIXED_TS }); + const ulid = monotonicFactory(() => context.globalThis.Math.random()); + const workflowStartedAt = context.globalThis.Date.now(); + const promiseQueueHolder = { current: Promise.resolve() }; + const ctxRef: { current?: WorkflowOrchestratorContext } = {}; + const ctx: WorkflowOrchestratorContext = { + suspensionGeneration: 0, + runId: 'wrun_test', + encryptionKey: undefined, + replayPayloadCache: new ReplayPayloadCache(undefined), + globalThis: context.globalThis, + eventsConsumer: new EventsConsumer(events, { + isDeliveryIdle: () => + ctxRef.current ? isDeliveryIdle(ctxRef.current) : true, + onUnconsumedEvent: (event) => { + ctxRef.current?.onWorkflowError( + new WorkflowRuntimeError( + `Unconsumed event: ${event.eventType} ${event.correlationId}` + ) + ); + }, + getPromiseQueue: () => promiseQueueHolder.current, + }), + invocationsQueue: new Map(), + generateUlid: () => ulid(workflowStartedAt), + generateNanoid: nanoid.customRandom(nanoid.urlAlphabet, 21, (size) => + new Uint8Array(size).map(() => 256 * context.globalThis.Math.random()) + ), + onWorkflowError: () => {}, + get promiseQueue() { + return promiseQueueHolder.current; + }, + set promiseQueue(value: Promise) { + promiseQueueHolder.current = value; + }, + pendingDeliveries: 0, + pendingDeliveryBarriers: new Map(), + }; + ctxRef.current = ctx; + return ctx; +} + +interface FixtureEvent { + s: number; + t: string; + k: string; + r: number; + n: string; +} + +const WATCHDOG = Symbol.for('storm-log-sweep:watchdog'); +const TOKEN = 'storm-log-replay-token'; + +async function buildEvents( + fixture: FixtureEvent[], + ulids: string[] +): Promise { + const ops: Promise[] = []; + const stepResult = await dehydrateStepReturnValue( + { ok: true }, + 'wrun_test', + undefined, + ops + ); + const hookPayload = await dehydrateStepReturnValue( + { round: 0, index: 0, sentAt: 1 }, + 'wrun_test', + undefined, + ops + ); + const events: Event[] = []; + for (const f of fixture) { + const eventId = `evnt_${String(f.s).padStart(26, '0')}`; + const createdAt = new Date(FIXED_TS + f.s); + if (f.t === 'attr_set') { + events.push({ + eventId, + runId: 'wrun_test', + eventType: 'attr_set', + correlationId: 'wrun_test', + eventData: { attributes: { settle: 'x' } }, + createdAt, + } as unknown as Event); + continue; + } + const correlationId = `${f.k}_${ulids[f.r]}`; + let eventData: Record; + switch (f.t) { + case 'hook_created': + eventData = { token: `${TOKEN}:poke`, isWebhook: false }; + break; + case 'hook_received': + eventData = { payload: hookPayload }; + break; + case 'wait_created': + case 'wait_completed': + eventData = { resumeAt: new Date(FIXED_TS + 1_000_000 + f.r) }; + break; + case 'step_created': + case 'step_started': + eventData = { stepName: f.n }; + break; + case 'step_completed': + eventData = { stepName: f.n, result: stepResult }; + break; + default: + throw new Error(`unhandled event type ${f.t}`); + } + events.push({ + eventId, + runId: 'wrun_test', + eventType: f.t, + correlationId, + eventData, + createdAt, + } as unknown as Event); + } + return events; +} + +function workflowBody(ctx: WorkflowOrchestratorContext) { + const useStep = createUseStep(ctx); + const sleep = createSleep(ctx); + const createHook = createCreateHook(ctx); + const settleStep = useStep('settleStep'); + const recoverStep = useStep('recoverStep'); + const finalizeStep = useStep('finalizeStep'); + const releaseStep = useStep('releaseStep'); + const reconcileStep = useStep('reconcileStep'); + const cfg = { + rounds: 6, + width: 8, + watchdogMs: 2500, + betweenRoundSleepMs: 1000, + reconcileBase: 2, + }; + return async () => { + const pokeHook = createHook({ token: `${TOKEN}:poke` }); + try { + for (let round = 0; round < cfg.rounds; round += 1) { + const branches = await Promise.all( + Array.from({ length: cfg.width }, (_, index) => + (async () => { + try { + const winner = await Promise.race([ + settleStep({ round, index }), + sleep(cfg.watchdogMs).then(() => WATCHDOG), + ]); + if (winner === WATCHDOG) { + await recoverStep({ round, index }); + await finalizeStep({ round, index, winner: 'watchdog' }); + return { winner: 'watchdog' as const }; + } + await finalizeStep({ round, index, winner: 'settled' }); + return { winner: 'settled' as const }; + } finally { + await releaseStep({ round, index }); + } + })() + ) + ); + const stragglers = branches.filter( + (b) => b.winner === 'watchdog' + ).length; + await Promise.all( + Array.from({ length: cfg.reconcileBase + stragglers }, (_, index) => + reconcileStep({ round, index, stragglers }) + ) + ); + if (cfg.betweenRoundSleepMs > 0) { + await sleep(cfg.betweenRoundSleepMs); + } + } + } finally { + pokeHook.dispose(); + } + }; +} + +// ~90s of replays; diagnostic tool rather than a regression test. Run with +// STORM_LOG_SWEEP=1 to reproduce the flip-point table in the PR description. +describe.skipIf(!process.env.STORM_LOG_SWEEP)( + 'prefix-length sweep over the corrupted storm log', + () => { + it('reports rank 197/198 bindings per prefix length', async () => { + const fixture: FixtureEvent[] = JSON.parse( + readFileSync( + join( + dirname(fileURLToPath(import.meta.url)), + '__fixtures__', + 'wrun-41KZYJ92TP-storm-log.json' + ), + 'utf8' + ) + ); + const maxRank = Math.max(...fixture.map((f) => f.r)); + const ulids = generateUlidSequence(maxRank + 8); + const rankOf = new Map(ulids.map((u, i) => [u, i])); + + const lines: string[] = []; + for (const len of [ + 605, 608, 610, 611, 612, 615, 617, 618, 619, 620, 621, 622, 623, 624, + 625, 630, 635, 640, 645, 650, 655, + ]) { + const prefix = fixture.filter((f) => f.s <= len); + const events = await buildEvents(prefix, ulids); + const ctx = setupWorkflowContext(events); + const discontinuation = withResolvers(); + ctx.onWorkflowError = discontinuation.reject; + let error: any; + try { + await Promise.race([workflowBody(ctx)(), discontinuation.promise]); + } catch (err) { + error = err; + } + const byRank: Record = {}; + for (const item of ctx.invocationsQueue.values()) { + if (item.type !== 'step') continue; + const r = rankOf.get(item.correlationId.split('_', 2)[1]); + if (r !== undefined && r >= 196 && r <= 200) { + byRank[r] = item.stepName; + } + } + lines.push( + `len=${len} outcome=${error?.constructor?.name ?? 'none'} ` + + `r196=${byRank[196] ?? '-'} r197=${byRank[197] ?? '-'} r198=${byRank[198] ?? '-'} r199=${byRank[199] ?? '-'} r200=${byRank[200] ?? '-'} ` + + `${error?.message?.slice(0, 110)?.replace(/\n/g, ' ') ?? ''}` + ); + } + // eslint-disable-next-line no-console + console.log(`\nSWEEP RESULTS\n${lines.join('\n')}`); + }, 240_000); + } +);