mirror of
https://github.com/qdrant/qdrant.git
synced 2026-09-29 01:17:56 -05:00
edge-docs-diff
367
Commits
| Author | SHA1 | Message | Date | |
|---|---|---|---|---|
|
|
b50c0fe012 |
test: wait for replicas Active before snapshot layout assert (#10739)
Avoid racing peer 0's cluster view after kill+recover, where remotes can still be Partial while the recovered peer's local shards are Active. |
||
|
|
7c8dc0e08d | test: skip three known consensus failures (#10712) | ||
|
|
c8c8895beb |
test(consensus): prove leader removal can stall (#10684)
* test(consensus): prove leader removal can stall A departing leader can clear peer addresses before its queued commit notification reaches the surviving voter. Control the removal append and its acknowledgement to reproduce that ordering without fixed sleeps or stopping the leader process. Check the actual notification failure and the survivor's commit index. Keep follower-removal and three-voter controls, and cover both the default removal wait and timeout=60. * test: timeout=60s with follower or leader Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com> * test: fail when leader removal loses its commit Require the surviving voter to receive the committed removal and remain operational. Fail explicitly when the commit notification never reaches its gate instead of treating the stuck cluster as success. Both two-node leader-removal cases now fail against the existing bug. The follower-removal and three-node controls pass. --------- Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com> |
||
|
|
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> |
||
|
|
c4b2d1cadd |
test: wait for applied CommitRead before cleanup in test_resharding_deferred[up] (#10702)
The test compared raft commit indices across peers before issuing shard cleanup. A follower advances its commit index before it applies the entry, so cleanup could still race the follower applying CommitRead, which invalidates running clean tasks and makes the endpoint return 500. Wait for every peer to report the applied `read_hash_ring_committed` resharding stage via telemetry instead. Reading it takes the same shard holder lock as the consensus handler, so an observed stage is fully applied. Repurpose the unused stage check helper to read the stage from telemetry, since the `comment` field it read no longer exists. Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> |
||
|
|
f5e75477a0 |
Cleanup ConsensusStateMachine validation and docs (#10658)
|
||
|
|
1d0c2c1bb2 |
test: wait for consensus catch-up before snapshot recovery on a new peer (#10642)
test_recover_from_snapshot_2 and test_upload_snapshot_2 start snapshot recovery on a freshly joined peer as soon as it lists the collection. The collection appears once the creation entry is applied, while the peer is still replaying the rest of the raft log, including the removal of the killed peer. Recovery then decides which other replicas to remove or mark dead from that stale local view and drops a healthy replica, leaving a shard with a single replica. Add a helper that waits until all peers share the same commit index and have no pending operations, and use it in both tests before recovering. Also fix a misleading comment in the recovery replica cleanup branch. Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> |
||
|
|
52d8e39451 |
Remove RocksDB references (#10561)
* Remove RocksDB specifics from shell.nix * Re-enable sparse benches, replace RocksDB structures * Rewrite congruence test, in-memory ID tracker vs mutable ID tracker * Remove RocksDB flag from test * Remove RocksDB tool * Remove RocksDB comments * Bump OpenAPI spec |
||
|
|
8185a2692e |
Validate consensus against ConsensusStateMachine (#10469)
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
c1098b652f |
test: wait for receiver Active after snapshot transfer (#10330)
Finish can apply on peer 0 before the destination peer, so asserting Active on peer 2 right after peer 0 reports all-Active was flaky. |
||
|
|
f59480ef0c |
test: de-flake test_recover_from_snapshot shard placement assertion (#10310)
Mirror the fix from #9938 for test_upload_snapshot: assert every shard has n_replicas active replicas across local+remote shards instead of assuming peer 0 always sees exactly 2*n_replicas remote shards. |
||
|
|
995083e123 |
tests: restart reinit peer on the same port in readyz test (#10255)
`test_reinit_removed_peer_readyz_ignores_old_cluster` restarted the reinitialized peer on a fresh port. A changed `--uri` makes the peer announce its new address to every address-book entry, including the injected old first peer, which re-adds it to the *old* cluster as a learner and starts replicating its log to it. Normally the restarted peer is a term ahead and ignores those messages, but when the reinit run's hard state save did not finish before the kill, it restarted at the old term, accepted the old leader, took its log (commit 13 > 12) and failed the guard. The scenario is a plain restart, so keep the URI; then nothing is announced and only the `/readyz` membership filter is exercised. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
774f052130 |
test: harden flaky consensus rejoin and JWT snapshot upload checks (#10158)
Tolerate not-yet-ready collection upserts in test_rejoin_cluster and give JWT snapshot uploads more headroom while still bounding auth-rejection hangs. |
||
|
|
4bc241aee2 |
Remove flaky consensus_lag consensus test (#10133)
The profiler consensus lag integration test has been flaky in CI; drop it while keeping the feature. |
||
|
|
9e8282fb6a |
test(resharding): crash a follower, not the leader, in the scale-down revert test (#10104)
The crashing peer was picked positionally, so it could be the raft leader. The staging crash exits inside the apply of the `Dead` entry while raft messages leave through an async send queue, so a crashing leader takes the append carrying the new commit index down with it. The other live peer is then left holding that entry appended but uncommitted and, as one voter out of three, can never commit it: it never aborts resharding, and the test times out waiting for its resharding state to clear — the restart that would restore quorum only comes after that wait. Pick the victim after the receiver is killed instead: wait for the two live peers to agree on a live leader, keep that leader as the survivor and crash the follower. The survivor then commits and applies the abort locally, with no commit index left to escape a dying process. Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
0805d3f422 |
Add /profiler/consensus_lag to measure apply lag between peers (#10090)
* Add /profiler/consensus_lag to measure apply lag between peers Raft commit index advances on a peer whose apply loop is stalled, so the existing signals - `raft_info.commit` and the `all_nodes_have_same_commit` test helper - report a stuck peer as healthy. Nothing exposes how long a peer has been behind at *applying* entries, which is what shard transfer's `await_consensus_sync` barrier actually waits on. Each peer now keeps a ring of the last 32 entries it applied, stamped with its own wall clock and the time that entry took to apply. The ring is in memory on ConsensusManager, not in Persistent, so the on-disk format is untouched. `/profiler/consensus_lag` collects those rings from every peer over a new internal RPC and lines them up on the entry indices they share. Each entry is measured from whichever peer applied it first, so a lag is never negative; the peer that is first can differ per entry, so the baseline is per entry rather than a single chosen peer. Entries only one peer still remembers are excluded, otherwise a peer would be measured against itself. A peer stalled part-way through an entry keeps healthy lag statistics - everything it did apply, it applied on time - so the report carries `behind_entries` and `newest_applied_age_ms` alongside, which is what actually exposes the stall. Peers that fail or time out are listed rather than failing the request: a partial answer is more useful than none when the point is to find a peer that stopped answering. The endpoint follows `/profiler/slow_requests`: manage access, and outside OpenAPI, so no endpoint-count or ACTION_ACCESS guard applies. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Move applied-entry log into its own module Keeps the new code out of files that are already large. The ring, its entry type and the snapshot served over RPC move to `content_manager/consensus/applied_log.rs`, alongside the other consensus internals; `ConsensusManager` is left with a field, an accessor and the one `record` call in the apply loop. The grpc encoding moves next to the decoding it mirrors, in `common/consensus_lag.rs`, leaving the internal service handler three lines instead of thirty. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Test that a consensus stall is still in the report after the peer catches up * Take each peer's applied index from consensus state, not its ring --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-authored-by: tellet-q <elena.dubrovina@qdrant.com> |
||
|
|
c65d7d2d63 | test(resharding): test resharding state clears before the replica revert (#9656) | ||
|
|
aa6c5d8403 |
Global quota API (#10035)
* feat: global quota API Memory and disk are node-wide resources, so configuring their thresholds per collection through strict mode makes little sense. Move them behind a single cluster-wide `QuotaManager`. The quota config is seeded from `storage.quotas` in the settings (and so from env vars), overridden by `quota.json` in the storage directory, and updated cluster-wide through a new `SetQuotaConfig` consensus operation which rewrites that file on every peer. Raft snapshots carry it too, so a peer that joins by snapshot picks it up. Quotas are enforced wherever the strict mode memory and disk checks used to run, but no longer gated behind `strict_mode.enabled`: a value set in an enabled strict mode config still wins per resource, the quota is the default. Rejections name both the condition that tripped and the config that governs it. `GET /quotas` reports the config plus current utilization to global read users; `PUT /quotas` replaces it for global manage users. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: cover the quota endpoints in the API consistency checks `test_all_rest_endpoints_are_covered` and the OpenAPI endpoint count both break on any new REST endpoint. Add `GET`/`PUT /quotas` to `ACTION_ACCESS` with their JWT access tests, and bump the expected API count. The quota endpoints stay out of `REST_ENDPOINT_WHITELIST`: that list is for data-plane endpoints reported per-endpoint in metrics. Also add a Raft snapshot CBOR compatibility test — snapshots are exchanged between peers of different versions during a rolling upgrade, so `quota_config` must be absent-tolerant in both directions. Review feedback: persist through `SaveOnDisk`, which already implements the write-before-swap protocol this was doing by hand; validate the config at both persistence boundaries, since a hand-edited quota file or a config arriving through consensus does not pass the REST handler's validation, and a `0%` limit would reject every update forever. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * test: assert seeding a quota from invalid settings persists nothing Follow-up to review feedback claiming `SaveOnDisk::load_or_init` writes the init value before it is validated. It does not — only `SaveOnDisk::new` persists — but the property matters: were seeding to persist first, invalid settings would leave a `quota.json` that fails validation on every subsequent start, and the node could only be recovered by deleting it by hand. Pin it down with a test. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: make QuotaManager the single reader of memory and disk The quota checks measured memory and disk themselves, while the optimizer and the WAL disk watcher each called `fs4::available_space` behind their own ad-hoc caches. Fold all of it into QuotaManager: it owns the readings, the freshness policy, and the limits they are compared against. Moves the module to `lib/shard`, since the optimizer sits below `storage` and has to reach it; `storage::quota` re-exports it, so consensus, the `/quotas` API and StorageConfig are unchanged. The manager is installed as a process singleton by TableOfContent, ahead of loading any collection. - Callers hand in QuotaLimits overrides instead of a StrictModeConfig, and an override can now only tighten. A collection-level admin could raise `max_disk_usage_percent` past a cluster-wide limit that needed global manage rights to set; ties resolve to the quota so the rejection names the knob that actually has to change. - Measurements are cached for 5s, but a reading at or above its limit is never reused: a rejected client retries, and freeing the resource has to take effect on the next request rather than a TTL later. - `fits_on_disk` sizes an optimization against physical free space only, never the configured limits. Optimizations are what free a full disk, so the quota must not be what stops one. - `percent_of` widens to u128 instead of saturating the multiply, which under-reported utilization (failing open) above ~184 PB. - StorageConfig::quotas is optional; absent means no quota is enforced. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: don't recover dead replicas onto a node at a resource limit Recovering a dead replica pulls a whole copy of its shard onto this node. If it is already at its memory or disk quota that transfer cannot finish, and starting it only pushes the node further past the limit. Skip it and reconsider on a later sync, once the resource frees up. Adds QuotaManager::check_capacity for work that lands bytes here without being an update. Unlike fits_on_disk the configured limits do apply: taking on a replica is not what frees a full node, so there is no deadlock to avoid by letting it through. The check is hoisted out of the per-shard loop because a node over its limit re-measures on every call, so checking per dead shard would cost a statvfs each. It is free when no quota is configured. Also trims the comments across the quota module, which had grown well past what the code needs. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * test: drop trivial and duplicated quota tests Six tests removed, ~140 lines, with no loss of coverage: - a_rejection_names_the_knob_that_has_to_change asserted that a format! contains its own literals; the message is covered end-to-end by the override test and by test_global_quota.py. - a_node_over_its_quota_has_no_capacity_to_take_on_a_replica was 30 lines for check_capacity, a one-line delegation to check_update the test above it already calls. - a_rejecting_measurement_is_never_served_from_the_cache duplicated the meter test, which proves the same rule with an injected reader instead of inferring it from the real filesystem. - free_space_is_reported_without_enforcing_anything covered a one-line accessor, and its point is what the fits_on_disk test is for. - The two resolve tests and the three meter tests each collapse into one. DiskFit::Unknown keeps its coverage as two lines inside the fits_on_disk test rather than its own fixture. The snapshot compat pair becomes one test: the second only asserted cluster_metadata.is_empty(), which says nothing about quotas — the real check was the deserialize, now an expect that states it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: re-measure free space as the disk fills, and drop a Windows-only assert Two CI failures, both from this branch. e2e test_low_disk: the DiskUsageWatcher I replaced escalated to checking on every call once free space fell below 512 MB. Folding it into the quota manager lost that — available_bytes passed no limit, so a reading was reused for the full 5s however little space was left. On a disk filling as fast as that test fills it, 5s blind is enough to actually run out and the WAL write dies instead of returning "No space left on device". available_bytes now takes a watch_below level and never reuses a reading under it, which is what the old ladder was expressing. The watcher passes max(min_free, 512 MB), so the escalation point is back; above it the 5s cache still costs fewer syscalls than the old 128-call ladder. fits_on_disk gets the same rule by passing required_bytes, so a merge that does not fit re-checks rather than sitting on a stale sample. Windows: fits_on_disk on a missing path was asserted to be Unknown, but GetDiskFreeSpaceEx resolves up to the containing drive and succeeds — as common::disk_usage's own test documents. Dropped; the branch is a two-line else and is not portably reachable. Also renames an_optimization_is_sized_against_the_disk_not_the_quota, which needed explaining to be understood. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: report the quota config in telemetry Reads it from the quota manager rather than the settings, so it is the config the node is actually enforcing: a peer that missed a consensus update reports what it is applying, not what the cluster agreed on. Gated on global access, the same access `GET /quotas` requires, and left out of `PeerTelemetry` — a quota is per-node state, so each peer reports its own. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: regenerate OpenAPI, and cover the quota in the telemetry key sets Two CI failures from the previous commit. Referencing QuotaConfig from TelemetryData moves its definition earlier in `components/schemas`, because TelemetryData is generated ahead of QuotaStatus. Regenerated rather than hand-patched, so the schema is a pure move. test_telemetry_detail asserts the exact set of top-level telemetry keys. The quota is reported at every level, including 0 — it is three scalars, it is the default the endpoint serves, and it is what explains an update being rejected — so both key sets gain it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: remove `max_disk_usage_percent` from strict mode Disk is a node-wide resource, so a per-collection percentage of it never meant anything a caller could act on: the limit describes how full the *node* is, and which collection the write happens to target has nothing to do with it. The global quota is where it belongs. It shipped in 1.18.2 without documentation, so this drops it outright rather than deprecating. Removal is soft in every direction: StrictModeConfig has no `deny_unknown_fields`, so a client still sending it gets it ignored rather than a 400, and the same struct deserializes the persisted collection config, so collections created on 1.18.2+ keep loading. Proto field 22 is reserved so the number is never reused. The e2e test becomes a quota test — the fixture and the timing are the interesting parts and they carry over unchanged; only how the threshold is configured differs. `max_resident_memory_percent` was documented and stays for now. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: enforce the strict mode memory limit outside the quota `max_resident_memory_percent` was folded into the quota as an override, which meant the quota check had to know about strict mode, and retiring the setting would mean unpicking `EffectiveLimit` and `LimitSource` from the resolution logic. It is now a check of its own in `verification/mod.rs`, next to the strict mode checks it belongs with, borrowing only the measurement from the quota manager — which stays the node's single reader of process memory, so both checks still share one reading. Deleting the setting later is deleting one function and its one caller. `QuotaManager::check_update` takes no arguments and consults the quota alone. A collection can still only tighten the limit for itself, because its own check runs in addition rather than in place of the quota's, and each rejection now names the config that has to change without having to carry a `LimitSource` to say so. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: enforce the quota on the update path, not in strict mode The quota check sat inside `check_strict_mode_toc_batch` only because that was the one place holding the collection's strict mode config. It doesn't need one any more, and the placement had a real cost: coverage depended on each handler remembering to ask for a strict mode check, and four of the internal update RPCs do — `sync_internal`, which moves the most bytes onto a node, does not. It now runs in `Collection::update_from_client` and `update_from_peer`, which every update passes through. `update_from_client` checks ahead of the shard split, so an operation is accepted or refused whole rather than landing on some shards and being refused by others. Classification moves with it, from ~10 `consumes_memory` impls on request DTOs to one exhaustive `CollectionUpdateOperations::consumes_quota`. The internal enum has variants — raw upserts, conditional upserts, the syncs — that have no client-facing request type, so per-DTO impls structurally could not classify them. Shard-transfer syncs stay excluded, as they are today: a transfer is sized up once before it starts, and refusing its batches partway abandons work that is nearly done only for it to restart from the beginning. Index and named-vector creation reach shards through consensus, past this check — a peer must not refuse what the cluster agreed to — so they keep their pre-consensus check, now against the quota manager directly. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: deprecate `max_resident_memory_percent` in strict mode Same reason the disk threshold went: memory is node-wide, so a per-collection percentage of it caps how full the *node* is, which has nothing to do with which collection is being written to. The node-wide quota caps it once for everything. Unlike the disk threshold this one shipped documented, in 1.18.0, so it keeps working — as a limit a collection can tighten for itself, never lift — and gets the usual markers: `#[deprecated]` on both Rust structs, `[deprecated = true]` on proto field 21, and `deprecated: true` in the OpenAPI schema, which schemars derives from the attribute. The note names 1.21 as the removal. Recording a version matters here: the audit in docs/plans/overdue-deprecations.md found that this repo has never written a removal deadline down, and members of the 1.15.0 deprecation batch are still in tree. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: reconcile the quota readers with #9891 #9891 landed effective (cgroup) figures in telemetry while this branch was making QuotaManager the single reader of memory and disk. Two collisions, neither of which git sees. `segment::utils::mem::total_memory_bytes` is now a shared accessor with a 5s TTL, so a cgroup resize is picked up. The quota module had its own `OnceLock` copy that froze the value at startup — exactly what #9891 set out to fix — so it delegates to the shared one instead. Telemetry's new `disk_size` called `common::disk_usage::disk_usage` directly. That reader lost its TTL cache on this branch when the caching moved into the quota manager's meter, so it would have taken an uncached `statvfs` on every telemetry request, and it put a second disk reader back in the tree. It goes through `QuotaManager::disk_capacity_bytes` now, sharing the reading the quota check already takes. Verified it still reports the storage filesystem, matching `df`. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: split the quota manager by what each half does `manager.rs` had grown to 450 lines holding three separate jobs: owning the config and its file, taking the readings, and comparing one against the other. - `manager/store.rs` — the `Store` enum, `QUOTA_CONFIG_FILE`, and config validation, which is now the store's own business rather than something every caller has to remember to do first. - `manager/measure.rs` — every reading, and `DiskFit`. The "nothing else calls `statvfs` or reads process RSS" claim is now checkable by looking at one file. - `manager/enforce.rs` — `check_update` / `check_capacity` and the threshold comparison. - `manager/mod.rs` — the struct, its construction, and the config accessors: what a reader needs to see first. Tests move with their subject. No behaviour change: `set_config` used to validate before delegating to the store, and now the store validates on write, which is the same order of operations from the outside. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: count copy-on-write deletes toward the quota Dropping a vector or a payload key does not free anything on its own: copy-on-write rewrites the point to produce the version without that field, so storage grows first and is only reclaimed once the optimizer gets to it. Gating those as if they were reclaiming space let a full node keep taking writes that make it fuller. Deleting whole points stays exempt. That is the one operation that has to work on a node at its limit, or there is no way back under it. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * docs: say that the quota's reported usage is per node `GET /quotas` returns one cluster-wide config and one set of utilization figures, which reads as though both describe the cluster. They do not: memory and disk are node-local, so `usage` is whatever the peer that served the request is seeing, and a peer under its limit says nothing about the others. Also corrects `resident_memory_percent`, which claimed to be a share of total system memory. It is a share of the memory available to the process, which under a cgroup is the limit rather than the host's RAM. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: treat a node over its quota as a failed replica, not a bad request A quota rejection described the request as invalid (400) and was classified non-transient, which is how the replica set recognises errors that every replica would produce alike. A quota is the opposite: the input is fine and the answer depends on which machine you ask. On the default `wait=false` path that combination silently dropped the write — `update.rs` only deactivates transient failures when nothing completed — leaving the replica Active and permanently missing data its co-replicas had. It is now `InsufficientStorage`, transient, HTTP 507 / gRPC `ResourceExhausted`. So a node that is out of room is handled like one that is offline: - last active replica, or every replica over quota: nothing could take the write, and the client is told the cluster is out of room. - more than one replica: the full node is deactivated through the same path a dead peer takes, and the update stands if enough replicas accepted it. `check_capacity` already keeps recovery off that node until it has room. The check also moves off `update_from_client`, which applied the coordinator's own limit to the whole operation even when it held no replica of the shards being written. Each replica set now gates its own local write and records the refusal as a failure of this peer, so a node only ever answers for itself. `ResourceExhausted` is shared with rate limiting, and the reverse conversion mapped it straight to `RateLimitExceeded` — a forwarded rejection came back as 429. Statuses now carry a marker so the two stay distinguishable across the wire. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: report quota pressure per node, and across the cluster A quota is node-local, so finding out which node has hit one meant asking each of them in turn — and nothing at all showed up in monitoring. `/metrics` gains a `quota_exceeded` gauge for the local node. It is emitted only while the quota is enabled: with it off the value would be a constant 0 that says nothing about the node, and an alert built on it would go quiet rather than fire if someone disabled the quota. Telemetry's `quota` field carries the same verdict alongside the config, since that is where the metric is derived from. `GET /quotas` now answers for the whole cluster. A new `GetQuotaUsage` RPC on the internal `QdrantInternal` service returns what one peer is using, and the handler fans it out to every known peer in parallel. Peers that do not answer are left out rather than failing the request — the nodes that are out of room are exactly the ones most likely to time out, and a partial answer still names them. Outside distributed mode the field is absent rather than a map of one. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: report the quota metric per resource `quota_exceeded` was one flag for the whole node, which does not say what to go and fix — disk is freed by deleting or optimizing, memory by unloading. It now carries a `resource` label: quota_exceeded{resource="memory"} 0 quota_exceeded{resource="disk"} 1 A resource with no limit gets no series at all, for the same reason the metric is absent while the quota is disabled: a series that can never reach 1 reads as healthy and would quietly carry an alert that cannot fire. `QuotaManager::exceeded` returns the per-resource verdict, with `None` for a resource this node does not cap. Telemetry reports the same breakdown, since the metric is derived from it. The peer usage RPC keeps a single flag — it sits next to both percentages, so it only has to answer "is this peer refusing writes". Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix: drop a no-op error conversion the linter caught `check_global_access` already returns a `StorageError`, so mapping it through `StorageError::from` converted the type to itself and tripped `clippy::useless_conversion`. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: hold a tripped quota until usage clears a release margin A resource resting on its limit crosses it in both directions on the noise between two readings, and each crossing is expensive: the node refuses a write, its replica is deactivated, usage dips, recovery starts sending a whole shard copy back, and the arriving data pushes it over again. The loop sustains itself, and every lap costs a shard transfer. A limit now trips at its configured value but only clears once usage has fallen 5 percentage points below it, so the crossing has to be real. The margin is floored at 1%, since a limit smaller than the margin would otherwise be impossible to fall back under and would strand the node. The verdict is carried on the manager rather than recomputed, which makes it the thing reporting shows: expect `exceeded` to be set while the utilization next to it is already back under the limit. Rejections say so too, rather than claiming a limit that is no longer exceeded: Disk usage is at 87% of total capacity. It reached the configured limit of 90% and has to fall below 85% before this node takes writes again. Changing the config clears the verdicts. New limits are a deliberate act, and should not be held back by the margin of a limit that no longer exists. Both resources are now evaluated on every check instead of stopping at the first failure, so a verdict is never left behind reporting a reading that has since been superseded. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * feat: make the quota release margin configurable 5 points is a guess about how noisy a deployment's usage is, which is not something one number can be right about: a node whose disk moves in gigabyte steps needs a wider margin than one that creeps, and an operator who wants the old flip-on-every-reading behaviour should be able to ask for it. `release_margin_percent` joins the rest of the quota config, so it seeds from `QDRANT__STORAGE__QUOTAS__RELEASE_MARGIN_PERCENT`, replicates through consensus, and changes with `PUT /quotas`. Defaults to 5 and is filled in when a request omits it, so it always answers with the margin actually in force rather than leaving the caller to assume one. `0` releases as soon as usage is back under the limit. `QuotaConfig` grows a hand-written `Default` for it, since deriving one would have quietly defaulted the margin to 0 and disabled the hysteresis for anyone constructing a config in code. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * refactor: leave the release margin unset by default, and hold verdicts in atomics `release_margin_percent` is `null` unless someone sets it, rather than materialising 5 into every config. A quota written today then does not pin a number a later release may want to revise, and `{"enabled": false}` still round-trips as itself. `QuotaConfig::limits` resolves it, next to `enabled`, so enforcement never sees the unset case. The verdicts move from a `Mutex<QuotaExceeded>` to one `AtomicBool` per resource. They are judged independently and nothing reads them as a pair, so the lock only added contention to the path every update takes; a verdict that races a concurrent check is re-decided by the next one from a fresh reading. That also drops the tri-state. Only "was this over its limit" has to survive between checks — whether a resource is enforced at all follows from the config and the reading, so it is derived when reporting rather than stored. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * docs: drop a comment arguing with a design that was never here Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
cb18bd7e4c |
test: wait for raft leader before collection recovery (#9972)
POST /cluster/recover can return 200 while raft silently drops the snapshot request when no leader is known yet, leaving the test stuck on a missing collection until timeout. Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
69186edec9 |
test: fix flaky test_partial_snapshot optimizer race (#9951)
Wait for green on write (and read after recover_read) so collection and partial snapshots are not taken mid-indexing. Otherwise a leftover appendable segment survives partial merge and breaks manifest equality. Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
6ba5e11aa5 |
test: tolerate transfer race in corrupted snapshot recovery (#9947)
#9013 skipped the manual replicate_shard when a transfer was already visible, but the recovery loop can still start one between that check and the POST. Accept 400 "already involved in transfer" as success so the remaining wait assertions still cover recovery. Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
c94c5fa4bc |
test: harden test_routing_token_sticky_reads against post-recovery flakiness (#9937)
Right after the no_sync snapshot recovery, the recovered replica serves local
reads immediately, but a remote read to it can transiently fail for a short
window. The read path then falls back to the other replica in hash order on
just the requesting peer, so a single routing token momentarily resolves to
different replicas across peers (observed as {A, B, B}), failing the
determinism assertion.
Wait until token-routed reads are stable across all peers for every token the
test asserts on before measuring, so the transient post-recovery fallback
window is passed. Pure test-side change; routing behaviour is unchanged.
Co-authored-by: Cursor <cursoragent@cursor.com>
|
||
|
|
e64c20a48a |
test: de-flake test_upload_snapshot (robust to shard placement) (#9938)
* test: make test_upload_snapshot robust to shard placement balance The final assertion in recover_from_uploaded_snapshot assumed a perfectly balanced shard placement (peer 0 having exactly 2*n_replicas remote shards). Shard placement across peers is not guaranteed to be balanced, so this made the test flaky (e.g. peer 0 ended up hosting all shards locally, leaving only 3 remote replicas instead of 4). Instead, verify the full replica layout is healthy: peer 0 observes every replica through its local + remote shards, so assert that all replicas are Active and every shard has exactly n_replicas copies across the cluster. Co-authored-by: Cursor <cursoragent@cursor.com> * test: fetch cluster info once for shard validation Read local and remote shards from a single /cluster response so both lists come from the same cluster revision, per review feedback. Co-authored-by: Cursor <cursoragent@cursor.com> --------- Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
b98443c2f9 |
Clean up stale shard transfers when applying consensus snapshot (#9928)
A transfer source that misses the transfer abort (e.g. while partitioned
or paused) keeps its local shard wrapped in a proxy. When such a peer can
only catch up via consensus snapshot, snapshot application re-creates
payload indexes with an update operation that the stale forward proxy
forwards to a transfer target which may no longer have the shard. The
resulting precondition error fails snapshot application and stops the
consensus thread ("No target shard N found for update"), leaving the
peer unable to ever catch up.
Snapshot application now explicitly cleans up transfers that are no
longer registered in consensus: the transfer task is stopped and the
proxy is reverted via the new `ShardReplicaSet::discard_proxy_local`,
which is infallible, never contacts the remote, and forgets queued
updates (replica states in the same snapshot already reflect the
transfer outcome).
The consensus test reproduces the incident: pause the transfer source
mid-transfer, restart the other peers so the aborted transfer can only
be learned via snapshot, and verify the source recovers.
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
|
||
|
|
53dfa5b022 |
Fix resharding, on queries filter shards on all shard selectors (#9882)
* Fix resharding, on queries filter shards on all shard selectors * Add failing consensus test: search during resharding with shard keys (#9880) Reproduces a known bug: after resharding is initialized on a custom sharded collection with a shard key, searches (with and without the shard key selector) fail with "does not have enough active replicas", because the new resharding shard is included in reads before it has an active replica. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> * Exempt explicit shard id selection from resharding read filter Explicit shard id selection is only used by internal per-shard operations (local shard API, internal gRPC reads), including the resharding driver reading back migrated points from the new shard. These must reach the resharding shard before it becomes visible to user-facing selectors, and filtering them also made per-shard reads return silently empty results on peers lagging on hashring commits. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Explicitly set resharding filtering per match branch --------- Co-authored-by: Andrey Vasnetsov <andrey@vasnetsov.com> Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
db2a135203 | raw vector grpc send (#9843) | ||
|
|
2735d40ecc |
Reset first_voter and prune address book on first-peer --reinit (#9785)
* Reset first_voter and prune address book on first-peer --reinit
A peer removed from consensus and killed after `RemoveNode(self)` was
committed but before it was applied keeps the old cluster's
`first_voter` and peer addresses in `raft_state.json`. First-peer
`--reinit` reset `conf_state` to a single voter (itself) but left both
untouched (`first_voter` is in fact never reset, even when the removal
is applied).
Both values are served to peers bootstrapping onto the reinitialized
cluster. A joining peer seeds its initial `conf_state` with the
advertised `first_voter`, and Raft conf-changes are deltas on top of
that base - so a stale `first_voter` permanently corrupts the joining
peer's voter set: it ends up with {old first peer, itself}, missing the
actual leader. Its `/readyz` then treats the still-alive old peer as a
cluster member and waits for the old cluster's commit index, which its
own consensus never reaches.
This only manifests when the reinitialized leader replicates its log as
plain entries (nothing applied before the kill, so the log anchor is
index 0). If the leader sends a snapshot instead, the snapshot's full
`conf_state` heals the corrupted seed - which is why
test_reinit_removed_peer only failed sporadically on CI.
Fix first-peer `--reinit` to behave like founding a fresh cluster:
reset `first_voter` to this peer (not `None`, or `recover_first_voter`
would re-derive the old first voter from the retained Raft log) and
prune `peer_address_by_id` to this peer only.
The kill-before-apply state is now injected deterministically into
test_reinit_removed_peer, reproducing the exact CI failure against the
unfixed binary. Since `--reinit` now prunes the address book, the
stale-address injection in
test_reinit_removed_peer_readyz_ignores_old_cluster moved to a restart
without `--reinit`, so it keeps exercising the /readyz `conf_state`
membership filter from #9688.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* Fix test_reinit_consensus expecting stale address book after --reinit
The test waited for cluster size 2 right after starting the reinitialized
first peer, before the second peer was even started. That only passed
because first-peer --reinit used to keep the old cluster's addresses in
the address book - the stale state the previous commit removes. A
reinitialized first peer is a fresh single-member cluster; the other
peers re-join and re-register their new addresses right after.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
|
||
|
|
506ccc7f5c | test: fix flaky test_recover_from_snapshot version comparison (#9818) | ||
|
|
46f6e7da31 |
tests: stop uploaders cleanly before consistency check in WAL delta tests (#9817)
The end-of-test teardown killed the uploader processes and slept for one second before scrolling all peers for the consistency check. Killing the client does not cancel an already-sent upsert server-side: on slow CI the last PUT can take longer than the sleep, so its replication is still propagating while the peers are scrolled at slightly different times, making the scrolls diverge by the last batch of points. Observed in test_shard_wal_delta_transfer_abort_and_retry: peer 0 was scrolled at 19.089s, received the forwarded batch at 19.179s, while peer 1 (which had already applied it locally) was scrolled at 19.267s. Use stop_update_process() (introduced in #8713 for the pre-peer-kill case) so the uploader exits between requests. Since uploads use wait=true, once the last PUT returns all active replicas have applied it, and the scrolls can no longer race. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
19fe868387 |
test(consensus): de-flake replace-peer-same-uri tests (#9758)
Wait for the collection metadata to propagate to the newly added extra peer before querying its collection cluster info. Being online and present in consensus does not guarantee the peer has already applied the collection-creation Raft entry locally, so get_collection_cluster_info could race and return 404. Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
bc7207b230 |
test(consensus): de-flake replicate_points_stream_transfer_updates override case (#9755)
* test(consensus): de-flake replicate_points_stream_transfer_updates override case With override_points=True the background writer re-upserts points 9990-9999, re-rolling their city payload. Points flipping away from "London" legitimately drop out of the filtered count on both shards, so asserting dest_filtered_count >= original snapshot count is not a valid invariant. On a slow CI runner the sleep(1)+kill() stopped the writer right after the overrides, before new inserts could compensate, making a net-negative flip likely (observed: 4954 >= 4959 failure). Replace the blind kill with a bounded workload (60 points) joined cleanly before the ~10s transfer of 10k points can finish, assert the writer exit code, and allow the filtered count to drop by up to the number of overridden points. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * test(consensus): guard that writer finishes while transfer is running The exact count consistency check requires every concurrent write to go through the transfer proxy. Make that precondition explicit: if the transfer ever finishes before the writer, fail with a clear message instead of a confusing count mismatch. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * test(consensus): make replicate_points update consistency checks exact Assign the city payload deterministically by point ID parity so filter membership can never change under concurrent overwrites. All assertions become exact ID-set comparisons with no slack: random city re-rolls made count-based checks unsound, since the forward proxy filters forwarded updates by post-update state and a point flipping out of the filter legitimately goes stale or missing on the destination. Replace the background writer process (sleep/kill/join choreography) with synchronous wait=true upserts issued while the transfer streams the initial points. Leave low point IDs unoccupied and insert into them during the transfer: the stream cursor passes them immediately, so these points can only reach the destination through live update forwarding, which the previous layout (writes at the stream tail) never verified. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
7b339f4643 |
Resolve filter-based update operations to point ids before WAL write (#9678)
* Resolve filter-based update operations to point ids before WAL write Filter/condition-resolving operations (delete-by-filter, conditional upsert, the *-by-filter payload/vector operations) stored their filter in the WAL and re-resolved it against live segment state on every apply. Replay-time state can differ from the original apply-time state (the optimizer drops deleted points and their version records during compaction), so WAL replay was not a deterministic function of the log and could resurrect filter-deleted points. Resolve such operations into concrete point ids at submit time, under a fence that guarantees the resolution sees exactly the operations that precede it in WAL order. The WAL now only ever contains id-based operations (pre-existing variants only — no format change), so replay applies the exact same point set as the original run. Fixes #9575 Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01DfVDMFQcy9Ww791x8sHobW * Fix rustfmt Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01DfVDMFQcy9Ww791x8sHobW * Drop coordinator-side resolution: every replica resolves locally Replicas holding the same data resolve the same filter to the same point set, and replicas that already diverged would not become consistent by agreeing on a filter's resolution. Forward the original filter operation as usual and let each replica's submit fallback resolve it under its own fence — one uniform path regardless of where the update lands. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01DfVDMFQcy9Ww791x8sHobW * Guard against is_filter_resolving / resolve_operation drift A resolved operation must never still classify as filter-resolving, otherwise a filter-carrying record could reach the WAL again (#9575). Catch one direction of drift between the gate and the rewriter with a debug assertion right after resolution. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Dedup points-vs-filter precedence into resolve_points_or_filter The "explicit id list wins over the filter" rule was written twice on the resolver side (DeletePayload arm and resolve_set_payload); a future tweak landing in one copy only would make SetPayload and DeletePayload silently diverge in what gets persisted to the WAL. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Assert rewritten WAL record reuses the incoming clock tag The single-record-reuses-the-tag property is what WAL-delta recovery and replica dedup rely on, but no test asserted it: submit the delete-by-filter with a real clock tag and check the resolved DeletePoints record carries it. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Test replay of old-style filter records left in the WAL Upgraded nodes can still hold WALs with unresolved filter operations; the by-filter apply paths are kept so they replay one final time with the old semantics. No test covered that path (the new submit flow can no longer produce such WALs), so append a raw DeletePointsByFilter record at the WAL layer, reload, and assert the matched points are gone. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Add consensus test for per-replica filter-op resolution Exercises the replicated path for filter/condition-resolving updates: the coordinator forwards the original filter op and each replica resolves it locally (delete-by-filter, insert-only and update-filter conditional upserts, set-payload-by-filter, including per-shard empty resolutions on a 2-shard collection). Asserts both replicas hold identical state (reads prefer the local replica), then restarts the whole cluster and asserts each replica replays its id-based WAL to the same state. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> Co-authored-by: Arnaud Gourlay <arnaud.gourlay@gmail.com> |
||
|
|
03f97d4c06 |
Fix /readyz of reinitialized peer waiting on a foreign consensus (#9688)
Applying `RemoveNode(self)` prunes all other peers from the removed peer's persisted address book, but the process may be stopped after the removal is committed and before the entry is applied. The old first peer's address then survives `--reinit`, and the readiness checker - which treated every `peer_address_by_id` entry as a cluster member - would wait for the reinitialized peer to reach the *old* cluster's commit index: a foreign consensus it can never catch up with, so `/readyz` never passed. Filter the address book by current `conf_state` membership instead, falling back to all known addresses while `conf_state` is still empty (a bootstrapping node that has not applied any configuration change yet). After `--reinit` the `conf_state` is reset to a single voter, so the readiness check correctly ignores peers of the old cluster. Fixes flaky `test_reinit_removed_peer`, which hit this race when the removed peer was killed before applying its own `RemoveNode`. The new regression test simulates that state and fails with the exact CI error without the fix. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
0346ea66cd |
fix: unblock optimizer after deleting a named vector (#9641)
* fix: unblock optimizer after deleting a named vector Deleting a named vector could permanently block the config-mismatch optimizer. The source-superset check in SegmentBuilder::update cancelled every rebuild that found the deleted vector still in old segment files, and each retry cancelled again, so optimizations got stuck forever. Removing the check (as in #9609) would fix delete but reintroduce data loss for the CreateVectorName race. Instead, tell the two cases apart with the live collection schema: prune a source vector that is gone from the schema (a real deletion), but cancel when it is still present (a freshly created vector this optimizer has not yet seen). This is safe because the schema is persisted before the op reaches segments, and the live schema is read after the source segments are frozen. The live set covers dense and sparse vectors, since a segment stores both together. When no live source is wired in, the conservative always-cancel behavior is kept. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * fix: wire live vector names into edge optimizers Deleting a named vector left the edge path with the pre-fix behavior: segment_optimizer_config hardcoded live_vector_names to None, so a merge touching a segment that still carried the deleted vector cancelled, and EdgeShard::optimize() propagated the cancellation as a hard error forever. Share the shard config behind an Arc and hand the blocking optimizers a provider that reads the current vector names on every call. Same safety argument as the server wiring: update() holds the segments read guard across both the segment application and the config update, so any name a frozen source segment carries is visible to the live read. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * refactor: share vector-name enumeration via CollectionParams::vector_names The optimizer's live-schema set and the WAL-recovery valid-name set are the same dense+sparse enumeration and must stay in lockstep; a drift between them would reintroduce a wrong prune/cancel decision. Replace the private helper in optimizers_builder and the inline block in WAL recovery with a single CollectionParams method. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * refactor: drop SegmentOptimizer::live_vector_names forwarding hop The default trait method only forwarded to the config getter and had a single caller; ShardOptimizationStrategy now reads the config directly, removing one layer of indirection. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> |
||
|
|
45634336ba |
Fix flaky count check in snapshot transfer missing-point test (#9685)
Killing the background load processes does not cancel requests already executing server-side: a wait=true upsert accepted just before the kill can still be propagating to the second replica while the test counts points, so peers transiently observe different totals (e.g. [20120, 20118, 20118] on CI). Poll the exact counts until they converge instead of asserting on the first sample. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
b590f929ce |
Fix reinit failing with "removed all voters" on a removed peer (#9654)
* Fix reinit failing with "removed all voters" on a removed peer When `--reinit` is run on the first peer, `conf_state` is reset to a single voter (this peer), but the Raft log is left untouched. If this peer had been removed from consensus before reinit, its log still holds a committed-but-unapplied `RemoveNode(self)` conf-change. On startup that entry is replayed on top of the freshly reset single-voter config, and Raft aborts with "removed all voters", so the node can never start. Resetting only the apply-progress queue is not enough: a fresh single-node leader re-commits and re-applies any log entries still physically present beyond `commit`, re-triggering the failure. The stale tail has to be physically dropped from the WAL. On the first-peer reinit path, discard committed-but-unapplied entries inherited from the previous cluster: truncate the WAL to the last applied index, pin `commit` to it and clear the apply-progress queue. The entry at `commit` is retained as the snapshot anchor, so the first peer can still serve snapshots to bootstrapping peers. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Add regression test for reinit of a removed peer Starts a 2-node cluster, gracefully removes the second peer from consensus (so it commits `RemoveNode(self)` into its own WAL), then reinitializes it as a fresh first peer. Before the fix this panicked on startup with "removed all voters"; the test asserts the peer comes back online, elects itself leader, and can still seed a fresh bootstrapping peer. Verified the test fails against the pre-fix binary with exactly: Failed to apply configuration change entry Caused by: Error in Raft consensus: removed all voters Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> |
||
|
|
f3bb3f43f4 |
Restart consensus thread on failure instead of stopping permanently (#9662)
* Restart consensus thread on failure instead of stopping permanently If the consensus loop failed (e.g. due to a transient I/O error such as running out of disk space), the consensus thread stopped permanently and the only way to recover was a full service restart. Now the consensus thread rebuilds the Raft node from persisted state and retries with exponential backoff (1s initial, doubled up to 5 min cap, retrying indefinitely). Rebuilding from persisted state is equivalent to a process restart, so this introduces no new recovery semantics. The `reinit` logic is never repeated on restart, and initial startup remains fail-fast. The consensus message channel is created outside of `Consensus` so that internal gRPC handlers and the forward-proposals thread keep their senders across restarts, and messages buffered during the outage are drained after recovery. While in restart backoff, the node reports the existing `StoppedWithErr` cluster status with restart attempt info appended, and flips back to `Working` after a successful restart. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Add staging TestTransientError op and consensus restart integration test New staging-only cluster operation `test_transient_error` (mirroring `test_slow_down`): applying it fails with the given probability on the targeted peer, stopping its consensus thread with a service error. The peer re-applies the entry on every consensus thread restart, rolling the probability again, so it exercises the consensus restart loop end to end. Probability 1.0 simulates a permanently failing operation. The integration test runs a 3-peer cluster, poisons a follower through its own API (so the simulated failure is reported deterministically in the response), and asserts that the restart loop engages, the remaining peers keep serving consensus operations with a quorum, and the failed peer eventually recovers and catches up. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Fix import order to satisfy rustfmt Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
8c8a72d120 |
Remove read_multi_iter to fix macOS linker symbol overflow (#9643)
* remove unused iter_offsets
* Replace MultivectorOffsetsStorage::iter_offsets with callback-based for_each_offset
First step of removing the iterator-returning read API (whose deep,
composable generic types blow up mangled symbol size). Convert the
offsets read from an iterator to a callback the caller pushes into:
- trait method iter_offsets -> for_each_offset(ids, FnMut(usize, MultivectorOffset))
returning common::universal_io::Result<()>
- Mmap impl now uses the callback read_batch (drops one read_iter use)
- Ram / Chunked impls push into the callback; Chunked still goes through
iter_vectors for now (converted in a later step)
- the single caller (for_each_in_multi_batch) passes a closure
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Implement read_batch directly on ReadPipeline, not via read_iter
read_batch now drives the pipeline itself (refill-then-wait loop, like
read_multi_iter) and invokes the callback per result, instead of
consuming the iterator returned by read_iter. A step toward removing the
iterator-returning read API.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Remove the ReadMulti RPC from the StorageRead gRPC service
ReadMulti was the only real consumer of UniversalRead::read_multi (which
itself relies on read_multi_iter). Removing the RPC end-to-end clears the
path to dropping that read API. StorageReadService keeps all its other
RPCs (ListFiles, FileExists, FileLength, ReadBytes, ReadBytesStream,
ReadWhole, ReadBatch).
- proto: drop `rpc ReadMulti` + ReadMulti{Entry,Request,Response}
- regenerated lib/api + uio-client generated code; drop ReadMulti
validation rules in lib/api/build.rs
- tonic: delete the read_multi handler + its 2 tests
- uio-client: delete Client::read_multi, the mock-server impl, and 2 tests
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Remove UniversalRead::read_multi
Its only real consumer was the StorageRead ReadMulti gRPC handler (removed
in the previous commit); the two wrapper forwarders had no callers. Drop
the trait method and both forwarders (typed/read_only), and remove the
io_uring test that only existed to compare read_multi vs read_multi_iter
(read_multi_iter stays covered by the other tests). Another step toward
removing the iterator-returning read API.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Drive ReadPipeline directly in gridstore read_from_pages
Replace the read_multi_iter call in Pages::read_from_pages with a direct
pipeline loop (refill-then-drain), scheduling each multi-page read on its
own page file. Behavior unchanged; another step toward removing the
iterator-returning read_multi_iter.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Drive ReadPipeline directly in gridstore read_batch_from_pages
Replace the second read_multi_iter call (in Pages::read_batch_from_pages)
with a direct pipeline loop, scheduling each (ReadMeta, page, range) on its
own page file and propagating errors via GridstoreError. Single/multi-page
buffering and out-of-order reassembly are unchanged. No more read_multi_iter
in gridstore.
Measured overhead on warm mmap (both paths are zero-copy borrows): ~0.3 ns
per read of fixed control cost, flat across read sizes — well under 0.1% of
a real payload read.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Drive ReadPipeline directly in on-disk postings with_posting_views
Replace read_iter in OnDiskPostings::with_posting_views with a direct
pipeline loop. wait_bytemuck yields a file-borrowed Cow, so postings are
still stored zero-copy in raw_postings (a read_batch swap would have forced
an owned copy of every posting list per query on the mmap backend).
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Read on-disk posting headers via read_batch, drop the HeadersBatch iterator
headers_iter now reads headers with the callback read_batch API: each header
is parsed (copied) out of the read bytes, so nothing borrows the file past the
read — no pipeline needed. Since the read is now eager, HeadersBatch holds the
collected Vec<HeaderResult> directly instead of a Box<dyn Iterator>, dropping
the boxing, the dynamic dispatch, and the struct's lifetime parameter.
with_posting_views takes the Vec and still pipelines the posting reads.
Removes the last read_iter use in this file.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Use read_batch in simple_disk_cache populate_from
populate_from reads one byte per block only to fault blocks into the local
cache, discarding the bytes — a no-op-callback read_batch fits exactly. Drives
the same DiskCachePipeline as before; one less read_iter caller.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Implement read_iter directly on ReadPipeline, not via read_multi_iter
read_iter now drives the pipeline itself (refill-then-wait loop, mirroring
read_bytes_iter) instead of mapping its ranges onto self and calling
read_multi_iter. Same signature and iterator contract, so all callers are
unchanged. Leaves iter_vectors as read_multi_iter's only remaining caller.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Revert "Drive ReadPipeline directly in on-disk postings with_posting_views"
This reverts commit
|
||
|
|
0fd361571a | Fix abort resharding live-lock (#7849) | ||
|
|
463a305404 |
Add routing token for deterministic read routes (#9338)
* Add routing token structure * Implement routing token in read operation executor as per design doc * Add TODO to glue routing token to user requests * Implement routing header for REST API * Source routing token from request, not from JWT token * Implement routing token in gRPC API * Add test * Review remarks Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> * Use lower case header name to prevent panic * Rename header to X-Qdrant-Route-Affinity * Assert routing consistency in test on all peers --------- Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> |
||
|
|
45fb36323a |
test: fix flaky test_partial_snapshot_empty (#9493)
The test asserts that creating a partial snapshot between two in-sync peers returns 304 (empty diff). It only waited for the write peer to become green, but read the read peer's manifest for the comparison. An async optimization reshaping the read peer's segments after collection-snapshot recovery makes its manifest diverge from the write peer's files, producing 200 instead of 304 (assert 200 == 304). Wait for the read peer to become green as well before comparing manifests, mirroring the earlier flaky-test fixes (#7358, #7360). Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
3053ca6f8a |
test(consensus): de-flake replicate_points_stream_transfer_updates (#9263)
The test relied on concurrent background upserts adding *new* matching points to the destination shard during the streamed transfer, asserting the destination count was strictly greater than the original snapshot. Whether new matching points land in the destination before the transfer completes is timing-dependent (especially under pytest-xdist load), so the strict `>` assertion is racy and occasionally fails with equal counts (e.g. `assert 5031 > 5031`). Relax to `>=` and keep the strict `dest == src` consistency check, which is the actual invariant being verified. Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
c4a22aa124 |
Add timeout to streaming shard snapshot writer (#9239)
* Add block level timeout to snapshot stream writer * Add tooling to exercise streaming snapshot stalls - tests/manual/slow_snapshot_download.py: slowly / partially download a streaming shard snapshot from a URL to exercise sender-side backpressure. Supports hold / rst / fin / blackhole termination to simulate a stalled, killed, or offline (network-partitioned) consumer against a remote node. Stdlib only; read-only against the target. - tests/consensus_tests/test_streaming_snapshot_receiver_kill.py: throttled receiver killed mid-flight + a second receiver, to observe whether the sender releases the SegmentHolder lock and recovers. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> --------- Co-authored-by: timvisee <tim@visee.me> Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> |
||
|
|
031547e820 |
test(consensus): batch initial upsert in snapshot-transfer missing-point test (#9261)
The initial 20k-point insert was sent as a single HTTP request (no batch_size), saturating all cores long enough to starve the consensus thread (cascading leader elections) and to exceed the 2000ms per-shard update healthcheck deadline, returning a flaky 408 before the test's actual snapshot-transfer logic even started. Batch the insert into 1k-point requests, matching the convention used by every other large initial upsert in the suite. Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> |
||
|
|
29b933f66e |
Fix REST auth whitelist, resolve route before authorizing (#9254)
* Don't whitelist endpoints on user provided path, but on endpoint pattern * Add test * Update comment |
||
|
|
5480878db8 |
Batch initial load in resharding abort-crash test to fix flakiness (#9240)
The test fired a single 20k-point `wait=true` upsert into a fresh 3-shard x 2-replica collection. Under CI load that one op keeps a replica busy past the hardcoded 2s inter-node health-check (transport_channel_pool HEALTH_CHECK_TIMEOUT), so the coordinator fails the forward with a transient "Healthcheck timeout 2000ms exceeded" 408 before resharding even starts. It's the only resharding test doing a 20k single-shot upsert (others do ~1k), which is why it flakes and they don't. Batch the load at 1000 points so each op stays well under the health-check window, matching the proven-stable pattern. Data and transfer behavior are unchanged. Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> |
||
|
|
74071c4891 |
fix(tests): wait for abort target to apply transfer-start before aborting (#9238)
test_resharding_down_abort_converges_when_killed_mid_abort fired the fire-and-forget abort_transfer request at alive_uri after waiting only for the *victim* to apply the transfer-start. When alive_uri lagged in replicating that consensus entry, the abort handler's local check_transfer_exists returned 404 (swallowed by the fire-and-forget thread); the transfer then completed naturally, resharding was never aborted, and _victim_in_window never opened, timing out after 30s. Wait until both alive_uri (the abort target that gates the 404) and the victim have applied the transfer-start before firing the abort. Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> |
||
|
|
a06c2109e6 |
Fix abort transfer resharding idempotency (#9215)
* Abort resharding before we abort transfer * Test resharding-down abort converges when a peer is killed mid-abort --------- Co-authored-by: tellet-q <elena.dubrovina@qdrant.com> |
||
|
|
df0a09d692 |
test: surface peer startup crashes in dirty-shard test (#9161)
* test: wait for WAL flock release on peer restart in dirty-shard test The flake addressed by #9124 (bumping `wait_for_peer_online` to 60s) was misdiagnosed as CPU contention. Logs show the restarted peer panics within ~1s of startup at `consensus_wal.rs:36` with: Wal error: Can't init WAL: Kind(WouldBlock) Panic: Can't open consensus WAL: Kind(WouldBlock) `wal::Wal::open` calls `fs4::FileExt::try_lock` (non-blocking `flock`) on the WAL directory fd. After `p.kill()` (SIGKILL + waitpid) the kernel normally releases the killed peer's flock immediately, but under pytest-xdist load there is a small window where it lags. The fresh peer's startup then races and panics. After the panic `/readyz` never returns 200, so neither 30s nor 60s rescues the test. Fix: add a `wait_for_wal_unlocked` helper that polls both the consensus WAL directory and the local-shard WAL directory with the same exclusive non-blocking flock that qdrant uses, and call it after every `p.kill()` that is followed by a `start_peer` on the same `peer_dir`. The two 60s timeouts are restored to the default 30s now that the underlying race is gone. Co-authored-by: Cursor <cursoragent@cursor.com> * test: fail fast if sync restart crashes in dirty-shard test The sync restart in `test_dirty_shard_survives_update_collection` used plain `wait_for_peer_online(sync_uri)`, which only polls `/readyz`. If the freshly started peer panics on startup (e.g. WAL `WouldBlock`), the test waits the full 30s timeout and then reports a `/readyz` timeout instead of the actual panic message and exit code. Switch the sync restart to `wait_for_peer_online_or_crash(...)` (same helper already used for the dirty restart). On crash it dumps the peer log tail so the next CI failure shows the real reason instead of a generic timeout. Co-authored-by: Cursor <cursoragent@cursor.com> * test: drop speculative WAL flock wait, keep crash-detection switch Reverts the `wait_for_wal_unlocked` helper that polled the WAL directories with non-blocking flock. The kernel-side flock race it was guarding against could not be reproduced in isolation (0/300 iters of SIGKILL+wait+re-flock on bare Linux), so it was speculative. Kept: - Sync restart now uses `wait_for_peer_online_or_crash(...)` instead of plain `wait_for_peer_online`, so any startup panic surfaces fast with the actual log tail instead of hiding behind a 30s `/readyz` timeout. - The two 60s timeouts bumped in #9124 are restored to the default 30s. If the flake recurs in CI, the new failure output will tell us the actual cause (panic message + exit code), which is more useful than papering over it. Co-authored-by: Cursor <cursoragent@cursor.com> --------- Co-authored-by: Cursor Agent <agent@cursor.com> Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
8e277d2bd3 |
tests: use uniform from .utils import * in consensus tests (#9155)
* tests: use uniform `from .utils import *` in consensus tests Eight consensus tests imported specific symbols from `utils` instead of using `from .utils import *`. The most consequential side effect was missing the `every_test` autouse fixture, which is responsible for cleaning up leaked `processes` and resetting the port-slice allocator between tests. Without it, leaks from one test can carry into the next on the same xdist worker, causing `processes.pop(target_idx)` to return the wrong peer in tests like `test_dirty_shard_crash_loop` and triggering WAL-lock conflicts when the intended target is left running. Make the imports uniform across the package so the autouse fixture is always in scope. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * tests: keep `test_issues_api` on explicit imports This module uses a `@pytest.fixture(scope="module")` cluster setup shared across all tests in the file. With `from .utils import *` the function-scoped `every_test` autouse fixture comes into scope and runs after the module-scoped `setup`, so on the first test it sees the already-populated `processes` and calls `kill_all_processes()` — killing the shared cluster before the test body runs. The test then fails with connection refused on the cluster's port. Keep explicit imports here so `every_test` is not autoloaded; the module teardown already cleans up via `kill_all_processes()`. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com> |