diff --git a/design/BulkAPI.md b/design/BulkAPI.md index 3c6a2c92d7..5dc0540589 100644 --- a/design/BulkAPI.md +++ b/design/BulkAPI.md @@ -60,8 +60,7 @@ progress of the job. "metadata_profile": "cluster-metadata-local-monitoring", "measurement_duration": "15min", "experiment_types": [ - "container", - "namespace" + "container" ], "cluster_name": "prod-cluster", "model_settings": { @@ -94,7 +93,9 @@ progress of the job. - **datasource:** The data source, e.g., `"Cbank1Xyz"`. -- **experiment_types:** Specifies the type(s) of experiments to run, e.g., `"container"` or `"namespace"`. +- **experiment_types:** Specifies the type of experiment to create for this bulk job. Currently, only a single value is + supported per request — either `"container"` (default) or `"namespace"`. Support for specifying multiple experiment + types in one request will be added in a future release. - **webhook:** The `webhook` parameter allows the system to notify an external service or consumer about the completion status of diff --git a/src/main/java/com/autotune/analyzer/serviceObjects/BulkInput.java b/src/main/java/com/autotune/analyzer/serviceObjects/BulkInput.java index e8b344cd45..3c965334ff 100644 --- a/src/main/java/com/autotune/analyzer/serviceObjects/BulkInput.java +++ b/src/main/java/com/autotune/analyzer/serviceObjects/BulkInput.java @@ -17,6 +17,7 @@ import com.autotune.analyzer.kruizeObject.ModelSettings; import com.autotune.analyzer.kruizeObject.TermSettings; +import com.autotune.analyzer.utils.AnalyzerConstants; import com.fasterxml.jackson.annotation.JsonIgnore; import java.util.List; import java.util.Map; @@ -32,26 +33,33 @@ public class BulkInput { private String metadata_profile; private String measurement_duration; private String requestId; //TODO: to be used for the Kafka consumer case to map requestID with jobID - + /** * Cluster name to use for all experiments in this bulk job. * If provided, overrides cluster name from datasource metadata. * If not provided, cluster name from metadata will be used. */ private String cluster_name; - + /** * Optional model settings to customize which recommendation models to generate. * If not provided, all models will be generated. */ private ModelSettings model_settings; - + /** * Optional term settings to customize which recommendation terms to generate. * If not provided, all terms will be generated. */ private TermSettings term_settings; + /** + * Experiment types to create in this bulk job (e.g. CONTAINER, NAMESPACE). + * If provided, only experiments of the specified type(s) will be created. + * If not provided or empty, defaults to container experiments. + */ + private List experiment_types; + // Getters and Setters public String getRequestId() { @@ -68,7 +76,8 @@ public BulkInput() { @JsonIgnore public boolean isEmpty() { return (filter == null && time_range == null && measurement_duration == null && metadata_profile == null - && datasource == null && cluster_name == null && model_settings == null && term_settings == null); + && datasource == null && cluster_name == null && model_settings == null && term_settings == null + && experiment_types == null); } public TimeRange getTime_range() { @@ -135,6 +144,17 @@ public void setTerm_settings(TermSettings term_settings) { this.term_settings = term_settings; } + + public List getExperiment_types() { + return experiment_types; + } + + public void setExperiment_types(List experiment_types) { + if (experiment_types != null && !experiment_types.isEmpty()) { + this.experiment_types = experiment_types; + } + } + // Nested class for FilterWrapper that contains 'exclude' and 'include' public static class FilterWrapper { private Filter exclude; diff --git a/src/main/java/com/autotune/analyzer/utils/AnalyzerConstants.java b/src/main/java/com/autotune/analyzer/utils/AnalyzerConstants.java index 829bbc4bf2..046f37c2f5 100644 --- a/src/main/java/com/autotune/analyzer/utils/AnalyzerConstants.java +++ b/src/main/java/com/autotune/analyzer/utils/AnalyzerConstants.java @@ -17,6 +17,7 @@ import com.autotune.utils.KruizeConstants; import com.autotune.utils.Utils; +import com.fasterxml.jackson.annotation.JsonCreator; import java.util.Map; import java.util.*; @@ -326,7 +327,17 @@ public enum ExperimentType { CONTAINER, // For container-level experiments NAMESPACE, // For namespace-level experiments CLUSTER, // For cluster-wide experiments - WORKLOAD // For application-specific experiments + WORKLOAD; // For application-specific experiments + + @JsonCreator + public static ExperimentType fromString(String value) { + if (value == null || value.trim().isEmpty()) return null; + try { + return ExperimentType.valueOf(value.trim().toUpperCase()); + } catch (IllegalArgumentException e) { + return null; + } + } } /** diff --git a/src/main/java/com/autotune/analyzer/utils/AnalyzerErrorConstants.java b/src/main/java/com/autotune/analyzer/utils/AnalyzerErrorConstants.java index 198014c969..b01d016ec6 100644 --- a/src/main/java/com/autotune/analyzer/utils/AnalyzerErrorConstants.java +++ b/src/main/java/com/autotune/analyzer/utils/AnalyzerErrorConstants.java @@ -215,6 +215,7 @@ public static final class CreateExperimentAPI { public static final String MISSING_NAMESPACE_DATA = "Missing NamespaceData for experimentType: %s"; public static final String MISSING_NAMESPACE = "Missing namespace for experimentType: %s"; public static final String INVALID_EXPERIMENT_TYPE = "Invalid experiment_type : %s"; + public static final String BULK_INVALID_EXPERIMENT_TYPES = "Invalid experiment_types value(s): %s. Valid values are: container, namespace"; private CreateExperimentAPI() { diff --git a/src/main/java/com/autotune/analyzer/workerimpl/BulkJobManager.java b/src/main/java/com/autotune/analyzer/workerimpl/BulkJobManager.java index a3df491bd0..adaefabbda 100644 --- a/src/main/java/com/autotune/analyzer/workerimpl/BulkJobManager.java +++ b/src/main/java/com/autotune/analyzer/workerimpl/BulkJobManager.java @@ -189,7 +189,7 @@ public void run() { setFinalJobStatus(COMPLETED, String.valueOf(HttpURLConnection.HTTP_OK), NOTHING_INFO, datasource); } else { jobData.setMetadata(metadataInfo); - Map createExperimentAPIObjectMap = getExperimentMap(labelString, jobData, metadataInfo, datasource); //Todo Store this map in buffer and use it if BulkAPI pods restarts and support experiment_type + Map createExperimentAPIObjectMap = getExperimentMap(labelString, jobData, metadataInfo, datasource); //Todo Store this map in buffer and use it if BulkAPI pods restarts // TODO: Remove getExperimentMap and instead collect all metadata, process it, and create experiments dynamically during metadata iteration. jobData.getSummary().setTotal_experiments(createExperimentAPIObjectMap.size()); jobData.getSummary().setProcessed_experiments(0); @@ -511,6 +511,11 @@ Map getExperimentMap(String labelString, Bulk String statusValue = "failure"; Timer.Sample timerGetExpMap = Timer.start(MetricsConfig.meterRegistry()); try { + List experimentTypes = this.bulkInput.getExperiment_types(); + AnalyzerConstants.ExperimentType experimentType = (experimentTypes == null || experimentTypes.isEmpty()) + ? AnalyzerConstants.ExperimentType.CONTAINER + : experimentTypes.get(0); + Map createExperimentAPIObjectMap = new HashMap<>(); Collection dataSourceCollection = metadataInfo.getDatasources().values(); for (DataSource ds : dataSourceCollection) { @@ -521,19 +526,29 @@ Map getExperimentMap(String labelString, Bulk : dsc.getDataSourceClusterName(); HashMap namespaceHashMap = dsc.getNamespaces(); for (DataSourceNamespace namespace : namespaceHashMap.values()) { - HashMap dataSourceWorkloadHashMap = namespace.getWorkloads(); - if (dataSourceWorkloadHashMap != null) { - for (DataSourceWorkload dsw : dataSourceWorkloadHashMap.values()) { - HashMap dataSourceContainerHashMap = dsw.getContainers(); - if (dataSourceContainerHashMap != null) { - for (DataSourceContainer dc : dataSourceContainerHashMap.values()) { - // Experiment name - dynamically constructed - String experiment_name = frameExperimentName(labelString, clusterName, namespace, dsw, dc); - // create JSON to be passed in the createExperimentAPI - List createExperimentAPIObjectList = new ArrayList<>(); - CreateExperimentAPIObject apiObject = prepareCreateExperimentJSONInput(dc, clusterName, dsw, namespace, - experiment_name, createExperimentAPIObjectList); - createExperimentAPIObjectMap.put(experiment_name, apiObject); + if (experimentType == AnalyzerConstants.ExperimentType.NAMESPACE) { + // One namespace experiment per namespace + String experiment_name = frameNamespaceExperimentName(labelString, clusterName, namespace); + List createExperimentAPIObjectList = new ArrayList<>(); + CreateExperimentAPIObject apiObject = prepareNamespaceExperimentJSONInput(clusterName, namespace, + experiment_name, createExperimentAPIObjectList); + createExperimentAPIObjectMap.put(experiment_name, apiObject); + } else { + // Default: one container experiment per container + HashMap dataSourceWorkloadHashMap = namespace.getWorkloads(); + if (dataSourceWorkloadHashMap != null) { + for (DataSourceWorkload dsw : dataSourceWorkloadHashMap.values()) { + HashMap dataSourceContainerHashMap = dsw.getContainers(); + if (dataSourceContainerHashMap != null) { + for (DataSourceContainer dc : dataSourceContainerHashMap.values()) { + // Experiment name - dynamically constructed + String experiment_name = frameExperimentName(labelString, clusterName, namespace, dsw, dc); + // create JSON to be passed in the createExperimentAPI + List createExperimentAPIObjectList = new ArrayList<>(); + CreateExperimentAPIObject apiObject = prepareCreateExperimentJSONInput(dc, clusterName, dsw, namespace, + experiment_name, createExperimentAPIObjectList); + createExperimentAPIObjectMap.put(experiment_name, apiObject); + } } } } @@ -725,11 +740,11 @@ private CreateExperimentAPIObject prepareCreateExperimentJSONInput(DataSourceCon kubernetesAPIObject.setNamespace(namespace.getNamespace()); kubernetesAPIObjectList.add(kubernetesAPIObject); createExperimentAPIObject.setKubernetesObjects(kubernetesAPIObjectList); - + // Create recommendation settings with threshold RecommendationSettings rs = new RecommendationSettings(); rs.setThreshold(CREATE_EXPERIMENT_CONFIG_BEAN.getThreshold()); - + // Pass through model_settings and term_settings from bulk payload if provided if (bulkInput.getModel_settings() != null) { rs.setModelSettings(bulkInput.getModel_settings()); @@ -737,7 +752,7 @@ private CreateExperimentAPIObject prepareCreateExperimentJSONInput(DataSourceCon if (bulkInput.getTerm_settings() != null) { rs.setTermSettings(bulkInput.getTerm_settings()); } - + createExperimentAPIObject.setRecommendationSettings(rs); TrialSettings trialSettings = new TrialSettings(); trialSettings.setMeasurement_durationMinutes(CREATE_EXPERIMENT_CONFIG_BEAN.getMeasurementDurationStr()); @@ -752,9 +767,100 @@ private CreateExperimentAPIObject prepareCreateExperimentJSONInput(DataSourceCon return createExperimentAPIObject; } + + /** + * Builds a CreateExperimentAPIObject for a namespace-level experiment. + * The kubernetes_objects entry contains only a namespaces block (no + * workload name/type or containers), matching the payload expected by + * CreateExperiment for experiment_type "namespace". + * + * @param clusterName resolved cluster name (user override or data source cluster name) + * @param namespace DataSourceNamespace whose namespace is being tracked + * @param experiment_name pre-framed experiment name + * @param createExperimentAPIObjects accumulator list + * @return the constructed CreateExperimentAPIObject + */ + private CreateExperimentAPIObject prepareNamespaceExperimentJSONInput(String clusterName, DataSourceNamespace namespace, + String experiment_name, List createExperimentAPIObjects) throws IOException { + CreateExperimentAPIObject createExperimentAPIObject = new CreateExperimentAPIObject(); + createExperimentAPIObject.setMode(CREATE_EXPERIMENT_CONFIG_BEAN.getMode()); + createExperimentAPIObject.setTargetCluster(CREATE_EXPERIMENT_CONFIG_BEAN.getTarget()); + createExperimentAPIObject.setApiVersion(CREATE_EXPERIMENT_CONFIG_BEAN.getVersion()); + createExperimentAPIObject.setExperimentName(experiment_name); + createExperimentAPIObject.setDatasource(this.bulkInput.getDatasource()); + createExperimentAPIObject.setClusterName(clusterName); + createExperimentAPIObject.setPerformanceProfile(CREATE_EXPERIMENT_CONFIG_BEAN.getPerformanceProfile()); + createExperimentAPIObject.setMetadataProfile(CREATE_EXPERIMENT_CONFIG_BEAN.getMetadataProfile()); + + // Namespace experiment: kubernetes_objects has only a namespaces block, no containers + List kubernetesAPIObjectList = new ArrayList<>(); + KubernetesAPIObject kubernetesAPIObject = new KubernetesAPIObject(); + NamespaceAPIObject namespaceAPIObject = new NamespaceAPIObject(namespace.getNamespace(), null, null); + kubernetesAPIObject.setNamespaceAPIObject(namespaceAPIObject); + kubernetesAPIObjectList.add(kubernetesAPIObject); + createExperimentAPIObject.setKubernetesObjects(kubernetesAPIObjectList); + + // Recommendation settings + RecommendationSettings rs = new RecommendationSettings(); + rs.setThreshold(CREATE_EXPERIMENT_CONFIG_BEAN.getThreshold()); + if (this.bulkInput.getModel_settings() != null) { + rs.setModelSettings(this.bulkInput.getModel_settings()); + } + if (this.bulkInput.getTerm_settings() != null) { + rs.setTermSettings(this.bulkInput.getTerm_settings()); + } + createExperimentAPIObject.setRecommendationSettings(rs); + + TrialSettings trialSettings = new TrialSettings(); + trialSettings.setMeasurement_durationMinutes(CREATE_EXPERIMENT_CONFIG_BEAN.getMeasurementDurationStr()); + createExperimentAPIObject.setTrialSettings(trialSettings); + + createExperimentAPIObject.setExperiment_id(Utils.generateID(createExperimentAPIObject.toString())); + createExperimentAPIObject.setStatus(AnalyzerConstants.ExperimentStatus.IN_PROGRESS); + createExperimentAPIObject.setExperimentType(AnalyzerConstants.ExperimentType.NAMESPACE); + + createExperimentAPIObjects.add(createExperimentAPIObject); + return createExperimentAPIObject; + } + + /** + * Frames the experiment name for a namespace-level experiment. + * Uses datasource, cluster name, and namespace — workload/container + * segments are not meaningful for namespace experiments. + * + * @param labelString label filter string (may be null) + * @param clusterName resolved cluster name (user override or data source cluster name) + * @param namespace namespace metadata + * @return framed experiment name + */ + public String frameNamespaceExperimentName(String labelString, String clusterName, + DataSourceNamespace namespace) { + String datasource = this.bulkInput.getDatasource(); + String namespaceName = namespace.getNamespace(); + + // Namespace experiment name: datasource|clustername|namespace + String experimentName = KruizeDeploymentInfo.namespace_experiment_name_format + .replace("%datasource%", datasource) + .replace("%clustername%", clusterName) + .replace("%namespace%", namespaceName); + + if (null != labelString) { + Map labelsMap = parseLabelString(labelString); + Pattern labelPattern = Pattern.compile("%label:([a-zA-Z0-9_]+)%"); + Matcher matcher = labelPattern.matcher(experimentName); + while (matcher.find()) { + String labelKey = matcher.group(1); + String labelValue = labelsMap.getOrDefault(labelKey, "unknown" + labelKey); + experimentName = experimentName.replace(matcher.group(), labelValue != null ? labelValue : "unknown" + labelKey); + } + } + LOGGER.debug("Namespace experiment name: {}", experimentName); + return experimentName; + } + /** * @param labelString - * @param dataSourceCluster + * @param clusterName * @param dataSourceNamespace * @param dataSourceWorkload * @param dataSourceContainer diff --git a/src/main/java/com/autotune/common/bulk/BulkServiceValidation.java b/src/main/java/com/autotune/common/bulk/BulkServiceValidation.java index e2352b6d8c..980c7111ea 100644 --- a/src/main/java/com/autotune/common/bulk/BulkServiceValidation.java +++ b/src/main/java/com/autotune/common/bulk/BulkServiceValidation.java @@ -18,19 +18,22 @@ import com.autotune.analyzer.kruizeObject.ModelSettings; import com.autotune.analyzer.kruizeObject.TermSettings; import com.autotune.analyzer.serviceObjects.BulkInput; +import com.autotune.analyzer.utils.AnalyzerErrorConstants; import com.autotune.common.data.ValidationOutputData; import com.autotune.common.datasource.DataSourceInfo; import com.autotune.common.datasource.DataSourceOperatorImpl; import com.autotune.common.utils.CommonUtils; import com.autotune.database.service.ExperimentDBService; +import com.autotune.analyzer.utils.AnalyzerConstants; import com.autotune.utils.KruizeConstants; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.time.OffsetDateTime; import java.time.format.DateTimeParseException; -import java.util.Arrays; import java.util.List; +import java.util.Set; +import java.util.Arrays; /** * Utility class that performs validation for bulk service requests. @@ -49,6 +52,13 @@ public class BulkServiceValidation { private static final Logger LOGGER = LoggerFactory.getLogger(BulkServiceValidation.class); + // Cluster name validation constants + private static final int MAX_CLUSTER_NAME_LENGTH = 253; + + // Valid experiment types supported by the bulk engine + private static final Set VALID_EXPERIMENT_TYPES = Set.of( + AnalyzerConstants.ExperimentType.CONTAINER, AnalyzerConstants.ExperimentType.NAMESPACE); + // Valid model and term names (case-insensitive) private static final List VALID_MODELS = Arrays.asList( KruizeConstants.JSONKeys.PERFORMANCE, @@ -102,6 +112,10 @@ public static ValidationOutputData validate(BulkInput payload, String jobID) thr validationOutputData = buildErrorOutput(validateTermSettings(payload.getTerm_settings()), jobID); if (validationOutputData != null) return validationOutputData; + // Validate experiment_types if provided + validationOutputData = buildErrorOutput(validateExperimentTypes(payload.getExperiment_types()), jobID); + if (validationOutputData != null) return validationOutputData; + if (payload.getDatasource() != null) { validationOutputData = buildErrorOutput(validateDatasourceConnection(payload.getDatasource()), jobID); } @@ -292,4 +306,42 @@ public static String validateTermSettings(TermSettings termSettings) { return ""; } + /** + * Validates the experiment_types field if provided. + * Checks for: + *
    + *
  • Exactly one entry — multiple types are not supported; the bulk engine + * creates experiments of a single type per job
  • + *
  • Non-null, non-empty entry value
  • + *
  • Valid type name (container, namespace)
  • + *
+ * + * @param experimentTypes the list of experiment types to validate (can be null) + * @return an error message if validation fails; otherwise an empty string + */ + public static String validateExperimentTypes(List experimentTypes) { + if (experimentTypes == null || experimentTypes.isEmpty()) { + return ""; // null/empty is valid; defaults to container experiments + } + + if (experimentTypes.size() > 1) { + return "experiment_types accepts at most one value per bulk job. " + + "Provided: " + experimentTypes; + } + + AnalyzerConstants.ExperimentType type = experimentTypes.get(0); + if (type == null) { + return "experiment_types contains a null value"; + } + + if (!VALID_EXPERIMENT_TYPES.contains(type)) { + return String.format( + AnalyzerErrorConstants.APIErrors.CreateExperimentAPI.BULK_INVALID_EXPERIMENT_TYPES, + List.of(type) + ); + } + + return ""; + } + } diff --git a/src/main/java/com/autotune/operator/KruizeDeploymentInfo.java b/src/main/java/com/autotune/operator/KruizeDeploymentInfo.java index b9e5eb61ff..668cc3407c 100644 --- a/src/main/java/com/autotune/operator/KruizeDeploymentInfo.java +++ b/src/main/java/com/autotune/operator/KruizeDeploymentInfo.java @@ -85,6 +85,7 @@ public class KruizeDeploymentInfo { public static int generate_recommendations_date_range_limit_in_days = 15; public static Integer delete_partition_threshold_in_days = DELETE_PARTITION_THRESHOLD_IN_DAYS; public static String experiment_name_format = "%datasource%|%clustername%|%namespace%|%workloadname%(%workloadtype%)|%containername%"; + public static String namespace_experiment_name_format = "%datasource%|%clustername%|%namespace%"; private static Hashtable tunableLayerPair; //private static KubernetesClient kubernetesClient; private static KubeEventLogger kubeEventLogger;