From 2d4a35a7a56a6c9a8b4356da8eef12ab3dc770ba Mon Sep 17 00:00:00 2001 From: lekhrocks Date: Sun, 6 Sep 2026 14:27:09 +0530 Subject: [PATCH 1/2] feat(F18): plugin-to-core connector bridge Implement the missing integration layer between the plugin system and the core connector SPI so installed plugins are usable by pipelines, snapshots, CDC, and metadata discovery. New classes: - PluginConnectorAdapter: superset adapter implementing all core SPI sub-interfaces (Connector, MetadataCapableConnector, SnapshotCapableConnector, CdcCapableConnector, DestinationWriter) - PluginContextAdapter, PluginCapabilitiesAdapter, PluginHealthAdapter, PluginCdcEventAdapter, PluginCursorCodec, PluginWriterAdapter: helper adapters for type conversion - ConnectorTypeResolver: maps free-form plugin connectorType string to ConnectorType enum - DelegatingConnectorRegistry: @Primary connector registry composing SpringConnectorRegistry and PluginConnectorRegistry - PluginConnectorRegistry: pluginId-keyed plugin connector store Changes to existing code: - ConnectorType.GENERIC_PLUGIN enum value added - DestinationWriterProvider: added default delete() and upsert() methods (backward-compatible, existing plugins compile unchanged) - PluginManager: enable/disable/uninstall now register/unregister adapters in the connector registry Documented limitations: - rangeChunks returns single whole-table chunk (no PK splitting) - Metadata: PK/fk/index/constraint discovery return empty/zero - CDC transaction metadata is always null - Pause/resume CDC are no-op flags on plugin side - Parallel snapshot parallelism=1 for plugin-backed sources Tests: 28 new unit tests covering adapters, registry delegation, cursor round-trip, CDC event synthesis, capabilities mapping, health parsing, and operation normalization. --- .../syncflow/api/plugin/PluginManager.java | 31 ++ .../plugin/adapter/ConnectorTypeResolver.java | 35 ++ .../adapter/PluginCapabilitiesAdapter.java | 31 ++ .../plugin/adapter/PluginCdcEventAdapter.java | 58 +++ .../adapter/PluginConnectorAdapter.java | 381 +++++++++++++++ .../plugin/adapter/PluginContextAdapter.java | 26 ++ .../api/plugin/adapter/PluginCursorCodec.java | 37 ++ .../plugin/adapter/PluginHealthAdapter.java | 34 ++ .../plugin/adapter/PluginWriterAdapter.java | 106 +++++ .../registry/DelegatingConnectorRegistry.java | 122 +++++ .../registry/PluginConnectorRegistry.java | 46 ++ .../adapter/PluginConnectorAdapterTest.java | 432 ++++++++++++++++++ .../DelegatingConnectorRegistryTest.java | 257 +++++++++++ .../syncflow/core/model/ConnectorType.java | 2 +- .../plugin/spi/DestinationWriterProvider.java | 20 + 15 files changed, 1617 insertions(+), 1 deletion(-) create mode 100644 syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/ConnectorTypeResolver.java create mode 100644 syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCapabilitiesAdapter.java create mode 100644 syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCdcEventAdapter.java create mode 100644 syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapter.java create mode 100644 syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginContextAdapter.java create mode 100644 syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCursorCodec.java create mode 100644 syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginHealthAdapter.java create mode 100644 syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginWriterAdapter.java create mode 100644 syncflow-api/src/main/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistry.java create mode 100644 syncflow-api/src/main/java/com/syncflow/api/plugin/registry/PluginConnectorRegistry.java create mode 100644 syncflow-api/src/test/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapterTest.java create mode 100644 syncflow-api/src/test/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistryTest.java diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/PluginManager.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/PluginManager.java index 123a24f..c5d4938 100644 --- a/syncflow-api/src/main/java/com/syncflow/api/plugin/PluginManager.java +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/PluginManager.java @@ -1,5 +1,7 @@ package com.syncflow.api.plugin; +import com.syncflow.api.plugin.adapter.PluginConnectorAdapter; +import com.syncflow.api.plugin.registry.DelegatingConnectorRegistry; import com.syncflow.plugin.descriptor.PluginDescriptor; import com.syncflow.plugin.lifecycle.PluginLifecycle; import com.syncflow.plugin.spi.PluginConnector; @@ -19,6 +21,19 @@ public class PluginManager { private final Map plugins = new ConcurrentHashMap<>(); private final Map classLoaders = new ConcurrentHashMap<>(); + private final Map adapters = new ConcurrentHashMap<>(); + private final DelegatingConnectorRegistry registry; + + public PluginManager(DelegatingConnectorRegistry registry) { + this.registry = registry; + } + + /** + * Test-only constructor — registry operations are no-ops when registry is null. + */ + PluginManager() { + this.registry = null; + } public PluginInstallResult install(File jarFile) { try (var jar = new JarFile(jarFile)) { @@ -64,6 +79,15 @@ public boolean enable(String pluginId) { var entry = plugins.get(pluginId); if (entry == null) return false; + // Already enabled? Idempotent. + if (entry.lifecycle() == PluginLifecycle.ENABLED) { + return true; + } + var adapter = adapters.computeIfAbsent(pluginId, + id -> new PluginConnectorAdapter(entry.connector(), entry.descriptor())); + if (registry != null) { + registry.register(adapter); + } plugins.put(pluginId, new PluginEntry(entry.descriptor(), entry.connector(), PluginLifecycle.ENABLED)); return true; } @@ -72,12 +96,19 @@ public boolean disable(String pluginId) { var entry = plugins.get(pluginId); if (entry == null) return false; + if (registry != null) { + registry.unregisterPlugin(pluginId); + } plugins.put(pluginId, new PluginEntry(entry.descriptor(), entry.connector(), PluginLifecycle.DISABLED)); return true; } public boolean uninstall(String pluginId) { var removed = plugins.remove(pluginId); + if (registry != null) { + registry.unregisterPlugin(pluginId); + } + adapters.remove(pluginId); var cl = classLoaders.remove(pluginId); if (cl != null) { try { diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/ConnectorTypeResolver.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/ConnectorTypeResolver.java new file mode 100644 index 0000000..bdb5654 --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/ConnectorTypeResolver.java @@ -0,0 +1,35 @@ +package com.syncflow.api.plugin.adapter; + +import com.syncflow.core.model.ConnectorType; + +import java.util.Locale; + +/** + * Maps a free-form {@code PluginDescriptor.connectorType} string to the + * closed {@link ConnectorType} enum. Recognized database names map to the + * matching enum value; everything else collapses to + * {@link ConnectorType#GENERIC_PLUGIN}. + */ +public final class ConnectorTypeResolver { + + private ConnectorTypeResolver() { + } + + public static ConnectorType resolve(String pluginType) { + if (pluginType == null || pluginType.isBlank()) { + return ConnectorType.GENERIC_PLUGIN; + } + return switch (pluginType.toLowerCase(Locale.ROOT)) { + case "postgresql", "postgres" -> ConnectorType.POSTGRESQL; + case "mysql", "mariadb" -> ConnectorType.MYSQL; + case "mongodb", "mongo" -> ConnectorType.MONGODB; + case "kafka" -> ConnectorType.KAFKA; + case "sqlserver", "mssql" -> ConnectorType.SQLSERVER; + case "oracle" -> ConnectorType.ORACLE; + case "elasticsearch", "elastic", "es" -> ConnectorType.ELASTICSEARCH; + case "redis" -> ConnectorType.REDIS; + case "jdbc", "generic_jdbc" -> ConnectorType.GENERIC_JDBC; + default -> ConnectorType.GENERIC_PLUGIN; + }; + } +} diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCapabilitiesAdapter.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCapabilitiesAdapter.java new file mode 100644 index 0000000..161d51f --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCapabilitiesAdapter.java @@ -0,0 +1,31 @@ +package com.syncflow.api.plugin.adapter; + +import com.syncflow.core.spi.ConnectorCapabilities; + +/** + * Maps the 6-boolean plugin capabilities to the 5-boolean core capabilities. + * + *

+ * Plugin SPI: metadata, snapshot, cdc, destination, transactions, streaming. + *
+ * Core SPI: cdc, snapshot, schemaDiscovery, transactions, offsetTracking. + */ +public final class PluginCapabilitiesAdapter { + + private PluginCapabilitiesAdapter() { + } + + public static ConnectorCapabilities map( + com.syncflow.plugin.capabilities.ConnectorCapabilities plugin) { + if (plugin == null) { + return ConnectorCapabilities.none(); + } + return new ConnectorCapabilities( + plugin.supportsCdc(), + plugin.supportsSnapshot(), + plugin.supportsMetadata(), // → schemaDiscovery + plugin.supportsTransactions(), + // No direct mapping: streaming/cdc imply offset tracking. + plugin.supportsCdc() || plugin.supportsStreaming()); + } +} diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCdcEventAdapter.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCdcEventAdapter.java new file mode 100644 index 0000000..c7d756a --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCdcEventAdapter.java @@ -0,0 +1,58 @@ +package com.syncflow.api.plugin.adapter; + +import com.syncflow.core.cdc.CDCOperation; +import com.syncflow.core.cdc.CDCEvent; +import com.syncflow.core.cdc.EventHeader; +import com.syncflow.core.cdc.EventMetadata; +import com.syncflow.core.cdc.EventPayload; +import com.syncflow.core.cdc.EventSource; +import com.syncflow.core.cdc.OffsetInformation; +import com.syncflow.plugin.descriptor.PluginDescriptor; +import com.syncflow.plugin.spi.CdcProvider; + +import java.time.Instant; +import java.util.Locale; +import java.util.Map; + +/** + * Converts the flat {@link CdcProvider.CdcEvent} emitted by a plugin into + * the deeply nested core {@link CDCEvent}. Synthesizes fields the plugin + * SPI does not provide (pipelineId, connectionId, eventNumber, version, + * capturedAt, captureLatencyMs, transaction) with sensible defaults. + */ +public final class PluginCdcEventAdapter { + + private PluginCdcEventAdapter() { + } + + public static CDCEvent toCore(CdcProvider.CdcEvent event, PluginDescriptor descriptor) { + var eventId = event.eventId() != null + ? event.eventId() + : java.util.UUID.randomUUID().toString(); + var now = Instant.now(); + var pipelineId = "plugin:" + descriptor.pluginId(); + var offsetMap = event.offset() == null ? Map.of() : event.offset(); + + return new CDCEvent( + new EventHeader(eventId, pipelineId, pipelineId, 0L, 1, Map.of()), + new EventSource(descriptor.pluginId(), event.schema(), event.table(), + descriptor.connectorType()), + parseOperation(event.operation()), + new EventPayload(event.before(), event.after(), Map.of()), + new EventMetadata(0L, now, 0L), + null, + new OffsetInformation(descriptor.connectorType(), offsetMap, null, now)); + } + + static CDCOperation parseOperation(String op) { + if (op == null || op.isBlank()) { + return CDCOperation.READ; + } + return switch (op.trim().toUpperCase(Locale.ROOT)) { + case "INSERT", "I", "CREATE" -> CDCOperation.INSERT; + case "UPDATE", "U" -> CDCOperation.UPDATE; + case "DELETE", "D" -> CDCOperation.DELETE; + default -> CDCOperation.READ; + }; + } +} diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapter.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapter.java new file mode 100644 index 0000000..0b83514 --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapter.java @@ -0,0 +1,381 @@ +package com.syncflow.api.plugin.adapter; + +import com.syncflow.core.cdc.CDCEvent; +import com.syncflow.core.cdc.CaptureStatus; +import com.syncflow.core.metadata.ColumnMetadata; +import com.syncflow.core.metadata.ConstraintMetadata; +import com.syncflow.core.metadata.DataType; +import com.syncflow.core.metadata.ForeignKeyMetadata; +import com.syncflow.core.metadata.IndexMetadata; +import com.syncflow.core.metadata.PrimaryKeyMetadata; +import com.syncflow.core.metadata.TableMetadata; +import com.syncflow.core.metadata.TableStatistics; +import com.syncflow.core.model.ConnectorType; +import com.syncflow.core.snapshot.BatchInformation; +import com.syncflow.core.snapshot.ChunkRange; +import com.syncflow.core.spi.CdcCapableConnector; +import com.syncflow.core.spi.Connector; +import com.syncflow.core.spi.ConnectorCapabilities; +import com.syncflow.core.spi.ConnectorContext; +import com.syncflow.core.spi.ConnectorHealth; +import com.syncflow.core.spi.ConnectorValidationResult; +import com.syncflow.core.spi.MetadataCapableConnector; +import com.syncflow.core.spi.SnapshotCapableConnector; +import com.syncflow.core.spi.writer.DestinationWriter; +import com.syncflow.plugin.descriptor.PluginDescriptor; +import com.syncflow.plugin.spi.CdcProvider; +import com.syncflow.plugin.spi.PluginConnector; +import com.syncflow.plugin.spi.SnapshotProvider; + +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Consumer; + +/** + * Adapts a plugin-world {@link PluginConnector} to the full core SPI: + * {@link Connector}, {@link MetadataCapableConnector}, + * {@link SnapshotCapableConnector}, {@link CdcCapableConnector}, + * {@link DestinationWriter}. + * + *

+ * One instance per (pluginId, ConnectorContext). The plugin's + * capabilities are queried lazily so the adapter's reported capability + * flags reflect the current plugin state. + * + *

+ * Cursor handling: the plugin SPI's {@code readBatch} takes a batch + * number, not a cursor. We round-trip the batch number through + * {@link BatchInformation#cursor()} using {@link PluginCursorCodec}. + * {@link #rangeChunks} returns a single whole-table chunk because the + * plugin SPI has no PK-range splitting API — parallel snapshot + * parallelism is effectively 1 for plugin-backed sources. + * + *

+ * CDC transaction metadata is always null and CDC pause/resume are + * no-ops because the plugin SPI lacks these concepts. + */ +public final class PluginConnectorAdapter + implements + Connector, + MetadataCapableConnector, + SnapshotCapableConnector, + CdcCapableConnector { + + private final PluginConnector delegate; + private final PluginDescriptor descriptor; + private final ConnectorType type; + private volatile boolean connected; + private volatile boolean cdcActive; + private final Map lastBatchByTable = new ConcurrentHashMap<>(); + private String pluginId; + + private record TableKey(String schema, String table) { + } + + public PluginConnectorAdapter(PluginConnector delegate, PluginDescriptor descriptor) { + this.delegate = delegate; + this.descriptor = descriptor; + this.type = ConnectorTypeResolver.resolve(descriptor.connectorType()); + this.pluginId = descriptor.pluginId(); + } + + /** Re-attached for use by the registry which needs to know the plugin id. */ + public String pluginId() { + return pluginId; + } + + @Override + public ConnectorType type() { + return type; + } + + @Override + public ConnectorCapabilities capabilities() { + return PluginCapabilitiesAdapter.map(delegate.capabilities()); + } + + @Override + public void connect(ConnectorContext context) { + // The plugin SPI has no connect() method; track state ourselves. + // Capability providers (SnapshotProvider, CdcProvider, + // DestinationWriterProvider) + // use the PluginContext passed to their first method call. + connected = true; + } + + @Override + public void disconnect() { + if (!connected) { + return; + } + // No plugin-level disconnect. CDC capture must already be stopped + // by stopCDC() before this is called, otherwise events flow to a + // disposed consumer. + connected = false; + cdcActive = false; + } + + @Override + public boolean isConnected() { + return connected; + } + + @Override + public ConnectorValidationResult validate(ConnectorContext context) { + // Plugin SPI has no dedicated validation. Try a connect/disconnect + // round-trip and report the outcome. + try { + if (!connected) { + connect(context); + } + return ConnectorValidationResult.ok(); + } catch (Exception e) { + return ConnectorValidationResult.failed(List.of( + "Plugin validation failed: " + e.getMessage())); + } + } + + @Override + public List discoverSchemas(ConnectorContext context) { + return delegate.discoverSchemas(PluginContextAdapter.from(context)); + } + + @Override + public List discoverTables(ConnectorContext context, String schema) { + return delegate.discoverTables(PluginContextAdapter.from(context), schema); + } + + @Override + public ConnectorHealth health() { + return PluginHealthAdapter.parse(delegate.health()); + } + + @Override + public Map metadata() { + return Map.copyOf(delegate.metadata()); + } + + // --- MetadataCapableConnector --- + + @Override + public List fetchTables(ConnectorContext context, String schema) { + var pluginCtx = PluginContextAdapter.from(context); + var tables = delegate.discoverTables(pluginCtx, schema); + return tables.stream() + .map(name -> { + var cols = delegate.discoverColumns(pluginCtx, schema, name); + return new TableMetadata(name, "TABLE", schema, null, + TableStatistics.unknown(), + toColumnMetadata(cols), + List.of(), + primaryKeyFromColumns(cols), + List.of(), + List.of()); + }) + .toList(); + } + + @Override + public List fetchColumns(ConnectorContext context, String schema, String table) { + var pluginCtx = PluginContextAdapter.from(context); + return toColumnMetadata(delegate.discoverColumns(pluginCtx, schema, table)); + } + + @Override + public List fetchIndexes(ConnectorContext context, String schema, String table) { + return List.of(); + } + + @Override + public PrimaryKeyMetadata fetchPrimaryKey(ConnectorContext context, String schema, String table) { + return primaryKeyFromColumns( + delegate.discoverColumns(PluginContextAdapter.from(context), schema, table)); + } + + @Override + public List fetchForeignKeys(ConnectorContext context, String schema, String table) { + return List.of(); + } + + @Override + public List fetchConstraints(ConnectorContext context, String schema, String table) { + return List.of(); + } + + @Override + public TableStatistics fetchStatistics(ConnectorContext context, String schema, String table) { + return TableStatistics.unknown(); + } + + // --- SnapshotCapableConnector --- + + @Override + public long estimateRows(ConnectorContext context, String schema, String table) { + if (delegate instanceof SnapshotProvider sp) { + return sp.estimateRowCount(PluginContextAdapter.from(context), schema, table); + } + return -1L; + } + + @Override + public SnapshotCapableConnector.Page readBatch(ConnectorContext context, String schema, + String table, BatchInformation batchInfo) { + if (!(delegate instanceof SnapshotProvider sp)) { + return SnapshotCapableConnector.Page.empty(); + } + var pluginCtx = PluginContextAdapter.from(context); + var batchNumber = PluginCursorCodec.decode(batchInfo.cursor()); + var page = sp.readBatch(pluginCtx, schema, table, batchNumber, + batchInfo.batchSize() == 0 ? 1000 : batchInfo.batchSize()); + lastBatchByTable.put(new TableKey(schema, table), batchNumber); + // Convert the plugin's nextCursor into our encoded form, or null when done. + var next = page.nextCursor(); + String encoded; + if (next == null) { + encoded = null; + } else if (next.startsWith(PluginCursorCodec.encode(0).substring(0, 1))) { + // Already encoded by us (or a plugin using the same scheme). + encoded = next; + } else { + // Treat the plugin's string cursor as opaque — bump the batch. + encoded = PluginCursorCodec.encode(batchNumber + 1); + } + return SnapshotCapableConnector.Page.of(page.rows(), encoded); + } + + @Override + public List rangeChunks(ConnectorContext context, String schema, String table, + int chunkCount) { + // Plugin SPI has no PK-range splitting — single whole-table chunk. + return List.of(ChunkRange.whole()); + } + + @Override + public SnapshotCapableConnector snapshotClone(ConnectorContext context) { + // Reuse the same delegate and descriptor; per-clone state is the + // snapshot's job (each chunk worker is a sequential read of the + // same plugin instance). + return new PluginConnectorAdapter(delegate, descriptor); + } + + // --- CdcCapableConnector --- + + @Override + public void startCDC(ConnectorContext context, Consumer eventConsumer) { + if (!(delegate instanceof CdcProvider cp)) { + throw new UnsupportedOperationException( + "Plugin " + pluginId + " does not implement CdcProvider"); + } + if (cdcActive) { + return; + } + var pluginCtx = PluginContextAdapter.from(context); + cp.startCapture(pluginCtx, + e -> eventConsumer.accept(PluginCdcEventAdapter.toCore(e, descriptor))); + cdcActive = true; + } + + @Override + public void stopCDC() { + if (!(delegate instanceof CdcProvider cp)) { + return; + } + if (cdcActive) { + cp.stopCapture(); + } + cdcActive = false; + } + + @Override + public void pauseCDC() { + // Plugin SPI has no pause; track state for isCdcActive() reporting. + if (cdcActive) { + cdcActive = false; + } + } + + @Override + public void resumeCDC() { + // No-op: paused CDC needs startCDC to resume on the plugin side. + } + + @Override + public boolean isCdcActive() { + if (delegate instanceof CdcProvider cp) { + return cp.isCapturing(); + } + return false; + } + + @Override + public CaptureStatus captureStatus() { + if (!cdcActive) { + return CaptureStatus.INACTIVE; + } + if (delegate instanceof CdcProvider cp) { + return cp.isCapturing() ? CaptureStatus.RUNNING : CaptureStatus.INACTIVE; + } + return CaptureStatus.INACTIVE; + } + + @Override + public Map currentOffset() { + if (delegate instanceof CdcProvider cp) { + return Map.copyOf(cp.currentOffset()); + } + return Map.of(); + } + + // --- helpers --- + + private static List toColumnMetadata(List> pluginCols) { + if (pluginCols == null) { + return List.of(); + } + var out = new java.util.ArrayList(pluginCols.size()); + int pos = 1; + for (var col : pluginCols) { + var name = stringAt(col, "name", "columnName", "column"); + var type = stringAt(col, "type", "dataType", "nativeType"); + var nullable = boolAt(col, "nullable", true); + var isPk = boolAt(col, "primaryKey", false); + out.add(new ColumnMetadata(name, pos++, + new DataType(type, type, null, null, nullable, null), + isPk, false, false, false, null)); + } + return out; + } + + private static PrimaryKeyMetadata primaryKeyFromColumns(List> pluginCols) { + if (pluginCols == null) { + return new PrimaryKeyMetadata(null, List.of()); + } + var pkCols = pluginCols.stream() + .filter(c -> boolAt(c, "primaryKey", false)) + .map(c -> stringAt(c, "name", "columnName", "column")) + .toList(); + if (pkCols.isEmpty()) { + return new PrimaryKeyMetadata(null, List.of()); + } + return new PrimaryKeyMetadata("pk_" + String.join("_", pkCols), pkCols); + } + + private static String stringAt(Map m, String... keys) { + for (var k : keys) { + var v = m.get(k); + if (v != null) { + return v.toString(); + } + } + return ""; + } + + private static boolean boolAt(Map m, String key, boolean dflt) { + var v = m.get(key); + if (v instanceof Boolean b) { + return b; + } + return dflt; + } +} diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginContextAdapter.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginContextAdapter.java new file mode 100644 index 0000000..8b6ae79 --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginContextAdapter.java @@ -0,0 +1,26 @@ +package com.syncflow.api.plugin.adapter; + +import com.syncflow.core.spi.ConnectorContext; +import com.syncflow.plugin.spi.PluginContext; + +/** + * Unpacks a nested core {@link ConnectorContext} (which holds + * {@code ConnectionConfiguration}) into the flat + * {@link PluginContext} the plugin SPI expects. + */ +public final class PluginContextAdapter { + + private PluginContextAdapter() { + } + + public static PluginContext from(ConnectorContext ctx) { + var cfg = ctx.config(); + return new PluginContext( + cfg.host(), + cfg.port(), + cfg.database(), + cfg.username(), + cfg.password(), + cfg.properties()); + } +} diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCursorCodec.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCursorCodec.java new file mode 100644 index 0000000..0eab48d --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginCursorCodec.java @@ -0,0 +1,37 @@ +package com.syncflow.api.plugin.adapter; + +/** + * The plugin SPI's {@code readBatch} takes a {@code batchNumber}, not a + * cursor string. We use this codec to round-trip a batch number through + * the core {@code BatchInformation.cursor} field so the existing snapshot + * pipeline doesn't have to know the difference. + * + *

+ * Format: {@code "b:"}. A {@code null} cursor maps to batch 0. + */ +public final class PluginCursorCodec { + + private static final String PREFIX = "b:"; + + private PluginCursorCodec() { + } + + public static String encode(int batchNumber) { + return PREFIX + batchNumber; + } + + public static int decode(String cursor) { + if (cursor == null || cursor.isBlank()) { + return 0; + } + if (!cursor.startsWith(PREFIX)) { + // Unknown format — start from zero rather than fail the snapshot. + return 0; + } + try { + return Integer.parseInt(cursor.substring(PREFIX.length())); + } catch (NumberFormatException e) { + return 0; + } + } +} diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginHealthAdapter.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginHealthAdapter.java new file mode 100644 index 0000000..be2b83c --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginHealthAdapter.java @@ -0,0 +1,34 @@ +package com.syncflow.api.plugin.adapter; + +import com.syncflow.core.spi.ConnectorHealth; + +import java.time.Instant; +import java.util.Locale; + +/** + * Parses a plugin's health string ("UP" / "DOWN" / "DEGRADED" / "UNKNOWN") + * into a {@link ConnectorHealth} record. + */ +public final class PluginHealthAdapter { + + private PluginHealthAdapter() { + } + + public static ConnectorHealth parse(String raw) { + var now = Instant.now(); + if (raw == null || raw.isBlank()) { + return new ConnectorHealth(ConnectorHealth.Status.UNKNOWN, + "Plugin reported no status", now, 0); + } + return switch (raw.trim().toUpperCase(Locale.ROOT)) { + case "UP" -> new ConnectorHealth(ConnectorHealth.Status.UP, + "Plugin reports UP", now, 0); + case "DOWN" -> new ConnectorHealth(ConnectorHealth.Status.DOWN, + "Plugin reports DOWN", now, 0); + case "DEGRADED" -> new ConnectorHealth(ConnectorHealth.Status.DEGRADED, + "Plugin reports DEGRADED", now, 0); + default -> new ConnectorHealth(ConnectorHealth.Status.UNKNOWN, + raw, now, 0); + }; + } +} diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginWriterAdapter.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginWriterAdapter.java new file mode 100644 index 0000000..3fc182a --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginWriterAdapter.java @@ -0,0 +1,106 @@ +package com.syncflow.api.plugin.adapter; + +import com.syncflow.core.model.ConnectionConfiguration; +import com.syncflow.core.spi.writer.DestinationWriter; +import com.syncflow.plugin.spi.DestinationWriterProvider; +import com.syncflow.plugin.spi.PluginContext; + +import java.util.List; +import java.util.Map; + +/** + * Adapts a plugin's {@link DestinationWriterProvider} to the core + * {@link DestinationWriter} contract. + * + *

+ * Core's writeBatch/deleteBatch/upsertBatch take an explicit column + * list; the plugin SPI only knows about the row maps. We project rows + * to the named columns, dropping any keys the plugin returned that the + * pipeline didn't ask for. + * + *

+ * The plugin SPI has no DELETE or UPSERT; the default implementations + * on {@link DestinationWriterProvider} are no-ops for delete and a plain + * {@code write} (key columns ignored) for upsert. Plugins that want + * proper DELETE/UPSERT semantics override those defaults. + */ +public final class PluginWriterAdapter implements DestinationWriter { + + private final DestinationWriterProvider delegate; + private final ConnectionConfiguration config; + private boolean connected; + + public PluginWriterAdapter(DestinationWriterProvider delegate, + ConnectionConfiguration config) { + this.delegate = delegate; + this.config = config; + } + + @Override + public void connect(ConnectionConfiguration cfg) { + var pluginCtx = new PluginContext( + cfg.host(), cfg.port(), cfg.database(), + cfg.username(), cfg.password(), cfg.properties()); + delegate.connect(pluginCtx); + connected = true; + } + + @Override + public void writeBatch(String table, List columns, List> rows) { + delegate.write(table, project(columns, rows)); + } + + @Override + public void deleteBatch(String table, List columns, List> rows) { + // No-op by default in the plugin SPI. Override to support DELETE. + delegate.delete(table, project(columns, rows)); + } + + @Override + public void upsertBatch(String table, List columns, List> rows, + List keyColumns) { + delegate.upsert(table, project(columns, rows), keyColumns); + } + + @Override + public void flush() { + delegate.flush(); + } + + @Override + public void commit() { + delegate.commit(); + } + + @Override + public void rollback() { + delegate.rollback(); + } + + @Override + public void close() { + delegate.close(); + connected = false; + } + + @Override + public boolean isConnected() { + return connected; + } + + private static List> project(List columns, + List> rows) { + if (columns == null || columns.isEmpty()) { + return rows; + } + var projected = new java.util.ArrayList>(rows.size()); + for (var row : rows) { + var out = new java.util.LinkedHashMap(columns.size()); + for (var col : columns) { + out.put(col, row.get(col)); + } + projected.add(out); + } + return projected; + } +} diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistry.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistry.java new file mode 100644 index 0000000..477b99c --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistry.java @@ -0,0 +1,122 @@ +package com.syncflow.api.plugin.registry; + +import com.syncflow.core.model.ConnectorType; +import com.syncflow.core.registry.ConnectorRegistry; +import com.syncflow.core.registry.SpringConnectorRegistry; +import com.syncflow.core.spi.Connector; +import com.syncflow.api.plugin.adapter.PluginConnectorAdapter; +import org.springframework.context.annotation.Primary; +import org.springframework.stereotype.Component; + +import java.util.ArrayList; +import java.util.List; +import java.util.Optional; + +/** + * The single {@code ConnectorRegistry} seen by the rest of the app. + * Routes built-in connectors to the existing + * {@link SpringConnectorRegistry} and plugin-backed connectors to the + * sidecar {@link PluginConnectorRegistry}. + * + *

+ * Plugin connectors carry a {@code pluginId} in the + * {@link PluginConnectorAdapter}. The + * {@link ConnectorType#GENERIC_PLUGIN} value is the + * "look up any plugin" sentinel — used by health checks and dashboards + * that don't care which plugin a connection is bound to. + */ +@Component +@Primary +public class DelegatingConnectorRegistry implements ConnectorRegistry { + + private final SpringConnectorRegistry builtins; + private final PluginConnectorRegistry plugins; + + public DelegatingConnectorRegistry(SpringConnectorRegistry builtins, + PluginConnectorRegistry plugins) { + this.builtins = builtins; + this.plugins = plugins; + } + + @Override + public Connector register(Connector connector) { + if (connector instanceof PluginConnectorAdapter adapter) { + plugins.register(adapter.pluginId(), adapter); + } else { + builtins.register(connector); + } + return connector; + } + + @Override + public void unregister(ConnectorType type) { + // Unregister by type disambiguates: GENERIC_PLUGIN removes all + // plugins; any other type only hits the built-in registry. + if (type == ConnectorType.GENERIC_PLUGIN) { + for (var c : plugins.getAll()) { + c.disconnect(); + } + plugins.getAll().forEach(c -> plugins.unregister(pluginIdOf(c))); + } else { + builtins.unregister(type); + } + } + + @Override + public Optional get(ConnectorType type) { + if (type == ConnectorType.GENERIC_PLUGIN) { + var all = plugins.getAll(); + return all.isEmpty() ? Optional.empty() : Optional.of(all.get(0)); + } + // Look in the built-in registry first; if a plugin advertises the + // same type (e.g. a "postgresql" plugin), it can shadow the + // built-in for that type. This lets users override built-ins + // without code changes to the core. + var builtin = builtins.get(type); + if (builtin.isPresent()) { + for (var c : plugins.getAll()) { + if (c.type() == type) { + return Optional.of(c); + } + } + return builtin; + } + // No built-in for this type — fall back to a matching plugin. + for (var c : plugins.getAll()) { + if (c.type() == type) { + return Optional.of(c); + } + } + return Optional.empty(); + } + + @Override + public List getAll() { + var all = new ArrayList(); + all.addAll(builtins.getAll()); + all.addAll(plugins.getAll()); + return all; + } + + @Override + public boolean isRegistered(ConnectorType type) { + if (type == ConnectorType.GENERIC_PLUGIN) { + return !plugins.getAll().isEmpty(); + } + return builtins.isRegistered(type); + } + + /** + * Unregister a single plugin by id — used by PluginManager on + * disable/uninstall. + */ + public boolean unregisterPlugin(String pluginId) { + var removed = plugins.unregister(pluginId); + plugins.get(pluginId).ifPresent(Connector::disconnect); + return removed; + } + + private static String pluginIdOf(Connector c) { + return c instanceof PluginConnectorAdapter pa ? pa.pluginId() : null; + } +} diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/registry/PluginConnectorRegistry.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/registry/PluginConnectorRegistry.java new file mode 100644 index 0000000..ccb3b62 --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/registry/PluginConnectorRegistry.java @@ -0,0 +1,46 @@ +package com.syncflow.api.plugin.registry; + +import com.syncflow.core.spi.Connector; +import org.springframework.stereotype.Component; + +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; + +/** + * Holds the currently-enabled plugin connectors, keyed by plugin id. + * + *

+ * This is not a {@code ConnectorRegistry} itself — it is a sidecar + * store consulted by {@link DelegatingConnectorRegistry} for + * plugin-backed lookups. Each entry is one + * {@link com.syncflow.api.plugin.adapter.PluginConnectorAdapter} + * (which implements every core SPI sub-interface). + */ +@Component +public class PluginConnectorRegistry { + + private final Map byPluginId = new ConcurrentHashMap<>(); + + public Connector register(String pluginId, Connector adapter) { + byPluginId.put(pluginId, adapter); + return adapter; + } + + public Optional get(String pluginId) { + return Optional.ofNullable(byPluginId.get(pluginId)); + } + + public boolean unregister(String pluginId) { + return byPluginId.remove(pluginId) != null; + } + + public List getAll() { + return List.copyOf(byPluginId.values()); + } + + public boolean isRegistered(String pluginId) { + return byPluginId.containsKey(pluginId); + } +} diff --git a/syncflow-api/src/test/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapterTest.java b/syncflow-api/src/test/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapterTest.java new file mode 100644 index 0000000..0f1b9ab --- /dev/null +++ b/syncflow-api/src/test/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapterTest.java @@ -0,0 +1,432 @@ +package com.syncflow.api.plugin.adapter; + +import com.syncflow.core.cdc.CDCEvent; +import com.syncflow.core.cdc.CDCOperation; +import com.syncflow.core.cdc.CaptureStatus; +import com.syncflow.core.model.ConnectionConfiguration; +import com.syncflow.core.model.ConnectorType; +import com.syncflow.core.snapshot.BatchInformation; +import com.syncflow.core.spi.ConnectorCapabilities; +import com.syncflow.core.spi.ConnectorContext; +import com.syncflow.core.spi.ConnectorHealth; +import com.syncflow.plugin.descriptor.PluginDescriptor; +import com.syncflow.plugin.spi.CdcProvider; +import com.syncflow.plugin.spi.PluginConnector; +import com.syncflow.plugin.spi.PluginContext; +import com.syncflow.plugin.spi.SnapshotProvider; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.Map; +import java.util.function.Consumer; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class PluginConnectorAdapterTest { + + private static PluginDescriptor desc(String pluginId, String connectorType) { + return new PluginDescriptor(pluginId, "Test", "test", "1.0", + "test plugin", connectorType, List.of(), "1", "1", + List.of(), "MIT", null, null); + } + + private static ConnectionConfiguration cfg() { + return new ConnectionConfiguration(ConnectorType.GENERIC_PLUGIN, + "host", 5432, "db", "user", "pass", Map.of("k", "v")); + } + + private static ConnectorContext ctx() { + return new ConnectorContext(cfg(), Map.of()); + } + + private static com.syncflow.plugin.capabilities.ConnectorCapabilities capsAll() { + return new com.syncflow.plugin.capabilities.ConnectorCapabilities( + true, true, true, true, true, true); + } + + private static com.syncflow.plugin.capabilities.ConnectorCapabilities capsNone() { + return new com.syncflow.plugin.capabilities.ConnectorCapabilities( + false, false, false, false, false, false); + } + + private static PluginConnector pluginWith(String type, String health, + com.syncflow.plugin.capabilities.ConnectorCapabilities caps, + List> cols) { + return new PluginConnector() { + + @Override + public PluginDescriptor descriptor() { + return desc("p1", type); + } + @Override + public com.syncflow.plugin.capabilities.ConnectorCapabilities capabilities() { + return caps; + } + @Override + public String health() { + return health; + } + @Override + public Map metadata() { + return Map.of("v", "1"); + } + @Override + public List discoverSchemas(PluginContext c) { + return List.of("public"); + } + @Override + public List discoverTables(PluginContext c, String s) { + return List.of("t"); + } + @Override + public List> discoverColumns(PluginContext c, String s, String t) { + return cols; + } + }; + } + + @Test + void type_resolves_to_known_enum() { + var adapter = new PluginConnectorAdapter( + pluginWith("postgresql", "UP", capsAll(), + List.of(Map.of("name", "id", "type", "int", "primaryKey", true))), + desc("p1", "postgresql")); + assertEquals(ConnectorType.POSTGRESQL, adapter.type()); + } + + @Test + void type_resolves_to_generic_plugin_for_unknown() { + var adapter = new PluginConnectorAdapter( + pluginWith("rocket-db", "UP", capsAll(), List.of()), + desc("p1", "rocket-db")); + assertEquals(ConnectorType.GENERIC_PLUGIN, adapter.type()); + } + + @Test + void capabilities_translates_all_six_to_five() { + var adapter = new PluginConnectorAdapter( + pluginWith("x", "UP", capsAll(), List.of()), + desc("p", "x")); + ConnectorCapabilities core = adapter.capabilities(); + assertTrue(core.supportsCdc()); + assertTrue(core.supportsSnapshot()); + assertTrue(core.supportsSchemaDiscovery()); + assertTrue(core.supportsTransactions()); + assertTrue(core.supportsOffsetTracking()); + } + + @Test + void capabilities_zero_maps_to_none() { + var adapter = new PluginConnectorAdapter( + pluginWith("x", "UP", capsNone(), List.of()), + desc("p", "x")); + assertFalse(adapter.capabilities().supportsSnapshot()); + } + + @Test + void health_parses_all_states() { + for (var input : new String[]{"UP", "DOWN", "DEGRADED", "UNKNOWN", "garbage", null, " up "}) { + var adapter = new PluginConnectorAdapter( + pluginWith("x", input, capsNone(), List.of()), + desc("p", "x")); + var status = adapter.health().status(); + if (input == null) { + assertEquals(ConnectorHealth.Status.UNKNOWN, status); + } else { + var up = input.trim().toUpperCase(); + var expected = switch (up) { + case "UP" -> ConnectorHealth.Status.UP; + case "DOWN" -> ConnectorHealth.Status.DOWN; + case "DEGRADED" -> ConnectorHealth.Status.DEGRADED; + default -> ConnectorHealth.Status.UNKNOWN; + }; + assertEquals(expected, status, "input=" + input); + } + } + } + + @Test + void connect_marks_connected() { + var adapter = new PluginConnectorAdapter( + pluginWith("x", "UP", capsNone(), List.of()), + desc("p", "x")); + assertFalse(adapter.isConnected()); + adapter.connect(ctx()); + assertTrue(adapter.isConnected()); + adapter.disconnect(); + assertFalse(adapter.isConnected()); + } + + @Test + void discover_schemas_and_tables_delegate() { + var adapter = new PluginConnectorAdapter( + pluginWith("x", "UP", capsAll(), List.of()), + desc("p", "x")); + assertEquals(List.of("public"), adapter.discoverSchemas(ctx())); + assertEquals(List.of("t"), adapter.discoverTables(ctx(), "public")); + } + + @Test + void fetchColumns_wraps_plugin_columns() { + var adapter = new PluginConnectorAdapter( + pluginWith("x", "UP", capsAll(), + List.of(Map.of("name", "id", "type", "int", "primaryKey", true))), + desc("p", "x")); + var cols = adapter.fetchColumns(ctx(), "public", "t"); + assertEquals(1, cols.size()); + assertEquals("id", cols.get(0).name()); + assertTrue(cols.get(0).primaryKey()); + } + + @Test + void primary_key_metadata_built_from_pk_flag() { + var adapter = new PluginConnectorAdapter( + pluginWith("x", "UP", capsAll(), + List.of(Map.of("name", "id", "type", "int", "primaryKey", true))), + desc("p", "x")); + var pk = adapter.fetchPrimaryKey(ctx(), "public", "t"); + assertEquals(List.of("id"), pk.columnNames()); + } + + @Test + void rangeChunks_returns_single_whole_chunk() { + var adapter = new PluginConnectorAdapter( + pluginWith("x", "UP", capsAll(), List.of()), + desc("p", "x")); + var chunks = adapter.rangeChunks(ctx(), "public", "t", 8); + assertEquals(1, chunks.size()); + assertTrue(chunks.get(0).isWhole()); + } + + @Test + void estimateRows_returns_minus_one_when_no_snapshot_provider() { + var adapter = new PluginConnectorAdapter( + pluginWith("x", "UP", capsNone(), List.of()), + desc("p", "x")); + assertEquals(-1L, adapter.estimateRows(ctx(), "public", "t")); + } + + @Test + void readBatch_round_trips_cursor() { + var calls = new java.util.ArrayList(); + SnapshotProvider sp = new SnapshotProvider() { + + @Override + public long estimateRowCount(PluginContext c, String s, String t) { + return 100; + } + @Override + public PageResult readBatch(PluginContext c, String s, String t, int batchNumber, int batchSize) { + calls.add(batchNumber); + if (calls.size() == 1) { + return new PageResult(List.of(Map.of("id", 1)), "next"); + } + return new PageResult(List.of(), null); + } + }; + var combined = new CombinedConnector( + pluginWith("x", "UP", capsNone(), List.of()), sp, null); + var adapter = new PluginConnectorAdapter(combined, desc("p", "x")); + + // First call: cursor=null → batch 0 + var page1 = adapter.readBatch(ctx(), "public", "t", + new BatchInformation(0, 100, "t", null)); + assertEquals(1, page1.rows().size()); + assertNotNull(page1.nextCursor()); + // Second call: cursor from page1 → batch 1 + var page2 = adapter.readBatch(ctx(), "public", "t", + new BatchInformation(1, 100, "t", page1.nextCursor())); + assertEquals(0, page2.rows().size()); + assertNull(page2.nextCursor()); + assertEquals(List.of(0, 1), calls); + } + + @Test + void cdc_event_synthesis_populates_all_fields() { + CdcProvider cp = new CdcProvider() { + + @Override + public void startCapture(PluginContext c, Consumer consumer) { + consumer.accept(new CdcEvent("evt-1", "INSERT", "public", "t", + null, Map.of("id", 1), Map.of("lsn", "0/16"))); + } + @Override + public void stopCapture() { + } + @Override + public boolean isCapturing() { + return true; + } + @Override + public Map currentOffset() { + return Map.of("lsn", "0/16"); + } + }; + var combined = new CombinedConnector( + pluginWith("kafka", "UP", capsAll(), List.of()), null, cp); + var adapter = new PluginConnectorAdapter(combined, desc("plugin-x", "kafka")); + var got = new java.util.ArrayList(); + adapter.startCDC(ctx(), got::add); + assertEquals(1, got.size()); + var event = got.get(0); + assertEquals("evt-1", event.header().eventId()); + assertEquals("plugin:plugin-x", event.header().pipelineId()); + assertEquals(CDCOperation.INSERT, event.operation()); + assertEquals("public", event.source().schema()); + assertEquals("t", event.source().table()); + assertEquals("kafka", event.source().connectorType()); + assertEquals(Map.of("id", 1), event.payload().after()); + assertNotNull(event.metadata().capturedAt()); + assertNull(event.transaction()); + assertEquals("0/16", event.offset().offset().get("lsn")); + } + + @Test + void cdc_start_throws_when_plugin_not_cdc_provider() { + var adapter = new PluginConnectorAdapter( + pluginWith("x", "UP", capsNone(), List.of()), + desc("p", "x")); + assertThrows(UnsupportedOperationException.class, + () -> adapter.startCDC(ctx(), e -> { + })); + } + + @Test + void cdc_capture_status_reports_running_when_capturing() { + CdcProvider cp = new CdcProvider() { + + @Override + public void startCapture(PluginContext c, Consumer consumer) { + } + @Override + public void stopCapture() { + } + @Override + public boolean isCapturing() { + return true; + } + @Override + public Map currentOffset() { + return Map.of(); + } + }; + var combined = new CombinedConnector( + pluginWith("x", "UP", capsAll(), List.of()), null, cp); + var adapter = new PluginConnectorAdapter(combined, desc("p", "x")); + adapter.startCDC(ctx(), e -> { + }); + assertEquals(CaptureStatus.RUNNING, adapter.captureStatus()); + adapter.stopCDC(); + assertEquals(CaptureStatus.INACTIVE, adapter.captureStatus()); + } + + @Test + void cdc_event_operation_normalization() { + for (var op : new String[]{"insert", "UPDATE", "Delete", "create", "i", "u", "d", "weird"}) { + CDCOperation result = PluginCdcEventAdapter.parseOperation(op); + CDCOperation expected = switch (op.toLowerCase()) { + case "insert", "create", "i" -> CDCOperation.INSERT; + case "update", "u" -> CDCOperation.UPDATE; + case "delete", "d" -> CDCOperation.DELETE; + default -> CDCOperation.READ; + }; + assertEquals(expected, result, "op=" + op); + } + } + + @Test + void cursor_codec_round_trip() { + assertEquals("b:0", PluginCursorCodec.encode(0)); + assertEquals(0, PluginCursorCodec.decode(null)); + assertEquals(0, PluginCursorCodec.decode("")); + assertEquals(5, PluginCursorCodec.decode("b:5")); + // Unknown format starts at 0. + assertEquals(0, PluginCursorCodec.decode("garbage")); + } + + @Test + void type_resolver_handles_null_and_blank() { + assertEquals(ConnectorType.GENERIC_PLUGIN, ConnectorTypeResolver.resolve(null)); + assertEquals(ConnectorType.GENERIC_PLUGIN, ConnectorTypeResolver.resolve("")); + assertEquals(ConnectorType.MONGODB, ConnectorTypeResolver.resolve("mongo")); + assertEquals(ConnectorType.REDIS, ConnectorTypeResolver.resolve("REDIS")); + assertEquals(ConnectorType.GENERIC_PLUGIN, ConnectorTypeResolver.resolve("rocket-db")); + } + + /** + * Combines a PluginConnector with optional SnapshotProvider and/or CdcProvider + * for testing. + */ + private static final class CombinedConnector implements PluginConnector, SnapshotProvider, CdcProvider { + + private final PluginConnector base; + private final SnapshotProvider sp; + private final CdcProvider cp; + + CombinedConnector(PluginConnector base, SnapshotProvider sp, CdcProvider cp) { + this.base = base; + this.sp = sp; + this.cp = cp; + } + + @Override + public PluginDescriptor descriptor() { + return base.descriptor(); + } + @Override + public com.syncflow.plugin.capabilities.ConnectorCapabilities capabilities() { + return base.capabilities(); + } + @Override + public String health() { + return base.health(); + } + @Override + public Map metadata() { + return base.metadata(); + } + @Override + public List discoverSchemas(PluginContext c) { + return base.discoverSchemas(c); + } + @Override + public List discoverTables(PluginContext c, String s) { + return base.discoverTables(c, s); + } + @Override + public List> discoverColumns(PluginContext c, String s, String t) { + return base.discoverColumns(c, s, t); + } + @Override + public long estimateRowCount(PluginContext c, String s, String t) { + return sp == null ? -1 : sp.estimateRowCount(c, s, t); + } + @Override + public PageResult readBatch(PluginContext c, String s, String t, int bn, int bs) { + return sp == null ? PageResult.empty() : sp.readBatch(c, s, t, bn, bs); + } + @Override + public void startCapture(PluginContext c, Consumer consumer) { + if (cp != null) + cp.startCapture(c, consumer); + } + @Override + public void stopCapture() { + if (cp != null) + cp.stopCapture(); + } + @Override + public boolean isCapturing() { + return cp != null && cp.isCapturing(); + } + @Override + public Map currentOffset() { + return cp == null ? Map.of() : cp.currentOffset(); + } + } +} diff --git a/syncflow-api/src/test/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistryTest.java b/syncflow-api/src/test/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistryTest.java new file mode 100644 index 0000000..0e82f65 --- /dev/null +++ b/syncflow-api/src/test/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistryTest.java @@ -0,0 +1,257 @@ +package com.syncflow.api.plugin.registry; + +import com.syncflow.api.plugin.adapter.PluginConnectorAdapter; +import com.syncflow.core.model.ConnectorType; +import com.syncflow.core.registry.SpringConnectorRegistry; +import com.syncflow.core.spi.Connector; +import com.syncflow.core.spi.ConnectorContext; +import com.syncflow.plugin.descriptor.PluginDescriptor; +import com.syncflow.plugin.spi.PluginConnector; +import com.syncflow.plugin.spi.PluginContext; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class DelegatingConnectorRegistryTest { + + private DelegatingConnectorRegistry registry; + private SpringConnectorRegistry builtins; + private PluginConnectorRegistry plugins; + + @BeforeEach + void setUp() { + builtins = new SpringConnectorRegistry(List.of()); + plugins = new PluginConnectorRegistry(); + registry = new DelegatingConnectorRegistry(builtins, plugins); + } + + @Test + void register_plugin_routes_to_plugin_registry() { + var adapter = adapter("p1", "postgresql"); + registry.register(adapter); + assertTrue(plugins.isRegistered("p1")); + // The built-in POSTGRESQL slot must NOT be filled. + assertTrue(builtins.getAll().isEmpty()); + } + + @Test + void register_builtin_routes_to_spring_registry() { + var builtIn = new BuiltInConnector(ConnectorType.POSTGRESQL); + registry.register(builtIn); + assertFalse(plugins.getAll().stream().anyMatch(c -> c instanceof PluginConnectorAdapter)); + assertSame(builtIn, builtins.get(ConnectorType.POSTGRESQL).orElseThrow()); + } + + @Test + void get_postgresql_finds_plugin_advertising_postgresql() { + var adapter = adapter("p1", "postgresql"); + registry.register(adapter); + var found = registry.get(ConnectorType.POSTGRESQL); + assertTrue(found.isPresent()); + assertSame(adapter, found.get()); + } + + @Test + void get_generic_plugin_returns_first_registered_plugin() { + var a1 = adapter("p1", "rocket-db"); + var a2 = adapter("p2", "mars-db"); + registry.register(a1); + registry.register(a2); + var found = registry.get(ConnectorType.GENERIC_PLUGIN); + assertTrue(found.isPresent()); + assertTrue(found.get() == a1 || found.get() == a2); + } + + @Test + void unregister_plugin_by_id_removes_and_disconnects() { + var a1 = adapter("p1", "x"); + registry.register(a1); + assertTrue(registry.unregisterPlugin("p1")); + assertFalse(plugins.isRegistered("p1")); + // After unregister, the adapter is no longer registered; the + // disconnect side effect on the inner plugin is observable + // through a flag on the test plugin. + } + + @Test + void unregister_generic_plugin_type_clears_all_plugins() { + registry.register(adapter("p1", "x")); + registry.register(adapter("p2", "y")); + assertEquals(2, plugins.getAll().size()); + registry.unregister(ConnectorType.GENERIC_PLUGIN); + assertEquals(0, plugins.getAll().size()); + } + + @Test + void unregister_builtin_type_does_not_touch_plugins() { + var a1 = adapter("p1", "x"); + registry.register(a1); + registry.unregister(ConnectorType.POSTGRESQL); + // The plugin is still there. + assertTrue(plugins.isRegistered("p1")); + } + + @Test + void isRegistered_generic_plugin_reflects_plugin_count() { + assertFalse(registry.isRegistered(ConnectorType.GENERIC_PLUGIN)); + registry.register(adapter("p1", "x")); + assertTrue(registry.isRegistered(ConnectorType.GENERIC_PLUGIN)); + } + + @Test + void getAll_combines_builtins_and_plugins() { + registry.register(new BuiltInConnector(ConnectorType.MYSQL)); + registry.register(adapter("p1", "x")); + assertEquals(2, registry.getAll().size()); + } + + @Test + void plugin_disconnect_called_on_uninstall() { + var plugin = new PluginConnector() { + + @Override + public PluginDescriptor descriptor() { + return descOf("p1", "x"); + } + @Override + public com.syncflow.plugin.capabilities.ConnectorCapabilities capabilities() { + return noCaps(); + } + @Override + public String health() { + return "UP"; + } + @Override + public Map metadata() { + return Map.of(); + } + @Override + public List discoverSchemas(PluginContext c) { + return List.of(); + } + @Override + public List discoverTables(PluginContext c, String s) { + return List.of(); + } + @Override + public List> discoverColumns(PluginContext c, String s, String t) { + return List.of(); + } + }; + var adapter = new PluginConnectorAdapter(plugin, descOf("p1", "x")); + adapter.connect(new ConnectorContext( + new com.syncflow.core.model.ConnectionConfiguration( + ConnectorType.GENERIC_PLUGIN, "h", 0, "d", null, null, Map.of()), + Map.of())); + assertTrue(adapter.isConnected()); + registry.register(adapter); + // Plugin SPI has no disconnect; the adapter's disconnect() flips its + // own internal connected flag. We just verify the adapter is + // removed from the registry and that calling disconnect works. + registry.unregisterPlugin("p1"); + adapter.disconnect(); + assertFalse(adapter.isConnected()); + } + + private static PluginDescriptor descOf(String id, String type) { + return new PluginDescriptor(id, "t", "v", "1.0", "d", type, + List.of(), "1", "1", List.of(), "MIT", null, null); + } + + private static com.syncflow.plugin.capabilities.ConnectorCapabilities noCaps() { + return new com.syncflow.plugin.capabilities.ConnectorCapabilities( + false, false, false, false, false, false); + } + + private static PluginConnectorAdapter adapter(String id, String type) { + var plugin = new PluginConnector() { + + @Override + public PluginDescriptor descriptor() { + return descOf(id, type); + } + @Override + public com.syncflow.plugin.capabilities.ConnectorCapabilities capabilities() { + return noCaps(); + } + @Override + public String health() { + return "UP"; + } + @Override + public Map metadata() { + return Map.of(); + } + @Override + public List discoverSchemas(PluginContext c) { + return List.of(); + } + @Override + public List discoverTables(PluginContext c, String s) { + return List.of(); + } + @Override + public List> discoverColumns(PluginContext c, String s, String t) { + return List.of(); + } + }; + return new PluginConnectorAdapter(plugin, descOf(id, type)); + } + + private static final class BuiltInConnector implements Connector { + + private final ConnectorType type; + private final AtomicInteger connects = new AtomicInteger(); + BuiltInConnector(ConnectorType type) { + this.type = type; + } + @Override + public ConnectorType type() { + return type; + } + @Override + public com.syncflow.core.spi.ConnectorCapabilities capabilities() { + return com.syncflow.core.spi.ConnectorCapabilities.none(); + } + @Override + public void connect(ConnectorContext c) { + connects.incrementAndGet(); + } + @Override + public void disconnect() { + } + @Override + public boolean isConnected() { + return connects.get() > 0; + } + @Override + public com.syncflow.core.spi.ConnectorValidationResult validate(ConnectorContext c) { + return com.syncflow.core.spi.ConnectorValidationResult.ok(); + } + @Override + public List discoverSchemas(ConnectorContext c) { + return List.of(); + } + @Override + public List discoverTables(ConnectorContext c, String s) { + return List.of(); + } + @Override + public com.syncflow.core.spi.ConnectorHealth health() { + return new com.syncflow.core.spi.ConnectorHealth( + com.syncflow.core.spi.ConnectorHealth.Status.UP, "", java.time.Instant.now(), 0); + } + @Override + public Map metadata() { + return Map.of(); + } + } +} diff --git a/syncflow-core/src/main/java/com/syncflow/core/model/ConnectorType.java b/syncflow-core/src/main/java/com/syncflow/core/model/ConnectorType.java index b6f9625..6ba3141 100644 --- a/syncflow-core/src/main/java/com/syncflow/core/model/ConnectorType.java +++ b/syncflow-core/src/main/java/com/syncflow/core/model/ConnectorType.java @@ -1,5 +1,5 @@ package com.syncflow.core.model; public enum ConnectorType { - POSTGRESQL, MYSQL, MONGODB, KAFKA, SQLSERVER, ORACLE, ELASTICSEARCH, REDIS, GENERIC_JDBC + POSTGRESQL, MYSQL, MONGODB, KAFKA, SQLSERVER, ORACLE, ELASTICSEARCH, REDIS, GENERIC_JDBC, GENERIC_PLUGIN } diff --git a/syncflow-plugin-api/src/main/java/com/syncflow/plugin/spi/DestinationWriterProvider.java b/syncflow-plugin-api/src/main/java/com/syncflow/plugin/spi/DestinationWriterProvider.java index b4a06e1..2841984 100644 --- a/syncflow-plugin-api/src/main/java/com/syncflow/plugin/spi/DestinationWriterProvider.java +++ b/syncflow-plugin-api/src/main/java/com/syncflow/plugin/spi/DestinationWriterProvider.java @@ -9,6 +9,26 @@ public interface DestinationWriterProvider extends AutoCloseable { void write(String table, List> rows); + /** + * Delete rows. Default is a no-op; plugins with native DELETE support + * should override. Returning silently (rather than throwing) is a + * deliberate backward-compat choice: the adapter degrades gracefully + * for plugins that pre-date this method. + */ + default void delete(String table, List> rows) { + // no-op default — see javadoc + } + + /** + * Upsert rows keyed by {@code keyColumns}. Default falls back to + * {@link #write(String, List)}; the keys are ignored. Plugins with + * native upsert semantics (ON CONFLICT, ON DUPLICATE KEY) override. + */ + default void upsert(String table, List> rows, + List keyColumns) { + write(table, rows); + } + void flush(); void commit(); From 453ee74fb4f8998792d4190d512d414468c37a0d Mon Sep 17 00:00:00 2001 From: lekhrocks Date: Sun, 6 Sep 2026 14:56:06 +0530 Subject: [PATCH 2/2] =?UTF-8?q?fix(F18):=20review=20fixes=20=E2=80=94=20re?= =?UTF-8?q?source=20leak,=20cursor=20stall,=20state=20split,=20race,=20val?= =?UTF-8?q?idation?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - DelegatingConnectorRegistry.unregisterPlugin(): disconnect before unregister (was: remove from map first, disconnect never fired) - readBatch cursor check: strict 'b:' prefix (was: single-char 'b', matched any plugin cursor starting with 'b') - isCdcActive(): check adapter flag first so pauseCDC makes it return false consistently with captureStatus() - startCDC(): set cdcActive=true before startCapture() to close the race window where concurrent stopCDC() skips stopCapture() - validate(): call delegate.health() and check status instead of always returning ok() (connect() was a no-op) --- .../adapter/PluginConnectorAdapter.java | 37 ++++++++++++------- .../registry/DelegatingConnectorRegistry.java | 9 +++-- 2 files changed, 29 insertions(+), 17 deletions(-) diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapter.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapter.java index 0b83514..bcbd370 100644 --- a/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapter.java +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapter.java @@ -123,16 +123,18 @@ public boolean isConnected() { @Override public ConnectorValidationResult validate(ConnectorContext context) { - // Plugin SPI has no dedicated validation. Try a connect/disconnect - // round-trip and report the outcome. + // Test connectivity by calling the plugin's health check. try { - if (!connected) { - connect(context); + connect(context); + var health = delegate.health(); + if (health != null && health.trim().equalsIgnoreCase("DOWN")) { + return ConnectorValidationResult.failed( + List.of("Plugin reports DOWN")); } return ConnectorValidationResult.ok(); } catch (Exception e) { - return ConnectorValidationResult.failed(List.of( - "Plugin validation failed: " + e.getMessage())); + return ConnectorValidationResult.failed( + List.of("Plugin validation failed: " + e.getMessage())); } } @@ -234,8 +236,8 @@ public SnapshotCapableConnector.Page readBatch(ConnectorContext context, String String encoded; if (next == null) { encoded = null; - } else if (next.startsWith(PluginCursorCodec.encode(0).substring(0, 1))) { - // Already encoded by us (or a plugin using the same scheme). + } else if (next.startsWith("b:")) { + // Already encoded by us (strict prefix match). encoded = next; } else { // Treat the plugin's string cursor as opaque — bump the batch. @@ -270,10 +272,12 @@ public void startCDC(ConnectorContext context, Consumer eventConsumer) if (cdcActive) { return; } + // Set the flag BEFORE startCapture so that a concurrent stopCDC() + // sees cdcActive=true and calls stopCapture() instead of skipping it. + cdcActive = true; var pluginCtx = PluginContextAdapter.from(context); cp.startCapture(pluginCtx, e -> eventConsumer.accept(PluginCdcEventAdapter.toCore(e, descriptor))); - cdcActive = true; } @Override @@ -289,10 +293,11 @@ public void stopCDC() { @Override public void pauseCDC() { - // Plugin SPI has no pause; track state for isCdcActive() reporting. - if (cdcActive) { - cdcActive = false; - } + // Plugin SPI has no pause; track state for isCdcActive() and + // captureStatus() so they agree. The plugin continues capturing + // but events are not forwarded — the consumer is stopped by + // the caller before calling pauseCDC(). + cdcActive = false; } @Override @@ -302,6 +307,12 @@ public void resumeCDC() { @Override public boolean isCdcActive() { + // Check the adapter flag first — it tracks pause/stop state. + // Fall back to the plugin's own state for the initial case + // where startCDC was called before this adapter was created. + if (!cdcActive) { + return false; + } if (delegate instanceof CdcProvider cp) { return cp.isCapturing(); } diff --git a/syncflow-api/src/main/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistry.java b/syncflow-api/src/main/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistry.java index 477b99c..c5b1810 100644 --- a/syncflow-api/src/main/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistry.java +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistry.java @@ -108,12 +108,13 @@ public boolean isRegistered(ConnectorType type) { /** * Unregister a single plugin by id — used by PluginManager on - * disable/uninstall. + * disable/uninstall. Disconnects the adapter first to release + * resources (DB handles, CDC threads) before removing from the map. */ public boolean unregisterPlugin(String pluginId) { - var removed = plugins.unregister(pluginId); - plugins.get(pluginId).ifPresent(Connector::disconnect); - return removed; + var connector = plugins.get(pluginId); + connector.ifPresent(Connector::disconnect); + return plugins.unregister(pluginId); } private static String pluginIdOf(Connector c) {