Files
qdrant/tests/consensus_tests/test_peer_proxy_cluster.py
Anton Antonov 4c4c3d967e test: improve consensus test suite (#10671)
* test: add a gated proxy for peer RPCs

Pause one selected internal request while other peer traffic continues.
Preserve payloads, metadata, deadlines, and cancellation so consensus
tests can control transfer timing without blocking unrelated requests.

Cover forwarding, independent gates, and cleanup with socket tests.

* test: connect peer proxies to consensus clusters

Let consensus tests route internal RPCs through request gates. Keep each
proxy alive across peer restarts so advertised addresses remain stable,
and close all proxies during test cleanup.

Wait for the upstream gRPC connection before returning from proxied
startup. Verify consensus progress during a held WAL-delta request,
recovery data, and restart behavior with both URI configuration modes.

* test: fix potentially misleading peer proxy method names

Explicitly state the guarantees, or lack of.

Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com>

* test: add support for hold_snapshot_download

Removes flakiness from snapshot-related consensus tests too

Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com>

* test: improve asserts when force deleting peer

Actually verify survivors recover and retain the expected data.
Making sure no data loss happens.

Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com>

* test: add OsError socket handling + explicit wal_delta tests

Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com>

* test: reject zero as a defined consensus leader

* test: recheck leader agreement on each poll

After a restart, the leader can change during election. Let the cluster
wait resample the leader on each poll and require agreement on a nonzero
leader before the snapshot test starts its transfer.

Keep explicit leader checks for existing callers, membership-size checks,
and the existing timeout. Cover election changes and offline peers.

* test: verify independent snapshot download gates

* test: use a positive peer connection deadline

* test: cover recovery after the removed source exits

* chore: add clarifying comment on timeout=0 usage

It's not obvious at first why it's like so.

Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com>

* test: share consensus response gates

Move response gates, their tests, and Raft decoding from the leader
removal proof into the base test infrastructure. Both removal scenarios
can then use the same successful-response check.

* test: support selective RPC blocking

Keep a removed source unaware of membership changes while its transfer
continues. Block its Raft traffic in both directions so election
attempts cannot disrupt survivor recovery.

* test: make source removal scenarios deterministic

Separate recovery after source exit from late data sent by a removed
source. Require a successful receiver response in the late scenario,
and retain complete data and replica-state checks in both cases.

* test: refactor timeouts and deadlines

* Cancellation happens after observing the intended phase, without an RPC deadline.
* Separate deadline tests cover held requests, upstream work, and held responses.

Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com>

* test: bound peer probes and removal requests

Give cluster probes and peer removal finite client timeouts so a stalled
HTTP request cannot leave the test waiting indefinitely.

* test: separate RPC release from termination

Keep the upstream handler blocked until the test releases it or the RPC
terminates. Use a separate termination event for cancellation assertions,
and release the handler during teardown instead of racing a fixture timer.

* test: use monotonic polling deadlines

Measure elapsed polling time with a monotonic clock so system clock
adjustments cannot shorten or extend the wait.

* test: allow more time to observe proxy events

Allow ten seconds for proxy observations and ordinary test requests.
Event and future waits still return as soon as they complete. Keep the
one-second expiry tests and document the HTTP deadline setup race.

* test: bound leader and replication requests

Limit how long leader lookup and transfer submission wait for an HTTP
response. A stalled submission must fail so the test can release its
transfer gates and clean up the peers.

* test: preserve readiness failures in diagnostics

Catch request failures while collecting cluster diagnostics, including
read timeouts. Report the original readiness failure instead of replacing
it with a diagnostic error.

* test: assert points calls for the correct collection

Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com>

* test: make sure check_cluster_size and check_leader cannot stall

Have an explicit timeout.

Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com>

* test: retry timeouts during initial leader lookup

Treat request timeouts as retryable while discovering the expected leader,
matching the subsequent leader and membership checks. Keep polling after
a transient timeout instead of aborting the cluster-status wait.

* test: ensure batch data is different

Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com>

---------

Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com>
2026-09-21 12:04:51 +03:00

151 lines
6.9 KiB
Python

import pytest
import requests
from .fixtures import create_collection, upsert_points
from .utils import *
COLLECTION_NAME = "test_collection"
RECOVERY_POINT = "/qdrant.CollectionsInternal/GetShardRecoveryPoint"
def peer_has_metadata_value(uri, metadata_path, expected_value):
response = requests.get(f"{uri}{metadata_path}", timeout=5)
assert_http_ok(response)
return response.json()["result"] == expected_value
@pytest.mark.parametrize("uris_in_env", [False, True], ids=["cli-uri", "env-uri"])
def test_peer_proxy_cluster_transfer_and_restart(tmp_path, uris_in_env):
peer_uris, peer_dirs, bootstrap_uri = start_cluster(
tmp_path, 3, uris_in_env=uris_in_env, use_peer_proxy=True,
)
peers = list(processes)
peer_ids = [get_cluster_info(uri)["peer_id"] for uri in peer_uris]
expected_addresses = {
str(peer_id): peer.proxy.uri for peer_id, peer in zip(peer_ids, peers)
}
assert bootstrap_uri == peers[0].proxy.uri
for uri in peer_uris:
assert {
peer_id: peer["uri"].rstrip("/")
for peer_id, peer in get_cluster_info(uri)["peers"].items()
} == expected_addresses
create_collection(peer_uris[0], shard_number=1, replication_factor=3, write_consistency_factor=3)
wait_collection_exists_and_active_on_all_peers(COLLECTION_NAME, peer_uris)
points = [
{"id": index, "vector": [1.0, 0.0, 0.0, 0.0], "payload": {"index": index}}
for index in range(100)
]
assert_http_ok(upsert_points(peer_uris[0], points))
with peers[0].proxy.hold_rpc(RECOVERY_POINT) as gate:
replicate_shard(peer_uris[2], COLLECTION_NAME, 0, peer_ids[2], peer_ids[0], method="wal_delta")
gate.wait_for_request()
transfer = get_collection_cluster_info(peer_uris[0], COLLECTION_NAME)["shard_transfers"]
assert len(transfer) == 1
assert (transfer[0]["from"], transfer[0]["to"], transfer[0]["shard_id"]) == (
peer_ids[2], peer_ids[0], 0,
)
assert transfer[0]["method"] == "wal_delta"
# This requires Raft progress while the selected transfer RPC is held.
metadata_path = "/cluster/metadata/keys/proxy-check"
assert_http_ok(requests.put(f"{peer_uris[1]}{metadata_path}?wait=true", json="while-held", timeout=10))
for uri in peer_uris:
wait_for(peer_has_metadata_value, uri, metadata_path, "while-held")
assert not gate.cancelled.is_set()
for uri in peer_uris:
wait_for_collection_shard_transfers_count(uri, COLLECTION_NAME, 0)
wait_for_all_replicas_active(uri, COLLECTION_NAME, min_local_replicas=1)
# A restart must retain its advertised address even without repeating the
# option. Otherwise surviving peers can bypass the proxy after the restart.
restarted_peer = peers[2]
restarted_peer.kill()
processes.remove(restarted_peer)
peer_uris[2] = start_peer(
peer_dirs[2], "peer_2_restarted.log", bootstrap_uri,
port=restarted_peer.p2p_port, uris_in_env=uris_in_env,
)
assert processes[-1].proxy is restarted_peer.proxy
wait_for_peer_online(peer_uris[2])
wait_collection_exists_and_active_on_all_peers(COLLECTION_NAME, peer_uris)
for uri in peer_uris:
assert get_cluster_info(uri)["peers"][str(peer_ids[2])]["uri"].rstrip("/") == restarted_peer.proxy.uri
local_shards = get_collection_cluster_info(uri, COLLECTION_NAME)["local_shards"]
assert len(local_shards) == 1
assert local_shards[0]["points_count"] == len(points)
response = requests.post(
f"{uri}/collections/{COLLECTION_NAME}/points/scroll?consistency=all",
json={"limit": 100, "with_payload": True, "with_vector": True},
timeout=10,
)
assert_http_ok(response)
assert response.json()["result"]["points"] == points
proxies = [peer.proxy for peer in peers]
kill_all_processes()
assert not peer_proxies
assert all(not proxy._thread.is_alive() for proxy in proxies)
assert all(proxy.port not in busy_ports for proxy in proxies)
assert all(proxy.http_port not in busy_ports for proxy in proxies)
@pytest.mark.parametrize("restart_receiver", [False, True], ids=["initial-start", "after-restart"])
def test_snapshot_download_gate_pauses_after_receiver_is_cleared(tmp_path, restart_receiver):
peer_uris, peer_dirs, _ = start_cluster(tmp_path, 3, use_peer_proxy=True)
peers = list(processes)
peer_ids = [get_cluster_info(uri)["peer_id"] for uri in peer_uris]
create_collection(peer_uris[0], shard_number=1, replication_factor=3, write_consistency_factor=3)
wait_collection_exists_and_active_on_all_peers(COLLECTION_NAME, peer_uris)
points = [{"id": index, "vector": [1.0, 0.0, 0.0, 0.0]} for index in range(100)]
assert_http_ok(upsert_points(peer_uris[0], points))
if restart_receiver:
peers[0].kill()
processes.remove(peers[0])
peer_uris[0] = start_peer(
peer_dirs[0], "receiver_restarted.log", peers[1].proxy.uri, port=peers[0].p2p_port,
)
assert processes[-1].proxy is peers[0].proxy
wait_for_peer_online(peer_uris[0])
wait_collection_exists_and_active_on_all_peers(COLLECTION_NAME, peer_uris)
# Readiness does not mean the restarted peer has learned the new leader.
wait_for_uniform_cluster_status(peer_uris)
with peers[0].proxy.hold_snapshot_download(peer_uris[2], COLLECTION_NAME, 0) as gate:
replicate_shard(peer_uris[2], COLLECTION_NAME, 0, peer_ids[2], peer_ids[0], method="snapshot")
gate.wait_for_request()
receiver = get_collection_cluster_info(peer_uris[0], COLLECTION_NAME)
transfer, = receiver["shard_transfers"]
assert (transfer["from"], transfer["to"], transfer["shard_id"], transfer["method"]) == (
peer_ids[2], peer_ids[0], 0, "snapshot",
)
shard, = receiver["local_shards"]
assert shard["state"] == "Recovery"
assert shard["points_count"] == 0
# Consensus must still apply new operations while download is held.
metadata_path = "/cluster/metadata/keys/snapshot-proxy-check"
assert_http_ok(requests.put(f"{peer_uris[1]}{metadata_path}?wait=true", json="while-held", timeout=10))
for uri in peer_uris:
wait_for(peer_has_metadata_value, uri, metadata_path, "while-held")
assert not gate.cancelled.is_set()
for uri in peer_uris:
wait_for_collection_shard_transfers_count(uri, COLLECTION_NAME, 0)
wait_for_all_replicas_active(uri, COLLECTION_NAME, min_local_replicas=1)
shard, = get_collection_cluster_info(uri, COLLECTION_NAME)["local_shards"]
assert shard["points_count"] == len(points)
response = requests.post(
f"{uri}/collections/{COLLECTION_NAME}/points/scroll?consistency=all",
json={"limit": 100, "with_payload": False, "with_vector": True}, timeout=10,
)
assert_http_ok(response)
assert response.json()["result"]["points"] == points
kill_all_processes()