This is an automated email from the ASF dual-hosted git repository.
numinnex pushed a change to branch kafka_proxy_auth
in repository https://gitbox.apache.org/repos/asf/iggy.git
from fb48d08e9 address review
add 362bd7829 fix(server): note root password reset in generated password
log (#4217)
add ad8951f57 chore(deps): Bump github.com/onsi/ginkgo/v2 from 2.32.1 to
2.32.2 in /bdd/go in the go group across 1 directory (#4227)
add 2ad0b1900 refactor(gateways): prepare the Kafka gateway for parallel
handler work (#4228)
add 92aee83c4 fix(server): skip NUMA memory binding on a single-node host
(#4243)
add 5d33f9d11 docs: refresh the README (#4208)
add ee226f971 refactor(server): split responses.rs by role and unify reply
framing (#4089)
add 968a23e20 chore(java): start 0.9.1-SNAPSHOT and document the 0.9.0 SDK
(#4247)
add 9a59f7857 ci(rust): run clippy on macOS so cfg(target_os) code gets
linted (#4248)
add a69b4e36f fix(connectors): meilisearch_sink URL scheme check is
case-sensitive (#4158)
add 41c6956bb fix(ci): retry flaky downloads and stop dry runs promising
new tags (#4216)
add d1cb5aee7 fix(configs): check server env vars once and refuse boot
only in debug (#4200)
add 7bd69670f refactor(partitions): make purge recovery testable with
SimStorage (#4240)
add f0bd70283 perf(message_bus): batch small replica frames per socket
read (#4224)
add da4f3188d feat(cpp): add functions related to consumers and offsets to
high level client (#4180)
add ddbe222c2 chore(deps): Bump the csharp group with 11 updates (#4254)
add 74614ed05 feat(server): warn when listener binds loopback inside a
container (#4218)
add fa8070494 fix(cluster): stop partition repair asking for an inverted
op range (#4246)
add 2cacea5b1 fix(partitions): sync checkpoints through original writers
(#4253)
add 7bae8ec84 fix(connectors): tag sink batches with the payload's schema
(#4204)
add 2894368e1 feat(gateways): map Kafka records to Iggy messages and back
(#4229)
add 3739da8be merge master
No new revisions were added by this update.
Summary of changes:
.claude/skills/connector-runtime/SKILL.md | 5 +-
.claude/skills/connector-sdk/SKILL.md | 17 +-
.claude/skills/connector-sink/SKILL.md | 2 +-
.config/nextest.toml | 20 +-
.../actions/csharp-dotnet/post-merge/action.yml | 2 +-
.github/actions/java-gradle/post-merge/action.yml | 2 +-
.../actions/python-maturin/pre-merge/action.yml | 4 +-
.github/actions/rust/pre-merge/action.yml | 9 +-
.../actions/utils/setup-node-with-cache/action.yml | 2 +-
.../actions/utils/setup-rust-with-cache/action.yml | 17 +-
.../utils/validate-third-party-licenses/action.yml | 10 +-
.github/config/components.yml | 1 +
.github/workflows/_build_python_wheels.yml | 46 +-
.github/workflows/_build_rust_artifacts.yml | 4 +-
.github/workflows/_common.yml | 9 +-
.github/workflows/_detect.yml | 2 +-
.github/workflows/_publish_rust_crates.yml | 2 +-
.github/workflows/_test.yml | 2 +-
.github/workflows/coverage-baseline.yml | 4 +-
.github/workflows/post-merge.yml | 2 +-
.github/workflows/pre-merge.yml | 54 +-
.github/workflows/publish.yml | 127 +-
Cargo.lock | 38 +-
Cargo.toml | 7 +-
README.md | 104 +-
assets/cli.png | Bin 231092 -> 151099
bytes
assets/server.png | Bin 1176428 -> 756489
bytes
assets/web_ui.png | Bin 145001 -> 136258
bytes
.../features/step_definitions/background_steps.cpp | 3 +-
bdd/go/go.mod | 2 +-
bdd/go/go.sum | 4 +-
core/binary_protocol/src/consensus/header.rs | 153 ++
core/binary_protocol/src/consensus/mod.rs | 5 +-
core/binary_protocol/src/consensus/reply_result.rs | 38 +-
core/binary_protocol/src/lib.rs | 13 +-
core/configs/src/configs_impl/file_provider.rs | 101 +-
.../configs/src/configs_impl/typed_env_provider.rs | 100 +-
core/configs/src/server_config/defaults.rs | 1 +
core/configs/src/server_config/displays.rs | 8 +-
core/configs/src/server_config/message_bus.rs | 131 +-
core/configs/src/server_config/server.rs | 36 +-
core/configs/src/server_config/validators.rs | 293 ++-
core/connectors/runtime/Cargo.toml | 4 +
core/connectors/runtime/README.md | 5 +-
core/connectors/runtime/src/benchmark.rs | 5 +-
core/connectors/runtime/src/metrics.rs | 34 +-
core/connectors/runtime/src/sink.rs | 669 +++++-
core/connectors/sdk/Cargo.toml | 2 +-
core/connectors/sdk/README.md | 15 +-
core/connectors/sdk/src/lib.rs | 334 ++-
core/connectors/sdk/src/sink.rs | 149 +-
.../sdk/src/transforms/flatbuffer_convert.rs | 23 +-
.../connectors/sdk/src/transforms/proto_convert.rs | 65 +-
core/connectors/sinks/README.md | 8 +
core/connectors/sinks/clickhouse_sink/README.md | 8 +-
core/connectors/sinks/clickhouse_sink/src/body.rs | 110 +-
core/connectors/sinks/delta_sink/Cargo.toml | 3 +
core/connectors/sinks/delta_sink/src/sink.rs | 81 +-
core/connectors/sinks/doris_sink/README.md | 3 +-
core/connectors/sinks/doris_sink/src/lib.rs | 87 +-
.../connectors/sinks/elasticsearch_sink/src/lib.rs | 80 +-
core/connectors/sinks/http_sink/README.md | 3 +-
core/connectors/sinks/http_sink/src/lib.rs | 37 +-
.../iceberg_sink/src/router/dynamic_router.rs | 58 +-
.../sinks/iceberg_sink/src/router/mod.rs | 43 +-
.../connectors/sinks/influxdb_sink/dependencies.md | 2 +-
core/connectors/sinks/meilisearch_sink/README.md | 13 +-
core/connectors/sinks/meilisearch_sink/src/lib.rs | 107 +-
core/connectors/sinks/quickwit_sink/README.md | 2 +-
core/connectors/sinks/quickwit_sink/src/lib.rs | 22 +-
core/connectors/sinks/s3_sink/README.md | 2 +-
core/connectors/sinks/s3_sink/src/formatter.rs | 28 +-
core/connectors/sinks/surrealdb_sink/README.md | 11 +-
core/connectors/sinks/surrealdb_sink/src/lib.rs | 17 +
.../sources/influxdb_source/dependencies.md | 2 +-
core/consensus/src/plane_helpers.rs | 61 +-
core/integration/Cargo.toml | 1 +
core/integration/tests/cluster/mod.rs | 1 +
.../tests/cluster/replica_read_batching.rs | 150 ++
.../integration/tests/connectors/clickhouse/mod.rs | 2 +
.../tests/connectors/clickhouse/proto_text.rs | 113 +
.../sink.toml => clickhouse/proto_text.toml} | 2 +-
.../proto_text_config}/clickhouse_sink.toml | 43 +-
core/integration/tests/connectors/doris/mod.rs | 1 +
.../meilisearch_sink.rs => doris/proto_text.rs} | 56 +-
.../{delta/sink.toml => doris/proto_text.toml} | 2 +-
.../proto_text_config/doris_sink.toml} | 35 +-
.../tests/connectors/elasticsearch/mod.rs | 1 +
.../tests/connectors/elasticsearch/proto_text.rs | 126 ++
.../sink.toml => elasticsearch/proto_text.toml} | 2 +-
.../proto_text_config/elasticsearch_sink.toml} | 23 +-
core/integration/tests/connectors/runtime/mod.rs | 1 +
.../tests/connectors/runtime/schema_tagging.rs | 134 ++
.../sink.toml => runtime/schema_tagging.toml} | 2 +-
.../stdout_sink.toml} | 15 +-
core/journal/Cargo.toml | 1 +
core/journal/src/durable_storage.rs | 192 +-
core/message_bus/src/config.rs | 8 +
core/message_bus/src/framing.rs | 95 +-
core/message_bus/src/installer/replica.rs | 12 +-
core/message_bus/src/lib.rs | 83 +
core/message_bus/src/transports/tcp.rs | 296 ++-
core/metadata/src/stm/result.rs | 13 +-
core/metadata/src/stm/stream.rs | 15 +-
core/metadata/src/stm/user.rs | 2 +-
core/partitions/Cargo.toml | 9 +-
core/partitions/src/iggy_index_writer.rs | 4 +
core/partitions/src/iggy_partition.rs | 245 +-
core/partitions/src/lib.rs | 2 +
core/partitions/src/offset_storage.rs | 247 ++-
core/partitions/src/persistence.rs | 49 +-
core/server/config.toml | 19 +
core/server/src/boot/credentials.rs | 4 +-
core/server/src/boot/mod.rs | 2 +-
core/server/src/consumer_group.rs | 2 +-
core/server/src/dispatch/authz.rs | 7 +-
core/server/src/dispatch/failure.rs | 463 ++--
core/server/src/dispatch/mod.rs | 2 +-
core/server/src/dispatch/partition.rs | 6 +-
core/server/src/dispatch/reads.rs | 7 +-
core/server/src/dispatch/session_ops.rs | 2 +-
core/server/src/dispatch/submit.rs | 4 +-
core/server/src/http/handlers.rs | 5 +-
core/server/src/http/reads.rs | 26 +-
core/server/src/http/reply.rs | 9 +-
core/server/src/http/submit.rs | 2 +-
core/server/src/http/wire.rs | 2 +-
core/server/src/lib.rs | 8 +-
core/server/src/namespace.rs | 362 +++
core/server/src/offset_recovery.rs | 179 +-
core/server/src/partition_helpers.rs | 105 +-
core/server/src/partition_reconciler.rs | 67 +
core/server/src/reply_frame.rs | 875 ++++++++
core/server/src/responses.rs | 1343 +----------
core/server/src/session_manager.rs | 2 +-
core/server/src/shard_allocator.rs | 44 +-
core/server/src/sysinfo_probe.rs | 151 ++
core/shard/src/lib.rs | 32 +-
core/shard/src/metrics.rs | 27 +
core/simulator/src/storage.rs | 10 +-
core/simulator/src/storage/purge.rs | 469 ++++
core/simulator/src/storage/tests.rs | 95 +-
.../Iggy_SDK.Examples.Basic.Consumer.csproj | 8 +-
.../Iggy_SDK.Examples.Basic.Producer.csproj | 10 +-
...ggy_SDK.Examples.GettingStarted.Consumer.csproj | 4 +-
...ggy_SDK.Examples.GettingStarted.Producer.csproj | 4 +-
.../Iggy_SDK.Examples.Shared.csproj | 2 +-
...gy_SDK.Examples.MessageEnvelope.Consumer.csproj | 10 +-
...gy_SDK.Examples.MessageEnvelope.Producer.csproj | 10 +-
...ggy_SDK.Examples.MessageHeaders.Consumer.csproj | 4 +-
...ggy_SDK.Examples.MessageHeaders.Producer.csproj | 4 +-
.../Iggy_SDK.Examples.NewSdk.Consumer.csproj | 10 +-
.../Iggy_SDK.Examples.NewSdk.Producer.csproj | 10 +-
.../Iggy_SDK.Examples.TcpTls.Consumer.csproj | 4 +-
.../Iggy_SDK.Examples.TcpTls.Producer.csproj | 4 +-
foreign/cpp/.clang-tidy | 49 +
foreign/cpp/BUILD.bazel | 1 +
foreign/cpp/MODULE.bazel | 1 +
foreign/cpp/include/iggy.hpp | 668 +++++-
foreign/cpp/src/client.cpp | 97 +
foreign/cpp/src/type_conversions.cpp | 68 +-
foreign/cpp/tests/e2e/client.cpp | 2 -
foreign/cpp/tests/e2e/consumer_group.cpp | 2337 ++++++++++++++------
foreign/cpp/tests/unit/unit_tests.cpp | 35 +-
foreign/csharp/Directory.Packages.props | 12 +-
foreign/java/README.md | 32 +-
foreign/java/gradle.properties | 2 +-
gateways/kafka/Cargo.toml | 16 +-
gateways/kafka/README.md | 1 +
gateways/kafka/docs/BRIDGE_MAPPING.md | 195 +-
gateways/kafka/src/bridge/config.rs | 4 +-
gateways/kafka/src/bridge/iggy_bridge.rs | 529 -----
.../kafka/src/bridge/iggy_bridge/fetch.rs | 11 +-
gateways/kafka/src/bridge/iggy_bridge/mod.rs | 187 ++
gateways/kafka/src/bridge/iggy_bridge/offsets.rs | 176 ++
.../kafka/src/bridge/iggy_bridge/produce.rs | 11 +-
gateways/kafka/src/bridge/iggy_bridge/topics.rs | 220 ++
gateways/kafka/src/lib.rs | 1 +
gateways/kafka/src/main.rs | 92 +-
gateways/kafka/src/protocol/api.rs | 706 ++----
.../kafka/src/protocol/handlers/api_versions.rs | 124 ++
.../kafka/src/protocol/handlers/create_topics.rs | 130 ++
gateways/kafka/src/protocol/handlers/fetch.rs | 118 +
.../kafka/src/protocol/handlers/list_offsets.rs | 115 +
gateways/kafka/src/protocol/handlers/metadata.rs | 171 ++
gateways/kafka/src/protocol/handlers/mod.rs | 179 ++
gateways/kafka/src/protocol/handlers/produce.rs | 150 ++
gateways/kafka/src/protocol/mod.rs | 2 +-
gateways/kafka/src/protocol/responses.rs | 372 ----
gateways/kafka/src/protocol/sasl.rs | 63 +-
gateways/kafka/src/records.rs | 1830 +++++++++++++++
gateways/kafka/src/server.rs | 57 +-
gateways/kafka/tests/api_handler_tests.rs | 200 +-
.../kafka/tests/bridge_iggy_integration_tests.rs | 398 +---
gateways/kafka/tests/broker_advertise_tests.rs | 5 +-
gateways/kafka/tests/common/iggy_server.rs | 406 ++++
gateways/kafka/tests/golden_wire_fixtures_tests.rs | 15 +-
gateways/kafka/tests/listener_robustness_tests.rs | 1 +
gateways/kafka/tests/response_negative_tests.rs | 42 +-
gateways/kafka/tests/server_e2e_tests.rs | 2 +-
gateways/kafka/tests/version_firewall_tests.rs | 148 +-
licenserc.toml | 1 +
scripts/prepare-release.sh | 1 +
203 files changed, 14546 insertions(+), 5308 deletions(-)
create mode 100644 core/integration/tests/cluster/replica_read_batching.rs
create mode 100644 core/integration/tests/connectors/clickhouse/proto_text.rs
copy core/integration/tests/connectors/{delta/sink.toml =>
clickhouse/proto_text.toml} (93%)
copy core/{connectors/runtime/example_config/connectors =>
integration/tests/connectors/clickhouse/proto_text_config}/clickhouse_sink.toml
(54%)
copy core/integration/tests/connectors/{meilisearch/meilisearch_sink.rs =>
doris/proto_text.rs} (55%)
copy core/integration/tests/connectors/{delta/sink.toml =>
doris/proto_text.toml} (93%)
copy
core/integration/tests/connectors/{runtime/sink_transform_error_config/stdout.toml
=> doris/proto_text_config/doris_sink.toml} (60%)
create mode 100644
core/integration/tests/connectors/elasticsearch/proto_text.rs
copy core/integration/tests/connectors/{delta/sink.toml =>
elasticsearch/proto_text.toml} (92%)
copy core/{connectors/sinks/elasticsearch_sink/config.toml =>
integration/tests/connectors/elasticsearch/proto_text_config/elasticsearch_sink.toml}
(60%)
create mode 100644 core/integration/tests/connectors/runtime/schema_tagging.rs
copy core/integration/tests/connectors/{delta/sink.toml =>
runtime/schema_tagging.toml} (92%)
copy
core/integration/tests/connectors/runtime/{sink_invalid_config/stdout_valid.toml
=> schema_tagging_config/stdout_sink.toml} (76%)
create mode 100644 core/server/src/namespace.rs
create mode 100644 core/server/src/reply_frame.rs
create mode 100644 core/server/src/sysinfo_probe.rs
create mode 100644 core/simulator/src/storage/purge.rs
create mode 100644 foreign/cpp/.clang-tidy
delete mode 100644 gateways/kafka/src/bridge/iggy_bridge.rs
copy bdd/rust/tests/steps/mod.rs =>
gateways/kafka/src/bridge/iggy_bridge/fetch.rs (85%)
create mode 100644 gateways/kafka/src/bridge/iggy_bridge/mod.rs
create mode 100644 gateways/kafka/src/bridge/iggy_bridge/offsets.rs
copy bdd/rust/tests/steps/mod.rs =>
gateways/kafka/src/bridge/iggy_bridge/produce.rs (85%)
create mode 100644 gateways/kafka/src/bridge/iggy_bridge/topics.rs
create mode 100644 gateways/kafka/src/protocol/handlers/api_versions.rs
create mode 100644 gateways/kafka/src/protocol/handlers/create_topics.rs
create mode 100644 gateways/kafka/src/protocol/handlers/fetch.rs
create mode 100644 gateways/kafka/src/protocol/handlers/list_offsets.rs
create mode 100644 gateways/kafka/src/protocol/handlers/metadata.rs
create mode 100644 gateways/kafka/src/protocol/handlers/mod.rs
create mode 100644 gateways/kafka/src/protocol/handlers/produce.rs
delete mode 100644 gateways/kafka/src/protocol/responses.rs
create mode 100644 gateways/kafka/src/records.rs
create mode 100644 gateways/kafka/tests/common/iggy_server.rs