reifydb_core/value/column/view/
group_by.rs1use std::{
5 iter::{Enumerate, FilterMap},
6 vec::IntoIter as VecIntoIter,
7};
8
9use indexmap::IndexMap;
10use reifydb_codec::key::{encoded::EncodedKey, serializer::KeySerializer};
11use reifydb_value::{Result, error::Error, value::Value};
12
13use crate::{
14 error::CoreError,
15 value::column::{ColumnBuffer, columns::Columns},
16};
17
18pub type GroupKey = Vec<Value>;
19
20#[repr(transparent)]
21#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
22pub struct GroupId(pub u32);
23
24impl GroupId {
25 pub fn index(self) -> usize {
26 self.0 as usize
27 }
28}
29
30pub type GroupRows = Vec<(GroupId, Vec<usize>)>;
31
32#[derive(Debug, Clone)]
33pub struct GroupSlots<T> {
34 slots: Vec<Option<T>>,
35 occupied: usize,
36}
37
38impl<T> Default for GroupSlots<T> {
39 fn default() -> Self {
40 Self::new()
41 }
42}
43
44impl<T> GroupSlots<T> {
45 pub fn new() -> Self {
46 Self {
47 slots: Vec::new(),
48 occupied: 0,
49 }
50 }
51
52 pub fn len(&self) -> usize {
53 self.occupied
54 }
55
56 pub fn is_empty(&self) -> bool {
57 self.occupied == 0
58 }
59
60 pub fn get(&self, group: GroupId) -> Option<&T> {
61 self.slots.get(group.index()).and_then(Option::as_ref)
62 }
63
64 pub fn insert(&mut self, group: GroupId, value: T) {
65 self.reserve_for(group);
66 if self.slots[group.index()].replace(value).is_none() {
67 self.occupied += 1;
68 }
69 }
70
71 pub fn get_or_insert_with(&mut self, group: GroupId, default: impl FnOnce() -> T) -> &mut T {
72 self.reserve_for(group);
73 if self.slots[group.index()].is_none() {
74 self.occupied += 1;
75 }
76 self.slots[group.index()].get_or_insert_with(default)
77 }
78
79 pub fn or_insert(&mut self, group: GroupId, default: T) -> &mut T {
80 self.get_or_insert_with(group, || default)
81 }
82
83 pub fn remove(&mut self, group: GroupId) -> Option<T> {
84 let removed = self.slots.get_mut(group.index()).and_then(Option::take);
85 if removed.is_some() {
86 self.occupied -= 1;
87 }
88 removed
89 }
90
91 pub fn iter(&self) -> impl Iterator<Item = (GroupId, &T)> {
92 self.slots
93 .iter()
94 .enumerate()
95 .filter_map(|(index, slot)| slot.as_ref().map(|value| (GroupId(index as u32), value)))
96 }
97
98 pub fn drain(&mut self) -> impl Iterator<Item = (GroupId, T)> + '_ {
99 self.occupied = 0;
100 self.slots
101 .drain(..)
102 .enumerate()
103 .filter_map(|(index, slot)| slot.map(|value| (GroupId(index as u32), value)))
104 }
105
106 fn reserve_for(&mut self, group: GroupId) {
107 if self.slots.len() <= group.index() {
108 self.slots.resize_with(group.index() + 1, || None);
109 }
110 }
111}
112
113fn occupied_slot<T>((index, slot): (usize, Option<T>)) -> Option<(GroupId, T)> {
114 slot.map(|value| (GroupId(index as u32), value))
115}
116
117impl<T> IntoIterator for GroupSlots<T> {
118 type Item = (GroupId, T);
119 type IntoIter = FilterMap<Enumerate<VecIntoIter<Option<T>>>, fn((usize, Option<T>)) -> Option<(GroupId, T)>>;
120
121 fn into_iter(self) -> Self::IntoIter {
122 self.slots
123 .into_iter()
124 .enumerate()
125 .filter_map(occupied_slot as fn((usize, Option<T>)) -> Option<(GroupId, T)>)
126 }
127}
128
129#[derive(Debug, Default, Clone)]
130pub struct GroupKeyDict {
131 entries: IndexMap<EncodedKey, GroupKey>,
132}
133
134impl GroupKeyDict {
135 pub fn new() -> Self {
136 Self {
137 entries: IndexMap::new(),
138 }
139 }
140
141 pub fn len(&self) -> usize {
142 self.entries.len()
143 }
144
145 pub fn is_empty(&self) -> bool {
146 self.entries.is_empty()
147 }
148
149 pub fn values(&self, group: GroupId) -> Option<&GroupKey> {
150 self.entries.get_index(group.index()).map(|(_, values)| values)
151 }
152
153 pub fn iter(&self) -> impl Iterator<Item = (GroupId, &GroupKey)> {
154 self.entries.values().enumerate().map(|(index, values)| (GroupId(index as u32), values))
155 }
156
157 fn intern(&mut self, encoded: &EncodedKey, materialize: impl FnOnce() -> GroupKey) -> GroupId {
158 if let Some(index) = self.entries.get_index_of(encoded) {
159 return GroupId(index as u32);
160 }
161 let (index, _) = self.entries.insert_full(encoded.clone(), materialize());
162 GroupId(index as u32)
163 }
164}
165
166impl Columns {
167 pub fn group_by_ids(&self, keys: &[&str], dict: &mut GroupKeyDict) -> Result<GroupRows> {
168 let row_count = self.columns.first().map_or(0, |c| c.len());
169 let key_columns = self.key_columns(keys)?;
170
171 let mut rows_by_group: IndexMap<GroupId, Vec<usize>> = IndexMap::new();
172
173 for row in 0..row_count {
174 let mut serializer = KeySerializer::new();
175 for column in &key_columns {
176 column.extend_key(row, &mut serializer);
177 }
178 let encoded = serializer.to_encoded_key();
179
180 let group = dict
181 .intern(&encoded, || key_columns.iter().map(|column| column.get_value(row)).collect());
182 rows_by_group.entry(group).or_default().push(row);
183 }
184
185 Ok(rows_by_group.into_iter().collect())
186 }
187
188 fn key_columns(&self, keys: &[&str]) -> Result<Vec<&ColumnBuffer>> {
189 let mut key_columns: Vec<&ColumnBuffer> = Vec::with_capacity(keys.len());
190 for &key in keys {
191 let pos = self.names.iter().position(|n| n.text() == key).ok_or_else(|| {
192 Error::from(CoreError::FrameError {
193 message: format!("Column '{}' not found", key),
194 })
195 })?;
196 key_columns.push(&self.columns[pos]);
197 }
198 Ok(key_columns)
199 }
200}