mattp5657 commented on PR #4121:
URL: https://github.com/apache/iggy/pull/4121#issuecomment-5622724301
# Connector Review - opensearch_sink (f3d9688d)
Target: `3873` (apache/iggy, open, unmerged)
Reviewed commit: `f3d9688d072e6b3337c7233a3fd2c5c353138abb`
Timestamp: 2026-09-10T00:00:00Z (America/New_York)
Roles: plugin, runtime, sdk, testing
Validators: 1 shard validator (all 5 claims) + 1 sweep validator
Contested: 0
## Review: add OpenSearch sink connector
### Confirmed (expert + clean-room validator)
- [nit] opensearch_sink::index_documents
(core/connectors/sinks/opensearch_sink/src/lib.rs:494) - multi-chunk
`consume()` keeps the *first* chunk error, not the last, unlike
`http_sink`/`mongodb_sink`/`surrealdb_sink`'s last-err convention. Fix: none
required, consistency-awareness only; `indexed`/`failed` counts are correct
regardless of which error message is kept. (raised: plugin; validated: PASS)
- [nit] OpenSearchSink::retry_on_open
(core/connectors/sinks/opensearch_sink/src/lib.rs:276) - a non-transient
`open()`-time failure is returned raw, not wrapped in `Error::InitError`,
unlike
`postgres_sink`/`mongodb_sink`/`http_sink`/`doris_sink`/`redshift_sink`/`surrealdb_sink`/`clickhouse_sink`
convention. Fix: wrap in `Error::InitError` at lib.rs:276 for consistency;
behaviorally inert - the FFI boundary (`sdk/src/sink.rs:89`,
`runtime/src/sink.rs:459-475`) collapses any `Result` to a generic
`i32`/message regardless of variant, so even the error text is lost either way.
(raised: plugin; validated: FIX - anchor corrected from the claimed
`:280`/`:338` (already-correct code) to the actual unwrapped-return site `:276`)
- [nit] OpenSearchSinkConfig
(core/connectors/sinks/opensearch_sink/src/lib.rs:74-75) - comment falsely
claims no in-tree redacted `SecretString` serializer exists;
`iggy_common::serde_secret::serialize_optional_redacted`
(core/common/src/utils/serde_secret.rs:85-93) predates this PR and does exactly
that. Fix: reword the comment to state the real reason (simplicity of omitting
`Serialize` entirely) rather than a false tooling claim; security posture is
unaffected either way. (raised: sdk; validated: PASS)
- [nit] core/integration/tests/connectors/opensearch/{opensearch_sink.rs:33,
opensearch_sink_failures.rs:132} - new integration tests use unit-test
`given_when_should` naming instead of the `connector-testing` skill's mandated
declarative `<subject>_<action>_<observation>` style. Fix: rename to match
`postgres_sink.rs::json_messages_sink_stores_as_bytea` /
`elasticsearch_sink.rs::elasticsearch_sink_stores_json_messages`. (raised:
testing; validated: FIX - origin corrected `pre-surfaced` -> `intro` since both
files are brand-new (`new file mode 100644`); framing corrected from "repo-wide
drift" to "widens the one prior instance" (`meilisearch_sink`), since the
declarative style is actually the dominant repo convention)
- [nit] OpenSearchSink::sleep_before_retry / index_documents
(core/connectors/sinks/opensearch_sink/src/lib.rs:669-672, also
:318,333,372,385,782,787) - retry/lifecycle/index-count logs omit connector ID
and index name, unattributable when several plugin instances share one
`LOG_CALLBACK` (exactly the setup the PR's own 3-instance failure test uses).
Sibling arm at lib.rs:798-801 does carry `self.id`; the success-count arm at
:782 drops it. Fix: add `connector ID: {self.id}` to every log site where
`&self` is in scope. (sweep, unvalidated)
- [nit] OpenSearchSinkConfig
(core/connectors/sinks/opensearch_sink/src/lib.rs:76) - derived `Debug` prints
`url` verbatim; embedded credentials in `url` are rejected only later, in
`normalize_url`, which never runs before `From<OpenSearchSinkConfig>`.
`ResolvedOpenSearchSinkConfig` hand-redacts this field but the raw config does
not, and the existing redaction test covers only `password`, not `url`. Latent
- no in-tree caller `Debug`-formats the raw config today. Fix: hand-write
`Debug` for `OpenSearchSinkConfig` through the existing
`redact_url_credentials` helper, or drop the derive. (sweep, unvalidated)
- [nit] opensearch_sink::index_chunk
(core/connectors/sinks/opensearch_sink/src/lib.rs:588) - the `?` on
`build_bulk_body` sits after the partial-merge step, so a failure there would
discard the `indexed` count already credited from earlier attempts, causing
already-indexed documents to be reported as failed. Unreachable today
(`serde_json::to_writer` over this `Value` cannot fail), so latent rather than
live. Fix: merge `outcome` into a returned `Ok` before propagating, or make
body construction infallible. (sweep, unvalidated, conf:M)
### Contested
None.
### Retracted (validator REMOVE)
None.
### Pre-existing (origin pre-*, not blocking)
None as findings against this PR. Two repo-wide gaps were independently
checked by the runtime and plugin experts and confirmed correctly disclosed in
the PR body rather than hidden: `runtime/src/sink.rs:740` discards the FFI
`consume` return value (so a plugin-side failure never reaches connector
status/`/stats`), and `runtime/src/sink.rs:522` commits offsets via
`AutoCommit::When(PollingMessages)` before `consume` runs. Both are
`pre-untouched` by this diff and affect every sink identically - not something
this PR introduces or is responsible for fixing.
### Simplification opportunities (non-blocking)
- [simplify] sanitize_url_for_log
(core/connectors/sinks/opensearch_sink/src/lib.rs:1296-1308) -
credential-stripping branches are dead in practice: the sole call site always
passes an already-sanitized `normalized_url`. Simpler: reduce to a
`to_string()` alias. Saves: ~6 lines - but the code already self-documents this
exact tradeoff as intentional defense-in-depth (lib.rs:1293-1295), and removing
it orphans the unit test at :2861, so marginal value is low. (raised: plugin;
validated: PASS, low marginal value)
- [simplify] opensearch_sink::documents_at
(core/connectors/sinks/opensearch_sink/src/lib.rs:928-936) - out-of-range
filter is dead: `parse_bulk_response` already rejects any response whose
`items.len() != expected` before `documents_at` runs. Same dead-defense class
as the item above, different function. Fix: none required; if kept, correct the
doc comment's panic claim to reflect that the state it guards against is
already excluded upstream. (sweep, unvalidated)
- [simplify] cluster_health / index_chunk
(core/connectors/sinks/opensearch_sink/src/lib.rs:298, lib.rs:545) - per-call
`tokio::time::timeout` duplicates the transport-level timeout already installed
(lib.rs:248-249), which covers strictly more (including body reads);
`index_exists`/`create_index` rely on the transport timeout alone, so the
codebase is internally inconsistent about which layer owns the bound. Fix: drop
the per-call wrappers, or document why two layers are intentional. (sweep,
unvalidated, conf:M)
### Verdict: APPROVE
Four independent experts (plugin, runtime, sdk, testing) plus clean-room
validation and an adversarial sweep found zero critical or warning findings.
Every confirmed item is a nit-level naming/logging/comment-accuracy issue or a
latent (not live) gap; two are already self-documented by the author as
intentional tradeoffs. `runtime/src/` has no changes, no STOP-tripwire was
triggered, and the mandatory connector checklist (SecretString, no `Serialize`,
FFI codes, idempotency, retry/transient classification, drop accounting,
forward-compat config, `&self` receivers, no blocking calls) verified clean
across all four passes.
Counts: critical 0, warning 0, nit 7, simplify 3
---
## Raw findings per expert
### plugin
```markdown
# Plugin review: PR 3873 (opensearch_sink) — role: plugin
Exemplars compared: `meilisearch_sink` (closest — same batch_size=1000
default,
same `Payload` dispatch shape, same `InvalidRecordValue` on unsupported
payload), `elasticsearch_sink` (closest sibling search-engine sink, same
`owned_value_to_serde_json` usage, same `close()` shape).
## Findings
[nit] opensearch_sink::index_documents
(core/connectors/sinks/opensearch_sink/src/lib.rs:494) - keeps the *first*
chunk error across a multi-chunk consume(), not the last. Every other
multi-batch sink (`http_sink:595`, `mongodb_sink:241`, `surrealdb_sink:596`)
uses a last-err accumulator. Deliberate and commented ("kept, not the last,
matching BulkOutcome::merge... the one an operator can act on"), and every
chunk is still processed and counted regardless of which error surfaces —
counts (`indexed`/`failed`) are correct either way, only the reported message
text differs. Fix: none required; flagging for consistency awareness only.
(intro, conf:H)
[nit] OpenSearchSink::retry_on_open
(core/connectors/sinks/opensearch_sink/src/lib.rs:263-290, call sites
:280,:338) - a non-transient open()-time failure returns the raw
`PermanentHttpError`/`Connection` error directly instead of wrapping it in
`Error::InitError`, unlike every other sink's connectivity-failure convention
(`postgres_sink:171`, `mongodb_sink:179-214`, `http_sink:347`,
`doris_sink:253`, `redshift_sink:146`, `surrealdb_sink:362`,
`clickhouse_sink:176` all wrap in `InitError`). Confirmed intentional and
unit-tested (`given_permanent_open_failure_should_not_retry`, lib.rs:2259-2274,
asserts `PermanentHttpError` survives untouched). Behaviorally inert:
`manager/sink.rs` flips `ConnectorStatus::Error` on any `Err` from `open()`
regardless of variant (sink.rs:106), so this is pure error-message
classification, not a functional gap. (intro, conf:M)
[simplify] sanitize_url_for_log
(core/connectors/sinks/opensearch_sink/src/lib.rs:1296-1308) - the
username/password-stripping branches are dead in practice: its only call site
(`open()`, lib.rs:685) always passes `normalized_url`, and `normalize_url`
(lib.rs:1221) already rejects any URL with embedded credentials before it can
reach here. Already documented as intentional defense-in-depth via the comment
above the fn. Simpler: could be removed or reduced to a `to_string()` alias
since the branch never fires today, but doing so would remove the safety net
for any future caller that bypasses `normalize_url`. Not worth changing. Saves:
~6 lines, but conf:L this is worth doing. (intro, conf:L)
## Verified correct (no finding, checked because charter mandates it)
- `SecretString` on `password`; no `Serialize` on `OpenSearchSinkConfig`
(comment at lib.rs:74 explains why — matches convention, e.g. no sink derives
`Serialize` on its config type).
- Custom `Debug` impls on both `OpenSearchSink` and
`ResolvedOpenSearchSinkConfig` redact the password and URL credentials
(lib.rs:109-165) — needed because `opensearch::auth::Credentials::Basic`
derives an unredacted `Debug`; verified via
`given_debug_formatted_*_should_redact_password` tests (lib.rs:2988-3040).
- `consume(&self, ...)` — no interior `&mut self`, no locks, no
`tokio::spawn`/`block_on`.
- `BTreeMap<HeaderKey, HeaderValue>` headers handled via `headers_to_json`
(lib.rs:1113), and unlike the pre-existing gap the PR documents in
`elasticsearch_sink`/`meilisearch_sink` (headers silently dropped because
`serde_json` can't serialize a `HeaderKey`-keyed map), this sink actually
converts and indexes them — a real fix, not a regression.
- `try_to_bytes`/no-clone JSON path: `prepare_document` uses `mem::replace`
to take `message.payload` by value (lib.rs:408), then
`owned_value_to_serde_json(&value)` (a reference-taking conversion helper
already used identically by `elasticsearch_sink:195`) — matches SDK convention,
not a clone-then-replace anti-pattern.
- `Vec::with_capacity(messages_count)` at lib.rs:741.
- `[lib] crate-type = ["cdylib", "lib"]` present in Cargo.toml.
- `verbose_logging` mirrors `debug!`→`info!` correctly in both `consume()`
receive-log and indexed-count log (lib.rs:709-733, 780-790).
- 3-attempt cap (`DEFAULT_MAX_RETRIES = 3`), exponential backoff + jitter
shared via SDK `retry` module, matches sibling sinks.
- Idempotency: dedup-on-write via deterministic `_id` (hash of
stream/topic/partition/offset/message-id, or user `document_id_field`),
verified end-to-end by `connectors::opensearch::opensearch_sink` (resends under
existing `order_id`, asserts count unchanged) — correctly distinguished from
the ES-auto-`_id` non-idempotent case the PR body calls out.
- No `unwrap()`/`expect()` on external-I/O `Result`s in production code
(verified via grep over lib.rs:1-1350; all `unwrap()`/`expect()` usage is in
`#[cfg(test)] mod tests`).
- Config forward-compat: every `OpenSearchSinkConfig` field is `Option<T>`
with defaults resolved in `From<OpenSearchSinkConfig> for
ResolvedOpenSearchSinkConfig` (lib.rs:167-203), conflict (retry_delay >
max_retry_delay) fixed with `warn!` + swap in `new()`-time conversion, never in
`consume()`.
- `Permanent*` vs transient classification:
`is_transient_error`/`is_transient_client_error`/`map_status_error`/`map_client_error`
all correctly route through `Error::HttpRequestFailed` (transient) vs
`Error::PermanentHttpError` (permanent), including the per-item `_bulk`
response classification (429/5xx retryable, everything else including mapping
errors permanent) — unit-tested extensively (lib.rs:1971-2194).
- Drop accounting: `invalid_records`/`preparation_errors` both roll into
`errors_count` (lib.rs:765-770); a single bad message in a batch does not lose
the rest (comment at lib.rs:754 states this explicitly and
`given_permanently_failing_chunk_should_not_abandon_later_chunks` proves the
chunk-boundary case at the unit level, integration test
`opensearch_sink_failures.rs` proves it against a live server).
- FFI/runtime consume()-swallow gap and
`elasticsearch_sink`/`meilisearch_sink` header-drop gap: both correctly
classified in the PR body as `pre-untouched` (verified against current `master`
citations, e.g. `sink.rs:740-748`, `elasticsearch_sink/src/lib.rs:205-219`) —
out of scope for this PR, not introduced by it.
- No `#[repr(C)]`, `Schema` variant, consumer-group, or `state.rs` changes —
no STOP-tripwire scope.
- `.config/nextest.toml` OpenSearch test-group addition mirrors the existing
Elasticsearch group exactly (serialized, shared reusable container by fixed
name) — consistent pattern, not scope creep.
## Simplifications: none (beyond the low-priority dead-branch note above,
not worth acting on)
## Verdict: APPROVE
```
### runtime
```markdown
# PR 3873 review — runtime lane (FFI host, `runtime/src/`)
Scope check first: `git diff --stat` / `files.txt` touch zero files under
`core/connectors/runtime/src/`. Only
`runtime/example_config/connectors/opensearch_sink.toml`
(new example, no code). This PR is a new sink plugin
(`core/connectors/sinks/opensearch_sink/`) plus its integration tests,
README/registry
entries, and workspace/nextest/bump-version plumbing. Nothing in my ownership
(FFI dispatch, plugin lifecycle, `plugin_id`, `LogCallback`, container `Arc`,
`DashMap`, state atomic-rename, metrics/config plumbing) is modified.
## Verified against runtime source (not changed by this PR)
- PR body claims: `core/connectors/runtime/src/sink.rs:740-748` invokes the
FFI
`consume` callback as a bare statement, discarding its `i32`, so a plugin
failure never reaches `ConnectorStatus`/`last_error`/`/stats`, and combined
with `AutoCommit::When(PollingMessages)` (`sink.rs:522`) the offset is
already committed by the time `consume` runs. Read `sink.rs:480-756`
directly: confirmed accurate on both counts — `process_messages` calls
`(consume)(plugin_id, ...)` at `sink.rs:740` with no binding, and
`setup_sink_consumers` builds the consumer with
`.auto_commit(AutoCommit::When(AutoCommitWhen::PollingMessages))` at
`sink.rs:522`. Pre-existing, untouched by this diff, correctly scoped out
of the PR rather than silently relied on. (pre-untouched, conf:H)
- `opensearch_sink/README.md`'s "Delivery Semantics" section states the same
gap and additionally claims it verified the failure mode against a live
server (mapping-conflict batch logged as `PermanentHttpError`, connector
stayed `Running`). Consistent with the runtime code read above — a
plugin-level `Err` from `consume()` has no path back into the host's status
machinery, so "stays Running" is the only possible outcome given current
`sink.rs`. No fix expected or requested here; this is the same repo-wide
limitation the PR states affects every sink.
## Connector-mandatory checks, in scope for this plugin
- `SecretString` on `password` (`lib.rs:591`), no `Serialize` derive on
`OpenSearchSinkConfig`/`ResolvedOpenSearchSinkConfig` — both have hand
comments explaining why (`lib.rs:584-585`, `lib.rs:651`). `Debug` impls are
hand-written to redact the URL and never print `password` in plaintext via
a derived path (`lib.rs:619-630`, `652-675`). Matches `postgres_sink`/
`http_sink` convention.
- `consume(&self, ...)` (`lib.rs:1211`) — immutable receiver, matches the
`Sink` trait and every other sink; no interior mutability via `Mutex`, only
`AtomicU64` counters, so no lock-across-`.await` risk.
- No `tokio::spawn` / `block_on` anywhere in the plugin (grepped the whole
file) — retries use `tokio::time::sleep`/`tokio::time::timeout` inline in
the `&self` async fns, which is how the runtime expects a plugin to behave
under its own executor.
- `[lib] crate-type = ["cdylib", "lib"]` present (`Cargo.toml:221-222`).
- `BTreeMap<HeaderKey, HeaderValue>` used for headers (`lib.rs:1623`,
matches `ConsumedMessage.headers` shape), converted explicitly to JSON
since `HeaderKey`/`HeaderValue` aren't `Serialize` as map keys — same
workaround shape the PR body says `elasticsearch_sink`/`meilisearch_sink`
skip (verified true against master's `serialize_headers` remark; this sink
does the conversion by hand rather than reusing that helper, but does not
drop headers, which is the actual bug in the other two).
- `Vec::with_capacity(messages_count)` at `lib.rs:1251`.
- FFI panic-safety: `documents_at` (`lib.rs:1438-1446`) explicitly drops
out-of-range positions from a server-echoed `items` array instead of
indexing/panicking, with a comment citing exactly the right reason — a
panic across the `extern "C"` plugin boundary aborts the whole connectors
runtime process. Backed by a dedicated regression test
(`given_out_of_range_positions_should_drop_them_instead_of_panicking`,
`lib.rs:2731-2741`). This is the one place in the plugin that actually
touches a runtime-host invariant, and it's handled correctly.
- Logging: `info!`/`debug!` pair gated on `verbose_logging`
(`lib.rs:1219-1243`, `1290-1300`) duplicates the exact same format string
under both branches. Checked whether this is a plugin-specific tell:
`mongodb_sink/src/lib.rs:263-273` and `redshift_sink/src/lib.rs:417-427` do
the identical `if self.verbose { info!(...) } else { debug!(...) }`
duplication already on master. Established repo-wide pattern, not
introduced here — not flagging as a new simplify target for this PR alone.
(pre-untouched convention, closest exemplar: `mongodb_sink`)
- Example config path has no platform extension
(`path = "target/release/libiggy_connector_opensearch_sink"`,
`example_config/connectors/opensearch_sink.toml:9`) — matches
`delta_sink.toml`/`clickhouse_sink.toml`/`iceberg_sink.toml`, the runtime
resolves the `.so`/`.dylib`/`.dll` suffix itself. Correct.
## STOP tripwires
None triggered: no `#[repr(C)]` touched, no `Schema` variant change, no
transient↔`Permanent*` promotion in runtime code (the plugin's own
classification in `is_transient_error`/`map_client_error`/`map_status_error`
is plugin-local, not a runtime reclassification), no consumer-group default
rename (`opensearch_sink_connector` is this plugin's own new default, not a
rename of an existing one), no plugin path resolution change, no
`iggy_connector_sdk` trait change.
## Findings
None in the runtime lane — no `runtime/src/` files changed, and everything
this plugin does that touches a runtime-owned invariant (FFI boundary panic
safety, crate-type, `&self` receiver, no blocking calls, config
forward-compat via `Option<T>`) is handled correctly.
Simplifications: none (no runtime-owned code changed; the one duplication
noted above — verbose/info vs debug logging — is a pre-existing repo-wide
pattern, not something this PR should be asked to fix unilaterally).
Verdict: APPROVE - no runtime/src/ changes; the plugin's own runtime-facing
behavior (FFI panic safety, crate-type, blocking-call discipline, config
forward-compat) is correct, and its two disclosed pre-existing runtime gaps
(consume() return value discarded, offset auto-committed before consume())
were independently verified against current `sink.rs` and are accurately
scoped out of this PR rather than papered over.
```
### sdk
```markdown
# SDK role review — PR #3873 (opensearch_sink) @ f3d9688d
Scope note: no files under `core/connectors/sdk/src/` are touched by this
diff.
Review covers whether the new plugin correctly *uses* the existing SDK
contract (`Sink` trait, `Error` enum, `retry.rs`, `convert.rs`,
`sink_connector!`
macro) rather than a contract change. Closest exemplar per PR body:
`elasticsearch_sink` (also compared against `http_sink`, `meilisearch_sink`
for
config-field conventions).
## Findings
[nit] opensearch_sink::OpenSearchSinkConfig
(core/connectors/sinks/opensearch_sink/src/lib.rs:74-75) - comment "the only
in-tree helper for a `SecretString` field writes the credential in plaintext"
is factually wrong: `iggy_common::serde_secret::serialize_optional_redacted`
(core/common/src/utils/serde_secret.rs:85-93) exists precisely for a struct
that must serialize but must not expose the secret, and predates this PR (added
before f3d9688d; the module's doc header even warns against exactly this kind
of inverted claim after a prior fix, commit f900e7abe "fix the inverted
serde_secret redaction claim", #3803). Omitting `Serialize` entirely is a
defensible choice on its own merits (elasticsearch_sink instead derives
`Serialize` with `serialize_optional_secret`, i.e. plaintext-on-purpose for the
control API; clickhouse_sink/delta_sink/doris_sink/redshift_sink also omit
`Serialize` like this PR does — split convention, not a deviation), but the
justification comment misstates th
e available tools and could steer a future maintainer wrong. Fix: reword to
something like "No `Serialize`: nothing in-tree serializes this type, and
`serialize_optional_redacted` is unused here because dropping the impl entirely
is simpler than remembering to route password through it." (intro, conf:H)
Simplifications: none. The bulk-retry state machine
(`BulkOutcome`/`BulkAttempt`/`index_chunk`) looks large but every branch is
exercised by a named test (`given_*`) and each piece (partial-chunk accounting,
transient vs permanent classification, first-failure retention) maps to a real
OpenSearch `_bulk` behavior called out in the PR body; collapsing it would drop
correctness, not just lines. `AtomicU64` counters instead of
`elasticsearch_sink`'s `Mutex<State>` is already a simplification (no lock, no
lock-across-await concern).
## SDK-contract compliance checklist
- `Sink` trait (core/connectors/sdk/src/lib.rs:133-147): `open(&mut self)`,
`consume(&self, ...)`, `close(&mut self)` signatures match exactly; `consume`
takes `&self` as required (no interior `Mutex` needed, uses `AtomicU64`).
- `Send + Sync` on `Sink`: satisfied automatically — every field
(`AtomicU64`, `Option<OpenSearch>`, `ResolvedOpenSearchSinkConfig` of
`String`/`bool`/`Duration`/`Option<Value>`/`Option<Refresh>`/`SecretString`) is
`Send + Sync`. No manual impl needed or attempted.
- `sink_connector!(OpenSearchSink)` (lib.rs:53) — plain macro invocation,
matches the macro's expansion in core/connectors/sdk/src/sink.rs:239-318
(duplicate-ID guard, `0`/`-1`/`1` FFI codes, `INSTANCES` global, `version()`
static) — all handled by the macro, plugin does nothing FFI-level itself.
`crate-type = ["cdylib", "lib"]` set correctly in Cargo.toml:32-33.
- `Error` variants used (`InvalidConfigValue`, `Connection`, `InitError`,
`HttpRequestFailed`, `InvalidRecordValue`, `PermanentHttpError`,
`Serialization`) — all pre-existing variants in
core/connectors/sdk/src/lib.rs:389+, used per their documented retry semantics:
`HttpRequestFailed` = transient (checked by
`is_transient_error`/`retry_on_open`), `PermanentHttpError` = non-retryable
data/schema issue (matches the variant's own doc comment about circuit breakers
not tripping on bad data). No new `Error` variant added — correct, none of
these needed distinct handling.
- `retry.rs` helpers (`exponential_backoff`, `jitter`, `parse_duration`,
`is_transient_status`) used with correct signatures.
`is_transient_status(status: reqwest::StatusCode)` called with
`opensearch::http::StatusCode` — verified same type:
`opensearch::http::mod.rs:41` is `pub use reqwest::StatusCode;`, and
`reqwest::StatusCode` is itself `pub use http::{StatusCode, ...}`
(reqwest-0.12.28/src/lib.rs:279) — no type mismatch, compiles.
- `max_retries` semantics ("N retries after the initial attempt", not "N
total attempts") differs from `retry.rs::HttpRetryMiddleware`'s self-documented
"total attempt count" convention, but matches the established sink-level
convention: http_sink's own doc comment (`batch_length * (max_retries + 1) *
...`, http_sink/src/lib.rs:1191) uses the same "+1 for initial attempt"
semantics. Not a deviation — two different sub-conventions coexist pre-PR
(middleware-level vs bespoke-retry-loop sinks), and this plugin correctly
follows the bespoke-loop one it actually implements.
- `convert::owned_value_to_serde_json`
(core/connectors/sdk/src/convert.rs:27) called with `&simd_json::OwnedValue`
from a destructured `Payload::Json` — signature matches, same usage pattern as
elasticsearch_sink.
- `Payload` match in `prepare_document` (lib.rs:410-423) handles
`Json`/`Raw`/`Text` explicitly, `Proto`/`FlatBuffer`/`Avro` fall to an explicit
`Error::InvalidRecordValue` (not a silent drop) — acceptable scope limit, no
SDK contract violation.
- `SecretString` on `password` (config field, lib.rs:81) with `ExposeSecret`
used only inside `create_client` (Basic auth header) — never logged, never in a
`Display`/`Debug` path (manual `Debug` impls on both `OpenSearchSink` and
`ResolvedOpenSearchSinkConfig` explicitly avoid deriving through the
client/password, with an accurate comment explaining why on lib.rs:105-108 and
lib.rs:141).
- No `tokio::spawn` / `block_on` / `std::sync::Mutex` anywhere in the plugin
(verified by grep across lib.rs); all `.unwrap()`/`.expect()` calls are
confined to `#[cfg(test)] mod tests` (first non-test line is 1351, all
unwrap/expect hits are ≥1405).
- BDD test naming: all unit tests are `#[test] fn given_..._should_...` —
consistent 3-part naming, matches CLAUDE.md convention.
- Config forward-compat: every `OpenSearchSinkConfig` field is `Option<T>`
except the two required ones (`url`, `index`), consistent with
`#[serde(default)]`-equivalent forward compatibility (no `#[serde(default)]`
needed since there's no `Default` derive requirement beyond what's already
`Option`-wrapped; struct does derive `Default`).
- `Refresh` config field: `Option<Refresh>` where
`opensearch::params::Refresh` derives `Deserialize` with `#[serde(rename =
"true"/"false"/"wait_for")]` (opensearch-2.4.0/src/params.rs:152-159) —
`refresh = "false"` in config.toml/README/example_config all deserialize
correctly.
- Workspace `Cargo.toml` diff (lines 89-101 of diff.patch) only adds the new
workspace member and the `opensearch = "2.4.0"` dependency — no version bump to
`iggy_connector_sdk` or any STOP-tripwire crate in this PR's actual diff (the
version-bump lines visible via a stale `git diff HEAD~1` are from the unrelated
`Merge branch 'master'` commit at HEAD, not from this PR's content — confirmed
against diff.patch directly).
## Verdict: APPROVE - clean SDK-contract usage; one nit-level comment
inaccuracy (not a correctness or security issue — the credential itself is
still never leaked, since `Serialize` is correctly omitted either way)
```
### testing
```markdown
# Testing review - PR #3873 (opensearch_sink)
[nit]
core/integration/tests/connectors/opensearch/opensearch_sink.rs::given_json_messages_when_sink_consumes_should_index_documents_and_upsert_by_natural_key
(opensearch_sink.rs:33) - integration test uses unit-test BDD
`given_when_should` naming instead of `connector-testing` SKILL.md's mandated
declarative `<subject>_<action>_<observation>` for integration tests ("never
given_should there"). Same pattern in opensearch_sink_failures.rs:132
(`given_missing_index_and_mapping_conflict_should_isolate_failures_from_healthy_sibling`).
Fix: rename to e.g. `json_messages_sink_indexes_and_upserts_by_natural_key` /
`missing_index_and_mapping_conflict_isolate_from_healthy_sibling`, matching
`postgres_sink.rs::json_messages_sink_stores_as_bytea` or
`elasticsearch_sink.rs::elasticsearch_sink_stores_json_messages`. Low priority:
`meilisearch_sink.rs::given_json_messages_when_sink_consumes_should_index_documents`
(near-identical name/shape to this PR's test) and `doris_sink.rs`,
`quickwit_sink
.rs`, `runtime/benchmark.rs`, `runtime/http_state.rs` all already use the same
`given_should` style in integration tests, so this is repo-wide drift the PR
inherited rather than introduced. (origin: pre-surfaced, conf: H)
Simplifications: none. Unit test file
(`core/connectors/sinks/opensearch_sink/src/lib.rs:1351-3085`, ~1730 lines) is
large but every test targets a distinct branch with a comment justifying why
(e.g. the four separate `given_debug_formatted_*_should_redact_password` tests
each catch a different leak surface: config Debug, sink Debug with no client,
sink Debug with a real opened `OpenSearch` client - which has its own
un-redacted `Debug` derive down through `Credentials::Basic` - and URL-embedded
credentials). No duplicate/mocked-should-be-real-infra tests, no fixed sleeps
outside legitimate poll loops or backoff-under-test.
## What was verified
- **BDD naming, unit tests**: consistent `given_X_should_Y` /
`given_X_when_Y_should_Z` throughout `opensearch_sink/src/lib.rs::tests`
(lib.rs:1408-3084). No `test_foo`/`does_bar` names.
- **Config helper**: `base_config()` (lib.rs:1386) matches the naming
convention of the closest exemplar `meilisearch_sink::tests::base_config` (not
`test_config`, which is fine - `meilisearch_sink`, `mongodb_sink`
(`given_default_config`), and `http_sink` (`given_default_config`) all name
this helper differently; no single fixed name is enforced repo-wide).
- **`#[tokio::test]` in `src/lib.rs`**: used for all async unit tests
(lib.rs:2241 onward) rather than local `Runtime::new()`. This looks like a
skill-guideline deviation at first read, but the closest exemplars for
wiremock-backed sink tests, `http_sink/src/lib.rs` and
`mongodb_sink/src/lib.rs`, both already use `#[tokio::test]` directly in
`src/lib.rs`. The `Runtime::new()` guidance in the skill is written for source
state round-trip tests (`random_source`-style); sinks with real async HTTP
mocking already establish `#[tokio::test]` as the norm. Not a defect. (origin:
pre-surfaced, conf: H)
- **Sink pure-logic coverage** (mandatory checklist item): config defaults +
fallback/clamping (lib.rs:1409-1442), all three `Payload` variants incl.
non-object JSON wrapping and raw-non-JSON base64 fallback (lib.rs:1494-1715),
header encoding incl. the dynamic-mapping-collision regression test
(lib.rs:1648-1684), `_bulk` response/query building and parsing incl.
malformed/short/missing-items edge cases (lib.rs:1923-2231), and
transient-vs-permanent classification for both client errors and per-item bulk
statuses (lib.rs:2845-2856, 1170-1195, plus the wiremock-driven retry tests at
lib.rs:2394-2653). This is the most thorough sink test suite of the set I
compared against (`elasticsearch_sink`, `meilisearch_sink`, `mongodb_sink`,
`http_sink`).
- **No four canonical source-state tests required**: this is a sink, not a
source - correctly out of scope.
- **Integration layout**:
`core/integration/tests/connectors/opensearch/{mod.rs,opensearch_sink.rs,opensearch_sink_failures.rs,sink.toml,failure_states.toml,failure_states/*.toml}`
plus `fixtures/opensearch/{mod.rs,container.rs,sink.rs,failure.rs}` mirrors
`postgres`/`elasticsearch` exactly: `TestFixture` impl, `ConfigEnv`
env-injection (all three documented forms: plugin-config leaf, indexed
`streams_N_*`, top-level `path`), `iggy-test-opensearch` fixed name +
`ReuseDirective::Always` matching the `elasticsearch`/`doris` reuse pattern
(skill's "as of now only elasticsearch and doris share a container" note is now
stale text, not a PR defect), polling helpers (`wait_for_document_count`,
`wait_for_status`) instead of fixed sleeps, `iggy-test-` container prefix
honored.
- **`.config/nextest.toml`**: new `[test-groups.opensearch]` + `max-threads
= 1` override correctly serializes the reused-container test group, matching
the `elasticsearch`/`doris` blocks immediately above it (nextest.toml:56-69).
- **Two integration tests cover exactly what's claimed**: happy path
(natural-key upsert idempotency proven end-to-end by resending `order_id: A-1`
at a new offset, headers round-trip, both `Payload::Raw` branches via a second
stream/topic on the `raw` schema) and a real failure-isolation test
(missing-index `open()` failure vs. a live `mapper_parsing_exception`
mid-`consume()`, three-sink isolation, chunk-boundary survival cross-checked
against the unit-level
`given_permanently_failing_chunk_should_not_abandon_later_chunks`). The module
doc in `opensearch_sink_failures.rs:18-65` explicitly traces the runtime code
path (`sink.rs::process_messages` discarding the FFI `consume` return) backing
its "not visible anywhere" claim rather than asserting it blind.
- **Doc/config sync**: `opensearch_sink/README.md` config table,
`opensearch_sink/config.toml`, and
`runtime/example_config/connectors/opensearch_sink.toml` all list the identical
`[plugin_config]` field set (defaults match `OpenSearchSinkConfig`'s
`unwrap_or` values in lib.rs). `core/connectors/README.md` and
`core/connectors/sinks/README.md` both gained the new sink row in alphabetical
position. No `max_connections` field applies here (sink has no connection
pool); retry knobs (`max_retries`, `retry_delay`, `max_retry_delay`,
`max_open_retries`) and `verbose_logging` are all present in both TOMLs and
documented.
- **Investigated and ruled out**: an apparent "PR deletes `rabbitmq_sink`"
signal from a raw `git diff master -- ...` was a false positive from the local
`master` ref having advanced past this PR's actual merge-base (post-merge
commits landed on master after this branch's base). The actual `diff.patch`
supplied for review contains zero references to `rabbitmq_sink`; confirmed via
direct grep. Not a real finding.
Verdict: APPROVE
```
## Validation record
- Shard validator (validate-1.md): 5/5 claims PASS on substance; 2 FIX (C2
line-anchor correction, C5 origin-tag correction); 0 REMOVE.
- Sweep validator (sweep.md): 5 additions (S1-S3 nit, S4-S5 simplify), all
conf H/M, none critical/warning; 0 corrections that reverse a PASS (C5's
"repo-wide drift" framing narrowed to "one prior instance, widened by this PR"
- doesn't change verdict or severity).
- Contested items: 0 (no REMOVE, no critical/warning downgrade, no sweep
addition contradicts a PASS).
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]