mirror of
https://github.com/qdrant/qdrant.git
synced 2026-09-27 16:37:42 -05:00
edge-docs
1014
Commits
| Author | SHA1 | Message | Date | |
|---|---|---|---|---|
|
|
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> |
||
|
|
242a8d3b1e |
Warn when distributed mode has API key but enforce_internal_auth is off (#10718)
* Warn when distributed mode has API key but enforce_internal_auth is off Make the insecure internal (p2p) gRPC configuration visible at startup so operators enable enforcement after rolling upgrades complete. * Shorten enforce_internal_auth warning message Drop the rolling-upgrade guidance from the log line. |
||
|
|
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 |
||
|
|
f5e75477a0 |
Cleanup ConsensusStateMachine validation and docs (#10658)
|
||
|
|
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> |
||
|
|
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>
|
||
|
|
b01f6d841e | Perf: Move BM25 sparse embedding on a blocking thread (#10610) | ||
|
|
b81ad12c20 |
Don't re-apply committed entries already applied during consensus start (#10277)
Co-authored-by: Roman Titov <ffuugoo@users.noreply.github.com> |
||
|
|
9224014e93 |
Optimizations and improvements for ConsensusStateMachine (#10533)
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
8185a2692e |
Validate consensus against ConsensusStateMachine (#10469)
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
9687da6c41 |
Use more direct calls (#10347)
* Direct call for shard transfer method and keys * Reuse cardinality estimate in sparse plain search * Avoid recounting available points in segment size info * Avoid cloning segment config when updating quantization * Avoid cloning search request for load profile * Direct call for counting read-only segments * Avoid re-reading point range for values count * Direct call to check replica states when initializing collection * Direct call to look up transfer on restart * Direct call for shard replicas after snapshot recovery * Direct call for local replica states in health check * Direct call for payload index schema keys when applying state * Direct calls for sharding method and key mapping when creating shard key * Direct call to check if peer has shards * Direct call for sharding method and keys when dropping shard key * Avoid cloning collection params for group by ordering * Avoid cloning collection params in local shard search * Direct call for peer address when sending Raft messages * Direct call for peer address in who_is * Avoid cloning remote query batch request * Avoid cloning operation in queue proxy update * Avoid cloning gRPC search groups request * Fetch cluster status once in cluster telemetry * Direct call to validate transfer exists on finish * Direct call for sharding method when dropping shard key * Avoid cloning peer address map when listing peers * Avoid cloning peer address map when adding peer to known * Avoid cloning shard key mapping when routing writes with fallback * Avoid cloning shard key mapping when checking resharding start * Avoid cloning gRPC recommend groups request * Avoid cloning operation when retaining forwarded point IDs * Direct call for counting collections in telemetry * Direct call to validate transfer exists on recovery * Direct call for shard IDs by shard key * Direct call for shard keys * Direct call to check if peer has shards in consensus * Direct call for replica state on transfer recovery * Direct call to check for active replicas when routing writes with fallback * Direct call to validate transfer exists on abort |
||
|
|
e0110f3fe8 |
feat: optional dial9 Tokio telemetry behind a dial9 feature (#10442)
* Add optional dial9 Tokio telemetry behind a `dial9` feature Integrate dial9 so storage runtimes can emit production-friendly Tokio traces. Recording is off unless the crate is built with `--features dial9` and DIAL9_ENABLED=true is set at runtime; with the feature off, runtime construction is byte-for-byte unchanged. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Y7X6MkjY3P7wpP2MfHdTDY * Enable dial9 CPU and schedule profiling Turn on cpu-profiling and sched events behind the same `dial9` feature, add the DIAL9_CPU_* / DIAL9_SCHEDULE_* env knobs, and document the frame pointer rustflags the stack unwinder needs. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Y7X6MkjY3P7wpP2MfHdTDY * Harden dial9 env parsing and the writer-failure path - Reset Cargo.lock to the branch point and re-resolve, so the diff is additive instead of re-resolving unrelated packages. This drops the heck 0.5.0 -> 0.4.1 downgrade, which sat in the default build graph and would have changed proto codegen identifier casing. The remaining non-additive entry, toml_parser 1.0.9 -> 1.1.3, is forced by proc-macro-crate via dial9-trace-format-derive. - Parse DIAL9_* booleans the way dial9 does, accepting 1/y/yes/on and 0/n/no/off and warning on anything else. `str::parse::<bool>` took only exact lowercase true/false, so DIAL9_CPU_PROFILE_ENABLED=0 silently left 99 Hz sampling on and DIAL9_ENABLED=1 silently left recording off. - Require the numeric knobs to be positive. A zero disk budget made dial9 evict everything and stop recording within seconds while the log still reported telemetry enabled. - Treat a set-but-empty DIAL9_TRACE_DIR as unset. It skipped the /tmp fallback and wrote up to the full budget into the working directory, which is /qdrant next to storage/ in the official image. - Return a disabled guard as soon as the trace writer fails, before with_cpu_profiling and with_sched_events run. Those start their profilers eagerly, opening a perf event per thread and installing a process-global signal handler that build() would then discard. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Y7X6MkjY3P7wpP2MfHdTDY * Correct the dial9 docs and give them their own section - `--cfg tokio_unstable` is required for any task data at all, not merely for fuller coverage: dial9's poll, spawn and terminate hooks are all `#[cfg(tokio_unstable)]`, and nothing in the repo sets the flag. Without it there is no task timeline and DIAL9_TASK_TRACKING_ENABLED does nothing. - Document `-C debuginfo=2`. `[profile.perf]` inherits `release` and sets no `debug` key, so the documented build symbolized off the ELF symtab with inlined callees collapsed and no file or line, unlike `[profile.bench]` which sets `debug = true` for this reason. - Move the dial9 material out from between the feature list and the prose that belongs to it. Those paragraphs describe `tracing` instrumentation and read as dial9's when the example is wedged in front of them, which points readers at `#[tracing::instrument]` for a tool that records Tokio runtime events and no tracing spans. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Y7X6MkjY3P7wpP2MfHdTDY * Use cfg_select! --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-authored-by: timvisee <tim@visee.me> |
||
|
|
b1a1c00059 |
Fix proxied changes dropped when an optimization fails (#10364)
* Add SetFlushInterval op to the model tester Changes the collection's flush_interval_sec mid-run through the same path update_collection takes (persist the optimizer-config diff, then recreate the optimizers in the background). The model is untouched: what it perturbs is the flush cadence, so how much of the workload is still WAL-only when a restart hits, plus the worker stop/start race in on_optimizer_config_update. Kept in FORCE_OFF for now: with the optimizer on it makes stale point state visible within a few ops of the config change. Narrowed to recreate_optimizers_background, see the comment on Swarm::BASE. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011AjmS5GFeGztfP3JqutnXj * Keep SetFlushInterval enabled in the swarm Drops it from FORCE_OFF so the divergence it surfaces is reachable without --enable-force-off (which would also enable the broken vector-name ops). The evidence moves from the FORCE_OFF comment onto the op's own doc. The two optimizer-on harness gates now fail whenever the swarm draws the op. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011AjmS5GFeGztfP3JqutnXj * Propagate proxied changes when unwrapping proxies on optimization failure unwrap_proxy puts the wrapped segments back into the segment holder, so the changes recorded on the proxy while the optimization ran (deleted points, index and vector-name changes) have to reach the wrapped segment first. They did not, so every point deleted or overwritten during the optimization kept its pre-optimization copy live next to the new copy in the write segment, and reads saw both: counts too high, scroll and search returning the stale copy. The snapshot unproxy path already does this; the optimizer failure path was the only place putting a wrapped segment back without it. It is reachable whenever the shard outlives the cancellation, in particular an update_collection that recreates the optimizers while an optimization is in flight. Lock order is holder-then-updates, matching try_unproxy_segment: updates-then- holder-write deadlocks against the snapshot path. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011AjmS5GFeGztfP3JqutnXj * Drop the stale failure note from the SetFlushInterval doc The divergence it described is fixed in this branch. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_011AjmS5GFeGztfP3JqutnXj * test as well with 0s as flushing interval * Fail optimization unwrapping when proxy propagation fails Losing proxied deletes and index changes is data corruption, so return the error instead of logging it: no proxy is unwrapped and the changes stay served by the proxies. The cancelled-segment cleanup moves ahead of unwrap_proxy so the orphan is still removed when that error fires. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01CpoaAtGbAQAEuxEi8ScHc5 * Drop the model tester --flush-interval-sec flag SetFlushInterval covers the interval now, so the run starts at the shipped 5s default (fixture::INITIAL_FLUSH_INTERVAL_SEC, still traced in the header) and the ops move it from there. Also documents what 0 does now that it is a generated value. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01CpoaAtGbAQAEuxEi8ScHc5 * Fail snapshot unproxying when proxy propagation fails Both paths logged the error and unwrapped anyway, dropping the deletes and index changes that never reached the wrapped segment. Same reasoning as unwrap_proxy in the optimizer. try_unproxy_segment hands the lock back and leaves the proxy installed, the failure mode its doc already describes: the caller keeps it in `proxies` and unproxy_all_segments retries the propagation right after. unproxy_all_segments returns before touching the holder, so the temp segment the surviving proxies write into stays in place (remove_segment_if_not_needed only checks whether it is empty and appendable, not whether a proxy still references it). Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01CpoaAtGbAQAEuxEi8ScHc5 --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> |
||
|
|
55d59b0501 |
fix(grpc): raise http2 max pending accept reset streams on public API (#10306)
Mirror the internal gRPC server setting to avoid GOAWAY/ENHANCE_YOUR_CALM errors when clients multiplex many short-lived streams (e.g. coach high_concurrency drill). See #1907. |
||
|
|
ce2f31c773 |
Detect container runtime beyond Docker in telemetry (#10132)
* telemetry: detect container runtime, not just Docker `is_docker` only checks `/.dockerenv`, a marker the Docker daemon creates. containerd/CRI-O (Kubernetes), Podman, etc. don't, so a containerized node — notably every Qdrant Cloud pod — reported `is_docker: false`, indistinguishable from bare metal. Add a `container_runtime` field (enum: none/docker/kubernetes/other) detected from well-known markers, most specific first: KUBERNETES_SERVICE_HOST → /.dockerenv → other container markers → none. Only docker and kubernetes are enumerated (the modes that matter for Qdrant); the rest fold into `other`. `is_docker` is kept but derived (== docker). Detected once via LazyLock. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Regenerate OpenAPI spec Add the ContainerRuntime schema and container_runtime field. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> |
||
|
|
896deaeeef |
Fix clippy warnings from Rust 1.98 beta (#10265)
* Fix clippy warnings from Rust 1.98 beta
* drop a redundant trait import in an io_bridge test module, it already
arrives through `use super::*`
* rewrite two `chunks_exact(CONST)` sites as `as_chunks::<{ CONST }>()`
for the new `chunks_exact_to_as_chunks` lint
* return `bool` from `wait_for_consensus_commit` instead of
`Result<(), ()>`, which `result_unit_err` now flags on `async fn`. Its
only caller did `.is_ok()` on it
* allow `result_large_err` on `QueueProxyShard::new_from_version`, which
hands the `LocalShard` back to the caller on failure. Mirrors the allow
already on `ForwardProxyShard::new`
* migrate three `Atomic::fetch_update` calls to `try_update`, the name it
is renamed to in 1.99. The new name already exists at our 1.97 MSRV
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* keep guarantee on caller
---------
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
|
||
|
|
7b2dc8819b |
Expose cpu_cores_used in Prometheus /metrics (#10243)
Surface the existing process CPU usage from telemetry as a gauge so operators can scrape average cores used without polling /telemetry. |
||
|
|
00b17c3e74 |
Replace sys-info dependency with already-present sysinfo (#10216)
The app telemetry was the only user of the sys-info crate, while segment already depends on sysinfo for cgroup-aware memory accounting. Read the distribution id/version via sysinfo statics, and the disk size fallback via common::disk_usage, so the whole sys-info crate (and its bundled C sources) drops out of the build. sysinfo is hoisted to a workspace dependency, shared by the root crate and segment. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
06ffcb881f |
Add CachedBlobFile: cached reads + write-through appends for object stores (#10206)
* Add CachedBlobFile: cached reads + write-through appends for object stores Combine a DiskCache mirror (reads) with a BlobFile remote handle (appends) into CachedBlobFile/CachedBlobFs, the appendable universal-IO citizen for object stores. Appends perform the remote mutation inline and are durable at Ok: a native write-offset append in AppendMode::Native (with a soft limit on appends per object), or a whole-object rewrite in AppendMode::Rewrite for stores without native append. After a successful append the mirror length is advanced without extra IO; appended blocks fault in from the remote on first read. The multipart UploadPartCopy rewrite path (prefix >= 5 MiB) and the rewrite-required error classification are left as todo!() pending the AsyncRewrite backend capability. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Backend-advertised AppendMethod; reactive appended-block cap recovery Replace CachedBlobFile's stored AppendMode with AsyncAppend::supported_append: the backend advertises Native or PartialUpload, and append takes a matching AppendRequest variant, rejecting the ones it does not support. The multipart UploadPartCopy todo moves into the S3 backend's PartialUpload arm. Drop the native_appends soft-limit counter: it is per-handle in-memory state that resets on every restart, so it can never be the correctness mechanism and persisting it would not make it authoritative either. The store is the authority: hitting its appended-block cap now surfaces as the new UniversalIoError::AppendRewriteRequired (S3 400 TooManyParts), and CachedBlobFile recovers with a whole-object rewrite. Unrecognized errors stay hard errors instead of silently triggering rewrites. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Per-store append strategies; server-side rewrites for plain S3 and GCS Replace the single AppendContext struct with an enum of strategy objects, one per store capability, each owning its append logic: - NativeAppend: the signed write-offset PutObject (S3 Express, MinIO AiStor; AwsConfig::native_append declares it for AiStor-like endpoints, s3_express implies it). - PartCopyAppend: plain S3 — appends land as one atomic multipart rewrite whose prefix parts are server-side UploadPartCopy requests; nothing but the appended data crosses the network. object_store keeps such provider-specific calls out of its portable surface, so the requests are hand-signed like the native append. - ComposeAppend: GCS — the appended data is uploaded as a temporary neighbor object and composed onto the destination server-side, conditional on the observed generation (a real compare-and-swap). AppendMethod is replaced by AppendSupport, which tells the caller the only thing it needs: when the store takes a direct append. Always (native, and compose: no part minimums, no block cap), AboveThreshold (part-copy: the copied prefix lands as non-last multipart parts, >= 5 MiB each), or Never. CachedBlobFile drops its hardcoded MIN_COPY_PREFIX and rewrites locally only below the backend-advertised threshold; AppendRequest::Rewrite now means only "append and rebuild as a single blob" — the appended-block cap recovery. The append module is split one file per strategy, with a shared SignedRequestContext transport and a test-only HTTP stub. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * DiskCache tracks the remote object's etag Seeded from the new known_etag open extra (OpenExtra::with_known_etag), refreshed from FileInfo on schedule_reopen, and settable directly for callers that mutate the remote out of band. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Remove AppendRequest enum; appended-block cap recovery moves into the backend AsyncAppend::append takes plain (path, offset, data). A native S3 store that rejects an append with TooManyParts now falls back to the part-copy rewrite inside the dispatcher, instead of surfacing AppendRewriteRequired to CachedBlobFile for a second Rewrite request. The Rewrite variant was handled identically to Append everywhere except that one native path. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Escalate to download+rewrite when the store rejects a part-copy rewrite The cap-recovery rewrite is chosen by the store's returned error, not a client-side threshold: a part-copy attempt rejected with EntityTooSmall (typed as UniversalIoError::AppendEntityTooSmall, parsed from the S3 error <Code>) falls back to downloading the sub-part-minimum prefix and PUTting the whole object back, guarded by a prefix-length offset check. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Fix S3 Express appends: zonal endpoint + s3express SigV4 service Hand-issued appends targeted the standard endpoint and signed as "s3", so every append to a directory bucket got 404 NoSuchBucket, masked as AppendOffsetConflict by the 404 mapping. Derive the zonal {bucket}.s3express-{az}.{region} base from the mandatory --{az}--x-s3 bucket suffix (mirroring object_store's private derivation), carry the SigV4 service name in SignedRequestContext, and treat a 404 as a conflict only for NoSuchKey or bodiless responses — NoSuchBucket stays a loud error guarding the endpoint derivation. extract_xml_tag moves up to the context module and now tolerates tag attributes and pretty-printed bodies. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Server-side etag precondition on appends; BlobFile loses UniversalAppend AsyncAppend::append carries an expected_etag that S3 part-copy rewrites attach as x-amz-copy-source-if-match (412 -> AppendEtagMismatch, a new typed error) and download_rewrite checks against the GET's own etag; native write-offset PUTs and GCS compose ignore it. BlobFile appends only through the inherent etag-aware append_bytes now — CachedBlobFile calls it directly with its DiskCache-tracked etag — and BlobFs's mutating ops become inherent, delegated from CachedBlobFs, per the standing TODOs. The append conformance battery runs over the CachedBlobFs stack, via new direct constructors that share one backend. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Drop unfulfilled too_many_arguments expectation rewrite_parts has exactly seven parameters — at the clippy threshold, not over it — so the lint never fires and the expect fails CI under -D unfulfilled-lint-expectations. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
500eed65b1 | add etag to FileInfo (#10190) | ||
|
|
8f5e83b13d |
[CachedFs] Unchanged file open is no-op (#10050)
* add `UniversalIoError::UnchangedOpen` * add & impl `CachedReadFs::reschedule_prefetch` * clippy * add `OkUnchanged` helper * propagate scheduling errors immediately * drop lock before reacquiring it * clear prefetched_files on new snapshot * only avoid prefetch on full FileInfo equality * `UnchangedOpen` maps to `Cancelled` * match-all match |
||
|
|
0805d3f422 |
Add /profiler/consensus_lag to measure apply lag between peers (#10090)
* Add /profiler/consensus_lag to measure apply lag between peers Raft commit index advances on a peer whose apply loop is stalled, so the existing signals - `raft_info.commit` and the `all_nodes_have_same_commit` test helper - report a stuck peer as healthy. Nothing exposes how long a peer has been behind at *applying* entries, which is what shard transfer's `await_consensus_sync` barrier actually waits on. Each peer now keeps a ring of the last 32 entries it applied, stamped with its own wall clock and the time that entry took to apply. The ring is in memory on ConsensusManager, not in Persistent, so the on-disk format is untouched. `/profiler/consensus_lag` collects those rings from every peer over a new internal RPC and lines them up on the entry indices they share. Each entry is measured from whichever peer applied it first, so a lag is never negative; the peer that is first can differ per entry, so the baseline is per entry rather than a single chosen peer. Entries only one peer still remembers are excluded, otherwise a peer would be measured against itself. A peer stalled part-way through an entry keeps healthy lag statistics - everything it did apply, it applied on time - so the report carries `behind_entries` and `newest_applied_age_ms` alongside, which is what actually exposes the stall. Peers that fail or time out are listed rather than failing the request: a partial answer is more useful than none when the point is to find a peer that stopped answering. The endpoint follows `/profiler/slow_requests`: manage access, and outside OpenAPI, so no endpoint-count or ACTION_ACCESS guard applies. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Move applied-entry log into its own module Keeps the new code out of files that are already large. The ring, its entry type and the snapshot served over RPC move to `content_manager/consensus/applied_log.rs`, alongside the other consensus internals; `ConsensusManager` is left with a field, an accessor and the one `record` call in the apply loop. The grpc encoding moves next to the decoding it mirrors, in `common/consensus_lag.rs`, leaving the internal service handler three lines instead of thirty. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Test that a consensus stall is still in the report after the peer catches up * Take each peer's applied index from consensus state, not its ring --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-authored-by: tellet-q <elena.dubrovina@qdrant.com> |
||
|
|
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>
|
||
|
|
c07d57bd8f | chore: fix dead code lints on macos (#10045) | ||
|
|
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> |
||
|
|
0a16a62f99 |
feat: io_uring setting to control which components use the io_uring backend (#10008)
* feat: `io_uring` setting to control which components use the io_uring backend A few components have both an mmap and an io_uring variant reading the very same files: the immutable dense vector storages, the single-file TurboQuant storage, and the mmap payload storage. Until now the choice was a side effect of `async_scorer` — a vector-search knob — plus, for the payload storage, a feature flag that was parked off because io_uring is ~2x slower than mmap when the data fits the page cache (#9310, #9409). Add `storage.performance.io_uring`, optional, with two modes: - unset (default): unchanged behaviour. The vector storages keep following `async_scorer`; the payload storage stays on mmap. - `disabled`: no component uses io_uring. - `auto`: a component uses io_uring when its memory placement is `cold` (data is left on disk, so reads hit the disk and there is something to gain), its feature flag allows it, and the kernel supports io_uring. Components meant to sit in RAM keep using mmap. The decision lives in one place, `segment::common::io_uring::use_io_uring`, so the openers no longer each reach for the async-scorer global. Kernel support is now probed up front through `is_io_uring_supported()` instead of opening a file and falling back on error. `async_payload_storage` now defaults to on: it no longer decides anything by itself, it only lifts the ban, and the payload storage no longer follows `async_scorer` at all — so turning it on cannot silently move an existing `async_scorer: true` deployment onto the slower path. Which backend a component ended up on depends on the config, the placement and the kernel at once, so report it in `SegmentInfo`: `vector_data[name].io_backend` and `payload_storage_io_backend`, both `"mmap" | "io_uring"`, absent for components that have no such choice. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Trim comments, drop trivial tests Two tests were only restating their own implementation: `test_mode_round_trip` round-tripped the encode/decode pair next to it, and `test_io_uring_config` checked that serde deserializes a two-variant enum. The mode matrix test stays, it is the one that pins the semantics. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Flatten `IoBackend` in OpenAPI, derive `JsonSchema` for `IoUringMode` Per-variant doc comments on a plain string enum make schemars emit a `oneOf` of anonymous single-value objects instead of a flat `enum`. Move the variant descriptions into the enum doc, as `Memory` and friends already do. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * Update lib/segment/src/vector_storage/turbo/turbo_vector_storage.rs Co-authored-by: Roman Titov <ffuugoo@users.noreply.github.com> * Update lib/segment/src/types.rs Co-authored-by: Roman Titov <ffuugoo@users.noreply.github.com> * Update lib/segment/src/types.rs Co-authored-by: Roman Titov <ffuugoo@users.noreply.github.com> * Update lib/segment/src/types.rs Co-authored-by: Roman Titov <ffuugoo@users.noreply.github.com> * Update lib/segment/src/types.rs Co-authored-by: Roman Titov <ffuugoo@users.noreply.github.com> * upd openapi schema * Update lib/common/common/src/flags.rs Co-authored-by: Roman Titov <ffuugoo@users.noreply.github.com> * Require kernel io_uring support in the async-scorer fallback `use_io_uring` returned `get_async_scorer()` verbatim when the `io_uring` setting is unset, so an enabled async scorer on a kernel without io_uring opened the io_uring storage, failed, and fell back to mmap with an error log per segment. Gate that branch on `is_io_uring_supported()` too, like `Auto` already is, so the component just stays on mmap. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * upd openapi schema --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com> Co-authored-by: Roman Titov <ffuugoo@users.noreply.github.com> |
||
|
|
8a7325ad5d |
Remove deprecated search endpoints from OpenAPI, deprecate them in gRPC (#9982)
* Remove deprecated search/recommend/discover endpoints from OpenAPI Remove deprecated REST API endpoint definitions from the OpenAPI generator. These endpoints were deprecated in v1.13.3 (`f4ced2567`, #5907, 2025-01-30) in favor of the universal `/points/query` endpoint: - POST /points/search - POST /points/search/batch - POST /points/search/groups - POST /points/recommend - POST /points/recommend/batch - POST /points/recommend/groups - POST /points/discover - POST /points/discover/batch Also removes the corresponding request types from the schema generator and updates the expected API count in the consistency check. Co-authored-by: Cursor <cursoragent@cursor.com> * Migrate OpenAPI integration tests to /points/query The deprecated /points/search, /points/recommend and /points/discover endpoints (along with their /batch and /groups variants) were removed from the OpenAPI spec, which caused validation failures in the Python integration test harness. This commit migrates the affected tests to the universal /points/query endpoint: - Delete tests dedicated to the deprecated endpoints: test_recommend.py, test_discover.py, test_multicollection_reco.py, test_recommendation_multivector.py - Refactor remaining tests to call /points/query (and /query/batch, /query/groups), translating request bodies (vector -> query / using, positive/negative -> query.recommend, target/context -> query.discover) and unwrapping the new result.points response shape. - Drop equivalence assertions against the now-removed legacy endpoints. Co-authored-by: Cursor <cursoragent@cursor.com> * Relax non-empty assertions in migrated recommend/discover tests The previous migration added `len(...) > 0` assertions to tests that previously only checked equivalence between the deprecated and new API. These assertions are too strict because the parametrized `query_filter` cases legitimately produce empty result sets. Drop the `> 0` assertion and rely on `request_with_validation` to verify the response is well-formed and HTTP OK. Co-authored-by: Cursor <cursoragent@cursor.com> * Migrate remaining OpenAPI tests off deprecated search endpoints Tests added to dev after the original migration was written still call /points/search and /points/recommend/groups through `request_with_validation`, which resolves the endpoint against the OpenAPI spec and therefore breaks once the endpoint is not in the spec: - test_turbo4_storage.py, test_sparse_idf_corpus.py, test_validation.py: translate /points/search to /points/query (vector{name,vector} -> query + using, result -> result.points). - test_group.py: drop the /points/recommend/groups half of the lookup_from validation test in favour of the query equivalent. test_sparse_idf_corpus.py's test_query_api_supports_idf_corpus goes away: with the helper on /points/query every test in the file now exercises what it asserted. Also record why test_recommend_group cannot assert on its groups: it uses every point in the collection as a recommend example, so all of them are excluded and the result is legitimately empty. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * Regenerate openapi.json without the deprecated search endpoints Drops the 8 deprecated paths and the request schemas that only they referenced: Search/Recommend/Discover request (+Batch, +Groups) types and their exclusive dependencies (NamedVector, NamedSparseVector, NamedVectorStruct, UsingVector, RecommendExample, ContextExamplePair). Regenerated output is a strict subset of the previous spec, and every remaining $ref still resolves. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * Deprecate the search/recommend/discover RPCs in gRPC The REST counterparts have carried `deprecated: true` since v1.13.3 and are now gone from the OpenAPI spec, while the gRPC RPCs never got any deprecation annotation at all. Mark all 8 with `option deprecated = true` so generated clients warn, and point each doc comment at its `Query` replacement. tonic puts `#[deprecated]` on the generated client methods only; the server trait gets the doc comment alone, so our own `impl` is unaffected. The RPCs keep serving traffic — this is annotation only. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> * Restore the deleted recommend/discover suites on /points/query The earlier migration deleted these four files outright, but the query-side tests it left behind are all shallow smoke tests (`len(result) > 0`, `"points" in result[0]`). The deleted ones carried invariants with no query-API equivalent anywhere, so deleting them was a real loss of coverage rather than de-duplication: - test_recommend.py: default strategy equals average_vector; batch results identical to sequential singles across six request shapes; best_score with only negatives yields all-negative scores; best_score with a single positive orders identically to a nearest query; raw vectors as examples equal ids as examples. - test_discover.py: context-only scores are all <= 0; target-only orders identically to a nearest query but scores differently; with a fixed context the integer part of the score is stable while the decimal part moves, and vice versa with a fixed target; batch equals singles; lookup_from by id equals by vector. - test_multicollection_reco.py: cross-collection lookup_from, plus wrong-vector-size, unknown-collection and unknown-vector rejections. - test_recommendation_multivector.py: the same recommend invariants over a max_sim multivector collection, which the query suite never covered. Only test_recommend_missing_lookup_from_collection_with_raw_vector is dropped as genuinely redundant — test_query.py's test_query_missing_lookup_from_collection covers query, query/batch and prefetch. Two request-shape differences the translation had to absorb: - Giving no examples at all is 422 (a RecommendInput validation rule), where the legacy API reported 400 from the query itself. A malformed example, such as an empty vector, is still 400. - DiscoverInput requires the `context` key and accepts only an explicit null to mean "no context", so target-only discover must spell it out. The legacy API let it be omitted. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> --------- Co-authored-by: Cursor <cursoragent@cursor.com> Co-authored-by: Claude Opus 5 <noreply@anthropic.com> |
||
|
|
94fdd0e746 |
Support memory placement in service-level storage config defaults (#9950)
* Support memory placement in service-level storage config defaults Follow-up to #9684: `storage.payload.memory` and `storage.collection.vectors.memory` set service-wide placement defaults for newly created collections, deprecating `storage.on_disk_payload` and `storage.collection.vectors.on_disk`. Defaults resolve as: request `memory` > request legacy flag > service `memory` > service legacy flag; exactly one level is filled to avoid spurious memory-vs-legacy mismatch warnings. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Mark hnsw_index.on_disk deprecated in config.yaml, document memory option Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
48e710ce0b |
Cleanup UniversalRead methods interface (#9934)
* use common::generic_consts::{Random, Sequential};
* UniversalRead::read_batch: generic over E
* UniversalRead::read_batch: pass `AccessPattern` as ZST arg
* UniversalRead::read_bytes_iter: pass `AccessPattern` as ZST arg
* UniversalRead::read_iter: pass `AccessPattern` as ZST arg
* UniversalRead::read: pass `AccessPattern` as ZST arg
* UniversalRead::read_bytes: pass `AccessPattern` as ZST arg
|
||
|
|
7e99cdd86a | UioResult (#9933) | ||
|
|
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> |
||
|
|
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 |
||
|
|
71e6c70458 |
Universal IO: append to file (#9720)
* universal_io: add UniversalAppend for atomic single-operation appends Growing a file previously took a separate set_len + reopen + write dance (bypassing universal_io and leaving a zero-filled window on crash), and there was no way to express appends for backends without random-offset writes. UniversalAppend::append grows the file by writing at the current end of file in one atomic grow+write operation and returns the offset at which the data landed; append_batch lands multiple buffers contiguously in as few operations as the backend allows. flusher() moves from UniversalWrite into a new UniversalFlush supertrait so append-only handles can require it without duplicating the method. Local backends, both single-syscall: - MmapFile appends through a dedicated O_APPEND fd (every write(2) / writev(2) is an atomic grow+write at EOF), then remaps via reopen(). Its flusher also fdatasyncs after appends, since msync alone does not persist file-size metadata. - IoUringFile appends via pwritev2(RWF_APPEND). O_APPEND is not an option there: on Linux, pwrite on an O_APPEND fd appends regardless of the given offset, which would break positioned writes on clones sharing the fd. Concurrent appenders are out of contract (single logical writer); object-store backends surface the new AppendOffsetConflict error and recover via reopen() + retry. * io_bridge: add AsyncWrite/AsyncAppend and appendable BlobFile Add the write-side backend traits reserved next to AsyncRead: AsyncWrite (create/remove/save) powers a UniversalWriteFileOps impl on BlobFs (create-or-truncate put, delete, atomic whole-object save; directory ops are no-ops), and AsyncAppend — a single-request append where the offset must equal the current object size, acting as a compare-and-swap token — powers UniversalAppend on BlobFile. BlobFile caches the object size across appends (one HEAD for N appends; a missing object counts as empty so the first append creates it), concatenates batches into a single request, and drops the cache on reopen() — the documented recovery path after AppendOffsetConflict. Its flusher is a no-op: appends are durable once the backend acknowledges them. * io_bridge_object_store: native single-request S3 append object_store has no append support, so issue the PutObject + x-amz-write-offset-bytes request ourselves, reusing the store's credential chain (AmazonS3::credentials) and object_store's SigV4 AwsAuthorizer, which signs every header present on the request — no hand-rolled signing and no direct reqwest dependency. The offset doubles as a compare-and-swap token: a mismatch (400 InvalidWriteOffset, or 412 on some S3-compatibles) maps to AppendOffsetConflict. The write-offset append API exists on AWS S3 Express One Zone directory buckets and compatible stores (e.g. MinIO AiStor) — plain S3 Standard buckets reject it, and real Express zonal endpoints / session auth are not verified yet; MinIO-AiStor-compatible stores are the primary target for now. GCS and Azure sources simply do not implement AsyncAppend. ObjectStoreSource carries an AppendContext (HTTP client + object URL base + signing region) built per backend from its config, and gains a generic AsyncWrite impl (single-put create/save, delete). A test-only multi-request CAS emulation over InMemory exercises the BlobFile append stack hermetically; an end-to-end flow against a real append-capable store is gated behind S3_APPEND_INTEGRATION_TEST=1. * simple_disk_cache: write-through UniversalAppend for DiskCache Append to the remote (the single grow+write operation), then write the same bytes through into the local mirror so tail reads do not re-fetch what was just uploaded. LocalState::append_local keeps the fetched bitmap accurate: blocks fully covered by the appended range are marked fetched, and the pre-append partial tail block — which resize() drops because set_len zero-fills its gap — is re-marked only when its prefix was already fetched. If the mirror turns out stale (the remote grew behind our back), append falls back to resize-only and lazy fetches heal the gap on the next read. Writeable opens are now allowed on DiskCacheFs solely to enable append; DiskCache still never implements UniversalWrite. The writeable flag propagates to the remote handle, which is opened buffered instead of O_DIRECT: appends write through the page cache, which O_DIRECT reads on the same fd would fight (and IoUringFile rejects appends on prevent_caching handles). The remote-immutability docs are relaxed to append-only with an immutable prefix, matching what reopen() already assumed. Includes a full-stack composition test: DiskCache write-through over BlobFile offset tracking over an in-memory object store. * io_bridge_object_store: build the append HTTP client lazily Opening a source from an AwsConfig eagerly built the reqwest client (TLS setup, connection pool) even when append was never used. Keep the AppendContext construction to pure config (allow_http flag, object URL base, signing region) and build the client on first append instead, cached in an Arc<OnceLock> shared across clones of the source — and thus across the file handles opened from it. Sources that never append now pay nothing; client-construction errors surface on the first append instead of at open. * universal_io: test that append grows the regular file on disk The conformance suite reads appended bytes back through universal-io handles; also assert the underlying regular file itself — created outside universal_io, verified with plain fs reads — for both local backends. * Mention why we use custom HTTP client, object_store crate has no support * io_bridge_object_store: reject appends unconfirmed by the size header A store without write-offset support may accept the signed PutObject as a plain put — replacing the object with just the appended bytes — and return 2xx (community MinIO did exactly this before 2025-05, commit minio/minio@6d18dba9). The old success path fabricated the new length when x-amz-object-size was missing, so the destruction stayed invisible while every subsequent append repeated it. Require the x-amz-object-size response header (returned by AWS and MinIO AiStor appends) for any append at offset > 0 and fail loudly without it. Offset-0 appends are equivalent to a whole-object write, so they remain valid either way — a misconfigured store now fails on the second append instead of never. * io_bridge_object_store: honor endpoint/region env vars for appends With AwsCredentials::Default the store is built via AmazonS3Builder::from_env, which honors AWS_ENDPOINT_URL_S3, AWS_ENDPOINT_URL, AWS_ENDPOINT, AWS_REGION and AWS_DEFAULT_REGION — but append_context derived the append URL and SigV4 region only from the typed config fields. An env-configured deployment would read from one host while signing and sending appends to https://{bucket}.s3.us-east-1.amazonaws.com. Resolve the append endpoint and region the same way build_store does: explicit config first, then (default credential chain only) the same environment variables, with AWS_ENDPOINT_URL_S3 taking precedence as in from_env. The resolution is a pure function over an injected env lookup so the test does not touch process-global environment state. * simple_disk_cache: delegate the append flusher to the remote DiskCache's UniversalFlush impl was an unconditional no-op, justified by object-store appends being durable on acknowledgement — but the impl is generic over any appendable remote, and for local remotes (MmapFile, IoUringFile, exactly the compositions the tests instantiate) that silently dropped the fdatasync the UniversalAppend contract requires: append, flush Ok, power loss, appended bytes gone. Delegate to the remote's flusher once the cache is materialized: local remotes get their sync, object-store flushers remain no-ops, and a never-materialized cache has made no appends so a no-op stays correct. * io_bridge_object_store: retry transient append failures The append RPC was a single unretried HTTP attempt, while every other request in this stack goes through object_store's retry layer — a routine transient 503 SlowDown or connection reset failed the append hard where a concurrent read would have silently recovered. Retry connection errors, 5xx and 429 up to three attempts with a short linear backoff, re-signing per attempt (the SigV4 signature embeds the request date). Retrying is safe because the write offset is a compare-and-swap; the one ambiguity — an attempt that landed but whose acknowledgement was lost — surfaces as a write-offset conflict on the retry, which is reconciled with a HEAD: under the single-writer contract, an object size of exactly offset + data_len proves the tail is ours, so the append reports success instead of a spurious conflict (whose reopen-and-retry recovery would duplicate the record). * universal_io: forward TypedStorage::flusher for any UniversalFlush The flusher forwarding lived in TypedStorage's S: UniversalWrite impl block, so append-only storages (DiskCache, BlobFile — UniversalAppend + UniversalFlush but not UniversalWrite) offered append through the wrapper while the durability flusher the append contract mandates was unreachable without going through .inner. Move it to an S: UniversalFlush block: UniversalWrite implies UniversalFlush, so existing callers resolve unchanged, and duplicating the method instead would have hit E0592 on backends implementing both — the very ambiguity UniversalFlush was extracted to avoid. * io_bridge_object_store: surface unbuildable append requests as errors Request building could panic on two reachable paths: url accepts URIs the http crate rejects (IPv6 zone identifiers, URIs beyond u16::MAX bytes), and AppendContext::new is public so the object URL base is not guaranteed to be a base URL. Both expects become S3Config errors, so a configuration edge case fails the append instead of panicking the thread driving the bridge runtime. * simple_disk_cache: don't fail appends the remote already committed append_impl committed to the remote first and returned Err when the subsequent local-mirror update failed (e.g. ENOSPC on the cache volume) — indistinguishable from "nothing was appended", so a retrying caller would duplicate the record on the remote. The mirror is cache maintenance, not part of the append: on a failed write-through, log and degrade to bare growth so lazy fetches heal the unmarked blocks (safe — blocks are only marked fetched after their bytes landed). Only an unresizable mirror still surfaces an error, and the UniversalAppend contract now documents that an append Err does not guarantee nothing was appended: reopen() and re-check the length before retrying. * universal_io: bounds-check positioned io_uring writes against EOF The UniversalAppend contract states that UniversalWrite::write beyond the end-of-file fails and append is the only growth path — mmap enforces it, but IoUringFile's write/write_batch/write_multi were unchecked pwrites that silently extended the file with a zero-filled hole, inflating subsequent append offsets. Check every positioned write against the file length (fstat once per call), matching mmap's OutOfBounds semantics, and generalize the regression test to run on both local backends including the batched path. * io_bridge: reject appends on handles opened without writeable BlobFs::open dropped OpenOptions entirely, so a BlobFile opened with writeable: false still accepted appends — mmap and DiskCache enforce the writeable requirement, the blob backend silently didn't, and a stray append through a nominally read-only handle would mutate a shared object. Thread OpenOptions::writeable into BlobFile and reject appends with PermissionDenied when it is unset, mirroring the other backends. Directly-constructed handles (BlobFile::new/open, which take no OpenOptions) remain writeable. With every backend now enforcing the flag, the UniversalAppend contract drops its 'where the backend enforces open modes' hedge. * simple_disk_cache: answer empty appends from the mirror An empty append still went through remote.append_batch, so it returned the remote's live end-of-file — which can diverge from what this handle's len() and reads observe when the remote grew behind our back — while leaving the stale mirror unhealed (unlike a non-empty append in the same state, which resizes). Accept empty appends early: return the mirror length without touching the remote at all, keeping the answer consistent with the handle's own view. The trait contract now spells out that empty appends return the handle's view of the end of file without growth I/O. * universal_io: grow the mmap in place after appends Every mmap append ended in a full reopen(): an open+fstat+close by path just to learn the new length, and — on the non-Linux fallback, which rebuilds the mapping with the open-time populate flag — a re-population of the ENTIRE file per append, making appends O(file size) for handles opened with Populate::Blocking. Populating after an append is pointless anyway: we just touched the data we wrote. Extract the remap machinery into remap_to() (reopen() keeps its exact semantics, populate included) and add grow_mapping(): a stat-free grow that never re-populates. Appends learn the new length from a single fstat on the already-open O_APPEND fd — kept rather than trusting the mapping length, which is stale exactly in the externally-grown-remote scenario the disk cache heals through lazy fetches (the foreign-growth tests catch the difference). The mirror's resize() passes the length it just set_len'd, dropping its stat round-trip entirely. Per small append this is write+fstat+mremap, down from write+open+fstat+close+mremap, with no populate anywhere. * universal_io: share the mmap append fd across clones The flusher captured the per-clone append_file at flusher-creation time, so a flusher obtained from a sibling clone — or created before the handle's first append — msynced the shared mapping's data pages but skipped the fdatasync that persists the appended file size: a half-persist where a crash loses the acknowledged tail even though a flusher ran after the appends (writer thread + long-lived flush-worker clone is exactly the natural WAL shape). Store the fd in an Arc<OnceLock> shared by all clones and read it at flush time instead of capture time: any clone's append makes every handle's flusher sync the size metadata, whatever the clone/flusher creation order. Initialization races between clones keep exactly one fd. No hot-path cost: reads and positioned writes never touch the cell, and the append path pays one atomic load next to its syscalls. The interior mutability also lets append_fd take &self. * universal_io: document the clone remap hazard truthfully The remap SAFETY comment claimed moving is safe "since we are holding &mut self" — which says nothing about clones: they share the mapping but keep their own raw ptr/len copies, so after a moving (or, on non-Linux, replacing) remap a sibling clone's next read dereferences an unmapped address. The trait contract understated the same hazard as a concurrent-read constraint, while the UB persists after append returns. State the contract once on MmapFile (clones must reopen before reading after any growth; a stale clone read is undefined behavior, not a stale view), correct the SAFETY argument to rely on it explicitly, annotate the as_bytes unsafe blocks that depend on it, and sharpen the UniversalAppend contract bullet accordingly. Making clones structurally safe (resolving ptr/len through the shared Arc) is deliberately left as a separate change. * universal_io: share the vectored append machinery between backends The IOV_MAX-chunking / EINTR-retry / WriteZero / advance_slices loop existed twice — as local_file_ops::write_all_vectored (mmap) and inlined around pwritev2 in IoUringFile::append_slices — along with a verbatim collect/cast/filter-empties preamble in both append_batch impls. Two copies of subtle short-write handling introduced by one branch will diverge the first time only one of them gets a fix. Add an io::Write adapter whose write_vectored issues pwritev2(RWF_APPEND), letting the io_uring append delegate to the shared write_all_vectored, and hoist the slice collection into local_file_ops::collect_append_slices. IOV_MAX becomes private to the one function enforcing it. No behavior change; the existing conformance tests (including the beyond-IOV_MAX batch) cover both backends through the shared path. * io_bridge_object_store: add s3_express to the test config helper The AwsConfig struct gained the s3_express field; update the resolve-endpoint test helper accordingly. * universal_io: run the append conformance suite over the S3 stack Promote the backend-generic UniversalAppend battery (offsets, batches across IOV_MAX, empty appends, read-after-append, reopen visibility, flusher) from a private test into universal_io::conformance, exposed under the testing feature so backend crates can run the identical suite. mmap and io_uring keep running it as before; the object-store bridge now runs it too, over BlobFs/BlobFile with the in-memory offset-CAS append emulation — so local file system and S3 append behavior are asserted by the same test. The real write-offset RPC remains covered by the gated test_native_append_flow integration test. * Swap order * universal_io: disambiguate the io_uring crate import The import reorder dropped the leading `::`, making `io_uring` ambiguous with this very module (pulled into scope by the `use super::*` glob) and breaking the build. * io_bridge: don't materialize the mock object on rejected appends MutableMockSource::append called get_or_insert_with before validating the offset, so a rejected stale append against a missing object left an empty entry behind (exists() flipping true) — a fidelity gap versus the real backends, where a rejected append has no side effects. Check the offset against the current length first and only materialize the buffer on a match. * universal_io: disambiguate the io_uring crate import Restore the leading `::` on the io_uring crate import — without it the name is ambiguous with this very module, which the `use super::*` glob pulls into scope, and the crate fails to compile. Matches the sibling files (pool.rs, runtime.rs), which already import via `::io_uring`. * io_bridge_object_store: treat 404 under a nonzero append offset as a conflict A missing object while the handle expected a nonzero end-of-file is a stale view (the object was deleted behind our back) — the same situation as an offset mismatch, with the same reopen-and-retry recovery, and it is exactly what the in-memory emulation and the io_bridge mock already report. The RPC path mapped every 404 to NotFound instead, so the three implementations disagreed on the same logical case. Keep NotFound for offset-0 appends, where a 404 is a genuine missing-target error (e.g. a missing bucket) that retrying cannot heal. * Preallocate vector * universal_io: disallow appends through the disk cache Appends must go directly to the backing storage (mmap, io_uring, S3) — the disk cache is strictly read-only again. Remove DiskCache's UniversalAppend/UniversalFlush impls and the mirror write-through machinery (LocalState::append_local), reject writeable opens at DiskCacheFs::open, and drop the writeable/prevent_caching plumbing that existed solely for cached appends, restoring read-only remote handles. Attempting to append through the cache is now a compile-time error (the trait impl no longer exists), and opening a cached handle writeable is rejected at runtime, covered by a test in each backend variant. * universal_io: make append idempotent via caller-supplied offset Append now takes the byte offset where the data must land: append(offset, data) -> Result<()>. Every backend validates that the offset equals the current end of file before writing (mmap and io_uring fstat the fd, object stores validate server-side via x-amz-write-offset-bytes); on mismatch nothing is written and the append fails with AppendOffsetConflict. Retrying an already-landed append therefore conflicts instead of appending twice, and recovery is re-deriving the offset from len(). BlobFile no longer tracks the object length locally; the store's own offset check is the compare-and-swap. * Review remarks * Validate file length in S3 append response * Fix linting, we don't mind a large enum variant on index builder * universal_io: conformance-test stale-handle append conflict recovery Promote the two-handle conflict scenario from the in-memory BlobFile test into the backend-generic conformance battery: a second writeable handle grows the file, the stale handle's append conflicts cleanly (the offset check runs against the file, not the handle's view), and the contract's documented recovery — reopen, re-check the length, append at the real end — lands the data exactly once. Now exercised over mmap, io_uring, and the object-store stack instead of only the in-memory emulation. * io_bridge_object_store: stub-server tests for append response handling The native append's HTTP state machine was only exercised by the gated live-store integration test (S3_APPEND_INTEGRATION_TEST=1), so none of its branches ran in CI. Cover them hermetically against a minimal local HTTP stub — one connection per canned response, no new dependencies: - the signed write-offset PUT, and the new-size validation on success (matching, mismatching, unparseable, and absent size headers — the absent case at offset zero and past it); - conflict mapping for 400 InvalidWriteOffset, 412, and 404 under a nonzero offset, with 404 at offset zero staying NotFound, and a 400 without the conflict code staying a plain error; - 429/5xx retries re-sending the same offset, giving up after MAX_ATTEMPTS, and the lost-acknowledgement reconciliation via HEAD (accepted when the object ends at offset + len, rejected otherwise); - the status + body excerpt on unexpected failures. * simple_disk_cache: statically assert the cache stays read-only Disallowing appends through the disk cache made them a compile-time error by removing the impls; pin that with assert_not_impl_any so the UniversalAppend/UniversalFlush/UniversalWrite impls cannot quietly return. Runtime rejection of writeable opens stays covered per backend variant. |
||
|
|
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>
|
||
|
|
5593bc5564 |
Add optional last-modified timestamp to ListedFile and CachedFs FileInfo (#9803)
* Add optional last-modified timestamp to ListedFile and CachedFs FileInfo Filled where the listing backend exposes one: local filesystems (entry metadata) and object stores (ObjectMeta::last_modified). The uio-grpc backend reports None since the RPC does not carry mtimes. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Carry last-modified over the StorageRead ListFiles RPC Extend ListFilesEntry with an optional google.protobuf.Timestamp, fill it on the server from the listing metadata, and convert it back to SystemTime in the uio-grpc client. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Update lib/common/io_bridge_object_store/src/source.rs Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> * Fix missing SystemTime import in io_bridge_object_store The CodeRabbit-suggested last_modified mapping used SystemTime::from without importing std::time::SystemTime, breaking compile and CI. Co-authored-by: Cursor <cursoragent@cursor.com> * Retrigger CI after flaky integration-tests-consensus timeout The compile fix is in; the prior run failed on an unrelated 30s timeout in test_collection_recovery, not on PR changes. Co-authored-by: Cursor <cursoragent@cursor.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> Co-authored-by: root <111755117+qdrant-cloud-bot@users.noreply.github.com> Co-authored-by: Cursor <cursoragent@cursor.com> |
||
|
|
8f54076cbe |
Add uio-grpc backend and key-triggered live-reload to edge-shard-query (#9795)
* Add uio-grpc backend and key-triggered live-reload to edge-shard-query --backend uio-grpc opens the shard directly over a running Qdrant peer's StorageRead gRPC service (public gRPC endpoint, api-key aware), addressed by --collection/--shard-id — no object storage involved. The UioGrpcSource backend existed since #9634 but was never wired into the tool's CLI. --bucket is now per-backend optional (required for aws/gcs), and the default cache dir is scoped by collection/shard for uio-grpc, whose mirror has no distinguishing key prefix. --live-reload-key runs the same watch loop as --live-reload with each reload triggered by pressing Enter instead of a timer — easier when stepping through a debug scenario. Timer mode is unchanged; the two flags are mutually exclusive. Closed stdin ends the loop gracefully, so the flag cannot busy-loop on piped input (and is documented as incompatible with @- stdin request arguments). Verified end-to-end against a live instance started with QDRANT__FEATURE_FLAGS__WRITE_SEGMENT_MANIFEST=true: scroll over uio-grpc returns all points with payloads, timer mode picks up an upsert + delete as +/- diff lines, and key mode fires one reload per Enter and exits cleanly on EOF. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Log StorageRead not-found gRPC failures at debug level A read-only follower routinely probes files the writer creates lazily (e.g. the mutable id tracker's mappings and versions before the first flush), so every uio-grpc follower poll spammed INFO logs like: gRPC /qdrant.StorageRead/FileLength failed with NotFound "File not found: .../mutable_id_tracker.versions" NotFound on /qdrant.StorageRead/* now logs at debug; NotFound on all other services (missing collection etc.) stays at info. The service is matched by parsing the path's service component and comparing it to the tonic-generated storage_read_server::SERVICE_NAME constant rather than a hard-coded path string. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
dd4bb1eada |
Fix /readyz false-positive on freshly bootstrapped peer (#9774)
A bootstrapping peer seeds `conf_state` with just the first voter of the cluster (see `Consensus::init`) until it applies the configuration change entries of the Raft log. Since #9688, `member_peer_addresses()` filters known peer addresses by `conf_state` membership, so during that catch-up window the health checker saw a single member, took the single-node short-circuit in `cluster_commit_index()`, and latched /readyz to ready before the peer reached the cluster commit index. Only apply the `conf_state` filter when this peer is a member itself. This keeps the #9688 behavior for reinitialized peers (whose conf_state is reset to the peer itself) and genuine single-node clusters, while a catching-up peer falls back to all known peer addresses and properly waits for the cluster commit index. Fixes flaky test_replace_running_peer_without_shards_same_uri, where the test queried collection state right after /readyz passed, before the peer applied the CreateCollection entry. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
a7eb45931d |
fix(cli): parse Windows snapshot mappings (#9723)
* fix(cli): parse Windows snapshot mappings Parse CLI snapshot mappings from the final colon so Windows drive-letter paths keep their drive prefix. Keep full-snapshot recovery on typed mappings instead of formatting paths back into the CLI string protocol. Co-authored-by: chatgpt-codex-connector[bot] <199175422+chatgpt-codex-connector[bot]@users.noreply.github.com> * fix(cli): reject path-like snapshot collection names Reject path-like collection names when parsing CLI snapshot mappings so a Windows path without a target collection suffix reports the intended missing-collection error instead of producing a bogus mapping. Co-authored-by: chatgpt-codex-connector[bot] <199175422+chatgpt-codex-connector[bot]@users.noreply.github.com> Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> * Merge tests --------- Co-authored-by: chatgpt-codex-connector[bot] <199175422+chatgpt-codex-connector[bot]@users.noreply.github.com> Co-authored-by: coderabbitai[bot] <136622811+coderabbitai[bot]@users.noreply.github.com> Co-authored-by: timvisee <tim@visee.me> |
||
|
|
7d4f70bb19 |
Add time to gRPC responses (#9733)
* Add missing time response in some gRPC APIs, make consistent with REST * Don't use destructor |
||
|
|
b4e9326248 | Create concurrent snapshot in model testing (#9611) | ||
|
|
cad112bb1c |
Fix Clippy 1.97 (#9716)
* Remove from_iter_instead_of_collect from workspace lints The lint was removed from clippy (beta) and now triggers renamed_and_removed_lints warnings in every crate. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Fix clippy::chunks_exact_to_as_chunks Replace chunks_exact with a constant chunk size by as_chunks. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Fix clippy::needless_late_init Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Fix clippy::useless_borrows_in_formatting Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Fix clippy::uninlined_format_args Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Fix clippy::for_kv_map Iterate map values directly instead of discarding keys. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Allow clippy::result_large_err on QueueProxyShard::new_from_version The Err variant intentionally hands the LocalShard back to the caller. Same pattern as the existing allow on ForwardProxyShard::new. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Allow clippy::result_unit_err on wait_for_consensus_commit Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
ab0d3ecc62 |
Add unified memory: cold|cached|pinned placement parameter for collection components (#9684)
* Add unified `memory: cold|cached|pinned` placement parameter for collection components
Introduce a single `memory` parameter that controls how each collection
component's data is held in RAM, replacing the inconsistent zoo of
`on_disk` / `always_ram` / `on_disk_payload` flags:
- `cold`: not pre-loaded from disk, cached with usage
- `cached`: pre-populated into page cache on load, evictable under pressure
- `pinned`: materialized on heap, never evicted by cache pressure
The parameter is available on dense vectors, HNSW config, all quantization
configs, the sparse index, all payload field index types, and payload
storage (as a new `payload: { memory }` sub-object on collection params).
When set, it overrides the deprecated legacy flag; when unset, behavior is
unchanged. Legacy flags are marked deprecated (Rust + proto) but keep
working; conflicts are resolved in favor of `memory` with a warning.
New capabilities enabled by the tri-state model:
- HNSW graph links can be pinned (first production caller of the existing
`GraphLinksResidency::Pinned`)
- sparse mmap index, quantized vectors and on-disk payload field indexes
gain a `cached` tier (mmap + populate on open)
`pinned` is rejected by API validation for components without a heap
variant (dense vector storage, payload storage). Low-memory mode degrades
placements at load time via `Memory::clamp_to_low_memory`, matching the
existing `prefer_disk`/`skip_populate` behavior. Effective-placement
comparison in the config-mismatch optimizer avoids spurious rebuilds when
the same placement is expressed through the new parameter.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
* Fix gpu-gated tests for the new `memory` field
CI clippy runs with --all-features, which compiles the gpu-gated tests
that were missed locally: add the `memory` field to config literals and
allow deprecated placement params, same as in the rest of the tests.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* Add OpenAPI tests for memory placement, keep sparse config downgrade-clean
- OpenAPI tests: create/update collections with `memory` on every component,
assert the parameters are echoed in collection info, assert legacy-only
collections expose no new fields, and assert `pinned` is rejected (422)
for dense vector storage and payload storage on both create and update.
- Persist only the explicitly requested `memory` parameter in
`sparse_index_config.json` instead of the legacy-resolved placement, so
configurations using only the deprecated `on_disk` flag keep byte-identical
files that older Qdrant versions load without unknown fields.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* Validate collection meta ops at construction, not only in the API layer
The `memory: pinned` rejection for dense vectors and payload storage
lived in `Validate` impls on the internal request types, which only ran
through the REST actix extractor. gRPC validates just the proto message,
so a gRPC client could persist `pinned` where it is not supported and
have it silently treated as `cached`.
Run the derived validation in `CreateCollectionOperation::new` and
`UpdateCollectionOperation::new` instead: the constructors are the
common chokepoint for all API paths, before the operation is proposed
to consensus. This covers every validator on these types, not just the
`memory` checks, and keeps consensus-apply unaffected so mixed-version
clusters never reject already-committed operations.
`UpdateCollectionOperation::new` becomes fallible; `remove_replica` now
uses `new_empty` since it carries no user config. Regression tests drive
the gRPC conversion path and assert `InvalidArgument` for `pinned` on
create and update, with `cold`/`cached` accepted as a control.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
|
||
|
|
03f97d4c06 |
Fix /readyz of reinitialized peer waiting on a foreign consensus (#9688)
Applying `RemoveNode(self)` prunes all other peers from the removed peer's persisted address book, but the process may be stopped after the removal is committed and before the entry is applied. The old first peer's address then survives `--reinit`, and the readiness checker - which treated every `peer_address_by_id` entry as a cluster member - would wait for the reinitialized peer to reach the *old* cluster's commit index: a foreign consensus it can never catch up with, so `/readyz` never passed. Filter the address book by current `conf_state` membership instead, falling back to all known addresses while `conf_state` is still empty (a bootstrapping node that has not applied any configuration change yet). After `--reinit` the `conf_state` is reset to a single voter, so the readiness check correctly ignores peers of the old cluster. Fixes flaky `test_reinit_removed_peer`, which hit this race when the removed peer was killed before applying its own `RemoveNode`. The new regression test simulates that state and fails with the exact CI error without the fix. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
80c9454141 |
[UIO] Include file size in UniversalReadFileOps::list_files (#9675)
* [AI + manual] Include file size in `UniversalReadFileOps::list_files` * Use dedicated `ListedFile` struct instead of `(PathBuf, u64)` pair Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> --------- Co-authored-by: generall <andrey@vasnetsov.com> Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> |
||
|
|
df8ab9f360 |
Trace-log 'Missing request message.' internal gRPC errors (#9663)
Since the tonic 0.14 upgrade (hyper 1.x), a client cancelling a unary call between HEADERS and DATA surfaces as Internal "Missing request message." instead of Cancelled, because hyper hides stream resets during request body read (hyperium/hyper#3681). Cluster read fan-out generates these routinely, flooding logs with spurious ERROR lines. Log them at trace level, matching the existing Cancelled handling. Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |
||
|
|
b590f929ce |
Fix reinit failing with "removed all voters" on a removed peer (#9654)
* Fix reinit failing with "removed all voters" on a removed peer When `--reinit` is run on the first peer, `conf_state` is reset to a single voter (this peer), but the Raft log is left untouched. If this peer had been removed from consensus before reinit, its log still holds a committed-but-unapplied `RemoveNode(self)` conf-change. On startup that entry is replayed on top of the freshly reset single-voter config, and Raft aborts with "removed all voters", so the node can never start. Resetting only the apply-progress queue is not enough: a fresh single-node leader re-commits and re-applies any log entries still physically present beyond `commit`, re-triggering the failure. The stale tail has to be physically dropped from the WAL. On the first-peer reinit path, discard committed-but-unapplied entries inherited from the previous cluster: truncate the WAL to the last applied index, pin `commit` to it and clear the apply-progress queue. The entry at `commit` is retained as the snapshot anchor, so the first peer can still serve snapshots to bootstrapping peers. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * Add regression test for reinit of a removed peer Starts a 2-node cluster, gracefully removes the second peer from consensus (so it commits `RemoveNode(self)` into its own WAL), then reinitializes it as a fresh first peer. Before the fix this panicked on startup with "removed all voters"; the test asserts the peer comes back online, elects itself leader, and can still seed a fresh bootstrapping peer. Verified the test fails against the pre-fix binary with exactly: Failed to apply configuration change entry Caused by: Error in Raft consensus: removed all voters Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 4.8 (1M context) <noreply@anthropic.com> |
||
|
|
f3bb3f43f4 |
Restart consensus thread on failure instead of stopping permanently (#9662)
* Restart consensus thread on failure instead of stopping permanently If the consensus loop failed (e.g. due to a transient I/O error such as running out of disk space), the consensus thread stopped permanently and the only way to recover was a full service restart. Now the consensus thread rebuilds the Raft node from persisted state and retries with exponential backoff (1s initial, doubled up to 5 min cap, retrying indefinitely). Rebuilding from persisted state is equivalent to a process restart, so this introduces no new recovery semantics. The `reinit` logic is never repeated on restart, and initial startup remains fail-fast. The consensus message channel is created outside of `Consensus` so that internal gRPC handlers and the forward-proposals thread keep their senders across restarts, and messages buffered during the outage are drained after recovery. While in restart backoff, the node reports the existing `StoppedWithErr` cluster status with restart attempt info appended, and flips back to `Working` after a successful restart. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Add staging TestTransientError op and consensus restart integration test New staging-only cluster operation `test_transient_error` (mirroring `test_slow_down`): applying it fails with the given probability on the targeted peer, stopping its consensus thread with a service error. The peer re-applies the entry on every consensus thread restart, rolling the probability again, so it exercises the consensus restart loop end to end. Probability 1.0 simulates a permanently failing operation. The integration test runs a 3-peer cluster, poisons a follower through its own API (so the simulated failure is reported deterministically in the response), and asserts that the restart loop engages, the remaining peers keep serving consensus operations with a quorum, and the failed peer eventually recovers and catches up. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> * Fix import order to satisfy rustfmt Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> --------- Co-authored-by: Claude Fable 5 <noreply@anthropic.com> |