* Split async IO into extension traits; only async-capable backends implement them
Move `read_bytes_async` / `open_async` off the universal `UniversalRead` /
`UniversalReadFs` traits into dedicated extension traits, `UniversalReadAsync`
and `UniversalReadFsAsync` (traits/async_io.rs). Only backends with a genuine
async story implement them — the blob family, the disk caches layered over it,
and a trivial ready-impl for mmap (tests and the mmap lookup path) — each in a
dedicated async_io.rs next to its sync impl.
`CachedFs` now requires its inner filesystem to be `UniversalReadFsAsync`; the
requirement reaches segment code through one supertrait bound on
`UniversalReadExt`. io_uring implements no async surface anymore: the
tokio_uring bridge thread, its tests, the musl-gated tokio-uring dependency,
and the `IoUringFile` read-only-segment wiring (`UniversalReadExt` impl and
the *RoIoUring condition-checker variants) are deleted — io_uring is not a
read-only-segment backend.
The payoff for live reload: `CachedFs::resolve_prefetched` awaits every parked
prefetch, and the edge refresh flow now runs preload -> resolve -> reload, so
the per-segment write locks never wait on IO.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* Decouple UniversalReadExt from the async filesystem requirement
UniversalReadExt is condition-checker dispatch; it never consumed the async
surface itself. Drop its `Fs: UniversalReadFsAsync` supertrait bound and relax
CachedFs's struct-level bound back to `UniversalReadFs` — the async requirement
now lives on the one impl that consumes it, `CachedReadFs for CachedFs`
(schedule_open parks the inner filesystem's `open_async` futures).
The bound then surfaces only on the lifecycle/preload impl blocks that go
through CachedReadFs (segment open, live-preload/reload, config reload, edge
load/refresh); the search path carries no async bounds at all.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
* existing segments: wait for IO outside of search pool
* new segments: wait for IO outside of search pool
* extract reload into separate function
* Update lib/edge/Cargo.toml
---------
Co-authored-by: Tim Visée <tim+github@visee.me>
* rename `reopen`->`live_reload` and `schedule_reopen`->`live_preload`
* `UniversalRead::live_preload` returns a shared future
* assert snapshot-miss eagerly on `live_preload`
`live_reload` cannot see the failed preload: its blocking fallback
re-resolves the length from the remote and succeeds. The error
surfaces at preload time, as callers (`ok_not_found`) expect.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
* [CachedFs] new `schedule` and `wait_all` primitives
* [AppendableIdTracker] don't reopen if just opened
* eager NotFound in `schedule_open`
* add traces for async reads
* finish `preopen`/`preload` with `wait_all`
* lock all segments in parallel for `live_reload`
* LIST before everything
to do: we don't have whole-fetch in async mode. to prevent sequential
`len`, we won't overlap static files with LIST.
* `wait_all` returns nothing
* `UniversalReadFs::open_async`
* `schedule_open` polls once
Scheduled opens must start eagerly: sync backends complete their
`open_async` on the first poll, preserving the prefetch contract
(handles outlive later file deletions/replacements). Moved down from
the integration branch so this PR stays green.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
* Read vector runs that straddle a chunk boundary
Resolve a run into per-chunk parts instead of a single range, borrowing
when it lands in one chunk and copying when it spans two. The read
pipeline schedules one range per read, so a straddling run is read
outside it.
No writer produces such a run yet, so this changes nothing on its own.
It is what a reader needs before one does — including edge and
live-reload readers, which read files a different version wrote.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Place multivector runs without regard to chunk boundaries
Writers appended a multivector's inner vectors at the end of the row
space unless the run would cross a chunk boundary, in which case they
skipped the chunk tail — the batch writers padding the skipped rows with
explicit zero rows. That made chunk geometry part of the interface every
multivector storage had to reuse.
Runs now go at the end unconditionally and the chunked storage splits
the write across chunks, as it already did for a batch of single
vectors.
What is left of the geometry is a size cap: a multivector may not exceed
one chunk. It is fill-independent, so it constrains nothing about
placement, and it is what the volatile storage needs anyway — that one
returns a plain slice and so cannot serve a straddling run.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Split a run at chunk boundaries in one place
Reading, writing in place and appending each derived the split from
`remaining_chunk_capacity`, so every one of them had to know that a run
does not necessarily fit where it starts.
`split_run` hands out the parts instead: one per chunk the run covers,
each carrying where it goes and how much of the run it takes. Nothing
asks how much room is left any more, and `get_chunk_offset` goes with it.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Keep straddling runs on the read pipeline
Reading a straddling run outside the pipeline blocked the scheduling
loop on one read, which costs a round trip on a backend that fetches
remotely and drops the batch back to sequential.
A run is now scheduled as one read per chunk it covers. Parts complete
in any order, so each run holds what has landed until the last part
does, then hands the callback the stitched vectors. Runs taking a single
read carry the caller's data in the tag and never touch that table.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Stop capping a multivector at one chunk
The cap outlived its reason on disk, but the volatile storage still
needed it: its `get_many` handed out a slice of one chunk, so a run that
crossed a boundary had nowhere to come from. And since a volatile
storage is a target of the batched copy that builds a segment, dropping
the cap only on disk would have turned a rejected write into a failed
merge.
So the volatile storage splits and stitches too. Both are a few lines
each, and placing a run no longer skips a chunk tail, so `extend` is now
`insert_many` at the end of the storage.
Nothing user-facing moves: `MAX_MULTIVECTOR_FLATTENED_LEN` caps a
multivector at 1M elements, far inside a 32 MiB chunk, so the storages
only ever rejected what reached them unvalidated.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Schedule a single-read run without the queue
The scheduling loop resolved every run into the queue and then took it
straight back out, so the overwhelmingly common run — one that fits a
chunk — paid a push and a pop for nothing. It now goes to the pipeline
directly, and the queue holds only what a straddling run leaves behind.
Worth ~10% on the multivector read benchmark, and it collapses the
"top up, then take" pair into one decision. Extracting that bookkeeping
into helpers instead was measured and is much worse: the mmap pipeline
alternates one schedule with one wait, so the loop body is a few dozen
nanoseconds, and a helper carrying the cold map and stitching paths is
too big for the compiler to inline back into it.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Repoint the multivector WAL-replay test at a live rejection
The test upserted a multivector too large for a storage chunk, which no
longer fails: the storages stopped capping one at a chunk. Nothing else
covered a multivector operation that only the apply path rejects.
A raw blob that is not a whole number of quantized records still does,
so the test now uses that, alongside its dense and sparse siblings.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Test reading multivectors with legacy chunk-tail padding
Locks the compatibility contract that pre-straddle files — runs that
skip a chunk's leftover slots — still reopen as single-chunk borrows.
* chore: retrigger CI after flaky test-consensus-compose
* Move ReadTag into for_each_vector, its only user
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01AN4Hgbd65gDhesthJk5bUY
---------
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
Co-authored-by: qdrant-cloud-bot <111755117+qdrant-cloud-bot@users.noreply.github.com>
Unindexed Match::Text and Match::Phrase previously shared a String::contains
arm, so phrase order was ignored and queries matched across token boundaries.
Use the default Word tokenizer for best-effort parity with indexed fields.
Fixes#10182
Add stop checks during the pre-HNSW build setup phase so cancellation
is observed promptly, and widen timing tolerance for noisy Windows debug
builds where post-stop delays can exceed 1s.
* Drop obsolete clippy large-error-threshold override
The 256 threshold was pinned for clippy 1.87 while tonic's `Status` was a
large error type. Upstream boxed its contents in `5de7bad` (hyperium/tonic#2253),
which is in the pinned 0.14.6 fork, so `Status` is now a single `Box` and the
default threshold of 128 passes.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Remove stale clippy allows
These 11 allows no longer suppress anything under any of the three CI clippy
configurations (default, --all-targets, --all-targets --all-features).
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
Drops 6 crates from the release build and 7 from the workspace test
build, with no source changes.
- geo: no triangulation, only Contains/Intersects/Haversine (spade, earcut)
- jsonwebtoken: HS256 from_secret only, no PEM keys (pem, simple_asn1)
- tar: nothing sets unpack_xattrs, which defaults to false (xattr)
- duplicate: every duplicate_item names its module (proc-macro2-diagnostics)
- pprof: no C++ frames to demangle (cpp_demangle)
Also promotes duplicate to a workspace dependency.
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
* perf: skip external-id resolution in search post-processing
process_search_result hands its scored points to retrieve(), which resolves
every external id back into an internal offset — although the offsets are
already known there: they come straight from the vector index as
ScoredPointOffset. Every request therefore performs top_k x segments
redundant external->internal lookups.
Whether that costs anything depends on the id tracker being mutable: a
freshly optimized segment looks ids up through a BTree, while a restarted
node maps an immutable tracker and resolves in constant time. So the effect
shows up right after ingestion or optimization and disappears after a
restart, which is what makes it easy to miss — a benchmark that starts from
a freshly booted node never sees it.
Split retrieve() into resolution + retrieve_resolved() and pass the offsets
from process_search_result directly, applying the deferred cutoff by offset
instead of resolving ids just to filter them. retrieve() behaviour is
unchanged for all other callers; retrieve_resolved() is private and states
that deferred filtering is the caller's responsibility.
Measured on glove-100-angular (1.18M points, 21 segments, top_k=10) at a
fixed request rate, on a node that had just finished ingesting: median
latency ~15-24% lower, ~8% less CPU per request. Both figures come from the
same node before and after the change.
* review fix
* review fix
---------
Co-authored-by: Ivan Dashchinskiy <iadashchinskiy@sbertech.ru>
* Optimization: Skipping items before pushing into PriorityQueue.
* Apply suggestion from top-k-update branch
* Stabilize order of equal scored items in tests
* Replace cgroups-rs with direct cgroup memory file reads
We used cgroups-rs in exactly one place, to read the memory limit and
usage of our own cgroup, so read those files directly instead. Drops 34
crates from the lockfile, including the zbus stack that carries
RUSTSEC-2026-0221.
Also fixes two latent cgroup v1 bugs (the LONG_MAX unlimited sentinel
reported ~9 EB of total memory, an unreadable limit file reported 0
bytes) and the hierarchy mix-up on hybrid hosts.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Decline cgroup memory reporting when the usage read fails
Reporting a usage of 0 made available_memory_bytes claim the whole cgroup
limit as free. Fall back to sysinfo when the usage file cannot be read at
init, and keep the last known value on a failed refresh.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Keep the last known memory limit when its read fails
A transient read failure cleared the cached limit and silently fell back
to host memory while the process was still capped, the same direction of
over-reporting as the usage read. Both now keep their last known value,
and a limit lifted at runtime still clears.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Treat a malformed memory limit as an error, not as unlimited
Parse failures returned Ok(None), so garbage in the limit file cleared a
valid cached limit on refresh and read as unlimited at init. Reserve
Ok(None) for "max" and the v1 sentinel, and report anything else as
InvalidData so the last known limit survives.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
Follow-up to #10287, which added `max`. Expressing a minimum still
required spelling out `(a + b - |a - b|) / 2`, the sign flip of the max
identity — drop the `neg` and you silently get a maximum instead. It also
only works for two operands and mentions each one twice, so the scorer
walks every sub-tree twice per candidate point.
The pair is what makes clamping expressible:
{"max": [0.0, {"min": [1.0, "$score"]}]}
`min` mirrors `max` throughout, and both guard helpers introduced in
#10287 already took an `operator: &str`, so they are reused unchanged: an
empty operand list is rejected at parse time rather than folding to
+infinity, and the Edge FFI rejects it at construction time. The result
needs no `is_finite` check, since `min` cannot produce a non-finite value
from finite inputs.
The unindexed-field walker shares one arm for `Max | Min` as the bodies
are identical, with a test pinning `min` separately so a later split
cannot silently drop it.
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
11 names across 7 files. Renames that did not reach the comment above them
(further_searches for further_results, query_context for segment_query_context,
block_ranges for local_block_ranges, op for operation twice, request for
requests, max_threads for max_kmeans_threads), and 3 arguments that were
removed from a signature and left documented (is_on_disk, collection_params,
search_runtime_handle with timeout).
Documentation only, no behaviour change.
* feat: add a dedicated max operator to score formulas
Expressing a maximum in a score formula required spelling out the
arithmetic identity `(a + b + |a - b|) / 2`. That is easy to get wrong
(the `/ 2` is load-bearing), only works for two operands, and mentions
each operand twice, so the scorer evaluates every sub-tree twice per
candidate point.
`max` is variadic, mirroring `sum` and `mult`:
{"max": ["$score", {"mult": [0.5, "popularity"]}]}
Unlike `sum` and `mult`, `max` has no identity element for the empty
case, so an empty operand list is rejected at parse time rather than
folding to -infinity and scoring every point with a non-finite value.
The check lives in `ExpressionInternal::parse_and_convert`, which every
entry point passes through, and the Edge FFI additionally rejects it at
construction time to match how that crate validates elsewhere.
The result needs no `is_finite` check: unlike `log10`, `exp`, `div`,
`sqrt` and `pow`, `max` cannot produce a non-finite value from finite
inputs, so it follows the existing `sum` convention.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
* test: cover max error propagation and datetime operands
An operand that fails must fail the whole expression rather than being
passed over in favour of a finite sibling. Covered with the failure both
before and after the finite operand: `mult` short-circuits on zero and
so can skip evaluating later operands, and this pins down that `max`
must not grow a similar shortcut that would swallow an error.
Also covers `max` over datetime operands, which reach the scorer through
a separate conversion to seconds, so that "score by whichever timestamp
is newer" is verified rather than assumed.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
* impl live_preload for payload indexes
enable live_preload for bool and null indexes
* (not) impl live_preload for `VectorIndexReadEnum`
* impl live_preload for `ReadOnlyPayloadStorage`
Extend the Windows-only ignore pattern to the remaining long-pole
unit/integration tests that dominate the Windows nextest phase:
- custom_query_scorer_equivalency::compare_scoring_equivalency
(~130s from product_x4 cases alone; rstest via test_attr)
- test_appendable_multi_turbo_vector_storage::{congruent_upsert_read_all_distances,
congruent_random_ops_{dot,cosine}} (~80s)
- test_appendable_turbo_vector_storage::{upsert_flush_reload_in_ram_matches_independent_oracle,
turbo_model_test_random_ops_{dot,cosine}} (~55s)
- scroll_filtering_test::test_filtering_context_consistency (~28s)
None of these exercise Windows-specific behavior; Linux/macOS keep
full coverage. Roughly ~300s of Windows test time removed on top of
the multivector rstest fix.
`multivector_filtrable_hnsw_test::test_multi_filterable_hnsw` and
`multivector_quantization_test::test_multivector_quantization_hnsw`
were meant to be ignored on Windows, but the
#[cfg_attr(target_os = "windows", ignore = "...")]
#[rstest]
pattern places the attribute BEFORE `#[rstest]`, so it never reaches
the per-case functions rstest generates — the cases were still running
on Windows CI (visible in the streaming test log).
Switch to the pattern that already works for
`byte_storage_quantization_test.rs`:
#[rstest]
#[cfg_attr(target_os = "windows",
test_attr(ignore = "..."))]
which uses rstest's `test_attr(...)` forwarding, so the `#[ignore]` is
applied to each generated per-case test.
Removes ~160 s of Windows test time (multivector_filtrable_hnsw ~100 s,
multivector_quantization ~64 s), which with ~3.6× test parallelism
should trim the Windows job wall-clock by ~45 s. Coverage on
Linux/macOS is unchanged.
* 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>
The `average_vector` recommend strategy folds the examples into one query
vector. That fold lived in `collection::recommendations`, though it only
touches segment/sparse types (`VectorRef`, `VectorInternal`,
`SparseVector::combine_aggregate`, `TypedMultiDenseVector`); edge-based
consumers (the serverless search worker) had to re-implement it because the
`collection` crate is the whole node layer.
Move `avg_vector_for_recommendation` (+ its private `avg_vectors` /
`merge_positive_and_negative_avg`) next to `RecoQuery` in
`segment::vector_storage::query`, returning `OperationResult` (validation
errors map to `CollectionError::BadInput` as before), keep the two collection
call sites on it, and re-export it — plus `VectorRef`, needed to call it —
from edge. Tests move along, plus one for the fold itself.
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
* Add `FlagsMode::from_feature_flags`, the mode for newly created flags
Compact in serverless-compatible deployments, dynamic otherwise. Only
creation consults it; opening existing flags detects their mode from
disk.
* Support the compact mode in the read-only flags types
Add `ReadOnlyFlags`, the mode-dispatching union of the two read-only
counterparts, serving the shared `RoaringFlagsRead` surface. Teach
`InMemoryBitvecFlags` to detect the mode it opens; its compact
`reload_appended` decodes the whole (small) file, as the format has no
random access.
* Create flags through mode selection in storages and indexes
Vector storage deleted flags and the bool/null indexes now open through
`open_or_create` with the mode from the feature flags: serverless
deployments create compact flags, dedicated ones keep creating dynamic
flags, and existing flags are opened in their detected mode either way.
* Read flags in either mode in the read-only bool and null indexes
`ReadOnlyFlags` shares the `RoaringFlagsRead` surface and the lifecycle
signatures of the roaring type it replaces, so the swap is a type
rename.
* Add TODO to not lock bitmask structure during flush
* `MutableStoredBitmask::save` returns the number of bytes written
Zero when the skip-clean save wrote nothing. Lets wrappers charge the
actual write to a hardware counter.
* Refuse to open compact flags in a dynamic-mode directory
Creating the compact file next to dynamic files would leave a directory
of both modes behind, which every later open rejects — refuse up front
instead. Both production callers already rule the case out through
`FlagsMode::detect`, so this only removes a foot-gun for future callers.
The open-or-eagerly-create logic moves into
`open_or_create_compact_mask`, shared with the update-only writer next.
* Rewrite `UpdateOnlyStoredFlags` onto the compact bitmask
The update-only flags writer now writes the compact mode — a single
roaring-encoded `compact_flags.dat` through `MutableStoredBitmask` —
instead of rewriting the whole padded dynamic file pair every batch. A
flush with no effective changes now writes nothing at all, where the old
writer rewrote the full mask on any `set`.
This also fixes opening serverless-created segments: the old open
eagerly wrote a `status.dat` into directories the writable side had
created in the compact mode, leaving files of both modes behind and
poisoning the directory for every later open.
A directory already holding dynamic-mode flags is refused loudly rather
than kept current or migrated; rebuild the segment to migrate its flags.
Migration may come later.
Drops the now-dead `InMemoryBitvecFlags::into_bitvec` and
`DynamicFlagsStatus::new`, and demotes `file_size_for` to private.
* Run edge tests with serverless feature flags
The edge fixtures ran with default feature flags, building leader shards
with dynamic-mode flags — a configuration edge never serves in
production, and one the update-only flags writer now refuses. It also
hid that the writer poisoned compact directories: no test exercised
update-only writes over a serverless-created shard.
Feature flags are process-global and first-init-wins, so every fixture
in the binary initializes the same serverless set; the manifest test
folds into it, since serverless implies `write_segment_manifest`.
* Don't use sequencial mode for one shot reads
* Add `CompactStoredFlags`, segment wrapper over the mutable bitmask
RAM-resident flags with a Flusher (skips the write when clean, cancels after drop) and files lister, backed by one compact stored-bitmask file rewritten whole on flush. Not integrated yet.
* Add `FlagsMode`, detecting the storage mode of a flags directory
`Dynamic` is the existing mmap stack for dedicated deployments, `Compact`
the compact stored-bitmask file for serverless ones; detection probes
which files are present. Also add the clippy allow the compact flags
tests were missing.
* Support the compact storage mode in `BitvecFlags` and `RoaringFlags`
The wrappers keep their in-memory read state in both modes; the new
`FlagsStorage` dispatches the write side between `BufferedDynamicFlags`
and `CompactStoredFlags`. `open_or_create` opens existing flags in their
detected mode and only applies `mode_if_create` to fresh ones — existing
call sites keep constructing the dynamic stack through `new`.
* Add `ReadOnlyCompactFlags`, read-only counterpart of compact flags
Bound to `UniversalRead`: opens on the bitmask header alone,
materializes the bitmap lazily on first query, and never creates a
missing file. Implements `RoaringFlagsRead` for the shared query
surface; `live_reload` reopens a fresh handle, as flushes replace the
file whole but cached handles keep serving the bytes they were opened
on. Not integrated yet.
* Skip compact live-reload tests on Windows, which forbids the rename
Both tests replace the compact flags file behind a reader whose disk
cache keeps the "remote" file mapped. On Unix the rename-over succeeds
and the mapping serves the old inode — the staleness under test — but
Windows forbids renaming over a mapped file, failing the writer's flush
with access denied. A limitation of the local-mmap remote stand-in, not
of the reload logic, which stays covered on the other targets.
* Don't check legacy flag file
The on-disk numeric index ended check_values_any with .unwrap_or(false),
so a real mmap read error was downgraded to "point does not match" and a
filtered search/scroll could silently return an incomplete result set.
Promote NumericIndexRead::check_values_any to OperationResult<bool> and
propagate the error through the dispatcher and the range condition checker,
which already returns OperationResult<bool>. Resolves the two FIXMEs in
numeric_index_read.rs and storage/read_ops.rs.
Unary inverse hyperbolic cosine, parallel to sqrt/ln/exp/log10, in REST,
gRPC, and edge (FFI + Python) interfaces. Inputs below 1 produce the same
NonFiniteNumber error as an invalid sqrt or ln.
Closes#10186
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
* Drop the vestigial UniversalWrite bound from the update-only writer
Neither writer kind performs in-place writes: AppendableSegment is built on
UniversalAppend, and DeleteOnlySegment tombstones via whole-mask atomic_save
(UniversalWriteFileOps), which UniversalAppend's supertrait already carries.
The bound is a leftover from the DiskIdTracker-based iterations that mutated
the deleted mask in place.
With it gone, UpdateOnlyEdgeShard::apply_batch is instantiable with the
object-store-appendable CachedBlobFile.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* edge-shard-update: --apply writes the batch, over object storage too
Open the object-storage backends through CachedBlobFs/CachedBlobFile
instead of the read-only DiskCacheFs handle, so the shard is appendable in
both modes, and add --apply: generate the same schema-derived batch and run
apply_batch instead of preview_batch. Dry run stays the default and the
generation is shared, so the preview cannot drift from what an apply would
do. AwsConfig::native_append is exposed as --native-append for
AiStor/RustFS-style endpoints; the Cached* types join io_bridge_object_store's
re-export of the io_bridge stack.
Applying to a leader-produced shard currently fails with a clean refusal —
its appendable segment's payload storage was created in mutable mode, which
the append-only writer rejects — the known segment-bootstrap gap, next in
line.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* CachedBlobFile: latency tracing for append_bytes
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* UpdateOnlyEdgeShard: sequential batches through one writer
Writers open once at shard open, next to the lookup segments they resume
from. apply_batch hands the writer back on success, live-reloading the
lookup half of every segment the batch wrote to (new
LookupSegment::live_reload, mirroring the read-only segment's); on error
the writer is consumed, since its lookups may no longer describe the
durable state.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* edge-shard-update: --interactive mode, sequential batches on one writer
After each applied batch, prompt on stdin for the next round's ids and
apply them through the writer apply_batch handed back — no shard
re-open — with op-num (and seed) incremented per round.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* CachedBlobFile: create the missing object on an offset-0 rewrite append
The caller-side rewrite path (part-copy S3 stores below the direct-append
threshold) validated the offset against the mirror length, whose
initialization HEAD-requests the remote and surfaced NotFound for an
object that does not exist yet. Direct-append backends (GCS compose,
native append) already create the object on an offset-0 append; the
rewrite now reads a missing remote as length zero so its whole-object PUT
does the same, and a non-zero offset against a missing object reports an
offset conflict.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* UpdateBatchOutcome: per-point records of retired slots
Each applied point now carries a PointApplyRecord: what happened to it
(stored/deleted/skipped/missing) and which slots it vacated where —
tombstoned per segment, or superseded in place for the old write-target
copy of a stored point. Built in the same loop that decides
tombstone-vs-supersede, so the report cannot drift from the writes.
edge-shard-update logs one line per point after the applied summary,
telling a fresh insert from an overwrite and naming the segments the old
copies were deleted from.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
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>
* 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>
A raw point can carry its payload as the byte blob it is stored as, mirroring
`PointStructRaw.raw_payload` on the internal gRPC API. The blob travels from the
sending node into the receiving node's WAL untouched, so the sender never parses
the payload it read and neither node builds a protobuf value tree for it.
It is parsed exactly once, where the operation is unpacked for apply
(`process_point_operation`), because that is the first place the parsed form is
actually needed: `set_full_payload` goes through the payload index, which cannot
be updated from bytes. The gRPC boundary therefore only checks the encoding tag
and rejects a point that sets both payload fields, the way the enclosing request
already rejects both `points` and `raw_points`.
Moving the parse onto the apply path makes its error classification load-bearing,
so a malformed blob is reported as `OperationError::MalformedPayloadBlob` — the
payload sibling of `MalformedVectorBlob`, mapped to `CollectionError::BadInput`
for the same reason: a bad blob that reached the WAL has to be skipped on replay
instead of crash-looping recovery.
Three consequences of the blob living that long are handled explicitly rather
than by convention:
- `decode_payload_raw` takes the blob only once it has parsed, so a failure
leaves the point holding it instead of holding neither representation.
- `upsert_points_raw` and `sync_points_raw` refuse a point that still carries a
blob. They read the parsed payload, so such a point would otherwise be stored
with no payload at all, and a `debug_assert!` would not catch it in release.
- `is_equal_to` compares blob to stored blob as bytes. A differing encoding costs
a redundant upsert on sync, never a skipped one.
The `raw_payload_transfer` bench measures the trade, per 100-point batch (one
transfer batch) at payloads of ~200 B / ~700 B / ~7 KB:
- Sender, storage bytes to wire: 16x / 37x / 113x faster. This is where the whole
win is — no parse of the blob that was read, no value tree built.
- WAL encode: 5x / 11x / 25x faster, writing a byte string instead of a map.
- Receiver, wire to applicable point: 1.09x / 1.10x / 1.06x. Near neutral, as it
swaps walking a prost value tree for a JSON parse.
- Wire bytes: ~6% smaller. WAL bytes: 10-32% *larger*, because the blob is JSON
while a parsed payload is written as a compact CBOR map.
The WAL growth is accepted rather than fixed: decoding earlier to win those bytes
back costs a second full deserialization, and would leave the receiving side with
a `payload_raw` that is never populated. Making the blob itself compact belongs in
the payload storage encoding (`RawPayloadEncoding` is the extension point for it),
not here.
Two flags, both off by default and both sender-only (nodes accept raw points and
raw payloads regardless), read where the transfer batch is prepared:
- `transfer_raw_points` transfers every collection as raw points, not only those
whose vector storage would drift in a decode-encode round-trip.
- `transfer_raw_payloads` ships the blob a raw read hands out; without it the
prepared batch decodes it back into the parsed payload, and the wire message is
exactly what it is today.
Neither is enabled by `all`: a node only accepts them once it runs a version that
understands them, so they can only be switched on a release later. Nothing
enforces that yet — the transfer has no peer-version gate.
Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
* Move delete-only tombstoning into per-format update-only id trackers
DeleteOnlySegment::tombstone_points wrote the deleted mask file directly,
hard-coding that both immutable id-tracker formats store it the same way.
Give each immutable format (in-RAM and disk-resident) its own update-only
tracker that owns the decision of where its tombstones go, and dispatch
through DeleteOnlyIdTrackerEnum, which lives next to ReadOnlyIdTrackerEnum.
Both formats share one deleted-mask file today, so the trackers delegate to
a shared writer in deleted_storage.rs; a format that diverges later changes
only its own update_only module.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
* Enforce empty-batch guard in the shared tombstone writer
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
---------
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>