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