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	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}