diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardActionClaim.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardActionClaim.java new file mode 100644 index 000000000..988bd3ba7 --- /dev/null +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardActionClaim.java @@ -0,0 +1,9 @@ +package com.bencodez.advancedcore.core.reward; + +/** Result of an atomic, durable attempt to claim one native reward action. */ +public enum SharedRewardActionClaim { + /** This caller persisted the pending action and may invoke it once. */ + STARTED, + /** A prior caller may have invoked the action; automatic replay is unsafe. */ + INDETERMINATE +} diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardContext.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardContext.java index 6f24331ab..12efb01e6 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardContext.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardContext.java @@ -10,6 +10,7 @@ public final class SharedRewardContext { private final UUID userId; private final String playerName; private final HashMap placeholders; + private volatile boolean keyedExecution; public SharedRewardContext(UUID userId, String playerName, Map placeholders) { this.userId = Objects.requireNonNull(userId, "userId"); @@ -32,4 +33,12 @@ public String playerName() { public HashMap placeholders() { return placeholders; } + + void markKeyedExecution() { + keyedExecution = true; + } + + boolean isKeyedExecution() { + return keyedExecution; + } } diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardIndeterminateException.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardIndeterminateException.java new file mode 100644 index 000000000..3d848b68f --- /dev/null +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardIndeterminateException.java @@ -0,0 +1,10 @@ +package com.bencodez.advancedcore.core.reward; + +/** A native action may have taken effect but has no durable completion acknowledgement. */ +public final class SharedRewardIndeterminateException extends IllegalStateException { + private static final long serialVersionUID = 1L; + + public SharedRewardIndeterminateException(String executionPath, int stepIndex) { + super("Native reward action requires reconciliation: " + executionPath + " step " + stepIndex); + } +} diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardKeyedDurability.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardKeyedDurability.java new file mode 100644 index 000000000..74fc80fa5 --- /dev/null +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardKeyedDurability.java @@ -0,0 +1,30 @@ +package com.bencodez.advancedcore.core.reward; + +import java.util.concurrent.CompletionStage; + +/** + * Durable action admission for one logical reward occurrence. The injected owner + * serializes the occurrence across processes. Its database is authoritative. + * Plans used with this adapter must have flat native steps; nested composite + * execution has no durable parent/child admission contract yet. + */ +public interface SharedRewardKeyedDurability extends SharedRewardDurability { + /** + * Atomically persist a pending action before any native side effect. Return + * STARTED to exactly one caller. If a prior pending action or a stale cursor + * exists, return INDETERMINATE without granting execution. An acknowledgement + * lost after committing STARTED must also return INDETERMINATE on retry. + * Validate the fingerprint and step index against the stored occurrence. + */ + CompletionStage claimAction(String executionPath, String fingerprint, + int stepIndex, SharedRewardContext context); + + /** + * {@inheritDoc} For keyed execution this must atomically advance the cursor + * and clear the matching pending action. A failure after the native effect + * leaves the action pending for explicit reconciliation, never auto-replay. + */ + @Override + CompletionStage checkpoint(String executionPath, String fingerprint, int completedSteps, + SharedRewardContext context); +} diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardOrchestrator.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardOrchestrator.java index cc3d7670f..be3b7320b 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardOrchestrator.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardOrchestrator.java @@ -27,7 +27,27 @@ public CompletionStage execute(SharedRewardPlan plan, Shared Objects.requireNonNull(plan, "plan"); Objects.requireNonNull(context, "context"); SharedRewardDurability replay = durability == null ? SharedRewardDurability.NONE : durability; - return execute(plan, context, replay, pathSegment(plan.id())); + return execute(plan, context, replay, pathSegment(plan.id()), null); + } + + /** + * Execute a prepared plan for one stable logical occurrence. The caller-owned + * durability adapter must atomically claim each native action before it runs; + * an uncertain earlier attempt fails closed for reconciliation. Native steps + * must be flattened; keyed nested composite plans are not supported. + */ + public CompletionStage executeKeyed(SharedRewardPlan plan, SharedRewardContext context, + SharedRewardKeyedDurability durability, String occurrenceKey) { + Objects.requireNonNull(plan, "plan"); + Objects.requireNonNull(context, "context"); + Objects.requireNonNull(durability, "durability"); + Objects.requireNonNull(occurrenceKey, "occurrenceKey"); + if (occurrenceKey.isBlank() || occurrenceKey.length() > 256) { + throw new IllegalArgumentException("Occurrence key must contain 1..256 characters"); + } + if (!durability.durable()) throw new IllegalArgumentException("Keyed rewards require durable action admission"); + context.markKeyedExecution(); + return execute(plan, context, durability, pathSegment(occurrenceKey) + "/" + pathSegment(plan.id()), durability); } public CompletionStage executeNested(SharedRewardPlan plan, SharedRewardContext context, @@ -37,18 +57,22 @@ public CompletionStage executeNested(SharedRewardPlan plan, Objects.requireNonNull(parentPath, "parentPath"); String segment = pathSegment(plan.id()); String path = parentPath.isBlank() ? segment : parentPath + "/" + segment; - return execute(plan, context, durability == null ? SharedRewardDurability.NONE : durability, path); + SharedRewardDurability replay = durability == null ? SharedRewardDurability.NONE : durability; + if (context.isKeyedExecution()) { + return failed("Keyed nested plans require an explicit composite-step contract; flatten native steps instead"); + } + return execute(plan, context, replay, path, null); } private CompletionStage execute(SharedRewardPlan plan, SharedRewardContext context, - SharedRewardDurability durability, String executionPath) { + SharedRewardDurability durability, String executionPath, SharedRewardKeyedDurability keyed) { try { if (platform.isShuttingDown()) return failed("Reward platform is shutting down"); String fingerprint = durability.durable() ? plan.fingerprint() : "non-durable"; SharedRewardProgress saved = durability.durable() ? durability.loadProgress(executionPath) : null; if (saved != null) { validateProgress(plan, fingerprint, saved, executionPath); - return continueFromProgress(plan, context, durability, executionPath, fingerprint, saved); + return continueFromProgress(plan, context, durability, executionPath, fingerprint, saved, keyed); } int legacyCursor = durability.completedSteps(executionPath); @@ -90,7 +114,7 @@ private CompletionStage execute(SharedRewardPlan plan, Share } return persisted.thenCompose(progress -> { validateProgress(plan, fingerprint, progress, executionPath); - return continueFromProgress(plan, context, durability, executionPath, fingerprint, progress); + return continueFromProgress(plan, context, durability, executionPath, fingerprint, progress, keyed); }); }); } catch (Throwable failure) { @@ -99,11 +123,12 @@ private CompletionStage execute(SharedRewardPlan plan, Share } private CompletionStage continueFromProgress(SharedRewardPlan plan, SharedRewardContext context, - SharedRewardDurability durability, String executionPath, String fingerprint, SharedRewardProgress progress) { + SharedRewardDurability durability, String executionPath, String fingerprint, SharedRewardProgress progress, + SharedRewardKeyedDurability keyed) { context.placeholders().clear(); context.placeholders().putAll(progress.placeholders()); if (!progress.eligible()) return CompletableFuture.completedFuture(SharedRewardResult.NOT_ELIGIBLE); - return executeEligible(plan, context, durability, executionPath, fingerprint, progress); + return executeEligible(plan, context, durability, executionPath, fingerprint, progress, keyed); } private void validateProgress(SharedRewardPlan plan, String fingerprint, SharedRewardProgress progress, @@ -121,10 +146,11 @@ private void validateProgress(SharedRewardPlan plan, String fingerprint, SharedR } private CompletionStage executeEligible(SharedRewardPlan plan, SharedRewardContext context, - SharedRewardDurability durability, String executionPath, String fingerprint, SharedRewardProgress progress) { + SharedRewardDurability durability, String executionPath, String fingerprint, SharedRewardProgress progress, + SharedRewardKeyedDurability keyed) { int resume = progress.completedSteps(); java.util.function.Supplier> work = - () -> executeSteps(plan, context, durability, executionPath, fingerprint, resume); + () -> executeSteps(plan, context, durability, executionPath, fingerprint, resume, keyed); Duration delay = Duration.ZERO; if (resume == 0 && !plan.delay().isZero()) { delay = durability.durable() ? Duration.between(platform.now(), progress.notBefore()) : plan.delay(); @@ -158,13 +184,14 @@ private CompletionStage evaluateRequirements(List executeSteps(SharedRewardPlan plan, SharedRewardContext context, - SharedRewardDurability durability, String executionPath, String fingerprint, int index) { + SharedRewardDurability durability, String executionPath, String fingerprint, int index, + SharedRewardKeyedDurability keyed) { CompletionStage chain = CompletableFuture.completedFuture(SharedRewardResult.COMPLETED); for (int current = index; current < plan.steps().size(); current++) { final int stepIndex = current; chain = chain.thenCompose(previous -> previous == SharedRewardResult.DEFERRED ? CompletableFuture.completedFuture(SharedRewardResult.DEFERRED) - : executeStep(plan.steps().get(stepIndex), context, durability, executionPath, fingerprint, stepIndex)); + : executeStep(plan.steps().get(stepIndex), context, durability, executionPath, fingerprint, stepIndex, keyed)); } return chain.thenCompose(result -> result != SharedRewardResult.DEFERRED && platform.isShuttingDown() ? failed("Reward platform shut down before execution completed") @@ -172,9 +199,29 @@ private CompletionStage executeSteps(SharedRewardPlan plan, } private CompletionStage executeStep(SharedRewardStep step, SharedRewardContext context, - SharedRewardDurability durability, String executionPath, String fingerprint, int index) { + SharedRewardDurability durability, String executionPath, String fingerprint, int index, + SharedRewardKeyedDurability keyed) { + if (keyed != null) { + if (!step.requiresOnlinePlayer()) { + return executeStepOnNative(step, context, durability, executionPath, fingerprint, index, keyed, false); + } + try { + CompletionStage availability = platform.checkActionAvailability(context.userId()); + if (availability == null) return failed("Reward platform returned null availability stage"); + return availability.thenCompose(online -> executeStepOnNative(step, context, durability, + executionPath, fingerprint, index, keyed, Boolean.TRUE.equals(online))); + } catch (Throwable failure) { + return CompletableFuture.failedFuture(failure); + } + } + return executeStepOnNative(step, context, durability, executionPath, fingerprint, index, null, false); + } + + private CompletionStage executeStepOnNative(SharedRewardStep step, SharedRewardContext context, + SharedRewardDurability durability, String executionPath, String fingerprint, int index, + SharedRewardKeyedDurability keyed, boolean onlineChecked) { if (platform.isShuttingDown()) return failed("Reward platform shut down before execution completed"); - if (step.requiresOnlinePlayer() && !platform.isOnline(context.userId())) { + if (step.requiresOnlinePlayer() && (keyed != null ? !onlineChecked : !platform.isOnline(context.userId()))) { if (!durability.durable()) return failed("Player became unavailable during non-durable reward step " + step.id()); CompletionStage deferred; try { @@ -188,6 +235,42 @@ private CompletionStage executeStep(SharedRewardStep step, S } String stepPath = executionPath + "/" + pathSegment(step.id()) + ":" + index; + if (keyed == null) return executeClaimedStep(step, context, durability, executionPath, fingerprint, index, stepPath, false); + CompletionStage claim; + try { + claim = keyed.claimAction(executionPath, fingerprint, index, context); + if (claim == null) return failed("Keyed reward adapter returned null action claim stage"); + } catch (Throwable failure) { + return CompletableFuture.failedFuture(failure); + } + return claim.thenCompose(result -> { + if (result == SharedRewardActionClaim.INDETERMINATE) { + return CompletableFuture.failedFuture(new SharedRewardIndeterminateException(executionPath, index)); + } + if (result != SharedRewardActionClaim.STARTED) return failed("Invalid keyed reward action claim"); + try { + CompletionStage dispatched = platform.runClaimedAction(context.userId(), + step.requiresOnlinePlayer(), () -> { + // Admission may have awaited storage while the server disabled + // or the player disconnected. Leave the claim for reconciliation. + if (platform.isShuttingDown() + || (step.requiresOnlinePlayer() && !platform.isOnline(context.userId()))) { + return CompletableFuture.failedFuture( + new SharedRewardIndeterminateException(executionPath, index)); + } + return executeClaimedStep(step, context, durability, executionPath, fingerprint, index, + stepPath, true); + }); + return dispatched == null ? failed("Reward platform returned null claimed action stage") : dispatched; + } catch (Throwable failure) { + return CompletableFuture.failedFuture(failure); + } + }); + } + + private CompletionStage executeClaimedStep(SharedRewardStep step, SharedRewardContext context, + SharedRewardDurability durability, String executionPath, String fingerprint, int index, String stepPath, + boolean claimedStrict) { CompletionStage action; try { action = step.action().execute(context, stepPath); @@ -200,7 +283,12 @@ private CompletionStage executeStep(SharedRewardStep step, S return action.thenCompose(result -> { if (result == null) return CompletableFuture.failedFuture( new IllegalStateException("Reward step returned null result: " + step.id())); - if (result == SharedRewardResult.DEFERRED) return CompletableFuture.completedFuture(SharedRewardResult.DEFERRED); + if (result == SharedRewardResult.DEFERRED) { + if (claimedStrict) { + return failed("Keyed reward action deferred after it was claimed: " + step.id()); + } + return CompletableFuture.completedFuture(SharedRewardResult.DEFERRED); + } CompletionStage checkpoint; try { checkpoint = durability.checkpoint(executionPath, fingerprint, index + 1, context); diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardPlatform.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardPlatform.java index f8bea6a17..cf63c383d 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardPlatform.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/reward/SharedRewardPlatform.java @@ -3,6 +3,7 @@ import java.time.Duration; import java.time.Instant; import java.util.UUID; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; import java.util.function.Supplier; @@ -25,4 +26,27 @@ CompletionStage delay(Duration delay, Supplier> operation); boolean isShuttingDown(); + + /** + * Check online availability on the player's owner thread before claiming a + * durable action. The returned stage must finish without waiting for the + * action or its claim. A false result permits durable offline deferral. + */ + default CompletionStage checkActionAvailability(UUID userId) { + return CompletableFuture.failedFuture(new UnsupportedOperationException( + "Keyed online rewards require an owner-thread availability check")); + } + + /** + * Run the post-claim state check and native action on the platform's owner + * thread or scheduler. When {@code requiresOnlinePlayer} is true, route to + * that player's entity/region owner. Keyed execution fails closed until an + * adapter provides this hook; legacy execution does not use it. The returned + * stage must cover the whole supplied operation, not just its scheduling. + */ + default CompletionStage runClaimedAction( + UUID userId, boolean requiresOnlinePlayer, Supplier> operation) { + return CompletableFuture.failedFuture(new UnsupportedOperationException( + "Keyed reward execution requires a native action scheduler")); + } } diff --git a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/rewards/SharedRewardKeyedRecoveryTest.java b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/rewards/SharedRewardKeyedRecoveryTest.java new file mode 100644 index 000000000..89ae384b0 --- /dev/null +++ b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/rewards/SharedRewardKeyedRecoveryTest.java @@ -0,0 +1,361 @@ +package com.bencodez.advancedcore.tests.rewards; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.time.Duration; +import java.time.Instant; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.CompletionStage; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Supplier; + +import org.junit.jupiter.api.Test; + +import com.bencodez.advancedcore.core.reward.SharedRewardActionClaim; +import com.bencodez.advancedcore.core.reward.SharedRewardContext; +import com.bencodez.advancedcore.core.reward.SharedRewardIndeterminateException; +import com.bencodez.advancedcore.core.reward.SharedRewardKeyedDurability; +import com.bencodez.advancedcore.core.reward.SharedRewardOrchestrator; +import com.bencodez.advancedcore.core.reward.SharedRewardPlan; +import com.bencodez.advancedcore.core.reward.SharedRewardPlatform; +import com.bencodez.advancedcore.core.reward.SharedRewardProgress; +import com.bencodez.advancedcore.core.reward.SharedRewardResult; +import com.bencodez.advancedcore.core.reward.SharedRewardStep; + +class SharedRewardKeyedRecoveryTest { + private static final UUID USER = UUID.fromString("fef273b7-aa45-42f9-ac14-cb047533afde"); + + @Test + void pendingClaimPrecedesNativeActionAndCheckpointClearsIt() { + Store store = new Store(); + AtomicInteger actions = new AtomicInteger(); + SharedRewardPlan plan = plan(new SharedRewardStep("command", false, (context, path) -> { + assertEquals(0, store.cursor("vote-a/vote")); + assertEquals(0, store.pending.get("vote-a/vote")); + actions.incrementAndGet(); + return done(); + })); + + assertEquals(SharedRewardResult.COMPLETED, execute(store, plan, "vote-a").join()); + assertEquals(1, actions.get()); + assertEquals(1, store.cursor("vote-a/vote")); + assertNull(store.pending.get("vote-a/vote")); + assertEquals(SharedRewardResult.COMPLETED, execute(store, plan, "vote-a").join()); + assertEquals(1, actions.get()); + } + + @Test + void concurrentSubmissionCannotRunAnAlreadyClaimedAction() { + Store store = new Store(); + AtomicInteger actions = new AtomicInteger(); + CompletableFuture effect = new CompletableFuture<>(); + SharedRewardPlan plan = plan(new SharedRewardStep("command", false, (context, path) -> { + actions.incrementAndGet(); + return effect; + })); + + CompletableFuture first = execute(store, plan, "vote-a"); + assertEquals(1, actions.get()); + CompletionException failure = assertThrows(CompletionException.class, + () -> execute(store, plan, "vote-a").join()); + assertInstanceOf(SharedRewardIndeterminateException.class, failure.getCause()); + assertEquals(1, actions.get()); + + effect.complete(SharedRewardResult.COMPLETED); + assertEquals(SharedRewardResult.COMPLETED, first.join()); + assertEquals(SharedRewardResult.COMPLETED, execute(store, plan, "vote-a").join()); + assertEquals(1, actions.get()); + } + + @Test + void keyedNestedPlanIsRejectedBeforeAnyChildSideEffect() { + Store store = new Store(); + Platform platform = new Platform(); + SharedRewardOrchestrator orchestrator = new SharedRewardOrchestrator(platform); + AtomicInteger nestedActions = new AtomicInteger(); + SharedRewardPlan nested = new SharedRewardPlan("child", 1, Duration.ZERO, List.of(), + List.of(step("item", nestedActions)), "child-config-v1"); + SharedRewardPlan parent = plan(new SharedRewardStep("nested", false, + (context, path) -> orchestrator.executeNested(nested, context, store, path))); + + assertThrows(CompletionException.class, () -> execute(platform, store, parent, "vote-a").join()); + assertEquals(0, nestedActions.get()); + assertEquals(0, store.completedSteps("vote-a/vote/nested:0/child")); + assertNull(store.pending.get("vote-a/vote/nested:0/child")); + } + + @Test + void legacyNestedCallAcceptsAKeyedCapableAdapter() { + Store store = new Store(); + SharedRewardPlan empty = new SharedRewardPlan("child", 1, Duration.ZERO, List.of(), List.of(), "empty"); + + assertEquals(SharedRewardResult.COMPLETED, new SharedRewardOrchestrator(new Platform()) + .executeNested(empty, new SharedRewardContext(USER, "Ben", Map.of()), store, "legacy-parent") + .toCompletableFuture().join()); + } + + @Test + void disconnectWhileClaimIsPendingCannotStartOnlineAction() { + Store store = new Store(); + Platform platform = new Platform(); + AtomicInteger actions = new AtomicInteger(); + CompletableFuture claim = new CompletableFuture<>(); + store.delayedClaim = claim; + SharedRewardPlan plan = plan(new SharedRewardStep("item", true, (context, path) -> { + assertTrue(platform.nativeDispatch); + actions.incrementAndGet(); + return done(); + })); + + CompletableFuture execution = execute(platform, store, plan, "vote-a"); + platform.online = false; + claim.complete(SharedRewardActionClaim.STARTED); + + CompletionException failure = assertThrows(CompletionException.class, execution::join); + assertInstanceOf(SharedRewardIndeterminateException.class, failure.getCause()); + assertEquals(0, actions.get()); + assertEquals(0, store.pending.get("vote-a/vote")); + } + + @Test + void shutdownWhileClaimIsPendingCannotStartNativeAction() { + Store store = new Store(); + Platform platform = new Platform(); + AtomicInteger actions = new AtomicInteger(); + CompletableFuture claim = new CompletableFuture<>(); + store.delayedClaim = claim; + SharedRewardPlan plan = plan(step("command", actions)); + + CompletableFuture execution = execute(platform, store, plan, "vote-a"); + platform.shuttingDown = true; + claim.complete(SharedRewardActionClaim.STARTED); + + CompletionException failure = assertThrows(CompletionException.class, execution::join); + assertInstanceOf(SharedRewardIndeterminateException.class, failure.getCause()); + assertEquals(0, actions.get()); + assertEquals(0, store.pending.get("vote-a/vote")); + assertEquals(USER, platform.lastDispatchUser); + assertEquals(false, platform.lastDispatchRequiresOnline); + } + + @Test + void completedNativeEffectWithFailedCheckpointIsIndeterminateOnRetry() { + Store store = new Store(); + store.failCheckpoint = true; + AtomicInteger actions = new AtomicInteger(); + SharedRewardPlan plan = plan(step("money", actions)); + + assertThrows(CompletionException.class, () -> execute(store, plan, "vote-a").join()); + assertEquals(1, actions.get()); + assertEquals(0, store.cursor("vote-a/vote")); + assertEquals(0, store.pending.get("vote-a/vote")); + + Store restarted = store.reopen(); + CompletionException failure = assertThrows(CompletionException.class, + () -> execute(restarted, plan, "vote-a").join()); + assertInstanceOf(SharedRewardIndeterminateException.class, failure.getCause()); + assertEquals(1, actions.get()); + } + + @Test + void lostCheckpointAcknowledgementRecoversCommittedPrefixWithoutRepeatingAction() { + Store store = new Store(); + store.loseCheckpointAck = true; + AtomicInteger actions = new AtomicInteger(); + SharedRewardPlan plan = plan(step("command", actions)); + + assertThrows(CompletionException.class, () -> execute(store, plan, "vote-a").join()); + assertEquals(1, actions.get()); + assertEquals(1, store.cursor("vote-a/vote")); + assertNull(store.pending.get("vote-a/vote")); + assertEquals(SharedRewardResult.COMPLETED, execute(store.reopen(), plan, "vote-a").join()); + assertEquals(1, actions.get()); + } + + @Test + void lostClaimAcknowledgementFailsClosedWithoutRunningAction() { + Store store = new Store(); + store.loseClaimAck = true; + AtomicInteger actions = new AtomicInteger(); + SharedRewardPlan plan = plan(step("item", actions)); + + assertThrows(CompletionException.class, () -> execute(store, plan, "vote-a").join()); + assertEquals(0, actions.get()); + assertEquals(0, store.pending.get("vote-a/vote")); + CompletionException failure = assertThrows(CompletionException.class, + () -> execute(store.reopen(), plan, "vote-a").join()); + assertInstanceOf(SharedRewardIndeterminateException.class, failure.getCause()); + assertEquals(0, actions.get()); + } + + @Test + void offlineDeferralDoesNotClaimActionAndDistinctOccurrencesDoNotCollide() { + Store store = new Store(); + Platform platform = new Platform(); + platform.online = false; + AtomicInteger actions = new AtomicInteger(); + SharedRewardPlan plan = plan(new SharedRewardStep("item", true, (context, path) -> { + assertTrue(platform.nativeDispatch); + actions.incrementAndGet(); + return done(); + })); + + assertEquals(SharedRewardResult.DEFERRED, execute(platform, store, plan, "vote-a").join()); + assertNull(store.pending.get("vote-a/vote")); + platform.online = true; + assertEquals(SharedRewardResult.COMPLETED, execute(platform, store, plan, "vote-a").join()); + assertEquals(SharedRewardResult.COMPLETED, execute(platform, store, plan, "vote-b").join()); + assertEquals(2, actions.get()); + assertEquals(5, platform.nativeDispatches); + assertEquals(2, platform.claimedDispatches); + assertEquals(USER, platform.lastDispatchUser); + assertTrue(platform.lastDispatchRequiresOnline); + } + + @Test + void completedPrefixIsSkippedAndChangedPlanFailsBeforeSideEffects() { + Store store = new Store(); + AtomicInteger first = new AtomicInteger(), second = new AtomicInteger(); + SharedRewardPlan original = plan(step("first", first), step("second", second)); + store.failAtStep = 1; + assertThrows(CompletionException.class, () -> execute(store, original, "vote-a").join()); + assertEquals(1, first.get()); + assertEquals(0, second.get()); + assertEquals(1, store.cursor("vote-a/vote")); + store.failAtStep = -1; + SharedRewardPlan changed = new SharedRewardPlan("vote", 1, Duration.ZERO, List.of(), + List.of(step("first", first), step("second", second)), "changed-config"); + assertThrows(CompletionException.class, () -> execute(store, changed, "vote-a").join()); + assertEquals(0, second.get()); + assertEquals(SharedRewardResult.COMPLETED, execute(store, original, "vote-a").join()); + assertEquals(1, first.get()); + assertEquals(1, second.get()); + } + + private static SharedRewardPlan plan(SharedRewardStep... steps) { + return new SharedRewardPlan("vote", 1, Duration.ZERO, List.of(), List.of(steps), "config-v1"); + } + + private static SharedRewardStep step(String id, AtomicInteger calls) { + return new SharedRewardStep(id, false, (context, path) -> { + calls.incrementAndGet(); + return done(); + }); + } + + private static CompletionStage done() { + return CompletableFuture.completedFuture(SharedRewardResult.COMPLETED); + } + + private static CompletableFuture execute(Store store, SharedRewardPlan plan, String key) { + return execute(new Platform(), store, plan, key); + } + + private static CompletableFuture execute(Platform platform, Store store, + SharedRewardPlan plan, String key) { + return new SharedRewardOrchestrator(platform) + .executeKeyed(plan, new SharedRewardContext(USER, "Ben", Map.of()), store, key).toCompletableFuture(); + } + + private static final class Platform implements SharedRewardPlatform { + boolean online = true; + boolean shuttingDown; + boolean nativeDispatch; + int nativeDispatches; + int claimedDispatches; + UUID lastDispatchUser; + boolean lastDispatchRequiresOnline; + public Instant now() { return Instant.EPOCH; } + public boolean isOnline(UUID uuid) { + assertTrue(nativeDispatch); + return online; + } + public boolean isShuttingDown() { return shuttingDown; } + public double nextChanceRoll() { return 0; } + public CompletionStage delay(Duration delay, + Supplier> work) { return work.get(); } + public CompletionStage checkActionAvailability(UUID userId) { + assertFalse(nativeDispatch); + nativeDispatches++; + nativeDispatch = true; + try { return CompletableFuture.completedFuture(isOnline(userId)); } + finally { nativeDispatch = false; } + } + public CompletionStage runClaimedAction( + UUID userId, boolean requiresOnlinePlayer, Supplier> operation) { + assertFalse(nativeDispatch); + nativeDispatches++; + claimedDispatches++; + lastDispatchUser = userId; + lastDispatchRequiresOnline = requiresOnlinePlayer; + nativeDispatch = true; + try { return operation.get(); } + finally { nativeDispatch = false; } + } + } + + /** Simulates a caller-owned durable row that survives constructing a new adapter. */ + private static final class Store implements SharedRewardKeyedDurability { + final Map progress; + final Map pending; + boolean failCheckpoint, loseClaimAck, loseCheckpointAck; + CompletableFuture delayedClaim; + int failAtStep = -1; + + Store() { this(new HashMap<>(), new HashMap<>()); } + Store(Map progress, Map pending) { + this.progress = progress; + this.pending = pending; + } + Store reopen() { return new Store(progress, pending); } + synchronized int cursor(String path) { return progress.get(path).completedSteps(); } + public boolean durable() { return true; } + public synchronized int completedSteps(String path) { return progress.containsKey(path) ? cursor(path) : 0; } + public synchronized SharedRewardProgress loadProgress(String path) { return progress.get(path); } + public synchronized CompletionStage begin(String path, SharedRewardProgress proposed) { + progress.putIfAbsent(path, proposed); + return CompletableFuture.completedFuture(progress.get(path)); + } + public synchronized CompletionStage claimAction(String path, String fingerprint, + int index, SharedRewardContext context) { + SharedRewardProgress state = progress.get(path); + if (state == null || !fingerprint.equals(state.planFingerprint()) || state.completedSteps() != index) { + return CompletableFuture.completedFuture(SharedRewardActionClaim.INDETERMINATE); + } + if (failAtStep == index) return CompletableFuture.failedFuture(new IllegalStateException("claim failed")); + if (pending.putIfAbsent(path, index) != null) { + return CompletableFuture.completedFuture(SharedRewardActionClaim.INDETERMINATE); + } + if (delayedClaim != null) return delayedClaim; + if (loseClaimAck) return CompletableFuture.failedFuture(new IllegalStateException("lost claim acknowledgement")); + return CompletableFuture.completedFuture(SharedRewardActionClaim.STARTED); + } + public synchronized CompletionStage checkpoint(String path, String fingerprint, + int nextStep, SharedRewardContext context) { + if (failCheckpoint) return CompletableFuture.failedFuture(new IllegalStateException("checkpoint failed")); + if (pending.get(path) == null || pending.get(path) != nextStep - 1) { + return CompletableFuture.failedFuture(new IllegalStateException("missing pending action")); + } + progress.put(path, progress.get(path).advance(nextStep, context)); + pending.remove(path); + if (loseCheckpointAck) return CompletableFuture.failedFuture(new IllegalStateException("lost checkpoint acknowledgement")); + return CompletableFuture.completedFuture(null); + } + public CompletionStage checkpoint(String path, int nextStep, SharedRewardContext context) { + return checkpoint(path, progress.get(path).planFingerprint(), nextStep, context); + } + public CompletionStage defer(String path, int nextStep, SharedRewardContext context) { + return CompletableFuture.completedFuture(null); + } + } +}