New CLI flags have been added to both TrafficReplayer.Parameters and KafkaConfig.KafkaParameters:
kafkaAuthType/kafkaTrafficAuthType— values:none,msk-iam,scram-sha-512kafkaListenerName/kafkaTrafficListenerName— orchestration metadatakafkaSecretName/kafkaTrafficSecretName— orchestration metadatakafkaUserName/kafkaTrafficUserName— SCRAM username
Problem: Everything downstream still passes a boolean enableMSKAuth — the scram-sha-512 auth type is accepted at the CLI but silently ignored when building Kafka client properties.
Wire the auth type string through the entire call chain so that scram-sha-512 (and any future types) actually configure the Kafka client. Internal signatures change freely; the CLI interface is already stable.
Add kafka password arg names and update POSSIBLE_CREDENTIALS_ARG_FLAG_NAMES and CENSORED_ARGS so the password is redacted in logs and can be picked up from env vars without prefix.
File: coreUtilities/src/main/java/org/opensearch/migrations/arguments/ArgNameConstants.java
Add kafkaPassword @Parameter field. EnvVarParameterPuller will automatically map CAPTURE_PROXY_KAFKA_PASSWORD from the environment (K8s secret mounted as env var).
File: TrafficCapture/captureKafkaOffloader/src/main/java/org/opensearch/migrations/trafficcapture/kafkaoffloader/KafkaConfig.java
Change the 4-arg overload from boolean mskAuthEnabled → String authType, String kafkaUserName, String kafkaPassword. Add SCRAM-SHA-512 handling that constructs the JAAS config from username + password.
Update the convenience buildKafkaProperties(KafkaParameters) overload to pass getEffectiveKafkaAuthType(), kafkaUserName, kafkaPassword.
Merge order:
- Hardcoded defaults (serializers, timeouts)
- Property file (if provided) — overrides defaults
- CLI/env-derived properties (brokers, client ID, auth config) — wins over property file
File: TrafficCapture/captureKafkaOffloader/src/main/java/org/opensearch/migrations/trafficcapture/kafkaoffloader/KafkaConfig.java
Add kafkaTrafficPassword @Parameter field. EnvVarParameterPuller maps TRAFFIC_REPLAYER_KAFKA_TRAFFIC_PASSWORD.
Add getEffectiveKafkaAuthType():
String getEffectiveKafkaAuthType() {
validateKafkaAuthFlags();
if (kafkaTrafficAuthType != null && !kafkaTrafficAuthType.isBlank()) {
return kafkaTrafficAuthType;
}
return Boolean.TRUE.equals(kafkaTrafficEnableMSKAuth) ? KAFKA_AUTH_TYPE_MSK_IAM : KAFKA_AUTH_TYPE_NONE;
}Simplify isKafkaTrafficEnableMSKAuth() to delegate to it.
File: TrafficCapture/trafficReplayer/src/main/java/org/opensearch/migrations/replay/TrafficReplayer.java
Change signature from boolean enableMSKAuth → String authType, String kafkaUserName, String kafkaPassword. Same SCRAM-SHA-512 JAAS construction. Same merge order (property file as defaults, CLI wins).
File: TrafficCapture/trafficReplayer/src/main/java/org/opensearch/migrations/replay/kafka/KafkaTrafficCaptureSource.java
Change boolean enableMSKAuth → String authType, String kafkaUserName, String kafkaPassword and pass through.
File: TrafficCapture/trafficReplayer/src/main/java/org/opensearch/migrations/replay/kafka/KafkaTrafficCaptureSource.java
return KafkaTrafficCaptureSource.buildKafkaSource(
ctx, appParams.kafkaTrafficBrokers, appParams.kafkaTrafficTopic,
appParams.kafkaTrafficGroupId,
appParams.getEffectiveKafkaAuthType(),
appParams.kafkaTrafficUserName,
appParams.kafkaTrafficPassword,
appParams.kafkaTrafficPropertyFile,
Clock.systemUTC(), new KafkaBehavioralPolicy()
);File: TrafficCapture/trafficReplayer/src/main/java/org/opensearch/migrations/replay/TrafficCaptureSourceFactory.java
Change boolean mskAuth → String authType, String kafkaUserName, String kafkaPassword and pass through to buildKafkaProperties.
File: TrafficCapture/trafficReplayer/src/main/java/org/opensearch/migrations/replay/kafka/KafkaTopicDumper.java
runner.runDumpFromKafka(params.mode, params.kafkaTrafficBrokers, params.kafkaTrafficTopic,
params.getEffectiveKafkaAuthType(), params.kafkaTrafficUserName, params.kafkaTrafficPassword,
params.kafkaTrafficPropertyFile, ...);File: TrafficCapture/trafficReplayer/src/main/java/org/opensearch/migrations/replay/TrafficReplayer.java
Update all test call sites that use the old boolean signatures:
| Test File | Change |
|---|---|
CaptureProxySetupTest.java |
No change — calls buildKafkaProperties(KafkaParameters) |
KafkaTrafficCaptureSourceTest.java |
false → "none", null, null; true → "msk-iam", null, null |
KafkaCommitsWorkBetweenLongPollsTest.java |
false → "none", null, null |
KafkaKeepAliveTests.java |
false → "none", null, null |
KafkaTrafficCaptureSourceLongTermTest.java |
false → "none", null, null |
KafkaRestartingTrafficReplayerTest.java |
false → "none", null, null |
authType |
security.protocol |
sasl.mechanism |
JAAS config source |
|---|---|---|---|
none |
(not set) | (not set) | N/A |
msk-iam |
SASL_SSL |
AWS_MSK_IAM |
Hardcoded IAM module |
scram-sha-512 |
SASL_SSL |
SCRAM-SHA-512 |
Built from kafkaUserName + kafkaPassword params/env |
The SCRAM password is injected as an environment variable by K8s (from the secret named by --kafkaSecretName). EnvVarParameterPuller maps it into the @Parameter field automatically:
- Replayer:
TRAFFIC_REPLAYER_KAFKA_TRAFFIC_PASSWORD→kafkaTrafficPassword - Proxy:
CAPTURE_PROXY_KAFKA_PASSWORD→kafkaPassword
The JAAS config is constructed programmatically:
org.apache.kafka.common.security.scram.ScramLoginModule required username="<user>" password="<pass>";
No --kafkaPropertyFile required. If one is provided, it supplies defaults that CLI/env params override.