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
Expand Up @@ -680,6 +680,7 @@ private void metrics() {

@Override
public void onPostLoad() {
ensureCommunicationSecret();
// auto conversion for Shop.yml
if (plugin.getShopFile().isJustCreated()) {
if (!plugin.getGui().isJustCreated() && !getServerData().isVoteShopConverted()) {
Expand Down Expand Up @@ -881,6 +882,24 @@ public void run() {

}

private void ensureCommunicationSecret() {
try {
boolean created = com.bencodez.votingplugin.proxy.security.SharedSecretKeyFile
.ensure(getDataFolder().toPath().resolve("secretkey.key"));
if (created) getLogger().info("Created secretkey.key for VotingPlugin communication security");
if (!bungeeSettings.isCommunicationEncryption()) getLogger().warning(
"CommunicationEncryption is disabled. Copy the proxy secretkey.key to every VotingPlugin node, enable CommunicationEncryption everywhere, and restart (recommended).");
} catch (java.io.IOException failure) {
boolean required = bungeeSettings.isCommunicationEncryption()
|| com.bencodez.votingplugin.proxy.security.SharedTransportEnvelopeAuthenticator.Mode
.parse(bungeeSettings.getSharedTransportAuthentication())
== com.bencodez.votingplugin.proxy.security.SharedTransportEnvelopeAuthenticator.Mode.REQUIRED;
if (required) throw new IllegalStateException(
"Unable to prepare required VotingPlugin communication secretkey.key", failure);
getLogger().warning("Unable to create optional secretkey.key; continuing with legacy plaintext/unsigned communication. Fix data-folder permissions before enabling communication security.");
}
}

private void startBackendHostedControl() {
HostedControlManager.HostConfiguration configuration = readBackendHostedControlConfiguration();
try {
Expand Down Expand Up @@ -2056,11 +2075,11 @@ private void registerEvents() {
*/
@Override
public void reload() {
reloadPlugin(false, true);
reloadPlugin(false, true, true);
}

public void reloadAll() {
reloadPlugin(true, true);
reloadPlugin(true, true, true);
}

/** Captures Bukkit presence while the lifecycle caller owns platform access. */
Expand All @@ -2082,10 +2101,13 @@ static UUID placeholderStorageUuid(Player player, boolean onlineMode) {

/** Reloads configuration applied by Control before its result is acknowledged. */
public void reloadFromControl() {
reloadPlugin(false, false);
// Control publishes a separately prepared handler only after validation and
// handoff. Keep the predecessor and its transport policy stable until then.
reloadPlugin(false, false, false);
}

private void reloadPlugin(boolean userStorage, boolean reconcileHostedControl) {
private void reloadPlugin(boolean userStorage, boolean reconcileHostedControl,
boolean updateActiveBackendRuntime) {
configFile.reloadData();
configFile.loadValues();

Expand All @@ -2106,20 +2128,7 @@ private void reloadPlugin(boolean userStorage, boolean reconcileHostedControl) {
// Re-evaluate after storage has reloaded; UserManager keeps this lifecycle task unique.
getVotingPluginUserManager().startSharedPointTransferRecovery();

if (bungeeSettings.isUseBungeecoord()) {
BackendProxyHandler handler = getBackendProxyHandler();
if (handler == null) {
loadBungeeHandler();
handler = getBackendProxyHandler();
} else {
handler.reloadPresenceReporting();
}
if (userStorage && handler != null) {
handler.loadGlobalMysql();
}
} else if (getBackendProxyHandler() != null) {
getBackendProxyHandler().disablePresenceReporting();
}
reloadBackendProxyRuntime(updateActiveBackendRuntime, userStorage);
checkYMLError();

plugin.loadVoteSites();
Expand Down Expand Up @@ -2151,6 +2160,34 @@ private void reloadPlugin(boolean userStorage, boolean reconcileHostedControl) {
setUpdate(true);
}

void reloadBackendProxyRuntime(boolean updateActiveRuntime, boolean userStorage) {
if (!updateActiveRuntime) return;
if (bungeeSettings.isUseBungeecoord()) {
BackendProxyHandler handler = getBackendProxyHandler();
if (handler == null) {
loadBungeeHandler();
handler = getBackendProxyHandler();
} else {
reloadActiveBackendTransportSecurity(handler);
handler.reloadPresenceReporting();
}
if (userStorage && handler != null) {
handler.loadGlobalMysql();
}
} else if (getBackendProxyHandler() != null) {
getBackendProxyHandler().disablePresenceReporting();
}
}

void reloadActiveBackendTransportSecurity(BackendProxyHandler handler) {
try {
handler.reloadSharedTransportSecurity();
} catch (RuntimeException failure) {
getLogger().warning("Backend transport security settings were not applied; the previous policy remains active");
debug(failure);
}
}

private void loadVoteBroadcast() {
ConfigurationSection sec = getConfigFile().getData().getConfigurationSection("VoteBroadcast");
BroadcastSettings settings = BroadcastSettings.load(sec);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,11 @@
import com.bencodez.votingplugin.backendproxy.voteparty.BackendVotePartySync;
import com.bencodez.votingplugin.proxy.BungeeMethod;
import com.bencodez.votingplugin.proxy.VotingPluginWire;
import com.bencodez.votingplugin.proxy.security.SharedTransportEnvelopeAuthenticator;
import com.bencodez.votingplugin.proxy.security.SharedTransportEnvelopeAuthenticator.Mode;
import com.bencodez.votingplugin.proxy.security.TransportEnvelopeEncryption;
import com.bencodez.votingplugin.proxy.security.TransportEnvelopeEncryption.Decryption;
import com.bencodez.votingplugin.proxy.security.TransportEnvelopeEncryption.Domain;

import lombok.Getter;

Expand Down Expand Up @@ -68,6 +73,10 @@ public class BackendProxyHandler implements Listener {
private BackendVotePartySync votePartySync;
private boolean persistVotePartyOnClose = true;
private BackendProxyMessageRouter messageRouter;
private volatile TransportEnvelopeEncryption communicationEncryption;
private Mode sharedTransportMode;
private boolean communicationEncryptionEnabled;
private final AtomicBoolean encryptionFailureLogged = new AtomicBoolean();

@Getter
private BungeeMethod method;
Expand Down Expand Up @@ -107,26 +116,28 @@ private void load(boolean activatePresenceReporting) {
plugin.debug("Loading backend proxy handler");
method = BungeeMethod.getByName(plugin.getBungeeSettings().getBungeeMethod());
plugin.getLogger().info("Using BungeeMethod: " + method.toString());
try {
communicationEncryption = TransportEnvelopeEncryption.load(
plugin.getDataFolder().toPath().resolve("secretkey.key"), Domain.PROXY_BACKEND,
plugin.getBungeeSettings().isCommunicationEncryption());
} catch (java.io.IOException failure) {
throw new IllegalStateException("Proxy communication encryption initialization failed", failure);
}
transportManager.setHttpEncryption(communicationEncryption);

globalDataSync.load();
globalMessageHandler = new GlobalMessageHandler() {
@Override
public void onMessage(JsonEnvelope envelope) {
BackendProxyHandler.this.dispatchIncomingAfterPublication(envelope, () -> super.onMessage(envelope));
}

@Override
public void sendMessage(JsonEnvelope envelope) {
transportManager.send(envelope);
}
};
globalMessageHandler = new EncryptedGlobalMessageHandler();

presenceManager = new BackendPresenceManager(plugin, method, globalMessageHandler);
votePartySync = new BackendVotePartySync(plugin);
messageRouter = new BackendProxyMessageRouter(plugin, presenceManager, globalDataSync, votePartySync,
processedVoteCache);
messageRouter.register(globalMessageHandler, method);
transportManager.start(method, globalMessageHandler, activatePresenceReporting);
if (method == BungeeMethod.REDIS || method == BungeeMethod.MQTT) {
sharedTransportMode = Mode.parse(plugin.getBungeeSettings().getSharedTransportAuthentication());
communicationEncryptionEnabled = plugin.getBungeeSettings().isCommunicationEncryption();
}

if (plugin.getOptions().getServer().equalsIgnoreCase("pleaseset")) {
plugin.getLogger().warning("Server name for bungee voting is not set, please set it");
Expand All @@ -137,6 +148,45 @@ public void sendMessage(JsonEnvelope envelope) {
}
}

private final class EncryptedGlobalMessageHandler extends GlobalMessageHandler {
@Override
public void onMessage(JsonEnvelope envelope) {
if (transportHandlesEncryption()) {
acceptDecrypted(envelope);
return;
}
Decryption decrypted = communicationEncryption.decrypt(envelope);
if (!decrypted.accepted()) {
if (encryptionFailureLogged.compareAndSet(false, true)) plugin.getLogger().warning(
"Proxy communication message rejected by encryption policy (" + decrypted.reason() + ")");
return;
}
acceptDecrypted(decrypted.envelope());
}

private void acceptDecrypted(JsonEnvelope accepted) {
BackendProxyHandler.this.dispatchIncomingAfterPublication(accepted, () -> super.onMessage(accepted));
}

@Override
public void sendMessage(JsonEnvelope envelope) {
transportManager.send(transportHandlesEncryption() ? envelope : communicationEncryption.encrypt(envelope));
}
}

private boolean transportHandlesEncryption() {
return usesSharedBrokerSecurity() || method == BungeeMethod.HTTP;
}

private boolean usesSharedBrokerSecurity() {
return method == BungeeMethod.REDIS || method == BungeeMethod.MQTT;
}

private void acceptAlreadyDecrypted(JsonEnvelope envelope) {
if (globalMessageHandler instanceof EncryptedGlobalMessageHandler handler) handler.acceptDecrypted(envelope);
else if (globalMessageHandler != null) globalMessageHandler.onMessage(envelope);
}

/** Starts presence only after a staged handler reaches the atomic publication boundary. */
public void activatePresenceReporting() {
if (presenceManager != null && !presenceReportingActivated) {
Expand Down Expand Up @@ -182,16 +232,31 @@ public void activateInboundMessages() {
activateOrderedVoteDispatch();
}

/** Routes an already accepted staged callback through the restored predecessor on rollback. */
/**
* Routes a staged callback through the restored predecessor only when both
* runtimes enforced the exact same inbound cryptographic policy. A callback
* accepted under a different replacement policy must not bypass the restored
* runtime's authentication or encryption boundary as plaintext.
*/
public void abortStagedInboundTo(BackendProxyHandler previous) {
BackendProxyHandler safeRollbackTarget = hasEquivalentInboundSecurity(previous) ? previous : null;
synchronized (inboundPublication) {
if (inboundPublished || inboundAborted) return;
inboundRollbackTarget = previous;
inboundRollbackTarget = safeRollbackTarget;
inboundAborted = true;
inboundPublication.notifyAll();
}
}

private boolean hasEquivalentInboundSecurity(BackendProxyHandler previous) {
if (previous == null || method != previous.method) return false;
if (communicationEncryption == null || previous.communicationEncryption == null) {
if (communicationEncryption != previous.communicationEncryption) return false;
} else if (!communicationEncryption.hasEquivalentInboundPolicy(previous.communicationEncryption)) return false;
if (method != BungeeMethod.REDIS && method != BungeeMethod.MQTT) return true;
return transportManager.hasEquivalentSharedInboundPolicy(previous.transportManager);
}

void dispatchIncomingAfterPublication(JsonEnvelope envelope, Runnable localDispatch) {
BackendProxyHandler rollbackTarget;
synchronized (inboundPublication) {
Expand All @@ -207,8 +272,7 @@ void dispatchIncomingAfterPublication(JsonEnvelope envelope, Runnable localDispa
if (!inboundPublished && rollbackTarget == null) return;
}
if (rollbackTarget != null) {
GlobalMessageHandler rollbackHandler = rollbackTarget.globalMessageHandler;
if (rollbackHandler != null) rollbackHandler.onMessage(envelope);
rollbackTarget.acceptAlreadyDecrypted(envelope);
Comment thread
BenCodez marked this conversation as resolved.
return;
}
// Reward-bearing proxy votes construct an asynchronous PlayerVoteEvent. The
Expand Down Expand Up @@ -861,6 +925,32 @@ public void reloadPresenceReporting() {
}
}

/** Refreshes policies captured when a Redis or MQTT transport first started. */
public void reloadSharedTransportSecurity() {
if (method != BungeeMethod.REDIS && method != BungeeMethod.MQTT) return;
Mode requestedMode = Mode.parse(plugin.getBungeeSettings().getSharedTransportAuthentication());
boolean requestedEncryption = plugin.getBungeeSettings().isCommunicationEncryption();
try {
java.nio.file.Path keyFile = plugin.getDataFolder().toPath().resolve("secretkey.key");
TransportEnvelopeEncryption candidateEncryption = TransportEnvelopeEncryption.load(
keyFile, Domain.PROXY_BACKEND, requestedEncryption);
SharedTransportEnvelopeAuthenticator candidateAuthenticator =
SharedTransportEnvelopeAuthenticator.load(keyFile, requestedMode);
boolean authenticatorChanged = !transportManager.hasEquivalentSharedTransportAuthenticator(
candidateAuthenticator);
boolean encryptionChanged = !transportManager.hasEquivalentSharedTransportEncryption(candidateEncryption);
if (!authenticatorChanged && !encryptionChanged) return;
transportManager.updateSharedTransportSecurity(authenticatorChanged ? candidateAuthenticator : null,
encryptionChanged ? candidateEncryption : null);
if (encryptionChanged) communicationEncryption = candidateEncryption;
encryptionFailureLogged.set(false);
sharedTransportMode = requestedMode;
communicationEncryptionEnabled = requestedEncryption;
} catch (java.io.IOException failure) {
throw new IllegalStateException("Shared backend transport security reload failed", failure);
}
}

public void disablePresenceReporting() {
if (presenceManager != null && presenceReportingActivated) {
presenceManager.stop();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
import com.bencodez.votingplugin.backendproxy.cache.ProcessedVoteCache;
import com.bencodez.votingplugin.proxy.BungeeMethod;
import com.bencodez.votingplugin.proxy.VotingPluginWire;
import com.bencodez.votingplugin.proxy.security.SharedTransportEnvelopeAuthenticator;

/**
* Selects and owns the active backend-to-proxy transport.
Expand All @@ -22,6 +23,7 @@ public class BackendProxyTransportManager {

private final VotingPluginMain plugin;
private final ProcessedVoteCache processedVoteCache;
private com.bencodez.votingplugin.proxy.security.TransportEnvelopeEncryption httpEncryption;
private BackendProxyTransport transport;
private BackendProxyTransport preparedTransport;
private BackendProxyTransport retiredTransport;
Expand Down Expand Up @@ -51,6 +53,11 @@ public BackendProxyTransportManager(VotingPluginMain plugin, ProcessedVoteCache
this.processedVoteCache = processedVoteCache;
}

public void setHttpEncryption(
com.bencodez.votingplugin.proxy.security.TransportEnvelopeEncryption httpEncryption) {
this.httpEncryption = httpEncryption;
}

public void start(BungeeMethod method, GlobalMessageHandler messageHandler) {
start(method, messageHandler, true);
}
Expand All @@ -68,7 +75,7 @@ public void start(BungeeMethod method, GlobalMessageHandler messageHandler, bool
transport = new SocketBackendProxyTransport(plugin);
break;
case HTTP:
transport = new HttpBackendProxyTransport(plugin);
transport = new HttpBackendProxyTransport(plugin, httpEncryption);
break;
case REDIS:
transport = new RedisBackendProxyTransport(plugin, processedVoteCache);
Expand Down Expand Up @@ -111,6 +118,41 @@ public synchronized void send(JsonEnvelope envelope) {
}
}

public synchronized void updateSharedTransportSecurity(SharedTransportEnvelopeAuthenticator authenticator,
com.bencodez.votingplugin.proxy.security.TransportEnvelopeEncryption encryption) {
if (transport instanceof RedisBackendProxyTransport redis) redis.updateSecurity(authenticator, encryption);
else if (transport instanceof MqttBackendProxyTransport mqtt) mqtt.updateSecurity(authenticator, encryption);
else throw new IllegalStateException("No active shared backend transport to update");
}

public synchronized boolean hasEquivalentSharedTransportAuthenticator(
SharedTransportEnvelopeAuthenticator authenticator) {
SharedInboundPolicy policy = sharedInboundPolicySnapshot();
return policy != null && policy.authenticator() != null
&& policy.authenticator().hasEquivalentInboundPolicy(authenticator);
}

public synchronized boolean hasEquivalentSharedTransportEncryption(
com.bencodez.votingplugin.proxy.security.TransportEnvelopeEncryption encryption) {
SharedInboundPolicy policy = sharedInboundPolicySnapshot();
return policy != null && policy.encryption() != null
&& policy.encryption().hasEquivalentInboundPolicy(encryption);
}

private synchronized SharedInboundPolicy sharedInboundPolicySnapshot() {
if (transport instanceof RedisBackendProxyTransport redis) return redis.sharedInboundPolicySnapshot();
if (transport instanceof MqttBackendProxyTransport mqtt) return mqtt.sharedInboundPolicySnapshot();
return null;
}

/** Compares broker authentication and the destination bound into its MAC without nesting manager locks. */
public boolean hasEquivalentSharedInboundPolicy(BackendProxyTransportManager other) {
if (other == null) return false;
SharedInboundPolicy current = sharedInboundPolicySnapshot();
SharedInboundPolicy restored = other.sharedInboundPolicySnapshot();
return current != null && current.hasEquivalentPolicy(restored);
}

private void acceptPreparedSend(JsonEnvelope envelope) {
if (preparedSends.size() < MAX_PREPARED_SENDS) {
preparedSends.addLast(envelope);
Expand Down
Loading
Loading