mirror of
https://github.com/qdrant/qdrant.git
synced 2026-09-28 00:47:32 -05:00
* 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>
151 lines
6.9 KiB
Python
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()
|