diff --git a/.github/workflows/ci-tests-and-coverage.yml b/.github/workflows/ci-tests-and-coverage.yml index f91d7d3..6c6691b 100644 --- a/.github/workflows/ci-tests-and-coverage.yml +++ b/.github/workflows/ci-tests-and-coverage.yml @@ -30,20 +30,31 @@ jobs: cache: maven - name: Run Unit and Backend Integration Tests + id: tests run: mvn -U -B -ntp -Pintegration-tests verify + - name: Print failing test reports + if: failure() && steps.tests.outcome == 'failure' + shell: bash + run: | + find . -type f \( -path '*/target/surefire-reports/*.txt' -o -path '*/target/failsafe-reports/*.txt' \) -print0 \ + | sort -z \ + | xargs -0 -r -n1 sh -c 'echo "::group::$0"; cat "$0"; echo "::endgroup::"' + - name: Upload JaCoCo Report if: always() uses: actions/upload-artifact@v7 with: name: jacoco-report - if-no-files-found: error + if-no-files-found: warn path: "**/target/site/jacoco" - - name: Upload Integration Test Reports + - name: Upload Test Reports if: always() uses: actions/upload-artifact@v7 with: - name: integration-test-reports - if-no-files-found: error - path: "**/target/failsafe-reports" + name: test-reports + if-no-files-found: warn + path: | + **/target/surefire-reports + **/target/failsafe-reports diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/api/DataProviderAPI.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/api/DataProviderAPI.java index 7338f8f..9b09add 100644 --- a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/api/DataProviderAPI.java +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/api/DataProviderAPI.java @@ -1,33 +1,22 @@ package nl.hauntedmc.dataprovider.api; +import nl.hauntedmc.dataprovider.api.orm.ORMContext; import nl.hauntedmc.dataprovider.database.DataAccess; import nl.hauntedmc.dataprovider.database.DatabaseProvider; import nl.hauntedmc.dataprovider.database.DatabaseType; -import nl.hauntedmc.dataprovider.api.orm.ORMContext; +import nl.hauntedmc.dataprovider.exception.DataProviderFailureContext; +import nl.hauntedmc.dataprovider.exception.DataProviderRegistrationException; +import nl.hauntedmc.dataprovider.exception.ExecutionOutcome; +import nl.hauntedmc.dataprovider.exception.RetryAdvice; import nl.hauntedmc.dataprovider.logging.LoggerAdapter; import javax.sql.DataSource; import java.util.Objects; import java.util.Optional; -/** - * Public, platform-neutral facade for plugin-scoped database registrations. - * - *

The API artifact intentionally contains contracts only. Platform modules provide the - * runtime implementation and expose an instance through their native service mechanism.

- */ +/** Public, platform-neutral facade for plugin-scoped database registrations. */ public interface DataProviderAPI { - /** - * Creates an isolated ORM context owned by the calling plugin. - * - * @param pluginName plugin name used for ORM diagnostics - * @param dataSource relational data source obtained from a registered provider - * @param logger logger that receives ORM lifecycle diagnostics - * @param schemaMode Hibernate schema mode: validate, none, update, or create - * @param entityClasses annotated entity classes to register - * @return a new, initialized ORM context - */ ORMContext createOrmContext( String pluginName, DataSource dataSource, @@ -36,8 +25,31 @@ ORMContext createOrmContext( Class... entityClasses ); + /** Legacy nullable registration method retained for compatibility. */ DatabaseProvider registerDatabase(DatabaseType databaseType, String connectionIdentifier); + /** + * Registers a database or throws a structured public exception retaining the failure category. + * Implementations should override this method to preserve backend-specific failure details. + */ + default DatabaseProvider registerDatabaseOrThrow(DatabaseType databaseType, String connectionIdentifier) { + DatabaseProvider provider = registerDatabase(databaseType, connectionIdentifier); + if (provider != null) { + return provider; + } + throw new DataProviderRegistrationException( + "Database registration failed.", + DataProviderFailureContext.of( + databaseType, + connectionIdentifier, + "registerDatabase", + RetryAdvice.CONDITIONAL, + ExecutionOutcome.NOT_STARTED + ), + null + ); + } + DataProviderScope scope(OwnerScope ownerScope); void unregisterDatabase(DatabaseType databaseType, String connectionIdentifier); @@ -46,9 +58,28 @@ ORMContext createOrmContext( void unregisterAllDatabasesForPlugin(); + /** Legacy nullable lookup retained for compatibility. */ DatabaseProvider getRegisteredDatabase(DatabaseType databaseType, String connectionIdentifier); - /** Creates an isolated ownership scope for independently managed plugin components. */ + /** Returns a registered provider or throws a structured registration-state failure. */ + default DatabaseProvider requireRegisteredDatabase(DatabaseType databaseType, String connectionIdentifier) { + DatabaseProvider provider = getRegisteredDatabase(databaseType, connectionIdentifier); + if (provider != null) { + return provider; + } + throw new DataProviderRegistrationException( + "No active database registration exists for the requested connection.", + DataProviderFailureContext.of( + databaseType, + connectionIdentifier, + "requireRegisteredDatabase", + RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED + ), + null + ); + } + default DataProviderScope scope(String ownerScope) { return scope(OwnerScope.of(ownerScope)); } diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/api/DataProviderScope.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/api/DataProviderScope.java index 4a196da..3a4cae3 100644 --- a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/api/DataProviderScope.java +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/api/DataProviderScope.java @@ -3,16 +3,17 @@ import nl.hauntedmc.dataprovider.database.DataAccess; import nl.hauntedmc.dataprovider.database.DatabaseProvider; import nl.hauntedmc.dataprovider.database.DatabaseType; +import nl.hauntedmc.dataprovider.exception.DataProviderFailureContext; +import nl.hauntedmc.dataprovider.exception.DataProviderRegistrationException; +import nl.hauntedmc.dataprovider.exception.ExecutionOutcome; +import nl.hauntedmc.dataprovider.exception.RetryAdvice; import java.util.Objects; import java.util.Optional; -/** - * Isolated lifecycle boundary for a logical component within one plugin. - */ +/** Isolated lifecycle boundary for a logical component within one plugin. */ public interface DataProviderScope extends AutoCloseable { - /** Lifecycle states for a scope. A closed scope cannot be reopened. */ enum LifecycleState { OPEN, CLOSING, @@ -21,29 +22,56 @@ enum LifecycleState { OwnerScope ownerScope(); - /** - * Returns this scope's current lifecycle state. - * Implementations created by DataProvider transition from OPEN to CLOSING to CLOSED on close. - */ default LifecycleState lifecycleState() { return LifecycleState.OPEN; } DatabaseProvider registerDatabase(DatabaseType databaseType, String connectionIdentifier); + default DatabaseProvider registerDatabaseOrThrow(DatabaseType databaseType, String connectionIdentifier) { + DatabaseProvider provider = registerDatabase(databaseType, connectionIdentifier); + if (provider != null) { + return provider; + } + throw new DataProviderRegistrationException( + "Scoped database registration failed.", + DataProviderFailureContext.of( + databaseType, + connectionIdentifier, + "scope.registerDatabase", + RetryAdvice.CONDITIONAL, + ExecutionOutcome.NOT_STARTED + ).withDiagnostics(java.util.Map.of("ownerScope", ownerScope().value())), + null + ); + } + void unregisterDatabase(DatabaseType databaseType, String connectionIdentifier); void unregisterAllDatabases(); - /** - * Retrieves a provider registered by this scope. - * - * @throws UnsupportedOperationException if the scope implementation does not support scoped lookup - */ default DatabaseProvider getRegisteredDatabase(DatabaseType databaseType, String connectionIdentifier) { throw new UnsupportedOperationException("Scoped provider lookup is not supported by this implementation."); } + default DatabaseProvider requireRegisteredDatabase(DatabaseType databaseType, String connectionIdentifier) { + DatabaseProvider provider = getRegisteredDatabase(databaseType, connectionIdentifier); + if (provider != null) { + return provider; + } + throw new DataProviderRegistrationException( + "No active scoped database registration exists for the requested connection.", + DataProviderFailureContext.of( + databaseType, + connectionIdentifier, + "scope.requireRegisteredDatabase", + RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED + ).withDiagnostics(java.util.Map.of("ownerScope", ownerScope().value())), + null + ); + } + default Optional registerDatabaseOptional( DatabaseType databaseType, String connectionIdentifier diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/BackendAuthenticationException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/BackendAuthenticationException.java new file mode 100644 index 0000000..c38fbc8 --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/BackendAuthenticationException.java @@ -0,0 +1,8 @@ +package nl.hauntedmc.dataprovider.exception; + +/** Backend rejected configured authentication. */ +public final class BackendAuthenticationException extends DataProviderException { + public BackendAuthenticationException(String message, DataProviderFailureContext context, Throwable cause) { + super(DataProviderErrorCode.AUTHENTICATION_FAILED, message, context, cause); + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/BackendUnavailableException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/BackendUnavailableException.java new file mode 100644 index 0000000..d4dc5ed --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/BackendUnavailableException.java @@ -0,0 +1,13 @@ +package nl.hauntedmc.dataprovider.exception; + +/** Backend is disabled, unreachable, or otherwise unavailable. */ +public final class BackendUnavailableException extends DataProviderException { + public BackendUnavailableException(DataProviderErrorCode code, String message, + DataProviderFailureContext context, Throwable cause) { + super(code, message, context, cause); + if (code != DataProviderErrorCode.BACKEND_DISABLED + && code != DataProviderErrorCode.BACKEND_UNAVAILABLE) { + throw new IllegalArgumentException("Unsupported backend availability error code: " + code); + } + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataConflictException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataConflictException.java new file mode 100644 index 0000000..7e55589 --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataConflictException.java @@ -0,0 +1,8 @@ +package nl.hauntedmc.dataprovider.exception; + +/** A uniqueness, optimistic-locking, or compare-and-set conflict occurred. */ +public final class DataConflictException extends DataProviderException { + public DataConflictException(String message, DataProviderFailureContext context, Throwable cause) { + super(DataProviderErrorCode.CONFLICT, message, context, cause); + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderConfigurationException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderConfigurationException.java new file mode 100644 index 0000000..0c14845 --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderConfigurationException.java @@ -0,0 +1,13 @@ +package nl.hauntedmc.dataprovider.exception; + +/** Invalid or missing DataProvider configuration. */ +public final class DataProviderConfigurationException extends DataProviderException { + public DataProviderConfigurationException(DataProviderErrorCode code, String message, + DataProviderFailureContext context, Throwable cause) { + super(code, message, context, cause); + if (code != DataProviderErrorCode.CONFIGURATION_INVALID + && code != DataProviderErrorCode.CONFIGURATION_MISSING) { + throw new IllegalArgumentException("Unsupported configuration error code: " + code); + } + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderErrorCode.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderErrorCode.java new file mode 100644 index 0000000..1aba041 --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderErrorCode.java @@ -0,0 +1,18 @@ +package nl.hauntedmc.dataprovider.exception; + +/** Stable machine-readable failure codes exposed by the DataProvider API. */ +public enum DataProviderErrorCode { + CONFIGURATION_INVALID, + CONFIGURATION_MISSING, + REGISTRATION_FAILED, + BACKEND_DISABLED, + BACKEND_UNAVAILABLE, + AUTHENTICATION_FAILED, + OPERATION_FAILED, + OPERATION_TIMED_OUT, + QUEUE_SATURATED, + SERIALIZATION_FAILED, + CONFLICT, + TRANSACTION_FAILED, + PROVIDER_CLOSED +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderException.java new file mode 100644 index 0000000..888023b --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderException.java @@ -0,0 +1,160 @@ +package nl.hauntedmc.dataprovider.exception; + +import nl.hauntedmc.dataprovider.database.DatabaseType; + +import java.util.LinkedHashMap; +import java.util.Locale; +import java.util.Map; +import java.util.Objects; +import java.util.UUID; +import java.util.regex.Pattern; + +/** Base type for safe, structured failures exposed by DataProvider. */ +public abstract class DataProviderException extends RuntimeException { + + private static final Pattern DIAGNOSTIC_KEY_PATTERN = Pattern.compile("[A-Za-z][A-Za-z0-9_.-]{0,63}"); + private static final Pattern DIAGNOSTIC_ID_PATTERN = Pattern.compile("[A-Za-z0-9][A-Za-z0-9_.:-]{0,127}"); + + private final DataProviderErrorCode errorCode; + private final DatabaseType backendType; + private final String connectionIdentifier; + private final String operationName; + private final RetryAdvice retryAdvice; + private final ExecutionOutcome executionOutcome; + private final Map diagnostics; + private final String diagnosticId; + + protected DataProviderException( + DataProviderErrorCode errorCode, + String safeMessage, + DataProviderFailureContext context, + Throwable safeCause + ) { + this( + errorCode, + safeMessage, + Objects.requireNonNull(context, "Failure context cannot be null.").backendType(), + context.connectionIdentifier(), + context.operationName(), + context.retryAdvice(), + context.executionOutcome(), + context.diagnostics(), + context.diagnosticId(), + safeCause + ); + } + + protected DataProviderException( + DataProviderErrorCode errorCode, + String safeMessage, + DatabaseType backendType, + String connectionIdentifier, + String operationName, + RetryAdvice retryAdvice, + ExecutionOutcome executionOutcome, + Map diagnostics, + String diagnosticId, + Throwable safeCause + ) { + super(requireSafeText(safeMessage, "safeMessage"), safeCause); + this.errorCode = Objects.requireNonNull(errorCode, "Error code cannot be null."); + this.backendType = backendType; + this.connectionIdentifier = normalizeNullable(connectionIdentifier); + this.operationName = normalizeNullable(operationName); + this.retryAdvice = Objects.requireNonNull(retryAdvice, "Retry advice cannot be null."); + this.executionOutcome = Objects.requireNonNull(executionOutcome, "Execution outcome cannot be null."); + this.diagnostics = sanitizeDiagnostics(diagnostics); + this.diagnosticId = normalizeDiagnosticId(diagnosticId); + } + + public final DataProviderErrorCode errorCode() { + return errorCode; + } + + public final DatabaseType backendType() { + return backendType; + } + + public final String connectionIdentifier() { + return connectionIdentifier; + } + + public final String operationName() { + return operationName; + } + + public final RetryAdvice retryAdvice() { + return retryAdvice; + } + + public final boolean retryable() { + return retryAdvice != RetryAdvice.NEVER; + } + + public final ExecutionOutcome executionOutcome() { + return executionOutcome; + } + + public final Map diagnostics() { + return diagnostics; + } + + public final String diagnosticId() { + return diagnosticId; + } + + private static Map sanitizeDiagnostics(Map source) { + if (source == null || source.isEmpty()) { + return Map.of(); + } + LinkedHashMap safe = new LinkedHashMap<>(); + source.forEach((key, value) -> { + String normalizedKey = requireSafeText(key, "diagnostic key"); + if (!DIAGNOSTIC_KEY_PATTERN.matcher(normalizedKey).matches()) { + throw new IllegalArgumentException("Unsupported diagnostic key: " + normalizedKey); + } + if (isSensitiveKey(normalizedKey)) { + throw new IllegalArgumentException("Sensitive diagnostic keys are not allowed: " + normalizedKey); + } + String normalizedValue = requireSafeText(value, "diagnostic value"); + if (normalizedValue.length() > 256) { + throw new IllegalArgumentException("Diagnostic values cannot exceed 256 characters."); + } + safe.put(normalizedKey, normalizedValue); + }); + return Map.copyOf(safe); + } + + private static boolean isSensitiveKey(String key) { + String lower = key.toLowerCase(Locale.ROOT); + return lower.contains("password") || lower.contains("secret") || lower.contains("token") + || lower.contains("credential") || lower.contains("authorization") || lower.contains("payload") + || lower.contains("query") || lower.contains("parameter") || lower.contains("url"); + } + + private static String normalizeDiagnosticId(String diagnosticId) { + if (diagnosticId == null || diagnosticId.isBlank()) { + return UUID.randomUUID().toString(); + } + String normalized = diagnosticId.trim(); + if (!DIAGNOSTIC_ID_PATTERN.matcher(normalized).matches()) { + throw new IllegalArgumentException("Unsupported diagnostic identifier."); + } + return normalized; + } + + private static String requireSafeText(String value, String field) { + if (value == null || value.isBlank()) { + throw new IllegalArgumentException(field + " cannot be null or blank."); + } + String normalized = value.trim(); + if (normalized.indexOf('\0') >= 0) { + throw new IllegalArgumentException(field + " cannot contain null characters."); + } + return normalized; + } + + private static String normalizeNullable(String value) { + return value == null || value.isBlank() ? null : value.trim(); + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderFailureContext.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderFailureContext.java new file mode 100644 index 0000000..f118590 --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderFailureContext.java @@ -0,0 +1,64 @@ +package nl.hauntedmc.dataprovider.exception; + +import nl.hauntedmc.dataprovider.database.DatabaseType; + +import java.util.Map; + +/** Safe context attached to a structured DataProvider failure. */ +public record DataProviderFailureContext( + DatabaseType backendType, + String connectionIdentifier, + String operationName, + RetryAdvice retryAdvice, + ExecutionOutcome executionOutcome, + Map diagnostics, + String diagnosticId +) { + public DataProviderFailureContext { + retryAdvice = retryAdvice == null ? RetryAdvice.NEVER : retryAdvice; + executionOutcome = executionOutcome == null ? ExecutionOutcome.UNKNOWN : executionOutcome; + diagnostics = diagnostics == null ? Map.of() : Map.copyOf(diagnostics); + } + + public static DataProviderFailureContext of( + DatabaseType backendType, + String connectionIdentifier, + String operationName, + RetryAdvice retryAdvice, + ExecutionOutcome executionOutcome + ) { + return new DataProviderFailureContext( + backendType, + connectionIdentifier, + operationName, + retryAdvice, + executionOutcome, + Map.of(), + null + ); + } + + public DataProviderFailureContext withDiagnostics(Map diagnostics) { + return new DataProviderFailureContext( + backendType, + connectionIdentifier, + operationName, + retryAdvice, + executionOutcome, + diagnostics, + diagnosticId + ); + } + + public DataProviderFailureContext withDiagnosticId(String diagnosticId) { + return new DataProviderFailureContext( + backendType, + connectionIdentifier, + operationName, + retryAdvice, + executionOutcome, + diagnostics, + diagnosticId + ); + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderOperationException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderOperationException.java new file mode 100644 index 0000000..19cb9e6 --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderOperationException.java @@ -0,0 +1,13 @@ +package nl.hauntedmc.dataprovider.exception; + +/** Backend operation failed without matching a more specific public failure category. */ +public final class DataProviderOperationException extends DataProviderException { + + public DataProviderOperationException( + String message, + DataProviderFailureContext context, + Throwable cause + ) { + super(DataProviderErrorCode.OPERATION_FAILED, message, context, cause); + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderRegistrationException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderRegistrationException.java new file mode 100644 index 0000000..98b79a5 --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderRegistrationException.java @@ -0,0 +1,8 @@ +package nl.hauntedmc.dataprovider.exception; + +/** Failure in registration ownership, publication, or lookup state. */ +public final class DataProviderRegistrationException extends DataProviderException { + public DataProviderRegistrationException(String message, DataProviderFailureContext context, Throwable cause) { + super(DataProviderErrorCode.REGISTRATION_FAILED, message, context, cause); + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderTimeoutException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderTimeoutException.java new file mode 100644 index 0000000..4b37745 --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataProviderTimeoutException.java @@ -0,0 +1,8 @@ +package nl.hauntedmc.dataprovider.exception; + +/** Operation exceeded a configured or backend timeout. */ +public final class DataProviderTimeoutException extends DataProviderException { + public DataProviderTimeoutException(String message, DataProviderFailureContext context, Throwable cause) { + super(DataProviderErrorCode.OPERATION_TIMED_OUT, message, context, cause); + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataSerializationException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataSerializationException.java new file mode 100644 index 0000000..41f1ff0 --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataSerializationException.java @@ -0,0 +1,8 @@ +package nl.hauntedmc.dataprovider.exception; + +/** Data could not be serialized or deserialized safely. */ +public final class DataSerializationException extends DataProviderException { + public DataSerializationException(String message, DataProviderFailureContext context, Throwable cause) { + super(DataProviderErrorCode.SERIALIZATION_FAILED, message, context, cause); + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataTransactionException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataTransactionException.java new file mode 100644 index 0000000..bcdfff4 --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/DataTransactionException.java @@ -0,0 +1,19 @@ +package nl.hauntedmc.dataprovider.exception; + +import java.util.Objects; + +/** Transaction failure retaining the phase and primary cause. */ +public final class DataTransactionException extends DataProviderException { + + private final TransactionPhase phase; + + public DataTransactionException(String message, TransactionPhase phase, + DataProviderFailureContext context, Throwable cause) { + super(DataProviderErrorCode.TRANSACTION_FAILED, message, context, cause); + this.phase = Objects.requireNonNull(phase, "Transaction phase cannot be null."); + } + + public TransactionPhase phase() { + return phase; + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/ExecutionOutcome.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/ExecutionOutcome.java new file mode 100644 index 0000000..61a021b --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/ExecutionOutcome.java @@ -0,0 +1,9 @@ +package nl.hauntedmc.dataprovider.exception; + +/** What is known about whether a failed operation reached the backend. */ +public enum ExecutionOutcome { + NOT_STARTED, + NOT_APPLIED, + MAY_HAVE_APPLIED, + UNKNOWN +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/ProviderClosedException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/ProviderClosedException.java new file mode 100644 index 0000000..11261ea --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/ProviderClosedException.java @@ -0,0 +1,8 @@ +package nl.hauntedmc.dataprovider.exception; + +/** A runtime-scoped provider or registration scope is already closed. */ +public final class ProviderClosedException extends DataProviderException { + public ProviderClosedException(String message, DataProviderFailureContext context, Throwable cause) { + super(DataProviderErrorCode.PROVIDER_CLOSED, message, context, cause); + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/QueueSaturatedException.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/QueueSaturatedException.java new file mode 100644 index 0000000..6877cad --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/QueueSaturatedException.java @@ -0,0 +1,8 @@ +package nl.hauntedmc.dataprovider.exception; + +/** Shared execution capacity rejected an operation before it started. */ +public final class QueueSaturatedException extends DataProviderException { + public QueueSaturatedException(String message, DataProviderFailureContext context, Throwable cause) { + super(DataProviderErrorCode.QUEUE_SATURATED, message, context, cause); + } +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/RetryAdvice.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/RetryAdvice.java new file mode 100644 index 0000000..cfacefd --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/RetryAdvice.java @@ -0,0 +1,8 @@ +package nl.hauntedmc.dataprovider.exception; + +/** Guidance for callers considering a retry. */ +public enum RetryAdvice { + NEVER, + SAFE, + CONDITIONAL +} diff --git a/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/TransactionPhase.java b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/TransactionPhase.java new file mode 100644 index 0000000..a79f061 --- /dev/null +++ b/dataprovider-api/src/main/java/nl/hauntedmc/dataprovider/exception/TransactionPhase.java @@ -0,0 +1,10 @@ +package nl.hauntedmc.dataprovider.exception; + +/** Transaction stage at which a failure occurred. */ +public enum TransactionPhase { + BEGIN, + CALLBACK, + COMMIT, + ROLLBACK, + CLEANUP +} diff --git a/dataprovider-api/src/test/java/nl/hauntedmc/dataprovider/exception/DataProviderExceptionTest.java b/dataprovider-api/src/test/java/nl/hauntedmc/dataprovider/exception/DataProviderExceptionTest.java new file mode 100644 index 0000000..becd236 --- /dev/null +++ b/dataprovider-api/src/test/java/nl/hauntedmc/dataprovider/exception/DataProviderExceptionTest.java @@ -0,0 +1,68 @@ +package nl.hauntedmc.dataprovider.exception; + +import nl.hauntedmc.dataprovider.database.DatabaseType; +import org.junit.jupiter.api.Test; + +import java.util.HashMap; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; + +class DataProviderExceptionTest { + + @Test + void contextIsImmutableAndMetadataIsStable() { + Map diagnostics = new HashMap<>(); + diagnostics.put("sqlState", "23000"); + DataProviderRegistrationException exception = new DataProviderRegistrationException( + "Registration failed safely.", + new DataProviderFailureContext( + DatabaseType.MYSQL, + "main", + "registerDatabase", + RetryAdvice.CONDITIONAL, + ExecutionOutcome.NOT_STARTED, + diagnostics, + null + ), + null + ); + diagnostics.put("sqlState", "changed"); + + assertEquals(DataProviderErrorCode.REGISTRATION_FAILED, exception.errorCode()); + assertEquals("23000", exception.diagnostics().get("sqlState")); + assertThrows(UnsupportedOperationException.class, + () -> exception.diagnostics().put("other", "value")); + assertNotNull(exception.diagnosticId()); + } + + @Test + void sensitiveDiagnosticKeysAreRejected() { + DataProviderFailureContext context = context(Map.of("password", "must-not-appear"), null); + + assertThrows(IllegalArgumentException.class, + () -> new DataProviderRegistrationException("Safe message.", context, null)); + } + + @Test + void malformedDiagnosticIdentifiersAreRejected() { + DataProviderFailureContext context = context(Map.of(), "invalid identifier with spaces"); + + assertThrows(IllegalArgumentException.class, + () -> new DataProviderRegistrationException("Safe message.", context, null)); + } + + private static DataProviderFailureContext context(Map diagnostics, String diagnosticId) { + return new DataProviderFailureContext( + DatabaseType.REDIS, + "cache", + "redis.getKey", + RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED, + diagnostics, + diagnosticId + ); + } +} diff --git a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/DataProviderHandler.java b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/DataProviderHandler.java index cce80e9..6990d5b 100644 --- a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/DataProviderHandler.java +++ b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/DataProviderHandler.java @@ -4,12 +4,18 @@ import nl.hauntedmc.dataprovider.core.concurrent.DataProviderExecutionRuntime; import nl.hauntedmc.dataprovider.core.concurrent.ExecutionRuntimeConfig; import nl.hauntedmc.dataprovider.core.config.ConfigHandler; +import nl.hauntedmc.dataprovider.core.exception.DataProviderExceptionMapper; import nl.hauntedmc.dataprovider.core.identity.CallerContext; import nl.hauntedmc.dataprovider.core.identity.CallerContextResolver; import nl.hauntedmc.dataprovider.core.identity.StackCallerClassLoaderResolver; import nl.hauntedmc.dataprovider.database.DatabaseConnectionKey; import nl.hauntedmc.dataprovider.database.DatabaseProvider; import nl.hauntedmc.dataprovider.database.DatabaseType; +import nl.hauntedmc.dataprovider.exception.DataProviderFailureContext; +import nl.hauntedmc.dataprovider.exception.DataProviderRegistrationException; +import nl.hauntedmc.dataprovider.exception.ExecutionOutcome; +import nl.hauntedmc.dataprovider.exception.ProviderClosedException; +import nl.hauntedmc.dataprovider.exception.RetryAdvice; import nl.hauntedmc.dataprovider.logging.LoggerAdapter; import java.nio.file.Path; @@ -42,8 +48,10 @@ public DataProviderHandler( Objects.requireNonNull(resourceClassLoader, "Resource class loader cannot be null."); Objects.requireNonNull(configHandler, "Config handler cannot be null."); this.logger = Objects.requireNonNull(logger, "Logger cannot be null."); - this.callerContextResolver = Objects.requireNonNull(callerContextResolver, - "Caller context resolver cannot be null."); + this.callerContextResolver = Objects.requireNonNull( + callerContextResolver, + "Caller context resolver cannot be null." + ); ownClassLoader = resourceClassLoader; DatabaseConfigMap configMap = new DatabaseConfigMap(dataPath, this.logger, resourceClassLoader); executionRuntime = new DataProviderExecutionRuntime(ExecutionRuntimeConfig.from(configHandler.getConfig())); @@ -59,19 +67,35 @@ public DataProviderHandler( ) { this.logger = Objects.requireNonNull(logger, "Logger cannot be null."); this.registry = Objects.requireNonNull(registry, "Registry cannot be null."); - this.callerContextResolver = Objects.requireNonNull(callerContextResolver, - "Caller context resolver cannot be null."); + this.callerContextResolver = Objects.requireNonNull( + callerContextResolver, + "Caller context resolver cannot be null." + ); this.ownClassLoader = Objects.requireNonNull(ownClassLoader, "Own class loader cannot be null."); executionRuntime = null; } public DatabaseProvider registerDatabase(DatabaseType databaseType, String connectionIdentifier) { - requireOpen(); - Objects.requireNonNull(databaseType, "Database type cannot be null"); - PluginId pluginId = PluginId.of(resolveCallerContext().pluginId()); - ConnectionIdentifier identifier = ConnectionIdentifier.of(connectionIdentifier); - return DatabaseFactory.withCreationPlugin(pluginId, () -> registry.registerDatabase( - pluginId, OwnerScopeId.of(pluginId.value()), databaseType, identifier)); + requireLegacyOpen(); + PluginId pluginId = resolvePluginId(); + return registerLegacy( + pluginId, + OwnerScopeId.of(pluginId.value()), + requireType(databaseType), + ConnectionIdentifier.of(connectionIdentifier) + ); + } + + public DatabaseProvider registerDatabaseOrThrow(DatabaseType databaseType, String connectionIdentifier) { + requireStructuredOpen("registerDatabase"); + PluginId pluginId = resolvePluginId(); + return registerStrict( + pluginId, + OwnerScopeId.of(pluginId.value()), + requireType(databaseType), + ConnectionIdentifier.of(connectionIdentifier), + "registerDatabase" + ); } public DatabaseProvider registerDatabaseForScope( @@ -87,28 +111,44 @@ public DatabaseProvider registerDatabaseForScope( DatabaseType databaseType, String connectionIdentifier ) { - requireOpen(); - Objects.requireNonNull(databaseType, "Database type cannot be null"); - Objects.requireNonNull(ownerScope, "Owner scope cannot be null."); - PluginId pluginId = PluginId.of(resolveCallerContext().pluginId()); - ConnectionIdentifier identifier = ConnectionIdentifier.of(connectionIdentifier); - return DatabaseFactory.withCreationPlugin(pluginId, () -> registry.registerDatabase( - pluginId, OwnerScopeId.from(ownerScope), databaseType, identifier)); - } - - public void unregisterDatabase(DatabaseType databaseType, String connectionIdentifier) { - requireOpen(); - Objects.requireNonNull(databaseType, "Database type cannot be null"); - PluginId pluginId = PluginId.of(resolveCallerContext().pluginId()); - registry.unregisterDatabase(pluginId, OwnerScopeId.of(pluginId.value()), databaseType, - ConnectionIdentifier.of(connectionIdentifier)); + requireLegacyOpen(); + PluginId pluginId = resolvePluginId(); + return registerLegacy( + pluginId, + OwnerScopeId.from(Objects.requireNonNull(ownerScope, "Owner scope cannot be null.")), + requireType(databaseType), + ConnectionIdentifier.of(connectionIdentifier) + ); } - public void unregisterDatabaseForScope( - String ownerScope, + public DatabaseProvider registerDatabaseForScopeOrThrow( + OwnerScope ownerScope, DatabaseType databaseType, String connectionIdentifier ) { + requireStructuredOpen("scope.registerDatabase"); + PluginId pluginId = resolvePluginId(); + return registerStrict( + pluginId, + OwnerScopeId.from(Objects.requireNonNull(ownerScope, "Owner scope cannot be null.")), + requireType(databaseType), + ConnectionIdentifier.of(connectionIdentifier), + "scope.registerDatabase" + ); + } + + public void unregisterDatabase(DatabaseType databaseType, String connectionIdentifier) { + requireLegacyOpen(); + PluginId pluginId = resolvePluginId(); + registry.unregisterDatabase( + pluginId, + OwnerScopeId.of(pluginId.value()), + requireType(databaseType), + ConnectionIdentifier.of(connectionIdentifier) + ); + } + + public void unregisterDatabaseForScope(String ownerScope, DatabaseType databaseType, String connectionIdentifier) { unregisterDatabaseForScope(OwnerScope.of(ownerScope), databaseType, connectionIdentifier); } @@ -117,17 +157,18 @@ public void unregisterDatabaseForScope( DatabaseType databaseType, String connectionIdentifier ) { - requireOpen(); - Objects.requireNonNull(databaseType, "Database type cannot be null"); - Objects.requireNonNull(ownerScope, "Owner scope cannot be null."); - PluginId pluginId = PluginId.of(resolveCallerContext().pluginId()); - registry.unregisterDatabase(pluginId, OwnerScopeId.from(ownerScope), databaseType, - ConnectionIdentifier.of(connectionIdentifier)); + requireLegacyOpen(); + registry.unregisterDatabase( + resolvePluginId(), + OwnerScopeId.from(Objects.requireNonNull(ownerScope, "Owner scope cannot be null.")), + requireType(databaseType), + ConnectionIdentifier.of(connectionIdentifier) + ); } public void unregisterAllDatabases() { - requireOpen(); - PluginId pluginId = PluginId.of(resolveCallerContext().pluginId()); + requireLegacyOpen(); + PluginId pluginId = resolvePluginId(); registry.unregisterAllDatabases(pluginId, OwnerScopeId.of(pluginId.value())); } @@ -136,15 +177,16 @@ public void unregisterAllDatabasesForScope(String ownerScope) { } public void unregisterAllDatabasesForScope(OwnerScope ownerScope) { - requireOpen(); - Objects.requireNonNull(ownerScope, "Owner scope cannot be null."); - PluginId pluginId = PluginId.of(resolveCallerContext().pluginId()); - registry.unregisterAllDatabases(pluginId, OwnerScopeId.from(ownerScope)); + requireLegacyOpen(); + registry.unregisterAllDatabases( + resolvePluginId(), + OwnerScopeId.from(Objects.requireNonNull(ownerScope, "Owner scope cannot be null.")) + ); } public void unregisterAllDatabasesForPlugin() { - requireOpen(); - registry.unregisterAllDatabasesForPlugin(PluginId.of(resolveCallerContext().pluginId())); + requireLegacyOpen(); + registry.unregisterAllDatabasesForPlugin(resolvePluginId()); } public void shutdownAllDatabases() { @@ -159,10 +201,27 @@ public void shutdownAllDatabases() { } public DatabaseProvider getRegisteredDatabase(DatabaseType databaseType, String connectionIdentifier) { - requireOpen(); - Objects.requireNonNull(databaseType, "Database type cannot be null"); - return registry.getDatabase(PluginId.of(resolveCallerContext().pluginId()), databaseType, - ConnectionIdentifier.of(connectionIdentifier)); + requireLegacyOpen(); + return registry.getDatabase( + resolvePluginId(), + requireType(databaseType), + ConnectionIdentifier.of(connectionIdentifier) + ); + } + + public DatabaseProvider requireRegisteredDatabase(DatabaseType databaseType, String connectionIdentifier) { + requireStructuredOpen("requireRegisteredDatabase"); + DatabaseType type = requireType(databaseType); + String identifier = ConnectionIdentifier.of(connectionIdentifier).value(); + DatabaseProvider provider = registry.getDatabase( + resolvePluginId(), + type, + ConnectionIdentifier.of(identifier) + ); + if (provider != null) { + return provider; + } + throw missingRegistration(type, identifier, "requireRegisteredDatabase"); } public DatabaseProvider getRegisteredDatabaseForScope( @@ -170,53 +229,133 @@ public DatabaseProvider getRegisteredDatabaseForScope( DatabaseType databaseType, String connectionIdentifier ) { - requireOpen(); - Objects.requireNonNull(ownerScope, "Owner scope cannot be null."); - Objects.requireNonNull(databaseType, "Database type cannot be null"); - return registry.getDatabase(PluginId.of(resolveCallerContext().pluginId()), OwnerScopeId.from(ownerScope), - databaseType, ConnectionIdentifier.of(connectionIdentifier)); + requireLegacyOpen(); + return registry.getDatabase( + resolvePluginId(), + OwnerScopeId.from(Objects.requireNonNull(ownerScope, "Owner scope cannot be null.")), + requireType(databaseType), + ConnectionIdentifier.of(connectionIdentifier) + ); + } + + public DatabaseProvider requireRegisteredDatabaseForScope( + OwnerScope ownerScope, + DatabaseType databaseType, + String connectionIdentifier + ) { + requireStructuredOpen("scope.requireRegisteredDatabase"); + DatabaseType type = requireType(databaseType); + ConnectionIdentifier identifier = ConnectionIdentifier.of(connectionIdentifier); + DatabaseProvider provider = registry.getDatabase( + resolvePluginId(), + OwnerScopeId.from(Objects.requireNonNull(ownerScope, "Owner scope cannot be null.")), + type, + identifier + ); + if (provider != null) { + return provider; + } + throw missingRegistration(type, identifier.value(), "scope.requireRegisteredDatabase"); } public ConcurrentMap getActiveDatabases() { - requireOpen(); + requireLegacyOpen(); requireInternalCaller(); return registry.getActiveDatabases(); } public Map getActiveDatabaseReferenceCounts() { - requireOpen(); + requireLegacyOpen(); requireInternalCaller(); return registry.getActiveDatabaseReferenceCounts(); } public Map getCachedDatabaseHealth() { - requireOpen(); + requireLegacyOpen(); requireInternalCaller(); return registry.getCachedHealthSnapshots(); } public CompletableFuture probeDatabaseHealthAsync() { - requireOpen(); + requireLegacyOpen(); requireInternalCaller(); return registry.probeRemoteHealthAsync(); } public Map getConfiguredDatabaseTypeStates() { - requireOpen(); + requireLegacyOpen(); requireInternalCaller(); return registry.getConfiguredDatabaseTypeStates(); } public String getConfiguredOrmSchemaMode() { - requireOpen(); + requireLegacyOpen(); requireInternalCaller(); return registry.getOrmSchemaMode(); } public void reloadConfiguration() { - requireOpen(); + requireLegacyOpen(); requireInternalCaller(); - registry.reloadConfiguration(); + try { + registry.reloadConfiguration(); + } catch (RuntimeException failure) { + throw DataProviderExceptionMapper.configurationFailure(failure, "reloadConfiguration"); + } + } + + private DatabaseProvider registerLegacy( + PluginId pluginId, + OwnerScopeId ownerScope, + DatabaseType type, + ConnectionIdentifier identifier + ) { + return DatabaseFactory.withCreationPlugin( + pluginId, + () -> registry.registerDatabase(pluginId, ownerScope, type, identifier) + ); + } + + private DatabaseProvider registerStrict( + PluginId pluginId, + OwnerScopeId ownerScope, + DatabaseType type, + ConnectionIdentifier identifier, + String operation + ) { + DatabaseProvider provider = registerLegacy(pluginId, ownerScope, type, identifier); + if (provider != null) { + return provider; + } + if (!registry.getConfiguredDatabaseTypeStates().getOrDefault(type, true)) { + throw DataProviderExceptionMapper.backendDisabled(type, identifier.value()); + } + DatabaseConnectionKey key = new DatabaseConnectionKey(pluginId.value(), type, identifier.value()); + ProviderLifecycleSnapshot snapshot = registry.getProviderLifecycleSnapshots().get(key); + Throwable failure = snapshot == null ? null : snapshot.failure(); + throw DataProviderExceptionMapper.registrationFailure(failure, type, identifier.value(), operation); + } + + private static DataProviderRegistrationException missingRegistration( + DatabaseType type, + String identifier, + String operation + ) { + return new DataProviderRegistrationException( + "No active database registration exists for the requested connection.", + DataProviderFailureContext.of( + type, + identifier, + operation, + RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED + ), + null + ); + } + + private PluginId resolvePluginId() { + return PluginId.of(resolveCallerContext().pluginId()); } private CallerContext resolveCallerContext() { @@ -230,16 +369,37 @@ private CallerContext resolveCallerContext() { private void requireInternalCaller() { ClassLoader callerLoader = StackCallerClassLoaderResolver.resolveNearestCallerOutsidePackage( - INTERNAL_PACKAGE_PREFIX); + INTERNAL_PACKAGE_PREFIX + ); if (callerLoader == null || callerLoader != ownClassLoader) { logger.error("Rejected privileged operation from non-internal caller."); throw new SecurityException("Privileged DataProvider operation is restricted to internal callers."); } } - private void requireOpen() { + private void requireLegacyOpen() { if (registry.isClosed()) { throw new IllegalStateException(CLOSED_MESSAGE); } } + + private void requireStructuredOpen(String operation) { + if (registry.isClosed()) { + throw new ProviderClosedException( + CLOSED_MESSAGE, + DataProviderFailureContext.of( + null, + null, + operation, + RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED + ), + null + ); + } + } + + private static DatabaseType requireType(DatabaseType databaseType) { + return Objects.requireNonNull(databaseType, "Database type cannot be null"); + } } diff --git a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/DatabaseFactory.java b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/DatabaseFactory.java index 5bd2d93..61db6ef 100644 --- a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/DatabaseFactory.java +++ b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/DatabaseFactory.java @@ -1,11 +1,13 @@ package nl.hauntedmc.dataprovider.core; +import nl.hauntedmc.dataprovider.core.concurrent.ContextualExecutionHandle; import nl.hauntedmc.dataprovider.core.concurrent.DataProviderExecutionRuntime; import nl.hauntedmc.dataprovider.core.concurrent.ExecutionHandle; import nl.hauntedmc.dataprovider.core.database.document.impl.mongodb.MongoDBDatabase; import nl.hauntedmc.dataprovider.core.database.keyvalue.impl.redis.RedisDatabase; import nl.hauntedmc.dataprovider.core.database.messaging.impl.redis.RedisMessagingDatabase; import nl.hauntedmc.dataprovider.core.database.relational.impl.mysql.MySQLDatabase; +import nl.hauntedmc.dataprovider.core.exception.DataProviderExceptionMapper; import nl.hauntedmc.dataprovider.database.DatabaseType; import nl.hauntedmc.dataprovider.logging.LoggerAdapter; import org.spongepowered.configurate.CommentedConfigurationNode; @@ -74,11 +76,17 @@ protected ManagedDatabaseProvider createDatabaseProvider( CommentedConfigurationNode connectionConfig = configMap.getConfig(type, connectionIdentifier); if (connectionConfig == null) { logger.error("Could not load configuration for " + connectionIdentifier.value() + " (" + type.name() + ")"); - return null; + throw DataProviderExceptionMapper.missingConfigurationFailure(); } - ExecutionHandle execution = executionRuntime == null + ExecutionHandle rawExecution = executionRuntime == null ? ExecutionHandle.direct() : executionRuntime.openScope(pluginId.value(), type, connectionIdentifier.value()); + ExecutionHandle execution = new ContextualExecutionHandle( + rawExecution, + pluginId.value(), + type, + connectionIdentifier.value() + ); try { return switch (type) { case MYSQL -> new MySQLDatabase(connectionConfig, logger, execution); diff --git a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/api/DefaultDataProviderApi.java b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/api/DefaultDataProviderApi.java index f95c0d7..4e7a7a9 100644 --- a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/api/DefaultDataProviderApi.java +++ b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/api/DefaultDataProviderApi.java @@ -4,10 +4,10 @@ import nl.hauntedmc.dataprovider.api.DataProviderScope; import nl.hauntedmc.dataprovider.api.OwnerScope; import nl.hauntedmc.dataprovider.api.orm.ORMContext; - +import nl.hauntedmc.dataprovider.core.DataProviderHandler; import nl.hauntedmc.dataprovider.database.DataAccess; -import nl.hauntedmc.dataprovider.database.DatabaseType; import nl.hauntedmc.dataprovider.database.DatabaseProvider; +import nl.hauntedmc.dataprovider.database.DatabaseType; import nl.hauntedmc.dataprovider.database.document.DocumentDataAccess; import nl.hauntedmc.dataprovider.database.document.DocumentDatabaseProvider; import nl.hauntedmc.dataprovider.database.keyvalue.KeyValueDataAccess; @@ -17,31 +17,15 @@ import nl.hauntedmc.dataprovider.database.relational.RelationalDataAccess; import nl.hauntedmc.dataprovider.database.relational.RelationalDatabaseProvider; import nl.hauntedmc.dataprovider.database.relational.schema.SchemaManager; -import nl.hauntedmc.dataprovider.core.DataProviderHandler; import javax.sql.DataSource; import java.util.Objects; -/** - * DataProviderAPI is the public facade that exposes safe, read-only database handles - * for third-party plugins. Internally, it delegates to a DataProviderHandler, but it does not - * expose lifecycle-sensitive methods (like shutdownAllDatabases or getActiveDatabases). - * - * For most integrations, the primary lifecycle is: - * register -> use provider/data access -> unregister. - * Optional scoped ownership is available through {@link #scope(String)} for advanced cases - * where one plugin/software process needs isolated ownership domains for independently - * managed components. - */ +/** Public read-only facade for plugin-scoped DataProvider access. */ public final class DefaultDataProviderApi implements DataProviderAPI { private final DataProviderHandler handler; - /** - * Constructs the API wrapper. - * - * @param handler the internal DataProviderHandler instance. - */ public DefaultDataProviderApi(DataProviderHandler handler) { this.handler = Objects.requireNonNull(handler, "DataProviderHandler cannot be null"); } @@ -55,70 +39,49 @@ public ORMContext createOrmContext( Class... entityClasses ) { return new nl.hauntedmc.dataprovider.core.orm.ORMContext( - pluginName, - dataSource, - logger, - schemaMode, - entityClasses - ); - } - - /** - * Registers a database connection for the resolved caller plugin. - * This is the default path for most integrations. - * - * @param databaseType the type of database (e.g. MYSQL, MONGODB, etc.) - * @param connectionIdentifier a unique identifier for the connection - * @return the registered read-only {@link DatabaseProvider} handle. - */ + pluginName, dataSource, logger, schemaMode, entityClasses); + } + + @Override public DatabaseProvider registerDatabase(DatabaseType databaseType, String connectionIdentifier) { return wrapProvider(handler.registerDatabase(databaseType, connectionIdentifier)); } - /** - * Creates an optional scoped lifecycle facade using a typed owner scope. - */ + @Override + public DatabaseProvider registerDatabaseOrThrow(DatabaseType databaseType, String connectionIdentifier) { + return wrapProvider(handler.registerDatabaseOrThrow(databaseType, connectionIdentifier)); + } + + @Override public DataProviderScope scope(OwnerScope ownerScope) { return new DefaultDataProviderScope(handler, ownerScope); } - /** - * Unregisters a specific database connection for the resolved caller plugin. - * This is the default path for most integrations. - * - * @param databaseType the type of database. - * @param connectionIdentifier the connection identifier. - */ + @Override public void unregisterDatabase(DatabaseType databaseType, String connectionIdentifier) { handler.unregisterDatabase(databaseType, connectionIdentifier); } - /** - * Unregisters all database connections for the resolved caller plugin default owner scope. - */ + @Override public void unregisterAllDatabases() { handler.unregisterAllDatabases(); } - /** - * Unregisters all database connections for the caller plugin across all caller scopes. - * Use this for deterministic full-plugin shutdown cleanup. - */ + @Override public void unregisterAllDatabasesForPlugin() { handler.unregisterAllDatabasesForPlugin(); } - /** - * Retrieves a registered database connection for the resolved caller plugin. - * - * @param databaseType the type of database. - * @param connectionIdentifier the connection identifier. - * @return the {@link DatabaseProvider} instance, or null if not registered. - */ + @Override public DatabaseProvider getRegisteredDatabase(DatabaseType databaseType, String connectionIdentifier) { return wrapProvider(handler.getRegisteredDatabase(databaseType, connectionIdentifier)); } + @Override + public DatabaseProvider requireRegisteredDatabase(DatabaseType databaseType, String connectionIdentifier) { + return wrapProvider(handler.requireRegisteredDatabase(databaseType, connectionIdentifier)); + } + static DatabaseProvider wrapProvider(DatabaseProvider provider) { if (provider == null || provider instanceof WrappedDatabaseProvider) { return provider; @@ -145,21 +108,9 @@ private record DatabaseProviderView(DatabaseProvider delegate) implements Wrappe private DatabaseProviderView { Objects.requireNonNull(delegate, "Delegate database provider cannot be null."); } - - @Override - public boolean isConnected() { - return delegate.isConnected(); - } - - @Override - public DataAccess getDataAccess() { - return delegate.getDataAccess(); - } - - @Override - public DataSource getDataSource() { - return delegate.getDataSource(); - } + @Override public boolean isConnected() { return delegate.isConnected(); } + @Override public DataAccess getDataAccess() { return delegate.getDataAccess(); } + @Override public DataSource getDataSource() { return delegate.getDataSource(); } } private record RelationalDatabaseProviderView(RelationalDatabaseProvider delegate) @@ -167,26 +118,10 @@ private record RelationalDatabaseProviderView(RelationalDatabaseProvider delegat private RelationalDatabaseProviderView { Objects.requireNonNull(delegate, "Delegate relational database provider cannot be null."); } - - @Override - public boolean isConnected() { - return delegate.isConnected(); - } - - @Override - public RelationalDataAccess getDataAccess() { - return delegate.getDataAccess(); - } - - @Override - public DataSource getDataSource() { - return delegate.getDataSource(); - } - - @Override - public SchemaManager getSchemaManager() { - return delegate.getSchemaManager(); - } + @Override public boolean isConnected() { return delegate.isConnected(); } + @Override public RelationalDataAccess getDataAccess() { return delegate.getDataAccess(); } + @Override public DataSource getDataSource() { return delegate.getDataSource(); } + @Override public SchemaManager getSchemaManager() { return delegate.getSchemaManager(); } } private record DocumentDatabaseProviderView(DocumentDatabaseProvider delegate) @@ -194,16 +129,8 @@ private record DocumentDatabaseProviderView(DocumentDatabaseProvider delegate) private DocumentDatabaseProviderView { Objects.requireNonNull(delegate, "Delegate document database provider cannot be null."); } - - @Override - public boolean isConnected() { - return delegate.isConnected(); - } - - @Override - public DocumentDataAccess getDataAccess() { - return delegate.getDataAccess(); - } + @Override public boolean isConnected() { return delegate.isConnected(); } + @Override public DocumentDataAccess getDataAccess() { return delegate.getDataAccess(); } } private record KeyValueDatabaseProviderView(KeyValueDatabaseProvider delegate) @@ -211,16 +138,8 @@ private record KeyValueDatabaseProviderView(KeyValueDatabaseProvider delegate) private KeyValueDatabaseProviderView { Objects.requireNonNull(delegate, "Delegate key-value database provider cannot be null."); } - - @Override - public boolean isConnected() { - return delegate.isConnected(); - } - - @Override - public KeyValueDataAccess getDataAccess() { - return delegate.getDataAccess(); - } + @Override public boolean isConnected() { return delegate.isConnected(); } + @Override public KeyValueDataAccess getDataAccess() { return delegate.getDataAccess(); } } private record MessagingDatabaseProviderView(MessagingDatabaseProvider delegate) @@ -228,15 +147,7 @@ private record MessagingDatabaseProviderView(MessagingDatabaseProvider delegate) private MessagingDatabaseProviderView { Objects.requireNonNull(delegate, "Delegate messaging database provider cannot be null."); } - - @Override - public boolean isConnected() { - return delegate.isConnected(); - } - - @Override - public MessagingDataAccess getDataAccess() { - return delegate.getDataAccess(); - } + @Override public boolean isConnected() { return delegate.isConnected(); } + @Override public MessagingDataAccess getDataAccess() { return delegate.getDataAccess(); } } } diff --git a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/api/DefaultDataProviderScope.java b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/api/DefaultDataProviderScope.java index 3c64f6c..2aaacd1 100644 --- a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/api/DefaultDataProviderScope.java +++ b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/api/DefaultDataProviderScope.java @@ -2,23 +2,18 @@ import nl.hauntedmc.dataprovider.api.DataProviderScope; import nl.hauntedmc.dataprovider.api.OwnerScope; - +import nl.hauntedmc.dataprovider.core.DataProviderHandler; import nl.hauntedmc.dataprovider.database.DatabaseProvider; import nl.hauntedmc.dataprovider.database.DatabaseType; -import nl.hauntedmc.dataprovider.core.DataProviderHandler; +import nl.hauntedmc.dataprovider.exception.DataProviderFailureContext; +import nl.hauntedmc.dataprovider.exception.ExecutionOutcome; +import nl.hauntedmc.dataprovider.exception.ProviderClosedException; +import nl.hauntedmc.dataprovider.exception.RetryAdvice; +import java.util.Map; import java.util.Objects; -/** - * Optional scoped lifecycle helper for advanced integrations that need isolated ownership domains - * inside one plugin/software process. - * - * Typical use: - * - create one scope per logical component - * - register and use connections through this scope - * - release the scope's registrations via {@link #unregisterAllDatabases()} or terminate the - * scope via {@link #close()} - */ +/** Optional scoped lifecycle helper for independently managed plugin components. */ public final class DefaultDataProviderScope implements DataProviderScope { private static final String CLOSED_MESSAGE = "DataProvider scope is closed."; @@ -33,9 +28,7 @@ public final class DefaultDataProviderScope implements DataProviderScope { this.ownerScope = Objects.requireNonNull(ownerScope, "Owner scope cannot be null."); } - /** - * Returns the normalized scope identifier used for ownership tracking. - */ + @Override public OwnerScope ownerScope() { return ownerScope; } @@ -45,34 +38,38 @@ public LifecycleState lifecycleState() { return lifecycleState; } - /** - * Registers a database connection under this scope. - */ + @Override public DatabaseProvider registerDatabase(DatabaseType databaseType, String connectionIdentifier) { synchronized (lifecycleMonitor) { - requireOpen(); + requireLegacyOpen(); return DefaultDataProviderApi.wrapProvider( handler.registerDatabaseForScope(ownerScope, databaseType, connectionIdentifier) ); } } - /** - * Releases one scoped registration reference. - */ + @Override + public DatabaseProvider registerDatabaseOrThrow(DatabaseType databaseType, String connectionIdentifier) { + synchronized (lifecycleMonitor) { + requireStructuredOpen("scope.registerDatabase"); + return DefaultDataProviderApi.wrapProvider( + handler.registerDatabaseForScopeOrThrow(ownerScope, databaseType, connectionIdentifier) + ); + } + } + + @Override public void unregisterDatabase(DatabaseType databaseType, String connectionIdentifier) { synchronized (lifecycleMonitor) { - requireOpen(); + requireLegacyOpen(); handler.unregisterDatabaseForScope(ownerScope, databaseType, connectionIdentifier); } } - /** - * Releases all registrations held by this scope. - */ + @Override public void unregisterAllDatabases() { synchronized (lifecycleMonitor) { - requireOpen(); + requireLegacyOpen(); handler.unregisterAllDatabasesForScope(ownerScope); } } @@ -80,13 +77,23 @@ public void unregisterAllDatabases() { @Override public DatabaseProvider getRegisteredDatabase(DatabaseType databaseType, String connectionIdentifier) { synchronized (lifecycleMonitor) { - requireOpen(); + requireLegacyOpen(); return DefaultDataProviderApi.wrapProvider( handler.getRegisteredDatabaseForScope(ownerScope, databaseType, connectionIdentifier) ); } } + @Override + public DatabaseProvider requireRegisteredDatabase(DatabaseType databaseType, String connectionIdentifier) { + synchronized (lifecycleMonitor) { + requireStructuredOpen("scope.requireRegisteredDatabase"); + return DefaultDataProviderApi.wrapProvider( + handler.requireRegisteredDatabaseForScope(ownerScope, databaseType, connectionIdentifier) + ); + } + } + @Override public void close() { synchronized (lifecycleMonitor) { @@ -102,9 +109,25 @@ public void close() { } } - private void requireOpen() { + private void requireLegacyOpen() { if (lifecycleState != LifecycleState.OPEN) { throw new IllegalStateException(CLOSED_MESSAGE); } } + + private void requireStructuredOpen(String operation) { + if (lifecycleState != LifecycleState.OPEN) { + throw new ProviderClosedException( + CLOSED_MESSAGE, + DataProviderFailureContext.of( + null, + null, + operation, + RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED + ).withDiagnostics(Map.of("ownerScope", ownerScope.value())), + null + ); + } + } } diff --git a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/concurrent/AsyncTaskSupport.java b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/concurrent/AsyncTaskSupport.java index 3a46749..9d4283d 100644 --- a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/concurrent/AsyncTaskSupport.java +++ b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/concurrent/AsyncTaskSupport.java @@ -1,11 +1,15 @@ package nl.hauntedmc.dataprovider.core.concurrent; +import nl.hauntedmc.dataprovider.core.exception.DataProviderExceptionMapper; +import nl.hauntedmc.dataprovider.exception.DataProviderException; + import java.util.Objects; +import java.util.concurrent.CancellationException; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executor; import java.util.concurrent.RejectedExecutionException; -/** Shared helpers for queue-backed async execution with rejection-safe futures. */ +/** Shared helpers for queue-backed async execution with rejection-safe structured futures. */ public final class AsyncTaskSupport { private AsyncTaskSupport() { @@ -59,14 +63,18 @@ public void run() { future.complete(supplier.get()); } catch (Throwable throwable) { failed = true; - future.completeExceptionally(throwable); + future.completeExceptionally(mapFailure(throwable, executor, operationName)); } } @Override public void reject(RejectedExecutionException rejection) { failed = true; - future.completeExceptionally(rejection); + future.completeExceptionally(DataProviderExceptionMapper.translate( + rejection, + executor, + operationName + )); } @Override @@ -76,17 +84,39 @@ public boolean failed() { }; try { executor.execute(task); - } catch (ExecutionRejectedException e) { - future.completeExceptionally(e); - } catch (RejectedExecutionException e) { - future.completeExceptionally(new ExecutionRejectedException( + } catch (ExecutionRejectedException rejection) { + future.completeExceptionally(DataProviderExceptionMapper.translate( + rejection, + executor, + operationName + )); + } catch (RejectedExecutionException rejection) { + ExecutionRejectedException structuredRejection = new ExecutionRejectedException( ExecutionRejectedException.Reason.LANE_QUEUE_FULL, "Rejected async operation '" + operationName + "'.", - e + rejection + ); + future.completeExceptionally(DataProviderExceptionMapper.translate( + structuredRejection, + executor, + operationName )); - } catch (RuntimeException e) { - future.completeExceptionally(e); + } catch (RuntimeException failure) { + future.completeExceptionally(mapFailure(failure, executor, operationName)); } return future; } + + private static Throwable mapFailure(Throwable failure, Executor executor, String operationName) { + if (failure instanceof DataProviderException + || failure instanceof IllegalArgumentException + || failure instanceof NullPointerException + || failure instanceof UnsupportedOperationException + || failure instanceof SecurityException + || failure instanceof CancellationException + || failure instanceof Error) { + return failure; + } + return DataProviderExceptionMapper.translate(failure, executor, operationName); + } } diff --git a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/concurrent/ContextualExecutionHandle.java b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/concurrent/ContextualExecutionHandle.java new file mode 100644 index 0000000..031b1a4 --- /dev/null +++ b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/concurrent/ContextualExecutionHandle.java @@ -0,0 +1,83 @@ +package nl.hauntedmc.dataprovider.core.concurrent; + +import nl.hauntedmc.dataprovider.database.DatabaseType; + +import java.util.Objects; + +/** Adds immutable backend identity to a runtime execution scope. */ +public final class ContextualExecutionHandle implements ExecutionHandle { + + private final ExecutionHandle delegate; + private final String pluginId; + private final DatabaseType backendType; + private final String connectionIdentifier; + + public ContextualExecutionHandle( + ExecutionHandle delegate, + String pluginId, + DatabaseType backendType, + String connectionIdentifier + ) { + this.delegate = Objects.requireNonNull(delegate, "Delegate execution handle cannot be null."); + this.pluginId = requireText(pluginId, "pluginId"); + this.backendType = Objects.requireNonNull(backendType, "Backend type cannot be null."); + this.connectionIdentifier = requireText(connectionIdentifier, "connectionIdentifier"); + } + + @Override + public void execute(Runnable command) { + delegate.execute(command); + } + + @Override + public ExecutionMetricsSnapshot metrics() { + return delegate.metrics(); + } + + @Override + public boolean isClosed() { + return delegate.isClosed(); + } + + @Override + public DatabaseType backendType() { + return backendType; + } + + @Override + public String connectionIdentifier() { + return connectionIdentifier; + } + + @Override + public String pluginId() { + return pluginId; + } + + @Override + public boolean tryAcquireSubscription() { + return delegate.tryAcquireSubscription(); + } + + @Override + public void releaseSubscription() { + delegate.releaseSubscription(); + } + + @Override + public void recordDroppedMessages(long count) { + delegate.recordDroppedMessages(count); + } + + @Override + public void close() { + delegate.close(); + } + + private static String requireText(String value, String field) { + if (value == null || value.isBlank()) { + throw new IllegalArgumentException(field + " cannot be null or blank."); + } + return value.trim(); + } +} diff --git a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/concurrent/ExecutionHandle.java b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/concurrent/ExecutionHandle.java index d4f4ffb..3ef6159 100644 --- a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/concurrent/ExecutionHandle.java +++ b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/concurrent/ExecutionHandle.java @@ -1,5 +1,7 @@ package nl.hauntedmc.dataprovider.core.concurrent; +import nl.hauntedmc.dataprovider.database.DatabaseType; + import java.util.concurrent.Executor; /** Connection-scoped execution handle backed by the shared runtime. */ @@ -9,6 +11,18 @@ public interface ExecutionHandle extends Executor, AutoCloseable { boolean isClosed(); + default DatabaseType backendType() { + return null; + } + + default String connectionIdentifier() { + return null; + } + + default String pluginId() { + return null; + } + default boolean tryAcquireSubscription() { return true; } diff --git a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/database/messaging/impl/redis/RedisMessagingDataAccess.java b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/database/messaging/impl/redis/RedisMessagingDataAccess.java index 31a3d50..afba33e 100644 --- a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/database/messaging/impl/redis/RedisMessagingDataAccess.java +++ b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/database/messaging/impl/redis/RedisMessagingDataAccess.java @@ -2,6 +2,9 @@ import nl.hauntedmc.dataprovider.core.concurrent.AsyncTaskSupport; import nl.hauntedmc.dataprovider.core.concurrent.ExecutionHandle; +import nl.hauntedmc.dataprovider.core.concurrent.ExecutionRejectedException; +import nl.hauntedmc.dataprovider.core.exception.DataProviderExceptionMapper; +import nl.hauntedmc.dataprovider.core.exception.StructuredFailures; import nl.hauntedmc.dataprovider.database.messaging.MessagingDataAccess; import nl.hauntedmc.dataprovider.database.messaging.api.EventMessage; import nl.hauntedmc.dataprovider.database.messaging.api.MessageRegistry; @@ -96,12 +99,21 @@ public CompletableFuture publish(String destinati String validatedDestination = validateDestination(destination); Objects.requireNonNull(message, "Message cannot be null"); if (shuttingDown.get()) { - return CompletableFuture.failedFuture(new IllegalStateException("Messaging provider is shutting down")); + return CompletableFuture.failedFuture(StructuredFailures.closed(workers, "redis.messaging.publish")); + } + final String json; + try { + json = messageRegistry.toJson(message); + } catch (RuntimeException failure) { + return CompletableFuture.failedFuture( + StructuredFailures.serialization(failure, workers, "redis.messaging.serialize")); } - String json = messageRegistry.toJson(message); if (json.length() > maxPayloadChars) { - return CompletableFuture.failedFuture(new IllegalArgumentException( - "Message payload exceeds maxPayloadChars (" + maxPayloadChars + ")")); + return CompletableFuture.failedFuture(StructuredFailures.serialization( + new IllegalArgumentException("Serialized message exceeds configured size limit."), + workers, + "redis.messaging.serialize" + )); } return AsyncTaskSupport.runAsync(workers, "redis.messaging.publish", () -> { try (Jedis jedis = pool.getResource()) { @@ -120,7 +132,7 @@ public Subscription subscribe( Objects.requireNonNull(type, "Type cannot be null"); Objects.requireNonNull(handler, "Handler cannot be null"); if (shuttingDown.get()) { - throw new IllegalStateException("Messaging provider is shutting down"); + throw StructuredFailures.closed(workers, "redis.messaging.subscribe"); } ChannelSubscription channelSubscription; @@ -129,11 +141,22 @@ public Subscription subscribe( channelSubscription = channelSubscriptions.get(validatedDestination); if (channelSubscription == null) { if (channelSubscriptions.size() >= maxSubscriptions) { - throw new IllegalStateException( - "Maximum active Redis subscriptions reached (" + maxSubscriptions + ")"); + throw DataProviderExceptionMapper.translate( + new ExecutionRejectedException( + ExecutionRejectedException.Reason.SUBSCRIPTION_LIMIT, + "Connection subscription limit reached."), + workers, + "redis.messaging.subscribe" + ); } if (executionBudget != null && !executionBudget.tryAcquireSubscription()) { - throw new IllegalStateException("DataProvider messaging subscription budget exhausted"); + throw DataProviderExceptionMapper.translate( + new ExecutionRejectedException( + ExecutionRejectedException.Reason.SUBSCRIPTION_LIMIT, + "Runtime subscription limit reached."), + workers, + "redis.messaging.subscribe" + ); } channelSubscription = new ChannelSubscription(validatedDestination, executionBudget != null); channelSubscriptions.put(validatedDestination, channelSubscription); diff --git a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/database/relational/impl/mysql/MySQLDataAccess.java b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/database/relational/impl/mysql/MySQLDataAccess.java index aeecf22..6ec08eb 100644 --- a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/database/relational/impl/mysql/MySQLDataAccess.java +++ b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/database/relational/impl/mysql/MySQLDataAccess.java @@ -1,8 +1,12 @@ package nl.hauntedmc.dataprovider.core.database.relational.impl.mysql; import nl.hauntedmc.dataprovider.core.concurrent.AsyncTaskSupport; +import nl.hauntedmc.dataprovider.core.exception.DataProviderExceptionMapper; import nl.hauntedmc.dataprovider.database.relational.RelationalDataAccess; import nl.hauntedmc.dataprovider.database.relational.TransactionCallback; +import nl.hauntedmc.dataprovider.exception.DataTransactionException; +import nl.hauntedmc.dataprovider.exception.ExecutionOutcome; +import nl.hauntedmc.dataprovider.exception.TransactionPhase; import javax.sql.DataSource; import java.sql.Connection; @@ -22,6 +26,8 @@ /** MySQL implementation of RelationalDataAccess. */ public class MySQLDataAccess implements RelationalDataAccess { + private static final String TRANSACTION_OPERATION = "mysql.executeTransactionally"; + private final DataSource dataSource; private final Executor executor; private final int queryTimeoutSeconds; @@ -47,8 +53,6 @@ public CompletableFuture executeUpdate(String query, Object... params) { applyStatementTuning(stmt); setParameters(stmt, params); stmt.executeUpdate(); - } catch (SQLException e) { - throw new RuntimeException("Failed to execute update", e); } }); } @@ -64,8 +68,6 @@ public CompletableFuture> queryForSingle(String query, Objec try (ResultSet rs = stmt.executeQuery()) { return rs.next() ? mapRow(rs) : null; } - } catch (SQLException e) { - throw new RuntimeException("Failed to execute queryForSingle", e); } }); } @@ -85,8 +87,6 @@ public CompletableFuture>> queryForList(String query, O } } return result; - } catch (SQLException e) { - throw new RuntimeException("Failed to execute queryForList", e); } }); } @@ -102,8 +102,6 @@ public CompletableFuture queryForSingleValue(String query, Object... par try (ResultSet rs = stmt.executeQuery()) { return rs.next() ? rs.getObject(1) : null; } - } catch (SQLException e) { - throw new RuntimeException("Failed to execute queryForSingleValue", e); } }); } @@ -127,14 +125,16 @@ public CompletableFuture executeBatchUpdate(String query, List b } stmt.executeBatch(); connection.commit(); - } catch (SQLException e) { - connection.rollback(); - throw e; + } catch (SQLException primary) { + try { + connection.rollback(); + } catch (SQLException rollbackFailure) { + primary.addSuppressed(rollbackFailure); + } + throw primary; } finally { connection.setAutoCommit(oldAutoCommit); } - } catch (SQLException e) { - throw new RuntimeException("Failed to execute batch update", e); } }); } @@ -142,24 +142,159 @@ public CompletableFuture executeBatchUpdate(String query, List b @Override public CompletableFuture executeTransactionally(TransactionCallback callback) { Objects.requireNonNull(callback, "Transaction callback cannot be null."); - return AsyncTaskSupport.supplyAsync(executor, "mysql.executeTransactionally", () -> { - try (Connection connection = dataSource.getConnection()) { - boolean oldAutoCommit = connection.getAutoCommit(); - connection.setAutoCommit(false); - try { - T result = callback.doInTransaction(connection); - connection.commit(); - return result; - } catch (Exception e) { - connection.rollback(); - throw new RuntimeException("Transaction failed, rolled back.", e); - } finally { - connection.setAutoCommit(oldAutoCommit); - } - } catch (SQLException e) { - throw new RuntimeException("Failed to execute transactionally", e); - } - }); + return AsyncTaskSupport.supplyAsync(executor, TRANSACTION_OPERATION, () -> executeTransaction(callback)); + } + + private T executeTransaction(TransactionCallback callback) { + Connection connection = null; + boolean oldAutoCommit; + try { + connection = dataSource.getConnection(); + oldAutoCommit = connection.getAutoCommit(); + connection.setAutoCommit(false); + } catch (Error fatal) { + closeAfterFatal(connection, fatal); + throw fatal; + } catch (Exception beginFailure) { + DataTransactionException structured = transactionFailure( + beginFailure, TransactionPhase.BEGIN, ExecutionOutcome.NOT_STARTED); + closeConnection(connection, structured, ExecutionOutcome.NOT_STARTED); + throw structured; + } + + T result; + try { + result = callback.doInTransaction(connection); + } catch (Error fatal) { + cleanupAfterFatal(connection, oldAutoCommit, fatal); + throw fatal; + } catch (Exception callbackFailure) { + DataTransactionException structured = transactionFailure( + callbackFailure, TransactionPhase.CALLBACK, ExecutionOutcome.NOT_APPLIED); + rollback(connection, structured); + restoreAutoCommit(connection, oldAutoCommit, structured, ExecutionOutcome.NOT_APPLIED); + closeConnection(connection, structured, ExecutionOutcome.NOT_APPLIED); + throw structured; + } + + try { + connection.commit(); + } catch (Error fatal) { + cleanupAfterFatal(connection, oldAutoCommit, fatal); + throw fatal; + } catch (Exception commitFailure) { + DataTransactionException structured = transactionFailure( + commitFailure, TransactionPhase.COMMIT, ExecutionOutcome.MAY_HAVE_APPLIED); + rollback(connection, structured); + restoreAutoCommit(connection, oldAutoCommit, structured, ExecutionOutcome.MAY_HAVE_APPLIED); + closeConnection(connection, structured, ExecutionOutcome.MAY_HAVE_APPLIED); + throw structured; + } + + try { + connection.setAutoCommit(oldAutoCommit); + } catch (Error fatal) { + closeAfterFatal(connection, fatal); + throw fatal; + } catch (Exception restoreFailure) { + DataTransactionException structured = transactionFailure( + restoreFailure, TransactionPhase.CLEANUP, ExecutionOutcome.MAY_HAVE_APPLIED); + closeConnection(connection, structured, ExecutionOutcome.MAY_HAVE_APPLIED); + throw structured; + } + + try { + connection.close(); + } catch (Error fatal) { + throw fatal; + } catch (Exception closeFailure) { + throw transactionFailure(closeFailure, TransactionPhase.CLEANUP, ExecutionOutcome.MAY_HAVE_APPLIED); + } + return result; + } + + private DataTransactionException transactionFailure( + Throwable failure, + TransactionPhase phase, + ExecutionOutcome outcome + ) { + return DataProviderExceptionMapper.transactionFailure( + failure, executor, TRANSACTION_OPERATION, phase, outcome); + } + + private void rollback(Connection connection, DataTransactionException primary) { + try { + connection.rollback(); + } catch (Error fatal) { + fatal.addSuppressed(primary); + closeAfterFatal(connection, fatal); + throw fatal; + } catch (Exception rollbackFailure) { + primary.addSuppressed(transactionFailure( + rollbackFailure, TransactionPhase.ROLLBACK, ExecutionOutcome.UNKNOWN)); + } + } + + private void restoreAutoCommit( + Connection connection, + boolean oldAutoCommit, + DataTransactionException primary, + ExecutionOutcome outcome + ) { + try { + connection.setAutoCommit(oldAutoCommit); + } catch (Error fatal) { + fatal.addSuppressed(primary); + closeAfterFatal(connection, fatal); + throw fatal; + } catch (Exception restoreFailure) { + primary.addSuppressed(transactionFailure( + restoreFailure, TransactionPhase.CLEANUP, outcome)); + } + } + + private void closeConnection( + Connection connection, + DataTransactionException primary, + ExecutionOutcome outcome + ) { + if (connection == null) { + return; + } + try { + connection.close(); + } catch (Error fatal) { + fatal.addSuppressed(primary); + throw fatal; + } catch (Exception closeFailure) { + primary.addSuppressed(transactionFailure( + closeFailure, TransactionPhase.CLEANUP, outcome)); + } + } + + private static void cleanupAfterFatal(Connection connection, boolean oldAutoCommit, Error fatal) { + try { + connection.rollback(); + } catch (Throwable cleanupFailure) { + fatal.addSuppressed(cleanupFailure); + } + try { + connection.setAutoCommit(oldAutoCommit); + } catch (Throwable cleanupFailure) { + fatal.addSuppressed(cleanupFailure); + } + closeAfterFatal(connection, fatal); + } + + private static void closeAfterFatal(Connection connection, Error fatal) { + if (connection == null) { + return; + } + try { + connection.close(); + } catch (Throwable closeFailure) { + fatal.addSuppressed(closeFailure); + } } @Override @@ -179,8 +314,6 @@ public CompletableFuture executeInsert(String query, Object... params) { } throw new SQLException("Insert succeeded but no generated key was returned."); } - } catch (SQLException e) { - throw new RuntimeException("Failed to execute insert", e); } }); } diff --git a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/exception/DataProviderExceptionMapper.java b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/exception/DataProviderExceptionMapper.java new file mode 100644 index 0000000..c4c75fc --- /dev/null +++ b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/exception/DataProviderExceptionMapper.java @@ -0,0 +1,353 @@ +package nl.hauntedmc.dataprovider.core.exception; + +import com.mongodb.MongoSecurityException; +import com.mongodb.MongoSocketException; +import com.mongodb.MongoTimeoutException; +import com.mongodb.MongoWriteException; +import nl.hauntedmc.dataprovider.core.concurrent.ExecutionHandle; +import nl.hauntedmc.dataprovider.core.concurrent.ExecutionRejectedException; +import nl.hauntedmc.dataprovider.database.DatabaseType; +import nl.hauntedmc.dataprovider.exception.BackendAuthenticationException; +import nl.hauntedmc.dataprovider.exception.BackendUnavailableException; +import nl.hauntedmc.dataprovider.exception.DataConflictException; +import nl.hauntedmc.dataprovider.exception.DataProviderConfigurationException; +import nl.hauntedmc.dataprovider.exception.DataProviderErrorCode; +import nl.hauntedmc.dataprovider.exception.DataProviderException; +import nl.hauntedmc.dataprovider.exception.DataProviderFailureContext; +import nl.hauntedmc.dataprovider.exception.DataProviderOperationException; +import nl.hauntedmc.dataprovider.exception.DataProviderRegistrationException; +import nl.hauntedmc.dataprovider.exception.DataProviderTimeoutException; +import nl.hauntedmc.dataprovider.exception.DataSerializationException; +import nl.hauntedmc.dataprovider.exception.DataTransactionException; +import nl.hauntedmc.dataprovider.exception.ExecutionOutcome; +import nl.hauntedmc.dataprovider.exception.ProviderClosedException; +import nl.hauntedmc.dataprovider.exception.QueueSaturatedException; +import nl.hauntedmc.dataprovider.exception.RetryAdvice; +import nl.hauntedmc.dataprovider.exception.TransactionPhase; +import org.bson.codecs.configuration.CodecConfigurationException; +import redis.clients.jedis.exceptions.JedisAccessControlException; +import redis.clients.jedis.exceptions.JedisConnectionException; + +import java.net.SocketTimeoutException; +import java.sql.SQLIntegrityConstraintViolationException; +import java.sql.SQLException; +import java.sql.SQLTimeoutException; +import java.util.LinkedHashMap; +import java.util.Locale; +import java.util.Map; +import java.util.concurrent.CompletionException; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.Executor; +import java.util.concurrent.TimeoutException; + +/** Internal classification and redaction boundary for public DataProvider failures. */ +public final class DataProviderExceptionMapper { + + private DataProviderExceptionMapper() { + } + + public static DataProviderException translate(Throwable failure, Executor executor, String operationName) { + Throwable root = unwrapAsync(failure); + if (root instanceof DataProviderException structured) { + return structured; + } + ExecutionHandle execution = executor instanceof ExecutionHandle handle ? handle : null; + DatabaseType backend = execution == null ? inferBackend(operationName) : execution.backendType(); + String connection = execution == null ? null : execution.connectionIdentifier(); + + if (root instanceof ExecutionRejectedException rejected) { + Map diagnostics = rejectionDiagnostics(rejected, execution); + return switch (rejected.reason()) { + case RUNTIME_SHUTTING_DOWN, SCOPE_CLOSED -> new ProviderClosedException( + "The DataProvider execution scope is closed.", + context(backend, connection, operationName, RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED, diagnostics), safeCause(root)); + case LANE_QUEUE_FULL, PLUGIN_QUEUE_LIMIT, CONNECTION_QUEUE_LIMIT, SUBSCRIPTION_LIMIT -> + new QueueSaturatedException( + "DataProvider execution capacity is currently exhausted.", + context(backend, connection, operationName, RetryAdvice.SAFE, + ExecutionOutcome.NOT_STARTED, diagnostics), safeCause(root)); + }; + } + if (isConflict(root)) { + return new DataConflictException( + "The operation conflicted with existing backend state.", + context(backend, connection, operationName, RetryAdvice.NEVER, + ExecutionOutcome.NOT_APPLIED, diagnosticsFor(root)), safeCause(root)); + } + if (isAuthenticationFailure(root)) { + return new BackendAuthenticationException( + "The backend rejected DataProvider authentication.", + context(backend, connection, operationName, RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED, diagnosticsFor(root)), safeCause(root)); + } + if (isTimeout(root)) { + boolean readOperation = isReadOperation(operationName); + return new DataProviderTimeoutException( + "The backend operation timed out.", + context(backend, connection, operationName, + readOperation ? RetryAdvice.SAFE : RetryAdvice.CONDITIONAL, + readOperation ? ExecutionOutcome.NOT_APPLIED : ExecutionOutcome.MAY_HAVE_APPLIED, + diagnosticsFor(root)), safeCause(root)); + } + if (isSerializationFailure(root)) { + return new DataSerializationException( + "Data serialization or deserialization failed.", + context(backend, connection, operationName, RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED, diagnosticsFor(root)), safeCause(root)); + } + if (isUnavailable(root)) { + return unavailable(backend, connection, operationName, root); + } + return new DataProviderOperationException( + "The backend operation failed.", + context(backend, connection, operationName, RetryAdvice.CONDITIONAL, + ExecutionOutcome.UNKNOWN, diagnosticsFor(root)), safeCause(root)); + } + + public static DataProviderException registrationFailure( + Throwable failure, DatabaseType backend, String connectionIdentifier, String operationName) { + Throwable root = unwrapRegistration(failure); + if (root instanceof DataProviderException structured) { + return structured; + } + if (root == null) { + return registrationException(backend, connectionIdentifier, operationName, Map.of(), null); + } + if (root instanceof MissingConfigurationFailure) { + return new DataProviderConfigurationException( + DataProviderErrorCode.CONFIGURATION_MISSING, + "No configuration exists for the requested database connection.", + context(backend, connectionIdentifier, operationName, RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED, diagnosticsFor(root)), safeCause(root)); + } + DataProviderException mapped = translate( + root, new RegistrationExecutionHandle(backend, connectionIdentifier), operationName); + if (mapped instanceof DataProviderOperationException) { + return registrationException( + backend, connectionIdentifier, operationName, + Map.of("causeCode", mapped.errorCode().name()), mapped); + } + return mapped; + } + + public static BackendUnavailableException backendDisabled(DatabaseType backend, String connectionIdentifier) { + return new BackendUnavailableException( + DataProviderErrorCode.BACKEND_DISABLED, + "The requested database backend is disabled.", + context(backend, connectionIdentifier, "registerDatabase", RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED, Map.of()), null); + } + + public static DataProviderConfigurationException configurationFailure(Throwable failure, String operationName) { + Throwable root = unwrapAsync(failure); + return new DataProviderConfigurationException( + DataProviderErrorCode.CONFIGURATION_INVALID, + "DataProvider configuration is invalid.", + context(null, null, operationName, RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED, diagnosticsFor(root)), safeCause(root)); + } + + public static DataTransactionException transactionFailure( + Throwable failure, + Executor executor, + String operationName, + TransactionPhase phase, + ExecutionOutcome outcome + ) { + DataProviderException mapped = translate(failure, executor, operationName); + return new DataTransactionException( + "The database transaction failed during " + phase.name().toLowerCase(Locale.ROOT) + ".", + phase, + context(mapped.backendType(), mapped.connectionIdentifier(), operationName, + phase == TransactionPhase.COMMIT ? RetryAdvice.CONDITIONAL : RetryAdvice.NEVER, + outcome, Map.of("phase", phase.name(), "causeCode", mapped.errorCode().name())), + mapped + ); + } + + public static MissingConfigurationFailure missingConfigurationFailure() { + return new MissingConfigurationFailure(); + } + + private static DataProviderRegistrationException registrationException( + DatabaseType backend, + String connectionIdentifier, + String operationName, + Map diagnostics, + Throwable cause + ) { + return new DataProviderRegistrationException( + "Database registration failed.", + context(backend, connectionIdentifier, operationName, RetryAdvice.CONDITIONAL, + ExecutionOutcome.NOT_STARTED, diagnostics), cause); + } + + private static BackendUnavailableException unavailable( + DatabaseType backend, String connection, String operationName, Throwable root) { + return new BackendUnavailableException( + DataProviderErrorCode.BACKEND_UNAVAILABLE, + "The configured backend is unavailable.", + context(backend, connection, operationName, RetryAdvice.CONDITIONAL, + ExecutionOutcome.UNKNOWN, diagnosticsFor(root)), safeCause(root)); + } + + private static DataProviderFailureContext context( + DatabaseType backend, + String connection, + String operation, + RetryAdvice retry, + ExecutionOutcome outcome, + Map diagnostics + ) { + return new DataProviderFailureContext(backend, connection, operation, retry, outcome, diagnostics, null); + } + + private static Throwable unwrapAsync(Throwable failure) { + if (failure == null) { + return null; + } + Throwable current = failure; + while (current.getCause() != null && (current instanceof CompletionException + || current instanceof ExecutionException + || current.getClass() == RuntimeException.class)) { + current = current.getCause(); + } + return current; + } + + private static Throwable unwrapRegistration(Throwable failure) { + Throwable current = unwrapAsync(failure); + while (current instanceof IllegalStateException && current.getCause() != null) { + current = unwrapAsync(current.getCause()); + } + return current; + } + + private static boolean isConflict(Throwable failure) { + return failure instanceof SQLIntegrityConstraintViolationException + || sqlStateStartsWith(failure, "23") + || failure instanceof MongoWriteException write && write.getError().getCode() == 11000; + } + + private static boolean isAuthenticationFailure(Throwable failure) { + if (failure instanceof MongoSecurityException || failure instanceof JedisAccessControlException) { + return true; + } + if (failure instanceof SQLException sql) { + return startsWith(sql.getSQLState(), "28") || sql.getErrorCode() == 1045; + } + String name = failure.getClass().getName(); + return name.contains("Authentication") || name.contains("AuthException"); + } + + private static boolean isTimeout(Throwable failure) { + return failure instanceof SQLTimeoutException || failure instanceof MongoTimeoutException + || failure instanceof SocketTimeoutException || failure instanceof TimeoutException + || failure.getClass().getSimpleName().contains("Timeout"); + } + + private static boolean isSerializationFailure(Throwable failure) { + String name = failure.getClass().getName(); + return failure instanceof CodecConfigurationException || name.startsWith("com.google.gson.") + || name.contains("JsonProcessingException") || name.contains("SerializationException"); + } + + private static boolean isUnavailable(Throwable failure) { + if (failure instanceof MongoSocketException || failure instanceof JedisConnectionException) { + return true; + } + if (failure instanceof SQLException sql) { + return startsWith(sql.getSQLState(), "08"); + } + String name = failure.getClass().getSimpleName(); + return name.contains("Connection") || name.contains("Socket") || name.contains("ServerSelection"); + } + + private static boolean sqlStateStartsWith(Throwable failure, String prefix) { + return failure instanceof SQLException sql && startsWith(sql.getSQLState(), prefix); + } + + private static boolean startsWith(String value, String prefix) { + return value != null && value.startsWith(prefix); + } + + private static Map rejectionDiagnostics( + ExecutionRejectedException rejection, ExecutionHandle execution) { + LinkedHashMap diagnostics = new LinkedHashMap<>(); + diagnostics.put("reason", rejection.reason().name()); + if (execution != null && execution.pluginId() != null) { + diagnostics.put("plugin", execution.pluginId()); + } + return Map.copyOf(diagnostics); + } + + private static Map diagnosticsFor(Throwable failure) { + LinkedHashMap diagnostics = new LinkedHashMap<>(); + if (failure != null) { + diagnostics.put("causeType", failure.getClass().getName()); + } + if (failure instanceof SQLException sql) { + if (sql.getSQLState() != null && !sql.getSQLState().isBlank()) { + diagnostics.put("sqlState", sql.getSQLState()); + } + diagnostics.put("vendorCode", Integer.toString(sql.getErrorCode())); + } + return Map.copyOf(diagnostics); + } + + private static Throwable safeCause(Throwable failure) { + return failure == null ? null : new SafeBackendCause(failure.getClass().getName()); + } + + private static boolean isReadOperation(String operationName) { + if (operationName == null) { + return false; + } + String lower = operationName.toLowerCase(Locale.ROOT); + return lower.contains("get") || lower.contains("find") || lower.contains("query") + || lower.contains("scan") || lower.contains("range") || lower.contains("health"); + } + + private static DatabaseType inferBackend(String operationName) { + if (operationName == null) { + return null; + } + if (operationName.startsWith("mysql.")) { + return DatabaseType.MYSQL; + } + if (operationName.startsWith("mongodb.")) { + return DatabaseType.MONGODB; + } + if (operationName.startsWith("redis.messaging.")) { + return DatabaseType.REDIS_MESSAGING; + } + if (operationName.startsWith("redis.")) { + return DatabaseType.REDIS; + } + return null; + } + + private record RegistrationExecutionHandle(DatabaseType backendType, String connectionIdentifier) + implements ExecutionHandle { + @Override public void execute(Runnable command) { command.run(); } + @Override public nl.hauntedmc.dataprovider.core.concurrent.ExecutionMetricsSnapshot metrics() { + return new nl.hauntedmc.dataprovider.core.concurrent.ExecutionMetricsSnapshot( + 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0); + } + @Override public boolean isClosed() { return false; } + @Override public void close() { } + } + + public static final class MissingConfigurationFailure extends RuntimeException { + private MissingConfigurationFailure() { + super("Missing database configuration", null, false, false); + } + } + + private static final class SafeBackendCause extends RuntimeException { + private SafeBackendCause(String causeType) { + super("Backend failure type: " + causeType, null, false, false); + } + } +} diff --git a/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/exception/StructuredFailures.java b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/exception/StructuredFailures.java new file mode 100644 index 0000000..7eb2073 --- /dev/null +++ b/dataprovider-core/src/main/java/nl/hauntedmc/dataprovider/core/exception/StructuredFailures.java @@ -0,0 +1,77 @@ +package nl.hauntedmc.dataprovider.core.exception; + +import nl.hauntedmc.dataprovider.core.concurrent.ExecutionHandle; +import nl.hauntedmc.dataprovider.database.DatabaseType; +import nl.hauntedmc.dataprovider.exception.DataProviderFailureContext; +import nl.hauntedmc.dataprovider.exception.DataSerializationException; +import nl.hauntedmc.dataprovider.exception.ExecutionOutcome; +import nl.hauntedmc.dataprovider.exception.ProviderClosedException; +import nl.hauntedmc.dataprovider.exception.RetryAdvice; + +import java.util.Map; +import java.util.concurrent.Executor; + +/** Constructors for structured failures that occur before async submission. */ +public final class StructuredFailures { + + private StructuredFailures() { + } + + public static DataSerializationException serialization( + Throwable failure, + Executor executor, + String operationName + ) { + Context context = context(executor); + return new DataSerializationException( + "Data serialization or deserialization failed.", + new DataProviderFailureContext( + context.backendType, + context.connectionIdentifier, + operationName, + RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED, + Map.of("causeType", failure.getClass().getName()), + null + ), + redactedCause(failure) + ); + } + + public static ProviderClosedException closed(Executor executor, String operationName) { + Context context = context(executor); + return new ProviderClosedException( + "The DataProvider provider is closed.", + new DataProviderFailureContext( + context.backendType, + context.connectionIdentifier, + operationName, + RetryAdvice.NEVER, + ExecutionOutcome.NOT_STARTED, + Map.of(), + null + ), + null + ); + } + + private static Context context(Executor executor) { + if (executor instanceof ExecutionHandle handle) { + return new Context(handle.backendType(), handle.connectionIdentifier()); + } + return new Context(null, null); + } + + private static Throwable redactedCause(Throwable failure) { + return new RedactedCause(failure.getClass().getName()); + } + + private record Context(DatabaseType backendType, String connectionIdentifier) { + } + + private static final class RedactedCause extends RuntimeException { + private RedactedCause(String type) { + super("Backend failure type: " + type, null, false, false); + } + } +} diff --git a/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/DatabaseFactoryTest.java b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/DatabaseFactoryTest.java index f3706ae..9a31283 100644 --- a/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/DatabaseFactoryTest.java +++ b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/DatabaseFactoryTest.java @@ -1,17 +1,16 @@ package nl.hauntedmc.dataprovider.core; -import nl.hauntedmc.dataprovider.database.DatabaseProvider; -import nl.hauntedmc.dataprovider.database.DatabaseType; import nl.hauntedmc.dataprovider.core.database.document.impl.mongodb.MongoDBDatabase; import nl.hauntedmc.dataprovider.core.database.keyvalue.impl.redis.RedisDatabase; import nl.hauntedmc.dataprovider.core.database.messaging.impl.redis.RedisMessagingDatabase; import nl.hauntedmc.dataprovider.core.database.relational.impl.mysql.MySQLDatabase; +import nl.hauntedmc.dataprovider.core.exception.DataProviderExceptionMapper; import nl.hauntedmc.dataprovider.core.testutil.RecordingLoggerAdapter; +import nl.hauntedmc.dataprovider.database.DatabaseType; import org.junit.jupiter.api.Test; import org.spongepowered.configurate.CommentedConfigurationNode; import static org.junit.jupiter.api.Assertions.assertInstanceOf; -import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; @@ -29,15 +28,14 @@ void constructorValidatesArguments() { } @Test - void returnsNullAndLogsWhenConfigurationIsMissing() { + void retainsTypedFailureAndLogsWhenConfigurationIsMissing() { RecordingLoggerAdapter logger = new RecordingLoggerAdapter(); DatabaseConfigMap configMap = mock(DatabaseConfigMap.class); when(configMap.getConfig(DatabaseType.MYSQL, ConnectionIdentifier.of("missing"))).thenReturn(null); DatabaseFactory factory = new DatabaseFactory(configMap, logger); - DatabaseProvider provider = factory.createDatabaseProvider(DatabaseType.MYSQL, "missing"); - - assertNull(provider); + assertThrows(DataProviderExceptionMapper.MissingConfigurationFailure.class, + () -> factory.createDatabaseProvider(DatabaseType.MYSQL, "missing")); assertTrue(logger.errorMessages().stream().anyMatch(m -> m.contains("Could not load configuration"))); } @@ -56,6 +54,7 @@ void createsProviderImplementationForEachDatabaseType() { assertInstanceOf(MySQLDatabase.class, factory.createDatabaseProvider(DatabaseType.MYSQL, "default")); assertInstanceOf(MongoDBDatabase.class, factory.createDatabaseProvider(DatabaseType.MONGODB, "default")); assertInstanceOf(RedisDatabase.class, factory.createDatabaseProvider(DatabaseType.REDIS, "default")); - assertInstanceOf(RedisMessagingDatabase.class, factory.createDatabaseProvider(DatabaseType.REDIS_MESSAGING, "default")); + assertInstanceOf(RedisMessagingDatabase.class, + factory.createDatabaseProvider(DatabaseType.REDIS_MESSAGING, "default")); } } diff --git a/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/HandlerLifecycleExceptionTest.java b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/HandlerLifecycleExceptionTest.java new file mode 100644 index 0000000..630dfff --- /dev/null +++ b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/HandlerLifecycleExceptionTest.java @@ -0,0 +1,41 @@ +package nl.hauntedmc.dataprovider.core; + +import nl.hauntedmc.dataprovider.core.identity.CallerContext; +import nl.hauntedmc.dataprovider.core.testutil.RecordingLoggerAdapter; +import nl.hauntedmc.dataprovider.database.DatabaseType; +import nl.hauntedmc.dataprovider.exception.ProviderClosedException; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +class HandlerLifecycleExceptionTest { + + @Test + void strictMethodsUseStructuredClosureWhileLegacyMethodsRemainCompatible() { + DataProviderRegistry registry = mock(DataProviderRegistry.class); + when(registry.isClosed()).thenReturn(true); + ClassLoader pluginLoader = new ClassLoader() { }; + DataProviderHandler handler = new DataProviderHandler( + registry, + () -> new CallerContext("plugin", pluginLoader), + new RecordingLoggerAdapter(), + getClass().getClassLoader() + ); + + assertThrows(IllegalStateException.class, + () -> handler.registerDatabase(DatabaseType.MYSQL, "default")); + assertThrows(IllegalStateException.class, + () -> handler.getRegisteredDatabase(DatabaseType.MYSQL, "default")); + + ProviderClosedException registration = assertThrows(ProviderClosedException.class, + () -> handler.registerDatabaseOrThrow(DatabaseType.MYSQL, "default")); + assertEquals("registerDatabase", registration.operationName()); + + ProviderClosedException lookup = assertThrows(ProviderClosedException.class, + () -> handler.requireRegisteredDatabase(DatabaseType.MYSQL, "default")); + assertEquals("requireRegisteredDatabase", lookup.operationName()); + } +} diff --git a/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/concurrent/AsyncTaskSupportTest.java b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/concurrent/AsyncTaskSupportTest.java index dc7166b..ad6f2c2 100644 --- a/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/concurrent/AsyncTaskSupportTest.java +++ b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/concurrent/AsyncTaskSupportTest.java @@ -1,5 +1,9 @@ package nl.hauntedmc.dataprovider.core.concurrent; +import nl.hauntedmc.dataprovider.exception.DataProviderOperationException; +import nl.hauntedmc.dataprovider.exception.ExecutionOutcome; +import nl.hauntedmc.dataprovider.exception.QueueSaturatedException; +import nl.hauntedmc.dataprovider.exception.RetryAdvice; import org.junit.jupiter.api.Test; import java.util.concurrent.CompletableFuture; @@ -8,56 +12,69 @@ import java.util.concurrent.RejectedExecutionException; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertThrows; -import static org.junit.jupiter.api.Assertions.assertTrue; class AsyncTaskSupportTest { @Test void supplyAsyncRunsOnExecutorAndReturnsValue() { Executor directExecutor = Runnable::run; - CompletableFuture result = AsyncTaskSupport.supplyAsync( - directExecutor, - "unit.supply", - () -> 42 - ); - + directExecutor, "unit.supply", () -> 42); assertEquals(42, result.join()); } @Test - void runAsyncReturnsFailedFutureWhenExecutorRejects() { + void runAsyncReturnsStructuredFailureWhenExecutorRejects() { Executor rejectingExecutor = command -> { - throw new RejectedExecutionException("full"); + throw new RejectedExecutionException("full internal queue detail"); }; + CompletableFuture future = AsyncTaskSupport.runAsync( + rejectingExecutor, "unit.reject", () -> { }); + + CompletionException completion = assertThrows(CompletionException.class, future::join); + QueueSaturatedException rejection = assertInstanceOf( + QueueSaturatedException.class, completion.getCause()); + assertEquals("unit.reject", rejection.operationName()); + assertEquals(RetryAdvice.SAFE, rejection.retryAdvice()); + assertEquals(ExecutionOutcome.NOT_STARTED, rejection.executionOutcome()); + } + @Test + void runAsyncRedactsAndStructuresUnclassifiedFailures() { + Executor directExecutor = Runnable::run; CompletableFuture future = AsyncTaskSupport.runAsync( - rejectingExecutor, - "unit.reject", + directExecutor, + "unit.failure", () -> { + throw new IllegalStateException("password=boom"); } ); - CompletionException ex = assertThrows(CompletionException.class, future::join); - assertTrue(ex.getCause() instanceof RejectedExecutionException); - assertTrue(ex.getCause().getMessage().contains("unit.reject")); + CompletionException completion = assertThrows(CompletionException.class, future::join); + DataProviderOperationException failure = assertInstanceOf( + DataProviderOperationException.class, completion.getCause()); + assertEquals("unit.failure", failure.operationName()); + assertEquals("java.lang.IllegalStateException", failure.diagnostics().get("causeType")); } @Test - void runAsyncPropagatesTaskFailures() { + void supplyAsyncPreservesCallerValidationFailures() { Executor directExecutor = Runnable::run; - CompletableFuture future = AsyncTaskSupport.runAsync( directExecutor, - "unit.failure", + "unit.validation", () -> { - throw new IllegalStateException("boom"); + throw new IllegalArgumentException("invalid identifier"); } ); - CompletionException ex = assertThrows(CompletionException.class, future::join); - assertTrue(ex.getCause() instanceof IllegalStateException); - assertEquals("boom", ex.getCause().getMessage()); + CompletionException completion = assertThrows(CompletionException.class, future::join); + IllegalArgumentException validation = assertInstanceOf( + IllegalArgumentException.class, + completion.getCause() + ); + assertEquals("invalid identifier", validation.getMessage()); } } diff --git a/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/concurrent/DataProviderExecutionRuntimeHardeningTest.java b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/concurrent/DataProviderExecutionRuntimeHardeningTest.java index 9ef674d..cd78281 100644 --- a/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/concurrent/DataProviderExecutionRuntimeHardeningTest.java +++ b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/concurrent/DataProviderExecutionRuntimeHardeningTest.java @@ -1,6 +1,7 @@ package nl.hauntedmc.dataprovider.core.concurrent; import nl.hauntedmc.dataprovider.database.DatabaseType; +import nl.hauntedmc.dataprovider.exception.QueueSaturatedException; import org.junit.jupiter.api.Test; import java.time.Duration; @@ -70,14 +71,10 @@ void exceptionalFuturePreservesStructuredRejectionReason() throws Exception { CompletableFuture rejected = AsyncTaskSupport.runAsync(scope, "rejected", () -> { }); CompletionException completion = org.junit.jupiter.api.Assertions.assertThrows( - CompletionException.class, - rejected::join - ); - ExecutionRejectedException rejection = assertInstanceOf( - ExecutionRejectedException.class, - completion.getCause() - ); - assertEquals(ExecutionRejectedException.Reason.CONNECTION_QUEUE_LIMIT, rejection.reason()); + CompletionException.class, rejected::join); + QueueSaturatedException rejection = assertInstanceOf( + QueueSaturatedException.class, completion.getCause()); + assertEquals("CONNECTION_QUEUE_LIMIT", rejection.diagnostics().get("reason")); release.countDown(); } } @@ -126,10 +123,7 @@ void interruptedWorkerDoesNotLeakInterruptFlagToNextTask() throws Exception { first.close(); CompletableFuture next = AsyncTaskSupport.supplyAsync( - second, - "check-interrupt", - () -> Thread.currentThread().isInterrupted() - ); + second, "check-interrupt", () -> Thread.currentThread().isInterrupted()); assertFalse(next.get(2, TimeUnit.SECONDS)); } } @@ -144,24 +138,12 @@ private static DataProviderExecutionRuntime runtime( long scopeGraceMs ) { ExecutionRuntimeConfig.LaneConfig lane = new ExecutionRuntimeConfig.LaneConfig( - workers, - queueCapacity, - pluginActive, - pluginQueue, - connectionActive, - connectionQueue - ); + workers, queueCapacity, pluginActive, pluginQueue, connectionActive, connectionQueue); EnumMap lanes = new EnumMap<>(ExecutionLane.class); for (ExecutionLane executionLane : ExecutionLane.values()) { lanes.put(executionLane, lane); } return new DataProviderExecutionRuntime(new ExecutionRuntimeConfig( - Map.copyOf(lanes), - Duration.ofMillis(scopeGraceMs), - Duration.ofMillis(250), - 16, - 8, - 4 - )); + Map.copyOf(lanes), Duration.ofMillis(scopeGraceMs), Duration.ofMillis(250), 16, 8, 4)); } } diff --git a/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/database/relational/impl/mysql/MySQLDataAccessTest.java b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/database/relational/impl/mysql/MySQLDataAccessTest.java index 2fd89f0..cf21163 100644 --- a/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/database/relational/impl/mysql/MySQLDataAccessTest.java +++ b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/database/relational/impl/mysql/MySQLDataAccessTest.java @@ -1,6 +1,11 @@ package nl.hauntedmc.dataprovider.core.database.relational.impl.mysql; import nl.hauntedmc.dataprovider.core.testutil.DirectExecutorService; +import nl.hauntedmc.dataprovider.exception.BackendUnavailableException; +import nl.hauntedmc.dataprovider.exception.DataProviderOperationException; +import nl.hauntedmc.dataprovider.exception.DataTransactionException; +import nl.hauntedmc.dataprovider.exception.ExecutionOutcome; +import nl.hauntedmc.dataprovider.exception.TransactionPhase; import org.junit.jupiter.api.Test; import javax.sql.DataSource; @@ -12,15 +17,18 @@ import java.sql.Statement; import java.util.List; import java.util.Map; -import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.doNothing; +import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -51,7 +59,6 @@ void queryForSingleMapsFirstRow() throws Exception { PreparedStatement statement = mock(PreparedStatement.class); ResultSet resultSet = mock(ResultSet.class); ResultSetMetaData metaData = mock(ResultSetMetaData.class); - when(dataSource.getConnection()).thenReturn(connection); when(connection.prepareStatement(anyString())).thenReturn(statement); when(statement.executeQuery()).thenReturn(resultSet); @@ -63,9 +70,8 @@ void queryForSingleMapsFirstRow() throws Exception { when(resultSet.getObject(1)).thenReturn(7); when(resultSet.getObject(2)).thenReturn("Remy"); - MySQLDataAccess access = new MySQLDataAccess(dataSource, new DirectExecutorService()); - Map row = access.queryForSingle("SELECT * FROM players WHERE id=?", 7).join(); - + Map row = new MySQLDataAccess(dataSource, new DirectExecutorService()) + .queryForSingle("SELECT * FROM players WHERE id=?", 7).join(); assertEquals(7, row.get("id")); assertEquals("Remy", row.get("name")); } @@ -80,9 +86,8 @@ void queryForSingleReturnsNullWhenNoRowsFound() throws Exception { when(connection.prepareStatement(anyString())).thenReturn(statement); when(statement.executeQuery()).thenReturn(resultSet); when(resultSet.next()).thenReturn(false); - - MySQLDataAccess access = new MySQLDataAccess(dataSource, new DirectExecutorService()); - assertNull(access.queryForSingle("SELECT 1").join()); + assertNull(new MySQLDataAccess(dataSource, new DirectExecutorService()) + .queryForSingle("SELECT 1").join()); } @Test @@ -92,7 +97,6 @@ void queryForListMapsAllRows() throws Exception { PreparedStatement statement = mock(PreparedStatement.class); ResultSet resultSet = mock(ResultSet.class); ResultSetMetaData metaData = mock(ResultSetMetaData.class); - when(dataSource.getConnection()).thenReturn(connection); when(connection.prepareStatement(anyString())).thenReturn(statement); when(statement.executeQuery()).thenReturn(resultSet); @@ -104,11 +108,9 @@ void queryForListMapsAllRows() throws Exception { when(resultSet.getObject(1)).thenReturn(1, 2); when(resultSet.getObject(2)).thenReturn("a", "b"); - MySQLDataAccess access = new MySQLDataAccess(dataSource, new DirectExecutorService()); - List> rows = access.queryForList("SELECT * FROM players").join(); - + List> rows = new MySQLDataAccess(dataSource, new DirectExecutorService()) + .queryForList("SELECT * FROM players").join(); assertEquals(2, rows.size()); - assertEquals(1, rows.get(0).get("id")); assertEquals("b", rows.get(1).get("name")); } @@ -118,7 +120,6 @@ void queryForSingleValueReturnsFirstColumnAndHandlesEmptyResult() throws Excepti Connection connection = mock(Connection.class); PreparedStatement statement = mock(PreparedStatement.class); ResultSet resultSet = mock(ResultSet.class); - when(dataSource.getConnection()).thenReturn(connection); when(connection.prepareStatement(anyString())).thenReturn(statement); when(statement.executeQuery()).thenReturn(resultSet); @@ -126,13 +127,9 @@ void queryForSingleValueReturnsFirstColumnAndHandlesEmptyResult() throws Excepti when(resultSet.getObject(1)).thenReturn("value"); MySQLDataAccess access = new MySQLDataAccess(dataSource, new DirectExecutorService()); - - Object value = access.queryForSingleValue("SELECT value FROM test").join(); - assertEquals("value", value); - + assertEquals("value", access.queryForSingleValue("SELECT value FROM test").join()); when(resultSet.next()).thenReturn(false); - Object missing = access.queryForSingleValue("SELECT value FROM test WHERE id=999").join(); - assertNull(missing); + assertNull(access.queryForSingleValue("SELECT value FROM test WHERE id=999").join()); } @Test @@ -143,8 +140,7 @@ void executeBatchUpdateAddsEachBatchEntry() throws Exception { when(dataSource.getConnection()).thenReturn(connection); when(connection.prepareStatement(anyString())).thenReturn(statement); - MySQLDataAccess access = new MySQLDataAccess(dataSource, new DirectExecutorService()); - access.executeBatchUpdate( + new MySQLDataAccess(dataSource, new DirectExecutorService()).executeBatchUpdate( "INSERT INTO test(a,b) VALUES (?,?)", List.of(new Object[]{1, "x"}, new Object[]{2, "y"}) ).join(); @@ -157,39 +153,117 @@ void executeBatchUpdateAddsEachBatchEntry() throws Exception { } @Test - void executeTransactionallyCommitsOnSuccessAndRollsBackOnFailure() throws Exception { + void executeTransactionallyCommitsRestoresAndClosesConnection() throws Exception { DataSource dataSource = mock(DataSource.class); Connection connection = mock(Connection.class); when(dataSource.getConnection()).thenReturn(connection); when(connection.getAutoCommit()).thenReturn(true); - MySQLDataAccess access = new MySQLDataAccess(dataSource, new DirectExecutorService()); - String result = access.executeTransactionally(conn -> "done").join(); + String result = new MySQLDataAccess(dataSource, new DirectExecutorService()) + .executeTransactionally(conn -> "done").join(); assertEquals("done", result); verify(connection).setAutoCommit(false); verify(connection).commit(); verify(connection).setAutoCommit(true); + verify(connection).close(); + } - DataSource failingDataSource = mock(DataSource.class); - Connection failingConnection = mock(Connection.class); - when(failingDataSource.getConnection()).thenReturn(failingConnection); - when(failingConnection.getAutoCommit()).thenReturn(false); - MySQLDataAccess failingAccess = new MySQLDataAccess(failingDataSource, new DirectExecutorService()); - - CompletionException ex = assertThrows( - CompletionException.class, - () -> failingAccess.executeTransactionally(conn -> { - throw new IllegalStateException("boom"); - }).join() - ); - - assertInstanceOf(RuntimeException.class, ex.getCause()); - verify(failingConnection).rollback(); - verify(failingConnection, times(2)).setAutoCommit(false); + @Test + void failedTransactionSetupClosesAcquiredConnection() throws Exception { + DataSource dataSource = mock(DataSource.class); + Connection connection = mock(Connection.class); + when(dataSource.getConnection()).thenReturn(connection); + when(connection.getAutoCommit()).thenReturn(true); + doThrow(new SQLException("setup secret")).when(connection).setAutoCommit(false); + + CompletionException completion = assertThrows(CompletionException.class, + () -> new MySQLDataAccess(dataSource, new DirectExecutorService()) + .executeTransactionally(conn -> "unused").join()); + DataTransactionException transaction = assertInstanceOf( + DataTransactionException.class, completion.getCause()); + assertEquals(TransactionPhase.BEGIN, transaction.phase()); + assertEquals(ExecutionOutcome.NOT_STARTED, transaction.executionOutcome()); + verify(connection).close(); } @Test - void executeInsertReturnsGeneratedKeyAndFailsWhenNoneReturned() throws Exception { + void callbackRollbackAndCloseFailuresRemainStructuredAndRedacted() throws Exception { + DataSource dataSource = mock(DataSource.class); + Connection connection = mock(Connection.class); + when(dataSource.getConnection()).thenReturn(connection); + when(connection.getAutoCommit()).thenReturn(false); + doThrow(new SQLException("rollback secret")).when(connection).rollback(); + doThrow(new SQLException("close secret")).when(connection).close(); + + CompletionException completion = assertThrows(CompletionException.class, + () -> new MySQLDataAccess(dataSource, new DirectExecutorService()) + .executeTransactionally(conn -> { + throw new IllegalStateException("callback secret"); + }).join()); + DataTransactionException transaction = assertInstanceOf( + DataTransactionException.class, completion.getCause()); + assertEquals(TransactionPhase.CALLBACK, transaction.phase()); + assertEquals(ExecutionOutcome.NOT_APPLIED, transaction.executionOutcome()); + assertEquals(2, transaction.getSuppressed().length); + for (Throwable suppressed : transaction.getSuppressed()) { + DataTransactionException structured = assertInstanceOf(DataTransactionException.class, suppressed); + assertFalse(structured.getMessage().contains("secret")); + assertFalse(structured.getCause().getMessage().contains("secret")); + } + assertEquals(TransactionPhase.ROLLBACK, + ((DataTransactionException) transaction.getSuppressed()[0]).phase()); + assertEquals(TransactionPhase.CLEANUP, + ((DataTransactionException) transaction.getSuppressed()[1]).phase()); + verify(connection).rollback(); + verify(connection, times(2)).setAutoCommit(false); + verify(connection).close(); + } + + @Test + void commitFailureReportsUnknownWriteOutcome() throws Exception { + DataSource dataSource = mock(DataSource.class); + Connection connection = mock(Connection.class); + when(dataSource.getConnection()).thenReturn(connection); + when(connection.getAutoCommit()).thenReturn(true); + doThrow(new SQLException("commit lost", "08006")).when(connection).commit(); + + CompletionException completion = assertThrows(CompletionException.class, + () -> new MySQLDataAccess(dataSource, new DirectExecutorService()) + .executeTransactionally(conn -> "done").join()); + DataTransactionException transaction = assertInstanceOf( + DataTransactionException.class, completion.getCause()); + assertEquals(TransactionPhase.COMMIT, transaction.phase()); + assertEquals(ExecutionOutcome.MAY_HAVE_APPLIED, transaction.executionOutcome()); + assertTrue(transaction.retryable()); + verify(connection).rollback(); + verify(connection).setAutoCommit(true); + verify(connection).close(); + } + + @Test + void postCommitRestoreFailureUsesCleanupPhaseAndAppliedOutcome() throws Exception { + DataSource dataSource = mock(DataSource.class); + Connection connection = mock(Connection.class); + when(dataSource.getConnection()).thenReturn(connection); + when(connection.getAutoCommit()).thenReturn(true); + doNothing().doThrow(new SQLException("restore secret")) + .when(connection).setAutoCommit(anyBoolean()); + + CompletionException completion = assertThrows(CompletionException.class, + () -> new MySQLDataAccess(dataSource, new DirectExecutorService()) + .executeTransactionally(conn -> "done").join()); + DataTransactionException transaction = assertInstanceOf( + DataTransactionException.class, completion.getCause()); + assertEquals(TransactionPhase.CLEANUP, transaction.phase()); + assertEquals(ExecutionOutcome.MAY_HAVE_APPLIED, transaction.executionOutcome()); + assertFalse(transaction.retryable()); + assertFalse(transaction.getCause().getMessage().contains("secret")); + verify(connection).commit(); + verify(connection).close(); + } + + @Test + void executeInsertReturnsGeneratedKeyAndStructuresFailure() throws Exception { DataSource dataSource = mock(DataSource.class); Connection connection = mock(Connection.class); PreparedStatement statement = mock(PreparedStatement.class); @@ -203,28 +277,25 @@ void executeInsertReturnsGeneratedKeyAndFailsWhenNoneReturned() throws Exception when(generatedKeys.getObject(1)).thenReturn(42L); MySQLDataAccess access = new MySQLDataAccess(dataSource, new DirectExecutorService()); - Object key = access.executeInsert("INSERT INTO players(name) VALUES (?)", "test").join(); - assertEquals(42L, key); + assertEquals(42L, access.executeInsert("INSERT INTO players(name) VALUES (?)", "test").join()); when(statement.executeUpdate()).thenReturn(0); - CompletionException ex = assertThrows( - CompletionException.class, - () -> access.executeInsert("INSERT INTO players(name) VALUES (?)", "test").join() - ); - assertTrue(ex.getCause().getMessage().contains("Failed to execute insert")); + CompletionException completion = assertThrows(CompletionException.class, + () -> access.executeInsert("INSERT INTO players(name) VALUES (?)", "test").join()); + assertInstanceOf(DataProviderOperationException.class, completion.getCause()); } @Test - void wrapsSqlExceptionsInRuntimeExceptions() throws Exception { + void connectionSqlExceptionsCompleteWithUnavailableFailure() throws Exception { DataSource dataSource = mock(DataSource.class); - when(dataSource.getConnection()).thenThrow(new SQLException("no connection")); - MySQLDataAccess access = new MySQLDataAccess(dataSource, new DirectExecutorService()); - - CompletionException ex = assertThrows( - CompletionException.class, - () -> access.executeUpdate("UPDATE test SET value=1").join() - ); - assertInstanceOf(RuntimeException.class, ex.getCause()); - assertTrue(ex.getCause().getMessage().contains("Failed to execute update")); + when(dataSource.getConnection()).thenThrow(new SQLException("password=secret", "08001")); + + CompletionException completion = assertThrows(CompletionException.class, + () -> new MySQLDataAccess(dataSource, new DirectExecutorService()) + .executeUpdate("UPDATE test SET value=1").join()); + BackendUnavailableException failure = assertInstanceOf( + BackendUnavailableException.class, completion.getCause()); + assertEquals("08001", failure.diagnostics().get("sqlState")); + assertFalse(failure.getMessage().contains("secret")); } } diff --git a/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/exception/DataProviderExceptionMapperTest.java b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/exception/DataProviderExceptionMapperTest.java new file mode 100644 index 0000000..ddb6959 --- /dev/null +++ b/dataprovider-core/src/test/java/nl/hauntedmc/dataprovider/core/exception/DataProviderExceptionMapperTest.java @@ -0,0 +1,131 @@ +package nl.hauntedmc.dataprovider.core.exception; + +import nl.hauntedmc.dataprovider.core.concurrent.ContextualExecutionHandle; +import nl.hauntedmc.dataprovider.core.concurrent.ExecutionHandle; +import nl.hauntedmc.dataprovider.core.concurrent.ExecutionRejectedException; +import nl.hauntedmc.dataprovider.database.DatabaseType; +import nl.hauntedmc.dataprovider.exception.BackendAuthenticationException; +import nl.hauntedmc.dataprovider.exception.DataConflictException; +import nl.hauntedmc.dataprovider.exception.DataProviderException; +import nl.hauntedmc.dataprovider.exception.DataProviderOperationException; +import nl.hauntedmc.dataprovider.exception.DataProviderTimeoutException; +import nl.hauntedmc.dataprovider.exception.ExecutionOutcome; +import nl.hauntedmc.dataprovider.exception.ProviderClosedException; +import nl.hauntedmc.dataprovider.exception.QueueSaturatedException; +import nl.hauntedmc.dataprovider.exception.RetryAdvice; +import org.junit.jupiter.api.Test; + +import java.sql.SQLIntegrityConstraintViolationException; +import java.sql.SQLException; +import java.sql.SQLTimeoutException; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNotNull; + +class DataProviderExceptionMapperTest { + + private final ExecutionHandle execution = new ContextualExecutionHandle( + ExecutionHandle.direct(), "test-plugin", DatabaseType.MYSQL, "main"); + + @Test + void queueRejectionIsSafeAndRetryable() { + DataProviderException mapped = DataProviderExceptionMapper.translate( + new ExecutionRejectedException( + ExecutionRejectedException.Reason.PLUGIN_QUEUE_LIMIT, + "internal detail"), + execution, + "mysql.queryForList" + ); + + QueueSaturatedException saturated = assertInstanceOf(QueueSaturatedException.class, mapped); + assertEquals(DatabaseType.MYSQL, saturated.backendType()); + assertEquals("main", saturated.connectionIdentifier()); + assertEquals("mysql.queryForList", saturated.operationName()); + assertEquals(RetryAdvice.SAFE, saturated.retryAdvice()); + assertEquals(ExecutionOutcome.NOT_STARTED, saturated.executionOutcome()); + assertEquals("PLUGIN_QUEUE_LIMIT", saturated.diagnostics().get("reason")); + assertNotNull(saturated.diagnosticId()); + assertFalse(saturated.getMessage().contains("internal detail")); + } + + @Test + void closureMapsToNonRetryableProviderFailure() { + ProviderClosedException closed = assertInstanceOf( + ProviderClosedException.class, + DataProviderExceptionMapper.translate( + new ExecutionRejectedException( + ExecutionRejectedException.Reason.SCOPE_CLOSED, + "scope detail"), + execution, + "mysql.executeUpdate") + ); + assertEquals(RetryAdvice.NEVER, closed.retryAdvice()); + assertEquals(ExecutionOutcome.NOT_STARTED, closed.executionOutcome()); + } + + @Test + void sqlFailuresAreClassifiedWithoutLeakingMessages() { + SQLException authentication = new SQLException( + "password=top-secret jdbc:mysql://secret-host", "28000", 1045); + BackendAuthenticationException auth = assertInstanceOf( + BackendAuthenticationException.class, + DataProviderExceptionMapper.translate(authentication, execution, "mysql.connect") + ); + assertFalse(auth.getMessage().contains("top-secret")); + assertFalse(auth.getCause().getMessage().contains("top-secret")); + assertEquals("28000", auth.diagnostics().get("sqlState"), + "Authentication diagnostics should retain safe SQL state when available."); + + DataConflictException conflict = assertInstanceOf( + DataConflictException.class, + DataProviderExceptionMapper.translate( + new SQLIntegrityConstraintViolationException("duplicate secret", "23000", 1062), + execution, + "mysql.executeInsert") + ); + assertEquals(ExecutionOutcome.NOT_APPLIED, conflict.executionOutcome()); + + DataProviderTimeoutException writeTimeout = assertInstanceOf( + DataProviderTimeoutException.class, + DataProviderExceptionMapper.translate( + new SQLTimeoutException("write timed out", "HYT00"), + execution, + "mysql.executeUpdate") + ); + assertEquals(ExecutionOutcome.MAY_HAVE_APPLIED, writeTimeout.executionOutcome()); + assertEquals(RetryAdvice.CONDITIONAL, writeTimeout.retryAdvice()); + } + + @Test + void readTimeoutIsSafeToRetryAndCannotHaveAppliedData() { + DataProviderTimeoutException timeout = assertInstanceOf( + DataProviderTimeoutException.class, + DataProviderExceptionMapper.translate( + new SQLTimeoutException("read timed out", "HYT00"), + execution, + "mysql.queryForList") + ); + + assertEquals(RetryAdvice.SAFE, timeout.retryAdvice()); + assertEquals(ExecutionOutcome.NOT_APPLIED, timeout.executionOutcome()); + } + + @Test + void unclassifiedDriverFailureUsesGenericOperationCategory() { + DataProviderOperationException failure = assertInstanceOf( + DataProviderOperationException.class, + DataProviderExceptionMapper.translate( + new SQLException("syntax near password=secret", "42000", 1064), + execution, + "mysql.executeUpdate") + ); + + assertEquals(RetryAdvice.CONDITIONAL, failure.retryAdvice()); + assertEquals(ExecutionOutcome.UNKNOWN, failure.executionOutcome()); + assertEquals("42000", failure.diagnostics().get("sqlState")); + assertFalse(failure.getMessage().contains("secret")); + assertFalse(failure.getCause().getMessage().contains("secret")); + } +} diff --git a/docs/EXCEPTIONS.md b/docs/EXCEPTIONS.md new file mode 100644 index 0000000..85db394 --- /dev/null +++ b/docs/EXCEPTIONS.md @@ -0,0 +1,57 @@ +# Structured exceptions + +DataProvider exposes unchecked structured failures from `nl.hauntedmc.dataprovider.exception`. + +## Strict and compatibility APIs + +Use `registerDatabaseOrThrow(...)` when startup must distinguish missing configuration, disabled backends, authentication failure, timeout, or backend unavailability. Use `requireRegisteredDatabase(...)` when absence is exceptional. + +The original `registerDatabase(...)` and optional helpers remain available for compatibility. They continue returning `null` or `Optional.empty()` and intentionally discard the failure category after DataProvider records registration failures internally. + +Legacy methods retain their previous lifecycle behavior, including `IllegalStateException` after their API or scope has closed. The new strict registration and lookup methods report closure as `ProviderClosedException`. + +Caller input validation remains distinct from backend failure classification. Invalid identifiers, unsupported document values, null arguments, and similar programming errors continue to use standard validation exceptions such as `IllegalArgumentException` and `NullPointerException`. + +## Common handling + +```java +try { + DatabaseProvider provider = api.registerDatabaseOrThrow(DatabaseType.MYSQL, "main"); +} catch (BackendAuthenticationException exception) { + // Configuration intervention is required; retrying unchanged credentials is not useful. +} catch (BackendUnavailableException exception) { + // The backend is disabled or unreachable. +} catch (DataProviderOperationException exception) { + // The operation failed without matching a more specific public category. +} +``` + +All structured exceptions expose: + +- `errorCode()` — stable machine-readable category +- `backendType()` — backend involved, when applicable +- `connectionIdentifier()` — safe logical identifier, never a connection URL +- `operationName()` — stable operation identifier +- `retryAdvice()` — `NEVER`, `SAFE`, or `CONDITIONAL` +- `executionOutcome()` — whether the operation started or may already have applied +- `diagnostics()` — immutable allowlisted metadata +- `diagnosticId()` — validated correlation identifier for operational support + +## Retry safety + +`retryable()` is a convenience method. Prefer `retryAdvice()` and `executionOutcome()` for writes: + +- `SAFE` + `NOT_STARTED` or `NOT_APPLIED`: retrying is normally safe. +- `CONDITIONAL` + `MAY_HAVE_APPLIED`: do not retry blindly; use an idempotency key or verify backend state. +- `CONDITIONAL` + `UNKNOWN`: retry only when the operation is idempotent or backend state has been checked. +- `NEVER`: correct configuration, ownership, authentication, lifecycle, or cleanup state first. + +Read timeouts are reported as `SAFE` with outcome `NOT_APPLIED`. Write timeouts remain `CONDITIONAL` with outcome `MAY_HAVE_APPLIED`. + +A transaction commit timeout can mean the commit succeeded but its acknowledgement was lost. DataProvider reports this as `DataTransactionException` with phase `COMMIT` and outcome `MAY_HAVE_APPLIED`. Failures while restoring or closing a connection after a successful commit use phase `CLEANUP`, outcome `MAY_HAVE_APPLIED`, and are not retryable. + +## Redaction + +DataProvider-generated public exception messages, diagnostics, causes, and suppressed cleanup failures do not include passwords, tokens, payloads, query parameter values, raw configuration, or credential-bearing URLs. Public causes preserve the original failure type through a redacted surrogate. Registration lifecycle failures retain their original internal cause for diagnostics while exposing only redacted public metadata. + +Rollback and cleanup failures are attached as suppressed structured exceptions without replacing the primary transaction failure. JVM-fatal errors remain primary and are never converted into ordinary DataProvider failures.