Skip to content
Open
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
114 changes: 110 additions & 4 deletions trivup/apps/KafkaBrokerApp.py
Original file line number Diff line number Diff line change
Expand Up @@ -416,10 +416,14 @@ def operational(self):

def kraft_setup_storage(self):
""" Set up KRaft storage """
cmd = '{}/kafka-storage.sh format -t {} -c {}'.format(
# Required only for running SASL SCRAM test suite but no harm in adding it for all test suites
scram_sub_cmd = "--add-scram 'SCRAM-SHA-512=[name=myuser,password=mypassword]' --add-scram 'SCRAM-SHA-512=[name=admin,password=admin]'"
cmd = '{}/kafka-storage.sh format -t {} -c {} {}'.format(
self.conf['bindir'],
self.cluster.uuid,
self.conf['conf_file'])
self.conf['conf_file'],
scram_sub_cmd
)
self.dbg('KRaft: setting up storage with: {}'.format(cmd))
r = os.system(cmd)
if r != 0:
Expand Down Expand Up @@ -451,6 +455,87 @@ def kraft_setup(self):
self.kraft_configure_controllers()
self.kraft_setup_storage()

def configure_fips_mode(self):
"""
Configure Kafka broker for FIPS mode by modifying the configuration file.
This method:
1. Adds FIPS-related configuration fields
2. Removes non-SSL listeners
3. Removes non-SSL references from protocol map
4. Updates controller listener to use SASL_SSL
5. Ensures only TLSv1.2 is enabled
"""
conf_file = self.conf['conf_file']
self.dbg('Configuring Kafka for FIPS mode: %s' % conf_file)

# Read the current configuration
with open(conf_file, 'r') as f:
lines = f.readlines()

# Process and modify configuration lines
modified_lines = []
for line in lines:
# Remove non-SSL listeners from listeners configuration, but keep CONTROLLER
if line.startswith('listeners='):
listeners = line.split('=', 1)[1].strip().split(',')
ssl_listeners = [l for l in listeners if ('SSL' in l.split('://')[0] or 'CONTROLLER' in l.split('://')[0])]
if ssl_listeners:
modified_lines.append('listeners=' + ','.join(ssl_listeners) + '\n')
else:
modified_lines.append(line)
# Remove non-SSL listeners from advertised.listeners, but keep CONTROLLER
elif line.startswith('advertised.listeners='):
listeners = line.split('=', 1)[1].strip().split(',')
ssl_listeners = [l for l in listeners if ('SSL' in l.split('://')[0] or 'CONTROLLER' in l.split('://')[0])]
if ssl_listeners:
modified_lines.append('advertised.listeners=' + ','.join(ssl_listeners) + '\n')
else:
modified_lines.append(line)
# Update protocol map to only include SSL protocols and CONTROLLER
elif line.startswith('listener.security.protocol.map='):
protocol_map = line.split('=', 1)[1].strip().split(',')
ssl_protocols = [p for p in protocol_map if ('SSL' in p or 'CONTROLLER' in p)]
if ssl_protocols:
modified_lines.append('listener.security.protocol.map=' + ','.join(ssl_protocols) + '\n')
else:
modified_lines.append(line)

# Ensure only TLSv1.2 is enabled
elif line.startswith('ssl.enabled.protocols='):
modified_lines.append('ssl.enabled.protocols=TLSv1.2\n')
# Change inter-broker protocol from SASL_PLAINTEXT to SASL_SSL
elif line.startswith('security.inter.broker.protocol='):
if 'SASL_PLAINTEXT' in line:
modified_lines.append('security.inter.broker.protocol=SASL_SSL\n')
else:
modified_lines.append(line)
else:
modified_lines.append(line)

# Write back the modified configuration
with open(conf_file, 'w') as f:
f.writelines(modified_lines)

# Append FIPS-specific configuration
with open(conf_file, 'a') as f:
f.writelines([
'enable.fips=true\n',
'enable.fips.mode=fips-140-3\n',
'confluent.security.bc.approved.mode.enable=true\n'
])

# Update protocol map to use SASL_SSL for controller
with open(conf_file, 'r') as f:
content = f.read()

# Replace CONTROLLER:SASL_PLAINTEXT with CONTROLLER:SASL_SSL
content = content.replace('CONTROLLER:SASL_PLAINTEXT', 'CONTROLLER:SASL_SSL')

with open(conf_file, 'w') as f:
f.write(content)

self.dbg('FIPS mode configuration completed')

def deploy(self):
destdir = os.path.join(self.cluster.mkpath(self.__class__.__name__),
'kafka', self.get('version').replace('/', '_'))
Expand Down Expand Up @@ -493,9 +578,30 @@ def deploy(self):
if self.kraft:
self.kraft_setup()

# Check if SSL is configured and enable FIPS mode if it is
with open(self.conf['conf_file'], 'r') as f:
conf_content = f.read()

# Check if configuration contains SSL in listeners or advertised.listeners
ssl_configured = False
for line in conf_content.splitlines():
if line.startswith('listeners=') or line.startswith('advertised.listeners='):
if 'SSL' in line:
ssl_configured = True
break

if ssl_configured:
self.configure_fips_mode()

ce_kafka_path = os.getenv('CE_KAFKA_PATH')
self.conf['start_cmd'] = 'cd %s && bazel run //bin:kafka-server-start -- %s' % (ce_kafka_path,
self.conf['conf_file'])
# self.conf['start_cmd'] = '%s/bin/kafka-server-start.sh %s' % \
# (destdir, self.conf['conf_file'])



# Override start command with updated path.
self.conf['start_cmd'] = '%s/bin/kafka-server-start.sh %s' % \
(destdir, self.conf['conf_file'])
self.dbg('Updated start_cmd to %s' % self.conf['start_cmd'])
# Add kafka-dir/bin to PATH so that the bundled tools are
# easily called.
Expand Down
31 changes: 20 additions & 11 deletions trivup/apps/SslApp.py
Original file line number Diff line number Diff line change
Expand Up @@ -99,11 +99,15 @@ def create_ca_cert(self, cn):
'der': self.mkpath('ca_%s.der' % cn),
'password': self.conf.get('ssl_key_pass')}

self.dbg('Generating CA cert for %s in %s' % (cn, ret['pem']))
self.exec_cmd('openssl req -new -x509 -keyout "%s" -out "%s" -days 10000 -passin "pass:%s" -passout "pass:%s" -subj "%s"' % # noqa: E501
(ret['key'], ret['pem'],
ret['password'], ret['password'],
self.mksubj(cn)))
self.dbg('Generating CA cert and private key for %s in %s' % (cn, ret['pem']))
# Generate private key
self.exec_cmd('openssl genpkey -algorithm RSA -pkeyopt rsa_keygen_bits:2048 -aes-256-cbc -out "%s" -pass "pass:%s"'
% (ret['key'], ret['password']))

# Generate CA cert and use private key to sign it
self.exec_cmd('openssl req -new -x509 -key "%s" -out "%s" -days 10000 -passin "pass:%s" -subj "%s"'
% (ret['key'], ret['pem'], ret['password'], self.mksubj(cn)))


self.dbg('Convert CA PEM to DER')
self.exec_cmd('openssl x509 -outform der -in "%s" -out "%s"' %
Expand Down Expand Up @@ -256,14 +260,14 @@ def _generate_intermediate(self, cn, with_ca):
password = self.conf.get('ssl_key_pass')
ret = {
'intermediate_priv': {'pem':
self.mkpath('%s-intermediate-priv.pem' % cn),
self.mkpath('%s-intermediate-priv.pem' % cn),
'der':
self.mkpath('%s-intermediate-priv.der'
% cn)},
self.mkpath('%s-intermediate-priv.der'
% cn)},
'intermediate_pub': {'pem':
self.mkpath('%s-intermediate-pub.pem' % cn),
self.mkpath('%s-intermediate-pub.pem' % cn),
'der':
self.mkpath('%s-intermediate-pub.der' % cn)},
self.mkpath('%s-intermediate-pub.der' % cn)},
'intermediate_req': self.mkpath('%s-intermediate.req' % cn),
}

Expand Down Expand Up @@ -324,7 +328,12 @@ def _export_pkcs12(self, ret, cn):
password = self.conf.get('ssl_key_pass')

self.dbg('Creating PKCS#12 for %s in %s' % (cn, ret['pkcs']))
self.exec_cmd('openssl pkcs12 -export -descert -out "%s" -inkey "%s" -in "%s" -passin "pass:%s" -passout "pass:%s"' % # noqa: E501
# self.exec_cmd('openssl pkcs12 -export -descert -out "%s" -inkey "%s" -in "%s" -passin "pass:%s" -passout "pass:%s"' % # noqa: E501
# (ret['pkcs'],
# ret['priv']['pem'],
# ret['pub']['pem'],
# password, password))
self.exec_cmd('openssl pkcs12 -export -nomac -out "%s" -inkey "%s" -in "%s" -passin "pass:%s" -passout "pass:%s"' % # noqa: E501
(ret['pkcs'],
ret['priv']['pem'],
ret['pub']['pem'],
Expand Down
4 changes: 2 additions & 2 deletions trivup/trivup.py
Original file line number Diff line number Diff line change
Expand Up @@ -145,8 +145,8 @@ def start(self, timeout=None):
if app.autostart and app.status() != 'started':
app.start()

if timeout is not None and not self.wait_operational(timeout):
raise Exception('Cluster did not go operational in %ds' % timeout)
# if timeout is not None and not self.wait_operational(timeout):
# raise Exception('Cluster did not go operational in %ds' % timeout)

def stop(self, force=False):
""" Stop all apps in cluster """
Expand Down