Skip to content

Commit 36e8cdc

Browse files
committed
Address comments
1 parent 2601a85 commit 36e8cdc

8 files changed

Lines changed: 190 additions & 75 deletions

File tree

bindings/cpp/test/test_sasl_auth.cpp

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,9 @@ class SaslAuthTest : public ::testing::Test {
2626
const std::string& sasl_servers() {
2727
return fluss_test::FlussTestEnvironment::Instance()->GetSaslBootstrapServers();
2828
}
29+
const std::string& plaintext_servers() {
30+
return fluss_test::FlussTestEnvironment::Instance()->GetBootstrapServers();
31+
}
2932
};
3033

3134
TEST_F(SaslAuthTest, SaslConnectWithValidCredentials) {
@@ -55,6 +58,26 @@ TEST_F(SaslAuthTest, SaslConnectWithValidCredentials) {
5558
ASSERT_OK(admin.DropDatabase(db_name, true, true));
5659
}
5760

61+
TEST_F(SaslAuthTest, SaslConnectWithSecondUser) {
62+
fluss::Configuration config;
63+
config.bootstrap_servers = sasl_servers();
64+
config.security_protocol = "sasl";
65+
config.security_sasl_mechanism = "PLAIN";
66+
config.security_sasl_username = "alice";
67+
config.security_sasl_password = "alice-secret";
68+
69+
fluss::Connection conn;
70+
ASSERT_OK(fluss::Connection::Create(config, conn));
71+
72+
fluss::Admin admin;
73+
ASSERT_OK(conn.GetAdmin(admin));
74+
75+
// Basic operation to confirm functional connection
76+
bool exists = false;
77+
ASSERT_OK(admin.DatabaseExists("some_nonexistent_db_alice", exists));
78+
ASSERT_FALSE(exists);
79+
}
80+
5881
TEST_F(SaslAuthTest, SaslConnectWithWrongPassword) {
5982
fluss::Configuration config;
6083
config.bootstrap_servers = sasl_servers();
@@ -87,3 +110,16 @@ TEST_F(SaslAuthTest, SaslConnectWithUnknownUser) {
87110
EXPECT_NE(result.error_message.find("Authentication failed"), std::string::npos)
88111
<< "Expected 'Authentication failed' in: " << result.error_message;
89112
}
113+
114+
TEST_F(SaslAuthTest, SaslClientToPlaintextServer) {
115+
fluss::Configuration config;
116+
config.bootstrap_servers = plaintext_servers();
117+
config.security_protocol = "sasl";
118+
config.security_sasl_mechanism = "PLAIN";
119+
config.security_sasl_username = "admin";
120+
config.security_sasl_password = "admin-secret";
121+
122+
fluss::Connection conn;
123+
auto result = fluss::Connection::Create(config, conn);
124+
ASSERT_FALSE(result.Ok()) << "SASL client connecting to plaintext server should fail";
125+
}

bindings/cpp/test/test_utils.h

Lines changed: 61 additions & 49 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,29 @@ static constexpr int kPlainClientTabletPort = 9224;
6363
/// Execute a shell command and return its exit code.
6464
inline int RunCommand(const std::string& cmd) { return system(cmd.c_str()); }
6565

66+
/// Join property lines with the escaped newline separator used by `printf` in docker commands.
67+
inline std::string JoinProps(const std::vector<std::string>& lines) {
68+
std::string result;
69+
for (size_t i = 0; i < lines.size(); ++i) {
70+
if (i > 0) result += "\\n";
71+
result += lines[i];
72+
}
73+
return result;
74+
}
75+
76+
/// Build a `docker run` command with FLUSS_PROPERTIES.
77+
inline std::string DockerRunCmd(const std::string& name, const std::string& props,
78+
const std::vector<std::string>& port_mappings,
79+
const std::string& server_type) {
80+
std::string cmd = "docker run -d --rm --name " + name + " --network " + kNetworkName;
81+
for (const auto& pm : port_mappings) {
82+
cmd += " -p " + pm;
83+
}
84+
cmd += " -e FLUSS_PROPERTIES=\"$(printf '" + props + "')\"";
85+
cmd += " " + std::string(kFlussImage) + ":" + kFlussVersion + " " + server_type;
86+
return cmd;
87+
}
88+
6689
/// Wait until a TCP port is accepting connections, or timeout.
6790
inline bool WaitForPort(const std::string& host, int port, int timeout_seconds = 60) {
6891
auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(timeout_seconds);
@@ -129,28 +152,24 @@ class FlussTestCluster {
129152
std::string sasl_jaas =
130153
"org.apache.fluss.security.auth.sasl.plain.PlainLoginModule required"
131154
" user_admin=\"admin-secret\" user_alice=\"alice-secret\";";
132-
std::string coord_props =
133-
"zookeeper.address: " + std::string(kZookeeperName) +
134-
":2181\\n"
135-
"bind.listeners: INTERNAL://" +
136-
std::string(kCoordinatorName) + ":0, CLIENT://" + std::string(kCoordinatorName) +
137-
":9123, PLAIN_CLIENT://" + std::string(kCoordinatorName) +
138-
":9223\\n"
139-
"advertised.listeners: CLIENT://localhost:9123, PLAIN_CLIENT://localhost:9223\\n"
140-
"internal.listener.name: INTERNAL\\n"
141-
"security.protocol.map: CLIENT:sasl\\n"
142-
"security.sasl.enabled.mechanisms: plain\\n"
143-
"security.sasl.plain.jaas.config: " +
144-
sasl_jaas +
145-
"\\n"
146-
"netty.server.num-network-threads: 1\\n"
147-
"netty.server.num-worker-threads: 3";
148-
149-
std::string coord_cmd = std::string("docker run -d --rm") + " --name " + kCoordinatorName +
150-
" --network " + kNetworkName + " -p 9123:9123" + " -p 9223:9223" +
151-
" -e FLUSS_PROPERTIES=\"$(printf '" + coord_props + "')\"" + " " +
152-
std::string(kFlussImage) + ":" + kFlussVersion +
153-
" coordinatorServer";
155+
156+
std::string coord = std::string(kCoordinatorName);
157+
std::string zk = std::string(kZookeeperName);
158+
std::string coord_props = JoinProps({
159+
"zookeeper.address: " + zk + ":2181",
160+
"bind.listeners: INTERNAL://" + coord + ":0, CLIENT://" + coord +
161+
":9123, PLAIN_CLIENT://" + coord + ":9223",
162+
"advertised.listeners: CLIENT://localhost:9123, PLAIN_CLIENT://localhost:9223",
163+
"internal.listener.name: INTERNAL",
164+
"security.protocol.map: CLIENT:sasl",
165+
"security.sasl.enabled.mechanisms: plain",
166+
"security.sasl.plain.jaas.config: " + sasl_jaas,
167+
"netty.server.num-network-threads: 1",
168+
"netty.server.num-worker-threads: 3",
169+
});
170+
171+
std::string coord_cmd = DockerRunCmd(kCoordinatorName, coord_props,
172+
{"9123:9123", "9223:9223"}, "coordinatorServer");
154173
if (RunCommand(coord_cmd) != 0) {
155174
std::cerr << "Failed to start Coordinator Server" << std::endl;
156175
Stop();
@@ -165,33 +184,26 @@ class FlussTestCluster {
165184
}
166185

167186
// Start Tablet Server (dual listeners: CLIENT=SASL on 9123, PLAIN_CLIENT=plaintext on 9223)
168-
std::string ts_props =
169-
"zookeeper.address: " + std::string(kZookeeperName) +
170-
":2181\\n"
171-
"bind.listeners: INTERNAL://" +
172-
std::string(kTabletServerName) + ":0, CLIENT://" + std::string(kTabletServerName) +
173-
":9123, PLAIN_CLIENT://" + std::string(kTabletServerName) +
174-
":9223\\n"
175-
"advertised.listeners: CLIENT://localhost:" +
176-
std::to_string(kTabletServerPort) +
177-
", PLAIN_CLIENT://localhost:" + std::to_string(kPlainClientTabletPort) +
178-
"\\n"
179-
"internal.listener.name: INTERNAL\\n"
180-
"security.protocol.map: CLIENT:sasl\\n"
181-
"security.sasl.enabled.mechanisms: plain\\n"
182-
"security.sasl.plain.jaas.config: " +
183-
sasl_jaas +
184-
"\\n"
185-
"tablet-server.id: 0\\n"
186-
"netty.server.num-network-threads: 1\\n"
187-
"netty.server.num-worker-threads: 3";
188-
189-
std::string ts_cmd = std::string("docker run -d --rm") + " --name " + kTabletServerName +
190-
" --network " + kNetworkName + " -p " +
191-
std::to_string(kTabletServerPort) + ":9123" + " -p " +
192-
std::to_string(kPlainClientTabletPort) + ":9223" +
193-
" -e FLUSS_PROPERTIES=\"$(printf '" + ts_props + "')\"" + " " +
194-
std::string(kFlussImage) + ":" + kFlussVersion + " tabletServer";
187+
std::string ts = std::string(kTabletServerName);
188+
std::string ts_props = JoinProps({
189+
"zookeeper.address: " + zk + ":2181",
190+
"bind.listeners: INTERNAL://" + ts + ":0, CLIENT://" + ts + ":9123, PLAIN_CLIENT://" +
191+
ts + ":9223",
192+
"advertised.listeners: CLIENT://localhost:" + std::to_string(kTabletServerPort) +
193+
", PLAIN_CLIENT://localhost:" + std::to_string(kPlainClientTabletPort),
194+
"internal.listener.name: INTERNAL",
195+
"security.protocol.map: CLIENT:sasl",
196+
"security.sasl.enabled.mechanisms: plain",
197+
"security.sasl.plain.jaas.config: " + sasl_jaas,
198+
"tablet-server.id: 0",
199+
"netty.server.num-network-threads: 1",
200+
"netty.server.num-worker-threads: 3",
201+
});
202+
203+
std::string ts_cmd = DockerRunCmd(kTabletServerName, ts_props,
204+
{std::to_string(kTabletServerPort) + ":9123",
205+
std::to_string(kPlainClientTabletPort) + ":9223"},
206+
"tabletServer");
195207
if (RunCommand(ts_cmd) != 0) {
196208
std::cerr << "Failed to start Tablet Server" << std::endl;
197209
Stop();

bindings/python/src/config.rs

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -95,21 +95,21 @@ impl Config {
9595
}
9696
};
9797
}
98-
"client.connect-timeout" => {
98+
"connect-timeout" => {
9999
config.connect_timeout_ms = value.parse::<u64>().map_err(|e| {
100100
FlussError::new_err(format!("Invalid value '{value}' for '{key}': {e}"))
101101
})?;
102102
}
103-
"client.security.protocol" => {
103+
"security.protocol" => {
104104
config.security_protocol = value;
105105
}
106-
"client.security.sasl.mechanism" => {
106+
"security.sasl.mechanism" => {
107107
config.security_sasl_mechanism = value;
108108
}
109-
"client.security.sasl.username" => {
109+
"security.sasl.username" => {
110110
config.security_sasl_username = value;
111111
}
112-
"client.security.sasl.password" => {
112+
"security.sasl.password" => {
113113
config.security_sasl_password = value;
114114
}
115115
_ => {

bindings/python/test/conftest.py

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -159,6 +159,13 @@ def sasl_bootstrap_servers(fluss_cluster):
159159
return sasl_addr
160160

161161

162+
@pytest.fixture(scope="session")
163+
def plaintext_bootstrap_servers(fluss_cluster):
164+
"""Bootstrap servers for the plaintext (non-SASL) listener."""
165+
plaintext_addr, _sasl_addr = fluss_cluster
166+
return plaintext_addr
167+
168+
162169
@pytest_asyncio.fixture(scope="session")
163170
async def admin(connection):
164171
"""Session-scoped admin client."""

bindings/python/test/test_sasl_auth.py

Lines changed: 42 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -29,10 +29,10 @@ async def test_sasl_connect_with_valid_credentials(sasl_bootstrap_servers):
2929
"""Verify that a client with correct SASL credentials can connect and perform operations."""
3030
config = fluss.Config({
3131
"bootstrap.servers": sasl_bootstrap_servers,
32-
"client.security.protocol": "sasl",
33-
"client.security.sasl.mechanism": "PLAIN",
34-
"client.security.sasl.username": "admin",
35-
"client.security.sasl.password": "admin-secret",
32+
"security.protocol": "sasl",
33+
"security.sasl.mechanism": "PLAIN",
34+
"security.sasl.username": "admin",
35+
"security.sasl.password": "admin-secret",
3636
})
3737
conn = await fluss.FlussConnection.create(config)
3838
admin = await conn.get_admin()
@@ -48,14 +48,31 @@ async def test_sasl_connect_with_valid_credentials(sasl_bootstrap_servers):
4848
conn.close()
4949

5050

51+
async def test_sasl_connect_with_second_user(sasl_bootstrap_servers):
52+
"""Verify that a second user can also authenticate successfully."""
53+
config = fluss.Config({
54+
"bootstrap.servers": sasl_bootstrap_servers,
55+
"security.protocol": "sasl",
56+
"security.sasl.mechanism": "PLAIN",
57+
"security.sasl.username": "alice",
58+
"security.sasl.password": "alice-secret",
59+
})
60+
conn = await fluss.FlussConnection.create(config)
61+
admin = await conn.get_admin()
62+
63+
# Basic operation to confirm functional connection
64+
assert not await admin.database_exists("some_nonexistent_db_alice")
65+
conn.close()
66+
67+
5168
async def test_sasl_connect_with_wrong_password(sasl_bootstrap_servers):
5269
"""Verify that wrong credentials are rejected with AUTHENTICATE_EXCEPTION."""
5370
config = fluss.Config({
5471
"bootstrap.servers": sasl_bootstrap_servers,
55-
"client.security.protocol": "sasl",
56-
"client.security.sasl.mechanism": "PLAIN",
57-
"client.security.sasl.username": "admin",
58-
"client.security.sasl.password": "wrong-password",
72+
"security.protocol": "sasl",
73+
"security.sasl.mechanism": "PLAIN",
74+
"security.sasl.username": "admin",
75+
"security.sasl.password": "wrong-password",
5976
})
6077
with pytest.raises(fluss.FlussError) as exc_info:
6178
await fluss.FlussConnection.create(config)
@@ -67,12 +84,25 @@ async def test_sasl_connect_with_unknown_user(sasl_bootstrap_servers):
6784
"""Verify that a nonexistent user is rejected with AUTHENTICATE_EXCEPTION."""
6885
config = fluss.Config({
6986
"bootstrap.servers": sasl_bootstrap_servers,
70-
"client.security.protocol": "sasl",
71-
"client.security.sasl.mechanism": "PLAIN",
72-
"client.security.sasl.username": "nonexistent_user",
73-
"client.security.sasl.password": "some-password",
87+
"security.protocol": "sasl",
88+
"security.sasl.mechanism": "PLAIN",
89+
"security.sasl.username": "nonexistent_user",
90+
"security.sasl.password": "some-password",
7491
})
7592
with pytest.raises(fluss.FlussError) as exc_info:
7693
await fluss.FlussConnection.create(config)
7794

7895
assert exc_info.value.error_code == fluss.ErrorCode.AUTHENTICATE_EXCEPTION
96+
97+
98+
async def test_sasl_client_to_plaintext_server(plaintext_bootstrap_servers):
99+
"""Verify that a SASL-configured client fails when connecting to a plaintext server."""
100+
config = fluss.Config({
101+
"bootstrap.servers": plaintext_bootstrap_servers,
102+
"security.protocol": "sasl",
103+
"security.sasl.mechanism": "PLAIN",
104+
"security.sasl.username": "admin",
105+
"security.sasl.password": "admin-secret",
106+
})
107+
with pytest.raises(fluss.FlussError):
108+
await fluss.FlussConnection.create(config)

crates/fluss/tests/integration/fluss_cluster.rs

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -378,6 +378,11 @@ impl FlussTestingCluster {
378378
&self.sasl_users
379379
}
380380

381+
/// Returns the plaintext (non-SASL) bootstrap servers address.
382+
pub fn plaintext_bootstrap_servers(&self) -> &str {
383+
&self.bootstrap_servers
384+
}
385+
381386
pub async fn get_fluss_connection(&self) -> FlussConnection {
382387
let config = Config {
383388
writer_acks: "all".to_string(),

crates/fluss/tests/integration/sasl_auth.rs

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,8 @@
1818
#[cfg(test)]
1919
mod sasl_auth_test {
2020
use crate::integration::utils::get_shared_cluster;
21+
use fluss::client::FlussConnection;
22+
use fluss::config::Config;
2123
use fluss::error::FlussError;
2224
use fluss::metadata::DatabaseDescriptorBuilder;
2325

@@ -102,6 +104,29 @@ mod sasl_auth_test {
102104
);
103105
}
104106

107+
/// Verify that a SASL-configured client fails when connecting to a plaintext server.
108+
#[tokio::test]
109+
async fn test_sasl_client_to_plaintext_server() {
110+
let cluster = get_shared_cluster();
111+
let plaintext_addr = cluster.plaintext_bootstrap_servers().to_string();
112+
113+
let config = Config {
114+
writer_acks: "all".to_string(),
115+
bootstrap_servers: plaintext_addr,
116+
security_protocol: "sasl".to_string(),
117+
security_sasl_mechanism: "PLAIN".to_string(),
118+
security_sasl_username: SASL_USERNAME.to_string(),
119+
security_sasl_password: SASL_PASSWORD.to_string(),
120+
..Default::default()
121+
};
122+
123+
let result = FlussConnection::new(config).await;
124+
assert!(
125+
result.is_err(),
126+
"SASL client connecting to plaintext server should fail"
127+
);
128+
}
129+
105130
/// Verify that a nonexistent user is rejected with a typed error.
106131
#[tokio::test]
107132
async fn test_sasl_connect_with_unknown_user() {

website/docs/user-guide/python/example/configuration.md

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -32,11 +32,11 @@ with await fluss.FlussConnection.create(config) as conn:
3232
| `scanner.remote-log.prefetch-num` | Number of remote log segments to prefetch | `4` |
3333
| `remote-file.download-thread-num` | Number of threads for remote log downloads | `3` |
3434
| `scanner.log.max-poll-records` | Max records returned in a single poll() | `500` |
35-
| `client.connect-timeout` | TCP connect timeout in milliseconds | `120000` |
36-
| `client.security.protocol` | `PLAINTEXT` (default) or `sasl` for SASL auth | `PLAINTEXT` |
37-
| `client.security.sasl.mechanism` | SASL mechanism (only `PLAIN` is supported) | `PLAIN` |
38-
| `client.security.sasl.username` | SASL username (required when protocol is `sasl`) | (empty) |
39-
| `client.security.sasl.password` | SASL password (required when protocol is `sasl`) | (empty) |
35+
| `connect-timeout` | TCP connect timeout in milliseconds | `120000` |
36+
| `security.protocol` | `PLAINTEXT` (default) or `sasl` for SASL auth | `PLAINTEXT` |
37+
| `security.sasl.mechanism` | SASL mechanism (only `PLAIN` is supported) | `PLAIN` |
38+
| `security.sasl.username` | SASL username (required when protocol is `sasl`) | (empty) |
39+
| `security.sasl.password` | SASL password (required when protocol is `sasl`) | (empty) |
4040

4141
## SASL Authentication
4242

@@ -45,10 +45,10 @@ To connect to a Fluss cluster with SASL/PLAIN authentication enabled:
4545
```python
4646
config = fluss.Config({
4747
"bootstrap.servers": "127.0.0.1:9123",
48-
"client.security.protocol": "sasl",
49-
"client.security.sasl.mechanism": "PLAIN",
50-
"client.security.sasl.username": "admin",
51-
"client.security.sasl.password": "admin-secret",
48+
"security.protocol": "sasl",
49+
"security.sasl.mechanism": "PLAIN",
50+
"security.sasl.username": "admin",
51+
"security.sasl.password": "admin-secret",
5252
})
5353
conn = await fluss.FlussConnection.create(config)
5454
```

0 commit comments

Comments
 (0)