diff --git a/build.gradle.kts b/build.gradle.kts index 738436e..8a5956e 100644 --- a/build.gradle.kts +++ b/build.gradle.kts @@ -174,6 +174,8 @@ dependencies { implementation("org.semver4j:semver4j:5.3.0") + implementation("org.eclipse.ditto:ditto-wot-model:3.6.0") + testImplementation("io.javaoperatorsdk:operator-framework-spring-boot-starter-test:5.5.0") { exclude(group = "org.apache.logging.log4j", module = "log4j-slf4j2-impl") } diff --git a/src/main/helm/templates/clusterrole.yaml b/src/main/helm/templates/clusterrole.yaml index 408b672..2941510 100644 --- a/src/main/helm/templates/clusterrole.yaml +++ b/src/main/helm/templates/clusterrole.yaml @@ -16,4 +16,19 @@ rules: {{- range .Values.customResources }} - "{{ . }}" {{- end }} - verbs: ["get", "list", "create", "update", "delete", "watch", "patch"] \ No newline at end of file + verbs: ["get", "list", "create", "update", "delete", "watch", "patch"] + + # Add permission to access deployments in the apps API group + - apiGroups: [ "apps" ] + resources: [ "deployments" ] + verbs: ["get", "list", "watch", "create", "update", "patch", "delete"] + + # Add permission to access pods + - apiGroups: [ "" ] + resources: [ "pods" ] + verbs: [ "get", "list", "watch" ] + + # Add permission to access services + - apiGroups: [ "" ] + resources: [ "services" ] + verbs: [ "get", "list", "watch" ] \ No newline at end of file diff --git a/src/main/java/ai/ancf/lmos/operator/reconciler/AgentReconciler.java b/src/main/java/ai/ancf/lmos/operator/reconciler/AgentReconciler.java index 0ede0eb..c9a06d2 100644 --- a/src/main/java/ai/ancf/lmos/operator/reconciler/AgentReconciler.java +++ b/src/main/java/ai/ancf/lmos/operator/reconciler/AgentReconciler.java @@ -6,28 +6,63 @@ package ai.ancf.lmos.operator.reconciler; -import io.javaoperatorsdk.operator.api.reconciler.Context; -import io.javaoperatorsdk.operator.api.reconciler.ControllerConfiguration; -import io.javaoperatorsdk.operator.api.reconciler.Reconciler; -import io.javaoperatorsdk.operator.api.reconciler.UpdateControl; -import ai.ancf.lmos.operator.resources.agent.AgentResource; +import ai.ancf.lmos.operator.service.AgentDeploymentStatusService; +import ai.ancf.lmos.operator.service.AgentServiceQuery; +import ai.ancf.lmos.operator.service.KubernetesResourceManager; +import io.fabric8.kubernetes.api.model.apps.Deployment; +import io.javaoperatorsdk.operator.api.reconciler.*; +import org.eclipse.ditto.wot.model.ThingDescription; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.stereotype.Component; +import java.util.concurrent.TimeUnit; +/** + * Reconciles Deployment resources by watching associated Pods and registering Services. + */ @Component -@ControllerConfiguration -public class AgentReconciler implements Reconciler { +@ControllerConfiguration(labelSelector = "wot-agent=true") +public class AgentReconciler implements Reconciler, Cleaner { private static final Logger LOG = LoggerFactory.getLogger(AgentReconciler.class); + private final AgentServiceQuery agentServiceQuery; + private final AgentDeploymentStatusService agentDeploymentStatusService; + private final KubernetesResourceManager kubernetesResourceManager; + + public AgentReconciler(AgentServiceQuery agentServiceQuery, AgentDeploymentStatusService agentDeploymentStatusService, KubernetesResourceManager kubernetesResourceManager) { + this.agentServiceQuery = agentServiceQuery; + this.agentDeploymentStatusService = agentDeploymentStatusService; + this.kubernetesResourceManager = kubernetesResourceManager; + } @Override - public UpdateControl reconcile(AgentResource agentResource, Context context) { - // TODO: fill in logic - LOG.debug("Agent reconcile"); + public UpdateControl reconcile(Deployment deployment, Context context) { + + boolean isDeploymentReady = agentDeploymentStatusService.isDeploymentReady(deployment); + + LOG.info("is Deployment {} ready: {}", deployment.getMetadata().getName(), isDeploymentReady); - return UpdateControl.noUpdate(); + if(isDeploymentReady) { + try { + String serviceUrl = kubernetesResourceManager.getServiceUrl(deployment); + ThingDescription thingDescription = agentServiceQuery.queryAgentService(serviceUrl); + kubernetesResourceManager.createOrUpdateAgentResource(thingDescription, deployment); + return UpdateControl.noUpdate(); + } catch (Exception e) { + LOG.error("Error processing td for deployment: {}", deployment.getMetadata().getName(), e); + return UpdateControl.noUpdate().rescheduleAfter(10, TimeUnit.SECONDS); + } + } + + return UpdateControl.noUpdate().rescheduleAfter(10, TimeUnit.SECONDS); + } + + @Override + public DeleteControl cleanup(Deployment deployment, Context context) { + LOG.info("Trigger AgentResource Deletion for deployment: {}", deployment.getMetadata().getName()); + kubernetesResourceManager.deleteAgentResource(deployment); + return DeleteControl.defaultDelete(); } } \ No newline at end of file diff --git a/src/main/java/ai/ancf/lmos/operator/resources/agent/AgentSpec.java b/src/main/java/ai/ancf/lmos/operator/resources/agent/AgentSpec.java index 0a82c28..bbb45ef 100644 --- a/src/main/java/ai/ancf/lmos/operator/resources/agent/AgentSpec.java +++ b/src/main/java/ai/ancf/lmos/operator/resources/agent/AgentSpec.java @@ -14,6 +14,7 @@ public class AgentSpec { private Set supportedTenants; private Set supportedChannels; private Set providedCapabilities; + private String thingDescription; public AgentSpec() { } @@ -58,6 +59,14 @@ public void setProvidedCapabilities(Set providedCapabilities this.providedCapabilities = providedCapabilities; } + public String getThingDescription() { + return thingDescription; + } + + public void setThingDescription(String thingDescription) { + this.thingDescription = thingDescription; + } + @Override public boolean equals(Object o) { if (this == o) return true; diff --git a/src/main/java/ai/ancf/lmos/operator/service/AgentDeploymentStatusService.java b/src/main/java/ai/ancf/lmos/operator/service/AgentDeploymentStatusService.java new file mode 100644 index 0000000..9e07bcd --- /dev/null +++ b/src/main/java/ai/ancf/lmos/operator/service/AgentDeploymentStatusService.java @@ -0,0 +1,31 @@ +/* + * SPDX-FileCopyrightText: 2024 Deutsche Telekom AG + * + * SPDX-License-Identifier: Apache-2.0 + */ + +package ai.ancf.lmos.operator.service; + +import io.fabric8.kubernetes.api.model.apps.Deployment; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +@org.springframework.stereotype.Service +public class AgentDeploymentStatusService { + + private static final Logger LOG = LoggerFactory.getLogger(AgentDeploymentStatusService.class); + + public boolean isDeploymentReady(Deployment deployment) { + String deploymentName = deployment.getMetadata().getName(); + String deploymentNamespace = deployment.getMetadata().getNamespace(); + Integer replicas = deployment.getStatus().getReplicas(); + Integer availableReplicas = deployment.getStatus().getAvailableReplicas(); + Integer desiredReplicas = deployment.getSpec().getReplicas(); + + LOG.info( + "Reconciling Deployment: {} in namespace: {}, Replicas, availableReplicas: {}, desiredReplicas: {}", + deploymentName, deploymentNamespace, availableReplicas, desiredReplicas); + + return (replicas != null && availableReplicas != null && replicas.equals(desiredReplicas) && availableReplicas.equals(desiredReplicas)); + } +} \ No newline at end of file diff --git a/src/main/java/ai/ancf/lmos/operator/service/AgentServiceQuery.java b/src/main/java/ai/ancf/lmos/operator/service/AgentServiceQuery.java new file mode 100644 index 0000000..c6388b3 --- /dev/null +++ b/src/main/java/ai/ancf/lmos/operator/service/AgentServiceQuery.java @@ -0,0 +1,53 @@ +/* + * SPDX-FileCopyrightText: 2024 Deutsche Telekom AG + * + * SPDX-License-Identifier: Apache-2.0 + */ + +package ai.ancf.lmos.operator.service; + +import org.eclipse.ditto.json.JsonObject; +import org.eclipse.ditto.wot.model.ThingDescription; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.http.MediaType; +import org.springframework.web.reactive.function.client.WebClient; + +@org.springframework.stereotype.Service +public class AgentServiceQuery { + + private static final Logger LOG = LoggerFactory.getLogger(AgentServiceQuery.class); + + private final WebClient webClient; + + public AgentServiceQuery() { + this.webClient = WebClient.builder().build(); + } + + public ThingDescription queryAgentService(String serviceUrl) { + LOG.info("Querying Agent for TD: {}", serviceUrl); + String thingDescriptionJson = webClient.get() + .uri(serviceUrl) + .accept(MediaType.APPLICATION_JSON) + .retrieve() + .bodyToMono(String.class).block(); + + if (thingDescriptionJson == null || thingDescriptionJson.isEmpty()) { + throw new IllegalStateException("TD Response body from agent is empty"); + } + + ThingDescription thingDescription = ThingDescription.fromJson(JsonObject.of(thingDescriptionJson)); + validateTD(thingDescription); + return thingDescription; + } + + private void validateTD(ThingDescription thingDescription) { + if (thingDescription == null) { + throw new RuntimeException("ThingDescription is null"); + } + + thingDescription.getDescription() + .filter(desc -> !desc.isEmpty()) + .orElseThrow(() -> new RuntimeException("Description of TD is not present or is empty")); + } +} \ No newline at end of file diff --git a/src/main/java/ai/ancf/lmos/operator/service/KubernetesResourceManager.java b/src/main/java/ai/ancf/lmos/operator/service/KubernetesResourceManager.java new file mode 100644 index 0000000..ffabc52 --- /dev/null +++ b/src/main/java/ai/ancf/lmos/operator/service/KubernetesResourceManager.java @@ -0,0 +1,97 @@ +/* + * SPDX-FileCopyrightText: 2024 Deutsche Telekom AG + * + * SPDX-License-Identifier: Apache-2.0 + */ + +package ai.ancf.lmos.operator.service; + +import ai.ancf.lmos.operator.resources.agent.AgentResource; +import ai.ancf.lmos.operator.resources.agent.AgentSpec; +import io.fabric8.kubernetes.api.model.*; +import io.fabric8.kubernetes.api.model.apps.Deployment; +import io.fabric8.kubernetes.client.KubernetesClient; +import org.eclipse.ditto.wot.model.ThingDescription; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.List; +import java.util.Map; + +@org.springframework.stereotype.Service +public class KubernetesResourceManager { + + private static final Logger LOG = LoggerFactory.getLogger(KubernetesResourceManager.class); + + private final KubernetesClient kubernetesClient; + + public KubernetesResourceManager(KubernetesClient kubernetesClient) { + this.kubernetesClient = kubernetesClient; + } + + public void createOrUpdateAgentResource(ThingDescription thingDescription, Deployment deployment) { + AgentResource agentResource = new AgentResource(); + agentResource.setMetadata(new ObjectMetaBuilder() + .withName(deployment.getMetadata().getName()) + .withNamespace(deployment.getMetadata().getNamespace()) + .build()); + + AgentSpec spec = new AgentSpec(); + spec.setDescription(thingDescription.getDescription().get() + .toString()); + spec.setThingDescription(thingDescription.toJsonString()); + + agentResource.setSpec(spec); + + AgentResource agentResourceCreated = kubernetesClient.resources(AgentResource.class) + .inNamespace(deployment.getMetadata().getNamespace()) + .withName(deployment.getMetadata().getName()) + .createOrReplace(agentResource); + + LOG.info("AgentResource {} created/updated for deployment: {}", agentResourceCreated.getFullResourceName(), deployment.getMetadata().getName()); + } + + public Service findService(Deployment deployment) { + Map selectorLabels = deployment.getSpec().getSelector().getMatchLabels(); + String deploymentNamespace = deployment.getMetadata().getNamespace(); + ServiceList serviceList = kubernetesClient.services() + .inNamespace(deploymentNamespace) + .withLabels(selectorLabels) + .list(); + if (serviceList.getItems().size() != 1) { + LOG.error("Expected exactly one service, but got {}, {}", serviceList.getItems().size(), serviceList.getItems()); + throw new IllegalStateException("Expected exactly one service, but got " + serviceList.getItems().size()); + } + return serviceList + .getItems() + .getFirst(); + } + + public String getBaseServiceUrl(Service service) { + ServicePort servicePort = service.getSpec().getPorts().getFirst(); + String port = servicePort.getPort().toString(); + boolean isHttps = servicePort.getPort() == 443 || "https".equalsIgnoreCase(servicePort.getName()); + String protocol = isHttps ? "https://" : "http://"; + String url = protocol + service.getMetadata().getName() + ":" + port; + LOG.info("Service URL is: {}", url); + return url; + } + + public String getServiceUrl(Deployment deployment) { + String agentPath = deployment.getMetadata().getAnnotations().getOrDefault("wot.w3.org/td-endpoint", ".well-known/wot"); + String baseServiceUrl = getBaseServiceUrl(findService(deployment)); + if (agentPath.startsWith("/")) { + return baseServiceUrl + agentPath; + } else { + return baseServiceUrl + "/" + agentPath; + } + } + + public void deleteAgentResource(Deployment deployment) { + List deleteStatus = kubernetesClient.resources(AgentResource.class) + .inNamespace(deployment.getMetadata().getNamespace()) + .withName(deployment.getMetadata().getName()) + .delete(); + LOG.info("AgentResource {} deleted for deployment: {}", deleteStatus, deployment.getMetadata().getName()); + } +} \ No newline at end of file diff --git a/src/test/java/ai/ancf/lmos/operator/reconciler/AgentReconcilerTest.java b/src/test/java/ai/ancf/lmos/operator/reconciler/AgentReconcilerTest.java new file mode 100644 index 0000000..54bc18a --- /dev/null +++ b/src/test/java/ai/ancf/lmos/operator/reconciler/AgentReconcilerTest.java @@ -0,0 +1,113 @@ +/* + * SPDX-FileCopyrightText: 2024 Deutsche Telekom AG + * + * SPDX-License-Identifier: Apache-2.0 + */ + +package ai.ancf.lmos.operator.reconciler; + +import ai.ancf.lmos.operator.service.AgentDeploymentStatusService; +import ai.ancf.lmos.operator.service.AgentServiceQuery; +import ai.ancf.lmos.operator.service.KubernetesResourceManager; +import io.fabric8.kubernetes.api.model.ObjectMeta; +import io.fabric8.kubernetes.api.model.apps.Deployment; +import io.javaoperatorsdk.operator.api.reconciler.DeleteControl; +import io.javaoperatorsdk.operator.api.reconciler.UpdateControl; +import io.javaoperatorsdk.operator.springboot.starter.test.EnableMockOperator; +import org.eclipse.ditto.wot.model.ThingDescription; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.mock.mockito.MockBean; +import org.springframework.context.annotation.Import; + +import java.util.Objects; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.*; + +@SpringBootTest +@EnableMockOperator +@Import(AgentReconciler.class) +public class AgentReconcilerTest { + + @MockBean + private AgentDeploymentStatusService agentDeploymentStatusService; + + @MockBean + private AgentServiceQuery agentServiceQuery; + + @MockBean + private KubernetesResourceManager kubernetesResourceManager; + + private Deployment deployment; + private AgentReconciler agentReconciler; + + @BeforeEach + public void setUp() { + // Initialize a mock deployment with necessary metadata and status + deployment = new Deployment(); + ObjectMeta meta = new ObjectMeta(); + meta.setName("test-deployment"); + meta.setNamespace("test-namespace"); + deployment.setMetadata(meta); + + agentReconciler = new AgentReconciler(agentServiceQuery, agentDeploymentStatusService, kubernetesResourceManager); + } + + @Test + public void testReconcile_whenDeploymentNotReady() { + // Pre-conditions + when(agentDeploymentStatusService.isDeploymentReady(any(Deployment.class))).thenReturn(false); + + // Steps & Expected Result + UpdateControl updateControl = agentReconciler.reconcile(deployment, null); + verify(agentDeploymentStatusService).isDeploymentReady(deployment); + verify(kubernetesResourceManager, never()).createOrUpdateAgentResource(any(), any()); + Assertions.assertTrue(updateControl.isNoUpdate()); + Assertions.assertEquals(updateControl.getScheduleDelay().get(), 10000); + } + + @Test + public void testReconcile_whenDeploymentReady() { + // Pre-conditions + when(agentDeploymentStatusService.isDeploymentReady(any(Deployment.class))).thenReturn(true); + when(kubernetesResourceManager.getServiceUrl(any(Deployment.class))).thenReturn("http://example.com"); + when(agentServiceQuery.queryAgentService(anyString())).thenReturn(ThingDescription.newBuilder().build()); + + // Steps & Expected Result + UpdateControl updateControl = agentReconciler.reconcile(deployment, null); + verify(agentDeploymentStatusService).isDeploymentReady(deployment); + verify(kubernetesResourceManager).getServiceUrl(deployment); + verify(agentServiceQuery).queryAgentService("http://example.com"); + verify(kubernetesResourceManager).createOrUpdateAgentResource(any(), eq(deployment)); + Assertions.assertTrue(updateControl.isNoUpdate()); + Assertions.assertFalse(updateControl.getScheduleDelay().isPresent()); + } + + @Test + public void testReconcile_onQueryFailure() { + // Pre-conditions + when(agentDeploymentStatusService.isDeploymentReady(any(Deployment.class))).thenReturn(true); + when(kubernetesResourceManager.getServiceUrl(any(Deployment.class))).thenReturn("http://example.com"); + when(agentServiceQuery.queryAgentService(anyString())).thenThrow(new RuntimeException("Service query failed")); + + // Steps & Expected Result + UpdateControl updateControl = agentReconciler.reconcile(deployment, null); + verify(agentDeploymentStatusService).isDeploymentReady(deployment); + verify(kubernetesResourceManager).getServiceUrl(deployment); + verify(agentServiceQuery).queryAgentService("http://example.com"); + verify(kubernetesResourceManager, never()).createOrUpdateAgentResource(any(), any()); + Assertions.assertTrue(updateControl.isNoUpdate()); + Assertions.assertEquals(updateControl.getScheduleDelay().get(), 10000); + } + + @Test + public void testCleanup_deletesAgentResource() { + // Steps & Expected Result + DeleteControl deleteControl = agentReconciler.cleanup(deployment, null); + verify(kubernetesResourceManager).deleteAgentResource(deployment); + assert Objects.equals(deleteControl.isRemoveFinalizer(), true); + } +} \ No newline at end of file