Execution: pull rows ahead while the window has room, judge each in a scoped task, emit in the child's order.
10use pgrx::pg_sys::panic::ErrorReport; 11use pgrx::{FromDatum, IntoDatum, PgLogLevel, PgMemoryContexts, PgSqlErrorCode, pg_guard, pg_sys}; 12 13use super::ffi::{ 14 append, copy_minimal_tuple, exec_clear_tuple, exec_proc_node, ints, invoke_function_execute_hook, is_a, ptrs, 15 tup_is_null, walk, 16}; 17use super::plan::tables; 18use super::tuple::{Expr, ScanTuple}; 19use super::backend; 20use super::reach; 21use super::guard::scan_estimate; 22use super::{Arg, Call, Field, Judge, Out, Refusal, Returns, Shown, Statement, Value}; 23use crate::executor::{Scope, block_on}; 24use crate::row::{CanonicalOutput, RowJson};
Postgres holds this as its CustomScanState: the base must come
first. It is Rust-heap memory that Postgres never frees; the query's
memory context drops it (pgrx's leak_and_drop_on_delete), which
also covers the queries that error and never reach EndCustomScan.
Everything a running scan owns. Absent under EXPLAIN without ANALYZE.
One per call column, in scan-tuple order.
42 calls: Vec<Planned>,
Writes row arguments as JSON.
44 json: RefCell<RowJson>,
No call is referenced: rows pass straight through.
46 passthrough: bool,
Receives a row's copy when it is emitted, and frees it at the next.
48 holder: *mut pg_sys::TupleTableSlot,
Where waiting rows are copied; reset on rescan.
50 rows: pg_sys::MemoryContext,
Where storing answers allocates (SPI's argument datums); reset after each store.
The blessed descriptor of a [Returns::Record] function's result
type, in the query's context; null otherwise.
66 record: pg_sys::TupleDesc,
Each argument, and the type it evaluates to.
68 args: Vec<(Arg, Expr, pg_sys::Oid, i32)>,
Referenced by the quals or the projection; others are left NULL.
70 used: bool,
When the row does not satisfy it, Postgres never reads the call,
which is left NULL and not sent (reach.rs); None is always.
Where tasks report to the waiting scan.
89impl Board { 90 fn wake(&self) { 91 if let Some(waker) = self.waker.take() { 92 waker.wake(); 93 } 94 } 95} 96 97pub(super) fn methods<J: Judge>() -> pg_sys::CustomExecMethods { 98 pg_sys::CustomExecMethods { 99 CustomName: J::NAME.as_ptr(), 100 BeginCustomScan: Some(begin::<J>), 101 ExecCustomScan: Some(exec::<J>), 102 EndCustomScan: Some(end::<J>), 103 ReScanCustomScan: Some(rescan::<J>), 104 ExplainCustomScan: Some(explain::<J>), 105 ..Default::default() 106 } 107} 108 109#[pg_guard] 110pub(super) unsafe extern "C-unwind" fn create_state<J: Judge>(_plan: *mut pg_sys::CustomScan) -> *mut pg_sys::Node { 111 let state = State::<J> { 112 base: pg_sys::CustomScanState { 113 ss: pg_sys::ScanState { 114 ps: pg_sys::PlanState { type_: pg_sys::NodeTag::T_CustomScanState, ..Default::default() }, 115 ..Default::default() 116 }, 117 methods: &tables().exec, 118 ..Default::default() 119 }, 120 run: None, 121 }; 122 PgMemoryContexts::CurrentMemoryContext.leak_and_drop_on_delete(state).cast() 123}
Safety
node came from [create_state::<J>].
131#[pg_guard] 132unsafe extern "C-unwind" fn begin<J: Judge>(node: *mut pg_sys::CustomScanState, estate: *mut pg_sys::EState, eflags: c_int) { 133 // SAFETY: ExecInitCustomScan has set up the node for this plan. 134 unsafe { 135 let plan = (*node).ss.ps.plan.cast::<pg_sys::CustomScan>(); 136 let child_plan = ptrs::<pg_sys::Plan>((*plan).custom_plans).next().expect("one child plan"); 137 let child = pg_sys::ExecInitNode(child_plan, estate, eflags); 138 (*node).custom_ps = append(std::ptr::null_mut(), [child]); 139 140 let mut private = ints((*plan).custom_private); 141 let prefix = private.next().expect("the child's column count") as usize; 142 let functions: Vec<usize> = private.map(|i| i as usize).collect(); 143 let funcids: Vec<pg_sys::Oid> = ptrs::<pg_sys::TargetEntry>((*plan).custom_scan_tlist) 144 .skip(prefix) 145 .map(|entry| (*(*entry).expr.cast::<pg_sys::FuncExpr>()).funcid) 146 .collect(); 147 148 // After setrefs, a call the node computes is an INDEX_VAR column 149 // past the prefix wherever the quals or projection use it. 150 let mut used = vec![false; functions.len()]; 151 for list in [(*plan).scan.plan.targetlist, (*plan).scan.plan.qual] { 152 walk(list.cast(), &mut |n| { 153 if is_a(n.cast(), pg_sys::NodeTag::T_Var) { 154 let var = n.cast::<pg_sys::Var>(); 155 let column = (*var).varattno as usize; 156 if (*var).varno == pg_sys::INDEX_VAR && column > prefix { 157 used[column - prefix - 1] = true; 158 } 159 } 160 true 161 }); 162 } 163 164 // The calls are never initialized as expressions, which is where 165 // Postgres checks EXECUTE; so check it here, as ExecInitFunc does. 166 for (&funcid, _) in funcids.iter().zip(&used).filter(|&(_, &u)| u) { 167 let acl = pg_sys::object_aclcheck(pg_sys::ProcedureRelationId, funcid, pg_sys::GetUserId(), pg_sys::ACL_EXECUTE as _); 168 if acl != pg_sys::AclResult::ACLCHECK_OK { 169 pg_sys::aclcheck_error(acl, pg_sys::ObjectType::OBJECT_FUNCTION, pg_sys::get_func_name(funcid)); 170 } 171 invoke_function_execute_hook(funcid); 172 } 173 174 if eflags as u32 & pg_sys::EXEC_FLAG_EXPLAIN_ONLY != 0 { 175 return; 176 } 177 178 let parent = &raw mut (*node).ss.ps; 179 let mut exprs = ptrs::<pg_sys::Expr>((*plan).custom_exprs); 180 let calls: Vec<Planned> = functions 181 .iter() 182 .zip(used) 183 .zip(&funcids) 184 .map(|((&function, used), &funcid)| { 185 let spec = &J::FUNCTIONS[function]; 186 let record = match spec.returns { 187 Returns::Record => { 188 pg_sys::BlessTupleDesc(pg_sys::lookup_rowtype_tupdesc_copy(pg_sys::get_func_rettype(funcid), -1)) 189 } 190 _ => std::ptr::null_mut(), 191 }; 192 let args = spec 193 .args 194 .iter() 195 .map(|&kind| { 196 let expr = exprs.next().expect("an argument per parameter"); 197 let (typid, typmod) = (pg_sys::exprType(expr.cast()), pg_sys::exprTypmod(expr.cast())); 198 (kind, Expr::init(expr, parent), typid, typmod) 199 }) 200 .collect(); 201 Planned { function, returns: spec.returns, record, args, used, reach: None } 202 }) 203 .collect(); 204 let mut calls = calls; 205 for planned in &mut calls { 206 let reach = exprs.next().expect("a reach per call"); 207 if !reach::is_always(reach.cast()) { 208 planned.reach = Some(Expr::init(reach, parent)); 209 } 210 } 211 assert!(exprs.next().is_none(), "custom_exprs is each call's arguments, then each call's reach"); 212 213 let judge = J::begin(&Statement::new(estate)).unwrap_or_else(|e| raise(*e)); 214 let rows = pg_sys::AllocSetContextCreateInternal( 215 pg_sys::CurrentMemoryContext, 216 c"jev scan rows".as_ptr(), 217 pg_sys::ALLOCSET_DEFAULT_MINSIZE as _, 218 pg_sys::ALLOCSET_DEFAULT_INITSIZE as _, 219 pg_sys::ALLOCSET_DEFAULT_MAXSIZE as _, 220 ); 221 let settling = pg_sys::AllocSetContextCreateInternal( 222 pg_sys::CurrentMemoryContext, 223 c"jev scan settle".as_ptr(), 224 pg_sys::ALLOCSET_DEFAULT_MINSIZE as _, 225 pg_sys::ALLOCSET_DEFAULT_INITSIZE as _, 226 pg_sys::ALLOCSET_DEFAULT_MAXSIZE as _, 227 ); 228 let passthrough = !calls.iter().any(|c| c.used); 229 state::<J>(node).run = Some(Run { 230 judge, 231 child, 232 tuple: ScanTuple::new(&raw mut (*node).ss, prefix), 233 calls, 234 json: RefCell::new(RowJson::new()), 235 passthrough, 236 holder: pg_sys::ExecInitExtraTupleSlot(estate, pg_sys::ExecGetResultType(child), &pg_sys::TTSOpsMinimalTuple), 237 rows, 238 settling, 239 waiting: VecDeque::new(), 240 child_done: false, 241 board: Rc::default(), 242 scope: Scope::new(), 243 }); 244 } 245} 246 247#[pg_guard] 248unsafe extern "C-unwind" fn exec<J: Judge>(node: *mut pg_sys::CustomScanState) -> *mut pg_sys::TupleTableSlot { 249 // SAFETY: ExecScan applies the quals and projection to what `next` 250 // returns, as for any scan. 251 unsafe { pg_sys::ExecScan(&raw mut (*node).ss, Some(next::<J>), Some(recheck::<J>)) } 252} 253 254#[pg_guard] 255unsafe extern "C-unwind" fn next<J: Judge>(ss: *mut pg_sys::ScanState) -> *mut pg_sys::TupleTableSlot { 256 // SAFETY: `ss` is this node's (ExecScan passes it back); begin ran. 257 unsafe { state::<J>(ss.cast()).run.as_mut().expect("the scan began").next() } 258}
An EvalPlanQual recheck (DML or a row lock meeting a concurrent
update). The scan's scanrelid is 0, so ExecScan asks this to fill the
scan slot: the child, a scan of the relation in the recheck's executor,
returns the row's current version once, and it is judged as any row.
That executor's scans begin with [Statement::rechecking], so the
judge answers from the statement's judgments or the cache and never
sends.
267#[pg_guard] 268unsafe extern "C-unwind" fn recheck<J: Judge>(ss: *mut pg_sys::ScanState, _slot: *mut pg_sys::TupleTableSlot) -> bool { 269 // SAFETY: `ss` is this node's (ExecScan passes it back); begin ran, 270 // and `next` loads the scan slot, which is `_slot`. 271 unsafe { !tup_is_null(state::<J>(ss.cast()).run.as_mut().expect("the scan began").next()) } 272}
274#[pg_guard] 275unsafe extern "C-unwind" fn rescan<J: Judge>(node: *mut pg_sys::CustomScanState) { 276 // SAFETY: an initialized node; the child is the one begin created. 277 unsafe { 278 if let Some(run) = state::<J>(node).run.as_mut() { 279 run.restart(); 280 } 281 pg_sys::ExecScanReScan(&raw mut (*node).ss); 282 let child = ptrs::<pg_sys::PlanState>((*node).custom_ps).next().expect("one child"); 283 // A changed parameter makes the next ExecProcNode rescan it. 284 if (*child).chgParam.is_null() { 285 pg_sys::ExecReScan(child); 286 } 287 } 288} 289 290#[pg_guard] 291unsafe extern "C-unwind" fn explain<J: Judge>( 292 node: *mut pg_sys::CustomScanState, 293 _ancestors: *mut pg_sys::List, 294 es: *mut pg_sys::ExplainState, 295) { 296 // SAFETY: an initialized node of the plan being explained. Under ANALYZE this runs after 297 // the executor finished and before it ends, so `run` is still there. 298 unsafe { 299 let estimate = scan_estimate((*node).ss.ps.plan.cast()); 300 let run = state::<J>(node).run.as_ref().filter(|_| (*es).analyze); 301 for property in J::explain((*es).costs.then_some(&estimate), run.map(|r| &r.judge)) { 302 let unit = property.unit.map_or(std::ptr::null(), CStr::as_ptr); 303 match property.value { 304 Shown::Integer(v) => pg_sys::ExplainPropertyInteger(property.label.as_ptr(), unit, v, es), 305 Shown::Float(v, digits) => { 306 pg_sys::ExplainPropertyFloat(property.label.as_ptr(), unit, v, c_int::from(digits), es) 307 } 308 } 309 } 310 } 311} 312 313#[pg_guard] 314unsafe extern "C-unwind" fn end<J: Judge>(node: *mut pg_sys::CustomScanState) { 315 // SAFETY: an initialized node, ended once. 316 unsafe { 317 // Keeps what was answered before the scope is dropped. 318 if let Some(run) = state::<J>(node).run.as_ref() { 319 run.judge.settle().unwrap_or_else(|e| raise(*e)); 320 } 321 // Drops the scope first: unanswered requests are reset now, not 322 // when the query's memory goes. 323 state::<J>(node).run = None; 324 for child in ptrs::<pg_sys::PlanState>((*node).custom_ps) { 325 pg_sys::ExecEndNode(child); 326 } 327 } 328} 329 330impl<J: Judge> Run<J> {
The next row, answered, in the scan slot; or the slot cleared at the end.
333 unsafe fn next(&mut self) -> *mut pg_sys::TupleTableSlot { 334 unsafe { 335 if self.passthrough { 336 let row = exec_proc_node(self.child); 337 if tup_is_null(row) { 338 return self.tuple.clear(); 339 } 340 return self.tuple.load(row, std::iter::repeat_n(None, self.calls.len())); 341 } 342 loop { 343 self.fill(); 344 self.settle(); 345 if let Some(failure) = self.board.failure.take() { 346 raise(*failure); 347 } 348 let Some(head) = self.waiting.front() else { 349 return self.tuple.clear(); 350 }; 351 if head.answer.borrow().is_some() { 352 let row = self.waiting.pop_front().expect("the head"); 353 return self.emit(row); 354 } 355 let (board, answer) = (self.board.clone(), head.answer.clone()); 356 block_on(std::future::poll_fn(|cx| { 357 if board.failure.borrow().is_some() || answer.borrow().is_some() || backend::pending(cx.waker()) { 358 return Poll::Ready(()); 359 } 360 board.waker.set(Some(cx.waker().clone())); 361 Poll::Pending 362 })); 363 } 364 } 365 }
Pulls rows until the window is full, starting each one's judgement.
368 unsafe fn fill(&mut self) { 369 unsafe { 370 while self.waiting.len() < self.judge.window().max(1) && !self.child_done { 371 let slot = exec_proc_node(self.child); 372 if tup_is_null(slot) { 373 self.child_done = true; 374 break; 375 } 376 let tuple = PgMemoryContexts::For(self.rows).switch_to(|_| copy_minimal_tuple(slot)); 377 self.tuple.load(slot, std::iter::repeat_n(None, self.calls.len())); 378 let (columns, calls) = self.evaluate(); 379 let answer = Rc::new(RefCell::new(None)); 380 if calls.is_empty() { 381 *answer.borrow_mut() = Some(vec![None; self.calls.len()]); 382 } else { 383 let judgement = self.judge.judge(calls); 384 let (board, into, width) = (self.board.clone(), answer.clone(), self.calls.len()); 385 self.scope.spawn(async move { 386 match judgement.await { 387 Ok(outs) if outs.len() == columns.len() => { 388 let mut row = vec![None; width]; 389 for (column, out) in columns.into_iter().zip(outs) { 390 row[column] = out; 391 } 392 *into.borrow_mut() = Some(row); 393 } 394 Ok(outs) => { 395 let message = format!("{} answers for {} calls", outs.len(), columns.len()); 396 board.failure.borrow_mut().get_or_insert_with(|| Box::new(internal(message))); 397 } 398 Err(e) => { 399 board.failure.borrow_mut().get_or_insert(e); 400 } 401 } 402 board.wake(); 403 }); 404 } 405 self.waiting.push_back(Waiting { tuple, answer }); 406 } 407 } 408 }
The row's calls that Postgres reaches and that have no NULL argument (STRICT: those are NULL without asking), and the call column each answers.
413 unsafe fn evaluate(&self) -> (Vec<usize>, Vec<Call>) { 414 unsafe { 415 self.tuple.reset(); 416 // The JSON of a row's leaves depends on these; see crate::row. 417 let _output = CanonicalOutput::enter().unwrap_or_else(|e| raise(internal(e.to_string()))); 418 PgMemoryContexts::For(self.tuple.memory()).switch_to(|_| { 419 let mut columns = Vec::new(); 420 let mut calls = Vec::new(); 421 for (column, planned) in self.calls.iter().enumerate().filter(|(_, p)| p.used) { 422 if let Some(reach) = &planned.reach 423 && !self.tuple.eval(reach).is_some_and(|d| d.value() != 0) 424 { 425 continue; 426 } 427 let args: Option<Vec<Value>> = planned 428 .args 429 .iter() 430 .map(|&(kind, ref expr, typid, typmod)| { 431 let datum = self.tuple.eval(expr)?; 432 Some(match kind { 433 Arg::Row => { 434 let mut json = String::new(); 435 self.json.borrow_mut().write(&mut json, datum, typid, typmod); 436 Value::Json(json) 437 } 438 Arg::Text => Value::Text( 439 CStr::from_ptr(pg_sys::text_to_cstring(datum.cast_mut_ptr::<pg_sys::text>())).to_string_lossy().into_owned(), 440 ), 441 // SAFETY: the planner matched this argument to a float8 442 // parameter, and eval returned it non-NULL. 443 Arg::Float8 => Value::Float8(f64::from_datum(datum, false).expect("a non-NULL float8")), 444 // SAFETY: as above, a non-NULL text[]. The Array 445 // borrows the per-tuple context, and is copied out. 446 Arg::TextArray => Value::TextArray( 447 pgrx::Array::<String>::from_datum(datum, false) 448 .expect("a non-NULL text[]") 449 .iter() 450 .collect(), 451 ), 452 }) 453 }) 454 .collect(); 455 if let Some(args) = args { 456 columns.push(column); 457 calls.push(Call { function: planned.function, args }); 458 } 459 } 460 (columns, calls) 461 }) 462 } 463 }
Loads row into the scan tuple: the child's columns, then the
answers.
467 unsafe fn emit(&mut self, row: Waiting) -> *mut pg_sys::TupleTableSlot { 468 let answers = row.answer.take().expect("an answered row"); 469 // By-reference answers (text) are palloc'd in per-tuple memory, not 470 // the current context: that is the query's (ExecutorRun's), which 471 // would keep every row's until the statement ends. Per-tuple memory 472 // is reset only when the parent pulls again (ExecScan, then `fill`), 473 // so the row stays valid while the parent reads it, as a 474 // projection's results do. 475 unsafe { 476 // SAFETY: the node's own per-tuple context, live for the executor. 477 let datums: Vec<Option<pg_sys::Datum>> = PgMemoryContexts::For(self.tuple.memory()).switch_to(|_| { 478 answers 479 .into_iter() 480 .zip(&self.calls) 481 .map(|(answer, planned)| match (answer, planned.returns) { 482 (None, _) => None, 483 (Some(Out::Float8(p)), Returns::Float8) => p.into_datum(), 484 (Some(Out::Bool(b)), Returns::Bool) => b.into_datum(), 485 (Some(Out::Text(t)), Returns::Text) => t.into_datum(), 486 (Some(Out::Record(fields)), Returns::Record) => Some(record(planned.record, fields)), 487 (Some(out), returns) => raise(internal(format!("{out:?} answered a call returning {returns:?}"))), 488 }) 489 .collect() 490 }); 491 // The holder frees the previous row's copy; this one stays 492 // until the next emit, as long as the scan tuple borrows it. 493 pg_sys::ExecStoreMinimalTuple(row.tuple, self.holder, true); 494 self.tuple.load(self.holder, datums.into_iter()) 495 } 496 }
Runs the work judgements queued for the backend (see
[backend]), then stores what was answered since the last call.
Its SPI arguments
(each question with its labels, each answer) are palloc'd in the
current context, the query's, which would keep every row's until
the statement ends; a context of its own is reset once they are
stored.
505 unsafe fn settle(&self) { 506 unsafe { 507 // SAFETY: the node's own context, live for the executor. 508 PgMemoryContexts::For(self.settling) 509 .switch_to(|_| { 510 backend::run(); 511 self.judge.settle() 512 }) 513 .unwrap_or_else(|e| raise(*e)); 514 pg_sys::MemoryContextReset(self.settling); 515 } 516 }
Back to the first row: drop every judgement and waiting copy.
A composite datum of desc from fields, in the current context.
Safety
desc is a blessed descriptor, live for the call.
544unsafe fn record(desc: pg_sys::TupleDesc, fields: Vec<Option<Field>>) -> pg_sys::Datum { 545 unsafe { 546 if (*desc).natts as usize != fields.len() { 547 raise(internal(format!("{} fields for a record of {} columns", fields.len(), (*desc).natts))); 548 } 549 let (mut values, mut nulls): (Vec<pg_sys::Datum>, Vec<bool>) = fields 550 .into_iter() 551 .map(|field| match field { 552 None => (pg_sys::Datum::from(0), true), 553 Some(Field::Float8(v)) => (v.into_datum().expect("a float8"), false), 554 Some(Field::Int4(v)) => (v.into_datum().expect("an int4"), false), 555 Some(Field::Bool(v)) => (v.into_datum().expect("a bool"), false), 556 Some(Field::Text(v)) => (v.into_datum().expect("a text"), false), 557 Some(Field::Jsonb(v)) => { 558 let text = std::ffi::CString::new(v).expect("JSON text has no NUL"); 559 // SAFETY: jsonb_in reads a NUL-terminated cstring. 560 let datum = pgrx::fcinfo::direct_function_call_as_datum(pg_sys::jsonb_in, &[text.as_c_str().into_datum()]); 561 (datum.expect("jsonb_in returns a value"), false) 562 } 563 }) 564 .unzip(); 565 // SAFETY: one value and one null flag per column of `desc`. 566 let tuple = pg_sys::heap_form_tuple(desc, values.as_mut_ptr(), nulls.as_mut_ptr()); 567 pg_sys::HeapTupleHeaderGetDatum((*tuple).t_data) 568 } 569}