Skip to content

Commit ccf8064

Browse files
committed
Fix SQL schema and insert follow-ups
1 parent 765d4e1 commit ccf8064

7 files changed

Lines changed: 245 additions & 343 deletions

File tree

‎AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/JdbcSqlUserStorage.java‎

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -123,8 +123,8 @@ private enum BooleanStorage { TEXT, NATIVE, NUMERIC }
123123
boolean transactionEnded = false;
124124
Throwable transactionFailure = null;
125125
try {
126-
ensureRow(connection, updates);
127-
updateValues(connection, updates);
126+
boolean updateExistingRow = ensureRow(connection, updates);
127+
if (updateExistingRow) updateValues(connection, updates);
128128
connection.commit(); committed = true; transactionEnded = true;
129129
} catch (SQLException | RuntimeException | Error e) {
130130
transactionFailure = e;
@@ -165,8 +165,9 @@ private void committedCleanupFailure(String operation, Exception error) {
165165
}
166166
private static void suppress(Throwable primary, Throwable secondary) { if (primary != secondary) primary.addSuppressed(secondary); }
167167

168-
private void ensureRow(Connection connection, Map<String, DataValue> updates) throws SQLException {
169-
if (dialect != Dialect.SQLITE && rowExists(connection)) return;
168+
/** @return true when the row already existed and still needs the batch UPDATE. */
169+
private boolean ensureRow(Connection connection, Map<String, DataValue> updates) throws SQLException {
170+
if (dialect != Dialect.SQLITE && rowExists(connection)) return true;
170171
StringBuilder names = new StringBuilder(quote(SqlUserSchema.UUID_COLUMN));
171172
StringBuilder parameters = new StringBuilder("?");
172173
for (String key : updates.keySet()) { names.append(", ").append(quote(key)); parameters.append(", ?"); }
@@ -180,10 +181,12 @@ private void ensureRow(Connection connection, Map<String, DataValue> updates) th
180181
for (Map.Entry<String, DataValue> entry : updates.entrySet()) bind(statement, index++, entry.getValue(), schema.column(entry.getKey()));
181182
inserted = statement.executeUpdate();
182183
} catch (SQLException insertFailure) {
183-
if (dialect == Dialect.MYSQL && isDuplicateKey(insertFailure) && rowExists(connection)) return;
184+
if (dialect == Dialect.MYSQL && isDuplicateKey(insertFailure) && rowExists(connection)) return true;
184185
throw insertFailure;
185186
}
186-
if (inserted == 0 && !rowExists(connection)) throw new SQLException("SQL user row was not created");
187+
if (inserted > 0) return false;
188+
if (!rowExists(connection)) throw new SQLException("SQL user row was not created");
189+
return true;
187190
}
188191

189192
private boolean isDuplicateKey(SQLException failure) { return failure.getErrorCode() == 1062 || "23000".equals(failure.getSQLState()); }

‎AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/MysqlUserBackend.java‎

Lines changed: 35 additions & 74 deletions
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,7 @@ public MysqlUserBackend(String baseTableName, MysqlConfig config, SqlUserSchema
5050

5151
@Override
5252
public SqlUserStorage user(UUID uuid) {
53-
requireOpen();
53+
requireAdmissionOpen();
5454
JdbcSqlUserStorage.Dialect dialect = JdbcSqlUserStorage.Dialect.fromDbType(table.getMysql().getConnectionManager().getDbType());
5555
SqlUserStorage delegate = new JdbcSqlUserStorage(UserStorage.MYSQL, uuid, table.getTableName(), schema,
5656
() -> table.getMysql().getConnectionManager().getConnection(), dialect, logger);
@@ -63,46 +63,37 @@ public SqlUserStorage user(UUID uuid) {
6363
};
6464
}
6565

66-
@Override
67-
public List<UUID> enumerateUsers() {
66+
@Override public List<UUID> enumerateUsers() {
6867
ArrayList<UUID> users = new ArrayList<>();
6968
forEachUser(uuid -> {
70-
if (users.size() >= MAX_MATERIALIZED_USERS) throw new IllegalStateException("User enumeration exceeds "
71-
+ MAX_MATERIALIZED_USERS + " entries; use forEachUser for streaming access");
69+
if (users.size() >= MAX_MATERIALIZED_USERS) throw new IllegalStateException("User enumeration exceeds " + MAX_MATERIALIZED_USERS + " entries; use forEachUser for streaming access");
7270
users.add(uuid);
7371
});
7472
return users;
7573
}
7674

77-
@Override
78-
public void forEachUser(Consumer<UUID> consumer) {
75+
@Override public void forEachUser(Consumer<UUID> consumer) {
7976
Objects.requireNonNull(consumer, "consumer");
8077
withOperation(() -> {
8178
String cursor = null;
8279
while (true) {
8380
List<UserPageEntry> page = readUserPage(cursor);
8481
if (page.isEmpty()) return null;
8582
cursor = page.get(page.size() - 1).cursor();
86-
for (UserPageEntry entry : page) {
87-
if (entry.uuid() != null) consumer.accept(entry.uuid());
88-
}
83+
for (UserPageEntry entry : page) if (entry.uuid() != null) consumer.accept(entry.uuid());
8984
if (page.size() < USER_PAGE_SIZE) return null;
9085
}
9186
});
9287
}
9388

9489
private List<UserPageEntry> readUserPage(String cursor) {
9590
String uuidColumn = table.quote(SqlUserSchema.UUID_COLUMN);
96-
String sql = "SELECT " + uuidColumn + " FROM " + table.quote(table.getTableName())
97-
+ (cursor == null ? "" : " WHERE " + uuidColumn + " > ?")
98-
+ " ORDER BY " + uuidColumn + " ASC LIMIT ?";
91+
String sql = "SELECT " + uuidColumn + " FROM " + table.quote(table.getTableName()) + (cursor == null ? "" : " WHERE " + uuidColumn + " > ?") + " ORDER BY " + uuidColumn + " ASC LIMIT ?";
9992
JdbcSqlUserStorage.Dialect dialect = JdbcSqlUserStorage.Dialect.fromDbType(table.getDbType());
100-
try (Connection connection = table.getMysql().getConnectionManager().getConnection();
101-
PreparedStatement statement = connection.prepareStatement(sql)) {
93+
try (Connection connection = table.getMysql().getConnectionManager().getConnection(); PreparedStatement statement = connection.prepareStatement(sql)) {
10294
int index = 1;
10395
if (cursor != null) {
104-
if (dialect == JdbcSqlUserStorage.Dialect.POSTGRESQL) dialect.bindUuid(statement, index++, UUID.fromString(cursor));
105-
else statement.setString(index++, cursor);
96+
if (dialect == JdbcSqlUserStorage.Dialect.POSTGRESQL) dialect.bindUuid(statement, index++, UUID.fromString(cursor)); else statement.setString(index++, cursor);
10697
}
10798
statement.setInt(index, USER_PAGE_SIZE);
10899
ArrayList<UserPageEntry> page = new ArrayList<>(USER_PAGE_SIZE);
@@ -117,48 +108,39 @@ private List<UserPageEntry> readUserPage(String cursor) {
117108
}
118109
}
119110
return page;
120-
} catch (IllegalArgumentException invalidCursor) {
121-
throw new IllegalStateException("Failed to advance SQL user enumeration cursor", invalidCursor);
122-
} catch (SQLException failure) {
123-
throw new IllegalStateException("Failed to enumerate MySQL users", failure);
124-
}
111+
} catch (IllegalArgumentException invalidCursor) { throw new IllegalStateException("Failed to advance SQL user enumeration cursor", invalidCursor); }
112+
catch (SQLException failure) { throw new IllegalStateException("Failed to enumerate MySQL users", failure); }
125113
}
126114

127115
private record UserPageEntry(String cursor, UUID uuid) {}
128116

129117
@Override public boolean isOpen() { return open.get(); }
130118

131-
@Override
132-
public void close() {
119+
@Override public void close() {
133120
if (operations.getReadHoldCount() != 0) throw new IllegalStateException("Cannot close MySQL from inside an active storage operation");
134121
open.set(false);
135122
operations.writeLock().lock();
136123
try {
137-
if (!tableClosed) {
138-
table.close();
139-
tableClosed = true;
140-
}
141-
} finally {
142-
operations.writeLock().unlock();
143-
}
124+
if (!tableClosed) { table.close(); tableClosed = true; }
125+
} finally { operations.writeLock().unlock(); }
144126
}
145127

146128
private <T> T withOperation(Supplier<T> operation) {
147-
requireOpen();
129+
requireAdmissionOpen();
148130
operations.readLock().lock();
149-
try { requireOpen(); return operation.get(); }
150-
finally { operations.readLock().unlock(); }
131+
try {
132+
requireAdmissionOpen();
133+
return operation.get();
134+
} finally { operations.readLock().unlock(); }
151135
}
152136

153137
private void ensureRegisteredColumns() {
154138
table.ensureUuidType();
155-
for (SqlUserSchema.ColumnDefinition column : schema.columns()) {
156-
if (!SqlUserSchema.UUID_COLUMN.equalsIgnoreCase(column.name())) table.ensureColumn(column);
157-
}
139+
for (SqlUserSchema.ColumnDefinition column : schema.columns()) if (!SqlUserSchema.UUID_COLUMN.equalsIgnoreCase(column.name())) table.ensureColumn(column);
158140
}
159141

160-
private void requireOpen() {
161-
if (!open.get() && operations.getReadHoldCount() == 0) throw new IllegalStateException("MySQL user backend is closed");
142+
private void requireAdmissionOpen() {
143+
if (!open.get()) throw new IllegalStateException("MySQL user backend is closed");
162144
}
163145

164146
private static final class HeadlessUserTable extends AbstractSqlTable {
@@ -178,14 +160,11 @@ private static final class HeadlessUserTable extends AbstractSqlTable {
178160
}
179161

180162
@Override public String getPrimaryKeyColumn() { return SqlUserSchema.UUID_COLUMN; }
181-
182-
@Override
183-
public String buildCreateTableSql(DbType dbType) {
163+
@Override public String buildCreateTableSql(DbType dbType) {
184164
StringBuilder sql = new StringBuilder("CREATE TABLE IF NOT EXISTS ").append(quote(tableName)).append(" (");
185165
boolean first = true;
186166
for (SqlUserSchema.ColumnDefinition column : schema.columns()) {
187-
if (!first) sql.append(", ");
188-
first = false;
167+
if (!first) sql.append(", "); first = false;
189168
String type = SqlUserSchema.UUID_COLUMN.equalsIgnoreCase(column.name()) ? bestUuidType() : normaliseTypeForDb(column.sqlType());
190169
sql.append(quote(column.name())).append(' ').append(type);
191170
}
@@ -204,10 +183,8 @@ void ensureUuidType() {
204183
String uuidType = bestUuidType();
205184
if (!columnNeedsAlter(SqlUserSchema.UUID_COLUMN, uuidType)) return;
206185
String uuidColumn = quote(SqlUserSchema.UUID_COLUMN);
207-
String sql = "ALTER TABLE " + quote(tableName) + " ALTER COLUMN " + uuidColumn
208-
+ " TYPE " + uuidType + " USING NULLIF(" + uuidColumn + ", '')::uuid;";
209-
try (Connection connection = getMysql().getConnectionManager().getConnection();
210-
PreparedStatement statement = connection.prepareStatement(sql)) { statement.executeUpdate(); }
186+
String sql = "ALTER TABLE " + quote(tableName) + " ALTER COLUMN " + uuidColumn + " TYPE " + uuidType + " USING NULLIF(" + uuidColumn + ", '')::uuid;";
187+
try (Connection connection = getMysql().getConnectionManager().getConnection(); PreparedStatement statement = connection.prepareStatement(sql)) { statement.executeUpdate(); }
211188
catch (SQLException ddlFailure) {
212189
try { if (columnNeedsAlter(SqlUserSchema.UUID_COLUMN, uuidType)) throw ddlFailure; }
213190
catch (SQLException inspectionFailure) {
@@ -224,13 +201,10 @@ void ensureColumn(SqlUserSchema.ColumnDefinition column) {
224201
String storedName = findRegisteredColumn(column.name());
225202
if (storedName != null) {
226203
if (getDbType() == DbType.POSTGRESQL && !storedName.equals(column.name())) renamePostgresColumn(storedName, column.name());
227-
rememberColumn(column);
228-
return;
204+
rememberColumn(column); return;
229205
}
230-
String sql = "ALTER TABLE " + quote(tableName) + " ADD COLUMN " + quote(column.name())
231-
+ " " + normaliseTypeForDb(column.sqlType()) + ";";
232-
try (Connection connection = getMysql().getConnectionManager().getConnection();
233-
PreparedStatement statement = connection.prepareStatement(sql)) { statement.executeUpdate(); }
206+
String sql = "ALTER TABLE " + quote(tableName) + " ADD COLUMN " + quote(column.name()) + " " + normaliseTypeForDb(column.sqlType()) + ";";
207+
try (Connection connection = getMysql().getConnectionManager().getConnection(); PreparedStatement statement = connection.prepareStatement(sql)) { statement.executeUpdate(); }
234208
catch (SQLException ddlFailure) {
235209
if (!isDuplicateColumn(ddlFailure)) throw ddlFailure;
236210
try {
@@ -249,27 +223,18 @@ void ensureColumn(SqlUserSchema.ColumnDefinition column) {
249223

250224
private void renamePostgresColumn(String storedName, String requestedName) throws SQLException {
251225
String sql = "ALTER TABLE " + quote(tableName) + " RENAME COLUMN " + quote(storedName) + " TO " + quote(requestedName) + ";";
252-
try (Connection connection = getMysql().getConnectionManager().getConnection(); PreparedStatement statement = connection.prepareStatement(sql)) {
253-
statement.executeUpdate();
254-
} catch (SQLException renameFailure) {
255-
String current = findRegisteredColumn(requestedName);
256-
if (!requestedName.equals(current)) throw renameFailure;
257-
}
226+
try (Connection connection = getMysql().getConnectionManager().getConnection(); PreparedStatement statement = connection.prepareStatement(sql)) { statement.executeUpdate(); }
227+
catch (SQLException renameFailure) { String current = findRegisteredColumn(requestedName); if (!requestedName.equals(current)) throw renameFailure; }
258228
}
259229

260230
private void rememberColumn(SqlUserSchema.ColumnDefinition column) {
261-
columns.removeIf(existing -> existing.equalsIgnoreCase(column.name()));
262-
columns.add(column.name());
263-
intColumns.removeIf(existing -> existing.equalsIgnoreCase(column.name()));
264-
if (column.dataType() == DataType.INTEGER) intColumns.add(column.name());
231+
columns.removeIf(existing -> existing.equalsIgnoreCase(column.name())); columns.add(column.name());
232+
intColumns.removeIf(existing -> existing.equalsIgnoreCase(column.name())); if (column.dataType() == DataType.INTEGER) intColumns.add(column.name());
265233
}
266234

267235
private String findRegisteredColumn(String name) throws SQLException {
268-
try (Connection connection = getMysql().getConnectionManager().getConnection();
269-
PreparedStatement statement = connection.prepareStatement("SELECT * FROM " + quote(tableName) + " WHERE 1=0");
270-
ResultSet result = statement.executeQuery()) {
271-
ResultSetMetaData metadata = result.getMetaData();
272-
String foldedMatch = null;
236+
try (Connection connection = getMysql().getConnectionManager().getConnection(); PreparedStatement statement = connection.prepareStatement("SELECT * FROM " + quote(tableName) + " WHERE 1=0"); ResultSet result = statement.executeQuery()) {
237+
ResultSetMetaData metadata = result.getMetaData(); String foldedMatch = null;
273238
for (int i = 1; i <= metadata.getColumnCount(); i++) {
274239
String storedName = metadata.getColumnName(i);
275240
if (name.equals(storedName)) return storedName;
@@ -279,11 +244,7 @@ private String findRegisteredColumn(String name) throws SQLException {
279244
}
280245
}
281246

282-
private boolean isDuplicateColumn(SQLException failure) {
283-
return getDbType() == DbType.POSTGRESQL ? "42701".equals(failure.getSQLState())
284-
: failure.getErrorCode() == 1060 && "42S21".equals(failure.getSQLState());
285-
}
286-
247+
private boolean isDuplicateColumn(SQLException failure) { return getDbType() == DbType.POSTGRESQL ? "42701".equals(failure.getSQLState()) : failure.getErrorCode() == 1060 && "42S21".equals(failure.getSQLState()); }
287248
String quote(String identifier) { return qi(identifier); }
288249
}
289250
}

‎AdvancedCore/src/main/java/com/bencodez/advancedcore/core/user/storage/sql/SqlUserSchema.java‎

Lines changed: 15 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -22,12 +22,8 @@ public record ColumnDefinition(String name, String sqlType, DataType dataType) {
2222
Objects.requireNonNull(name, "name");
2323
Objects.requireNonNull(sqlType, "sqlType");
2424
Objects.requireNonNull(dataType, "dataType");
25-
if (name.isBlank()) {
26-
throw new IllegalArgumentException("Column name cannot be blank");
27-
}
28-
if (sqlType.isBlank()) {
29-
throw new IllegalArgumentException("SQL type cannot be blank");
30-
}
25+
if (name.isBlank()) throw new IllegalArgumentException("Column name cannot be blank");
26+
if (sqlType.isBlank()) throw new IllegalArgumentException("SQL type cannot be blank");
3127
}
3228
}
3329

@@ -37,45 +33,32 @@ private SqlUserSchema(Map<String, ColumnDefinition> columnsByLowerName) {
3733
this.columnsByLowerName = Collections.unmodifiableMap(new LinkedHashMap<>(columnsByLowerName));
3834
}
3935

40-
public static Builder builder() {
41-
return new Builder();
42-
}
36+
public static Builder builder() { return new Builder(); }
4337

4438
public static SqlUserSchema fromKeys(Collection<? extends UserDataKey> keys) {
4539
Builder builder = builder();
4640
for (UserDataKey key : Objects.requireNonNull(keys, "keys")) {
4741
DataType type = DataType.STRING;
48-
if (key instanceof UserDataKeyInt) {
49-
type = DataType.INTEGER;
50-
} else if (key instanceof UserDataKeyBoolean) {
51-
type = DataType.BOOLEAN;
52-
}
42+
if (key instanceof UserDataKeyInt) type = DataType.INTEGER;
43+
else if (key instanceof UserDataKeyBoolean) type = DataType.BOOLEAN;
5344
builder.column(key.getKey(), key.getColumnType(), type);
5445
}
5546
return builder.build();
5647
}
5748

58-
public List<ColumnDefinition> columns() {
59-
return new ArrayList<>(columnsByLowerName.values());
60-
}
49+
public List<ColumnDefinition> columns() { return new ArrayList<>(columnsByLowerName.values()); }
6150

6251
public ColumnDefinition column(String name) {
63-
if (name == null) {
64-
return null;
65-
}
52+
if (name == null) return null;
6653
return columnsByLowerName.get(name.toLowerCase(Locale.ROOT));
6754
}
6855

69-
public boolean contains(String name) {
70-
return column(name) != null;
71-
}
56+
public boolean contains(String name) { return column(name) != null; }
7257

7358
public static final class Builder {
7459
private final Map<String, ColumnDefinition> columns = new LinkedHashMap<>();
7560

7661
private Builder() {
77-
// Identity spelling/type is invariant. PostgreSQL quotes identifiers and
78-
// therefore cannot tolerate a custom "UUID" replacing canonical "uuid".
7962
columns.put(UUID_COLUMN, new ColumnDefinition(UUID_COLUMN, "VARCHAR(37)", DataType.STRING));
8063
}
8164

@@ -84,13 +67,15 @@ public Builder column(String name, String sqlType, DataType dataType) {
8467
if (UUID_COLUMN.equalsIgnoreCase(name)) {
8568
throw new IllegalArgumentException("Column name 'uuid' is reserved for user identity");
8669
}
87-
ColumnDefinition definition = new ColumnDefinition(name, sqlType, dataType);
88-
columns.put(name.toLowerCase(Locale.ROOT), definition);
70+
String canonical = name.toLowerCase(Locale.ROOT);
71+
ColumnDefinition existing = columns.get(canonical);
72+
if (existing != null) {
73+
throw new IllegalArgumentException("Duplicate SQL column name: " + name + " conflicts with " + existing.name());
74+
}
75+
columns.put(canonical, new ColumnDefinition(name, sqlType, dataType));
8976
return this;
9077
}
9178

92-
public SqlUserSchema build() {
93-
return new SqlUserSchema(columns);
94-
}
79+
public SqlUserSchema build() { return new SqlUserSchema(columns); }
9580
}
9681
}

0 commit comments

Comments
 (0)