Skip to content
Draft
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
40 changes: 20 additions & 20 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[workspace]
members = [
"crates/asap-physical-operators",
"crates/executor",
"crates/types",
"crates/frontend-common",
"crates/sql-function-catalog",
Expand Down
Original file line number Diff line number Diff line change
@@ -1,8 +1,12 @@
[package]
name = "asap-physical-operators"
name = "asap-executor"
version = "0.1.0"
edition = "2021"

# #509 Stage 4 reference executor. No stage crate (asap-logical-optimizer,
# asap-physical-optimizer, asap-plan-selection) depends on it; their manifest
# guards reject it. Its tests use the stage crates as dev-dependencies.

[dependencies]
chrono = { version = "0.4.39", default-features = false, features = ["std"] }
futures = "0.3"
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
# ASAP physical operators
# ASAP executor

An independent Rust physical operator DAG runtime shared by ingestion time and
The #509 Stage 4 reference executor: an independent Rust physical operator DAG runtime shared by ingestion time and
query time execution. The library requires neither backend engine, a server,
a storage implementation, Arrow nor DataFusion. DataFusion informed the design;
it is not the execution framework.
Expand All @@ -20,14 +20,14 @@ operator is needed. Summary construction updates state batch by batch. End of
input means the supplied query range or ingestion window is complete.

```rust
use asap_physical_operators::{
use asap_executor::{
expressions::Expression,
operators::Operator,
values::Value,
plan::PhysicalDAG,
runtime::{Limits, RunContext, Scope},
};
use asap_physical_operators::planner::ir::schema::DataType;
use asap_executor::planner::ir::schema::DataType;
use futures::{executor::block_on, StreamExt};

let source = Operator::scalar(Value::Int64(7), DataType::Int64)?;
Expand All @@ -44,7 +44,7 @@ let run = RunContext::new(
let mut output = plan.execute(&[1], run)?.remove(0);
let batch = block_on(output.next()).unwrap()?;
assert!(matches!(batch.rows()[0][0], Value::Int64(-7)));
# Ok::<(), asap_physical_operators::dag::Error>(())
# Ok::<(), asap_executor::dag::Error>(())
```

`physical_planner::compile` accepts a logical Post-ASAP DAG (`PhysicalASAPDAG`) and typed input contracts.
Expand Down
File renamed without changes.
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
//! Blocking operators enforce resources before returning their first batch.
use asap_physical_operators::{
use asap_executor::{
operators::Operator,
plan::{PhysicalDAG, PhysicalOperator},
runtime::{Limits, RunContext, Scope},
Expand Down Expand Up @@ -107,7 +107,7 @@ fn join_yields_during_computation_and_observes_cancellation() {
// Sorting and grouping yield even for one large batch.
#[test]
fn blocking_reductions_yield_and_release_memory_on_cancellation() {
use asap_physical_operators::{
use asap_executor::{
operators::{Reduction, SortKey},
plan::PhysicalOperator,
};
Expand Down Expand Up @@ -143,7 +143,7 @@ fn blocking_reductions_yield_and_release_memory_on_cancellation() {
// Merge-sort rounds preserve input order for tied keys across chunk boundaries.
#[test]
fn cooperative_sort_preserves_ties_across_chunks() {
use asap_physical_operators::operators::SortKey;
use asap_executor::operators::SortKey;
let batch = Batch::try_new(
schema(2),
(0..1025)
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
//! Spatial heap weights come from a fresh instant vector, never sample history.
mod common;
use asap_physical_operators::{
use asap_executor::{
operators::Operator,
physical_planner::{
promql_rows::{decode_series_identity, series_row, SERIES_IDENTITY_COLUMN},
Expand Down Expand Up @@ -234,7 +234,7 @@ fn spatial_heap_ranks_latest_values_in_independent_runs() {
// Blocking membership selection shares the run's cancellation and byte budget.
#[test]
fn current_series_observes_resource_limits() {
use asap_physical_operators::Error;
use asap_executor::Error;
let plan = snapshot_plan();
for cancelled in [false, true] {
let data = input(&[("one", 50_000, 1.)]);
Expand Down Expand Up @@ -273,7 +273,7 @@ fn current_series_observes_resource_limits() {

#[test]
fn identity_encoding_is_lossless_and_rejects_noncanonical_inputs() {
use asap_physical_operators::physical_planner::promql_rows::encode_series_identity;
use asap_executor::physical_planner::promql_rows::encode_series_identity;
let labels = BTreeMap::from([
("a".into(), "quote\"slash\\".into()),
("other".into(), "".into()),
Expand All @@ -296,7 +296,7 @@ fn identity_encoding_is_lossless_and_rejects_noncanonical_inputs() {
// test does not manually assemble the computation or its dependency edges.
#[test]
fn planner_current_series_candidate_compiles_with_dynamic_identity() {
use asap_physical_operators::physical_planner::{compile, promql_rows::with_series_identity};
use asap_executor::physical_planner::{compile, promql_rows::with_series_identity};
use planner_types::{types::AccuracyTarget, workload::*};
use std::rc::Rc;
let workload = PlanningWorkload {
Expand Down Expand Up @@ -334,7 +334,7 @@ fn planner_current_series_candidate_compiles_with_dynamic_identity() {
.candidate(&open_root)
.unwrap();
let snapshot_program =
asap_physical_operators::physical_planner::promql_rows::compile_current_series_evaluation(
asap_executor::physical_planner::promql_rows::compile_current_series_evaluation(
&open_selected,
)
.unwrap();
Expand Down
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
//! Exercise the public library without a backend server, store, or scheduler.
use asap_physical_operators::planner::ir::{
use asap_executor::planner::ir::{
scalar::ColumnRef,
schema::{
FieldDataType, GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, SummaryUpdate,
},
};
use asap_physical_operators::{factory::create_planner_accumulator, AggregateCore};
use asap_executor::{factory::create_planner_accumulator, AggregateCore};

fn family(k: u32) -> FieldDataType {
FieldDataType::Sketch(
Expand All @@ -28,9 +28,7 @@ fn build(values: &[f64]) -> Box<dyn AggregateCore> {
}
fn read(state: &dyn AggregateCore) -> f64 {
state
.estimate(
&asap_physical_operators::planner::ir::schema::SketchStatistic::Quantile { q: 0.5 },
)
.estimate(&asap_executor::planner::ir::schema::SketchStatistic::Quantile { q: 0.5 })
.unwrap()
}

Expand Down Expand Up @@ -64,8 +62,8 @@ fn invalid_kll_parameters_are_rejected_at_binding() {
// a packed-wire column-bit budget must not be imposed on this constructor.
#[test]
fn native_count_sketch_dimensions_are_not_packed_wire_dimensions() {
use asap_physical_operators::planner::ir::schema::SummaryInputExpr;
use asap_physical_operators::KeyByLabelValues;
use asap_executor::planner::ir::schema::SummaryInputExpr;
use asap_executor::KeyByLabelValues;
let family = FieldDataType::Sketch(
SketchKind::new(
SketchAlgorithm::CountSketchWithHeap,
Expand All @@ -85,7 +83,7 @@ fn native_count_sketch_dimensions_are_not_packed_wire_dimensions() {
let state = operator.into_accumulator();
let state = state
.as_any()
.downcast_ref::<asap_physical_operators::summary_kernels::CountSketchWithHeapAccumulator>()
.downcast_ref::<asap_executor::summary_kernels::CountSketchWithHeapAccumulator>()
.unwrap();
assert_eq!(state.query_key(&key), 7.0);
}
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
//! Planner-selected PromQL computation compiles from the timed DAG alone;
//! the deployment supplies only raw rows at the ingestion frontier.
mod common;
use asap_physical_operators::{
use asap_executor::{
operators::Operator,
physical_planner::{compile, promql_rows, CompiledPhysicalDAG, InputContract, Source},
runtime::{Limits, RunContext, Scope},
Expand Down Expand Up @@ -115,7 +115,7 @@ fn execute(
dag: &PhysicalASAPDAG,
samples: &[Sample],
end: i64,
) -> Result<Vec<asap_physical_operators::runtime::SharedValue<Batch>>, String> {
) -> Result<Vec<asap_executor::runtime::SharedValue<Batch>>, String> {
execute_relabeled(dag, samples, end, &BTreeMap::new())
}

Expand All @@ -126,7 +126,7 @@ fn execute_relabeled(
samples: &[Sample],
end: i64,
relabel: &BTreeMap<&str, (&str, &str)>,
) -> Result<Vec<asap_physical_operators::runtime::SharedValue<Batch>>, String> {
) -> Result<Vec<asap_executor::runtime::SharedValue<Batch>>, String> {
let inputs = raw_inputs(dag);
let program = compile(
dag,
Expand Down Expand Up @@ -624,8 +624,8 @@ fn population_sums_and_averages_are_compensated() {
// returns the sketch's total update weight, including colliding items.
#[test]
fn stored_count_min_bare_count_compiles_to_a_evaluation() {
use asap_executor::summary_kernels::CountMinSketchAccumulator;
use asap_logical_optimizer::{Replacement, ReplacementStrategy, TargetSubDAG};
use asap_physical_operators::summary_kernels::CountMinSketchAccumulator;
let root = lower_with("count(up)", AccuracyTarget::Epsilon(0.02));
let dag = asap_logical_optimizer::ASAPStrategies::default()
.replacements(&TargetSubDAG::new(&root))
Expand Down
Loading