#!/usr/bin/env bash set -euo pipefail repo_root="$(git rev-parse --show-toplevel 2>/dev/null || pwd)" env_name="${ENVIRONMENT:-$(basename "$(dirname "${TG_ROOT:-live/prod/kafka-topics}")")}" infra_path="${INFRASTRUCTURE_CONFIG_PATH:-${repo_root}/infrastructure.yaml}" python3 - "${infra_path}" "${env_name}" <<'PY' import os import sys try: import yaml except Exception as exc: raise SystemExit(f"pyyaml is required for Kafka topics validation: {exc}") infra_path, env = sys.argv[1:3] def load_yaml(path): if not os.path.exists(path): raise SystemExit(f"File not found: {path}") with open(path, "r", encoding="utf-8") as fh: return yaml.safe_load(fh) or {} def env_value(value, env, default=None): if value is None: return default if isinstance(value, dict): if env in value: return value[env] if "_default" in value: return value["_default"] return default return value infra = load_yaml(infra_path) env_cfg = (infra.get("environments") or {}).get(env) if not isinstance(env_cfg, dict): raise SystemExit(f"Environment '{env}' is not defined in {infra_path}") cluster_refs = env_cfg.get("kafka_cluster_refs") or {} policy = env_cfg.get("kafka_policy") or {} topics = ((env_cfg.get("kafka") or {}).get("topics")) if topics is None: topics = [] if not isinstance(topics, list): raise SystemExit(f"environments.{env}.kafka.topics must be an array") allowed_cleanup = {"delete", "compact"} allowed_deletion = {"orphan", "delete"} errors = [] seen_names = set() for idx, topic in enumerate(topics): prefix = f"environments.{env}.kafka.topics[{idx}]" if not isinstance(topic, dict): errors.append(f"{prefix}: must be an object") continue name = topic.get("name") owner = topic.get("owner") cluster_ref = topic.get("clusterRef") if not name: errors.append(f"{prefix}.name is required") elif name in seen_names: errors.append(f"{prefix}.name '{name}' is duplicated") else: seen_names.add(name) if not owner: errors.append(f"{prefix}.owner is required") if not cluster_ref: errors.append(f"{prefix}.clusterRef is required") continue if cluster_ref not in cluster_refs: errors.append(f"{prefix}.clusterRef '{cluster_ref}' is not defined in infrastructure.yaml for env '{env}'") continue cluster = cluster_refs[cluster_ref] or {} partitions = env_value(topic.get("partitions"), env, cluster.get("default_partitions")) replication = env_value(topic.get("replicationFactor"), env, cluster.get("default_replication_factor")) if partitions is None: errors.append(f"{prefix}.partitions is required or cluster default_partitions must be set") if replication is None: errors.append(f"{prefix}.replicationFactor is required or cluster default_replication_factor must be set") try: partitions_num = int(partitions) if partitions_num < 1: errors.append(f"{prefix}.partitions must be >= 1") except Exception: errors.append(f"{prefix}.partitions must be a number") try: replication_num = int(replication) if replication_num < 1: errors.append(f"{prefix}.replicationFactor must be >= 1") except Exception: replication_num = None errors.append(f"{prefix}.replicationFactor must be a number") max_replication = cluster.get("max_replication_factor") if replication_num is not None and max_replication is not None and replication_num > int(max_replication): errors.append(f"{prefix}.replicationFactor={replication_num} exceeds cluster max_replication_factor={max_replication}") deletion_policy = topic.get("deletionPolicy", "orphan") if deletion_policy not in allowed_deletion: errors.append(f"{prefix}.deletionPolicy must be one of {sorted(allowed_deletion)}") if deletion_policy == "delete" and not bool(policy.get("allow_delete", False)): errors.append(f"{prefix}.deletionPolicy=delete is forbidden by kafka_policy.allow_delete=false") merged_config = {} merged_config.update(cluster.get("default_topic_config") or {}) merged_config.update(topic.get("config") or {}) cleanup_policy = env_value(merged_config.get("cleanup.policy"), env) if cleanup_policy is not None and cleanup_policy not in allowed_cleanup: errors.append(f"{prefix}.config.cleanup.policy must be one of {sorted(allowed_cleanup)}") min_isr = env_value(merged_config.get("min.insync.replicas"), env) if min_isr is not None and replication_num is not None: try: min_isr_num = int(min_isr) if min_isr_num > replication_num: errors.append(f"{prefix}.config.min.insync.replicas={min_isr_num} is greater than replicationFactor={replication_num}") except Exception: errors.append(f"{prefix}.config.min.insync.replicas must be a number") if errors: raise SystemExit("Kafka topics validation failed:\n- " + "\n- ".join(errors)) print(f"Kafka topics declaration is valid: {infra_path} env={env} topics={len(topics)}") PY