mod.rsannotatedmod.rssource191 lines · 6.3 KB · raw

The batch scan (contract Execution and §1): one plan node evaluates every call to the extension's functions over one relation, many rows at a time, instead of the executor calling them row by row.

Planning. set_rel_pathlist_hook finds the calls in a relation's WHERE clauses and, when the relation is the whole query, its select list. It wraps every path of the relation in a CustomPath whose child is the original path, so index and parameterized paths survive. With ORDER BY, Postgres computes volatile select-list expressions in a projection above the Sort; create_upper_paths_hook makes that projection the scan too, with the Sort as its child.

The plan puts each call in custom_scan_tlist after the child's own columns; setrefs then rewrites every occurrence of the call in the node's quals and projection into a column of the scan tuple, which the node fills. A call it could not claim (under an aggregate, in a join's select list) stays an ordinary function call. guard.rs refuses such a plan in ExecutorStart, before any row is judged.

Execution. The node pulls rows from its child while fewer than [Judge::window] are unanswered, evaluates each row's call arguments, and spawns [Judge::judge] for the row in a task Scope. It emits rows in the child's order, waiting on the executor only when the oldest row is still unanswered. A LIMIT therefore stops the judging: nothing past the window is ever sent.

The safe surface is [Judge]; everything that touches pg_sys is in this module.

32use std::ffi::CStr;
33use std::future::Future;
34use std::pin::Pin;
36use pgrx::pg_sys::panic::ErrorReport;
37
38mod backend;
39mod exec;
40mod ffi;
41mod guard;
42mod lift;
43mod plan;
44mod reach;
45mod statement;
46mod tuple;
47
48pub use backend::on_backend;
49pub use plan::{register, support};
50pub use statement::Statement;

A function the scan evaluates, as declared in SQL.

53pub struct Function {
54    pub name: &'static str,
55    pub args: &'static [Arg],
56    pub returns: Returns,
57}

How an argument reaches [Judge::judge].

60#[derive(Clone, Copy, Debug, PartialEq, Eq)]
61pub enum Arg {

anyelement, as JSON written by [crate::row]'s RowJson.

63    Row,

text.

65    Text,

float8.

67    Float8,

text[], in element order; a NULL element is None.

69    TextArray,
70}
72#[derive(Clone, Copy, Debug, PartialEq, Eq)]
73pub enum Returns {
74    Float8,
75    Bool,
76    Text,

The function's composite result type, from [Out::Record]'s fields in column order.

79    Record,
80}

An evaluated argument.

83#[derive(Clone, Debug, PartialEq)]
84pub enum Value {
85    Json(String),
86    Text(String),
87    Float8(f64),
88    TextArray(Vec<Option<String>>),
89}

One call in one row: which of [Judge::FUNCTIONS], and its arguments in declaration order. Calls with a NULL argument never get here: every function is STRICT, so the scan answers them NULL itself.

94#[derive(Clone, Debug, PartialEq)]
95pub struct Call {
96    pub function: usize,
97    pub args: Vec<Value>,
98}

A call's result; it must match the function's [Returns].

101#[derive(Clone, Debug, PartialEq)]
102pub enum Out {
103    Float8(f64),
104    Bool(bool),
105    Text(String),

One per column of the result type; None is NULL.

107    Record(Vec<Option<Field>>),
108}

A column of an [Out::Record].

111#[derive(Clone, Debug, PartialEq)]
112pub enum Field {
113    Float8(f64),
114    Int4(i32),
115    Bool(bool),
116    Text(String),

jsonb, from its text.

118    Jsonb(String),
119}

A refusal the scan raises as an ERROR (boxed: a report is large).

122pub type Refusal = Box<ErrorReport>;

A row's answers, one per call; None is SQL NULL (a judgment the judge chose not to make, such as jev.on_error = unsure).

126pub type Judgement = Pin<Box<dyn Future<Output = Result<Vec<Option<Out>>, Refusal>>>>;

One line a scan adds to EXPLAIN, as label: value unit.

129#[derive(Clone, Copy, Debug, PartialEq)]
130pub struct Property {
131    pub label: &'static CStr,
132    pub unit: Option<&'static CStr>,
133    pub value: Shown,
134}
136#[derive(Clone, Copy, Debug, PartialEq)]
137pub enum Shown {
138    Integer(i64),

With this many digits after the point.

140    Float(f64, u8),
141}

What a statement's jev scans are estimated to judge, from the planner's row counts after every SQL filter.

145#[derive(Clone, Copy, Debug, Default, PartialEq)]
146pub struct Estimate {

Rows judged, summed over the statement's jev scans.

148    pub rows: f64,

Their width in the child's own layout, summed.

150    pub bytes: f64,
151}

What the scan asks of the extension. begin runs once per scan, in the executor's startup, so settings are read once per statement (a generic plan outlives a SET). Its [Statement] holds what the statement's scans share.

157pub trait Judge: Sized + 'static {

Shown in EXPLAIN as Custom Scan (NAME).

159    const NAME: &'static CStr;

The extension whose schema holds [Self::FUNCTIONS].

161    const EXTENSION: &'static CStr;
162    const FUNCTIONS: &'static [Function];
164    fn begin(statement: &Statement) -> Result<Self, Refusal>;

Refuses a statement whose estimate is over budget. Runs in ExecutorStart, before any scan begins, so nothing has been sent.

168    fn afford(estimate: &Estimate) -> Result<(), Refusal>;

What EXPLAIN shows for one scan: estimate is its child's, absent under COSTS OFF; run is the scan that ran, absent without ANALYZE, where [Self::begin] is never called.

173    fn explain(estimate: Option<&Estimate>, run: Option<&Self>) -> Vec<Property>;

Rows judged at once. Asked before each row is pulled, so it may follow what the connection learns mid-scan.

177    fn window(&self) -> usize;

Every call one row makes, answered in the same order. The future runs on the backend's executor and must not call Postgres.

181    fn judge(&self, calls: Vec<Call>) -> Judgement;

Runs on the backend between waits and when the scan ends, where Postgres may be called: work a judgement queued for the backend (storing answers in the cache) happens here. It runs after the closures judgements queued through [on_backend], which is how a judgement of several rounds reads the cache between them.

188    fn settle(&self) -> Result<(), Refusal> {
189        Ok(())
190    }
191}