Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 11 additions & 14 deletions docs/en/antalya/protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand All @@ -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.

Expand All @@ -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.
2 changes: 1 addition & 1 deletion docs/en/interfaces/specs/NativeProtocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
2 changes: 1 addition & 1 deletion src/Client/Connection.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
7 changes: 7 additions & 0 deletions src/Core/AntalyaProtocol.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,13 @@ UInt64 parseMarker(std::string_view name)
return std::min<UInt64>(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));
}

}

}
3 changes: 3 additions & 0 deletions src/Core/AntalyaProtocol.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);

}

}
6 changes: 6 additions & 0 deletions src/Core/tests/gtest_antalya_protocol.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<UInt64>(DBMS_ANTALYA_PROTOCOL_VERSION));
Expand Down
4 changes: 3 additions & 1 deletion src/Server/TCPHandler.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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 {}",
Expand Down Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions src/Server/TCPHandler.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
5 changes: 5 additions & 0 deletions tests/integration/test_antalya_protocol/test.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
)

NEGOTIATED = "Antalya protocol: "
MARKED_CLIENT = "(antalya:[0-9]*) version"


@pytest.fixture(scope="module")
Expand All @@ -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):
Expand All @@ -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
Loading