fix(connectors): Harden Elasticsearch sink test readiness#3729
fix(connectors): Harden Elasticsearch sink test readiness#3729ryerraguntla wants to merge 5 commits into
Conversation
Add a timeout to the Elasticsearch sink transport, make connectors-runtime readiness check `/health`, and strengthen the Elasticsearch test fixture startup path with retry, cluster-health polling, and stale-index cleanup. Also refresh the sink index before counting documents so test assertions see newly indexed data sooner.
Use short-timeout probes without retries for cluster health and stale index sweep, and cap docker rm so a wedged daemon cannot stall setup.
|
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✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## master #3729 +/- ##
=============================================
- Coverage 74.34% 40.21% -34.13%
Complexity 950 950
=============================================
Files 1303 1301 -2
Lines 148712 129057 -19655
Branches 124225 104570 -19655
=============================================
- Hits 110562 51905 -58657
- Misses 34673 74406 +39733
+ Partials 3477 2746 -731
🚀 New features to boost your workflow:
|
|
/request-review @user-or-team |
There was a problem hiding this comment.
start() force-removes the shared container after any try_start() error, not only after confirming it is wedged. Transient health, port-inspection, or Docker errors can therefore remove an Elasticsearch container still used by another process. The fixed name limits which container is removed but does not provide ownership or concurrency safety. Could this recovery path verify the container is unhealthy and coordinate removal with a cross-process lock before running docker rm -f?
There was a problem hiding this comment.
Let me add the logic to remove when required.
Only docker rm the shared reuse container after inspect shows a removable state, and serialize recovery with a cross-process lock so transient start flakes cannot yank a healthy instance from a peer nextest worker.
|
/ready |
hubcio
left a comment
There was a problem hiding this comment.
not in this diff, same class of problem:
- elasticsearch source
create_client()(core/connectors/sources/elasticsearch_source/src/lib.rs:304) still builds its transport with no timeout - sameopen()-hang failure mode this PR fixes for the sink, so source tests keep the #3728 flake class. worth a follow-up. mcp.rswait_ready()still uses the pattern this PR fixed for connectors: GET/, anyOk(_)counts as ready, no process-crash fast-fail. consistency follow-up.
| .unwrap_or(DEFAULT_TIMEOUT_SECONDS) | ||
| .max(1); | ||
| let mut transport_builder = | ||
| TransportBuilder::new(conn_pool).timeout(Duration::from_secs(timeout_seconds)); |
There was a problem hiding this comment.
this timeout is client-global in elasticsearch-rs, so it bounds bulk requests in consume() too, not just open(). that changes the failure mode: before, a bulk on degraded ES stalled until it eventually succeeded; now anything over timeout_seconds returns Err, which the runtime discards (#2927) with the offset already committed at poll (#2928) - the batch is silently dropped. mechanism is pre-existing and tracked there, but this PR is what arms it for slow-ES bulk. worth an upgrade note in the PR description: raise timeout_seconds for slow bulk workloads, and note the tradeoff until #2927/#2928 land.
separately, the comment slightly oversells the timeout as the flake fix - with the 30s default the harness readiness budget (~20s) expires first, so what fixes #3728 is the fixture readiness gate; this is the infinite-hang backstop. maybe say that instead.
| .config | ||
| .timeout_seconds | ||
| .unwrap_or(DEFAULT_TIMEOUT_SECONDS) | ||
| .max(1); |
There was a problem hiding this comment.
.max(1) silently turns timeout_seconds: 0 into 1s. defensible - Duration::ZERO would instant-fail every request - but undocumented, and clickhouse_sink passes 0 through unclamped, so the two siblings now disagree on the same config field. document the clamp here; clickhouse should probably get the same guard.
| // Prefer /health over `/` so readiness means the connectors API is up, | ||
| // not some other listener that happened to bind the reserved port. | ||
| let health_url = format!("{}/health", self.http_url()); | ||
| let client = reqwest::Client::new(); |
There was a problem hiding this comment.
reqwest::Client::new() has no request timeout, so a black-holed GET (connection accepted, no response) hangs send() forever - the loop never gets back to the pid check above and the retry budget never fires. low risk on localhost (dead process means ECONNREFUSED, which fast-fails), but a 1-2s timeout like the fixture probe client makes both guards effective.
| match client.get(&health_url).send().await { | ||
| Ok(response) if response.status().is_success() => { | ||
| let body = response.text().await.unwrap_or_default(); | ||
| if body.contains("\"timed_out\":true") { |
There was a problem hiding this comment.
?timeout=1s without a wait_for_* param makes /_cluster/health return immediately with timed_out: false, so this check never fires - and any 200 counts as ready even with cluster status red. ?wait_for_status=yellow&timeout=1s makes both the timeout param and this check do real work (and then parsing the bool beats string-matching the body).
| } | ||
|
|
||
| fn acquire_recovery_lock() -> Result<RecoveryLock, String> { | ||
| let path = std::env::temp_dir().join(RECOVERY_LOCK_FILE_NAME); |
There was a problem hiding this comment.
cross-process serialization only holds if every nextest process resolves the same temp_dir() - with a per-process TMPDIR (some sandbox/CI setups) each process locks its own file and recovery races silently. probably fine to accept, but worth a comment stating the assumption.
| .truncate(false) | ||
| .open(&path) | ||
| .map_err(|error| format!("open {}: {error}", path.display()))?; | ||
| file.lock() |
There was a problem hiding this comment.
File::lock() blocks with no acquire timeout - worst case a peer holds this through its whole recovery (two 120s container starts plus readiness waits), and since the test runtime is current-thread, the block parks the entire runtime for that long. bounded and failure-path-only, so nice-to-have: try_lock() with a deadline.
| ) || health == "unhealthy" | ||
| } | ||
|
|
||
| fn docker_command_stdout(args: &[&str], timeout_secs: u64) -> Option<String> { |
There was a problem hiding this comment.
docker_command_stdout and force_remove_container duplicate the spawn / try_wait / deadline / kill loop, and both block the async caller (thread::sleep on a current-thread runtime). tokio::process::Command + tokio::time::timeout(dur, cmd.output()) + kill_on_drop(true) collapses both to a few lines with the same 5s/15s kill semantics - delta/fixture.rs:254 is the in-repo precedent and tokio full is already enabled here. it also removes a latent footgun: stdout is only drained after exit, so docker output over the pipe buffer (~64KB) would deadlock until the timeout kill. fine today with one-line inspect -f output, but the helper reads as generic.
|
|
||
| for _ in 0..POLL_ATTEMPTS { | ||
| // Refresh so near-real-time search/count sees recently indexed docs. | ||
| if let Err(error) = self.refresh_index().await { |
There was a problem hiding this comment.
refresh_index() goes through create_http_client() (30s timeout + 3 retries) inside a 50ms poll loop - on the failure path that is up to 100 refreshes, each with retry backoff. happy path exits in a few iterations so this only inflates time-to-failure; the short-timeout probe-client pattern from container.rs would fit here too.
Summary
iggy-connectorsnever becomes healthy because sink
open()can hang on Elasticsearch HTTPwith no client timeout, while the connectors HTTP API only binds after
plugin open.
127.0.0.1, alwaysre-check
/_cluster/healthon reuse attach, self-heal a wedgediggy-test-elasticsearchcontainer only, and sweep stale harness indicesolder than 30 minutes.
/healthand fail fast if thechild process dies during retries; refresh the index before document-count
polls in the sink fixture.
Closes #3728
Test plan
cargo fmt --allcargo sort --no-format --workspacecargo clippy -p iggy_connector_elasticsearch_sink -p integration --all-features --all-targets -- -D warningscargo test -p integration -- connectors::elasticsearch::elasticsearch_sink::elasticsearch_sink_stores_json_messagesdocker rm -frecovery targets onlyiggy-test-elasticsearch(no other Elasticsearch containers)
iggy-test-elasticsearchrunning and re-runthe sink test to verify reuse + readiness path
Notes
iggy-test-elasticsearchonly.tests sharing the reused container are not disrupted.