From 6406ccf39f367096a2b2fad7998bdfe031566532 Mon Sep 17 00:00:00 2001 From: poushalibanik Date: Thu, 6 Mar 2025 13:05:05 +0530 Subject: [PATCH 1/4] Added a new functionality which now enables us to add broker ports of our own --- trivup/apps/KafkaBrokerApp.py | 21 +++++++--- trivup/clusters/KafkaCluster.py | 73 ++++++++++++++++++++++++++++++--- 2 files changed, 83 insertions(+), 11 deletions(-) diff --git a/trivup/apps/KafkaBrokerApp.py b/trivup/apps/KafkaBrokerApp.py index b5e2712..20080ec 100644 --- a/trivup/apps/KafkaBrokerApp.py +++ b/trivup/apps/KafkaBrokerApp.py @@ -41,6 +41,7 @@ class KafkaBrokerApp (trivup.App): """ Kafka broker app Depends on ZookeeperApp (unless KRaft mode) """ + def __init__(self, cluster, conf=None, on=None): """ @param cluster Current cluster @@ -73,6 +74,7 @@ def __init__(self, cluster, conf=None, on=None): * fdlimit - RLIMIT_NOFILE (or "max") (default: max) * conf - arbitary server.properties config as a list of strings. * realm - Kerberos realm to use when sasl_mechanisms contains GSSAPI + * """ super(KafkaBrokerApp, self).__init__(cluster, conf=conf, on=on) @@ -171,6 +173,7 @@ def __init__(self, cluster, conf=None, on=None): conf_blob.append(listener_map) + def sort_listener(a): """ Sort listener_types list so that the PLAINTEXTs are first, since the first listener is used by operational(). """ @@ -178,12 +181,19 @@ def sort_listener(a): return 0 else: return 1 - - # Allocate a port for each listener type - ports = [(x, trivup.TcpPortAllocator(self.cluster).next( + + # If brokers ports are specified those would be used, if not it would fallback to TcpPortAllocator logic + if 'user_port' in self.conf and self.conf['user_port']: + ports = [(x, self.conf.get('user_port')) for x in sorted(set(listener_types), key=sort_listener)] + self.conf['port'] = self.conf.get('user_port') + self.conf['address'] = '%s:%d' % (listener_host, self.conf.get('user_port')) + else : + ports = [(x, trivup.TcpPortAllocator(self.cluster).next( self, self.conf.get('port_base', self.conf.get('port', None)))) for x in sorted(set(listener_types), key=sort_listener)] - self.conf['port'] = ports[0][1] # "Default" port + self.conf['port'] = ports[0][1] # "Default" port + self.conf['address'] = '%s:%d' % (listener_host, self.conf['port']) + if can_docker: # Add docker listener to allow services (e.g, SchemaRegistry) in @@ -191,7 +201,6 @@ def sort_listener(a): docker_port = trivup.TcpPortAllocator(self.cluster).next(self) docker_host = '%s:%d' % (cluster.get_docker_host(), docker_port) - self.conf['address'] = '%s:%d' % (listener_host, self.conf['port']) # Create a listener for each port listeners = ['%s://%s:%d' % (x[0], "0.0.0.0", x[1]) for x in ports] if can_docker: @@ -462,4 +471,4 @@ def _add_simple_authorizer(self, conf_blob): if self.version[0] >= 3: conf_blob.append('authorizer.class.name=kafka.security.authorizer.AclAuthorizer') # noqa: E501 else: - conf_blob.append('authorizer.class.name=kafka.security.auth.SimpleAclAuthorizer') # noqa: E501 + conf_blob.append('authorizer.class.name=kafka.security.auth.SimpleAclAuthorizer') # noqa: E501 \ No newline at end of file diff --git a/trivup/clusters/KafkaCluster.py b/trivup/clusters/KafkaCluster.py index 8cae55c..31dd4da 100755 --- a/trivup/clusters/KafkaCluster.py +++ b/trivup/clusters/KafkaCluster.py @@ -64,9 +64,10 @@ import argparse import subprocess import copy - +import socket class KafkaCluster(object): + # conf dict structure with defaults: # commented-out fields are not defaults but show what is available. default_conf = { @@ -92,8 +93,10 @@ class KafkaCluster(object): 'oidc': False, # Additional broker server.properties configuration # 'broker_conf': ['connections.max.idle.ms=1234', ..] + # 'broker_ports' : Comma-separated list of Kafka broker ports. If not provided, random ports will be used. } + def __init__(self, **kwargs): """ Create and start a KafkaCluster. See default_conf above for parameters. """ @@ -108,6 +111,27 @@ def __init__(self, **kwargs): self.version_num = [int(x) for x in self.version.split('.')][:3] self.kraft = self.conf.get('kraft') + + # Checking if ports being passed are available + if 'broker_ports' in conf and conf['broker_ports']: + self.broker_ports_list = [int(port) for port in conf.get('broker_ports').split(',')] + self._check_ports_availability() + else: + self.broker_ports_list = [] + + # Scenario where no. of ports being passed doesn't match broker count + if 'broker_cnt' in conf and conf['broker_cnt'] and 'broker_ports' in conf and conf['broker_ports']: + if len(self.broker_ports_list) != self.conf.get('broker_cnt'): + raise ValueError(f"The number of ports :({len(self.broker_ports_list)}) does not match broker_cnt : ({conf['broker_cnt']}).") + else: + if 'broker_ports' in conf and conf['broker_ports']: + self.conf['broker_cnt'] = len(self.broker_ports_list) + + # Broker count's default value has been set to None so has to be overwritten + if 'broker_cnt' not in self.conf or not self.conf['broker_cnt']: + self.conf['broker_cnt'] = self.default_conf.get('broker_cnt') + + # Create trivup Cluster self.cluster = Cluster( self.__class__.__name__, @@ -173,6 +197,9 @@ def __init__(self, **kwargs): self.brokers = dict() for n in range(0, broker_cnt): bconf = copy.deepcopy(self.broker_conf) + if self.broker_ports_list: + bconf['user_port'] = self.broker_ports_list[n] + if self.version_num >= [2, 4, 0]: # Configure rack & replica selector if broker supports # fetch-from-follower @@ -438,7 +465,8 @@ def interactive(self, cmd=None): retcode, fullcmd)) return retcode - + + def client_conf(self): """ Get a dict copy of the client configuration """ return deepcopy(self._client_conf) @@ -451,7 +479,37 @@ def write_client_conf(self, path, additional_blob=None): if additional_blob is not None: f.write(str('#\n# Additional configuration:')) f.write(str(additional_blob)) - + + + def _check_ports_availability(self): + """ Check availability of the broker ports and exit if any are unavailable. """ + unavailable_ports = [] + + for port in self.broker_ports_list: + if not self._is_port_available(port): + unavailable_ports.append(port) + + if unavailable_ports: + print(f"Error: The following broker ports are unavailable: {', '.join(map(str, unavailable_ports))}") + print("Closing application due to unavailable ports.") + sys.exit(1) + + print(f"All broker ports are available: {', '.join(map(str, self.broker_ports_list))}") + + + def _is_port_available(self, port): + """ Check if a port is available by trying to bind to it. """ + s = None + try: + s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + s.bind(('', port)) + return True + except socket.error: + return False + finally: + if s: + s.close() + if __name__ == '__main__': @@ -472,7 +530,7 @@ def write_client_conf(self, path, additional_blob=None): default=KafkaCluster.default_conf['with_sr'], help='Enable SchemaRegistry') parser.add_argument('--brokers', dest='broker_cnt', type=int, - default=KafkaCluster.default_conf['broker_cnt'], + default=None, help='Number of Kafka brokers') parser.add_argument('--version', dest='version', type=str, default=KafkaCluster.default_conf['version'], @@ -494,6 +552,9 @@ def write_client_conf(self, path, additional_blob=None): action='store_true', default=KafkaCluster.default_conf['oidc'], help='Enable Oauthbearer OIDC JWT server') + parser.add_argument('--broker-ports', dest='broker_ports', type=str, default=None, + help='Comma-separated list of Kafka broker ports. If not provided, default ports will be used.') + args = parser.parse_args() @@ -507,7 +568,9 @@ def write_client_conf(self, path, additional_blob=None): 'broker_cnt': args.broker_cnt, 'kafka_path': args.kafka_src, 'cleanup': not args.no_cleanup, - 'oidc': args.oidc} + 'oidc': args.oidc, + 'broker_ports': args.broker_ports + } kc = KafkaCluster(**conf) From c84ea0fc1240742bb57322405737b0e861859932 Mon Sep 17 00:00:00 2001 From: poushalibanik Date: Thu, 6 Nov 2025 15:55:16 +0530 Subject: [PATCH 2/4] added changes that brings in a new option with trivup command i.e a daemon flag(--d), also displays the process id, broker port and schema port in a single place --- trivup/apps/KafkaBrokerApp.py | 3 +- trivup/clusters/KafkaCluster.py | 150 ++++++++++++++++++++++++-------- 2 files changed, 117 insertions(+), 36 deletions(-) diff --git a/trivup/apps/KafkaBrokerApp.py b/trivup/apps/KafkaBrokerApp.py index 20080ec..8e2ca5f 100644 --- a/trivup/apps/KafkaBrokerApp.py +++ b/trivup/apps/KafkaBrokerApp.py @@ -207,7 +207,8 @@ def sort_listener(a): listeners.append('%s://%s:%d' % ('DOCKER', "0.0.0.0", docker_port)) self.conf['listeners'] = ','.join(listeners) if 'advertised_hostname' not in self.conf: - self.conf['advertised_hostname'] = self.conf['nodename'] + # self.conf['advertised_hostname'] = self.conf['nodename'] + self.conf['advertised_hostname'] = socket.gethostname()+ ".hursley.ibm.com" advertised_listeners = ['%s://%s:%d' % (x[0], self.conf['advertised_hostname'], x[1]) for x in ports if x[0] != 'CONTROLLER'] diff --git a/trivup/clusters/KafkaCluster.py b/trivup/clusters/KafkaCluster.py index 31dd4da..99fffda 100755 --- a/trivup/clusters/KafkaCluster.py +++ b/trivup/clusters/KafkaCluster.py @@ -8,11 +8,11 @@ # modification, are permitted provided that the following conditions are met: # # * Redistributions of source code must retain the above copyright notice, this -# list of conditions and the following disclaimer. +#   list of conditions and the following disclaimer. # # * Redistributions in binary form must reproduce the above copyright notice, -# this list of conditions and the following disclaimer in the documentation -# and/or other materials provided with the distribution. +#   this list of conditions and the following disclaimer in the documentation +#   and/or other materials provided with the distribution. # # THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" # AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE @@ -28,24 +28,24 @@ # Provides a Kafka cluster with the following components: -# * ZooKeeperApp (1) -# * KafkaBrokerApp brokers (broker_cnt=3) -# * SslApp (optional, if with_ssl=True) -# * KerberosKdcApp (optional, sasl.mechanism=GSSAPI, -# cross-realm if realm_cnt=2) -# * SchemaRegistryApp (optional, if with_sr=True) -# * OauthbearerOIDCApp (optional, if oidc=True) +#  * ZooKeeperApp (1) +#  * KafkaBrokerApp brokers (broker_cnt=3) +#  * SslApp (optional, if with_ssl=True) +#  * KerberosKdcApp (optional, sasl.mechanism=GSSAPI, +#                     cross-realm if realm_cnt=2) +#  * SchemaRegistryApp (optional, if with_sr=True) +#  * OauthbearerOIDCApp (optional, if oidc=True) # # cluster.env (dict) will contain: -# TRIVUP_ROOT -# ZK_ADDRESS (unless --kraft) -# BROKERS -# BROKER_PID_ -# KAFKA_PATH (path to kafka package root directory) -# SR_URL (if with_sr) -# SSL_(ca|pub|priv)_.. (if with_ssl, paths to certs and keys) -# SSL_password (if with_ssl, key password) -# KRB5CCNAME, KRB5_COFNIG, KRB5_KDC_PROFILE (if GSSAPI enabled) +#      TRIVUP_ROOT +#      ZK_ADDRESS  (unless --kraft) +#      BROKERS +#      BROKER_PID_ +#      KAFKA_PATH  (path to kafka package root directory) +#      SR_URL      (if with_sr) +#      SSL_(ca|pub|priv)_.. (if with_ssl, paths to certs and keys) +#      SSL_password          (if with_ssl, key password) +#      KRB5CCNAME, KRB5_COFNIG, KRB5_KDC_PROFILE (if GSSAPI enabled) # # See conf dict structure below. @@ -65,15 +65,16 @@ import subprocess import copy import socket +import urllib.parse; class KafkaCluster(object): - + # conf dict structure with defaults: # commented-out fields are not defaults but show what is available. default_conf = { 'version': '2.8.0', # Apache Kafka version 'cp_version': '6.1.0', # Confluent Platform version (for SR) - 'broker_cnt': 3, + 'broker_cnt': 1, 'sasl_mechanism': '', # GSSAPI, PLAIN, SCRAM-.., ... 'realm_cnt': 1, 'krb_renew_lifetime': 30, @@ -94,9 +95,10 @@ class KafkaCluster(object): # Additional broker server.properties configuration # 'broker_conf': ['connections.max.idle.ms=1234', ..] # 'broker_ports' : Comma-separated list of Kafka broker ports. If not provided, random ports will be used. + # 'd' : Runs the trivup in the background and returns the main process ID along with broker port and schema port. } - + def __init__(self, **kwargs): """ Create and start a KafkaCluster. See default_conf above for parameters. """ @@ -111,7 +113,7 @@ def __init__(self, **kwargs): self.version_num = [int(x) for x in self.version.split('.')][:3] self.kraft = self.conf.get('kraft') - + # Checking if ports being passed are available if 'broker_ports' in conf and conf['broker_ports']: self.broker_ports_list = [int(port) for port in conf.get('broker_ports').split(',')] @@ -199,7 +201,7 @@ def __init__(self, **kwargs): bconf = copy.deepcopy(self.broker_conf) if self.broker_ports_list: bconf['user_port'] = self.broker_ports_list[n] - + if self.version_num >= [2, 4, 0]: # Configure rack & replica selector if broker supports # fetch-from-follower @@ -238,7 +240,15 @@ def __init__(self, **kwargs): def __del__(self): """ Destructor: forcibly stop the cluster """ - self.stop(force=True) + # Only stop if cleanup is enabled. Otherwise, assume the user handles it. + # This try/except prevents the TypeError/AttributeError if self.conf is cleared by interpreter on exit. + try: + if self.conf.get('cleanup', True): + self.stop(force=True) + except AttributeError: + # Ignore the error that happens when attributes are already gone (e.g., in daemon mode) + pass + def _setup_env(self): """ Set up convenience envs """ @@ -391,6 +401,7 @@ def stop(self, cleanup=True, keeptypes=['log'], force=False, timeout=0): if timeout > 0: self.cluster.wait_stopped(timeout) if cleanup: + # self.cluster.cleanup is a method, not a boolean, this is correct. self.cluster.cleanup(keeptypes) def stopped(self): @@ -431,7 +442,7 @@ def interactive(self, cmd=None): print("# - Waiting for cluster to go operational in {}/{}".format( self.cluster.root_path, self.cluster.instance)) - kc.wait_operational() + self.wait_operational() env = self.env.copy() @@ -465,8 +476,8 @@ def interactive(self, cmd=None): retcode, fullcmd)) return retcode - - + + def client_conf(self): """ Get a dict copy of the client configuration """ return deepcopy(self._client_conf) @@ -479,7 +490,7 @@ def write_client_conf(self, path, additional_blob=None): if additional_blob is not None: f.write(str('#\n# Additional configuration:')) f.write(str(additional_blob)) - + def _check_ports_availability(self): """ Check availability of the broker ports and exit if any are unavailable. """ @@ -492,7 +503,7 @@ def _check_ports_availability(self): if unavailable_ports: print(f"Error: The following broker ports are unavailable: {', '.join(map(str, unavailable_ports))}") print("Closing application due to unavailable ports.") - sys.exit(1) + sys.exit(1) print(f"All broker ports are available: {', '.join(map(str, self.broker_ports_list))}") @@ -502,14 +513,14 @@ def _is_port_available(self, port): s = None try: s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - s.bind(('', port)) - return True + s.bind(('', port)) + return True except socket.error: - return False + return False finally: if s: s.close() - + if __name__ == '__main__': @@ -554,6 +565,9 @@ def _is_port_available(self, port): help='Enable Oauthbearer OIDC JWT server') parser.add_argument('--broker-ports', dest='broker_ports', type=str, default=None, help='Comma-separated list of Kafka broker ports. If not provided, default ports will be used.') + parser.add_argument('--d', dest='d', action='store_true', + help='Runs the trivup in the background and returns the main process ID along with broker port and schema port') + args = parser.parse_args() @@ -569,11 +583,77 @@ def _is_port_available(self, port): 'kafka_path': args.kafka_src, 'cleanup': not args.no_cleanup, 'oidc': args.oidc, - 'broker_ports': args.broker_ports + 'broker_ports': args.broker_ports, + 'd': args.d } kc = KafkaCluster(**conf) + if args.d: + # Disabling cleanup via the configuration dictionary + kc.conf['cleanup'] = False + + # Wait for the cluster to be fully operational. + print("# Waiting for cluster to go operational...") + try: + # Increased timeout for Schema Registry + kc.wait_operational(timeout=180) + except Exception: + print("FATAL ERROR: Cluster failed to become operational.") + sys.exit(1) + + # Getting SR Listener + sr_listeners = None + sr_app = None + + # Checking if Schema Registry was requested and started + sr_apps = kc.cluster.find_apps(SchemaRegistryApp, 'started') + + if sr_apps: + # Getting the running Schema Registry App instance + sr_app = sr_apps[0] + + # ACCESS THE CORRECT VARIABLE: 'url' from the configuration dictionary + sr_listeners = sr_app.conf.get('url') + + # HOSTNAME EXTRACTION + # Parsing the bootstrap server string to get the netloc (host:port) + parsed_broker = urllib.parse.urlparse(kc.bootstrap_servers) + + # Extract the hostname + if parsed_broker.netloc: + host = parsed_broker.netloc.split(':')[0] + else: + host = parsed_broker.path.split(':')[0] + + # APPLYING THE HOSTNAME TO THE SR LISTENER + # Only run replacement if the SR listener URL was successfully retrieved + if sr_listeners is not None: + if 'localhost' in sr_listeners or '0.0.0.0' in sr_listeners: + sr_listeners = sr_listeners.replace('localhost', host).replace('0.0.0.0', host) + + + # Robustly getting the first running broker + running_brokers = kc.cluster.find_apps(KafkaBrokerApp, 'started') + + if running_brokers and running_brokers[0].proc: + pid = running_brokers[0].proc.pid + + print(f"Bootstrap Servers: {kc.bootstrap_servers}") + + # Printing the Schema Registry port only if it was successfully found + if sr_listeners: + print(f"Schema Registry port: {sr_listeners}") + + # Output the PID and exit immediately, leaving the subprocesses running. + print(f"Cluster started in background. Main broker PID: {pid}") + # Exit cleanly + sys.exit(0) + else: + print("ERROR: Cluster started, but could not find main broker PID.") + sys.exit(1) + + # Normal interactive mode / Command execution ret = kc.interactive(args.cmd) print("# Stopping cluster in {}/{}".format(kc.cluster.root_path, From 6432dd3e7715801ae5a194aaab55d4ea103e78b5 Mon Sep 17 00:00:00 2001 From: poushalibanik Date: Thu, 4 Dec 2025 15:14:57 +0530 Subject: [PATCH 3/4] programatically fetching the domain name along with hostname --- trivup/apps/KafkaBrokerApp.py | 23 +++++++++++++++++++++-- 1 file changed, 21 insertions(+), 2 deletions(-) diff --git a/trivup/apps/KafkaBrokerApp.py b/trivup/apps/KafkaBrokerApp.py index 8e2ca5f..bfb2cbc 100644 --- a/trivup/apps/KafkaBrokerApp.py +++ b/trivup/apps/KafkaBrokerApp.py @@ -37,7 +37,6 @@ import socket import time - class KafkaBrokerApp (trivup.App): """ Kafka broker app Depends on ZookeeperApp (unless KRaft mode) """ @@ -208,7 +207,27 @@ def sort_listener(a): self.conf['listeners'] = ','.join(listeners) if 'advertised_hostname' not in self.conf: # self.conf['advertised_hostname'] = self.conf['nodename'] - self.conf['advertised_hostname'] = socket.gethostname()+ ".hursley.ibm.com" + + def get_fqdn_from_resolvconf(): + domain = None + with open("/etc/resolv.conf") as f: + for line in f: + if line.startswith("search") or line.startswith("domain"): + parts = line.split() + if len(parts) > 1: + domain = parts[1] + break + + hostname = socket.gethostname() + + if domain: + return f"{hostname}.{domain}" + else: + return hostname + + + fqdn = get_fqdn_from_resolvconf() + self.conf['advertised_hostname'] = fqdn advertised_listeners = ['%s://%s:%d' % (x[0], self.conf['advertised_hostname'], x[1]) for x in ports if x[0] != 'CONTROLLER'] From d3d9d3e35440cf7fb36d861c02c5766a95f9c035 Mon Sep 17 00:00:00 2001 From: Poushali-Banik Date: Thu, 4 Dec 2025 18:03:06 +0530 Subject: [PATCH 4/4] programatically fetching the domain name along with hostname --- trivup/apps/KafkaBrokerApp.py | 1 - 1 file changed, 1 deletion(-) diff --git a/trivup/apps/KafkaBrokerApp.py b/trivup/apps/KafkaBrokerApp.py index bfb2cbc..f14ab9a 100644 --- a/trivup/apps/KafkaBrokerApp.py +++ b/trivup/apps/KafkaBrokerApp.py @@ -225,7 +225,6 @@ def get_fqdn_from_resolvconf(): else: return hostname - fqdn = get_fqdn_from_resolvconf() self.conf['advertised_hostname'] = fqdn advertised_listeners = ['%s://%s:%d' %