Files
qdrant/src/model_testing.rs
T
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

266 lines
12 KiB
Rust

use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::{Duration, Instant};
use clap::Parser;
use collection::profiling::interface::init_requests_profile_collector;
use common::flags::{FeatureFlags, init_feature_flags};
use common::mmap::MULTI_MMAP_SUPPORT_CHECK_RESULT;
#[cfg(all(
not(target_env = "msvc"),
any(target_arch = "x86_64", target_arch = "aarch64")
))]
use tikv_jemallocator::Jemalloc;
#[cfg(all(
not(target_env = "msvc"),
any(target_arch = "x86_64", target_arch = "aarch64")
))]
#[global_allocator]
static GLOBAL: Jemalloc = Jemalloc;
#[derive(Parser, Debug)]
#[command(
version,
about = "Long-running soak test: random ops against a Collection, verified live against an in-memory model"
)]
struct Args {
/// RNG seed — same seed reproduces the exact op sequence.
#[clap(long, default_value_t = 0)]
seed: u64,
/// Number of randomized operations to apply. Ignored as the stop condition when
/// `--duration` is set (the run is then bounded by wall-clock time instead).
#[clap(long, default_value_t = 10_000, value_parser = clap::value_parser!(u64).range(1..))]
op_num: u64,
/// Run continuously for this many wall-clock seconds instead of stopping after `--op-num` ops.
/// When set, `--op-num` is ignored as the stop condition and the run ends at the deadline (or
/// on Ctrl-C, whichever comes first). The post-run live verification + final close/reopen
/// reload check run as usual on whatever ops were applied. Use this for time-boxed soak runs
/// (e.g. an overnight or per-CI-slot budget) where the interesting variable is "how long"
/// rather than "how many ops".
#[clap(long, value_parser = clap::value_parser!(u64).range(1..))]
duration_sec: Option<u64>,
/// Number of shards in the test collection.
#[clap(long, default_value_t = 3, value_parser = clap::value_parser!(u32).range(1..))]
shard_count: u32,
/// Size of the point ID space (ids are drawn uniformly from a pool of this many ids,
/// precomputed at startup; see `--uuid-id-fraction` for the numeric/UUID split).
#[clap(long, default_value_t = 500, value_parser = clap::value_parser!(u64).range(1..))]
id_pool: u64,
/// Fraction (0.0..=1.0) of the id pool backed by UUID point ids instead of numeric ones,
/// exercising the UUID side of the id tracker, WAL round-tripping and shard routing. The
/// UUID slots are well-formed v4 ids drawn once at startup from the seeded rng, so the
/// pool is stable for the run, `--seed` reproduces it exactly, and ids keep their
/// upsert-overwrite / delete-hits-live-point reuse semantics. 0.0 consumes no rng draws
/// and reproduces the numeric-only op stream of builds without UUID support.
#[clap(long, default_value_t = 0.5)]
uuid_id_fraction: f64,
/// Where to write the collection + snapshots data. Wiped at the start of each run;
/// copy the directory out beforehand to preserve a previous run for post-mortem.
#[clap(long, default_value = "./storage-model")]
storage_path: PathBuf,
/// Disable the segment optimizer entirely (sets `max_optimization_threads=0` in the
/// fixture). Useful for isolating optimizer-race bugs: if a failure stops reproducing
/// with this flag on, the bug depends on concurrent optimization.
#[clap(long, default_value_t = false)]
disable_optimizer: bool,
/// Optimizer `max_segment_size` in KB. Smaller values produce more segments and more
/// optimizer churn (more chances to trip races); larger values mimic production
/// cadence. The fixture's default (10 KB) is well below production thresholds because
/// the soak collection only holds a few MB total — at production defaults
/// (~200 MB) the optimizer would never fire on this workload.
#[clap(long, default_value_t = 10, value_parser = clap::value_parser!(u64).range(1..))]
max_segment_size_kb: u64,
/// Optimizer `indexing_threshold` in KB. Same rationale as `max_segment_size_kb`:
/// scaled down from production (~20 MB) to fire on the soak's small workload.
#[clap(long, default_value_t = 5, value_parser = clap::value_parser!(u64).range(1..))]
indexing_threshold_kb: u64,
/// Per-iteration probability (0.0..=1.0) of restarting the collection mid-run:
/// close + reopen + full model verification. Surfaces reload/WAL-replay bugs at the
/// op where they're introduced rather than only at end-of-run. 0.0 (default) skips
/// restarts entirely; small positive values (e.g. 0.005) are enough to expose the
/// known engine bugs without dominating the run with restart overhead.
#[clap(long, default_value_t = 0.0)]
restart_probability: f64,
/// Recompute the swarm config (the random subset of enabled ops) every N ops, turning one
/// long run into a sequence of swarm sub-tests for broader feature-interaction coverage
/// (Groce et al., "Swarm Testing"). Each redraw is logged to the trace. Set
/// `>= op_num` for a single config for the whole run. Default 10000 — long enough for a
/// config to build up meaningful state (e.g. a no-delete epoch grows segments) before the
/// next redraw; smaller values trade that depth for more config breadth.
#[clap(long, default_value_t = 10_000, value_parser = clap::value_parser!(u64).range(1..))]
swarm_interval: u64,
/// Set `on_disk` on the dense vector params. Vectors are persisted to disk either way; this
/// chooses mmap-on-demand storage (`VectorStorageType::Mmap`, reads hit the page cache/disk)
/// over a RAM-resident copy (`InRamMmap`, mmap populated into memory). Required for
/// `--async-scorer` to actually exercise io_uring: with an in-RAM copy, reads are served from
/// memory and the async read path is bypassed even though the io_uring storage is opened.
#[clap(long, default_value_t = false)]
on_disk: bool,
/// Enable the io_uring async scorer for mmap dense vector storage (Linux only). Stays on
/// plain mmap without a word if the kernel does not support io_uring, so read the storage's
/// `io_backend` telemetry to confirm it actually engaged. Pair with `--on-disk` (otherwise
/// reads come from the RAM-resident copy); has no effect on sparse or multivector storage.
#[clap(long, default_value_t = false)]
async_scorer: bool,
/// Run the live pre-close verification on a mid-run restart (the full model check before
/// `stop_gracefully`). Off by default: it's slow, and scrolling the whole collection first
/// warms caches / forces lazy loads that can mask a reload divergence, so the default cold
/// close+reopen surfaces more bugs. Enable it to attribute a `restart at op:N` failure to the
/// apply path vs. the reload path.
#[clap(long, default_value_t = false)]
pre_restart_check: bool,
/// Fix the number of Tokio worker threads instead of defaulting to the CPU core count.
/// Doesn't make thread interleaving reproducible (work-stealing timing, OS scheduling and
/// I/O readiness still vary), but pins the runtime shape so a `--seed` repro shares the same
/// worker count across machines. Leave unset for full parallelism on soak runs.
#[clap(long, value_parser = clap::value_parser!(u64).range(1..))]
worker_threads: Option<u64>,
/// Promote the always-disabled (`FORCE_OFF`) ops — DeleteByFilter, CreateVectorName,
/// DeleteVectorName — to forced-on, so they're enabled in every swarm config and guaranteed
/// to fire. These ops are masked off by default because they trip known engine bugs (see the
/// `FORCE_OFF` comment in `op/mod.rs`); enable this to deliberately reproduce them. The rng-draw
/// count is unchanged vs. the default, so the non-broken op stream stays reproducible per seed.
#[clap(long, default_value_t = false)]
enable_force_off: bool,
/// Disable the `CreateSnapshot` op entirely — it's masked out of every swarm config. Off by
/// default, so a normal soak takes snapshots in the background concurrently with the workload
/// (the archive is discarded — this stresses snapshot *creation* under concurrent writes, not
/// recovery). The op draws no rng, so `--seed` reproduces identically whether or not snapshots
/// are enabled; disable them only to avoid the background snapshot IO. The rng-draw count for
/// the swarm config is unchanged either way.
#[clap(long, default_value_t = false)]
disable_snapshots: bool,
}
fn main() {
let args = Args::parse();
// Built by hand (not `#[tokio::main]`) to seed Tokio's RNG off `--seed`, pinning in-poll
// draws like `select!` branch selection. `rng_seed` needs `--cfg tokio_unstable`; without
// it the binary still builds, just unseeded.
let mut builder = tokio::runtime::Builder::new_multi_thread();
builder.enable_all();
if let Some(worker_threads) = args.worker_threads {
builder.worker_threads(worker_threads as usize);
}
#[cfg(tokio_unstable)]
{
builder.rng_seed(tokio::runtime::RngSeed::from_bytes(
&args.seed.to_le_bytes(),
));
}
let runtime = builder.build().expect("failed to build tokio runtime");
runtime.block_on(run_main(args));
}
async fn run_main(args: Args) {
env_logger::init();
init_feature_flags(FeatureFlags::default());
let _ = MULTI_MMAP_SUPPORT_CHECK_RESULT.set(true);
init_requests_profile_collector(tokio::runtime::Handle::current());
// Ctrl-C handler: first signal flags shutdown so the run loop exits cleanly between
// ops; a second Ctrl-C falls through to the default handler and terminates. Wired
// BEFORE the run starts so an early signal still gets caught.
let shutdown = Arc::new(AtomicBool::new(false));
{
let s = shutdown.clone();
ctrlc::set_handler(move || {
if s.swap(true, Ordering::Relaxed) {
eprintln!("\nmodel_testing: second Ctrl-C, terminating");
std::process::exit(130);
} else {
eprintln!("\nmodel_testing: shutdown requested, draining current op...");
}
})
.expect("failed to install ctrl-c handler");
}
if !(0.0..=1.0).contains(&args.restart_probability) {
eprintln!(
"model_testing: --restart-probability must be in [0.0, 1.0], got {}",
args.restart_probability,
);
std::process::exit(2);
}
if !(0.0..=1.0).contains(&args.uuid_id_fraction) {
eprintln!(
"model_testing: --uuid-id-fraction must be in [0.0, 1.0], got {}",
args.uuid_id_fraction,
);
std::process::exit(2);
}
// Process-global flag, read when each on-disk dense vector storage is opened — set it before
// the fixture builds any segments.
segment::vector_storage::common::set_async_scorer(args.async_scorer);
// Show the active stop condition: the duration when time-bounded, else the op count.
let stop = match args.duration_sec {
Some(s) => format!("duration_sec={s}"),
None => format!("op_num={}", args.op_num),
};
println!(
"model_testing: seed={} {stop} shard_count={} id_pool={} uuid_id_fraction={} \
storage_path={} disable_optimizer={} max_segment_size_kb={} indexing_threshold_kb={} \
restart_probability={} swarm_interval={} \
on_disk={} async_scorer={} pre_restart_check={} enable_force_off={} \
disable_snapshots={}",
args.seed,
args.shard_count,
args.id_pool,
args.uuid_id_fraction,
args.storage_path.display(),
args.disable_optimizer,
args.max_segment_size_kb,
args.indexing_threshold_kb,
args.restart_probability,
args.swarm_interval,
args.on_disk,
args.async_scorer,
args.pre_restart_check,
args.enable_force_off,
args.disable_snapshots,
);
let start = Instant::now();
collection::model_testing::run(
args.seed,
args.op_num as usize,
args.shard_count,
args.id_pool,
args.uuid_id_fraction,
&args.storage_path,
args.disable_optimizer,
args.max_segment_size_kb as usize,
args.indexing_threshold_kb as usize,
args.restart_probability,
args.swarm_interval as usize,
args.on_disk,
args.pre_restart_check,
args.enable_force_off,
args.disable_snapshots,
args.duration_sec.map(Duration::from_secs),
shutdown,
)
.await;
println!("model_testing: ok in {:.1?}", start.elapsed());
}