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;