Loading src/main/java/org/etsi/osl/hypo/registry/kafka/outgoing/RegistryKafkaProducer.java +13 −0 Original line number Diff line number Diff line Loading @@ -86,6 +86,19 @@ public class RegistryKafkaProducer { logger.debugf(ApplicationProperties.PAYLOAD_S, responseMessage.toStringAsJson()); } public void sendKubeConfigResponseMessage( VaultSecretResponseMessage responseMessage, Headers headers) { OutgoingKafkaRecordMetadata<Object> kafkaMetadata = OutgoingKafkaRecordMetadata.builder().withHeaders(headers).build(); vaultSecretResponseMessageEmitter.send(Message.of(responseMessage, Metadata.of(kafkaMetadata))); logger.debugf( "Kube Conf Response data message with headers sent to KAFKA topic: %s", REGISTRY_RESPONSE_KUBE_CONF_CHANNEL); logger.debugf(ApplicationProperties.PAYLOAD_S, responseMessage.toStringAsJson()); } public void sendOrganizationCredentialsResponseMessage( CloudEventRegistryOrganizationCredentialsResult responseMessage) { organizationCredentialsResultEmitter.send(responseMessage); Loading src/main/java/org/etsi/osl/hypo/registry/service/RegistryService.java +13 −3 Original line number Diff line number Diff line Loading @@ -284,6 +284,9 @@ public class RegistryService { } public Uni<Void> retrieveKubernetesConfigFromKafkaMessage(Message<String> message) { String serviceId = message.getPayload(); Headers headers = extractKafkaHeaders(message); return Uni.createFrom() .voidItem() .invoke( Loading @@ -291,9 +294,16 @@ public class RegistryService { VaultSecretResponseMessage response = getKubernetesConfigSecretResponseMessage(message); logger.debugf( "Successfully retrieved KubeConfig for ID: %s. Sending response.", message.getPayload()); registryKafkaProducer.sendKubeConfigResponseMessage(response); "Successfully retrieved KubeConfig for ID: %s. Sending response.", serviceId); registryKafkaProducer.sendKubeConfigResponseMessage(response, headers); }) .onFailure() .invoke( failure -> { VaultSecretResponseMessage failureResponse = buildKubernetesConfigSecretResponseMessage( null, serviceId, Status.FAILURE, failure.getMessage()); registryKafkaProducer.sendKubeConfigResponseMessage(failureResponse, headers); }); } Loading src/main/resources/application.properties +2 −0 Original line number Diff line number Diff line Loading @@ -133,6 +133,8 @@ mp.messaging.outgoing.registry-store-kubernetes-config-result.value.deserializer ## Test Environment %test.quarkus.oidc.enabled=false %test.quarkus.otel.enabled=false %test.quarkus.kafka.devservices.enabled=false %test.quarkus.vault.devservices.enabled=false %test.quarkus.rest-client.vault-api.url=http://localhost:8200 # %test.kubernetes.namespace-sa.vault.secret-kv-path=test/kube-sa/path Loading src/test/java/org/etsi/osl/hypo/registry/kafka/outgoing/RegistryKafkaProducerTest.java +31 −0 Original line number Diff line number Diff line Loading @@ -110,6 +110,37 @@ class RegistryKafkaProducerTest { () -> assertThat(vaultSecretResponseMessageInMemorySink.received()).hasSize(1)); } @Test void sendKubeConfigResponseMessage_withHeaders() { VaultSecretResponseMessage response = new VaultSecretResponseMessage(); response.setId("vault-id"); response.setStatus(Status.SUCCESS); String value = "value"; RecordHeader recordHeader = new RecordHeader("correlationId", value.getBytes()); RecordHeaders recordHeaders = new RecordHeaders(); recordHeaders.add(recordHeader); producer.sendKubeConfigResponseMessage(response, recordHeaders); await() .atMost(Duration.ofSeconds(5)) .untilAsserted( () -> assertThat(vaultSecretResponseMessageInMemorySink.received()).hasSize(1)); Message<VaultSecretResponseMessage> message = vaultSecretResponseMessageInMemorySink.received().get(0); assertThat(message.getPayload().getId()).isNotBlank().isEqualTo("vault-id"); Headers headers = message.getMetadata(OutgoingKafkaRecordMetadata.class).orElseThrow().getHeaders(); assertThat(headers.toArray()).isNotEmpty().hasSize(1); Header header = headers.lastHeader("correlationId"); assertThat(header).isNotNull(); String stringValue = new String(header.value(), StandardCharsets.UTF_8); assertThat(stringValue).isEqualTo("value"); } @Test void sendOrganizationCredentialsResponseMessage_Success() { CloudEventRegistryOrganizationCredentialsResult response = Loading src/test/java/org/etsi/osl/hypo/registry/service/RegistryServiceTest.java +80 −2 Original line number Diff line number Diff line Loading @@ -630,12 +630,14 @@ class RegistryServiceTest { await() .atMost(Duration.ofSeconds(5)) .untilAsserted(() -> verify(registryKafkaProducer).sendKubeConfigResponseMessage(any())); .untilAsserted( () -> verify(registryKafkaProducer).sendKubeConfigResponseMessage(any(), any())); verify(registryService, times(1)).getKubernetesConfigSecretResponseMessage(messageWithToken); ArgumentCaptor<VaultSecretResponseMessage> responseMessageCaptor = ArgumentCaptor.captor(); verify(registryKafkaProducer).sendKubeConfigResponseMessage(responseMessageCaptor.capture()); verify(registryKafkaProducer) .sendKubeConfigResponseMessage(responseMessageCaptor.capture(), any()); VaultSecretResponseMessage actual = responseMessageCaptor.getValue(); Loading @@ -644,6 +646,82 @@ class RegistryServiceTest { assertThat(actual.getStatus()).isEqualTo(Status.SUCCESS); } @Test void retrieveKubernetesConfigFromKafkaMessage_EchoesIncomingHeadersOnResponse() { String serviceId = "my-service-id"; IncomingKafkaRecordMetadata<?, ?> metadata = Mockito.mock(IncomingKafkaRecordMetadata.class); RecordHeaders incomingHeaders = new RecordHeaders(); incomingHeaders.add( new RecordHeader( "Authorization", ("Bearer " + KEYCLOAK_TOKEN_WITH_GROUP).getBytes(StandardCharsets.UTF_8))); incomingHeaders.add( new RecordHeader("correlationId", "corr-123".getBytes(StandardCharsets.UTF_8))); when(metadata.getHeaders()).thenReturn(incomingHeaders); Message<String> messageWithHeaders = Message.of(serviceId).addMetadata(metadata); registryService .retrieveKubernetesConfigFromKafkaMessage(messageWithHeaders) .await() .indefinitely(); ArgumentCaptor<org.apache.kafka.common.header.Headers> headersCaptor = ArgumentCaptor.captor(); await() .atMost(Duration.ofSeconds(5)) .untilAsserted( () -> verify(registryKafkaProducer) .sendKubeConfigResponseMessage(any(), headersCaptor.capture())); org.apache.kafka.common.header.Headers echoedHeaders = headersCaptor.getValue(); assertThat(echoedHeaders.lastHeader("correlationId")).isNotNull(); assertThat( new String(echoedHeaders.lastHeader("correlationId").value(), StandardCharsets.UTF_8)) .isEqualTo("corr-123"); } @Test void retrieveKubernetesConfigFromKafkaMessage_OnFailure_SendsFailureResponseWithHeaders() { String serviceId = "my-service-id"; IncomingKafkaRecordMetadata<?, ?> metadata = Mockito.mock(IncomingKafkaRecordMetadata.class); RecordHeaders incomingHeaders = new RecordHeaders(); incomingHeaders.add( new RecordHeader( "Authorization", "Bearer invalid-token-string".getBytes(StandardCharsets.UTF_8))); incomingHeaders.add( new RecordHeader("correlationId", "corr-456".getBytes(StandardCharsets.UTF_8))); when(metadata.getHeaders()).thenReturn(incomingHeaders); Message<String> messageWithInvalidToken = Message.of(serviceId).addMetadata(metadata); assertThrows( Exception.class, () -> registryService .retrieveKubernetesConfigFromKafkaMessage(messageWithInvalidToken) .await() .indefinitely()); ArgumentCaptor<VaultSecretResponseMessage> responseCaptor = ArgumentCaptor.captor(); ArgumentCaptor<org.apache.kafka.common.header.Headers> headersCaptor = ArgumentCaptor.captor(); await() .atMost(Duration.ofSeconds(5)) .untilAsserted( () -> verify(registryKafkaProducer) .sendKubeConfigResponseMessage( responseCaptor.capture(), headersCaptor.capture())); VaultSecretResponseMessage failureResponse = responseCaptor.getValue(); assertThat(failureResponse.getStatus()).isEqualTo(Status.FAILURE); assertThat(failureResponse.getId()).isEqualTo(serviceId); org.apache.kafka.common.header.Headers echoedHeaders = headersCaptor.getValue(); assertThat(echoedHeaders.lastHeader("correlationId")).isNotNull(); assertThat( new String(echoedHeaders.lastHeader("correlationId").value(), StandardCharsets.UTF_8)) .isEqualTo("corr-456"); } @Test void persistTmfOrganizationCredentialsSecretFromKafkaMessage() { Message<OrganizationAuthCredentials> messageWithToken = Loading Loading
src/main/java/org/etsi/osl/hypo/registry/kafka/outgoing/RegistryKafkaProducer.java +13 −0 Original line number Diff line number Diff line Loading @@ -86,6 +86,19 @@ public class RegistryKafkaProducer { logger.debugf(ApplicationProperties.PAYLOAD_S, responseMessage.toStringAsJson()); } public void sendKubeConfigResponseMessage( VaultSecretResponseMessage responseMessage, Headers headers) { OutgoingKafkaRecordMetadata<Object> kafkaMetadata = OutgoingKafkaRecordMetadata.builder().withHeaders(headers).build(); vaultSecretResponseMessageEmitter.send(Message.of(responseMessage, Metadata.of(kafkaMetadata))); logger.debugf( "Kube Conf Response data message with headers sent to KAFKA topic: %s", REGISTRY_RESPONSE_KUBE_CONF_CHANNEL); logger.debugf(ApplicationProperties.PAYLOAD_S, responseMessage.toStringAsJson()); } public void sendOrganizationCredentialsResponseMessage( CloudEventRegistryOrganizationCredentialsResult responseMessage) { organizationCredentialsResultEmitter.send(responseMessage); Loading
src/main/java/org/etsi/osl/hypo/registry/service/RegistryService.java +13 −3 Original line number Diff line number Diff line Loading @@ -284,6 +284,9 @@ public class RegistryService { } public Uni<Void> retrieveKubernetesConfigFromKafkaMessage(Message<String> message) { String serviceId = message.getPayload(); Headers headers = extractKafkaHeaders(message); return Uni.createFrom() .voidItem() .invoke( Loading @@ -291,9 +294,16 @@ public class RegistryService { VaultSecretResponseMessage response = getKubernetesConfigSecretResponseMessage(message); logger.debugf( "Successfully retrieved KubeConfig for ID: %s. Sending response.", message.getPayload()); registryKafkaProducer.sendKubeConfigResponseMessage(response); "Successfully retrieved KubeConfig for ID: %s. Sending response.", serviceId); registryKafkaProducer.sendKubeConfigResponseMessage(response, headers); }) .onFailure() .invoke( failure -> { VaultSecretResponseMessage failureResponse = buildKubernetesConfigSecretResponseMessage( null, serviceId, Status.FAILURE, failure.getMessage()); registryKafkaProducer.sendKubeConfigResponseMessage(failureResponse, headers); }); } Loading
src/main/resources/application.properties +2 −0 Original line number Diff line number Diff line Loading @@ -133,6 +133,8 @@ mp.messaging.outgoing.registry-store-kubernetes-config-result.value.deserializer ## Test Environment %test.quarkus.oidc.enabled=false %test.quarkus.otel.enabled=false %test.quarkus.kafka.devservices.enabled=false %test.quarkus.vault.devservices.enabled=false %test.quarkus.rest-client.vault-api.url=http://localhost:8200 # %test.kubernetes.namespace-sa.vault.secret-kv-path=test/kube-sa/path Loading
src/test/java/org/etsi/osl/hypo/registry/kafka/outgoing/RegistryKafkaProducerTest.java +31 −0 Original line number Diff line number Diff line Loading @@ -110,6 +110,37 @@ class RegistryKafkaProducerTest { () -> assertThat(vaultSecretResponseMessageInMemorySink.received()).hasSize(1)); } @Test void sendKubeConfigResponseMessage_withHeaders() { VaultSecretResponseMessage response = new VaultSecretResponseMessage(); response.setId("vault-id"); response.setStatus(Status.SUCCESS); String value = "value"; RecordHeader recordHeader = new RecordHeader("correlationId", value.getBytes()); RecordHeaders recordHeaders = new RecordHeaders(); recordHeaders.add(recordHeader); producer.sendKubeConfigResponseMessage(response, recordHeaders); await() .atMost(Duration.ofSeconds(5)) .untilAsserted( () -> assertThat(vaultSecretResponseMessageInMemorySink.received()).hasSize(1)); Message<VaultSecretResponseMessage> message = vaultSecretResponseMessageInMemorySink.received().get(0); assertThat(message.getPayload().getId()).isNotBlank().isEqualTo("vault-id"); Headers headers = message.getMetadata(OutgoingKafkaRecordMetadata.class).orElseThrow().getHeaders(); assertThat(headers.toArray()).isNotEmpty().hasSize(1); Header header = headers.lastHeader("correlationId"); assertThat(header).isNotNull(); String stringValue = new String(header.value(), StandardCharsets.UTF_8); assertThat(stringValue).isEqualTo("value"); } @Test void sendOrganizationCredentialsResponseMessage_Success() { CloudEventRegistryOrganizationCredentialsResult response = Loading
src/test/java/org/etsi/osl/hypo/registry/service/RegistryServiceTest.java +80 −2 Original line number Diff line number Diff line Loading @@ -630,12 +630,14 @@ class RegistryServiceTest { await() .atMost(Duration.ofSeconds(5)) .untilAsserted(() -> verify(registryKafkaProducer).sendKubeConfigResponseMessage(any())); .untilAsserted( () -> verify(registryKafkaProducer).sendKubeConfigResponseMessage(any(), any())); verify(registryService, times(1)).getKubernetesConfigSecretResponseMessage(messageWithToken); ArgumentCaptor<VaultSecretResponseMessage> responseMessageCaptor = ArgumentCaptor.captor(); verify(registryKafkaProducer).sendKubeConfigResponseMessage(responseMessageCaptor.capture()); verify(registryKafkaProducer) .sendKubeConfigResponseMessage(responseMessageCaptor.capture(), any()); VaultSecretResponseMessage actual = responseMessageCaptor.getValue(); Loading @@ -644,6 +646,82 @@ class RegistryServiceTest { assertThat(actual.getStatus()).isEqualTo(Status.SUCCESS); } @Test void retrieveKubernetesConfigFromKafkaMessage_EchoesIncomingHeadersOnResponse() { String serviceId = "my-service-id"; IncomingKafkaRecordMetadata<?, ?> metadata = Mockito.mock(IncomingKafkaRecordMetadata.class); RecordHeaders incomingHeaders = new RecordHeaders(); incomingHeaders.add( new RecordHeader( "Authorization", ("Bearer " + KEYCLOAK_TOKEN_WITH_GROUP).getBytes(StandardCharsets.UTF_8))); incomingHeaders.add( new RecordHeader("correlationId", "corr-123".getBytes(StandardCharsets.UTF_8))); when(metadata.getHeaders()).thenReturn(incomingHeaders); Message<String> messageWithHeaders = Message.of(serviceId).addMetadata(metadata); registryService .retrieveKubernetesConfigFromKafkaMessage(messageWithHeaders) .await() .indefinitely(); ArgumentCaptor<org.apache.kafka.common.header.Headers> headersCaptor = ArgumentCaptor.captor(); await() .atMost(Duration.ofSeconds(5)) .untilAsserted( () -> verify(registryKafkaProducer) .sendKubeConfigResponseMessage(any(), headersCaptor.capture())); org.apache.kafka.common.header.Headers echoedHeaders = headersCaptor.getValue(); assertThat(echoedHeaders.lastHeader("correlationId")).isNotNull(); assertThat( new String(echoedHeaders.lastHeader("correlationId").value(), StandardCharsets.UTF_8)) .isEqualTo("corr-123"); } @Test void retrieveKubernetesConfigFromKafkaMessage_OnFailure_SendsFailureResponseWithHeaders() { String serviceId = "my-service-id"; IncomingKafkaRecordMetadata<?, ?> metadata = Mockito.mock(IncomingKafkaRecordMetadata.class); RecordHeaders incomingHeaders = new RecordHeaders(); incomingHeaders.add( new RecordHeader( "Authorization", "Bearer invalid-token-string".getBytes(StandardCharsets.UTF_8))); incomingHeaders.add( new RecordHeader("correlationId", "corr-456".getBytes(StandardCharsets.UTF_8))); when(metadata.getHeaders()).thenReturn(incomingHeaders); Message<String> messageWithInvalidToken = Message.of(serviceId).addMetadata(metadata); assertThrows( Exception.class, () -> registryService .retrieveKubernetesConfigFromKafkaMessage(messageWithInvalidToken) .await() .indefinitely()); ArgumentCaptor<VaultSecretResponseMessage> responseCaptor = ArgumentCaptor.captor(); ArgumentCaptor<org.apache.kafka.common.header.Headers> headersCaptor = ArgumentCaptor.captor(); await() .atMost(Duration.ofSeconds(5)) .untilAsserted( () -> verify(registryKafkaProducer) .sendKubeConfigResponseMessage( responseCaptor.capture(), headersCaptor.capture())); VaultSecretResponseMessage failureResponse = responseCaptor.getValue(); assertThat(failureResponse.getStatus()).isEqualTo(Status.FAILURE); assertThat(failureResponse.getId()).isEqualTo(serviceId); org.apache.kafka.common.header.Headers echoedHeaders = headersCaptor.getValue(); assertThat(echoedHeaders.lastHeader("correlationId")).isNotNull(); assertThat( new String(echoedHeaders.lastHeader("correlationId").value(), StandardCharsets.UTF_8)) .isEqualTo("corr-456"); } @Test void persistTmfOrganizationCredentialsSecretFromKafkaMessage() { Message<OrganizationAuthCredentials> messageWithToken = Loading