use std::path::Path;
use crate::artifacts::Loaded;
use crate::ids::CaseId;
use crate::perf::{JourneyCatalogue, Measurement, PerfClass, PerformanceCase};
use crate::perf_run::client::PerfClient;
use crate::perf_run::corpus::SeededCorpus;
use crate::perf_run::pack::JourneyPack;
use crate::pipeline::{Error, load_clean_root, load_ixit, load_party_json, to_json_document};
use crate::probe::AqlProbeReport;
use crate::schema::results_schema;
use crate::stress::StressReport;
#[derive(Debug, Clone, Copy)]
pub struct SustainedWindow(u64);
impl SustainedWindow {
pub const LADDER: &'static [u64] = &[1, 2, 4, 6, 8, 12];
#[must_use]
pub fn hours(hours: u64) -> Option<Self> {
Self::LADDER.contains(&hours).then_some(Self(hours))
}
#[must_use]
pub fn seconds(self) -> u64 {
self.0.saturating_mul(3600)
}
}
impl Default for SustainedWindow {
fn default() -> Self {
Self(1)
}
}
#[derive(Debug, Clone, Copy)]
pub enum SeedStage {
BeforeScale,
AfterScale,
AfterWard,
}
#[derive(Debug)]
pub enum MeasuredEvent<'a> {
Progress(String),
CaseStarted {
case: &'a PerformanceCase,
source: &'a Path,
},
Measured(&'a Measurement),
PrunedOrphan(&'a CaseId),
Merged(&'a Path),
}
#[derive(Debug)]
pub struct MeasuredRun {
pub measurements: Vec<Measurement>,
pub earned_all: bool,
}
#[derive(Debug)]
pub struct MeasuredRequest<'a> {
pub root: &'a Path,
pub ixit: &'a Path,
pub results: &'a Path,
pub class: &'a str,
pub seed_workers: usize,
pub window: SustainedWindow,
}
#[derive(Debug)]
pub struct StressRequest<'a> {
pub root: &'a Path,
pub ixit: &'a Path,
pub corpus_class: &'a str,
pub seed_workers: usize,
pub step_secs: u64,
pub bisections: u32,
pub max_rate: f64,
}
#[derive(Debug)]
pub struct ProbeRequest<'a> {
pub root: &'a Path,
pub ixit: &'a Path,
pub corpus_class: &'a str,
pub seed_workers: usize,
pub requests: u32,
}
pub fn performance_case_of_class<'a>(
loaded: &'a Loaded,
class: PerfClass,
token: &str,
) -> Result<(&'a Path, &'a PerformanceCase), Error> {
loaded
.set
.performance
.iter()
.find(|(_, c)| c.class == class)
.map(|(path, case)| (path.as_path(), case))
.ok_or_else(|| {
Error::Missing(format!(
"no performance case of class {token} in the catalogue"
))
})
}
pub fn scale_opt_xml(loaded: &Loaded) -> Result<String, Error> {
let corpus_dir = loaded
.set
.corpus_dir
.as_deref()
.ok_or_else(|| Error::Missing("artifact set has no corpus directory".to_owned()))?;
let key = crate::ids::CorpusKey::parse("cnf.opt.blood_pressure")
.map_err(|e| Error::Instrument(e.to_string()))?;
let source = loaded
.set
.corpus
.as_ref()
.and_then(|(_, m)| m.get(&key))
.and_then(|entry| entry.source.clone())
.ok_or_else(|| {
Error::Missing("corpus manifest has no cnf.opt.blood_pressure fixture".to_owned())
})?;
std::fs::read_to_string(corpus_dir.join(&source))
.map_err(|e| Error::Instrument(format!("cannot read OPT fixture {source}: {e}")))
}
pub fn journey_context(loaded: &Loaded) -> Result<(JourneyCatalogue, JourneyPack), Error> {
let catalogue = loaded
.set
.journeys
.as_ref()
.map(|(_, catalogue)| catalogue.clone())
.ok_or_else(|| {
Error::Missing("artifact set has no vocab/journey_catalogue.yaml".to_owned())
})?;
let corpus_dir = loaded
.set
.corpus_dir
.as_deref()
.ok_or_else(|| Error::Missing("artifact set has no corpus directory".to_owned()))?;
let manifest = loaded
.set
.corpus
.as_ref()
.map(|(_, manifest)| manifest)
.ok_or_else(|| Error::Missing("artifact set has no corpus manifest".to_owned()))?;
let pack = JourneyPack::load(corpus_dir, manifest, &catalogue).map_err(Error::Instrument)?;
Ok((catalogue, pack))
}
pub fn seed_corpus(
client: &PerfClient,
corpus_key: &str,
opt_xml: &str,
journey_pack: &JourneyPack,
seed_workers: usize,
progress: &(dyn Fn(String) + Sync),
stage: &mut dyn FnMut(SeedStage),
) -> Result<SeededCorpus, Error> {
use crate::perf_run::corpus;
let (ehrs, versions) = corpus::scale_shape(corpus_key).map_err(Error::Instrument)?;
stage(SeedStage::BeforeScale);
let mut seeded = corpus::seed_scale_ladder(
client,
corpus_key,
opt_xml,
ehrs,
versions,
seed_workers,
progress,
)
.map_err(|e| Error::Instrument(format!("seeding failed: {e}")))?;
stage(SeedStage::AfterScale);
corpus::seed_ward(client, &mut seeded, journey_pack, seed_workers, progress)
.map_err(|e| Error::Instrument(format!("ward seeding failed: {e}")))?;
stage(SeedStage::AfterWard);
Ok(seeded)
}
pub fn run_stress(
request: &StressRequest<'_>,
progress: &(dyn Fn(String) + Sync),
) -> Result<StressReport, Error> {
use crate::perf_run;
let class = PerfClass::parse(request.corpus_class).map_err(Error::Selector)?;
let loaded = load_clean_root(request.root)?;
let (ixit, _) = load_ixit(request.ixit)?;
let (principals, environment) =
perf_run::window::measured_run_context(&ixit).map_err(Error::Instrument)?;
let client = principals.primary().clone();
let (_, case) = performance_case_of_class(&loaded, class, request.corpus_class)?;
let opt_xml = scale_opt_xml(&loaded)?;
let (catalogue, journey_pack) = journey_context(&loaded)?;
let corpus = seed_corpus(
&client,
case.corpus.as_str(),
&opt_xml,
&journey_pack,
request.seed_workers,
progress,
&mut |_| {},
)?;
let options = crate::stress::StressOptions {
step_hold_s: request.step_secs.max(10),
bisections: request.bisections,
max_rate: request.max_rate,
..crate::stress::StressOptions::default()
};
let workload = perf_run::schedule::JourneyWorkload {
catalogue: &catalogue,
shares: &case.workload.journeys,
pack: &journey_pack,
curve: crate::perf::ArrivalCurve::Uniform,
principals: &principals,
};
let report = crate::stress::run_stress(
&principals,
&corpus,
&workload,
environment,
ixit.containers.as_ref(),
&options,
progress,
)
.map_err(|e| Error::Instrument(format!("stress run failed: {e}")))?;
if perf_run::rate_limited_observed() {
return Err(Error::Instrument(perf_run::rate_limited_refusal("stress")));
}
Ok(report)
}
pub fn run_aql_probe(
request: &ProbeRequest<'_>,
progress: &(dyn Fn(String) + Sync),
) -> Result<AqlProbeReport, Error> {
use crate::perf_run;
let class = PerfClass::parse(request.corpus_class).map_err(Error::Selector)?;
let loaded = load_clean_root(request.root)?;
let (ixit, _) = load_ixit(request.ixit)?;
let (principals, environment) =
perf_run::window::measured_run_context(&ixit).map_err(Error::Instrument)?;
let client = principals.primary().clone();
let (_, case) = performance_case_of_class(&loaded, class, request.corpus_class)?;
let opt_xml = scale_opt_xml(&loaded)?;
let (_, journey_pack) = journey_context(&loaded)?;
let corpus = seed_corpus(
&client,
case.corpus.as_str(),
&opt_xml,
&journey_pack,
request.seed_workers,
progress,
&mut |_| {},
)?;
let options = crate::probe::ProbeOptions {
requests: request.requests,
};
crate::probe::run_probe(
&client,
&corpus,
environment,
ixit.containers.as_ref(),
&options,
progress,
)
.map_err(|e| Error::Instrument(format!("probe run failed: {e}")))
}
#[expect(
clippy::too_many_lines,
reason = "the measured window is one sequence: seed, settle, drive, attach, merge"
)]
pub fn run_measured(
request: &MeasuredRequest<'_>,
observe: &(dyn Fn(MeasuredEvent<'_>) + Sync),
) -> Result<MeasuredRun, Error> {
use crate::perf_run;
let class = PerfClass::parse(request.class).map_err(Error::Selector)?;
let loaded = load_clean_root(request.root)?;
let (ixit, _) = load_ixit(request.ixit)?;
let (principals, environment) =
perf_run::window::measured_run_context(&ixit).map_err(Error::Instrument)?;
let client = principals.primary().clone();
let selected: Vec<_> = loaded
.set
.performance
.iter()
.filter(|(_, c)| c.class == class)
.collect();
if selected.is_empty() {
return Err(Error::Missing(format!(
"no performance case of class {} in the catalogue",
request.class
)));
}
let opt_xml = scale_opt_xml(&loaded)?;
let (catalogue, journey_pack) = journey_context(&loaded)?;
let progress = |message: String| observe(MeasuredEvent::Progress(message));
let containers = ixit.containers.clone();
if containers.is_none() {
progress("resources: not sampled (ixit declares no `containers` block)".to_owned());
}
let mut run = MeasuredRun {
measurements: Vec::new(),
earned_all: true,
};
for (path, case) in selected {
observe(MeasuredEvent::CaseStarted { case, source: path });
let mut disk = crate::perf::DiskAnchors {
before_scale_seed_bytes: None,
after_scale_seed_bytes: None,
after_ward_seed_bytes: None,
after_window_bytes: None,
seed_compositions: perf_run::corpus::scale_shape(case.corpus.as_str())
.ok()
.and_then(|(ehrs, versions)| u64::try_from(ehrs.saturating_mul(versions)).ok()),
};
let probe_volume = |label: &str| -> Option<u64> {
let db = &containers.as_ref()?.db;
match perf_run::resources::db_volume_bytes(db) {
Ok(bytes) => {
progress(format!("disk anchor {label}: {bytes} bytes"));
Some(bytes)
}
Err(e) => {
progress(format!("disk anchor {label} unavailable: {e}"));
None
}
}
};
let corpus = seed_corpus(
&client,
case.corpus.as_str(),
&opt_xml,
&journey_pack,
request.seed_workers,
&progress,
&mut |milestone| match milestone {
SeedStage::BeforeScale => {
disk.before_scale_seed_bytes = probe_volume("before scale seed");
}
SeedStage::AfterScale => {
disk.after_scale_seed_bytes = probe_volume("after scale seed");
}
SeedStage::AfterWard => {
disk.after_ward_seed_bytes = probe_volume("after preflight + ward seed");
}
},
)?;
if let Some(c) = &containers {
progress(
"settling maintenance before the measured window (vacuumdb --analyze)".to_owned(),
);
if let Err(e) = perf_run::resources::settle_maintenance(&c.db) {
progress(format!("maintenance not settled: {e}"));
}
}
let warmup_s = case.workload.warmup.0;
let duration_s = case.workload.duration.0.max(request.window.seconds());
let sampler = containers
.as_ref()
.map(|c| perf_run::resources::ResourceSampler::start(c, warmup_s, duration_s));
let mut measurement = perf_run::window::drive_case(
case,
&principals,
&corpus,
&journey_pack,
&catalogue,
environment,
warmup_s,
duration_s,
&progress,
)
.map_err(|e| Error::Instrument(format!("measured run failed: {e}")))?;
if let Some(sampler) = sampler {
let (series, notes) = sampler.stop();
for note in notes {
progress(note);
}
disk.after_window_bytes = probe_volume("after measured window");
let sampled_any = series.iter().any(|s| !s.samples.is_empty());
let anchored_any = disk.before_scale_seed_bytes.is_some()
|| disk.after_scale_seed_bytes.is_some()
|| disk.after_ward_seed_bytes.is_some()
|| disk.after_window_bytes.is_some();
if sampled_any || anchored_any {
measurement.resources = Some(crate::perf::ResourcesRecord {
sample_interval_s: perf_run::resources::SAMPLE_INTERVAL.as_secs(),
containers: series,
disk: Some(disk),
});
} else {
progress(
"resources: not sampled (container runtime unreachable for the whole run)"
.to_owned(),
);
}
}
observe(MeasuredEvent::Measured(&measurement));
if measurement.verdict != crate::perf::ClassVerdict::Earned {
run.earned_all = false;
}
if perf_run::rate_limited_observed() {
return Err(Error::Instrument(perf_run::rate_limited_refusal("perf")));
}
merge_measurement(request, &loaded, measurement.clone(), observe)?;
run.measurements.push(measurement);
}
Ok(run)
}
fn merge_measurement(
request: &MeasuredRequest<'_>,
loaded: &Loaded,
measurement: Measurement,
observe: &(dyn Fn(MeasuredEvent<'_>) + Sync),
) -> Result<(), Error> {
let mut results: crate::party::Results =
load_party_json(request.results, &results_schema(), "results.schema.json")?;
results.measurements.retain(|m| m.case != measurement.case);
results.measurements.retain(|m| {
let known = loaded.set.performance.iter().any(|(_, c)| c.id == m.case);
if !known {
observe(MeasuredEvent::PrunedOrphan(&m.case));
}
known
});
results.measurements.push(measurement);
results
.measurements
.sort_by(|a, b| a.case.as_str().cmp(b.case.as_str()));
let document = to_json_document(&results, "serialize")?;
crate::pipeline::write_file(request.results, &document)?;
observe(MeasuredEvent::Merged(request.results));
Ok(())
}