Conversation
Adds a JDBC source connector that polls a configured SQL query against any JDBC-compliant database (PostgreSQL, MySQL, Oracle, SQL Server, H2) through an embedded JVM and produces each row as a JSON message. Supports bulk and incremental (offset-tracked) modes with persisted offset state. Also adds the shared connector integration-test harness and a small TCP listener change so that harness can read the server's written runtime config when bound to a fixed port. The matching JDBC sink connector follows in a separate PR.
Aligns the JDBC source crate with the other source connectors: crate version/type and publish flag, and poll_interval parsing via the workspace humantime crate rather than a direct humantime-serde pin. Removes a redundant per-column JNI getObject probe (using wasNull after the typed getter instead) and caps the per-poll liveness check at five seconds so a dead connection cannot stall a shared worker. Drops the unused partition_id field from the example configs and docs.
# Conflicts: # Cargo.lock
Derives the incremental cursor from the last ordered row rather than a Rust-side max, so the cursor always matches the database's own ORDER BY and degrades to re-reads rather than skips if the ordering is imperfect. Fails loudly when the tracking column is absent from the result set, and rejects an empty query, a blank tracking_column, or a blank initial_offset at open(). Routes date/time columns through the null-safe string path so a NULL value no longer fails the poll, and preserves companion WHERE conditions (such as an operator-added IS NOT NULL) when stripping the offset predicate at cold start.
…C examples The incremental examples used timestamp tracking columns (updated_at, OrderDate) without noting that a non-unique tracking value can skip rows when a same-value group is split across a batch boundary. Add that caveat to the affected examples and point to the unique/strictly-increasing requirement.
Incremental-mode validation and row reading disagreed about what a
tracking-column name refers to, so a config that reads rows correctly
could be rejected at open(), and a config that skips rows could be
accepted.
Read-time tracking_column_matches compares against both the raw driver
label and the snake_cased output key, but the ORDER BY check compared a
literal lowercased string. With snake_case_columns on, a snake_cased
tracking_column against a CamelCase ORDER BY column (order_date vs
OrderDate) failed validation despite matching at runtime. The check now
normalizes the ORDER BY key the same way and routes both paths through
tracking_column_matches, and it strips a surrounding identifier quote so
a case-preserving column such as PostgreSQL's ORDER BY "OrderDate"
validates against its unquoted driver label.
Validation also derived the outer ordering with rfind("order by"), which
matched a window function's internal ORDER BY (ROW_NUMBER() OVER (ORDER
BY id) with no outer clause). That orders values within the frame, not
the emitted ResultSet, so setMaxRows still returns an arbitrary subset
and the advancing cursor skips rows, reopening the skip-row bug at
validate time. Ordering detection now scans at parenthesis depth zero,
ignoring ORDER BY inside a subquery, CTE, or OVER (...), and skips
single-quoted string literals so a parenthesis in a literal cannot
perturb the depth. A query whose only ordering sits inside such a
construct is rejected.
Update the README ORDER BY caveat to match and drop a stray line-ending
backslash in the bulk example query.
open() booted the JVM and ran DriverManager.getConnection synchronously and unbounded, so a single unreachable or slow-DNS database hung the whole connectors runtime at startup (sources open sequentially, and the FFI drives open() via block_on), taking every other source, sink, and the HTTP control API with it. Wrap the JVM init and connection in block_in_place like poll()/close(), and add login_timeout_ms (DriverManager.setLoginTimeout) so a stuck connect fails instead of blocking forever. BIGINT was emitted as a JSON number, silently rounding values above 2^53 in consumers that parse JSON numbers as f64. Emit it as a string, the same as NUMERIC/DECIMAL already do. The incremental cursor comment claimed it "degrades to re-reads rather than skips", which is false when a run of rows sharing one tracking value is split across a setMaxRows batch boundary: the remainder is skipped. Correct the comment and warn at validate time that the tracking column must be unique or batch_size must exceed the largest tie group. Hoist the jni dependency to the workspace so the pin is shared with the upcoming JDBC sink rather than duplicated. Document that SecretString redaction is not end-to-end: the runtime's generic configs/plugin control endpoint and trace-level config logging emit the raw plugin_config for every connector.
# Conflicts: # Cargo.lock # core/integration/tests/connectors/mod.rs
A proactive claims-vs-code sweep found two README statements the code does not back. `table_name` is a reserved field hardcoded to null, never derived from the query, yet the metadata was advertised as including the table name; document that it is always null and that operation_type is always SELECT. "Prevents duplicate reads" contradicted the connector's own at-least-once guarantee; reword to "avoids re-reading rows below the tracked offset (at-least-once, not exactly-once)".
Harden outer_order_by() to skip SQL line (--) and block (/* */) comments, so a commented-out ORDER BY can no longer satisfy the incremental ordering requirement and let an unordered query pass validation and skip rows at runtime. Complete JNI exception clearing in throwable_string_method(): a Java exception raised by getMessage()/getSQLState() or the string conversion is now cleared before returning, so the next JNI call on the thread is not aborted for running with an exception pending. Drop the Serialize derive (and the serialize_secret helpers) from JdbcSourceConfig: the runtime only deserializes it, and serializing would risk writing the jdbc_url/password SecretStrings in plaintext. This also removes the now-unused iggy_common dependency. Redacted logging still goes through the manual Debug impl. Read column names with ResultSetMetaData.getColumnLabel() instead of getColumnName() so a SELECT expr AS alias yields the alias the caller configured rather than the base-table column (or an empty string). Remove the unused connectors_api_address() test helper that failed integration clippy, and make the JDBC integration tests panic on a Postgres/testcontainers setup failure instead of returning early and reporting a false pass.
Correct "unparseable" -> "unparsable" in poll_interval doc comments and code comments (flagged by the typos check), and reflow jdbc_oracle.toml to taplo's canonical format (a stray blank line).
The JDBC Postgres tests downloaded the driver JAR at runtime and cached
it without validating it. When the download returned a truncated body or
an error page (seen in CI), the file was still cached and put on the JVM
classpath; the JVM booted fine but Class.forName("org.postgresql.Driver")
then failed, so the connector never initialized and every test expecting
messages saw zero, while the bad jar was reused by every later test.
Download from Maven Central, verify the bytes are a real JAR (ZIP magic
plus a minimum size) before persisting, and write via a temp file and
atomic rename so a partial or corrupt download can never be cached. A
previously cached invalid jar now self-heals, and a genuinely bad
download fails setup with a clear message instead of a later opaque
ClassNotFoundException.
# Conflicts: # Cargo.toml # core/server/src/tcp/tcp_listener.rs
The incremental cursor advanced as soon as rows were fetched, so a batch whose send later failed was skipped permanently: the next poll queried past it, and its further-advanced offset then overwrote the checkpoint, removing any chance of replay. The connector could not fix this before, because the runtime handed a source its state only at open() with no per-batch delivery signal. Upstream now supplies that signal. Stage the advanced cursor in poll() and resolve it in on_batch_result: Ack (batch sent and checkpoint persisted) commits it, Nack discards it so the same range is re-polled. That makes at-least-once hold for an ordinary in-process send failure, not just across restarts, and follows the staging pattern documented for sources and implemented by random_source. An empty poll now checkpoints nothing rather than rewriting an unchanged cursor, which on the HTTP state backend costs a round-trip per interval. Also repair the shared connectors test harness for the create_topic and send_messages signature changes that came with master, and rename the JDBC integration test file to the sibling convention.
Only a status code crosses the FFI boundary when a plugin is opened, and both containers discarded the error behind `result.is_ok()`. The runtime could then report no more than "plugin initialization failed", so an operator saw a connector silently skipped at startup with nothing to say why, and the plugin's own diagnosis was thrown away. A JDBC source whose driver JAR could not be read cost a full CI investigation for exactly this reason: the connector had already produced "Failed to load driver class 'org.postgresql.Driver'" and it never reached a log. Log the cause before collapsing it to a code, for sinks and sources alike. The signature is unchanged, so pre-built plugins keep working.
The driver-JAR integrity check accepted any file of at least 500 KB starting with the ZIP local-file-header magic. A download truncated past that size still satisfies both conditions while its central directory is gone, so it is not readable as an archive at all. The check's own comment claimed such a download would fail it. That file was then cached and reused by every later JDBC test in the job, and the JVM could not load the driver class from it. `Class.forName` threw, the connector's `open()` failed, and the runtime skipped the source, so tests failed having received no messages rather than pointing at the JAR. One bad download made this deterministic across all four nextest retries and took out every JDBC test sharing the runner. Verify instead that the bytes open as an archive containing `org/postgresql/Driver.class`, which is the property the JVM needs; a cached file that fails is already deleted and re-downloaded. Also stop comparing the recovery assertion against the raw message list. Bulk mode re-runs its query every poll interval, so it re-delivers the same rows indefinitely and two batches arriving in one client poll turned a healthy run into an id mismatch. Match on the distinct ids seen, and report the received count so "nothing was delivered" is distinguishable from "delivered without the expected data.id".
The shared connectors test client polled in batches of 10, and the JDBC collect-until-N loop gave up after 30 attempts. Those multiply into a ceiling: the loop could never observe more than 300 messages however many the source delivered. large_result_set_streams_without_crashing waited for 150 of that 300 and locally spent ~16 attempts draining 10 at a time, so a slower runner exhausted the budget before the count was reached and the test failed on every retry rather than flaking. The five tests that pass are all bulk mode, which re-runs its query each interval and so retries for free; the two that failed were the only ones whose expectations could not absorb a slow runner. Let the caller choose the batch count, request more per poll than any test waits for, and bound the wait with a wall-clock deadline so runner speed costs time instead of correctness. Assert the incremental cursor on the distinct ids seen rather than on a fixed-length prefix of the received messages. The claim under test is unchanged, that the offset advanced past 3, but a re-delivered batch no longer fails it, which matters because the connector's delivery contract is at-least-once. Both panics now report the ids seen and the raw message count, so a future failure distinguishes a stalled source from a replay.
# Conflicts: # .config/nextest.toml # .github/workflows/_build_rust_artifacts.yml # Cargo.lock # Cargo.toml
|
Thanks for the PR. It is labeled Slash commands (own line, regular comment) move it around the queue:
See CONTRIBUTING.md for details. |
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## master #4272 +/- ##
============================================
+ Coverage 87.56% 87.70% +0.14%
- Complexity 1575 1576 +1
============================================
Files 1284 1285 +1
Lines 225457 227541 +2084
Branches 188821 190905 +2084
============================================
+ Hits 197413 199570 +2157
+ Misses 23309 23154 -155
- Partials 4735 4817 +82
🚀 New features to boost your workflow:
|
|
you dont need to constantly rebase. maintainers will do that :) |
# Conflicts: # .config/nextest.toml
|
was I not clear enough? 🥲 please, dont rebase, your PR is merge-able without conflicts, it's good for review now, expect review this week. if it won't rebase cleanly we'll ping you. |
My bad, there was a conflict earlier and thus had to merge and resolve conflicts. |
|
/skill team-review-slim |
There was a problem hiding this comment.
Summary: The JDBC source connector's staged-cursor protocol, JNI error clearing and SQL validation hold up under tracing; fifteen findings remain, none critical. The strongest are a zero poll_interval that busy-loops the database, an incremental query without a cursor predicate that re-reads the same rows forever, a tie group that can skip rows, and two credential leaks around the logs - a URL sanitizer that stops at the first @, and the new SDK open-error log that can print the full jdbc_url; claims this checkout cannot settle, such as jni-rs method lookups and non-MySQL backslash escaping, were dropped.
Counts: critical 0, warning 8, nit 3, simplification 4
Findings without an anchor on a changed line:
core/connectors/sources/jdbc_source/src/lib.rs:1536warning:poll_interval = "0s"parses, so the sleep never runs andpoll()issues back-to-back JDBC queries against the database. Reject a zero interval invalidate_config, or clamp it to a 1ms floor.core/connectors/sources/jdbc_source/src/lib.rs:968warning:read_rowstakes the last row's tracking value as the cursor, so a tie group split bysetMaxRowsloses its remaining rows on the next strict>poll. Probe one row pastbatch_sizein incremental mode and fail closed when the extra row repeats the last in-batch value.core/connectors/sources/jdbc_source/src/lib.rs:1450warning: Incremental mode accepts a query without{last_offset}, so every poll re-reads the same firstbatch_sizerows while the cursor never enters the query. Require the placeholder in the query atopen(), or error the poll when the built query holds no cursor comparison.core/connectors/sources/jdbc_source/src/lib.rs:1027warning:read_single_rowtakes the cursor from the first result column whose label matchestracking_column, whilequery_orders_by_tracking_columnstrips theORDER BYqualifier, so a join that returns twoupdated_atlabels can set the cursor from the column the query does not order by. Reject a result set with more than one matching label, or take the cursor from the ordered column.core/connectors/sources/jdbc_source/src/lib.rs:143warning:sanitize_jdbc_urlstops at the first@, so ajdbc_urlwhose password contains one, such asuser:p@ss@host, still logsss@hostat lines 444 and 1503. Mask the credential up to the last@before the host.core/connectors/sources/jdbc_source/src/lib.rs:1006warning:read_single_rowrebuilds each column's output key and repeats the collision warning for every row, so one collapsed key pair logs once per row per poll. Compute the key and warn once per column, inread_column_metadata.core/connectors/sources/jdbc_source/src/lib.rs:1702warning:get_or_create_jvmreturns the live JVM without comparingdriver_jar_pathorjvm_options, so a second JDBC source silently uses the first connector's classpath. Fail atopen()when the requested classpath differs, becausevalidate_configalready accepted the unused JAR path.core/connectors/sources/jdbc_source/src/lib.rs:367simplification:stateis anArc<Mutex<State>>that nothing clones, because the source is its only owner. Simpler: use a plainMutex<State>, which removes one heap allocation and one level of indirection.core/connectors/sources/jdbc_source/src/lib.rs:334simplification:State.last_poll_timeis written into every checkpoint and never read. Simpler: delete the field, since the connector is new and no older checkpoint depends on its position.core/connectors/sources/jdbc_source/src/lib.rs:2044simplification:java::sql::TypesdefinesCHAR,VARCHAR,LONGVARCHARandNULL, and nothing reads them, so two#[allow(dead_code)]attributes hide the gap. Simpler: delete the four constants and both attributes.
This review was generated by Claude Code 2.1.284 on deepseek-flash[1m]. Review the output before you act on it.
| match result { | ||
| Ok(()) => 0, | ||
| Err(error) => { | ||
| error!("Failed to open source connector with ID: {id}. {error}"); |
There was a problem hiding this comment.
warning: The new log line writes the connector's open error verbatim, and a driver rejection such as JDBC's No suitable driver found for <url> echoes the full jdbc_url, credentials included. Sanitize the URL out of that error before the connector returns it. Also at core/connectors/sdk/src/sink.rs:101.
| ```sql | ||
| -- Configuration | ||
| tracking_column = "id" | ||
| query = "SELECT * FROM users WHERE id > {last_offset} ORDER BY id" |
There was a problem hiding this comment.
nit: The example config has no initial_offset, and strip_offset_predicate removes only the literal WHERE {tracking_column} > {last_offset} form, so this query fails open() on the first poll. Add initial_offset = "0" to the example.
| password = "" | ||
|
|
||
| # Simple query for testing | ||
| query = "SELECT * FROM users WHERE id > {last_offset} ORDER BY id" |
There was a problem hiding this comment.
nit: The in-memory H2 database starts empty, and nothing creates the users table the query selects, so this example errors on every poll. Add an INIT clause that creates and seeds the table, as test_jdbc_h2.toml does.
| uuid = { workspace = true, features = ["v4"] } | ||
|
|
||
| [dev-dependencies] | ||
| toml = { workspace = true } |
There was a problem hiding this comment.
nit: The crate omits [lints] workspace = true, so the workspace warnings = "deny" gate does not cover it. Add the section, as the five sibling source crates do.
| impl ConnectorsIggyClient { | ||
| /// Send messages to the configured stream/topic (used by sink connector tests). | ||
| #[allow(dead_code)] | ||
| async fn send_messages( |
There was a problem hiding this comment.
simplification: ConnectorsIggyClient::send_messages is never called, and its doc comment claims a sink test uses it, so #[allow(dead_code)] hides 15 lines of dead test scaffolding. Simpler: delete the method.
Which issue does this PR address?
Relates to #2500.
This supersedes #3588, which was closed by the stale bot. It includes the review follow-ups from that PR and merges the latest
master.Rationale
Iggy currently has database-specific connectors. A generic JDBC source lets Iggy read from PostgreSQL, MySQL, Oracle, SQL Server, H2, and other JDBC-compliant databases through their existing JDBC drivers.
What changed?
java.sqlAPI.{last_offset}substitution.master, including persisted topic durability required by source connectors.How to review
core/connectors/sources/jdbc_source/src/lib.rscore/connectors/sources/jdbc_source/README.mdandconfig.tomlcore/integration/tests/connectors/jdbc/core/integration/tests/connectors/mod.rsVerification
cargo fmt --allcargo sort --no-format --workspacecargo check -p iggy_connector_jdbc_sourcecargo clippy -p iggy_connector_jdbc_source --all-features --all-targets -- -D warningscargo test -p iggy_connector_jdbc_source --all-features(65 passed)cargo clippy -p iggy_connector_sdk --all-targets --all-features -- -D warningscargo test -p iggy_connector_sdk --all-features(220 passed)cargo clippy -p integration --all-targets --all-features -- -D warningsIGGY_TEST_CLUSTER_NODES=1 cargo test -p integration -- connectors::jdbc::(7 passed)The local three-node integration harness did not reach VSR mesh readiness on this machine, before any JDBC test ran. The JDBC suite passes in the harness-supported single-node mode above.
Limitations
AI Usage
Claude Code was used substantially for implementation, tests, documentation, and review hardening. The changes were verified with the commands above and can be explained line by line.