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;

Reply via email to