diff --git a/docs/en/antalya/protocol.md b/docs/en/antalya/protocol.md index f7ca781159e7..be10b9014a31 100644 --- a/docs/en/antalya/protocol.md +++ b/docs/en/antalya/protocol.md @@ -10,16 +10,17 @@ doc_type: 'reference' # Antalya protocol version {#antalya-protocol-version} Antalya versions its own wire-protocol changes with `DBMS_ANTALYA_PROTOCOL_VERSION`, a counter that -upstream ClickHouse cannot reach, defined in `src/Core/AntalyaProtocol.h`. A server advertises it in -the `ServerHello` name string, on every connection: +upstream ClickHouse cannot reach, defined in `src/Core/AntalyaProtocol.h`. Both sides advertise it in +the name string of their `Hello`, on every connection: ```text +client -> server "ClickHouse client (antalya:1)" server -> client "ClickHouse (antalya:1)" ``` -The client parses the suffix, caps the value with `min(own, server)` and keeps the result. `0` means -the peer is not an Antalya build. Negotiation is per hop and not transitive: initiator to worker and -worker to worker negotiate independently. +Each side parses the peer's suffix, caps the value with `min(own, peer)` and keeps the result. `0` +means the peer is not an Antalya build. Negotiation is per hop and not transitive: initiator to +worker and worker to worker negotiate independently. Version 1 is the advertisement itself. Nothing is gated on it yet. @@ -30,8 +31,6 @@ Version 1 is the advertisement itself. Nothing is gated on it yet. `DBMS_CLUSTER_PROCESSING_PROTOCOL_VERSION` for a feature upstream does not have. - Keep the counter cumulative. A backport takes the whole contiguous range up to the value it needs, or does not bump at all - the `min(own, server)` cap is only sound for a cumulative feature set. -- Gate only what the *client* decides to do. The server never learns the client's version, because - only the server advertises. - Update this page, and update `docs/en/interfaces/specs/NativeProtocol.md` when the change alters a packet layout described there. @@ -40,11 +39,9 @@ Version 1 is the advertisement itself. Nothing is gated on it yet. An upstream rebase can reuse the next value of an upstream protocol counter for a different feature. Keeping the Antalya counter separate prevents the same version from describing two wire layouts. -## Why the marker rides in `ServerHello` {#why-the-marker-rides-in-serverhello} +## Why the marker rides in the `Hello` name {#why-the-marker-rides-in-the-hello-name} -The client `Hello` cannot advertise the version because it is sent before the peer is known. Its -`client_name` is also stored and validated against the Query packet, so changing it can raise -`CLIENT_INFO_DOES_NOT_MATCH` on an upstream peer. - -The server advertises through `server_name`, which is display text. The marker stays inside that -existing string because adding a field would make older peers read it as the next packet. +The marker stays inside the existing name string because adding a field would make older peers read +it as the next packet. The client writes its `Hello` before it knows the peer, so it always marks +`client_name`. With `validate_tcp_client_information` enabled, an upstream server rejects an initial +query from an Antalya client with `CLIENT_INFO_DOES_NOT_MATCH`; an Antalya server ignores the marker. diff --git a/docs/en/interfaces/specs/NativeProtocol.md b/docs/en/interfaces/specs/NativeProtocol.md index 02eb626f2598..8a2c8f6ae40b 100644 --- a/docs/en/interfaces/specs/NativeProtocol.md +++ b/docs/en/interfaces/specs/NativeProtocol.md @@ -426,7 +426,7 @@ Client → Server. The first message after the TCP connection opens. | # | Field | Type | Role | Description | |---|------------------|---------|-----------|-------------| -| 1 | client_name | String | universal | Client identifier (e.g., `"clickhouse-client"`) | +| 1 | client_name | String | universal | Client identifier (e.g., `"clickhouse-client"`). An Altinity Antalya build appends `" (antalya:N)"`, where `N` is its Antalya protocol version. See [Antalya protocol version](/antalya/protocol). | | 2 | version_major | VarUInt | universal | Client major version | | 3 | version_minor | VarUInt | universal | Client minor version | | 4 | protocol_version | VarUInt | universal | Client's max supported protocol version | diff --git a/src/Client/Connection.cpp b/src/Client/Connection.cpp index beb525678813..71b11aded41c 100644 --- a/src/Client/Connection.cpp +++ b/src/Client/Connection.cpp @@ -504,7 +504,7 @@ void Connection::sendHello() "Parameters 'default_database', 'user' and 'password' must not contain ASCII control characters"); writeVarUInt(Protocol::Client::Hello, *out); - writeStringBinary(std::string(VERSION_NAME) + " " + client_name, *out); + writeStringBinary(AntalyaProtocol::appendMarker(std::string(VERSION_NAME) + " " + client_name), *out); writeVarUInt(VERSION_MAJOR, *out); writeVarUInt(VERSION_MINOR, *out); // NOTE For backward compatibility of the protocol, client cannot send its version_patch. diff --git a/src/Core/AntalyaProtocol.cpp b/src/Core/AntalyaProtocol.cpp index b697a52e2c3a..2d5a2f129d56 100644 --- a/src/Core/AntalyaProtocol.cpp +++ b/src/Core/AntalyaProtocol.cpp @@ -55,6 +55,13 @@ UInt64 parseMarker(std::string_view name) return std::min(version, DBMS_ANTALYA_PROTOCOL_VERSION); } +std::string_view removeMarker(std::string_view name) +{ + if (parseMarker(name) == 0) + return name; + return name.substr(0, name.rfind(MARKER_PREFIX)); +} + } } diff --git a/src/Core/AntalyaProtocol.h b/src/Core/AntalyaProtocol.h index c1bb7ece6776..53cf8d1ceb68 100644 --- a/src/Core/AntalyaProtocol.h +++ b/src/Core/AntalyaProtocol.h @@ -18,6 +18,9 @@ String appendMarker(std::string_view name); /// Returns the negotiated version, or `0` if there is no marker. UInt64 parseMarker(std::string_view name); +/// Returns `name` without the marker, or `name` itself if there is no marker. +std::string_view removeMarker(std::string_view name); + } } diff --git a/src/Core/tests/gtest_antalya_protocol.cpp b/src/Core/tests/gtest_antalya_protocol.cpp index 57fdbe5ef035..2363aa9408e9 100644 --- a/src/Core/tests/gtest_antalya_protocol.cpp +++ b/src/Core/tests/gtest_antalya_protocol.cpp @@ -28,6 +28,12 @@ TEST(AntalyaProtocol, RejectsInvalidMarkers) EXPECT_EQ(parseMarker(name), 0u) << "should not have parsed: " << name; } +TEST(AntalyaProtocol, RemoveMarkerRemovesOnlyAValidMarker) +{ + EXPECT_EQ(removeMarker(appendMarker("ClickHouse client")), "ClickHouse client"); + EXPECT_EQ(removeMarker("ClickHouse client (antalya:0)"), "ClickHouse client (antalya:0)"); +} + TEST(AntalyaProtocol, ParsesMarkerAndCapsVersion) { EXPECT_EQ(parseMarker("ClickHouse server (antalya:999999999)"), static_cast(DBMS_ANTALYA_PROTOCOL_VERSION)); diff --git a/src/Server/TCPHandler.cpp b/src/Server/TCPHandler.cpp index c85dc0847448..e890daed3a7a 100644 --- a/src/Server/TCPHandler.cpp +++ b/src/Server/TCPHandler.cpp @@ -235,7 +235,8 @@ void validateClientInfo(const ClientInfo & session_client_info, const ClientInfo if (session_client_info.interface == ClientInfo::Interface::TCP) { - if (session_client_info.client_name != client_info.client_name) + /// The Antalya marker is only in the Hello `client_name`, not in the Query packet. + if (AntalyaProtocol::removeMarker(session_client_info.client_name) != client_info.client_name) throw Exception( DB::ErrorCodes::CLIENT_INFO_DOES_NOT_MATCH, "Client info's client_name does not match: {} not equal to {}", @@ -1938,6 +1939,7 @@ void TCPHandler::receiveHello() } readStringBinary(client_name, *in, MAX_HELLO_STRING_SIZE); + client_antalya_protocol_version = AntalyaProtocol::parseMarker(client_name); readVarUInt(client_version_major, *in); readVarUInt(client_version_minor, *in); // NOTE For backward compatibility of the protocol, client cannot send its version_patch. diff --git a/src/Server/TCPHandler.h b/src/Server/TCPHandler.h index bdec8f161e33..95c83195cdfc 100644 --- a/src/Server/TCPHandler.h +++ b/src/Server/TCPHandler.h @@ -206,6 +206,7 @@ class TCPHandler : public Poco::Net::TCPServerConnection UInt64 client_version_patch = 0; UInt32 client_tcp_protocol_version = 0; UInt32 client_parallel_replicas_protocol_version = 0; + UInt64 client_antalya_protocol_version = 0; String proto_send_chunked_cl = "notchunked"; String proto_recv_chunked_cl = "notchunked"; String quota_key; diff --git a/tests/integration/test_antalya_protocol/test.py b/tests/integration/test_antalya_protocol/test.py index 81404c264cae..436c65466cbe 100644 --- a/tests/integration/test_antalya_protocol/test.py +++ b/tests/integration/test_antalya_protocol/test.py @@ -15,6 +15,7 @@ ) NEGOTIATED = "Antalya protocol: " +MARKED_CLIENT = "(antalya:[0-9]*) version" @pytest.fixture(scope="module") @@ -32,10 +33,12 @@ def count_in_log(node, substring): def test_remote_function_negotiates(started_cluster): initiator_before = count_in_log(node1, NEGOTIATED) + worker_before = count_in_log(node2, MARKED_CLIENT) assert node1.query("SELECT count() FROM remote('node2', system.one)") == "1\n" assert count_in_log(node1, NEGOTIATED) > initiator_before + assert count_in_log(node2, MARKED_CLIENT) > worker_before def test_new_initiator_against_an_unmarked_worker(started_cluster): @@ -45,4 +48,6 @@ def test_new_initiator_against_an_unmarked_worker(started_cluster): def test_unmarked_initiator_against_a_marked_server(started_cluster): + before = count_in_log(node1, MARKED_CLIENT) assert node_old.query("SELECT count() FROM remote('node1', numbers(10))") == "10\n" + assert count_in_log(node1, MARKED_CLIENT) == before