Skip to content

Commit 3c4ad7c

Browse files
committed
feat: add programmatic extension configuration for BigDataExtensions
- Introduced support for programmatic extension declarations in `BigDataExtensions` via `BigDataExtensionsConfigurer`. - Enabled dynamic extension setup using Kotlin companion objects, explicit configurer classes, or test instance implementations. - Added `BigDataExtensionsBuilder` for simplified and declarative extension configurations (e.g., Kafka Avro, S3 JCEKS). - Updated examples, tests, and documentation to demonstrate programmatic configurations. - Enhanced JUnit extension to combine TOML resources with programmatic declarations for runtime flexibility.
1 parent dad923f commit 3c4ad7c

6 files changed

Lines changed: 292 additions & 2 deletions

File tree

README.md

Lines changed: 43 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -54,7 +54,7 @@ val kit = BigDataTestKit.builder()
5454
The built-in extensions currently support:
5555

5656
- `s3Jceks`: creates an HDFS-backed JCEKS file from the LocalStack S3 endpoint credentials and exposes `s3-jceks.credential-provider.path`.
57-
- `kafkaAvro`: creates Kafka topics and produces Avro records through Schema Registry from inline JSON records or a records resource.
57+
- `kafkaAvro`: creates Kafka topics and produces Avro records through Schema Registry from inline TOML records or a records resource.
5858

5959
JUnit usage is declaration-driven. The test declares the services with `@BigDataTest`, then points `@BigDataExtensions` at one or more TOML resources:
6060

@@ -94,6 +94,48 @@ records = [
9494
]
9595
```
9696

97+
The same setup can be declared programmatically when names, records, or options need to be generated dynamically:
98+
99+
```kotlin
100+
@BigDataExtensions
101+
@BigDataTest(kafka = true, schemaRegistry = true, localStackS3 = true, hdfs = true)
102+
class MyIntegrationTest {
103+
companion object : BigDataExtensionsConfigurer {
104+
override fun configure(extensions: BigDataExtensionsBuilder) {
105+
val suffix = System.nanoTime()
106+
extensions.s3Jceks {
107+
hdfsDir = "/bigdata-test/$suffix"
108+
}
109+
extensions.kafkaAvro {
110+
topic("events-$suffix", "classpath:schemas/event.avsc") {
111+
record("alpha", mapOf("id" to 1, "name" to "alpha"))
112+
}
113+
}
114+
}
115+
}
116+
}
117+
```
118+
119+
Java tests can use the explicit configurer form:
120+
121+
```java
122+
@BigDataExtensions(configurer = MyExtensionsConfigurer.class)
123+
@BigDataTest(kafka = true, schemaRegistry = true)
124+
class MyJavaIntegrationTest {
125+
}
126+
127+
public final class MyExtensionsConfigurer implements BigDataExtensionsConfigurer {
128+
@Override
129+
public void configure(BigDataExtensionsBuilder extensions) {
130+
extensions.kafkaAvro(kafka -> kafka.topic(
131+
"events",
132+
"classpath:schemas/event.avsc",
133+
topic -> topic.record("alpha", Map.of("id", 1, "name", "alpha"))
134+
));
135+
}
136+
}
137+
```
138+
97139
For future extension modules, implement `BigDataExtension` directly for programmatic use or publish a `BigDataExtensionProvider` via `ServiceLoader`. Providers are selected from config entries under `extensions` by `type`, and extensions can hook lifecycle events such as `AFTER_KIT_START`, `BEFORE_TEST_EXECUTION`, and `AFTER_ALL`.
98140

99141
## JUnit 5

doc/user-guide.adoc

Lines changed: 63 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -324,6 +324,69 @@ class MyIntegrationTest {
324324

325325
Extension output is exposed through `BigDataExtensionResult`.
326326

327+
=== Programmatic Extension Declaration
328+
329+
Use programmatic declaration when test data, topic names, paths, or custom extensions need to be generated at runtime. The JUnit extension combines TOML resources first, then programmatic declarations.
330+
331+
[source,kotlin]
332+
----
333+
@BigDataExtensions
334+
@BigDataTest(
335+
hdfs = true,
336+
kafka = true,
337+
schemaRegistry = true,
338+
localStackS3 = true,
339+
)
340+
class MyIntegrationTest {
341+
companion object : BigDataExtensionsConfigurer {
342+
override fun configure(extensions: BigDataExtensionsBuilder) {
343+
val suffix = System.nanoTime()
344+
345+
extensions.s3Jceks {
346+
hdfsDir = "/bigdata-test/$suffix"
347+
fileName = "s3.jceks"
348+
}
349+
350+
extensions.kafkaAvro {
351+
topic("events-$suffix", "classpath:schemas/event.avsc") {
352+
record("alpha", mapOf("id" to 1, "name" to "alpha"))
353+
record("beta", mapOf("id" to 2, "name" to "beta"))
354+
}
355+
}
356+
}
357+
}
358+
}
359+
----
360+
361+
Programmatic declarations can be supplied in three ways:
362+
363+
* A Kotlin companion object implementing `BigDataExtensionsConfigurer`.
364+
* The test instance implementing `BigDataExtensionsConfigurer`.
365+
* An explicit no-arg configurer class with `@BigDataExtensions(configurer = MyConfigurer::class)`.
366+
367+
Java tests can use the same explicit configurer form. Kotlin's `KClass` annotation parameter is exposed to Java as a normal `.class` value.
368+
369+
[source,java]
370+
----
371+
@BigDataExtensions(configurer = MyExtensionsConfigurer.class)
372+
@BigDataTest(kafka = true, schemaRegistry = true)
373+
class MyJavaIntegrationTest {
374+
}
375+
376+
public final class MyExtensionsConfigurer implements BigDataExtensionsConfigurer {
377+
@Override
378+
public void configure(BigDataExtensionsBuilder extensions) {
379+
extensions.kafkaAvro(kafka -> kafka.topic(
380+
"events",
381+
"classpath:schemas/event.avsc",
382+
topic -> topic.record("alpha", Map.of("id", 1, "name", "alpha"))
383+
));
384+
}
385+
}
386+
----
387+
388+
Use `extensions.extension(...)` to add a custom `BigDataExtension` directly when the built-in builder methods are not enough.
389+
327390
=== S3 JCEKS Extension
328391

329392
The `s3Jceks` extension creates a Hadoop credential provider file in HDFS and stores LocalStack S3 credentials in it.
Lines changed: 141 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,141 @@
1+
package org.openprojectx.bigdata.test.extensions.config
2+
3+
import kotlinx.serialization.json.JsonArray
4+
import kotlinx.serialization.json.JsonElement
5+
import kotlinx.serialization.json.JsonNull
6+
import kotlinx.serialization.json.JsonObject
7+
import kotlinx.serialization.json.JsonPrimitive
8+
import org.openprojectx.bigdata.test.extensions.core.BigDataExtension
9+
import org.openprojectx.bigdata.test.extensions.hadoop.S3JceksExtension
10+
import org.openprojectx.bigdata.test.extensions.kafka.KafkaAvroRecordSeed
11+
import org.openprojectx.bigdata.test.extensions.kafka.KafkaAvroSeedExtension
12+
import org.openprojectx.bigdata.test.extensions.kafka.KafkaAvroTopicSeed
13+
import java.util.function.Consumer
14+
15+
class BigDataExtensionsBuilder {
16+
private val extensions = mutableListOf<BigDataExtension>()
17+
18+
fun extension(extension: BigDataExtension) {
19+
extensions += extension
20+
}
21+
22+
fun s3Jceks(configure: S3JceksBuilder.() -> Unit = {}) {
23+
extensions += S3JceksBuilder().apply(configure).build()
24+
}
25+
26+
fun s3Jceks(configure: Consumer<S3JceksBuilder>) {
27+
extensions += S3JceksBuilder().also { configure.accept(it) }.build()
28+
}
29+
30+
fun kafkaAvro(configure: KafkaAvroBuilder.() -> Unit) {
31+
extensions += KafkaAvroBuilder().apply(configure).build()
32+
}
33+
34+
fun kafkaAvro(configure: Consumer<KafkaAvroBuilder>) {
35+
extensions += KafkaAvroBuilder().also { configure.accept(it) }.build()
36+
}
37+
38+
internal fun build(): List<BigDataExtension> = extensions.toList()
39+
}
40+
41+
class S3JceksBuilder {
42+
var id: String = "s3-jceks"
43+
var hdfsDir: String = "/bigdata-test/config"
44+
var fileName: String = "s3.jceks"
45+
var accessKeyAlias: String = "fs.s3a.access.key"
46+
var secretKeyAlias: String = "fs.s3a.secret.key"
47+
48+
internal fun build(): S3JceksExtension =
49+
S3JceksExtension(
50+
id = id,
51+
hdfsDir = hdfsDir,
52+
fileName = fileName,
53+
accessKeyAlias = accessKeyAlias,
54+
secretKeyAlias = secretKeyAlias,
55+
)
56+
}
57+
58+
class KafkaAvroBuilder {
59+
var id: String = "kafka-avro-seed"
60+
private val topics = mutableListOf<KafkaAvroTopicSeed>()
61+
62+
fun topic(
63+
name: String,
64+
schema: String,
65+
partitions: Int = 1,
66+
replicationFactor: Short = 1,
67+
configure: KafkaAvroTopicBuilder.() -> Unit,
68+
) {
69+
topics += KafkaAvroTopicBuilder(name, schema, partitions, replicationFactor).apply(configure).build()
70+
}
71+
72+
fun topic(
73+
name: String,
74+
schema: String,
75+
configure: Consumer<KafkaAvroTopicBuilder>,
76+
) {
77+
topics += KafkaAvroTopicBuilder(name, schema, 1, 1).also { configure.accept(it) }.build()
78+
}
79+
80+
fun topic(
81+
name: String,
82+
schema: String,
83+
partitions: Int,
84+
replicationFactor: Short,
85+
configure: Consumer<KafkaAvroTopicBuilder>,
86+
) {
87+
topics += KafkaAvroTopicBuilder(name, schema, partitions, replicationFactor).also { configure.accept(it) }.build()
88+
}
89+
90+
internal fun build(): KafkaAvroSeedExtension =
91+
KafkaAvroSeedExtension(id = id, topics = topics.toList())
92+
}
93+
94+
class KafkaAvroTopicBuilder internal constructor(
95+
private val name: String,
96+
private val schema: String,
97+
private val partitions: Int,
98+
private val replicationFactor: Short,
99+
) {
100+
private val records = mutableListOf<KafkaAvroRecordSeed>()
101+
102+
fun record(key: String, value: JsonObject) {
103+
records += KafkaAvroRecordSeed(key = key, value = value)
104+
}
105+
106+
fun record(key: String, value: Map<String, Any?>) {
107+
record(key, value.toJsonObject())
108+
}
109+
110+
internal fun build(): KafkaAvroTopicSeed =
111+
KafkaAvroTopicSeed(
112+
name = name,
113+
schema = schema,
114+
records = records.toList(),
115+
partitions = partitions,
116+
replicationFactor = replicationFactor,
117+
)
118+
}
119+
120+
private fun Map<String, Any?>.toJsonObject(): JsonObject =
121+
JsonObject(mapValues { (_, value) -> value.toJsonElement() })
122+
123+
private fun Any?.toJsonElement(): JsonElement = when (this) {
124+
null -> JsonNull
125+
is JsonElement -> this
126+
is String -> JsonPrimitive(this)
127+
is Boolean -> JsonPrimitive(this)
128+
is Int -> JsonPrimitive(this)
129+
is Long -> JsonPrimitive(this)
130+
is Float -> JsonPrimitive(this)
131+
is Double -> JsonPrimitive(this)
132+
is Map<*, *> -> JsonObject(
133+
entries.associate { (key, value) ->
134+
require(key is String) { "Kafka Avro record map keys must be strings" }
135+
key to value.toJsonElement()
136+
},
137+
)
138+
is Iterable<*> -> JsonArray(map { it.toJsonElement() })
139+
is Array<*> -> JsonArray(map { it.toJsonElement() })
140+
else -> error("Unsupported Kafka Avro record value type: ${this::class}")
141+
}
Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,9 @@
1+
package org.openprojectx.bigdata.test.extensions.config
2+
3+
interface BigDataExtensionsConfigurer {
4+
fun configure(extensions: BigDataExtensionsBuilder)
5+
}
6+
7+
class NoopBigDataExtensionsConfigurer : BigDataExtensionsConfigurer {
8+
override fun configure(extensions: BigDataExtensionsBuilder) = Unit
9+
}
Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,14 @@
11
package org.openprojectx.bigdata.test.extensions.junit5
22

33
import org.junit.jupiter.api.extension.ExtendWith
4+
import org.openprojectx.bigdata.test.extensions.config.BigDataExtensionsConfigurer
5+
import org.openprojectx.bigdata.test.extensions.config.NoopBigDataExtensionsConfigurer
6+
import kotlin.reflect.KClass
47

58
@Target(AnnotationTarget.CLASS)
69
@Retention(AnnotationRetention.RUNTIME)
710
@ExtendWith(BigDataExtensionsExtension::class)
811
annotation class BigDataExtensions(
912
vararg val value: String,
13+
val configurer: KClass<out BigDataExtensionsConfigurer> = NoopBigDataExtensionsConfigurer::class,
1014
)

extensions/src/main/kotlin/org/openprojectx/bigdata/test/extensions/junit5/BigDataExtensionsExtension.kt

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,10 @@ import org.junit.jupiter.api.extension.BeforeTestExecutionCallback
55
import org.junit.jupiter.api.extension.ExtensionContext
66
import org.junit.jupiter.api.extension.ParameterContext
77
import org.junit.jupiter.api.extension.ParameterResolver
8+
import org.openprojectx.bigdata.test.extensions.config.BigDataExtensionsBuilder
89
import org.openprojectx.bigdata.test.extensions.config.BigDataExtensionsConfigLoader
10+
import org.openprojectx.bigdata.test.extensions.config.BigDataExtensionsConfigurer
11+
import org.openprojectx.bigdata.test.extensions.config.NoopBigDataExtensionsConfigurer
912
import org.openprojectx.bigdata.test.extensions.core.BigDataExtensionEvent
1013
import org.openprojectx.bigdata.test.extensions.core.BigDataExtensionResult
1114
import org.openprojectx.bigdata.test.extensions.core.BigDataExtensionRunner
@@ -36,7 +39,8 @@ class BigDataExtensionsExtension : BeforeTestExecutionCallback, AfterAllCallback
3639
val annotation = context.requiredTestClass.getAnnotation(BigDataExtensions::class.java)
3740
?: error("@BigDataExtensions is missing")
3841
val resources = BigDataExtensionResourceLoader(context.requiredTestClass.classLoader)
39-
val extensions = BigDataExtensionsConfigLoader(resources).load(annotation.value.asIterable())
42+
val extensions = BigDataExtensionsConfigLoader(resources).load(annotation.value.asIterable()) +
43+
programmaticExtensions(annotation, context)
4044
val runner = BigDataExtensionRunner(extensions, resources)
4145
val kit = BigDataTestKitStore.get(context)
4246
val result = runner.fire(BigDataExtensionEvent.AFTER_KIT_START, kit)
@@ -46,6 +50,33 @@ class BigDataExtensionsExtension : BeforeTestExecutionCallback, AfterAllCallback
4650
return result
4751
}
4852

53+
private fun programmaticExtensions(
54+
annotation: BigDataExtensions,
55+
context: ExtensionContext,
56+
): List<org.openprojectx.bigdata.test.extensions.core.BigDataExtension> {
57+
val builder = BigDataExtensionsBuilder()
58+
explicitConfigurer(annotation)?.configure(builder)
59+
companionConfigurer(context)?.configure(builder)
60+
instanceConfigurer(context)?.configure(builder)
61+
return builder.build()
62+
}
63+
64+
private fun explicitConfigurer(annotation: BigDataExtensions): BigDataExtensionsConfigurer? {
65+
val type = annotation.configurer.java
66+
if (type == NoopBigDataExtensionsConfigurer::class.java) return null
67+
return type.getDeclaredConstructor().also { it.isAccessible = true }.newInstance()
68+
}
69+
70+
private fun companionConfigurer(context: ExtensionContext): BigDataExtensionsConfigurer? =
71+
runCatching {
72+
context.requiredTestClass.getDeclaredField("Companion")
73+
.also { it.isAccessible = true }
74+
.get(null) as? BigDataExtensionsConfigurer
75+
}.getOrNull()
76+
77+
private fun instanceConfigurer(context: ExtensionContext): BigDataExtensionsConfigurer? =
78+
context.testInstance.orElse(null) as? BigDataExtensionsConfigurer
79+
4980
private val ExtensionContext.store: ExtensionContext.Store
5081
get() = getStore(ExtensionContext.Namespace.create(BigDataExtensionsExtension::class.java, requiredTestClass))
5182

0 commit comments

Comments
 (0)