1use std::sync::Arc;
10use uqa_core::{
11 memory::{
12 Budgeted, BudgetedVec, MemoryBudget, MemoryError, MemoryReservation, Produced,
13 ProductionControl,
14 },
15 CancellationToken, QueryCancelled, Value, ValueRetentionError,
16};
17
18use crate::ast::{ColumnDef, ColumnType, Expr, FunctionBinding};
19
20mod columns;
21mod expressions;
22
23#[derive(Debug, thiserror::Error)]
24pub enum CatalogRetentionError {
25 #[error(transparent)]
26 Memory(#[from] MemoryError),
27 #[error(transparent)]
28 Cancelled(#[from] QueryCancelled),
29 #[error("validated column expression contains a subquery")]
30 UnexpectedSubquery,
31 #[error("malformed {kind} value: {reason}")]
32 Malformed { kind: &'static str, reason: String },
33}
34
35impl From<ValueRetentionError> for CatalogRetentionError {
36 fn from(error: ValueRetentionError) -> Self {
37 match error {
38 ValueRetentionError::Memory(error) => Self::Memory(error),
39 ValueRetentionError::Cancelled(error) => Self::Cancelled(error),
40 ValueRetentionError::Malformed { kind, reason } => Self::Malformed { kind, reason },
41 }
42 }
43}
44
45impl From<CatalogRetentionError> for crate::SQLError {
46 fn from(error: CatalogRetentionError) -> Self {
47 match error {
48 CatalogRetentionError::Memory(error) => Self::Routine {
49 sqlstate: "53200".into(),
50 message: error.to_string(),
51 },
52 CatalogRetentionError::Cancelled(error) => Self::Cancelled(error),
53 CatalogRetentionError::UnexpectedSubquery => Self::Internal(error.to_string()),
54 CatalogRetentionError::Malformed { .. } => Self::Routine {
55 sqlstate: "XX001".into(),
56 message: error.to_string(),
57 },
58 }
59 }
60}
61
62type Result<T> = std::result::Result<T, CatalogRetentionError>;
63
64#[derive(Debug, Clone)]
66pub struct RetainedColumns(Arc<Budgeted<Arc<Vec<ColumnDef>>>>);
67
68impl RetainedColumns {
69 pub fn capture(
71 columns: &Arc<Vec<ColumnDef>>,
72 budget: &MemoryBudget,
73 cancellation: &CancellationToken,
74 ) -> Result<Self> {
75 let mut walker = Walker::new(budget, cancellation);
76 walker.charge(size_of::<Vec<ColumnDef>>())?;
77 walker.children(columns, Node::Column)?;
78 let memory = walker.finish()?;
79 cancellation.check()?;
80 Ok(Self(
81 Budgeted::new(Arc::clone(columns), memory).into_shared()?,
82 ))
83 }
84
85 pub fn as_slice(&self) -> &[ColumnDef] {
86 &self.0
87 }
88
89 pub fn reserved_bytes(&self) -> usize {
90 self.0.reserved_bytes()
91 }
92}
93
94impl std::ops::Deref for RetainedColumns {
95 type Target = [ColumnDef];
96
97 fn deref(&self) -> &Self::Target {
98 self.as_slice()
99 }
100}
101
102impl ColumnDef {
103 pub fn reserve_retained_payload(
105 &self,
106 budget: &MemoryBudget,
107 cancellation: &CancellationToken,
108 ) -> Result<MemoryReservation> {
109 Walker::new(budget, cancellation).root(Node::Column(self))
110 }
111}
112
113impl ColumnType {
114 pub(crate) fn retain_external_with_control(
116 self,
117 control: &ProductionControl<'_>,
118 ) -> Result<Produced<Self>> {
119 struct Handoff {
120 value: ColumnType,
121 memory: Option<MemoryReservation>,
122 }
123 let mut output = Handoff {
124 value: self,
125 memory: control.empty_reservation(),
126 };
127 control.check()?;
128 if let Some(budget) = control.budget() {
129 let mut walker = Walker::with_control(budget, *control);
130 let result = walker
131 .visit(Node::Type(&output.value))
132 .and_then(|()| walker.drain());
133 let Walker {
134 memory, pending, ..
135 } = walker;
136 drop(pending);
137 output
138 .memory
139 .as_mut()
140 .expect("controlled catalog handoff")
141 .absorb(memory);
142 result?;
143 }
144 Ok(control.finish(output.value, output.memory)?)
145 }
146
147 pub fn reserve_retained_payload(
149 &self,
150 budget: &MemoryBudget,
151 cancellation: &CancellationToken,
152 ) -> Result<MemoryReservation> {
153 Walker::new(budget, cancellation).root(Node::Type(self))
154 }
155}
156
157impl Expr {
158 pub fn reserve_column_payload(
160 &self,
161 budget: &MemoryBudget,
162 cancellation: &CancellationToken,
163 ) -> Result<MemoryReservation> {
164 Walker::new(budget, cancellation).root(Node::Expr(self))
165 }
166}
167
168enum Node<'a> {
169 Column(&'a ColumnDef),
170 Type(&'a ColumnType),
171 Expr(&'a Expr),
172 Binding(&'a FunctionBinding),
173}
174
175struct Walker<'a> {
176 memory: MemoryReservation,
177 pending: BudgetedVec<Node<'a>>,
178 control: ProductionControl<'a>,
179}
180
181impl<'a> Walker<'a> {
182 fn new(budget: &'a MemoryBudget, cancellation: &'a CancellationToken) -> Self {
183 Self::with_control(
184 budget,
185 ProductionControl::new(budget, cancellation, cancellation),
186 )
187 }
188
189 fn with_control(budget: &'a MemoryBudget, control: ProductionControl<'a>) -> Self {
190 Self {
191 memory: budget.empty_reservation(),
192 pending: BudgetedVec::new(budget),
193 control,
194 }
195 }
196
197 fn root(mut self, root: Node<'a>) -> Result<MemoryReservation> {
198 self.visit(root)?;
199 self.finish()
200 }
201
202 fn finish(mut self) -> Result<MemoryReservation> {
203 self.drain()?;
204 Ok(self.memory)
205 }
206
207 fn drain(&mut self) -> Result<()> {
208 self.control.check_cancellation()?;
209 while let Some(node) = self.pending.pop() {
210 self.visit(node)?;
211 }
212 Ok(())
213 }
214
215 fn visit(&mut self, node: Node<'a>) -> Result<()> {
216 self.control.check_cancellation()?;
217 match node {
218 Node::Column(column) => self.column(column),
219 Node::Type(ty) => self.ty(ty),
220 Node::Expr(expr) => self.expr(expr),
221 Node::Binding(binding) => self.binding(binding),
222 }
223 }
224
225 fn node(&mut self, node: Node<'a>) -> Result<()> {
226 self.control.check_cancellation()?;
227 self.pending.push(node)?;
228 Ok(())
229 }
230
231 fn charge(&mut self, bytes: usize) -> Result<()> {
232 self.control.check_cancellation()?;
233 self.memory.grow(bytes)?;
234 Ok(())
235 }
236
237 fn buffer<T>(&mut self, capacity: usize) -> Result<()> {
238 self.charge(
239 capacity
240 .checked_mul(size_of::<T>())
241 .ok_or(MemoryError::SizeOverflow)?,
242 )
243 }
244
245 fn children<T>(&mut self, items: &'a Vec<T>, node: fn(&'a T) -> Node<'a>) -> Result<()> {
246 self.buffer::<T>(items.capacity())?;
247 for item in items {
248 self.node(node(item))?;
249 }
250 Ok(())
251 }
252
253 fn boxed<T>(&mut self, item: &'a T, node: fn(&'a T) -> Node<'a>) -> Result<()> {
254 self.charge(size_of::<T>())?;
255 self.node(node(item))
256 }
257
258 fn text(&mut self, text: &String) -> Result<()> {
259 self.charge(text.capacity())
260 }
261
262 fn optional_text(&mut self, text: Option<&String>) -> Result<()> {
263 if let Some(text) = text {
264 self.text(text)?;
265 }
266 Ok(())
267 }
268
269 fn texts(&mut self, texts: &Vec<String>) -> Result<()> {
270 self.buffer::<String>(texts.capacity())?;
271 for text in texts {
272 self.text(text)?;
273 }
274 Ok(())
275 }
276
277 fn value(&mut self, value: &Value) -> Result<()> {
278 let memory = value.reserve_retained_payload_with_check(self.memory.budget(), || {
279 self.control.check_cancellation()
280 })?;
281 self.memory.absorb(memory);
282 Ok(())
283 }
284
285 fn optional_expr(&mut self, expr: Option<&'a Expr>) -> Result<()> {
286 if let Some(expr) = expr {
287 self.node(Node::Expr(expr))?;
288 }
289 Ok(())
290 }
291
292 fn optional_boxed_expr(&mut self, expr: Option<&'a Expr>) -> Result<()> {
293 if let Some(expr) = expr {
294 self.boxed(expr, Node::Expr)?;
295 }
296 Ok(())
297 }
298}
299
300#[cfg(test)]
301mod tests;