From b0328c32e0c92270b23ab5a4e95a821d031209c0 Mon Sep 17 00:00:00 2001 From: BenCodez <17074231+BenCodez@users.noreply.github.com> Date: Sat, 19 Sep 2026 22:23:43 -0600 Subject: [PATCH 1/7] Add atomic user transaction extension for caller SQL --- .../api/user/usercache/UserDataManager.java | 16 + .../user/runtime/BukkitUserCacheOwner.java | 5 + .../user/storage/BukkitSqlUserBackend.java | 39 ++ .../user/runtime/SharedUserDataRuntime.java | 58 +++ .../core/user/runtime/UserCacheOwner.java | 2 + .../core/user/storage/SqlUserStorage.java | 48 ++- .../user/storage/sql/JdbcSqlUserStorage.java | 57 ++- .../user/storage/sql/MysqlUserBackend.java | 4 + .../storage/sql/SqlUserBackendFactory.java | 20 + .../user/storage/sql/SqliteUserBackend.java | 4 + .../storage/SqliteRequiredColumnTest.java | 18 + .../storage/SqliteUserTransactionTest.java | 399 ++++++++++++++++++ 12 files changed, 664 insertions(+), 6 deletions(-) create mode 100644 AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/usercache/UserDataManager.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/usercache/UserDataManager.java index 77a38c4e6..5a7fd394b 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/usercache/UserDataManager.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/usercache/UserDataManager.java @@ -3,6 +3,7 @@ import java.util.ArrayList; import java.util.Objects; import java.util.Map.Entry; +import java.util.Map; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; @@ -36,6 +37,7 @@ import com.bencodez.advancedcore.core.user.runtime.SharedUserDataRuntime; import com.bencodez.simpleapi.debug.DebugLevel; import com.bencodez.simpleapi.array.ArrayUtils; +import com.bencodez.simpleapi.sql.data.DataValue; import lombok.Getter; @@ -325,6 +327,20 @@ public final boolean hasSharedRuntime() { return runtime != null && !runtime.isClosed(); } + /** + * Cache-safe transaction entry for callers adding their own durable SQL + * record to a user mutation. Initial values are row prerequisites and may + * commit before queued cache work; put operation-specific mutation only in + * the callback. Blocking work must run off the game thread. + */ + public final T withAtomicUserTransaction(UUID uuid, UserStorage storage, + Map initialValues, + SqlUserStorage.TransactionWork work) { + SharedUserDataRuntime runtime = sharedRuntime; + if (runtime == null) throw new IllegalStateException("Shared user runtime is unavailable"); + return runtime.transaction(uuid, storage, initialValues, work); + } + /** * Replace the active shared backend on this manager's worker. The runtime * flushes and retires the old cache generation before publishing the new diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/runtime/BukkitUserCacheOwner.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/runtime/BukkitUserCacheOwner.java index 052c1cfe1..1e3c626f5 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/runtime/BukkitUserCacheOwner.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/runtime/BukkitUserCacheOwner.java @@ -145,6 +145,11 @@ private void bind(UserDataCache cache, UUID uuid) { @Override public boolean isCached(UUID uuid) { return manager.isCached(uuid); } + @Override public boolean hasPendingChanges(UUID uuid) { + UserDataCache cache = manager.getUserDataCache().get(uuid); + return cache != null && cache.hasChangesToProcess(); + } + @Override public DataValue getIfPresent(UUID uuid, String key) { UserDataCache cache = manager.getUserDataCache().get(uuid); if (cache == null) return null; diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java index 1b845cad6..f4f4298da 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java @@ -7,13 +7,18 @@ import java.util.UUID; import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Consumer; +import java.sql.DriverManager; +import java.sql.SQLException; import com.bencodez.advancedcore.AdvancedCorePlugin; import com.bencodez.advancedcore.api.user.UserStorage; import com.bencodez.advancedcore.api.user.userstorage.mysql.MySQL; import com.bencodez.advancedcore.api.user.userstorage.sql.UserTable; import com.bencodez.advancedcore.core.user.storage.SqlUserStorage; +import com.bencodez.advancedcore.core.user.storage.sql.SqlBackendLogger; import com.bencodez.advancedcore.core.user.storage.sql.SqlUserBackend; +import com.bencodez.advancedcore.core.user.storage.sql.SqlUserBackendFactory; +import com.bencodez.advancedcore.core.user.storage.sql.SqlUserSchema; import com.bencodez.simpleapi.sql.Column; import com.bencodez.simpleapi.sql.data.DataValue; import com.bencodez.simpleapi.sql.data.DataValueString; @@ -74,6 +79,12 @@ public SqlUserStorage user(UUID uuid) { @Override public void writeValues(UserStorage storage, HashMap values) { BukkitSqlUserBackend.this.writeValues(storage, uuid, values); } + @Override public T transaction(UserStorage storage, TransactionWork work) { + return BukkitSqlUserBackend.this.transaction(storage, uuid, java.util.Map.of(), work); + } + @Override public T transaction(UserStorage storage, java.util.Map initialValues, TransactionWork work) { + return BukkitSqlUserBackend.this.transaction(storage, uuid, initialValues, work); + } }; } @@ -157,6 +168,34 @@ private void writeValues(UserStorage storage, UUID uuid, HashMap T transaction(UserStorage storage, UUID uuid, java.util.Map initialValues, SqlUserStorage.TransactionWork work) { + requireOpen(); + requireStorage(storage); + Objects.requireNonNull(work, "work"); + SqlUserSchema schema = SqlUserSchema.fromKeys(plugin.getUserManager().getDataManager().getKeys()); + SqlBackendLogger logger = new SqlBackendLogger() { + @Override public void info(String message) { plugin.getLogger().info(message); } + @Override public void warn(String message, Throwable error) { + plugin.getLogger().warning(message + (error == null ? "" : ": " + error.getMessage())); + } + }; + if (storage == UserStorage.MYSQL) { + var manager = mysql().getMysql().getConnectionManager(); + return SqlUserBackendFactory.existingUser(storage, uuid, mysql().getTableName(), schema, + manager::getConnection, manager.getDbType(), logger).transaction(storage, initialValues, work); + } + synchronized (sqliteOperations) { + // The legacy SQLite provider retains one shared connection. Obtain + // its actual database URL, then open a separate transaction-owned + // connection so AdvancedCore can close it after commit/rollback. + String url; + try { url = table().getSqLite().getSQLConnection().getMetaData().getURL(); } + catch (SQLException failure) { throw new IllegalStateException("Failed to locate Bukkit SQLite user database", failure); } + return SqlUserBackendFactory.existingUser(storage, uuid, table().getName(), schema, + () -> DriverManager.getConnection(url), null, logger).transaction(storage, initialValues, work); + } + } + private Column primary(UUID uuid) { return new Column("uuid", new DataValueString(uuid.toString())); } private MySQL mysql() { return requireMysql(); } private UserTable table() { return requireTable(); } diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/SharedUserDataRuntime.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/SharedUserDataRuntime.java index 38a57a7db..bbd03e013 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/SharedUserDataRuntime.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/SharedUserDataRuntime.java @@ -2,6 +2,7 @@ import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.Objects; import java.util.Set; import java.util.UUID; @@ -16,7 +17,9 @@ import java.util.function.Supplier; import com.bencodez.advancedcore.api.user.UserDataFetchMode; +import com.bencodez.advancedcore.api.user.UserStorage; import com.bencodez.advancedcore.core.user.storage.SqlUserDataAccess; +import com.bencodez.advancedcore.core.user.storage.SqlUserStorage; import com.bencodez.advancedcore.core.user.storage.sql.SqlUserBackend; import com.bencodez.simpleapi.sql.Column; import com.bencodez.simpleapi.sql.data.DataValue; @@ -119,6 +122,61 @@ public void flush(UUID uuid) { private void flushInternal(UUID uuid) { cacheOwner.flush(uuid, backend.storageType(), backend.user(uuid)); } + /** + * Extend the active SQL backend's user transaction with caller-owned SQL. + * Pending cache changes are flushed first. Only after commit is the old + * cache generation retired, so a later load observes committed values and + * a rolled-back callback cannot publish speculative values. + */ + public T transaction(UUID uuid, SqlUserStorage.TransactionWork work) { + return transactionInternal(uuid, null, Map.of(), work); + } + + /** + * Require the selected physical store. Initial values are prerequisites + * for creating a row with required columns, not part of the caller's atomic + * mutation. When a cache already exists, they can commit before its queued + * changes are flushed; callers must write operation-specific values only + * inside the transaction callback. + */ + public T transaction(UUID uuid, UserStorage expectedStorage, Map initialValues, + SqlUserStorage.TransactionWork work) { + Objects.requireNonNull(expectedStorage, "expectedStorage"); + return transactionInternal(uuid, expectedStorage, initialValues, work); + } + + private T transactionInternal(UUID uuid, UserStorage expectedStorage, Map initialValues, + SqlUserStorage.TransactionWork work) { + Objects.requireNonNull(uuid, "uuid"); + Objects.requireNonNull(initialValues, "initialValues"); + Objects.requireNonNull(work, "work"); + cacheOwner.requireBlockingAllowed(); + try { + return userExclusiveAccess(uuid, () -> { + if (expectedStorage != null && backend.storageType() != expectedStorage) { + throw new IllegalStateException("Cannot access " + expectedStorage + + " user storage while the shared runtime owns " + backend.storageType()); + } + cacheOwner.beginRemoval(uuid); + try { + if (cacheOwner.hasPendingChanges(uuid) && !initialValues.isEmpty()) { + // Existing queued changes may need a required column on + // first write. Establish only row prerequisites before + // their ordinary, independently durable cache flush. + backend.user(uuid).transaction(backend.storageType(), initialValues, scope -> null); + } + flushInternal(uuid); + T result = backend.user(uuid).transaction(backend.storageType(), initialValues, work); + cacheOwner.remove(uuid); + return result; + } catch (RuntimeException | Error failure) { + cacheOwner.cancelRemoval(uuid); + throw failure; + } + }); + } finally { cacheOwner.dispatchNotifications(uuid); } + } + public void flushAll() { try { storageAccess(() -> { for (UUID uuid : Set.copyOf(cacheOwner.cachedUsers())) userAccess(uuid, () -> { flushInternal(uuid); return null; }); diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/UserCacheOwner.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/UserCacheOwner.java index 81038cf0b..050489f28 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/UserCacheOwner.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/UserCacheOwner.java @@ -14,6 +14,8 @@ /** Port for the existing cache/queue owner. No parallel cache is allocated. */ public interface UserCacheOwner { boolean isCached(UUID uuid); + /** True when a cached user has queued values that need a pre-transaction flush. */ + default boolean hasPendingChanges(UUID uuid) { return false; } DataValue getIfPresent(UUID uuid, String key); void populate(UUID uuid, HashMap values); diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/SqlUserStorage.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/SqlUserStorage.java index b7bca5caf..228d803a5 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/SqlUserStorage.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/SqlUserStorage.java @@ -1,7 +1,10 @@ package com.bencodez.advancedcore.core.user.storage; +import java.sql.Connection; +import java.sql.SQLException; import java.util.HashMap; import java.util.List; +import java.util.Map; import com.bencodez.advancedcore.api.user.UserStorage; import com.bencodez.simpleapi.sql.Column; @@ -15,7 +18,8 @@ * * These operations retain the backing provider's error/commit semantics. A void * return is not an additional durable reward acknowledgement. This boundary does - * not add retries, transactions, or a second write queue. + * does not add retries or a second write queue. JDBC implementations also + * expose a caller extension of their existing user-row transaction. */ public interface SqlUserStorage { List readRow(UserStorage storage); @@ -27,4 +31,46 @@ public interface SqlUserStorage { void write(UserStorage storage, String key, DataValue value); void writeValues(UserStorage storage, HashMap values); + + /** + * Run caller-owned SQL and user writes in one backend-owned transaction. + * The user row exists and is locked before the callback runs. A callback + * must use only this scope for user writes; ordinary methods open another + * connection. Cache-backed callers must enter through + * UserDataManager.withAtomicUserTransaction so cached writes are flushed and the + * committed snapshot is reconciled. This blocking method belongs on a + * storage worker. Deadlock/busy retries are the caller's responsibility. + */ + default T transaction(UserStorage storage, TransactionWork work) { + return transaction(storage, Map.of(), work); + } + + /** + * Initial values are inserted only when the user row does not yet exist. + * Supply any required non-default columns here; existing rows retain their + * values until the callback explicitly writes them. + */ + default T transaction(UserStorage storage, Map initialValues, TransactionWork work) { + throw new UnsupportedOperationException("Atomic JDBC user transactions are unavailable"); + } + + @FunctionalInterface + interface TransactionWork { + T run(TransactionScope scope) throws SQLException; + } + + interface TransactionScope { + /** + * Active connection for caller-owned tables. Never commit, roll back, + * close, or change auto-commit. AdvancedCore owns its lifecycle. Do not + * retain the connection or scope after the callback returns. + */ + Connection connection(); + + /** Read the locked user row on the active transaction connection. */ + List readRow() throws SQLException; + + /** Write registered user values on the active transaction connection. */ + void writeValues(Map values) throws SQLException; + } } diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/JdbcSqlUserStorage.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/JdbcSqlUserStorage.java index 3fde49890..63c6da2d5 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/JdbcSqlUserStorage.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/JdbcSqlUserStorage.java @@ -69,8 +69,14 @@ private enum BooleanStorage { TEXT, NATIVE, NUMERIC, POSTGRES_BIT } @Override public List readRow(UserStorage requestedStorage) { requireStorage(requestedStorage); + try (Connection connection = connections.open()) { + return readRow(connection); + } catch (SQLException e) { throw failure("read user row", e); } + } + + private List readRow(Connection connection) throws SQLException { String sql = "SELECT * FROM " + quote(tableName) + " WHERE " + quote(SqlUserSchema.UUID_COLUMN) + "=?"; - try (Connection connection = connections.open(); PreparedStatement statement = connection.prepareStatement(sql)) { + try (PreparedStatement statement = connection.prepareStatement(sql)) { dialect.bindUuid(statement, 1, uuid); try (ResultSet result = statement.executeQuery()) { if (!result.next()) return new ArrayList<>(); @@ -86,7 +92,7 @@ private enum BooleanStorage { TEXT, NATIVE, NUMERIC, POSTGRES_BIT } } return columns; } - } catch (SQLException e) { throw failure("read user row", e); } + } } @Override public boolean contains(UserStorage requestedStorage) { @@ -117,15 +123,53 @@ private enum BooleanStorage { TEXT, NATIVE, NUMERIC, POSTGRES_BIT } Objects.requireNonNull(values, "values"); Map updates = canonicalize(values); if (updates.isEmpty()) return; + inTransaction("write user values", connection -> { + writeValues(connection, updates); + return null; + }); + } + + @Override public T transaction(UserStorage requestedStorage, Map initialValues, TransactionWork work) { + requireStorage(requestedStorage); + Objects.requireNonNull(work, "work"); + Map seed = canonicalize(Objects.requireNonNull(initialValues, "initialValues")); + return inTransaction("run user transaction", connection -> { + ensureRow(connection, seed); + Scope scope = new Scope(connection); + try { return work.run(scope); } + finally { scope.active = false; } + }); + } + + private final class Scope implements TransactionScope { + private final Connection connection; + private boolean active = true; + private Scope(Connection connection) { this.connection = connection; } + private void requireActive() { if (!active) throw new IllegalStateException("SQL user transaction scope has ended"); } + @Override public Connection connection() { requireActive(); return connection; } + @Override public List readRow() throws SQLException { requireActive(); return JdbcSqlUserStorage.this.readRow(connection); } + @Override public void writeValues(Map values) throws SQLException { + requireActive(); + JdbcSqlUserStorage.this.writeValues(connection, canonicalize(Objects.requireNonNull(values, "values"))); + } + } + + private void writeValues(Connection connection, Map updates) throws SQLException { + if (updates.isEmpty()) return; + boolean updateExistingRow = ensureRow(connection, updates); + if (updateExistingRow) updateValues(connection, updates); + } + + private T inTransaction(String operation, TransactionWorkOnConnection work) { boolean committed = false; + T result = null; try (Connection connection = connections.open()) { boolean autoCommit = connection.getAutoCommit(); connection.setAutoCommit(false); boolean transactionEnded = false; Throwable transactionFailure = null; try { - boolean updateExistingRow = ensureRow(connection, updates); - if (updateExistingRow) updateValues(connection, updates); + result = work.run(connection); connection.commit(); committed = true; transactionEnded = true; } catch (SQLException | RuntimeException | Error e) { transactionFailure = e; @@ -141,12 +185,15 @@ private enum BooleanStorage { TEXT, NATIVE, NUMERIC, POSTGRES_BIT } } } } catch (SQLException e) { - if (committed) committedCleanupFailure("close SQL connection", e); else throw failure("write user values", e); + if (committed) committedCleanupFailure("close SQL connection", e); else throw failure(operation, e); } catch (RuntimeException e) { if (committed) committedCleanupFailure("close SQL connection", e); else throw e; } + return result; } + @FunctionalInterface private interface TransactionWorkOnConnection { T run(Connection connection) throws SQLException; } + private Map canonicalize(Map values) { Map updates = new LinkedHashMap<>(); for (Map.Entry entry : values.entrySet()) { diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/MysqlUserBackend.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/MysqlUserBackend.java index 24cd41307..3c2837aea 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/MysqlUserBackend.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/MysqlUserBackend.java @@ -62,6 +62,10 @@ public SqlUserStorage user(UUID uuid) { @Override public void delete(UserStorage storage) { withOperation(() -> { delegate.delete(storage); return null; }); } @Override public void write(UserStorage storage, String key, DataValue value) { withOperation(() -> { delegate.write(storage, key, value); return null; }); } @Override public void writeValues(UserStorage storage, HashMap values) { withOperation(() -> { delegate.writeValues(storage, values); return null; }); } + @Override public T transaction(UserStorage storage, TransactionWork work) { return withOperation(() -> delegate.transaction(storage, work)); } + @Override public T transaction(UserStorage storage, java.util.Map initialValues, TransactionWork work) { + return withOperation(() -> delegate.transaction(storage, initialValues, work)); + } }; } diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/SqlUserBackendFactory.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/SqlUserBackendFactory.java index 114ff26de..354a38506 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/SqlUserBackendFactory.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/SqlUserBackendFactory.java @@ -1,10 +1,16 @@ package com.bencodez.advancedcore.core.user.storage.sql; +import java.sql.Connection; +import java.sql.SQLException; import java.nio.file.Path; import java.util.Collection; import java.util.Objects; +import java.util.UUID; +import com.bencodez.advancedcore.api.user.UserStorage; import com.bencodez.advancedcore.api.user.usercache.keys.UserDataKey; +import com.bencodez.advancedcore.core.user.storage.SqlUserStorage; +import com.bencodez.simpleapi.sql.mysql.DbType; import com.bencodez.simpleapi.sql.mysql.config.MysqlConfig; /** @@ -27,4 +33,18 @@ public static MysqlUserBackend mysql(String baseTableName, MysqlConfig config, Objects.requireNonNull(keys, "keys"); return new MysqlUserBackend(baseTableName, config, SqlUserSchema.fromKeys(keys), logger); } + + @FunctionalInterface + public interface ConnectionOpener { Connection open() throws SQLException; } + + /** Adapt an existing platform-owned SQL provider to the shared JDBC user transaction. */ + public static SqlUserStorage existingUser(UserStorage storage, UUID uuid, String tableName, + SqlUserSchema schema, ConnectionOpener connections, DbType dbType, SqlBackendLogger logger) { + Objects.requireNonNull(storage, "storage"); + JdbcSqlUserStorage.Dialect dialect = storage == UserStorage.SQLITE + ? JdbcSqlUserStorage.Dialect.SQLITE + : JdbcSqlUserStorage.Dialect.fromDbType(Objects.requireNonNull(dbType, "dbType")); + return new JdbcSqlUserStorage(storage, uuid, tableName, schema, + connections::open, dialect, logger); + } } diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/SqliteUserBackend.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/SqliteUserBackend.java index 5a7ed8b44..d42abbdba 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/SqliteUserBackend.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/SqliteUserBackend.java @@ -60,6 +60,10 @@ public SqlUserStorage user(UUID uuid) { @Override public void delete(UserStorage storage) { withOperation(() -> { delegate.delete(storage); return null; }); } @Override public void write(UserStorage storage, String key, DataValue value) { withOperation(() -> { delegate.write(storage, key, value); return null; }); } @Override public void writeValues(UserStorage storage, HashMap values) { withOperation(() -> { delegate.writeValues(storage, values); return null; }); } + @Override public T transaction(UserStorage storage, TransactionWork work) { return withOperation(() -> delegate.transaction(storage, work)); } + @Override public T transaction(UserStorage storage, java.util.Map initialValues, TransactionWork work) { + return withOperation(() -> delegate.transaction(storage, initialValues, work)); + } }; } diff --git a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteRequiredColumnTest.java b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteRequiredColumnTest.java index 5b685a2ef..846fc9e40 100644 --- a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteRequiredColumnTest.java +++ b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteRequiredColumnTest.java @@ -5,6 +5,7 @@ import java.nio.file.Path; import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.UUID; import org.junit.jupiter.api.Test; @@ -50,6 +51,23 @@ class SqliteRequiredColumnTest { } } + @Test void firstAtomicTransactionCanSupplyRequiredColumnWithoutChangingExistingRows() { + UUID uuid = UUID.randomUUID(); + try (SqliteUserBackend backend = open()) { + var user = backend.user(uuid); + assertThrows(IllegalStateException.class, () -> user.transaction(UserStorage.SQLITE, scope -> null)); + assertFalse(user.contains(UserStorage.SQLITE)); + user.transaction(UserStorage.SQLITE, Map.of("PlayerName", new DataValueString("Ben")), scope -> { + scope.writeValues(Map.of("Points", new DataValueInt(12))); + return null; + }); + user.transaction(UserStorage.SQLITE, Map.of("PlayerName", new DataValueString("Other")), scope -> null); + var values = new SqlUserDataAccess(user).getValues(UserStorage.SQLITE); + assertEquals("Ben", values.get("PlayerName").getString()); + assertEquals(12, values.get("Points").getInt()); + } + } + @Test void missingRequiredValueCannotReportASuccessfulWriteToANonexistentRow() { UUID uuid = UUID.randomUUID(); try (SqliteUserBackend backend = open()) { diff --git a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java new file mode 100644 index 000000000..86862f6a2 --- /dev/null +++ b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java @@ -0,0 +1,399 @@ +package com.bencodez.advancedcore.tests.storage; + +import static org.junit.jupiter.api.Assertions.*; +import static org.mockito.Mockito.*; + +import java.nio.file.Path; +import java.sql.Connection; +import java.sql.DriverManager; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.logging.Logger; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.junit.jupiter.api.io.TempDir; + +import com.bencodez.advancedcore.api.user.UserStorage; +import com.bencodez.advancedcore.AdvancedCorePlugin; +import com.bencodez.advancedcore.api.user.UserManager; +import com.bencodez.advancedcore.api.user.usercache.UserDataManager; +import com.bencodez.advancedcore.api.user.usercache.keys.UserDataKeyInt; +import com.bencodez.advancedcore.api.user.userstorage.sql.UserTable; +import com.bencodez.advancedcore.bukkit.user.storage.BukkitSqlUserBackend; +import com.bencodez.advancedcore.api.user.UserDataFetchMode; +import com.bencodez.advancedcore.core.user.runtime.SharedUserDataRuntime; +import com.bencodez.advancedcore.core.user.runtime.UserCacheOwner; +import com.bencodez.advancedcore.core.user.storage.SqlUserStorage; +import com.bencodez.advancedcore.core.user.storage.sql.SqlBackendLogger; +import com.bencodez.advancedcore.core.user.storage.sql.SqlUserSchema; +import com.bencodez.advancedcore.core.user.storage.sql.SqliteUserBackend; +import com.bencodez.simpleapi.sql.Column; +import com.bencodez.simpleapi.sql.DataType; +import com.bencodez.simpleapi.sql.data.DataValue; +import com.bencodez.simpleapi.sql.data.DataValueInt; +import com.bencodez.simpleapi.sql.data.DataValueString; +import com.bencodez.simpleapi.sql.sqlite.db.SQLite; + +@Timeout(20) +class SqliteUserTransactionTest { + @TempDir Path tempDir; + private static final UserStorage TYPE = UserStorage.SQLITE; + + private SqliteUserBackend backend() { + return new SqliteUserBackend(tempDir, "Users", "Users", SqlUserSchema.builder() + .column("Points", "INTEGER", DataType.INTEGER) + .column("Votes", "INTEGER", DataType.INTEGER).build(), SqlBackendLogger.NO_OP); + } + + private void createReceipts(Path file) throws SQLException { + try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + file); + PreparedStatement statement = connection.prepareStatement( + "CREATE TABLE Receipts (operation_key TEXT PRIMARY KEY, intent TEXT NOT NULL)")) { + statement.executeUpdate(); + } + } + + private static int value(List row, String key) { + return row.stream().filter(column -> key.equalsIgnoreCase(column.getName())) + .findFirst().orElseThrow().getValue().getInt(); + } + + private static void insert(Connection connection, String key) throws SQLException { + try (PreparedStatement statement = connection.prepareStatement( + "INSERT INTO Receipts (operation_key, intent) VALUES (?, ?)")) { + statement.setString(1, key); + statement.setString(2, "immutable-plan"); + statement.executeUpdate(); + } + } + + private static boolean receipt(Connection connection, String key) throws SQLException { + try (PreparedStatement statement = connection.prepareStatement( + "SELECT intent FROM Receipts WHERE operation_key=?")) { + statement.setString(1, key); + try (ResultSet result = statement.executeQuery()) { + return result.next() && "immutable-plan".equals(result.getString(1)); + } + } + } + + private static HashMap values(int points, int votes) { + HashMap values = new HashMap<>(); + values.put("Points", new DataValueInt(points)); + values.put("Votes", new DataValueInt(votes)); + return values; + } + + @Test void commitsUserValuesAndCallerReceiptTogetherAcrossRestart() throws Exception { + UUID uuid = UUID.randomUUID(); + try (SqliteUserBackend backend = backend()) { + createReceipts(backend.databaseFile()); + SqlUserStorage user = backend.user(uuid); + assertEquals("done", user.transaction(TYPE, scope -> { + scope.writeValues(values(7, 2)); + insert(scope.connection(), "one"); + return "done"; + })); + } + try (SqliteUserBackend reopened = backend(); + Connection connection = DriverManager.getConnection("jdbc:sqlite:" + reopened.databaseFile())) { + assertEquals(7, value(reopened.user(uuid).readRow(TYPE), "Points")); + assertEquals(2, value(reopened.user(uuid).readRow(TYPE), "Votes")); + assertTrue(receipt(connection, "one")); + } + } + + @Test void rollsBackBothFailureOrders() throws Exception { + UUID uuid = UUID.randomUUID(); + try (SqliteUserBackend backend = backend()) { + createReceipts(backend.databaseFile()); + SqlUserStorage user = backend.user(uuid); + user.writeValues(TYPE, values(3, 4)); + assertThrows(IllegalStateException.class, () -> user.transaction(TYPE, scope -> { + scope.writeValues(values(9, 10)); + insert(scope.connection(), "after"); + throw new SQLException("fail after both writes"); + })); + assertThrows(IllegalStateException.class, () -> user.transaction(TYPE, scope -> { + insert(scope.connection(), "before"); + throw new SQLException("fail before user write"); + })); + UUID absent = UUID.randomUUID(); + assertThrows(IllegalStateException.class, () -> backend.user(absent).transaction(TYPE, scope -> { + insert(scope.connection(), "new-user"); + throw new SQLException("fail before creating user values"); + })); + assertFalse(backend.user(absent).contains(TYPE)); + assertThrows(IllegalStateException.class, () -> user.transaction(TYPE, scope -> { + scope.writeValues(values(20, 21)); + insert(scope.connection(), "same"); + insert(scope.connection(), "same"); + return null; + })); + try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + backend.databaseFile())) { + assertEquals(3, value(user.readRow(TYPE), "Points")); + assertEquals(4, value(user.readRow(TYPE), "Votes")); + assertFalse(receipt(connection, "after")); + assertFalse(receipt(connection, "before")); + assertFalse(receipt(connection, "new-user")); + assertFalse(receipt(connection, "same")); + } + } + } + + @Test void concurrentDuplicateAndLostAcknowledgementApplyOnce() throws Exception { + UUID uuid = UUID.randomUUID(); + try (SqliteUserBackend backend = backend(); SqliteUserBackend otherBackend = backend()) { + createReceipts(backend.databaseFile()); + SqlUserStorage user = backend.user(uuid); + // Independent backend instances have no shared JVM operation lock. + SqlUserStorage otherUser = otherBackend.user(uuid); + CountDownLatch ready = new CountDownLatch(2); + CountDownLatch start = new CountDownLatch(1); + ExecutorService workers = Executors.newFixedThreadPool(2); + try { + java.util.function.Function> attempt = candidate -> () -> { + ready.countDown(); + assertTrue(start.await(5, TimeUnit.SECONDS)); + return applyOnce(candidate, "duplicate"); + }; + Future first = workers.submit(attempt.apply(user)); + Future second = workers.submit(attempt.apply(otherUser)); + assertTrue(ready.await(5, TimeUnit.SECONDS)); + start.countDown(); + assertNotEquals(first.get(10, TimeUnit.SECONDS), second.get(10, TimeUnit.SECONDS)); + // The first success acknowledgement is deliberately discarded. + assertFalse(applyOnce(user, "duplicate")); + assertEquals(1, value(user.readRow(TYPE), "Points")); + assertEquals(1, value(user.readRow(TYPE), "Votes")); + try (Connection connection = DriverManager.getConnection("jdbc:sqlite:" + backend.databaseFile()); + PreparedStatement statement = connection.prepareStatement("SELECT COUNT(*) FROM Receipts")) { + try (ResultSet result = statement.executeQuery()) { + assertTrue(result.next()); + assertEquals(1, result.getInt(1)); + } + } + } finally { + start.countDown(); + workers.shutdownNow(); + assertTrue(workers.awaitTermination(5, TimeUnit.SECONDS)); + } + } + } + + private static boolean applyOnce(SqlUserStorage user, String key) { + return user.transaction(TYPE, scope -> { + if (receipt(scope.connection(), key)) return false; + List row = scope.readRow(); + scope.writeValues(values(value(row, "Points") + 1, value(row, "Votes") + 1)); + insert(scope.connection(), key); + return true; + }); + } + + @Test void runtimePublishesOnlyCommittedValuesToItsCacheOwner() throws Exception { + UUID uuid = UUID.randomUUID(); + try (SqliteUserBackend backend = backend()) { + backend.user(uuid).writeValues(TYPE, values(2, 3)); + SimpleCacheOwner cache = new SimpleCacheOwner(); + SharedUserDataRuntime runtime = new SharedUserDataRuntime(backend, cache); + UserDataManager manager = new UserDataManager(mock(AdvancedCorePlugin.class)); + manager.bindSharedRuntime(runtime); + try { + assertEquals(2, runtime.populate(uuid).get("Points").getInt()); + assertThrows(IllegalStateException.class, () -> manager.withAtomicUserTransaction(uuid, TYPE, Map.of(), scope -> { + scope.writeValues(values(8, 9)); + throw new SQLException("abort"); + })); + assertEquals(2, runtime.read(uuid, "Points", UserDataFetchMode.CACHE_ONLY, + null, new DataValueInt(-1)).getInt()); + manager.withAtomicUserTransaction(uuid, TYPE, Map.of(), scope -> { + scope.writeValues(values(10, 11)); + return null; + }); + assertFalse(cache.isCached(uuid)); + assertEquals(10, runtime.read(uuid, "Points", UserDataFetchMode.DEFAULT, + null, new DataValueInt(-1)).getInt()); + assertEquals(11, value(backend.user(uuid).readRow(TYPE), "Votes")); + } finally { manager.getTimer().shutdownNow(); } + } + } + + @Test void productionBukkitSqliteRouteCanCommitCallerRecordWithUserValues() throws Exception { + UUID uuid = UUID.randomUUID(); + try (SqliteUserBackend standalone = backend()) { + createReceipts(standalone.databaseFile()); + String url = "jdbc:sqlite:" + standalone.databaseFile(); + try (Connection legacyConnection = DriverManager.getConnection(url)) { + AdvancedCorePlugin plugin = mock(AdvancedCorePlugin.class); + UserTable table = mock(UserTable.class); + SQLite sqlite = mock(SQLite.class); + UserManager users = mock(UserManager.class); + UserDataManager manager = mock(UserDataManager.class); + when(plugin.getStorageType()).thenReturn(TYPE); + when(plugin.getSQLiteUserTable()).thenReturn(table); + when(plugin.getUserManager()).thenReturn(users); + when(plugin.getLogger()).thenReturn(Logger.getAnonymousLogger()); + when(users.getDataManager()).thenReturn(manager); + when(manager.getKeys()).thenReturn(new java.util.ArrayList<>(List.of( + new UserDataKeyInt("Points"), new UserDataKeyInt("Votes")))); + when(table.getSqLite()).thenReturn(sqlite); + when(table.getName()).thenReturn("Users"); + when(sqlite.getSQLConnection()).thenReturn(legacyConnection); + + SharedUserDataRuntime runtime = new SharedUserDataRuntime(new BukkitSqlUserBackend(plugin), + new SimpleCacheOwner()); + assertEquals("saved", runtime.transaction(uuid, scope -> { + scope.writeValues(values(4, 5)); + insert(scope.connection(), "bukkit"); + return "saved"; + })); + assertEquals(4, value(standalone.user(uuid).readRow(TYPE), "Points")); + try (Connection check = DriverManager.getConnection(url)) { + assertTrue(receipt(check, "bukkit")); + } + } + } + } + + @Test void transactionFencesBypassPublishersAndReopensCacheAfterRollback() throws Exception { + UUID uuid = UUID.randomUUID(); + try (SqliteUserBackend backend = backend()) { + backend.user(uuid).writeValues(TYPE, values(1, 1)); + FencedCacheOwner cache = new FencedCacheOwner(); + SharedUserDataRuntime runtime = new SharedUserDataRuntime(backend, cache); + runtime.populate(uuid); + CountDownLatch inCallback = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + ExecutorService worker = Executors.newSingleThreadExecutor(); + try { + Future transaction = worker.submit(() -> runtime.transaction(uuid, scope -> { + inCallback.countDown(); + try { if (!release.await(5, TimeUnit.SECONDS)) throw new SQLException("timed out"); } + catch (InterruptedException interrupted) { + Thread.currentThread().interrupt(); + throw new SQLException("interrupted", interrupted); + } + scope.writeValues(values(2, 2)); + return null; + })); + assertTrue(inCallback.await(5, TimeUnit.SECONDS)); + assertFalse(cache.tryPublishBypass()); + release.countDown(); + transaction.get(5, TimeUnit.SECONDS); + assertEquals(2, value(backend.user(uuid).readRow(TYPE), "Points")); + runtime.populate(uuid); + assertThrows(IllegalStateException.class, () -> runtime.transaction(uuid, scope -> { + throw new SQLException("rollback"); + })); + assertTrue(cache.tryPublishBypass()); + } finally { + release.countDown(); + worker.shutdownNow(); + assertTrue(worker.awaitTermination(5, TimeUnit.SECONDS)); + } + } + } + + @Test void requiredRowPrerequisiteAllowsQueuedCacheFlushBeforeAtomicWork() throws Exception { + UUID uuid = UUID.randomUUID(); + SqlUserSchema schema = SqlUserSchema.builder() + .column("PlayerName", "VARCHAR(30) NOT NULL", DataType.STRING) + .column("Points", "INTEGER DEFAULT 0", DataType.INTEGER).build(); + try (SqliteUserBackend backend = new SqliteUserBackend(tempDir, "Required", "Users", schema, + SqlBackendLogger.NO_OP)) { + createReceipts(backend.databaseFile()); + PendingCacheOwner cache = new PendingCacheOwner(); + SharedUserDataRuntime runtime = new SharedUserDataRuntime(backend, cache); + UUID pristine = UUID.randomUUID(); + assertThrows(IllegalStateException.class, () -> runtime.transaction(pristine, TYPE, + Map.of("PlayerName", new DataValueString("Ben")), scope -> { + throw new SQLException("abort before prior cache work"); + })); + assertFalse(backend.user(pristine).contains(TYPE)); + runtime.queueChange(uuid, "Points", new DataValueInt(5)); + assertThrows(IllegalStateException.class, () -> runtime.transaction(uuid, TYPE, + Map.of("PlayerName", new DataValueString("Ben")), scope -> { + assertEquals(5, value(scope.readRow(), "Points")); + scope.writeValues(Map.of("Points", new DataValueInt(99))); + insert(scope.connection(), "failed-vote"); + throw new SQLException("abort accepted operation"); + })); + List row = backend.user(uuid).readRow(TYPE); + assertEquals(5, value(row, "Points")); + assertEquals("Ben", row.stream().filter(c -> c.getName().equals("PlayerName")) + .findFirst().orElseThrow().getValue().getString()); + try (Connection check = DriverManager.getConnection("jdbc:sqlite:" + backend.databaseFile())) { + assertFalse(receipt(check, "failed-vote")); + } + runtime.transaction(uuid, TYPE, Map.of("PlayerName", new DataValueString("Ben")), scope -> { + scope.writeValues(Map.of("Points", new DataValueInt(6))); + insert(scope.connection(), "accepted-vote"); + return null; + }); + assertEquals(6, value(backend.user(uuid).readRow(TYPE), "Points")); + try (Connection check = DriverManager.getConnection("jdbc:sqlite:" + backend.databaseFile())) { + assertTrue(receipt(check, "accepted-vote")); + } + } + } + + private static class SimpleCacheOwner implements UserCacheOwner { + private final Map> cached = new java.util.concurrent.ConcurrentHashMap<>(); + @Override public boolean isCached(UUID uuid) { return cached.containsKey(uuid); } + @Override public DataValue getIfPresent(UUID uuid, String key) { + HashMap values = cached.get(uuid); + return values == null ? null : values.get(key); + } + @Override public void populate(UUID uuid, HashMap values) { cached.put(uuid, new HashMap<>(values)); } + @Override public void queueChange(UUID uuid, String key, DataValue value) { cached.get(uuid).put(key, value); } + @Override public void flush(UUID uuid, SqlUserStorage storage) {} + @Override public Set cachedUsers() { return Set.copyOf(cached.keySet()); } + @Override public void remove(UUID uuid) { cached.remove(uuid); } + @Override public void clearAfterFlush() { cached.clear(); } + @Override public void shutdown() {} + } + + private static final class FencedCacheOwner extends SimpleCacheOwner { + private boolean removing; + private boolean queued; + synchronized boolean tryPublishBypass() { + if (removing) return false; + queued = true; + return true; + } + @Override public synchronized void beginRemoval(UUID uuid) { removing = true; } + @Override public synchronized void cancelRemoval(UUID uuid) { removing = false; } + @Override public synchronized void flush(UUID uuid, SqlUserStorage storage) { queued = false; } + @Override public synchronized void remove(UUID uuid) { + if (queued) throw new IllegalStateException("cache has unflushed work"); + super.remove(uuid); + } + } + + private static final class PendingCacheOwner extends SimpleCacheOwner { + private final HashMap pending = new HashMap<>(); + @Override public boolean hasPendingChanges(UUID uuid) { return !pending.isEmpty(); } + @Override public void queueChange(UUID uuid, String key, DataValue value) { + super.queueChange(uuid, key, value); + pending.put(key, value); + } + @Override public void flush(UUID uuid, SqlUserStorage storage) { + if (pending.isEmpty()) return; + storage.writeValues(TYPE, new HashMap<>(pending)); + pending.clear(); + } + } +} From 29ff179c23bac35f203fb1dbdc45eed68995db14 Mon Sep 17 00:00:00 2001 From: BenCodez <17074231+BenCodez@users.noreply.github.com> Date: Sun, 20 Sep 2026 10:08:58 -0600 Subject: [PATCH 2/7] Handle post-commit cache retirement and publish MySQL user indexes --- .../api/user/usercache/UserDataManager.java | 5 ++ .../api/user/userstorage/mysql/MySQL.java | 6 ++ .../user/runtime/BukkitUserCacheOwner.java | 4 + .../user/storage/BukkitSqlUserBackend.java | 18 ++++- .../user/runtime/SharedUserDataRuntime.java | 27 ++++++- .../core/user/runtime/UserCacheOwner.java | 5 +- .../storage/SqliteUserTransactionTest.java | 79 +++++++++++++++++++ 7 files changed, 138 insertions(+), 6 deletions(-) diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/usercache/UserDataManager.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/usercache/UserDataManager.java index 5a7fd394b..8204df656 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/usercache/UserDataManager.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/usercache/UserDataManager.java @@ -985,6 +985,11 @@ private void reportDeferredStorageFailure(Throwable failure) { } } + /** Retain a cache reconciliation failure after its SQL transaction committed. */ + public final void recordSharedStorageFailure(Throwable failure) { + reportDeferredStorageFailure(Objects.requireNonNull(failure, "failure")); + } + /** Last asynchronous cache-cleanup failure, retained for diagnosis and recovery. */ public Throwable getLastDeferredStorageFailure() { return lastDeferredStorageFailure.get(); } diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java index f899fb516..d23723655 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java @@ -360,6 +360,12 @@ public void forEachUser(java.util.function.BiConsumer> p // Keep existing methods (getUuids / getUUID / etc.) // ------------------------- + /** Publish a user committed through the shared JDBC transaction route. */ + public void recordCommittedUser(UUID uuid, String playerName) { + uuids.add(uuid.toString()); + if (playerName != null && !playerName.isEmpty()) names.add(playerName); + } + public Set getUuids() { if (uuids == null || uuids.isEmpty()) { uuids.clear(); diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/runtime/BukkitUserCacheOwner.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/runtime/BukkitUserCacheOwner.java index 1e3c626f5..d842dca02 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/runtime/BukkitUserCacheOwner.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/runtime/BukkitUserCacheOwner.java @@ -258,6 +258,10 @@ private record CachePopulation(UUID uuid, UserDataCache cache, long version) imp for (UUID uuid : Set.copyOf(pendingNotifications.keySet())) dispatchNotifications(uuid); } + @Override public void reportCommittedFailure(UUID uuid, Throwable failure) { + manager.recordSharedStorageFailure(failure); + } + @Override public void discardAllNotifications() { pendingNotifications.clear(); } @Override public Set cachedUsers() { return new HashSet<>(manager.getUserDataCache().keySet()); } diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java index f4f4298da..230fe4ff8 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java @@ -181,8 +181,22 @@ private T transaction(UserStorage storage, UUID uuid, java.util.Map { + T value = work.run(scope); + for (Column column : scope.readRow()) { + if ("PlayerName".equalsIgnoreCase(column.getName()) && column.getValue() != null) { + committedName[0] = column.getValue().getString(); + break; + } + } + return value; + }); + // The native enumerations are populated by its own write path. This + // transaction bypasses that path, so publish only after JDBC commit. + mysql().recordCommittedUser(uuid, committedName[0]); + return result; } synchronized (sqliteOperations) { // The legacy SQLite provider retains one shared connection. Obtain diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/SharedUserDataRuntime.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/SharedUserDataRuntime.java index bbd03e013..50ce28ff7 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/SharedUserDataRuntime.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/SharedUserDataRuntime.java @@ -166,13 +166,34 @@ private T transactionInternal(UUID uuid, UserStorage expectedStorage, Map null); } flushInternal(uuid); - T result = backend.user(uuid).transaction(backend.storageType(), initialValues, work); - cacheOwner.remove(uuid); - return result; } catch (RuntimeException | Error failure) { cacheOwner.cancelRemoval(uuid); throw failure; } + // SQL has committed. Cache retirement can fail, but reporting a + // transaction failure here would invite a duplicate caller retry. + T result; + try { + result = backend.user(uuid).transaction(backend.storageType(), initialValues, work); + } catch (RuntimeException | Error failure) { + cacheOwner.cancelRemoval(uuid); + throw failure; + } + try { + cacheOwner.remove(uuid); + } catch (RuntimeException | Error failure) { + try { + cacheOwner.populate(uuid, SqlUserDataAccess.convert(readStorageRow(uuid))); + } catch (RuntimeException | Error recoveryFailure) { + failure.addSuppressed(recoveryFailure); + } finally { + try { cacheOwner.cancelRemoval(uuid); } + catch (RuntimeException | Error recoveryFailure) { failure.addSuppressed(recoveryFailure); } + } + try { cacheOwner.reportCommittedFailure(uuid, failure); } + catch (RuntimeException | Error reportingFailure) { failure.addSuppressed(reportingFailure); } + } + return result; }); } finally { cacheOwner.dispatchNotifications(uuid); } } diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/UserCacheOwner.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/UserCacheOwner.java index 050489f28..948fd8123 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/UserCacheOwner.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/runtime/UserCacheOwner.java @@ -59,7 +59,10 @@ default void bindLifecycle(SqlUserBackend backend, Consumer gate, default void requireBlockingAllowed() {} /** Deliver callbacks accumulated by a flush after its user admission is released. */ - default void dispatchNotifications(UUID uuid) {} + default void dispatchNotifications(UUID uuid) {} + + /** Record a cache failure after SQL has committed without making it retryable. */ + default void reportCommittedFailure(UUID uuid, Throwable failure) {} /** Deliver callbacks accumulated by a lifecycle-wide flush after all admission is released. */ default void dispatchAllNotifications() {} diff --git a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java index 86862f6a2..24201c452 100644 --- a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java +++ b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java @@ -8,6 +8,7 @@ import java.sql.DriverManager; import java.sql.PreparedStatement; import java.sql.ResultSet; +import java.sql.ResultSetMetaData; import java.sql.SQLException; import java.util.HashMap; import java.util.List; @@ -30,6 +31,8 @@ import com.bencodez.advancedcore.api.user.UserManager; import com.bencodez.advancedcore.api.user.usercache.UserDataManager; import com.bencodez.advancedcore.api.user.usercache.keys.UserDataKeyInt; +import com.bencodez.advancedcore.api.user.usercache.keys.UserDataKeyString; +import com.bencodez.advancedcore.api.user.userstorage.mysql.MySQL; import com.bencodez.advancedcore.api.user.userstorage.sql.UserTable; import com.bencodez.advancedcore.bukkit.user.storage.BukkitSqlUserBackend; import com.bencodez.advancedcore.api.user.UserDataFetchMode; @@ -45,6 +48,7 @@ import com.bencodez.simpleapi.sql.data.DataValueInt; import com.bencodez.simpleapi.sql.data.DataValueString; import com.bencodez.simpleapi.sql.sqlite.db.SQLite; +import com.bencodez.simpleapi.sql.mysql.DbType; @Timeout(20) class SqliteUserTransactionTest { @@ -268,6 +272,75 @@ private static boolean applyOnce(SqlUserStorage user, String key) { } } + @Test void committedReceiptIsReturnedWhenCacheRetirementFails() throws Exception { + UUID uuid = UUID.randomUUID(); + try (SqliteUserBackend backend = backend()) { + createReceipts(backend.databaseFile()); + backend.user(uuid).writeValues(TYPE, values(2, 3)); + FailingRetirementCacheOwner cache = new FailingRetirementCacheOwner(); + SharedUserDataRuntime runtime = new SharedUserDataRuntime(backend, cache); + runtime.populate(uuid); + boolean applied = runtime.transaction(uuid, scope -> { + if (receipt(scope.connection(), "accepted")) return false; + scope.writeValues(values(3, 4)); + insert(scope.connection(), "accepted"); + return true; + }); + assertTrue(applied); + assertEquals(3, runtime.read(uuid, "Points", UserDataFetchMode.CACHE_ONLY, + null, new DataValueInt(-1)).getInt()); + assertNotNull(cache.committedFailure); + boolean repeated = runtime.transaction(uuid, scope -> { + if (receipt(scope.connection(), "accepted")) return false; + scope.writeValues(values(4, 5)); + return true; + }); + assertFalse(repeated); + assertEquals(3, value(backend.user(uuid).readRow(TYPE), "Points")); + } + } + + @Test void productionMysqlRoutePublishesIdentityOnlyAfterCommit() throws Exception { + UUID uuid = UUID.randomUUID(); + AdvancedCorePlugin plugin = mock(AdvancedCorePlugin.class); + MySQL mysql = mock(MySQL.class, RETURNS_DEEP_STUBS); + UserManager users = mock(UserManager.class); + UserDataManager manager = mock(UserDataManager.class); + Connection connection = mock(Connection.class); + PreparedStatement exists = mock(PreparedStatement.class); + PreparedStatement row = mock(PreparedStatement.class); + ResultSet existsResult = mock(ResultSet.class); + ResultSet rowResult = mock(ResultSet.class); + ResultSetMetaData metadata = mock(ResultSetMetaData.class); + when(plugin.getUserManager()).thenReturn(users); + when(plugin.getLogger()).thenReturn(Logger.getAnonymousLogger()); + when(users.getDataManager()).thenReturn(manager); + when(manager.getKeys()).thenReturn(new java.util.ArrayList<>(List.of(new UserDataKeyString("PlayerName")))); + when(mysql.getTableName()).thenReturn("Users"); + when(mysql.getMysql().getConnectionManager().getConnection()).thenReturn(connection); + when(mysql.getMysql().getConnectionManager().getDbType()).thenReturn(DbType.MYSQL); + when(connection.prepareStatement(org.mockito.ArgumentMatchers.contains("FOR UPDATE"))).thenReturn(exists); + when(connection.prepareStatement(org.mockito.ArgumentMatchers.startsWith("SELECT *"))).thenReturn(row); + when(exists.executeQuery()).thenReturn(existsResult); + when(existsResult.next()).thenReturn(true); + when(row.executeQuery()).thenReturn(rowResult); + when(rowResult.next()).thenReturn(true); + when(rowResult.getMetaData()).thenReturn(metadata); + when(metadata.getColumnCount()).thenReturn(1); + when(metadata.getColumnLabel(1)).thenReturn("PlayerName"); + when(rowResult.getString(1)).thenReturn("Ben"); + BukkitSqlUserBackend backend = new BukkitSqlUserBackend(plugin, UserStorage.MYSQL, mysql, null); + assertEquals("accepted", backend.user(uuid).transaction(UserStorage.MYSQL, scope -> "accepted")); + org.mockito.InOrder order = inOrder(connection, mysql); + order.verify(connection).commit(); + order.verify(mysql).recordCommittedUser(uuid, "Ben"); + clearInvocations(connection, mysql); + assertThrows(IllegalStateException.class, () -> backend.user(uuid).transaction(UserStorage.MYSQL, + scope -> { throw new SQLException("receipt failed"); })); + verify(connection).rollback(); + verify(mysql, never()).recordCommittedUser(any(), any()); + } + @Test void transactionFencesBypassPublishersAndReopensCacheAfterRollback() throws Exception { UUID uuid = UUID.randomUUID(); try (SqliteUserBackend backend = backend()) { @@ -366,6 +439,12 @@ private static class SimpleCacheOwner implements UserCacheOwner { @Override public void shutdown() {} } + private static final class FailingRetirementCacheOwner extends SimpleCacheOwner { + private Throwable committedFailure; + @Override public void remove(UUID uuid) { throw new IllegalStateException("cache retirement failed"); } + @Override public void reportCommittedFailure(UUID uuid, Throwable failure) { committedFailure = failure; } + } + private static final class FencedCacheOwner extends SimpleCacheOwner { private boolean removing; private boolean queued; From ffc5e39dd9bf7b08b3d294cc71626c5001210e48 Mon Sep 17 00:00:00 2001 From: BenCodez <17074231+BenCodez@users.noreply.github.com> Date: Sun, 20 Sep 2026 10:20:06 -0600 Subject: [PATCH 3/7] Propagate shared cache flush failures before receipt transactions --- .../user/storage/BukkitSqlUserBackend.java | 62 ++++++++----- .../storage/SqliteUserTransactionTest.java | 89 +++++++++++++++++++ 2 files changed, 128 insertions(+), 23 deletions(-) diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java index 230fe4ff8..fe52b44e6 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java @@ -159,19 +159,49 @@ private void writeValues(UserStorage storage, UUID uuid, HashMap columns = new ArrayList<>(); + java.util.Map updates = new HashMap<>(); values.forEach((key, value) -> { - if (!"uuid".equalsIgnoreCase(key) && value != null) columns.add(new Column(key, value)); + if (!"uuid".equalsIgnoreCase(key) && value != null) updates.put(key, value); }); - if (columns.isEmpty()) return; - if (storage == UserStorage.MYSQL) mysql().update(uuid.toString(), columns, false); - else synchronized (sqliteOperations) { table().update(primary(uuid), columns); } + if (updates.isEmpty()) return; + // Shared cache flushes must surface SQL failures. The legacy update + // methods log and swallow them, which could acknowledge a later caller + // transaction while its prerequisite queued values were never stored. + withSqlUser(storage, uuid, user -> { + user.writeValues(storage, new HashMap<>(updates)); + return null; + }); + if (storage == UserStorage.MYSQL) { + DataValue name = updates.entrySet().stream() + .filter(entry -> "PlayerName".equalsIgnoreCase(entry.getKey())) + .map(java.util.Map.Entry::getValue).findFirst().orElse(null); + mysql().recordCommittedUser(uuid, name != null && name.isString() ? name.getString() : null); + } } private T transaction(UserStorage storage, UUID uuid, java.util.Map initialValues, SqlUserStorage.TransactionWork work) { requireOpen(); requireStorage(storage); Objects.requireNonNull(work, "work"); + return withSqlUser(storage, uuid, user -> { + if (storage != UserStorage.MYSQL) return user.transaction(storage, initialValues, work); + String[] committedName = new String[1]; + T result = user.transaction(storage, initialValues, scope -> { + T value = work.run(scope); + for (Column column : scope.readRow()) { + if ("PlayerName".equalsIgnoreCase(column.getName()) && column.getValue() != null) { + committedName[0] = column.getValue().getString(); + break; + } + } + return value; + }); + mysql().recordCommittedUser(uuid, committedName[0]); + return result; + }); + } + + private T withSqlUser(UserStorage storage, UUID uuid, java.util.function.Function work) { SqlUserSchema schema = SqlUserSchema.fromKeys(plugin.getUserManager().getDataManager().getKeys()); SqlBackendLogger logger = new SqlBackendLogger() { @Override public void info(String message) { plugin.getLogger().info(message); } @@ -181,22 +211,8 @@ private T transaction(UserStorage storage, UUID uuid, java.util.Map { - T value = work.run(scope); - for (Column column : scope.readRow()) { - if ("PlayerName".equalsIgnoreCase(column.getName()) && column.getValue() != null) { - committedName[0] = column.getValue().getString(); - break; - } - } - return value; - }); - // The native enumerations are populated by its own write path. This - // transaction bypasses that path, so publish only after JDBC commit. - mysql().recordCommittedUser(uuid, committedName[0]); - return result; + return work.apply(SqlUserBackendFactory.existingUser(storage, uuid, mysql().getTableName(), schema, + manager::getConnection, manager.getDbType(), logger)); } synchronized (sqliteOperations) { // The legacy SQLite provider retains one shared connection. Obtain @@ -205,8 +221,8 @@ private T transaction(UserStorage storage, UUID uuid, java.util.Map DriverManager.getConnection(url), null, logger).transaction(storage, initialValues, work); + return work.apply(SqlUserBackendFactory.existingUser(storage, uuid, table().getName(), schema, + () -> DriverManager.getConnection(url), null, logger)); } } diff --git a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java index 24201c452..23b7dfa80 100644 --- a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java +++ b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java @@ -300,6 +300,95 @@ private static boolean applyOnce(SqlUserStorage user, String key) { } } + @Test void productionCacheFlushFailureStopsCallerTransaction() throws Exception { + UUID uuid = UUID.randomUUID(); + try (SqliteUserBackend standalone = backend()) { + createReceipts(standalone.databaseFile()); + standalone.user(uuid).writeValues(TYPE, values(2, 3)); + String url = "jdbc:sqlite:" + standalone.databaseFile(); + try (Connection legacyConnection = DriverManager.getConnection(url); + Connection setup = DriverManager.getConnection(url); + PreparedStatement trigger = setup.prepareStatement( + "CREATE TRIGGER reject_points BEFORE UPDATE ON Users " + + "BEGIN SELECT RAISE(ABORT, 'write rejected'); END")) { + trigger.executeUpdate(); + AdvancedCorePlugin plugin = mock(AdvancedCorePlugin.class); + UserTable table = mock(UserTable.class); + SQLite sqlite = mock(SQLite.class); + UserManager users = mock(UserManager.class); + UserDataManager manager = mock(UserDataManager.class); + when(plugin.getSQLiteUserTable()).thenReturn(table); + when(plugin.getUserManager()).thenReturn(users); + when(plugin.getLogger()).thenReturn(Logger.getAnonymousLogger()); + when(users.getDataManager()).thenReturn(manager); + when(manager.getKeys()).thenReturn(new java.util.ArrayList<>(List.of( + new UserDataKeyInt("Points"), new UserDataKeyInt("Votes")))); + when(table.getSqLite()).thenReturn(sqlite); + when(table.getName()).thenReturn("Users"); + when(sqlite.getSQLConnection()).thenReturn(legacyConnection); + PendingCacheOwner cache = new PendingCacheOwner(); + SharedUserDataRuntime runtime = new SharedUserDataRuntime( + new BukkitSqlUserBackend(plugin, TYPE, null, table), cache); + cache.populate(uuid, values(2, 3)); + cache.queueChange(uuid, "Points", new DataValueInt(7)); + java.util.concurrent.atomic.AtomicBoolean callbackRan = new java.util.concurrent.atomic.AtomicBoolean(); + assertThrows(IllegalStateException.class, () -> runtime.transaction(uuid, scope -> { + callbackRan.set(true); + insert(scope.connection(), "should-not-commit"); + return null; + })); + assertFalse(callbackRan.get()); + assertEquals(2, value(standalone.user(uuid).readRow(TYPE), "Points")); + try (Connection check = DriverManager.getConnection(url)) { + assertFalse(receipt(check, "should-not-commit")); + } + assertTrue(cache.hasPendingChanges(uuid)); + } + } + } + + @Test void productionFirstCacheFlushSuppliesRequiredColumns() throws Exception { + UUID uuid = UUID.randomUUID(); + String url = "jdbc:sqlite:" + tempDir.resolve("required-users.db"); + try (Connection setup = DriverManager.getConnection(url); + PreparedStatement create = setup.prepareStatement( + "CREATE TABLE Users (uuid TEXT PRIMARY KEY, PlayerName TEXT NOT NULL, Points INTEGER)")) { + create.executeUpdate(); + } + try (Connection legacyConnection = DriverManager.getConnection(url)) { + AdvancedCorePlugin plugin = mock(AdvancedCorePlugin.class); + UserTable table = mock(UserTable.class); + SQLite sqlite = mock(SQLite.class); + UserManager users = mock(UserManager.class); + UserDataManager manager = mock(UserDataManager.class); + when(plugin.getUserManager()).thenReturn(users); + when(plugin.getLogger()).thenReturn(Logger.getAnonymousLogger()); + when(users.getDataManager()).thenReturn(manager); + when(manager.getKeys()).thenReturn(new java.util.ArrayList<>(List.of( + new UserDataKeyString("PlayerName"), new UserDataKeyInt("Points")))); + when(table.getSqLite()).thenReturn(sqlite); + when(table.getName()).thenReturn("Users"); + when(sqlite.getSQLConnection()).thenReturn(legacyConnection); + PendingCacheOwner cache = new PendingCacheOwner(); + SharedUserDataRuntime runtime = new SharedUserDataRuntime( + new BukkitSqlUserBackend(plugin, TYPE, null, table), cache); + cache.populate(uuid, new HashMap<>()); + cache.queueChange(uuid, "PlayerName", new DataValueString("Ben")); + cache.queueChange(uuid, "Points", new DataValueInt(7)); + runtime.flush(uuid); + assertFalse(cache.hasPendingChanges(uuid)); + try (Connection check = DriverManager.getConnection(url); + PreparedStatement select = check.prepareStatement("SELECT PlayerName, Points FROM Users WHERE uuid=?")) { + select.setString(1, uuid.toString()); + try (ResultSet row = select.executeQuery()) { + assertTrue(row.next()); + assertEquals("Ben", row.getString(1)); + assertEquals(7, row.getInt(2)); + } + } + } + } + @Test void productionMysqlRoutePublishesIdentityOnlyAfterCommit() throws Exception { UUID uuid = UUID.randomUUID(); AdvancedCorePlugin plugin = mock(AdvancedCorePlugin.class); From b78cb029fb32c2e10a03ab4e859e572e8a648691 Mon Sep 17 00:00:00 2001 From: BenCodez <17074231+BenCodez@users.noreply.github.com> Date: Sun, 20 Sep 2026 10:32:26 -0600 Subject: [PATCH 4/7] Retain dynamic user columns in strict shared cache flushes --- .../user/storage/BukkitSqlUserBackend.java | 24 +++++++++++++++---- .../storage/SqliteUserTransactionTest.java | 13 +++++++++- 2 files changed, 32 insertions(+), 5 deletions(-) diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java index fe52b44e6..c72985a14 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java @@ -167,7 +167,7 @@ private void writeValues(UserStorage storage, UUID uuid, HashMap { + withSqlUser(storage, uuid, updates, user -> { user.writeValues(storage, new HashMap<>(updates)); return null; }); @@ -183,7 +183,7 @@ private T transaction(UserStorage storage, UUID uuid, java.util.Map { + return withSqlUser(storage, uuid, java.util.Map.of(), user -> { if (storage != UserStorage.MYSQL) return user.transaction(storage, initialValues, work); String[] committedName = new String[1]; T result = user.transaction(storage, initialValues, scope -> { @@ -201,8 +201,22 @@ private T transaction(UserStorage storage, UUID uuid, java.util.Map T withSqlUser(UserStorage storage, UUID uuid, java.util.function.Function work) { - SqlUserSchema schema = SqlUserSchema.fromKeys(plugin.getUserManager().getDataManager().getKeys()); + private T withSqlUser(UserStorage storage, UUID uuid, java.util.Map updates, + java.util.function.Function work) { + SqlUserSchema registered = SqlUserSchema.fromKeys(plugin.getUserManager().getDataManager().getKeys()); + SqlUserSchema.Builder builder = SqlUserSchema.builder(); + for (SqlUserSchema.ColumnDefinition column : registered.columns()) { + if (!"uuid".equalsIgnoreCase(column.name())) builder.column(column.name(), column.sqlType(), column.dataType()); + } + ArrayList dynamic = new ArrayList<>(); + for (java.util.Map.Entry entry : updates.entrySet()) { + if (registered.contains(entry.getKey())) continue; + com.bencodez.simpleapi.sql.DataType type = entry.getValue().getType(); + builder.column(entry.getKey(), type == com.bencodez.simpleapi.sql.DataType.STRING ? "TEXT" : + type == com.bencodez.simpleapi.sql.DataType.INTEGER ? "INTEGER" : "BOOLEAN", type); + dynamic.add(new Column(entry.getKey(), type)); + } + SqlUserSchema schema = builder.build(); SqlBackendLogger logger = new SqlBackendLogger() { @Override public void info(String message) { plugin.getLogger().info(message); } @Override public void warn(String message, Throwable error) { @@ -210,11 +224,13 @@ private T withSqlUser(UserStorage storage, UUID uuid, java.util.function.Fun } }; if (storage == UserStorage.MYSQL) { + for (Column column : dynamic) mysql().checkColumn(column.getName(), column.getDataType()); var manager = mysql().getMysql().getConnectionManager(); return work.apply(SqlUserBackendFactory.existingUser(storage, uuid, mysql().getTableName(), schema, manager::getConnection, manager.getDbType(), logger)); } synchronized (sqliteOperations) { + for (Column column : dynamic) table().checkColumn(column); // The legacy SQLite provider retains one shared connection. Obtain // its actual database URL, then open a separate transaction-owned // connection so AdvancedCore can close it after commit/rollback. diff --git a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java index 23b7dfa80..dbfbc7e82 100644 --- a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java +++ b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java @@ -369,21 +369,32 @@ private static boolean applyOnce(SqlUserStorage user, String key) { when(table.getSqLite()).thenReturn(sqlite); when(table.getName()).thenReturn("Users"); when(sqlite.getSQLConnection()).thenReturn(legacyConnection); + doAnswer(call -> { + Column column = call.getArgument(0); + if ("RepeatSpecial".equals(column.getName())) { + try (PreparedStatement alter = legacyConnection.prepareStatement( + "ALTER TABLE Users ADD COLUMN RepeatSpecial INTEGER")) { alter.executeUpdate(); } + } + return null; + }).when(table).checkColumn(any(Column.class)); PendingCacheOwner cache = new PendingCacheOwner(); SharedUserDataRuntime runtime = new SharedUserDataRuntime( new BukkitSqlUserBackend(plugin, TYPE, null, table), cache); cache.populate(uuid, new HashMap<>()); cache.queueChange(uuid, "PlayerName", new DataValueString("Ben")); cache.queueChange(uuid, "Points", new DataValueInt(7)); + cache.queueChange(uuid, "RepeatSpecial", new DataValueInt(3)); runtime.flush(uuid); assertFalse(cache.hasPendingChanges(uuid)); try (Connection check = DriverManager.getConnection(url); - PreparedStatement select = check.prepareStatement("SELECT PlayerName, Points FROM Users WHERE uuid=?")) { + PreparedStatement select = check.prepareStatement( + "SELECT PlayerName, Points, RepeatSpecial FROM Users WHERE uuid=?")) { select.setString(1, uuid.toString()); try (ResultSet row = select.executeQuery()) { assertTrue(row.next()); assertEquals("Ben", row.getString(1)); assertEquals(7, row.getInt(2)); + assertEquals(3, row.getInt(3)); } } } From a90d9e397ace14149920faee133380a756d2dfde Mon Sep 17 00:00:00 2001 From: BenCodez <17074231+BenCodez@users.noreply.github.com> Date: Sun, 20 Sep 2026 10:39:25 -0600 Subject: [PATCH 5/7] Invalidate MySQL name cache after committed user changes --- .../api/user/userstorage/mysql/MySQL.java | 19 ++++++----- .../user/storage/BukkitSqlUserBackend.java | 21 ++---------- .../storage/SqliteUserTransactionTest.java | 34 ++++++++++++------- 3 files changed, 34 insertions(+), 40 deletions(-) diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java index d23723655..338845b6f 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java @@ -361,9 +361,11 @@ public void forEachUser(java.util.function.BiConsumer> p // ------------------------- /** Publish a user committed through the shared JDBC transaction route. */ - public void recordCommittedUser(UUID uuid, String playerName) { + public void recordCommittedUser(UUID uuid) { uuids.add(uuid.toString()); - if (playerName != null && !playerName.isEmpty()) names.add(playerName); + // A prior name may belong to this UUID. Reload the complete set from + // committed SQL on its next read instead of retaining that stale name. + synchronized (names) { names.clear(); } } public Set getUuids() { @@ -391,13 +393,12 @@ public ArrayList getUuidsQuery() { return out; } - public Set getNames() { - if (names == null || names.isEmpty()) { - names.clear(); - names.addAll(getNamesQuery()); - } - return names; - } + public Set getNames() { + synchronized (names) { + if (names.isEmpty()) names.addAll(getNamesQuery()); + return names; + } + } public ArrayList getNamesQuery() { ArrayList out = new ArrayList<>(); diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java index c72985a14..12807972d 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java @@ -171,12 +171,7 @@ private void writeValues(UserStorage storage, UUID uuid, HashMap(updates)); return null; }); - if (storage == UserStorage.MYSQL) { - DataValue name = updates.entrySet().stream() - .filter(entry -> "PlayerName".equalsIgnoreCase(entry.getKey())) - .map(java.util.Map.Entry::getValue).findFirst().orElse(null); - mysql().recordCommittedUser(uuid, name != null && name.isString() ? name.getString() : null); - } + if (storage == UserStorage.MYSQL) mysql().recordCommittedUser(uuid); } private T transaction(UserStorage storage, UUID uuid, java.util.Map initialValues, SqlUserStorage.TransactionWork work) { @@ -185,18 +180,8 @@ private T transaction(UserStorage storage, UUID uuid, java.util.Map { if (storage != UserStorage.MYSQL) return user.transaction(storage, initialValues, work); - String[] committedName = new String[1]; - T result = user.transaction(storage, initialValues, scope -> { - T value = work.run(scope); - for (Column column : scope.readRow()) { - if ("PlayerName".equalsIgnoreCase(column.getName()) && column.getValue() != null) { - committedName[0] = column.getValue().getString(); - break; - } - } - return value; - }); - mysql().recordCommittedUser(uuid, committedName[0]); + T result = user.transaction(storage, initialValues, work); + mysql().recordCommittedUser(uuid); return result; }); } diff --git a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java index dbfbc7e82..589a68c99 100644 --- a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java +++ b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java @@ -8,7 +8,6 @@ import java.sql.DriverManager; import java.sql.PreparedStatement; import java.sql.ResultSet; -import java.sql.ResultSetMetaData; import java.sql.SQLException; import java.util.HashMap; import java.util.List; @@ -408,10 +407,7 @@ private static boolean applyOnce(SqlUserStorage user, String key) { UserDataManager manager = mock(UserDataManager.class); Connection connection = mock(Connection.class); PreparedStatement exists = mock(PreparedStatement.class); - PreparedStatement row = mock(PreparedStatement.class); ResultSet existsResult = mock(ResultSet.class); - ResultSet rowResult = mock(ResultSet.class); - ResultSetMetaData metadata = mock(ResultSetMetaData.class); when(plugin.getUserManager()).thenReturn(users); when(plugin.getLogger()).thenReturn(Logger.getAnonymousLogger()); when(users.getDataManager()).thenReturn(manager); @@ -420,25 +416,37 @@ private static boolean applyOnce(SqlUserStorage user, String key) { when(mysql.getMysql().getConnectionManager().getConnection()).thenReturn(connection); when(mysql.getMysql().getConnectionManager().getDbType()).thenReturn(DbType.MYSQL); when(connection.prepareStatement(org.mockito.ArgumentMatchers.contains("FOR UPDATE"))).thenReturn(exists); - when(connection.prepareStatement(org.mockito.ArgumentMatchers.startsWith("SELECT *"))).thenReturn(row); when(exists.executeQuery()).thenReturn(existsResult); when(existsResult.next()).thenReturn(true); - when(row.executeQuery()).thenReturn(rowResult); - when(rowResult.next()).thenReturn(true); - when(rowResult.getMetaData()).thenReturn(metadata); - when(metadata.getColumnCount()).thenReturn(1); - when(metadata.getColumnLabel(1)).thenReturn("PlayerName"); - when(rowResult.getString(1)).thenReturn("Ben"); BukkitSqlUserBackend backend = new BukkitSqlUserBackend(plugin, UserStorage.MYSQL, mysql, null); assertEquals("accepted", backend.user(uuid).transaction(UserStorage.MYSQL, scope -> "accepted")); org.mockito.InOrder order = inOrder(connection, mysql); order.verify(connection).commit(); - order.verify(mysql).recordCommittedUser(uuid, "Ben"); + order.verify(mysql).recordCommittedUser(uuid); clearInvocations(connection, mysql); assertThrows(IllegalStateException.class, () -> backend.user(uuid).transaction(UserStorage.MYSQL, scope -> { throw new SQLException("receipt failed"); })); verify(connection).rollback(); - verify(mysql, never()).recordCommittedUser(any(), any()); + verify(mysql, never()).recordCommittedUser(any()); + } + + @Test void committedMysqlIdentityInvalidatesStaleNameCache() throws Exception { + MySQL mysql = mock(MySQL.class, CALLS_REAL_METHODS); + Set names = java.util.concurrent.ConcurrentHashMap.newKeySet(); + Set uuids = java.util.concurrent.ConcurrentHashMap.newKeySet(); + java.lang.reflect.Field namesField = MySQL.class.getDeclaredField("names"); + java.lang.reflect.Field uuidsField = MySQL.class.getDeclaredField("uuids"); + namesField.setAccessible(true); + uuidsField.setAccessible(true); + namesField.set(mysql, names); + uuidsField.set(mysql, uuids); + doReturn(new java.util.ArrayList<>(List.of("New"))).when(mysql).getNamesQuery(); + names.add("Old"); + UUID uuid = UUID.randomUUID(); + mysql.recordCommittedUser(uuid); + assertTrue(uuids.contains(uuid.toString())); + assertFalse(names.contains("Old")); + assertEquals(Set.of("New"), mysql.getNames()); } @Test void transactionFencesBypassPublishersAndReopensCacheAfterRollback() throws Exception { From 985ae094390f8e2f0317b7d758340fd0fa246ace Mon Sep 17 00:00:00 2001 From: BenCodez <17074231+BenCodez@users.noreply.github.com> Date: Sun, 20 Sep 2026 10:47:43 -0600 Subject: [PATCH 6/7] Invalidate name snapshots only for committed name changes --- .../api/user/userstorage/mysql/MySQL.java | 10 ++++----- .../user/storage/BukkitSqlUserBackend.java | 22 ++++++++++++++++--- .../core/user/storage/SqlUserStorage.java | 3 +++ .../user/storage/sql/JdbcSqlUserStorage.java | 11 +++++++--- .../storage/SqliteUserTransactionTest.java | 22 ++++++++++++++++--- 5 files changed, 53 insertions(+), 15 deletions(-) diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java index 338845b6f..57e697617 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java @@ -360,12 +360,10 @@ public void forEachUser(java.util.function.BiConsumer> p // Keep existing methods (getUuids / getUUID / etc.) // ------------------------- - /** Publish a user committed through the shared JDBC transaction route. */ - public void recordCommittedUser(UUID uuid) { + /** Publish a committed UUID; refresh names only when PlayerName may have changed. */ + public void recordCommittedUser(UUID uuid, boolean nameMayHaveChanged) { uuids.add(uuid.toString()); - // A prior name may belong to this UUID. Reload the complete set from - // committed SQL on its next read instead of retaining that stale name. - synchronized (names) { names.clear(); } + if (nameMayHaveChanged) synchronized (names) { names.clear(); } } public Set getUuids() { @@ -396,7 +394,7 @@ public ArrayList getUuidsQuery() { public Set getNames() { synchronized (names) { if (names.isEmpty()) names.addAll(getNamesQuery()); - return names; + return new java.util.HashSet<>(names); } } diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java index 12807972d..943b93f3f 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java @@ -171,7 +171,7 @@ private void writeValues(UserStorage storage, UUID uuid, HashMap(updates)); return null; }); - if (storage == UserStorage.MYSQL) mysql().recordCommittedUser(uuid); + if (storage == UserStorage.MYSQL) mysql().recordCommittedUser(uuid, containsPlayerName(updates)); } private T transaction(UserStorage storage, UUID uuid, java.util.Map initialValues, SqlUserStorage.TransactionWork work) { @@ -180,12 +180,28 @@ private T transaction(UserStorage storage, UUID uuid, java.util.Map { if (storage != UserStorage.MYSQL) return user.transaction(storage, initialValues, work); - T result = user.transaction(storage, initialValues, work); - mysql().recordCommittedUser(uuid); + boolean[] nameTouched = {false}; + T result = user.transaction(storage, initialValues, scope -> { + nameTouched[0] = scope.createdUserRow() && containsPlayerName(initialValues); + return work.run(new SqlUserStorage.TransactionScope() { + @Override public boolean createdUserRow() { return scope.createdUserRow(); } + @Override public java.sql.Connection connection() { return scope.connection(); } + @Override public List readRow() throws SQLException { return scope.readRow(); } + @Override public void writeValues(java.util.Map values) throws SQLException { + scope.writeValues(values); + if (containsPlayerName(values)) nameTouched[0] = true; + } + }); + }); + mysql().recordCommittedUser(uuid, nameTouched[0]); return result; }); } + private static boolean containsPlayerName(java.util.Map values) { + return values.keySet().stream().anyMatch("PlayerName"::equalsIgnoreCase); + } + private T withSqlUser(UserStorage storage, UUID uuid, java.util.Map updates, java.util.function.Function work) { SqlUserSchema registered = SqlUserSchema.fromKeys(plugin.getUserManager().getDataManager().getKeys()); diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/SqlUserStorage.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/SqlUserStorage.java index 228d803a5..2c4612b2e 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/SqlUserStorage.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/SqlUserStorage.java @@ -60,6 +60,9 @@ interface TransactionWork { } interface TransactionScope { + /** True only when this transaction inserted the user row. */ + default boolean createdUserRow() { return false; } + /** * Active connection for caller-owned tables. Never commit, roll back, * close, or change auto-commit. AdvancedCore owns its lifecycle. Do not diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/JdbcSqlUserStorage.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/JdbcSqlUserStorage.java index 63c6da2d5..e472ff238 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/JdbcSqlUserStorage.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/JdbcSqlUserStorage.java @@ -134,8 +134,8 @@ private List readRow(Connection connection) throws SQLException { Objects.requireNonNull(work, "work"); Map seed = canonicalize(Objects.requireNonNull(initialValues, "initialValues")); return inTransaction("run user transaction", connection -> { - ensureRow(connection, seed); - Scope scope = new Scope(connection); + boolean createdUserRow = !ensureRow(connection, seed); + Scope scope = new Scope(connection, createdUserRow); try { return work.run(scope); } finally { scope.active = false; } }); @@ -143,9 +143,14 @@ private List readRow(Connection connection) throws SQLException { private final class Scope implements TransactionScope { private final Connection connection; + private final boolean createdUserRow; private boolean active = true; - private Scope(Connection connection) { this.connection = connection; } + private Scope(Connection connection, boolean createdUserRow) { + this.connection = connection; + this.createdUserRow = createdUserRow; + } private void requireActive() { if (!active) throw new IllegalStateException("SQL user transaction scope has ended"); } + @Override public boolean createdUserRow() { requireActive(); return createdUserRow; } @Override public Connection connection() { requireActive(); return connection; } @Override public List readRow() throws SQLException { requireActive(); return JdbcSqlUserStorage.this.readRow(connection); } @Override public void writeValues(Map values) throws SQLException { diff --git a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java index 589a68c99..3bcdcc81b 100644 --- a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java +++ b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java @@ -118,6 +118,14 @@ private static HashMap values(int points, int votes) { } } + @Test void transactionScopeReportsWhetherItCreatedTheUserRow() { + UUID uuid = UUID.randomUUID(); + try (SqliteUserBackend backend = backend()) { + assertTrue(backend.user(uuid).transaction(TYPE, SqlUserStorage.TransactionScope::createdUserRow)); + assertFalse(backend.user(uuid).transaction(TYPE, SqlUserStorage.TransactionScope::createdUserRow)); + } + } + @Test void rollsBackBothFailureOrders() throws Exception { UUID uuid = UUID.randomUUID(); try (SqliteUserBackend backend = backend()) { @@ -422,12 +430,16 @@ private static boolean applyOnce(SqlUserStorage user, String key) { assertEquals("accepted", backend.user(uuid).transaction(UserStorage.MYSQL, scope -> "accepted")); org.mockito.InOrder order = inOrder(connection, mysql); order.verify(connection).commit(); - order.verify(mysql).recordCommittedUser(uuid); + order.verify(mysql).recordCommittedUser(uuid, false); + clearInvocations(connection, mysql); + assertEquals("named", backend.user(uuid).transaction(UserStorage.MYSQL, + Map.of("PlayerName", new DataValueString("New")), scope -> "named")); + verify(mysql).recordCommittedUser(uuid, false); clearInvocations(connection, mysql); assertThrows(IllegalStateException.class, () -> backend.user(uuid).transaction(UserStorage.MYSQL, scope -> { throw new SQLException("receipt failed"); })); verify(connection).rollback(); - verify(mysql, never()).recordCommittedUser(any()); + verify(mysql, never()).recordCommittedUser(any(), anyBoolean()); } @Test void committedMysqlIdentityInvalidatesStaleNameCache() throws Exception { @@ -443,9 +455,13 @@ private static boolean applyOnce(SqlUserStorage user, String key) { doReturn(new java.util.ArrayList<>(List.of("New"))).when(mysql).getNamesQuery(); names.add("Old"); UUID uuid = UUID.randomUUID(); - mysql.recordCommittedUser(uuid); + Set priorSnapshot = mysql.getNames(); + mysql.recordCommittedUser(uuid, false); + assertTrue(names.contains("Old")); + mysql.recordCommittedUser(uuid, true); assertTrue(uuids.contains(uuid.toString())); assertFalse(names.contains("Old")); + assertEquals(Set.of("Old"), priorSnapshot); assertEquals(Set.of("New"), mysql.getNames()); } From 2b94f74d8beb8411d8a860226d78fe9beac90774 Mon Sep 17 00:00:00 2001 From: BenCodez <17074231+BenCodez@users.noreply.github.com> Date: Sun, 20 Sep 2026 10:55:17 -0600 Subject: [PATCH 7/7] Coordinate native name refresh and include dynamic transaction seeds --- .../api/user/userstorage/mysql/MySQL.java | 25 +++++---- .../user/storage/BukkitSqlUserBackend.java | 2 +- .../storage/SqliteUserTransactionTest.java | 54 ++++++++++++++++++- 3 files changed, 68 insertions(+), 13 deletions(-) diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java index 57e697617..19f5288b4 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/api/user/userstorage/mysql/MySQL.java @@ -125,16 +125,20 @@ public void loadData() { uuids.clear(); uuids.addAll(getUuidsQuery()); - names.clear(); - names.addAll(getNamesQuery()); + synchronized (names) { + names.clear(); + names.addAll(getNamesQuery()); + } } public void clearCacheBasic() { clearCaches(); uuids.clear(); uuids.addAll(getUuidsQuery()); - names.clear(); - names.addAll(getNamesQuery()); + synchronized (names) { + names.clear(); + names.addAll(getNamesQuery()); + } } // ------------------------- @@ -581,7 +585,7 @@ public void deletePlayer(String uuid) { } uuids.remove(uuid); - names.remove(UuidLookup.getInstance().getCachedName(uuid)); + synchronized (names) { names.remove(UuidLookup.getInstance().getCachedName(uuid)); } clearCacheBasic(); } @@ -821,12 +825,11 @@ public void insertQuery(String index, List cols) { playerName = col.getValue().toString(); } } - if (playerName == null || playerName.isEmpty()) { - names.add(UuidLookup.getInstance().getPlayerName( - plugin.getUserManager().getUser(java.util.UUID.fromString(index), false), index, false)); - } else { - names.add(playerName); - } + String committedName = playerName == null || playerName.isEmpty() + ? UuidLookup.getInstance().getPlayerName( + plugin.getUserManager().getUser(java.util.UUID.fromString(index), false), index, false) + : playerName; + synchronized (names) { names.add(committedName); } uuids.add(index); plugin.devDebug("Inserting " + index + " into database"); diff --git a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java index 943b93f3f..bef301a5d 100644 --- a/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java +++ b/AdvancedCore/src/main/java/com/bencodez/advancedcore/bukkit/user/storage/BukkitSqlUserBackend.java @@ -178,7 +178,7 @@ private T transaction(UserStorage storage, UUID uuid, java.util.Map { + return withSqlUser(storage, uuid, initialValues, user -> { if (storage != UserStorage.MYSQL) return user.transaction(storage, initialValues, work); boolean[] nameTouched = {false}; T result = user.transaction(storage, initialValues, scope -> { diff --git a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java index 3bcdcc81b..bcf5963f3 100644 --- a/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java +++ b/AdvancedCore/src/test/java/com/bencodez/advancedcore/tests/storage/SqliteUserTransactionTest.java @@ -376,9 +376,10 @@ private static boolean applyOnce(SqlUserStorage user, String key) { when(table.getSqLite()).thenReturn(sqlite); when(table.getName()).thenReturn("Users"); when(sqlite.getSQLConnection()).thenReturn(legacyConnection); + java.util.concurrent.atomic.AtomicBoolean dynamicCreated = new java.util.concurrent.atomic.AtomicBoolean(); doAnswer(call -> { Column column = call.getArgument(0); - if ("RepeatSpecial".equals(column.getName())) { + if ("RepeatSpecial".equals(column.getName()) && dynamicCreated.compareAndSet(false, true)) { try (PreparedStatement alter = legacyConnection.prepareStatement( "ALTER TABLE Users ADD COLUMN RepeatSpecial INTEGER")) { alter.executeUpdate(); } } @@ -404,6 +405,12 @@ private static boolean applyOnce(SqlUserStorage user, String key) { assertEquals(3, row.getInt(3)); } } + UUID seeded = UUID.randomUUID(); + int seededRepeat = runtime.transaction(seeded, TYPE, Map.of( + "PlayerName", new DataValueString("Other"), + "RepeatSpecial", new DataValueInt(9)), + scope -> value(scope.readRow(), "RepeatSpecial")); + assertEquals(9, seededRepeat); } } @@ -465,6 +472,51 @@ private static boolean applyOnce(SqlUserStorage user, String key) { assertEquals(Set.of("New"), mysql.getNames()); } + @Test void nativeNameRefreshCannotReinsertOldNameAfterCommit() throws Exception { + MySQL mysql = mock(MySQL.class, CALLS_REAL_METHODS); + Set names = java.util.concurrent.ConcurrentHashMap.newKeySet(); + Set uuids = java.util.concurrent.ConcurrentHashMap.newKeySet(); + java.lang.reflect.Field namesField = MySQL.class.getDeclaredField("names"); + java.lang.reflect.Field uuidsField = MySQL.class.getDeclaredField("uuids"); + namesField.setAccessible(true); + uuidsField.setAccessible(true); + namesField.set(mysql, names); + uuidsField.set(mysql, uuids); + doNothing().when(mysql).clearCaches(); + doReturn(new java.util.ArrayList()).when(mysql).getUuidsQuery(); + CountDownLatch oldQueryStarted = new CountDownLatch(1); + CountDownLatch finishOldQuery = new CountDownLatch(1); + java.util.concurrent.atomic.AtomicInteger queries = new java.util.concurrent.atomic.AtomicInteger(); + doAnswer(call -> { + if (queries.getAndIncrement() == 0) { + oldQueryStarted.countDown(); + assertTrue(finishOldQuery.await(5, TimeUnit.SECONDS)); + return new java.util.ArrayList<>(List.of("Old")); + } + return new java.util.ArrayList<>(List.of("New")); + }).when(mysql).getNamesQuery(); + ExecutorService workers = Executors.newFixedThreadPool(2); + try { + Future refresh = workers.submit(mysql::clearCacheBasic); + assertTrue(oldQueryStarted.await(5, TimeUnit.SECONDS)); + CountDownLatch commitStarted = new CountDownLatch(1); + Future committed = workers.submit(() -> { + commitStarted.countDown(); + mysql.recordCommittedUser(UUID.randomUUID(), true); + }); + assertTrue(commitStarted.await(5, TimeUnit.SECONDS)); + assertThrows(java.util.concurrent.TimeoutException.class, () -> committed.get(250, TimeUnit.MILLISECONDS)); + finishOldQuery.countDown(); + refresh.get(5, TimeUnit.SECONDS); + committed.get(5, TimeUnit.SECONDS); + assertTrue(names.isEmpty()); + assertEquals(Set.of("New"), mysql.getNames()); + } finally { + finishOldQuery.countDown(); + workers.shutdownNow(); + } + } + @Test void transactionFencesBypassPublishersAndReopensCacheAfterRollback() throws Exception { UUID uuid = UUID.randomUUID(); try (SqliteUserBackend backend = backend()) {