Loading src/main/java/org/etsi/osl/hypo/api/peering/clients/RegistryClient.java 0 → 100644 +23 −0 Original line number Diff line number Diff line package org.etsi.osl.hypo.api.peering.clients; import jakarta.ws.rs.Consumes; import jakarta.ws.rs.GET; import jakarta.ws.rs.HeaderParam; import jakarta.ws.rs.Path; import jakarta.ws.rs.PathParam; import jakarta.ws.rs.Produces; import jakarta.ws.rs.core.MediaType; import org.eclipse.microprofile.rest.client.inject.RegisterRestClient; import org.etsi.osl.hypo.api.peering.model.PeeringInfoSecret; @Path("registry-api") @RegisterRestClient(configKey = "registry-api") public interface RegistryClient { @GET @Path("organization/secret/{secret-name}") @Consumes(MediaType.APPLICATION_JSON) @Produces(MediaType.APPLICATION_JSON) PeeringInfoSecret getSecretWithVaultToken( @HeaderParam("X-Vault-Token") String vaultToken, @PathParam("secret-name") String secretName); } src/main/java/org/etsi/osl/hypo/api/peering/entity/OssClientData.java +0 −6 Original line number Diff line number Diff line Loading @@ -48,12 +48,6 @@ public class OssClientData { @Column(name = "oauth2TokenUri") private String oauth2TokenUri; @Column(name = "username") private String username; @Column(name = "password") private String password; @Column(name = "healthCheckCounter") private int healthCheckCounter; Loading src/main/java/org/etsi/osl/hypo/api/peering/model/PeeringInfoSecret.java 0 → 100644 +14 −0 Original line number Diff line number Diff line package org.etsi.osl.hypo.api.peering.model; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; import lombok.ToString; @Data @NoArgsConstructor @AllArgsConstructor public class PeeringInfoSecret { private String username; @ToString.Exclude private String password; } src/main/java/org/etsi/osl/hypo/api/peering/services/HealthCheckService.java +23 −14 Original line number Diff line number Diff line Loading @@ -12,12 +12,13 @@ import java.util.List; import java.util.Map; import org.apache.kafka.common.errors.TimeoutException; import org.eclipse.microprofile.config.inject.ConfigProperty; import org.eclipse.microprofile.rest.client.inject.RestClient; import org.etsi.osl.hypo.api.peering.clients.RegistryClient; import org.etsi.osl.hypo.api.peering.entity.OssClientData; import org.etsi.osl.hypo.api.peering.kafka.outgoing.PeeringProducer; import org.etsi.osl.hypo.api.peering.model.PeeringInfoSecret; import org.etsi.osl.hypo.api.peering.oss.OssRestClientImpl; import org.etsi.osl.hypo.api.peering.util.HelperFunctions; import org.etsi.osl.hypo.api.tmf.alarm.schema.Alarm; import org.etsi.osl.hypo.api.tmf.alarm.schema.model.AlarmType; import org.etsi.osl.hypo.api.peering.util.FileReader; import org.etsi.osl.hypo.api.tmf.peering.PeeredOrganizationData; import org.etsi.osl.hypo.api.tmf.services.catalog.schema.ServiceSpecificationEntity; import org.hibernate.exception.JDBCConnectionException; Loading @@ -30,6 +31,7 @@ public class HealthCheckService { private final OssClientDataService ossClientDataService; private final PeeringProducer peeringProducer; private final RegistryClient registryClient; private static final String MAXIMUM_HEARTBEAT_ATTEMPTS = "Heartbeat failure exceeded maximum attempts for organization"; Loading @@ -37,14 +39,20 @@ public class HealthCheckService { @ConfigProperty(name = "failed.attempts") private int failedAttempts; private final String vaultTokenPath; @Inject public HealthCheckService( OssRestClientImpl ossRestClient, OssClientDataService ossClientDataService, PeeringProducer peeringProducer) { PeeringProducer peeringProducer, @RestClient RegistryClient registryClient, @ConfigProperty(name = "peering.vault.token-path") String vaultTokenPath) { this.ossRestClientImpl = ossRestClient; this.ossClientDataService = ossClientDataService; this.peeringProducer = peeringProducer; this.registryClient = registryClient; this.vaultTokenPath = vaultTokenPath; } @Scheduled(cron = "{cron.expr}") Loading @@ -61,7 +69,11 @@ public class HealthCheckService { for (OssClientData ossClientDatum : ossClientData) { try { checkPeeredOrganizationStatus(ossClientDatum); String keycloakToken = FileReader.getFileContent(vaultTokenPath); PeeringInfoSecret peeringInfoSecret = registryClient.getSecretWithVaultToken( keycloakToken, ossClientDatum.getOrganizationId()); checkPeeredOrganizationStatus(ossClientDatum, peeringInfoSecret); if (ossClientDatum.getHealthCheckCounter() > 0) { safeUpdateHeartbeat(ossClientDatum, 0); Loading @@ -73,7 +85,8 @@ public class HealthCheckService { } } public void checkPeeredOrganizationStatus(OssClientData ossClientData) { public void checkPeeredOrganizationStatus( OssClientData ossClientData, PeeringInfoSecret peeringInfoSecret) { List<ServiceSpecificationEntity> serviceSpecificationEntities = ossRestClientImpl.callOssToGetServiceSpecifications( Loading @@ -81,8 +94,8 @@ public class HealthCheckService { ossClientData.getBaseUrl(), ossClientData.getOauth2ClientId(), ossClientData.getOauth2ClientSecret(), ossClientData.getUsername(), ossClientData.getPassword(), peeringInfoSecret.getUsername(), peeringInfoSecret.getPassword(), ossClientData.getOrganizationId()); List<ServiceSpecificationEntity> updatedServiceSpecificationEntities = new ArrayList<>(); Loading @@ -96,8 +109,8 @@ public class HealthCheckService { ossClientData.getBaseUrl(), ossClientData.getOauth2ClientId(), ossClientData.getOauth2ClientSecret(), ossClientData.getUsername(), ossClientData.getPassword(), peeringInfoSecret.getUsername(), peeringInfoSecret.getPassword(), entry.getKey(), ossClientData.getOrganizationId()); serviceSpecificationEntityList.add(serviceSpecificationEntity); Loading Loading @@ -158,10 +171,6 @@ public class HealthCheckService { logger.errorf(MAXIMUM_HEARTBEAT_ATTEMPTS + " %s", ossClientDatum.getOrganizationId()); try { peeringProducer.sendHeartbeatFailure(ossClientDatum.getOrganizationId()); Alarm alarm = HelperFunctions.prepareAlarm( ossClientDatum.getOrganizationId(), e.getMessage(), AlarmType.COMMUNICATIONSALARM); } catch (Exception kafkaEx) { logger.error("Failed to send Kafka alarm. Cluster might be unstable.", kafkaEx); } Loading src/main/java/org/etsi/osl/hypo/api/peering/util/FileReader.java 0 → 100644 +26 −0 Original line number Diff line number Diff line package org.etsi.osl.hypo.api.peering.util; import java.io.IOException; import java.nio.file.Files; import java.nio.file.Path; public class FileReader { // 👇 Prevent instantiation FileReader() { throw new UnsupportedOperationException("Utility class — instantiation not allowed"); } public static String getFileContent(String filePath) { Path path = Path.of(filePath); if (!Files.exists(path)) { throw new IllegalArgumentException(String.format("File not found in path: %s", path)); } try { return Files.readString(path); } catch (IOException e) { throw new ReadFileException(String.format("Failed to read file in path: %s", path)); } } } Loading
src/main/java/org/etsi/osl/hypo/api/peering/clients/RegistryClient.java 0 → 100644 +23 −0 Original line number Diff line number Diff line package org.etsi.osl.hypo.api.peering.clients; import jakarta.ws.rs.Consumes; import jakarta.ws.rs.GET; import jakarta.ws.rs.HeaderParam; import jakarta.ws.rs.Path; import jakarta.ws.rs.PathParam; import jakarta.ws.rs.Produces; import jakarta.ws.rs.core.MediaType; import org.eclipse.microprofile.rest.client.inject.RegisterRestClient; import org.etsi.osl.hypo.api.peering.model.PeeringInfoSecret; @Path("registry-api") @RegisterRestClient(configKey = "registry-api") public interface RegistryClient { @GET @Path("organization/secret/{secret-name}") @Consumes(MediaType.APPLICATION_JSON) @Produces(MediaType.APPLICATION_JSON) PeeringInfoSecret getSecretWithVaultToken( @HeaderParam("X-Vault-Token") String vaultToken, @PathParam("secret-name") String secretName); }
src/main/java/org/etsi/osl/hypo/api/peering/entity/OssClientData.java +0 −6 Original line number Diff line number Diff line Loading @@ -48,12 +48,6 @@ public class OssClientData { @Column(name = "oauth2TokenUri") private String oauth2TokenUri; @Column(name = "username") private String username; @Column(name = "password") private String password; @Column(name = "healthCheckCounter") private int healthCheckCounter; Loading
src/main/java/org/etsi/osl/hypo/api/peering/model/PeeringInfoSecret.java 0 → 100644 +14 −0 Original line number Diff line number Diff line package org.etsi.osl.hypo.api.peering.model; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; import lombok.ToString; @Data @NoArgsConstructor @AllArgsConstructor public class PeeringInfoSecret { private String username; @ToString.Exclude private String password; }
src/main/java/org/etsi/osl/hypo/api/peering/services/HealthCheckService.java +23 −14 Original line number Diff line number Diff line Loading @@ -12,12 +12,13 @@ import java.util.List; import java.util.Map; import org.apache.kafka.common.errors.TimeoutException; import org.eclipse.microprofile.config.inject.ConfigProperty; import org.eclipse.microprofile.rest.client.inject.RestClient; import org.etsi.osl.hypo.api.peering.clients.RegistryClient; import org.etsi.osl.hypo.api.peering.entity.OssClientData; import org.etsi.osl.hypo.api.peering.kafka.outgoing.PeeringProducer; import org.etsi.osl.hypo.api.peering.model.PeeringInfoSecret; import org.etsi.osl.hypo.api.peering.oss.OssRestClientImpl; import org.etsi.osl.hypo.api.peering.util.HelperFunctions; import org.etsi.osl.hypo.api.tmf.alarm.schema.Alarm; import org.etsi.osl.hypo.api.tmf.alarm.schema.model.AlarmType; import org.etsi.osl.hypo.api.peering.util.FileReader; import org.etsi.osl.hypo.api.tmf.peering.PeeredOrganizationData; import org.etsi.osl.hypo.api.tmf.services.catalog.schema.ServiceSpecificationEntity; import org.hibernate.exception.JDBCConnectionException; Loading @@ -30,6 +31,7 @@ public class HealthCheckService { private final OssClientDataService ossClientDataService; private final PeeringProducer peeringProducer; private final RegistryClient registryClient; private static final String MAXIMUM_HEARTBEAT_ATTEMPTS = "Heartbeat failure exceeded maximum attempts for organization"; Loading @@ -37,14 +39,20 @@ public class HealthCheckService { @ConfigProperty(name = "failed.attempts") private int failedAttempts; private final String vaultTokenPath; @Inject public HealthCheckService( OssRestClientImpl ossRestClient, OssClientDataService ossClientDataService, PeeringProducer peeringProducer) { PeeringProducer peeringProducer, @RestClient RegistryClient registryClient, @ConfigProperty(name = "peering.vault.token-path") String vaultTokenPath) { this.ossRestClientImpl = ossRestClient; this.ossClientDataService = ossClientDataService; this.peeringProducer = peeringProducer; this.registryClient = registryClient; this.vaultTokenPath = vaultTokenPath; } @Scheduled(cron = "{cron.expr}") Loading @@ -61,7 +69,11 @@ public class HealthCheckService { for (OssClientData ossClientDatum : ossClientData) { try { checkPeeredOrganizationStatus(ossClientDatum); String keycloakToken = FileReader.getFileContent(vaultTokenPath); PeeringInfoSecret peeringInfoSecret = registryClient.getSecretWithVaultToken( keycloakToken, ossClientDatum.getOrganizationId()); checkPeeredOrganizationStatus(ossClientDatum, peeringInfoSecret); if (ossClientDatum.getHealthCheckCounter() > 0) { safeUpdateHeartbeat(ossClientDatum, 0); Loading @@ -73,7 +85,8 @@ public class HealthCheckService { } } public void checkPeeredOrganizationStatus(OssClientData ossClientData) { public void checkPeeredOrganizationStatus( OssClientData ossClientData, PeeringInfoSecret peeringInfoSecret) { List<ServiceSpecificationEntity> serviceSpecificationEntities = ossRestClientImpl.callOssToGetServiceSpecifications( Loading @@ -81,8 +94,8 @@ public class HealthCheckService { ossClientData.getBaseUrl(), ossClientData.getOauth2ClientId(), ossClientData.getOauth2ClientSecret(), ossClientData.getUsername(), ossClientData.getPassword(), peeringInfoSecret.getUsername(), peeringInfoSecret.getPassword(), ossClientData.getOrganizationId()); List<ServiceSpecificationEntity> updatedServiceSpecificationEntities = new ArrayList<>(); Loading @@ -96,8 +109,8 @@ public class HealthCheckService { ossClientData.getBaseUrl(), ossClientData.getOauth2ClientId(), ossClientData.getOauth2ClientSecret(), ossClientData.getUsername(), ossClientData.getPassword(), peeringInfoSecret.getUsername(), peeringInfoSecret.getPassword(), entry.getKey(), ossClientData.getOrganizationId()); serviceSpecificationEntityList.add(serviceSpecificationEntity); Loading Loading @@ -158,10 +171,6 @@ public class HealthCheckService { logger.errorf(MAXIMUM_HEARTBEAT_ATTEMPTS + " %s", ossClientDatum.getOrganizationId()); try { peeringProducer.sendHeartbeatFailure(ossClientDatum.getOrganizationId()); Alarm alarm = HelperFunctions.prepareAlarm( ossClientDatum.getOrganizationId(), e.getMessage(), AlarmType.COMMUNICATIONSALARM); } catch (Exception kafkaEx) { logger.error("Failed to send Kafka alarm. Cluster might be unstable.", kafkaEx); } Loading
src/main/java/org/etsi/osl/hypo/api/peering/util/FileReader.java 0 → 100644 +26 −0 Original line number Diff line number Diff line package org.etsi.osl.hypo.api.peering.util; import java.io.IOException; import java.nio.file.Files; import java.nio.file.Path; public class FileReader { // 👇 Prevent instantiation FileReader() { throw new UnsupportedOperationException("Utility class — instantiation not allowed"); } public static String getFileContent(String filePath) { Path path = Path.of(filePath); if (!Files.exists(path)) { throw new IllegalArgumentException(String.format("File not found in path: %s", path)); } try { return Files.readString(path); } catch (IOException e) { throw new ReadFileException(String.format("Failed to read file in path: %s", path)); } } }