diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteIdentity.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteIdentity.java new file mode 100644 index 000000000..b5f835317 --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteIdentity.java @@ -0,0 +1,14 @@ +package com.bencodez.votingplugin.core.vote; + +import java.util.Objects; +import java.util.UUID; + +public record SharedVoteIdentity(UUID uuid, String playerName, boolean online) { + public SharedVoteIdentity { + Objects.requireNonNull(uuid, "uuid"); + Objects.requireNonNull(playerName, "playerName"); + if (playerName.isBlank()) { + throw new IllegalArgumentException("playerName cannot be blank"); + } + } +} diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteIdentityResolver.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteIdentityResolver.java new file mode 100644 index 000000000..bbb95d2f4 --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteIdentityResolver.java @@ -0,0 +1,9 @@ +package com.bencodez.votingplugin.core.vote; + +import java.util.concurrent.CompletionStage; + +/** Adapter to the existing AdvancedCore/user identity services. */ +@FunctionalInterface +public interface SharedVoteIdentityResolver { + CompletionStage resolve(SharedVoteInput input); +} diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteInput.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteInput.java new file mode 100644 index 000000000..946bc13de --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteInput.java @@ -0,0 +1,54 @@ +package com.bencodez.votingplugin.core.vote; + +import java.util.Objects; +import java.util.UUID; + +/** + * An already-accepted vote entering the shared processing path. Native ingress, + * proxy/global duplicate filtering, service-site validation and security checks + * remain upstream and are intentionally not reimplemented here. + */ +public record SharedVoteInput(UUID voteId, String playerName, String serviceSite, long voteTime, + boolean realVote, boolean addTotals, boolean proxyVote, boolean forceProxyRouting, boolean wasOnline) { + public SharedVoteInput { + Objects.requireNonNull(voteId, "voteId"); + Objects.requireNonNull(playerName, "playerName"); + Objects.requireNonNull(serviceSite, "serviceSite"); + if (playerName.isBlank()) throw new IllegalArgumentException("playerName cannot be blank"); + if (serviceSite.isBlank()) throw new IllegalArgumentException("serviceSite cannot be blank"); + if (voteTime < 0) throw new IllegalArgumentException("voteTime cannot be negative"); + } + + /** + * Compatibility constructor for callers that historically had one proxy bit. + * Native adapters should use the full constructor and map isBungee() and + * isForceBungee() independently. + */ + public SharedVoteInput(UUID voteId, String playerName, String serviceSite, long voteTime, + boolean realVote, boolean addTotals, boolean proxyVote, boolean wasOnline) { + this(voteId, playerName, serviceSite, voteTime, realVote, addTotals, + proxyVote, proxyVote, wasOnline); + } + + /** PlayerVoteEvent uses zero as "now"; normalize before durable mutation/receipt creation. */ + public SharedVoteInput normalizedVoteTime(long nowEpochMillis) { + if (voteTime != 0) return this; + if (nowEpochMillis <= 0) throw new IllegalArgumentException("normalized vote time must be positive"); + return new SharedVoteInput(voteId, playerName, serviceSite, nowEpochMillis, + realVote, addTotals, proxyVote, forceProxyRouting, wasOnline); + } + + /** A retry carrying the zero sentinel still identifies its already-normalized persisted receipt. */ + public boolean matchesPersisted(SharedVoteInput persisted) { + if (persisted == null) return false; + return voteId.equals(persisted.voteId()) + && playerName.equals(persisted.playerName()) + && serviceSite.equals(persisted.serviceSite()) + && (voteTime == 0 || voteTime == persisted.voteTime()) + && realVote == persisted.realVote() + && addTotals == persisted.addTotals() + && proxyVote == persisted.proxyVote() + && forceProxyRouting == persisted.forceProxyRouting() + && wasOnline == persisted.wasOnline(); + } +} diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteMutation.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteMutation.java new file mode 100644 index 000000000..bec276506 --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteMutation.java @@ -0,0 +1,16 @@ +package com.bencodez.votingplugin.core.vote; + +import java.util.Objects; +import java.util.UUID; + +/** + * Describes the logical vote mutation. The AdvancedCore-facing storage adapter + * owns actual keys, cache/queue ordering, point hooks/caps, and persistence. + */ +public record SharedVoteMutation(UUID voteId, String serviceSite, long voteTime, + boolean countTotals, boolean awardConfiguredPoints) { + public SharedVoteMutation { + Objects.requireNonNull(voteId, "voteId"); + Objects.requireNonNull(serviceSite, "serviceSite"); + } +} diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVotePolicy.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVotePolicy.java new file mode 100644 index 000000000..824bcf9e2 --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVotePolicy.java @@ -0,0 +1,29 @@ +package com.bencodez.votingplugin.core.vote; + +/** + * Platform-neutral subset of the existing vote-processing configuration. + * Native adapters should populate it from the current VotingPlugin config. + */ +public record SharedVotePolicy(boolean countFakeVotes, boolean addTotals, + boolean addTotalsOffline, boolean processRewards, boolean giveOfflineRewards) { + + boolean shouldApplyConfiguredVoteMutation(SharedVoteInput input) { + return input.addTotals() && (input.realVote() || countFakeVotes); + } + + boolean shouldCountTotals(SharedVoteInput input, boolean online) { + return shouldApplyConfiguredVoteMutation(input) && addTotals && (addTotalsOffline || online); + } + + boolean shouldAwardConfiguredPoints(SharedVoteInput input) { + // Existing Bukkit behavior awards configured vote points when the accepted + // vote is countable even if Config.AddTotals itself is disabled. + return shouldApplyConfiguredVoteMutation(input); + } + + boolean shouldExecuteRewardsNow(SharedVoteInput input, boolean online) { + // Proxy votes preserve the existing force-processing behavior. For native + // votes, ProcessRewards and per-site offline eligibility remain authoritative. + return input.proxyVote() || (processRewards && (online || giveOfflineRewards)); + } +} diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessingResult.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessingResult.java new file mode 100644 index 000000000..4262fc005 --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessingResult.java @@ -0,0 +1,17 @@ +package com.bencodez.votingplugin.core.vote; + +import java.util.Objects; + +public record SharedVoteProcessingResult(SharedVoteIdentity identity, SharedVoteUserSnapshot persistedState, + RewardDisposition rewardDisposition) { + public SharedVoteProcessingResult { + Objects.requireNonNull(identity, "identity"); + Objects.requireNonNull(persistedState, "persistedState"); + Objects.requireNonNull(rewardDisposition, "rewardDisposition"); + } + + public enum RewardDisposition { + EXECUTED, + DEFERRED + } +} diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessor.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessor.java new file mode 100644 index 000000000..32ddcac44 --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessor.java @@ -0,0 +1,129 @@ +package com.bencodez.votingplugin.core.vote; + +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Objects; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionStage; +import java.util.function.Supplier; + +import com.bencodez.votingplugin.core.vote.SharedVoteProcessingResult.RewardDisposition; + +public final class SharedVoteProcessor { + private static final int MAX_RECOVERY_BATCH = 100; + private final SharedVoteIdentityResolver identities; + private final SharedVoteUserServices users; + private final SharedVoteRewardServices rewards; + + public SharedVoteProcessor(SharedVoteIdentityResolver identities, SharedVoteUserServices users, + SharedVoteRewardServices rewards) { + this.identities = Objects.requireNonNull(identities, "identities"); + this.users = Objects.requireNonNull(users, "users"); + this.rewards = Objects.requireNonNull(rewards, "rewards"); + } + + public CompletionStage process(SharedVoteInput input, SharedVotePolicy policy) { + Objects.requireNonNull(input, "input"); + Objects.requireNonNull(policy, "policy"); + return call(() -> users.findVote(input.voteId()), "receipt lookup").thenCompose(existing -> { + if (existing != null) { + existing.requireInput(input); + return deliver(existing); + } + SharedVoteInput normalized = input.normalizedVoteTime(System.currentTimeMillis()); + return call(() -> identities.resolve(normalized), "identity resolution").thenCompose(identity -> { + if (identity == null) return failed("Identity resolver returned null identity"); + boolean currentOnline = identity.online(); + boolean rewardOnline = normalized.proxyVote() ? normalized.wasOnline() : currentOnline; + SharedVoteMutation mutation = new SharedVoteMutation(normalized.voteId(), normalized.serviceSite(), + normalized.voteTime(), policy.shouldCountTotals(normalized, currentOnline), + policy.shouldAwardConfiguredPoints(normalized)); + boolean executeNow = policy.shouldExecuteRewardsNow(normalized, rewardOnline); + return call(() -> rewards.prepareVoteRewards(normalized, identity, executeNow), "reward preparation") + .thenCompose(rewardPlan -> { + if (rewardPlan == null) return failed("Reward services returned null prepared reward plan"); + // Preserve the ingress sentinel for the atomic uniqueness check. The + // mutation carries this caller's normalized candidate; if another + // concurrent caller wins, the returned receipt is authoritative. + return call(() -> users.persistVoteWithReward(input, identity, mutation, executeNow, rewardPlan), + "atomic persistence").thenCompose(receipt -> { + if (receipt == null) return failed("User services returned null vote receipt"); + receipt.requireInput(input); + if (!receipt.identity().uuid().equals(identity.uuid())) { + return failed("Vote receipt belongs to a different resolved identity"); + } + return deliver(receipt); + }); + }); + }); + }); + } + + public CompletionStage recover(UUID voteId) { + Objects.requireNonNull(voteId, "voteId"); + return call(() -> users.findVote(voteId), "receipt lookup").thenCompose(receipt -> { + if (receipt == null || !voteId.equals(receipt.input().voteId())) return failed("Pending vote receipt was not found"); + return deliver(receipt); + }); + } + + public CompletionStage> recoverPending(int limit) { + if (limit < 1 || limit > MAX_RECOVERY_BATCH) throw new IllegalArgumentException("Recovery limit must be 1..100"); + return call(() -> users.pendingVotes(limit), "pending receipt scan").thenCompose(receipts -> { + if (receipts == null || receipts.size() > limit) return failed("Invalid pending receipt batch"); + List batch = List.copyOf(receipts); + HashSet ids = new HashSet<>(); + for (SharedVoteReceipt receipt : batch) { + if (!receipt.pending() || !ids.add(receipt.input().voteId())) return failed("Invalid pending receipt entry"); + } + List results = new ArrayList<>(); + List failures = new ArrayList<>(); + CompletionStage chain = CompletableFuture.completedFuture(null); + for (SharedVoteReceipt receipt : batch) { + chain = chain.thenCompose(ignored -> deliver(receipt).handle((result, failure) -> { + if (failure == null) results.add(result); else failures.add(failure); + return null; + })); + } + return chain.thenCompose(ignored -> { + if (failures.isEmpty()) return CompletableFuture.completedFuture(List.copyOf(results)); + IllegalStateException failure = new IllegalStateException("Some pending vote rewards could not be recovered"); + failures.forEach(failure::addSuppressed); + return CompletableFuture.failedFuture(failure); + }); + }); + } + + private CompletionStage deliver(SharedVoteReceipt receipt) { + if (!receipt.pending()) return CompletableFuture.completedFuture(result(receipt)); + return call(() -> rewards.deliverOnce(receipt), "keyed reward delivery").thenCompose(disposition -> { + if (disposition == null) return failed("Reward services returned null delivery disposition"); + return call(() -> users.markRewardCompleted(receipt.input().voteId(), disposition), "reward acknowledgement") + .thenApply(completed -> { + if (completed == null) throw new IllegalStateException("Missing acknowledged vote receipt"); + completed.requireSameOrigin(receipt); + if (completed.completedDisposition() != disposition) { + throw new IllegalStateException("Vote reward acknowledgement did not preserve the delivery result"); + } + return result(completed); + }); + }); + } + + private static SharedVoteProcessingResult result(SharedVoteReceipt receipt) { + return new SharedVoteProcessingResult(receipt.identity(), receipt.persistedState(), receipt.completedDisposition()); + } + + private static CompletionStage call(Supplier> operation, String name) { + try { + CompletionStage stage = operation.get(); + return stage == null ? failed("Adapter returned null stage for " + name) : stage; + } catch (Throwable failure) { return CompletableFuture.failedFuture(failure); } + } + + private static CompletionStage failed(String message) { + return CompletableFuture.failedFuture(new IllegalStateException(message)); + } +} diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteReceipt.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteReceipt.java new file mode 100644 index 000000000..492b25666 --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteReceipt.java @@ -0,0 +1,53 @@ +package com.bencodez.votingplugin.core.vote; + +import java.util.Objects; + +import com.bencodez.votingplugin.core.vote.SharedVoteProcessingResult.RewardDisposition; + +/** + * Immutable receipt keyed by voteId in the existing persistence owner. Its initial + * form is committed atomically with totals/points, pending reward intent and the + * immutable prepared reward version. Across retries only the acknowledgement changes. + */ +public record SharedVoteReceipt(SharedVoteInput input, SharedVoteIdentity identity, SharedVoteMutation mutation, + SharedVoteUserSnapshot persistedState, boolean executeRewardsNow, SharedVoteRewardPlan rewardPlan, + RewardDisposition completedDisposition) { + public SharedVoteReceipt { + Objects.requireNonNull(input, "input"); + Objects.requireNonNull(identity, "identity"); + Objects.requireNonNull(mutation, "mutation"); + Objects.requireNonNull(persistedState, "persistedState"); + Objects.requireNonNull(rewardPlan, "rewardPlan"); + if (input.voteTime() <= 0) throw new IllegalArgumentException("Persisted vote time must be normalized"); + if (!input.voteId().equals(mutation.voteId()) || !input.serviceSite().equals(mutation.serviceSite()) + || input.voteTime() != mutation.voteTime()) { + throw new IllegalArgumentException("Vote receipt input and mutation do not match"); + } + } + + public boolean pending() { return completedDisposition == null; } + + public SharedVoteReceipt completed(RewardDisposition disposition) { + Objects.requireNonNull(disposition, "disposition"); + if (!pending() && disposition != completedDisposition) { + throw new IllegalStateException("A completed vote receipt cannot change its reward result"); + } + return new SharedVoteReceipt(input, identity, mutation, persistedState, executeRewardsNow, rewardPlan, disposition); + } + + public void requireInput(SharedVoteInput expected) { + Objects.requireNonNull(expected, "expected"); + if (!expected.matchesPersisted(input)) { + throw new IllegalStateException("voteId is already bound to different vote input"); + } + } + + public void requireSameOrigin(SharedVoteReceipt expected) { + requireInput(expected.input()); + if (!identity.equals(expected.identity()) || !mutation.equals(expected.mutation()) + || !persistedState.equals(expected.persistedState()) || executeRewardsNow != expected.executeRewardsNow() + || !rewardPlan.equals(expected.rewardPlan())) { + throw new IllegalStateException("Vote receipt changed while acknowledging reward delivery"); + } + } +} diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteRewardPlan.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteRewardPlan.java new file mode 100644 index 000000000..840f94db3 --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteRewardPlan.java @@ -0,0 +1,17 @@ +package com.bencodez.votingplugin.core.vote; + +import java.util.Objects; + +/** + * Immutable reference to the exact reward definition prepared for one vote. + * The version reference must continue resolving to the same archived/snapshotted + * definition after configuration reloads and process restarts. + */ +public record SharedVoteRewardPlan(String planId, String versionReference) { + public SharedVoteRewardPlan { + Objects.requireNonNull(planId, "planId"); + Objects.requireNonNull(versionReference, "versionReference"); + if (planId.isBlank()) throw new IllegalArgumentException("planId cannot be blank"); + if (versionReference.isBlank()) throw new IllegalArgumentException("versionReference cannot be blank"); + } +} diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteRewardServices.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteRewardServices.java new file mode 100644 index 000000000..3fb9c6a3e --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteRewardServices.java @@ -0,0 +1,37 @@ +package com.bencodez.votingplugin.core.vote; + +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionStage; + +import com.bencodez.votingplugin.core.vote.SharedVoteProcessingResult.RewardDisposition; + +/** Adapter for AdvancedCore reward orchestration and its existing durable replay owner. */ +public interface SharedVoteRewardServices { + CompletionStage executeVoteRewards(SharedVoteInput input, SharedVoteIdentity identity, + SharedVoteUserSnapshot persistedState); + + CompletionStage deferVoteRewards(SharedVoteInput input, SharedVoteIdentity identity, + SharedVoteUserSnapshot persistedState); + + /** + * Resolve and freeze the exact reward configuration BEFORE the vote transaction. + * versionReference must resolve to this same definition after reload/restart. + */ + default CompletionStage prepareVoteRewards(SharedVoteInput input, + SharedVoteIdentity identity, boolean executeRewardsNow) { + return CompletableFuture.failedFuture(new UnsupportedOperationException( + "Shared vote rewards require a persistable prepared reward version")); + } + + /** + * Admit/deduplicate by receipt.input().voteId() in the existing reward owner. + * Use receipt.rewardPlan() rather than current configuration. Serialize concurrent + * delivery/recovery for that occurrence, preserve step checkpoints and durably + * remember the terminal disposition. A repeated call resumes pending work or + * returns the same result without redoing completed work. + */ + default CompletionStage deliverOnce(SharedVoteReceipt receipt) { + return CompletableFuture.failedFuture(new UnsupportedOperationException( + "Shared vote rewards require keyed durable replay support")); + } +} diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteUserServices.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteUserServices.java new file mode 100644 index 000000000..899bbc677 --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteUserServices.java @@ -0,0 +1,44 @@ +package com.bencodez.votingplugin.core.vote; + +import java.util.List; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionStage; + +import com.bencodez.votingplugin.core.vote.SharedVoteProcessingResult.RewardDisposition; + +/** Port for AdvancedCore's shared user/cache/storage owner. */ +public interface SharedVoteUserServices { + CompletionStage persistVote(SharedVoteIdentity identity, SharedVoteMutation mutation); + CompletionStage load(UUID uuid); + + default CompletionStage findVote(UUID voteId) { return unsupported(); } + + /** + * One transaction atomically applies the mutation and inserts a pending receipt, + * including the prepared immutable reward version, or returns the existing receipt + * without applying totals/points again. Enforce unique voteId in that transaction. + */ + default CompletionStage persistVoteWithReward(SharedVoteInput input, + SharedVoteIdentity identity, SharedVoteMutation mutation, boolean executeRewardsNow, + SharedVoteRewardPlan rewardPlan) { + return unsupported(); + } + + /** Unsafe legacy shape intentionally has no mutation-only fallback. */ + default CompletionStage persistVoteWithReward(SharedVoteInput input, + SharedVoteIdentity identity, SharedVoteMutation mutation, boolean executeRewardsNow) { + return unsupported(); + } + + default CompletionStage markRewardCompleted(UUID voteId, RewardDisposition disposition) { + return unsupported(); + } + + default CompletionStage> pendingVotes(int limit) { return unsupported(); } + + private static CompletionStage unsupported() { + return CompletableFuture.failedFuture(new UnsupportedOperationException( + "Shared vote persistence requires atomic vote/receipt recovery support")); + } +} diff --git a/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteUserSnapshot.java b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteUserSnapshot.java new file mode 100644 index 000000000..c1ac9bace --- /dev/null +++ b/VotingPlugin/src/main/java/com/bencodez/votingplugin/core/vote/SharedVoteUserSnapshot.java @@ -0,0 +1,6 @@ +package com.bencodez.votingplugin.core.vote; + +/** Durable state returned after the vote mutation has been applied. */ +public record SharedVoteUserSnapshot(int allTimeTotal, int monthTotal, int weeklyTotal, + int dailyTotal, int points) { +} diff --git a/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessorEndToEndTest.java b/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessorEndToEndTest.java new file mode 100644 index 000000000..26c35b645 --- /dev/null +++ b/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessorEndToEndTest.java @@ -0,0 +1,125 @@ +package com.bencodez.votingplugin.core.vote; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.CompletionStage; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +class SharedVoteProcessorEndToEndTest { + @TempDir Path tempDir; + + @Test void votePersistsThenExecutesCommandAndMessageAcrossRestartAndOfflineVote() { + UUID uuid = UUID.randomUUID(); + Path storeFile = tempDir.resolve("users.properties"); + ArrayList rewards = new ArrayList<>(); + SharedVotePolicy policy = new SharedVotePolicy(false, true, true, true, true); + SharedVoteTestStore firstStore = new SharedVoteTestStore(storeFile, 10); + SharedVoteProcessor firstRuntime = new SharedVoteProcessor( + input -> CompletableFuture.completedFuture(new SharedVoteIdentity(uuid, "Ben", true)), firstStore, recording(firstStore, rewards)); + SharedVoteProcessingResult first = firstRuntime.process( + new SharedVoteInput(UUID.randomUUID(), "Ben", "ExampleSite", 1000L, true, true, false, true), policy).toCompletableFuture().join(); + assertEquals(new SharedVoteUserSnapshot(1, 1, 1, 1, 10), first.persistedState()); + assertEquals(SharedVoteProcessingResult.RewardDisposition.EXECUTED, first.rewardDisposition()); + assertTrue(Files.isRegularFile(storeFile)); + assertEquals(List.of("command:say Thanks Ben:1", "message:Thanks Ben:10"), rewards); + + SharedVoteTestStore restartedStore = new SharedVoteTestStore(storeFile, 10); + assertEquals(first.persistedState(), restartedStore.load(uuid).toCompletableFuture().join()); + SharedVoteProcessor restartedRuntime = new SharedVoteProcessor( + input -> CompletableFuture.completedFuture(new SharedVoteIdentity(uuid, "Ben", false)), restartedStore, recording(restartedStore, rewards)); + SharedVoteProcessingResult second = restartedRuntime.process( + new SharedVoteInput(UUID.randomUUID(), "Ben", "ExampleSite", 2000L, true, true, false, false), policy).toCompletableFuture().join(); + assertEquals(new SharedVoteUserSnapshot(2, 2, 2, 2, 20), second.persistedState()); + assertEquals(SharedVoteProcessingResult.RewardDisposition.EXECUTED, second.rewardDisposition()); + assertEquals(List.of("command:say Thanks Ben:1", "message:Thanks Ben:10", + "command:say Thanks Ben:2", "message:Thanks Ben:20"), rewards); + } + + @Test void offlineIneligibleRewardIsDurablyDelegatedAfterPersistence() { + UUID uuid = UUID.randomUUID(); + ArrayList rewards = new ArrayList<>(); + SharedVoteTestStore store = new SharedVoteTestStore(tempDir.resolve("deferred.properties"), 3); + SharedVoteProcessor runtime = new SharedVoteProcessor( + input -> CompletableFuture.completedFuture(new SharedVoteIdentity(uuid, "Ben", false)), store, recording(store, rewards)); + SharedVoteProcessingResult result = runtime.process( + new SharedVoteInput(UUID.randomUUID(), "Ben", "ExampleSite", 3000L, true, true, false, false), + new SharedVotePolicy(false, true, true, true, false)).toCompletableFuture().join(); + assertEquals(new SharedVoteUserSnapshot(1, 1, 1, 1, 3), result.persistedState()); + assertEquals(SharedVoteProcessingResult.RewardDisposition.DEFERRED, result.rewardDisposition()); + assertEquals(List.of("defer:ExampleSite:1"), rewards); + } + + @Test void proxyTotalsUseLiveStateInsteadOfHistoricalWasOnline() { + UUID uuid = UUID.randomUUID(); + SharedVotePolicy policy = new SharedVotePolicy(false, true, false, true, true); + SharedVoteTestStore offlineStore = new SharedVoteTestStore(tempDir.resolve("proxy-offline.properties"), 2); + SharedVoteProcessingResult offline = new SharedVoteProcessor( + input -> CompletableFuture.completedFuture(new SharedVoteIdentity(uuid, "Ben", false)), offlineStore, recording(offlineStore, new ArrayList<>())).process( + new SharedVoteInput(UUID.randomUUID(), "Ben", "ExampleSite", 3100L, true, true, true, true), policy).toCompletableFuture().join(); + assertEquals(new SharedVoteUserSnapshot(0, 0, 0, 0, 2), offline.persistedState()); + SharedVoteTestStore onlineStore = new SharedVoteTestStore(tempDir.resolve("proxy-online.properties"), 2); + SharedVoteProcessingResult online = new SharedVoteProcessor( + input -> CompletableFuture.completedFuture(new SharedVoteIdentity(uuid, "Ben", true)), onlineStore, recording(onlineStore, new ArrayList<>())).process( + new SharedVoteInput(UUID.randomUUID(), "Ben", "ExampleSite", 3200L, true, true, true, false), policy).toCompletableFuture().join(); + assertEquals(new SharedVoteUserSnapshot(1, 1, 1, 1, 2), online.persistedState()); + } + + @Test void persistenceFailurePreventsRewardExecution() { + ArrayList rewards = new ArrayList<>(); + SharedVoteTestStore failing = new SharedVoteTestStore(tempDir.resolve("failure.properties"), 1) { + @Override public CompletionStage persistVoteWithReward(SharedVoteInput input, + SharedVoteIdentity identity, SharedVoteMutation mutation, boolean execute, SharedVoteRewardPlan rewardPlan) { + return CompletableFuture.failedFuture(new IllegalStateException("storage failed")); + } + }; + SharedVoteProcessor runtime = new SharedVoteProcessor( + input -> CompletableFuture.completedFuture(new SharedVoteIdentity(UUID.randomUUID(), "Ben", true)), failing, recording(failing, rewards)); + CompletionException failure = assertThrows(CompletionException.class, () -> runtime.process( + new SharedVoteInput(UUID.randomUUID(), "Ben", "ExampleSite", 4000L, true, true, false, true), + new SharedVotePolicy(false, true, true, true, true)).toCompletableFuture().join()); + assertEquals("storage failed", failure.getCause().getMessage()); + assertTrue(rewards.isEmpty()); + } + + @Test void fakeVoteAndAddTotalsPolicyPreserveExistingCountingRules() { + UUID uuid = UUID.randomUUID(); + SharedVoteTestStore store = new SharedVoteTestStore(tempDir.resolve("policy.properties"), 5); + SharedVoteProcessor runtime = new SharedVoteProcessor( + input -> CompletableFuture.completedFuture(new SharedVoteIdentity(uuid, "Ben", true)), store, recording(store, new ArrayList<>())); + SharedVoteProcessingResult ignoredFake = runtime.process( + new SharedVoteInput(UUID.randomUUID(), "Ben", "ExampleSite", 5000L, false, true, false, true), + new SharedVotePolicy(false, true, true, true, true)).toCompletableFuture().join(); + assertEquals(new SharedVoteUserSnapshot(0, 0, 0, 0, 0), ignoredFake.persistedState()); + SharedVoteProcessingResult countedFake = runtime.process( + new SharedVoteInput(UUID.randomUUID(), "Ben", "ExampleSite", 6000L, false, true, false, true), + new SharedVotePolicy(true, false, true, true, true)).toCompletableFuture().join(); + assertEquals(new SharedVoteUserSnapshot(0, 0, 0, 0, 5), countedFake.persistedState()); + } + + private static SharedVoteRewardServices recording(SharedVoteTestStore store, List events) { + return new SharedVoteRewardServices() { + @Override public CompletionStage prepareVoteRewards(SharedVoteInput input, + SharedVoteIdentity identity, boolean executeRewardsNow) { + return store.prepareVoteRewards(input, identity, executeRewardsNow); + } + @Override public CompletionStage deliverOnce(SharedVoteReceipt receipt) { + return store.deliverOnce(receipt).thenApply(result -> { events.clear(); events.addAll(store.events()); return result; }); + } + @Override public CompletionStage executeVoteRewards(SharedVoteInput input, SharedVoteIdentity identity, + SharedVoteUserSnapshot state) { throw new AssertionError("Unkeyed execution"); } + @Override public CompletionStage deferVoteRewards(SharedVoteInput input, SharedVoteIdentity identity, + SharedVoteUserSnapshot state) { throw new AssertionError("Unkeyed deferral"); } + }; + } +} diff --git a/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessorRecoveryTest.java b/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessorRecoveryTest.java new file mode 100644 index 000000000..d9548d17a --- /dev/null +++ b/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteProcessorRecoveryTest.java @@ -0,0 +1,177 @@ +package com.bencodez.votingplugin.core.vote; + +import static org.junit.jupiter.api.Assertions.*; + +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.CompletionStage; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.api.io.TempDir; + +import com.bencodez.votingplugin.core.vote.SharedVoteProcessingResult.RewardDisposition; + +/** Restart and acknowledgement boundaries with test-only file transaction/replay owners. */ +@Timeout(15) +class SharedVoteProcessorRecoveryTest { + @TempDir Path directory; + private static final UUID USER = UUID.fromString("9837d441-a4d6-461f-aa86-17958c01bc8c"); + private static final SharedVotePolicy POLICY = new SharedVotePolicy(false, true, true, true, true); + + @Test void crashAfterCommitRecoversPendingRewardWithoutAnotherIngressVote() { + SharedVoteTestStore store = store(); + store.failAfterCommit = true; + SharedVoteInput input = input(); + assertThrows(CompletionException.class, () -> runtime(store).process(input, POLICY).toCompletableFuture().join()); + assertEquals(new SharedVoteUserSnapshot(1, 1, 1, 1, 5), store.load(USER).toCompletableFuture().join()); + assertTrue(store.events().isEmpty()); + SharedVoteTestStore reopened = store(); + SharedVoteProcessor recovered = new SharedVoteProcessor(vote -> { throw new AssertionError("Do not resolve a committed vote again"); }, reopened, reopened); + assertEquals(1, recovered.recoverPending(10).toCompletableFuture().join().size()); + assertEquals(2, reopened.events().size()); + assertEquals(1, reopened.load(USER).toCompletableFuture().join().allTimeTotal()); + assertTrue(reopened.pendingVotes(10).toCompletableFuture().join().isEmpty()); + } + + @Test void lostCommitAckRetryUsesOriginalIdentityPolicyAndSnapshot() { + SharedVoteTestStore store = store(); + store.failAfterCommit = true; + SharedVoteInput input = input(); + assertThrows(CompletionException.class, () -> runtime(store).process(input, POLICY).toCompletableFuture().join()); + SharedVoteTestStore reopened = store(); + SharedVoteProcessor retry = new SharedVoteProcessor(vote -> { throw new AssertionError("Do not re-resolve"); }, reopened, reopened); + SharedVoteProcessingResult result = retry.process(input, + new SharedVotePolicy(false, false, false, false, false)).toCompletableFuture().join(); + assertEquals(RewardDisposition.EXECUTED, result.rewardDisposition()); + assertEquals(new SharedVoteUserSnapshot(1, 1, 1, 1, 5), result.persistedState()); + } + + @Test void failureBeforeCommitCreatesNeitherMutationNorRewardWork() { + SharedVoteTestStore store = store(); + store.failBeforeCommit = true; + SharedVoteInput input = input(); + assertThrows(CompletionException.class, () -> runtime(store).process(input, POLICY).toCompletableFuture().join()); + assertEquals(new SharedVoteUserSnapshot(0, 0, 0, 0, 0), store.load(USER).toCompletableFuture().join()); + assertEquals(null, store.findVote(input.voteId()).toCompletableFuture().join()); + assertTrue(store.events().isEmpty()); + runtime(store).process(input, POLICY).toCompletableFuture().join(); + assertEquals(1, store.mutationCount); + } + + @Test void failedRewardDeliveryRemainsPendingForRestart() { + SharedVoteTestStore store = store(); + store.failBeforeDelivery = true; + SharedVoteInput input = input(); + assertThrows(CompletionException.class, () -> runtime(store).process(input, POLICY).toCompletableFuture().join()); + assertTrue(store.findVote(input.voteId()).toCompletableFuture().join().pending()); + SharedVoteTestStore reopened = store(); + runtime(reopened).recover(input.voteId()).toCompletableFuture().join(); + assertEquals(2, reopened.events().size()); + assertEquals(1, reopened.load(USER).toCompletableFuture().join().allTimeTotal()); + } + + @Test void lostRewardOwnerAcknowledgementDoesNotRepeatCompletedEffects() { + lostAck(true, false, false); + } + @Test void failedReceiptAcknowledgementDoesNotRepeatCompletedEffects() { + lostAck(false, true, false); + } + @Test void lostTerminalReceiptAcknowledgementDoesNotRepeatCompletedEffects() { + lostAck(false, false, true); + } + private void lostAck(boolean reward, boolean beforeMark, boolean afterMark) { + SharedVoteTestStore store = store(); + store.failAfterDelivery = reward; + store.failBeforeMark = beforeMark; + store.failAfterMark = afterMark; + SharedVoteInput input = input(); + assertThrows(CompletionException.class, () -> runtime(store).process(input, POLICY).toCompletableFuture().join()); + assertEquals(2, store.events().size()); + SharedVoteTestStore reopened = store(); + SharedVoteProcessingResult result = runtime(reopened).process(input, POLICY).toCompletableFuture().join(); + assertEquals(RewardDisposition.EXECUTED, result.rewardDisposition()); + assertEquals(2, reopened.events().size()); + assertEquals(1, reopened.load(USER).toCompletableFuture().join().allTimeTotal()); + } + + @Test void duplicateVoteIdWithDifferentInputIsRejectedBeforeAnotherMutation() { + SharedVoteTestStore store = store(); + SharedVoteInput input = input(); + runtime(store).process(input, POLICY).toCompletableFuture().join(); + SharedVoteInput conflict = new SharedVoteInput(input.voteId(), "Other", "Example", input.voteTime(), true, true, false, true); + assertThrows(CompletionException.class, () -> runtime(store).process(conflict, POLICY).toCompletableFuture().join()); + assertEquals(1, store.mutationCount); + assertEquals(2, store.events().size()); + } + + @Test void concurrentDuplicateSubmissionsShareThePersistenceAndRewardOwners() throws Exception { + SharedVoteTestStore store = store(); + SharedVoteInput input = input(); + var workers = Executors.newFixedThreadPool(4); + try { + List> tasks = new ArrayList<>(); + for (int i = 0; i < 12; i++) tasks.add(workers.submit(() -> runtime(store).process(input, POLICY).toCompletableFuture().join())); + for (var task : tasks) task.get(5, TimeUnit.SECONDS); + assertEquals(1, store.mutationCount); + assertEquals(2, store.events().size()); + } finally { + workers.shutdownNow(); + assertTrue(workers.awaitTermination(5, TimeUnit.SECONDS)); + } + } + + @Test void unacknowledgedAtomicPersistenceDoesNotTriggerTheOriginatingRewardCall() { + SharedVoteTestStore store = store(); + store.commitAck = new CompletableFuture<>(); + var result = runtime(store).process(input(), POLICY).toCompletableFuture(); + assertFalse(result.isDone()); + assertTrue(store.events().isEmpty()); + store.commitAck.complete(null); + assertEquals(RewardDisposition.EXECUTED, result.join().rewardDisposition()); + } + + @Test void offlineHandoffIsIdempotentAcrossRestartAndLostAcknowledgement() { + SharedVoteTestStore store = store(); + store.failBeforeMark = true; + SharedVoteInput input = input(); + SharedVoteProcessor offline = new SharedVoteProcessor(vote -> CompletableFuture.completedFuture(new SharedVoteIdentity(USER, "Ben", false)), store, store); + SharedVotePolicy defer = new SharedVotePolicy(false, true, true, true, false); + assertThrows(CompletionException.class, () -> offline.process(input, defer).toCompletableFuture().join()); + SharedVoteTestStore reopened = store(); + assertEquals(RewardDisposition.DEFERRED, runtime(reopened).recover(input.voteId()).toCompletableFuture().join().rewardDisposition()); + assertEquals(List.of("defer:Example:1"), reopened.events()); + } + + @Test void recoveryRejectsInvalidBoundsWithoutScanningStorage() { + SharedVoteTestStore store = store(); + assertThrows(IllegalArgumentException.class, () -> runtime(store).recoverPending(0)); + assertThrows(IllegalArgumentException.class, () -> runtime(store).recoverPending(101)); + } + + @Test void aFailedRecoveryEntryDoesNotBlockOtherPendingEntries() { + SharedVoteTestStore store = store(); + store.failBeforeDelivery = true; + SharedVoteInput first = input(); + assertThrows(CompletionException.class, () -> runtime(store).process(first, POLICY).toCompletableFuture().join()); + store.failBeforeDelivery = true; + SharedVoteInput second = input(); + assertThrows(CompletionException.class, () -> runtime(store).process(second, POLICY).toCompletableFuture().join()); + store.failBeforeDelivery = true; + assertThrows(CompletionException.class, () -> runtime(store).recoverPending(10).toCompletableFuture().join()); + assertEquals(1, store.pendingVotes(10).toCompletableFuture().join().size()); + assertEquals(2, store.events().size()); + } + + private SharedVoteTestStore store() { return new SharedVoteTestStore(directory.resolve("users.properties"), 5); } + private static SharedVoteProcessor runtime(SharedVoteTestStore store) { + return new SharedVoteProcessor(vote -> CompletableFuture.completedFuture(new SharedVoteIdentity(USER, "Ben", true)), store, store); + } + private static SharedVoteInput input() { return new SharedVoteInput(UUID.randomUUID(), "Ben", "Example", 1000, true, true, false, true); } +} diff --git a/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteReviewFollowupTest.java b/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteReviewFollowupTest.java new file mode 100644 index 000000000..c081ad233 --- /dev/null +++ b/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteReviewFollowupTest.java @@ -0,0 +1,126 @@ +package com.bencodez.votingplugin.core.vote; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.nio.file.Path; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +class SharedVoteReviewFollowupTest { + @TempDir Path directory; + private static final UUID USER = UUID.fromString("9837d441-a4d6-461f-aa86-17958c01bc8c"); + private static final SharedVotePolicy POLICY = new SharedVotePolicy(false, true, true, true, true); + + @Test + void zeroTimestampIsNormalizedOnceAndZeroSentinelRetryMatchesTheReceipt() { + SharedVoteTestStore store = new SharedVoteTestStore(directory.resolve("zero.properties"), 5); + SharedVoteProcessor processor = runtime(store); + UUID voteId = UUID.randomUUID(); + SharedVoteInput sentinel = new SharedVoteInput(voteId, "Ben", "Example", 0, true, true, false, true); + processor.process(sentinel, POLICY).toCompletableFuture().join(); + SharedVoteReceipt receipt = store.findVote(voteId).toCompletableFuture().join(); + assertTrue(receipt.input().voteTime() > 0); + long normalized = receipt.input().voteTime(); + processor.process(sentinel, POLICY).toCompletableFuture().join(); + assertEquals(1, store.mutationCount); + assertEquals(normalized, store.findVote(voteId).toCompletableFuture().join().input().voteTime()); + } + + @Test + void concurrentZeroSentinelCandidatesUseTheFirstAtomicReceipt() throws Exception { + SharedVoteTestStore store = new SharedVoteTestStore(directory.resolve("zero-race.properties"), 5); + UUID voteId = UUID.randomUUID(); + SharedVoteInput sentinel = new SharedVoteInput(voteId, "Ben", "Example", 0, true, true, false, true); + SharedVoteIdentity identity = new SharedVoteIdentity(USER, "Ben", true); + SharedVoteRewardPlan firstPlan = new SharedVoteRewardPlan("vote-site:Example", "v1"); + SharedVoteRewardPlan losingPlan = new SharedVoteRewardPlan("vote-site:Example", "v2"); + CountDownLatch ready = new CountDownLatch(2); + CountDownLatch start = new CountDownLatch(1); + var workers = Executors.newFixedThreadPool(2); + try { + CompletableFuture firstAttempt = CompletableFuture.supplyAsync(() -> { + awaitStart(ready, start); + return store.persistVoteWithReward(sentinel, identity, + new SharedVoteMutation(voteId, "Example", 1001, true, true), true, firstPlan) + .toCompletableFuture().join(); + }, workers); + CompletableFuture retryAttempt = CompletableFuture.supplyAsync(() -> { + awaitStart(ready, start); + return store.persistVoteWithReward(sentinel, identity, + new SharedVoteMutation(voteId, "Example", 2002, true, true), true, losingPlan) + .toCompletableFuture().join(); + }, workers); + assertTrue(ready.await(5, TimeUnit.SECONDS)); + start.countDown(); + SharedVoteReceipt first = firstAttempt.join(); + SharedVoteReceipt retry = retryAttempt.join(); + assertEquals(first, retry); + SharedVoteRewardPlan winningPlan = first.input().voteTime() == 1001 ? firstPlan : losingPlan; + assertEquals(winningPlan, first.rewardPlan()); + assertTrue(first.input().voteTime() == 1001 || first.input().voteTime() == 2002); + assertEquals(1, store.mutationCount); + } finally { + start.countDown(); + workers.shutdownNow(); + } + } + + @Test + void restartUsesThePersistedPreparedRewardVersionInsteadOfCurrentConfiguration() { + Path file = directory.resolve("plan.properties"); + SharedVoteTestStore first = new SharedVoteTestStore(file, 5); + first.preparedPlanVersion = "config-v1"; + first.failAfterCommit = true; + SharedVoteInput input = new SharedVoteInput(UUID.randomUUID(), "Ben", "Example", 1000, true, true, false, true); + assertThrows(CompletionException.class, () -> runtime(first).process(input, POLICY).toCompletableFuture().join()); + assertEquals("config-v1", first.findVote(input.voteId()).toCompletableFuture().join().rewardPlan().versionReference()); + + SharedVoteTestStore reopened = new SharedVoteTestStore(file, 5); + reopened.preparedPlanVersion = "config-v2"; + runtime(reopened).recover(input.voteId()).toCompletableFuture().join(); + assertEquals("config-v1", reopened.lastDeliveredPlanVersion); + assertEquals("config-v1", reopened.findVote(input.voteId()).toCompletableFuture().join().rewardPlan().versionReference()); + } + + @Test + void proxyOriginAndForcedRoutingRemainIndependentAcrossPersistence() { + SharedVoteTestStore store = new SharedVoteTestStore(directory.resolve("proxy-flags.properties"), 5); + UUID voteId = UUID.randomUUID(); + SharedVoteInput input = new SharedVoteInput(voteId, "Ben", "Example", 1234, + true, true, true, false, true); + SharedVoteIdentity identity = new SharedVoteIdentity(USER, "Ben", true); + SharedVoteReceipt receipt = store.persistVoteWithReward(input, identity, + new SharedVoteMutation(voteId, "Example", 1234, true, true), true, + new SharedVoteRewardPlan("vote-site:Example", "v1")).toCompletableFuture().join(); + assertTrue(receipt.input().proxyVote()); + assertFalse(receipt.input().forceProxyRouting()); + SharedVoteReceipt reopened = new SharedVoteTestStore(store.file, 5).findVote(voteId).toCompletableFuture().join(); + assertTrue(reopened.input().proxyVote()); + assertFalse(reopened.input().forceProxyRouting()); + } + + private static SharedVoteProcessor runtime(SharedVoteTestStore store) { + return new SharedVoteProcessor( + vote -> CompletableFuture.completedFuture(new SharedVoteIdentity(USER, "Ben", true)), store, store); + } + + private static void awaitStart(CountDownLatch ready, CountDownLatch start) { + ready.countDown(); + try { + if (!start.await(5, TimeUnit.SECONDS)) throw new IllegalStateException("Timed out awaiting concurrent start"); + } catch (InterruptedException interrupted) { + Thread.currentThread().interrupt(); + throw new IllegalStateException("Interrupted awaiting concurrent start", interrupted); + } + } +} diff --git a/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteTestStore.java b/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteTestStore.java new file mode 100644 index 000000000..aad9f531d --- /dev/null +++ b/VotingPlugin/src/test/java/com/bencodez/votingplugin/core/vote/SharedVoteTestStore.java @@ -0,0 +1,226 @@ +package com.bencodez.votingplugin.core.vote; + +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.StandardCopyOption; +import java.util.ArrayList; +import java.util.List; +import java.util.Properties; +import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionStage; + +import com.bencodez.votingplugin.core.vote.SharedVoteProcessingResult.RewardDisposition; + +/** TEST ONLY atomic file snapshot/reward owner. */ +class SharedVoteTestStore implements SharedVoteUserServices, SharedVoteRewardServices { + final Path file; + final Path rewardFile; + final int pointsPerVote; + boolean failBeforeCommit, failAfterCommit, failBeforeMark, failAfterMark; + boolean failBeforeDelivery, failAfterDelivery; + CompletableFuture commitAck = CompletableFuture.completedFuture(null); + int mutationCount; + String preparedPlanVersion = "fixture-v1"; + String lastDeliveredPlanVersion; + + SharedVoteTestStore(Path file, int pointsPerVote) { + this.file = file; + this.rewardFile = file.resolveSibling(file.getFileName() + ".rewards"); + this.pointsPerVote = pointsPerVote; + } + + @Override public CompletionStage persistVote(SharedVoteIdentity user, SharedVoteMutation mutation) { + return CompletableFuture.failedFuture(new AssertionError("Unsafe mutation-only API was called")); + } + + @Override public synchronized CompletionStage load(UUID uuid) { + try { return CompletableFuture.completedFuture(snapshot(read(file), uuid + ".")); } + catch (Exception failure) { return CompletableFuture.failedFuture(failure); } + } + + @Override public synchronized CompletionStage findVote(UUID id) { + try { return CompletableFuture.completedFuture(receipt(read(file), id)); } + catch (Exception failure) { return CompletableFuture.failedFuture(failure); } + } + + @Override + public CompletionStage prepareVoteRewards(SharedVoteInput input, + SharedVoteIdentity identity, boolean executeRewardsNow) { + return CompletableFuture.completedFuture( + new SharedVoteRewardPlan("vote-site:" + input.serviceSite(), preparedPlanVersion)); + } + + @Override public synchronized CompletionStage persistVoteWithReward(SharedVoteInput input, + SharedVoteIdentity identity, SharedVoteMutation mutation, boolean execute, SharedVoteRewardPlan rewardPlan) { + try { + Properties properties = read(file); + SharedVoteReceipt existing = receipt(properties, input.voteId()); + if (existing != null) { + existing.requireInput(input); + if (!existing.identity().uuid().equals(identity.uuid())) throw new IllegalStateException("Conflicting user"); + return CompletableFuture.completedFuture(existing); + } + if (!input.voteId().equals(mutation.voteId()) || !input.serviceSite().equals(mutation.serviceSite())) { + throw new IllegalStateException("Mutation origin does not match vote input"); + } + if (input.voteTime() != 0 && input.voteTime() != mutation.voteTime()) { + throw new IllegalStateException("Explicit vote timestamp changed before persistence"); + } + SharedVoteInput persistedInput = input.voteTime() == 0 + ? input.normalizedVoteTime(mutation.voteTime()) : input; + if (failBeforeCommit) { failBeforeCommit = false; throw new IOException("before commit"); } + String prefix = identity.uuid() + "."; + SharedVoteUserSnapshot previous = snapshot(properties, prefix); + int increment = mutation.countTotals() ? 1 : 0; + SharedVoteUserSnapshot next = new SharedVoteUserSnapshot(previous.allTimeTotal() + increment, + previous.monthTotal() + increment, previous.weeklyTotal() + increment, previous.dailyTotal() + increment, + previous.points() + (mutation.awardConfiguredPoints() ? pointsPerVote : 0)); + putSnapshot(properties, prefix, next); + SharedVoteReceipt created = new SharedVoteReceipt(persistedInput, identity, mutation, next, execute, rewardPlan, null); + putReceipt(properties, created); + write(file, properties); + mutationCount++; + if (failAfterCommit) { failAfterCommit = false; throw new IOException("lost commit acknowledgement"); } + return commitAck.thenApply(ignored -> created); + } catch (Exception failure) { return CompletableFuture.failedFuture(failure); } + } + + @Override public synchronized CompletionStage markRewardCompleted(UUID id, RewardDisposition disposition) { + try { + if (failBeforeMark) { failBeforeMark = false; throw new IOException("before receipt acknowledgement"); } + Properties properties = read(file); + SharedVoteReceipt completed = receipt(properties, id).completed(disposition); + putReceipt(properties, completed); + write(file, properties); + if (failAfterMark) { failAfterMark = false; throw new IOException("lost receipt acknowledgement"); } + return CompletableFuture.completedFuture(completed); + } catch (Exception failure) { return CompletableFuture.failedFuture(failure); } + } + + @Override public synchronized CompletionStage> pendingVotes(int limit) { + try { + Properties properties = read(file); + List pending = new ArrayList<>(); + for (String key : properties.stringPropertyNames().stream().sorted().toList()) { + if (key.startsWith("receipt.") && key.endsWith(".done") && properties.getProperty(key).isEmpty()) { + UUID id = UUID.fromString(key.substring("receipt.".length(), key.length() - ".done".length())); + pending.add(receipt(properties, id)); + if (pending.size() == limit) break; + } + } + return CompletableFuture.completedFuture(pending); + } catch (Exception failure) { return CompletableFuture.failedFuture(failure); } + } + + @Override public synchronized CompletionStage deliverOnce(SharedVoteReceipt receipt) { + try { + if (failBeforeDelivery) { failBeforeDelivery = false; throw new IOException("reward owner unavailable"); } + Properties properties = read(rewardFile); + String prefix = receipt.input().voteId() + "."; + String origin = receipt.input() + "|" + receipt.identity() + "|" + receipt.mutation() + + "|" + receipt.persistedState() + "|" + receipt.executeRewardsNow() + "|" + receipt.rewardPlan(); + if (properties.containsKey(prefix + "done")) { + if (!origin.equals(properties.getProperty(prefix + "origin"))) throw new IllegalStateException("Conflicting reward receipt"); + lastDeliveredPlanVersion = receipt.rewardPlan().versionReference(); + return CompletableFuture.completedFuture(RewardDisposition.valueOf(properties.getProperty(prefix + "done"))); + } + lastDeliveredPlanVersion = receipt.rewardPlan().versionReference(); + RewardDisposition disposition = receipt.executeRewardsNow() ? RewardDisposition.EXECUTED : RewardDisposition.DEFERRED; + if (receipt.executeRewardsNow()) { + effect(properties, "command:say Thanks " + receipt.identity().playerName() + ":" + receipt.persistedState().allTimeTotal()); + effect(properties, "message:Thanks " + receipt.identity().playerName() + ":" + receipt.persistedState().points()); + } else { + effect(properties, "defer:" + receipt.input().serviceSite() + ":" + receipt.persistedState().allTimeTotal()); + } + properties.setProperty(prefix + "origin", origin); + properties.setProperty(prefix + "done", disposition.name()); + write(rewardFile, properties); + if (failAfterDelivery) { failAfterDelivery = false; throw new IOException("lost durable reward acknowledgement"); } + return CompletableFuture.completedFuture(disposition); + } catch (Exception failure) { return CompletableFuture.failedFuture(failure); } + } + + @Override public CompletionStage executeVoteRewards(SharedVoteInput input, SharedVoteIdentity identity, SharedVoteUserSnapshot snapshot) { + return CompletableFuture.failedFuture(new AssertionError("Unkeyed execution API was called")); + } + @Override public CompletionStage deferVoteRewards(SharedVoteInput input, SharedVoteIdentity identity, SharedVoteUserSnapshot snapshot) { + return CompletableFuture.failedFuture(new AssertionError("Unkeyed deferral API was called")); + } + + synchronized List events() { + try { + Properties properties = read(rewardFile); + List events = new ArrayList<>(); + for (int i = 0; i < integer(properties, "effects"); i++) events.add(properties.getProperty("effect." + i)); + return events; + } catch (IOException failure) { throw new IllegalStateException(failure); } + } + private static void effect(Properties p, String value) { + int size = integer(p, "effects"); + p.setProperty("effect." + size, value); + p.setProperty("effects", Integer.toString(size + 1)); + } + private static Properties read(Path file) throws IOException { + Properties properties = new Properties(); + if (Files.exists(file)) try (InputStream input = Files.newInputStream(file)) { properties.load(input); } + return properties; + } + private static void write(Path file, Properties properties) throws IOException { + Files.createDirectories(file.getParent()); + Path temporary = file.resolveSibling(file.getFileName() + ".tmp"); + try (OutputStream out = Files.newOutputStream(temporary)) { properties.store(out, "test snapshot"); } + Files.move(temporary, file, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING); + } + private static int integer(Properties p, String key) { return Integer.parseInt(p.getProperty(key, "0")); } + private static boolean bool(Properties p, String key) { return Boolean.parseBoolean(p.getProperty(key)); } + private static SharedVoteUserSnapshot snapshot(Properties p, String prefix) { + return new SharedVoteUserSnapshot(integer(p, prefix + "all"), integer(p, prefix + "month"), integer(p, prefix + "week"), + integer(p, prefix + "day"), integer(p, prefix + "points")); + } + private static void putSnapshot(Properties p, String prefix, SharedVoteUserSnapshot s) { + p.setProperty(prefix + "all", Integer.toString(s.allTimeTotal())); + p.setProperty(prefix + "month", Integer.toString(s.monthTotal())); + p.setProperty(prefix + "week", Integer.toString(s.weeklyTotal())); + p.setProperty(prefix + "day", Integer.toString(s.dailyTotal())); + p.setProperty(prefix + "points", Integer.toString(s.points())); + } + private static void putReceipt(Properties p, SharedVoteReceipt r) { + String k = "receipt." + r.input().voteId() + "."; + p.setProperty(k + "name", r.input().playerName()); + p.setProperty(k + "site", r.input().serviceSite()); + p.setProperty(k + "time", Long.toString(r.input().voteTime())); + p.setProperty(k + "real", Boolean.toString(r.input().realVote())); + p.setProperty(k + "add", Boolean.toString(r.input().addTotals())); + p.setProperty(k + "proxy", Boolean.toString(r.input().proxyVote())); + p.setProperty(k + "forceProxy", Boolean.toString(r.input().forceProxyRouting())); + p.setProperty(k + "wasOnline", Boolean.toString(r.input().wasOnline())); + p.setProperty(k + "uuid", r.identity().uuid().toString()); + p.setProperty(k + "resolvedName", r.identity().playerName()); + p.setProperty(k + "online", Boolean.toString(r.identity().online())); + p.setProperty(k + "count", Boolean.toString(r.mutation().countTotals())); + p.setProperty(k + "award", Boolean.toString(r.mutation().awardConfiguredPoints())); + p.setProperty(k + "execute", Boolean.toString(r.executeRewardsNow())); + p.setProperty(k + "planId", r.rewardPlan().planId()); + p.setProperty(k + "planVersion", r.rewardPlan().versionReference()); + p.setProperty(k + "done", r.pending() ? "" : r.completedDisposition().name()); + putSnapshot(p, k + "state.", r.persistedState()); + } + private static SharedVoteReceipt receipt(Properties p, UUID id) { + String k = "receipt." + id + "."; + if (!p.containsKey(k + "done")) return null; + SharedVoteInput input = new SharedVoteInput(id, p.getProperty(k + "name"), p.getProperty(k + "site"), + Long.parseLong(p.getProperty(k + "time")), bool(p, k + "real"), bool(p, k + "add"), + bool(p, k + "proxy"), bool(p, k + "forceProxy"), bool(p, k + "wasOnline")); + SharedVoteIdentity identity = new SharedVoteIdentity(UUID.fromString(p.getProperty(k + "uuid")), + p.getProperty(k + "resolvedName"), bool(p, k + "online")); + SharedVoteMutation mutation = new SharedVoteMutation(id, input.serviceSite(), input.voteTime(), bool(p, k + "count"), bool(p, k + "award")); + SharedVoteRewardPlan plan = new SharedVoteRewardPlan(p.getProperty(k + "planId"), p.getProperty(k + "planVersion")); + String done = p.getProperty(k + "done"); + return new SharedVoteReceipt(input, identity, mutation, snapshot(p, k + "state."), bool(p, k + "execute"), plan, + done.isEmpty() ? null : RewardDisposition.valueOf(done)); + } +}