feat(connectors): add OpenSearch sink connector - #3873
Conversation
|
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 #3873 +/- ##
============================================
+ Coverage 83.85% 83.91% +0.05%
Complexity 1358 1358
============================================
Files 1212 1213 +1
Lines 166843 168701 +1858
Branches 134304 136288 +1984
============================================
+ Hits 139905 141559 +1654
- Misses 23298 23320 +22
- Partials 3640 3822 +182
🚀 New features to boost your workflow:
|
|
Getting some errors with a 503 returned when they passed locally, for example for Typo. Will try running again later: Connecting to github.com (github.com)|140.82.112.4|:443... connected.
HTTP request sent, awaiting response... 503 Service Unavailable
2026-08-12 19:30:09 ERROR 503: Service Unavailable.Update: These seem to have resolved. |
|
/ready |
| // The top-level flag lets a clean batch skip the per-item scan entirely. | ||
| if !response | ||
| .get("errors") | ||
| .and_then(Value::as_bool) | ||
| .unwrap_or(true) | ||
| { | ||
| return Ok(BulkAttempt { | ||
| indexed: items.len(), | ||
| ..BulkAttempt::default() |
There was a problem hiding this comment.
This fast path still accepts malformed item entries without validating them. For example, {"errors":false,"items":[{}]} passes the length check and is reported as one indexed document, even though the response contains no recognizable index result or successful status. That leaves the malformed-response fix incomplete, and the runtime will commit the offset for a document that was never actually accounted for. Please validate that every item contains the expected operation result and a 2xx status before returning success, and add an errors: false regression case with a malformed item.
|
|
||
| [package] | ||
| name = "iggy_connector_opensearch_sink" | ||
| version = "0.5.0-edge.1" |
| impl std::fmt::Debug for OpenSearchSink { | ||
| fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { | ||
| f.debug_struct("OpenSearchSink") | ||
| .field("id", &self.id) | ||
| .field("config", &self.config) | ||
| .field("client", &self.client.is_some()) | ||
| .field("invocations_count", &self.invocations_count) | ||
| .field("documents_indexed", &self.documents_indexed) | ||
| .field("errors_count", &self.errors_count) | ||
| .finish() | ||
| } | ||
| } |
There was a problem hiding this comment.
This hides the authenticated client and relies on SecretString for the dedicated password field, but it still prints self.config.url verbatim through self.config. Since embedded URL credentials are rejected only later in open(), formatting a newly constructed sink from https://admin:hunter2@host exposes hunter2. Let's redact or sanitize the URL in the resolved config's Debug output.
| secrecy = { workspace = true } | ||
| serde = { workspace = true } | ||
| serde_json = { workspace = true } | ||
| simd-json = { workspace = true } |
There was a problem hiding this comment.
simd-json is used only by test fixtures inside the #[cfg(test)] module. It should be under [dev-dependencies].
Which issue does this PR address?
Closes #3504
Rationale
Iggy connectors ship sinks for several external systems but not for OpenSearch, a widely used search/analytics backend. This adds one, modeled on the existing
elasticsearch_sinkshape but closing a retry gap that sink still has (see Known trade-offs).What changed?
Adds
core/connectors/sinks/opensearch_sink/, following the sink lifecycle end to end:open(): validates config (URL shape, credential pairing,document_id_fieldconstraints), thenretry_on_open-wraps a cluster health check, an index-exists check, and, whencreate_index_if_not_exists(defaulttrue), index creation with an optional custom mapping. Capped atmax_open_retrieswith exponential backoff and jitter (sharedretryhelpers).consume(): batches incoming messages intobatch_sizechunks, builds aniggy_*-enriched document per message (hashed or field-derived_id), and hands each chunk toindex_chunk.index_chunk(): POSTs a_bulkrequest and loops up tomax_retrieswith the same backoff._bulkanswers200even when individual documents fail, so each response is parsed per item rather than trusting the top-level status: permanent failures (4xx, e.g. a mapping conflict) are recorded immediately, while transient ones (429/5xx) shrink the pending set to just the rejected documents and get resent, so a partial rejection under load doesn't re-index or lose the rest of the chunk. Counts from earlier attempts are merged into the outcome so a later failure doesn't erase already-indexed documents from the tally.close(): drops the client, no special teardown.Integration tests (
core/integration/tests/connectors/opensearch/) run against a real container (testcontainers-modules, reused across tests viaReuseDirective::Always, per-test-unique index names for isolation), covering the happy path plus a static mapping conflict, a missing index with index-creation disabled, and confirming a failing chunk doesn't block chunks queued behind it.Credentials: HTTP Basic auth only (
username/password, both-or-neither validated at config time);passwordis aSecretString, never logged or serialized. AWS SigV4 (AWS-managed OpenSearch / Serverless) is not supported.Known trade-offs, deliberately out of scope here, verified against current
master:elasticsearch_sinkhas the same gap this PR fixes for OpenSearch:bulk_index_documents(elasticsearch_sink/src/lib.rs:205-219) tallies per-item_bulkfailures intoerrors_countbut never retries the transient subset (e.g. 429es_rejected_execution_exception). Worth a follow-up issue rather than folding into this PR.consume()return value:core/connectors/runtime/src/sink.rs:740-748invokes the FFIconsumecallback as a bare statement, never binding itsi32result, soprocess_messagesalways returnsOk. Combined with offsets auto-committing at poll time (sink.rs:522), a plugin-level failure never reaches connector status,last_error, or/stats, and the batch is never redelivered. Pre-existing, repo-wide, affects every sink.meilisearch_sink(lib.rs:451-454) andelasticsearch_sink(lib.rs:319-328) both silently dropiggy_headers/_iggy_headers:BTreeMap<HeaderKey, HeaderValue>can't serialize as a JSON object (serde_jsonrequires string keys), and both sinks swallow that error viaif let Ok(...)instead of surfacing it.core/commoneven ships aserialize_headersworkaround for this exact case that neither sink uses.elasticsearch_source, the client is built onSingleNodeConnectionPool(lib.rs:208) with no cluster sniffing or multi-node failover. A dead configured node fails every request rather than routing around it.Local Execution
AI Usage