spetz commented on code in PR #4121: URL: https://github.com/apache/iggy/pull/4121#discussion_r3981857140
########## .claude/skills/connector-review/SKILL.md: ########## @@ -0,0 +1,151 @@ +--- +name: connector-review +description: Adversarial 4-expert review of a connectors PR, branch, or ref range (sinks, sources, runtime, SDK, transforms) with clean-room validation of every finding. Experts work alone, no peer debate. Expensive, one run spawns ~10 subagents. Load connectors-overview first. +argument-hint: "[PR number | branch | ref range]" +disable-model-invocation: true +--- + +# Apache Iggy Connector Review + +`<TARGET>` = `$ARGUMENTS`: a PR number, a branch, or a ref range. Empty means `origin/master..HEAD`. Scope: anything under `core/connectors/` plus connector integration tests under `core/integration/tests/connectors/`. + +Prereq: load `connectors-overview` first as router (sibling skill in this directory; repo-wide rules in `<repo>/AGENTS.md`). This file owns **review mechanics only**. + +You = **moderator**. You never open the diff or a source file: you route paths, merge claims, synthesize. Every token you load rides along every later turn. Reviewers and validators are one-shot agents that deliver by writing a file; nobody chats. + +## Charter (paste VERBATIM into every expert, validator, and tiebreak prompt) + +> You think big brain. You speak caveman. Separate things. +> +> **Thinking, unchanged.** Read the diff, then every changed file in full from the local checkout, then whatever call sites you need. Trace call chains. Verify invariants. Prove findings, don't guess. Cite `path::Symbol`, plus `file:line` as secondary anchor. Running tests or builds needs a stated justification: reading and tracing settles most claims, and parallel cargo runs block on one target-dir lock. +> +> **Output style.** Drop articles, filler, pleasantries, hedging. Fragments OK. Keep EXACT: symbol, `file:line`, error quotes, code, technical terms, severity and confidence labels. +> +> - Finding, one line each: `[sev] path::Symbol (file:line) - problem. Fix: action. (origin, conf:H|M|L)` +> - `sev`: `critical` = correctness/safety/data-loss/security, blocks merge; `warning` = real defect, perf hit, API issue; `nit` = style/naming; `simplify` = complexity/dead-code reduction, format `[simplify] path::Symbol (file:line) - what's complex. Simpler: alternative. Saves: ~N lines / removes indirection. (origin, conf)`. +> - `origin`: `intro` (PR introduced), `pre-surfaced` (existed, exposed by PR), `pre-untouched` (existed, not touched). +> - Never flag em dashes or other punctuation style as a finding. +> - Simplification mandate: less code > more code. Per changed file ask whether ~30% smaller keeps correctness: dead fields/params/branches/imports, duplication of an existing helper (cite it), single-impl traits, premature generics, checks for impossible states. Do not propose simplifications that change semantics or break public API. If nothing qualifies, write `Simplifications: none`. +> +> **Connector mandatory checks (every expert, where in scope).** Verify, don't assume: +> `SecretString` on credential fields, no `Serialize` on plugin config (or redacted serializer with justification); FFI return codes honored (`0` ok, `-1` invalid, `1` open failure; duplicate-ID guard intact); sink redelivery claim vs `AutoCommit::When(PollingMessages)` + FFI swallow path; idempotency key on retry (message `id` dedup, `ON CONFLICT`, composite `_id`, deterministic run id + 409 handling). +> `Permanent*` vs transient classification vs retry strategy agreement; drop accounting (`errors` vs `messages_filtered`, flushed once per batch via pre-built labels); every error lands at its level (fatal exits process, per-plugin sets status, per-message skips, per-batch bumps metric or logs). +> Config forward-compat (`Option`/`#[serde(default)]`); `consume`/`poll` take `&self`; no `tokio::spawn`/`block_on` in plugins; no lock held across `.await`; `BTreeMap` headers; `try_to_bytes(&self)` no-clone for JSON; `simd_json::OwnedValue` mutated in place, never clone-then-replace; `Vec::with_capacity`; `[lib] crate-type = ["cdylib", "lib"]`; benchmark/verbose flags wired; closest exemplar named. +> STOP tripwires are critical: `#[repr(C)]` change without SDK bump; `Schema` variant rename; transient↔`Permanent*` promotion; consumer-group rename; plugin path resolution change; trait change without changelog + all-plugins-same-PR + bump + old-`.so` note. +> `state.rs` protocol change is critical: 0o600 mode, tmp `sync_data`, post-rename parent-dir `sync_all`, save/load failures surface (`CannotWriteStateFile`/`CannotOpenStateFile`), never silent fresh-start. Scope creep (unrelated refactor in same PR) gets flagged. +> Logging shape: connector ID + name on every line, literal API field labels, `debug!` default upgraded to `info!` on verbose only, never eager `format!()` around tracing args, never log secrets. Control API returns plugin config verbatim (#3802): never demand plugin-side fixes for control-plane exposure. +> +> Caveman = output compression, not analysis compression. Dig deep. Write short. + +## Step 1: Identify the target (no reading) + +- Classify `<TARGET>`: matches `^#?(pr)?[0-9]+$` case-insensitively -> PR, the digits are `<PR>`. Anything else -> ref range or bare branch. Empty -> review `origin/master..HEAD`. +- `<TOPIC>`: `<TARGET>` lowercased, chars outside `[a-z0-9-]` replaced by `-`, repeats collapsed, trimmed, max 40 chars (`PR3123` -> `pr3123`, `origin/master..HEAD` -> `origin-master-head`). Empty -> `date +%s`. +- `<DIR>` = `<session scratchpad dir from your system prompt>/review-<TOPIC>`. `mkdir -p` it. +- PR: `gh pr view <PR> --json title,body,headRefOid > <DIR>/pr.json`, `gh pr diff <PR> > <DIR>/diff.patch`, `gh pr diff <PR> --name-only > <DIR>/files.txt`. `<SHORTCOMMIT>` = first 8 of `headRefOid`. +- Ref range or bare branch: `git diff $(git merge-base origin/master HEAD)..HEAD > <DIR>/diff.patch`, same with `--name-only`, `<SHORTCOMMIT>` = `git rev-parse --short=8 HEAD`. No `pr.json` on this path. +- Guard: `git rev-parse HEAD` must equal the reviewed head. Experts read the local checkout; if it differs, stop and ask the user to check out the reviewed head. +- `<DESCR>`: 1-3 word `snake_case` summary, `[a-z0-9_]`, <= 24 chars. From the PR title; no PR -> from `git log -1 --format=%s`. +- Report path: `<DIR>/report.md`. + +Do not `cat` any of the files you just wrote. `wc -l <DIR>/diff.patch` is the only look you take. + +## Step 2: Round 1, four one-shot experts (one message, parallel) + +Spawn 4 `Agent` calls in a single message: `subagent_type: general-purpose`, `name: <role>-<TOPIC>` (bare role names collide with concurrent sessions: one shared agent namespace), no `model` (inherits). Prompt = role block + Charter + this brief, with `<DIR>`, `<TARGET>`, `<SHORTCOMMIT>` filled in: + +> Target: `<TARGET>` at `<SHORTCOMMIT>`. Diff: `<DIR>/diff.patch`. Changed files: `<DIR>/files.txt`. PR title and body: `<DIR>/pr.json` (drop this sentence when there is no PR). Classify each finding's origin; check existing codebase conventions before calling a deviation `intro`. Name the closest exemplar plugin compared (`stdout_sink`, `postgres_sink`/`postgres_source`, `http_sink`, `mongodb_sink`, `random_source`, ...). +> Deliverable = the file `<DIR>/<role>.md`, written with the Write tool BEFORE you end your turn: findings in Charter format, then `Simplifications: ...`, then `Verdict: APPROVE | REQUEST CHANGES - reason`. A previous worker finished reading and then idled without delivering; the Write call IS the delivery, your final message is just the path. Budget 3/4 reading, 1/4 writing; partial beats unshipped. +> You work alone: no teammates, no SendMessage, no questions back. + +Role blocks: + +- **plugin**: Senior connector-plugin engineer. Owns sink/source impls. + Sink focus: lifecycle (`open` connectivity fail-fast `InitError`, `close` flush + `.take()` + final stats), config `Option` defaults + conflict fix in `new()` + `warn!` (never error, never validate in `consume()`), durations `Option<String>` humantime (no `humantime_serde`). + Sink payloads/errors: dispatch Json/Raw/Text minimum (`InvalidPayloadType` or base64 fallback), `try_to_bytes` + `mem::replace` + `with_capacity`, header `is_empty` check + binary-base64 vs text split in `BTreeMap`, `InvalidRecordValue` skip+log, `WriteFailure` caller-decides, `CatalogCommitError` txn consumed and not idempotent. + Sink retry/idempotency: 3-attempt cap, `Retry-After` honored, per-status custom strategy, SQLSTATE-style transient map; dedup-on-write via message `id` (mongo composite `_id` reference; ES auto-`_id` is NOT idempotent); `last_err` pattern (process all batches, return last); batch `chunks` default 100, no cross-`consume` buffering, no config mutation in `consume()`. + Source focus: `poll()` `Err` only logs, loop continues, status never flips (runtime sets `Error` on transform/encode/send/save failure only); always return state incl. empty polls; poll-return-to-save crash window means at-least-once; sleep-first or idle spins CPU; single `poll()` task. + Source state/ids: `ProducedMessage.id` natural ID + `origin_timestamp` nanos, `timestamp`/`checksum` left `None`; brief-lock cursor pattern; state helpers `Option` + log, non-fatal; `State` small and bounded; `SOURCE_SENDERS` cleaned on close; sync poll blocks a worker, `handle` with no timeout stalls one. + Transform focus: three `mod.rs` edits (module, `TransformType`, `from_config`) + `snake_case` TOML match; `InvalidConfigValue` with context; JSON-only with non-JSON passthrough (format converters exempt); compile regex/paths once in `new()`, plain fields, no locks; sync stays sync; stateless `Send + Sync`. + Transform errors: missing field prefers passthrough over `Err` unless config says fail; `Ok(None)` filter logs at `debug!`; no `SystemTime::now()`. + Logging: redaction at construction, `verbose_logging` mirror. + Simplify: dead payload arms, skip branches that never fire, single-variant enums, near-duplicate path/config helpers. +- **runtime**: Connectors runtime + FFI host engineer. Owns `runtime/src/`. + Focus: FFI pointer lifetimes (call-duration only), `plugin_id` monotonic never reused, `LogCallback` static, container `Arc` outlives tasks (no unload mid-call), `DashMap` keyed by TOML key + status counters on transitions. + Loops: sink (autocommit timing, `consumer.next()` Err auto-committed drop, decode/transform drop-and-continue, postcard encode, nonzero-return, flush on offset gap, failed batch still emits); source (flume handoff, close-then-`cleanup_sender`-then-drain, state save after Iggy send, transform `Err` never flips status); state atomic-rename protocol; `restart_guard.try_lock()`. Review Comment: “Transform `Err` never flips status” contradicts the runtime and line 67. `source_forwarding_loop` rejects batches with processing errors, calls `context.sources.set_error`, and returns NACK. Update this instruction to match that behavior. ########## .claude/skills/connector-review/SKILL.md: ########## @@ -0,0 +1,151 @@ +--- +name: connector-review +description: Adversarial 4-expert review of a connectors PR, branch, or ref range (sinks, sources, runtime, SDK, transforms) with clean-room validation of every finding. Experts work alone, no peer debate. Expensive, one run spawns ~10 subagents. Load connectors-overview first. +argument-hint: "[PR number | branch | ref range]" +disable-model-invocation: true +--- + +# Apache Iggy Connector Review + +`<TARGET>` = `$ARGUMENTS`: a PR number, a branch, or a ref range. Empty means `origin/master..HEAD`. Scope: anything under `core/connectors/` plus connector integration tests under `core/integration/tests/connectors/`. + +Prereq: load `connectors-overview` first as router (sibling skill in this directory; repo-wide rules in `<repo>/AGENTS.md`). This file owns **review mechanics only**. + +You = **moderator**. You never open the diff or a source file: you route paths, merge claims, synthesize. Every token you load rides along every later turn. Reviewers and validators are one-shot agents that deliver by writing a file; nobody chats. + +## Charter (paste VERBATIM into every expert, validator, and tiebreak prompt) + +> You think big brain. You speak caveman. Separate things. +> +> **Thinking, unchanged.** Read the diff, then every changed file in full from the local checkout, then whatever call sites you need. Trace call chains. Verify invariants. Prove findings, don't guess. Cite `path::Symbol`, plus `file:line` as secondary anchor. Running tests or builds needs a stated justification: reading and tracing settles most claims, and parallel cargo runs block on one target-dir lock. +> +> **Output style.** Drop articles, filler, pleasantries, hedging. Fragments OK. Keep EXACT: symbol, `file:line`, error quotes, code, technical terms, severity and confidence labels. +> +> - Finding, one line each: `[sev] path::Symbol (file:line) - problem. Fix: action. (origin, conf:H|M|L)` +> - `sev`: `critical` = correctness/safety/data-loss/security, blocks merge; `warning` = real defect, perf hit, API issue; `nit` = style/naming; `simplify` = complexity/dead-code reduction, format `[simplify] path::Symbol (file:line) - what's complex. Simpler: alternative. Saves: ~N lines / removes indirection. (origin, conf)`. +> - `origin`: `intro` (PR introduced), `pre-surfaced` (existed, exposed by PR), `pre-untouched` (existed, not touched). +> - Never flag em dashes or other punctuation style as a finding. +> - Simplification mandate: less code > more code. Per changed file ask whether ~30% smaller keeps correctness: dead fields/params/branches/imports, duplication of an existing helper (cite it), single-impl traits, premature generics, checks for impossible states. Do not propose simplifications that change semantics or break public API. If nothing qualifies, write `Simplifications: none`. +> +> **Connector mandatory checks (every expert, where in scope).** Verify, don't assume: +> `SecretString` on credential fields, no `Serialize` on plugin config (or redacted serializer with justification); FFI return codes honored (`0` ok, `-1` invalid, `1` open failure; duplicate-ID guard intact); sink redelivery claim vs `AutoCommit::When(PollingMessages)` + FFI swallow path; idempotency key on retry (message `id` dedup, `ON CONFLICT`, composite `_id`, deterministic run id + 409 handling). +> `Permanent*` vs transient classification vs retry strategy agreement; drop accounting (`errors` vs `messages_filtered`, flushed once per batch via pre-built labels); every error lands at its level (fatal exits process, per-plugin sets status, per-message skips, per-batch bumps metric or logs). +> Config forward-compat (`Option`/`#[serde(default)]`); `consume`/`poll` take `&self`; no `tokio::spawn`/`block_on` in plugins; no lock held across `.await`; `BTreeMap` headers; `try_to_bytes(&self)` no-clone for JSON; `simd_json::OwnedValue` mutated in place, never clone-then-replace; `Vec::with_capacity`; `[lib] crate-type = ["cdylib", "lib"]`; benchmark/verbose flags wired; closest exemplar named. +> STOP tripwires are critical: `#[repr(C)]` change without SDK bump; `Schema` variant rename; transient↔`Permanent*` promotion; consumer-group rename; plugin path resolution change; trait change without changelog + all-plugins-same-PR + bump + old-`.so` note. +> `state.rs` protocol change is critical: 0o600 mode, tmp `sync_data`, post-rename parent-dir `sync_all`, save/load failures surface (`CannotWriteStateFile`/`CannotOpenStateFile`), never silent fresh-start. Scope creep (unrelated refactor in same PR) gets flagged. +> Logging shape: connector ID + name on every line, literal API field labels, `debug!` default upgraded to `info!` on verbose only, never eager `format!()` around tracing args, never log secrets. Control API returns plugin config verbatim (#3802): never demand plugin-side fixes for control-plane exposure. +> +> Caveman = output compression, not analysis compression. Dig deep. Write short. + +## Step 1: Identify the target (no reading) + +- Classify `<TARGET>`: matches `^#?(pr)?[0-9]+$` case-insensitively -> PR, the digits are `<PR>`. Anything else -> ref range or bare branch. Empty -> review `origin/master..HEAD`. +- `<TOPIC>`: `<TARGET>` lowercased, chars outside `[a-z0-9-]` replaced by `-`, repeats collapsed, trimmed, max 40 chars (`PR3123` -> `pr3123`, `origin/master..HEAD` -> `origin-master-head`). Empty -> `date +%s`. +- `<DIR>` = `<session scratchpad dir from your system prompt>/review-<TOPIC>`. `mkdir -p` it. +- PR: `gh pr view <PR> --json title,body,headRefOid > <DIR>/pr.json`, `gh pr diff <PR> > <DIR>/diff.patch`, `gh pr diff <PR> --name-only > <DIR>/files.txt`. `<SHORTCOMMIT>` = first 8 of `headRefOid`. +- Ref range or bare branch: `git diff $(git merge-base origin/master HEAD)..HEAD > <DIR>/diff.patch`, same with `--name-only`, `<SHORTCOMMIT>` = `git rev-parse --short=8 HEAD`. No `pr.json` on this path. Review Comment: This command ignores `<TARGET>` and always compares the current checkout against `origin/master`. Supplying another branch or revision range therefore reviews the wrong changes. Construct the diff from the supplied target and verify the checkout against that target’s head. ########## .claude/skills/connector-review/SKILL.md: ########## @@ -0,0 +1,151 @@ +--- +name: connector-review +description: Adversarial 4-expert review of a connectors PR, branch, or ref range (sinks, sources, runtime, SDK, transforms) with clean-room validation of every finding. Experts work alone, no peer debate. Expensive, one run spawns ~10 subagents. Load connectors-overview first. +argument-hint: "[PR number | branch | ref range]" +disable-model-invocation: true +--- + +# Apache Iggy Connector Review + +`<TARGET>` = `$ARGUMENTS`: a PR number, a branch, or a ref range. Empty means `origin/master..HEAD`. Scope: anything under `core/connectors/` plus connector integration tests under `core/integration/tests/connectors/`. + +Prereq: load `connectors-overview` first as router (sibling skill in this directory; repo-wide rules in `<repo>/AGENTS.md`). This file owns **review mechanics only**. + +You = **moderator**. You never open the diff or a source file: you route paths, merge claims, synthesize. Every token you load rides along every later turn. Reviewers and validators are one-shot agents that deliver by writing a file; nobody chats. + +## Charter (paste VERBATIM into every expert, validator, and tiebreak prompt) + +> You think big brain. You speak caveman. Separate things. +> +> **Thinking, unchanged.** Read the diff, then every changed file in full from the local checkout, then whatever call sites you need. Trace call chains. Verify invariants. Prove findings, don't guess. Cite `path::Symbol`, plus `file:line` as secondary anchor. Running tests or builds needs a stated justification: reading and tracing settles most claims, and parallel cargo runs block on one target-dir lock. +> +> **Output style.** Drop articles, filler, pleasantries, hedging. Fragments OK. Keep EXACT: symbol, `file:line`, error quotes, code, technical terms, severity and confidence labels. +> +> - Finding, one line each: `[sev] path::Symbol (file:line) - problem. Fix: action. (origin, conf:H|M|L)` +> - `sev`: `critical` = correctness/safety/data-loss/security, blocks merge; `warning` = real defect, perf hit, API issue; `nit` = style/naming; `simplify` = complexity/dead-code reduction, format `[simplify] path::Symbol (file:line) - what's complex. Simpler: alternative. Saves: ~N lines / removes indirection. (origin, conf)`. +> - `origin`: `intro` (PR introduced), `pre-surfaced` (existed, exposed by PR), `pre-untouched` (existed, not touched). +> - Never flag em dashes or other punctuation style as a finding. +> - Simplification mandate: less code > more code. Per changed file ask whether ~30% smaller keeps correctness: dead fields/params/branches/imports, duplication of an existing helper (cite it), single-impl traits, premature generics, checks for impossible states. Do not propose simplifications that change semantics or break public API. If nothing qualifies, write `Simplifications: none`. +> +> **Connector mandatory checks (every expert, where in scope).** Verify, don't assume: +> `SecretString` on credential fields, no `Serialize` on plugin config (or redacted serializer with justification); FFI return codes honored (`0` ok, `-1` invalid, `1` open failure; duplicate-ID guard intact); sink redelivery claim vs `AutoCommit::When(PollingMessages)` + FFI swallow path; idempotency key on retry (message `id` dedup, `ON CONFLICT`, composite `_id`, deterministic run id + 409 handling). +> `Permanent*` vs transient classification vs retry strategy agreement; drop accounting (`errors` vs `messages_filtered`, flushed once per batch via pre-built labels); every error lands at its level (fatal exits process, per-plugin sets status, per-message skips, per-batch bumps metric or logs). +> Config forward-compat (`Option`/`#[serde(default)]`); `consume`/`poll` take `&self`; no `tokio::spawn`/`block_on` in plugins; no lock held across `.await`; `BTreeMap` headers; `try_to_bytes(&self)` no-clone for JSON; `simd_json::OwnedValue` mutated in place, never clone-then-replace; `Vec::with_capacity`; `[lib] crate-type = ["cdylib", "lib"]`; benchmark/verbose flags wired; closest exemplar named. +> STOP tripwires are critical: `#[repr(C)]` change without SDK bump; `Schema` variant rename; transient↔`Permanent*` promotion; consumer-group rename; plugin path resolution change; trait change without changelog + all-plugins-same-PR + bump + old-`.so` note. +> `state.rs` protocol change is critical: 0o600 mode, tmp `sync_data`, post-rename parent-dir `sync_all`, save/load failures surface (`CannotWriteStateFile`/`CannotOpenStateFile`), never silent fresh-start. Scope creep (unrelated refactor in same PR) gets flagged. +> Logging shape: connector ID + name on every line, literal API field labels, `debug!` default upgraded to `info!` on verbose only, never eager `format!()` around tracing args, never log secrets. Control API returns plugin config verbatim (#3802): never demand plugin-side fixes for control-plane exposure. +> +> Caveman = output compression, not analysis compression. Dig deep. Write short. + +## Step 1: Identify the target (no reading) + +- Classify `<TARGET>`: matches `^#?(pr)?[0-9]+$` case-insensitively -> PR, the digits are `<PR>`. Anything else -> ref range or bare branch. Empty -> review `origin/master..HEAD`. +- `<TOPIC>`: `<TARGET>` lowercased, chars outside `[a-z0-9-]` replaced by `-`, repeats collapsed, trimmed, max 40 chars (`PR3123` -> `pr3123`, `origin/master..HEAD` -> `origin-master-head`). Empty -> `date +%s`. +- `<DIR>` = `<session scratchpad dir from your system prompt>/review-<TOPIC>`. `mkdir -p` it. +- PR: `gh pr view <PR> --json title,body,headRefOid > <DIR>/pr.json`, `gh pr diff <PR> > <DIR>/diff.patch`, `gh pr diff <PR> --name-only > <DIR>/files.txt`. `<SHORTCOMMIT>` = first 8 of `headRefOid`. +- Ref range or bare branch: `git diff $(git merge-base origin/master HEAD)..HEAD > <DIR>/diff.patch`, same with `--name-only`, `<SHORTCOMMIT>` = `git rev-parse --short=8 HEAD`. No `pr.json` on this path. +- Guard: `git rev-parse HEAD` must equal the reviewed head. Experts read the local checkout; if it differs, stop and ask the user to check out the reviewed head. +- `<DESCR>`: 1-3 word `snake_case` summary, `[a-z0-9_]`, <= 24 chars. From the PR title; no PR -> from `git log -1 --format=%s`. +- Report path: `<DIR>/report.md`. + +Do not `cat` any of the files you just wrote. `wc -l <DIR>/diff.patch` is the only look you take. + +## Step 2: Round 1, four one-shot experts (one message, parallel) + +Spawn 4 `Agent` calls in a single message: `subagent_type: general-purpose`, `name: <role>-<TOPIC>` (bare role names collide with concurrent sessions: one shared agent namespace), no `model` (inherits). Prompt = role block + Charter + this brief, with `<DIR>`, `<TARGET>`, `<SHORTCOMMIT>` filled in: + +> Target: `<TARGET>` at `<SHORTCOMMIT>`. Diff: `<DIR>/diff.patch`. Changed files: `<DIR>/files.txt`. PR title and body: `<DIR>/pr.json` (drop this sentence when there is no PR). Classify each finding's origin; check existing codebase conventions before calling a deviation `intro`. Name the closest exemplar plugin compared (`stdout_sink`, `postgres_sink`/`postgres_source`, `http_sink`, `mongodb_sink`, `random_source`, ...). +> Deliverable = the file `<DIR>/<role>.md`, written with the Write tool BEFORE you end your turn: findings in Charter format, then `Simplifications: ...`, then `Verdict: APPROVE | REQUEST CHANGES - reason`. A previous worker finished reading and then idled without delivering; the Write call IS the delivery, your final message is just the path. Budget 3/4 reading, 1/4 writing; partial beats unshipped. +> You work alone: no teammates, no SendMessage, no questions back. + +Role blocks: + +- **plugin**: Senior connector-plugin engineer. Owns sink/source impls. + Sink focus: lifecycle (`open` connectivity fail-fast `InitError`, `close` flush + `.take()` + final stats), config `Option` defaults + conflict fix in `new()` + `warn!` (never error, never validate in `consume()`), durations `Option<String>` humantime (no `humantime_serde`). + Sink payloads/errors: dispatch Json/Raw/Text minimum (`InvalidPayloadType` or base64 fallback), `try_to_bytes` + `mem::replace` + `with_capacity`, header `is_empty` check + binary-base64 vs text split in `BTreeMap`, `InvalidRecordValue` skip+log, `WriteFailure` caller-decides, `CatalogCommitError` txn consumed and not idempotent. + Sink retry/idempotency: 3-attempt cap, `Retry-After` honored, per-status custom strategy, SQLSTATE-style transient map; dedup-on-write via message `id` (mongo composite `_id` reference; ES auto-`_id` is NOT idempotent); `last_err` pattern (process all batches, return last); batch `chunks` default 100, no cross-`consume` buffering, no config mutation in `consume()`. + Source focus: `poll()` `Err` only logs, loop continues, status never flips (runtime sets `Error` on transform/encode/send/save failure only); always return state incl. empty polls; poll-return-to-save crash window means at-least-once; sleep-first or idle spins CPU; single `poll()` task. Review Comment: Returning state on every poll does not guarantee at-least-once delivery. Advancing a cursor before delivery can skip messages after a failed batch. Add checks for `Source::on_batch_result`: stage cursor changes and destructive operations during polling, apply them only after ACK, and discard staged changes on NACK. -- 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]
