Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,7 @@
*/
package org.apache.camel.quarkus.support.debezium.deployment;

import java.util.function.BooleanSupplier;

import io.debezium.connector.base.DefaultQueueProvider;
import io.debezium.connector.common.BaseSourceTask;
import io.debezium.embedded.async.ConvertingAsyncEngineBuilderFactory;
import io.debezium.engine.DebeziumEngine;
Expand All @@ -38,35 +37,27 @@
import io.debezium.snapshot.mode.ConfigurationBasedSnapshotter;
import io.debezium.snapshot.mode.InitialOnlySnapshotter;
import io.debezium.snapshot.mode.InitialSnapshotter;
import io.debezium.snapshot.mode.NeverSnapshotter;
import io.debezium.snapshot.mode.NoDataSnapshotter;
import io.debezium.snapshot.mode.RecoverySnapshotter;
import io.debezium.snapshot.mode.WhenNeededNoDataSnapshotter;
import io.debezium.snapshot.mode.WhenNeededSnapshotter;
import io.debezium.snapshot.spi.SnapshotLock;
import io.debezium.storage.file.history.FileSchemaHistory;
import io.quarkus.arc.deployment.AdditionalBeanBuildItem;
import io.quarkus.deployment.annotations.BuildProducer;
import io.quarkus.deployment.annotations.BuildStep;
import io.quarkus.deployment.builditem.BytecodeTransformerBuildItem;
import io.quarkus.deployment.builditem.CombinedIndexBuildItem;
import io.quarkus.deployment.builditem.IndexDependencyBuildItem;
import io.quarkus.deployment.builditem.nativeimage.NativeImageResourceBuildItem;
import io.quarkus.deployment.builditem.nativeimage.ReflectiveClassBuildItem;
import io.quarkus.deployment.builditem.nativeimage.ServiceProviderBuildItem;
import io.quarkus.gizmo.Gizmo;
import org.apache.camel.quarkus.support.debezium.DebeziumComponentObserver;
import org.apache.kafka.connect.source.SourceTask;
import org.jboss.jandex.DotName;
import org.jboss.jandex.IndexView;
import org.jboss.logging.Logger;
import org.objectweb.asm.ClassVisitor;
import org.objectweb.asm.MethodVisitor;
import org.objectweb.asm.Opcodes;

public class DebeziumSupportProcessor {

private static final Logger LOG = Logger.getLogger(DebeziumSupportProcessor.class);

@BuildStep
void addDependencies(BuildProducer<IndexDependencyBuildItem> indexDependency) {
indexDependency.produce(new IndexDependencyBuildItem("org.apache.kafka", "connect-json"));
Expand Down Expand Up @@ -105,7 +96,8 @@ void reflectiveClasses(CombinedIndexBuildItem combinedIndex, BuildProducer<Refle
"io.debezium.storage.kafka.history.KafkaSchemaHistory",
"io.debezium.relational.history.FileDatabaseHistory",
"io.debezium.embedded.ConvertingEngineBuilderFactory",
"io.debezium.processors.PostProcessorRegistry")
"io.debezium.processors.PostProcessorRegistry",
"io.debezium.relational.ConcurrentMapTableMappingStorage")
.build());

reflectiveClasses.produce(ReflectiveClassBuildItem.builder(
Expand All @@ -125,7 +117,7 @@ void reflectiveClasses(CombinedIndexBuildItem combinedIndex, BuildProducer<Refle
NoDataSnapshotter.class,
RecoverySnapshotter.class,
WhenNeededSnapshotter.class,
NeverSnapshotter.class,
WhenNeededNoDataSnapshotter.class,
ConfigurationBasedSnapshotter.class,
SourceSignalChannel.class,
KafkaSignalChannel.class,
Expand All @@ -134,7 +126,8 @@ void reflectiveClasses(CombinedIndexBuildItem combinedIndex, BuildProducer<Refle
InProcessSignalChannel.class,
StandardActionProvider.class,
SourceTask.class,
FileSchemaHistory.class)
FileSchemaHistory.class,
DefaultQueueProvider.class)
.build());

}
Expand Down Expand Up @@ -164,99 +157,9 @@ void registerNativeImageResources(BuildProducer<NativeImageResourceBuildItem> re

resources.produce(new NativeImageResourceBuildItem("META-INF/services/io.debezium.snapshot.spi.SnapshotLock"));
resources.produce(new NativeImageResourceBuildItem("META-INF/services/io.debezium.snapshot.spi.SnapshotQuery"));
resources.produce(new NativeImageResourceBuildItem("META-INF/services/io.debezium.connector.base.QueueProvider"));
resources
.produce(new NativeImageResourceBuildItem("META-INF/services/org.apache.kafka.connect.source.SourceConnector"));
}

// TODO: Remove this - https://github.com/apache/camel-quarkus/issues/8530
@BuildStep(onlyIf = KafkaClients42IsPresent.class)
BytecodeTransformerBuildItem patchConfigInfos() {
// Patch ConfigInfos to add values() method as duplicate of configs()
// This provides backward compatibility for Debezium with kafka-clients 4.2.0
return new BytecodeTransformerBuildItem.Builder()
.setClassToTransform("org.apache.kafka.connect.runtime.rest.entities.ConfigInfos")
.setCacheable(true)
.setVisitorFunction((className, classVisitor) -> new ConfigInfosClassVisitor(classVisitor))
.build();
}

// TODO: Remove this - https://github.com/apache/camel-quarkus/issues/8530
static final class KafkaClients42IsPresent implements BooleanSupplier {
@Override
public boolean getAsBoolean() {
try {
// Check if ConfigInfos.values() is present. If it's not, then kafka-clients >= 4.2.0 is on the classpath
Class<?> configInfos = Thread.currentThread().getContextClassLoader()
.loadClass("org.apache.kafka.connect.runtime.rest.entities.ConfigInfos");
configInfos.getDeclaredMethod("values");
return false;
} catch (ClassNotFoundException e) {
throw new RuntimeException(e);
} catch (NoSuchMethodException e) {
return true;
}
}
}

/**
* Adds a values() method to ConfigInfos that duplicates the configs() method.
* This provides backward compatibility with older Debezium versions.
*/
static class ConfigInfosClassVisitor extends ClassVisitor {

private String configsFieldDescriptor = null;
private String configsMethodSignature = null;

protected ConfigInfosClassVisitor(ClassVisitor classVisitor) {
super(Gizmo.ASM_API_VERSION, classVisitor);
}

@Override
public org.objectweb.asm.FieldVisitor visitField(int access, String name, String descriptor, String signature,
Object value) {
// Track the configs field descriptor
if ("configs".equals(name)) {
configsFieldDescriptor = descriptor;
}
return super.visitField(access, name, descriptor, signature, value);
}

@Override
public MethodVisitor visitMethod(int access, String name, String descriptor, String signature,
String[] exceptions) {
// Track the signature of configs() method
if ("configs".equals(name) && "()Ljava/util/List;".equals(descriptor)) {
configsMethodSignature = signature;
}
return super.visitMethod(access, name, descriptor, signature, exceptions);
}

@Override
public void visitEnd() {
// Add values() method that duplicates configs()
LOG.debug("Adding values() method to ConfigInfos as duplicate of configs()");

MethodVisitor mv = cv.visitMethod(
Opcodes.ACC_PUBLIC,
"values",
"()Ljava/util/List;",
configsMethodSignature, // Same generic signature as configs()
null);

if (mv != null) {
mv.visitCode();
// Method body: return this.configs;
mv.visitVarInsn(Opcodes.ALOAD, 0); // Load 'this'
mv.visitFieldInsn(Opcodes.GETFIELD,
"org/apache/kafka/connect/runtime/rest/entities/ConfigInfos",
"configs",
configsFieldDescriptor != null ? configsFieldDescriptor : "Ljava/util/List;");
mv.visitInsn(Opcodes.ARETURN); // Return the field value
mv.visitMaxs(1, 1); // Max stack=1, max locals=1
mv.visitEnd();
}

super.visitEnd();
}
}
}
2 changes: 1 addition & 1 deletion pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@
<camel.docs.branch>camel-${camel.major.minor}.x</camel.docs.branch><!-- The stable camel branch on which our Antora docs depends -->
<camel-kamelets.version>4.21.0</camel-kamelets.version>
<cassandra-quarkus.version>1.4.1</cassandra-quarkus.version><!-- This should be in sync with quarkus-platform https://repo1.maven.org/maven2/com/datastax/oss/quarkus/cassandra-quarkus-bom/ -->
<debezium.version>3.5.2.Final</debezium.version> <!-- This should be in sync with quarkus-platform https://github.com/quarkusio/quarkus-platform-->
<debezium.version>3.6.0.Final</debezium.version> <!-- This should be in sync with quarkus-platform https://github.com/quarkusio/quarkus-platform-->
<rocketmq.version>${rocketmq-version}</rocketmq.version>
<optaplanner.version>10.0.0</optaplanner.version><!-- This should be in sync with quarkus-platform https://repo1.maven.org/maven2/org/optaplanner/optaplanner-quarkus/ -->
<quarkiverse-amazonservices.version>3.21.0</quarkiverse-amazonservices.version><!-- This should be in sync with quarkus-platform https://repo1.maven.org/maven2/io/quarkiverse/amazonservices/quarkus-amazon-services-parent/ -->
Expand Down
Loading
Loading