This is an automated email from the ASF dual-hosted git repository.
slbotbm pushed a change to branch connectors-team-review-skill
in repository https://gitbox.apache.org/repos/asf/iggy.git
discard e0bf857e6 add connector-review skill
omit f3d9688d0 Merge branch 'master' into feat/connectors-opensearch-sink
omit 751caed3e Merge branch 'master' into feat/connectors-opensearch-sink
omit 058a0fed9 Merge branch 'master' into feat/connectors-opensearch-sink
omit 027dea8f3 Merge branch 'master' into feat/connectors-opensearch-sink
omit 32373e9ff chore(deps): sync Cargo.lock after merging master
omit 07cbe3d80 Merge branch 'master' into feat/connectors-opensearch-sink
omit aedfdb1c2 docs(connectors): clarify OpenSearch sink retry backoff scope
omit ee210faeb fix(connectors): report first failing chunk in OpenSearch
batches
omit a1ffd8d77 fix(connectors): move simd-json to dev-dependencies in
OpenSearch sink
omit 7fa604f87 fix(connectors): redact OpenSearch sink URL credentials in
Debug
omit 7427e8324 chore(connectors): bump OpenSearch sink to 0.5.0-edge.4
omit eef41feb3 fix(connectors): reject malformed items in OpenSearch bulk
fast path
omit 8fbad8fef fix(connectors): retry OpenSearch test container
create-or-attach race
omit 5fbcf8471 fix(connectors): make OpenSearch sink URL scheme detection
case-insensitive
omit d9946e2f9 fix(connectors): correct OpenSearch sink retry backoff and
jitter cap
omit 7adad928f perf(connectors): stop cloning every raw payload in
OpenSearch sink
omit 91ba9b0cd fix(connectors): stop OpenSearch sink Debug from leaking the
password
omit 2817a4e43 fix(connectors): don't count unparsable OpenSearch bulk
responses as success
omit fa4cc2448 Merge branch 'master' into feat/connectors-opensearch-sink
omit 6f078e2ed Merge branch 'master' into feat/connectors-opensearch-sink
omit d92ea8058 feat(connectors): add OpenSearch sink connector
add add0bd967 feat(python): add QuicConfig transport configuration (#3991)
add 03ca396bb feat(python): add stream listing, update, delete, and purge
(#3701)
add b87b6a6f7 fix(java): size wire strings by UTF-8 byte length, not char
count (#4068)
add ab734f17f fix(connectors): validate connector key before it becomes a
path (#4083)
add 06a12e306 chore: add Iggy banner to server startup (#4085)
add bd1a873fe feat(python): add HttpConfig transport configuration (#3992)
add 1004eee2d ci(connectors): run a plugin's integration suite when it
changes (#4077)
add be3663805 feat(python): support message partitioning strategies (#3927)
add 3c1a8ccf6 feat(python): add WebSocketConfig transport configuration
(#4000)
add 6c07954a0 chore(deps): Bump the csharp group with 1 update (#4087)
add a71ed7420 fix(metadata): roll deleted partitions and topics out of
parent stats (#4046)
add 2e3823196 feat(connectors): add RabbitMQ sink (#3973)
add a2a4a2cb9 fix(consensus): keep a replica off a hole in its committed
prefix (#4073)
add 8dde7e7f6 feat(java): let TCP clients run one I/O thread or share a
group (#4096)
add 422810f82 chore(deps): Bump the java group across 3 directories with 1
update (#4098)
add 6d894501b chore(deps): Bump the github-actions group with 3 updates
(#4099)
add 1290e4fff chore(deps): remove Microsoft.SourceLink.GitHub package
reference (#4113)
add 7a77f4592 fix(connectors): use DLL_EXTENSION for plugin path suffix
(#4114)
add 62f945a34 chore: add justinmclean to .asf.yaml collaborators (#4115)
add 46a907596 fix(connectors): bring quickwit_sink up to convention (#3523)
add df4c75d34 feat(python): expose partition management (#4017)
add 3f0646fce chore(ci): add team-review-slim, cheaper team-review and
plain english (#4095)
new 420a66fc0 add connector-review skill
This update added new revisions after undoing existing revisions.
That is to say, some revisions that were in the old version of the
branch are not in the new version. This situation occurs
when a user --force pushes a change and generates a repository
containing something like this:
* -- * -- B -- O -- O -- O (e0bf857e6)
\
N -- N -- N refs/heads/connectors-team-review-skill (420a66fc0)
You should already have received notification emails for all of the O
revisions, and so the following emails describe only the N revisions
from the common base, B.
Any revisions marked "omit" are not gone; other references still
refer to them. Any revisions marked "discard" are gone forever.
The 1 revisions listed above as "new" are entirely new to this
repository and will be described in separate emails. The revisions
listed as "add" were already present in the repository and have only
been added to this reference.
Summary of changes:
.asf.yaml | 3 +
.claude/skills/team-review-slim/SKILL.md | 230 ++
.claude/skills/team-review/SKILL.md | 184 +-
.config/nextest.toml | 14 -
.../actions/python-maturin/pre-merge/action.yml | 16 +-
.github/actions/rust/pre-merge/action.yml | 51 +-
.github/workflows/_build_rust_artifacts.yml | 2 +-
.github/workflows/_common.yml | 4 +-
.github/workflows/coverage-baseline.yml | 19 +-
.github/workflows/edge-release.yml | 3 +-
.github/workflows/post-merge.yml | 2 +-
Cargo.lock | 1535 ++++++----
Cargo.toml | 10 +-
bdd/java/build.gradle.kts | 2 +-
bdd/python/pylock.toml | 346 +--
bdd/python/uv.lock | 2 -
.../websocket_config/websocket_client_config.rs | 75 +-
core/common/src/types/streaming_stats.rs | 367 ++-
core/connectors/README.md | 2 +-
core/connectors/runtime/README.md | 15 +
.../example_config/connectors/quickwit_sink.toml | 22 +-
.../{opensearch_sink.toml => rabbitmq_sink.toml} | 35 +-
core/connectors/runtime/src/api/error.rs | 15 +-
.../mod.rs => connectors/runtime/src/api/key.rs} | 35 +-
core/connectors/runtime/src/api/mod.rs | 1 +
core/connectors/runtime/src/api/sink.rs | 44 +-
core/connectors/runtime/src/api/source.rs | 44 +-
core/connectors/runtime/src/configs/connectors.rs | 182 +-
.../src/configs/connectors/http_provider.rs | 12 +-
.../src/configs/connectors/local_provider.rs | 184 +-
core/connectors/runtime/src/error.rs | 6 +
core/connectors/runtime/src/main.rs | 42 +-
core/connectors/sinks/README.md | 2 +-
core/connectors/sinks/opensearch_sink/README.md | 205 --
core/connectors/sinks/opensearch_sink/src/lib.rs | 3085 --------------------
core/connectors/sinks/quickwit_sink/Cargo.toml | 10 +-
core/connectors/sinks/quickwit_sink/README.md | 64 +-
core/connectors/sinks/quickwit_sink/config.toml | 9 +
core/connectors/sinks/quickwit_sink/src/lib.rs | 957 +++++-
.../{opensearch_sink => rabbitmq_sink}/Cargo.toml | 23 +-
core/connectors/sinks/rabbitmq_sink/README.md | 50 +
.../{opensearch_sink => rabbitmq_sink}/config.toml | 30 +-
core/connectors/sinks/rabbitmq_sink/src/lib.rs | 783 +++++
core/consensus/src/impls.rs | 440 ++-
core/consensus/src/plane_helpers.rs | 169 +-
core/integration/Cargo.toml | 5 +
core/integration/tests/connectors/api/endpoints.rs | 104 +
.../api/{config.toml => key_validation.toml} | 4 +-
core/integration/tests/connectors/fixtures/mod.rs | 13 +-
.../connectors/fixtures/opensearch/container.rs | 238 --
.../connectors/fixtures/opensearch/failure.rs | 129 -
.../tests/connectors/fixtures/opensearch/mod.rs | 24 -
.../tests/connectors/fixtures/opensearch/sink.rs | 161 -
.../connectors/fixtures/quickwit/container.rs | 86 +-
.../tests/connectors/fixtures/quickwit/mod.rs | 5 +-
.../connectors/fixtures/rabbitmq/container.rs | 327 +++
.../fixtures/{clickhouse => rabbitmq}/mod.rs | 5 +-
.../tests/connectors/fixtures/rabbitmq/sink.rs | 273 ++
core/integration/tests/connectors/mod.rs | 2 +-
.../connectors/opensearch/failure_states.toml | 20 -
.../failure_states/opensearch_healthy.toml | 43 -
.../opensearch_mapping_conflict.toml | 54 -
.../failure_states/opensearch_missing_index.toml | 42 -
.../integration/tests/connectors/opensearch/mod.rs | 19 -
.../tests/connectors/opensearch/opensearch_sink.rs | 325 ---
.../opensearch/opensearch_sink_failures.rs | 401 ---
.../tests/connectors/opensearch/sink.toml | 20 -
.../tests/connectors/quickwit/quickwit_sink.rs | 128 +-
.../tests/connectors/{quickwit => rabbitmq}/mod.rs | 2 +-
.../tests/connectors/rabbitmq/rabbitmq_sink.rs | 429 +++
.../tests/connectors/{delta => rabbitmq}/sink.toml | 2 +-
.../data_integrity/verify_after_server_restart.rs | 52 +
core/integration/tests/server/general.rs | 14 +-
.../scenarios/delete_stats_rollback_scenario.rs | 225 ++
core/integration/tests/server/scenarios/mod.rs | 137 +-
.../server/scenarios/purge_delete_scenario.rs | 26 +-
.../scenarios/stream_size_validation_scenario.rs | 140 +-
core/metadata/src/impls/metadata.rs | 18 +
core/metadata/src/stm/stream.rs | 899 +++++-
core/partitions/src/iggy_partition.rs | 309 +-
core/sdk/src/prelude.rs | 6 +-
core/server/src/banner.rs | 80 +
core/server/src/boot/mod.rs | 16 -
core/server/src/boot/recovery.rs | 2 +-
core/server/src/http/metrics.rs | 71 +-
core/server/src/http/reads.rs | 10 +-
core/server/src/main.rs | 2 +
core/server/src/partition_reconciler.rs | 331 ++-
core/server/src/responses.rs | 28 +-
core/shard/src/lib.rs | 997 ++++++-
core/shard/src/router.rs | 9 +
core/simulator/src/bin/workload-fuzz.rs | 44 +-
core/simulator/src/lib.rs | 158 +-
core/simulator/src/workload/invariants.rs | 164 +-
core/simulator/src/workload/oracle.rs | 14 +
core/simulator/src/workload/state_checker.rs | 170 +-
examples/java/build.gradle.kts | 2 +-
.../apache/iggy/examples/async/AsyncConsumer.java | 3 +-
examples/python/README.md | 2 +-
examples/python/basic/producer.py | 8 +-
examples/python/getting-started/consumer.py | 42 +-
examples/python/getting-started/producer.py | 42 +-
examples/python/pylock.toml | 81 +-
examples/python/uv.lock | 2 -
foreign/csharp/Directory.Build.props | 4 -
foreign/csharp/Directory.Packages.props | 3 +-
foreign/java/README.md | 30 +
foreign/java/gradle/libs.versions.toml | 2 +-
.../apache/iggy/client/async/MessagesClient.java | 2 +-
.../iggy/client/async/tcp/AsyncIggyTcpClient.java | 32 +-
.../async/tcp/AsyncIggyTcpClientBuilder.java | 79 +-
.../iggy/client/async/tcp/AsyncTcpConnection.java | 88 +-
.../client/async/tcp/ConsumerGroupsTcpClient.java | 10 +-
.../async/tcp/PersonalAccessTokensTcpClient.java | 14 +-
.../iggy/client/async/tcp/StreamsTcpClient.java | 22 +-
.../iggy/client/async/tcp/TopicsTcpClient.java | 4 +-
.../iggy/client/async/tcp/UsersTcpClient.java | 46 +-
.../iggy/client/async/tcp/vsr/VsrLoginCodec.java | 58 +-
.../client/blocking/tcp/IggyTcpClientBuilder.java | 31 +
.../org/apache/iggy/identifier/Identifier.java | 32 +-
.../main/java/org/apache/iggy/message/Message.java | 3 +-
.../java/org/apache/iggy/message/Partitioning.java | 42 +-
.../org/apache/iggy/serde/BytesSerializer.java | 37 +-
.../async/tcp/AsyncIggyTcpClientBuilderTest.java | 94 +
.../tcp/AsyncIggyTcpClientEventLoopGroupTest.java | 464 +++
.../tcp/AsyncTcpConnectionConcurrencyTest.java | 4 +
.../client/async/tcp/vsr/VsrLoginCodecTest.java | 159 +
.../client/blocking/tcp/BytesSerializerTest.java | 109 +-
.../client/blocking/tcp/MessagesTcpClientTest.java | 35 +
.../client/blocking/tcp/StreamTcpClientTest.java | 24 +
.../org/apache/iggy/identifier/IdentifierTest.java | 49 +
.../java/org/apache/iggy/message/MessageTest.java | 9 +
.../org/apache/iggy/message/PartitioningTest.java | 60 +
foreign/python/README.md | 38 +-
foreign/python/apache_iggy.pyi | 727 ++++-
foreign/python/docker-compose.test.yml | 2 +
foreign/python/pylock.toml | 37 +-
foreign/python/pyproject.toml | 2 -
foreign/python/src/client.rs | 468 ++-
foreign/python/src/config.rs | 1162 +++++++-
foreign/python/src/duration.rs | 33 +
foreign/python/src/lib.rs | 17 +-
foreign/python/src/options.rs | 2 +
foreign/python/src/partitioning.rs | 105 +
foreign/python/src/stream.rs | 78 +-
foreign/python/src/topic.rs | 19 +
foreign/python/tests/conftest.py | 27 +-
foreign/python/tests/test_client_config.py | 4 +-
foreign/python/tests/test_connectivity.py | 2 +-
foreign/python/tests/test_http_config.py | 363 +++
foreign/python/tests/test_message_operations.py | 137 +
foreign/python/tests/test_partition.py | 228 ++
foreign/python/tests/test_quic_config.py | 494 ++++
foreign/python/tests/test_stream.py | 444 ++-
foreign/python/tests/test_websocket_config.py | 522 ++++
foreign/python/tests/utils.py | 55 +-
foreign/python/uv.lock | 26 -
licenserc.toml | 1 +
scripts/bump-version.sh | 2 +-
159 files changed, 15975 insertions(+), 6862 deletions(-)
create mode 100644 .claude/skills/team-review-slim/SKILL.md
rename core/connectors/runtime/example_config/connectors/{opensearch_sink.toml
=> rabbitmq_sink.toml} (68%)
copy core/{integration/tests/connectors/postgres/mod.rs =>
connectors/runtime/src/api/key.rs} (50%)
delete mode 100644 core/connectors/sinks/opensearch_sink/README.md
delete mode 100644 core/connectors/sinks/opensearch_sink/src/lib.rs
rename core/connectors/sinks/{opensearch_sink => rabbitmq_sink}/Cargo.toml
(78%)
create mode 100644 core/connectors/sinks/rabbitmq_sink/README.md
rename core/connectors/sinks/{opensearch_sink => rabbitmq_sink}/config.toml
(71%)
create mode 100644 core/connectors/sinks/rabbitmq_sink/src/lib.rs
copy core/integration/tests/connectors/api/{config.toml =>
key_validation.toml} (83%)
delete mode 100644
core/integration/tests/connectors/fixtures/opensearch/container.rs
delete mode 100644
core/integration/tests/connectors/fixtures/opensearch/failure.rs
delete mode 100644 core/integration/tests/connectors/fixtures/opensearch/mod.rs
delete mode 100644
core/integration/tests/connectors/fixtures/opensearch/sink.rs
create mode 100644
core/integration/tests/connectors/fixtures/rabbitmq/container.rs
copy core/integration/tests/connectors/fixtures/{clickhouse =>
rabbitmq}/mod.rs (77%)
create mode 100644 core/integration/tests/connectors/fixtures/rabbitmq/sink.rs
delete mode 100644
core/integration/tests/connectors/opensearch/failure_states.toml
delete mode 100644
core/integration/tests/connectors/opensearch/failure_states/opensearch_healthy.toml
delete mode 100644
core/integration/tests/connectors/opensearch/failure_states/opensearch_mapping_conflict.toml
delete mode 100644
core/integration/tests/connectors/opensearch/failure_states/opensearch_missing_index.toml
delete mode 100644 core/integration/tests/connectors/opensearch/mod.rs
delete mode 100644
core/integration/tests/connectors/opensearch/opensearch_sink.rs
delete mode 100644
core/integration/tests/connectors/opensearch/opensearch_sink_failures.rs
delete mode 100644 core/integration/tests/connectors/opensearch/sink.toml
copy core/integration/tests/connectors/{quickwit => rabbitmq}/mod.rs (97%)
create mode 100644 core/integration/tests/connectors/rabbitmq/rabbitmq_sink.rs
copy core/integration/tests/connectors/{delta => rabbitmq}/sink.toml (94%)
create mode 100644
core/integration/tests/server/scenarios/delete_stats_rollback_scenario.rs
create mode 100644 core/server/src/banner.rs
create mode 100644
foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientEventLoopGroupTest.java
create mode 100644
foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/vsr/VsrLoginCodecTest.java
create mode 100644 foreign/python/src/partitioning.rs
create mode 100644 foreign/python/tests/test_http_config.py
create mode 100644 foreign/python/tests/test_partition.py
create mode 100644 foreign/python/tests/test_quic_config.py
create mode 100644 foreign/python/tests/test_websocket_config.py