import json import re import argparse import os from html import escape as _xml_escape from typing import List, Dict, Any, Optional def extract_payload(value_str): """Extract the `payload={...}` section, parse key=value pairs safely, return JSON. Handles values containing spaces, parentheses, commas, colons and URLs by using a regex that looks ahead for the next ", =" pattern instead of naively splitting on commas. Commas that are part of the value (e.g. inside parentheses) are preserved. """ match = re.search(r'payload=({.*})\s*$', value_str) if not match: return value_str body = match.group(1) inner = body[1:-1] # strip outer braces # Pattern: key= (value up to next ', =' sequence or end of string) pattern = re.compile(r'(\w+)=((?:(?!, \w+=).)*)') result: Dict[str, Any] = {} for m in pattern.finditer(inner): key = m.group(1) raw_val = m.group(2).strip() if raw_val.endswith(','): raw_val = raw_val[:-1].rstrip() # Normalize nulls if raw_val == 'null': val: Any = None else: # Try int if re.fullmatch(r'-?\d+', raw_val): try: val = int(raw_val) except ValueError: # pragma: no cover val = raw_val else: val = raw_val result[key] = val return json.dumps(result, ensure_ascii=False) def generate_changelog( records: List[Dict[str, Any]], override_author: Optional[str] = None, id_pattern: Optional[str] = None, ) -> str: """Generate changelog XML string from list of record dicts. Parameters: records: list of record dicts. override_author: global author (record author ignored) if provided. id_pattern: id template. Supports {n}/{i}/{index}. If no placeholder and multiple records, an index suffix (:n) is appended. """ xml_header = ( '\n' '\n' ) xml_footer = '\n' changesets: List[str] = [] # Base required fields; may relax id/author if overridden # Only topic is required in raw messages; id & author come from CLI. # Only the base id (via CLI) is required now; topic may be absent in # DLQ export. required_fields = ["id"] optional_defaults = [ ("kafka-cluster-name", "managed-kafka-cdc-destination") ] kafka_defaults = [("key", "null"), ("tombstone", False)] reserved_non_attr = {"value"} kafka_msg_fields = {k for k, _ in kafka_defaults} if override_author is None: raise ValueError("override_author (CLI --author) must be provided") for idx, record in enumerate(records, start=1): # Determine effective required fields per record effective_required = [f for f in required_fields] if id_pattern: effective_required = [ f for f in effective_required if f != "id" ] missing = [f for f in effective_required if not record.get(f)] if missing: raise ValueError( f"Record {idx} missing required field(s): {', '.join(missing)}" ) # Resolve author/id with overrides def resolve_id() -> str: if id_pattern: if any(ph in id_pattern for ph in ("{n}", "{i}", "{index}")): return id_pattern.format(n=idx, i=idx, index=idx) if len(records) > 1: return f"{id_pattern}:{idx}" return id_pattern return record.get("id") def resolve_author() -> str: # always CLI provided return override_author attr_items = [ ("id", resolve_id()), ("author", resolve_author()), ] topic_val = record.get("topic") if topic_val is not None: attr_items.append(("topic", topic_val)) for k, default in optional_defaults: attr_items.append((k, record.get(k, default))) extra_keys = sorted( k for k in record.keys() if k not in {a for a, _ in attr_items} and k not in kafka_msg_fields and k not in reserved_non_attr ) for k in extra_keys: attr_items.append((k, record[k])) kafka_items = [] for k, default in kafka_defaults: val = record.get(k, default) if k == "tombstone": val = str(bool(val)).lower() kafka_items.append((k, val)) payload_json = extract_payload(record.get("value", "")) def serialize(items): parts = [] for a_key, a_val in items: if isinstance(a_val, bool): a_val = str(a_val).lower() elif a_val is None: a_val = "null" parts.append( f"{_xml_escape(str(a_key))}=\"{_xml_escape(str(a_val))}\"" ) return " ".join(parts) changeset_attrs_str = serialize(attr_items) kafka_attrs_str = serialize(kafka_items) changeset = ( f" \n" f" \n" f" {payload_json}\n" f" \n" f" \n" ) changesets.append(changeset) return f"{xml_header}{''.join(changesets)}{xml_footer}" def main(): parser = argparse.ArgumentParser( description="Convert JSON to changelog XML." ) parser.add_argument("input_json", help="Path to the input JSON file") # Output path is derived automatically from ID unless --output overrides parser.add_argument( "--author", required=True, help="Author to apply to all changesets (overrides any record author)", ) parser.add_argument( "--id", required=True, help=( "Base ID for changesets and output filename. If multiple records, " ":n suffix appended to changeset ids." ), ) parser.add_argument( "--output-dir", default=os.path.expanduser( "~/code/database/kafka-db-deploy/kafka/build/changelog/dml" ), help=( "Directory to place the derived output file. Created if missing. " "Default: %(default)s" ), ) parser.add_argument( "--suffix", default="_redrive_dql_messages.xml", help=( "Filename suffix appended to ID for output file name. Default: " "%(default)s" ), ) args = parser.parse_args() with open(args.input_json, "r") as f: records = json.load(f) # Determine output path output_dir = os.path.abspath(os.path.expanduser(args.output_dir)) os.makedirs(output_dir, exist_ok=True) output_filename = ( f"{args.id}{args.suffix}" if not args.suffix.startswith(args.id) else args.suffix ) output_path = os.path.join(output_dir, output_filename) xml = generate_changelog( records, override_author=args.author, id_pattern=args.id, ) with open(output_path, "w") as f: f.write(xml) print(f"Wrote changelog: {output_path}") if __name__ == "__main__": # pragma: no cover main()