1use std::cell::Cell;
29use std::sync::Arc;
30
31use rudb_catalog::{Catalog, QualifiedName};
32use rudb_common::{Cancel, Memory, Result};
33use rudb_functions::TableFunction;
34use rudb_metrics::{Counters, Report};
35use rudb_pipeline::{Source, Watched};
36use rudb_plan::{Node, NodeRef, Pipelines, Plan};
37
38use crate::adapt::{Broken, Fed, Paired, Pulled, Streamed};
39use crate::cancel::Guarded;
40use crate::gather::{Gather, Keep};
41use crate::group::{Aggregate, Distinct};
42use crate::join::{CrossProduct, Gathered, Join};
43use crate::operator::Operator;
44use crate::schema::Schema;
45use crate::setop::SetOp;
46use crate::sort::Sort;
47use crate::source::{Dummy, FileScan, Scan, Series, Values};
48use crate::strategies::Strategies;
49use crate::stream::{Filter, Limit, Project};
50use crate::topn::TopN;
51
52pub fn build<'a>(plan: &'a Plan, catalog: &'a Catalog) -> Result<Box<dyn Operator + 'a>> {
59 build_with(plan, catalog, &Cancel::new(), &Memory::unlimited())
60}
61
62pub fn build_with<'a>(
82 plan: &'a Plan,
83 catalog: &'a Catalog,
84 cancel: &Cancel,
85 memory: &Memory,
86) -> Result<Box<dyn Operator + 'a>> {
87 build_measured(plan, catalog, cancel, memory, &Report::new())
88}
89
90pub fn build_measured<'a>(
100 plan: &'a Plan,
101 catalog: &'a Catalog,
102 cancel: &Cancel,
103 memory: &Memory,
104 report: &Report,
105) -> Result<Box<dyn Operator + 'a>> {
106 let pipelines = Pipelines::of(plan);
107 for pipeline in pipelines.all() {
108 report.pipeline(pipeline);
109 for waits_for in pipelines.waits_for(pipeline) {
110 report.depends(pipeline, *waits_for);
111 }
112 }
113 let building =
114 Building { plan, catalog, cancel, memory, report, pipelines, next_operator: Cell::new(0) };
115 building.node(plan.root())
116}
117
118struct Building<'a, 'b> {
124 plan: &'a Plan,
125 catalog: &'a Catalog,
126 cancel: &'b Cancel,
127 memory: &'b Memory,
128 report: &'b Report,
129 pipelines: Pipelines,
130 next_operator: Cell<u32>,
131}
132
133impl<'a> Building<'a, '_> {
134 fn id(&self) -> u32 {
137 let id = self.next_operator.get();
138 self.next_operator.set(id + 1);
139 id
140 }
141
142 fn watch(&self, id: u32, pipeline: u32, kind: &str, detail: Option<&str>) -> Arc<Counters> {
149 let counters = Counters::new(id, pipeline, kind).reference();
150 let counters = match detail {
151 Some(detail) => counters.detailed(detail),
152 None => counters,
153 };
154 self.report.watch(counters)
155 }
156
157 fn node(&self, reference: NodeRef) -> Result<Box<dyn Operator + 'a>> {
158 let plan = self.plan;
159 let memory = self.memory;
160 let pipeline = self.pipelines.pipeline(reference);
161 let inner: Box<dyn Operator + 'a> = match *plan.node(reference) {
162 Node::Get { catalog: database, schema, table, index, columns, .. } => {
163 let id = self.id();
164 let name = QualifiedName::new(
165 plan.string(database),
166 plan.string(schema),
167 plan.string(table),
168 );
169 let scan = Scan::new(plan, self.catalog.table(&name)?, index, columns)?;
170 let schema = scan.schema().clone();
171 let counters = self.watch(id, pipeline, "Scan", Some(plan.string(table)));
172 pulled(Watched::new(scan, counters), schema)
173 }
174 Node::Dummy => {
175 let id = self.id();
176 let dummy = Dummy::new();
177 let schema = dummy.schema().clone();
178 pulled(Watched::new(dummy, self.watch(id, pipeline, "Dummy", None)), schema)
179 }
180 Node::Values { index, columns, rows } => {
181 let id = self.id();
182 let values = Values::new(plan, index, columns, rows)?;
183 let schema = values.schema().clone();
184 pulled(Watched::new(values, self.watch(id, pipeline, "Values", None)), schema)
185 }
186 Node::TableFunction { index, function, args, options, settings, columns } => {
187 let id = self.id();
188 let name = plan.string(function);
189 match TableFunction::lookup(name) {
190 Some(function @ (TableFunction::ReadParquet | TableFunction::ReadCsv)) => {
191 let scan =
192 FileScan::new(plan, index, function, args, options, settings, columns)?;
193 let schema = scan.schema().clone();
194 let counters = self.watch(id, pipeline, "FileScan", Some(name));
195 pulled(Watched::new(scan, counters), schema)
196 }
197 Some(TableFunction::RudbStrategies) => {
198 let table = Strategies::new(plan, index, columns)?;
199 let schema = table.schema().clone();
200 let counters = self.watch(id, pipeline, "Strategies", None);
201 pulled(Watched::new(table, counters), schema)
202 }
203 _ => {
204 let series = Series::new(plan, index, name, args)?;
205 let schema = series.schema().clone();
206 let counters = self.watch(id, pipeline, "Series", Some(name));
207 pulled(Watched::new(series, counters), schema)
208 }
209 }
210 }
211 Node::Filter { input, predicate } => {
212 let id = self.id();
213 let input = self.node(input)?;
214 let schema = input.schema().clone();
215 let filter = Filter::new(plan, predicate, &schema)?;
216 let counters = self.watch(id, pipeline, "Filter", None);
217 Box::new(Streamed::new(input, Watched::new(filter, counters), schema))
218 }
219 Node::Project { input, index, exprs, names } => {
220 let id = self.id();
221 let input = self.node(input)?;
222 let project = Project::new(plan, input.schema(), index, exprs, names)?;
223 let schema = project.schema().clone();
224 let counters = self.watch(id, pipeline, "Project", None);
225 Box::new(Streamed::new(input, Watched::new(project, counters), schema))
226 }
227 Node::Aggregate { input, index, groups, aggregates } => {
228 let id = self.id();
229 let input = self.node(input)?;
230 let (aggregate, out) =
231 Aggregate::new(plan, input.schema(), index, groups, aggregates, memory)?;
232 let schema = aggregate.schema().clone();
233 let counters = self.watch(id, pipeline, "Aggregate", None);
234 Box::new(Broken::new(input, Watched::new(aggregate, counters), out, schema))
235 }
236 Node::Sort { input, keys } => {
237 let id = self.id();
238 let input = self.node(input)?;
239 let schema = input.schema().clone();
240 let (sort, out) = Sort::new(plan, &schema, keys, memory)?;
241 let counters = self.watch(id, pipeline, "Sort", None);
242 Box::new(Broken::new(input, Watched::new(sort, counters), out, schema))
243 }
244 Node::Limit { input, count, offset } => {
245 let id = self.id();
246 let input = self.node(input)?;
247 let schema = input.schema().clone();
248 let limit = Limit::new(count, offset);
249 let counters = self.watch(id, pipeline, "Limit", None);
250 Box::new(Streamed::new(input, Watched::new(limit, counters), schema))
251 }
252 Node::TopN { input, keys, count, offset } => {
253 let id = self.id();
254 let input = self.node(input)?;
255 let schema = input.schema().clone();
256 let (top, out) = TopN::new(plan, &schema, keys, count, offset, memory)?;
257 let counters = self.watch(id, pipeline, "TopN", None);
258 Box::new(Broken::new(input, Watched::new(top, counters), out, schema))
259 }
260 Node::Distinct { input, on } => {
261 let id = self.id();
262 let input = self.node(input)?;
263 let schema = input.schema().clone();
264 let (distinct, out) = Distinct::new(plan, &schema, on, memory)?;
265 let counters = self.watch(id, pipeline, "Distinct", None);
266 Box::new(Broken::new(input, Watched::new(distinct, counters), out, schema))
267 }
268 Node::Join { left, right, kind, conditions } => {
269 let id = self.id();
270 let gather_id = self.id();
271 let gathering = self.pipelines.pipeline(right);
277 let right = self.node(right)?;
278 let left = self.node(left)?;
279 let (gather, gathered) = Gather::new(memory);
280 let side = Gathered { schema: right.schema(), rows: gathered };
281 let (join, out) =
282 Join::new(plan, left.schema(), side, kind, conditions, self.cancel, memory);
283 let schema = join.schema().clone();
284 let kept = self.watch(gather_id, gathering, "Gather", None);
285 let counters = self.watch(id, pipeline, "Join", None);
286 Box::new(Paired::new(
287 right,
288 Watched::new(gather, kept),
289 left,
290 Watched::new(join, counters),
291 out,
292 schema,
293 ))
294 }
295 Node::CrossProduct { left, right } => {
296 let id = self.id();
297 let keep_id = self.id();
298 let aside = self.pipelines.pipeline(right);
303 let right = self.node(right)?;
304 let left = self.node(left)?;
305 let (keep, kept) = Keep::new(memory);
306 let cross = CrossProduct::new(left.schema(), right.schema(), kept);
307 let schema = cross.schema().clone();
308 let held = self.watch(keep_id, aside, "Keep", None);
309 let counters = self.watch(id, pipeline, "CrossProduct", None);
310 Box::new(Fed::new(
311 right,
312 Watched::new(keep, held),
313 Streamed::new(left, Watched::new(cross, counters), schema),
314 ))
315 }
316 Node::SetOp { left, right, kind, all, index } => {
317 let id = self.id();
318 let gather_id = self.id();
319 let counting = self.pipelines.pipeline(right);
322 let right = self.node(right)?;
323 let left = self.node(left)?;
324 let (gather, gathered) = Gather::new(memory);
325 let (setop, out) = SetOp::new(left.schema(), gathered, kind, all, index, memory);
326 let schema = setop.schema().clone();
327 let kept = self.watch(gather_id, counting, "Gather", None);
328 let counters = self.watch(id, pipeline, "SetOp", None);
329 Box::new(Paired::new(
330 right,
331 Watched::new(gather, kept),
332 left,
333 Watched::new(setop, counters),
334 out,
335 schema,
336 ))
337 }
338 };
339 Ok(Box::new(Guarded::new(inner, self.cancel.clone())))
340 }
341}
342
343fn pulled<'a, S: Source + 'a>(source: S, schema: Schema) -> Box<dyn Operator + 'a> {
350 Box::new(Pulled::new(source, schema))
351}