Files
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

29 lines
1.1 KiB
Python

import pytest
from .points_messages import UpdateBatchInternal, is_upsert_batch_for
@pytest.mark.parametrize("routes, expected", [
([("test_collection", 0)], True),
([("test_collection", 0), ("test_collection", 0)], True),
([("other_collection", 0)], False),
([("test_collection", 1)], False),
([("test_collection", None)], False),
([("test_collection", 0), ("test_collection", 1)], False),
([], False),
])
def test_upsert_batch_matches_collection_and_explicit_shard(routes, expected):
batch = UpdateBatchInternal()
for collection_name, shard_id in routes:
upsert = batch.operations.add().upsert
upsert.upsert_points.collection_name = collection_name
if shard_id is not None:
upsert.shard_id = shard_id
assert is_upsert_batch_for(batch.SerializeToString(), "test_collection", 0) is expected
def test_upsert_batch_rejects_unrecognized_operation():
# An operation containing only unknown fields must not match shard 0.
assert not is_upsert_batch_for(b"\x0a\x02\x1a\x00", "test_collection", 0)