1use std::cmp::Ordering;
22use std::collections::HashSet;
23
24use rudb_common::{Error, LogicalType, Result, Value};
25use rudb_vector::{Buffer, Data, Live, Validity, Vector, interleave};
26
27use crate::aggregate::{Accumulator, NOWHERE, finish_run, update_runs, update_scattered};
28use crate::compare::order;
29use crate::datetime;
30use crate::number::integral;
31
32const NULL_PICK: &str = "NULLs are not allowed as list elements in the second input parameter.";
34
35pub(crate) fn before_nulls(
37 name: &str,
38 args: &[Value],
39 returns: &LogicalType,
40) -> Option<Result<Value>> {
41 Some(match (name, args) {
42 ("list_position", [Value::Null, _]) => Ok(Value::Null),
43 ("list_position", [Value::List { values, .. }, needle]) => position(values, needle),
44 ("list_resize" | "list_intersect", [Value::Null, ..]) => Ok(Value::Null),
45 ("list_intersect", [Value::List { .. }, Value::Null]) => listed(Vec::new(), returns),
46 ("list_resize", [Value::List { values, .. }, size, filler @ ..]) => {
47 resize(values, size, filler.first().unwrap_or(&Value::Null), returns)
48 }
49 _ => return None,
50 })
51}
52
53pub(crate) fn value(name: &str, args: &[Value], returns: &LogicalType) -> Option<Result<Value>> {
55 let list = |values| listed(values, returns);
56 Some(match (name, args) {
57 ("list_contains", [Value::List { values, .. }, needle]) => {
58 found(values, needle).map(|at| Value::Boolean(at.is_some()))
59 }
60 ("list_has_any", [Value::List { values, .. }, Value::List { values: wanted, .. }]) => {
61 has_any(values, wanted).map(Value::Boolean)
62 }
63 ("list_has_all", [Value::List { values, .. }, Value::List { values: wanted, .. }]) => {
64 has_any_missing(values, wanted).map(|missing| Value::Boolean(!missing))
65 }
66 ("list_distinct", [Value::List { values, .. }]) => distinct(values).and_then(&list),
67 ("list_unique", [Value::List { values, .. }]) => {
68 distinct(values).map(|kept| Value::UBigInt(kept.len() as u64))
69 }
70 ("list_intersect", [Value::List { values, .. }, Value::List { values: other, .. }]) => {
71 intersect(values, other).and_then(&list)
72 }
73 ("list_where", [Value::List { values, .. }, Value::List { values: mask, .. }]) => {
74 masked(values, mask).and_then(&list)
75 }
76 ("list_select", [Value::List { values, .. }, Value::List { values: indexes, .. }]) => {
77 selected(values, indexes).and_then(&list)
78 }
79 ("list_sort", [Value::List { values, .. }, spelled @ ..]) => {
80 let order = spelled.first().map(spelled_order).transpose();
81 let nulls = spelled.get(1).map(spelled_nulls).transpose();
82 match (order, nulls) {
83 (Ok(order), Ok(nulls)) => {
84 sort(values, order.unwrap_or(false), nulls.unwrap_or(false)).and_then(list)
85 }
86 (Err(error), _) | (_, Err(error)) => Err(error),
87 }
88 }
89 ("range" | "generate_series", _) => ranged(name == "generate_series", args).and_then(list),
90 ("list_grade_up", [Value::List { values, .. }, spelled @ ..]) => {
91 let order = spelled.first().map(spelled_order).transpose();
92 let nulls = spelled.get(1).map(spelled_nulls).transpose();
93 match (order, nulls) {
94 (Ok(order), Ok(nulls)) => {
95 graded(values, order.unwrap_or(false), nulls.unwrap_or(false)).and_then(list)
96 }
97 (Err(error), _) | (_, Err(error)) => Err(error),
98 }
99 }
100 ("list_reverse_sort", [Value::List { values, .. }, spelled @ ..]) => {
101 match spelled.first().map(spelled_nulls).transpose() {
102 Ok(nulls) => sort(values, true, nulls.unwrap_or(false)).and_then(list),
103 Err(error) => Err(error),
104 }
105 }
106 ("list_reverse", [Value::List { values, .. }]) => {
107 list(values.iter().rev().cloned().collect())
108 }
109 ("flatten", [Value::List { values, .. }]) => {
110 let mut flat = Vec::new();
111 for inner in values {
112 if let Value::List { values: held, .. } = inner {
113 flat.extend(held.iter().cloned());
114 }
115 }
116 list(flat)
117 }
118 _ => return None,
119 })
120}
121
122fn listed(values: Vec<Value>, returns: &LogicalType) -> Result<Value> {
124 let LogicalType::List(element) = returns else {
125 return Err(Error::internal(format!("a list function returning {returns}")));
126 };
127 Ok(Value::List { element: (**element).clone(), values })
128}
129
130fn ranged(inclusive: bool, args: &[Value]) -> Result<Vec<Value>> {
136 if let [start, stop, Value::Interval { months, days, micros }] = args {
137 return stepped(inclusive, start, stop, (*months, *days, *micros));
138 }
139 let whole = |value: &Value| {
140 integral(value)
141 .and_then(|held| i64::try_from(held).ok())
142 .ok_or_else(|| Error::internal(format!("a range over a {}", value.logical_type())))
143 };
144 let (start, stop, step) = match args {
145 [stop] => (0, whole(stop)?, 1),
146 [start, stop] => (whole(start)?, whole(stop)?, 1),
147 [start, stop, step] => (whole(start)?, whole(stop)?, whole(step)?),
148 _ => return Err(Error::internal(format!("a range over {} arguments", args.len()))),
149 };
150 let count = series_length(start, stop, step, inclusive)?;
151 let mut at = start;
154 let mut values = Vec::with_capacity(count);
155 for _ in 0..count {
158 values.push(Value::BigInt(at));
159 at = at.wrapping_add(step);
160 }
161 Ok(values)
162}
163
164fn series_length(start: i64, stop: i64, step: i64, inclusive: bool) -> Result<usize> {
167 if step == 0 || (start > stop && step > 0) || (start < stop && step < 0) {
168 return Ok(0);
169 }
170 let apart = stop.abs_diff(start);
171 let by = step.unsigned_abs();
172 let (whole, over) =
174 if by == 1 { (apart, false) } else { (apart / by, !apart.is_multiple_of(by)) };
175 let count = u128::from(whole) + u128::from(inclusive || over);
176 usize::try_from(count).ok().filter(|&count| count <= MAX_SERIES).ok_or_else(too_long)
177}
178
179fn stepped(
182 inclusive: bool,
183 start: &Value,
184 stop: &Value,
185 interval: (i32, i32, i64),
186) -> Result<Vec<Value>> {
187 let moment = |value: &Value| match value {
188 Value::Timestamp(stamp) | Value::TimestampTz(stamp) => Ok(*stamp),
189 other => Err(Error::internal(format!("a range from a {}", other.logical_type()))),
190 };
191 let mut stamps = Vec::new();
192 step_stamps(inclusive, moment(start)?, moment(stop)?, interval, &mut stamps)?;
193 let zoned = matches!(start, Value::TimestampTz(_));
194 Ok(stamps
195 .into_iter()
196 .map(|stamp| if zoned { Value::TimestampTz(stamp) } else { Value::Timestamp(stamp) })
197 .collect())
198}
199
200fn step_stamps(
205 inclusive: bool,
206 start: i64,
207 end: i64,
208 (months, days, micros): (i32, i32, i64),
209 out: &mut Vec<i64>,
210) -> Result<usize> {
211 let forward = months > 0 || days > 0 || micros > 0;
212 let backward = months < 0 || days < 0 || micros < 0;
213 if forward && backward {
214 return Err(Error::invalid_input(
215 "Interval with mix of negative/positive entries not supported",
216 ));
217 }
218 if [start, end].iter().any(|&stamp| stamp == i64::MAX || stamp == -i64::MAX) {
220 return Err(Error::invalid_input("Interval infinite bounds not supported"));
221 }
222 let from = out.len();
223 if !forward && !backward {
224 return Ok(0);
225 }
226 let whole = i128::from(days) * i128::from(datetime::MICROS_PER_DAY) + i128::from(micros);
229 if let (0, Ok(step)) = (months, i64::try_from(whole)) {
230 let count = series_length(start, end, step, inclusive)?;
231 let steps = i64::try_from(count).map_err(|_| too_long())?;
232 out.reserve(count);
233 out.extend((0..steps).map(|at| start.wrapping_add(at.wrapping_mul(step))));
234 return Ok(count);
235 }
236 let (months, days, micros) = (i64::from(months), i64::from(days), i128::from(micros));
237 let mut at = start;
238 loop {
239 let past = if forward { at > end } else { at < end };
240 if past || (at == end && !inclusive) {
241 return Ok(out.len() - from);
242 }
243 if out.len() - from == MAX_SERIES {
244 return Err(too_long());
245 }
246 let next = datetime::shifted_stamp(at, months, days, micros)?;
247 out.push(at);
248 at = next;
249 }
250}
251
252#[derive(Debug, Clone, PartialEq, Eq)]
254pub enum Stepping {
255 Even {
257 start: i64,
259 step: i64,
261 count: usize,
263 },
264 Listed(Vec<i64>),
266}
267
268impl Stepping {
269 #[must_use]
271 pub fn len(&self) -> usize {
272 match self {
273 Self::Even { count, .. } => *count,
274 Self::Listed(stamps) => stamps.len(),
275 }
276 }
277
278 #[must_use]
280 pub fn is_empty(&self) -> bool {
281 self.len() == 0
282 }
283
284 #[must_use]
286 pub fn at(&self, position: usize) -> i64 {
287 match self {
288 Self::Even { start, step, .. } => {
289 let steps = i64::try_from(position).unwrap_or(i64::MAX);
290 start.wrapping_add(steps.wrapping_mul(*step))
291 }
292 Self::Listed(stamps) => stamps[position],
293 }
294 }
295}
296
297pub fn moment_steps(
304 inclusive: bool,
305 start: i64,
306 end: i64,
307 interval: (i32, i32, i64),
308) -> Result<Stepping> {
309 let (months, days, micros) = interval;
310 let whole = i128::from(days) * i128::from(datetime::MICROS_PER_DAY) + i128::from(micros);
311 if let (0, Ok(step)) = (months, i64::try_from(whole)) {
312 return Ok(Stepping::Even {
313 start,
314 step,
315 count: series_length(start, end, step, inclusive)?,
316 });
317 }
318 let mut stamps = Vec::new();
319 step_stamps(inclusive, start, end, interval, &mut stamps)?;
320 Ok(Stepping::Listed(stamps))
321}
322
323const MAX_SERIES: usize = u32::MAX as usize;
325
326fn too_long() -> Error {
328 Error::invalid_input("Lists larger than 2^32 elements are not supported")
329}
330
331fn same(left: &Value, right: &Value) -> Result<bool> {
333 Ok(match (left.is_null(), right.is_null()) {
334 (true, true) => true,
335 (true, false) | (false, true) => false,
336 (false, false) => order(left, right)? == Ordering::Equal,
337 })
338}
339
340fn found(values: &[Value], needle: &Value) -> Result<Option<usize>> {
342 for (at, value) in values.iter().enumerate() {
343 if same(value, needle)? {
344 return Ok(Some(at));
345 }
346 }
347 Ok(None)
348}
349
350fn position(values: &[Value], needle: &Value) -> Result<Value> {
352 Ok(match found(values, needle)? {
353 Some(at) => Value::Integer(i32::try_from(at + 1).map_err(|_| {
354 Error::out_of_range(format!("a list position of {} does not fit in INTEGER", at + 1))
355 })?),
356 None => Value::Null,
357 })
358}
359
360fn sorted(values: &[Value]) -> Result<Vec<&Value>> {
362 let mut held: Vec<&Value> = values.iter().filter(|value| !value.is_null()).collect();
363 let mut failed = None;
364 held.sort_by(|left, right| {
365 order(left, right).unwrap_or_else(|error| {
366 failed.get_or_insert(error);
367 Ordering::Equal
368 })
369 });
370 match failed {
371 Some(error) => Err(error),
372 None => Ok(held),
373 }
374}
375
376fn contains(haystack: &[&Value], needle: &Value) -> Result<bool> {
378 let mut failed = None;
379 let hit = haystack
380 .binary_search_by(|probe| {
381 order(probe, needle).unwrap_or_else(|error| {
382 failed.get_or_insert(error);
383 Ordering::Equal
384 })
385 })
386 .is_ok();
387 match failed {
388 Some(error) => Err(error),
389 None => Ok(hit),
390 }
391}
392
393fn has_any(values: &[Value], wanted: &[Value]) -> Result<bool> {
395 let haystack = sorted(values)?;
396 for value in wanted.iter().filter(|value| !value.is_null()) {
397 if contains(&haystack, value)? {
398 return Ok(true);
399 }
400 }
401 Ok(false)
402}
403
404fn has_any_missing(values: &[Value], wanted: &[Value]) -> Result<bool> {
407 let haystack = sorted(values)?;
408 for value in wanted.iter().filter(|value| !value.is_null()) {
409 if !contains(&haystack, value)? {
410 return Ok(true);
411 }
412 }
413 Ok(false)
414}
415
416fn distinct(values: &[Value]) -> Result<Vec<Value>> {
418 let mut at: Vec<usize> = (0..values.len()).filter(|&at| !values[at].is_null()).collect();
419 let mut failed = None;
420 at.sort_by(|&left, &right| match order(&values[left], &values[right]) {
423 Ok(Ordering::Equal) => left.cmp(&right),
424 Ok(ordering) => ordering,
425 Err(error) => {
426 failed.get_or_insert(error);
427 Ordering::Equal
428 }
429 });
430 if let Some(error) = failed {
431 return Err(error);
432 }
433 let mut kept = Vec::with_capacity(at.len());
434 for (index, &here) in at.iter().enumerate() {
435 if index == 0 || order(&values[at[index - 1]], &values[here])? != Ordering::Equal {
436 kept.push(here);
437 }
438 }
439 kept.sort_unstable();
440 Ok(kept.into_iter().map(|at| values[at].clone()).collect())
441}
442
443fn intersect(values: &[Value], other: &[Value]) -> Result<Vec<Value>> {
445 let haystack = sorted(other)?;
446 let mut kept = Vec::new();
447 for value in distinct(values)? {
448 if contains(&haystack, &value)? {
449 kept.push(value);
450 }
451 }
452 Ok(kept)
453}
454
455fn masked(values: &[Value], mask: &[Value]) -> Result<Vec<Value>> {
458 let mut kept = Vec::new();
459 for (at, flag) in mask.iter().enumerate() {
460 match flag {
461 Value::Boolean(true) => kept.push(values.get(at).cloned().unwrap_or(Value::Null)),
462 Value::Boolean(false) => {}
463 Value::Null => return Err(Error::invalid_input(NULL_PICK)),
464 other => {
465 return Err(Error::internal(format!(
466 "list_where with a {} mask",
467 other.logical_type()
468 )));
469 }
470 }
471 }
472 Ok(kept)
473}
474
475fn selected(values: &[Value], indexes: &[Value]) -> Result<Vec<Value>> {
478 let mut kept = Vec::with_capacity(indexes.len());
479 for index in indexes {
480 if index.is_null() {
481 return Err(Error::invalid_input(NULL_PICK));
482 }
483 let picked = index
484 .as_i64()
485 .and_then(|index| usize::try_from(index).ok())
486 .and_then(|index| index.checked_sub(1))
487 .and_then(|at| values.get(at));
488 kept.push(picked.cloned().unwrap_or(Value::Null));
489 }
490 Ok(kept)
491}
492
493fn resize(values: &[Value], size: &Value, filler: &Value, returns: &LogicalType) -> Result<Value> {
496 let size = match size {
497 Value::Null => 0,
498 Value::UBigInt(size) => usize::try_from(*size).map_err(|_| {
499 Error::out_of_range(format!("a list of {size} elements is too long to build"))
500 })?,
501 other => {
502 return Err(Error::internal(format!("list_resize to a {}", other.logical_type())));
503 }
504 };
505 let element = match returns {
506 LogicalType::List(element) => element,
507 _ => &LogicalType::Null,
508 };
509 let width = element.physical().size().max(1);
510 if (size as u128) * (width as u128) > MAX_VECTOR_BYTES {
511 return Err(Error::out_of_range(format!(
512 "Cannot resize vector to {size} rows: maximum allowed vector size is 128.0 GiB"
513 )));
514 }
515 let mut kept: Vec<Value> = values.iter().take(size).cloned().collect();
516 kept.resize(size, filler.clone());
517 listed(kept, returns)
518}
519
520const MAX_VECTOR_BYTES: u128 = 1 << 37;
523
524fn spelled_order(spelled: &Value) -> Result<bool> {
526 let spelled = spelled.to_string().to_uppercase();
527 match spelled.as_str() {
528 "ASC" | "ASCENDING" | "DEFAULT" | "ORDER_DEFAULT" => Ok(false),
529 "DESC" | "DESCENDING" => Ok(true),
530 _ => Err(unrecognized(&spelled, "OrderType")),
531 }
532}
533
534fn spelled_nulls(spelled: &Value) -> Result<bool> {
536 let spelled = spelled.to_string().to_uppercase();
537 match spelled.as_str() {
538 "NULLS FIRST" | "NULLS_FIRST" => Ok(true),
539 "NULLS LAST" | "NULLS_LAST" | "DEFAULT" | "ORDER_DEFAULT" => Ok(false),
540 _ => Err(unrecognized(&spelled, "OrderByNullType")),
541 }
542}
543
544fn unrecognized(spelled: &str, kind: &str) -> Error {
547 Error::not_implemented(format!(
548 "Enum value: unrecognized value \"{spelled}\" for enum \"{kind}\""
549 ))
550}
551
552fn sort(values: &[Value], descending: bool, nulls_first: bool) -> Result<Vec<Value>> {
554 Ok(grade(values, descending, nulls_first)?.into_iter().map(|at| values[at].clone()).collect())
555}
556
557fn graded(values: &[Value], descending: bool, nulls_first: bool) -> Result<Vec<Value>> {
559 grade(values, descending, nulls_first)?
560 .into_iter()
561 .map(|at| {
562 Ok(Value::BigInt(
563 i64::try_from(at + 1).map_err(|error| Error::internal(error.to_string()))?,
564 ))
565 })
566 .collect()
567}
568
569fn grade(values: &[Value], descending: bool, nulls_first: bool) -> Result<Vec<usize>> {
575 let (mut held, nulls): (Vec<usize>, Vec<usize>) =
576 (0..values.len()).partition(|&at| !values[at].is_null());
577 let mut failed = None;
578 held.sort_by(|&left, &right| {
579 let ordering = order(&values[left], &values[right]).unwrap_or_else(|error| {
580 failed.get_or_insert(error);
581 Ordering::Equal
582 });
583 if descending { ordering.reverse() } else { ordering }
584 });
585 if let Some(error) = failed {
586 return Err(error);
587 }
588 Ok(if nulls_first { [nulls, held].concat() } else { [held, nulls].concat() })
589}
590
591pub(crate) fn vectorized<V: AsRef<Vector>>(
599 name: &str,
600 args: &[V],
601 returns: &LogicalType,
602 rows: usize,
603) -> Result<Option<Vector>> {
604 match (name, args) {
605 ("list_value", [_, ..]) => built(args, returns, rows),
606 ("range" | "generate_series", [_, ..]) => {
607 series(name == "generate_series", args, returns, rows)
608 }
609 ("list_aggr", [list, aggregate]) => aggregated(list.as_ref(), aggregate.as_ref(), returns),
610 ("list_reverse", [list]) => reversed(list.as_ref()),
611 ("length" | "array_length", [list]) => counted(list.as_ref()),
612 ("list_distinct", [list]) => deduplicated(false, list.as_ref()),
613 ("list_unique", [list]) => deduplicated(true, list.as_ref()),
614 ("list_contains" | "list_position", [list, needle]) => {
615 searched(name == "list_position", list.as_ref(), needle.as_ref())
616 }
617 ("list_sort" | "list_grade_up", [list, spelled @ ..]) => ordered(
618 list.as_ref(),
619 spelled.first().map(AsRef::as_ref),
620 spelled.get(1).map(AsRef::as_ref),
621 false,
622 name == "list_grade_up",
623 ),
624 ("list_reverse_sort", [list, spelled @ ..]) => {
625 ordered(list.as_ref(), None, spelled.first().map(AsRef::as_ref), true, false)
626 }
627 _ => Ok(None),
628 }
629}
630
631fn built<V: AsRef<Vector>>(
636 args: &[V],
637 returns: &LogicalType,
638 rows: usize,
639) -> Result<Option<Vector>> {
640 let LogicalType::List(element) = returns else {
641 return Ok(None);
642 };
643 if nested_or_null(element) || args.iter().any(|arg| arg.as_ref().logical_type() != &**element) {
644 return Ok(None);
645 }
646 let pieces = args.iter().map(|arg| arg.as_ref().flatten()).collect::<Result<Vec<_>>>()?;
649 let width = args.len();
650 let order: Vec<usize> =
651 (0..rows).flat_map(|row| (0..width).map(move |at| at * rows + row)).collect();
652 let child = interleave(element, &pieces, &order)?;
653 let count = entry(width)?;
654 let entries = (0..rows).map(|row| Ok((entry(row * width)?, count))).collect::<Result<_>>()?;
655 Vector::list(entries, child).map(Some)
656}
657
658fn series<V: AsRef<Vector>>(
661 inclusive: bool,
662 args: &[V],
663 returns: &LogicalType,
664 rows: usize,
665) -> Result<Option<Vector>> {
666 if let [start, stop, step] = args
667 && step.as_ref().logical_type() == &LogicalType::Interval
668 {
669 return timed(inclusive, [start.as_ref(), stop.as_ref(), step.as_ref()], returns, rows);
670 }
671 if args.iter().any(|arg| arg.as_ref().logical_type() != &LogicalType::BigInt) {
672 return Ok(None);
673 }
674 let flat: Vec<Vector> = args.iter().map(|arg| arg.as_ref().flatten()).collect::<Result<_>>()?;
677 let mut columns = Vec::with_capacity(flat.len());
678 for vector in &flat {
679 let Some(Data::Int64(values)) = vector.data() else {
680 return Ok(None);
681 };
682 columns.push((values.as_slice(), vector.validity().live()));
683 }
684 let mut runs = Vec::with_capacity(rows);
687 let mut live = vec![true; rows];
688 let mut total = 0_usize;
689 for (row, live) in live.iter_mut().enumerate() {
690 let mut held = [0_i64; 3];
691 for (at, (values, valid)) in columns.iter().enumerate() {
692 match values.get(row) {
693 Some(&value) if valid.at(row) => held[at] = value,
694 _ => *live = false,
695 }
696 }
697 let (start, stop, step) = match columns.len() {
698 1 => (0, held[0], 1),
699 2 => (held[0], held[1], 1),
700 _ => (held[0], held[1], held[2]),
701 };
702 let count = if *live { series_length(start, stop, step, inclusive)? } else { 0 };
703 runs.push((start, step, count));
704 total += count;
705 }
706 let mut entries = Vec::with_capacity(rows);
707 let mut child = Vec::with_capacity(total);
708 for &(start, step, count) in &runs {
709 entries.push((entry(child.len())?, entry(count)?));
710 let steps = i64::try_from(count).map_err(|_| too_long())?;
713 child.extend((0..steps).map(|at| start.wrapping_add(at.wrapping_mul(step))));
714 }
715 let child = Vector::flat(LogicalType::BigInt, Data::Int64(Buffer::from(child)))?;
716 let validity = Validity::from_iter(rows, |row| live[row]).normalize(rows);
717 Ok(Some(Vector::list(entries, child)?.with_validity(validity)))
718}
719
720fn timed(
723 inclusive: bool,
724 args: [&Vector; 3],
725 returns: &LogicalType,
726 rows: usize,
727) -> Result<Option<Vector>> {
728 let LogicalType::List(element) = returns else {
729 return Ok(None);
730 };
731 if args[..2].iter().any(|arg| arg.logical_type() != &**element) {
732 return Ok(None);
733 }
734 let [start, stop, step] = args.map(Vector::flatten);
735 let (start, stop, step) = (start?, stop?, step?);
736 let (Some(Data::Int64(starts)), Some(Data::Int64(stops)), Some(Data::Interval(steps))) =
737 (start.data(), stop.data(), step.data())
738 else {
739 return Ok(None);
740 };
741 let (starts, stops, steps) = (starts.as_slice(), stops.as_slice(), steps.as_slice());
742 let lives = [start.validity().live(), stop.validity().live(), step.validity().live()];
743 let mut entries = Vec::with_capacity(rows);
744 let mut child = Vec::new();
745 let mut live = vec![true; rows];
746 for row in 0..rows {
747 let at = entry(child.len())?;
748 let held = (starts.get(row), stops.get(row), steps.get(row));
749 let (Some(&from), Some(&to), Some(&interval)) = held else {
750 return Ok(None);
751 };
752 if lives.iter().any(|live| !live.at(row)) {
753 entries.push((at, 0));
754 live[row] = false;
755 continue;
756 }
757 let count = step_stamps(inclusive, from, to, interval, &mut child)?;
758 entries.push((at, entry(count)?));
759 }
760 let child = Vector::flat((**element).clone(), Data::Int64(Buffer::from(child)))?;
761 let validity = Validity::from_iter(rows, |row| live[row]).normalize(rows);
762 Ok(Some(Vector::list(entries, child)?.with_validity(validity)))
763}
764
765fn counted(list: &Vector) -> Result<Option<Vector>> {
767 let Some((entries, _)) = list.list_parts() else {
768 return Ok(None);
769 };
770 if !matches!(list.logical_type(), LogicalType::List(_)) {
771 return Ok(None);
772 }
773 let lengths: Vec<i64> = entries.iter().map(|&(_, len)| i64::from(len)).collect();
774 let answer = Vector::flat(LogicalType::BigInt, Data::Int64(Buffer::from(lengths)))?;
775 Ok(Some(answer.with_validity(list.validity().clone())))
776}
777
778fn aggregated(list: &Vector, aggregate: &Vector, returns: &LogicalType) -> Result<Option<Vector>> {
787 let (Some((entries, child)), Some(Value::Varchar(aggregate))) =
788 (list.list_parts(), aggregate.constant_value())
789 else {
790 return Ok(None);
791 };
792 let live = list.validity().live();
793 let mut slots = vec![NOWHERE; child.len()];
794 for (row, &(start, len)) in entries.iter().enumerate() {
795 if !live.at(row) {
796 continue;
797 }
798 let (start, len) = (start as usize, len as usize);
799 let Some(held) = slots.get_mut(start..start + len) else {
800 return Err(Error::internal("a list entry past the end of its child"));
801 };
802 if held.iter().any(|&slot| slot != NOWHERE) {
803 return Ok(None);
804 }
805 held.fill(row);
806 }
807 let rows = entries.len();
808 let mut states = vec![Accumulator::new(aggregate, returns)?; rows];
809 let mut runs: Vec<(usize, usize)> = Vec::new();
810 for (at, &slot) in slots.iter().enumerate() {
811 match runs.last_mut() {
812 Some((held, end)) if *held == slot => *end = at + 1,
813 _ => runs.push((slot, at + 1)),
814 }
815 }
816 if !update_runs(&mut states, &runs, 1, 0, Some(child), child.len())? {
817 update_scattered(&mut states, &slots, 1, 0, Some(child), child.len())?;
818 }
819 let every: Vec<usize> = (0..rows).collect();
820 let answer = match finish_run(&states, &every, 1, 0, returns)? {
821 Some(answer) => answer,
822 None => {
823 let values = states.iter().map(Accumulator::finish).collect::<Result<Vec<_>>>()?;
824 Vector::from_values(returns.clone(), &values)?
825 }
826 };
827 let validity = answer.validity().and(list.validity(), rows);
828 Ok(Some(answer.with_validity(validity)))
829}
830
831fn reversed(list: &Vector) -> Result<Option<Vector>> {
833 let Some((entries, child)) = list.list_parts() else {
834 return Ok(None);
835 };
836 if !matches!(list.logical_type(), LogicalType::List(_)) {
837 return Ok(None);
838 }
839 let live = list.validity().live();
840 let mut indices = Vec::with_capacity(child.len());
841 let mut placed = Vec::with_capacity(entries.len());
842 for (row, &(start, len)) in entries.iter().enumerate() {
843 let at = entry(indices.len())?;
844 if live.at(row) {
845 indices.extend((start..start + len).rev());
846 placed.push((at, len));
847 } else {
848 placed.push((at, 0));
849 }
850 }
851 let child = child.gather(&indices)?;
852 Ok(Some(Vector::list(placed, child)?.with_validity(list.validity().clone())))
853}
854
855fn searched(position: bool, list: &Vector, needle: &Vector) -> Result<Option<Vector>> {
862 let (Some((entries, child)), Some(wanted)) = (list.list_parts(), needle.constant_value())
863 else {
864 return Ok(None);
865 };
866 let plain = plain(child.logical_type());
867 let Some(wanted) =
868 integral(wanted).filter(|_| plain && needle.logical_type() == child.logical_type())
869 else {
870 return Ok(None);
871 };
872 let elements = child.validity().live();
873 macro_rules! scan {
874 ($($variant:ident),+) => {
875 match child.data() {
876 $(Some(Data::$variant(values)) => {
877 first_places(entries, values.as_slice(), elements, wanted)
878 })+
879 _ => return Ok(None),
880 }
881 };
882 }
883 let found = scan!(Int8, Int16, Int32, Int64, UInt8, UInt16, UInt32, UInt64);
884 let rows = list.validity().live();
885 if position {
886 let validity = Validity::from_iter(entries.len(), |row| rows.at(row) && found[row] != 0);
887 let data =
888 Data::Int32(Buffer::from(found.iter().map(|&place| place as i32).collect::<Vec<_>>()));
889 let answer = Vector::flat(LogicalType::Integer, data)?;
890 return Ok(Some(answer.with_validity(validity.normalize(entries.len()))));
891 }
892 let data = Data::Bool(Buffer::from(found.iter().map(|&place| place != 0).collect::<Vec<_>>()));
893 let answer = Vector::flat(LogicalType::Boolean, data)?;
894 Ok(Some(answer.with_validity(list.validity().clone())))
895}
896
897fn ordered(
906 list: &Vector,
907 order: Option<&Vector>,
908 nulls: Option<&Vector>,
909 reverse: bool,
910 grade: bool,
911) -> Result<Option<Vector>> {
912 let Some((entries, child)) = list.list_parts() else {
913 return Ok(None);
914 };
915 if !plain(child.logical_type()) || list.validity().count_valid(list.len()) == 0 {
916 return Ok(None);
917 }
918 let spelled = |arg: Option<&Vector>| match arg.map(Vector::constant_value) {
919 None => Some(None),
920 Some(Some(value @ Value::Varchar(_))) => Some(Some(value.clone())),
921 Some(_) => None,
922 };
923 let (Some(order), Some(nulls)) = (spelled(order), spelled(nulls)) else {
924 return Ok(None);
925 };
926 let descending = reverse || order.as_ref().map(spelled_order).transpose()?.unwrap_or(false);
927 let nulls_first = nulls.as_ref().map(spelled_nulls).transpose()?.unwrap_or(false);
928 let rows = list.validity().live();
929 let elements = child.validity().live();
930 macro_rules! permute {
931 ($($variant:ident),+) => {
932 match child.data() {
933 $(Some(Data::$variant(values)) => {
934 let values = values.as_slice();
935 permutation(entries, rows, elements, nulls_first, |left, right| {
936 let ordering = values[left as usize].cmp(&values[right as usize]);
937 if descending { ordering.reverse() } else { ordering }
938 })?
939 })+
940 _ => return Ok(None),
941 }
942 };
943 }
944 let (placed, indices) = permute!(Int8, Int16, Int32, Int64, UInt8, UInt16, UInt32, UInt64);
945 let child = if grade {
946 let mut places = Vec::with_capacity(indices.len());
948 for (&(at, len), &(start, _)) in placed.iter().zip(entries) {
949 let run = &indices[at as usize..(at + len) as usize];
950 places.extend(run.iter().map(|&index| i64::from(index - start) + 1));
951 }
952 Vector::flat(LogicalType::BigInt, Data::Int64(Buffer::from(places)))?
953 } else {
954 child.gather(&indices)?
955 };
956 Ok(Some(Vector::list(placed, child)?.with_validity(list.validity().clone())))
957}
958
959fn deduplicated(unique: bool, list: &Vector) -> Result<Option<Vector>> {
962 let Some((entries, child)) = list.list_parts() else {
963 return Ok(None);
964 };
965 if !plain(child.logical_type()) {
966 return Ok(None);
967 }
968 let rows = list.validity().live();
969 let elements = child.validity().live();
970 macro_rules! keep {
971 ($($variant:ident),+) => {
972 match child.data() {
973 $(Some(Data::$variant(values)) => {
974 let values = values.as_slice();
975 firsts(entries, rows, elements, |at| i128::from(values[at as usize]))?
976 })+
977 _ => return Ok(None),
978 }
979 };
980 }
981 let (placed, indices) = keep!(Int8, Int16, Int32, Int64, UInt8, UInt16, UInt32, UInt64);
982 if unique {
983 let counts: Vec<u64> = placed.iter().map(|&(_, len)| u64::from(len)).collect();
984 let answer = Vector::flat(LogicalType::UBigInt, Data::UInt64(Buffer::from(counts)))?;
985 return Ok(Some(answer.with_validity(list.validity().clone())));
986 }
987 let child = child.gather(&indices)?;
988 Ok(Some(Vector::list(placed, child)?.with_validity(list.validity().clone())))
989}
990
991fn firsts(
997 entries: &[(u32, u32)],
998 rows: Live<'_>,
999 elements: Live<'_>,
1000 key: impl Fn(u32) -> i128,
1001) -> Result<Permuted> {
1002 const SHORT: u32 = 32;
1003 let mut indices = Vec::new();
1004 let mut placed = Vec::with_capacity(entries.len());
1005 let mut kept: Vec<i128> = Vec::new();
1006 let mut seen: HashSet<i128> = HashSet::new();
1007 for (row, &(start, len)) in entries.iter().enumerate() {
1008 let at = entry(indices.len())?;
1009 if !rows.at(row) {
1010 placed.push((at, 0));
1011 continue;
1012 }
1013 let from = indices.len();
1014 kept.clear();
1015 seen.clear();
1016 for index in start..start + len {
1017 if !elements.at(index as usize) {
1018 continue;
1019 }
1020 let value = key(index);
1021 let fresh = if len <= SHORT {
1022 let fresh = !kept.contains(&value);
1023 if fresh {
1024 kept.push(value);
1025 }
1026 fresh
1027 } else {
1028 seen.insert(value)
1029 };
1030 if fresh {
1031 indices.push(index);
1032 }
1033 }
1034 placed.push((at, entry(indices.len() - from)?));
1035 }
1036 Ok((placed, indices))
1037}
1038
1039type Permuted = (Vec<(u32, u32)>, Vec<u32>);
1042
1043fn permutation(
1048 entries: &[(u32, u32)],
1049 rows: Live<'_>,
1050 elements: Live<'_>,
1051 nulls_first: bool,
1052 compare: impl Fn(u32, u32) -> Ordering,
1053) -> Result<Permuted> {
1054 let mut indices = Vec::new();
1055 let mut placed = Vec::with_capacity(entries.len());
1056 let mut nulls = Vec::new();
1057 for (row, &(start, len)) in entries.iter().enumerate() {
1058 let at = entry(indices.len())?;
1059 if !rows.at(row) {
1060 placed.push((at, 0));
1061 continue;
1062 }
1063 nulls.clear();
1064 let from = indices.len();
1065 for index in start..start + len {
1066 if elements.at(index as usize) {
1067 indices.push(index);
1068 } else {
1069 nulls.push(index);
1070 }
1071 }
1072 indices[from..].sort_by(|&left, &right| compare(left, right));
1073 if nulls_first {
1074 indices.splice(from..from, nulls.iter().copied());
1075 } else {
1076 indices.extend_from_slice(&nulls);
1077 }
1078 placed.push((at, len));
1079 }
1080 Ok((placed, indices))
1081}
1082
1083fn plain(ty: &LogicalType) -> bool {
1086 matches!(
1087 ty,
1088 LogicalType::TinyInt
1089 | LogicalType::SmallInt
1090 | LogicalType::Integer
1091 | LogicalType::BigInt
1092 | LogicalType::UTinyInt
1093 | LogicalType::USmallInt
1094 | LogicalType::UInteger
1095 | LogicalType::UBigInt
1096 )
1097}
1098
1099fn first_places<T: Copy + PartialEq + TryFrom<i128>>(
1106 entries: &[(u32, u32)],
1107 values: &[T],
1108 elements: Live<'_>,
1109 wanted: i128,
1110) -> Vec<u32> {
1111 let Ok(wanted) = T::try_from(wanted) else {
1112 return vec![0; entries.len()];
1113 };
1114 let place = |at: Option<usize>| at.map_or(0, |at| at as u32 + 1);
1115 entries
1116 .iter()
1117 .map(|&(start, len)| {
1118 let start = start as usize;
1119 let run = &values[start..start + len as usize];
1120 match elements {
1121 Live::All => place(run.iter().position(|&value| value == wanted)),
1122 _ => place(
1123 run.iter()
1124 .enumerate()
1125 .position(|(at, &value)| value == wanted && elements.at(start + at)),
1126 ),
1127 }
1128 })
1129 .collect()
1130}
1131
1132fn nested_or_null(element: &LogicalType) -> bool {
1134 matches!(
1135 element,
1136 LogicalType::Null
1137 | LogicalType::List(_)
1138 | LogicalType::Array(..)
1139 | LogicalType::Struct(_)
1140 | LogicalType::Map(..)
1141 | LogicalType::Union(_)
1142 )
1143}
1144
1145fn entry(offset: usize) -> Result<u32> {
1147 u32::try_from(offset)
1148 .map_err(|_| Error::out_of_range(format!("a list child of {offset} elements")))
1149}