Skip to content

Commit dc72803

Browse files
committed
feat(extensions): add support for custom JCEKS aliases in S3JceksExtension
- Introduced `aliases` parameter to define custom JCEKS entries, alongside default Hadoop 3.4 S3A encryption options. - Updated `BigDataExtensionsConfigLoader` and `BigDataExtensionsBuilder` to support `aliases` from TOML configurations. - Ensured empty aliases are stored as single spaces to maintain compatibility with S3AUtils. - Added comprehensive tests for default and custom alias handling logic. - Expanded user documentation with configuration and usage details for `aliases`.
1 parent 6abf7f4 commit dc72803

7 files changed

Lines changed: 109 additions & 2 deletions

File tree

doc/user-guide.adoc

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1481,8 +1481,16 @@ hdfsDir = "/bigdata-test/spark"
14811481
fileName = "s3.jceks"
14821482
accessKeyAlias = "fs.s3a.access.key"
14831483
secretKeyAlias = "fs.s3a.secret.key"
1484+
1485+
[s3Jceks.aliases]
1486+
"fs.s3a.encryption.algorithm" = ""
1487+
"fs.s3a.server-side-encryption-algorithm" = ""
1488+
"fs.s3a.encryption.key" = ""
1489+
"fs.s3a.server-side-encryption.key" = ""
14841490
----
14851491

1492+
`aliases` writes additional entries into the same JCEKS file. The extension writes empty defaults for Hadoop 3.4 S3A encryption aliases because `S3AUtils.lookupPassword()` reads those options through the credential provider even when encryption is disabled. The JCEKS provider cannot store a zero-length secret, so empty alias values are stored as a single space; S3A trims the value back to blank while reading. Override them only when your test needs S3A encryption options.
1493+
14861494
Outputs:
14871495

14881496
* `s3-jceks.hdfs.path`

example/spark/src/test/resources/spark-bigdata-extensions.toml

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,10 @@ enabled = true
33
hdfsDir = "/bigdata-test/spark"
44
fileName = "s3.jceks"
55

6+
[s3Jceks.aliases]
7+
"fs.s3a.encryption.algorithm" = ""
8+
"fs.s3a.server-side-encryption-algorithm" = ""
9+
610
[kafkaAvro]
711
enabled = true
812

extensions/src/main/kotlin/org/openprojectx/bigdata/test/extensions/config/BigDataExtensionsBuilder.kt

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -164,6 +164,15 @@ class S3JceksBuilder {
164164
var fileName: String = "s3.jceks"
165165
var accessKeyAlias: String = "fs.s3a.access.key"
166166
var secretKeyAlias: String = "fs.s3a.secret.key"
167+
private val aliases = linkedMapOf<String, String>()
168+
169+
fun alias(name: String, value: String) {
170+
aliases[name] = value
171+
}
172+
173+
fun aliases(values: Map<String, String>) {
174+
aliases += values
175+
}
167176

168177
internal fun build(): S3JceksExtension =
169178
S3JceksExtension(
@@ -172,6 +181,7 @@ class S3JceksBuilder {
172181
fileName = fileName,
173182
accessKeyAlias = accessKeyAlias,
174183
secretKeyAlias = secretKeyAlias,
184+
aliases = S3JceksExtension.defaultAliases + aliases,
175185
)
176186
}
177187

extensions/src/main/kotlin/org/openprojectx/bigdata/test/extensions/config/BigDataExtensionsConfigLoader.kt

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@ class BigDataExtensionsConfigLoader(
5555
fileName = config.string("fileName", "s3.jceks"),
5656
accessKeyAlias = config.string("accessKeyAlias", "fs.s3a.access.key"),
5757
secretKeyAlias = config.string("secretKeyAlias", "fs.s3a.secret.key"),
58+
aliases = S3JceksExtension.defaultAliases + config.stringMap("aliases"),
5859
)
5960
}
6061
}
@@ -163,6 +164,24 @@ class BigDataExtensionsConfigLoader(
163164
)
164165
}.orEmpty()
165166

167+
private fun JsonObject.stringMap(name: String): Map<String, String> =
168+
this[name]?.jsonObject?.flattenStringMap(errorPath = name).orEmpty()
169+
170+
private fun JsonObject.flattenStringMap(prefix: String = "", errorPath: String): Map<String, String> =
171+
flatMap { (key, value) ->
172+
val outputKey = if (prefix.isBlank()) key else "$prefix.$key"
173+
val fullErrorPath = "$errorPath.$key"
174+
when (value) {
175+
is JsonObject -> value.flattenStringMap(prefix = outputKey, errorPath = fullErrorPath).toList()
176+
else -> listOf(
177+
outputKey to (
178+
value.jsonPrimitive.contentOrNull
179+
?: error("Extension config field '$fullErrorPath' must be a string")
180+
),
181+
)
182+
}
183+
}.toMap()
184+
166185
private fun JsonObject.string(name: String): String =
167186
this[name]?.jsonPrimitive?.contentOrNull ?: error("Missing required extension config field '$name'")
168187

extensions/src/main/kotlin/org/openprojectx/bigdata/test/extensions/hadoop/HadoopCredentialProviders.kt

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,7 @@ object HadoopCredentialProviders {
4444
if (provider.getCredentialEntry(alias) != null) {
4545
provider.deleteCredentialEntry(alias)
4646
}
47-
provider.createCredentialEntry(alias, value.toCharArray())
47+
provider.createCredentialEntry(alias, value.ifEmpty { " " }.toCharArray())
4848
}
4949
provider.flush()
5050
}

extensions/src/main/kotlin/org/openprojectx/bigdata/test/extensions/hadoop/S3JceksExtension.kt

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ data class S3JceksExtension(
1111
val fileName: String = "s3.jceks",
1212
val accessKeyAlias: String = "fs.s3a.access.key",
1313
val secretKeyAlias: String = "fs.s3a.secret.key",
14+
val aliases: Map<String, String> = defaultAliases,
1415
) : BigDataExtension {
1516
override val requiredServices: Set<BigDataService> = setOf(BigDataService.HDFS, BigDataService.LOCALSTACK_S3)
1617
override val events: Set<BigDataExtensionEvent> = setOf(BigDataExtensionEvent.AFTER_KIT_START)
@@ -27,9 +28,21 @@ data class S3JceksExtension(
2728
credentials = mapOf(
2829
accessKeyAlias to s3.property("aws.accessKeyId"),
2930
secretKeyAlias to s3.property("aws.secretAccessKey"),
30-
),
31+
) + aliases,
3132
)
3233
context.putOutput("$id.credential-provider.path", providerPath)
3334
context.putOutput("$id.hdfs.path", "${hdfsDir.trimEnd('/')}/$fileName")
3435
}
36+
37+
companion object {
38+
val defaultAliases: Map<String, String> = mapOf(
39+
// Hadoop 3.4 S3A loads these through S3AUtils.lookupPassword().
40+
// Empty aliases keep credential-provider-only configurations from
41+
// failing when encryption is not enabled.
42+
"fs.s3a.encryption.algorithm" to "",
43+
"fs.s3a.server-side-encryption-algorithm" to "",
44+
"fs.s3a.encryption.key" to "",
45+
"fs.s3a.server-side-encryption.key" to "",
46+
)
47+
}
3548
}
Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,53 @@
1+
package org.openprojectx.bigdata.test.extensions.hadoop
2+
3+
import org.junit.jupiter.api.Assertions.assertEquals
4+
import org.junit.jupiter.api.Assertions.assertTrue
5+
import org.junit.jupiter.api.Test
6+
import org.openprojectx.bigdata.test.extensions.config.BigDataExtensionsConfigLoader
7+
8+
class S3JceksExtensionTest {
9+
@Test
10+
fun `loads default and custom jceks aliases from toml`() {
11+
val extension = BigDataExtensionsConfigLoader().load(
12+
"""
13+
[s3Jceks]
14+
hdfsDir = "/config"
15+
fileName = "s3.jceks"
16+
17+
[s3Jceks.aliases]
18+
"fs.s3a.encryption.algorithm" = "SSE-S3"
19+
"fs.s3a.encryption.key" = ""
20+
"custom.alias" = "custom-value"
21+
""".trimIndent(),
22+
).single() as S3JceksExtension
23+
24+
assertEquals("/config", extension.hdfsDir)
25+
assertEquals("s3.jceks", extension.fileName)
26+
assertEquals("SSE-S3", extension.aliases["fs.s3a.encryption.algorithm"])
27+
assertEquals("", extension.aliases["fs.s3a.server-side-encryption-algorithm"])
28+
assertEquals("", extension.aliases["fs.s3a.server-side-encryption.key"])
29+
assertEquals("custom-value", extension.aliases["custom.alias"])
30+
}
31+
32+
@Test
33+
fun `loads unquoted dotted jceks aliases from toml`() {
34+
val extension = BigDataExtensionsConfigLoader().load(
35+
"""
36+
[s3Jceks.aliases]
37+
fs.s3a.encryption.algorithm = "SSE-KMS"
38+
""".trimIndent(),
39+
).single() as S3JceksExtension
40+
41+
assertEquals("SSE-KMS", extension.aliases["fs.s3a.encryption.algorithm"])
42+
}
43+
44+
@Test
45+
fun `default jceks aliases include s3a encryption options`() {
46+
val defaults = S3JceksExtension.defaultAliases
47+
48+
assertTrue(defaults.containsKey("fs.s3a.encryption.algorithm"))
49+
assertTrue(defaults.containsKey("fs.s3a.server-side-encryption-algorithm"))
50+
assertTrue(defaults.containsKey("fs.s3a.encryption.key"))
51+
assertTrue(defaults.containsKey("fs.s3a.server-side-encryption.key"))
52+
}
53+
}

0 commit comments

Comments
 (0)