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


The following commit(s) were added to refs/heads/kafka_proxy_acls by this push:
     new a8c72d2cf address review
a8c72d2cf is described below

commit a8c72d2cf7582fbbf3147548e371cd4d05df3a4f
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Wed Sep 23 14:28:42 2026 +0200

    address review
---
 .github/actions/rust/pre-merge/action.yml      |  7 +--
 gateways/kafka/README.md                       |  6 ++-
 gateways/kafka/docs/TEST_SUITE.md              | 11 ++++-
 gateways/kafka/src/auth.rs                     | 46 ++++++++++++++-----
 gateways/kafka/src/protocol/acl.rs             |  5 +++
 gateways/kafka/tests/kafka_client_e2e_tests.rs | 62 ++++++++++++++++++++++++--
 gateways/kafka/tests/sasl_tests.rs             | 60 +++++++++++++++++++++++++
 7 files changed, 176 insertions(+), 21 deletions(-)

diff --git a/.github/actions/rust/pre-merge/action.yml 
b/.github/actions/rust/pre-merge/action.yml
index 1789d9024..29bc5547f 100644
--- a/.github/actions/rust/pre-merge/action.yml
+++ b/.github/actions/rust/pre-merge/action.yml
@@ -401,8 +401,9 @@ runs:
           # Pull the client images here rather than letting the first test 
pull them. A cold
           # runner fetches ~400 MB for the Kafka distribution, and inside a 
test that competes
           # with nextest's per-test kill budget, so a slow registry surfaces 
as a test timeout
-          # naming the feature under test instead of naming the pull. Failure 
stays a warning:
-          # the test still tries, and the guard above still fails the job if 
it cannot run.
+          # naming the feature under test instead of naming the pull. The 
suite runs its
+          # containers with --pull never, so a failed pull here makes it skip, 
which the guard
+          # above turns into a job failure that names the missing image.
           # Read the tags out of the suite itself; hardcoding them here lets a 
bump in the test
           # desync silently, which pre-pulls the superseded image and puts the 
real one back
           # inside the budget this block exists to avoid.
@@ -413,7 +414,7 @@ runs:
           fi
           for image in $E2E_IMAGES; do
             timeout 300 docker pull -q "$image" \
-              || echo "::warning::docker pull $image failed; the real-client 
suite will retry it inside the test"
+              || echo "::warning::docker pull $image failed; the real-client 
suite will fail naming it"
           done
         fi
 
diff --git a/gateways/kafka/README.md b/gateways/kafka/README.md
index e9af99f9c..6e74597b2 100644
--- a/gateways/kafka/README.md
+++ b/gateways/kafka/README.md
@@ -133,8 +133,10 @@ Four things to know before switching it on:
 `kafka-acls.sh --list` works against the gateway. It is read only: 
`CreateAcls` and `DeleteAcls`
 are not implemented and not advertised.
 
-A principal sees its own permissions and nobody else's, because the gateway 
holds no administrative
-credentials. Only global permissions are rendered, as wildcard bindings, and 
the view is a snapshot
+`--list` here is a snapshot of the caller's own grants, not the broker-wide 
dump it is against a
+Kafka cluster. A principal sees its own permissions and nobody else's, because 
the gateway holds no
+administrative credentials, so filtering on another `User:` returns an empty 
listing whatever that
+user holds, and root's listing is root's grants rather than a catalog of every 
binding. Only global permissions are rendered, as wildcard bindings, and the 
view is a snapshot
 taken when the connection authenticated, so a permission changed afterwards is 
invisible until the
 client reconnects. [docs/ACL_MAPPING.md](docs/ACL_MAPPING.md) has the mapping 
table and what is
 deliberately left out.
diff --git a/gateways/kafka/docs/TEST_SUITE.md 
b/gateways/kafka/docs/TEST_SUITE.md
index 206c3b36c..38c89acb9 100644
--- a/gateways/kafka/docs/TEST_SUITE.md
+++ b/gateways/kafka/docs/TEST_SUITE.md
@@ -81,16 +81,23 @@ already-built `iggy-server` in the same target directory. 
Missing either makes i
 printed reason rather than fail, which is what lets `cargo test -p 
iggy-gateway-kafka` stay usable
 without either.
 
+Client containers run with `--pull never`, so pull the two images first. 
Without them the suite
+skips and names the missing one.
+
 ```bash
 cargo build --bin iggy-server
+docker pull edenhill/kcat:1.7.1
+docker pull apache/kafka:3.9.0
 KAFKA_E2E_REQUIRED=1 cargo test -p iggy-gateway-kafka --test 
kafka_client_e2e_tests
 ```
 
 `KAFKA_E2E_REQUIRED=1` turns a skip into a failure, mirroring 
`KAFKA_FIXTURES_REQUIRED`, so a CI
 job that means to run these cannot report a pass over zero assertions. Set it 
there.
 
-The suite shares the `kafka_bridge` nextest group with the bridge tests, so 
its spawned servers are
-serialized against them rather than competing for cores and ports.
+The suite runs in its own `kafka_client_e2e` nextest group, capped at one 
thread, so its spawned
+servers are serialized against each other. It is kept apart from the 
`kafka_bridge` group so the
+container-driven tests do not queue behind the bridge tests, and both groups 
cap their servers'
+shard pools, so the two can run alongside each other.
 
 It automates categories S and T of [`MANUAL_TESTING.md`](MANUAL_TESTING.md). 
Those procedures stay,
 because they cover cases a test does not assert, but the load-bearing ones now 
run in CI.
diff --git a/gateways/kafka/src/auth.rs b/gateways/kafka/src/auth.rs
index 7f04d775b..e89791cf1 100644
--- a/gateways/kafka/src/auth.rs
+++ b/gateways/kafka/src/auth.rs
@@ -62,7 +62,8 @@ const VERIFY_RECONNECTION_RETRIES: u32 = 1;
 /// caller bounds the whole exchange at its pre-authentication budget, which 
must also absorb the
 /// wait for an authentication slot, so every second spent here is a second 
that wait does not get.
 /// At one second the inner worst case is 12s against a 15s outer budget, 
leaving the queue three
-/// seconds rather than one.
+/// seconds rather than one. A read that times out skips the teardown wait, 
which could not succeed
+/// behind the SDK's still-running request, so the common slow path costs 11s, 
not 12s.
 ///
 /// It also bounds only *this* future, not the SDK's work. A cancelled call 
leaves the SDK's own
 /// read running on a detached task that holds its connection lock until that 
task's own deadline,
@@ -298,6 +299,20 @@ impl SaslAuthenticator for IggyAuthenticator {
             Ok(Ok(())) => fetch_permissions(&client, 
&credentials.username).await,
         };
 
+        // A timed-out permission read leaves the SDK's detached request task 
holding the stream
+        // lock until its own response deadline, and `shutdown` needs that 
lock first, so waiting
+        // on it could only burn `TEARDOWN_TIMEOUT` with the caller's 
authentication slot held.
+        // Dropping the client instead still aborts the heartbeat, which is 
the part that would
+        // reconnect. The socket itself lives on in that task until its 
deadline, which nothing on
+        // this side of the SDK can shorten.
+        if matches!(outcome, Ok((_, PermissionRead::TimedOut))) {
+            return outcome.map(|(permissions, _)| AuthenticatedPrincipal {
+                username: credentials.username.clone(),
+                permissions,
+                permissions_known: false,
+            });
+        }
+
         // Shut down on both paths, and not `disconnect`: only `shutdown` 
stops the heartbeat task,
         // which would otherwise keep pinging, observe the dropped transport, 
and reconnect using
         // the very credentials this call was only meant to check.
@@ -316,14 +331,25 @@ impl SaslAuthenticator for IggyAuthenticator {
             }
         }
 
-        outcome.map(|(permissions, permissions_known)| AuthenticatedPrincipal {
+        outcome.map(|(permissions, read)| AuthenticatedPrincipal {
             username: credentials.username.clone(),
             permissions,
-            permissions_known,
+            permissions_known: read == PermissionRead::Resolved,
         })
     }
 }
 
+/// How the permission read after a login ended.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+enum PermissionRead {
+    /// The record resolved, so the permissions are a real answer, possibly an 
empty one.
+    Resolved,
+    /// The read answered but gave no usable record.
+    Failed,
+    /// The read did not answer in time, and its request is still in flight 
inside the SDK.
+    TimedOut,
+}
+
 /// Reads the just-authenticated user's own record and projects its global 
permissions.
 ///
 /// A missing record or absent permissions both yield an empty set rather than 
an error: the
@@ -332,23 +358,23 @@ impl SaslAuthenticator for IggyAuthenticator {
 async fn fetch_permissions(
     client: &impl Client,
     username: &str,
-) -> Result<(PrincipalPermissions, bool), AuthError> {
+) -> Result<(PrincipalPermissions, PermissionRead), AuthError> {
     // Unreachable in practice: a username that reached a successful login is 
already inside
     // Identifier's own length bounds. Propagating rather than degrading is 
still the wrong shape
     // for this function, so it degrades like every other failure below.
     let Ok(identifier) = Identifier::named(username) else {
         warn!("authenticated, but the principal's name is not a valid Iggy 
identifier");
-        return Ok((PrincipalPermissions::default(), false));
+        return Ok((PrincipalPermissions::default(), PermissionRead::Failed));
     };
     // Deliberately not propagated as a failure. The credentials were already 
accepted by the login
     // above, so turning a stumble on this second round trip into a rejection 
would answer a correct
     // password with `SASL_AUTHENTICATION_FAILED`, which a Kafka client treats 
as fatal and raises
     // to the application. Losing the ACL view is the lesser harm, and it 
degrades to an empty one.
     let fetched = tokio::time::timeout(PERMISSION_READ_TIMEOUT, 
client.get_user(&identifier)).await;
-    let mut known = true;
+    let mut read = PermissionRead::Resolved;
     let user = match fetched {
         Ok(Ok(None)) => {
-            known = false;
+            read = PermissionRead::Failed;
             // Distinct from an error: the login succeeded, so the account 
exists. A record that
             // resolves to nothing here means the read raced a deletion, and 
silently reporting an
             // empty ACL view for it would look identical to a principal with 
no grants.
@@ -357,12 +383,12 @@ async fn fetch_permissions(
         }
         Ok(Ok(user)) => user,
         Ok(Err(error)) => {
-            known = false;
+            read = PermissionRead::Failed;
             warn!(%error, "authenticated, but could not read the principal's 
permissions");
             None
         }
         Err(_elapsed) => {
-            known = false;
+            read = PermissionRead::TimedOut;
             warn!("authenticated, but timed out reading the principal's 
permissions");
             None
         }
@@ -375,7 +401,7 @@ async fn fetch_permissions(
             .as_ref()
             .map(PrincipalPermissions::from)
             .unwrap_or_default(),
-        known,
+        read,
     ))
 }
 
diff --git a/gateways/kafka/src/protocol/acl.rs 
b/gateways/kafka/src/protocol/acl.rs
index fddd4ba0d..e09547e19 100644
--- a/gateways/kafka/src/protocol/acl.rs
+++ b/gateways/kafka/src/protocol/acl.rs
@@ -73,6 +73,11 @@ pub const ANY_HOST: &str = "*";
 /// topic, and every Kafka topic lives inside one Iggy stream, so a stream 
grant is in practice a
 /// grant over the topics a Kafka client can reach.
 ///
+/// A description, never an authorization input. It holds only global flags, 
so a principal whose
+/// grants are per-stream or per-topic reads here as holding nothing, and an 
authorizer built on it
+/// would deny what Iggy allows. Enforcement belongs to Iggy's own 
`Permissions`, evaluated by the
+/// server.
+///
 /// The boolean count mirrors Iggy's own `GlobalPermissions`, which is a flat 
set of independent
 /// grants. Collapsing them into a bitfield would hide which grant is which at 
every call site for
 /// no gain, so the lint is allowed here the way it is elsewhere in this 
repository.
diff --git a/gateways/kafka/tests/kafka_client_e2e_tests.rs 
b/gateways/kafka/tests/kafka_client_e2e_tests.rs
index c4b3b4109..699d3ca52 100644
--- a/gateways/kafka/tests/kafka_client_e2e_tests.rs
+++ b/gateways/kafka/tests/kafka_client_e2e_tests.rs
@@ -56,6 +56,15 @@ const USER_PASSWORD: &str = "s3cretpass";
 /// build on a loaded machine is not quick.
 const SERVER_READY_TIMEOUT: Duration = Duration::from_secs(45);
 
+/// Wall-clock cap on one client container. The ACL test runs five in a row 
inside nextest's 300s
+/// kill budget, so a single wedged client must fail on its own terms, well 
before that budget
+/// kills the test and hides which step hung.
+const CLIENT_RUN_TIMEOUT: &str = "45s";
+
+/// Per-request and per-call budget for the Java admin tools. Their defaults 
(30s and 60s) let one
+/// unanswered request eat most of `CLIENT_RUN_TIMEOUT` retrying.
+const JAVA_CLIENT_TIMEOUT_MS: u32 = 10_000;
+
 /// Reports why the suite cannot run, and whether that is fatal.
 ///
 /// Returns `true` when the caller should skip. `KAFKA_E2E_REQUIRED=1` makes 
it panic instead, so a
@@ -83,6 +92,19 @@ fn docker_missing() -> bool {
     }
 }
 
+/// Whether `image` is already in the local Docker store.
+///
+/// Containers run with `--pull never`, so a missing image would otherwise 
surface as a client that
+/// failed to start, reported against whichever feature that test happened to 
cover.
+fn image_present(image: &str) -> bool {
+    Command::new("docker")
+        .args(["image", "inspect", image])
+        .stdout(Stdio::null())
+        .stderr(Stdio::null())
+        .status()
+        .is_ok_and(|status| status.success())
+}
+
 /// Locates the already-built `iggy-server` alongside this test binary. Does 
not build it.
 fn iggy_server_binary() -> Option<PathBuf> {
     let mut dir = std::env::current_exe().ok()?;
@@ -224,9 +246,22 @@ async fn spawn_gateway(iggy_address: &str) -> SocketAddr {
 /// This blocks the calling thread for the life of the container, which is why 
every test here uses
 /// a multi-threaded runtime: the gateway runs as a spawned task, and on the 
single-threaded runtime
 /// `#[tokio::test]` gives by default, this call would starve it and nothing 
would ever listen.
+///
+/// `--pull never` keeps a registry download out of the test's time budget: 
images are pulled up
+/// front, and `stack` skips when they are absent. `timeout` bounds the run 
itself, and a run it
+/// cuts short exits non-zero, so `expect_ran` reports it rather than 
asserting over partial output.
 fn run_client(image: &str, args: &[&str], mounts: &[(&str, &str)]) -> 
ClientRun {
-    let mut command = Command::new("docker");
-    command.args(["run", "--rm", "--network", "host"]);
+    let mut command = Command::new("timeout");
+    command.args(["--kill-after=10s", CLIENT_RUN_TIMEOUT]);
+    command.args([
+        "docker",
+        "run",
+        "--rm",
+        "--pull",
+        "never",
+        "--network",
+        "host",
+    ]);
     for (host, guest) in mounts {
         command.args(["-v", &format!("{host}:{guest}:ro")]);
     }
@@ -359,7 +394,9 @@ fn java_client_config(username: &str, password: &str) -> 
(tempfile::TempDir, Pat
             "security.protocol=SASL_PLAINTEXT\n\
              sasl.mechanism=PLAIN\n\
              
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule 
required \
-             username=\"{username}\" password=\"{password}\";\n"
+             username=\"{username}\" password=\"{password}\";\n\
+             request.timeout.ms={JAVA_CLIENT_TIMEOUT_MS}\n\
+             default.api.timeout.ms={JAVA_CLIENT_TIMEOUT_MS}\n"
         ),
     )
     .expect("write client config");
@@ -371,6 +408,15 @@ async fn stack() -> Option<(TestServer, SocketAddr)> {
     if docker_missing() {
         return None;
     }
+    if let Some(image) = [KCAT_IMAGE, KAFKA_IMAGE]
+        .into_iter()
+        .find(|image| !image_present(image))
+    {
+        skip(&format!(
+            "client image {image} is not pulled; run `docker pull {image}`"
+        ));
+        return None;
+    }
     let server = match TestServer::spawn() {
         Ok(server) => server,
         Err(reason) => {
@@ -526,9 +572,17 @@ async fn 
given_distinct_principals_when_listing_acls_should_describe_each_differ
         "every rendered binding names the authenticated principal, got: {root}"
     );
     assert!(
-        root.contains("resourceType=TOPIC") && 
root.contains("operation=WRITE"),
+        grants(&root, "TOPIC", "WRITE"),
         "root holds every permission, so it must be described as able to 
write, got: {root}"
     );
+    assert!(
+        grants(&root, "CLUSTER", "DESCRIBE"),
+        "root holds the server flags, so the cluster must render as 
describable, got: {root}"
+    );
+    assert!(
+        !grants(&root, "CLUSTER", "ALTER"),
+        "no Iggy flag gates a cluster mutation, so even root must not be shown 
one, got: {root}"
+    );
 
     let consumer = list_acls(gateway, "consumer-only", USER_PASSWORD);
     // Paired assertions: which operation sits under which resource is the 
whole claim.
diff --git a/gateways/kafka/tests/sasl_tests.rs 
b/gateways/kafka/tests/sasl_tests.rs
index 01f6da1a1..35c6a0d2e 100644
--- a/gateways/kafka/tests/sasl_tests.rs
+++ b/gateways/kafka/tests/sasl_tests.rs
@@ -56,6 +56,7 @@ const ERROR_UNSUPPORTED_SASL_MECHANISM: i16 = 33;
 const ERROR_ILLEGAL_SASL_STATE: i16 = 34;
 const ERROR_UNSUPPORTED_VERSION: i16 = 35;
 const ERROR_SASL_AUTHENTICATION_FAILED: i16 = 58;
+const ERROR_UNKNOWN_SERVER_ERROR: i16 = -1;
 
 const HANDSHAKE_VERSION: i16 = 1;
 const AUTHENTICATE_VERSION: i16 = 1;
@@ -1313,6 +1314,65 @@ async fn 
given_a_principal_with_no_permissions_should_report_an_empty_view_not_a
     );
 }
 
+/// Accepts `alice` but reports that her permissions could not be read, the 
shape
+/// `IggyAuthenticator` produces when the login succeeds and the follow-up 
`get_user` does not.
+#[derive(Debug)]
+struct UnreadPermissionsAuthenticator;
+
+#[async_trait]
+impl SaslAuthenticator for UnreadPermissionsAuthenticator {
+    async fn authenticate(
+        &self,
+        credentials: &PlainCredentials,
+    ) -> Result<AuthenticatedPrincipal, AuthError> {
+        Ok(AuthenticatedPrincipal {
+            username: credentials.username.clone(),
+            permissions: PrincipalPermissions::default(),
+            permissions_known: false,
+        })
+    }
+}
+
+#[tokio::test]
+async fn 
given_unread_permissions_when_describing_acls_should_answer_an_error_and_stay_open()
 {
+    // The fallback permissions are empty, so answering from them would render 
the same zero
+    // bindings as a principal that genuinely holds nothing. Every other stub 
reports its
+    // permissions as known, so without this case the two answers could swap 
unnoticed.
+    let (addr, shutdown) = spawn_test_server_with_authenticator(
+        sasl_config(),
+        Arc::new(UnreadPermissionsAuthenticator),
+    )
+    .await;
+    std::mem::forget(shutdown);
+    let mut stream = TcpStream::connect(addr).await.expect("connect");
+    authenticate(&mut stream).await;
+
+    let body = send(
+        &mut stream,
+        API_KEY_DESCRIBE_ACLS,
+        DESCRIBE_ACLS_VERSION,
+        3,
+        &any_acl_filter_body(),
+    )
+    .await;
+    assert_eq!(
+        acl_error_code(&body),
+        ERROR_UNKNOWN_SERVER_ERROR,
+        "a view that was never read must not be reported as an empty one"
+    );
+    assert!(
+        parse_acl_bindings(&body).is_empty(),
+        "an error answer carries no bindings"
+    );
+
+    // Kept open: the login was valid, and only the ACL view is missing.
+    let metadata = send(&mut stream, API_KEY_METADATA, 0, 4, &[0, 0, 0, 
0]).await;
+    assert!(
+        !metadata.is_empty(),
+        "the connection must keep serving after the ACL error"
+    );
+}
+
 #[tokio::test]
 async fn 
given_an_unauthenticated_connection_when_describing_acls_should_be_answered_then_closed()
 {
     let addr = spawn_gateway_for(PrincipalPermissions::default()).await;

Reply via email to