This is an automated email from the ASF dual-hosted git repository. numinnex pushed a commit to branch kafka_proxy_acls in repository https://gitbox.apache.org/repos/asf/iggy.git
commit 5c8bac1ce3ce3409b2563f593c52437845d747f0 Merge: df3a07d7a 11dd94c62 Author: Grzegorz Koszyk <[email protected]> AuthorDate: Mon Sep 28 12:00:29 2026 +0200 merge master .config/nextest.toml | 10 +- .github/actions/swift/pre-merge/action.yml | 75 ++ .github/actions/utils/setup-swift/action.yml | 56 ++ .github/config/components.yml | 14 + .github/workflows/_detect.yml | 11 +- .github/workflows/_test.yml | 8 + .github/workflows/pre-merge.yml | 21 +- .gitignore | 8 + AGENTS.md | 3 +- CONTRIBUTING.md | 9 + bdd/docker-compose.cluster.yml | 1 + bdd/go/go.mod | 4 +- bdd/go/go.sum | 8 +- core/binary_protocol/src/consensus/command.rs | 11 +- .../src/consensus/consumer_session.rs | 218 +++++ core/binary_protocol/src/consensus/header.rs | 6 + core/binary_protocol/src/consensus/mod.rs | 4 + core/binary_protocol/src/lib.rs | 22 +- core/cli/src/args/consumer_group.rs | 6 +- core/configs/src/common/defaults.rs | 10 + core/configs/src/common/displays.rs | 6 +- core/configs/src/common/server.rs | 7 + core/configs/src/server_config/validators.rs | 98 +- core/connectors/sinks/redshift_sink/README.md | 2 +- .../test_consumer_group_create_command.rs | 34 +- .../tests/cluster/client_table_restart.rs | 478 ++++++++- .../tests/cluster/partition_primary_routing.rs | 29 +- .../tests/connectors/fixtures/delta/fixture.rs | 158 +-- .../integration/tests/connectors/fixtures/floci.rs | 159 +++ .../tests/connectors/fixtures/iceberg/container.rs | 166 +--- core/integration/tests/connectors/fixtures/mod.rs | 1 + .../connectors/fixtures/redshift/container.rs | 73 +- .../tests/connectors/fixtures/redshift/mod.rs | 2 +- .../tests/connectors/fixtures/redshift/sink.rs | 203 +--- .../tests/connectors/fixtures/s3/fixture.rs | 145 +-- core/integration/tests/connectors/s3/s3_sink.rs | 33 +- .../stale_client_consumer_group_scenario.rs | 262 ++++- core/metadata/src/impls/metadata.rs | 427 +++++++- core/metadata/src/impls/recovery.rs | 53 +- core/metadata/src/stm/consumer_group.rs | 242 ++++- core/metadata/src/stm/mod.rs | 10 + core/metadata/src/stm/stream.rs | 81 +- core/server/config.toml | 19 +- core/server/src/boot/mod.rs | 47 +- core/server/src/boot/recovery.rs | 4 + core/server/src/boot/threads.rs | 12 +- core/server/src/consumer_group.rs | 5 + core/server/src/consumer_group/lease.rs | 160 +++ core/server/src/consumer_group/liveness.rs | 1034 ++++++++++++++++++++ core/server/src/dispatch/mod.rs | 14 +- core/server/src/dispatch/session_ops.rs | 2 +- core/server/src/dispatch/submit.rs | 24 +- core/server/src/dispatch/test_support.rs | 8 + core/server/src/partition_reconciler.rs | 1 + core/server/src/session_manager.rs | 70 +- core/server/src/shell.rs | 4 + core/server_common/src/consensus_message.rs | 16 +- core/server_common/src/sharding/mod.rs | 10 + core/shard/src/lib.rs | 75 +- core/shard/src/poll/completion_tests.rs | 1 + core/simulator/README.md | 1 + core/simulator/src/client.rs | 18 +- core/simulator/src/lib.rs | 274 +++++- core/simulator/src/packet.rs | 7 +- core/simulator/src/replica.rs | 12 +- foreign/cpp/MODULE.bazel | 2 +- foreign/cpp/MODULE.bazel.lock | 7 +- foreign/go/client/tcp/tcp_offset_management.go | 9 + foreign/go/tests/e2e_test.go | 49 + foreign/swift/.swift-format | 22 + .../redshift/mod.rs => foreign/swift/Package.swift | 31 +- foreign/swift/README.md | 87 ++ foreign/swift/Sources/Iggy/Errors/IggyError.swift | 78 ++ .../swift/Sources/Iggy/Errors/IggyErrorCode.swift | 515 ++++++++++ .../swift/Sources/Iggy/Iggy.swift | 17 +- .../Sources/Iggy/Utilities/UInt128Value.swift | 156 +++ foreign/swift/Sources/Iggy/Wire/ByteCodec.swift | 218 +++++ foreign/swift/Tests/IggyTests/ByteCodecTests.swift | 120 +++ foreign/swift/Tests/IggyTests/ErrorCodeTests.swift | 55 ++ gateways/kafka/README.md | 120 ++- gateways/kafka/docs/BRIDGE_MAPPING.md | 70 +- gateways/kafka/docs/MANUAL_TESTING.md | 30 +- gateways/kafka/docs/SCOPE.md | 122 ++- gateways/kafka/docs/TEST_SUITE.md | 9 +- gateways/kafka/src/bridge/config.rs | 47 + gateways/kafka/src/bridge/error.rs | 139 ++- gateways/kafka/src/bridge/iggy_bridge/mod.rs | 107 +- gateways/kafka/src/bridge/iggy_bridge/offsets.rs | 37 +- gateways/kafka/src/bridge/iggy_bridge/produce.rs | 183 +++- gateways/kafka/src/bridge/iggy_bridge/topics.rs | 330 ++++++- gateways/kafka/src/bridge/mod.rs | 8 +- gateways/kafka/src/bridge/topic_map.rs | 6 + gateways/kafka/src/protocol/api.rs | 81 +- .../kafka/src/protocol/handlers/create_topics.rs | 711 +++++++++++++- gateways/kafka/src/protocol/handlers/metadata.rs | 544 +++++++++- gateways/kafka/src/protocol/handlers/produce.rs | 895 ++++++++++++++++- gateways/kafka/src/records.rs | 976 +++++++++++++----- gateways/kafka/tests/api_handler_tests.rs | 23 +- .../kafka/tests/bridge_iggy_integration_tests.rs | 272 ++++- gateways/kafka/tests/common/iggy_server.rs | 3 +- .../kafka/tests/create_topics_real_bridge_tests.rs | 536 ++++++++++ gateways/kafka/tests/listener_robustness_tests.rs | 24 +- gateways/kafka/tests/metadata_real_bridge_tests.rs | 408 ++++++++ gateways/kafka/tests/produce_real_bridge_tests.rs | 750 ++++++++++++++ gateways/kafka/tests/response_negative_tests.rs | 9 +- gateways/kafka/tests/server_e2e_tests.rs | 26 +- gateways/kafka/tests/version_firewall_tests.rs | 19 +- licenserc.toml | 7 + 108 files changed, 11356 insertions(+), 1494 deletions(-) diff --cc .config/nextest.toml index bf6c1d4ad,c8428a4e4..4cb32eb3c --- a/.config/nextest.toml +++ b/.config/nextest.toml @@@ -55,22 -55,16 +55,28 @@@ test-group = "elasticsearch [test-groups.kafka_bridge] max-threads = 4 + # Every gateway suite that spawns a server belongs here. A new one left out does not fail, it + # just runs unbounded alongside these, which is the failure nobody notices until CI is slow. [[profile.default.overrides]] - filter = 'binary_id(iggy-gateway-kafka::bridge_iggy_integration_tests) or binary_id(iggy-gateway-kafka::list_offsets_real_bridge_tests)' + filter = ''' + binary_id(iggy-gateway-kafka::bridge_iggy_integration_tests) + + binary_id(iggy-gateway-kafka::list_offsets_real_bridge_tests) + + binary_id(iggy-gateway-kafka::produce_real_bridge_tests) + ''' test-group = "kafka_bridge" +# The real-client suite spawns its own server too, so it needs serializing for the same reason, but +# in its own group rather than queued behind the bridge tests: sharing one group would put six +# container-driven tests behind twenty-one unrelated ones, serially, for no isolation gained. Both +# groups cap their servers' shard pools themselves, so running the two alongside each other is +# bounded. +[test-groups.kafka_client_e2e] +max-threads = 1 + +[[profile.default.overrides]] +filter = 'binary_id(iggy-gateway-kafka::kafka_client_e2e_tests)' +test-group = "kafka_client_e2e" + [profile.default] slow-timeout = { period = "60s", terminate-after = 5 } diff --cc gateways/kafka/docs/TEST_SUITE.md index 38c89acb9,101932120..da897e3d0 --- a/gateways/kafka/docs/TEST_SUITE.md +++ b/gateways/kafka/docs/TEST_SUITE.md @@@ -64,11 -64,11 +64,12 @@@ file under `tests/` anymore | [`server_e2e_tests.rs`](../tests/server_e2e_tests.rs) | Full `KafkaGateway` TCP round-trips | Partial | | [`listener_robustness_tests.rs`](../tests/listener_robustness_tests.rs) | TCP listener robustness — framing, pipelining, concurrency, connection limits | No | | [`sasl_tests.rs`](../tests/sasl_tests.rs) | SASL/PLAIN over a socket — full handshake, every refusal path, and the disabled default. Drives a stub verifier implementing `SaslAuthenticator`, so no Iggy server is needed | No | +| [`kafka_client_e2e_tests.rs`](../tests/kafka_client_e2e_tests.rs) | **Real Kafka clients** against the whole stack: a spawned `iggy-server`, the gateway in-process with a real authenticator, and kcat / the Java tools from containers. The only suite that can catch a client-compatibility bug, since every other one hand-builds frames | No, but needs Docker and a built `iggy-server` | | [`bridge_iggy_integration_tests.rs`](../tests/bridge_iggy_integration_tests.rs) | `IggyBridge` against a real, spawned `iggy-server` — provisioning idempotency, high watermark, credential/connection edge cases | No (needs the `iggy-server` binary - see Prerequisites) | + | [`produce_real_bridge_tests.rs`](../tests/produce_real_bridge_tests.rs) | Produce (key 0) through the whole handler against a real, spawned `iggy-server` — records go in as Kafka wire bytes and come back through the Iggy SDK, plus one error code per partition | No (needs the `iggy-server` binary - see Prerequisites) | `tests/common/` holds shared helpers (`codec.rs`, `fixtures.rs`, `scope.rs`, `server.rs`, - `tcp.rs`, `wire.rs`), compiled per test binary via `#[path]`, not a test binary itself. `codec.rs` + `iggy_server.rs`, `tcp.rs`, `wire.rs`), compiled per test binary via `#[path]`, not a test binary itself. `codec.rs` is test-only primitive encode/decode scaffolding for hand-building legacy/adversarial wire shapes `kafka_protocol`'s spec-correct encoder cannot produce - it is not the gateway's production codec. diff --cc gateways/kafka/src/protocol/api.rs index 86690079c,b92274f92..0d9217e81 --- a/gateways/kafka/src/protocol/api.rs +++ b/gateways/kafka/src/protocol/api.rs @@@ -18,10 -18,9 +18,11 @@@ use std::sync::Arc; use bytes::Bytes; - + use kafka_protocol::error::ResponseError; -use kafka_protocol::messages::{SaslAuthenticateRequest, SaslHandshakeRequest}; +use kafka_protocol::messages::{ + DescribeAclsRequest, SaslAuthenticateRequest, SaslHandshakeRequest, +}; + use tokio::sync::Semaphore; use crate::bridge::IggyBridge; use crate::error::Result;
