handleError(
+ CR resource,
+ S status,
+ E exception
+ ) {
+ status.setPhase(CRPhase.ERROR)
+ .setMessage(exception.getMessage());
+
+ return UpdateControl.patchStatus(resource)
+ .rescheduleAfter(60, TimeUnit.SECONDS);
+ }
+}
diff --git a/src/main/java/it/aboutbits/postgresql/core/CRPhase.java b/src/main/java/it/aboutbits/postgresql/core/CRPhase.java
new file mode 100644
index 0000000..4c0e0df
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/core/CRPhase.java
@@ -0,0 +1,11 @@
+package it.aboutbits.postgresql.core;
+
+import org.jspecify.annotations.NullMarked;
+
+@NullMarked
+public enum CRPhase {
+ PENDING,
+ READY,
+ ERROR,
+ DELETING
+}
diff --git a/src/main/java/it/aboutbits/postgresql/core/CRStatus.java b/src/main/java/it/aboutbits/postgresql/core/CRStatus.java
new file mode 100644
index 0000000..db7ac62
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/core/CRStatus.java
@@ -0,0 +1,76 @@
+package it.aboutbits.postgresql.core;
+
+import lombok.AccessLevel;
+import lombok.Getter;
+import lombok.Setter;
+import lombok.experimental.Accessors;
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+
+import java.time.OffsetDateTime;
+import java.time.ZoneOffset;
+
+/**
+ * Status Object for the Custom Resources.
+ *
+ * This object captures the current state of a Custom Resource as observed by the reconciler.
+ */
+@NullMarked
+@Getter
+@Setter
+@Accessors(chain = true)
+public class CRStatus {
+ /**
+ * The Custom Resource name (may differ from metadata.name).
+ */
+ @Nullable
+ private String name = null;
+
+ /**
+ * Current lifecycle phase of the Bucket.
+ */
+ @Setter(AccessLevel.NONE)
+ private CRPhase phase = CRPhase.PENDING;
+
+ /**
+ * Human-readable message providing details about the current state.
+ */
+ @Nullable
+ private String message = null;
+
+ /**
+ * Last time the condition was probed/updated.
+ */
+ @Nullable
+ private OffsetDateTime lastProbeTime = null;
+
+ /**
+ * Last time the condition transitioned from one status to another.
+ */
+ @Nullable
+ @Setter(AccessLevel.NONE)
+ private OffsetDateTime lastPhaseTransitionTime = null;
+
+ /**
+ * Observed resource generation that the controller acted upon.
+ */
+ private long observedGeneration = 0;
+
+ /**
+ * Update the current phase. When the phase changes, the {@link #lastPhaseTransitionTime}
+ * is updated to the current UTC time and the message is set to {@code null}.
+ *
+ * @param newPhase the new phase
+ * @return this status instance
+ */
+ public CRStatus setPhase(CRPhase newPhase) {
+ if (this.phase == newPhase) {
+ return this;
+ }
+
+ this.phase = newPhase;
+ this.lastPhaseTransitionTime = OffsetDateTime.now(ZoneOffset.UTC);
+
+ return this;
+ }
+}
diff --git a/src/main/java/it/aboutbits/postgresql/core/ClusterReference.java b/src/main/java/it/aboutbits/postgresql/core/ClusterReference.java
new file mode 100644
index 0000000..262ace8
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/core/ClusterReference.java
@@ -0,0 +1,19 @@
+package it.aboutbits.postgresql.core;
+
+import io.fabric8.generator.annotation.Required;
+import lombok.Getter;
+import lombok.Setter;
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+
+@NullMarked
+@Getter
+@Setter
+public class ClusterReference {
+ @Required
+ private String name = "";
+
+ @Nullable
+ @io.fabric8.generator.annotation.Nullable
+ private String namespace;
+}
diff --git a/src/main/java/it/aboutbits/postgresql/core/Credentials.java b/src/main/java/it/aboutbits/postgresql/core/Credentials.java
new file mode 100644
index 0000000..dc376e4
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/core/Credentials.java
@@ -0,0 +1,11 @@
+package it.aboutbits.postgresql.core;
+
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+
+@NullMarked
+public record Credentials(
+ @Nullable String username,
+ String password
+) {
+}
diff --git a/src/main/java/it/aboutbits/postgresql/core/KubernetesService.java b/src/main/java/it/aboutbits/postgresql/core/KubernetesService.java
new file mode 100644
index 0000000..2252b1a
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/core/KubernetesService.java
@@ -0,0 +1,94 @@
+package it.aboutbits.postgresql.core;
+
+import io.fabric8.kubernetes.client.KubernetesClient;
+import it.aboutbits.postgresql.crd.clusterconnection.ClusterConnection;
+import jakarta.inject.Singleton;
+import org.jspecify.annotations.NullMarked;
+
+import java.nio.charset.Charset;
+import java.util.Base64;
+
+@NullMarked
+@Singleton
+public final class KubernetesService {
+ public static final String SECRET_TYPE_BASIC_AUTH = "kubernetes.io/basic-auth";
+ public static final String SECRET_DATA_BASIC_AUTH_USERNAME_KEY = "username";
+ public static final String SECRET_DATA_BASIC_AUTH_PASSWORD_KEY = "password";
+
+ public Credentials getSecretRefCredentials(
+ KubernetesClient kubernetesClient,
+ ClusterConnection clusterConnection
+ ) {
+ return getSecretRefCredentials(
+ kubernetesClient,
+ clusterConnection.getSpec().getAdminSecretRef(),
+ clusterConnection.getMetadata().getNamespace()
+ );
+ }
+
+ public Credentials getSecretRefCredentials(
+ KubernetesClient kubernetesClient,
+ SecretRef secretRef,
+ String defaultNamespace
+ ) {
+ var secretNamespace = secretRef.getNamespace() != null
+ ? secretRef.getNamespace()
+ : defaultNamespace;
+
+ var secretName = secretRef.getName();
+
+ var secret = kubernetesClient.secrets()
+ .inNamespace(secretNamespace)
+ .withName(secretName)
+ .get();
+
+ if (secret == null) {
+ throw new IllegalStateException("SecretRef not found [secret.namespace=%s, secret.name=%s]".formatted(
+ secretNamespace,
+ secretName
+ ));
+ }
+
+ if (!secret.getType().equals(SECRET_TYPE_BASIC_AUTH)) {
+ throw new IllegalArgumentException("The SecretRef is of the wrong type [secret.namespace=%s, secret.name=%s, expected.secret.type=%s, actual.secret.type=%s]".formatted(
+ secretNamespace,
+ secretName,
+ SECRET_TYPE_BASIC_AUTH,
+ secret.getType()
+ ));
+ }
+
+ var data = secret.getData();
+ if (data == null || data.isEmpty()) {
+ throw new IllegalStateException("The SecretRef has no data set [secret.namespace=%s, secret.name=%s]".formatted(
+ secretNamespace,
+ secretName
+ ));
+ }
+
+ var usernameBase64 = data.get(SECRET_DATA_BASIC_AUTH_USERNAME_KEY);
+ var username = usernameBase64 == null
+ ? null
+ : new String(
+ Base64.getDecoder().decode(usernameBase64),
+ Charset.defaultCharset()
+ );
+
+ var passwordBase64 = data.get(SECRET_DATA_BASIC_AUTH_PASSWORD_KEY);
+ if (passwordBase64 == null) {
+ throw new IllegalStateException("The SecretRef is missing required data password [secret.namespace=%s, secret.name=%s]".formatted(
+ secretNamespace,
+ secretName
+ ));
+ }
+ var password = new String(
+ Base64.getDecoder().decode(passwordBase64),
+ Charset.defaultCharset()
+ );
+
+ return new Credentials(
+ username,
+ password
+ );
+ }
+}
diff --git a/src/main/java/it/aboutbits/postgresql/core/Named.java b/src/main/java/it/aboutbits/postgresql/core/Named.java
new file mode 100644
index 0000000..1d90fe2
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/core/Named.java
@@ -0,0 +1,10 @@
+package it.aboutbits.postgresql.core;
+
+import com.fasterxml.jackson.annotation.JsonIgnore;
+import org.jspecify.annotations.NullMarked;
+
+@NullMarked
+public interface Named {
+ @JsonIgnore
+ String getName();
+}
diff --git a/src/main/java/it/aboutbits/postgresql/core/PostgreSQLAuthenticationService.java b/src/main/java/it/aboutbits/postgresql/core/PostgreSQLAuthenticationService.java
new file mode 100644
index 0000000..4f6c849
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/core/PostgreSQLAuthenticationService.java
@@ -0,0 +1,219 @@
+package it.aboutbits.postgresql.core;
+
+import com.ongres.scram.common.StringPreparation;
+import it.aboutbits.postgresql.crd.role.RoleSpec;
+import jakarta.inject.Singleton;
+import lombok.extern.slf4j.Slf4j;
+import org.jooq.DSLContext;
+import org.jspecify.annotations.NullMarked;
+
+import javax.crypto.Mac;
+import javax.crypto.SecretKeyFactory;
+import javax.crypto.spec.PBEKeySpec;
+import javax.crypto.spec.SecretKeySpec;
+import java.nio.charset.StandardCharsets;
+import java.security.MessageDigest;
+import java.security.NoSuchAlgorithmException;
+import java.util.Arrays;
+import java.util.Base64;
+import java.util.HexFormat;
+import java.util.Locale;
+
+import static it.aboutbits.postgresql.core.infrastructure.persistence.Tables.PG_AUTHID;
+
+@NullMarked
+@Slf4j
+@Singleton
+public final class PostgreSQLAuthenticationService {
+ private static final String MD5 = "MD5";
+ private static final String SHA_256 = "SHA-256";
+ private static final String HMAC_SHA_256 = "HmacSHA256";
+ private static final String PBKDF2_WITH_HMAC_SHA256 = "PBKDF2WithHmacSHA256";
+
+ public boolean passwordMatches(
+ DSLContext dsl,
+ RoleSpec spec,
+ String expectedPassword
+ ) {
+ var currentPasswordVerifier = dsl
+ .select(PG_AUTHID.ROLPASSWORD)
+ .from(PG_AUTHID)
+ .where(PG_AUTHID.ROLNAME.eq(spec.getName()))
+ .fetchSingle(PG_AUTHID.ROLPASSWORD);
+
+ if (currentPasswordVerifier == null || currentPasswordVerifier.isBlank()) {
+ return false;
+ }
+
+ // PostgreSQL stores either:
+ // - SCRAM verifier: SCRAM-SHA-256$:$:
+ // - or legacy md5: md5
+ if (currentPasswordVerifier.startsWith("SCRAM-SHA-256$")) {
+ return verifyPostgresScramSha256(
+ currentPasswordVerifier,
+ expectedPassword
+ );
+ }
+
+ if (currentPasswordVerifier.startsWith(MD5.toLowerCase(Locale.ROOT))) {
+ return verifyPostgresMd5(
+ currentPasswordVerifier,
+ expectedPassword,
+ spec.getName()
+ );
+ }
+
+ // Unknown format (or plain text, which PG should not store in rolpassword)
+ return false;
+ }
+
+ private static boolean verifyPostgresScramSha256(String postgresVerifier, String cleartextPassword) {
+ // Prepare the cleartext password with SASLprep
+ var preparedPassword = StringPreparation.POSTGRESQL_PREPARATION.normalize(
+ cleartextPassword.toCharArray()
+ );
+
+ // Format: SCRAM-SHA-256$:$:
+ var afterPrefix = postgresVerifier.substring("SCRAM-SHA-256$".length());
+ var dollar = afterPrefix.indexOf('$');
+ if (dollar < 0) {
+ return false;
+ }
+
+ // :
+ var iterationsAndSalt = afterPrefix.substring(0, dollar);
+ // :
+ var keys = afterPrefix.substring(dollar + 1);
+
+ var colonIterationsAndSalt = iterationsAndSalt.indexOf(':');
+ if (colonIterationsAndSalt < 0) {
+ return false;
+ }
+
+ int iterations;
+ try {
+ iterations = Integer.parseInt(iterationsAndSalt.substring(0, colonIterationsAndSalt));
+ } catch (NumberFormatException e) {
+ log.error("Invalid iterations format in PostgreSQL verifier: %s".formatted(postgresVerifier), e);
+ return false;
+ }
+ if (iterations <= 0) {
+ return false;
+ }
+
+ var saltB64 = iterationsAndSalt.substring(colonIterationsAndSalt + 1);
+
+ var colonKeys = keys.indexOf(':');
+ if (colonKeys < 0) {
+ return false;
+ }
+
+ var storedKeyB64 = keys.substring(0, colonKeys);
+
+ byte[] salt;
+ byte[] currentStoredKey;
+ try {
+ salt = Base64.getDecoder().decode(saltB64);
+ currentStoredKey = Base64.getDecoder().decode(storedKeyB64);
+ } catch (IllegalArgumentException e) {
+ log.error("Invalid salt or stored key format in PostgreSQL verifier: %s".formatted(postgresVerifier), e);
+ return false;
+ }
+
+ byte[] saltedPassword = null;
+ byte[] clientKey = null;
+ byte[] expectedStoredKey = null;
+ try {
+ // RFC 5802/7677:
+ // saltedPassword := Hi(password, salt, iterations) (PBKDF2-HMAC-SHA-256, 32 bytes)
+ // clientKey := HMAC(saltedPassword, "Client Key")
+ // storedKey := H(clientKey) (SHA-256)
+ saltedPassword = pbkdf2HmacSha256(preparedPassword, salt, iterations, 32);
+ clientKey = hmacSha256(saltedPassword, "Client Key".getBytes(StandardCharsets.UTF_8));
+ expectedStoredKey = sha256(clientKey);
+
+ return MessageDigest.isEqual(
+ currentStoredKey,
+ expectedStoredKey
+ );
+ } finally {
+ if (saltedPassword != null) {
+ Arrays.fill(saltedPassword, (byte) 0);
+ }
+ if (clientKey != null) {
+ Arrays.fill(clientKey, (byte) 0);
+ }
+ if (expectedStoredKey != null) {
+ Arrays.fill(expectedStoredKey, (byte) 0);
+ }
+ }
+ }
+
+ private static boolean verifyPostgresMd5(
+ String postgresMd5,
+ String expectedPassword,
+ String username
+ ) {
+ // PostgreSQL md5 is: "md5" + md5(password + username)
+ if (postgresMd5.length() != 3 + 32 || !postgresMd5.regionMatches(true, 0, MD5, 0, 3)) {
+ return false;
+ }
+
+ byte[] currentDigest;
+ try {
+ currentDigest = HexFormat.of().parseHex(
+ postgresMd5,
+ 3,
+ postgresMd5.length()
+ );
+ } catch (IllegalArgumentException e) {
+ log.error("Invalid MD5 format in PostgreSQL verifier: %s".formatted(postgresMd5), e);
+ return false; // not valid hex
+ }
+
+ MessageDigest md5;
+ try {
+ md5 = MessageDigest.getInstance(MD5);
+ } catch (NoSuchAlgorithmException e) {
+ throw new IllegalStateException("%s not available".formatted(MD5), e);
+ }
+
+ md5.update((expectedPassword + username).getBytes(StandardCharsets.UTF_8));
+ var expectedDigest = md5.digest();
+
+ return MessageDigest.isEqual(currentDigest, expectedDigest);
+ }
+
+ private static byte[] pbkdf2HmacSha256(
+ char[] password,
+ byte[] salt,
+ int iterations,
+ int keyLenBytes
+ ) {
+ try {
+ var secretKeyFactory = SecretKeyFactory.getInstance(PBKDF2_WITH_HMAC_SHA256);
+ var spec = new PBEKeySpec(password, salt, iterations, keyLenBytes * 8);
+ return secretKeyFactory.generateSecret(spec).getEncoded();
+ } catch (Exception e) {
+ throw new IllegalStateException("%s not available".formatted(PBKDF2_WITH_HMAC_SHA256), e);
+ }
+ }
+
+ private static byte[] hmacSha256(byte[] key, byte[] data) {
+ try {
+ var mac = Mac.getInstance(HMAC_SHA_256);
+ mac.init(new SecretKeySpec(key, HMAC_SHA_256));
+ return mac.doFinal(data);
+ } catch (Exception e) {
+ throw new IllegalStateException("%s not available".formatted(HMAC_SHA_256), e);
+ }
+ }
+
+ private static byte[] sha256(byte[] data) {
+ try {
+ return MessageDigest.getInstance(SHA_256).digest(data);
+ } catch (Exception e) {
+ throw new IllegalStateException("%s not available".formatted(SHA_256), e);
+ }
+ }
+}
diff --git a/src/main/java/it/aboutbits/postgresql/core/PostgreSQLContextFactory.java b/src/main/java/it/aboutbits/postgresql/core/PostgreSQLContextFactory.java
new file mode 100644
index 0000000..8e471a6
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/core/PostgreSQLContextFactory.java
@@ -0,0 +1,59 @@
+package it.aboutbits.postgresql.core;
+
+import io.fabric8.kubernetes.client.KubernetesClient;
+import it.aboutbits.postgresql.crd.clusterconnection.ClusterConnection;
+import jakarta.enterprise.context.ApplicationScoped;
+import lombok.RequiredArgsConstructor;
+import org.jooq.CloseableDSLContext;
+import org.jooq.impl.DSL;
+import org.jspecify.annotations.NullMarked;
+
+import java.util.Properties;
+
+@NullMarked
+@ApplicationScoped
+@RequiredArgsConstructor
+public class PostgreSQLContextFactory {
+ private static final String POSTGRESQL_AUTHENTICATION_USER_KEY = "user";
+ private static final String POSTGRESQL_AUTHENTICATION_PASSWORD_KEY = "password";
+
+ private final KubernetesService kubernetesService;
+ private final KubernetesClient kubernetesClient;
+
+ public CloseableDSLContext getDSLContext(ClusterConnection clusterConnection) {
+ var credentials = kubernetesService.getSecretRefCredentials(
+ kubernetesClient,
+ clusterConnection
+ );
+
+ var spec = clusterConnection.getSpec();
+
+ var jdbcUrl = "jdbc:postgresql://%s:%d/%s".formatted(
+ spec.getHost(),
+ spec.getPort(),
+ spec.getMaintenanceDatabase()
+ );
+
+ var properties = new Properties(2 + spec.getParameters().size());
+
+ properties.setProperty(
+ POSTGRESQL_AUTHENTICATION_USER_KEY,
+ credentials.username()
+ );
+ properties.setProperty(
+ POSTGRESQL_AUTHENTICATION_PASSWORD_KEY,
+ credentials.password()
+ );
+
+ if (!spec.getParameters().isEmpty()) {
+ properties.putAll(
+ spec.getParameters()
+ );
+ }
+
+ return DSL.using(
+ jdbcUrl,
+ properties
+ );
+ }
+}
diff --git a/src/main/java/it/aboutbits/postgresql/core/SQLUtil.java b/src/main/java/it/aboutbits/postgresql/core/SQLUtil.java
new file mode 100644
index 0000000..b07f347
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/core/SQLUtil.java
@@ -0,0 +1,53 @@
+package it.aboutbits.postgresql.core;
+
+import org.jooq.QueryPart;
+import org.jspecify.annotations.NullMarked;
+
+import java.util.List;
+
+import static org.jooq.impl.DSL.sql;
+
+@NullMarked
+public final class SQLUtil {
+ public static QueryPart concatenateQueryPartsWithSpaces(List extends QueryPart> parts) {
+ return concatenateQueryParts(parts, " ");
+ }
+
+ public static QueryPart concatenateQueryPartsWithComma(List extends QueryPart> parts) {
+ return concatenateQueryParts(parts, ", ");
+ }
+
+ /**
+ * Concatenate QueryParts with the requested separator
+ */
+ private static QueryPart concatenateQueryParts(
+ List extends QueryPart> items,
+ String separator
+ ) {
+ int size = items.size();
+
+ if (items.isEmpty()) {
+ return sql("");
+ } else if (size == 1) {
+ return items.getFirst();
+ }
+
+ var template = new StringBuilder();
+
+ // Add the first item without a separator
+ template.append('{').append(0).append('}');
+
+ // Add the rest of the items with the leading separator
+ for (int i = 1; i < size; i++) {
+ template.append(separator).append('{').append(i).append('}');
+ }
+
+ return sql(
+ template.toString(),
+ items.toArray(QueryPart[]::new)
+ );
+ }
+
+ private SQLUtil() {
+ }
+}
diff --git a/src/main/java/it/aboutbits/postgresql/core/SecretRef.java b/src/main/java/it/aboutbits/postgresql/core/SecretRef.java
new file mode 100644
index 0000000..ed8a784
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/core/SecretRef.java
@@ -0,0 +1,23 @@
+package it.aboutbits.postgresql.core;
+
+import io.fabric8.generator.annotation.Required;
+import lombok.Getter;
+import lombok.Setter;
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+
+@NullMarked
+@Getter
+@Setter
+public class SecretRef {
+ @Required
+ private String name = "";
+
+ /**
+ * The namespace where the Secret is located.
+ * If it is null, it means the Secret is in the same namespace as the resource referencing it.
+ */
+ @Nullable
+ @io.fabric8.generator.annotation.Nullable
+ private String namespace;
+}
diff --git a/src/main/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnection.java b/src/main/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnection.java
new file mode 100644
index 0000000..bbc6500
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnection.java
@@ -0,0 +1,76 @@
+package it.aboutbits.postgresql.crd.clusterconnection;
+
+import com.fasterxml.jackson.annotation.JsonIgnore;
+import io.fabric8.crd.generator.annotation.AdditionalPrinterColumn;
+import io.fabric8.kubernetes.api.model.Namespaced;
+import io.fabric8.kubernetes.client.CustomResource;
+import io.fabric8.kubernetes.model.annotation.Group;
+import io.fabric8.kubernetes.model.annotation.Version;
+import it.aboutbits.postgresql.core.CRStatus;
+import it.aboutbits.postgresql.core.Named;
+import org.jspecify.annotations.NullMarked;
+
+import java.net.URLEncoder;
+import java.nio.charset.StandardCharsets;
+import java.util.StringJoiner;
+
+@NullMarked
+@Version("v1")
+@Group("postgresql.aboutbits.it")
+@AdditionalPrinterColumn(
+ name = "Name",
+ jsonPath = ".status.name",
+ type = AdditionalPrinterColumn.Type.STRING
+)
+@AdditionalPrinterColumn(
+ name = "Phase",
+ jsonPath = ".status.phase",
+ type = AdditionalPrinterColumn.Type.STRING
+)
+@AdditionalPrinterColumn(
+ name = "Message",
+ jsonPath = ".status.message",
+ type = AdditionalPrinterColumn.Type.STRING
+)
+@AdditionalPrinterColumn(
+ name = "Since",
+ jsonPath = ".status.lastPhaseTransitionTime",
+ type = AdditionalPrinterColumn.Type.DATE
+)
+@AdditionalPrinterColumn(
+ name = "Age",
+ jsonPath = ".metadata.creationTimestamp",
+ type = AdditionalPrinterColumn.Type.DATE
+)
+public class ClusterConnection
+ extends CustomResource
+ implements Namespaced, Named {
+ @Override
+ @JsonIgnore
+ public String getName() {
+ var spec = getSpec();
+
+ var jdbcUrl = "jdbc:postgresql://%s:%d/%s".formatted(
+ spec.getHost(),
+ spec.getPort(),
+ spec.getMaintenanceDatabase()
+ );
+
+ if (spec.getParameters().isEmpty()) {
+ return jdbcUrl;
+ }
+
+ var stringJoiner = new StringJoiner("&", jdbcUrl + "?", "");
+
+ spec.getParameters().forEach((key, value) ->
+ stringJoiner.add("%s=%s".formatted(
+ URLEncoder.encode(key, StandardCharsets.UTF_8)
+ .replace("+", "%20"),
+ URLEncoder.encode(value, StandardCharsets.UTF_8)
+ .replace("+", "%20")
+ ))
+ );
+
+ return stringJoiner.toString();
+ }
+}
diff --git a/src/main/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionReconciler.java b/src/main/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionReconciler.java
new file mode 100644
index 0000000..1afc816
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionReconciler.java
@@ -0,0 +1,50 @@
+package it.aboutbits.postgresql.crd.clusterconnection;
+
+import io.javaoperatorsdk.operator.api.reconciler.Context;
+import io.javaoperatorsdk.operator.api.reconciler.Reconciler;
+import io.javaoperatorsdk.operator.api.reconciler.UpdateControl;
+import it.aboutbits.postgresql.core.BaseReconciler;
+import it.aboutbits.postgresql.core.CRPhase;
+import it.aboutbits.postgresql.core.CRStatus;
+import it.aboutbits.postgresql.core.PostgreSQLContextFactory;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.jspecify.annotations.NullMarked;
+
+@NullMarked
+@Slf4j
+@RequiredArgsConstructor
+public class ClusterConnectionReconciler
+ extends BaseReconciler
+ implements Reconciler {
+ private final PostgreSQLContextFactory contextFactory;
+
+ @Override
+ public UpdateControl reconcile(
+ ClusterConnection resource,
+ Context context
+ ) {
+ var status = initializeStatus(resource);
+
+ try (var dsl = contextFactory.getDSLContext(resource)) {
+ var version = dsl.fetchSingle("select version()").into(String.class);
+
+ status.setPhase(CRPhase.READY).setMessage(version);
+
+ return UpdateControl.patchStatus(resource);
+ } catch (Exception e) {
+ log.error("Failed to check database connectivity", e);
+
+ return handleError(
+ resource,
+ status,
+ e
+ );
+ }
+ }
+
+ @Override
+ protected CRStatus newStatus() {
+ return new CRStatus();
+ }
+}
diff --git a/src/main/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionSpec.java b/src/main/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionSpec.java
new file mode 100644
index 0000000..8b3bc9a
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionSpec.java
@@ -0,0 +1,30 @@
+package it.aboutbits.postgresql.crd.clusterconnection;
+
+import io.fabric8.generator.annotation.Required;
+import it.aboutbits.postgresql.core.SecretRef;
+import lombok.Getter;
+import lombok.Setter;
+import org.jspecify.annotations.NullMarked;
+
+import java.util.HashMap;
+import java.util.Map;
+
+@NullMarked
+@Getter
+@Setter
+public class ClusterConnectionSpec {
+ @Required
+ private String host = "";
+
+ @Required
+ private Integer port = -1;
+
+ @Required
+ private String maintenanceDatabase = "postgres";
+
+ @Required
+ private SecretRef adminSecretRef = new SecretRef();
+
+ @io.fabric8.generator.annotation.Nullable
+ private Map parameters = new HashMap<>();
+}
diff --git a/src/main/java/it/aboutbits/postgresql/crd/database/.gitkeep b/src/main/java/it/aboutbits/postgresql/crd/database/.gitkeep
new file mode 100644
index 0000000..e69de29
diff --git a/src/main/java/it/aboutbits/postgresql/crd/role/Role.java b/src/main/java/it/aboutbits/postgresql/crd/role/Role.java
new file mode 100644
index 0000000..c56df40
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/crd/role/Role.java
@@ -0,0 +1,49 @@
+package it.aboutbits.postgresql.crd.role;
+
+import com.fasterxml.jackson.annotation.JsonIgnore;
+import io.fabric8.crd.generator.annotation.AdditionalPrinterColumn;
+import io.fabric8.kubernetes.api.model.Namespaced;
+import io.fabric8.kubernetes.client.CustomResource;
+import io.fabric8.kubernetes.model.annotation.Group;
+import io.fabric8.kubernetes.model.annotation.Version;
+import it.aboutbits.postgresql.core.CRStatus;
+import it.aboutbits.postgresql.core.Named;
+import org.jspecify.annotations.NullMarked;
+
+@NullMarked
+@Version("v1")
+@Group("postgresql.aboutbits.it")
+@AdditionalPrinterColumn(
+ name = "Name",
+ jsonPath = ".status.name",
+ type = AdditionalPrinterColumn.Type.STRING
+)
+@AdditionalPrinterColumn(
+ name = "Phase",
+ jsonPath = ".status.phase",
+ type = AdditionalPrinterColumn.Type.STRING
+)
+@AdditionalPrinterColumn(
+ name = "Message",
+ jsonPath = ".status.message",
+ type = AdditionalPrinterColumn.Type.STRING
+)
+@AdditionalPrinterColumn(
+ name = "Since",
+ jsonPath = ".status.lastPhaseTransitionTime",
+ type = AdditionalPrinterColumn.Type.DATE
+)
+@AdditionalPrinterColumn(
+ name = "Age",
+ jsonPath = ".metadata.creationTimestamp",
+ type = AdditionalPrinterColumn.Type.DATE
+)
+public class Role
+ extends CustomResource
+ implements Namespaced, Named {
+ @Override
+ @JsonIgnore
+ public String getName() {
+ return getSpec().getName();
+ }
+}
diff --git a/src/main/java/it/aboutbits/postgresql/crd/role/RoleFlag.java b/src/main/java/it/aboutbits/postgresql/crd/role/RoleFlag.java
new file mode 100644
index 0000000..1731ed2
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/crd/role/RoleFlag.java
@@ -0,0 +1,45 @@
+package it.aboutbits.postgresql.crd.role;
+
+import lombok.Getter;
+import lombok.RequiredArgsConstructor;
+import lombok.experimental.Accessors;
+import org.jspecify.annotations.NullMarked;
+
+@NullMarked
+@Getter
+@Accessors(fluent = true)
+@RequiredArgsConstructor
+public enum RoleFlag {
+ SUPERUSER("SUPERUSER"),
+ NO_SUPERUSER("NOSUPERUSER"),
+
+ CREATEDB("CREATEDB"),
+ NO_CREATEDB("NOCREATEDB"),
+
+ CREATEROLE("CREATEROLE"),
+ NO_CREATEROLE("NOCREATEROLE"),
+
+ INHERIT("INHERIT"),
+ NO_INHERIT("NOINHERIT"),
+
+ LOGIN("LOGIN"),
+ NO_LOGIN("NOLOGIN"),
+
+ REPLICATION("REPLICATION"),
+ NO_REPLICATION("NOREPLICATION"),
+
+ BYPASSRLS("BYPASSRLS"),
+ NO_BYPASSRLS("NOBYPASSRLS"),
+
+ CONNECTION_LIMIT("CONNECTION LIMIT"),
+
+ PASSWORD("PASSWORD"),
+
+ VALID_UNTIL("VALID UNTIL"),
+
+ IN_ROLE("IN ROLE"),
+
+ ROLE("ROLE");
+
+ private final String flag;
+}
diff --git a/src/main/java/it/aboutbits/postgresql/crd/role/RoleReconciler.java b/src/main/java/it/aboutbits/postgresql/crd/role/RoleReconciler.java
new file mode 100644
index 0000000..76395a8
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/crd/role/RoleReconciler.java
@@ -0,0 +1,350 @@
+package it.aboutbits.postgresql.crd.role;
+
+import io.fabric8.kubernetes.api.model.Secret;
+import io.fabric8.kubernetes.client.KubernetesClient;
+import io.javaoperatorsdk.operator.api.config.informer.InformerEventSourceConfiguration;
+import io.javaoperatorsdk.operator.api.reconciler.Cleaner;
+import io.javaoperatorsdk.operator.api.reconciler.Context;
+import io.javaoperatorsdk.operator.api.reconciler.DeleteControl;
+import io.javaoperatorsdk.operator.api.reconciler.EventSourceContext;
+import io.javaoperatorsdk.operator.api.reconciler.Reconciler;
+import io.javaoperatorsdk.operator.api.reconciler.UpdateControl;
+import io.javaoperatorsdk.operator.processing.event.ResourceID;
+import io.javaoperatorsdk.operator.processing.event.source.EventSource;
+import io.javaoperatorsdk.operator.processing.event.source.SecondaryToPrimaryMapper;
+import io.javaoperatorsdk.operator.processing.event.source.informer.InformerEventSource;
+import it.aboutbits.postgresql.core.BaseReconciler;
+import it.aboutbits.postgresql.core.CRPhase;
+import it.aboutbits.postgresql.core.CRStatus;
+import it.aboutbits.postgresql.core.KubernetesService;
+import it.aboutbits.postgresql.core.PostgreSQLAuthenticationService;
+import it.aboutbits.postgresql.core.PostgreSQLContextFactory;
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+import org.jooq.DSLContext;
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+
+@NullMarked
+@Slf4j
+@RequiredArgsConstructor
+public class RoleReconciler
+ extends BaseReconciler
+ implements Reconciler, Cleaner {
+ private final RoleService roleService;
+ private final KubernetesService kubernetesService;
+ private final PostgreSQLAuthenticationService postgreSQLAuthenticationService;
+
+ private final KubernetesClient kubernetesClient;
+ private final PostgreSQLContextFactory contextFactory;
+
+ @Override
+ public UpdateControl reconcile(
+ Role resource,
+ Context context
+ ) {
+ var spec = resource.getSpec();
+ var status = initializeStatus(resource);
+
+ var name = resource.getMetadata().getName();
+ var namespace = resource.getMetadata().getNamespace();
+
+ log.info(
+ "Reconciling Role [resource={}/{}, status.phase={}]",
+ namespace,
+ name,
+ status.getPhase()
+ );
+
+ var clusterRef = spec.getClusterRef();
+ var expectedFlags = spec.getFlags();
+
+ var clusterConnectionOptional = getReferencedClusterConnection(
+ kubernetesClient,
+ resource,
+ clusterRef
+ );
+
+ if (clusterConnectionOptional.isEmpty()) {
+ status.setPhase(CRPhase.PENDING)
+ .setMessage("The specified ClusterConnection does not exist or is not ready yet [clusterRef=%s/%s]".formatted(
+ getResourceNamespaceOrOwn(resource, clusterRef.getNamespace()),
+ clusterRef.getName()
+ ));
+
+ return UpdateControl.patchStatus(resource)
+ .rescheduleAfter(60, TimeUnit.SECONDS);
+ }
+
+ var clusterConnection = clusterConnectionOptional.get();
+
+ // We need to case-insensitive sort the roles, as PostgreSQL will lowercase anything without quotes
+ expectedFlags.getRole().sort(String.CASE_INSENSITIVE_ORDER);
+ expectedFlags.getInRole().sort(String.CASE_INSENSITIVE_ORDER);
+
+ var passwordSecretRef = spec.getPasswordSecretRef();
+
+ String password;
+ if (passwordSecretRef != null) {
+ password = kubernetesService.getSecretRefCredentials(
+ kubernetesClient,
+ passwordSecretRef,
+ namespace
+ ).password();
+ } else {
+ password = null;
+ }
+
+ UpdateControl updateControl;
+
+ try (var dsl = contextFactory.getDSLContext(clusterConnection)) {
+ // Run everything in a single transaction
+ updateControl = dsl.transactionResult(
+ cfg -> reconcileInTransaction(
+ cfg.dsl(),
+ resource,
+ status,
+ password
+ )
+ );
+ } catch (Exception e) {
+ return handleError(
+ resource,
+ status,
+ e
+ );
+ }
+
+ return updateControl;
+ }
+
+ @Override
+ public DeleteControl cleanup(
+ Role resource,
+ Context context
+ ) {
+ var spec = resource.getSpec();
+ var status = initializeStatus(resource);
+
+ var name = resource.getMetadata().getName();
+ var namespace = resource.getMetadata().getNamespace();
+
+ log.info(
+ "Deleting Role [resource={}/{}, spec.name={}, status.phase={}]",
+ namespace,
+ name,
+ spec.getName(),
+ status.getPhase()
+ );
+
+ if (status.getPhase() != CRPhase.DELETING) {
+ status.setPhase(CRPhase.DELETING)
+ .setMessage("Role deletion in progress");
+ }
+
+ var clusterRef = spec.getClusterRef();
+
+ var clusterConnectionOptional = getReferencedClusterConnection(
+ kubernetesClient,
+ resource,
+ clusterRef
+ );
+
+ if (clusterConnectionOptional.isEmpty()) {
+ status.setMessage("The specified ClusterConnection no longer exists or is not ready yet [clusterRef=%s/%s]".formatted(
+ getResourceNamespaceOrOwn(resource, clusterRef.getNamespace()),
+ clusterRef.getName()
+ ));
+
+ return DeleteControl.noFinalizerRemoval()
+ .rescheduleAfter(60, TimeUnit.SECONDS);
+ }
+
+ var clusterConnection = clusterConnectionOptional.get();
+
+ try (var dsl = contextFactory.getDSLContext(clusterConnection)) {
+ roleService.dropRole(dsl, spec);
+
+ return DeleteControl.defaultDelete();
+ } catch (Exception e) {
+ log.error(
+ "Failed to delete Role [resource={}/{}, spec.name={}, status.phase={}]",
+ namespace,
+ name,
+ spec.getName(),
+ status.getPhase()
+ );
+
+ status.setMessage("Deletion failed: " + e.getMessage());
+
+ return DeleteControl.noFinalizerRemoval()
+ .rescheduleAfter(60, TimeUnit.SECONDS);
+ }
+ }
+
+ /**
+ * Watches for {@code Secret} changes to trigger reconciliation for dependent {@code Role} resources.
+ */
+ @Override
+ public List> prepareEventSources(EventSourceContext context) {
+ // 1. Define the Mapper
+ // We define how to find the Primary Resource (Role) when a Secret changes
+ // Filter Roles that reference this specific Secret
+ SecondaryToPrimaryMapper secretToRoleMapper = (Secret secret) -> context.getPrimaryCache()
+ .list()
+ .filter(role -> isReferencedBy(role, secret))
+ .map(ResourceID::fromResource)
+ .collect(Collectors.toSet());
+
+ // 2. Build the Event Source Configuration which binds the InformerConfig + Mapper
+ var eventSourceConfig = InformerEventSourceConfiguration.from(Secret.class, Role.class)
+ .withSecondaryToPrimaryMapper(secretToRoleMapper)
+ // or .withWatchAllNamespaces() if we want to have the secret in another namespace than the Role CR instance
+ .withNamespacesInheritedFromController()
+ .build();
+
+ // 3. Create the Event Source
+ // This will watch for Secret changes and run the mapper
+ var secretEventSource = new InformerEventSource<>(
+ eventSourceConfig,
+ context
+ );
+
+ return List.of(secretEventSource);
+ }
+
+ private UpdateControl reconcileInTransaction(
+ DSLContext tx,
+ Role resource,
+ CRStatus status,
+ @Nullable String password
+ ) {
+ var name = resource.getMetadata().getName();
+ var namespace = resource.getMetadata().getNamespace();
+
+ var spec = resource.getSpec();
+ var expectedFlags = spec.getFlags();
+
+ // Create and return the role if it doesn't exist yet
+ if (!roleService.roleExists(tx, spec)) {
+ log.info(
+ "Creating Role [resource={}/{}]",
+ namespace,
+ name
+ );
+
+ roleService.createRole(
+ tx,
+ spec,
+ password
+ );
+
+ status.setPhase(CRPhase.READY)
+ .setMessage(null);
+
+ return UpdateControl.patchStatus(resource);
+ }
+
+ // When there is NOLOGIN, we set no password
+ var passwordMatches = true;
+ var roleLoginMatches = roleService.roleLoginMatches(tx, spec);
+ var currentFlags = roleService.fetchCurrentFlags(tx, spec);
+ var flagsMatch = expectedFlags.equals(currentFlags);
+ var commentMatches = roleService.roleCommentMatches(tx, spec);
+
+ var passwordSecretRef = spec.getPasswordSecretRef();
+ var loginExpected = passwordSecretRef != null;
+
+ if (loginExpected && password != null) {
+ passwordMatches = postgreSQLAuthenticationService.passwordMatches(
+ tx,
+ spec,
+ password
+ );
+ }
+
+ if (roleLoginMatches && passwordMatches && flagsMatch && commentMatches) {
+ log.info(
+ "Role up-to-date [resource={}/{}]",
+ namespace,
+ name
+ );
+
+ return UpdateControl.noUpdate();
+ }
+
+ var changePassword = loginExpected && !passwordMatches;
+
+ log.info(
+ "Updating Role [resource={}/{}]",
+ namespace,
+ name
+ );
+
+ if (!roleLoginMatches || !passwordMatches || !flagsMatch) {
+ roleService.alterRole(
+ tx,
+ spec,
+ changePassword,
+ password
+ );
+ }
+
+ if (!flagsMatch) {
+ log.info(
+ "Updating Role membership [resource={}/{}]",
+ namespace,
+ name
+ );
+
+ roleService.reconcileRoleMembership(
+ tx,
+ spec,
+ expectedFlags,
+ currentFlags
+ );
+ }
+
+ if (!commentMatches) {
+ roleService.updateComment(
+ tx,
+ spec
+ );
+ }
+
+ status.setPhase(CRPhase.READY)
+ .setMessage(null);
+
+ return UpdateControl.patchStatus(resource);
+ }
+
+ @Override
+ protected CRStatus newStatus() {
+ return new CRStatus();
+ }
+
+ /**
+ * Checks if the given Role's spec.passwordSecretRef points to the changed Secret.
+ */
+ private boolean isReferencedBy(
+ Role role,
+ Secret secret
+ ) {
+ var spec = role.getSpec();
+
+ if (spec.getPasswordSecretRef() == null) {
+ return false;
+ }
+
+ var ref = spec.getPasswordSecretRef();
+ var refName = ref.getName();
+ var refNamespace = getResourceNamespaceOrOwn(role, ref.getNamespace());
+
+ return refName.equals(secret.getMetadata().getName())
+ && refNamespace.equals(secret.getMetadata().getNamespace());
+ }
+}
diff --git a/src/main/java/it/aboutbits/postgresql/crd/role/RoleService.java b/src/main/java/it/aboutbits/postgresql/crd/role/RoleService.java
new file mode 100644
index 0000000..740b387
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/crd/role/RoleService.java
@@ -0,0 +1,434 @@
+package it.aboutbits.postgresql.crd.role;
+
+import it.aboutbits.postgresql.core.SQLUtil;
+import it.aboutbits.postgresql.core.infrastructure.persistence.Routines;
+import jakarta.inject.Singleton;
+import org.jooq.DSLContext;
+import org.jooq.Query;
+import org.jooq.QueryPart;
+import org.jooq.Record1;
+import org.jooq.impl.DSL;
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.Objects;
+
+import static it.aboutbits.postgresql.core.infrastructure.persistence.Tables.PG_AUTHID;
+import static it.aboutbits.postgresql.core.infrastructure.persistence.Tables.PG_AUTH_MEMBERS;
+import static org.jooq.impl.DSL.field;
+import static org.jooq.impl.DSL.keyword;
+import static org.jooq.impl.DSL.multiset;
+import static org.jooq.impl.DSL.name;
+import static org.jooq.impl.DSL.query;
+import static org.jooq.impl.DSL.role;
+import static org.jooq.impl.DSL.select;
+import static org.jooq.impl.DSL.selectOne;
+import static org.jooq.impl.DSL.sql;
+import static org.jooq.impl.DSL.val;
+
+@NullMarked
+@Singleton
+public final class RoleService {
+ public boolean roleExists(
+ DSLContext tx,
+ RoleSpec spec
+ ) {
+ return tx.fetchExists(selectOne()
+ .from(PG_AUTHID)
+ .where(PG_AUTHID.ROLNAME.eq(spec.getName()))
+ );
+ }
+
+ public void createRole(
+ DSLContext tx,
+ RoleSpec spec,
+ @Nullable String password
+ ) {
+ var roleName = spec.getName();
+ var flags = spec.getFlags();
+ var comment = spec.getComment();
+
+ tx.execute(
+ buildCreateRole(
+ roleName,
+ flags,
+ password
+ )
+ );
+
+ // Optional comment
+ if (comment != null && !comment.isBlank()) {
+ tx.execute(
+ buildCommentOnRole(roleName, comment)
+ );
+ }
+ }
+
+ public void alterRole(
+ DSLContext tx,
+ RoleSpec spec,
+ boolean changePassword,
+ @Nullable String password
+ ) {
+ var roleName = spec.getName();
+ var flags = spec.getFlags();
+
+ tx.execute(
+ buildAlterRole(
+ roleName,
+ flags,
+ changePassword,
+ password
+ )
+ );
+ }
+
+ public void updateComment(
+ DSLContext tx,
+ RoleSpec spec
+ ) {
+ var roleName = spec.getName();
+ var expectedComment = normalizeComment(spec.getComment());
+
+ var currentComment = normalizeComment(
+ fetchCurrentRoleComment(tx, roleName)
+ );
+
+ if (!Objects.equals(currentComment, expectedComment)) {
+ tx.execute(
+ buildCommentOnRole(roleName, expectedComment)
+ );
+ }
+ }
+
+ public boolean roleCommentMatches(
+ DSLContext tx,
+ RoleSpec spec
+ ) {
+ var expectedComment = spec.getComment();
+
+ var currentComment = normalizeComment(
+ fetchCurrentRoleComment(tx, spec.getName())
+ );
+
+ return Objects.equals(currentComment, expectedComment);
+ }
+
+ public @Nullable String fetchCurrentRoleComment(
+ DSLContext tx,
+ String roleName
+ ) {
+
+ return tx
+ .select(Routines.shobjDescription(
+ PG_AUTHID.OID,
+ val(PG_AUTHID.getUnqualifiedName().last())
+ ))
+ .from(PG_AUTHID)
+ .where(PG_AUTHID.ROLNAME.eq(roleName))
+ .fetchOneInto(String.class);
+ }
+
+ public boolean roleLoginMatches(
+ DSLContext tx,
+ RoleSpec spec
+ ) {
+ var loginExpected = spec.getPasswordSecretRef() != null;
+
+ var canLogin = tx.fetchExists(selectOne()
+ .from(PG_AUTHID)
+ .where(PG_AUTHID.ROLNAME.eq(spec.getName()))
+ .and(PG_AUTHID.ROLCANLOGIN.isTrue())
+ );
+
+ return loginExpected == canLogin;
+ }
+
+ public RoleSpec.Flags fetchCurrentFlags(
+ DSLContext tx,
+ RoleSpec spec
+ ) {
+ var member = PG_AUTHID.as("member");
+ var parent = PG_AUTHID.as("parent");
+
+ return tx
+ .select(
+ PG_AUTHID.ROLSUPER.as("superuser"),
+ PG_AUTHID.ROLCREATEDB.as("createdb"),
+ PG_AUTHID.ROLCREATEROLE.as("createrole"),
+ PG_AUTHID.ROLINHERIT.as("inherit"),
+ PG_AUTHID.ROLREPLICATION.as("replication"),
+ PG_AUTHID.ROLBYPASSRLS.as("bypassrls"),
+ PG_AUTHID.ROLCONNLIMIT.as("connectionLimit"),
+ field("nullif({0}, 'infinity')", PG_AUTHID.ROLVALIDUNTIL.getDataType(), PG_AUTHID.ROLVALIDUNTIL).as("validUntil"),
+ multiset(
+ select(parent.ROLNAME)
+ .from(PG_AUTH_MEMBERS)
+ .join(member).on(member.OID.eq(PG_AUTH_MEMBERS.MEMBER))
+ .join(parent).on(parent.OID.eq(PG_AUTH_MEMBERS.ROLEID))
+ .where(member.OID.eq(PG_AUTHID.OID))
+ .orderBy(parent.ROLNAME)
+ ).as("inRole").convertFrom(result -> result.map(Record1::value1)),
+ multiset(
+ select(member.ROLNAME)
+ .from(PG_AUTH_MEMBERS)
+ .join(parent).on(parent.OID.eq(PG_AUTH_MEMBERS.ROLEID))
+ .join(member).on(member.OID.eq(PG_AUTH_MEMBERS.MEMBER))
+ .where(parent.OID.eq(PG_AUTHID.OID))
+ .orderBy(member.ROLNAME)
+ ).as("role").convertFrom(result -> result.map(Record1::value1))
+ )
+ .from(PG_AUTHID)
+ .where(PG_AUTHID.ROLNAME.eq(spec.getName()))
+ .fetchSingleInto(RoleSpec.Flags.class);
+ }
+
+ public void reconcileRoleMembership(
+ DSLContext tx,
+ RoleSpec spec,
+ RoleSpec.Flags expectedFlags,
+ RoleSpec.Flags currentFlags
+ ) {
+ var roleName = spec.getName();
+
+ // ROLE IN
+ var expectedInRole = new HashSet<>(expectedFlags.getInRole());
+ var currentInRole = new HashSet<>(currentFlags.getInRole());
+
+ var queries = new ArrayList();
+
+ var inRoleToGrant = new HashSet<>(expectedInRole);
+ inRoleToGrant.removeAll(currentInRole);
+
+ var inRoleToRevoke = new HashSet<>(currentInRole);
+ inRoleToRevoke.removeAll(expectedInRole);
+
+ for (var parentRole : inRoleToGrant) {
+ // GRANT parentRole TO roleName
+ queries.add(buildGrantRoleToMember(parentRole, roleName));
+ }
+ for (var parentRole : inRoleToRevoke) {
+ // REVOKE parentRole FROM roleName
+ queries.add(buildRevokeRoleFromMember(parentRole, roleName));
+ }
+
+ // ROLE
+ var expectedRoleMembers = new HashSet<>(expectedFlags.getRole());
+ var currentRoleMembers = new HashSet<>(currentFlags.getRole());
+
+ var roleMembersToGrant = new HashSet<>(expectedRoleMembers);
+ roleMembersToGrant.removeAll(currentRoleMembers);
+
+ var roleMembersToRevoke = new HashSet<>(currentRoleMembers);
+ roleMembersToRevoke.removeAll(expectedRoleMembers);
+
+ for (var member : roleMembersToGrant) {
+ // GRANT roleName TO member
+ queries.add(buildGrantRoleToMember(roleName, member));
+ }
+ for (var member : roleMembersToRevoke) {
+ // REVOKE roleName FROM member
+ queries.add(buildRevokeRoleFromMember(roleName, member));
+ }
+
+ if (!queries.isEmpty()) {
+ tx.batch(queries).execute();
+ }
+ }
+
+ public void dropRole(
+ DSLContext dsl,
+ RoleSpec spec
+ ) {
+ dsl.execute(
+ query("drop role if exists {0}", role(spec.getName()))
+ );
+ }
+
+ /**
+ * Build: CREATE ROLE [ [ WITH ] option [ ... ] ]
+ * See
+ * PostgreSQL: Documentation: CREATE ROLE
+ *
+ */
+ private static Query buildCreateRole(
+ String roleName,
+ RoleSpec.Flags flags,
+ @Nullable String password
+ ) {
+ var options = new ArrayList();
+
+ // Only allow the user to log in if a password is specified.
+ if (password != null) {
+ options.add(keyword(RoleFlag.LOGIN.flag()));
+ options.add(keyword(RoleFlag.PASSWORD.flag()));
+ options.add(val(password));
+ }
+
+ if (flags.isSuperuser()) {
+ options.add(keyword(RoleFlag.SUPERUSER.flag()));
+ }
+ if (flags.isCreatedb()) {
+ options.add(keyword(RoleFlag.CREATEDB.flag()));
+ }
+ if (flags.isCreaterole()) {
+ options.add(keyword(RoleFlag.CREATEROLE.flag()));
+ }
+ if (flags.isInherit()) {
+ options.add(keyword(RoleFlag.INHERIT.flag()));
+ }
+ if (flags.isReplication()) {
+ options.add(keyword(RoleFlag.REPLICATION.flag()));
+ }
+ if (flags.isBypassrls()) {
+ options.add(keyword(RoleFlag.BYPASSRLS.flag()));
+ }
+ if (flags.getConnectionLimit() >= 0) {
+ options.add(keyword(RoleFlag.CONNECTION_LIMIT.flag()));
+ options.add(val(flags.getConnectionLimit()));
+ }
+
+ var validUntil = flags.getValidUntil();
+ if (validUntil != null) {
+ options.add(keyword(RoleFlag.VALID_UNTIL.flag()));
+ options.add(val(validUntil.toString()));
+ }
+
+ if (!flags.getInRole().isEmpty()) {
+ options.add(keyword(RoleFlag.IN_ROLE.flag()));
+ options.add(SQLUtil.concatenateQueryPartsWithComma(
+ flags.getInRole()
+ .stream()
+ .map(DSL::role)
+ .toList()
+ ));
+ }
+ if (!flags.getRole().isEmpty()) {
+ options.add(keyword(RoleFlag.ROLE.flag()));
+ options.add(SQLUtil.concatenateQueryPartsWithComma(
+ flags.getRole()
+ .stream()
+ .map(DSL::role)
+ .toList()
+ ));
+ }
+
+ var optionsSql = options.isEmpty()
+ ? sql("") // nothing
+ : sql(" with {0}", SQLUtil.concatenateQueryPartsWithSpaces(options));
+
+ return query(
+ "create role {0}{1}",
+ role(roleName),
+ optionsSql
+ );
+ }
+
+ private static Query buildAlterRole(
+ String roleName,
+ RoleSpec.Flags flags,
+ boolean changePassword,
+ @Nullable String password
+ ) {
+ var options = new ArrayList();
+ var loginExpected = password != null;
+
+ // LOGIN / NOLOGIN
+ options.add(keyword(loginExpected
+ ? RoleFlag.LOGIN.flag()
+ : RoleFlag.NO_LOGIN.flag()
+ ));
+
+ // Password handling
+ // - if NOLOGIN, remove the password
+ // - if LOGIN and passwordChanged, set the new password
+ if (!loginExpected) {
+ options.add(keyword(RoleFlag.PASSWORD.flag()));
+ options.add(keyword("NULL"));
+ } else if (changePassword) {
+ options.add(keyword(RoleFlag.PASSWORD.flag()));
+ options.add(val(password));
+ }
+
+ // Explicitly set the expected state to make the statement idempotent
+ options.add(keyword(flags.isSuperuser()
+ ? RoleFlag.SUPERUSER.flag()
+ : RoleFlag.NO_SUPERUSER.flag()
+ ));
+ options.add(keyword(flags.isCreatedb()
+ ? RoleFlag.CREATEDB.flag()
+ : RoleFlag.NO_CREATEDB.flag()
+ ));
+ options.add(keyword(flags.isCreaterole()
+ ? RoleFlag.CREATEROLE.flag()
+ : RoleFlag.NO_CREATEROLE.flag()
+ ));
+ options.add(keyword(flags.isInherit()
+ ? RoleFlag.INHERIT.flag()
+ : RoleFlag.NO_INHERIT.flag()
+ ));
+ options.add(keyword(flags.isReplication()
+ ? RoleFlag.REPLICATION.flag()
+ : RoleFlag.NO_REPLICATION.flag()
+ ));
+ options.add(keyword(flags.isBypassrls()
+ ? RoleFlag.BYPASSRLS.flag()
+ : RoleFlag.NO_BYPASSRLS.flag()
+ ));
+
+ options.add(keyword(RoleFlag.CONNECTION_LIMIT.flag()));
+ options.add(val(flags.getConnectionLimit()));
+
+ var validUntil = flags.getValidUntil();
+ options.add(keyword(RoleFlag.VALID_UNTIL.flag()));
+ if (validUntil != null) {
+ options.add(val(validUntil.toString()));
+ } else {
+ options.add(val("infinity"));
+ }
+
+ return query(
+ "alter role {0} with {1}",
+ role(roleName),
+ SQLUtil.concatenateQueryPartsWithSpaces(options)
+ );
+ }
+
+ private static Query buildGrantRoleToMember(
+ String role,
+ String member
+ ) {
+ return query("grant {0} to {1}", role(role), role(member));
+ }
+
+ private static Query buildRevokeRoleFromMember(
+ String role,
+ String member
+ ) {
+ return query("revoke {0} from {1}", role(role), role(member));
+ }
+
+ /**
+ * Build: COMMENT ON ROLE IS
+ */
+ private static Query buildCommentOnRole(
+ String roleName,
+ @Nullable String comment
+ ) {
+ return query(
+ "comment on role {0} is {1}",
+ name(roleName),
+ val(comment)
+ );
+ }
+
+ private static @Nullable String normalizeComment(@Nullable String comment) {
+ if (comment == null || comment.isBlank()) {
+ return null;
+ }
+
+ return comment;
+ }
+}
diff --git a/src/main/java/it/aboutbits/postgresql/crd/role/RoleSpec.java b/src/main/java/it/aboutbits/postgresql/crd/role/RoleSpec.java
new file mode 100644
index 0000000..86c4bb6
--- /dev/null
+++ b/src/main/java/it/aboutbits/postgresql/crd/role/RoleSpec.java
@@ -0,0 +1,79 @@
+package it.aboutbits.postgresql.crd.role;
+
+import io.fabric8.generator.annotation.Required;
+import io.fabric8.generator.annotation.ValidationRule;
+import it.aboutbits.postgresql.core.ClusterReference;
+import it.aboutbits.postgresql.core.SecretRef;
+import lombok.EqualsAndHashCode;
+import lombok.Getter;
+import lombok.Setter;
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+
+import java.time.OffsetDateTime;
+import java.util.ArrayList;
+import java.util.List;
+
+@NullMarked
+@Getter
+@Setter
+public class RoleSpec {
+ @Required
+ @ValidationRule(
+ value = "self == oldSelf",
+ message = "The Role name must not be changed once it is created"
+ )
+ private String name = "";
+
+ @Nullable
+ @io.fabric8.generator.annotation.Nullable
+ private String comment;
+
+ @Required
+ private ClusterReference clusterRef = new ClusterReference();
+
+ @Nullable
+ @io.fabric8.generator.annotation.Nullable
+ private SecretRef passwordSecretRef;
+
+ @io.fabric8.generator.annotation.Nullable
+ private Flags flags = new Flags();
+
+ @Getter
+ @Setter
+ @EqualsAndHashCode
+ // The Fabric8 @Nullable annotation is relevant for generating nullable annotations in the resulting CRD YAML JSON Schema
+ @SuppressWarnings({"NullablePrimitive"})
+ public static class Flags {
+ @io.fabric8.generator.annotation.Nullable
+ private boolean superuser = false;
+
+ @io.fabric8.generator.annotation.Nullable
+ private boolean createdb = false;
+
+ @io.fabric8.generator.annotation.Nullable
+ private boolean createrole = false;
+
+ @io.fabric8.generator.annotation.Nullable
+ private boolean inherit = true;
+
+ @io.fabric8.generator.annotation.Nullable
+ private boolean replication = false;
+
+ @io.fabric8.generator.annotation.Nullable
+ private boolean bypassrls = false;
+
+ @io.fabric8.generator.annotation.Nullable
+ private int connectionLimit = -1;
+
+ @Nullable
+ @io.fabric8.generator.annotation.Nullable
+ private OffsetDateTime validUntil = null;
+
+ @io.fabric8.generator.annotation.Nullable
+ private List inRole = new ArrayList<>();
+
+ @io.fabric8.generator.annotation.Nullable
+ private List role = new ArrayList<>();
+ }
+}
diff --git a/src/main/resources/META-INF/branding/logo.png b/src/main/resources/META-INF/branding/logo.png
new file mode 100644
index 0000000..915e323
Binary files /dev/null and b/src/main/resources/META-INF/branding/logo.png differ
diff --git a/src/main/resources/application-dev.yml b/src/main/resources/application-dev.yml
new file mode 100644
index 0000000..96e49bc
--- /dev/null
+++ b/src/main/resources/application-dev.yml
@@ -0,0 +1,6 @@
+quarkus:
+ datasource:
+ devservices:
+ port: 5432
+ test:
+ continuous-testing: enabled
diff --git a/src/main/resources/application-test.yml b/src/main/resources/application-test.yml
new file mode 100644
index 0000000..d37b73f
--- /dev/null
+++ b/src/main/resources/application-test.yml
@@ -0,0 +1,8 @@
+quarkus:
+ operator-sdk:
+ activate-leader-election-for-profiles:
+ - test
+ log:
+ category:
+ "org.jooq":
+ level: DEBUG
diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml
new file mode 100644
index 0000000..b3783d2
--- /dev/null
+++ b/src/main/resources/application.yml
@@ -0,0 +1,158 @@
+quarkus:
+ kubernetes-client:
+ devservices:
+ enabled: true
+ # To not use our default kubeconfig from the home directory, which could lead to
+ # potential damage if we had the permissions to install the CRD in an existing configured cluster context.
+ # By setting this to true, Quarkus will use a temporary kubeconfig that will be removed after the test.
+ override-kubeconfig: true
+ flavor: k3s
+ # See https://github.com/dajudge/kindcontainer/blob/master/k8s-versions.json
+ api-version: 1.34.1
+ live-reload:
+ instrumentation: true
+ micrometer:
+ enabled: true
+ datasource:
+ devservices:
+ enabled: true
+ image-name: postgres:17.7
+ username: root
+ password: password
+ reuse: false
+ jdbc:
+ metrics:
+ enabled: true
+ operator-sdk:
+ crd:
+ generate: true
+ # NOTE that this option is only considered when *not* in production mode
+ # as applying the CRD to a production cluster could be dangerous if done automatically.
+ apply: true
+ # Whether controllers should only process events if the associated resource generation
+ # has increased since the last reconciliation, otherwise will process all events.
+ generation-aware: true
+ test:
+ hang-detection-timeout: PT1M
+ #log:
+ # category:
+ # "io.javaoperatorsdk":
+ # level: DEBUG
+ # "io.quarkiverse.operatorsdk":
+ # level: DEBUG
+
+ # Config for the generated Helm chart #
+ # Container Image config for Kubernetes Helm #
+ container-image:
+ registry: ghcr.io
+ group: aboutbits/postgresql-operator
+ name: app
+ tag: ${quarkus.application.version}
+ helm:
+ app-version: ${quarkus.application.version}
+ type: application
+ name: ${quarkus.kubernetes.name}
+ description: AboutBits PostgreSQL Operator Helm Chart
+ # Keep the map-system-properties flag below on false, else it will map all ${ENV_VARS} we define in the application.yml or application-prod.yml files.
+ # In Kubernetes the env entries take precedence over envFrom entries
+ map-system-properties: false
+ home: https://github.com/aboutbits/postgresql-operator
+ sources:
+ - https://github.com/aboutbits/postgresql-operator
+ annotations:
+ "catalog.cattle.io/os": linux
+ keywords:
+ - aboutbits
+ - postgresql
+ - operator
+ maintainers:
+ "AboutBits":
+ name: AboutBits
+ email: info@aboutbits.it
+ url: https://aboutbits.it/
+ create-tar-file: true
+ extension: tgz
+ values:
+ replicas:
+ property: replicas
+ value: 1
+ paths:
+ - (kind == Deployment).spec.replicas
+ image-pull-policy:
+ property: imagePullPolicy
+ value: IfNotPresent
+ paths:
+ - (kind == Deployment).spec.template.spec.containers.(name == ${quarkus.kubernetes.name}).imagePullPolicy
+ resource-requests-cpu:
+ property: resources.requests.cpu
+ value: ${quarkus.kubernetes.resources.requests.cpu}
+ paths:
+ - (kind == Deployment).spec.template.spec.containers.(name == ${quarkus.kubernetes.name}).resources.requests.cpu
+ resource-requests-memory:
+ property: resources.requests.memory
+ value: ${quarkus.kubernetes.resources.requests.memory}
+ paths:
+ - (kind == Deployment).spec.template.spec.containers.(name == ${quarkus.kubernetes.name}).resources.requests.memory
+ resource-limits-memory:
+ property: resources.limits.memory
+ value: ${quarkus.kubernetes.resources.limits.memory}
+ paths:
+ - (kind == Deployment).spec.template.spec.containers.(name == ${quarkus.kubernetes.name}).resources.limits.memory
+ image-pull-secret:
+ property: imagePullSecret
+ value: ${quarkus.kubernetes.image-pull-secrets[0]}
+ paths:
+ - (kind == Deployment).spec.template.spec.imagePullSecrets[0].name
+ expressions:
+ release-name-labels:
+ expression: "{{ .Release.Name }}"
+ path: metadata.labels.'app.kubernetes.io/name'
+ release-name-service-selector:
+ expression: "{{ .Release.Name }}"
+ path: (kind == Service).spec.selector.'app.kubernetes.io/name'
+ release-name-deployment-match-labels:
+ expression: "{{ .Release.Name }}"
+ path: (kind == Deployment).spec.selector.matchLabels.'app.kubernetes.io/name'
+ release-name-deployment-labels:
+ expression: "{{ .Release.Name }}"
+ path: (kind == Deployment).spec.template.metadata.labels.'app.kubernetes.io/name'
+ kubernetes:
+ name: postgresql-operator
+ version: ${quarkus.application.version}
+ add-version-to-label-selectors: false
+ image-pull-policy: IfNotPresent
+ image-pull-secrets:
+ - github-container-registry
+ replicas: 1
+ annotations:
+ "app.kubernetes.io/version": ${quarkus.application.version}
+ resources:
+ requests:
+ cpu: 50m
+ memory: 300Mi
+ limits:
+ memory: 512Mi
+ startup-probe:
+ http-action-port-name: http
+ initial-delay: PT2S
+ period: PT10S
+ timeout: PT3S
+ success-threshold: 1
+ failure-threshold: 3
+ readiness-probe:
+ http-action-port-name: http
+ initial-delay: PT0S
+ period: PT10S
+ timeout: PT3S
+ success-threshold: 1
+ failure-threshold: 3
+ liveness-probe:
+ http-action-port-name: http
+ initial-delay: PT10S
+ period: PT30S
+ timeout: PT10S
+ success-threshold: 1
+ failure-threshold: 3
+ env:
+ fields:
+ KUBERNETES_NODE_NAME: spec.nodeName
diff --git a/src/main/resources/default_banner.txt b/src/main/resources/default_banner.txt
new file mode 100644
index 0000000..a3b0c25
--- /dev/null
+++ b/src/main/resources/default_banner.txt
@@ -0,0 +1,6 @@
+ _ _ _ ____ _ _
+ / \ | |__ ___ _ _| |_| __ )(_) |_ ___
+ / _ \ | '_ \ / _ \| | | | __| _ \| | __/ __|
+ / ___ \| |_) | (_) | |_| | |_| |_) | | |_\__ \
+ /_/ \_\_.__/ \___/ \__,_|\__|____/|_|\__|___/
+ PostgreSQL Operator
\ No newline at end of file
diff --git a/src/test/java/it/aboutbits/postgresql/PostgreSQLInstanceReadinessCheckTest.java b/src/test/java/it/aboutbits/postgresql/PostgreSQLInstanceReadinessCheckTest.java
new file mode 100644
index 0000000..dd04b70
--- /dev/null
+++ b/src/test/java/it/aboutbits/postgresql/PostgreSQLInstanceReadinessCheckTest.java
@@ -0,0 +1,102 @@
+package it.aboutbits.postgresql;
+
+import io.fabric8.kubernetes.client.KubernetesClient;
+import io.quarkus.test.junit.QuarkusTest;
+import it.aboutbits.postgresql._support.testdata.persisted.Given;
+import it.aboutbits.postgresql.crd.clusterconnection.ClusterConnection;
+import jakarta.inject.Inject;
+import org.eclipse.microprofile.health.HealthCheckResponse;
+import org.eclipse.microprofile.health.Readiness;
+import org.jspecify.annotations.NullMarked;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import java.util.Objects;
+import java.util.concurrent.TimeUnit;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+@NullMarked
+@QuarkusTest
+class PostgreSQLInstanceReadinessCheckTest {
+ @Inject
+ Given given;
+
+ @Inject
+ @Readiness
+ PostgreSQLInstanceReadinessCheck readinessCheck;
+
+ @Inject
+ KubernetesClient kubernetesClient;
+
+ @BeforeEach
+ void cleanUp() {
+ kubernetesClient.resources(ClusterConnection.class)
+ .withTimeout(5, TimeUnit.SECONDS)
+ .delete();
+ }
+
+ @Test
+ void call_whenAllConnectionsUp_shouldReturnUp() {
+ given.one()
+ .clusterConnection()
+ .withName("test-db")
+ .returnFirst();
+
+ var response = readinessCheck.call();
+
+ assertThat(response.getStatus()).isEqualTo(
+ HealthCheckResponse.Status.UP
+ );
+
+ assertThat(response.getData())
+ .isPresent()
+ .get()
+ .satisfies(data -> {
+ assertThat(data).containsKey("test-db");
+
+ var dbStatus = Objects.requireNonNull(
+ data.get("test-db")
+ );
+
+ assertThat(
+ dbStatus.toString()
+ ).startsWith("UP (PostgreSQL");
+ });
+ }
+
+ @Test
+ void call_whenSomeConnectionsDown_shouldReturnDown() {
+ given.one()
+ .clusterConnection()
+ .withName("db-1")
+ .returnFirst();
+
+ given.one()
+ .clusterConnection()
+ .withName("db-2")
+ .withHost("non-existent-host")
+ .returnFirst();
+
+ var response = readinessCheck.call();
+
+ assertThat(response.getStatus()).isEqualTo(
+ HealthCheckResponse.Status.DOWN
+ );
+
+ assertThat(response.getData())
+ .isPresent()
+ .get()
+ .satisfies(data -> {
+ assertThat(data).containsKey("db-1");
+
+ var dbStatus = Objects.requireNonNull(
+ data.get("db-1")
+ );
+
+ assertThat(dbStatus.toString()).startsWith("UP (PostgreSQL");
+
+ assertThat(data).containsEntry("db-2", "DOWN");
+ });
+ }
+}
diff --git a/src/test/java/it/aboutbits/postgresql/_support/testdata/base/TestDataCreator.java b/src/test/java/it/aboutbits/postgresql/_support/testdata/base/TestDataCreator.java
new file mode 100644
index 0000000..e7e871b
--- /dev/null
+++ b/src/test/java/it/aboutbits/postgresql/_support/testdata/base/TestDataCreator.java
@@ -0,0 +1,69 @@
+package it.aboutbits.postgresql._support.testdata.base;
+
+import net.datafaker.Faker;
+import org.jspecify.annotations.NullMarked;
+
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+
+@NullMarked
+public abstract class TestDataCreator {
+ protected static final Faker FAKER = new Faker();
+
+ protected final int numberOfItems;
+
+ protected TestDataCreator(int numberOfItems) {
+ this.numberOfItems = numberOfItems;
+ }
+
+ public void apply() {
+ create();
+ }
+
+ public T returnFirst() {
+ return create().getFirst();
+ }
+
+ public List returnAll() {
+ return create();
+ }
+
+ public Set returnSet() {
+ return new HashSet<>(create());
+ }
+
+ protected List create() {
+ var result = new ArrayList();
+
+ for (var index = 0; index < numberOfItems; index++) {
+ result.add(
+ create(index)
+ );
+ }
+
+ return result;
+ }
+
+ protected abstract T create(int index);
+
+ public static String randomKubernetesNameSuffix(String name) {
+ var maxLength = 63; // Kubernetes hard enforces RFC-1123
+
+ if (name.length() > 60) {
+ throw new IllegalArgumentException(
+ "The name is too long (must be <= 60 to allow '-' + at least 2 random chars, max %d total) [name=%s, length=%d]".formatted(
+ maxLength,
+ name,
+ name.length()
+ )
+ );
+ }
+
+ var separator = "-";
+ var suffixLength = maxLength - name.length() - separator.length();
+
+ return name + separator + FAKER.regexify("[a-z0-9]{%d}".formatted(suffixLength));
+ }
+}
diff --git a/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/Given.java b/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/Given.java
new file mode 100644
index 0000000..112bb84
--- /dev/null
+++ b/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/Given.java
@@ -0,0 +1,110 @@
+package it.aboutbits.postgresql._support.testdata.persisted;
+
+import io.fabric8.kubernetes.client.KubernetesClient;
+import it.aboutbits.postgresql._support.testdata.persisted.creator.ClusterConnectionCreate;
+import it.aboutbits.postgresql._support.testdata.persisted.creator.RoleCreate;
+import it.aboutbits.postgresql._support.testdata.persisted.creator.SecretRefCreate;
+import jakarta.enterprise.context.ApplicationScoped;
+import lombok.AccessLevel;
+import lombok.RequiredArgsConstructor;
+import org.eclipse.microprofile.config.inject.ConfigProperty;
+import org.jspecify.annotations.NullMarked;
+
+import java.net.URI;
+
+@NullMarked
+@ApplicationScoped
+@RequiredArgsConstructor
+public class Given {
+ private final KubernetesClient kubernetesClient;
+
+ @SuppressWarnings("NullAway.Init")
+ @ConfigProperty(name = "quarkus.datasource.devservices.username")
+ String username;
+
+ @SuppressWarnings("NullAway.Init")
+ @ConfigProperty(name = "quarkus.datasource.devservices.password")
+ String password;
+
+ @SuppressWarnings("NullAway.Init")
+ @ConfigProperty(name = "quarkus.datasource.jdbc.url")
+ String jdbcUrl;
+
+ DBConnectionDetails dbConnectionDetails() {
+ return new DBConnectionDetails(
+ parsePortFromJdbcUrl(jdbcUrl),
+ username,
+ password
+ );
+ }
+
+ public One one() {
+ return new One(this);
+ }
+
+ public Many many(int numberOfItems) {
+ return new Many(numberOfItems, this);
+ }
+
+ public class One extends Item {
+ One(Given given) {
+ super(1, given);
+ }
+ }
+
+ public class Many extends Item {
+ Many(int numberOfItems, Given given) {
+ super(numberOfItems, given);
+ }
+ }
+
+ @RequiredArgsConstructor(access = AccessLevel.PACKAGE)
+ public abstract class Item {
+ private final int numberOfItems;
+ private final Given given;
+
+ @SuppressWarnings("unused")
+ public Item describedAs(String description) {
+ return this;
+ }
+
+ public SecretRefCreate secretRef() {
+ return new SecretRefCreate(
+ numberOfItems,
+ kubernetesClient
+ );
+ }
+
+ public ClusterConnectionCreate clusterConnection() {
+ return new ClusterConnectionCreate(
+ numberOfItems,
+ given,
+ kubernetesClient,
+ dbConnectionDetails()
+ );
+ }
+
+ public RoleCreate role() {
+ return new RoleCreate(
+ numberOfItems,
+ given,
+ kubernetesClient
+ );
+ }
+ }
+
+ public record DBConnectionDetails(
+ int port,
+ String username,
+ String password
+ ) {
+ }
+
+ private int parsePortFromJdbcUrl(String url) {
+ // Typical format: jdbc:postgresql://localhost:5432/db
+ // We strip "jdbc:" so URI.create can handle the "postgresql://..." part
+ return URI.create(
+ url.replace("jdbc:", "")
+ ).getPort();
+ }
+}
diff --git a/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/creator/ClusterConnectionCreate.java b/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/creator/ClusterConnectionCreate.java
new file mode 100644
index 0000000..824e451
--- /dev/null
+++ b/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/creator/ClusterConnectionCreate.java
@@ -0,0 +1,168 @@
+package it.aboutbits.postgresql._support.testdata.persisted.creator;
+
+import io.fabric8.kubernetes.api.model.ObjectMetaBuilder;
+import io.fabric8.kubernetes.client.KubernetesClient;
+import it.aboutbits.postgresql._support.testdata.base.TestDataCreator;
+import it.aboutbits.postgresql._support.testdata.persisted.Given;
+import it.aboutbits.postgresql.core.SecretRef;
+import it.aboutbits.postgresql.crd.clusterconnection.ClusterConnection;
+import it.aboutbits.postgresql.crd.clusterconnection.ClusterConnectionSpec;
+import lombok.AccessLevel;
+import lombok.Setter;
+import lombok.experimental.Accessors;
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+
+import java.util.Map;
+import java.util.Objects;
+import java.util.concurrent.TimeUnit;
+
+@NullMarked
+@Setter
+@Accessors(fluent = true, chain = true)
+public class ClusterConnectionCreate extends TestDataCreator {
+ private final Given given;
+ private final KubernetesClient kubernetesClient;
+ private final Given.DBConnectionDetails dbConnectionDetails;
+
+ @Nullable
+ private String withNamespace;
+ @Setter(AccessLevel.NONE)
+ private boolean withoutNamespace = false;
+
+ @Nullable
+ private String withName;
+
+ @Nullable
+ private String withHost;
+
+ @Nullable
+ private Integer withPort;
+
+ @Nullable
+ private String withMaintenanceDatabase;
+
+ @Nullable
+ private SecretRef withAdminSecretRef;
+
+ @Nullable
+ private String withApplicationName;
+
+ public ClusterConnectionCreate withoutNamespace() {
+ this.withoutNamespace = true;
+ return this;
+ }
+
+ public ClusterConnectionCreate(
+ int numberOfItems,
+ Given given, KubernetesClient kubernetesClient, Given.DBConnectionDetails dbConnectionDetails
+ ) {
+ super(numberOfItems);
+ this.given = given;
+ this.kubernetesClient = kubernetesClient;
+ this.dbConnectionDetails = dbConnectionDetails;
+ }
+
+ @Override
+ protected ClusterConnection create(int index) {
+ // given
+ var namespace = getNamespace();
+ var name = getName();
+
+ var item = new ClusterConnection();
+
+ item.setMetadata(new ObjectMetaBuilder()
+ .withName(name)
+ .withNamespace(namespace)
+ .build()
+ );
+
+ var spec = new ClusterConnectionSpec();
+ spec.setHost(getHost());
+ spec.setPort(getPort());
+ spec.setMaintenanceDatabase(getMaintenanceDatabase());
+ spec.setAdminSecretRef(getAdminSecretRef());
+ spec.setParameters(getParameters());
+
+ item.setSpec(spec);
+
+ kubernetesClient.resources(ClusterConnection.class)
+ .inNamespace(namespace)
+ .resource(item)
+ .serverSideApply();
+
+ //noinspection ConstantConditions
+ return kubernetesClient.resources(ClusterConnection.class)
+ .inNamespace(namespace)
+ .withName(name)
+ .waitUntilCondition(
+ clusterConnection -> clusterConnection.getStatus() != null,
+ 10,
+ TimeUnit.SECONDS
+ );
+ }
+
+ @Nullable
+ private String getNamespace() {
+ if (withoutNamespace) {
+ return null;
+ }
+
+ if (withNamespace != null) {
+ return withNamespace;
+ }
+
+ return kubernetesClient.getNamespace();
+ }
+
+ private String getName() {
+ if (withName != null) {
+ return withName;
+ }
+
+ return randomKubernetesNameSuffix("test-cluster-connection");
+ }
+
+ private String getHost() {
+ if (withHost != null) {
+ return withHost;
+ }
+
+ return "localhost";
+ }
+
+ private int getPort() {
+ if (withPort != null) {
+ return withPort;
+ }
+
+ return dbConnectionDetails.port();
+ }
+
+ private String getMaintenanceDatabase() {
+ return Objects.requireNonNullElse(
+ withMaintenanceDatabase,
+ "postgres"
+ );
+ }
+
+ private SecretRef getAdminSecretRef() {
+ if (withAdminSecretRef != null) {
+ return withAdminSecretRef;
+ }
+
+ return given.one()
+ .secretRef()
+ .withUsername(dbConnectionDetails.username())
+ .withPassword(dbConnectionDetails.password())
+ .returnFirst();
+ }
+
+ private Map getParameters() {
+ if (withApplicationName != null) {
+ return Map.of("ApplicationName", withApplicationName);
+ }
+
+ return Map.of();
+ }
+}
diff --git a/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/creator/RoleCreate.java b/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/creator/RoleCreate.java
new file mode 100644
index 0000000..a1884e2
--- /dev/null
+++ b/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/creator/RoleCreate.java
@@ -0,0 +1,155 @@
+package it.aboutbits.postgresql._support.testdata.persisted.creator;
+
+import io.fabric8.kubernetes.api.model.ObjectMetaBuilder;
+import io.fabric8.kubernetes.client.KubernetesClient;
+import it.aboutbits.postgresql._support.testdata.base.TestDataCreator;
+import it.aboutbits.postgresql._support.testdata.persisted.Given;
+import it.aboutbits.postgresql.core.ClusterReference;
+import it.aboutbits.postgresql.core.SecretRef;
+import it.aboutbits.postgresql.crd.role.Role;
+import it.aboutbits.postgresql.crd.role.RoleSpec;
+import lombok.AccessLevel;
+import lombok.Setter;
+import lombok.experimental.Accessors;
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+
+import java.util.Objects;
+import java.util.concurrent.TimeUnit;
+
+@NullMarked
+@Setter
+@Accessors(fluent = true, chain = true)
+public class RoleCreate extends TestDataCreator {
+ private final Given given;
+
+ private final KubernetesClient kubernetesClient;
+
+ @Nullable
+ private String withNamespace;
+ @Setter(AccessLevel.NONE)
+ private boolean withoutNamespace = false;
+
+ @Nullable
+ private String withName;
+
+ @Nullable
+ private String withComment;
+
+ @Nullable
+ private String withClusterConnectionName;
+
+ @Nullable
+ private String withClusterConnectionNamespace;
+
+ @Nullable
+ private SecretRef withPasswordSecretRef;
+
+ private RoleSpec.@Nullable Flags withFlags;
+
+ public RoleCreate withLogin(boolean login) {
+ if (!login) {
+ withPasswordSecretRef = null;
+ return this;
+ }
+
+ if (withPasswordSecretRef != null) {
+ return this;
+ }
+
+ withPasswordSecretRef = given.one()
+ .secretRef()
+ .returnFirst();
+
+ return this;
+ }
+
+ public RoleCreate withoutNamespace() {
+ withoutNamespace = true;
+ return this;
+ }
+
+ public RoleCreate(
+ int numberOfItems,
+ Given given,
+ KubernetesClient kubernetesClient
+ ) {
+ super(numberOfItems);
+ this.given = given;
+ this.kubernetesClient = kubernetesClient;
+ }
+
+ @Override
+ protected Role create(int index) {
+ var namespace = getNamespace();
+ var name = getName();
+
+ var item = new Role();
+
+ item.setMetadata(new ObjectMetaBuilder()
+ .withName(name)
+ .withNamespace(namespace)
+ .build()
+ );
+
+ var spec = new RoleSpec();
+ spec.setName(name);
+ spec.setComment(withComment);
+
+ var clusterRef = new ClusterReference();
+ clusterRef.setName(getClusterConnectionName());
+ clusterRef.setNamespace(withClusterConnectionNamespace);
+ spec.setClusterRef(clusterRef);
+
+ spec.setPasswordSecretRef(withPasswordSecretRef);
+
+ if (withFlags != null) {
+ spec.setFlags(withFlags);
+ }
+
+ item.setSpec(spec);
+
+ kubernetesClient.resources(Role.class)
+ .inNamespace(namespace)
+ .resource(item)
+ .serverSideApply();
+
+ //noinspection ConstantConditions
+ return kubernetesClient.resources(Role.class)
+ .inNamespace(namespace)
+ .withName(name)
+ .waitUntilCondition(
+ role -> role.getStatus() != null,
+ 10,
+ TimeUnit.SECONDS
+ );
+ }
+
+ @Nullable
+ private String getNamespace() {
+ if (withoutNamespace) {
+ return null;
+ }
+
+ if (withNamespace != null) {
+ return withNamespace;
+ }
+
+ return kubernetesClient.getNamespace();
+ }
+
+ private String getName() {
+ if (withName != null) {
+ return withName;
+ }
+
+ return randomKubernetesNameSuffix("test-role");
+ }
+
+ private String getClusterConnectionName() {
+ return Objects.requireNonNullElse(
+ withClusterConnectionName,
+ "test-cluster-connection"
+ );
+ }
+}
diff --git a/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/creator/SecretRefCreate.java b/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/creator/SecretRefCreate.java
new file mode 100644
index 0000000..7033537
--- /dev/null
+++ b/src/test/java/it/aboutbits/postgresql/_support/testdata/persisted/creator/SecretRefCreate.java
@@ -0,0 +1,140 @@
+package it.aboutbits.postgresql._support.testdata.persisted.creator;
+
+import io.fabric8.kubernetes.api.model.SecretBuilder;
+import io.fabric8.kubernetes.client.KubernetesClient;
+import it.aboutbits.postgresql._support.testdata.base.TestDataCreator;
+import it.aboutbits.postgresql.core.SecretRef;
+import lombok.AccessLevel;
+import lombok.Setter;
+import lombok.experimental.Accessors;
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+
+import static it.aboutbits.postgresql.core.KubernetesService.SECRET_DATA_BASIC_AUTH_PASSWORD_KEY;
+import static it.aboutbits.postgresql.core.KubernetesService.SECRET_DATA_BASIC_AUTH_USERNAME_KEY;
+import static it.aboutbits.postgresql.core.KubernetesService.SECRET_TYPE_BASIC_AUTH;
+
+@NullMarked
+@Setter
+@Accessors(fluent = true, chain = true)
+public class SecretRefCreate extends TestDataCreator {
+ private final KubernetesClient kubernetesClient;
+
+ @Nullable
+ private String withNamespace;
+ @Setter(AccessLevel.NONE)
+ private boolean withoutNamespace = false;
+
+ @Nullable
+ private String withName;
+
+ @Nullable
+ private String withUsername;
+ @Setter(AccessLevel.NONE)
+ private boolean withoutUsername = false;
+
+ @Nullable
+ private String withPassword;
+ @Setter(AccessLevel.NONE)
+ private boolean withoutPassword = false;
+
+ public SecretRefCreate(
+ int numberOfItems,
+ KubernetesClient kubernetesClient
+ ) {
+ super(numberOfItems);
+ this.kubernetesClient = kubernetesClient;
+ }
+
+ @SuppressWarnings("unused")
+ public SecretRefCreate withoutNamespace() {
+ withoutNamespace = true;
+ return this;
+ }
+
+ @SuppressWarnings("unused")
+ public SecretRefCreate withoutUsername() {
+ withoutUsername = true;
+ return this;
+ }
+
+ @SuppressWarnings("unused")
+ public SecretRefCreate withoutPassword() {
+ withoutPassword = true;
+ return this;
+ }
+
+ @Override
+ protected SecretRef create(int index) {
+ var namespace = getNamespace();
+ var name = getName();
+
+ var secret = new SecretBuilder()
+ .withNewMetadata()
+ .withNamespace(namespace)
+ .withName(name)
+ .endMetadata()
+ .withType(SECRET_TYPE_BASIC_AUTH)
+ .addToStringData(SECRET_DATA_BASIC_AUTH_USERNAME_KEY, getUsername())
+ .addToStringData(SECRET_DATA_BASIC_AUTH_PASSWORD_KEY, getPassword())
+ .build();
+
+ kubernetesClient.secrets()
+ .inNamespace(namespace)
+ .resource(secret)
+ .serverSideApply();
+
+ var secretRef = new SecretRef();
+ secretRef.setName(name);
+ secretRef.setNamespace(namespace);
+
+ return secretRef;
+ }
+
+ @Nullable
+ private String getNamespace() {
+ if (withoutNamespace) {
+ return null;
+ }
+
+ if (withNamespace != null) {
+ return withNamespace;
+ }
+
+ return kubernetesClient.getNamespace();
+ }
+
+ private String getName() {
+ if (withName != null) {
+ return withName;
+ }
+
+ return randomKubernetesNameSuffix("test-secret");
+ }
+
+ @Nullable
+ private String getUsername() {
+ if (withoutUsername) {
+ return null;
+ }
+
+ if (withUsername != null) {
+ return withUsername;
+ }
+
+ return FAKER.credentials().username();
+ }
+
+ @Nullable
+ private String getPassword() {
+ if (withoutPassword) {
+ return null;
+ }
+
+ if (withPassword != null) {
+ return withPassword;
+ }
+
+ return FAKER.credentials().username();
+ }
+}
diff --git a/src/test/java/it/aboutbits/postgresql/core/SQLUtilTest.java b/src/test/java/it/aboutbits/postgresql/core/SQLUtilTest.java
new file mode 100644
index 0000000..917bdc1
--- /dev/null
+++ b/src/test/java/it/aboutbits/postgresql/core/SQLUtilTest.java
@@ -0,0 +1,99 @@
+package it.aboutbits.postgresql.core;
+
+import org.jooq.QueryPart;
+import org.jooq.SQLDialect;
+import org.jooq.impl.DSL;
+import org.jspecify.annotations.NullMarked;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Nested;
+import org.junit.jupiter.api.Test;
+
+import java.util.List;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.jooq.impl.DSL.sql;
+
+@NullMarked
+class SQLUtilTest {
+ @Nested
+ class ConcatenateQueryPartsWithSpaces {
+ @Test
+ @DisplayName("when empty, should return empty string")
+ void whenEmpty_shouldReturnEmptyString() {
+ // given / when
+ var result = SQLUtil.concatenateQueryPartsWithSpaces(List.of());
+
+ // then
+ assertThat(render(result)).isEmpty();
+ }
+
+ @Test
+ @DisplayName("when single item, should return item")
+ void whenSingleItem_shouldReturnItem() {
+ // given
+ var part = sql("item1");
+
+ // when
+ var result = SQLUtil.concatenateQueryPartsWithSpaces(List.of(part));
+
+ // then
+ assertThat(render(result)).isEqualTo("item1");
+ }
+
+ @Test
+ @DisplayName("when multiple items, should join with spaces")
+ void whenMultipleItems_shouldJoinWithSpaces() {
+ // given
+ var parts = List.of(sql("item1"), sql("item2"), sql("item3"));
+
+ // when
+ var result = SQLUtil.concatenateQueryPartsWithSpaces(parts);
+
+ // then
+ assertThat(render(result)).isEqualTo("item1 item2 item3");
+ }
+ }
+
+ @Nested
+ class ConcatenateQueryPartsWithComma {
+ @Test
+ @DisplayName("when empty, should return empty string")
+ void whenEmpty_shouldReturnEmptyString() {
+ // given / when
+ var result = SQLUtil.concatenateQueryPartsWithComma(List.of());
+
+ // then
+ assertThat(render(result)).isEmpty();
+ }
+
+ @Test
+ @DisplayName("when single item, should return item")
+ void whenSingleItem_shouldReturnItem() {
+ // given
+ var part = sql("item1");
+
+ // when
+ var result = SQLUtil.concatenateQueryPartsWithComma(List.of(part));
+
+ // then
+ assertThat(render(result)).isEqualTo("item1");
+ }
+
+ @Test
+ @DisplayName("when multiple items, should join with comma")
+ void whenMultipleItems_shouldJoinWithComma() {
+ // given
+ var parts = List.of(sql("item1"), sql("item2"), sql("item3"));
+
+ // when
+ var result = SQLUtil.concatenateQueryPartsWithComma(parts);
+
+ // then
+ assertThat(render(result)).isEqualTo("item1, item2, item3");
+ }
+ }
+
+ private String render(QueryPart queryPart) {
+ return DSL.using(SQLDialect.POSTGRES).render(queryPart);
+ }
+}
diff --git a/src/test/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionReconcilerErrorTest.java b/src/test/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionReconcilerErrorTest.java
new file mode 100644
index 0000000..60c976f
--- /dev/null
+++ b/src/test/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionReconcilerErrorTest.java
@@ -0,0 +1,109 @@
+package it.aboutbits.postgresql.crd.clusterconnection;
+
+import io.fabric8.kubernetes.api.model.ObjectMeta;
+import io.javaoperatorsdk.operator.api.reconciler.Context;
+import io.quarkus.test.InjectMock;
+import io.quarkus.test.junit.QuarkusTest;
+import it.aboutbits.postgresql.core.PostgreSQLContextFactory;
+import jakarta.inject.Inject;
+import org.jooq.CloseableDSLContext;
+import org.jooq.exception.DataAccessException;
+import org.jspecify.annotations.NullMarked;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+@NullMarked
+@QuarkusTest
+class ClusterConnectionReconcilerErrorTest {
+ @SuppressWarnings("NullAway.Init")
+ @InjectMock
+ PostgreSQLContextFactory contextFactory;
+
+ @Inject
+ ClusterConnectionReconciler reconciler;
+
+ private ClusterConnection resource;
+ private Context context;
+
+ @BeforeEach
+ void setUp() {
+ resource = new ClusterConnection();
+
+ var metadata = new ObjectMeta();
+ metadata.setGeneration(1L);
+
+ // We mock the spec to ensure getName() works safely without throwing NPE,
+ // as the custom getName() implementation in ClusterConnection relies on spec fields.
+ var spec = mock(ClusterConnectionSpec.class);
+
+ when(spec.getHost()).thenReturn("localhost");
+ when(spec.getPort()).thenReturn(5432);
+ when(spec.getMaintenanceDatabase()).thenReturn("postgres");
+ when(spec.getParameters()).thenReturn(Collections.emptyMap());
+
+ resource.setSpec(spec);
+ resource.setMetadata(metadata);
+
+ //noinspection unchecked
+ context = mock(Context.class);
+ }
+
+ @Test
+ @DisplayName("Should handle SQLException during DSL context creation")
+ void reconcile_whenDslCreationFails_shouldReturnErrorStatus() {
+ // given
+ var errorMessage = "Connection refused to database";
+
+ when(contextFactory.getDSLContext(resource)).thenThrow(
+ new RuntimeException(errorMessage)
+ );
+
+ // when
+ var updateControl = reconciler.reconcile(resource, context);
+
+ // then
+ assertThat(updateControl.getResource())
+ .isPresent()
+ .get()
+ .extracting(ClusterConnection::getStatus)
+ .satisfies(status ->
+ assertThat(status.getMessage()).contains(errorMessage)
+ );
+ }
+
+ @Test
+ @DisplayName("Should handle DataAccessException during version check")
+ void reconcile_whenVersionQueryFails_shouldReturnErrorStatus() {
+ // given
+ var errorMessage = "Query execution failed";
+ var dslContext = mock(CloseableDSLContext.class);
+
+ when(contextFactory.getDSLContext(resource)).thenReturn(
+ dslContext
+ );
+
+ when(dslContext.fetchSingle(anyString())).thenThrow(
+ new DataAccessException(errorMessage)
+ );
+
+ // when
+ var updateControl = reconciler.reconcile(resource, context);
+
+ // then
+ assertThat(updateControl.getResource())
+ .isPresent()
+ .get()
+ .extracting(ClusterConnection::getStatus)
+ .satisfies(status ->
+ assertThat(status.getMessage()).contains(errorMessage)
+ );
+ }
+}
diff --git a/src/test/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionReconcilerTest.java b/src/test/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionReconcilerTest.java
new file mode 100644
index 0000000..02babfb
--- /dev/null
+++ b/src/test/java/it/aboutbits/postgresql/crd/clusterconnection/ClusterConnectionReconcilerTest.java
@@ -0,0 +1,102 @@
+package it.aboutbits.postgresql.crd.clusterconnection;
+
+import io.fabric8.kubernetes.client.KubernetesClient;
+import io.quarkus.test.junit.QuarkusTest;
+import it.aboutbits.postgresql._support.testdata.persisted.Given;
+import it.aboutbits.postgresql.core.CRPhase;
+import it.aboutbits.postgresql.core.CRStatus;
+import it.aboutbits.postgresql.core.PostgreSQLContextFactory;
+import lombok.RequiredArgsConstructor;
+import org.jooq.DSLContext;
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+
+import java.time.OffsetDateTime;
+import java.time.ZoneOffset;
+import java.time.temporal.ChronoUnit;
+import java.util.Objects;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatNoException;
+import static org.assertj.core.api.Assertions.within;
+
+@NullMarked
+@QuarkusTest
+@RequiredArgsConstructor
+class ClusterConnectionReconcilerTest {
+ private final Given given;
+
+ private final PostgreSQLContextFactory postgreSQLContextFactory;
+ private final KubernetesClient kubernetesClient;
+
+ @BeforeEach
+ void cleanUp() {
+ kubernetesClient.resources(ClusterConnection.class)
+ .withTimeout(5, TimeUnit.SECONDS)
+ .delete();
+ }
+
+ @Test
+ @DisplayName("When a ClusterConnection is created, the status should be ready")
+ void createsCustomResource_andReconcilerStatusIsReady() {
+ // given / when
+ var customResource = given.one()
+ .clusterConnection()
+ .withName("test-connection")
+ .returnFirst();
+
+ // then
+ AtomicReference<@Nullable DSLContext> dslAtomic = new AtomicReference<>();
+ assertThatNoException().isThrownBy(
+ () -> dslAtomic.set(postgreSQLContextFactory.getDSLContext(customResource))
+ );
+
+ var dsl = Objects.requireNonNull(dslAtomic.get());
+
+ var version = dsl.fetchSingle("select version()").into(String.class);
+
+ var expectedStatus = getInitialClusterConnectionStatus(customResource);
+ expectedStatus.setMessage(version);
+
+ assertThatClusterConnectionHasExpectedStatus(
+ customResource,
+ expectedStatus,
+ OffsetDateTime.now(ZoneOffset.UTC)
+ );
+ }
+
+ private static void assertThatClusterConnectionHasExpectedStatus(
+ ClusterConnection clusterConnection,
+ CRStatus expectedStatus,
+ OffsetDateTime now
+ ) {
+ assertThat(clusterConnection)
+ .isNotNull()
+ .extracting(ClusterConnection::getStatus)
+ .satisfies(status -> {
+ assertThat(status.getLastProbeTime()).isCloseTo(
+ now,
+ within(10, ChronoUnit.SECONDS)
+ );
+ assertThat(status.getLastPhaseTransitionTime()).isCloseTo(
+ now,
+ within(10, ChronoUnit.SECONDS)
+ );
+ })
+ .usingRecursiveComparison()
+ .ignoringFields("lastProbeTime", "lastPhaseTransitionTime")
+ .isEqualTo(expectedStatus);
+ }
+
+ private static CRStatus getInitialClusterConnectionStatus(ClusterConnection clusterConnection) {
+ return new CRStatus()
+ .setName(clusterConnection.getName())
+ .setPhase(CRPhase.READY)
+ .setObservedGeneration(1L);
+ }
+}
diff --git a/src/test/java/it/aboutbits/postgresql/crd/role/RoleReconcilerTest.java b/src/test/java/it/aboutbits/postgresql/crd/role/RoleReconcilerTest.java
new file mode 100644
index 0000000..d88e54a
--- /dev/null
+++ b/src/test/java/it/aboutbits/postgresql/crd/role/RoleReconcilerTest.java
@@ -0,0 +1,961 @@
+package it.aboutbits.postgresql.crd.role;
+
+import io.fabric8.kubernetes.api.model.SecretBuilder;
+import io.fabric8.kubernetes.client.KubernetesClient;
+import io.quarkus.test.junit.QuarkusTest;
+import it.aboutbits.postgresql._support.testdata.persisted.Given;
+import it.aboutbits.postgresql.core.CRPhase;
+import it.aboutbits.postgresql.core.CRStatus;
+import it.aboutbits.postgresql.core.PostgreSQLAuthenticationService;
+import it.aboutbits.postgresql.core.PostgreSQLContextFactory;
+import it.aboutbits.postgresql.core.SecretRef;
+import it.aboutbits.postgresql.crd.clusterconnection.ClusterConnection;
+import lombok.RequiredArgsConstructor;
+import org.jooq.DSLContext;
+import org.jooq.Field;
+import org.jspecify.annotations.NullMarked;
+import org.jspecify.annotations.Nullable;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import java.time.OffsetDateTime;
+import java.time.ZoneOffset;
+import java.time.temporal.ChronoUnit;
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+import java.util.function.BiConsumer;
+import java.util.function.Predicate;
+import java.util.stream.Stream;
+
+import static it.aboutbits.postgresql.core.KubernetesService.SECRET_DATA_BASIC_AUTH_PASSWORD_KEY;
+import static it.aboutbits.postgresql.core.infrastructure.persistence.Tables.PG_AUTHID;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.awaitility.Awaitility.await;
+import static org.jooq.impl.DSL.role;
+
+@NullMarked
+@QuarkusTest
+@RequiredArgsConstructor
+class RoleReconcilerTest {
+ private final Given given;
+
+ private final RoleService roleService;
+ private final PostgreSQLContextFactory postgreSQLContextFactory;
+ private final PostgreSQLAuthenticationService postgreSQLAuthenticationService;
+
+ private final KubernetesClient kubernetesClient;
+
+ @BeforeEach
+ void cleanUp() {
+ kubernetesClient.resources(Role.class)
+ .withTimeout(5, TimeUnit.SECONDS)
+ .delete();
+
+ kubernetesClient.resources(ClusterConnection.class)
+ .withTimeout(5, TimeUnit.SECONDS)
+ .delete();
+ }
+
+ @Test
+ @DisplayName("When a Role (LOGIN) is created, it should be reconciled to READY and present in pg_authid")
+ void createRole_withLogin_andStatusReady() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-connection-role-login")
+ .returnFirst();
+
+ var now = OffsetDateTime.now(ZoneOffset.UTC);
+ var roleName = "test-role-login";
+
+ // when
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .withPasswordSecretRef(clusterConnection.getSpec().getAdminSecretRef())
+ .returnFirst();
+
+ // then: assert READY
+ var expectedStatus = new CRStatus()
+ .setName(roleName)
+ .setPhase(CRPhase.READY)
+ .setObservedGeneration(1L);
+
+ assertThatRoleHasExpectedStatus(
+ role,
+ expectedStatus,
+ now
+ );
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ assertThat(roleService.roleExists(dsl, role.getSpec())).isTrue();
+ assertThat(roleService.roleLoginMatches(dsl, role.getSpec())).isTrue();
+ }
+
+ @Test
+ @DisplayName("When a Role (NOLOGIN) is created, it should be reconciled to READY and present with NOLOGIN")
+ void createRole_withoutLogin_andStatusReady() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-connection-role-nologin")
+ .returnFirst();
+
+ var now = OffsetDateTime.now(ZoneOffset.UTC);
+ var roleName = "test-role-nologin";
+
+ // when
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .returnFirst();
+
+ var expectedStatus = new CRStatus()
+ .setName(roleName)
+ .setPhase(CRPhase.READY)
+ .setObservedGeneration(1L);
+
+ assertThatRoleHasExpectedStatus(
+ role,
+ expectedStatus,
+ now
+ );
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ assertThat(roleService.roleExists(dsl, role.getSpec())).isTrue();
+ assertThat(roleService.roleLoginMatches(dsl, role.getSpec())).isTrue();
+ }
+
+ @Test
+ @DisplayName("When a Role login state is changed, it should be updated correctly in pg_authid")
+ void toggleRoleLogin_updatesCorrectly() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-connection-role-toggle-login")
+ .returnFirst();
+
+ var now = OffsetDateTime.now(ZoneOffset.UTC);
+ var roleName = "test-role-toggle-login";
+
+ // when
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .returnFirst();
+
+ var spec = role.getSpec();
+
+ // then
+ assertThatRoleHasExpectedStatus(
+ role,
+ new CRStatus()
+ .setName(roleName)
+ .setPhase(CRPhase.READY)
+ .setObservedGeneration(1L),
+ now
+ );
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ assertThat(
+ roleService.roleExists(dsl, role.getSpec())
+ ).isTrue();
+
+ assertThat(
+ getRoleFlagValue(dsl, roleName, PG_AUTHID.ROLCANLOGIN)
+ ).isFalse();
+
+ // 2. Add a passwordSecretRef to make it a login role
+ spec.setPasswordSecretRef(clusterConnection.getSpec().getAdminSecretRef());
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == 2L
+ );
+
+ // then
+ assertThat(getRoleFlagValue(dsl, roleName, PG_AUTHID.ROLCANLOGIN)).isTrue();
+
+ // 3. Remove passwordSecretRef again
+ spec.setPasswordSecretRef(null);
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == 3L
+ );
+
+ // then
+ assertThat(getRoleFlagValue(dsl, roleName, PG_AUTHID.ROLCANLOGIN)).isFalse();
+ }
+
+ @Test
+ @DisplayName("When a Role references a missing ClusterConnection, status should be PENDING with a helpful message")
+ void createRole_withMissingClusterConnection_setsPending() {
+ // given
+ var roleName = "test-role-missing-cc";
+ var missingClusterName = "non-existing-cc";
+
+ var now = OffsetDateTime.now(ZoneOffset.UTC);
+
+ var dummySecretRef = new SecretRef();
+ dummySecretRef.setName("dummy");
+
+ // when
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(missingClusterName)
+ .withClusterConnectionNamespace(kubernetesClient.getNamespace())
+ .withPasswordSecretRef(dummySecretRef)
+ .returnFirst();
+
+ // then
+ assertThat(role).isNotNull();
+ assertThat(role.getStatus()).isNotNull();
+
+ assertThat(role.getStatus().getPhase()).isEqualTo(CRPhase.PENDING);
+ assertThat(role.getStatus().getMessage()).startsWith(
+ "The specified ClusterConnection does not exist"
+ );
+ assertThat(role.getStatus().getLastProbeTime()).isAfter(
+ now
+ );
+ assertThat(role.getStatus().getLastPhaseTransitionTime()).isNull();
+ }
+
+ @Test
+ @DisplayName(
+ "When a Role (LOGIN) references a secret and that secret changes, it should trigger a re-reconciliation"
+ )
+ void secretChange_triggersReconciliation() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-connection-role-secret-change")
+ .returnFirst();
+
+ var roleName = "test-role-secret-change";
+
+ var initialPassword = "initial-password";
+ var newPassword = "new-password";
+
+ var secretRef = given.one()
+ .secretRef()
+ .withPassword(initialPassword)
+ .returnFirst();
+
+ var secret = kubernetesClient.secrets()
+ .inNamespace(kubernetesClient.getNamespace())
+ .withName(secretRef.getName())
+ .require();
+
+ // when: create Role
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .withPasswordSecretRef(secretRef)
+ .returnFirst();
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ // then: password should match the initial one
+ // Wait for password to match because reconciliation might take a bit
+ await().atMost(10, TimeUnit.SECONDS)
+ .pollInterval(500, TimeUnit.MILLISECONDS)
+ .until(() -> postgreSQLAuthenticationService.passwordMatches(
+ dsl,
+ role.getSpec(),
+ initialPassword
+ ));
+
+ // when: update secret
+ secret.getMetadata().setManagedFields(null);
+ secret = new SecretBuilder(secret)
+ .addToStringData(SECRET_DATA_BASIC_AUTH_PASSWORD_KEY, newPassword)
+ .build();
+
+ kubernetesClient.secrets()
+ .inNamespace(kubernetesClient.getNamespace())
+ .resource(secret)
+ .serverSideApply();
+
+ // then: password should eventually match the new one
+ await().atMost(10, TimeUnit.SECONDS)
+ .pollInterval(500, TimeUnit.MILLISECONDS)
+ .until(() -> postgreSQLAuthenticationService.passwordMatches(
+ dsl,
+ role.getSpec(),
+ newPassword
+ ));
+ }
+
+ @Test
+ @DisplayName(
+ "When a Role (LOGIN) changes its secret reference, it should trigger a re-reconciliation"
+ )
+ void secretRefChange_triggersReconciliation() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-connection-role-secret-ref-change")
+ .returnFirst();
+
+ var roleName = "test-role-secret-ref-change";
+
+ var initialPassword = "initial-password";
+ var newPassword = "new-password";
+
+ var initialSecretRef = given.one()
+ .secretRef()
+ .withPassword(initialPassword)
+ .returnFirst();
+
+ var newSecretRef = given.one()
+ .secretRef()
+ .withPassword(newPassword)
+ .returnFirst();
+
+ // when: create Role
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .withPasswordSecretRef(initialSecretRef)
+ .returnFirst();
+
+ var spec = role.getSpec();
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ // then: password should match the initial one
+ await().atMost(10, TimeUnit.SECONDS)
+ .pollInterval(500, TimeUnit.MILLISECONDS)
+ .until(() -> postgreSQLAuthenticationService.passwordMatches(
+ dsl,
+ role.getSpec(),
+ initialPassword
+ ));
+
+ // when: update secret reference in the Role
+ spec.setPasswordSecretRef(newSecretRef);
+
+ var updatedRole = applyRole(role);
+
+ // then: password should eventually match the new one
+ await().atMost(10, TimeUnit.SECONDS)
+ .pollInterval(500, TimeUnit.MILLISECONDS)
+ .until(() -> postgreSQLAuthenticationService.passwordMatches(
+ dsl,
+ updatedRole.getSpec(),
+ newPassword
+ ));
+ }
+
+ @Test
+ @DisplayName("When the comment is changed, it should be updated in the database")
+ void comment_updatesCorrectly() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-role-comment")
+ .returnFirst();
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ var roleName = "test-role-comment";
+
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .returnFirst();
+
+ var spec = role.getSpec();
+
+ // 1. Set a comment
+ var comment = "This is a test comment";
+ spec.setComment(comment);
+
+ // when
+ var reconciled = applyRole(role);
+ var initialGeneration = reconciled.getStatus().getObservedGeneration();
+
+ // then
+ assertThat(
+ roleService.fetchCurrentRoleComment(dsl, roleName)
+ ).isEqualTo(comment);
+
+ // 2. Change comment
+ var newComment = "Updated comment";
+ spec.setComment(newComment);
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == initialGeneration + 1
+ );
+
+ // then
+ assertThat(
+ roleService.fetchCurrentRoleComment(dsl, roleName)
+ ).isEqualTo(newComment);
+
+ // 3. Remove comment
+ spec.setComment(null);
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == initialGeneration + 2
+ );
+
+ // then
+ assertThat(
+ roleService.fetchCurrentRoleComment(dsl, roleName)
+ ).isNull();
+ }
+
+ @ParameterizedTest
+ @MethodSource("provideBooleanFlags")
+ @DisplayName("When a boolean Role flag is toggled, it should be updated in the database")
+ void roleFlag_togglesCorrectly(
+ Field field,
+ BiConsumer setter
+ ) {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-role-flags")
+ .returnFirst();
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ var roleName = "test-role-" + field.getName();
+
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .returnFirst();
+
+ var spec = role.getSpec();
+
+ // 1. Enable flag (true)
+ setter.accept(
+ spec.getFlags(),
+ true
+ );
+
+ // when
+ var reconciled = applyRole(role);
+ var initialGeneration = reconciled.getStatus().getObservedGeneration();
+
+ // then
+ assertThat(
+ getRoleFlagValue(
+ dsl,
+ roleName,
+ field
+ )
+ ).isTrue();
+
+ // 2. Disable flag (false)
+ setter.accept(
+ spec.getFlags(),
+ false
+ );
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == initialGeneration + 1
+ );
+
+ // then
+ assertThat(
+ getRoleFlagValue(
+ dsl,
+ roleName,
+ field
+ )
+ ).isFalse();
+ }
+
+ @Test
+ @DisplayName("When the CONNECTION LIMIT is changed, it should be updated in the database")
+ void connectionLimit_updatesCorrectly() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-role-conn-limit")
+ .returnFirst();
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ var roleName = "test-role-conn-limit";
+
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .returnFirst();
+
+ var spec = role.getSpec();
+
+ // 1. Set a connection limit
+ spec.getFlags().setConnectionLimit(10);
+
+ // when
+ var reconciled = applyRole(role);
+ var initialGeneration = reconciled.getStatus().getObservedGeneration();
+
+ // then
+ assertThat(
+ getRoleFlagValue(dsl, roleName, PG_AUTHID.ROLCONNLIMIT)
+ ).isEqualTo(10);
+
+ // 2. Change connection limit
+ spec.getFlags().setConnectionLimit(20);
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == initialGeneration + 1
+ );
+
+ // then
+ assertThat(
+ getRoleFlagValue(dsl, roleName, PG_AUTHID.ROLCONNLIMIT)
+ ).isEqualTo(20);
+
+ // 3. Reset connection limit to -1
+ spec.getFlags().setConnectionLimit(-1);
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == initialGeneration + 2
+ );
+
+ // then
+ assertThat(
+ getRoleFlagValue(dsl, roleName, PG_AUTHID.ROLCONNLIMIT)
+ ).isEqualTo(-1);
+ }
+
+ @Test
+ @DisplayName("When the VALID UNTIL is changed, it should be updated in the database")
+ void validUntil_updatesCorrectly() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-role-valid-until")
+ .returnFirst();
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ var roleName = "test-role-valid-until";
+
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .returnFirst();
+
+ var spec = role.getSpec();
+
+ var expiry = OffsetDateTime.now(ZoneOffset.UTC)
+ .plusDays(1)
+ .truncatedTo(ChronoUnit.SECONDS);
+
+ // 1. Set a valid until date
+ spec.getFlags().setValidUntil(expiry);
+
+ // when
+ var reconciled = applyRole(role);
+ var initialGeneration = reconciled.getStatus().getObservedGeneration();
+
+ var currentFlags = roleService.fetchCurrentFlags(dsl, spec);
+
+ // then
+ assertThat(
+ currentFlags.getValidUntil()
+ ).isEqualTo(expiry);
+
+ // 2. Change valid until date
+ var newExpiry = expiry.plusDays(1);
+ spec.getFlags().setValidUntil(newExpiry);
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == initialGeneration + 1
+ );
+
+ currentFlags = roleService.fetchCurrentFlags(dsl, spec);
+
+ // then
+ assertThat(
+ currentFlags.getValidUntil()
+ ).isEqualTo(newExpiry);
+
+ // 3. Reset valid until to null (infinity)
+ spec.getFlags().setValidUntil(null);
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == initialGeneration + 2
+ );
+
+ currentFlags = roleService.fetchCurrentFlags(dsl, spec);
+
+ // then
+ assertThat(
+ currentFlags.getValidUntil()
+ ).isNull();
+ }
+
+ @Test
+ @DisplayName("When IN ROLE membership is changed, it should be updated in the database")
+ void inRole_updatesCorrectly() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-role-in-role")
+ .returnFirst();
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ var parentRole1 = "parent_role_1";
+ var parentRole2 = "parent_role_2";
+
+ dsl.execute("create role {0}", role(parentRole1));
+ dsl.execute("create role {0}", role(parentRole2));
+
+ var roleName = "test-role-in-role";
+
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .returnFirst();
+
+ var spec = role.getSpec();
+
+ // 1. Add a parent role
+ spec.getFlags().setInRole(
+ List.of(parentRole1)
+ );
+
+ // when
+ var reconciled = applyRole(role);
+ var initialGeneration = reconciled.getStatus().getObservedGeneration();
+
+ // then
+ assertThat(
+ roleService.fetchCurrentFlags(dsl, spec).getInRole()
+ ).containsExactly(parentRole1);
+
+ // 2. Add another parent role and remove the first one
+ spec.getFlags().setInRole(
+ List.of(parentRole2)
+ );
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == initialGeneration + 1
+ );
+
+ // then
+ assertThat(
+ roleService.fetchCurrentFlags(dsl, spec).getInRole()
+ ).containsExactly(parentRole2);
+
+ // 3. Remove all parent roles
+ spec.getFlags().setInRole(
+ List.of()
+ );
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == initialGeneration + 2
+ );
+
+ // then
+ assertThat(
+ roleService.fetchCurrentFlags(dsl, spec).getInRole()
+ ).isEmpty();
+
+ // cleanup
+ dsl.execute("drop role if exists {0}", role(parentRole1));
+ dsl.execute("drop role if exists {0}", role(parentRole2));
+ }
+
+ @Test
+ @DisplayName("When ROLE membership is changed, it should be updated in the database")
+ void role_updatesCorrectly() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-role-role")
+ .returnFirst();
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ var memberRole1 = "member_role_1";
+ var memberRole2 = "member_role_2";
+
+ dsl.execute("create role {0}", role(memberRole1));
+ dsl.execute("create role {0}", role(memberRole2));
+
+ var roleName = "test-role-role";
+
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .returnFirst();
+
+ var spec = role.getSpec();
+
+ // 1. Add a member role
+ spec.getFlags().setRole(
+ List.of(memberRole1)
+ );
+
+ // when
+ var reconciled = applyRole(role);
+ var initialGeneration = reconciled.getStatus().getObservedGeneration();
+
+ // then
+ assertThat(
+ roleService.fetchCurrentFlags(dsl, spec).getRole()
+ ).containsExactly(memberRole1);
+
+ // 2. Add another member role and remove the first one
+ spec.getFlags().setRole(
+ List.of(memberRole2)
+ );
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == initialGeneration + 1
+ );
+
+ // then
+ assertThat(
+ roleService.fetchCurrentFlags(dsl, spec).getRole()
+ ).containsExactly(memberRole2);
+
+ // 3. Remove all member roles
+ spec.getFlags().setRole(
+ List.of()
+ );
+
+ // when
+ applyRole(
+ role,
+ r -> r.getStatus().getObservedGeneration() == initialGeneration + 2
+ );
+
+ // then
+ assertThat(
+ roleService.fetchCurrentFlags(dsl, spec).getRole()
+ ).isEmpty();
+
+ // cleanup
+ dsl.execute("drop role if exists {0}", role(memberRole1));
+ dsl.execute("drop role if exists {0}", role(memberRole2));
+ }
+
+ @Test
+ @DisplayName("When multiple ROLE memberships are added, they should be sorted and updated correctly")
+ void role_multipleMemberships_updatesCorrectly() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-role-multiple")
+ .returnFirst();
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ var roleA = "role_a";
+ var roleB = "role_b";
+ var roleC = "role_c";
+
+ dsl.execute("create role {0}", role(roleA));
+ dsl.execute("create role {0}", role(roleB));
+ dsl.execute("create role {0}", role(roleC));
+
+ var roleName = "test-role-multiple";
+
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .returnFirst();
+
+ var spec = role.getSpec();
+
+ // Add multiple roles out of order
+ spec.getFlags().setInRole(
+ List.of(roleC, roleA, roleB)
+ );
+
+ // when
+ applyRole(role);
+
+ // then
+ assertThat(
+ roleService.fetchCurrentFlags(dsl, spec).getInRole()
+ ).containsExactly(roleA, roleB, roleC);
+
+ // cleanup
+ dsl.execute("drop role if exists {0}", role(roleA));
+ dsl.execute("drop role if exists {0}", role(roleB));
+ dsl.execute("drop role if exists {0}", role(roleC));
+ }
+
+ @Test
+ @DisplayName("When a Role is deleted, it should be dropped from the database")
+ void deleteRole_removesFromDatabase() {
+ // given
+ var clusterConnection = given.one()
+ .clusterConnection()
+ .withName("test-connection-role-delete")
+ .returnFirst();
+
+ var roleName = "test-role-delete";
+
+ var role = given.one()
+ .role()
+ .withName(roleName)
+ .withClusterConnectionName(clusterConnection.getMetadata().getName())
+ .returnFirst();
+
+ var dsl = postgreSQLContextFactory.getDSLContext(clusterConnection);
+
+ // Verify it exists initially
+ assertThat(roleService.roleExists(dsl, role.getSpec())).isTrue();
+
+ // when
+ kubernetesClient.resources(Role.class)
+ .inNamespace(role.getMetadata().getNamespace())
+ .withName(role.getMetadata().getName())
+ .withTimeout(5, TimeUnit.SECONDS)
+ .delete();
+
+ // then
+ await().atMost(10, TimeUnit.SECONDS)
+ .pollInterval(500, TimeUnit.MILLISECONDS)
+ .until(() -> !roleService.roleExists(dsl, role.getSpec()));
+ }
+
+ private @Nullable T getRoleFlagValue(
+ DSLContext dsl,
+ String roleName,
+ Field field
+ ) {
+ return dsl.select(field)
+ .from(PG_AUTHID)
+ .where(PG_AUTHID.ROLNAME.eq(roleName))
+ .fetchSingle(field);
+ }
+
+ private static Stream provideBooleanFlags() {
+ return Stream.of(
+ Arguments.of(PG_AUTHID.ROLSUPER, (BiConsumer) RoleSpec.Flags::setSuperuser),
+ Arguments.of(PG_AUTHID.ROLCREATEDB, (BiConsumer) RoleSpec.Flags::setCreatedb),
+ Arguments.of(PG_AUTHID.ROLCREATEROLE, (BiConsumer) RoleSpec.Flags::setCreaterole),
+ Arguments.of(PG_AUTHID.ROLINHERIT, (BiConsumer) RoleSpec.Flags::setInherit),
+ Arguments.of(PG_AUTHID.ROLREPLICATION, (BiConsumer) RoleSpec.Flags::setReplication),
+ Arguments.of(PG_AUTHID.ROLBYPASSRLS, (BiConsumer) RoleSpec.Flags::setBypassrls)
+ );
+ }
+
+ private Role applyRole(Role role) {
+ var namespace = kubernetesClient.getNamespace();
+
+ role.getMetadata().setManagedFields(null);
+ role.getMetadata().setResourceVersion(null);
+
+ var applied = kubernetesClient.resources(Role.class)
+ .inNamespace(namespace)
+ .resource(role)
+ .serverSideApply();
+
+ var generation = applied.getMetadata().getGeneration();
+
+ //noinspection ConstantConditions
+ return kubernetesClient.resources(Role.class)
+ .inNamespace(namespace)
+ .withName(applied.getMetadata().getName())
+ .waitUntilCondition(
+ r -> r.getStatus() != null && r.getStatus().getObservedGeneration() >= generation,
+ 10,
+ TimeUnit.SECONDS
+ );
+ }
+
+ private Role applyRole(
+ Role role,
+ Predicate condition
+ ) {
+ var namespace = kubernetesClient.getNamespace();
+
+ role.getMetadata().setManagedFields(null);
+ role.getMetadata().setResourceVersion(null);
+
+ var applied = kubernetesClient.resources(Role.class)
+ .inNamespace(namespace)
+ .resource(role)
+ .serverSideApply();
+
+ return kubernetesClient.resources(Role.class)
+ .inNamespace(namespace)
+ .withName(applied.getMetadata().getName())
+ .waitUntilCondition(
+ condition,
+ 10,
+ TimeUnit.SECONDS
+ );
+ }
+
+ private static void assertThatRoleHasExpectedStatus(
+ Role role,
+ CRStatus expectedStatus,
+ OffsetDateTime now
+ ) {
+ assertThat(role)
+ .isNotNull()
+ .extracting(Role::getStatus)
+ .satisfies(status -> {
+ assertThat(status.getLastProbeTime()).isAfter(
+ now
+ );
+ assertThat(status.getLastPhaseTransitionTime()).isAfter(
+ now
+ );
+ })
+ .usingRecursiveComparison()
+ .ignoringFields("lastProbeTime", "lastPhaseTransitionTime")
+ .isEqualTo(expectedStatus);
+ }
+}