-
-
Notifications
You must be signed in to change notification settings - Fork 515
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Fix kafka internal docker connection #2490
Open
catinapoke
wants to merge
34
commits into
testcontainers:main
Choose a base branch
from
catinapoke:main
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
Changes from 19 commits
Commits
Show all changes
34 commits
Select commit
Hold shift + click to select a range
320e54d
fixed kafka internal docker connection
catinapoke 9b6537d
added new option for kafka listeners
catinapoke 15f4ca3
Merge remote-tracking branch 'original/main'
catinapoke c7552f4
fixed kafka docs with new options
catinapoke 8904938
Update docs/modules/kafka.md
catinapoke 7b38c28
Apply suggestions from code review for docs
catinapoke 98793d7
added new test, fixed docs from code review
catinapoke 62a68b5
Merge branch 'main' of https://github.com/catinapoke/testcontainers-go
catinapoke a7c4a9f
Merge remote-tracking branch 'Original/main'
catinapoke bfa9191
renamed docker folder to testdata and fixed test name
catinapoke f1b0d48
Merge remote-tracking branch 'Original/main'
catinapoke 128d31a
added test for input listeners validation and refactor from code review
catinapoke ead068a
added unit test for trimValidateListeners
catinapoke eb71110
Merge remote-tracking branch 'Original/main'
catinapoke a4d91e2
Merge branch 'main' into main
mdelapenya dd8fcec
Merge remote-tracking branch 'upstream/main'
catinapoke 5c147d4
go mod update
catinapoke b31e1af
todo for PR
catinapoke 9f23a4d
updated copyStarterScript for upstream
catinapoke e63a604
fix: typo
mdelapenya f9cb155
fix: use Run method
mdelapenya 8394529
chore: use PLAINTEXT and BROKER
mdelapenya 31950f4
fix: handle error in tests
mdelapenya 3aa2cbf
chore: make lint
mdelapenya d61004f
chore: use non deprecated APIs
mdelapenya df53491
chore: handle errors in testss
mdelapenya bf02852
fix: close sarama client
mdelapenya 152ef0c
fix: validation in test
mdelapenya efa0f7d
chore: refactor test
mdelapenya d669bfb
chore: add test using kcat
mdelapenya 4c44b97
Merge branch 'main' into catinapoke/main
mdelapenya ba701e1
Merge pull request #1 from mdelapenya/catinapoke/main
catinapoke b51be12
fixed basic kafka test but network one is broken
catinapoke c8c47f4
not working test with kcat
catinapoke File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -20,7 +20,7 @@ const ( | |
// starterScript { | ||
starterScriptContent = `#!/bin/bash | ||
source /etc/confluent/docker/bash-config | ||
export KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://%s:%d,BROKER://%s:9092 | ||
export KAFKA_ADVERTISED_LISTENERS=%s | ||
echo Starting Kafka KRaft mode | ||
sed -i '/KAFKA_ZOOKEEPER_CONNECT/d' /etc/confluent/docker/configure | ||
echo 'kafka-storage format --ignore-formatted -t "$(kafka-storage random-uuid)" -c /etc/kafka/kafka.properties' >> /etc/confluent/docker/configure | ||
|
@@ -34,6 +34,13 @@ echo '' > /etc/confluent/docker/ensure | |
type KafkaContainer struct { | ||
testcontainers.Container | ||
ClusterID string | ||
Listeners KafkaListener | ||
} | ||
|
||
type KafkaListener struct { | ||
Name string | ||
Host string | ||
Port string | ||
} | ||
|
||
// Deprecated: use Run instead | ||
|
@@ -49,10 +56,10 @@ func Run(ctx context.Context, img string, opts ...testcontainers.ContainerCustom | |
ExposedPorts: []string{string(publicPort)}, | ||
Env: map[string]string{ | ||
// envVars { | ||
"KAFKA_LISTENERS": "PLAINTEXT://0.0.0.0:9093,BROKER://0.0.0.0:9092,CONTROLLER://0.0.0.0:9094", | ||
"KAFKA_REST_BOOTSTRAP_SERVERS": "PLAINTEXT://0.0.0.0:9093,BROKER://0.0.0.0:9092,CONTROLLER://0.0.0.0:9094", | ||
"KAFKA_LISTENER_SECURITY_PROTOCOL_MAP": "BROKER:PLAINTEXT,PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT", | ||
"KAFKA_INTER_BROKER_LISTENER_NAME": "BROKER", | ||
"KAFKA_LISTENERS": "EXTERNAL://0.0.0.0:9093,INTERNAL://0.0.0.0:9092,CONTROLLER://0.0.0.0:9094", | ||
"KAFKA_REST_BOOTSTRAP_SERVERS": "EXTERNAL://0.0.0.0:9093,INTERNAL://0.0.0.0:9092,CONTROLLER://0.0.0.0:9094", | ||
"KAFKA_LISTENER_SECURITY_PROTOCOL_MAP": "INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT", | ||
"KAFKA_INTER_BROKER_LISTENER_NAME": "INTERNAL", | ||
"KAFKA_BROKER_ID": "1", | ||
"KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR": "1", | ||
"KAFKA_OFFSETS_TOPIC_NUM_PARTITIONS": "1", | ||
|
@@ -68,15 +75,43 @@ func Run(ctx context.Context, img string, opts ...testcontainers.ContainerCustom | |
Entrypoint: []string{"sh"}, | ||
// this CMD will wait for the starter script to be copied into the container and then execute it | ||
Cmd: []string{"-c", "while [ ! -f " + starterScript + " ]; do sleep 0.1; done; bash " + starterScript}, | ||
LifecycleHooks: []testcontainers.ContainerLifecycleHooks{ | ||
} | ||
|
||
genericContainerReq := testcontainers.GenericContainerRequest{ | ||
ContainerRequest: req, | ||
Started: true, | ||
} | ||
|
||
settings := defaultOptions() | ||
for _, opt := range opts { | ||
if apply, ok := opt.(Option); ok { | ||
apply(&settings) | ||
} | ||
if err := opt.Customize(&genericContainerReq); err != nil { | ||
return nil, err | ||
} | ||
} | ||
|
||
if err := trimValidateListeners(settings.Listeners); err != nil { | ||
return nil, fmt.Errorf("listeners validation: %w", err) | ||
} | ||
|
||
// apply envs for listeners | ||
envChange := editEnvsForListeners(settings.Listeners) | ||
for key, item := range envChange { | ||
genericContainerReq.Env[key] = item | ||
} | ||
|
||
genericContainerReq.ContainerRequest.LifecycleHooks = | ||
[]testcontainers.ContainerLifecycleHooks{ | ||
{ | ||
PostStarts: []testcontainers.ContainerHook{ | ||
// Use a single hook to copy the starter script and wait for | ||
// the Kafka server to be ready. This prevents the wait running | ||
// if the starter script fails to copy. | ||
func(ctx context.Context, c testcontainers.Container) error { | ||
// 1. copy the starter script into the container | ||
if err := copyStarterScript(ctx, c); err != nil { | ||
if err := copyStarterScript(ctx, c, &settings); err != nil { | ||
return fmt.Errorf("copy starter script: %w", err) | ||
} | ||
|
||
|
@@ -85,19 +120,7 @@ func Run(ctx context.Context, img string, opts ...testcontainers.ContainerCustom | |
}, | ||
}, | ||
}, | ||
}, | ||
} | ||
|
||
genericContainerReq := testcontainers.GenericContainerRequest{ | ||
ContainerRequest: req, | ||
Started: true, | ||
} | ||
|
||
for _, opt := range opts { | ||
if err := opt.Customize(&genericContainerReq); err != nil { | ||
return nil, err | ||
} | ||
} | ||
|
||
err := validateKRaftVersion(genericContainerReq.Image) | ||
if err != nil { | ||
|
@@ -116,32 +139,70 @@ func Run(ctx context.Context, img string, opts ...testcontainers.ContainerCustom | |
return &KafkaContainer{Container: container, ClusterID: clusterID}, nil | ||
} | ||
|
||
func trimValidateListeners(listeners []KafkaListener) error { | ||
// Trim | ||
for i := 0; i < len(listeners); i++ { | ||
listeners[i].Name = strings.ToUpper(strings.Trim(listeners[i].Name, " ")) | ||
listeners[i].Host = strings.Trim(listeners[i].Host, " ") | ||
mdelapenya marked this conversation as resolved.
Show resolved
Hide resolved
|
||
listeners[i].Port = strings.Trim(listeners[i].Port, " ") | ||
} | ||
|
||
// Validate | ||
var ports map[string]bool = make(map[string]bool, len(listeners)+2) | ||
var names map[string]bool = make(map[string]bool, len(listeners)+2) | ||
|
||
// check for default listeners | ||
ports["9094"] = true | ||
ports["9093"] = true | ||
|
||
// check for default listeners | ||
names["CONTROLLER"] = true | ||
names["EXTERNAL"] = true | ||
|
||
for _, item := range listeners { | ||
if names[item.Name] { | ||
return fmt.Errorf("duplicate of listener name: %s", item.Name) | ||
} | ||
names[item.Name] = true | ||
|
||
if ports[item.Port] { | ||
return fmt.Errorf("duplicate of listener port: %s", item.Port) | ||
} | ||
ports[item.Port] = true | ||
} | ||
|
||
return nil | ||
} | ||
|
||
// copyStarterScript copies the starter script into the container. | ||
func copyStarterScript(ctx context.Context, c testcontainers.Container) error { | ||
func copyStarterScript(ctx context.Context, c testcontainers.Container, settings *options) error { | ||
if err := wait.ForListeningPort(publicPort). | ||
SkipInternalCheck(). | ||
WaitUntilReady(ctx, c); err != nil { | ||
return fmt.Errorf("wait for exposed port: %w", err) | ||
} | ||
|
||
host, err := c.Host(ctx) | ||
if err != nil { | ||
return fmt.Errorf("host: %w", err) | ||
if len(settings.Listeners) == 0 { | ||
defaultInternal, err := internalListener(ctx, c) | ||
if err != nil { | ||
return fmt.Errorf("can't create default internal listener: %w", err) | ||
} | ||
settings.Listeners = append(settings.Listeners, defaultInternal) | ||
} | ||
|
||
inspect, err := c.Inspect(ctx) | ||
defaultExternal, err := externalListener(ctx, c) | ||
if err != nil { | ||
return fmt.Errorf("inspect: %w", err) | ||
return fmt.Errorf("can't create default external listener: %w", err) | ||
} | ||
|
||
hostname := inspect.Config.Hostname | ||
settings.Listeners = append(settings.Listeners, defaultExternal) | ||
|
||
port, err := c.MappedPort(ctx, publicPort) | ||
if err != nil { | ||
return fmt.Errorf("mapped port: %w", err) | ||
var advertised []string | ||
for _, item := range settings.Listeners { | ||
advertised = append(advertised, fmt.Sprintf("%s://%s:%s", item.Name, item.Host, item.Port)) | ||
} | ||
|
||
scriptContent := fmt.Sprintf(starterScriptContent, host, port.Int(), hostname) | ||
scriptContent := fmt.Sprintf(starterScriptContent, strings.Join(advertised, ",")) | ||
|
||
if err := c.CopyToContainer(ctx, []byte(scriptContent), starterScript, 0o755); err != nil { | ||
return fmt.Errorf("copy to container: %w", err) | ||
|
@@ -150,12 +211,43 @@ func copyStarterScript(ctx context.Context, c testcontainers.Container) error { | |
return nil | ||
} | ||
|
||
func WithClusterID(clusterID string) testcontainers.CustomizeRequestOption { | ||
return func(req *testcontainers.GenericContainerRequest) error { | ||
req.Env["CLUSTER_ID"] = clusterID | ||
func editEnvsForListeners(listeners []KafkaListener) map[string]string { | ||
catinapoke marked this conversation as resolved.
Show resolved
Hide resolved
|
||
if len(listeners) == 0 { | ||
// no change | ||
return map[string]string{} | ||
} | ||
|
||
return nil | ||
envs := map[string]string{ | ||
"KAFKA_LISTENERS": "CONTROLLER://0.0.0.0:9094, EXTERNAL://0.0.0.0:9093", | ||
"KAFKA_REST_BOOTSTRAP_SERVERS": "CONTROLLER://0.0.0.0:9094, EXTERNAL://0.0.0.0:9093", | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. kafka local provides a rest-proxy service but I don't see any test for it. It is also a good opportunity to add a test. There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. will take a look |
||
"KAFKA_LISTENER_SECURITY_PROTOCOL_MAP": "CONTROLLER:PLAINTEXT, EXTERNAL:PLAINTEXT", | ||
} | ||
|
||
// expect first listener has common network between kafka instances | ||
envs["KAFKA_INTER_BROKER_LISTENER_NAME"] = listeners[0].Name | ||
|
||
// expect small number of listeners, so joins is okay | ||
for _, item := range listeners { | ||
envs["KAFKA_LISTENERS"] = strings.Join( | ||
[]string{ | ||
envs["KAFKA_LISTENERS"], | ||
fmt.Sprintf("%s://0.0.0.0:%s", item.Name, item.Port), | ||
}, | ||
",", | ||
) | ||
|
||
envs["KAFKA_REST_BOOTSTRAP_SERVERS"] = envs["KAFKA_LISTENERS"] | ||
|
||
envs["KAFKA_LISTENER_SECURITY_PROTOCOL_MAP"] = strings.Join( | ||
[]string{ | ||
envs["KAFKA_LISTENER_SECURITY_PROTOCOL_MAP"], | ||
item.Name + ":" + "PLAINTEXT", | ||
}, | ||
",", | ||
) | ||
} | ||
|
||
return envs | ||
} | ||
|
||
// Brokers retrieves the broker connection strings from Kafka with only one entry, | ||
|
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
any reason why it should change?
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
More obvious way to show the purpose. It's also the way many people write listeners in articles
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I'd not change this, as it could break existing client code using for reasons these env vars and the values