From dfdc0cebae1eac0b98789abf0e3db2bf7fd49b60 Mon Sep 17 00:00:00 2001 From: Joseph Grogan Date: Fri, 5 Sep 2025 15:22:24 -0400 Subject: [PATCH 1/2] Expose configs via pipeline element table --- .../hoptimator/k8s/K8sPipelineElement.java | 11 +- .../hoptimator/k8s/K8sPipelineElementApi.java | 75 ++++- .../k8s/K8sPipelineElementMapApi.java | 2 +- .../k8s/K8sPipelineElementTable.java | 15 +- .../K8sPipelineElementStatusEstimator.java | 2 +- .../k8s/K8sPipelineElementApiTest.java | 303 ++++++++++++++++++ 6 files changed, 395 insertions(+), 13 deletions(-) create mode 100644 hoptimator-k8s/src/test/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApiTest.java diff --git a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElement.java b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElement.java index ac62d7d67..8ac80958f 100644 --- a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElement.java +++ b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElement.java @@ -2,6 +2,7 @@ import java.util.HashSet; import java.util.List; +import java.util.Map; import java.util.Objects; import java.util.Set; import java.util.stream.Collectors; @@ -11,13 +12,15 @@ /** Represents a pipeline element status and its associated pipelines. */ public class K8sPipelineElement { - private String name; + private final String name; private final K8sPipelineElementStatus status; + private final Map configs; private final Set pipelines = new HashSet<>(); - public K8sPipelineElement(V1alpha1Pipeline pipeline, K8sPipelineElementStatus status) { + public K8sPipelineElement(V1alpha1Pipeline pipeline, K8sPipelineElementStatus status, Map configs) { this.name = status.getName(); this.status = status; + this.configs = configs; this.pipelines.add(pipeline); } @@ -29,6 +32,10 @@ public K8sPipelineElementStatus status() { return status; } + public Map configs() { + return configs; + } + public void addPipeline(V1alpha1Pipeline pipeline) { pipelines.add(pipeline); } diff --git a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApi.java b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApi.java index 7a0d633c9..3f2843f99 100644 --- a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApi.java +++ b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApi.java @@ -1,6 +1,12 @@ package com.linkedin.hoptimator.k8s; +import com.google.gson.JsonElement; +import com.google.gson.JsonObject; +import io.kubernetes.client.util.generic.KubernetesApiResponse; +import io.kubernetes.client.util.generic.dynamic.DynamicKubernetesObject; +import io.kubernetes.client.util.generic.dynamic.Dynamics; import java.sql.SQLException; +import java.util.Arrays; import java.util.Collection; import java.util.HashMap; import java.util.List; @@ -11,9 +17,15 @@ import com.linkedin.hoptimator.k8s.status.K8sPipelineElementStatus; import com.linkedin.hoptimator.k8s.status.K8sPipelineElementStatusEstimator; import com.linkedin.hoptimator.util.Api; +import java.util.Objects; +import java.util.stream.Collectors; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + /** Provides all pipeline elements in a {@link com.linkedin.hoptimator.k8s.K8sContext} instance. */ public class K8sPipelineElementApi implements Api { + private final static Logger log = LoggerFactory.getLogger(K8sPipelineElementApi.class); private final K8sContext context; K8sPipelineElementApi(K8sContext context) { @@ -27,18 +39,73 @@ private Collection discoverAllElements(K8sContext context) t Map elements = new HashMap<>(); K8sPipelineElementStatusEstimator statusEstimator = new K8sPipelineElementStatusEstimator(context); for (V1alpha1Pipeline pipeline : pipelines) { - List elementStatuses = statusEstimator.estimateStatuses(pipeline); - for (K8sPipelineElementStatus elementStatus : elementStatuses) { + String namespace = Objects.requireNonNull(pipeline.getMetadata()).getNamespace(); + List elementYamls = getPipelineElements(pipeline); + for (String elementYaml : elementYamls) { + K8sPipelineElementStatus elementStatus = statusEstimator.estimateElementStatus(elementYaml, namespace); + Map configurations = getElementConfiguration(elementYaml, namespace); String key = elementStatus.getName(); if (!elements.containsKey(key)) { - elements.put(key, new K8sPipelineElement(pipeline, elementStatus)); + elements.put(key, new K8sPipelineElement(pipeline, elementStatus, configurations)); } - elements.get(key).addPipeline(pipeline); } } return elements.values(); } + /** + * Returns list of all elements specified in the given pipeline. + */ + List getPipelineElements(V1alpha1Pipeline pipeline) { + return Arrays.stream(Objects.requireNonNull(Objects.requireNonNull(pipeline.getSpec()).getYaml()) + .split("\n---\n")) + .map(String::trim) + .filter(x -> !x.isEmpty()) + .collect(Collectors.toList()); + } + + /** + * Returns spec.configs of the given element, if present. + */ + Map getElementConfiguration(String elementYaml, String pipelineNamespace) { + DynamicKubernetesObject obj = Dynamics.newFromYaml(elementYaml); + String name = obj.getMetadata().getName(); + String namespace = obj.getMetadata().getNamespace() == null ? pipelineNamespace : obj.getMetadata().getNamespace(); + String kind = obj.getKind(); + String nameWithKind = String.format("%s/%s", kind, name); + try { + KubernetesApiResponse existing = + context.dynamic(obj.getApiVersion(), K8sUtils.guessPlural(obj)).get(namespace, name); + String failureMessage = + String.format("Failed to fetch %s in namespace %s: %s.", nameWithKind, namespace, existing.toString()); + existing.onFailure((code, status) -> log.warn(failureMessage)); + if (!existing.isSuccess()) { + return new HashMap<>(); + } + DynamicKubernetesObject element = existing.getObject(); + if (element == null || element.getRaw() == null) { + return new HashMap<>(); + } + if (!element.getRaw().has("spec") || !element.getRaw().getAsJsonObject("spec").has("configs")) { + return new HashMap<>(); + } + + // Extract configs as a Map + JsonObject configs = element.getRaw().getAsJsonObject("spec").getAsJsonObject("configs"); + Map configMap = new HashMap<>(); + for (Map.Entry entry : configs.entrySet()) { + configMap.put(entry.getKey(), entry.getValue().getAsString()); + } + return configMap; + } catch (Exception e) { + String failureMessage = + String.format("Encountered exception while checking status of %s in namespace %s: %s", nameWithKind, + namespace, e); + log.error(failureMessage); + return new HashMap<>(); + } + } + /** * Lists all pipeline elements in the context. */ diff --git a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementMapApi.java b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementMapApi.java index 486d79c83..2328b8fdc 100644 --- a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementMapApi.java +++ b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementMapApi.java @@ -28,7 +28,7 @@ public Collection list() throws SQLException { } private static Stream mapEntriesFromElement(K8sPipelineElement element) { - String elementName = element.status().getName(); + String elementName = element.name(); return element.pipelineNames() .stream() .map(pipelineName -> new K8sPipelineElementMapEntry(elementName, pipelineName)); diff --git a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementTable.java b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementTable.java index ad5dde619..6eab71fc6 100644 --- a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementTable.java +++ b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementTable.java @@ -1,9 +1,9 @@ package com.linkedin.hoptimator.k8s; +import com.linkedin.hoptimator.util.RemoteTable; import org.apache.calcite.schema.Schema; +import org.json.simple.JSONObject; -import com.linkedin.hoptimator.k8s.status.K8sPipelineElementStatus; -import com.linkedin.hoptimator.util.RemoteTable; public class K8sPipelineElementTable extends RemoteTable { @@ -13,12 +13,14 @@ public static class Row { public boolean READY; public boolean FAILED; public String MESSAGE; + public String CONFIGS; - public Row(String name, boolean ready, boolean failed, String message) { + public Row(String name, boolean ready, boolean failed, String message, String configs) { this.NAME = name; this.READY = ready; this.FAILED = failed; this.MESSAGE = message; + this.CONFIGS = configs; } @Override @@ -34,8 +36,11 @@ public K8sPipelineElementTable(K8sPipelineElementApi pipelineElementApi) { @Override public Row toRow(K8sPipelineElement k8sPipelineElement) { - K8sPipelineElementStatus status = k8sPipelineElement.status(); - return new Row(status.getName(), status.isReady(), status.isFailed(), status.getMessage()); + return new Row(k8sPipelineElement.name(), + k8sPipelineElement.status().isReady(), + k8sPipelineElement.status().isFailed(), + k8sPipelineElement.status().getMessage(), + JSONObject.toJSONString(k8sPipelineElement.configs())); } @Override diff --git a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/status/K8sPipelineElementStatusEstimator.java b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/status/K8sPipelineElementStatusEstimator.java index 7661768f3..a860d816c 100644 --- a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/status/K8sPipelineElementStatusEstimator.java +++ b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/status/K8sPipelineElementStatusEstimator.java @@ -47,7 +47,7 @@ public List estimateStatuses(V1alpha1Pipeline pipeline /** * Estimates status of an element. If we can not retrieve it from K8s, we assume that it's not ready and not failed yet. */ - private K8sPipelineElementStatus estimateElementStatus(String elementYaml, String pipelineNamespace) { + public K8sPipelineElementStatus estimateElementStatus(String elementYaml, String pipelineNamespace) { DynamicKubernetesObject obj = Dynamics.newFromYaml(elementYaml); String name = obj.getMetadata().getName(); String namespace = obj.getMetadata().getNamespace() == null ? pipelineNamespace : obj.getMetadata().getNamespace(); diff --git a/hoptimator-k8s/src/test/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApiTest.java b/hoptimator-k8s/src/test/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApiTest.java new file mode 100644 index 000000000..599531a80 --- /dev/null +++ b/hoptimator-k8s/src/test/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApiTest.java @@ -0,0 +1,303 @@ +package com.linkedin.hoptimator.k8s; + +import com.google.gson.JsonObject; +import com.linkedin.hoptimator.k8s.models.V1alpha1Pipeline; +import com.linkedin.hoptimator.k8s.models.V1alpha1PipelineSpec; +import io.kubernetes.client.common.KubernetesType; +import io.kubernetes.client.openapi.models.V1ObjectMeta; +import io.kubernetes.client.util.generic.KubernetesApiResponse; +import io.kubernetes.client.util.generic.dynamic.DynamicKubernetesApi; +import io.kubernetes.client.util.generic.dynamic.DynamicKubernetesObject; +import java.util.List; +import java.util.Map; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.MockedStatic; +import org.mockito.junit.jupiter.MockitoExtension; +import org.mockito.junit.jupiter.MockitoSettings; +import org.mockito.quality.Strictness; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +/** + * Unit tests for K8sPipelineElementApi + */ +@ExtendWith(MockitoExtension.class) +@MockitoSettings(strictness = Strictness.LENIENT) +class K8sPipelineElementApiTest { + + @Mock + private K8sContext mockContext; + + @Mock + private KubernetesApiResponse mockApiResponse; + + @Mock + private DynamicKubernetesObject mockDynamicObject; + + @Mock + private MockedStatic k8sUtilsMockedStatic; + + private K8sPipelineElementApi api; + + @BeforeEach + void setUp() { + api = new K8sPipelineElementApi(mockContext); + } + + @Test + void testGetPipelineElementsSingleElement() { + // Given + V1alpha1Pipeline pipeline = createPipeline("test-pipeline", "default", + "apiVersion: apps/v1\n" + + "kind: Deployment\n" + + "metadata:\n" + + " name: test-deployment"); + + // When + List elements = api.getPipelineElements(pipeline); + + // Then + assertEquals(1, elements.size()); + assertEquals("apiVersion: apps/v1\nkind: Deployment\nmetadata:\n name: test-deployment", elements.get(0)); + } + + @Test + void testGetPipelineElementsMultipleElements() { + // Given + String yaml = "apiVersion: apps/v1\n" + + "kind: Deployment\n" + + "metadata:\n" + + " name: test-deployment\n" + + "---\n" + + "apiVersion: v1\n" + + "kind: Service\n" + + "metadata:\n" + + " name: test-service\n" + + "---\n" + + "apiVersion: v1\n" + + "kind: ConfigMap\n" + + "metadata:\n" + + " name: test-config"; + + V1alpha1Pipeline pipeline = createPipeline("test-pipeline", "default", yaml); + + // When + List elements = api.getPipelineElements(pipeline); + + // Then + assertEquals(3, elements.size()); + assertTrue(elements.get(0).contains("kind: Deployment")); + assertTrue(elements.get(1).contains("kind: Service")); + assertTrue(elements.get(2).contains("kind: ConfigMap")); + } + + @Test + void testGetPipelineElementsWithEmptyElements() { + // Given + String yaml = "apiVersion: apps/v1\n" + + "kind: Deployment\n" + + "metadata:\n" + + " name: test-deployment\n" + + "---\n" + + "\n" + + "---\n" + + "apiVersion: v1\n" + + "kind: Service\n" + + "metadata:\n" + + " name: test-service"; + + V1alpha1Pipeline pipeline = createPipeline("test-pipeline", "default", yaml); + + // When + List elements = api.getPipelineElements(pipeline); + + // Then + assertEquals(2, elements.size()); // Empty elements should be filtered out + assertTrue(elements.get(0).contains("kind: Deployment")); + assertTrue(elements.get(1).contains("kind: Service")); + } + + @Test + void testGetPipelineElementsNullYaml() { + // Given + V1alpha1Pipeline pipeline = createPipeline("test-pipeline", "default", null); + + // When & Then + assertThrows(NullPointerException.class, () -> api.getPipelineElements(pipeline)); + } + + @Test + void testGetElementConfigurationWithConfigs() { + // Given + String elementYaml = "apiVersion: hoptimator.linkedin.com/v1alpha1\n" + + "kind: SqlJob\n" + + "metadata:\n" + + " name: test-sqljob\n" + + " namespace: default"; + + JsonObject configsJson = new JsonObject(); + configsJson.addProperty("key1", "value1"); + configsJson.addProperty("key2", "value2"); + + JsonObject specJson = new JsonObject(); + specJson.add("configs", configsJson); + + JsonObject rootJson = new JsonObject(); + rootJson.add("spec", specJson); + + setupMockForElementConfiguration(rootJson, true); + + // When + Map result = api.getElementConfiguration(elementYaml, "default"); + + // Then + assertEquals(2, result.size()); + assertEquals("value1", result.get("key1")); + assertEquals("value2", result.get("key2")); + } + + @Test + void testGetElementConfigurationNoConfigs() { + // Given + String elementYaml = "apiVersion: apps/v1\n" + + "kind: Deployment\n" + + "metadata:\n" + + " name: test-deployment\n" + + " namespace: default"; + + JsonObject specJson = new JsonObject(); + JsonObject rootJson = new JsonObject(); + rootJson.add("spec", specJson); + + setupMockForElementConfiguration(rootJson, true); + + // When + Map result = api.getElementConfiguration(elementYaml, "default"); + + // Then + assertTrue(result.isEmpty()); + } + + @Test + void testGetElementConfigurationNoSpec() { + // Given + String elementYaml = "apiVersion: apps/v1\n" + + "kind: Deployment\n" + + "metadata:\n" + + " name: test-deployment\n" + + " namespace: default"; + + JsonObject rootJson = new JsonObject(); + + setupMockForElementConfiguration(rootJson, true); + + // When + Map result = api.getElementConfiguration(elementYaml, "default"); + + // Then + assertTrue(result.isEmpty()); + } + + @Test + void testGetElementConfigurationApiCallFails() { + // Given + String elementYaml = "apiVersion: apps/v1\n" + + "kind: Deployment\n" + + "metadata:\n" + + " name: test-deployment\n" + + " namespace: default"; + + setupMockForElementConfiguration(null, false); + + // When + Map result = api.getElementConfiguration(elementYaml, "default"); + + // Then + assertTrue(result.isEmpty()); + } + + @Test + void testGetElementConfigurationNullElement() { + // Given + String elementYaml = "apiVersion: apps/v1\n" + + "kind: Deployment\n" + + "metadata:\n" + + " name: test-deployment\n" + + " namespace: default"; + + k8sUtilsMockedStatic.when(() -> K8sUtils.guessPlural(any(DynamicKubernetesObject.class))).thenReturn("deployments"); + + when(mockApiResponse.isSuccess()).thenReturn(true); + when(mockApiResponse.getObject()).thenReturn(null); + + DynamicKubernetesApi mockDynamicApi = mock(DynamicKubernetesApi.class); + when(mockDynamicApi.get(anyString(), anyString())).thenReturn(mockApiResponse); + when(mockContext.dynamic(anyString(), anyString())).thenReturn(mockDynamicApi); + + // When + Map result = api.getElementConfiguration(elementYaml, "default"); + + // Then + assertTrue(result.isEmpty()); + } + + @Test + void testGetElementConfigurationExceptionThrown() { + // Given + String elementYaml = "apiVersion: apps/v1\n" + + "kind: Deployment\n" + + "metadata:\n" + + " name: test-deployment\n" + + " namespace: default"; + + k8sUtilsMockedStatic.when(() -> K8sUtils.guessPlural(any(DynamicKubernetesObject.class))).thenReturn("deployments"); + + when(mockContext.dynamic(anyString(), anyString())).thenThrow(new RuntimeException("API error")); + + // When + Map result = api.getElementConfiguration(elementYaml, "default"); + + // Then + assertTrue(result.isEmpty()); + } + + // Helper methods + + private V1alpha1Pipeline createPipeline(String name, String namespace, String yaml) { + V1alpha1Pipeline pipeline = new V1alpha1Pipeline(); + + V1ObjectMeta metadata = new V1ObjectMeta(); + metadata.setName(name); + metadata.setNamespace(namespace); + pipeline.setMetadata(metadata); + + V1alpha1PipelineSpec spec = new V1alpha1PipelineSpec(); + spec.setYaml(yaml); + pipeline.setSpec(spec); + + return pipeline; + } + + private void setupMockForElementConfiguration(JsonObject rootJson, boolean success) { + k8sUtilsMockedStatic.when(() -> K8sUtils.guessPlural(any(KubernetesType.class))).thenReturn("deployments"); + + when(mockApiResponse.isSuccess()).thenReturn(success); + if (success && rootJson != null) { + when(mockApiResponse.getObject()).thenReturn(mockDynamicObject); + when(mockDynamicObject.getRaw()).thenReturn(rootJson); + } + + DynamicKubernetesApi mockDynamicApi = mock(DynamicKubernetesApi.class); + when(mockDynamicApi.get(anyString(), anyString())).thenReturn(mockApiResponse); + when(mockContext.dynamic(anyString(), anyString())).thenReturn(mockDynamicApi); + } +} \ No newline at end of file From 16513578570a96d986d79f45ef6f3e002433df8b Mon Sep 17 00:00:00 2001 From: Joseph Grogan Date: Fri, 5 Sep 2025 17:51:31 -0400 Subject: [PATCH 2/2] Fix test failures --- .../java/com/linkedin/hoptimator/k8s/K8sPipelineElementApi.java | 1 + 1 file changed, 1 insertion(+) diff --git a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApi.java b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApi.java index 3f2843f99..0dc0b1041 100644 --- a/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApi.java +++ b/hoptimator-k8s/src/main/java/com/linkedin/hoptimator/k8s/K8sPipelineElementApi.java @@ -48,6 +48,7 @@ private Collection discoverAllElements(K8sContext context) t if (!elements.containsKey(key)) { elements.put(key, new K8sPipelineElement(pipeline, elementStatus, configurations)); } + elements.get(key).addPipeline(pipeline); } } return elements.values();