Commit 9d7a171d authored by Anastasios Poimenidis's avatar Anastasios Poimenidis
Browse files

fix: send correct termination messages to Sonata

parent 87ff4bea
Loading
Loading
Loading
Loading
+45 −36
Original line number Diff line number Diff line
@@ -11,7 +11,6 @@ import static org.etsi.osl.hypo.core.common.constants.SonataKafkaTriggersConstan
import io.micrometer.common.util.StringUtils;
import io.smallrye.reactive.messaging.kafka.api.OutgoingKafkaRecordMetadata;
import jakarta.enterprise.context.ApplicationScoped;
import jakarta.inject.Inject;
import org.eclipse.microprofile.reactive.messaging.Channel;
import org.eclipse.microprofile.reactive.messaging.Emitter;
import org.eclipse.microprofile.reactive.messaging.Message;
@@ -33,35 +32,39 @@ public class ServiceInventoryProducer {
            CLOUD_EVENT_CONTAINING_SONATA_SERVICE_TERMINATION_MESSAGE_S_HAS_BEEN_SENT_TO_KAFKA_TOPIC_S =
                    "Cloud Event containing Sonata service Termination message \n %s \n has been sent to kafka topic: %s";

    @Inject
    @Channel(TMF_SERVICE_ACTIVATION)
    Emitter<CloudEventServiceActivationState> cloudEventServiceActivationStateEmitter;
    private final Emitter<CloudEventServiceActivationState> cloudEventServiceActivationStateEmitter;
    private final Emitter<CloudEventServiceTerminationData> cloudEventServiceTerminationEmitter;
    private final Emitter<CloudEventServiceTerminationData> cloudEventK8sServiceTerminationEmitter;
    private final Emitter<CloudEventServiceTerminationData> ossServiceTermationEmitter;
    private final Emitter<CloudEventServiceOrder> cloudEventServiceUpdateEmitter;
    private final Emitter<CloudEventPlatformServiceUpdate> cloudEventPlatformServiceUpdateEmitter;
    private final Emitter<CloudEventKubernetesRegistration> cloudEventKubernetesRegistrationEmitter;

    @Inject
    @Channel(TMF_END_USER_SERVICE_TERMINATION)
    Emitter<CloudEventServiceTerminationData> cloudEventServiceTerminationEmitter;
    private static final Logger logger = Logger.getLogger(ServiceInventoryProducer.class);

    @Inject
    public ServiceInventoryProducer(
            @Channel(TMF_SERVICE_ACTIVATION)
                    Emitter<CloudEventServiceActivationState> cloudEventServiceActivationStateEmitter,
            @Channel(TMF_END_USER_SERVICE_TERMINATION)
                    Emitter<CloudEventServiceTerminationData> cloudEventServiceTerminationEmitter,
            @Channel(TMF_KUBERNETES_SERVICE_TERMINATION)
    Emitter<CloudEventServiceTerminationData> cloudEventServiceTerminationForKubernetesEmitter;

    @Inject
                    Emitter<CloudEventServiceTerminationData> cloudEventK8sServiceTerminationEmitter,
            @Channel(TMF_OSS_SERVICE_TERMINATION)
    Emitter<CloudEventServiceTerminationData> ossServiceTermationEmitter;

    @Inject
                    Emitter<CloudEventServiceTerminationData> ossServiceTermationEmitter,
            @Channel(SONATA_SERVICE_UPDATE)
    Emitter<CloudEventServiceOrder> cloudEventServiceUpdateEmitter;

    @Inject
                    Emitter<CloudEventServiceOrder> cloudEventServiceUpdateEmitter,
            @Channel(TMF_PLATFORM_SERVICE_UPDATE)
    Emitter<CloudEventPlatformServiceUpdate> cloudEventKubernetesServiceUpdateEmitter;

    @Inject
                    Emitter<CloudEventPlatformServiceUpdate> cloudEventPlatformServiceUpdateEmitter,
            @Channel(SONATA_REGISTER_KUBERNETES)
    Emitter<CloudEventKubernetesRegistration> cloudEventKubernetesRegistrationEmitter;

    private static final Logger logger = Logger.getLogger(ServiceInventoryProducer.class);
                    Emitter<CloudEventKubernetesRegistration> cloudEventKubernetesRegistrationEmitter) {
        this.cloudEventServiceActivationStateEmitter = cloudEventServiceActivationStateEmitter;
        this.cloudEventServiceTerminationEmitter = cloudEventServiceTerminationEmitter;
        this.cloudEventK8sServiceTerminationEmitter = cloudEventK8sServiceTerminationEmitter;
        this.ossServiceTermationEmitter = ossServiceTermationEmitter;
        this.cloudEventServiceUpdateEmitter = cloudEventServiceUpdateEmitter;
        this.cloudEventPlatformServiceUpdateEmitter = cloudEventPlatformServiceUpdateEmitter;
        this.cloudEventKubernetesRegistrationEmitter = cloudEventKubernetesRegistrationEmitter;
    }

    public void sendServiceActivationMessage(
            CloudEventServiceActivationState cloudEventServiceActivationState, String authHeader) {
@@ -98,15 +101,13 @@ public class ServiceInventoryProducer {

    public void sendKubernetesServiceTerminationMessage(
            CloudEventServiceTerminationData terminationMessage, String authHeader) {
        cloudEventServiceTerminationForKubernetesEmitter.send(terminationMessage);
        if (StringUtils.isNotBlank(authHeader)) {
            OutgoingKafkaRecordMetadata<String> metadata = Helpers.getMetadata(authHeader);

            Message<CloudEventServiceTerminationData> message =
                    Message.of(terminationMessage).addMetadata(metadata);
            cloudEventServiceTerminationForKubernetesEmitter.send(message);
            cloudEventK8sServiceTerminationEmitter.send(message);
        } else {
            cloudEventServiceTerminationForKubernetesEmitter.send(terminationMessage);
            cloudEventK8sServiceTerminationEmitter.send(terminationMessage);
        }
        logger.infof(
                CLOUD_EVENT_CONTAINING_SONATA_SERVICE_TERMINATION_MESSAGE_S_HAS_BEEN_SENT_TO_KAFKA_TOPIC_S,
@@ -114,11 +115,19 @@ public class ServiceInventoryProducer {
                TMF_KUBERNETES_SERVICE_TERMINATION);
    }

    public void sendOssServiceTerminationMessage(CloudEventServiceTerminationData message) {
    public void sendOssServiceTerminationMessage(
            CloudEventServiceTerminationData terminationMessage, String authHeader) {
        if (StringUtils.isNotBlank(authHeader)) {
            OutgoingKafkaRecordMetadata<String> metadata = Helpers.getMetadata(authHeader);
            Message<CloudEventServiceTerminationData> message =
                    Message.of(terminationMessage).addMetadata(metadata);
            ossServiceTermationEmitter.send(message);
        } else {
            ossServiceTermationEmitter.send(terminationMessage);
        }
        logger.infof(
                CLOUD_EVENT_CONTAINING_SONATA_SERVICE_TERMINATION_MESSAGE_S_HAS_BEEN_SENT_TO_KAFKA_TOPIC_S,
                message.toStringAsJson(),
                terminationMessage.toStringAsJson(),
                TMF_OSS_SERVICE_TERMINATION);
    }

@@ -143,7 +152,7 @@ public class ServiceInventoryProducer {

    public void sendPlatformServiceUpdateMessage(
            CloudEventPlatformServiceUpdate cloudEventPlatformServiceUpdate) {
        cloudEventKubernetesServiceUpdateEmitter.send(cloudEventPlatformServiceUpdate);
        cloudEventPlatformServiceUpdateEmitter.send(cloudEventPlatformServiceUpdate);
        logger.infof(
                CLOUD_EVENT_CONTAINING_SONATA_SERVICE_UPDATE_MESSAGE_S_HAS_BEEN_SENT_TO_KAFKA_TOPIC_S,
                cloudEventPlatformServiceUpdate.toStringAsJson(),
+1 −1
Original line number Diff line number Diff line
@@ -170,7 +170,7 @@ public class KafkaMessagesWrapper {
            case TMF_KUBERNETES_SERVICE_TERMINATION -> serviceInventoryProducer
                    .sendKubernetesServiceTerminationMessage(message, authHeader);
            case TMF_OSS_SERVICE_TERMINATION -> serviceInventoryProducer.sendOssServiceTerminationMessage(
                    message);
                    message, authHeader);
            default -> logger.errorf(
                    "No topic was determined for termination of service with id: %s", serviceId);
        }
+1 −1
Original line number Diff line number Diff line
@@ -1451,7 +1451,7 @@ class ServiceInventoryServiceTest {

        if (sonataTrigger.equals(TMF_OSS_SERVICE_TERMINATION)) {
            Mockito.verify(serviceInventoryProducer, Mockito.times(1))
                    .sendOssServiceTerminationMessage(argumentCaptor.capture());
                    .sendOssServiceTerminationMessage(argumentCaptor.capture(), authHeaderCaptor.capture());
        }

        Mockito.verify(serviceOrderAndServiceRepositories, Mockito.times(1)).storeServiceBackup(any());