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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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
}
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ public final class SharedRewardContext {
private final UUID userId;
private final String playerName;
private final HashMap<String, String> placeholders;
private volatile boolean keyedExecution;

public SharedRewardContext(UUID userId, String playerName, Map<String, String> placeholders) {
this.userId = Objects.requireNonNull(userId, "userId");
Expand All @@ -32,4 +33,12 @@ public String playerName() {
public HashMap<String, String> placeholders() {
return placeholders;
}

void markKeyedExecution() {
keyedExecution = true;
}

boolean isKeyedExecution() {
return keyedExecution;
}
}
Original file line number Diff line number Diff line change
@@ -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);
}
}
Original file line number Diff line number Diff line change
@@ -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<SharedRewardActionClaim> 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<Void> checkpoint(String executionPath, String fingerprint, int completedSteps,
SharedRewardContext context);
}
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,27 @@ public CompletionStage<SharedRewardResult> 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<SharedRewardResult> 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<SharedRewardResult> executeNested(SharedRewardPlan plan, SharedRewardContext context,
Expand All @@ -37,18 +57,22 @@ public CompletionStage<SharedRewardResult> 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<SharedRewardResult> 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);
Expand Down Expand Up @@ -90,7 +114,7 @@ private CompletionStage<SharedRewardResult> 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) {
Expand All @@ -99,11 +123,12 @@ private CompletionStage<SharedRewardResult> execute(SharedRewardPlan plan, Share
}

private CompletionStage<SharedRewardResult> 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,
Expand All @@ -121,10 +146,11 @@ private void validateProgress(SharedRewardPlan plan, String fingerprint, SharedR
}

private CompletionStage<SharedRewardResult> 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<CompletionStage<SharedRewardResult>> 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();
Expand Down Expand Up @@ -158,23 +184,44 @@ private CompletionStage<Outcome> evaluateRequirements(List<SharedRewardRequireme
}

private CompletionStage<SharedRewardResult> executeSteps(SharedRewardPlan plan, SharedRewardContext context,
SharedRewardDurability durability, String executionPath, String fingerprint, int index) {
SharedRewardDurability durability, String executionPath, String fingerprint, int index,
SharedRewardKeyedDurability keyed) {
CompletionStage<SharedRewardResult> 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")
: CompletableFuture.completedFuture(result));
}

private CompletionStage<SharedRewardResult> 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<Boolean> 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<SharedRewardResult> 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<Void> deferred;
try {
Expand All @@ -188,6 +235,42 @@ private CompletionStage<SharedRewardResult> 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<SharedRewardActionClaim> 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<SharedRewardResult> 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<SharedRewardResult> executeClaimedStep(SharedRewardStep step, SharedRewardContext context,
SharedRewardDurability durability, String executionPath, String fingerprint, int index, String stepPath,
boolean claimedStrict) {
CompletionStage<SharedRewardResult> action;
try {
action = step.action().execute(context, stepPath);
Expand All @@ -200,7 +283,12 @@ private CompletionStage<SharedRewardResult> 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<Void> checkpoint;
try {
checkpoint = durability.checkpoint(executionPath, fingerprint, index + 1, context);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -25,4 +26,27 @@ CompletionStage<SharedRewardResult> delay(Duration delay,
Supplier<CompletionStage<SharedRewardResult>> 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<Boolean> 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<SharedRewardResult> runClaimedAction(
UUID userId, boolean requiresOnlinePlayer, Supplier<CompletionStage<SharedRewardResult>> operation) {
return CompletableFuture.failedFuture(new UnsupportedOperationException(
"Keyed reward execution requires a native action scheduler"));
}
}
Loading
Loading