datafusion_expr/
physical_planning_context.rs1use std::fmt;
19use std::hash::{Hash, Hasher};
20use std::sync::{Arc, Mutex};
21
22use datafusion_common::{HashMap, Result, ScalarValue, TableReference, internal_err};
23
24#[derive(Clone, Debug, Default)]
50pub struct PhysicalPlanningContext {
51 indexes: Arc<HashMap<crate::logical_plan::Subquery, SubqueryIndex>>,
54 results: ScalarSubqueryResults,
55 lambda_variable_qualifier: HashMap<String, TableReference>,
58}
59
60impl PhysicalPlanningContext {
61 pub fn new(
65 indexes: HashMap<crate::logical_plan::Subquery, SubqueryIndex>,
66 results: ScalarSubqueryResults,
67 ) -> Self {
68 Self {
69 indexes: Arc::new(indexes),
70 results,
71 lambda_variable_qualifier: HashMap::new(),
72 }
73 }
74
75 pub fn index_of(
77 &self,
78 subquery: &crate::logical_plan::Subquery,
79 ) -> Option<SubqueryIndex> {
80 self.indexes.get(subquery).copied()
81 }
82
83 pub fn results(&self) -> &ScalarSubqueryResults {
85 &self.results
86 }
87
88 pub fn with_qualified_lambda_variables(
91 mut self,
92 qualifier: &TableReference,
93 variables: &[String],
94 ) -> Self {
95 for var in variables {
96 self.lambda_variable_qualifier
97 .entry_ref(var)
98 .insert(qualifier.clone());
99 }
100
101 self
102 }
103
104 pub fn lambda_variable_qualifier(&self, name: &str) -> Option<&TableReference> {
106 self.lambda_variable_qualifier.get(name)
107 }
108}
109
110#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
112pub struct SubqueryIndex(usize);
113
114impl SubqueryIndex {
115 pub const fn new(index: usize) -> Self {
117 Self(index)
118 }
119
120 pub const fn as_usize(self) -> usize {
122 self.0
123 }
124}
125
126#[derive(Clone, Default)]
133pub struct ScalarSubqueryResults {
134 slots: Arc<Vec<Mutex<Option<ScalarValue>>>>,
135}
136
137impl ScalarSubqueryResults {
138 pub fn new(n: usize) -> Self {
140 Self {
141 slots: Arc::new((0..n).map(|_| Mutex::new(None)).collect()),
142 }
143 }
144
145 pub fn get(&self, index: SubqueryIndex) -> Option<ScalarValue> {
147 let slot = self.slots.get(index.as_usize())?;
148 slot.lock().unwrap().clone()
149 }
150
151 pub fn set(&self, index: SubqueryIndex, value: ScalarValue) -> Result<()> {
153 let Some(slot) = self.slots.get(index.as_usize()) else {
154 return internal_err!(
155 "ScalarSubqueryResults: result index {} is out of bounds",
156 index.as_usize()
157 );
158 };
159
160 let mut slot = slot.lock().unwrap();
161 if slot.is_some() {
162 return internal_err!(
163 "ScalarSubqueryResults: result for index {} was already populated",
164 index.as_usize()
165 );
166 }
167 *slot = Some(value);
168
169 Ok(())
170 }
171
172 pub fn clear(&self) {
174 for slot in self.slots.iter() {
175 *slot.lock().unwrap() = None;
176 }
177 }
178
179 pub fn ptr_eq(this: &Self, other: &Self) -> bool {
181 Arc::ptr_eq(&this.slots, &other.slots)
182 }
183}
184
185impl fmt::Debug for ScalarSubqueryResults {
186 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
187 f.debug_list()
188 .entries(self.slots.iter().map(|slot| slot.lock().unwrap().clone()))
189 .finish()
190 }
191}
192
193impl PartialEq for ScalarSubqueryResults {
194 fn eq(&self, other: &Self) -> bool {
195 Self::ptr_eq(self, other)
196 }
197}
198
199impl Eq for ScalarSubqueryResults {}
200
201impl Hash for ScalarSubqueryResults {
202 fn hash<H: Hasher>(&self, state: &mut H) {
203 Arc::as_ptr(&self.slots).hash(state);
204 }
205}
206
207#[cfg(test)]
208mod tests {
209 use super::*;
210
211 #[test]
212 fn scalar_subquery_results_set_and_get() -> Result<()> {
213 let results = ScalarSubqueryResults::new(1);
214 assert_eq!(results.get(SubqueryIndex::new(0)), None);
215
216 results.set(SubqueryIndex::new(0), ScalarValue::Int32(Some(42)))?;
217 assert_eq!(
218 results.get(SubqueryIndex::new(0)),
219 Some(ScalarValue::Int32(Some(42)))
220 );
221 assert!(
222 results
223 .set(SubqueryIndex::new(0), ScalarValue::Int32(Some(7)))
224 .is_err()
225 );
226
227 Ok(())
228 }
229
230 #[test]
231 fn lambda_variables_shadow_outer_scope() {
232 let outer = TableReference::bare("lambda_1");
233 let inner = TableReference::bare("lambda_2");
234
235 let ctx = PhysicalPlanningContext::default()
236 .with_qualified_lambda_variables(&outer, &["x".to_string(), "y".to_string()])
237 .with_qualified_lambda_variables(&inner, &["y".to_string()]);
238
239 assert_eq!(ctx.lambda_variable_qualifier("x"), Some(&outer));
240 assert_eq!(ctx.lambda_variable_qualifier("y"), Some(&inner));
241 assert_eq!(ctx.lambda_variable_qualifier("z"), None);
242 }
243
244 #[test]
245 fn scalar_subquery_results_clear() -> Result<()> {
246 let results = ScalarSubqueryResults::new(1);
247 results.set(SubqueryIndex::new(0), ScalarValue::Int32(Some(42)))?;
248
249 results.clear();
250
251 assert_eq!(results.get(SubqueryIndex::new(0)), None);
252 results.set(SubqueryIndex::new(0), ScalarValue::Int32(Some(7)))?;
253 assert_eq!(
254 results.get(SubqueryIndex::new(0)),
255 Some(ScalarValue::Int32(Some(7)))
256 );
257
258 Ok(())
259 }
260}