Commit Graph
1014 Commits
Author SHA1 Message Date
Andrey VasnetsovandClaude Fable 5.1 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>
2026-09-22 11:55:19 +02:00
qdrant-cloud-bot 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.
2026-09-21 21:04:19 +02:00
Tim Visée 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
2026-09-15 15:45:37 +02:00
Roman Titov f5e75477a0 Cleanup ConsensusStateMachine validation and docs (#10658) 2026-09-15 18:37:21 +09:00
Andrey VasnetsovandClaude Opus 5 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>
2026-09-13 16:22:44 +02:00
Andrey VasnetsovandClaude Fable 5.1 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>
2026-09-12 12:28:24 +02:00
Arnaud Gourlay b01f6d841e Perf: Move BM25 sparse embedding on a blocking thread (#10610) 2026-09-11 14:55:23 +02:00
Tim ViséeandRoman Titov b81ad12c20 Don't re-apply committed entries already applied during consensus start (#10277)
Co-authored-by: Roman Titov <ffuugoo@users.noreply.github.com>
2026-09-09 18:54:45 +02:00
Roman TitovandClaude Opus 5 9224014e93 Optimizations and improvements for ConsensusStateMachine (#10533)
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-09 17:25:13 +02:00
Roman TitovandClaude Opus 5 8185a2692e Validate consensus against ConsensusStateMachine (#10469)
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-09-09 12:36:19 +02:00
Tim Visée 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
2026-09-03 17:11:40 +02:00
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>
2026-09-03 11:16:38 +02:00
Arnaud GourlayandClaude Opus 5 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>
2026-08-31 15:49:27 +02:00
qdrant-cloud-bot 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.
2026-08-24 15:23:43 +02:00
Kumar ShivenduandClaude Opus 4.8 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>
2026-08-21 17:16:30 +05:30
Arnaud GourlayandClaude Opus 5 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>
2026-08-19 16:10:47 +02:00
qdrant-cloud-bot 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.
2026-08-17 18:25:19 +02:00
Arnaud GourlayandClaude Fable 5 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>
2026-08-14 11:09:39 +02:00
Andrey VasnetsovandClaude Fable 5 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>
2026-08-13 18:26:49 +02:00
Luis Cossío 500eed65b1 add etag to FileInfo (#10190) 2026-08-12 11:34:25 +02:00
Luis Cossío 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
2026-08-11 17:41:44 -04:00
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>
2026-08-06 12:17:29 +02:00
Andrey VasnetsovandClaude Opus 5 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>
2026-08-03 11:20:34 +02:00
Kumar ShivenduandClaude Opus 4.8 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>
2026-08-02 14:31:06 +02:00
Luis Cossío c07d57bd8f chore: fix dead code lints on macos (#10045) 2026-07-31 14:32:32 -04:00
Andrey VasnetsovandClaude Opus 5 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>
2026-07-31 10:40:20 +02:00
xzfc 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)
2026-07-30 20:59:31 +00:00
Andrey VasnetsovandClaude Opus 5 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>
2026-07-29 12:56:31 +02:00
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>
2026-07-28 21:17:36 +02:00
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>
2026-07-27 18:09:43 +02:00
Andrey VasnetsovandClaude Fable 5 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>
2026-07-24 10:35:52 +02:00
xzfc 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
2026-07-22 15:28:07 +00:00
xzfc 7e99cdd86a UioResult (#9933) 2026-07-22 14:34:33 +00:00
Arnaud GourlayandClaude Fable 5 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>
2026-07-16 15:46:19 +02:00
Ivan Pleshkov 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
2026-07-15 15:10:47 +02:00
Tim Visée 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.
2026-07-15 13:55:08 +02:00
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>
2026-07-14 11:16:47 +02:00
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>
2026-07-12 09:18:14 +02:00
Andrey VasnetsovandClaude Fable 5 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>
2026-07-11 16:04:54 +02:00
Andrey VasnetsovandClaude Fable 5 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>
2026-07-10 11:44:41 +02:00
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>
2026-07-08 14:10:01 +02:00
Tim Visée 7d4f70bb19 Add time to gRPC responses (#9733)
* Add missing time response in some gRPC APIs, make consistent with REST

* Don't use destructor
2026-07-08 11:36:01 +02:00
Arnaud Gourlay b4e9326248 Create concurrent snapshot in model testing (#9611) 2026-07-07 15:25:38 +02:00
Arnaud GourlayandClaude Fable 5 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>
2026-07-07 14:53:15 +02:00
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>
2026-07-07 14:32:51 +02:00
Andrey VasnetsovandClaude Fable 5 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>
2026-07-06 17:24:00 +02:00
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>
2026-07-04 00:19:17 +02:00
Andrey VasnetsovandClaude Fable 5 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>
2026-07-03 10:27:17 +02:00
Andrey VasnetsovandClaude Opus 4.8 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>
2026-07-03 10:13:30 +02:00
Andrey VasnetsovandClaude Fable 5 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>
2026-07-03 09:53:39 +02:00