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
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -7,3 +7,6 @@ test.db
venv/
build/
dist/
certs/*
!/certs/generate-certs.sh
tests/.env.ssl
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
11 changes: 9 additions & 2 deletions docker-compose.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
1 change: 1 addition & 0 deletions requirements.txt
Original file line number Diff line number Diff line change
Expand Up @@ -20,3 +20,4 @@ isort==5.10.1
mypy==0.971
pytest==7.1.2
pytest-asyncio==0.19.0
python-dotenv
52 changes: 52 additions & 0 deletions scripts/generate-certs.sh
Original file line number Diff line number Diff line change
@@ -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
2 changes: 2 additions & 0 deletions scripts/start
Original file line number Diff line number Diff line change
Expand Up @@ -11,4 +11,6 @@ if [ -n "$1" ]; then
fi
fi

(cd ../certs && ./generate-certs.sh)

docker-compose up $1
2 changes: 1 addition & 1 deletion setup.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
13 changes: 13 additions & 0 deletions tests/test_broadcast.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import pytest

from broadcaster import Broadcast
from dotenv import load_dotenv


@pytest.mark.asyncio
Expand Down Expand Up @@ -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"