Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/integration.yml
Original file line number Diff line number Diff line change
Expand Up @@ -166,7 +166,7 @@ jobs:
# See: https://quarkus.io/guides/tests-with-coverage#coverage-for-integration-tests
#
mvn verify -B --no-transfer-progress -DskipSTs \
-Dquarkus.package.write-transformed-bytecode-to-build-output=true
-Dquarkus.package.write-transformed-bytecode-to-build-output=false

- name: Archive Failed Tests Results
uses: actions/upload-artifact@v7
Expand Down
3 changes: 2 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -36,4 +36,5 @@ release.properties
# Systemtests
systemtests/screenshots/
systemtests/config.yaml
systemtests/tracing/
systemtests/tracing/
/.playwright-mcp/
25 changes: 25 additions & 0 deletions api/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -306,6 +306,7 @@
<execution>
<goals>
<goal>build</goal>
<goal>generate-code</goal>
</goals>
<configuration>
<properties>
Expand All @@ -319,6 +320,30 @@
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-dependency-plugin</artifactId>
<executions>
<execution>
<id>unpack</id>
<phase>initialize</phase>
<goals>
<goal>unpack</goal>
</goals>
<configuration>
<artifactItems>
<artifactItem>
<!--
This artifact is not published to Maven Central and must be retrieve from
https://linkedin.jfrog.io/artifactory/release instead :-(
-->
<groupId>com.linkedin.cruisecontrol</groupId>
<artifactId>cruise-control</artifactId>
<version>2.5.146</version>
<type>jar</type>
<overWrite>true</overWrite>
<outputDirectory>${project.build.directory}/schema/cruise-control</outputDirectory>
<includes>yaml/**</includes>
</artifactItem>
</artifactItems>
</configuration>
</execution>
<execution>
<id>analyze</id>
<goals>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,7 @@ public Response listRebalances(
listParams,
KafkaRebalance::fromCursor);

var rebalanceList = rebalanceService.listRebalances(listSupport);
var rebalanceList = rebalanceService.listRebalances(fields, listSupport);
var responseEntity = new KafkaRebalance.RebalanceDataList(rebalanceList, listSupport);

return Response.ok(responseEntity).build();
Expand All @@ -156,7 +156,7 @@ public Response listRebalances(
@APIResponse(responseCode = "504", ref = "ServerTimeout")
@Authorized
@ResourcePrivilege(Privilege.GET)
public Response getRebalance(
public Response describeRebalance(
@Parameter(description = "Cluster identifier")
@PathParam("clusterId")
String clusterId,
Expand All @@ -176,6 +176,7 @@ public Response getRebalance(
KafkaRebalance.Fields.STATUS,
KafkaRebalance.Fields.MODE,
KafkaRebalance.Fields.BROKERS,
KafkaRebalance.Fields.BROKER_CAPACITY,
KafkaRebalance.Fields.GOALS,
KafkaRebalance.Fields.SKIP_HARD_GOAL_CHECK,
KafkaRebalance.Fields.REBALANCE_DISK,
Expand All @@ -187,6 +188,7 @@ public Response getRebalance(
KafkaRebalance.Fields.REPLICA_MOVEMENT_STRATEGIES,
KafkaRebalance.Fields.SESSION_ID,
KafkaRebalance.Fields.OPTIMIZATION_RESULT,
KafkaRebalance.Fields.OPTIMIZATION_PROPOSAL,
KafkaRebalance.Fields.CONDITIONS,
},
payload = ErrorCategory.InvalidQueryParameter.class)
Expand All @@ -203,6 +205,7 @@ public Response getRebalance(
KafkaRebalance.Fields.STATUS,
KafkaRebalance.Fields.MODE,
KafkaRebalance.Fields.BROKERS,
KafkaRebalance.Fields.BROKER_CAPACITY,
KafkaRebalance.Fields.GOALS,
KafkaRebalance.Fields.SKIP_HARD_GOAL_CHECK,
KafkaRebalance.Fields.REBALANCE_DISK,
Expand All @@ -214,13 +217,14 @@ public Response getRebalance(
KafkaRebalance.Fields.REPLICA_MOVEMENT_STRATEGIES,
KafkaRebalance.Fields.SESSION_ID,
KafkaRebalance.Fields.OPTIMIZATION_RESULT,
KafkaRebalance.Fields.OPTIMIZATION_PROPOSAL,
KafkaRebalance.Fields.CONDITIONS,
}))
List<String> fields) {

requestedFields.accept(fields);

var result = rebalanceService.getRebalance(rebalanceId);
var result = rebalanceService.getRebalance(rebalanceId, fields);
var responseEntity = new KafkaRebalance.RebalanceData(result);

return Response.ok(responseEntity).build();
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
package com.github.streamshub.console.api.model;

import java.math.BigDecimal;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Comparator;
Expand Down Expand Up @@ -64,6 +65,7 @@ public static class Fields {
public static final String STATUS = "status";
public static final String MODE = "mode";
public static final String BROKERS = "brokers";
public static final String BROKER_CAPACITY = "brokerCapacity";
public static final String GOALS = "goals";
public static final String SKIP_HARD_GOAL_CHECK = "skipHardGoalCheck";
public static final String REBALANCE_DISK = "rebalanceDisk";
Expand All @@ -75,6 +77,7 @@ public static class Fields {
public static final String REPLICA_MOVEMENT_STRATEGIES = "replicaMovementStrategies";
public static final String SESSION_ID = "sessionId";
public static final String OPTIMIZATION_RESULT = "optimizationResult";
public static final String OPTIMIZATION_PROPOSAL = "optimizationProposal";
public static final String CONDITIONS = "conditions";

static final Comparator<KafkaRebalance> ID_COMPARATOR =
Expand Down Expand Up @@ -196,6 +199,46 @@ static final class Meta extends JsonApiMeta {
String action;
}

public static final record BrokerCapacityOverride(
@JsonProperty
List<Integer> brokers,
@JsonProperty
String cpu,
@JsonProperty
String inboundNetwork,
@JsonProperty
String outboundNetwork
) {
}

public static final record BrokerCapacity(
@JsonProperty
String cpu,
@JsonProperty
String inboundNetwork,
@JsonProperty
String outboundNetwork,
@JsonProperty
List<BrokerCapacityOverride> overrides
) {
}

public static final record BrokerLoadImpact(
@JsonProperty
BigDecimal before,
@JsonProperty
BigDecimal after,
@JsonProperty
BigDecimal diff
) {
}

public static final record OptimizationProposal(
@JsonProperty
Map<String, Map<String, BrokerLoadImpact>> brokerImpact
) {
}

@JsonFilter("fieldFilter")
@Schema(name = "KafkaRebalanceAttributes")
static class Attributes extends KubeAttributes {
Expand All @@ -211,6 +254,9 @@ static class Attributes extends KubeAttributes {
@Schema(readOnly = true, nullable = true)
List<Integer> brokers;

@JsonProperty
BrokerCapacity brokerCapacity;

@JsonProperty
@Schema(readOnly = true, nullable = true)
List<String> goals;
Expand Down Expand Up @@ -253,7 +299,11 @@ static class Attributes extends KubeAttributes {

@JsonProperty
@Schema(readOnly = true)
Map<String, Object> optimizationResult = new HashMap<>(0);
Map<String, Object> optimizationResult = HashMap.newHashMap(0);

@JsonProperty
@Schema(readOnly = true)
OptimizationProposal optimizationProposal;

@JsonProperty
@Schema(readOnly = true)
Expand Down Expand Up @@ -317,6 +367,10 @@ public void brokers(List<Integer> brokers) {
attributes.brokers = brokers;
}

public void brokerCapacity(BrokerCapacity brokerCapacity) {
attributes.brokerCapacity = brokerCapacity;
}

public void goals(List<String> goals) {
attributes.goals = goals;
}
Expand Down Expand Up @@ -361,6 +415,10 @@ public Map<String, Object> optimizationResult() {
return attributes.optimizationResult;
}

public void optimizationProposal(OptimizationProposal optimizationProposal) {
attributes.optimizationProposal = optimizationProposal;
}

public void conditions(List<Condition> conditions) {
attributes.conditions = conditions;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
import java.nio.charset.StandardCharsets;
import java.util.Arrays;
import java.util.Base64;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
Expand All @@ -17,8 +18,11 @@

import org.jboss.logging.Logger;

import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.github.streamshub.console.api.model.Condition;
import com.github.streamshub.console.api.model.KafkaRebalance;
import com.github.streamshub.console.api.model.KafkaRebalance.BrokerLoadImpact;
import com.github.streamshub.console.api.security.PermissionService;
import com.github.streamshub.console.api.support.KafkaContext;
import com.github.streamshub.console.api.support.ListRequestContext;
Expand All @@ -31,6 +35,8 @@
import io.strimzi.api.ResourceAnnotations;
import io.strimzi.api.ResourceLabels;
import io.strimzi.api.kafka.model.kafka.Kafka;
import io.strimzi.api.kafka.model.kafka.KafkaSpec;
import io.strimzi.api.kafka.model.kafka.cruisecontrol.CruiseControlSpec;
import io.strimzi.api.kafka.model.rebalance.KafkaRebalanceMode;
import io.strimzi.api.kafka.model.rebalance.KafkaRebalanceSpec;
import io.strimzi.api.kafka.model.rebalance.KafkaRebalanceState;
Expand All @@ -45,6 +51,9 @@ public class KafkaRebalanceService {
@Inject
KubernetesClient client;

@Inject
ObjectMapper mapper;

@Inject
ConsoleConfig consoleConfig;

Expand All @@ -54,7 +63,7 @@ public class KafkaRebalanceService {
@Inject
PermissionService permissionService;

public List<KafkaRebalance> listRebalances(ListRequestContext<KafkaRebalance> listSupport) {
public List<KafkaRebalance> listRebalances(List<String> fields, ListRequestContext<KafkaRebalance> listSupport) {
final Map<String, Integer> statuses = new HashMap<>();
listSupport.meta().put("summary", Map.of("statuses", statuses));

Expand All @@ -65,7 +74,7 @@ public List<KafkaRebalance> listRebalances(ListRequestContext<KafkaRebalance> li
ResourceTypes.Kafka.REBALANCES,
Privilege.LIST,
r -> r.getMetadata().getName()))
.map(this::toKafkaRebalance)
.map(r -> toKafkaRebalance(r, fields))
.map(rebalance -> tallyStatus(statuses, rebalance))
.filter(listSupport.filter(KafkaRebalance.class))
.map(listSupport::tally)
Expand All @@ -77,9 +86,9 @@ public List<KafkaRebalance> listRebalances(ListRequestContext<KafkaRebalance> li
.toList();
}

public KafkaRebalance getRebalance(String id) {
public KafkaRebalance getRebalance(String id, List<String> fields) {
return findRebalance(id)
.map(this::toKafkaRebalance)
.map(r -> toKafkaRebalance(r, fields))
.map(permissionService.addPrivileges(ResourceTypes.Kafka.REBALANCES, KafkaRebalance::name))
.orElseThrow(() -> new NotFoundException("No such Kafka rebalance resource"));
}
Expand All @@ -98,11 +107,11 @@ public KafkaRebalance patchRebalance(String id, KafkaRebalance rebalance) {

return client.resource(resource).patch();
})
.map(this::toKafkaRebalance)
.map(r -> toKafkaRebalance(r, Collections.emptyList()))
.orElseThrow(() -> new NotFoundException("No such Kafka rebalance resource"));
}

KafkaRebalance toKafkaRebalance(io.strimzi.api.kafka.model.rebalance.KafkaRebalance resource) {
KafkaRebalance toKafkaRebalance(io.strimzi.api.kafka.model.rebalance.KafkaRebalance resource, List<String> fields) {
KafkaRebalanceSpec rebalanceSpec = resource.getSpec();
Optional<KafkaRebalanceStatus> rebalanceStatus = Optional.ofNullable(resource.getStatus());
Optional<KafkaRebalanceState> state = rebalanceStatus
Expand All @@ -115,13 +124,32 @@ KafkaRebalance toKafkaRebalance(io.strimzi.api.kafka.model.rebalance.KafkaRebala
.findFirst();

String id = Base64.getUrlEncoder().encodeToString(Cache.metaNamespaceKeyFunc(resource).getBytes(StandardCharsets.UTF_8));
String namespace = resource.getMetadata().getNamespace();
KafkaRebalance rebalance = new KafkaRebalance(id);
rebalance.name(resource.getMetadata().getName());
rebalance.namespace(resource.getMetadata().getNamespace());
rebalance.namespace(namespace);
rebalance.creationTimestamp(resource.getMetadata().getCreationTimestamp());
rebalance.status(state.map(Enum::name).orElse(null));
rebalance.mode(Optional.ofNullable(rebalanceSpec.getMode()).map(KafkaRebalanceMode::toValue).orElse(null));
rebalance.brokers(rebalanceSpec.getBrokers());
rebalance.brokerCapacity(Optional.ofNullable(kafkaContext.resource())
.map(Kafka::getSpec)
.map(KafkaSpec::getCruiseControl)
.map(CruiseControlSpec::getBrokerCapacity)
.map(capacity -> new KafkaRebalance.BrokerCapacity(
capacity.getCpu(),
capacity.getInboundNetwork(),
capacity.getOutboundNetwork(),
Optional.ofNullable(capacity.getOverrides())
.orElseGet(Collections::emptyList)
.stream()
.map(override -> new KafkaRebalance.BrokerCapacityOverride(
override.getBrokers(),
override.getCpu(),
override.getInboundNetwork(),
override.getOutboundNetwork()))
.toList()))
.orElse(null));
rebalance.goals(rebalanceSpec.getGoals());
rebalance.skipHardGoalCheck(rebalanceSpec.isSkipHardGoalCheck());
rebalance.rebalanceDisk(rebalanceSpec.isRebalanceDisk());
Expand All @@ -148,9 +176,43 @@ KafkaRebalance toKafkaRebalance(io.strimzi.api.kafka.model.rebalance.KafkaRebala
.map(allowed -> allowed.stream().map(Enum::name).toList())
.ifPresent(rebalance.allowedActions()::addAll);

if (fields.contains(KafkaRebalance.Fields.OPTIMIZATION_PROPOSAL)) {
rebalance.optimizationProposal(getOptimizationProposal(namespace, rebalanceStatus));
}

return rebalance;
}

private KafkaRebalance.OptimizationProposal getOptimizationProposal(String namespace, Optional<KafkaRebalanceStatus> rebalanceStatus) {
return rebalanceStatus
.map(KafkaRebalanceStatus::getOptimizationResult)
.map(result -> result.get("afterBeforeLoadConfigMap"))
.filter(String.class::isInstance)
.map(String.class::cast)
.map(configMapName -> client.configMaps().inNamespace(namespace).withName(configMapName).get())
.map(configMap -> {
var data = configMap.getData();
var qname = "%s/%s".formatted(configMap.getMetadata().getNamespace(), configMap.getMetadata().getName());
var brokerLoadImpact = Optional.ofNullable(data.get("brokerLoad.json"))
.map(json -> {
try {
return mapper.readValue(json, new TypeReference<Map<String, Map<String, BrokerLoadImpact>>>() {
// No implementation
});
} catch (Exception e) {
logger.warnf("""
Error reading 'brokerLoad.json' from rebalance \
afterBeforeLoadConfigMap ConfigMap[%s]: %s""", qname, e.getMessage());
return null;
}
})
.orElse(null);

return new KafkaRebalance.OptimizationProposal(brokerLoadImpact);
})
.orElse(null);
}

KafkaRebalance tallyStatus(Map<String, Integer> statuses, KafkaRebalance rebalance) {
String status = rebalance.status();
if (status != null) {
Expand Down Expand Up @@ -214,3 +276,4 @@ private boolean isTemplate(io.strimzi.api.kafka.model.rebalance.KafkaRebalance r
.orElse(false);
}
}

Loading
Loading