Commit a46fb321 authored by Vasilis Katopodis's avatar Vasilis Katopodis
Browse files

fix: harden kubeconfig retrieval and registry secrets listeners against lost/silent failures

parent 35ceebfc
Loading
Loading
Loading
Loading
Loading
+27 −21
Original line number Diff line number Diff line
@@ -6,6 +6,7 @@ import static org.etsi.osl.hypo.core.common.constants.SonataKafkaTriggersConstan
import static org.etsi.osl.hypo.registry.config.ApplicationProperties.*;

import io.smallrye.common.annotation.Blocking;
import io.smallrye.mutiny.Uni;
import io.smallrye.mutiny.infrastructure.Infrastructure;
import jakarta.enterprise.context.ApplicationScoped;
import java.util.concurrent.CompletionStage;
@@ -41,30 +42,33 @@ public class RegistrySecretsKafkaListener {
    public CompletionStage<Void> receiveRegistrySecretsMessage(Message<KafkaRequestPayload> message) {
        KafkaRequestPayload payload = message.getPayload();

        CompletionStage<Void> ack = message.ack();
        if (payload == null) {
            logger.warn(
                    "Received empty payload for Registry Secrets channel. Sending deserialization failure.");
            registryKafkaProducer.sendDeserializationFailure();
            return ack;
            return message.ack();
        }

        logger.infof("Kafka request received on Registry Secrets channel. Type: %s", payload.getType());

        registryService
        return registryService
                .handleReceiveRegistryKafkaMessage(message)
                .runSubscriptionOn(Infrastructure.getDefaultWorkerPool()) // worker thread
                .subscribe()
                .with(
                        unused -> logger.infof(THE_HANDLING_OF_KAFKA_MESSAGE_S_WAS_SUCCESSFUL, message),
                        failure ->
                .onItemOrFailure()
                .transformToUni(
                        (unused, failure) -> {
                            if (failure != null) {
                                logger.errorf(
                                        failure,
                                        "Kafka Handling failed for message %s with error: %s",
                                        message,
                                        failure.getMessage()));

        return ack;
                                        failure.getMessage());
                                return Uni.createFrom().completionStage(message.nack(failure));
                            }
                            logger.infof(THE_HANDLING_OF_KAFKA_MESSAGE_S_WAS_SUCCESSFUL, message);
                            return Uni.createFrom().completionStage(message.ack());
                        })
                .subscribeAsCompletionStage();
    }

    @Incoming(REGISTRY_STORE_KUBE_CONF_CHANNEL)
@@ -96,30 +100,32 @@ public class RegistrySecretsKafkaListener {
    @Incoming(REGISTRY_RETRIEVE_KUBE_CONF_CHANNEL)
    @Blocking
    public CompletionStage<Void> retrieveKubernetesConfig(Message<String> message) {
        CompletionStage<Void> ack = message.ack();

        String payload = message.getPayload();
        if (payload == null) {
            logger.warn("Received empty payload for KubeConfig retrieval channel. Ignoring message.");
            return ack;
            return message.ack();
        }

        logger.infof("Processing KubeConfig Retrieval request for ID: %s", payload);

        registryService
        return registryService
                .retrieveKubernetesConfigFromKafkaMessage(message)
                .runSubscriptionOn(Infrastructure.getDefaultWorkerPool()) // worker thread
                .subscribe()
                .with(
                        unused -> logger.infof(THE_HANDLING_OF_KAFKA_MESSAGE_S_WAS_SUCCESSFUL, message),
                        failure ->
                .onItemOrFailure()
                .transformToUni(
                        (unused, failure) -> {
                            if (failure != null) {
                                logger.errorf(
                                        failure,
                                        AlarmLogs.FAILED_TO_RETRIEVE_KUBERNETES_CONFIG_FOR_ID + ": %s. Reason: %s",
                                        payload,
                                        failure.getMessage()));

        return ack;
                                        failure.getMessage());
                                return Uni.createFrom().completionStage(message.nack(failure));
                            }
                            logger.infof(THE_HANDLING_OF_KAFKA_MESSAGE_S_WAS_SUCCESSFUL, message);
                            return Uni.createFrom().completionStage(message.ack());
                        })
                .subscribeAsCompletionStage();
    }

    @Incoming(CREATE_GROUP_POLICY_CHANNEL)
+2 −1
Original line number Diff line number Diff line
package org.etsi.osl.hypo.registry.model.registry;

import com.fasterxml.jackson.annotation.JsonProperty;
import java.util.Map;
import lombok.Data;

@Data
@@ -10,5 +11,5 @@ public class RegistryData {
    private String id;

    @JsonProperty("secretData")
    private Object secretData;
    private Map<String, Object> secretData;
}
+47 −7
Original line number Diff line number Diff line
@@ -8,6 +8,7 @@ import static org.etsi.osl.hypo.registry.config.ApplicationProperties.UNKNOWN;
import io.smallrye.mutiny.Uni;
import io.smallrye.reactive.messaging.kafka.api.IncomingKafkaRecordMetadata;
import jakarta.enterprise.context.ApplicationScoped;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
@@ -251,6 +252,10 @@ public class RegistryService {
            Message<String> message) {
        String serviceId = message.getPayload();

        if (StringUtils.isBlank(serviceId)) {
            throw new IllegalArgumentException("KubeConfig retrieval id must not be blank");
        }

        logger.infof("Processing KubeConfig Retrieval request for ServiceID: %s", serviceId);
        String keycloakJwt = KafkaUtils.extractTokenFromKafkaMessage(message);

@@ -269,11 +274,19 @@ public class RegistryService {

        if (secretData.isEmpty()) {
            logger.warnf("No secret data found in Vault for KubeConfig serviceId: %s", serviceId);

            RegistryData notFoundData = new RegistryData();
            notFoundData.setId(resolvedPath);
            return buildKubernetesConfigSecretResponseMessage(
                    notFoundData,
                    serviceId,
                    Status.FAILURE,
                    "No KubeConfig secret found for ID: " + serviceId);
        }

        RegistryData registryData = new RegistryData();
        registryData.setId(resolvedPath);
        registryData.setSecretData(secretData.orElse(null));
        registryData.setSecretData(secretData.get());

        return buildKubernetesConfigSecretResponseMessage(registryData, serviceId, Status.SUCCESS, "");
    }
@@ -283,16 +296,21 @@ public class RegistryService {
        return Uni.createFrom().voidItem().invoke(() -> handleKubernetesConfigKafkaMessage(message));
    }

    private static final Duration KUBE_CONFIG_RETRIEVAL_RETRY_BACKOFF = Duration.ofMillis(200);
    private static final int KUBE_CONFIG_RETRIEVAL_MAX_RETRIES = 2;

    public Uni<Void> retrieveKubernetesConfigFromKafkaMessage(Message<String> message) {
        String serviceId = message.getPayload();
        Headers headers = extractKafkaHeaders(message);

        return Uni.createFrom()
                .voidItem()
                .item(() -> getKubernetesConfigSecretResponseMessage(message))
                .onFailure(RegistryService::isTransientVaultFailure)
                .retry()
                .withBackOff(KUBE_CONFIG_RETRIEVAL_RETRY_BACKOFF)
                .atMost(KUBE_CONFIG_RETRIEVAL_MAX_RETRIES)
                .invoke(
                        () -> {
                            VaultSecretResponseMessage response =
                                    getKubernetesConfigSecretResponseMessage(message);
                        response -> {
                            logger.debugf(
                                    "Successfully retrieved KubeConfig for ID: %s. Sending response.", serviceId);
                            registryKafkaProducer.sendKubeConfigResponseMessage(response, headers);
@@ -304,7 +322,16 @@ public class RegistryService {
                                    buildKubernetesConfigSecretResponseMessage(
                                            null, serviceId, Status.FAILURE, failure.getMessage());
                            registryKafkaProducer.sendKubeConfigResponseMessage(failureResponse, headers);
                        });
                        })
                .replaceWithVoid();
    }

    // Retries only failures that look transient (a Vault server-side error). Not-found and
    // validation failures are deterministic and already turned into a Status.FAILURE response by
    // getKubernetesConfigSecretResponseMessage, so retrying them would be wasted work.
    private static boolean isTransientVaultFailure(Throwable failure) {
        return failure instanceof VaultException vaultException
                && vaultException.getStatusCode() >= 500;
    }

    public Uni<Void> persistFabricCredentialsSecretFromKafkaMessage(
@@ -317,8 +344,21 @@ public class RegistryService {
        return Uni.createFrom().voidItem().invoke(() -> deleteFabricCredentials(message));
    }

    private static final Duration REGISTRY_SECRETS_MESSAGE_TIMEOUT = Duration.ofSeconds(10);

    // Note: handleReceiveRegistry runs synchronously on the worker thread this Uni is subscribed
    // on, so a timeout here stops the listener from waiting on it (letting it nack for redelivery)
    // but cannot forcibly cancel the still-running call — handleReceiveRegistry already catches
    // its own failures and calls sendFailure/sendResponse internally, so this is a deadline on how
    // long we wait, not a guarantee the underlying work stops.
    public Uni<Void> handleReceiveRegistryKafkaMessage(Message<KafkaRequestPayload> message) {
        return Uni.createFrom().voidItem().invoke(() -> handleReceiveRegistry(message));
        return Uni.createFrom()
                .voidItem()
                .invoke(() -> handleReceiveRegistry(message))
                .ifNoItem()
                .after(REGISTRY_SECRETS_MESSAGE_TIMEOUT)
                .failWith(
                        () -> new TimeoutException("Timed out handling Registry Secrets message: " + message));
    }

    public Uni<Void> deleteOrganizationCredentialsSecretFromKafkaMessage(
+4 −0
Original line number Diff line number Diff line
@@ -19,6 +19,10 @@ kubernetes.config.vault.secret-kv-path=hypo/kubernetes-config
kubernetes.namespace-sa.vault.secret-kv-path=hypo/kubernetes-sa-config
service.order.vault.secret-kv-path=hypo/service-order-data

# Vault REST client resilience: bounds how long a hung/slow Vault can block a Kafka worker thread.
quarkus.rest-client.vault-api.connect-timeout=5000
quarkus.rest-client.vault-api.read-timeout=10000

#quarkus.log.category."org.jboss.resteasy.reactive.client.logging".level=DEBUG

quarkus.kafka.devservices.image-name=labs.etsi.org:5050/osl/hypo/code/org.etsi.osl.hypo.ops/cicd/integration.tests/vectorized/redpanda:v24.1.2
+70 −0
Original line number Diff line number Diff line
@@ -439,6 +439,76 @@ public class RegistrySecretsKafkaListenerTest {
                        organizationAuthCredentialsMessage);
    }

    @Test
    void retrieveKubernetesConfig_acksOnlyAfterSuccessfulCompletion() {
        Message<String> message = mock(Message.class);
        when(message.getPayload()).thenReturn("my-service-id");
        when(message.ack()).thenReturn(CompletableFuture.completedStage(null));

        when(registryService.retrieveKubernetesConfigFromKafkaMessage(message))
                .thenReturn(Uni.createFrom().voidItem());

        CompletionStage<Void> result = registrySecretsKafkaListener.retrieveKubernetesConfig(message);

        verify(message, times(1)).ack();
        verify(message, times(0)).nack(any());
        assertThat(result).isCompleted();
    }

    @Test
    void retrieveKubernetesConfig_nacksOnFailureInsteadOfAcking() {
        Message<String> message = mock(Message.class);
        when(message.getPayload()).thenReturn("my-service-id");
        when(message.nack(any())).thenReturn(CompletableFuture.completedStage(null));

        RuntimeException failure = new RuntimeException("Vault unreachable");
        when(registryService.retrieveKubernetesConfigFromKafkaMessage(message))
                .thenReturn(Uni.createFrom().failure(failure));

        CompletionStage<Void> result = registrySecretsKafkaListener.retrieveKubernetesConfig(message);

        verify(message, times(0)).ack();
        verify(message, times(1)).nack(failure);
        assertThat(result).isCompleted();
    }

    @Test
    void receiveRegistrySecretsMessage_acksOnlyAfterSuccessfulCompletion() {
        KubeConfigRequest payload = new KubeConfigRequest();
        Message<KafkaRequestPayload> message = mock(Message.class);
        when(message.getPayload()).thenReturn(payload);
        when(message.ack()).thenReturn(CompletableFuture.completedStage(null));

        when(registryService.handleReceiveRegistryKafkaMessage(message))
                .thenReturn(Uni.createFrom().voidItem());

        CompletionStage<Void> result =
                registrySecretsKafkaListener.receiveRegistrySecretsMessage(message);

        verify(message, times(1)).ack();
        verify(message, times(0)).nack(any());
        assertThat(result).isCompleted();
    }

    @Test
    void receiveRegistrySecretsMessage_nacksOnFailureInsteadOfAcking() {
        KubeConfigRequest payload = new KubeConfigRequest();
        Message<KafkaRequestPayload> message = mock(Message.class);
        when(message.getPayload()).thenReturn(payload);
        when(message.nack(any())).thenReturn(CompletableFuture.completedStage(null));

        RuntimeException failure = new RuntimeException("boom");
        when(registryService.handleReceiveRegistryKafkaMessage(message))
                .thenReturn(Uni.createFrom().failure(failure));

        CompletionStage<Void> result =
                registrySecretsKafkaListener.receiveRegistrySecretsMessage(message);

        verify(message, times(0)).ack();
        verify(message, times(1)).nack(failure);
        assertThat(result).isCompleted();
    }

    @Test
    void receiveRegistrySecretsMessage() {
        KubeConfigRequest kubeConfigRequest = new KubeConfigRequest();
Loading