diff --git a/.gitignore b/.gitignore index 013870b..63bab8a 100644 --- a/.gitignore +++ b/.gitignore @@ -7,3 +7,6 @@ test.db venv/ build/ dist/ +certs/* +!/certs/generate-certs.sh +tests/.env.ssl \ No newline at end of file diff --git a/README.md b/README.md index 0a52279..9f674c2 100644 --- a/README.md +++ b/README.md @@ -91,7 +91,7 @@ Python 3.7+ The following environment variables are exposed to allow SASL authentication with Kafka (along with their default assignment): ``` -KAFKA_SECURITY_PROTOCOL=PLAINTEXT # PLAINTEXT, SASL_PLAINTEXT, SASL_SSL +KAFKA_SECURITY_PROTOCOL=PLAINTEXT # PLAINTEXT, SASL_PLAINTEXT, SASL_SSL, SSL KAFKA_SASL_MECHANISM=PLAIN # PLAIN, SCRAM-SHA-256, SCRAM-SHA-512 KAFKA_PLAIN_USERNAME=None # any str KAFKA_PLAIN_PASSWORD=None # any str diff --git a/docker-compose.yaml b/docker-compose.yaml index 60073b4..64aaa6d 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -19,10 +19,17 @@ services: environment: - KAFKA_ZOOKEEPER_CONNECT=zookeeper:32181 - KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 - - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT - - KAFKA_ADVERTISED_LISTENERS=PLAINTEXT_HOST://localhost:29092,PLAINTEXT://localhost:9092 + - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=PLAINTEXT:PLAINTEXT,SSL:SSL + - KAFKA_ADVERTISED_LISTENERS=SSL://localhost:29092,PLAINTEXT://localhost:9092 - KAFKA_BROKER_ID=1 + - KAFKA_SSL_KEYSTORE_FILENAME=broadcaster-localhost.jks + - KAFKA_SSL_KEYSTORE_CREDENTIALS=password.txt + - KAFKA_SSL_KEY_CREDENTIALS=password.txt + - KAFKA_SSL_TRUSTSTORE_FILENAME=broadcaster-localhost.jks + - KAFKA_SSL_TRUSTSTORE_CREDENTIALS=password.txt - ALLOW_PLAINTEXT_LISTENER=yes + volumes: + - ./certs:/etc/kafka/secrets redis: image: "redis:alpine" ports: diff --git a/requirements.txt b/requirements.txt index 02846f0..f1b0a7c 100644 --- a/requirements.txt +++ b/requirements.txt @@ -20,3 +20,4 @@ isort==5.10.1 mypy==0.971 pytest==7.1.2 pytest-asyncio==0.19.0 +python-dotenv diff --git a/scripts/generate-certs.sh b/scripts/generate-certs.sh new file mode 100644 index 0000000..44f2dcb --- /dev/null +++ b/scripts/generate-certs.sh @@ -0,0 +1,52 @@ +#!/bin/bash +# Setup params +PASSWORD=guessme +VALIDITY=365 +PROJECT_PREFIX=broadcaster +BROKERS='localhost' +CLIENT_ALIAS=myclientname +CLIENT_KEYSTORE=$PROJECT_PREFIX.client.keystore.jks +CLIENT_CERT_FILE=$PROJECT_PREFIX-client-cert-file.txt +CLIENT_CERT_SIGNED=$PROJECT_PREFIX-client-cert-signed.crt +CA_ROOT_ALIAS=ca-root +CA_CERT_NAME=ca-cert.crt +CA_KEY=ca-key.key +BROKER_TRUSTSTORE=$PROJECT_PREFIX.truststore.jks +echo -e "OpenSSL based Keys/Cert generation for Kafka" + +rm $CLIENT_KEYSTORE +rm $BROKER_TRUSTSTORE + +# Generate CA +openssl req -new -x509 -keyout ca-key.key -out ca-cert.crt -days 365 -passin pass:$PASSWORD -subj "/CN=ca-root/OU=SomeUnit/O=SomeOrg/L=London/S=England/C=GB" -passout pass:$PASSWORD + +# Generate for all brokers +echo -e "\n\n###\n###Generating Keys for listed brokers = $BROKERS\n\n###\n###" +for BROKER in $BROKERS +do +keytool -genkeypair -keysize 2048 -keyalg RSA -keystore $PROJECT_PREFIX-$BROKER.jks -alias $BROKER -dname "CN=$BROKER,OU=SomeUnit,O=SomeOrg,L=London,S=England,C=GB" -ext SAN=DNS:$BROKER -validity $VALIDITY -keypass $PASSWORD -storepass $PASSWORD +echo -e "subjectAltName=DNS:$BROKER" > $PROJECT_PREFIX-x509v3-$BROKER.ext +done +echo -e "\n\n###\n###Signing and importing certificates using CA file $CA_CERT_NAME and CA keys file $CA_KEY\n\n###\n###" +for BROKER in $BROKERS +do +keytool -certreq -keystore $PROJECT_PREFIX-$BROKER.jks -alias $BROKER -ext SAN=DNS:$BROKER -file $PROJECT_PREFIX-$BROKER-cert-file.txt -storepass $PASSWORD -keypass $PASSWORD +openssl x509 -req -CA $CA_CERT_NAME -CAkey $CA_KEY -in $PROJECT_PREFIX-$BROKER-cert-file.txt -out $PROJECT_PREFIX-$BROKER-cert-signed.crt -days $VALIDITY -CAcreateserial -passin pass:$PASSWORD -extfile $PROJECT_PREFIX-x509v3-$BROKER.ext +done +echo -e "\n\n###\n###Importing CA root $CA_CERT_NAME and signed broker certs into keystoere\n\n###\n###" +for BROKER in $BROKERS +do +keytool -import -keystore $PROJECT_PREFIX-$BROKER.jks -alias $CA_ROOT_ALIAS -file $CA_CERT_NAME -storepass $PASSWORD -keypass $PASSWORD -noprompt +keytool -import -keystore $PROJECT_PREFIX-$BROKER.jks -alias $BROKER -file $PROJECT_PREFIX-$BROKER-cert-signed.crt -storepass $PASSWORD -keypass $PASSWORD -noprompt +done +echo -e "\n\n###\n###Preparing Client Certificates and keystores###\n\n###" +keytool -genkeypair -keysize 2048 -keyalg RSA -keystore $CLIENT_KEYSTORE -alias $CLIENT_ALIAS -dname "CN=$CLIENT_ALIAS,OU=SomeUnit,O=SomeOrg,L=London,S=England,C=GB" -validity $VALIDITY -storepass $PASSWORD -keypass $PASSWORD +keytool -certreq -keystore $CLIENT_KEYSTORE -alias $CLIENT_ALIAS -file $CLIENT_CERT_FILE -storepass $PASSWORD -keypass $PASSWORD +openssl x509 -req -CA $CA_CERT_NAME -CAkey $CA_KEY -in $CLIENT_CERT_FILE -out $CLIENT_CERT_SIGNED -days $VALIDITY -CAcreateserial -passin pass:$PASSWORD +keytool -import -keystore $CLIENT_KEYSTORE -alias $CA_ROOT_ALIAS -file $CA_CERT_NAME -storepass $PASSWORD -keypass $PASSWORD -noprompt +keytool -import -keystore $CLIENT_KEYSTORE -alias $CLIENT_ALIAS -file $CLIENT_CERT_SIGNED -storepass $PASSWORD -keypass $PASSWORD -noprompt +### +# Once everything is done - import CA into broker and client trust stores correctly + +#write password to password.txt for docker-compose kafka secrets +echo $PASSWORD > password.txt \ No newline at end of file diff --git a/scripts/start b/scripts/start index 188ffae..4f4a722 100755 --- a/scripts/start +++ b/scripts/start @@ -11,4 +11,6 @@ if [ -n "$1" ]; then fi fi +(cd ../certs && ./generate-certs.sh) + docker-compose up $1 diff --git a/setup.py b/setup.py index 4efa50b..7334a5a 100644 --- a/setup.py +++ b/setup.py @@ -49,7 +49,7 @@ def get_packages(package): "redis": ["asyncio-redis"], "postgres": ["asyncpg"], "kafka": ["aiokafka"], - "test": ["pytest", "pytest-asyncio"], + "test": ["pytest", "pytest-asyncio", "python-dotenv"], }, classifiers=[ "Development Status :: 3 - Alpha", diff --git a/tests/test_broadcast.py b/tests/test_broadcast.py index e3313bc..f450132 100644 --- a/tests/test_broadcast.py +++ b/tests/test_broadcast.py @@ -1,6 +1,7 @@ import pytest from broadcaster import Broadcast +from dotenv import load_dotenv @pytest.mark.asyncio @@ -44,3 +45,15 @@ async def test_kafka(): event = await subscriber.get() assert event.channel == "chatroom" assert event.message == "hello" + + +@pytest.mark.skip("Deadlock on `next_published`") +@pytest.mark.asyncio +async def test_kafka_ssl(): + load_dotenv(".env.ssl") + async with Broadcast("kafka://localhost:29092") as broadcast: + async with broadcast.subscribe("chatroom") as subscriber: + await broadcast.publish("chatroom", "hello") + event = await subscriber.get() + assert event.channel == "chatroom" + assert event.message == "hello" \ No newline at end of file