Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions crates/asap-physical-operators/src/operators/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,7 @@ enum Kind {
SeriesLabels {
kind: planner_types::pre_asap::VectorMatchKind,
labels: Vec<String>,
unique: bool,
},
SeriesBinary {
operator: planner_types::post_asap::BinaryOperator,
Expand Down
42 changes: 40 additions & 2 deletions crates/asap-physical-operators/src/operators/series_labels.rs
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,27 @@ impl Operator {
) -> Result<Self, Error> {
layout(&input)?;
Ok(Self {
kind: Kind::SeriesLabels { kind, labels },
kind: Kind::SeriesLabels {
kind,
labels,
unique: false,
},
inputs: vec![input.clone()],
output: input,
})
}

/// Drop the metric name from each row's label set, as PromQL arithmetic
/// with a literal does. Unlike matching labels, the result is itself a
/// vector, so two rows that become equal are an error, as in Prometheus.
pub fn series_without_name(input: Schema) -> Result<Self, Error> {
layout(&input)?;
Ok(Self {
kind: Kind::SeriesLabels {
kind: VectorMatchKind::Ignoring,
labels: vec![],
unique: true,
},
inputs: vec![input.clone()],
output: input,
})
Expand Down Expand Up @@ -134,7 +154,15 @@ pub(super) fn execute<'a>(
let (rows, _memory) = collect_rows(left, &context).await?;
let mut result = Vec::new();
match (&operator.kind, right) {
(Kind::SeriesLabels { kind, labels }, None) => {
(
Kind::SeriesLabels {
kind,
labels,
unique,
},
None,
) => {
let mut seen = std::collections::BTreeSet::new();
for mut row in rows {
work.checkpoint().await?;
let mut set = left_layout.read(&output, &row)?;
Expand All @@ -146,6 +174,16 @@ pub(super) fn execute<'a>(
}
left_layout.write(&output, &mut row, &set)?;
workspace.grow(row_bytes(&row))?;
// Matching and grouping sides may repeat a label set;
// their consumers decide whether that is an error.
if *unique {
workspace.grow(set.iter().map(|(k, v)| 64 + k.len() + v.len()).sum())?;
if !seen.insert(set) {
return Err(invalid(
"vector cannot contain metrics with the same labelset",
));
}
}
result.push(row);
}
}
Expand Down
5 changes: 4 additions & 1 deletion crates/asap-physical-operators/src/operators/unchecked.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,10 @@ impl TryFrom<UncheckedOperator> for Operator {
at_ms,
steps,
)?,
Kind::SeriesLabels { kind, labels } => {
// The rebuilt kind must equal the serialized one, which rejects
// a unique rewrite with other matching labels.
Kind::SeriesLabels { unique: true, .. } => Operator::series_without_name(input(0)?)?,
Kind::SeriesLabels { kind, labels, .. } => {
Operator::series_labels(input(0)?, kind, labels)?
}
Kind::SeriesBinary { operator } => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -222,6 +222,7 @@ impl Lowering {
};
let step = self.value(vector)?;
let input = self.schema(&step);
let step = self.add(Operator::series_without_name(input.clone())?, vec![step]);
let value = named_column(&input, &ColumnRef::SampleValue)?;
let literal = Expression::Literal {
value: crate::values::Value::Float64(literal),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ pub(super) fn series_scalar_binary(
literal_left: bool,
) -> Result<Vec<Operator>, Error> {
arithmetic(operator)?;
let relabel = Operator::series_labels(input.clone(), VectorMatchKind::Ignoring, vec![])?;
let relabel = Operator::series_without_name(input.clone())?;
let value = input
.fields
.iter()
Expand Down
31 changes: 31 additions & 0 deletions crates/asap-physical-operators/tests/deployment_computation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,17 @@ fn execute(
dag: &PostAsapDag,
samples: &[Sample],
end: i64,
) -> Result<Vec<asap_physical_operators::runtime::SharedValue<Batch>>, String> {
execute_relabeled(dag, samples, end, &BTreeMap::new())
}

/// [`execute`], supplying samples of each instance in `relabel` under its
/// `(__name__, instance)` instead.
fn execute_relabeled(
dag: &PostAsapDag,
samples: &[Sample],
end: i64,
relabel: &BTreeMap<&str, (&str, &str)>,
) -> Result<Vec<asap_physical_operators::runtime::SharedValue<Batch>>, String> {
let inputs = raw_inputs(dag);
let program = compile(
Expand All @@ -122,6 +133,8 @@ fn execute(
.iter()
.filter(|sample| sample.0 == metric)
.map(|(name, job, instance, at, value)| {
let (name, instance) =
relabel.get(instance).copied().unwrap_or((*name, *instance));
let labels = BTreeMap::from([
("__name__".to_string(), name.to_string()),
("job".into(), job.to_string()),
Expand Down Expand Up @@ -469,6 +482,24 @@ fn per_series_scalar_arithmetic_applies_to_stored_readouts() {
}
}

// A literal operand drops the metric name; series that then share a label set
// are an error, as in Prometheus, rather than duplicate output series.
#[test]
fn per_series_scalar_arithmetic_rejects_label_sets_equal_without_the_name() {
let dag = exact_dag("sum_over_time(m[5m]) * 2");
let samples = counter("m", "api", 10., 10.)
.chain(counter("m", "api", 1., 1.).map(|s| (s.0, s.1, "y", s.3, s.4)))
.collect::<Vec<_>>();
let distinct = BTreeMap::from([("y", ("n", "y"))]);
assert!(execute_relabeled(&dag, &samples, 300_000, &distinct).is_ok());
// Both series are then {instance="x",job="api"}, named m and n.
let equal = BTreeMap::from([("y", ("n", "x"))]);
let Err(error) = execute_relabeled(&dag, &samples, 300_000, &equal) else {
panic!("duplicate label sets were accepted");
};
assert!(error.contains("same labelset"), "{error}");
}

fn with_vector_match(
mut dag: PostAsapDag,
kind: planner_types::pre_asap::VectorMatchKind,
Expand Down
29 changes: 28 additions & 1 deletion crates/asap-physical-operators/tests/promql_fallback.rs
Original file line number Diff line number Diff line change
Expand Up @@ -128,7 +128,8 @@ fn evaluate(
.flat_map(|(_, samples)| samples.iter())
.map(|(spec, seconds, value)| {
let mut labels = labels(spec);
labels.insert("__name__".into(), name.clone());
// A sample may supply its own `__name__`, as a series of another metric.
labels.entry("__name__".into()).or_insert(name.clone());
promql_rows::series_row(&schema, &labels, seconds * 1000, *value).unwrap()
})
.collect();
Expand Down Expand Up @@ -604,6 +605,32 @@ fn without_grouping_drops_labels_and_the_name() {
vec![(String::new(), 4.)]
);
assert!(labeled("sum without (inst) (a)", &[], 60).is_empty());
// Series equal without the name share a group rather than colliding.
let named: &[Sample] = &[
("job=x,inst=1", 50, 1.),
("__name__=b,job=x,inst=1", 50, 2.),
];
assert_eq!(
labeled("sum without (inst) (a)", &[("a", named)], 60),
vec![("job=x".into(), 3.)]
);
}

// Arithmetic with a literal drops the metric name; series that then share a
// label set are an error, as in Prometheus.
#[test]
fn literal_arithmetic_drops_the_name_and_rejects_equal_label_sets() {
let a: &[Sample] = &[
("job=x,inst=1", 50, 1.),
("__name__=b,job=x,inst=2", 50, 2.),
];
assert_eq!(
labeled("a * 2", &[("a", a)], 60),
vec![("inst=1,job=x".into(), 2.), ("inst=2,job=x".into(), 4.)]
);
let equal: &[Sample] = &[("job=x", 50, 1.), ("__name__=b,job=x", 50, 2.)];
let error = evaluate("a * 2", &[("a", equal)], 60).unwrap_err();
assert!(error.contains("same labelset"), "{error}");
}

// An empty label value is an absent label, and an empty side yields an empty
Expand Down
2 changes: 1 addition & 1 deletion docs/develop_docs/physical-compile-coverage.md
Original file line number Diff line number Diff line change
Expand Up @@ -125,7 +125,7 @@ Totals are unchanged: 19 Supported, 5 Partial, 5 Missing, 2 Backend.
| Row | Change |
|---|---|
| 5 | Query-time `Binary` over rows with a series identity, such as per-series readouts of stored state, uses the Fallback's `series_labels` and `series_binary`. Examples: `avg_over_time` as stored sum/count, and `rate(a) / rate(b)`. Matching drops `__name__` and honors `on`/`ignoring` when the payload carries them. Only one-to-one arithmetic is covered; `group_left`/`group_right` stay rejected and comparisons are row 7. Now Supported. |
| 4, 8 | A literal operand also applies to per-series rows and drops `__name__`. |
| 4, 8 | A literal operand also applies to per-series rows and drops `__name__`, in the Fallback too. Series whose label sets become equal are an error, as in Prometheus. |

Grouped `sum`/`avg`, current-series `Sum`/`Average` readouts, and
`sum_over_time`/`avg_over_time` use Prometheus' Kahan-Neumaier summation. An
Expand Down
Loading