mirror of
https://github.com/qdrant/qdrant.git
synced 2026-09-25 07:27:41 -05:00
dev
4722
Commits
| Author | SHA1 | Message | Date | |
|---|---|---|---|---|
|
|
aefcb71b0f |
Meter remote requests and expose IO statistics to edge shard users (#10769)
* Meter remote requests and expose IO statistics to edge shard users Generalize the disk cache statistics from #10637 into a shared OpStats primitive and add a second layer counting every remote request issued through BlobFs and BlobFile, reads and writes alike. One helper at each call site drives the uio_trace request and the stats guard together, so the trace also gains the write operations. CachedBlobFs exposes both halves through stats(), DiskCacheFs exposes its remote filesystem, and both edge shards expose their filesystem, so an application can read cumulative totals. shard_update logs the IO counters moved by each phase; shard_query prints the combined dump. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Drop DiskCacheStats in favor of remote request stats Every disk-cache fetch is exactly one BlobFile read already metered by RemoteIoStats, so the cache half duplicated the remote Read/ReadFrom counters. CachedBlobFs::stats() now returns RemoteIoStats directly. The remote stats dump now includes per-op latency histograms, which were previously only printed for cache fetches. edge-shard-query and edge-shard-update initialize serverless_compatible feature flags instead of running with uninitialized ones. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * Sample uio-trace CPU usage every 10ms instead of 2ms Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * Address review: not-found bucket, empty-tail reads, strum Op - OpStats counts "not found" answers separately from errors, and the trace records them as a `not_found` outcome, so length and existence probes of absent objects no longer read as failures. - A read_from tail at or past EOF is settled after the len probe and counts as a completed empty read; a not-found error skips the probe, which could only repeat the answer. - Op derives strum EnumCount/EnumIter/IntoStaticStr instead of a hand-maintained ALL array. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> |
||
|
|
ee4981c280 |
Fix Edge query offset and group hit payload/vectors (#10778)
`PlannedQuery` fetches `limit + offset` points and leaves cutting the offset to the caller, which the server does at the collection level. Edge never did, so `query()` returned the first `offset + limit` hits. `query_groups` returned hits with only the `group_by` field as payload, ignoring the requested `with_payload`/`with_vector`. Hydrate the final hits once after grouping, as the collection-level group-by does. Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com> |
||
|
|
a8e0a9a4a5 |
[ReadOnlyShard] concurrent LIST files (#10771)
* list files concurrently Read-only shard load took each segment's LIST snapshot one after another, so cold start paid N sequential object-store round-trips. Take all segments' snapshots concurrently via a new `list_files_async` on `UniversalReadFsAsync`, then stage each open over its ready `CachedFs`. - `list_files_async` has no default, so no backend silently falls back to a blocking LIST; `BlobFs` spawns it on the `BridgeRuntime` and shares the traced, latency-logged future with the sync `list_files`. - `ReadOnlySegment::build_cached_fs_async` + `schedule_open_with_cached_fs` split `schedule_open` so the LIST can be awaited separately. - `CachedFs::schedule_open` polls once with `now_or_never` instead of a nested `block_on`, which would panic under the outer executor. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * use async LIST in live-reload --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com> |
||
|
|
1f0b3de20c |
[UIO] Fold *FileOps* into *Fs* traits (#10768)
* Merge UniversalReadFileOps into UniversalReadFs Every filesystem implemented both traits, so fold list/exists/from_context into UniversalReadFs and make UniversalWriteFileOps extend it directly. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * rename `UniversalWriteFileOps` -> `UniversalWriteFs` --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com> |
||
|
|
25be44d8eb |
Split edge shard_query and shard_update tools into modules (#10766)
Both tool binaries were single main.rs files of about 1100 and 840 lines. Move their items into per-concern modules without changing behavior: CLI definitions, argument parsing, backend configs, the prepared request and reporting for shard_query; CLI, backend configs, schema reading, point generation, dry run and apply for shard_update. Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> |
||
|
|
e376bc22b0 |
test(model_testing): make harness runs replayable and print the replay command (#10707)
Harness failures print their seed, but the seeded tests cannot be re-run with it and the soak binary needs hand-built flags, so a CI failure was effectively one-shot: - Honour MODEL_TESTING_SEED so the seed from a failing run replays the same op sequence in the same test; without it a fresh seed is still drawn. - Print the equivalent model_testing invocation for every run, so any failure (including the Linux-only harness tests on a non-Linux machine) becomes a one-line reproduction. - Bind the harness knobs once in smoke() and feed the same bindings to run() and to the printed command, so a run and its reproduction cannot disagree about them. - Cover both helpers with platform-independent unit tests: the helpers now compile and are tested on every platform, not only inside the Linux-gated harness module. Refs #10406, #10662, #10667. |
||
|
|
e695fc4487 |
Remove shard initializing flag under the local shard lock (#10752)
* Remove shard initializing flag under the local shard lock A late snapshot recovery from an aborted transfer and a new incoming transfer could race on the shard initializing flag: initiate_receiving removed the flag after releasing the local shard lock, deleting the flag a concurrent restore_local_replica_from had just created. The restore then failed with ENOENT on its own flag removal, and a crash in between would have left half-moved shard data without a dirty marker. Move the removal into init_empty_local_shard, under the local write lock and after the empty shard is built, tolerating a missing flag. This also fixes the try_exists(..).is_ok() check, which was true for absent files. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> * Install empty local shard only after removing the initializing flag If flag removal fails or the future is cancelled while it is pending, local keeps the dummy placeholder instead of a Local shard next to a leftover dirty marker, so a retried transfer re-initializes the shard and the documented cancel-safety holds. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com> |
||
|
|
19d4da9988 |
UIO tracing and visualizer (#10742)
* Add uio_trace module and visualizer * Handle object store, not just gRPC * Improve visualizer * Improve visualizer [2] |
||
|
|
878843e6e0 |
Unify the blobstore tracker read interface across storage modes (#10735)
One TrackerRead trait, in lib/blobstore/src/tracker/read.rs, now covers the Gridstore Tracker and ReadOnlyTracker and the Logstore AppendOnlyTracker. LogstoreView, LogstoreReader, and validate_consistency are generic over it, as GridstoreView already was over the Gridstore-only trait. The trait gains get_range and an access pattern on get, and its iter returns impl Iterator instead of the concrete Gridstore Iter, which drops the storage type parameter from the trait. The Gridstore trackers implement get_range through a new read_slots helper; nothing in Gridstore calls it yet. The lifecycle methods of LogstoreReader (open, preopen, files, live preload, live reload, clear cache) stay in an impl block bound to AppendOnlyTracker, so a future immutable in-RAM tracker plugs into the read path without inheriting reload semantics. No runtime change: Gridstore slot reads keep the Random access pattern, and PointerItem stays the iterator item type. Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> |
||
|
|
06433b7028 |
[BM25] Add a BM25-over-sparse baseline benchmark (#10677)
* Add a BM25-over-sparse baseline benchmark The text payload index is meant to score about as fast as BM25 over sparse vectors, so the number it has to match needs to exist before the scorer does. Measures a local shard end to end, no HTTP. - embeds the corpus through `lib/bm25` with its defaults, so the baseline is the route a user migrates from rather than a reimplementation - Zipf-like vocabulary. On a uniform one every term is equally selective, IDF is flat and pruning has nothing to prune, which would flatter any scorer measured against it - two shards rather than one shard before and after optimization: a shard that will optimize starts as soon as the upsert lands, so the first cut timed a half-converted index and called it fresh. Both states assert what they hold before anything is timed Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Measure the sparse BM25 baseline where the routes separate - 200k documents by default (BM25_SPARSE_DOCS overrides): at 20k every shape of both routes measures the same and half of a shard-level query is the shard; reachable through the shard since #10682 - a third shard with the sparse index on disk, which the optimized state never exercised - recall at 10 against BM25 by definition, printed per state: the default avg_len of 256 on a corpus averaging 110 tokens misses a quarter of the true top 10, so the optimized state is also timed with the corpus average - corpus, queries, reference and recall move to segment::fixtures::bm25_corpus, to be shared with the text-index bench and the comparison harness - module doc: the sparse shapes measure within a few percent of each other; the states exist for the text comparison Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Refuse an empty corpus in the sparse BM25 baseline BM25_SPARSE_DOCS=0 built empty shards, made the average length NaN and scored every empty truth as recall 1.0. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Measure what the sparse BM25 baseline states claim The fresh shard kept the default 10 MB indexing threshold and optimized itself in the background, so it was timed as a second optimized state. Disable indexing and re-check it after timing. Keep the shard storage under CARGO_TARGET_TMPDIR so the on-disk index is not read from a tmpfs. Fix doc comments that named missing files, a no-op IDF clamp and the wrong reason for the empty-corpus guard. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
87dc39bfdc |
[Diskcache] piggyback remote reads across pipelines (#10690)
* AI: Share reads across different diskcache pipelines * AI: Unify scheduled reads into single ScheduledRead struct * AI: Implement follower promotion on fetch abandonment * AI: self-promote if waiting for an abandoned placeholder * AI: simplify * AI: use Weak<Placeholder> instead of manually counting strong refs * AI: track abandoned waiters * AI + manual: park with timeout * AI: fix barrier in test * fmt with latest nightly * expect remote pipeline |
||
|
|
6a03f4288b |
Implement resharding operations for ConsensusStateMachine (#10697)
|
||
|
|
c296898355 |
[BM25] Add document length into the immutable and on-disk text indexes (#10639)
* Carry document length into the immutable and on-disk text indexes The mutable index records `doc_len`. The immutable and on-disk backends dropped it on conversion. They now carry it and persist it as a sidecar, so a length survives an optimizer run and a restart. - `ImmutableInvertedIndex` gains `point_to_doc_len`, parallel to `point_to_tokens_count` and zeroed wherever that vector is, so summing it never counts a deleted document. - `OnDiskInvertedIndex` writes `point_to_doc_len.dat`, only when the index records lengths. `files()` lists it only when it exists, so a snapshot carries it exactly when there is one. - Deletions are masked on load, not at build time. The file is written once and a point deleted later through the id-tracker is zeroed in `TryFrom<&OnDiskInvertedIndex>`. No segment total is stored, since it could only be summed after that masking. - A sidecar shorter than the counts is treated as absent. It can only come from a partially copied file set, and padding it would give every point past the truncation a zero that reads like a real length. - A missing sidecar on an index that should have one makes `new_mmap` report the index absent, which routes into the existing rebuild from payload. The check lives there rather than in `OnDiskInvertedIndex::open` because the read-only stack never builds and would drop the field instead. - `FullTextMmapIndexBuilder::add_many` reached below `index_str_tokens` and so declared no length at all. It now tokenizes through the shared helper and measures like the gridstore path. Recording is still gated behind `TextIndexParams::scoring()`, so none of this is written today. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Wipe the text index directory instead of its listed files `files()` lists the `doc_len` sidecar only when the index loaded it, so a truncated sidecar is omitted. `wipe` then deleted the listed files, failed to remove the now non-empty directory, discarded that error and returned success, leaving the directory and a stale sidecar behind. Remove the directory itself. Each field owns its own `{field}-text` directory, so this no longer depends on `files()` being an accurate inventory, and real removal errors propagate rather than being swallowed. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Drop the fully qualified visibility paths in the text index `pub(in crate::index::field_index::full_text_index)` is tedious to read and to keep correct, and it grants nothing here: `inverted_index` and `mutable_text_index` are private modules of `full_text_index`, so their types are not nameable from outside it whatever the field visibility says. Plain `pub` is no wider in practice. Widening `Storage` surfaced two `private_interfaces` warnings, since `ZerocopyPostingValue` and `PostingListHeader` in `types.rs` were still behind the long path; those move too. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Harden the document length sidecar after review - check for the sidecar before opening the index rather than after. The open populates the whole file set, so the first start after scoring is enabled would fault in every segment's postings only to discard them - treat a sidecar that covers more points than the index as untrustworthy, not just one that covers fewer, and warn with both counts the way `SortedBlockIndex::open` does. A longer one used to be accepted and then silently truncated when materialized - unlink a stale sidecar when a build records no lengths. It was the only file here that could outlive the build that wrote it, and `open` would have read it as this build's - say why an index is being rebuilt from payload instead of leaving a silent full re-index announced at debug level - read `phrase_matching` from the config in the mmap builder, like the other two callers of the same helper pair, so the two halves of the sentinel rule cannot drift apart - assert the length invariant the writers actually maintain, and correct what `files()` and the field comment claim Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Address review on the document length sidecar - shrink the document length vector when it is handed over to the immutable index, it arrives with the mutable index's doubling capacity - narrow the immutable index fields to pub(super), with a test-only accessor for the one reader outside the module - let wipe propagate a missing index directory instead of treating it as success, nothing reaches it with the directory already gone Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
c218fa67c0 |
Configurable key prefix for object storage snapshots (#10715)
* Configurable key prefix for object storage snapshots Snapshot object keys were always derived from the local snapshots path, so every deployment sharing a bucket wrote under the same `snapshots/` root. Each cloud config block gains an optional `prefix`, and objects become `<prefix>/snapshots/...`. The prefix is applied by wrapping the client in `PrefixStore`, so the snapshot operations and the names returned by the API are unchanged. Leading, trailing and repeated slashes are dropped, and an empty prefix is a no-op. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Test that snapshot operations cannot escape the storage prefix Hostile targets are handed to the cloud manager directly, past `validate_snapshot_name`: parent references, absolute paths, encoded slashes, backslashes and empty paths. Writes must stay under the prefix, and objects planted outside the prefix must be invisible to list, download, stream and delete. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Normalize empty prefix components Refactor prefix handling to remove empty components. Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --------- Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> Co-authored-by: Tim Visée <tim+github@visee.me> Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> |
||
|
|
6cd816acac |
Add match: { substring } filter condition (#10711)
* Add `match: { substring }` filter condition
Unindexed `text` and `text_any` matching became token-aware in #10341 and
#10593. Users who relied on the old raw substring behaviour get it back as
an explicit condition: `match: { "substring": "..." }` selects points with a
string value containing the given string, byte-wise and case-sensitive,
consistent with exact keyword and prefix matching.
Execution: a keyword index (with or without the `prefix` option) serves the
condition by scanning its value dictionary and uniting the postings of the
matching keys; cardinality reuses the prefix estimator, generalised into
`keys_union_cardinality`. The per-point checker goes through the forward
index. Without a keyword index the condition falls back to reading the
payload. Text, bool, integer and uuid indexes decline it.
Strict mode: the condition requires the `KeywordMatch` capability, so with
`unindexed_filtering_retrieve: false` it is rejected on unindexed and on
text-indexed fields and allowed on any keyword index.
API: `MatchSubstring` in the REST `Match` union with regenerated OpenAPI,
gRPC `Match.substring = 12`, edge python `MatchSubstring`, edge ffi
`Match::Substring`.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
* Test substring fallback on a text-indexed field, document estimator params
A text index cannot serve `substring`, so on a field that has only a text
index the condition runs through the payload fallback; only strict mode may
reject it. Pin that in the OpenAPI suite and reword the strict-mode unit
test comment, which read as if the text index itself blocked the query.
Also spell out what `keys` and `postings` mean in `keys_union_cardinality`.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
* Serve `match: { substring }` from the keyword key dictionary
The condition used to enumerate keys through `MapIndexRead::for_each_value`,
which on the on-disk variant drags the whole `value_to_points` file through
`for_each_entry`, plus one random read per matching key for its postings
count. Query planning paid that scan in full, before deciding whether to use
the clause at all.
Route it through the `prefix_index.bin` key dictionary instead: front-coded
keys with their postings counts inline, no postings. Estimation now reads
keys only and never touches `value_to_points`; filtering takes the matched
key list and resolves postings in one batched read, as prefix matching
already does.
This makes the `prefix` option a requirement: a keyword index without it has
no key dictionary, so it declines the condition and falls back to the payload
scan, the same as a text index. Strict mode follows — substring now infers
`KeywordPrefix`, so `unindexed_filtering_retrieve: false` names
`keyword (with prefix: true)` as the index to create.
`PrefixIndex::for_each_key` reads blocks in ~1 MiB chunks rather than the
whole key section at once: a substring cannot be pruned by the block index,
so the one-shot read would grow with the dictionary.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
* Sync MatchSubstring OpenAPI description with Rust docs
After routing substring matching through the prefix key dictionary,
the schema docstring required regenerating so OpenAPI stays consistent.
* Scan the map index keys for substring match without a dictionary
Without the keyword dictionary a substring condition was declined by the
field index and left to the per-point condition checker, which reads the
forward index for every candidate point. Enumerate the distinct keys of
`values_to_points` instead: the same one-pass-over-distinct-values shape as
the dictionary scan, only over a structure that interleaves keys with their
postings. Filtering and cardinality estimation are then always served, so
the condition can act as a primary clause on a plain keyword index.
Prefix matching keeps its per-point fallback: an ordered dictionary is what
makes a prefix a bounded range, and enumerating every key to answer one is
not a trade worth making implicitly.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
* Reject substring matching in strict mode
A substring condition is answered by looking at every distinct value of the
field: no index gives it a bounded access path, so there is no index a user
could create to make it affordable. Reject it under strict mode instead,
wherever a filter reaches verification — read and write filters, nested
sub-filters, and prefetch filters.
Filter limits are now checked before the unindexed-field check, so the
rejection is not reported as "create an index for this key", advice that
would lead to the same rejection afterwards.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
* Update the strict mode substring test to the new rejection
The test asserted that substring filtering under strict mode asks for a
keyword index with the `prefix` option. It is now rejected whatever index
the field carries, so every case in the test gets the same answer.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
* Estimate a substring condition without scanning
Counting the keys a substring matches costs the same scan as answering the
condition, and `filter` then repeats it to collect those keys. Report the
uninformed estimate instead — the one an unindexed condition has always
reported — and keep the primary clause, so the scan happens once, in
`filter`, and only when the planner picks the condition to drive iteration.
With no counts to collect, `substring_scan` collapses into `substring_keys`:
the in-RAM variants no longer look up a posting count per matched key.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
* Don't to parse everything as UTF-8
---------
Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
Co-authored-by: qdrant-cloud-bot <111755117+qdrant-cloud-bot@users.noreply.github.com>
Co-authored-by: timvisee <tim@visee.me>
|
||
|
|
baec6a3119 |
Support GCS and Azure Blob Storage for collection snapshots (#10714)
Snapshot storage was limited to S3 although the object_store dependency already ships the GCS and Azure backends. `snapshots_storage` now accepts `gcs` (alias `gcp`) and `azure`, configured through `gcs_config` and `azure_config` blocks next to the existing `s3_config`. The legacy S3 shape is unchanged. Client construction is split into one builder per backend, all sharing the Qdrant user agent and the plain-HTTP rule for `http://` endpoints. Startup warns when a config block for an unselected backend is present. The e2e snapshot recovery test is parameterized over the cloud backends. The GCS case is skipped because fake-gcs-server does not implement the XML multipart upload API that object_store uses for GCS uploads. Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> |
||
|
|
8f44ddd8ba | Combined storage tests (#10676) | ||
|
|
0a0fd87909 |
Count documents, not points, in the mutable text index (#10675)
A value that tokenizes to nothing is indexed and matches nothing. The mmap build already treats it as no document, folding it into the "no tokens" mask, while the mutable index counted it, so `points_count` for the same data depended on whether the segment had been optimized yet. - `index_tokens` counts a transition rather than incrementing, so a point becomes a document when it gains its first token and stops being one when rewritten to nothing. Re-indexing the same point no longer counts it twice - `remove` decrements only when the removed token set had tokens - the builder counts the same way `points_count` feeds `count_indexed_points` for cardinality estimation, and is `N` for the IDF of a text score. Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
5abf55f2c7 |
test/fix-flaky-building-cancellation (#10696)
* test(segment): replace timing-based building cancellation test `test_building_cancellation` compared wall-clock times of independent builds cancelled after fractions of a measured baseline. Its tolerances had to be widened repeatedly (#2039, #8243, #10346) and it still failed on Windows CI, e.g. when an early stop landed in a slow non-cancellable setup step (time_early: 969, time_later: 631, baseline 2488). Replace it with tests that target build phases explicitly and measure work instead of time: - test_building_cancelled_before_start: a build cancelled up front returns `Cancelled` without starting any vector index work. - test_building_cancelled_during_main_graph: a watcher cancels the build once the main HNSW graph (observed via the progress tracker) reaches 1000 of 10000 points. The build must return `Cancelled`, leave the graph unfinished, and insert at most 2 points per build thread after the flag is set (only in-flight insertions may complete). Both also check that a cancelled build leaves nothing behind in the segments or temp directory. Use random vectors with a fixed seed instead of identical zero vectors. Verified the main-graph test fails when the per-point stop check is removed, or only done every 64 points (102 points after cancel, limit 16). Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * test(segment): cover cancellation during HNSW graph healing it: an uncancellable heal phase is what blocked consensus when removing a replica during an optimization. Track progress of the `migrate` phase (healed `(point, level)` pairs out of the total), like `main_graph` already does. Besides showing healing progress in optimization telemetry, this lets a test tell "stopped right away" apart from "healed everything, then noticed the flag"; both return `Cancelled`, because the flag is checked again right after healing. Generalize the main-graph cancellation test into a helper that cancels any phase at 10% of its work and checks that at most 2 items per build thread complete afterwards, and add test_building_cancelled_during_heal: it builds an HNSW segment, deletes a quarter of its points (below the default healing_threshold of 0.3), and cancels the rebuild while healing. Verified the test fails with the heal stop check removed ("migrate was completed despite cancellation"). Measured: 8 items after cancel with 8 build threads, out of 3000-5700 to heal. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
173f9d3ed6 |
Split HNSWIndex::build into phase modules (#10709)
* Split additional links phase out of HNSWIndex::build Move the per-payload-block subgraph phase, together with its existing helpers condition_points and build_filtered_graph, from HNSWIndex::build into hnsw/build/additional_links.rs. The phase body is moved verbatim; the borrowed locals become explicit parameters, and the function returns the number of vectors indexed through the subgraphs. GPU insert context creation moves to gpu_build.rs next to the other GPU setup helpers, so build.rs no longer needs feature-gated GPU imports. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Split main graph CPU phase out of HNSWIndex::build Move old-index healing, the single-threaded warm-up and the parallel point insertion into hnsw/build/main_graph.rs as build_main_graph_on_cpu, mirroring build_main_graph_on_gpu. Level assignment and the GPU attempt stay in the orchestrator, since the GPU result decides whether the CPU path runs at all. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Extract config derivation, field subtasks and thread pool from HNSWIndex::build Share the kilobytes-to-vectors full scan threshold conversion between load_or_derive_config and build as derive_config, with the vector count as an explicit argument so both call sites keep their existing divisor. Move the per-field progress subtasks into additional_links_fields next to the phase that consumes them, and the rayon pool with its low-priority spawn handler into build_thread_pool. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Extract GPU attempt, graph saving and the deleted-points check from build upload_and_build_main_graph in gpu_build.rs owns the "any graph will be built" gate, the upload, the main graph attempt and its timing log, so build.rs keeps one cfg pair and adopts the GPU graph with an if-let. save_graph builds the links format param and calls into_graph_layers in one place, since the param borrows the inline vectors. The debug-only walk over deleted points becomes a predicate over links_empty under debug_assert!. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> |
||
|
|
6dd9b8002a |
Cancellable open and live_reload for read-only edge shards (#10699)
* Cancellable open and live_reload for read-only edge shards Add cooperative cancellation to ReadOnlyEdgeShard::open and ReadOnlyEdgeShard::live_reload, following the contract of EdgeShardReadWithCancellation: a shared stop flag is checked between stages and between segments, a set flag yields OperationError::Cancelled and never a partial result, and the flag is never set or reset by the callee. New entry points: ReadOnlyEdgeShard::open_with_cancellation and ReadOnlyEdgeShard::live_reload_with_cancellation. The existing open, live_reload and live_reload_with keep their signatures and delegate with a fresh flag. The flag is propagated into ReadOnlySegment::schedule_open. The staged handle carries it, so both the prefetch staging and finish check it between components. The edge loader propagates a Cancelled error instead of logging it as an unloadable segment. The holder swap and the config re-derivation that follows it are one indivisible step, so a cancellation never leaves a config lagging behind the segment set. Segment reloads stay atomic under their write lock, so a cancelled live_reload leaves the shard consistent and the next one continues from there. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Check cancellation between vector component preopens A preopen polls its open future once, so real work starts right away. Checking the flag between the storage, quantized vectors and index preopens of a dense vector, and between the storage and index preopens of a sparse vector, keeps the staging contract that the flag is observed between components. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> |
||
|
|
76a8647728 |
Follow-up tests for #10349: pinned findings on the merged code (#10661)
* test: pin persisted proxy segments follow-up findings * fix: tests * Respect `up_to` in flushing pending proxy changes * Fix unproxy divergence, centralize logic in single shared function * Don't pass locked segments we don't use * Close the post-swap window losing acknowledged proxy changes `finish_optimization` left a window between `swap_new` and the end of the function where the propagated proxy changes were durable in no reachable place. The proxies had left the holder, so a flush pass no longer saw them as unsaved work and acknowledged the WAL past what their pending changes logs persisted, while their source files were still on disk contradicting the newer state. Their `ack_pin` was only registered at the very end, and the optimized segment, the only durable home of those changes, had no version file yet, so `normalize_segment_dir` deleted it on the next load. A failure or crash in that window lost every change past the proxies' logs. Close it from both sides, reordering only: - Save the optimized segment's version file right after the flush that made the propagated changes durable, before the swap. That flush is what `SegmentBuilder::build` postponed the version file for, so the segment is loadable from the swap onwards. - Register the proxies' deferred destruction (and with it their `ack_pin`) under the same write lock that evicted them, so no flush pass can observe the holder without either the proxies or their pin. Collecting the deferred point ids moves up with the registration, as it borrows the swapped-out proxies. The deferred destruction still cannot run before the manifest is synced: `locked_proxies` holds the segments alive until the end of the function, so `try_drop_data` retries until then. * Rename the optimization test hook after the window it guards The hook no longer sits before the version save, it marks a failure anywhere in the window after the optimized segment was swapped in. Rename it and the test accordingly, and restate the test doc as the invariant that window must uphold rather than the bug it used to describe. * Pin WAL ack while creating snapshot * Patch test that was stuck * Initialize necessary feature flags in tests * Correctly propagate changes in two stages, lock updates on second stage * Use existing WAL ack pinning infrastructure * test: pin unproxy phase 2 propagation failure losing acknowledged changes * test: pin WAL ack pin at zero suppressing clock persistence * Fix propagate and unproxy data consistency error on failure * Store clocks before checking WAL ack pin * fix: wait for the flush worker when stopping it in tests * fix: linter --------- Co-authored-by: timvisee <tim@visee.me> |
||
|
|
651ca78d2d |
Fix small typos in comments across segment and storage crates (#10681)
* Fix small typos in comments across segment and storage crates - AddLearnerPeer -> AddLearnerNode (the ConfChangeType name used by the matched WAL entries) - 'adding adding' -> 'adding' in the progress_tracker debug assertion - CanellationToken -> CancellationToken in the snapshots recovery comment - 'any of of' -> 'any of' and 'valuse_set' -> 'value_set' in json_path comments - 'is not not changed' -> 'is not changed' in the mutable null index * Style: collapse the debug_assert to one line per rustfmt Nightly rustfmt prefers the single-line form now that the message fits (follow-up to the typo fix). |
||
|
|
be8fe9a678 |
Fix HNSW builds randomly running on fewer threads than they were given (#10686)
* Fix HNSW builds randomly running on fewer threads than they were given * Move HNSW_BUILD_MAX_PAR_LEN next to SINGLE_THREADED_HNSW_BUILD_THRESHOLD |
||
|
|
20229f99ba |
[combined-storage] Combined storage write (#10669)
* HNSWIndex::build(): add `inline_vectors` arg Let the caller decide whether to use `inline_vectors` format. * SegmentBuilder::build: finalize GraphInline vector storage Instead of old "graph-with-vectors plus regular vector storage", keep only the graph-with-vectors storage. * Gate the GraphInline segment build behind a feature flag |
||
|
|
c00e2aa7d5 |
Pin WAL acknowledgements from multiple places with pin guard (#10672)
* Rename `wal_keep_from` to `wal_ack_pin` The parameter pins the WAL acknowledge, the new name says so. * Support pinning the WAL acknowledge from multiple places at once The WAL acknowledge had a single pin slot, a shared `AtomicU64` any second user would have clobbered. Replace it with `WalAckPins`, holding any number of live pins. The flush worker never acknowledges at or past the lowest of them. Taking a pin hands out a `WalAckPinGuard` that releases when dropped, so the queue proxy shard no longer has to release it by hand. The registry holds its pins weakly, which means dropping the guard is all it takes and a queue proxy lost to an unwind can no longer stall the WAL acknowledge forever. * Tests for the WAL acknowledge pins * Debug assert that set does not move version backwards * Update tests |
||
|
|
de42b2c7b9 |
Implement CreateShardKey/RemoveShardKey for consensus state machine (#10666)
|
||
|
|
cfe3bd625c |
Fix stale UniversalMapIndex rustdoc links in the map index docs (#10680)
universal_map_index::UniversalMapIndex does not exist; the on-disk variant is on_disk_map_index::OnDiskMapIndex (the third MapIndex variant). Repointed the three rustdoc links accordingly. |
||
|
|
901fa8e865 |
Stop rewriting whole sparse posting lists on every upsert (#10682)
`PostingList::upsert` ends in `propagate_max_next_weight_to_the_left`, whose doc
comment states "If an entry has a weight larger than `max_next_weight`, the
propagation stops". It never stopped — the loop always walked the entire prefix.
Record ids ascend during an upload, so every insert lands at the end, walks the
whole list, and growing a posting list to length L costs O(L²). On the NeurIPS
2023 sparse base set (MS MARCO / SPLADE) the hottest dimension appears in ~66% of
documents, so at 1M points its posting list holds 660k elements and every new
point rewrites all of them.
Two changes to `PostingList`:
* restore the early exit — every entry satisfies the same recurrence the loop
walks, `max_next_weight[i] = max(max_next_weight[i + 1], weight[i + 1])`, so
once an entry already holds the value being written, every entry to its left is
correct too;
* add an append fast path — when the incoming record id is past the last stored
id, push directly instead of binary searching a list that can hold millions of
elements.
Neither changes what is stored. An index grown by `upsert` stays identical to one
built by `InvertedIndexBuilder`, `max_next_weight` included, which is what a
segment reload depends on; searches return identical results.
Building an `InvertedIndexRam` one point at a time, SPLADE vectors:
points before after
100k 2,216/s 124,016/s
1M 250/s 114,422/s
Uploading 1M of those points to a single node: 72.2s -> 23.7s with default
collection settings, and 839.8s -> 29.5s with `indexing_threshold: 0`, which
stops segments converting and is the documented way to speed up a bulk load.
Claude-Session: https://claude.ai/code/session_01HXjykBsMZaRNrP17o2fuEJ
Co-authored-by: Andrey Vasnetsov <andrey@qdrant.com>
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
||
|
|
3eefe44572 |
Fix never-ending optimization loop with vectors:{memory:cached} (#10664)
* Add `test_optimizers_should_settle` * Fix `"memory": "cached"` endless optimization loop |
||
|
|
88768e1b3d |
fix: expose ReshardingStage in telemetry (#10618)
* fix: expose ReshardingStage in internal telemetry Useful for cluster-manager operations. Non-breaking change for existing API. The current field values look stable enough, so it's not worrisome to keep them stable. Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com> * fix: address comment, also document uuid for symmetry Already returned, just document it in OpenAPI schema. Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com> --------- Signed-off-by: Anton Antonov <anton.synd.antonov@gmail.com> |
||
|
|
b2b509cf73 |
Tests for #10349: coverage gaps and pinned findings (#10445)
* shard: persist proxy pending changes on flush, stop holding back WAL ack Hook the pending changes component into the proxy segment's flush: the proxy flusher first persists the buffered operations into the pending changes log, then passes the flush along to the wrapped segment. The proxy's `persistent_version` now covers what the log durably holds on top of what the wrapped segment persisted itself. That is what lifts the WAL cap proxies imposed so far. `flush_all` compares each segment's version against its persistent version; a proxy used to report only the wrapped segment's persisted version while its own version climbed with every buffered operation, so the WAL could never be acknowledged past the point the proxy was created at, and a restart replayed all of it — potentially very expensive operations, such as an update by filter, all over again. With the buffered state durable on disk the generic rule acknowledges the full version, and a restart recovers it from the log instead. Dropping a proxy's data drops the component first, which waits for any in-flight pending changes flusher so it cannot append to the segment directory while that is being deleted. Update the proxy flush test to the new semantics, add a segment holder test asserting the acknowledged version advances past a proxied delete, and update the ack pin rationale in `finish_optimization`: the pin is still needed after the proxies leave the holder, it just snapshots a persistent version that now includes the log. * Persist wrapped segment before pending changes Prevents raising version of proxy segment too early * collection: test crash recovery through persisted proxy changes End-to-end test of the persisted pending changes: wrap every segment of a local shard in a proxy, delete points so the deletes are only buffered, flush, and assert the acknowledgeable version covers them. Then acknowledge the WAL up to that version, drop the shard without ever propagating the proxies, and load it again: the deletes are gone from the WAL and must come back through the pending changes logs. The delete under test is deliberately not the last WAL entry, as the acknowledge never passes the last entry and that one is always replayed. * test: cover untested pending proxy changes paths * test: pin persisted proxy changes findings * Fix bad merge * Remove corrupt length test, we cannot detect if last entry was corrupt Remove a test that asserts an entry that isn't the last cannot have a corrupt length. We cannot reliably detect whether the invalid length was the last entry or not, because nothing else tells us how many entries we expect in the file. At the same time we don't expect random bit flips. So I removed the test. * Simulate segments flush to clear pending changes log file * Update test, also assert proxy segment version * Fix bad merge * Fix blocked test * Enable necessary feature flags in tests (2/2) --------- Co-authored-by: timvisee <tim@visee.me> |
||
|
|
cc7a209c76 |
Persist proxy segment changes across restart, don't stall WAL ack's (#10349)
* segment: move proxy pending change types into segment crate Move the types describing the changes a proxy segment buffers — point deletes (`ProxyDeletedPoint`), payload index changes (`ProxyIndexChange`, `ProxyIndexChanges`) and vector name changes (`IntendedVector`, `ProxyVectorNameChanges`) — from `shard::proxy_segment` into a new `segment::pending_changes` module. Pure move, no behavior change: the proxy segment re-exports them from their old location. Having them in the segment crate lets both the proxy segment and the segment load path share them, in preparation for persisting pending proxy changes to disk and replaying them on restart. * segment: add PendingChange describing a persisted proxy operation Add the `PendingChange` enum with one variant per operation type a proxy segment buffers — point delete, payload index change, vector name change — each carrying the operation version it was issued with. This is the shape in which pending proxy changes are persisted to disk. Derive serde on it and on the buffered change types it embeds, so entries can be serialized into a log file and read back. `PartialEq` on those types lets a persisted batch be matched against the in-memory pending buffer after a flush. * segment: add PendingChanges component persisting proxy changes to a log Add `PendingChanges`, the component that manages the operations a proxy segment buffers for one proxy layer, and persists them to disk so they no longer only live in memory. It keeps the same per-type buffers the proxy segment served its reads from (point deletes, payload index changes, vector name changes), plus a single registration-ordered buffer of everything not yet persisted. `flusher()` writes that buffer into an append-only log file inside the wrapped segment's directory: `pending_changes.log` for the inner most proxy layer, with the layer number as a suffix for each layer above it. Appends follow the mutable ID tracker: all new entries are serialized into one buffer and written with a single call on an append-mode file, then fsynced, so a crash can only leave a torn entry at the very end. Loading truncates such an entry — its operations were never durable and thus never acknowledged in the WAL — but fails hard on a malformed entry in the middle, which cannot be explained by a torn append. The component tracks the highest operation version the log covers. Every registered operation at or below it is either durable in the log or was a no-op that does not need recovery; a flusher advances it to the proxy's version even when there is nothing to write. The pending buffer is deliberately not cleared when the proxy propagates its changes to the wrapped segment, as that only makes them durable once the wrapped segment flushes. Replaying an entry twice is a version-gated no-op. A log file left behind by a previous proxy on the same segment is adopted by `open()`: new entries are appended after it and its highest version is taken over, while its entries are not loaded into the buffers as they are already applied to the segment. `load()` also reconstructs the buffers, for callers that do want the buffered state. * segment: replay persisted pending proxy changes onto a segment on load Add `recover_pending_changes`, to be called when a segment is loaded on restart, before regular WAL replay. If the segment directory holds pending changes log files, the proxies that wrote them did not propagate their buffered state into the segment before the process stopped. Instead of reconstructing the proxies, replay all logged operations directly onto the segment: inner most proxy layer first, each file in append order, through the regular version-gated segment operations (`apply_change`). Entries the segment already applied are silently skipped, so a stale file is harmless. The segment is force-flushed before the files are removed; a crash in between merely replays the files once more. * segment: test PendingChanges component Cover the pending changes component: registering and flushing each operation type and reconstructing the buffers from the log, log file naming per proxy layer and gap-tolerant listing, covering the proxy version without entries, operations registered while a flusher is captured, flushers of a dropped component, torn-tail truncation versus mid-file corruption, adoption of an existing log, and replaying logs onto a real segment: fresh, stale (already applied), multi-layer, and vector name changes. * segment: include pending changes logs in segment snapshots Register the pending changes log files of a segment in its snapshot: add them to `snapshot_files` next to the segment state and version files, existence-guarded, and to the segment manifest as unversioned files. Full, partial and streamed snapshots therefore all carry them. The recovery side needs no changes: a restored segment is loaded like any other, which replays and removes the logs. * shard: back proxy segment pending changes by PendingChanges component Replace the proxy segment's separate `deleted_points`, `changed_indexes` and `changed_vector_names` fields with a single `PendingChanges` component. Reads keep going through the same per-type buffers, now behind accessors; writes go through the component's `register_*` methods, which additionally queue every operation for persistence. Opening the component is fallible, as it adopts a pending changes log a previous proxy may have left in the wrapped segment's directory, so `UnsyncedProxySegment::new` now returns a result. Wrapping another proxy opens the next proxy layer up, writing to its own dedicated log file. No behavior change yet: the proxy still flushes and reports persistence exactly as before, nothing is written to the log. * shard: persist proxy pending changes on flush, stop holding back WAL ack Hook the pending changes component into the proxy segment's flush: the proxy flusher first persists the buffered operations into the pending changes log, then passes the flush along to the wrapped segment. The proxy's `persistent_version` now covers what the log durably holds on top of what the wrapped segment persisted itself. That is what lifts the WAL cap proxies imposed so far. `flush_all` compares each segment's version against its persistent version; a proxy used to report only the wrapped segment's persisted version while its own version climbed with every buffered operation, so the WAL could never be acknowledged past the point the proxy was created at, and a restart replayed all of it — potentially very expensive operations, such as an update by filter, all over again. With the buffered state durable on disk the generic rule acknowledges the full version, and a restart recovers it from the log instead. Dropping a proxy's data drops the component first, which waits for any in-flight pending changes flusher so it cannot append to the segment directory while that is being deleted. Update the proxy flush test to the new semantics, add a segment holder test asserting the acknowledged version advances past a proxied delete, and update the ack pin rationale in `finish_optimization`: the pin is still needed after the proxies leave the holder, it just snapshots a persistent version that now includes the log. * shard: propagate proxy changes when unwrapping on optimizer cancel When an optimization is cancelled or fails, `unwrap_proxy` puts the wrapped segments back into the segment holder. Propagate the changes buffered in each proxy into its wrapped segment first, as the snapshot unproxy path already does, instead of dropping them with the proxy. The pending changes log is deliberately left in place when unwrapping: deleting it before the wrapped segment has flushed the propagated changes would not be crash safe. It is cleaned up on restart and when the segment directory is dropped, and a new proxy on the same segment adopts and appends to it; replaying a stale file is safe because all operations are version gated. * shard: test persisted proxy pending changes Test the proxy segment against its persisted pending changes: buffered changes survive dropping the proxy without propagation and are replayed onto the segment when it is loaded again; unwrapping leaves the log in place and a new proxy on the same segment adopts and appends to it; layered proxies each persist into their own log file and a restart replays both; and a persisted log is part of the segment manifest and snapshot. * collection, edge: recover persisted proxy changes on segment load Replay the pending changes logs left behind by proxy segments onto each segment when a shard loads its segments, right after consistency repair and before the payload index rebuild, vector name reconciliation and WAL replay. Proxy state that made it to disk no longer holds back the WAL acknowledge, so this is where it must be recovered from. Proxies are not reconstructed: the segment holder starts with plain segments carrying the replayed operations, and the logs are removed once the segment flushed them. * collection: test crash recovery through persisted proxy changes End-to-end test of the persisted pending changes: wrap every segment of a local shard in a proxy, delete points so the deletes are only buffered, flush, and assert the acknowledgeable version covers them. Then acknowledge the WAL up to that version, drop the shard without ever propagating the proxies, and load it again: the deletes are gone from the WAL and must come back through the pending changes logs. The delete under test is deliberately not the last WAL entry, as the acknowledge never passes the last entry and that one is always replayed. * segment: make replaying persisted proxy changes on load an explicit mode Add `PersistedProxyChanges` to state whether persisted pending proxy changes are replayed onto a segment when it is loaded. `Replay`, the default, recovers them and removes the logs as before. `Ignore` leaves both the segment and the log files untouched and logs at debug level that replaying was skipped; it is for segment files that mirror those of another writer, where replaying would make the local copy diverge from what the writer's manifest describes. All callers pass `Replay` for now, no behavior change. * collection: do not replay persisted proxy changes on partial snapshot recovery Partial snapshots are recovered by read replicas in a read/write segregation setup. A read replica must not mutate its segments, so it cannot replay the persisted proxy segment changes on load and must ignore them instead: its segment files are a local copy of the writer's that must stay a faithful mirror of them, as later partial snapshots are diffed against what the writer's manifest describes. Replaying would mutate the segment files and remove the logs, making the copy diverge. Thread the replay mode through `LocalShard::load` as a dedicated `PersistedProxyChanges` argument, derived from the recovery type: `RecoveryType::Full` replays as before, `RecoveryType::Partial` ignores the persisted changes and leaves the logs in place. Regular shard loads replay. Extend the crash recovery test with an ignoring load first: the delete under test must not come back and the logs must survive, before a replaying load recovers it. * Persist wrapped segment before pending changes Prevents raising version of proxy segment too early * Fix comment * Fix crash window, only ready optimized segment after propagating changes The optimizer renamed a newly built segment into segments_path and wrote its version file before finish_optimization propagated the proxies' buffered changes into it. A crash in that window left the segment restart-loadable but stale, permanently losing or resurrecting points. Defer the version file save until finish_optimization has fully reconciled proxy changes into the segment, including the post-swap dedup pass, so it stays invisible to restart and snapshot recovery until then. SegmentBuilder::build() gains a `ready` flag; load_segment gains `ignore_missing_version` for the one caller reloading before that point. Incidentally also closes the crash-unsafe cancellation-orphan cleanup gap noted in #9217, since a cancelled build is discarded on restart the same way. * Force flush optimized segment, otherwise we may lose proxy changes * Don't force flush after replay, defer deleting log files until flush * Include persisted proxy changes log file in segment manifest * Add random ID to proxy log files, prevent instance conflicts * Rename proxy log file, always include level * Delete proxy log file on unproxy, defer until next flush cycle * Fix truncation * Reformat * Lock persisted segments behind runtime feature flag * Enable necessary feature flags in tests * Fix linters |
||
|
|
31e97b5afc |
Persist doc_len in the mutable text index (#10616)
* Measure and persist document length in the mutable text index BM25 length normalization needs the total token count per point, which nothing stored: `point_to_tokens_count` is the distinct count, and its meaning is fixed by the user-visible `values_count` filter. - `doc_len` is measured in `add_many`, the only place that still sees every token, and persisted in the stored record alongside the tokens. - It is a parameter on `index_str_tokens` and `MutableInvertedIndexBuilder::add`, never derived. Without phrase matching the stored tokens are sorted and deduplicated, and the index is rebuilt from those records on every segment open, so a derived length would degrade to the distinct-term count on restart. - Array boundary sentinels are discounted by inserted count, not by value: `tokenize_doc` does not strip that character from user text the way `tokenize_query` does, so a payload containing it has those tokens indexed and they must be counted. - `MutableInvertedIndex` gains `point_to_doc_len` and a running `total_tokens`, maintained across add, overwrite and remove, so `avgdl` is a division rather than a scan. `set_doc_len` is the only writer, so an absent length means the same thing on every path: the slot is zeroed, never left stale. - Recording is gated behind `TextIndexParams::scoring()`, a private const for now. Nothing can ask for a ranked query yet, so recording on every text index would rebuild every collection to produce data no query can reach. State lives in the data rather than in a flag: `point_to_doc_len` is an `Option` the way `point_to_doc` already encodes positions, and `add_many` asks the index whether it records lengths rather than asking the config. - `StoredDocument::doc_len` is an `Option` skipped on write when absent, so a non-scoring index writes byte-identical records to today's and a legacy record reads back as `None`. A document whose tokens were all filtered is a real `Some(0)`, and stays distinguishable from one that was never measured. The immutable and on-disk backends drop the value for now and grow their own sidecar next. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Address text index document length review feedback --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-authored-by: generall <andrey@vasnetsov.com> |
||
|
|
f5e75477a0 |
Cleanup ConsensusStateMachine validation and docs (#10658)
|
||
|
|
0d2da625c4 |
Add cancellable reads for read-only edge shards (#10646)
* [AI] Add cancellable read trait for read-only edge shards * Move edge read cancellation tests into separate module |
||
|
|
f8c30a4a57 |
ci: terminate hung nextest tests after 5m and breadcrumb model_testing (#10644)
A hung harness_no_restarts run blocked ubuntu CI for ~50m with only SLOW markers and no failure dump. Kill tests after five slow-timeout periods, and print stage/op breadcrumbs to stdout so the timeout failure includes seed, storage path, and the last op/stage that never returned. |
||
|
|
39a9c6ae97 | [AI] Introduce QueryBatchRequest for edge batch queries (#10641) | ||
|
|
8f56a945f1 |
Add shared DiskCache statistics and latency histogram (#10637)
* [AI] Add shared disk cache statistics and latency histogram * [AI] Document disk cache statistics observer identity * [AI] Simplify disk cache statistics by removing pipeline error counters * [AI] Limit disk cache statistics to remote fetches and trim redundant tests * [AI] Close remote append handle before reload statistics snapshot |
||
|
|
1d0c2c1bb2 |
test: wait for consensus catch-up before snapshot recovery on a new peer (#10642)
test_recover_from_snapshot_2 and test_upload_snapshot_2 start snapshot recovery on a freshly joined peer as soon as it lists the collection. The collection appears once the creation entry is applied, while the peer is still replaying the rest of the raft log, including the removal of the killed peer. Recovery then decides which other replicas to remove or mark dead from that stale local view and drops a healthy replica, leaving a shard with a single replica. Add a helper that waits until all peers share the same commit index and have no pending operations, and use it in both tests before recovering. Also fix a misleading comment in the recovery replica cleanup branch. Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> |
||
|
|
0ff1c4e2b7 |
Send the queue proxy batch as a pre-encoded gRPC body (#10617)
* Send the queue proxy batch as a pre-encoded gRPC body The parent commits moved the WAL read, the operation clone and the request build off the async runtime. Two passes over the batch were left on it. Measured cost of each synchronous pass over a 26.2 MiB send batch (800 ops x 30 points x 256 dims), release build: pass 1 clone WAL operations 72.4 ms moved by the parent commits pass 2 build gRPC request 22.2 ms moved by the parent commits pass 3 clone request to send 65.7 ms on the async runtime pass 4 protobuf encode (tonic) 41.4 ms on the async runtime Pass 3 is there because `with_points_client` takes `impl Fn` and the channel pool calls that closure once per attempt, so each attempt needs its own owned message. Pass 4 runs inside `poll_next`: for a unary call tonic encodes the whole message in a single synchronous `encode_item`, so a worker is blocked for the full 41 ms, once per attempt. Encoding the batch up front removes both. The generated client cannot take a pre-encoded body, `update_batch` is typed `impl IntoRequest<UpdateBatchInternal>`, but all it does is pick a codec and a path and call `Grpc::unary`, and we already build the client ourselves from a pooled channel. `update_batch_pre_encoded` does the same three things with a codec that writes the encoded bytes through and decodes the response with prost. What stays on the runtime is the copy of the encoded body into tonic's send buffer: 14.9 ms for 26.2 MiB under jemalloc, nearly all of it faulting in freshly mapped pages rather than the copy itself (0.9 ms when the allocator hands back warm pages). Retries share the same refcounted bytes instead of cloning and re-encoding, so they drop with it. on the runtime before 107.2 ms per attempt on the runtime after 14.9 ms per attempt Bypassing the generated client means the RPC path and message types no longer follow the proto automatically, so a test checks them against the compiled descriptor set. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01WZjWYpKeYdGpKiLeEZU9oc * Hand the pre-encoded update batch a configured Grpc, take the service name from the generated code `update_batch_pre_encoded` took a bare channel plus a `max_decoding_message_size` argument, and its only caller passed `usize::MAX`. Every other internal client applies that limit inside its `with_*_client` helper, so do the same: `with_grpc` hands out the `tonic::client::Grpc` the generated clients wrap, already configured, and the argument goes away. The service half of the RPC identity now comes from the generated `points_internal_server::SERVICE_NAME` instead of a second literal. Only the method name and the path literal remain hand-written, still pinned to the descriptor set by the test. `PreEncodedMessage::encode` uses `encode_to_vec`: one pass instead of a separate `encoded_len` call, and no `expect`. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> * Adapt pre-encoding to the build/forward split from #10599 `forward_update_batch` takes `Arc<UpdateBatchInternal>` and encodes it once on the blocking pool, so the channel pool's attempts share the bytes. The queue proxy keeps the built request in that `Arc` across `BATCH_RETRIES` and for the per-operation isolation path. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> --------- Co-authored-by: generall <andrey@qdrant.com> Co-authored-by: Claude Opus 5 <noreply@anthropic.com> |
||
|
|
a6b3c7b4ca |
Share the text index post-tokenization indexing path (#10614)
* Share the text index post-tokenization indexing path Extract `MutableInvertedIndex::index_str_tokens` so the write path and read-only live reload cannot drift apart, and collapse the two identical on-disk file listings into one. The extracted helper gates the ordered document on `point_to_doc.is_some()` rather than re-reading `config.phrase_matching`. Equivalent, since `point_to_doc` is built from that flag, and it skips a needless clone if the two ever disagree. No behavior change. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Cover the text index live reload path `ReadOnlyAppendableFullTextIndex::live_reload` had no direct test: it replays stored documents through the same post-tokenization indexing as the write path, but nothing pinned that down. Asserts the incremental reload lands on the same state as a fresh `open_appendable` after a writer deletes one point and appends two, over both `phrase_matching` values. The phrase leg matters because adjacency depends on the ordered document being indexed, not just the token set. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Apply suggestion from @timvisee Co-authored-by: Tim Visée <tim+github@visee.me> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-authored-by: Tim Visée <tim+github@visee.me> |
||
|
|
f16b007daa |
fix(edge): tell manifest skew apart from a real fault when skipping a segment (#10627)
Every failure to open a segment was reported the same way: one warning, same wording, segment dropped, shard serves without it. Two very different things land there. The manifest is superset-biased, so it may list a segment the leader has not finalized yet or has already removed. Both arrive as `FileNotFound`, both fix themselves once the follower catches up, and both are routine. Anything else is a segment that should have loaded and did not. The shard opens without it and answers queries over a subset of its data, returning success to the client. In a recent load test this produced 667 warnings indistinguishable from ordinary leader churn. Report the first at debug and the second at error. The manifest does not need re-reading to tell them apart, so this costs nothing on the open path. Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
ddbc6cab0c |
fix(uio): carry the failing object in UniversalIoError::S3 (#10626)
`map_get_err` is handed the key it was reading, but only the `NotFound` arm kept it; every other error boxed the underlying failure and dropped the key on the floor. Callers that report such a failure are then unable to say what it was reading. The read-only segment open is the clearest case: it logs one warning per skipped segment, so a failure anywhere among a segment's objects — state file, id tracker, payload storage, per-vector storage and index, payload indexes — produces the same line, naming only the segment uuid. A recent load test hit exactly this: 667 warnings, all byte-identical, none of them saying which file failed. Give the `S3` variant an explicit `path` beside its source error, rather than folding the key into the message, so the object stays a field callers can read. It is optional because most construction sites are not about one particular object — a short or overlapping read from the scatter buffer, the append context's protocol errors — and those keep using `s3()` unchanged. `s3_at()` sets it, and the object-store read surface (every `map_get_err` caller, plus `list_files` and `exists`) now does. Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
e97bd7d1e5 |
fix(uio): don't fetch a zero-length object on a populating async open (#10625)
* fix(uio): don't fetch a zero-length object on a populating async open
The `PreferBackground`/`Blocking` arm of the disk cache's `open_async` built
its prefetch range by hand as `0..len`. For a zero-length object that is
`0..0`, which `object_store` rejects client-side with
`InvalidGetRange::Inconsistent` ("Range started at 0 and ended at 0") rather
than answering with an empty body.
The failure is not contained to the file: a single empty object anywhere
under a segment prefix fails the whole segment open, and a read-only
follower then logs "skipping unloadable segment" and serves the shard
without it — silently returning results computed over a subset of the data.
It is also deterministic, so the segment stays dropped on every retry.
The sync path already handles this correctly via `schedule_whole`
(`read_from_into_byte_buffer` disambiguates the unsatisfiable-range error
with a `len` call and yields an empty buffer); the async arm was the only
place assembling the range itself. Create the local mirror first and skip
the fetch entirely when there is nothing to read.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* ci: stop using the minio/mc image, which no longer exists
`minio/mc` has been withdrawn from Docker Hub — pulling it now fails with
"repository does not exist or may require 'docker login'". The readiness
loop ran it 60 times with output redirected to /dev/null, so the pull
error was invisible and the job failed as "rustfs did not become ready in
60s", pointing at the wrong component. Bucket creation used the same
image and would have failed next.
Use the AWS CLI that ships with the runner image instead: no third-party
container to pull for either step. The readiness probe stays an
authenticated call (`s3api list-buckets`), so it still waits out
credential setup and not just the port opening, and it now reports the
final failure instead of swallowing it.
Also pin the rustfs service image by digest. That was not the cause here
— rc.5 and rc.6 both work — but this workflow already pins its actions by
SHA, and a service tracking `latest` is how a green suite turns red with
no change to the repo.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
||
|
|
ca512c20fc |
Resolve filter ids once per filtered search (#10624)
External->internal id resolution is the expensive half of estimating a `has_id` filter, and the query API's rescore stage turns every prefetch into exactly such a filter. Two paths paid for it more than once: - `search_vectors_plain` re-estimated the filter its caller had already estimated to pick the plain strategy, so every segment resolved the whole id list twice. It now takes the estimation as an argument, and the dispatcher makes it once for every strategy. - Formula rescoring resolved the prefetched ids one at a time, which the disk-resident id tracker turns into an unpipelined block read per point, instead of the batched single pass it offers. Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
5e32ea89cb |
Reject snapshot upload without collection config, without exposing the temp path (#10556)
* Reject snapshot upload without collection config before loading it The raw IO error from `CollectionConfigInternal::load` embedded the server-side temporary path in the API response. Check for the file first and return a fixed bad-input error instead. Part of #10553 Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01A3o9eWSZNMa6WAs5F2HMZC * Fix missing-config check to run before restore_snapshot loads config The path-leak guard lived after Collection::restore_snapshot, but that function already calls CollectionConfigInternal::load and surfaced the temp path as a 500. Require a regular config.json file before loading, and cover a directory-shaped config entry in the openapi test. * Use a valid empty TAR in the missing-config snapshot upload test Avoid depending on malformed-archive handling; exercise the missing collection-config path with a real TAR that has no entries. --------- Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> Co-authored-by: qdrant-cloud-bot <111755117+qdrant-cloud-bot@users.noreply.github.com> |
||
|
|
17a0747d24 |
Move queue proxy WAL read and batch build off the async runtime (#10599)
* Move queue proxy WAL read off the async runtime `read_wal_batch()` read up to MAX_BATCH_BYTES (32 MiB) from the WAL and deserialized it synchronously on the async runtime. That runtime also serves all internal gRPC, including the health check the transport channel pool uses to decide whether a peer is alive, so during the queue-replay phase of a snapshot shard transfer the sender periodically stopped answering internal requests for the duration of a 32 MiB disk read plus decode. Measured on a 3-node cluster with no CPU/memory/IO limits and ~30% free RAM: the sender's health-check p99 went from 3-6 ms during the download phase to 81-93 ms during replay, across three separate transfers, while a peer probed at the same instant stayed flat at 3-5 ms. Take the lock with `lock_owned().await` before `spawn_blocking` rather than `blocking_lock()` inside it, so a blocking-pool thread is only ever occupied by the read and never by waiting for a concurrent writer. Lock scope is unchanged - the mutex was already held across the whole synchronous read. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01WZjWYpKeYdGpKiLeEZU9oc * remove unwanted comment * Build the queue proxy send batch off the async runtime too The previous commit moved the WAL read to the blocking pool, but measuring it showed no improvement: the sender's health-check p99 during queue replay stayed at ~90-100 ms. perf on the patched build explained why - optimizer/HNSW load is ambient (~90% of CPU in both the download and replay phases, so not what makes replay special), while one general-runtime worker burns 7.7% of all CPU during replay against ~0.5% during download. That worker is doing the *send* half of the loop, which is still synchronous: - transfer_operations_batch() clones every operation in the batch - forward_update_batch() converts each one into its gRPC representation Both are full passes over up to MAX_BATCH_BYTES (32 MiB) of point data, on the runtime that also answers internal health checks. Move both to the blocking pool. WalBatch now holds its operations behind an Arc so the batch can be shared into a blocking task without being copied first, and the gRPC request construction is split out of forward_update_batch into RemoteShard::build_update_batch_request so it can be spawned. The extracted function keeps the original body and indentation, so the diff is the move plus the plumbing rather than a reindent. forward_update_batch has exactly one caller (the queue proxy), so this adds no blocking-pool hop to the normal replication path. Still on the runtime and not addressed here: the per-attempt clone of the request inside with_points_client, and tonic's own protobuf encoding. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01WZjWYpKeYdGpKiLeEZU9oc * Use single clone, don't iterate manually * Methods are droppable, don't hang on spawn_blocking with full runtime * Build the WAL transfer batch once, reuse it across retries The batch was deep cloned on every send so the original operations stayed available for retries and for the one-by-one isolation path. Instead, move the operations into the gRPC request once and retry on the prebuilt request, which `with_points_client` already clones per attempt. Stripping WAL indices and setting the force flag now happens in the WAL read task, so no per-operation work runs on the async runtime. Drop the pre-1.14.1 fallback that transferred operations individually, all peers support batched updates by now. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> --------- Co-authored-by: generall <andrey@qdrant.com> Co-authored-by: Claude Opus 5 <noreply@anthropic.com> Co-authored-by: timvisee <tim@visee.me> |
||
|
|
95eb3c3f79 |
Support cached id tracker memory placement (#10598)
Thread a `Populate` through the disk-resident id tracker's open paths (`DiskMappingReader`, `DiskIdTracker::open`, `ReadOnlyDiskIdTracker`, `ReadOnlyIdTrackerEnum`) so a `cached` placement primes the page cache with the mapping files on load instead of leaving them to page in on demand. The populate is derived from the segment config's placement at load time, clamped by low-memory mode, in both the writable segment open and the read-only one. The update-only lookup path keeps its transfer-nothing policy, and the build-time open stays cold: the built segment is reloaded anyway. `cached` is no longer rejected by validation. Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com> |
||
|
|
81bb80a5d3 |
Expose id tracker memory placement in collection config (#10597)
Add `id_tracker: { memory: cold | pinned }` to CollectionParams,
CollectionParamsDiff and CreateCollection (REST + gRPC `IdTrackerParams`),
mirroring `payload: { memory }`. `cold` builds the disk-resident id tracker,
`pinned` the in-RAM immutable one. Unset keeps the current behavior: the
`serverless_compatible` feature flag decides.
The requested placement is persisted as an optional `id_tracker_memory` on
SegmentConfig (skipped when unset, so existing configs are unchanged); the
segment builder resolves it through `SegmentConfig::id_tracker_memory_placement`
instead of reading the feature flag directly.
The config mismatch optimizer rebuilds non-appendable segments whose effective
placement differs from the requested one. Appendable segments are skipped: they
always use the mutable tracker and get the current config when indexed.
`cached` is rejected by validation: the disk mapping reader has no
populate-on-open path.
Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
|