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..bcbd370 --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/adapter/PluginConnectorAdapter.java @@ -0,0 +1,392 @@ +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) { + // Test connectivity by calling the plugin's health check. + try { + 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())); + } + } + + @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("b:")) { + // Already encoded by us (strict prefix match). + 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; + } + // 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))); + } + + @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() 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 + public void resumeCDC() { + // No-op: paused CDC needs startCDC to resume on the plugin side. + } + + @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(); + } + 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..c5b1810 --- /dev/null +++ b/syncflow-api/src/main/java/com/syncflow/api/plugin/registry/DelegatingConnectorRegistry.java @@ -0,0 +1,123 @@ +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. Disconnects the adapter first to release + * resources (DB handles, CDC threads) before removing from the map. + */ + public boolean unregisterPlugin(String pluginId) { + var connector = plugins.get(pluginId); + connector.ifPresent(Connector::disconnect); + return plugins.unregister(pluginId); + } + + 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();