mirror of
https://github.com/qdrant/qdrant.git
synced 2026-09-28 17:07:49 -05:00
read_bytes_async_uring
2071
Commits
| Author | SHA1 | Message | Date | |
|---|---|---|---|---|
|
|
79c71a644b |
docs(collection): align SnapshotRecover.priority rustdoc with SnapshotPriority (#10233)
Replica prefers existing replica data, then this peer is synchronized from other replicas — not from the snapshot. Also fix the "if will" typo. |
||
|
|
dfca67f5ed |
feat(consensus): warn when applying a single entry stalls the consensus thread (#10215)
* feat(consensus): warn when applying a single entry stalls the consensus thread longer than a threshold |
||
|
|
e32d3fbf89 |
Add acosh expression to formula query (#10231)
Unary inverse hyperbolic cosine, parallel to sqrt/ln/exp/log10, in REST, gRPC, and edge (FFI + Python) interfaces. Inputs below 1 produce the same NonFiniteNumber error as an invalid sqrt or ln. Closes #10186 Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
86b9330628 |
transfer: send raw payloads, behind feature flags (#10066)
A raw point can carry its payload as the byte blob it is stored as, mirroring `PointStructRaw.raw_payload` on the internal gRPC API. The blob travels from the sending node into the receiving node's WAL untouched, so the sender never parses the payload it read and neither node builds a protobuf value tree for it. It is parsed exactly once, where the operation is unpacked for apply (`process_point_operation`), because that is the first place the parsed form is actually needed: `set_full_payload` goes through the payload index, which cannot be updated from bytes. The gRPC boundary therefore only checks the encoding tag and rejects a point that sets both payload fields, the way the enclosing request already rejects both `points` and `raw_points`. Moving the parse onto the apply path makes its error classification load-bearing, so a malformed blob is reported as `OperationError::MalformedPayloadBlob` — the payload sibling of `MalformedVectorBlob`, mapped to `CollectionError::BadInput` for the same reason: a bad blob that reached the WAL has to be skipped on replay instead of crash-looping recovery. Three consequences of the blob living that long are handled explicitly rather than by convention: - `decode_payload_raw` takes the blob only once it has parsed, so a failure leaves the point holding it instead of holding neither representation. - `upsert_points_raw` and `sync_points_raw` refuse a point that still carries a blob. They read the parsed payload, so such a point would otherwise be stored with no payload at all, and a `debug_assert!` would not catch it in release. - `is_equal_to` compares blob to stored blob as bytes. A differing encoding costs a redundant upsert on sync, never a skipped one. The `raw_payload_transfer` bench measures the trade, per 100-point batch (one transfer batch) at payloads of ~200 B / ~700 B / ~7 KB: - Sender, storage bytes to wire: 16x / 37x / 113x faster. This is where the whole win is — no parse of the blob that was read, no value tree built. - WAL encode: 5x / 11x / 25x faster, writing a byte string instead of a map. - Receiver, wire to applicable point: 1.09x / 1.10x / 1.06x. Near neutral, as it swaps walking a prost value tree for a JSON parse. - Wire bytes: ~6% smaller. WAL bytes: 10-32% *larger*, because the blob is JSON while a parsed payload is written as a compact CBOR map. The WAL growth is accepted rather than fixed: decoding earlier to win those bytes back costs a second full deserialization, and would leave the receiving side with a `payload_raw` that is never populated. Making the blob itself compact belongs in the payload storage encoding (`RawPayloadEncoding` is the extension point for it), not here. Two flags, both off by default and both sender-only (nodes accept raw points and raw payloads regardless), read where the transfer batch is prepared: - `transfer_raw_points` transfers every collection as raw points, not only those whose vector storage would drift in a decode-encode round-trip. - `transfer_raw_payloads` ships the blob a raw read hands out; without it the prepared batch decodes it back into the parsed payload, and the wire message is exactly what it is today. Neither is enabled by `all`: a node only accepts them once it runs a version that understands them, so they can only be switched on a release later. Nothing enforces that yet — the transfer has no peer-version gate. Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
e8c8eead24 |
Fix getting stuck on reshard down abort (#10205)
* test: parameterize collection fixture and share resharding consensus stub Let integration tests build a collection with a custom optimizers config, and move NoopReshardingConsensus from the consensus idempotency test into the shared test module. * test: reproduce scale-down resharding abort hanging on deferred points abort_resharding holds the shard holder write lock while scale_down_cleanup_points deletes migrated points with WaitUntil::Visible and no timeout. With prevent_unoptimized enabled that wait only resolves once the optimizer has cleared every deferred point of the shard, so a stalled optimizer wedges the write lock, and with it every shard holder reader and - in a cluster - the consensus apply thread driving the abort (SetShardReplicaState(Dead) on a ReshardingScaleDown replica, e.g. after that peer is killed). The test drives such an abort with the optimizer disabled and asserts it completes. It currently fails by hanging into its 30s timeout, and must pass once the cleanup delete no longer waits for visibility. * When cleaning up old points, don't wait until visible |
||
|
|
7f3912fb71 |
fix(collection): size WAL replay clamp by max_capacity, not available capacity (#9969)
load_from_wal sized the WAL-replay handoff from update_sender.capacity() (currently-available slots, which vary at runtime) but treated it as the total queue size, contradicting update_queue_length() in the same file which uses max_capacity() as the total. Use max_capacity(). This also removes the 'update_queue_size as u64 - 1' underflow risk, since max_capacity() is >= 1 while capacity() can be 0. |
||
|
|
6a4a1fc4d4 | test(model_testing): log run stages and fix nightly failure reporting (#10163) | ||
|
|
0a6cb3b4cf |
[Raw payloads]: read payload as stored bytes in retrieve_raw (#10040)
* segment: read payload as stored in retrieve_raw `retrieve_raw` already hands back vectors as stored; let the caller ask for the payload the same way, so a reader that only relocates a point parses nothing. `RawPayloadFormat` states what the caller wants — no payload, parsed, or as stored — and replaces the `WithPayload` argument, which could express a key selection that a raw read cannot serve anyway. [`MaybeRawPayload`] states what came back, which can differ from the request in one direction only: a payload storage that keeps payloads parsed cannot answer `Raw` with a blob, and now says so instead of encoding a payload for a reader that would parse it straight back. The raw path reaches the blobstore through `read_payloads_maybe_raw`, mirroring `read_payloads` down the payload storage and payload index traits, so it keeps the batched read. Every caller asks for `Parsed`, so this changes no behaviour: the copy-on-write move and the sync comparison need the parsed payload anyway, and the shard transfer switches over with the feature flag that ships the blob to another node. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * segment: always hand out the stored payload blob from retrieve_raw Review follow-up: instead of telling `retrieve_raw` in which form to return the payload, it always returns it as stored and a caller that needs the parsed form decodes it itself. - Drop `RawPayloadFormat` and the payload parameter it replaced: no production caller ever asked for anything but the whole payload, and a selector cannot be applied to an opaque blob anyway. - Drop `MaybeRawPayload` / `MaybeRawPayloadRef`: only `InMemoryPayloadStorage` could produce the parsed variant, and no segment can be built with that storage (`PayloadStorageType` is `Mmap` or `InRamMmap`, both blobstore-backed). `SegmentRecordRaw` carries a plain `Option<RawPayload>`. - `PayloadStorageRead::read_payloads_maybe_raw` becomes `read_payloads_raw` and hands out `Option<&[u8]>`. The in-memory storage keeps payloads parsed, so it encodes on read, producing the bytes an on-disk storage would have written. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * api: decode a received raw payload with the shared decoder `decode_payload` at the gRPC boundary matched on the encoding and parsed the blob itself, duplicating `RawPayload::decode`. Add the inbound conversion from the wire type and let the one decoder do the reading, so another encoding has a single place to be taught. The conversion also rejects an encoding number no variant maps to, which prost would otherwise hand out as the default encoding — a blob from a node that writes payloads some other way must not be read as JSON. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * simplification --------- Co-authored-by: Ivan Pleshkov <ivan.pleshkov@qdrant.com> Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
7364cc42ef |
feat(edge): add query_batch for batched planned queries (#10100)
* feat(edge): add query_batch for batched planned queries Expose the planned-query batch path as a public API so multiple independent queries can share one planning pass over leaf searches and scrolls. Wired through EdgeShardRead, FFI, and Python bindings. Co-authored-by: Cursor <cursoragent@cursor.com> * perf(edge): push batched query vectors down to segments `query_batch` planned the whole batch at once but then executed every leaf search on its own: one query context, one fan-out over all segments, and one single-vector `Segment::search_batch` call per leaf. Execute the batch as a batch instead: - `EdgeReadView::search_batch` builds the query context once, visits the segments once, and hands each segment the leaves that agree on everything but their query vector as a single multi-vector `search_batch` call. `search` is now a thin wrapper over a one-element batch. - Move `SearchType`/`BatchSearchParams` from `collection`'s segments searcher into `shard`, next to `CoreSearchRequest`, and add `group_search_batches` so both the collection and the edge read path share one grouping implementation. Edge computes the grouping once and reuses it per segment. - `search_matrix` now issues its per-sample nearest queries through `query_batch`; they share filter, limit and vector name, so the whole sample is scored in one batched search per segment instead of one full segment pass per sampled point. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Cursor <cursoragent@cursor.com> Co-authored-by: generall <andrey@vasnetsov.com> Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
c65d7d2d63 | test(resharding): test resharding state clears before the replica revert (#9656) | ||
|
|
2c40d5f2f1 |
Fix Transfer::Restart (#9786)
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> |
||
|
|
347fbcaf5b |
Fix SetRepicaState(Dead)/Transfer::Abort/Resharding::Abort` (#9760)
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> |
||
|
|
a21e0955ed |
Fix S3 snapshot path traversal vulnerability (#10085)
* Fix snapshot path traversal vulnerability * Add tests |
||
|
|
192ed73385 |
Make CreateShardKey idempotent (#10025)
[audit-K] Make create_shard_key crash-safe with a single commit point create_shard_key persisted shard_key_mapping.json incrementally, one add_shard per placement entry, so the operation was not atomic across a crash. The re-apply gate state.shards_key_mapping.contains_key(&shard_key) becomes true after the first loop iteration, so a crash mid-loop on a multi-shard placement left the peer permanently holding a subset of the key's shards while every other peer had all of them. Write the mapping exactly once instead, after every shard of the key exists on disk, so there is no partial state to observe. SaveOnDisk writes atomically, load_shards derives the shard id list from the mapping alone (unreferenced directories are invisible after restart), create_shard_dir wipes leftovers, and max_shard_id reads only the mapping - so a replay allocates exactly the ids the crashed attempt did, which are the ids every other peer allocated too. The contains_key gate then holds as intended: it fires only for an operation that already completed in full, or for a genuine duplicate. ShardHolder::add_shards registers a batch and persists the mapping once, with add_shard as a one-element wrapper; the mapping write moved ahead of the in-memory updates, being the only fallible step. No caller changes behavior: Collection::new and load_shards already hit the write_optional early return, and start_resharding_unchecked still adds a single id. An empty placement is now rejected rather than silently returning Ok without creating the key. Covered by create_shard_key_test.rs: id allocation, replay before and after the commit, duplicate rejection, and the empty placement guard. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
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> |
||
|
|
7743616b3b |
Report effective (cgroup) CPU, RAM and disk in telemetry (#9891)
* Report effective (cgroup) CPU, RAM and disk in telemetry
The `system` block reported host-level figures that ignore the limits the
kernel actually enforces on the process:
- `cores` <- sys_info::cpu_num() (host socket count)
- `ram_size` <- sys_info::mem_info() (host total RAM)
- `disk_size` <- sys_info::disk_info() (container root fs)
On any cgroup-limited deployment (containers, Kubernetes pods, systemd
slices) these overstate what Qdrant can use, are misleading for capacity /
oversubscription analysis, and don't match how Qdrant sizes itself.
Report the effective values instead, reusing existing helpers:
- `cores` -> common::cpu::get_num_cpus() (already drives sizing)
- `ram_size` -> segment::utils::mem::total_memory_bytes()
(cgroup limit via cgroups_rs, else sysinfo host total)
- `disk_size` -> common::disk_usage::disk_usage(storage_path)
(data-volume capacity, cached; sys_info host disk fallback)
`Mem::new()` is not free (it builds a sysinfo System and loads the cgroup
memory controller), so `total_memory_bytes()` caches with a 5s TTL — matching
the disk-usage cache — instead of recomputing per call. The TTL (rather than
caching once) means an in-place cgroup memory resize is reflected within a few
seconds, consistent with how `disk_size` and `cores` already behave. The
strict-mode helper in `collection` now delegates to this shared accessor
(dropping its own OnceLock), so it too becomes resize-aware.
ram_size/disk_size stay in KiB to match the previous unit.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Simplify comments and code
* Simplify openapi spec
* minor comment improve
* mem: cache total_memory_bytes like disk_usage (5s TTL)
`total_memory_bytes()` mirrors `common::disk_usage::disk_usage`: a small 5s
TTL cache over `Mem::new().total_memory_bytes()`. `Mem::new()` is not free
(builds a sysinfo System + loads the cgroup controller), and the short TTL
keeps the value in step with an in-place cgroup memory resize rather than
freezing at startup. Shared by telemetry and strict-mode, like the disk cache.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Regenerate OpenAPI spec
Add the ram_size / disk_size field descriptions produced by schemars from the
updated telemetry doc comments, keeping the generated spec consistent.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* mem: address review — parking_lot mutex, hold lock, saturating_duration_since
- Use parking_lot::Mutex (no poisoning; lock() returns the guard directly).
- Hold the lock across the whole method — single acquisition, simpler.
- saturating_duration_since instead of duration_since (no panic on a cached
timestamp spuriously ahead of now).
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
|
||
|
|
f44ae2939d |
Serialize snapshot recoveries of the same shard (#10026)
Two snapshot recoveries of one shard resolve as "last to finish wins".
`restore_local_replica_from` is destructive on its own - for a full recovery
it takes the `LocalShard::clear` + `move_data` branch - and nothing is atomic
across the download that precedes it.
That picks the wrong winner. A recovery gets abandoned because it stalled,
which is exactly why its caller timed out and retried, so it finishes *after*
the retry that replaced it and restores on top of it. It keeps running because
`recover_shard_snapshot_impl` is not cancel safe, so a caller that walked away
cannot stop it:
stalled: clear -> download ..........................-> restore -> rolls back retry
retry: clear -> download -> restore -> Ok to caller -> caller writes
The retry's `Ok` is what the sender of a snapshot shard transfer acts on: it
switches the transfer to `Partial` and flushes its queue proxy into the shard.
The stalled recovery then discards those writes - acknowledged, then lost.
Add a per-shard recovery lock on `ShardReplicaSet`, held across the whole
recovery (clear, download and restore), so what a caller asked for last is what
survives. Clearing before the download is kept, so recovery still never needs
disk space for two copies of the shard.
`ShardRecoveryGuard` already existed to track recovery progress for the whole
recovery, so the lock is folded into it rather than added alongside:
`ActiveRecoveries::start` takes ownership of the lock guard, making it
impossible to start a recovery without holding it, and `Collection::
start_shard_recovery` is the single call that acquires both.
Queueing behind another recovery is logged. An abandoned recovery is otherwise
invisible - no live caller, no request, no error - and shows up only as a
transfer sitting in `Recovering`.
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
||
|
|
75385df69f |
Remove dead code (#10030)
* Remove dead code * Remove unused dependencies * `allow(dead_code)` -> `expect(dead_code)` * ast-grep: rule-tests/*-test.yml => tests/*-test.yml For brevity. * ast-grep: forbid allow(dead_code) |
||
|
|
c124db6905 |
Make deprecated on_disk_payload optional in the REST schema (#10020)
* Make deprecated `on_disk_payload` optional in the API schema `CollectionParams::on_disk_payload` was the last place the deprecated flag was still a bare bool, in both the REST schema and the gRPC `CollectionParams` message. Clients generated from those schemas model it as a required bool, so removing the field in a future version would break them. Make it optional in both schemas while keeping it populated on every path that builds `CollectionParams`, so responses and the persisted collection config still carry a value and existing clients keep working until it is removed. The gRPC change is wire-compatible: proto3 `optional` only adds a synthetic oneof for presence tracking, the field number and wire type are unchanged. It does change what an absent field means, though, so decode absent as `false` when reading a remote peer's collection info -- peers that predate the change encode a plain bool, which is omitted from the wire when false. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Keep gRPC `on_disk_payload` a plain bool Adding proto3 `optional` would let a new client tell "unset" from `false`, but it also changes what an absent field means. A pre-upgrade server encodes a plain bool, which is omitted from the wire when false, so a client generated from the new proto would decode `false` as "unset" -- and upgrading the client first is the order we recommend. gRPC does not need the change anyway: protobuf tolerates a missing field by design, so a gRPC client does not break at runtime when we stop sending this one. The breakage this addresses is on the REST side, where a generated model with a required non-nullable bool fails to deserialize a response that omits the field. JSON encodes `false` explicitly, so the REST schema change carries no equivalent ambiguity. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
be543561e5 |
Add Logstore and Blobstore wrapper (#9673)
* Gridstore: introduce storage operating mode in config Add a mode field to the gridstore config, selecting between the dynamic mode (current behavior, the default) and the upcoming serverless mode. The mode is specified through StorageOptions on creation, persisted in config.json, and read back first when opening so the correct variant can be selected automatically. Configs written before this field existed deserialize as dynamic. For now, selecting the serverless mode returns an error; the variant itself is added in follow-up commits. * Gridstore: move dynamic implementation into dedicated module Mechanical move of the current Gridstore implementation into gridstore/dynamic.rs as DynamicGridstore. The public Gridstore struct becomes a thin wrapper holding a mode variant enum, propagating every call into the selected variant. For now the enum only has the dynamic variant; the serverless variant is added in follow-up commits. No logic changes to the dynamic implementation itself: only visibility, the config parameter now passed into open (the wrapper reads it first to select the mode), and open_or_create staying on the wrapper. * Gridstore: add serverless tracker Add the append-only mapping tracker for the serverless storage mode. The tracker file is a plain array of 16-byte mapping entries without any header: the number of mappings is defined by the exact file length, and the entry index is the point offset. The file starts empty and only ever grows by appending, existing bytes are never rewritten. Mappings must be set in monotonically increasing point offset order; skipped offsets are backfilled as zeroed entries which decode as None. New mappings are buffered in memory and appended with a single write per flush. A flush with a stale target is a no-op so bytes are never written twice. A torn trailing entry (file length not a multiple of the entry size) is ignored when reading and truncated away when opening writable. Unlike the dynamic tracker, the file is read and written directly with positional file IO instead of memory mapping, as serverless environments do not handle memory mapped files well. * Gridstore: add serverless storage variant Add the append-only gridstore variant for serverless deployments, which restrict IO to appending to files: existing bytes can never be rewritten, and IO is expensive so as few files as possible are used. The variant stores all value data in a single page file next to the serverless tracker and the storage config, three files in total. Both data files start empty and only ever grow by appending; there is no preallocation, no used-block bitmask and no gap/region bookkeeping. Values are appended at put time at the next block aligned offset, with the zero padding included in the write so it lands exactly at the end of the file. Mappings are buffered and appended to the tracker with a single write per flush, after the page file is synced, so a mapping on disk never points at data that is not durable. Values cannot be updated or deleted, and must be put at monotonically increasing point offsets; violations are rejected before any data is written. Files are read and written directly, never memory mapped. The mode is selected through StorageOptions on creation and picked up automatically from the persisted config when opening. * Gridstore: serverless support in reader and view Extend the read-only GridstoreReader and the GridstoreView with the serverless mode, keeping both public types unchanged: like the writable Gridstore they now hold a mode variant internally, selected automatically from the persisted config when opening. The serverless reader holds the tracker and page directly and reads the files positionally, without memory mapping. A live reload re-reads the mapping count from the exact tracker file length (there is no size header), ignoring a torn trailing entry, and never truncates as it is read-only. Value reads always go directly to the file, so newly appended data is readable without remapping anything. * Gridstore: document storage operating modes * Gridstore: review fixes for the serverless mode Hardening and cleanup from a review pass over the new serverless storage variant: - Batch the reader side iteration like the writer already did, instead of materializing tracker mappings for the full range in one go, which could transiently allocate gigabytes on large storages. - Recover the append cursors when a positional write fails partway: truncate the file back to the tracked length so a retried append or flush never rewrites bytes that already landed in the file. - Validate page addressability before appending value data, a rejected put must not grow the page file. - Cross-check tracker and page consistency when opening: mappings that reference value data past the end of the page file (e.g. after a partial copy or restore) now fail fast instead of surfacing as opaque read errors per point. - Reject value pointers into any page other than page 0 on the serverless read path with PageNotFound, matching the dynamic mode contract, instead of silently reading from a wrong location. - Refresh the reported storage size on reader live reload even when no new mappings were flushed, unflushed value data may have grown the page file already. - Validate configs read from disk: a corrupt config with zero sized blocks, pages or regions is now rejected when opening instead of panicking on a division by zero later. - Classify rejected serverless puts as UnsupportedOperation, consistent with rejected deletes, so they don't surface as user-facing validation errors at the segment level. - Deduplicate the compression dispatch into Compression::compress and Compression::decompress, and the serverless file create/open patterns into shared direct IO helpers, so the two modes and files can't silently drift apart. * Gridstore: cover both operating modes in mode-agnostic tests Parameterize the gridstore tests that exercise mode-agnostic behavior over both the dynamic and serverless mode with rstest, using a single and bulk put/get roundtrips, storage files, basic persistence, corrupt config rejection, batched read congruence, reader live reload, and the different block sizes. Mode specific expectations branch inside the tests: expected file names, storage size semantics (whole blocks vs exactly packed bytes), value pointer layout (page spill over vs a single packed page), and gaps (created by deletes in dynamic mode, by skipped puts in serverless mode). Dynamic-only internals assertions are kept behind a mode check. Tests around updates, deletes, page spanning, block reuse and other dynamic-only behavior intentionally stay dynamic; the serverless specific format invariants remain covered by the dedicated serverless tests. * Gridstore: port serverless specific tests from sibling branch Source the serverless specific test cases that the serverless-gridstore-updates branch added, adapted to the dedicated variant implemented here (distinct file names, headerless tracker with 16 byte entries, a single packed page without trailing padding, and rejected re-puts): - writes only ever append: tracker and page files only grow and previously written bytes stay byte-for-byte untouched - new mappings land exactly at the end of the tracker file, which always covers the exact number of mappings - mapping gaps are zero-padded on disk and survive reopening - values are packed back to back at block aligned offsets, the page file ends exactly at the last value - serverless mode never creates nor reports block flag files - a flusher persists exactly the mappings that existed at its creation, later puts stay pending - a config claiming the wrong mode fails loudly in both directions instead of loading the incompatible file format of the other mode Tests around their mode switching, page spanning and tolerated deletes don't apply to this design and are intentionally not ported. * Gridstore: test serverless production risk scenarios Add tests for the operational aspects that matter before serverless mode goes to production, each covering a scenario that wasn't evaluated yet: - Replayed puts of already persisted offsets (a WAL redo after a crash where the flush completed but was never acknowledged) are rejected without appending anything, and max_point_offset is the exact offset a replay must resume at. - The accepted crash case of a tracker file extended with zeroed bytes: the entries count as permanent None mappings, can never be put again, and the storage stays consistent and writable past them. - The read-only reader never modifies the files: opening over a torn tracker tail, reading, iterating and live reloading leave both files byte-for-byte untouched. - A multi-round put/flush/reopen cycle always exposes exactly the flushed prefix, with the mapping count matching the exact tracker file length and unflushed offsets reusable. - An append beyond the maximum addressable block offset is rejected before writing anything, keeping retried puts from growing the page file unboundedly. * Gridstore: rename serverless mode to append-only, split into module Rename the mode after its defining characteristic instead of its deployment target: files only ever grow, existing bytes are never rewritten. Renames Mode::Serverless to Mode::AppendOnly (persisted as "mode": "append_only") and the on-disk file names to append_only_tracker.dat and append_only_page_0.dat. The serverless deployment motivation stays in the documentation. Also split the single 2300 line serverless.rs into an append_only module with dedicated files for the storage, page, view, reader and tests. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Use universal IO in Gridstore * Include upstream preopen logic in new Gridstore variant * Gridstore: buffer append-only value writes until flush In append-only mode, put previously wrote the value data to the page file right away, one write operation per put, while mappings were already buffered and batch persisted on flush. Buffer value writes the same way: both the value and its mapping now only land on disk once a flush cycle executes. This batches all new value data into a single write operation per flush, which is significantly more efficient on S3 based storage where every write is a costly operation. A flush now performs exactly two writes: one appending all buffered value data to the page file, one appending all pending mappings to the tracker file, in that order, so a mapping on disk never points at value data that is not durable. The page mirrors the tracker's pending mechanism: an in-memory buffer that is byte for byte the next append (zero padding between block aligned values included), a watermark captured at flusher creation so puts made during a flush stay buffered, a stale-flush no-op guard so appended bytes are never written twice, and truncate-back recovery on failed writes. Reads transparently serve buffered values from memory. As a side effect, a crash between flushes now leaves nothing on disk at all, where the write-through approach left orphaned value bytes in the page file. The buffered data is held in memory until the next flush, bounded by the flush cadence. Universal IO filesystem handles are now required to be Send + Sync, so the flusher closure can carry one to grow the page file at flush time; all existing backends already satisfied this. * Gridstore: rename inner DynamicGridstore to Gridstore The dynamic variant keeps the Gridstore name; the outer dispatching type will be renamed to Blobstore in a follow-up. Until then the inner type is referred to as dynamic::Gridstore to distinguish it from the outer type. * Gridstore: rename append-only variant to Arenastore The append-only variant stores all value data in a single ever-growing page, allocating space by appending, hence: arena store. * Gridstore: rename outer storage type to Blobstore The outer type dispatching between the two storage variants is now called Blobstore, being more generic than Gridstore. This frees up the Gridstore name, which now exclusively refers to the dynamic mode variant, next to Arenastore for the append-only variant. Storage components keep using the outer type, so they now use Blobstore. The gridstore crate name, GridstoreError, and the persisted names (config.json mode, payload config storage_type) are unchanged. * Gridstore: split Gridstore and Arenastore into dedicated modules The outer module is now blobstore, matching the Blobstore type it defines. The two storage variants each get their own submodule: the dynamic Gridstore moves from dynamic.rs into gridstore/ with its reader and view extracted from the shared files, mirroring the arenastore/ module (previously append_only/) which already had this layout. * Rename gridstore crate to blobstore The crate is named after the outer Blobstore storage type it provides. The gridstore name lives on in the dynamic mode variant. GridstoreError and the persisted names (config.json mode, payload config storage_type) are unchanged. * Arenastore: pack values back to back across multiple pages Drop the block alignment from the append-only mode: values are packed byte to byte, without blocks, and the tracker offset is now a plain byte offset within the page. Blocks and regions are dynamic mode concepts; their page size constraints no longer apply to append-only configs. Bring back support for multiple pages. Once appending a value would grow the current page beyond the configured page size, a new page is started, bounding the size of and the number of appends to each file: object stores like S3 Express limit the number of appends per object. A value larger than the page size gets a page of its own; values never span pages. A rollover creates the new, empty page file at put time; the value data itself stays buffered until the next flush, which appends to each touched page with a single write, using per-page watermarks captured at flusher creation. The reader scans for consecutively numbered page files when opening, validates the most recent mappings against them, and adopts pages created since on a live reload. * Blobstore: rename dynamic mode to mutable Rename Mode::Dynamic to Mode::Mutable, and the persisted config value with it: config.json now writes "mode": "mutable". There is no compatibility alias for "dynamic", released versions never wrote the mode field (a missing field still defaults to mutable), only unreleased storages did. The Gridstore type and module names for the mutable variant are unchanged. * Fix Edge compilation due to package rename * Review remarks * Extract Gridstore preopen into module * Rename Arenastore files * Use universal IO for append operations * Rename GridstoreError to BlobstoreError The error type belongs to the Blobstore crate and is shared by both the Gridstore and Arenastore variants, so it follows the crate naming. Also update the user-facing error messages that referred to the old name. * Split config into per-variant types * Rename Arenastore to Logstore Rename the Arenastore type to Logstore, including the reader, view, config, module and variant names. The storage file names follow: log_page_{n}.dat and log_tracker.dat. The persisted mode tag stays "append_only". * Move bitmask module into the Gridstore variant The bitmask tracks free blocks, which only exists in the mutable mode. Move the module from the crate root into the Gridstore variant that owns it. It stays re-exported at the crate root because the bitmask benchmark needs a public path. * Move pages module into the Gridstore variant Like the bitmask, the block based pages module is only used by the mutable mode. Move it from the crate root into the Gridstore variant that owns it. The Logstore variant has its own page implementation. * Use universal IO for every Logstore operation Replace the direct_io module with universal IO in the append-only tracker, making the whole Logstore go through a universal IO backend bounded by UniversalRead and UniversalAppend: - The tracker is generic over the backend now. Reads go through UniversalRead with the caller's access pattern, flushes land as one atomic append with the same offset compare-and-swap recovery as the pages: a retried append after a lost acknowledgement is adopted instead of appended twice. A torn trailing entry is still truncated away on writable open, through a fresh handle since shrinking is not supported through an open one. - The reader now schedules a prefetch for the tracker file too, it no longer bypasses the backend. - The config write, clear and wipe use the backend file operations instead of local filesystem calls, matching the Gridstore variant. * Batch reads in Logstore read_values Apply the same batching logic as the Gridstore variant: resolve all mappings first, then fetch the value data, both through the backend's read pipeline so async backends can serve the reads in parallel. The tracker gains a batched lookup mirroring the mutable tracker's iter, serving pending mappings and out of range point offsets directly from memory. The pages gain a batched value read; unflushed values are served from the in-memory buffers, and since values never span pages each value is a single read without reassembly. Like in the Gridstore variant, the callback may now be invoked in a different order than the requested point offsets. * Better describe logstore live reload ordering * use enum for options, swap `*Options`<->`*Config` naming * don't wrap enum in struct * ditch unused `StorageConfig`, make deserialization more ergonomic * rename `*Options`->`*Config` * make `preopen` non-blocking * fixup! ditch unused `StorageConfig`, make deserialization more ergonomic * fixup! use enum for options, swap `*Options`<->`*Config` naming * fixup! don't wrap enum in struct * fix rebase * use `populate` param in Logstore * test: failing repro of stale page after live reload across rollover A reader that live-reloads between a page rollover and the following flush adopts the new, still empty page. The previous page is then no longer the last one and is never reloaded again, so the tail that the next flush appends to it stays invisible to the reader forever: value pointer at byte 100 with length 100 is out of range AppendOnlyPages::live_reload only reloads the last held page, assuming earlier pages never change once a newer page exists. But the rollover creates the new page file eagerly at put time, while the previous page's buffered tail only lands at the next flush (see test_rollover_writes_no_value_data_before_flush), so a page can keep growing on disk after its successor exists. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * fix: reload all pages that grew * use Fs in `open_or_create` * fix: publish tracker mappings only after the pages reload `AppendOnlyTracker::live_reload` observed the mapping count and made it visible in one step, before `LogstoreReader::live_reload` reloaded the pages. Every failure path in the page reload -- `list_files`, reopening a grown page, opening an adopted one, the truncation check -- therefore left the reader with mappings referencing value data it never loaded, so reads in the new offset range fail until a later reload happens to succeed. The edge refresh loop keeps a segment whose reload failed, expecting it to keep serving its pre-refresh state, which it then does not. Split observing from publishing: `reload_count` refreshes the handle and returns the count as a `PendingReload` token, `commit_reload` publishes it. The reader still observes the tracker first, as the writer persists pages before the mappings referencing them, but only commits once the pages are loaded. Reopening without committing is harmless: reads stay bounded by the unchanged count, and the bytes below it never change. A partial failure inside the page reload needs no unwinding, pages running ahead of the tracker is the safe direction. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * perf: batch the value reads in Logstore iteration `LogstoreView::iter_range`, the path behind `Logstore::iter` and `LogstoreReader::iter`, fetched the mappings for the whole range with a single read but then read the values themselves one at a time, serially. Gridstore routes its `iter` through `read_values` and pipelines both stages, so a full scan of an append-only storage was the one read path without batching -- one blocking round trip per value on the object store backends this variant exists for. It is reached by payload storage iteration and by the payload index build, which scans every payload. Feed the pointers into `read_batch_values` instead, keeping the single contiguous tracker read, which is better than the per-offset pipeline scheduling Gridstore does on that side. Values are now delivered through the read pipeline, so the callback may be invoked out of order, as it already could be for Gridstore's `iter` and for `read_values` in both variants. Both segment callers are order independent. Tests that happened to rely on the mmap backend completing reads in scheduling order now sort before comparing. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * test: don't run the failed-page-reload test on Windows The test shrinks a page file out of band to make the page reload fail, but Windows refuses to resize a file while the reader holds it mapped, which it does by construction here: "the requested operation cannot be performed on a file with a user-mapped section open". The panic is on the injection itself, the code under test never runs. There is no portable injection. Truncating a page the reader holds is what the check under test detects, so the mapping cannot be avoided; failing the adopted page open instead needs a listed but unopenable file, and `local_list_files` descends into matching directories rather than listing them; failing the directory listing needs the storage directory removed, which Windows also refuses while pages are mapped. The storage itself is fine on Windows, its append path grows mapped pages there and every other Logstore test passes. The logic under test is platform independent and stays covered elsewhere, with the tracker half of the guarantee pinned by `test_live_reload`, which runs on every target. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: generall <andrey@vasnetsov.com> Co-authored-by: Claude Fable 5 <noreply@anthropic.com> Co-authored-by: Luis Cossío <luis.cossio@outlook.com> |
||
|
|
eafabd267f |
docs: describe StartResharding fields in OpenAPI (#9946)
Add doc comments to `StartResharding` fields so the generated OpenAPI spec explains what a user has to pass, and regenerate the spec. Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com> |
||
|
|
21284fcef7 |
Honor applied_seq during WAL replay only under prevent_unoptimized (#9930)
`load_from_wal` splits WAL recovery in two: it replays `[first_index, applied_seq + APPLIED_SEQ_SAVE_INTERVAL + 1)` synchronously and hands the remaining tail to the update worker, which applies it in the background *after* `LocalShard::load` has returned and the shard has started serving reads. That split was introduced by #8008 and applies to every collection, so up to `update_queue_size - 1` operations already acknowledged to a client with `wait=true` can be missing from reads right after a restart, reappearing one by one as the worker catches up. Only `prevent_unoptimized` needs that routing: the update worker signals the optimizer per operation, and optimization is the only thing that makes deferred points visible. Everywhere else the synchronous replay is sufficient, so gate the use of `applied_seq` on the flag -- the same condition that already gates the worker's deferred-points wait -- and replay the whole WAL before load returns, as it did before #8008. Found by the crasher: after a crash-restart cycle it reported 72 missing points out of a confirmed 3202, with the shard counting 3202 points while 3930 had been acknowledged and 332 WAL entries were still queued across two shards. Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.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>
|
||
|
|
0745c36c8f | Rename WrongVectorBytesSize -> MalformedVectorBlob (#9929) | ||
|
|
c561e433c6 |
[TQDT] upsert raw malformed blob multi + sparse vectors (#9904)
* Fix malformed vector upsertion for multivecs and sparse-vecs too * Shorten comment * Fix after rebase |
||
|
|
1353d54eb5 | [TQDT] Fix/upsert raw malformed blob badinput (#9886) | ||
|
|
6fc3bcb124 |
test(model_testing): cover slice filter in scroll, count, and delete-by-filter (#9905)
* test(model_testing): add slice matcher to generated scroll filter Extend ScrollFilter with a Slice variant so paginated scroll exercises Condition::Slice. The generator draws small totals (1/2/3/4/5/8) and a valid index; the model verifier mirrors membership via Slice::check — the same hash contract the engine uses — so the existing paged-scroll id-set assertion covers sliced scroll under soak (optimizer, WAL reload, multi-shard, mixed UUID/numeric ids). Co-authored-by: Cursor <cursoragent@cursor.com> * test(model_testing): add CountBySlice verification op Exercise Condition::Slice through the exact count API under soak. Shares the slice generator with ScrollPaged; the model oracle uses Slice::check so engine and in-memory counts must agree. Co-authored-by: Cursor <cursoragent@cursor.com> * test(model_testing): compose slice with num on scroll and delete-by-filter - ScrollFilter::NumAndSlice: indexed num drives candidates; slice is a per-candidate check via Filter::merge. - DeleteByFilter { num, slice: Option<Slice> }: half the deletes also restrict by slice so submit-time filter resolution and WAL-replayed id lists exercise Condition::Slice. Co-authored-by: Cursor <cursoragent@cursor.com> --------- Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
446d140c2d |
Slice filtering condition: sliced scroll / deterministic sampling (#9899)
* feat: slice filtering condition for sliced scroll and deterministic sampling
Add a `slice` filter condition selecting points where
`stable_hash(point_id) % total == index`. The hash is SipHash-2-4 with a
zero key over canonical id bytes (8 LE bytes for numeric ids, 16 RFC 4122
bytes for UUIDs) — a frozen public contract, independent of the internal
resharding ring hash, reproducible by clients to predict membership.
For a fixed `total`, slices are disjoint and cover all points, enabling
parallel scroll streams (ES sliced-scroll style) and reproducible sampling
that composes with any other filter condition.
- REST: `{"slice": {"total": N, "index": R}}`; gRPC: `SliceCondition` in
the condition oneof (tag 8)
- Evaluated per point via id_tracker external-id lookup; no payload index
needed; cardinality estimated as `points / total` with no primary clause
- `total >= 1` enforced by NonZeroU32 at parse time, `index < total` by
validation in both REST and gRPC paths
- Hash contract locked by test vectors independently reproduced with a
reference SipHash-2-4 implementation
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* tests: minimal OpenAPI test for slice filter condition
Scrolls all slices of a fixed total over numeric + UUID ids asserting
disjointness and full coverage, checks must_not inversion, and pins the
two rejection paths (422 for index >= total, 400 for total = 0). Requests
and responses are validated against the regenerated OpenAPI spec by the
test harness.
Note: the spec cannot itself reject total = 0 client-side — the Condition
anyOf falls through to the permissive Filter schema, as with any invalid
condition — so rejection is asserted via the server response.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
|
||
|
|
c196d2eb1a |
Benches: use SmallRng instead of ChaCha12-based generators (#9887)
* Benches: use SmallRng instead of ChaCha12-based generators All benchmarks used StdRng or rand::rng() (ThreadRng), both backed by the ChaCha12 block cipher in rand 0.10. Benchmarks do not need crypto-strength randomness, and several draw random values inside the timed closure, so cipher work was included in the measurement itself. Switch every bench target to SmallRng (Xoshiro256++), and key the HNSW graph cache and sparse index cache by RNG algorithm so stale caches built from the old generator are not reused against newly generated vectors. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Benches: replace free-function rand::random with local SmallRng Addresses review: rand::random draws from the thread RNG (ChaCha12), including inside the timed loop of the pq score benchmark. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
1cd77fa5de |
Model tester: cover all quantization types (#9875)
* Model tester: cover all quantization types Add a `quantization` field to `VectorCandidate` so candidates can carry any `QuantizationConfig` variant, and materialize the configs in the fixture (`quantization_config`). The inline-storage vector "i" keeps its scalar Int8 config, now declared on the candidate instead of hard-coded in the fixture. New candidates: - "p" Dense(8) + Product x4 - "v" Dense(6) + Binary (non-byte-aligned dim, trailing-bit padding) - "r" Dense(8) + Turbo (search-side TQ over Float32 storage) Quantization x datatype combos: - "l" Dense(6) Float16 + Binary - "d" Dense(8) Turbo4 + Turbo default bits (keep-source-rotated branch of `should_keep_source_rotated`) - "g" Dense(8) Turbo4 + Turbo Bits1_5 (Padded rotation, rotate-back branch) This is model-safe: schema quantization keeps the original vectors, so read-back predictions are untouched; the approximate quantized scoring only feeds the membership-only Search/Query/Recommend checks. `assert_candidates_predictable` enforces the wiring constraints: quantization is dense-only (the fixture only wires the dense arm) and requires `initially_active` (CreateVectorName's `DenseVectorConfig` carries no quantization). Verified: seeds 1/2/3/7/42 soaks (5k ops, restarts, optimizer on) green; quantized codes confirmed on disk for every quantized candidate. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Model tester: make inline_storage a VectorCandidate knob, enable on "l" and "d" Replaces the fixture's name-based INLINE_STORAGE_VECTOR special case with an inline_storage field on VectorCandidate (requires quantization, enforced by the startup assert). Enables it on "l" (Float16 base + padded Binary links) and "d" (Turbo4 base + TQ links) to cover more (base layout, link encoding) pairs of the CompressedWithVectors format; "v", "r", "g" and "p" keep the non-inline paths. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- 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> |
||
|
|
eadf071a9c |
Clamp WAL replay target to the truncated prefix on shard load (#9851)
The WAL is only truncated past operations whose segment flush was confirmed, so first_index is a durable lower bound on the applied sequence. The persisted applied_seq can legitimately lag behind it by more than one save interval: it is saved every 64 update-worker calls from a counter that restarts at zero on process start, and synchronous WAL replay never feeds it. A replay target computed from such a stale applied_seq can then sit before first_index, tripping the debug_assert from #8454 (flaky model_testing gate, #9844) and, in release builds, enqueueing already-truncated indices that fail with spurious "Operation not found in WAL" errors. Clamp the replay target at first_index: nothing before it ever needs replay. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
e7f03d40c4 |
Model tester: support UUID point ids (#9866)
Every point id the workload draws now comes from an IdSpace pool precomputed at startup. A --uuid-id-fraction (default 0.5) of the --id-pool slots are well-formed v4 UUIDs built from the seeded rng via uuid::Builder::from_random_bytes, the rest stay numeric. Precomputing the pool keeps the id-reuse semantics (upserts overwrite live points, deletes and retrieves hit them) that fresh per-op random UUIDs would lose, and keeps runs seed-reproducible. Sampling consumes a single range draw per id, exactly like the previous NumId draw, so fraction 0 consumes no extra rng draws and reproduces the numeric-only op stream byte-for-byte. The harness smoke tests run with fraction 0.5, and the fraction is recorded in the trace header. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
561aa419f7 |
Add recoverable OutOfAppendableCapacity operation error (#9860)
Preparation for capping appendable segment growth in the update path (#9158): a dedicated error for "all appendable segments reached max_segment_size", so the update pipeline can recognize it and provision a fresh appendable segment before re-applying the operation. Maps to a transient service error at the collection level: if it ever escapes recovery, failed-operation recovery re-applies the operation. Part 1/5 of the appendable segment overflow fix. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
db2a135203 | raw vector grpc send (#9843) | ||
|
|
fd5ef26d59 |
Model tester: support Cosine, Euclid and Manhattan distance metrics (#9853)
Every vector candidate now carries a distance metric instead of the hardcoded Dot, threaded into the fixture schema, the CreateVectorName generator and the read-back prediction. The model predicts Cosine read-backs exactly by mirroring the engine's ingestion preprocessing: metric_preprocess follows NamedVectors::preprocess_dense_vector's per-datatype dispatch and calls Distance::preprocess_vector itself, so predictions track the engine by construction (including the identity preprocess of the byte metric, which stores Uint8 vectors un-normalized). Stored vectors are preprocessed exactly once (optimizer and CoW moves transfer raw bytes), so predictions stay exact across moves, including Cosine + Float16. New candidates: "e" (dense Cosine), "n" (multi-dense Cosine, per-row normalization), "x" (dense Cosine + Float16), "o" (dense Cosine + Turbo4, padding-free dim), "j" (dense Euclid), "k" (dense Manhattan). Euclid/Manhattan preprocess is an identity, so their value is engine side: Order::SmallBetter comparator coverage. The startup predictability check now also rejects sparse + non-Dot (sparse schemas carry no distance) and Turbo4 + Euclid/Manhattan (TQ's L1/L2 modes store lengths differently from Dot/Cosine and their copy-on-write re-quantization fixed point is not soak-validated yet). Soak-validated on seeds 1/2/4/5/6/7/8 (30k ops), including two restart runs (restart probability 0.002) with the optimizer enabled. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
f8a62ed423 | Fix flaky test_join_all_completes_sibling_restart_after_workers_stop (#9848) | ||
|
|
9a20fc49e4 |
[TQDT] TQ roundtrip: raw vectors in grpc, WAL and apply internal operation (#9813)
* raw vector grpc apply are you happy fmt clean up are you happy clippy review remarks review remarks * fix after rebase |
||
|
|
842ddfae10 |
fix: validate lookup_from collection for query and recommend APIs (#9531)
* fix: validate lookup_from collection in query APIs * Move validation to the bottom of the struct implementation * fix: preserve lookup_from missing collection error --------- Co-authored-by: timvisee <tim@visee.me> |
||
|
|
0d2f5c7e85 | Miscellaneous cleanups (#9849) | ||
|
|
355ac9fc2e |
Add Float16 and Uint8 storage datatypes to the model tester (#9815)
* Add Float16 and Uint8 storage datatypes to the model tester
Extend VectorCandidate with a datatype override and fold the DenseTurbo
kind into Dense + Some(Turbo4) so storage datatype has a single source
of truth. Two new initially-active candidates exercise half-precision
("h", dense 6) and unsigned-byte ("y", dense 4) storage; "c" carries an
explicit Some(Float32) to cover schema configs that spell the default
datatype out.
The model predicts lossy read-backs through the engine's own
PrimitiveVectorElement impls (as Turbo4 reuses turbo_storage_roundtrip)
and compares them exactly: both round-trips are deterministic and
idempotent, so they stay bit-stable across optimizer moves, WAL replay,
and reloads. Uint8 components are drawn from 0.0..256.0 since the
storage truncates with `x as u8` and unit-range draws would collapse to
zeros.
A compile-time assertion rejects datatype overrides on non-Dense
candidates: the fixture's sparse/multi-dense arms ignore the field and
multi-dense read-backs are compared without a round-trip prediction, so
a lossy multi-dense candidate would soak-panic with a false divergence.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* Plumb Float16 and Uint8 multi-dense support in the model tester
Multi-dense storage converts the flattened matrix component-wise
(from_float_multivector), so per-row round-trips through the same
PrimitiveVectorElement impls predict read-backs exactly. The fixture's
multi-dense arm now applies the candidate datatype (matching the
CreateVectorName path), model_vector predicts per-row, and two new
initially-active candidates exercise the combination: "w"
(MultiDense(5), Float16) and "z" (MultiDense(3), Uint8).
The compile-time candidate check narrows to the combinations that
remain unpredicted: Turbo4 multi-dense (the multivector quantization
path differs from the per-vector turbo_storage_roundtrip) and sparse
with any datatype override.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* Make datatype match exhaustive in random_dense_vec
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* Address review findings on datatype plumbing
- Start candidate "c" active so the explicit Float32 schema path runs in
default soaks (CreateVectorName is FORCE_OFF by default)
- Single Float16/Uint8 roundtrip dispatch shared by the dense and
multi-dense arms of model_vector
- Hoist shared fixture builder plumbing into dense_params_builder
- Build one DenseVectorConfig literal in the CreateVectorName generator
- Inline single-caller datatype_of wrapper
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* Replace const-eval candidate check with a startup assert
Const eval forbids iterators, forcing an index-based while loop. A plain
function called at the top of run() reads better, still fails before any
op is applied, and names the offending candidate in the panic message.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* Fold INITIAL_ACTIVE into an initially_active candidate field
The hand-maintained list duplicated ALL_CANDIDATES (11 of 12 names) and
had to be kept in sync when adding candidates; forgetting it was silent
since CreateVectorName is FORCE_OFF by default, so a forgotten name got
zero default-soak coverage. Each candidate now declares its activation
inline and the fixture and run() filter on it.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
|
||
|
|
55f4219a3c |
Fix flaky WAL replay tests by waiting for update worker idle (#9832)
An empty update channel only means the last operation was received, not that it finished applying in spawn_blocking. Use plunge_async as a barrier before asserting point counts. Fixes #9831 Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
8cbd218d7c |
Fix flaky test_wait_deferred_does_not_block_update_worker (#9816)
Drain the update worker queue after setup upserts with WaitUntil::Wal. Wal only waits for the WAL write, so on slow CI (notably Windows) B could time out while queued behind the setup backlog rather than because the worker was blocked on A's deferred wait. Fixes #9814 Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
1d4d6f02da |
Per-query IDF corpus for sparse vector search (#9661)
* Add per-query IDF corpus for sparse vector search
Let the caller choose, per query, which population sparse IDF statistics
are computed over. `params.idf` is either `"global"` (default, unchanged
behavior) or `{"corpus": <filter>}`, where the corpus filter is
independent of - and usually broader than - the retrieval filter.
Decoupling the two keeps the score scale stable when the retrieval
filter tightens: term importance is measured against a population the
user names, not against whatever subset the filter happens to select.
Design decisions:
- Corpus grammar is restricted to a conjunction (`must`) of `match`
conditions on payload fields; loosening later is backward compatible.
- Strict mode validates the corpus filter like a read filter
(unindexed fields rejected).
- `idf` on a vector without the IDF modifier is a validation error,
never silently ignored.
- An empty corpus yields degenerate but corpus-scoped scores (smoothed
IDF over N=0), never a fallback to global statistics - in multi-tenant
collections a fallback would leak term statistics across tenants.
Implementation:
- QueryContext IDF stats are keyed by corpus, so one batch can mix
requests with different corpora.
- Statistics come from the sparse index: df(term) is counted over the
query terms' posting lists only, never by scanning stored vectors.
Small corpora (under ~1/32 of the segment, by cardinality estimate)
are kept as a sorted id list galloping through posting lists via
skip_to; large ones as a dense membership mask filled streaming from
the filtered-points iterator. A misestimated small corpus degrades
into the mask.
- Exposed uniformly: REST (`params.idf`), gRPC (`IdfParams` message),
edge python bindings; OpenAPI schema regenerated.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Apply rustfmt
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Fix clippy manual_is_multiple_of in sparse IDF corpus test.
Co-authored-by: Cursor <cursoragent@cursor.com>
* Allow any filter as IDF corpus
Drop the must+match grammar restriction on the corpus filter. A
restriction enforced only as a validation step over the full Filter
type buys nothing; if a narrower corpus syntax is ever wanted, it
should be a dedicated API-level type instead.
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
* Fix build: add memory field to SparseIndexConfig in idf corpus test
Co-authored-by: Cursor <cursoragent@cursor.com>
---------
Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Co-authored-by: root <111755117+qdrant-cloud-bot@users.noreply.github.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
|
||
|
|
a647c807d4 | Cleanup | ||
|
|
b31d98d55a |
[audit-M] Propagate the force-abort error in drop_shard_key
drop_shard_key force-aborted any resharding on the key being dropped, but swallowed a failed abort_resharding with a log-only error and dropped the shards anyway. That could leave resharding_state.json referencing the just-dropped key — a latent inconsistent load-time state. Propagate the error with `?`. A ServiceError halts consensus and retries after restart, which is safe because the abort and everything else in drop_shard_key is replay-tolerant; a user error dismisses the entry before any shard is dropped, which is equally consistent. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> |
||
|
|
d7c7052eb5 |
[audit-F] Use bad_request instead of service_error in pre-write transfer validations
A ServiceError returned from apply halts consensus on every peer (the entry can never be applied), deterministically stalling the whole cluster. The transfer Start validations that report a missing source or destination shard are reachable — e.g. a committed Start racing a resharding-abort or shard-key-drop that removed the shard — and all run before any durable write, so dismissing the entry with a user error is safe and correct. Convert five such sites to CollectionError::bad_request: - validate_transfer: source shard missing, and destination shard missing in both the resharding and filtered branches (helpers.rs); - start_shard_transfer: the source and target get_shard lookups (shard_transfer.rs). The genuine "single node deployment" service_errors in collection_meta_ops.rs are left untouched — those are real misconfiguration guards. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> |
||
|
|
7b2799a394 |
[audit-D] Make transfer Start validation replay-tolerant
Start's first durable write registers the transfer record; its last write sets the destination replica state. On re-apply after a crash between the two, check_transfer_conflicts found the operation's own half-applied transfer and returned bad_request, which is dismissed and never retried — so the peer permanently lacked the destination replica entry. Worse, a later Finish on that peer then silently skips both the destination promotion and the source removal, pinning a replica-set divergence. Exclude the transfer's own key from the conflict scan so a replay falls through and re-runs the (idempotent) start: register_start_shard_transfer is a set-insert and the destination replica-state write is absolute, so re-running reconciles the partial state instead of dismissing it. A genuinely conflicting transfer (different key touching the same shard/peers) is still rejected. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> |
||
|
|
3fd4e187dd |
[audit-J] Stop swallowing config.save errors in resharding
start_resharding, finish_resharding and abort_resharding each saved the updated collection config with a log-only `if let Err(err) = config.save(..)` that swallowed the failure. A swallowed save let the operation report success with a stale shard_number persisted on disk, arming a shard-dir/loader panic (or a silently unloaded shard) at the next restart. Propagate the error with `config.save(&self.path)?;` instead. The resulting IO error is a ServiceError, so consensus halts and retries the entry after restart. That is safe because in all three functions the config save is value-idempotent (guarded by `shard_number != new_shard_number`) and every step is replay-tolerant, so the retry converges. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> |
||
|
|
58d55ba923 |
[audit-H] Persist decremented shard count before dropping shard in abort_resharding
In the up-direction of abort_resharding, ShardHolder::abort_resharding drops the new shard's directory, but the config.params.shard_number update was persisted only at the very end of the function — after the transfers abort. A crash anywhere in that wide window left shard_number pointing at an already-deleted shard directory, which makes the auto-sharding loader panic on the missing dir at startup (crash loop), before the consensus replay that would reconcile the state can run. Move the shard-count update block to before the shard_holder.abort_resharding call. The block keeps its value-idempotence guard, so replay converges. The config write lock is taken while holding the shard_holder write guard, matching the shard_holder -> config ordering already used in finish_resharding, and is released before abort_resharding. The reverse crash window (dir still present but count already decremented) is benign and reconciled by replay. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> |