Skip to content
Open
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
15 changes: 3 additions & 12 deletions crates/asap-physical-operators/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -116,15 +116,6 @@ the shared runtime with independent per-run state. Window coverage, revision and
maintenance-policy admission remain deployment/planning contracts; this compiler
does not discover storage or silently change a selected maintenance strategy.

`physical_planner::compile_temporal_pane_candidate` lowers a selected continuous
KLL lifecycle and Sliding/Tumbling framework into maintenance and query DAGs.
`TemporalPaneMaintenance` supplies pane geometry and a resolved complete entity
identity contract. The compiler inserts population guards, scan predicates,
pane construction, ordered state slots, a shared merge and quantile readouts.
Pane outputs have distinct physical identities from the logical whole-window
summary, and the returned candidate retains the maintenance contract for binding.
Each run checks phase, pane timestamps and duplicate entity states. The initial
realization uses complete bounded snapshots; partial edges, exponential
histograms and cross-run delta accumulation are unsupported. Storage identities,
revision selection, completeness/readiness evidence and scheduling stay with
deployment.
Pane construction and geometry belong to deployment. A deployment runs the
precompute DAG once per pane it constructs and binds the selected pane states to
query input slots; the query DAG merges and reads them out as computation.
18 changes: 3 additions & 15 deletions crates/asap-physical-operators/src/operators/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,8 +21,8 @@ mod current_series;
mod filter;
mod joins;
mod limit;
mod panes;
mod projection;
mod scope_timestamp;
mod sort;
mod source;
mod summary;
Expand All @@ -40,11 +40,6 @@ enum Kind {
value: Value,
dtype: DataType,
},
PaneInput {
coordinate: usize,
layout: planner_types::post_asap::PaneLayout,
offset_ms: Option<i64>,
},
ScopeTimestamp {
columns: Vec<Option<usize>>,
},
Expand Down Expand Up @@ -260,10 +255,7 @@ impl PhysicalOperator<Batch, Schema> for Operator {
};
PlanProperties {
boundedness,
emission: if matches!(
self.kind,
Kind::PaneInput { .. } | Kind::ScopeTimestamp { .. }
) {
emission: if matches!(self.kind, Kind::ScopeTimestamp { .. }) {
inputs
.first()
.map_or(Emission::Unknown, |input| input.emission)
Expand All @@ -279,7 +271,6 @@ impl PhysicalOperator<Batch, Schema> for Operator {
match self.kind {
Kind::Source(_) => "Source",
Kind::Constant { .. } => "Constant",
Kind::PaneInput { .. } => "PaneInput",
Kind::ScopeTimestamp { .. } => "ScopeTimestamp",
Kind::Union => "Union",
Kind::CurrentSeries { .. } => "CurrentSeries",
Expand All @@ -303,7 +294,6 @@ impl PhysicalOperator<Batch, Schema> for Operator {
}
}
fn validate_context(&self, context: &RunContext) -> Result<(), Error> {
panes::validate_context(self, context)?;
current_series::validate_context(self, context)?;
self.readout_range(context).map(|_| ())
}
Expand Down Expand Up @@ -332,9 +322,7 @@ impl PhysicalOperator<Batch, Schema> for Operator {
}
Kind::Project(_) => projection::execute(self, inputs, context),
Kind::CurrentSeries { .. } => current_series::execute(self, inputs, context),
Kind::PaneInput { .. } | Kind::ScopeTimestamp { .. } => {
panes::execute(self, inputs, context)
}
Kind::ScopeTimestamp { .. } => scope_timestamp::execute(self, inputs, context),
Kind::Filter(_) => filter::execute(self, inputs, context),
Kind::Limit { .. } => limit::execute(self, inputs, context),
Kind::Sort { .. } => sort::execute(self, inputs, context),
Expand Down
228 changes: 0 additions & 228 deletions crates/asap-physical-operators/src/operators/panes.rs

This file was deleted.

Loading
Loading