exec.rsannotatedexec.rssource569 lines · 24.2 KB · raw

Execution: pull rows ahead while the window has room, judge each in a scoped task, emit in the child's order.

4use std::cell::{Cell, RefCell};
5use std::collections::VecDeque;
6use std::ffi::{CStr, c_int};
7use std::rc::Rc;
8use std::task::{Poll, Waker};
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.

30#[repr(C)]
31pub(super) struct State<J: Judge> {
32    base: pg_sys::CustomScanState,
33    run: Option<Run<J>>,
34}

Everything a running scan owns. Absent under EXPLAIN without ANALYZE.

37struct Run<J: Judge> {
38    judge: J,
39    child: *mut pg_sys::PlanState,
40    tuple: ScanTuple,

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.

53    settling: pg_sys::MemoryContext,
54    waiting: VecDeque<Waiting>,
55    child_done: bool,
56    board: Rc<Board>,

Last, so the tasks go before the state they write into.

58    scope: Scope,
59}
61struct Planned {
62    function: usize,
63    returns: Returns,

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.

73    reach: Option<Expr>,
74}
76struct Waiting {
77    tuple: pg_sys::MinimalTuple,

Per call column: None is NULL.

79    answer: Rc<RefCell<Option<Vec<Option<Out>>>>>,
80}

Where tasks report to the waiting scan.

83#[derive(Default)]
84struct Board {
85    waker: Cell<Option<Waker>>,
86    failure: RefCell<Option<Refusal>>,
87}
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>].

127unsafe fn state<'a, J: Judge>(node: *mut pg_sys::CustomScanState) -> &'a mut State<J> {
128    unsafe { &mut *node.cast::<State<J>>() }
129}
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.

519    unsafe fn restart(&mut self) {
520        unsafe {
521            self.scope.cancel();
522            self.waiting.clear();
523            exec_clear_tuple(self.holder);
524            pg_sys::MemoryContextReset(self.rows);
525            self.board.failure.take();
526            self.child_done = false;
527        }
528    }
529}
531fn internal(message: String) -> ErrorReport {
532    ErrorReport::new(PgSqlErrorCode::ERRCODE_INTERNAL_ERROR, message, "jev scan")
533}
534
535fn raise(report: ErrorReport) -> ! {
536    report.report(PgLogLevel::ERROR);
537    unreachable!("an ERROR report does not return")
538}

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}