Skip to main content

reifydb_core/value/column/view/
group_by.rs

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