pgdog.git / summary / log / commit / refs
commit 2b42016b96a45cd8ce08df6bb6eea81fc8aeb313
Author: Lev Kokotov <levkk@users.noreply.github.com>
Commit: GitHub <noreply@github.com>
Date: Mon Aug 18 06:23:32 2025 +0000
Re-introducing plugins (#337)
* save
* save
* save
* safe
* save
* save
* docs
* fix tests
* save
* save
* readme
* save
* save
* save
* fmt
* Try to fix docsrs
* save
* save
* save
* save
* save
* save
* save
* upgrade
* save
Cargo.lock | 155 ++++--
Cargo.toml | 10 +-
dev/docs_redirect.html | 18 +
dev/sync_docs.sh | 11 +
examples/routing-plugin/Cargo.toml | 10 -
examples/routing-plugin/src/lib.rs | 29 -
integration/plugins/pgdog.toml | 11 +
integration/plugins/users.toml | 4 +
pgdog-macros/Cargo.toml | 15 +
pgdog-macros/src/lib.rs | 124 +++++
pgdog-plugin-build/Cargo.toml | 10 +
pgdog-plugin-build/src/lib.rs | 42 ++
pgdog-plugin/Cargo.toml | 9 +-
pgdog-plugin/LICENSE | 2 +-
pgdog-plugin/README.md | 21 +-
pgdog-plugin/build.rs | 9 +-
pgdog-plugin/include/plugin.h | 45 --
pgdog-plugin/include/types.h | 319 ++---------
pgdog-plugin/src/ast.rs | 132 +++++
pgdog-plugin/src/bindings.rs | 621 ++++++++++++----------
pgdog-plugin/src/c_api.rs | 24 -
pgdog-plugin/src/comp.rs | 8 +
pgdog-plugin/src/config.rs | 107 ----
pgdog-plugin/src/context.rs | 442 +++++++++++++++
pgdog-plugin/src/copy.rs | 241 ---------
pgdog-plugin/src/input.rs | 63 ---
pgdog-plugin/src/lib.rs | 148 +++++-
pgdog-plugin/src/order_by.rs | 42 --
pgdog-plugin/src/output.rs | 91 ----
pgdog-plugin/src/parameter.rs | 53 --
pgdog-plugin/src/plugin.rs | 170 +++---
pgdog-plugin/src/prelude.rs | 7 +
pgdog-plugin/src/query.rs | 80 ---
pgdog-plugin/src/route.rs | 166 ------
pgdog-plugin/src/string.rs | 80 +++
pgdog/src/frontend/router/mod.rs | 1 -
pgdog/src/frontend/router/parser/context.rs | 21 +
pgdog/src/frontend/router/parser/query/mod.rs | 41 +-
pgdog/src/frontend/router/parser/query/plugins.rs | 89 ++++
pgdog/src/frontend/router/parser/route.rs | 4 +
pgdog/src/frontend/router/request.rs | 46 --
pgdog/src/plugin/mod.rs | 52 +-
plugins/README.md | 13 +-
plugins/pgdog-example-plugin/Cargo.toml | 13 +
plugins/pgdog-example-plugin/src/lib.rs | 42 ++
plugins/pgdog-example-plugin/src/plugin.rs | 120 +++++
plugins/pgdog-routing/Cargo.toml | 27 -
plugins/pgdog-routing/build.rs | 6 -
plugins/pgdog-routing/postgres_hash/LICENSE | 23 -
plugins/pgdog-routing/postgres_hash/hashfn.c | 416 ---------------
plugins/pgdog-routing/src/comment.rs | 46 --
plugins/pgdog-routing/src/copy.rs | 143 -----
plugins/pgdog-routing/src/lib.rs | 137 -----
plugins/pgdog-routing/src/order_by.rs | 56 --
plugins/pgdog-routing/src/sharding_function.rs | 143 -----
55 files changed, 2009 insertions(+), 2749 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
index 6f0c81f9..ac86fc18 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -652,18 +652,6 @@ dependencies = [
"phf",
]
-[[package]]
-name = "csv"
-version = "1.3.1"
-source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "acdc4883a9c96732e4733212c01447ebd805833b7275a73ca3ee080fd77afdaf"
-dependencies = [
- "csv-core",
- "itoa",
- "ryu",
- "serde",
-]
-
[[package]]
name = "csv-core"
version = "0.1.12"
@@ -2277,7 +2265,7 @@ dependencies = [
"once_cell",
"parking_lot",
"pg_query",
- "pgdog-plugin",
+ "pgdog-plugin 0.1.6",
"pin-project",
"rand 0.8.5",
"ratatui",
@@ -2295,7 +2283,7 @@ dependencies = [
"tokio",
"tokio-rustls",
"tokio-util",
- "toml",
+ "toml 0.8.22",
"tracing",
"tracing-subscriber",
"url",
@@ -2303,30 +2291,68 @@ dependencies = [
]
[[package]]
-name = "pgdog-plugin"
+name = "pgdog-example-plugin"
+version = "0.1.0"
+dependencies = [
+ "once_cell",
+ "parking_lot",
+ "pgdog-plugin 0.1.6 (registry+https://github.com/rust-lang/crates.io-index)",
+ "thiserror 2.0.12",
+]
+
+[[package]]
+name = "pgdog-macros"
+version = "0.1.1"
+dependencies = [
+ "proc-macro2",
+ "quote",
+ "syn 2.0.101",
+]
+
+[[package]]
+name = "pgdog-macros"
version = "0.1.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "ce29ee22e1c763d01c2fbe68bf4ca512f3b000c8f6c5cf8a064e03ce197bc39f"
+dependencies = [
+ "proc-macro2",
+ "quote",
+ "syn 2.0.101",
+]
+
+[[package]]
+name = "pgdog-plugin"
+version = "0.1.6"
dependencies = [
"bindgen 0.71.1",
"libc",
"libloading",
+ "pg_query",
+ "pgdog-macros 0.1.1",
+ "toml 0.9.5",
"tracing",
]
[[package]]
-name = "pgdog-routing"
-version = "0.1.0"
+name = "pgdog-plugin"
+version = "0.1.6"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "9ef8d3679b846c077c274fbbd36f1e11f66d4984488e7201c0ac4b8d4a86d999"
dependencies = [
- "cc",
- "csv",
- "once_cell",
+ "bindgen 0.71.1",
+ "libc",
+ "libloading",
"pg_query",
- "pgdog-plugin",
- "postgres",
- "rand 0.8.5",
- "regex",
+ "pgdog-macros 0.1.1 (registry+https://github.com/rust-lang/crates.io-index)",
+ "toml 0.9.5",
"tracing",
- "tracing-subscriber",
- "uuid",
+]
+
+[[package]]
+name = "pgdog-plugin-build"
+version = "0.1.1"
+dependencies = [
+ "toml 0.9.5",
]
[[package]]
@@ -2446,20 +2472,6 @@ version = "1.11.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f84267b20a16ea918e43c6a88433c2d54fa145c92a811b5b047ccbe153674483"
-[[package]]
-name = "postgres"
-version = "0.19.10"
-source = "registry+https://github.com/rust-lang/crates.io-index"
-checksum = "363e6dfbdd780d3aa3597b6eb430db76bb315fa9bad7fae595bb8def808b8470"
-dependencies = [
- "bytes",
- "fallible-iterator",
- "futures-util",
- "log",
- "tokio",
- "tokio-postgres",
-]
-
[[package]]
name = "postgres-protocol"
version = "0.6.8"
@@ -2897,13 +2909,6 @@ dependencies = [
"serde",
]
-[[package]]
-name = "routing-plugin"
-version = "0.1.0"
-dependencies = [
- "pgdog-plugin",
-]
-
[[package]]
name = "rsa"
version = "0.9.8"
@@ -3191,6 +3196,15 @@ dependencies = [
"serde",
]
+[[package]]
+name = "serde_spanned"
+version = "1.0.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "40734c41988f7306bb04f0ecf60ec0f3f1caa34290e4e8ea471dcd3346483b83"
+dependencies = [
+ "serde",
+]
+
[[package]]
name = "serde_urlencoded"
version = "0.7.1"
@@ -3988,11 +4002,26 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "05ae329d1f08c4d17a59bed7ff5b5a769d062e64a62d34a3261b219e62cd5aae"
dependencies = [
"serde",
- "serde_spanned",
- "toml_datetime",
+ "serde_spanned 0.6.8",
+ "toml_datetime 0.6.9",
"toml_edit",
]
+[[package]]
+name = "toml"
+version = "0.9.5"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "75129e1dc5000bfbaa9fee9d1b21f974f9fbad9daec557a521ee6e080825f6e8"
+dependencies = [
+ "indexmap",
+ "serde",
+ "serde_spanned 1.0.0",
+ "toml_datetime 0.7.0",
+ "toml_parser",
+ "toml_writer",
+ "winnow",
+]
+
[[package]]
name = "toml_datetime"
version = "0.6.9"
@@ -4002,6 +4031,15 @@ dependencies = [
"serde",
]
+[[package]]
+name = "toml_datetime"
+version = "0.7.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "bade1c3e902f58d73d3f294cd7f20391c1cb2fbcb643b73566bc773971df91e3"
+dependencies = [
+ "serde",
+]
+
[[package]]
name = "toml_edit"
version = "0.22.26"
@@ -4010,18 +4048,33 @@ checksum = "310068873db2c5b3e7659d2cc35d21855dbafa50d1ce336397c666e3cb08137e"
dependencies = [
"indexmap",
"serde",
- "serde_spanned",
- "toml_datetime",
+ "serde_spanned 0.6.8",
+ "toml_datetime 0.6.9",
"toml_write",
"winnow",
]
+[[package]]
+name = "toml_parser"
+version = "1.0.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "b551886f449aa90d4fe2bdaa9f4a2577ad2dde302c61ecf262d80b116db95c10"
+dependencies = [
+ "winnow",
+]
+
[[package]]
name = "toml_write"
version = "0.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "bfb942dfe1d8e29a7ee7fcbde5bd2b9a25fb89aa70caea2eba3bee836ff41076"
+[[package]]
+name = "toml_writer"
+version = "1.0.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "fcc842091f2def52017664b53082ecbbeb5c7731092bad69d2c63050401dfd64"
+
[[package]]
name = "tower"
version = "0.5.2"
diff --git a/Cargo.toml b/Cargo.toml
index 4302c0e1..95cdffe8 100644
--- a/Cargo.toml
+++ b/Cargo.toml
@@ -1,10 +1,10 @@
[workspace]
-members = [ "examples/demo",
- "examples/routing-plugin", "integration/rust",
- "pgdog",
- "pgdog-plugin",
- "plugins/pgdog-routing",
+members = [
+ "examples/demo",
+ "integration/rust",
+ "pgdog", "pgdog-macros",
+ "pgdog-plugin", "pgdog-plugin-build", "plugins/pgdog-example-plugin",
]
resolver = "2"
diff --git a/dev/docs_redirect.html b/dev/docs_redirect.html
new file mode 100644
index 00000000..750d6cef
--- /dev/null
+++ b/dev/docs_redirect.html
@@ -0,0 +1,18 @@
+<!doctype html>
+<html>
+ <head>
+ <meta
+ http-equiv="refresh"
+ content="0; url=https://docsrs.pgdog.dev/pgdog_plugin/index.html"
+ />
+ <title>Redirecting...</title>
+ </head>
+ <body>
+ <p>
+ If you are not redirected automatically,
+ <a href="https://docsrs.pgdog.dev/pgdog_plugin/index.html"
+ >click here</a
+ >.
+ </p>
+ </body>
+</html>
diff --git a/dev/sync_docs.sh b/dev/sync_docs.sh
new file mode 100644
index 00000000..02361cd7
--- /dev/null
+++ b/dev/sync_docs.sh
@@ -0,0 +1,11 @@
+#!/bin/bash
+#
+SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
+set -e
+cd ${SCRIPT_DIR}/../
+cargo doc
+cd target/doc
+aws s3 sync pgdog_plugin s3://pgdog-docsrs/pgdog_plugin
+aws s3 sync pgdog s3://pgdog-docsrs/pgdog
+aws s3 sync pgdog_plugin_build s3://pgdog-docsrs/pgdog_plugin_build
+aws s3 cp ${SCRIPT_DIR}/docs_redirect.html s3://pgdog-docsrs/index.html
diff --git a/examples/routing-plugin/Cargo.toml b/examples/routing-plugin/Cargo.toml
deleted file mode 100644
index 84e4c29f..00000000
--- a/examples/routing-plugin/Cargo.toml
+++ /dev/null
@@ -1,10 +0,0 @@
-[package]
-name = "routing-plugin"
-version = "0.1.0"
-edition = "2021"
-
-[lib]
-crate-type = ["rlib", "cdylib"]
-
-[dependencies]
-pgdog-plugin = { path = "../../pgdog-plugin" }
diff --git a/examples/routing-plugin/src/lib.rs b/examples/routing-plugin/src/lib.rs
deleted file mode 100644
index 89ea0a3e..00000000
--- a/examples/routing-plugin/src/lib.rs
+++ /dev/null
@@ -1,29 +0,0 @@
-//! Simple routing plugin example using Rust.
-use pgdog_plugin::*;
-
-/// Route query.
-#[no_mangle]
-pub extern "C" fn pgdog_route_query(input: Input) -> Output {
- let is_read = input
- .query()
- .map(|query| query.query().to_lowercase().trim().starts_with("select"))
- .unwrap_or(false);
-
- // This is just an example of extracing a parameter from
- // the query. In the future, we'll use this to shard transactions.
- let _parameter = input.query().map(|query| {
- query.parameter(0).map(|parameter| {
- let id = parameter.as_str().map(|str| str.parse::<i64>());
- match id {
- Some(Ok(id)) => id,
- _ => i64::from_be_bytes(parameter.as_bytes().try_into().unwrap_or([0u8; 8])),
- }
- })
- });
-
- if is_read {
- Output::new_forward(Route::read_any())
- } else {
- Output::new_forward(Route::write_any())
- }
-}
diff --git a/integration/plugins/pgdog.toml b/integration/plugins/pgdog.toml
new file mode 100644
index 00000000..10a323e2
--- /dev/null
+++ b/integration/plugins/pgdog.toml
@@ -0,0 +1,11 @@
+[[plugins]]
+name = "pgdog_example_plugin"
+
+[[databases]]
+name = "pgdog"
+host = "127.0.0.1"
+
+[[databases]]
+name = "pgdog"
+host = "127.0.0.1"
+role = "replica"
diff --git a/integration/plugins/users.toml b/integration/plugins/users.toml
new file mode 100644
index 00000000..581cdb75
--- /dev/null
+++ b/integration/plugins/users.toml
@@ -0,0 +1,4 @@
+[[users]]
+name = "pgdog"
+database = "pgdog"
+password = "pgdog"
diff --git a/pgdog-macros/Cargo.toml b/pgdog-macros/Cargo.toml
new file mode 100644
index 00000000..4bf67124
--- /dev/null
+++ b/pgdog-macros/Cargo.toml
@@ -0,0 +1,15 @@
+[package]
+name = "pgdog-macros"
+version = "0.1.1"
+edition = "2024"
+authors = ["Lev Kokotov <lev@pgdog.dev>"]
+license = "MIT"
+description = "Macros used by the pgdog-plugin crate to generate safe FFI functions"
+
+[lib]
+proc-macro = true
+
+[dependencies]
+proc-macro2 = "1.0"
+quote = "1.0"
+syn = { version = "2.0", features = ["full"] }
diff --git a/pgdog-macros/src/lib.rs b/pgdog-macros/src/lib.rs
new file mode 100644
index 00000000..65df8ce1
--- /dev/null
+++ b/pgdog-macros/src/lib.rs
@@ -0,0 +1,124 @@
+//! Macros used by PgDog plugins.
+//!
+//! Required and exported by the `pgdog-plugin` crate. You don't have to add this crate separately.
+//!
+use proc_macro::TokenStream;
+use quote::quote;
+use syn::{ItemFn, parse_macro_input};
+
+/// Generates required methods for PgDog to run at plugin load time.
+///
+/// ### Methods
+///
+/// * `pgdog_rustc_version`: Returns the version of the Rust compiler used to build the plugin.
+/// * `pgdog_pg_query_version`: Returns the version of the pg_query library used by the plugin.
+/// * `pgdog_plugin_version`: Returns the version of the plugin itself, taken from Cargo.toml.
+///
+#[proc_macro]
+pub fn plugin(_input: TokenStream) -> TokenStream {
+ let expanded = quote! {
+ #[unsafe(no_mangle)]
+ pub unsafe extern "C" fn pgdog_rustc_version(output: *mut pgdog_plugin::PdStr) {
+ let version = pgdog_plugin::comp::rustc_version();
+ unsafe {
+ *output = version;
+ }
+ }
+
+ #[unsafe(no_mangle)]
+ pub unsafe extern "C" fn pgdog_pg_query_version(output: *mut pgdog_plugin::PdStr) {
+ let version: pgdog_plugin::PdStr = option_env!("PGDOG_PGQUERY_VERSION")
+ .unwrap_or_default()
+ .into();
+ unsafe {
+ *output = version;
+ }
+ }
+
+ #[unsafe(no_mangle)]
+ pub unsafe extern "C" fn pgdog_plugin_version(output: *mut pgdog_plugin::PdStr) {
+ let version: pgdog_plugin::PdStr = env!("CARGO_PKG_VERSION").into();
+ unsafe {
+ *output = version;
+ }
+ }
+ };
+ TokenStream::from(expanded)
+}
+
+/// Generate the `pgdog_init` method that's executed at plugin load time.
+#[proc_macro_attribute]
+pub fn init(_attr: TokenStream, item: TokenStream) -> TokenStream {
+ let input_fn = parse_macro_input!(item as ItemFn);
+ let fn_name = &input_fn.sig.ident;
+
+ let expanded = quote! {
+
+ #[unsafe(no_mangle)]
+ pub extern "C" fn pgdog_init() {
+ #input_fn
+
+ #fn_name();
+ }
+ };
+
+ TokenStream::from(expanded)
+}
+
+/// Generate the `pgdog_fini` method that runs at PgDog shutdown.
+#[proc_macro_attribute]
+pub fn fini(_attr: TokenStream, item: TokenStream) -> TokenStream {
+ let input_fn = parse_macro_input!(item as ItemFn);
+ let fn_name = &input_fn.sig.ident;
+
+ let expanded = quote! {
+ #[unsafe(no_mangle)]
+ pub extern "C" fn pgdog_fini() {
+ #input_fn
+
+ #fn_name();
+ }
+ };
+
+ TokenStream::from(expanded)
+}
+
+/// Generates the `pgdog_route` method for routing queries.
+#[proc_macro_attribute]
+pub fn route(_attr: TokenStream, item: TokenStream) -> TokenStream {
+ let input_fn = parse_macro_input!(item as ItemFn);
+ let fn_name = &input_fn.sig.ident;
+ let fn_inputs = &input_fn.sig.inputs;
+
+ // Extract the first parameter name and type for the pgdog_route function signature
+ let (first_param_name, _) = fn_inputs
+ .iter()
+ .filter_map(|input| {
+ if let syn::FnArg::Typed(pat_type) = input {
+ if let syn::Pat::Ident(pat_ident) = &*pat_type.pat {
+ Some((pat_ident.ident.clone(), pat_type.ty.clone()))
+ } else {
+ None
+ }
+ } else {
+ None
+ }
+ })
+ .next()
+ .expect("Route function must have at least one named parameter");
+
+ let expanded = quote! {
+ #[unsafe(no_mangle)]
+ pub unsafe extern "C" fn pgdog_route(#first_param_name: pgdog_plugin::PdRouterContext, output: *mut pgdog_plugin::PdRoute) {
+ #input_fn
+
+ let pgdog_context: pgdog_plugin::Context = #first_param_name.into();
+ let route: pgdog_plugin::PdRoute = #fn_name(pgdog_context).into();
+ unsafe {
+ *output = route;
+ }
+ }
+ };
+
+ TokenStream::from(expanded)
+}
diff --git a/pgdog-plugin-build/Cargo.toml b/pgdog-plugin-build/Cargo.toml
new file mode 100644
index 00000000..98f6cc75
--- /dev/null
+++ b/pgdog-plugin-build/Cargo.toml
@@ -0,0 +1,10 @@
+[package]
+name = "pgdog-plugin-build"
+version = "0.1.1"
+edition = "2024"
+authors = ["Lev Kokotov <lev@pgdog.dev> "]
+license = "MIT"
+description = "Build-time helpers for PgDog plugins"
+
+[dependencies]
+toml = "0.9"
diff --git a/pgdog-plugin-build/src/lib.rs b/pgdog-plugin-build/src/lib.rs
new file mode 100644
index 00000000..64ab7b18
--- /dev/null
+++ b/pgdog-plugin-build/src/lib.rs
@@ -0,0 +1,42 @@
+//! Build-time helpers for PgDog plugins.
+//!
+//! Include this package as a build dependency only.
+//!
+
+use std::{fs::File, io::Read};
+
+/// Extracts the `pg_query` crate version from `Cargo.toml`
+/// and sets it as an environment variable.
+///
+/// This should be used at build time only. It expects `Cargo.toml` to be present in the same
+/// folder as `build.rs`.
+///
+/// ### Note
+///
+/// You should have a strict version constraint on `pg_query`, for example:
+///
+/// ```toml
+/// pg_query = "6.1.0"
+/// ```
+///
+/// If the version in your plugin doesn't match what PgDog is using, your plugin won't be loaded.
+///
+pub fn pg_query_version() {
+ let mut contents = String::new();
+ if let Ok(mut file) = File::open("Cargo.toml") {
+ file.read_to_string(&mut contents).ok();
+ } else {
+ panic!("Cargo.toml not found");
+ }
+
+ let contents: Option<toml::Value> = toml::from_str(&contents).ok();
+ if let Some(contents) = contents {
+ if let Some(dependencies) = contents.get("dependencies") {
+ if let Some(pg_query) = dependencies.get("pg_query") {
+ if let Some(version) = pg_query.as_str() {
+ println!("cargo:rustc-env=PGDOG_PGQUERY_VERSION={}", version);
+ }
+ }
+ }
+ }
+}
diff --git a/pgdog-plugin/Cargo.toml b/pgdog-plugin/Cargo.toml
index 87ea5e6c..3a873ecd 100644
--- a/pgdog-plugin/Cargo.toml
+++ b/pgdog-plugin/Cargo.toml
@@ -1,13 +1,13 @@
[package]
name = "pgdog-plugin"
-version = "0.1.1"
+version = "0.1.6"
edition = "2021"
license = "MIT"
authors = ["Lev Kokotov <lev.kokotov@gmail.com>"]
readme = "README.md"
-repository = "https://github.com/levkk/pgdog"
+repository = "https://github.com/pgdogdev/pgdog"
homepage = "https://pgdog.dev"
-description = "pgDog plugin interface and helpers"
+description = "PgDog plugin interface and helpers"
include = ["src/", "include/", "build.rs", "LICENSE", "README.md"]
[lib]
@@ -17,6 +17,9 @@ crate-type = ["rlib", "cdylib"]
libloading = "0.8"
libc = "0.2"
tracing = "0.1"
+pg_query = "6.1.0"
+pgdog-macros = { path = "../pgdog-macros", version = "0.1.1" }
+toml = "0.9"
[build-dependencies]
bindgen = "0.71.0"
diff --git a/pgdog-plugin/LICENSE b/pgdog-plugin/LICENSE
index b1d87947..3859f508 100644
--- a/pgdog-plugin/LICENSE
+++ b/pgdog-plugin/LICENSE
@@ -1,4 +1,4 @@
-Copyright 2025 Lev Kokotov
+Copyright 2025 PgDog, Inc.
Permission is hereby granted, free of charge, to any person obtaining a copy of this software and
associated documentation files (the “Software”), to deal in the Software without restriction, including
diff --git a/pgdog-plugin/README.md b/pgdog-plugin/README.md
index 3c5327de..29ac8012 100644
--- a/pgdog-plugin/README.md
+++ b/pgdog-plugin/README.md
@@ -1,21 +1,24 @@
# PgDog plugins
-[](https://pgdog.dev)
+[](https://docsrs.pgdog.dev/pgdog_plugin/index.html)
[](https://crates.io/crates/pgdog-plugin)
-[](https://docs.rs/pgdog-plugin/)
-PgDog plugin system is based around shared libraries loaded at runtime.
-These libraries can be written in any language as long as they are compiled to `.so` (or `.dylib` on Mac),
-and can expose predefined C ABI functions.
+PgDog plugin system is based around shared libraries loaded at runtime. The plugins currently can only be
+written in Rust. This is because PgDog passes Rust-specific data types to plugin functions, and those cannot
+be easily made C ABI-compatible.
-This crate implements the bridge between the C ABI and PgDog, defines common C types and interface to use,
-and exposes internal PgDog configuration.
+This crate implements the bridge between PgDog and plugins, making sure data types can be safely passed through the FFI.
-This crate is a C (and Rust) library that should be linked at compile time against your plugins.
+Automatic checks include:
+
+- Rust compiler version check
+- `pg_query` version check
+
+This crate should be linked at compile time against your plugins.
## Writing plugins
-Examples of plugins written in C and Rust are available [here](https://github.com/levkk/pgdog/tree/main/examples).
+See [documentation](https://docsrs.pgdog.dev/pgdog_plugin/index.html) for examples. Example plugins are [available in GitHub](https://github.com/pgdogdev/pgdog/tree/main/plugins) as well.
## License
diff --git a/pgdog-plugin/build.rs b/pgdog-plugin/build.rs
index 7ab129c3..2015d87b 100644
--- a/pgdog-plugin/build.rs
+++ b/pgdog-plugin/build.rs
@@ -1,5 +1,6 @@
-use std::path::PathBuf;
+use std::{path::PathBuf, process::Command};
+#[cfg(not(docsrs))]
fn main() {
println!("cargo:rerun-if-changed=include/types.h");
@@ -16,4 +17,10 @@ fn main() {
let out_path = PathBuf::from("src");
let _ = bindings.write_to_file(out_path.join("bindings.rs"));
+
+ let rustc = std::env::var("RUSTC").unwrap();
+ let version = Command::new(rustc).arg("--version").output().unwrap();
+ let version_str = String::from_utf8(version.stdout).unwrap();
+
+ println!("cargo:rustc-env=RUSTC_VERSION={}", version_str.trim());
}
diff --git a/pgdog-plugin/include/plugin.h b/pgdog-plugin/include/plugin.h
deleted file mode 100644
index 43c0889d..00000000
--- a/pgdog-plugin/include/plugin.h
+++ /dev/null
@@ -1,45 +0,0 @@
-
-#include "types.h"
-
-/* Route query to a primary/replica and shard.
- *
- * Implementing this function is optional. If the plugin
- * implements it, the query router will use its decision
- * to route the query.
- *
- * ## Thread safety
- *
- * This function is not synchronized and can be called
- * for multiple queries at a time. If accessing global state,
- * make sure to protect access with a mutex.
- *
- * ## Performance
- *
- * This function is called for every transaction. It's a hot path,
- * so make sure to optimize for performance in the implementation.
- *
-*/
-Output pgdog_route_query(Input input);
-
-/*
- * Perform initialization at plugin loading time.
- *
- * Executed only once and execution is synchronized,
- * so it's safe to initialize sychroniziation primitives
- * like mutexes in this method.
- */
-void pgdog_init();
-
-/* Create new row.
-*
-* Implemented by pgdog_plugin library.
-* Make sure your plugin links with -lpgdog_plugin.
-*/
-extern Row pgdog_row_new(int num_columns);
-
-/* Free memory allocated for the row.
-*
-* Implemented by pgdog_plugin library.
-* Make sure your plugin links with -lpgdog_plugin.
-*/
-extern void pgdog_row_free(Row row);
diff --git a/pgdog-plugin/include/types.h b/pgdog-plugin/include/types.h
index 4831e4da..56a49163 100644
--- a/pgdog-plugin/include/types.h
+++ b/pgdog-plugin/include/types.h
@@ -1,280 +1,59 @@
-/**
- * Query parameter value.
- */
-typedef struct Parameter {
- int len;
- const char *data;
- int format;
-} Parameter;
-
-/* Query and parameters received by pgDog.
- *
- * The plugin is expected to parse the query and based on its
- * contents and the parameters, make a routing decision.
- */
-typedef struct Query {
- /* Length of the query */
- int len;
-
- /* The query text. */
- const char *query;
-
- /* Number of parameters. */
- int num_parameters;
-
- /* List of parameters. */
- const Parameter *parameters;
-} Query;
+#include <stddef.h>
+#include <stdint.h>
/**
- * The query is a read or a write.
- * In case the plugin isn't able to figure it out, it can return UNKNOWN and
- * pgDog will ignore the plugin's decision.
-*/
-typedef enum Affinity {
- READ = 1,
- WRITE = 2,
- TRANSACTION_START = 3,
- TRANSACTION_END = 4,
- UNKNOWN = -1,
-} Affinity;
-
-/**
- * In case the plugin doesn't know which shard to route the
- * the query, it can decide to route it to any shard or to all
- * shards. All shard queries return a result assembled by pgDog.
- *
-*/
-typedef enum Shard {
- ANY = -1,
- ALL = -2,
-} Shard;
-
-/*
- * Column sort direction.
-*/
-typedef enum OrderByDirection {
- ASCENDING,
- DESCENDING,
-} OrderByDirection;
-
-/*
- * Column sorting.
-*/
-typedef struct OrderBy {
- char *column_name;
- int column_index;
- OrderByDirection direction;
-} OrderBy;
-
-/**
- * Route the query should take.
- *
-*/
-typedef struct Route {
- Affinity affinity;
- int shard;
- int num_order_by;
- OrderBy *order_by;
-} Route;
-
-/**
- * The routing decision the plugin makes based on the query contents.
- *
- * FORWARD: The query is forwarded to a shard. Which shard (and whether it's a replica
- * or a primary) is decided by the plugin output.
- * REWRITE: The query text is rewritten. The plugin outputs new query text.
- * ERROR: The query is denied and the plugin returns an error instead. This error is sent
- * to the client.
- * INTERCEPT: The query is intercepted and the plugin returns rows instead. These rows
- are sent to the client and the original query is never sent to a backend server.
- * NO_DECISION: The plugin doesn't care about this query. The output is ignored by pgDog and the next
- plugin in the chain is attempted.
- * COPY: Client is sending over a COPY statement.
- *
-*/
-typedef enum RoutingDecision {
- FORWARD = 1,
- REWRITE = 2,
- ERROR = 3,
- INTERCEPT = 4,
- NO_DECISION = 5, /* The plugin doesn't want to make a decision. We'll try
- the next plugin in the chain. */
- COPY = 6, /* COPY */
- COPY_ROWS = 7, /* Copy rows. */
-} RoutingDecision;
+ * Wrapper around Rust's [`&str`], without allocating memory, unlike [`std::ffi::CString`].
+ * The caller must use it as a Rust string. This is not a C-string.
+ */
+typedef struct PdStr {
+ size_t len;
+ void *data;
+} RustString;
/*
- * Error returned by the router plugin.
- * This will be sent to the client and the transaction will be aborted.
-*/
-typedef struct Error {
- char *severity;
- char *code;
- char *message;
- char *detail;
-} Error;
-
-typedef struct RowColumn {
- int length;
- char *data;
-} RowColumn;
-
-typedef struct Row {
- int num_columns;
- RowColumn *columns;
-} Row;
-
-typedef struct RowDescriptionColumn {
- int len;
- char *name;
- int oid;
-} RowDescriptionColumn;
-
-typedef struct RowDescription {
- int num_columns;
- RowDescriptionColumn *columns;
-} RowDescription;
-
-typedef struct Intercept {
- RowDescription row_description;
- int num_rows;
- Row *rows;
-} Intercept;
-
-/**
- * Copy format. Currently supported:
- * - CSV
-*/
-typedef enum CopyFormat {
- INVALID,
- CSV,
-} CopyFormat;
-
-/**
- * Client requesting a COPY.
-*/
-typedef struct Copy {
- CopyFormat copy_format;
- char *table_name;
- int has_headers;
- char delimiter;
- int num_columns;
- char **columns;
-} Copy;
-
-/**
- * A copy row extracted from input,
- * with the shard it should go to.
- *
- * <div rustbindgen nodebug></div>
-*/
-typedef struct CopyRow {
- int len;
- char *data;
- int shard;
-} CopyRow;
+ * Wrapper around output by pg_query.
+ */
+typedef struct PdStatement {
+ int32_t version;
+ uint64_t len;
+ void *data;
+} PdStatement;
/**
- * Copy output.
- *
- * <div rustbindgen nodebug></div>
-*/
-typedef struct CopyOutput {
- int num_rows;
- CopyRow *rows;
- char *header;
-} CopyOutput;
-
-/*
- * Union of results a plugin can return.
- *
- * Route: FORWARD
- * Error: ERROR
- * Intercept: INTERCEPT
+ * Context on the database cluster configuration and the currently processed
+ * PostgreSQL statement.
*
+ * This struct is C FFI-safe and therefore uses C types. Use public methods to interact with it instead
+ * of reading the data directly.
*/
-typedef union RoutingOutput {
- Route route;
- Error error;
- Intercept intercept;
- Copy copy;
- CopyOutput copy_rows;
-} RoutingOutput;
-
-/*
- * Plugin output.
- *
- * This is returned by a plugin to communicate its routing decision.
+typedef struct PdRouterContext {
+ /** How many shards are configured. */
+ uint64_t shards;
+ /** Does the database cluster have replicas? `1` = `true`, `0` = `false`. */
+ uint8_t has_replicas;
+ /** Does the database cluster have a primary? `1` = `true`, `0` = `false`. */
+ uint8_t has_primary;
+ /** Is the query being executed inside a transaction? `1` = `true`, `0` = `false`. */
+ uint8_t in_transaction;
+ /** PgDog strongly believes this statement should go to a primary. `1` = `true`, `0` = `false`. */
+ uint8_t write_override;
+ /** pg_query generated Abstract Syntax Tree of the statement. */
+ PdStatement query;
+} PdRouterContext;
+
+/**
+ * Routing decision returned by the plugin.
*/
-typedef struct Output {
- RoutingDecision decision;
- RoutingOutput output;
-} Output;
-
-/**
- * Database role, e.g. primary or replica.
-*/
-typedef enum Role {
- PRIMARY = 1,
- REPLICA = 2,
-} Role;
-
-/**
- * Database configuration entry.
-*/
-typedef struct DatabaseConfig {
- int shard;
- Role role;
- char *host;
- int port;
-} DatabaseConfig;
-
-/**
- * Configuration for a database cluster
- * used to the serve a query passed to the plugin.
-*/
-typedef struct Config {
- int num_databases;
- DatabaseConfig *databases;
- /* Database name from pgdog.toml. */
- char *name;
- int shards;
-} Config;
-
-/**
- * Copy input.
-*/
-typedef struct CopyInput {
- int len;
- const char* data;
- char delimiter;
- int has_headers;
- int sharding_column;
-} CopyInput;
-
-/**
-* Routing input union passed to the plugin.
-*/
-typedef union RoutingInput {
- Query query;
- CopyInput copy;
-} RoutingInput;
-
-/**
- * Input type.
-*/
-typedef enum InputType {
- ROUTING_INPUT = 1,
- COPY_INPUT = 2,
-} InputType;
-
-/**
- * Plugin input.
-*/
-typedef struct Input {
- Config config;
- InputType input_type;
- RoutingInput input;
-} Input;
+ typedef struct PdRoute {
+ /** Which shard the query should go to.
+ *
+ * `-1` for all shards, `-2` for unknown, this setting is ignored.
+ */
+ int64_t shard;
+ /** Is the query a read and should go to a replica?
+ *
+ * `1` for `true`, `0` for `false`, `2` for unknown, this setting is ignored.
+ */
+ uint8_t read_write;
+ } PdRoute;
diff --git a/pgdog-plugin/src/ast.rs b/pgdog-plugin/src/ast.rs
new file mode 100644
index 00000000..567715ed
--- /dev/null
+++ b/pgdog-plugin/src/ast.rs
@@ -0,0 +1,132 @@
+//! Wrapper around `pg_query` protobuf-generated statement.
+//!
+//! This is passed through FFI safely by ensuring two conditions:
+//!
+//! 1. The version of the **Rust compiler** used to build the plugin is the same used to build PgDog
+//! 2. The version of the **`pg_query` library** used by the plugin is the same used by PgDog
+//!
+use std::{ffi::c_void, ops::Deref};
+
+use pg_query::protobuf::{ParseResult, RawStmt};
+
+use crate::bindings::PdStatement;
+
+impl PdStatement {
+ /// Create FFI binding from `pg_query` output.
+ ///
+ /// # Safety
+ ///
+ /// The reference must live for the entire time
+ /// this struct is used. This is _not_ checked by the compiler,
+ /// and is the responsibility of the caller.
+ ///
+ pub unsafe fn from_proto(value: &ParseResult) -> Self {
+ Self {
+ data: value.stmts.as_ptr() as *mut c_void,
+ version: value.version,
+ len: value.stmts.len() as u64,
+ }
+ }
+}
+
+/// Wrapper around [`pg_query::protobuf::ParseResult`], which allows
+/// the caller to use `pg_query` types and methods to inspect the statement.
+#[derive(Debug)]
+pub struct PdParseResult {
+ parse_result: Option<ParseResult>,
+ borrowed: bool,
+}
+
+impl Clone for PdParseResult {
+ /// Cloning the binding is safe. A new structure
+ /// will be created without any references to the original
+ /// reference.
+ fn clone(&self) -> Self {
+ Self {
+ parse_result: self.parse_result.clone(),
+ borrowed: false,
+ }
+ }
+}
+
+impl From<PdStatement> for PdParseResult {
+ /// Create the binding from a FFI-passed reference.
+ ///
+ /// SAFETY: Memory is owned by the caller.
+ ///
+ fn from(value: PdStatement) -> Self {
+ Self {
+ parse_result: Some(ParseResult {
+ version: value.version,
+ stmts: unsafe {
+ Vec::from_raw_parts(
+ value.data as *mut RawStmt,
+ value.len as usize,
+ value.len as usize,
+ )
+ },
+ }),
+ borrowed: true,
+ }
+ }
+}
+
+impl Drop for PdParseResult {
+ /// Drop the binding and forget the memory if the binding
+ /// is using the referenced struct. Otherwise, deallocate as normal.
+ fn drop(&mut self) {
+ if self.borrowed {
+ let parse_result = self.parse_result.take();
+ std::mem::forget(parse_result.unwrap().stmts);
+ }
+ }
+}
+
+impl Deref for PdParseResult {
+ type Target = ParseResult;
+
+ fn deref(&self) -> &Self::Target {
+ self.parse_result.as_ref().unwrap()
+ }
+}
+
+impl PdStatement {
+ /// Get the protobuf-wrapped PostgreSQL statement. The returned structure, [`PdParseResult`],
+ /// implements [`Deref`] to [`pg_query::protobuf::ParseResult`] and can be used to parse the
+ /// statement.
+ pub fn protobuf(&self) -> PdParseResult {
+ PdParseResult::from(*self)
+ }
+}
+
+#[cfg(test)]
+mod test {
+ use crate::pg_query::NodeEnum;
+
+ use super::*;
+
+ #[test]
+ fn test_ast() {
+ let ast = pg_query::parse("SELECT * FROM users WHERE id = $1").unwrap();
+ let ffi = unsafe { PdStatement::from_proto(&ast.protobuf) };
+ match ffi
+ .protobuf()
+ .stmts
+ .first()
+ .unwrap()
+ .stmt
+ .as_ref()
+ .unwrap()
+ .node
+ .as_ref()
+ .unwrap()
+ {
+ NodeEnum::SelectStmt(_) => (),
+ _ => {
+ panic!("not a select")
+ }
+ };
+
+ let _ = ffi.protobuf().clone();
+ }
+}
diff --git a/pgdog-plugin/src/bindings.rs b/pgdog-plugin/src/bindings.rs
index e30f41ec..7e86ce91 100644
--- a/pgdog-plugin/src/bindings.rs
+++ b/pgdog-plugin/src/bindings.rs
@@ -1,383 +1,418 @@
/* automatically generated by rust-bindgen 0.71.1 */
-#[doc = " Query parameter value."]
+pub const __WORDSIZE: u32 = 64;
+pub const __has_safe_buffers: u32 = 1;
+pub const __DARWIN_ONLY_64_BIT_INO_T: u32 = 1;
+pub const __DARWIN_ONLY_UNIX_CONFORMANCE: u32 = 1;
+pub const __DARWIN_ONLY_VERS_1050: u32 = 1;
+pub const __DARWIN_UNIX03: u32 = 1;
+pub const __DARWIN_64_BIT_INO_T: u32 = 1;
+pub const __DARWIN_VERS_1050: u32 = 1;
+pub const __DARWIN_NON_CANCELABLE: u32 = 0;
+pub const __DARWIN_SUF_EXTSN: &[u8; 14] = b"$DARWIN_EXTSN\0";
+pub const __DARWIN_C_ANSI: u32 = 4096;
+pub const __DARWIN_C_FULL: u32 = 900000;
+pub const __DARWIN_C_LEVEL: u32 = 900000;
+pub const __STDC_WANT_LIB_EXT1__: u32 = 1;
+pub const __DARWIN_NO_LONG_LONG: u32 = 0;
+pub const _DARWIN_FEATURE_64_BIT_INODE: u32 = 1;
+pub const _DARWIN_FEATURE_ONLY_64_BIT_INODE: u32 = 1;
+pub const _DARWIN_FEATURE_ONLY_VERS_1050: u32 = 1;
+pub const _DARWIN_FEATURE_ONLY_UNIX_CONFORMANCE: u32 = 1;
+pub const _DARWIN_FEATURE_UNIX_CONFORMANCE: u32 = 3;
+pub const __has_ptrcheck: u32 = 0;
+pub const USE_CLANG_TYPES: u32 = 0;
+pub const __PTHREAD_SIZE__: u32 = 8176;
+pub const __PTHREAD_ATTR_SIZE__: u32 = 56;
+pub const __PTHREAD_MUTEXATTR_SIZE__: u32 = 8;
+pub const __PTHREAD_MUTEX_SIZE__: u32 = 56;
+pub const __PTHREAD_CONDATTR_SIZE__: u32 = 8;
+pub const __PTHREAD_COND_SIZE__: u32 = 40;
+pub const __PTHREAD_ONCE_SIZE__: u32 = 8;
+pub const __PTHREAD_RWLOCK_SIZE__: u32 = 192;
+pub const __PTHREAD_RWLOCKATTR_SIZE__: u32 = 16;
+pub const INT8_MAX: u32 = 127;
+pub const INT16_MAX: u32 = 32767;
+pub const INT32_MAX: u32 = 2147483647;
+pub const INT64_MAX: u64 = 9223372036854775807;
+pub const INT8_MIN: i32 = -128;
+pub const INT16_MIN: i32 = -32768;
+pub const INT32_MIN: i32 = -2147483648;
+pub const INT64_MIN: i64 = -9223372036854775808;
+pub const UINT8_MAX: u32 = 255;
+pub const UINT16_MAX: u32 = 65535;
+pub const UINT32_MAX: u32 = 4294967295;
+pub const UINT64_MAX: i32 = -1;
+pub const INT_LEAST8_MIN: i32 = -128;
+pub const INT_LEAST16_MIN: i32 = -32768;
+pub const INT_LEAST32_MIN: i32 = -2147483648;
+pub const INT_LEAST64_MIN: i64 = -9223372036854775808;
+pub const INT_LEAST8_MAX: u32 = 127;
+pub const INT_LEAST16_MAX: u32 = 32767;
+pub const INT_LEAST32_MAX: u32 = 2147483647;
+pub const INT_LEAST64_MAX: u64 = 9223372036854775807;
+pub const UINT_LEAST8_MAX: u32 = 255;
+pub const UINT_LEAST16_MAX: u32 = 65535;
+pub const UINT_LEAST32_MAX: u32 = 4294967295;
+pub const UINT_LEAST64_MAX: i32 = -1;
+pub const INT_FAST8_MIN: i32 = -128;
+pub const INT_FAST16_MIN: i32 = -32768;
+pub const INT_FAST32_MIN: i32 = -2147483648;
+pub const INT_FAST64_MIN: i64 = -9223372036854775808;
+pub const INT_FAST8_MAX: u32 = 127;
+pub const INT_FAST16_MAX: u32 = 32767;
+pub const INT_FAST32_MAX: u32 = 2147483647;
+pub const INT_FAST64_MAX: u64 = 9223372036854775807;
+pub const UINT_FAST8_MAX: u32 = 255;
+pub const UINT_FAST16_MAX: u32 = 65535;
+pub const UINT_FAST32_MAX: u32 = 4294967295;
+pub const UINT_FAST64_MAX: i32 = -1;
+pub const INTPTR_MAX: u64 = 9223372036854775807;
+pub const INTPTR_MIN: i64 = -9223372036854775808;
+pub const UINTPTR_MAX: i32 = -1;
+pub const SIZE_MAX: i32 = -1;
+pub const RSIZE_MAX: i32 = -1;
+pub const WINT_MIN: i32 = -2147483648;
+pub const WINT_MAX: u32 = 2147483647;
+pub const SIG_ATOMIC_MIN: i32 = -2147483648;
+pub const SIG_ATOMIC_MAX: u32 = 2147483647;
+pub type wchar_t = ::std::os::raw::c_int;
+pub type max_align_t = f64;
+pub type int_least8_t = i8;
+pub type int_least16_t = i16;
+pub type int_least32_t = i32;
+pub type int_least64_t = i64;
+pub type uint_least8_t = u8;
+pub type uint_least16_t = u16;
+pub type uint_least32_t = u32;
+pub type uint_least64_t = u64;
+pub type int_fast8_t = i8;
+pub type int_fast16_t = i16;
+pub type int_fast32_t = i32;
+pub type int_fast64_t = i64;
+pub type uint_fast8_t = u8;
+pub type uint_fast16_t = u16;
+pub type uint_fast32_t = u32;
+pub type uint_fast64_t = u64;
+pub type __int8_t = ::std::os::raw::c_schar;
+pub type __uint8_t = ::std::os::raw::c_uchar;
+pub type __int16_t = ::std::os::raw::c_short;
+pub type __uint16_t = ::std::os::raw::c_ushort;
+pub type __int32_t = ::std::os::raw::c_int;
+pub type __uint32_t = ::std::os::raw::c_uint;
+pub type __int64_t = ::std::os::raw::c_longlong;
+pub type __uint64_t = ::std::os::raw::c_ulonglong;
+pub type __darwin_intptr_t = ::std::os::raw::c_long;
+pub type __darwin_natural_t = ::std::os::raw::c_uint;
+pub type __darwin_ct_rune_t = ::std::os::raw::c_int;
#[repr(C)]
-#[derive(Debug, Copy, Clone)]
-pub struct Parameter {
- pub len: ::std::os::raw::c_int,
- pub data: *const ::std::os::raw::c_char,
- pub format: ::std::os::raw::c_int,
+#[derive(Copy, Clone)]
+pub union __mbstate_t {
+ pub __mbstate8: [::std::os::raw::c_char; 128usize],
+ pub _mbstateL: ::std::os::raw::c_longlong,
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of Parameter"][::std::mem::size_of::<Parameter>() - 24usize];
- ["Alignment of Parameter"][::std::mem::align_of::<Parameter>() - 8usize];
- ["Offset of field: Parameter::len"][::std::mem::offset_of!(Parameter, len) - 0usize];
- ["Offset of field: Parameter::data"][::std::mem::offset_of!(Parameter, data) - 8usize];
- ["Offset of field: Parameter::format"][::std::mem::offset_of!(Parameter, format) - 16usize];
+ ["Size of __mbstate_t"][::std::mem::size_of::<__mbstate_t>() - 128usize];
+ ["Alignment of __mbstate_t"][::std::mem::align_of::<__mbstate_t>() - 8usize];
+ ["Offset of field: __mbstate_t::__mbstate8"]
+ [::std::mem::offset_of!(__mbstate_t, __mbstate8) - 0usize];
+ ["Offset of field: __mbstate_t::_mbstateL"]
+ [::std::mem::offset_of!(__mbstate_t, _mbstateL) - 0usize];
};
+pub type __darwin_mbstate_t = __mbstate_t;
+pub type __darwin_ptrdiff_t = ::std::os::raw::c_long;
+pub type __darwin_size_t = ::std::os::raw::c_ulong;
+pub type __darwin_va_list = __builtin_va_list;
+pub type __darwin_wchar_t = ::std::os::raw::c_int;
+pub type __darwin_rune_t = __darwin_wchar_t;
+pub type __darwin_wint_t = ::std::os::raw::c_int;
+pub type __darwin_clock_t = ::std::os::raw::c_ulong;
+pub type __darwin_socklen_t = __uint32_t;
+pub type __darwin_ssize_t = ::std::os::raw::c_long;
+pub type __darwin_time_t = ::std::os::raw::c_long;
+pub type __darwin_blkcnt_t = __int64_t;
+pub type __darwin_blksize_t = __int32_t;
+pub type __darwin_dev_t = __int32_t;
+pub type __darwin_fsblkcnt_t = ::std::os::raw::c_uint;
+pub type __darwin_fsfilcnt_t = ::std::os::raw::c_uint;
+pub type __darwin_gid_t = __uint32_t;
+pub type __darwin_id_t = __uint32_t;
+pub type __darwin_ino64_t = __uint64_t;
+pub type __darwin_ino_t = __darwin_ino64_t;
+pub type __darwin_mach_port_name_t = __darwin_natural_t;
+pub type __darwin_mach_port_t = __darwin_mach_port_name_t;
+pub type __darwin_mode_t = __uint16_t;
+pub type __darwin_off_t = __int64_t;
+pub type __darwin_pid_t = __int32_t;
+pub type __darwin_sigset_t = __uint32_t;
+pub type __darwin_suseconds_t = __int32_t;
+pub type __darwin_uid_t = __uint32_t;
+pub type __darwin_useconds_t = __uint32_t;
+pub type __darwin_uuid_t = [::std::os::raw::c_uchar; 16usize];
+pub type __darwin_uuid_string_t = [::std::os::raw::c_char; 37usize];
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct Query {
- pub len: ::std::os::raw::c_int,
- pub query: *const ::std::os::raw::c_char,
- pub num_parameters: ::std::os::raw::c_int,
- pub parameters: *const Parameter,
+pub struct __darwin_pthread_handler_rec {
+ pub __routine: ::std::option::Option<unsafe extern "C" fn(arg1: *mut ::std::os::raw::c_void)>,
+ pub __arg: *mut ::std::os::raw::c_void,
+ pub __next: *mut __darwin_pthread_handler_rec,
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of Query"][::std::mem::size_of::<Query>() - 32usize];
- ["Alignment of Query"][::std::mem::align_of::<Query>() - 8usize];
- ["Offset of field: Query::len"][::std::mem::offset_of!(Query, len) - 0usize];
- ["Offset of field: Query::query"][::std::mem::offset_of!(Query, query) - 8usize];
- ["Offset of field: Query::num_parameters"]
- [::std::mem::offset_of!(Query, num_parameters) - 16usize];
- ["Offset of field: Query::parameters"][::std::mem::offset_of!(Query, parameters) - 24usize];
+ ["Size of __darwin_pthread_handler_rec"]
+ [::std::mem::size_of::<__darwin_pthread_handler_rec>() - 24usize];
+ ["Alignment of __darwin_pthread_handler_rec"]
+ [::std::mem::align_of::<__darwin_pthread_handler_rec>() - 8usize];
+ ["Offset of field: __darwin_pthread_handler_rec::__routine"]
+ [::std::mem::offset_of!(__darwin_pthread_handler_rec, __routine) - 0usize];
+ ["Offset of field: __darwin_pthread_handler_rec::__arg"]
+ [::std::mem::offset_of!(__darwin_pthread_handler_rec, __arg) - 8usize];
+ ["Offset of field: __darwin_pthread_handler_rec::__next"]
+ [::std::mem::offset_of!(__darwin_pthread_handler_rec, __next) - 16usize];
};
-pub const Affinity_READ: Affinity = 1;
-pub const Affinity_WRITE: Affinity = 2;
-pub const Affinity_TRANSACTION_START: Affinity = 3;
-pub const Affinity_TRANSACTION_END: Affinity = 4;
-pub const Affinity_UNKNOWN: Affinity = -1;
-#[doc = " The query is a read or a write.\n In case the plugin isn't able to figure it out, it can return UNKNOWN and\n pgDog will ignore the plugin's decision."]
-pub type Affinity = ::std::os::raw::c_int;
-pub const Shard_ANY: Shard = -1;
-pub const Shard_ALL: Shard = -2;
-#[doc = " In case the plugin doesn't know which shard to route the\n the query, it can decide to route it to any shard or to all\n shards. All shard queries return a result assembled by pgDog."]
-pub type Shard = ::std::os::raw::c_int;
-pub const OrderByDirection_ASCENDING: OrderByDirection = 0;
-pub const OrderByDirection_DESCENDING: OrderByDirection = 1;
-pub type OrderByDirection = ::std::os::raw::c_uint;
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct OrderBy {
- pub column_name: *mut ::std::os::raw::c_char,
- pub column_index: ::std::os::raw::c_int,
- pub direction: OrderByDirection,
+pub struct _opaque_pthread_attr_t {
+ pub __sig: ::std::os::raw::c_long,
+ pub __opaque: [::std::os::raw::c_char; 56usize],
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of OrderBy"][::std::mem::size_of::<OrderBy>() - 16usize];
- ["Alignment of OrderBy"][::std::mem::align_of::<OrderBy>() - 8usize];
- ["Offset of field: OrderBy::column_name"]
- [::std::mem::offset_of!(OrderBy, column_name) - 0usize];
- ["Offset of field: OrderBy::column_index"]
- [::std::mem::offset_of!(OrderBy, column_index) - 8usize];
- ["Offset of field: OrderBy::direction"][::std::mem::offset_of!(OrderBy, direction) - 12usize];
+ ["Size of _opaque_pthread_attr_t"][::std::mem::size_of::<_opaque_pthread_attr_t>() - 64usize];
+ ["Alignment of _opaque_pthread_attr_t"]
+ [::std::mem::align_of::<_opaque_pthread_attr_t>() - 8usize];
+ ["Offset of field: _opaque_pthread_attr_t::__sig"]
+ [::std::mem::offset_of!(_opaque_pthread_attr_t, __sig) - 0usize];
+ ["Offset of field: _opaque_pthread_attr_t::__opaque"]
+ [::std::mem::offset_of!(_opaque_pthread_attr_t, __opaque) - 8usize];
};
-#[doc = " Route the query should take."]
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct Route {
- pub affinity: Affinity,
- pub shard: ::std::os::raw::c_int,
- pub num_order_by: ::std::os::raw::c_int,
- pub order_by: *mut OrderBy,
+pub struct _opaque_pthread_cond_t {
+ pub __sig: ::std::os::raw::c_long,
+ pub __opaque: [::std::os::raw::c_char; 40usize],
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of Route"][::std::mem::size_of::<Route>() - 24usize];
- ["Alignment of Route"][::std::mem::align_of::<Route>() - 8usize];
- ["Offset of field: Route::affinity"][::std::mem::offset_of!(Route, affinity) - 0usize];
- ["Offset of field: Route::shard"][::std::mem::offset_of!(Route, shard) - 4usize];
- ["Offset of field: Route::num_order_by"][::std::mem::offset_of!(Route, num_order_by) - 8usize];
- ["Offset of field: Route::order_by"][::std::mem::offset_of!(Route, order_by) - 16usize];
+ ["Size of _opaque_pthread_cond_t"][::std::mem::size_of::<_opaque_pthread_cond_t>() - 48usize];
+ ["Alignment of _opaque_pthread_cond_t"]
+ [::std::mem::align_of::<_opaque_pthread_cond_t>() - 8usize];
+ ["Offset of field: _opaque_pthread_cond_t::__sig"]
+ [::std::mem::offset_of!(_opaque_pthread_cond_t, __sig) - 0usize];
+ ["Offset of field: _opaque_pthread_cond_t::__opaque"]
+ [::std::mem::offset_of!(_opaque_pthread_cond_t, __opaque) - 8usize];
};
-pub const RoutingDecision_FORWARD: RoutingDecision = 1;
-pub const RoutingDecision_REWRITE: RoutingDecision = 2;
-pub const RoutingDecision_ERROR: RoutingDecision = 3;
-pub const RoutingDecision_INTERCEPT: RoutingDecision = 4;
-pub const RoutingDecision_NO_DECISION: RoutingDecision = 5;
-pub const RoutingDecision_COPY: RoutingDecision = 6;
-pub const RoutingDecision_COPY_ROWS: RoutingDecision = 7;
-#[doc = " The routing decision the plugin makes based on the query contents.\n\n FORWARD: The query is forwarded to a shard. Which shard (and whether it's a replica\n or a primary) is decided by the plugin output.\n REWRITE: The query text is rewritten. The plugin outputs new query text.\n ERROR: The query is denied and the plugin returns an error instead. This error is sent\n to the client.\n INTERCEPT: The query is intercepted and the plugin returns rows instead. These rows\nare sent to the client and the original query is never sent to a backend server.\n NO_DECISION: The plugin doesn't care about this query. The output is ignored by pgDog and the next\nplugin in the chain is attempted.\n COPY: Client is sending over a COPY statement."]
-pub type RoutingDecision = ::std::os::raw::c_uint;
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct Error {
- pub severity: *mut ::std::os::raw::c_char,
- pub code: *mut ::std::os::raw::c_char,
- pub message: *mut ::std::os::raw::c_char,
- pub detail: *mut ::std::os::raw::c_char,
+pub struct _opaque_pthread_condattr_t {
+ pub __sig: ::std::os::raw::c_long,
+ pub __opaque: [::std::os::raw::c_char; 8usize],
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of Error"][::std::mem::size_of::<Error>() - 32usize];
- ["Alignment of Error"][::std::mem::align_of::<Error>() - 8usize];
- ["Offset of field: Error::severity"][::std::mem::offset_of!(Error, severity) - 0usize];
- ["Offset of field: Error::code"][::std::mem::offset_of!(Error, code) - 8usize];
- ["Offset of field: Error::message"][::std::mem::offset_of!(Error, message) - 16usize];
- ["Offset of field: Error::detail"][::std::mem::offset_of!(Error, detail) - 24usize];
+ ["Size of _opaque_pthread_condattr_t"]
+ [::std::mem::size_of::<_opaque_pthread_condattr_t>() - 16usize];
+ ["Alignment of _opaque_pthread_condattr_t"]
+ [::std::mem::align_of::<_opaque_pthread_condattr_t>() - 8usize];
+ ["Offset of field: _opaque_pthread_condattr_t::__sig"]
+ [::std::mem::offset_of!(_opaque_pthread_condattr_t, __sig) - 0usize];
+ ["Offset of field: _opaque_pthread_condattr_t::__opaque"]
+ [::std::mem::offset_of!(_opaque_pthread_condattr_t, __opaque) - 8usize];
};
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct RowColumn {
- pub length: ::std::os::raw::c_int,
- pub data: *mut ::std::os::raw::c_char,
+pub struct _opaque_pthread_mutex_t {
+ pub __sig: ::std::os::raw::c_long,
+ pub __opaque: [::std::os::raw::c_char; 56usize],
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of RowColumn"][::std::mem::size_of::<RowColumn>() - 16usize];
- ["Alignment of RowColumn"][::std::mem::align_of::<RowColumn>() - 8usize];
- ["Offset of field: RowColumn::length"][::std::mem::offset_of!(RowColumn, length) - 0usize];
- ["Offset of field: RowColumn::data"][::std::mem::offset_of!(RowColumn, data) - 8usize];
+ ["Size of _opaque_pthread_mutex_t"][::std::mem::size_of::<_opaque_pthread_mutex_t>() - 64usize];
+ ["Alignment of _opaque_pthread_mutex_t"]
+ [::std::mem::align_of::<_opaque_pthread_mutex_t>() - 8usize];
+ ["Offset of field: _opaque_pthread_mutex_t::__sig"]
+ [::std::mem::offset_of!(_opaque_pthread_mutex_t, __sig) - 0usize];
+ ["Offset of field: _opaque_pthread_mutex_t::__opaque"]
+ [::std::mem::offset_of!(_opaque_pthread_mutex_t, __opaque) - 8usize];
};
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct Row {
- pub num_columns: ::std::os::raw::c_int,
- pub columns: *mut RowColumn,
+pub struct _opaque_pthread_mutexattr_t {
+ pub __sig: ::std::os::raw::c_long,
+ pub __opaque: [::std::os::raw::c_char; 8usize],
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of Row"][::std::mem::size_of::<Row>() - 16usize];
- ["Alignment of Row"][::std::mem::align_of::<Row>() - 8usize];
- ["Offset of field: Row::num_columns"][::std::mem::offset_of!(Row, num_columns) - 0usize];
- ["Offset of field: Row::columns"][::std::mem::offset_of!(Row, columns) - 8usize];
+ ["Size of _opaque_pthread_mutexattr_t"]
+ [::std::mem::size_of::<_opaque_pthread_mutexattr_t>() - 16usize];
+ ["Alignment of _opaque_pthread_mutexattr_t"]
+ [::std::mem::align_of::<_opaque_pthread_mutexattr_t>() - 8usize];
+ ["Offset of field: _opaque_pthread_mutexattr_t::__sig"]
+ [::std::mem::offset_of!(_opaque_pthread_mutexattr_t, __sig) - 0usize];
+ ["Offset of field: _opaque_pthread_mutexattr_t::__opaque"]
+ [::std::mem::offset_of!(_opaque_pthread_mutexattr_t, __opaque) - 8usize];
};
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct RowDescriptionColumn {
- pub len: ::std::os::raw::c_int,
- pub name: *mut ::std::os::raw::c_char,
- pub oid: ::std::os::raw::c_int,
+pub struct _opaque_pthread_once_t {
+ pub __sig: ::std::os::raw::c_long,
+ pub __opaque: [::std::os::raw::c_char; 8usize],
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of RowDescriptionColumn"][::std::mem::size_of::<RowDescriptionColumn>() - 24usize];
- ["Alignment of RowDescriptionColumn"][::std::mem::align_of::<RowDescriptionColumn>() - 8usize];
- ["Offset of field: RowDescriptionColumn::len"]
- [::std::mem::offset_of!(RowDescriptionColumn, len) - 0usize];
- ["Offset of field: RowDescriptionColumn::name"]
- [::std::mem::offset_of!(RowDescriptionColumn, name) - 8usize];
- ["Offset of field: RowDescriptionColumn::oid"]
- [::std::mem::offset_of!(RowDescriptionColumn, oid) - 16usize];
+ ["Size of _opaque_pthread_once_t"][::std::mem::size_of::<_opaque_pthread_once_t>() - 16usize];
+ ["Alignment of _opaque_pthread_once_t"]
+ [::std::mem::align_of::<_opaque_pthread_once_t>() - 8usize];
+ ["Offset of field: _opaque_pthread_once_t::__sig"]
+ [::std::mem::offset_of!(_opaque_pthread_once_t, __sig) - 0usize];
+ ["Offset of field: _opaque_pthread_once_t::__opaque"]
+ [::std::mem::offset_of!(_opaque_pthread_once_t, __opaque) - 8usize];
};
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct RowDescription {
- pub num_columns: ::std::os::raw::c_int,
- pub columns: *mut RowDescriptionColumn,
+pub struct _opaque_pthread_rwlock_t {
+ pub __sig: ::std::os::raw::c_long,
+ pub __opaque: [::std::os::raw::c_char; 192usize],
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of RowDescription"][::std::mem::size_of::<RowDescription>() - 16usize];
- ["Alignment of RowDescription"][::std::mem::align_of::<RowDescription>() - 8usize];
- ["Offset of field: RowDescription::num_columns"]
- [::std::mem::offset_of!(RowDescription, num_columns) - 0usize];
- ["Offset of field: RowDescription::columns"]
- [::std::mem::offset_of!(RowDescription, columns) - 8usize];
+ ["Size of _opaque_pthread_rwlock_t"]
+ [::std::mem::size_of::<_opaque_pthread_rwlock_t>() - 200usize];
+ ["Alignment of _opaque_pthread_rwlock_t"]
+ [::std::mem::align_of::<_opaque_pthread_rwlock_t>() - 8usize];
+ ["Offset of field: _opaque_pthread_rwlock_t::__sig"]
+ [::std::mem::offset_of!(_opaque_pthread_rwlock_t, __sig) - 0usize];
+ ["Offset of field: _opaque_pthread_rwlock_t::__opaque"]
+ [::std::mem::offset_of!(_opaque_pthread_rwlock_t, __opaque) - 8usize];
};
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct Intercept {
- pub row_description: RowDescription,
- pub num_rows: ::std::os::raw::c_int,
- pub rows: *mut Row,
+pub struct _opaque_pthread_rwlockattr_t {
+ pub __sig: ::std::os::raw::c_long,
+ pub __opaque: [::std::os::raw::c_char; 16usize],
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of Intercept"][::std::mem::size_of::<Intercept>() - 32usize];
- ["Alignment of Intercept"][::std::mem::align_of::<Intercept>() - 8usize];
- ["Offset of field: Intercept::row_description"]
- [::std::mem::offset_of!(Intercept, row_description) - 0usize];
- ["Offset of field: Intercept::num_rows"][::std::mem::offset_of!(Intercept, num_rows) - 16usize];
- ["Offset of field: Intercept::rows"][::std::mem::offset_of!(Intercept, rows) - 24usize];
+ ["Size of _opaque_pthread_rwlockattr_t"]
+ [::std::mem::size_of::<_opaque_pthread_rwlockattr_t>() - 24usize];
+ ["Alignment of _opaque_pthread_rwlockattr_t"]
+ [::std::mem::align_of::<_opaque_pthread_rwlockattr_t>() - 8usize];
+ ["Offset of field: _opaque_pthread_rwlockattr_t::__sig"]
+ [::std::mem::offset_of!(_opaque_pthread_rwlockattr_t, __sig) - 0usize];
+ ["Offset of field: _opaque_pthread_rwlockattr_t::__opaque"]
+ [::std::mem::offset_of!(_opaque_pthread_rwlockattr_t, __opaque) - 8usize];
};
-pub const CopyFormat_INVALID: CopyFormat = 0;
-pub const CopyFormat_CSV: CopyFormat = 1;
-#[doc = " Copy format. Currently supported:\n - CSV"]
-pub type CopyFormat = ::std::os::raw::c_uint;
-#[doc = " Client requesting a COPY."]
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct Copy {
- pub copy_format: CopyFormat,
- pub table_name: *mut ::std::os::raw::c_char,
- pub has_headers: ::std::os::raw::c_int,
- pub delimiter: ::std::os::raw::c_char,
- pub num_columns: ::std::os::raw::c_int,
- pub columns: *mut *mut ::std::os::raw::c_char,
-}
-#[allow(clippy::unnecessary_operation, clippy::identity_op)]
-const _: () = {
- ["Size of Copy"][::std::mem::size_of::<Copy>() - 40usize];
- ["Alignment of Copy"][::std::mem::align_of::<Copy>() - 8usize];
- ["Offset of field: Copy::copy_format"][::std::mem::offset_of!(Copy, copy_format) - 0usize];
- ["Offset of field: Copy::table_name"][::std::mem::offset_of!(Copy, table_name) - 8usize];
- ["Offset of field: Copy::has_headers"][::std::mem::offset_of!(Copy, has_headers) - 16usize];
- ["Offset of field: Copy::delimiter"][::std::mem::offset_of!(Copy, delimiter) - 20usize];
- ["Offset of field: Copy::num_columns"][::std::mem::offset_of!(Copy, num_columns) - 24usize];
- ["Offset of field: Copy::columns"][::std::mem::offset_of!(Copy, columns) - 32usize];
-};
-#[doc = " A copy row extracted from input,\n with the shard it should go to.\n\n <div rustbindgen nodebug></div>"]
-#[repr(C)]
-#[derive(Copy, Clone)]
-pub struct CopyRow {
- pub len: ::std::os::raw::c_int,
- pub data: *mut ::std::os::raw::c_char,
- pub shard: ::std::os::raw::c_int,
-}
-#[allow(clippy::unnecessary_operation, clippy::identity_op)]
-const _: () = {
- ["Size of CopyRow"][::std::mem::size_of::<CopyRow>() - 24usize];
- ["Alignment of CopyRow"][::std::mem::align_of::<CopyRow>() - 8usize];
- ["Offset of field: CopyRow::len"][::std::mem::offset_of!(CopyRow, len) - 0usize];
- ["Offset of field: CopyRow::data"][::std::mem::offset_of!(CopyRow, data) - 8usize];
- ["Offset of field: CopyRow::shard"][::std::mem::offset_of!(CopyRow, shard) - 16usize];
-};
-#[doc = " Copy output.\n\n <div rustbindgen nodebug></div>"]
-#[repr(C)]
-#[derive(Copy, Clone)]
-pub struct CopyOutput {
- pub num_rows: ::std::os::raw::c_int,
- pub rows: *mut CopyRow,
- pub header: *mut ::std::os::raw::c_char,
-}
-#[allow(clippy::unnecessary_operation, clippy::identity_op)]
-const _: () = {
- ["Size of CopyOutput"][::std::mem::size_of::<CopyOutput>() - 24usize];
- ["Alignment of CopyOutput"][::std::mem::align_of::<CopyOutput>() - 8usize];
- ["Offset of field: CopyOutput::num_rows"]
- [::std::mem::offset_of!(CopyOutput, num_rows) - 0usize];
- ["Offset of field: CopyOutput::rows"][::std::mem::offset_of!(CopyOutput, rows) - 8usize];
- ["Offset of field: CopyOutput::header"][::std::mem::offset_of!(CopyOutput, header) - 16usize];
-};
-#[repr(C)]
-#[derive(Copy, Clone)]
-pub union RoutingOutput {
- pub route: Route,
- pub error: Error,
- pub intercept: Intercept,
- pub copy: Copy,
- pub copy_rows: CopyOutput,
-}
-#[allow(clippy::unnecessary_operation, clippy::identity_op)]
-const _: () = {
- ["Size of RoutingOutput"][::std::mem::size_of::<RoutingOutput>() - 40usize];
- ["Alignment of RoutingOutput"][::std::mem::align_of::<RoutingOutput>() - 8usize];
- ["Offset of field: RoutingOutput::route"]
- [::std::mem::offset_of!(RoutingOutput, route) - 0usize];
- ["Offset of field: RoutingOutput::error"]
- [::std::mem::offset_of!(RoutingOutput, error) - 0usize];
- ["Offset of field: RoutingOutput::intercept"]
- [::std::mem::offset_of!(RoutingOutput, intercept) - 0usize];
- ["Offset of field: RoutingOutput::copy"][::std::mem::offset_of!(RoutingOutput, copy) - 0usize];
- ["Offset of field: RoutingOutput::copy_rows"]
- [::std::mem::offset_of!(RoutingOutput, copy_rows) - 0usize];
-};
-#[repr(C)]
-#[derive(Copy, Clone)]
-pub struct Output {
- pub decision: RoutingDecision,
- pub output: RoutingOutput,
+pub struct _opaque_pthread_t {
+ pub __sig: ::std::os::raw::c_long,
+ pub __cleanup_stack: *mut __darwin_pthread_handler_rec,
+ pub __opaque: [::std::os::raw::c_char; 8176usize],
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of Output"][::std::mem::size_of::<Output>() - 48usize];
- ["Alignment of Output"][::std::mem::align_of::<Output>() - 8usize];
- ["Offset of field: Output::decision"][::std::mem::offset_of!(Output, decision) - 0usize];
- ["Offset of field: Output::output"][::std::mem::offset_of!(Output, output) - 8usize];
+ ["Size of _opaque_pthread_t"][::std::mem::size_of::<_opaque_pthread_t>() - 8192usize];
+ ["Alignment of _opaque_pthread_t"][::std::mem::align_of::<_opaque_pthread_t>() - 8usize];
+ ["Offset of field: _opaque_pthread_t::__sig"]
+ [::std::mem::offset_of!(_opaque_pthread_t, __sig) - 0usize];
+ ["Offset of field: _opaque_pthread_t::__cleanup_stack"]
+ [::std::mem::offset_of!(_opaque_pthread_t, __cleanup_stack) - 8usize];
+ ["Offset of field: _opaque_pthread_t::__opaque"]
+ [::std::mem::offset_of!(_opaque_pthread_t, __opaque) - 16usize];
};
-pub const Role_PRIMARY: Role = 1;
-pub const Role_REPLICA: Role = 2;
-#[doc = " Database role, e.g. primary or replica."]
-pub type Role = ::std::os::raw::c_uint;
-#[doc = " Database configuration entry."]
+pub type __darwin_pthread_attr_t = _opaque_pthread_attr_t;
+pub type __darwin_pthread_cond_t = _opaque_pthread_cond_t;
+pub type __darwin_pthread_condattr_t = _opaque_pthread_condattr_t;
+pub type __darwin_pthread_key_t = ::std::os::raw::c_ulong;
+pub type __darwin_pthread_mutex_t = _opaque_pthread_mutex_t;
+pub type __darwin_pthread_mutexattr_t = _opaque_pthread_mutexattr_t;
+pub type __darwin_pthread_once_t = _opaque_pthread_once_t;
+pub type __darwin_pthread_rwlock_t = _opaque_pthread_rwlock_t;
+pub type __darwin_pthread_rwlockattr_t = _opaque_pthread_rwlockattr_t;
+pub type __darwin_pthread_t = *mut _opaque_pthread_t;
+pub type intmax_t = ::std::os::raw::c_long;
+pub type uintmax_t = ::std::os::raw::c_ulong;
+#[doc = " Wrapper around Rust's [`&str`], without allocating memory, unlike [`std::ffi::CString`].\n The caller must use it as a Rust string. This is not a C-string."]
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct DatabaseConfig {
- pub shard: ::std::os::raw::c_int,
- pub role: Role,
- pub host: *mut ::std::os::raw::c_char,
- pub port: ::std::os::raw::c_int,
+pub struct PdStr {
+ pub len: usize,
+ pub data: *mut ::std::os::raw::c_void,
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of DatabaseConfig"][::std::mem::size_of::<DatabaseConfig>() - 24usize];
- ["Alignment of DatabaseConfig"][::std::mem::align_of::<DatabaseConfig>() - 8usize];
- ["Offset of field: DatabaseConfig::shard"]
- [::std::mem::offset_of!(DatabaseConfig, shard) - 0usize];
- ["Offset of field: DatabaseConfig::role"]
- [::std::mem::offset_of!(DatabaseConfig, role) - 4usize];
- ["Offset of field: DatabaseConfig::host"]
- [::std::mem::offset_of!(DatabaseConfig, host) - 8usize];
- ["Offset of field: DatabaseConfig::port"]
- [::std::mem::offset_of!(DatabaseConfig, port) - 16usize];
+ ["Size of PdStr"][::std::mem::size_of::<PdStr>() - 16usize];
+ ["Alignment of PdStr"][::std::mem::align_of::<PdStr>() - 8usize];
+ ["Offset of field: PdStr::len"][::std::mem::offset_of!(PdStr, len) - 0usize];
+ ["Offset of field: PdStr::data"][::std::mem::offset_of!(PdStr, data) - 8usize];
};
-#[doc = " Configuration for a database cluster\n used to the serve a query passed to the plugin."]
+#[doc = " Wrapper around Rust's [`&str`], without allocating memory, unlike [`std::ffi::CString`].\n The caller must use it as a Rust string. This is not a C-string."]
+pub type RustString = PdStr;
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct Config {
- pub num_databases: ::std::os::raw::c_int,
- pub databases: *mut DatabaseConfig,
- pub name: *mut ::std::os::raw::c_char,
- pub shards: ::std::os::raw::c_int,
+pub struct PdStatement {
+ pub version: i32,
+ pub len: u64,
+ pub data: *mut ::std::os::raw::c_void,
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of Config"][::std::mem::size_of::<Config>() - 32usize];
- ["Alignment of Config"][::std::mem::align_of::<Config>() - 8usize];
- ["Offset of field: Config::num_databases"]
- [::std::mem::offset_of!(Config, num_databases) - 0usize];
- ["Offset of field: Config::databases"][::std::mem::offset_of!(Config, databases) - 8usize];
- ["Offset of field: Config::name"][::std::mem::offset_of!(Config, name) - 16usize];
- ["Offset of field: Config::shards"][::std::mem::offset_of!(Config, shards) - 24usize];
+ ["Size of PdStatement"][::std::mem::size_of::<PdStatement>() - 24usize];
+ ["Alignment of PdStatement"][::std::mem::align_of::<PdStatement>() - 8usize];
+ ["Offset of field: PdStatement::version"]
+ [::std::mem::offset_of!(PdStatement, version) - 0usize];
+ ["Offset of field: PdStatement::len"][::std::mem::offset_of!(PdStatement, len) - 8usize];
+ ["Offset of field: PdStatement::data"][::std::mem::offset_of!(PdStatement, data) - 16usize];
};
-#[doc = " Copy input."]
+#[doc = " Context on the database cluster configuration and the currently processed\n PostgreSQL statement.\n\n This struct is C FFI-safe and therefore uses C types. Use public methods to interact with it instead\n of reading the data directly."]
#[repr(C)]
#[derive(Debug, Copy, Clone)]
-pub struct CopyInput {
- pub len: ::std::os::raw::c_int,
- pub data: *const ::std::os::raw::c_char,
- pub delimiter: ::std::os::raw::c_char,
- pub has_headers: ::std::os::raw::c_int,
- pub sharding_column: ::std::os::raw::c_int,
+pub struct PdRouterContext {
+ #[doc = " How many shards are configured."]
+ pub shards: u64,
+ #[doc = " Does the database cluster have replicas? `1` = `true`, `0` = `false`."]
+ pub has_replicas: u8,
+ #[doc = " Does the database cluster have a primary? `1` = `true`, `0` = `false`."]
+ pub has_primary: u8,
+ #[doc = " Is the query being executed inside a transaction? `1` = `true`, `0` = `false`."]
+ pub in_transaction: u8,
+ #[doc = " PgDog strongly believes this statement should go to a primary. `1` = `true`, `0` = `false`."]
+ pub write_override: u8,
+ #[doc = " pg_query generated Abstract Syntax Tree of the statement."]
+ pub query: PdStatement,
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of CopyInput"][::std::mem::size_of::<CopyInput>() - 32usize];
- ["Alignment of CopyInput"][::std::mem::align_of::<CopyInput>() - 8usize];
- ["Offset of field: CopyInput::len"][::std::mem::offset_of!(CopyInput, len) - 0usize];
- ["Offset of field: CopyInput::data"][::std::mem::offset_of!(CopyInput, data) - 8usize];
- ["Offset of field: CopyInput::delimiter"]
- [::std::mem::offset_of!(CopyInput, delimiter) - 16usize];
- ["Offset of field: CopyInput::has_headers"]
- [::std::mem::offset_of!(CopyInput, has_headers) - 20usize];
- ["Offset of field: CopyInput::sharding_column"]
- [::std::mem::offset_of!(CopyInput, sharding_column) - 24usize];
+ ["Size of PdRouterContext"][::std::mem::size_of::<PdRouterContext>() - 40usize];
+ ["Alignment of PdRouterContext"][::std::mem::align_of::<PdRouterContext>() - 8usize];
+ ["Offset of field: PdRouterContext::shards"]
+ [::std::mem::offset_of!(PdRouterContext, shards) - 0usize];
+ ["Offset of field: PdRouterContext::has_replicas"]
+ [::std::mem::offset_of!(PdRouterContext, has_replicas) - 8usize];
+ ["Offset of field: PdRouterContext::has_primary"]
+ [::std::mem::offset_of!(PdRouterContext, has_primary) - 9usize];
+ ["Offset of field: PdRouterContext::in_transaction"]
+ [::std::mem::offset_of!(PdRouterContext, in_transaction) - 10usize];
+ ["Offset of field: PdRouterContext::write_override"]
+ [::std::mem::offset_of!(PdRouterContext, write_override) - 11usize];
+ ["Offset of field: PdRouterContext::query"]
+ [::std::mem::offset_of!(PdRouterContext, query) - 16usize];
};
-#[doc = " Routing input union passed to the plugin."]
+#[doc = " Routing decision returned by the plugin."]
#[repr(C)]
-#[derive(Copy, Clone)]
-pub union RoutingInput {
- pub query: Query,
- pub copy: CopyInput,
-}
-#[allow(clippy::unnecessary_operation, clippy::identity_op)]
-const _: () = {
- ["Size of RoutingInput"][::std::mem::size_of::<RoutingInput>() - 32usize];
- ["Alignment of RoutingInput"][::std::mem::align_of::<RoutingInput>() - 8usize];
- ["Offset of field: RoutingInput::query"][::std::mem::offset_of!(RoutingInput, query) - 0usize];
- ["Offset of field: RoutingInput::copy"][::std::mem::offset_of!(RoutingInput, copy) - 0usize];
-};
-pub const InputType_ROUTING_INPUT: InputType = 1;
-pub const InputType_COPY_INPUT: InputType = 2;
-#[doc = " Input type."]
-pub type InputType = ::std::os::raw::c_uint;
-#[doc = " Plugin input."]
-#[repr(C)]
-#[derive(Copy, Clone)]
-pub struct Input {
- pub config: Config,
- pub input_type: InputType,
- pub input: RoutingInput,
+#[derive(Debug, Copy, Clone)]
+pub struct PdRoute {
+ #[doc = " Which shard the query should go to.\n\n `-1` for all shards, `-2` for unknown, this setting is ignored."]
+ pub shard: i64,
+ #[doc = " Is the query a read and should go to a replica?\n\n `1` for `true`, `0` for `false`, `2` for unknown, this setting is ignored."]
+ pub read_write: u8,
}
#[allow(clippy::unnecessary_operation, clippy::identity_op)]
const _: () = {
- ["Size of Input"][::std::mem::size_of::<Input>() - 72usize];
- ["Alignment of Input"][::std::mem::align_of::<Input>() - 8usize];
- ["Offset of field: Input::config"][::std::mem::offset_of!(Input, config) - 0usize];
- ["Offset of field: Input::input_type"][::std::mem::offset_of!(Input, input_type) - 32usize];
- ["Offset of field: Input::input"][::std::mem::offset_of!(Input, input) - 40usize];
+ ["Size of PdRoute"][::std::mem::size_of::<PdRoute>() - 16usize];
+ ["Alignment of PdRoute"][::std::mem::align_of::<PdRoute>() - 8usize];
+ ["Offset of field: PdRoute::shard"][::std::mem::offset_of!(PdRoute, shard) - 0usize];
+ ["Offset of field: PdRoute::read_write"][::std::mem::offset_of!(PdRoute, read_write) - 8usize];
};
+pub type __builtin_va_list = *mut ::std::os::raw::c_char;
diff --git a/pgdog-plugin/src/c_api.rs b/pgdog-plugin/src/c_api.rs
deleted file mode 100644
index eac22b5c..00000000
--- a/pgdog-plugin/src/c_api.rs
+++ /dev/null
@@ -1,24 +0,0 @@
-use crate::bindings::*;
-use std::alloc::{alloc, dealloc, Layout};
-use std::ffi::c_int;
-
-/// Create new row.
-#[no_mangle]
-pub extern "C" fn pgdog_row_new(num_columns: c_int) -> Row {
- let layout = Layout::array::<RowColumn>(num_columns as usize).unwrap();
- let columns = unsafe { alloc(layout) };
-
- Row {
- num_columns,
- columns: columns as *mut RowColumn,
- }
-}
-
-/// Delete a row.
-#[no_mangle]
-pub extern "C" fn pgdog_row_free(row: Row) {
- let layout = Layout::array::<RowColumn>(row.num_columns as usize).unwrap();
- unsafe {
- dealloc(row.columns as *mut u8, layout);
- }
-}
diff --git a/pgdog-plugin/src/comp.rs b/pgdog-plugin/src/comp.rs
new file mode 100644
index 00000000..e49fe038
--- /dev/null
+++ b/pgdog-plugin/src/comp.rs
@@ -0,0 +1,8 @@
+//! Compatibility checks.
+
+use crate::PdStr;
+
+/// Rust compiler version used to build this library.
+pub fn rustc_version() -> PdStr {
+ env!("RUSTC_VERSION").into()
+}
diff --git a/pgdog-plugin/src/config.rs b/pgdog-plugin/src/config.rs
deleted file mode 100644
index 815928aa..00000000
--- a/pgdog-plugin/src/config.rs
+++ /dev/null
@@ -1,107 +0,0 @@
-//! Config helpers.
-
-use crate::bindings::*;
-use std::alloc::{alloc, dealloc, Layout};
-use std::ffi::{CStr, CString};
-use std::ptr::copy;
-
-impl DatabaseConfig {
- /// Create new database config.
- pub fn new(host: CString, port: u16, role: Role, shard: usize) -> Self {
- Self {
- shard: shard as i32,
- role,
- port: port as i32,
- host: host.into_raw(),
- }
- }
-
- /// Get host name.
- pub fn host(&self) -> &str {
- unsafe { CStr::from_ptr(self.host) }.to_str().unwrap()
- }
-
- /// Database port.
- pub fn port(&self) -> u16 {
- self.port as u16
- }
-
- /// Shard.
- pub fn shard(&self) -> usize {
- self.shard as usize
- }
-
- /// Is this a replica?
- pub fn replica(&self) -> bool {
- self.role == Role_REPLICA
- }
-
- /// Is this a primary?
- pub fn primary(&self) -> bool {
- !self.replica()
- }
-
- /// Deallocate this structure after use.
- ///
- /// # Safety
- ///
- /// This is not to be used by plugins.
- /// This is for internal pgDog usage only.
- pub(crate) unsafe fn deallocate(&self) {
- drop(unsafe { CString::from_raw(self.host) })
- }
-}
-
-impl Config {
- /// Create new config structure.
- pub fn new(name: CString, databases: &[DatabaseConfig], shards: usize) -> Self {
- let layout = Layout::array::<DatabaseConfig>(databases.len()).unwrap();
- let ptr = unsafe {
- let ptr = alloc(layout) as *mut DatabaseConfig;
- copy(databases.as_ptr(), ptr, databases.len());
- ptr
- };
-
- Self {
- num_databases: databases.len() as i32,
- databases: ptr,
- name: name.into_raw(),
- shards: shards as i32,
- }
- }
-
- /// Get database at index.
- pub fn database(&self, index: usize) -> Option<DatabaseConfig> {
- if index < self.num_databases as usize {
- Some(unsafe { *self.databases.add(index) })
- } else {
- None
- }
- }
-
- /// Get all databases in this configuration.
- pub fn databases(&self) -> Vec<DatabaseConfig> {
- (0..self.num_databases)
- .map(|i| self.database(i as usize).unwrap())
- .collect()
- }
-
- /// Number of shards.
- pub fn shards(&self) -> usize {
- self.shards as usize
- }
-
- /// Deallocate this structure.
- ///
- /// SAFETY: This is not to be used by plugins.
- /// # Safety
- ///
- /// This is for internal pgDog usage only.
- pub(crate) unsafe fn deallocate(&self) {
- self.databases().into_iter().for_each(|d| d.deallocate());
-
- let layout = Layout::array::<DatabaseConfig>(self.num_databases as usize).unwrap();
- unsafe { dealloc(self.databases as *mut u8, layout) };
- drop(unsafe { CString::from_raw(self.name) })
- }
-}
diff --git a/pgdog-plugin/src/context.rs b/pgdog-plugin/src/context.rs
new file mode 100644
index 00000000..fc76a8d6
--- /dev/null
+++ b/pgdog-plugin/src/context.rs
@@ -0,0 +1,442 @@
+//! Context passed to and from the plugins.
+
+use std::ops::Deref;
+
+use crate::{bindings::PdRouterContext, PdRoute, PdStatement};
+
+/// PostgreSQL statement, parsed by [`pg_query`].
+///
+/// Implements [`Deref`] on [`PdStatement`], which is passed
+/// in using the FFI interface.
+/// Use the [`PdStatement::protobuf`] method to obtain a reference
+/// to the Abstract Syntax Tree.
+///
+/// ### Example
+///
+/// ```no_run
+/// # use pgdog_plugin::Context;
+/// # let context = unsafe { Context::doc_test() };
+/// # let statement = context.statement();
+/// use pgdog_plugin::pg_query::NodeEnum;
+///
+/// let ast = statement.protobuf();
+/// let root = ast
+/// .stmts
+/// .first()
+/// .unwrap()
+/// .stmt
+/// .as_ref()
+/// .unwrap()
+/// .node
+/// .as_ref();
+///
+/// if let Some(NodeEnum::SelectStmt(stmt)) = root {
+/// println!("SELECT statement: {:#?}", stmt);
+/// }
+///
+/// ```
+pub struct Statement {
+ ffi: PdStatement,
+}
+
+impl Deref for Statement {
+ type Target = PdStatement;
+
+ fn deref(&self) -> &Self::Target {
+ &self.ffi
+ }
+}
+
+/// Context information provided by PgDog to the plugin at statement execution. It contains the actual statement and several metadata about
+/// the state of the database cluster:
+///
+/// - Number of shards
+/// - Does it have replicas
+/// - Does it have a primary
+///
+/// ### Example
+///
+/// ```
+/// use pgdog_plugin::{Context, Route, macros, Shard, ReadWrite};
+///
+/// #[macros::route]
+/// fn route(context: Context) -> Route {
+/// let shards = context.shards();
+/// let read_only = context.read_only();
+/// let ast = context.statement().protobuf();
+///
+/// println!("shards: {} (read_only: {})", shards, read_only);
+/// println!("ast: {:#?}", ast);
+///
+/// let read_write = if read_only {
+/// ReadWrite::Read
+/// } else {
+/// ReadWrite::Write
+/// };
+///
+/// Route::new(Shard::Direct(0), read_write)
+/// }
+/// ```
+///
+pub struct Context {
+ ffi: PdRouterContext,
+}
+
+impl From<PdRouterContext> for Context {
+ fn from(value: PdRouterContext) -> Self {
+ Self { ffi: value }
+ }
+}
+
+impl Context {
+ /// Returns a reference to the Abstract Syntax Tree (AST) created by [`pg_query`].
+ ///
+ /// # Example
+ ///
+ /// ```no_run
+ /// # use pgdog_plugin::Context;
+ /// # let context = unsafe { Context::doc_test() };
+ /// # let statement = context.statement();
+ /// let ast = context.statement().protobuf();
+ /// let nodes = ast.nodes();
+ /// ```
+ pub fn statement(&self) -> Statement {
+ Statement {
+ ffi: self.ffi.query,
+ }
+ }
+
+ /// Returns true if the database cluster doesn't have a primary database and can only serve
+ /// read queries.
+ ///
+ /// # Example
+ ///
+ /// ```
+ /// # use pgdog_plugin::Context;
+ /// # let context = unsafe { Context::doc_test() };
+ ///
+ /// let read_only = context.read_only();
+ ///
+ /// if read_only {
+ /// println!("Database cluster doesn't have a primary, only replicas.");
+ /// }
+ /// ```
+ pub fn read_only(&self) -> bool {
+ self.ffi.has_primary == 0
+ }
+
+ /// Returns true if the database cluster has replica databases.
+ ///
+ /// # Example
+ ///
+ /// ```
+ /// # use pgdog_plugin::Context;
+ /// # let context = unsafe { Context::doc_test() };
+ /// let has_replicas = context.has_replicas();
+ ///
+ /// if has_replicas {
+ /// println!("Database cluster can load balance read queries.")
+ /// }
+ /// ```
+ pub fn has_replicas(&self) -> bool {
+ self.ffi.has_replicas == 1
+ }
+
+ /// Returns true if the database cluster has a primary database and can serve write queries.
+ ///
+ /// # Example
+ ///
+ /// ```
+ /// # use pgdog_plugin::Context;
+ /// # let context = unsafe { Context::doc_test() };
+ /// let has_primary = context.has_primary();
+ ///
+ /// if has_primary {
+ /// println!("Database cluster can serve write queries.");
+ /// }
+ /// ```
+ pub fn has_primary(&self) -> bool {
+ !self.read_only()
+ }
+
+ /// Returns the number of shards in the database cluster.
+ ///
+ /// # Example
+ ///
+ /// ```
+ /// # use pgdog_plugin::Context;
+ /// # let context = unsafe { Context::doc_test() };
+ /// let shards = context.shards();
+ ///
+ /// if shards > 1 {
+ /// println!("Plugin should consider which shard to route the query to.");
+ /// }
+ /// ```
+ pub fn shards(&self) -> usize {
+ self.ffi.shards as usize
+ }
+
+ /// Returns true if the database cluster has more than one shard.
+ ///
+ /// # Example
+ ///
+ /// ```
+ /// # use pgdog_plugin::Context;
+ /// # let context = unsafe { Context::doc_test() };
+ /// let sharded = context.sharded();
+ /// let shards = context.shards();
+ ///
+ /// if sharded {
+ /// assert!(shards > 1);
+ /// } else {
+ /// assert_eq!(shards, 1);
+ /// }
+ /// ```
+ pub fn sharded(&self) -> bool {
+ self.shards() > 1
+ }
+
+ /// Returns true if PgDog strongly believes the statement should be sent to a primary. This indicates
+ /// that the statement is **not** a `SELECT` (e.g. `UPDATE`, `DELETE`, etc.), or a `SELECT` that is very likely to write data to the database, e.g.:
+ ///
+ /// ```sql
+ /// WITH users AS (
+ /// INSERT INTO users VALUES (1, 'test@acme.com') RETURNING *
+ /// )
+ /// SELECT * FROM users;
+ /// ```
+ ///
+ /// # Example
+ ///
+ /// ```
+ /// # use pgdog_plugin::Context;
+ /// # let context = unsafe { Context::doc_test() };
+ /// if context.write_override() {
+ /// println!("We should really send this query to the primary.");
+ /// }
+ /// ```
+ pub fn write_override(&self) -> bool {
+ self.ffi.write_override == 1
+ }
+}
+
+impl Context {
+ /// Used for doc tests only. **Do not use**.
+ ///
+ /// # Safety
+ ///
+ /// Not safe, don't use. We use it for doc tests only.
+ ///
+ pub unsafe fn doc_test() -> Context {
+ use std::{os::raw::c_void, ptr::null};
+
+ Context {
+ ffi: PdRouterContext {
+ shards: 1,
+ has_replicas: 1,
+ has_primary: 1,
+ in_transaction: 0,
+ write_override: 0,
+ query: PdStatement {
+ version: 1,
+ len: 0,
+ data: null::<c_void>() as *mut c_void,
+ },
+ },
+ }
+ }
+}
+
+/// What shard, if any, the statement should be sent to.
+///
+/// ### Example
+///
+/// ```
+/// use pgdog_plugin::Shard;
+///
+/// // Send query to shard 2.
+/// let direct = Shard::Direct(2);
+///
+/// // Send query to all shards.
+/// let cross_shard = Shard::All;
+///
+/// // Let PgDog handle sharding.
+/// let unknown = Shard::Unknown;
+/// ```
+#[derive(Debug, Copy, Clone, PartialEq, Eq)]
+pub enum Shard {
+ /// Direct-to-shard statement, sent to the specified shard only.
+ Direct(usize),
+ /// Send statement to all shards and let PgDog collect and transform the results.
+ All,
+ /// Not clear which shard it should go to, so let PgDog decide.
+ /// Use this if you don't want to handle sharding inside the plugin.
+ Unknown,
+}
+
+impl From<Shard> for i64 {
+ fn from(value: Shard) -> Self {
+ match value {
+ Shard::Direct(value) => value as i64,
+ Shard::All => -1,
+ Shard::Unknown => -2,
+ }
+ }
+}
+
+impl TryFrom<i64> for Shard {
+ type Error = ();
+ fn try_from(value: i64) -> Result<Self, Self::Error> {
+ Ok(if value == -1 {
+ Shard::All
+ } else if value == -2 {
+ Shard::Unknown
+ } else if value >= 0 {
+ Shard::Direct(value as usize)
+ } else {
+ return Err(());
+ })
+ }
+}
+
+impl TryFrom<u8> for ReadWrite {
+ type Error = ();
+
+ fn try_from(value: u8) -> Result<Self, Self::Error> {
+ Ok(if value == 0 {
+ ReadWrite::Write
+ } else if value == 1 {
+ ReadWrite::Read
+ } else if value == 2 {
+ ReadWrite::Unknown
+ } else {
+ return Err(());
+ })
+ }
+}
+
+/// Indicates if the statement is a read or a write. Read statements are sent to a replica, if one is configured.
+/// Write statements are sent to the primary.
+///
+/// ### Example
+///
+/// ```
+/// use pgdog_plugin::ReadWrite;
+///
+/// // The statement should go to a replica.
+/// let read = ReadWrite::Read;
+///
+/// // The statement should go the primary.
+/// let write = ReadWrite::Write;
+///
+/// // Skip and let PgDog decide.
+/// let unknown = ReadWrite::Unknown;
+/// ```
+#[derive(Debug, Copy, Clone, PartialEq, Eq)]
+pub enum ReadWrite {
+ /// Send the statement to a replica, if any are configured.
+ Read,
+ /// Send the statement to the primary.
+ Write,
+ /// Plugin doesn't know if the statement is a read or write. This let's PgDog decide.
+ /// Use this if you don't want to make this decision in the plugin.
+ Unknown,
+}
+
+impl From<ReadWrite> for u8 {
+ fn from(value: ReadWrite) -> Self {
+ match value {
+ ReadWrite::Write => 0,
+ ReadWrite::Read => 1,
+ ReadWrite::Unknown => 2,
+ }
+ }
+}
+
+impl Default for PdRoute {
+ fn default() -> Self {
+ Route::unknown().ffi
+ }
+}
+
+/// Statement route.
+///
+/// PgDog uses this to decide where a query should be sent to. Read statements are sent to a replica,
+/// while write ones are sent the primary. If the cluster has more than one shard, the statement can be
+/// sent to a specific database, or all of them.
+///
+/// ### Example
+///
+/// ```
+/// use pgdog_plugin::{Shard, ReadWrite, Route};
+///
+/// // This sends the query to the primary database of shard 0.
+/// let route = Route::new(Shard::Direct(0), ReadWrite::Write);
+///
+/// // This sends the query to all shards, routing them to a replica
+/// // of each shard, if any are configured.
+/// let route = Route::new(Shard::All, ReadWrite::Read);
+///
+/// // No routing information is available. PgDog will ignore it
+/// // and make its own decision.
+/// let route = Route::unknown();
+pub struct Route {
+ ffi: PdRoute,
+}
+
+impl Default for Route {
+ fn default() -> Self {
+ Self::unknown()
+ }
+}
+
+impl Deref for Route {
+ type Target = PdRoute;
+
+ fn deref(&self) -> &Self::Target {
+ &self.ffi
+ }
+}
+
+impl From<PdRoute> for Route {
+ fn from(value: PdRoute) -> Self {
+ Self { ffi: value }
+ }
+}
+
+impl From<Route> for PdRoute {
+ fn from(value: Route) -> Self {
+ value.ffi
+ }
+}
+
+impl Route {
+ /// Create new route.
+ ///
+ /// # Arguments
+ ///
+ /// * `shard`: Which shard the statement should be sent to.
+ /// * `read_write`: Does the statement read or write data. Read statements are sent to a replica. Write statements are sent to the primary.
+ ///
+ pub fn new(shard: Shard, read_write: ReadWrite) -> Route {
+ Self {
+ ffi: PdRoute {
+ shard: shard.into(),
+ read_write: read_write.into(),
+ },
+ }
+ }
+
+ /// Create new route with no sharding or read/write information.
+ /// Use this if you don't want your plugin to do query routing.
+ /// Plugins that do something else with queries, e.g., logging, metrics,
+ /// can return this route.
+ pub fn unknown() -> Route {
+ Self {
+ ffi: PdRoute {
+ shard: -2,
+ read_write: 2,
+ },
+ }
+ }
+}
diff --git a/pgdog-plugin/src/copy.rs b/pgdog-plugin/src/copy.rs
deleted file mode 100644
index 28f4e763..00000000
--- a/pgdog-plugin/src/copy.rs
+++ /dev/null
@@ -1,241 +0,0 @@
-//! Handle COPY commands.
-
-use libc::c_char;
-
-use crate::{
- bindings::{Copy, CopyInput, CopyOutput, CopyRow},
- CopyFormat_CSV, CopyFormat_INVALID,
-};
-use std::{
- alloc::{alloc, dealloc, Layout},
- ffi::{CStr, CString},
- ptr::{copy, null_mut},
- slice::from_raw_parts,
- str::from_utf8_unchecked,
-};
-
-impl Copy {
- /// Not a valid COPY statement. Will be ignored by the router.
- pub fn invalid() -> Self {
- Self {
- copy_format: CopyFormat_INVALID,
- table_name: null_mut(),
- has_headers: 0,
- delimiter: ',' as c_char,
- num_columns: 0,
- columns: null_mut(),
- }
- }
-
- /// Create new copy command.
- pub fn new(table_name: &str, headers: bool, delimiter: char, columns: &[&str]) -> Self {
- let mut cols = vec![];
- for column in columns {
- let cstr = CString::new(column.as_bytes()).unwrap();
- cols.push(cstr.into_raw());
- }
- let layout = Layout::array::<*mut i8>(columns.len()).unwrap();
- #[cfg(all(target_os = "linux", target_arch = "aarch64"))]
- let ptr = unsafe { alloc(layout) as *mut *mut u8 };
- #[cfg(not(all(target_os = "linux", target_arch = "aarch64")))]
- let ptr = unsafe { alloc(layout) as *mut *mut i8 };
- unsafe {
- copy(cols.as_ptr(), ptr, columns.len());
- }
-
- Self {
- table_name: CString::new(table_name).unwrap().into_raw(),
- has_headers: if headers { 1 } else { 0 },
- copy_format: CopyFormat_CSV,
- delimiter: delimiter as c_char,
- num_columns: columns.len() as i32,
- columns: ptr,
- }
- }
-
- /// Get table name.
- pub fn table_name(&self) -> &str {
- unsafe { CStr::from_ptr(self.table_name).to_str().unwrap() }
- }
-
- /// Does this COPY statement say to expect headers?
- pub fn has_headers(&self) -> bool {
- self.has_headers != 0
- }
-
- /// Columns specified by the caller.
- pub fn columns(&self) -> Vec<&str> {
- unsafe {
- (0..self.num_columns)
- .map(|s| {
- CStr::from_ptr(*self.columns.offset(s as isize))
- .to_str()
- .unwrap()
- })
- .collect()
- }
- }
-
- /// Get CSV delimiter.
- pub fn delimiter(&self) -> char {
- self.delimiter as u8 as char
- }
-
- /// Deallocate this structure.
- ///
- /// # Safety
- ///
- /// Call this only when finished with this.
- ///
- pub unsafe fn deallocate(&self) {
- unsafe { drop(CString::from_raw(self.table_name)) }
-
- (0..self.num_columns)
- .for_each(|i| drop(CString::from_raw(*self.columns.offset(i as isize))));
-
- let layout = Layout::array::<*mut i8>(self.num_columns as usize).unwrap();
- unsafe { dealloc(self.columns as *mut u8, layout) }
- }
-}
-
-impl CopyInput {
- /// Create new copy input.
- pub fn new(data: &[u8], sharding_column: usize, headers: bool, delimiter: char) -> Self {
- #[cfg(all(target_os = "linux", target_arch = "aarch64"))]
- let data_ptr = data.as_ptr() as *const u8;
- #[cfg(not(all(target_os = "linux", target_arch = "aarch64")))]
- let data_ptr = data.as_ptr() as *const i8;
- Self {
- len: data.len() as i32,
- data: data_ptr,
- sharding_column: sharding_column as i32,
- has_headers: if headers { 1 } else { 0 },
- delimiter: delimiter as c_char,
- }
- }
-
- /// Get data as slice.
- pub fn data(&self) -> &[u8] {
- unsafe { from_raw_parts(self.data as *const u8, self.len as usize) }
- }
-
- /// CSV delimiter.
- pub fn delimiter(&self) -> char {
- self.delimiter as u8 as char
- }
-
- /// Sharding column offset.
- pub fn sharding_column(&self) -> usize {
- self.sharding_column as usize
- }
-
- /// Does this input contain headers? Only the first one will.
- pub fn headers(&self) -> bool {
- self.has_headers != 0
- }
-}
-
-impl CopyRow {
- /// Create new row from data slice.
- pub fn new(data: &[u8], shard: i32) -> Self {
- #[cfg(all(target_os = "linux", target_arch = "aarch64"))]
- let data_ptr = data.as_ptr() as *mut u8;
- #[cfg(not(all(target_os = "linux", target_arch = "aarch64")))]
- let data_ptr = data.as_ptr() as *mut i8;
- Self {
- len: data.len() as i32,
- data: data_ptr,
- shard,
- }
- }
-
- /// Shard this row should go to.
- pub fn shard(&self) -> usize {
- self.shard as usize
- }
-
- /// Get data.
- pub fn data(&self) -> &[u8] {
- unsafe { from_raw_parts(self.data as *const u8, self.len as usize) }
- }
-}
-
-impl std::fmt::Debug for CopyRow {
- fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
- f.debug_struct("CopyRow")
- .field("len", &self.len)
- .field("shard", &self.shard)
- .field("data", &unsafe { from_utf8_unchecked(self.data()) })
- .finish()
- }
-}
-
-impl CopyOutput {
- /// Copy output from rows.
- pub fn new(rows: &[CopyRow]) -> Self {
- let layout = Layout::array::<CopyRow>(rows.len()).unwrap();
- unsafe {
- let ptr = alloc(layout) as *mut CopyRow;
- copy(rows.as_ptr(), ptr, rows.len());
- Self {
- num_rows: rows.len() as i32,
- rows: ptr,
- header: null_mut(),
- }
- }
- }
-
- /// Parse and give back the CSV header.
- pub fn with_header(mut self, header: Option<String>) -> Self {
- if let Some(header) = header {
- let ptr = CString::new(header).unwrap().into_raw();
- self.header = ptr;
- }
-
- self
- }
-
- /// Get rows.
- pub fn rows(&self) -> &[CopyRow] {
- unsafe { from_raw_parts(self.rows, self.num_rows as usize) }
- }
-
- /// Get header value, if any.
- pub fn header(&self) -> Option<&str> {
- unsafe {
- if !self.header.is_null() {
- CStr::from_ptr(self.header).to_str().ok()
- } else {
- None
- }
- }
- }
-
- /// Deallocate this structure.
- ///
- /// # Safety
- ///
- /// Don't use unless you don't need this data anymore.
- ///
- pub unsafe fn deallocate(&self) {
- let layout = Layout::array::<CopyRow>(self.num_rows as usize).unwrap();
- dealloc(self.rows as *mut u8, layout);
-
- if !self.header.is_null() {
- unsafe { drop(CString::from_raw(self.header)) }
- }
- }
-}
-
-impl std::fmt::Debug for CopyOutput {
- fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
- let rows = (0..self.num_rows)
- .map(|i| unsafe { *self.rows.offset(i as isize) })
- .collect::<Vec<_>>();
-
- f.debug_struct("CopyOutput")
- .field("num_rows", &self.num_rows)
- .field("rows", &rows)
- .finish()
- }
-}
diff --git a/pgdog-plugin/src/input.rs b/pgdog-plugin/src/input.rs
deleted file mode 100644
index 780aa66a..00000000
--- a/pgdog-plugin/src/input.rs
+++ /dev/null
@@ -1,63 +0,0 @@
-//! Plugin input helpers.
-#![allow(non_upper_case_globals)]
-use crate::bindings::{self, *};
-
-impl bindings::Input {
- /// Create new plugin input.
- pub fn new_query(config: Config, input: RoutingInput) -> Self {
- Self {
- config,
- input,
- input_type: InputType_ROUTING_INPUT,
- }
- }
-
- pub fn new_copy(config: Config, input: RoutingInput) -> Self {
- Self {
- config,
- input,
- input_type: InputType_COPY_INPUT,
- }
- }
-
- /// Deallocate memory.
- ///
- /// # Safety
- ///
- /// This is not to be used by plugins.
- /// # Safety
- ///
- /// This is for internal pgDog usage only.
- pub unsafe fn deallocate(&self) {
- self.config.deallocate();
- }
-
- /// Get query if this is a routing input.
- pub fn query(&self) -> Option<bindings::Query> {
- match self.input_type {
- InputType_ROUTING_INPUT => Some(unsafe { self.input.query }),
- _ => None,
- }
- }
-
- /// Get copy input, if any.
- pub fn copy(&self) -> Option<CopyInput> {
- if self.input_type == InputType_COPY_INPUT {
- Some(unsafe { self.input.copy })
- } else {
- None
- }
- }
-}
-
-impl RoutingInput {
- /// Create query routing input.
- pub fn query(query: bindings::Query) -> Self {
- Self { query }
- }
-
- /// Create copy routing input.
- pub fn copy(copy: CopyInput) -> Self {
- Self { copy }
- }
-}
diff --git a/pgdog-plugin/src/lib.rs b/pgdog-plugin/src/lib.rs
index 76af01e6..98e790c1 100644
--- a/pgdog-plugin/src/lib.rs
+++ b/pgdog-plugin/src/lib.rs
@@ -1,34 +1,136 @@
-//! pgDog plugin interface.
+//! PgDog plugins library.
+//!
+//! Implements data types and methods plugins can use to interact with PgDog at runtime.
+//!
+//! # Getting started
+//!
+//! Create a Rust library package with Cargo:
+//!
+//! ```bash
+//! cargo init --lib my_plugin
+//! ```
+//!
+//! The plugin needs to be built as a C ABI-compatible shared library. Add the following to Cargo.toml in the new plugin directory:
+//!
+//! ```toml
+//! [lib]
+//! crate-type = ["rlib", "cdylib"]
+//! ```
+//!
+//! ## Dependencies
+//!
+//! PgDog is using [`pg_query`] to parse SQL. It produces an Abstract Syntax Tree (AST) which plugins can use to inspect queries
+//! and make statement routing decisions.
+//!
+//! The AST is computed by PgDog at runtime. It then passes it down to plugins, using a FFI interface. To make this safe, plugins must follow the
+//! following 2 requirements:
+//!
+//! 1. Plugins must be compiled with the **same version of the Rust compiler** as PgDog. This is automatically checked at runtime and plugins that don't do this are not loaded.
+//! 2. Plugins must use the **same version of [`pg_query`] crate** as PgDog. This happens automatically when using `pg_query` structs re-exported by this crate.
+//!
+//!
+//! #### Configure dependencies
+//!
+//! Add the following to your plugin's `Cargo.toml`:
+//!
+//! ```toml
+//! [dependencies]
+//! pgdog-plugin = "0.1.6"
+//! ```
+//!
+//! # Required methods
+//!
+//! All plugins need to implement a set of functions that PgDog calls at runtime to load the plugin. You can implement them automatically
+//! using a macro. Inside the plugin's `src/lib.rs` file, add the following code:
+//!
+//! ```
+//! // src/lib.rs
+//! use pgdog_plugin::macros;
+//!
+//! macros::plugin!();
+//! ```
+//!
+//! # Routing queries
+//!
+//! Plugins are most commonly used to route queries. To do this, they need to implement a function that reads
+//! the [`Context`] passed in by PgDog, and returns a [`Route`] that indicates which database the query should be sent to.
+//!
+//! ### Example
+//!
+//! ```no_run
+//! use pgdog_plugin::prelude::*;
+//! use pg_query::{protobuf::{Node, RawStmt}, NodeEnum};
+//!
+//! #[route]
+//! fn route(context: Context) -> Route {
+//! let proto = context
+//! .statement()
+//! .protobuf();
+//! let root = proto.stmts.first();
+//! if let Some(root) = root {
+//! if let Some(ref stmt) = root.stmt {
+//! if let Some(ref node) = stmt.node {
+//! if let NodeEnum::SelectStmt(_) = node {
+//! return Route::new(Shard::Unknown, ReadWrite::Read);
+//! }
+//! }
+//! }
+//! }
+//!
+//! Route::new(Shard::Unknown, ReadWrite::Write)
+//! }
+//! ```
+//!
+//! The [`macros::route`] macro wraps the function into a safe FFI interface which PgDog calls at runtime.
+//!
+//! ### Errors
+//!
+//! Plugin functions cannot return errors or panic. To handle errors, you can log them to `stderr` and return a default route,
+//! which PgDog will ignore. Plugins currently cannot be used to block queries.
+//!
+//! # Enabling plugins
+//!
+//! Plugins are shared libraries, loaded by PgDog at runtime using `dlopen(3)`. If specifying only its name, make sure to place the plugin's shared library
+//! into one of the following locations:
+//!
+//! - Any of the system default paths, e.g.: `/lib`, `/usr/lib`, `/lib64`, `/usr/lib64`, etc.
+//! - Path specified by the `LD_LIBRARY_PATH` (on Linux) or `DYLD_LIBRARY_PATH` (Mac OS) environment variables.
+//!
+//! Alternatively, specify the relative or absolute path to the shared library as the plugin name. Plugins aren't loaded automatically. For each plugin you want to enable, add it to `pgdog.toml`:
+//!
+//! ```toml
+//! [[plugins]]
+//! # Plugin should be in /usr/lib or in LD_LIBRARY_PATH.
+//! name = "my_plugin"
+//!
+//! [[plugins]]
+//! # Plugin should be in $PWD/libmy_plugin.so
+//! name = "libmy_plugin.so"
+//!
+//! [[plugins]]
+//! # Absolute path to the plugin.
+//! name = "/usr/local/lib/libmy_plugin.so"
+//! ```
+//!
+/// Bindgen-generated FFI bindings.
#[allow(non_upper_case_globals)]
+#[allow(non_camel_case_types)]
+#[allow(non_snake_case)]
pub mod bindings;
-pub mod c_api;
-pub mod config;
-pub mod copy;
-pub mod input;
-pub mod order_by;
-pub mod output;
-pub mod parameter;
+pub mod ast;
+pub mod comp;
+pub mod context;
pub mod plugin;
-pub mod query;
-pub mod route;
+pub mod prelude;
+pub mod string;
pub use bindings::*;
-pub use c_api::*;
+pub use context::*;
pub use plugin::*;
pub use libloading;
-#[cfg(test)]
-mod test {
- use super::*;
- use std::ffi::CString;
-
- #[test]
- fn test_query() {
- let query = CString::new("SELECT 1").unwrap();
- let query = Query::new(query);
- assert_eq!(query.query(), "SELECT 1");
- }
-}
+pub use pg_query;
+pub use pgdog_macros as macros;
diff --git a/pgdog-plugin/src/order_by.rs b/pgdog-plugin/src/order_by.rs
index 4ee52988..8b137891 100644
--- a/pgdog-plugin/src/order_by.rs
+++ b/pgdog-plugin/src/order_by.rs
@@ -1,43 +1 @@
-use std::{
- ffi::{CStr, CString},
- ptr::null_mut,
-};
-use crate::{OrderBy, OrderByDirection};
-
-impl OrderBy {
- pub(crate) fn drop(&self) {
- if !self.column_name.is_null() {
- unsafe { drop(CString::from_raw(self.column_name)) }
- }
- }
-
- /// Order by column name.
- pub fn column_name(name: &str, direction: OrderByDirection) -> Self {
- let column_name = CString::new(name.as_bytes()).unwrap();
-
- Self {
- column_name: column_name.into_raw(),
- column_index: -1,
- direction,
- }
- }
-
- /// Order by column index.
- pub fn column_index(index: usize, direction: OrderByDirection) -> Self {
- Self {
- column_name: null_mut(),
- column_index: index as i32,
- direction,
- }
- }
-
- /// Get column name if any.
- pub fn name(&self) -> Option<&str> {
- if self.column_name.is_null() || self.column_index >= 0 {
- None
- } else {
- unsafe { CStr::from_ptr(self.column_name).to_str().ok() }
- }
- }
-}
diff --git a/pgdog-plugin/src/output.rs b/pgdog-plugin/src/output.rs
deleted file mode 100644
index 56b2972b..00000000
--- a/pgdog-plugin/src/output.rs
+++ /dev/null
@@ -1,91 +0,0 @@
-//! Plugin output helpers.
-#![allow(non_upper_case_globals)]
-use crate::bindings::*;
-
-impl std::fmt::Debug for Output {
- fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
- f.debug_struct("Output")
- .field("decision", &self.decision)
- .finish()
- }
-}
-
-impl Output {
- /// Plugin doesn't want to deal with the input.
- /// Router will skip it.
- pub fn skip() -> Self {
- Self {
- decision: RoutingDecision_NO_DECISION,
- output: RoutingOutput::new_route(Route::unknown()),
- }
- }
-
- /// Create new forward output.
- ///
- /// This means the query will be forwarded as-is to a destination
- /// specified in the route.
- pub fn new_forward(route: Route) -> Output {
- Output {
- decision: RoutingDecision_FORWARD,
- output: RoutingOutput::new_route(route),
- }
- }
-
- /// Create new copy statement.
- pub fn new_copy(copy: Copy) -> Output {
- Output {
- decision: RoutingDecision_COPY,
- output: RoutingOutput::new_copy(copy),
- }
- }
-
- /// Sharded copy rows.
- pub fn new_copy_rows(output: CopyOutput) -> Output {
- Output {
- decision: RoutingDecision_COPY_ROWS,
- output: RoutingOutput::new_copy_rows(output),
- }
- }
-
- /// Get route determined by the plugin.
- pub fn route(&self) -> Option<Route> {
- match self.decision {
- RoutingDecision_FORWARD => Some(unsafe { self.output.route }),
- _ => None,
- }
- }
-
- /// Get copy info determined by the plugin.
- pub fn copy(&self) -> Option<Copy> {
- if self.decision == RoutingDecision_COPY {
- Some(unsafe { self.output.copy })
- } else {
- None
- }
- }
-
- /// Get copy rows if any.
- pub fn copy_rows(&self) -> Option<CopyOutput> {
- if self.decision == RoutingDecision_COPY_ROWS {
- Some(unsafe { self.output.copy_rows })
- } else {
- None
- }
- }
-
- /// # Safety
- ///
- /// Don't use this function unless you're cleaning up plugin
- /// output.
- pub unsafe fn deallocate(&self) {
- if self.decision == RoutingDecision_FORWARD {
- self.output.route.deallocate();
- }
- if self.decision == RoutingDecision_COPY {
- self.output.copy.deallocate();
- }
- if self.decision == RoutingDecision_COPY_ROWS {
- self.output.copy_rows.deallocate();
- }
- }
-}
diff --git a/pgdog-plugin/src/parameter.rs b/pgdog-plugin/src/parameter.rs
deleted file mode 100644
index e952562a..00000000
--- a/pgdog-plugin/src/parameter.rs
+++ /dev/null
@@ -1,53 +0,0 @@
-use crate::bindings::Parameter;
-
-use libc::c_char;
-use std::alloc::{alloc, dealloc, Layout};
-use std::ptr::copy;
-use std::slice::from_raw_parts;
-use std::str::from_utf8;
-
-impl Parameter {
- /// Create new parameter from format code and raw data.
- pub fn new(format: i16, data: &[u8]) -> Self {
- let len = data.len() as i32;
- let layout = Layout::array::<u8>(len as usize).unwrap();
- let ptr = unsafe { alloc(layout) };
- unsafe {
- copy::<u8>(data.as_ptr(), ptr, len as usize);
- }
-
- Self {
- len,
- data: ptr as *const c_char,
- format: format as i32,
- }
- }
-
- /// Manually free memory allocated for this parameter.
- ///
- /// # Safety
- ///
- /// Call this after plugin finished executing to avoid memory leaks.
- pub unsafe fn deallocate(&mut self) {
- let layout = Layout::array::<u8>(self.len as usize).unwrap();
- unsafe {
- dealloc(self.data as *mut u8, layout);
- }
- }
-
- /// Get parameter value as a string if it's encoded as one.
- pub fn as_str(&self) -> Option<&str> {
- if self.format != 0 {
- return None;
- }
-
- from_utf8(self.as_bytes()).ok()
- }
-
- /// Get parameter value as bytes.
- pub fn as_bytes(&self) -> &[u8] {
- let slice = unsafe { from_raw_parts(self.data as *const u8, self.len as usize) };
-
- slice
- }
-}
diff --git a/pgdog-plugin/src/plugin.rs b/pgdog-plugin/src/plugin.rs
index ffe90076..44f09eec 100644
--- a/pgdog-plugin/src/plugin.rs
+++ b/pgdog-plugin/src/plugin.rs
@@ -1,48 +1,91 @@
-//! Plugin interface.
-use std::ops::Deref;
+//! PgDog's plugin interface.
+//!
+//! This loads the shared library using [`libloading`] and exposes
+//! a safe interface to the plugin's methods.
+//!
+
+use std::path::Path;
-use crate::bindings::{self, Input, Output};
use libloading::{library_filename, Library, Symbol};
+use crate::{PdRoute, PdRouterContext, PdStr};
+
/// Plugin interface.
+///
+/// Methods are loaded using `libloading`. If required methods aren't found,
+/// the plugin isn't loaded. All optional methods are checked first, before being
+/// executed.
+///
+/// Using this interface is reasonably safe.
+///
#[derive(Debug)]
pub struct Plugin<'a> {
+ /// Plugin name.
name: String,
/// Initialization routine.
init: Option<Symbol<'a, unsafe extern "C" fn()>>,
/// Shutdown routine.
fini: Option<Symbol<'a, unsafe extern "C" fn()>>,
- /// Route query to a shard.
- route: Option<Symbol<'a, unsafe extern "C" fn(bindings::Input) -> Output>>,
+ /// Route query.
+ route: Option<Symbol<'a, unsafe extern "C" fn(PdRouterContext, *mut PdRoute)>>,
+ /// Compiler version.
+ rustc_version: Option<Symbol<'a, unsafe extern "C" fn(*mut PdStr)>>,
+ /// Plugin version.
+ plugin_version: Option<Symbol<'a, unsafe extern "C" fn(*mut PdStr)>>,
}
impl<'a> Plugin<'a> {
- /// Load library using a cross-platform naming convention.
- pub fn library(name: &str) -> Result<Library, libloading::Error> {
- let name = library_filename(name);
- unsafe { Library::new(name) }
+ /// Load plugin's shared library using a cross-platform naming convention.
+ ///
+ /// Plugin has to be in `LD_LIBRARY_PATH`, in a standard location
+ /// for the operating system, or be provided as an absolute or relative path,
+ /// including the platform-specific extension.
+ ///
+ /// ### Example
+ ///
+ /// ```no_run
+ /// use pgdog_plugin::Plugin;
+ ///
+ /// let plugin_lib = Plugin::library("/home/pgdog/plugin.so").unwrap();
+ /// let plugin_lib = Plugin::library("plugin.so").unwrap();
+ /// ```
+ ///
+ pub fn library<P: AsRef<Path>>(name: P) -> Result<Library, libloading::Error> {
+ if name.as_ref().exists() {
+ let name = name.as_ref().display().to_string();
+ unsafe { Library::new(&name) }
+ } else {
+ let name = library_filename(name.as_ref());
+ unsafe { Library::new(name) }
+ }
}
- /// Load standard methods from the plugin library.
+ /// Load standard plugin methods from the plugin library.
+ ///
+ /// ### Arguments
+ ///
+ /// * `name`: Plugin name. Can be any name you want, it's only used for logging.
+ /// * `library`: `libloading::Library` reference. Must have the same, ideally static, lifetime as the plugin.
+ ///
pub fn load(name: &str, library: &'a Library) -> Self {
- let route = unsafe { library.get(b"pgdog_route_query\0") }.ok();
let init = unsafe { library.get(b"pgdog_init\0") }.ok();
let fini = unsafe { library.get(b"pgdog_fini\0") }.ok();
+ let route = unsafe { library.get(b"pgdog_route\0") }.ok();
+ let rustc_version = unsafe { library.get(b"pgdog_rustc_version\0") }.ok();
+ let plugin_version = unsafe { library.get(b"pgdog_plugin_version\0") }.ok();
Self {
name: name.to_owned(),
- route,
init,
fini,
+ route,
+ rustc_version,
+ plugin_version,
}
}
- /// Route query.
- pub fn route(&self, input: Input) -> Option<Output> {
- self.route.as_ref().map(|route| unsafe { route(input) })
- }
-
- /// Perform initialization.
+ /// Execute plugin's initialization routine.
+ /// Returns true if the route exists and was executed, false otherwise.
pub fn init(&self) -> bool {
if let Some(init) = &self.init {
unsafe {
@@ -54,71 +97,56 @@ impl<'a> Plugin<'a> {
}
}
+ /// Execute plugin's shutdown routine.
pub fn fini(&self) {
if let Some(ref fini) = &self.fini {
unsafe { fini() }
}
}
- /// Plugin name.
- pub fn name(&self) -> &str {
- &self.name
- }
-
- /// Check that we have the required methods.
- pub fn valid(&self) -> bool {
- self.route.is_some()
- }
-}
-
-pub struct PluginOutput {
- output: Output,
-}
-
-impl PluginOutput {
- pub fn new(output: Output) -> Self {
- Self { output }
- }
-}
-
-impl Deref for PluginOutput {
- type Target = Output;
-
- fn deref(&self) -> &Self::Target {
- &self.output
- }
-}
-
-impl Drop for PluginOutput {
- fn drop(&mut self) {
- unsafe {
- self.output.deallocate();
+ /// Execute plugin's route routine. Determines where a statement should be sent.
+ /// Returns a route if the routine is defined, or `None` if not.
+ ///
+ /// ### Arguments
+ ///
+ /// * `context`: Statement context created by PgDog's query router.
+ ///
+ pub fn route(&self, context: PdRouterContext) -> Option<PdRoute> {
+ if let Some(ref route) = &self.route {
+ let mut output = PdRoute::default();
+ unsafe {
+ route(context, &mut output as *mut PdRoute);
+ }
+ Some(output)
+ } else {
+ None
}
}
-}
-pub struct PluginInput {
- input: Input,
-}
-
-impl PluginInput {
- pub fn new(input: Input) -> Self {
- Self { input }
+ /// Returns plugin's name. This is the same name as what
+ /// is passed to [`Plugin::load`] function.
+ pub fn name(&self) -> &str {
+ &self.name
}
-}
-
-impl Deref for PluginInput {
- type Target = Input;
- fn deref(&self) -> &Self::Target {
- &self.input
+ /// Returns the Rust compiler version used to build the plugin.
+ /// This version must match the compiler version used to build
+ /// PgDog, or the plugin won't be loaded.
+ pub fn rustc_version(&self) -> Option<PdStr> {
+ let mut output = PdStr::default();
+ self.rustc_version.as_ref().map(|rustc_version| unsafe {
+ rustc_version(&mut output);
+ output
+ })
}
-}
-impl Drop for PluginInput {
- fn drop(&mut self) {
- unsafe {
- self.input.deallocate();
- }
+ /// Get plugin version. It's set in plugin's
+ /// `Cargo.toml`.
+ pub fn version(&self) -> Option<PdStr> {
+ let mut output = PdStr::default();
+ self.plugin_version.as_ref().map(|func| unsafe {
+ func(&mut output as *mut PdStr);
+ output
+ })
}
}
diff --git a/pgdog-plugin/src/prelude.rs b/pgdog-plugin/src/prelude.rs
new file mode 100644
index 00000000..e364a576
--- /dev/null
+++ b/pgdog-plugin/src/prelude.rs
@@ -0,0 +1,7 @@
+//! Commonly used structs and re-exports.
+
+pub use crate::pg_query;
+pub use crate::{
+ macros::{fini, init, route},
+ Context, ReadWrite, Route, Shard,
+};
diff --git a/pgdog-plugin/src/query.rs b/pgdog-plugin/src/query.rs
deleted file mode 100644
index 0721be7d..00000000
--- a/pgdog-plugin/src/query.rs
+++ /dev/null
@@ -1,80 +0,0 @@
-use crate::bindings::{Parameter, Query};
-
-use std::alloc::{alloc, dealloc, Layout};
-use std::ffi::{CStr, CString};
-use std::ptr::{copy, null};
-
-impl Query {
- /// Get query text.
- pub fn query(&self) -> &str {
- debug_assert!(!self.query.is_null());
- unsafe { CStr::from_ptr(self.query) }.to_str().unwrap()
- }
-
- /// Create new query to pass it over the FFI boundary.
- pub fn new(query: CString) -> Self {
- Self {
- len: query.as_bytes().len() as i32,
- query: query.into_raw(),
- num_parameters: 0,
- parameters: null(),
- }
- }
-
- /// Set parameters on this query. This is used internally
- /// by pgDog to construct this structure.
- pub fn set_parameters(&mut self, params: &[Parameter]) {
- let layout = Layout::array::<Parameter>(params.len()).unwrap();
- let parameters = unsafe { alloc(layout) };
-
- unsafe {
- copy(params.as_ptr(), parameters as *mut Parameter, params.len());
- }
- self.parameters = parameters as *const Parameter;
- self.num_parameters = params.len() as i32;
- }
-
- /// Get query parameters, if any.
- pub fn parameters(&self) -> Vec<Parameter> {
- (0..self.num_parameters)
- .map(|i| self.parameter(i as usize).unwrap())
- .collect()
- }
-
- /// Get parameter at offset if one exists.
- pub fn parameter(&self, index: usize) -> Option<Parameter> {
- if index < self.num_parameters as usize {
- unsafe { Some(*self.parameters.add(index)) }
- } else {
- None
- }
- }
-
- /// Free memory allocated for parameters, if any.
- ///
- /// # Safety
- ///
- /// This is not to be used by plugins.
- /// This is for internal pgDog usage only.
- pub unsafe fn deallocate(&mut self) {
- #[cfg(all(target_os = "linux", target_arch = "aarch64"))]
- let ptr = self.query as *mut u8;
- #[cfg(not(all(target_os = "linux", target_arch = "aarch64")))]
- let ptr = self.query as *mut i8;
-
- unsafe { drop(CString::from_raw(ptr)) }
-
- if !self.parameters.is_null() {
- for index in 0..self.num_parameters {
- if let Some(mut param) = self.parameter(index as usize) {
- param.deallocate();
- }
- }
- let layout = Layout::array::<Parameter>(self.num_parameters as usize).unwrap();
- unsafe {
- dealloc(self.parameters as *mut u8, layout);
- self.parameters = null();
- }
- }
- }
-}
diff --git a/pgdog-plugin/src/route.rs b/pgdog-plugin/src/route.rs
deleted file mode 100644
index 6de91173..00000000
--- a/pgdog-plugin/src/route.rs
+++ /dev/null
@@ -1,166 +0,0 @@
-//! Query routing helpers.
-#![allow(non_upper_case_globals)]
-
-use std::{
- alloc::{alloc, dealloc, Layout},
- ptr::{copy, null_mut},
-};
-
-use crate::bindings::*;
-
-impl RoutingOutput {
- /// Create new route.
- pub fn new_route(route: Route) -> RoutingOutput {
- RoutingOutput { route }
- }
-
- /// Create new copy statement.
- pub fn new_copy(copy: Copy) -> RoutingOutput {
- RoutingOutput { copy }
- }
-
- /// Create new copy rows output.
- pub fn new_copy_rows(copy_rows: CopyOutput) -> RoutingOutput {
- RoutingOutput { copy_rows }
- }
-}
-
-impl Route {
- /// The plugin has no idea what to do with this query.
- /// The router will ignore this and try another way.
- pub fn unknown() -> Route {
- Route {
- shard: Shard_ANY,
- affinity: Affinity_UNKNOWN,
- num_order_by: 0,
- order_by: null_mut(),
- }
- }
-
- /// Read from this shard.
- pub fn read(shard: usize) -> Route {
- Route {
- shard: shard as i32,
- affinity: Affinity_READ,
- num_order_by: 0,
- order_by: null_mut(),
- }
- }
-
- /// Write to this shard.
- pub fn write(shard: usize) -> Route {
- Route {
- shard: shard as i32,
- affinity: Affinity_WRITE,
- num_order_by: 0,
- order_by: null_mut(),
- }
- }
-
- /// Read from any shard.
- pub fn read_any() -> Self {
- Self {
- affinity: Affinity_READ,
- shard: Shard_ANY,
- num_order_by: 0,
- order_by: null_mut(),
- }
- }
-
- /// Read from all shards.
- pub fn read_all() -> Self {
- Self {
- affinity: Affinity_READ,
- shard: Shard_ALL,
- num_order_by: 0,
- order_by: null_mut(),
- }
- }
-
- /// Read from any shard.
- pub fn write_any() -> Self {
- Self {
- affinity: Affinity_WRITE,
- shard: Shard_ANY,
- num_order_by: 0,
- order_by: null_mut(),
- }
- }
-
- /// Write to all shards.
- pub fn write_all() -> Self {
- Self {
- affinity: Affinity_WRITE,
- shard: Shard_ALL,
- num_order_by: 0,
- order_by: null_mut(),
- }
- }
-
- /// Is this a read?
- pub fn is_read(&self) -> bool {
- self.affinity == Affinity_READ
- }
-
- /// Is this a write?
- pub fn is_write(&self) -> bool {
- self.affinity == Affinity_WRITE
- }
-
- /// This query indicates a transaction a starting, e.g. BEGIN.
- pub fn is_transaction_start(&self) -> bool {
- self.affinity == Affinity_TRANSACTION_START
- }
-
- /// This query indicates a transaction is ending, e.g. COMMIT/ROLLBACK.
- pub fn is_transaction_end(&self) -> bool {
- self.affinity == Affinity_TRANSACTION_END
- }
-
- /// Which shard, if any.
- pub fn shard(&self) -> Option<usize> {
- if self.shard < 0 {
- None
- } else {
- Some(self.shard as usize)
- }
- }
-
- /// Can send query to any shard.
- pub fn is_any_shard(&self) -> bool {
- self.shard == Shard_ANY
- }
-
- /// Send queries to all shards.
- pub fn is_all_shards(&self) -> bool {
- self.shard == Shard_ALL
- }
-
- /// The plugin has no idea where to route this query.
- pub fn is_unknown(&self) -> bool {
- self.shard == Shard_ANY && self.affinity == Affinity_UNKNOWN
- }
-
- /// Add order by columns to the route.
- pub fn order_by(&mut self, order_by: &[OrderBy]) {
- let num_order_by = order_by.len();
- let layout = Layout::array::<OrderBy>(num_order_by).unwrap();
- let ptr = unsafe { alloc(layout) as *mut OrderBy };
- unsafe { copy(order_by.as_ptr(), ptr, num_order_by) };
- self.num_order_by = num_order_by as i32;
- self.order_by = ptr;
- }
-
- /// Deallocate memory.
- ///
- /// # Safety
- ///
- /// Don't use this unless you're cleaning up plugin output.
- pub(crate) unsafe fn deallocate(&self) {
- if self.num_order_by > 0 {
- (0..self.num_order_by).for_each(|index| (*self.order_by.offset(index as isize)).drop());
- let layout = Layout::array::<OrderBy>(self.num_order_by as usize).unwrap();
- dealloc(self.order_by as *mut u8, layout);
- }
- }
-}
diff --git a/pgdog-plugin/src/string.rs b/pgdog-plugin/src/string.rs
new file mode 100644
index 00000000..aecdb47d
--- /dev/null
+++ b/pgdog-plugin/src/string.rs
@@ -0,0 +1,80 @@
+//! Wrapper around Rust's [`str`], a UTF-8 encoded slice.
+//!
+//! This is used to pass strings back and forth between the plugin and
+//! PgDog, without allocating memory as required by [`std::ffi::CString`].
+//!
+//! ### Example
+//!
+//! ```
+//! use pgdog_plugin::PdStr;
+//! use std::ops::Deref;
+//!
+//! let string = PdStr::from("hello world");
+//! assert_eq!(string.deref(), "hello world");
+//!
+//! let string = string.to_string(); // Owned version.
+//! ```
+//!
+use crate::bindings::PdStr;
+use std::{ops::Deref, os::raw::c_void, slice::from_raw_parts, str::from_utf8_unchecked};
+
+impl From<&str> for PdStr {
+ fn from(value: &str) -> Self {
+ PdStr {
+ data: value.as_ptr() as *mut c_void,
+ len: value.len(),
+ }
+ }
+}
+
+impl From<&String> for PdStr {
+ fn from(value: &String) -> Self {
+ PdStr {
+ data: value.as_ptr() as *mut c_void,
+ len: value.len(),
+ }
+ }
+}
+
+impl Deref for PdStr {
+ type Target = str;
+
+ fn deref(&self) -> &Self::Target {
+ unsafe {
+ let slice = from_raw_parts::<u8>(self.data as *mut u8, self.len);
+ from_utf8_unchecked(slice)
+ }
+ }
+}
+
+impl PartialEq for PdStr {
+ fn eq(&self, other: &Self) -> bool {
+ **self == **other
+ }
+}
+
+impl Default for PdStr {
+ fn default() -> Self {
+ Self {
+ len: 0,
+ data: "".as_ptr() as *mut c_void,
+ }
+ }
+}
+
+#[cfg(test)]
+mod test {
+ use super::*;
+
+ #[test]
+ fn test_pd_str() {
+ let s = "one_two_three";
+ let pd = PdStr::from(s);
+ assert_eq!(pd.deref(), "one_two_three");
+
+ let s = String::from("one_two");
+ let pd = PdStr::from(&s);
+ assert_eq!(pd.deref(), "one_two");
+ assert_eq!(&*pd, "one_two");
+ }
+}
diff --git a/pgdog/src/frontend/router/mod.rs b/pgdog/src/frontend/router/mod.rs
index 2629139c..399d00df 100644
--- a/pgdog/src/frontend/router/mod.rs
+++ b/pgdog/src/frontend/router/mod.rs
@@ -4,7 +4,6 @@ pub mod context;
pub mod copy;
pub mod error;
pub mod parser;
-pub mod request;
pub mod round_robin;
pub mod search_path;
pub mod sharding;
diff --git a/pgdog/src/frontend/router/parser/context.rs b/pgdog/src/frontend/router/parser/context.rs
index 4e3f093e..4809e6f4 100644
--- a/pgdog/src/frontend/router/parser/context.rs
+++ b/pgdog/src/frontend/router/parser/context.rs
@@ -1,5 +1,8 @@
//! Shortcut the parser given the cluster config.
+use pgdog_plugin::pg_query::protobuf::ParseResult;
+use pgdog_plugin::{PdRouterContext, PdStatement};
+
use crate::{
backend::ShardingSchema,
config::{config, MultiTenant, ReadWriteStrategy},
@@ -93,4 +96,22 @@ impl<'a> QueryParserContext<'a> {
pub(super) fn multi_tenant(&self) -> &Option<MultiTenant> {
self.multi_tenant
}
+
+ /// Create plugin context.
+ pub(super) fn plugin_context(&self, ast: &ParseResult) -> PdRouterContext {
+ PdRouterContext {
+ shards: self.shards as u64,
+ has_replicas: if self.read_only { 0 } else { 1 },
+ has_primary: if self.write_only { 0 } else { 1 },
+ in_transaction: if self.router_context.in_transaction {
+ 1
+ } else {
+ 0
+ },
+ // SAFETY: ParseResult lives for the entire time the plugin is executed.
+ // We could use lifetimes to guarantee this, but bindgen doesn't generate them.
+ query: unsafe { PdStatement::from_proto(ast) },
+ write_override: 0, // This is set inside `QueryParser::plugins`.
+ }
+ }
}
diff --git a/pgdog/src/frontend/router/parser/query/mod.rs b/pgdog/src/frontend/router/parser/query/mod.rs
index 245e635e..ea7c5e28 100644
--- a/pgdog/src/frontend/router/parser/query/mod.rs
+++ b/pgdog/src/frontend/router/parser/query/mod.rs
@@ -16,11 +16,13 @@ use crate::{
messages::{Bind, Vector},
parameter::ParameterValue,
},
+ plugin::plugins,
};
use super::*;
mod delete;
mod explain;
+mod plugins;
mod select;
mod set;
mod shared;
@@ -29,11 +31,13 @@ mod transaction;
mod update;
use multi_tenant::MultiTenantCheck;
-use pg_query::{
+use pgdog_plugin::pg_query::{
fingerprint,
protobuf::{a_const::Val, *},
NodeEnum,
};
+use plugins::PluginOutput;
+
use tracing::{debug, trace};
/// Query parser.
@@ -56,6 +60,8 @@ pub struct QueryParser {
write_override: bool,
// Currently calculated shard.
shard: Shard,
+ // Plugin read override.
+ plugin_output: PluginOutput,
}
impl Default for QueryParser {
@@ -64,6 +70,7 @@ impl Default for QueryParser {
in_transaction: false,
write_override: false,
shard: Shard::All,
+ plugin_output: PluginOutput::default(),
}
}
}
@@ -261,6 +268,16 @@ impl QueryParser {
_ => Ok(Command::Query(Route::write(None))),
}?;
+ // Run plugins, if any.
+ self.plugins(
+ context,
+ &statement,
+ match &command {
+ Command::Query(query) => query.is_read(),
+ _ => false,
+ },
+ )?;
+
// Overwrite shard using shard we got from a comment, if any.
if let Shard::Direct(shard) = self.shard {
if let Command::Query(ref mut route) = command {
@@ -268,6 +285,18 @@ impl QueryParser {
}
}
+ // Set plugin-specified route, if available.
+ // Plugins override what we calculated above.
+ if let Command::Query(ref mut route) = command {
+ if let Some(read) = self.plugin_output.read {
+ route.set_read_mut(read);
+ }
+
+ if let Some(ref shard) = self.plugin_output.shard {
+ route.set_shard_raw_mut(shard);
+ }
+ }
+
// If we only have one shard, set it.
//
// If the query parser couldn't figure it out,
@@ -280,12 +309,12 @@ impl QueryParser {
}
}
- // Last ditch attempt to route a query to a specific shard.
- //
- // Looking through manual queries to see if we have any
- // with the fingerprint.
- //
if let Command::Query(ref mut route) = command {
+ // Last ditch attempt to route a query to a specific shard.
+ //
+ // Looking through manual queries to see if we have any
+ // with the fingerprint.
+ //
if route.shard().all() {
let databases = databases();
// Only fingerprint the query if some manual queries are configured.
diff --git a/pgdog/src/frontend/router/parser/query/plugins.rs b/pgdog/src/frontend/router/parser/query/plugins.rs
new file mode 100644
index 00000000..e798f157
--- /dev/null
+++ b/pgdog/src/frontend/router/parser/query/plugins.rs
@@ -0,0 +1,89 @@
+use crate::frontend::router::parser::cache::CachedAst;
+use pgdog_plugin::{ReadWrite, Shard as PdShard};
+
+use super::*;
+
+/// Output by one of the plugins.
+#[derive(Default, Debug)]
+pub(super) struct PluginOutput {
+ pub(super) shard: Option<Shard>,
+ pub(super) read: Option<bool>,
+}
+
+impl PluginOutput {
+ fn provided(&self) -> bool {
+ self.shard.is_some() || self.read.is_some()
+ }
+}
+
+impl QueryParser {
+ /// Execute plugins, if any.
+ pub(super) fn plugins(
+ &mut self,
+ context: &QueryParserContext,
+ statement: &CachedAst,
+ read: bool,
+ ) -> Result<(), Error> {
+ // Don't run plugins on Parse only.
+ if context.router_context.bind.is_none() && statement.cached {
+ return Ok(());
+ }
+
+ let plugins = if let Some(plugins) = plugins() {
+ plugins
+ } else {
+ return Ok(());
+ };
+
+ if plugins.is_empty() {
+ return Ok(());
+ }
+
+ // Run plugins, if any.
+ // The first plugin to returns something, wins.
+ debug!("executing {} router plugins", plugins.len());
+
+ let mut context = context.plugin_context(&statement.ast().protobuf);
+ context.write_override = if self.write_override || !read { 1 } else { 0 };
+
+ for plugin in plugins {
+ if let Some(route) = plugin.route(context) {
+ match route.shard.try_into() {
+ Ok(shard) => match shard {
+ PdShard::All => self.plugin_output.shard = Some(Shard::All),
+ PdShard::Direct(shard) => {
+ self.plugin_output.shard = Some(Shard::Direct(shard))
+ }
+ PdShard::Unknown => self.plugin_output.shard = None,
+ },
+ Err(_) => self.plugin_output.shard = None,
+ }
+
+ match route.read_write.try_into() {
+ Ok(ReadWrite::Read) => self.plugin_output.read = Some(true),
+ Ok(ReadWrite::Write) => self.plugin_output.read = Some(false),
+ _ => self.plugin_output.read = None,
+ }
+
+ if self.plugin_output.provided() {
+ debug!(
+ "plugin \"{}\" returned route [{}, {}]",
+ plugin.name(),
+ match self.plugin_output.shard.as_ref() {
+ Some(shard) => format!("shard={}", shard),
+ None => format!("shard=unknown"),
+ },
+ match self.plugin_output.read {
+ Some(read) =>
+ format!("role={}", if read { "replica" } else { "primary" }),
+ None => format!("read=unknown"),
+ }
+ );
+ break;
+ }
+ }
+ }
+
+ Ok(())
+ }
+}
diff --git a/pgdog/src/frontend/router/parser/route.rs b/pgdog/src/frontend/router/parser/route.rs
index 2d7b9c89..e96bbc18 100644
--- a/pgdog/src/frontend/router/parser/route.rs
+++ b/pgdog/src/frontend/router/parser/route.rs
@@ -148,6 +148,10 @@ impl Route {
self
}
+ pub fn set_shard_raw_mut(&mut self, shard: &Shard) {
+ self.shard = shard.clone();
+ }
+
pub fn should_buffer(&self) -> bool {
!self.order_by().is_empty() || !self.aggregate().is_empty() || self.distinct().is_some()
}
diff --git a/pgdog/src/frontend/router/request.rs b/pgdog/src/frontend/router/request.rs
deleted file mode 100644
index bbe6f871..00000000
--- a/pgdog/src/frontend/router/request.rs
+++ /dev/null
@@ -1,46 +0,0 @@
-//! Memory-safe wrapper around the FFI binding to Query.
-use pgdog_plugin::Query;
-use std::{
- ffi::CString,
- ops::{Deref, DerefMut},
-};
-
-use super::Error;
-
-/// Memory-safe wrapper around the FFI binding to Query.
-pub struct Request {
- query: Query,
-}
-
-impl Deref for Request {
- type Target = Query;
- fn deref(&self) -> &Self::Target {
- &self.query
- }
-}
-
-impl DerefMut for Request {
- fn deref_mut(&mut self) -> &mut Self::Target {
- &mut self.query
- }
-}
-
-impl Request {
- /// New query request.
- pub fn new(query: &str) -> Result<Self, Error> {
- Ok(Self {
- query: Query::new(CString::new(query.as_bytes())?),
- })
- }
-
- /// Get constructed query.
- pub fn query(&self) -> Query {
- self.query
- }
-}
-
-impl Drop for Request {
- fn drop(&mut self) {
- unsafe { self.query.deallocate() }
- }
-}
diff --git a/pgdog/src/plugin/mod.rs b/pgdog/src/plugin/mod.rs
index 17e1481f..23cbbf6d 100644
--- a/pgdog/src/plugin/mod.rs
+++ b/pgdog/src/plugin/mod.rs
@@ -1,9 +1,11 @@
//! pgDog plugins.
+use std::ops::Deref;
+
use once_cell::sync::OnceCell;
-use pgdog_plugin::libloading;
use pgdog_plugin::libloading::Library;
use pgdog_plugin::Plugin;
+use pgdog_plugin::{comp, libloading};
use tokio::time::Instant;
use tracing::{debug, error, info, warn};
@@ -33,25 +35,43 @@ pub fn load(names: &[&str]) -> Result<(), libloading::Error> {
let _ = LIBS.set(libs);
+ let rustc_version = comp::rustc_version();
+
let mut plugins = vec![];
for (i, name) in names.iter().enumerate() {
if let Some(lib) = LIBS.get().unwrap().get(i) {
let now = Instant::now();
let plugin = Plugin::load(name, lib);
- if !plugin.valid() {
- warn!("plugin \"{}\" is missing required symbols, skipping", name);
- } else {
- if plugin.init() {
- debug!("plugin \"{}\" initialized", name);
+ // Check Rust compiler version.
+ if let Some(plugin_rustc) = plugin.rustc_version() {
+ if rustc_version != plugin_rustc {
+ warn!("skipping plugin \"{}\" because it was compiled with different compiler version ({})",
+ plugin.name(),
+ plugin_rustc.deref()
+ );
+ continue;
}
- plugins.push(plugin);
- info!(
- "loaded \"{}\" plugin [{:.4}ms]",
- name,
- now.elapsed().as_secs_f64() * 1000.0
+ } else {
+ warn!(
+ "skipping plugin \"{}\" because it doesn't expose its Rust compiler version",
+ plugin.name()
);
+ continue;
}
+
+ if plugin.init() {
+ debug!("plugin \"{}\" initialized", name);
+ }
+
+ info!(
+ "loaded \"{}\" plugin (v{}) [{:.4}ms]",
+ name,
+ plugin.version().unwrap_or_default().deref(),
+ now.elapsed().as_secs_f64() * 1000.0
+ );
+
+ plugins.push(plugin);
}
}
@@ -62,8 +82,10 @@ pub fn load(names: &[&str]) -> Result<(), libloading::Error> {
/// Shutdown plugins.
pub fn shutdown() {
- for plugin in plugins() {
- plugin.fini();
+ if let Some(plugins) = plugins() {
+ for plugin in plugins {
+ plugin.fini();
+ }
}
}
@@ -77,8 +99,8 @@ pub fn plugin(name: &str) -> Option<&Plugin<'_>> {
}
/// Get all loaded plugins.
-pub fn plugins() -> &'static Vec<Plugin<'static>> {
- PLUGINS.get().unwrap()
+pub fn plugins() -> Option<&'static Vec<Plugin<'static>>> {
+ PLUGINS.get()
}
/// Load plugins from config.
diff --git a/plugins/README.md b/plugins/README.md
index c8d86ca3..733ccec3 100644
--- a/plugins/README.md
+++ b/plugins/README.md
@@ -1,12 +1,13 @@
# PgDog plugins
-This directory contains (now and in the future) plugins that ship with PgDog and are built by original author(s)
-or the community. You can use these as-is or modify them to your needs.
+This directory contains plugins that ship with PgDog and are built by original author(s) or by the community. You can use these as-is or modify them to your needs.
## Plugins
-### `pgdog-routing`
+### `pgdog-example-plugin`
-The only plugin in here right now and the catch-all for routing traffic through PgDog. This plugin uses `pg_query.rs` (Rust bindings to `pg_query`)
-to parse queries using the PostgreSQL parser, and splits traffic between primary and replicas. This allows users of this plugin to deploy
-primaries and replicas in one PgDog configuration.
+Example plugin that can be used as reference by the community. It currently records
+when a write was made to a table and, for the next 5 seconds after the write, redirects
+all `SELECT` queries that touch table to the primary.
+
+It's a simple workaround for Postgres replica lag, if you're using batch writes.
diff --git a/plugins/pgdog-example-plugin/Cargo.toml b/plugins/pgdog-example-plugin/Cargo.toml
new file mode 100644
index 00000000..cc180455
--- /dev/null
+++ b/plugins/pgdog-example-plugin/Cargo.toml
@@ -0,0 +1,13 @@
+[package]
+name = "pgdog-example-plugin"
+version = "0.1.0"
+edition = "2024"
+
+[lib]
+crate-type = ["rlib", "cdylib"]
+
+[dependencies]
+pgdog-plugin = "0.1.6"
+once_cell = "1"
+parking_lot = "0.12"
+thiserror = "2"
diff --git a/plugins/pgdog-example-plugin/src/lib.rs b/plugins/pgdog-example-plugin/src/lib.rs
new file mode 100644
index 00000000..f15bbeca
--- /dev/null
+++ b/plugins/pgdog-example-plugin/src/lib.rs
@@ -0,0 +1,42 @@
+//! PgDog example plugin.
+//!
+//! All methods, except the ones generated by the `plugin!` macro, are optional.
+//!
+//! Implementing none of them produces a plugin that doesn't do anything, but it can
+//! still be loaded at runtime.
+//!
+
+pub mod plugin;
+
+use pgdog_plugin::{Context, Route, macros};
+
+// This identifies this library is a PgDog plugin and adds some
+// required methods automatically.
+macros::plugin!();
+
+/// Perform any plugin initialization routines here.
+/// These are running sync on boot, and will block startup util they are finished.
+#[macros::init]
+fn init() {}
+
+/// If defined, this function is called on every query going through PgDog.
+///
+/// It's provided with the AST generated by pg_query and context on how many databases
+/// PgDog is proxying.
+///
+/// N.B. Like all functions called via the FFI interface, it cannot return an error or panic.
+///
+#[macros::route]
+fn route(context: Context) -> Route {
+ crate::plugin::route_query(context).unwrap_or(Route::unknown())
+}
+
+/// Run any code before PgDog is shut down.
+///
+/// This allows for plugins to upload stats to some external service
+/// or perform some cleanup routines.
+///
+/// N.B. This is sync and will prevent PgDog from exiting if it gets stuck.
+///
+#[macros::fini]
+fn shutdown() {}
diff --git a/plugins/pgdog-example-plugin/src/plugin.rs b/plugins/pgdog-example-plugin/src/plugin.rs
new file mode 100644
index 00000000..4bb318b9
--- /dev/null
+++ b/plugins/pgdog-example-plugin/src/plugin.rs
@@ -0,0 +1,120 @@
+use std::{
+ collections::HashMap,
+ time::{Duration, Instant},
+};
+
+use once_cell::sync::Lazy;
+use parking_lot::Mutex;
+use pg_query::{NodeEnum, protobuf::RangeVar};
+use pgdog_plugin::prelude::*;
+use thiserror::Error;
+
+#[derive(Error, Debug)]
+pub enum PluginError {
+ #[error("{0}")]
+ PgQuery(#[from] pg_query::Error),
+
+ #[error("empty query")]
+ EmptyQuery,
+}
+
+static WRITE_TIMES: Lazy<Mutex<HashMap<String, Instant>>> =
+ Lazy::new(|| Mutex::new(HashMap::new()));
+
+/// Route query to a replica or a primary, depending on when was the last time
+/// we wrote to the table.
+pub(crate) fn route_query(context: Context) -> Result<Route, PluginError> {
+ // PgDog really thinks this should be a write.
+ // This could be because there is an INSERT statement in a CTE,
+ // or something else. You could override its decision here, but make
+ // sure you checked the AST first.
+ let write_override = context.write_override();
+
+ let proto = context.statement().protobuf();
+ let root = proto
+ .stmts
+ .first()
+ .ok_or(PluginError::EmptyQuery)?
+ .stmt
+ .as_ref()
+ .ok_or(PluginError::EmptyQuery)?;
+
+ match root.node.as_ref() {
+ Some(NodeEnum::SelectStmt(stmt)) => {
+ if write_override {
+ return Ok(Route::unknown());
+ }
+
+ let table_name = stmt
+ .from_clause
+ .first()
+ .ok_or(PluginError::EmptyQuery)?
+ .node
+ .as_ref()
+ .ok_or(PluginError::EmptyQuery)?;
+
+ if let NodeEnum::RangeVar(RangeVar { relname, .. }) = table_name {
+ // Got info on last write.
+ if let Some(last_write) = { WRITE_TIMES.lock().get(relname).cloned() }
+ && last_write.elapsed() > Duration::from_secs(5)
+ && context.has_replicas()
+ {
+ return Ok(Route::new(Shard::Unknown, ReadWrite::Read));
+ }
+ }
+ }
+ Some(NodeEnum::InsertStmt(stmt)) => {
+ if let Some(ref relation) = stmt.relation {
+ WRITE_TIMES
+ .lock()
+ .insert(relation.relname.clone(), Instant::now());
+ }
+ }
+ Some(NodeEnum::UpdateStmt(stmt)) => {
+ if let Some(ref relation) = stmt.relation {
+ WRITE_TIMES
+ .lock()
+ .insert(relation.relname.clone(), Instant::now());
+ }
+ }
+ Some(NodeEnum::DeleteStmt(stmt)) => {
+ if let Some(ref relation) = stmt.relation {
+ WRITE_TIMES
+ .lock()
+ .insert(relation.relname.clone(), Instant::now());
+ }
+ }
+ _ => {}
+ }
+
+ // Let PgDog decide.
+ Ok(Route::unknown())
+}
+
+#[cfg(test)]
+mod test {
+ use pgdog_plugin::PdStatement;
+
+ use super::*;
+
+ #[test]
+ fn test_routing_plugin() {
+ // Keep protobuf in memory.
+ let proto = pg_query::parse("SELECT * FROM users").unwrap().protobuf;
+ let query = unsafe { PdStatement::from_proto(&proto) };
+ let context = pgdog_plugin::PdRouterContext {
+ shards: 1,
+ has_replicas: 1,
+ has_primary: 1,
+ in_transaction: 0,
+ write_override: 0,
+ query,
+ };
+ let route = route_query(context.into()).unwrap();
+ let read_write: ReadWrite = route.read_write.try_into().unwrap();
+ let shard: Shard = route.shard.try_into().unwrap();
+
+ assert_eq!(read_write, ReadWrite::Read);
+ assert_eq!(shard, Shard::Unknown);
+ }
+}
diff --git a/plugins/pgdog-routing/Cargo.toml b/plugins/pgdog-routing/Cargo.toml
deleted file mode 100644
index 67d61fac..00000000
--- a/plugins/pgdog-routing/Cargo.toml
+++ /dev/null
@@ -1,27 +0,0 @@
-[package]
-name = "pgdog-routing"
-version = "0.1.0"
-edition = "2021"
-license = "AGPL-3.0"
-authors = ["Lev Kokotov <lev.kokotov@gmail.com>"]
-description = "De facto pgDog plugin for routing queries"
-
-[dependencies]
-pgdog-plugin = { path = "../../pgdog-plugin", version = "0.1.1" }
-pg_query = "6.0"
-tracing = "0.1"
-tracing-subscriber = { version = "0.3", features = ["env-filter", "std"] }
-rand = "0.8"
-once_cell = "1"
-regex = "1"
-uuid = { version = "1", features = ["v4"] }
-csv = "1"
-
-[lib]
-crate-type = ["rlib", "cdylib"]
-
-[build-dependencies]
-cc = "1"
-
-[dev-dependencies]
-postgres = {version = "0.19", features = ["with-uuid-1"] }
diff --git a/plugins/pgdog-routing/build.rs b/plugins/pgdog-routing/build.rs
deleted file mode 100644
index 79f815e1..00000000
--- a/plugins/pgdog-routing/build.rs
+++ /dev/null
@@ -1,6 +0,0 @@
-fn main() {
- println!("cargo:rerun-if-changed=postgres_hash/hashfn.c");
- cc::Build::new()
- .file("postgres_hash/hashfn.c")
- .compile("postgres_hash");
-}
diff --git a/plugins/pgdog-routing/postgres_hash/LICENSE b/plugins/pgdog-routing/postgres_hash/LICENSE
deleted file mode 100644
index be2d694b..00000000
--- a/plugins/pgdog-routing/postgres_hash/LICENSE
+++ /dev/null
@@ -1,23 +0,0 @@
-PostgreSQL Database Management System
-(formerly known as Postgres, then as Postgres95)
-
-Portions Copyright (c) 1996-2025, PostgreSQL Global Development Group
-
-Portions Copyright (c) 1994, The Regents of the University of California
-
-Permission to use, copy, modify, and distribute this software and its
-documentation for any purpose, without fee, and without a written agreement
-is hereby granted, provided that the above copyright notice and this
-paragraph and the following two paragraphs appear in all copies.
-
-IN NO EVENT SHALL THE UNIVERSITY OF CALIFORNIA BE LIABLE TO ANY PARTY FOR
-DIRECT, INDIRECT, SPECIAL, INCIDENTAL, OR CONSEQUENTIAL DAMAGES, INCLUDING
-LOST PROFITS, ARISING OUT OF THE USE OF THIS SOFTWARE AND ITS
-DOCUMENTATION, EVEN IF THE UNIVERSITY OF CALIFORNIA HAS BEEN ADVISED OF THE
-POSSIBILITY OF SUCH DAMAGE.
-
-THE UNIVERSITY OF CALIFORNIA SPECIFICALLY DISCLAIMS ANY WARRANTIES,
-INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY
-AND FITNESS FOR A PARTICULAR PURPOSE. THE SOFTWARE PROVIDED HEREUNDER IS
-ON AN "AS IS" BASIS, AND THE UNIVERSITY OF CALIFORNIA HAS NO OBLIGATIONS TO
-PROVIDE MAINTENANCE, SUPPORT, UPDATES, ENHANCEMENTS, OR MODIFICATIONS.
diff --git a/plugins/pgdog-routing/postgres_hash/hashfn.c b/plugins/pgdog-routing/postgres_hash/hashfn.c
deleted file mode 100644
index 22fc98ba..00000000
--- a/plugins/pgdog-routing/postgres_hash/hashfn.c
+++ /dev/null
@@ -1,416 +0,0 @@
-#include <stdint.h>
-
-/*
- * PostgreSQL Database Management System
- * (formerly known as Postgres, then as Postgres95)
- *
- * Portions Copyright (c) 1996-2025, PostgreSQL Global Development Group
- *
- * Portions Copyright (c) 1994, The Regents of the University of California
- *
- * Permission to use, copy, modify, and distribute this software and its
- * documentation for any purpose, without fee, and without a written agreement
- * is hereby granted, provided that the above copyright notice and this
- * paragraph and the following two paragraphs appear in all copies.
- *
- * IN NO EVENT SHALL THE UNIVERSITY OF CALIFORNIA BE LIABLE TO ANY PARTY FOR
- * DIRECT, INDIRECT, SPECIAL, INCIDENTAL, OR CONSEQUENTIAL DAMAGES, INCLUDING
- * LOST PROFITS, ARISING OUT OF THE USE OF THIS SOFTWARE AND ITS
- * DOCUMENTATION, EVEN IF THE UNIVERSITY OF CALIFORNIA HAS BEEN ADVISED OF THE
- * POSSIBILITY OF SUCH DAMAGE.
- *
- * THE UNIVERSITY OF CALIFORNIA SPECIFICALLY DISCLAIMS ANY WARRANTIES,
- * INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY
- * AND FITNESS FOR A PARTICULAR PURPOSE. THE SOFTWARE PROVIDED HEREUNDER IS
- * ON AN "AS IS" BASIS, AND THE UNIVERSITY OF CALIFORNIA HAS NO OBLIGATIONS TO
- * PROVIDE MAINTENANCE, SUPPORT, UPDATES, ENHANCEMENTS, OR MODIFICATIONS.
-*/
-
-#define uint64 uint64_t
-#define uint32 uint32_t
-#define int64 int64_t
-
-/*----------
- * mix -- mix 3 32-bit values reversibly.
- *
- * This is reversible, so any information in (a,b,c) before mix() is
- * still in (a,b,c) after mix().
- *
- * If four pairs of (a,b,c) inputs are run through mix(), or through
- * mix() in reverse, there are at least 32 bits of the output that
- * are sometimes the same for one pair and different for another pair.
- * This was tested for:
- * * pairs that differed by one bit, by two bits, in any combination
- * of top bits of (a,b,c), or in any combination of bottom bits of
- * (a,b,c).
- * * "differ" is defined as +, -, ^, or ~^. For + and -, I transformed
- * the output delta to a Gray code (a^(a>>1)) so a string of 1's (as
- * is commonly produced by subtraction) look like a single 1-bit
- * difference.
- * * the base values were pseudorandom, all zero but one bit set, or
- * all zero plus a counter that starts at zero.
- *
- * This does not achieve avalanche. There are input bits of (a,b,c)
- * that fail to affect some output bits of (a,b,c), especially of a. The
- * most thoroughly mixed value is c, but it doesn't really even achieve
- * avalanche in c.
- *
- * This allows some parallelism. Read-after-writes are good at doubling
- * the number of bits affected, so the goal of mixing pulls in the opposite
- * direction from the goal of parallelism. I did what I could. Rotates
- * seem to cost as much as shifts on every machine I could lay my hands on,
- * and rotates are much kinder to the top and bottom bits, so I used rotates.
- *----------
- */
-#define mix(a,b,c) \
-{ \
- a -= c; a ^= rot(c, 4); c += b; \
- b -= a; b ^= rot(a, 6); a += c; \
- c -= b; c ^= rot(b, 8); b += a; \
- a -= c; a ^= rot(c,16); c += b; \
- b -= a; b ^= rot(a,19); a += c; \
- c -= b; c ^= rot(b, 4); b += a; \
-}
-
-static inline uint32
-pg_rotate_left32(uint32 word, int n)
-{
- return (word << n) | (word >> (32 - n));
-}
-
-#define rot(x,k) pg_rotate_left32(x, k)
-
-#define UINT32_ALIGN_MASK (sizeof(uint32) - 1)
-
-/*----------
- * final -- final mixing of 3 32-bit values (a,b,c) into c
- *
- * Pairs of (a,b,c) values differing in only a few bits will usually
- * produce values of c that look totally different. This was tested for
- * * pairs that differed by one bit, by two bits, in any combination
- * of top bits of (a,b,c), or in any combination of bottom bits of
- * (a,b,c).
- * * "differ" is defined as +, -, ^, or ~^. For + and -, I transformed
- * the output delta to a Gray code (a^(a>>1)) so a string of 1's (as
- * is commonly produced by subtraction) look like a single 1-bit
- * difference.
- * * the base values were pseudorandom, all zero but one bit set, or
- * all zero plus a counter that starts at zero.
- *
- * The use of separate functions for mix() and final() allow for a
- * substantial performance increase since final() does not need to
- * do well in reverse, but is does need to affect all output bits.
- * mix(), on the other hand, does not need to affect all output
- * bits (affecting 32 bits is enough). The original hash function had
- * a single mixing operation that had to satisfy both sets of requirements
- * and was slower as a result.
- *----------
- */
-#define final(a,b,c) \
-{ \
- c ^= b; c -= rot(b,14); \
- a ^= c; a -= rot(c,11); \
- b ^= a; b -= rot(a,25); \
- c ^= b; c -= rot(b,16); \
- a ^= c; a -= rot(c, 4); \
- b ^= a; b -= rot(a,14); \
- c ^= b; c -= rot(b,24); \
-}
-
-#define UINT64CONST(x) UINT64_C(x)
-#define HASH_PARTITION_SEED UINT64CONST(0x7A5B22367996DCFD)
-
-/*
- * Combine two 64-bit hash values, resulting in another hash value, using the
- * same kind of technique as hash_combine(). Testing shows that this also
- * produces good bit mixing.
- */
-uint64
-hash_combine64(uint64 a, uint64 b)
-{
- /* 0x49a0f4dd15e5a8e3 is 64bit random data */
- a ^= b + UINT64CONST(0x49a0f4dd15e5a8e3) + (a << 54) + (a >> 7);
- return a;
-}
-
-/*
- * hash_bytes_extended() -- hash into a 64-bit value, using an optional seed
- * k : the key (the unaligned variable-length array of bytes)
- * len : the length of the key, counting by bytes
- * seed : a 64-bit seed (0 means no seed)
- *
- * Returns a uint64 value. Otherwise similar to hash_bytes.
- */
-uint64
-hash_bytes_extended(const unsigned char *k, int keylen)
-{
- uint32 a,
- b,
- c,
- len;
-
- uint64 seed = HASH_PARTITION_SEED;
-
- /* Set up the internal state */
- len = keylen;
- a = b = c = 0x9e3779b9 + len + 3923095;
-
- /* If the seed is non-zero, use it to perturb the internal state. */
- if (seed != 0)
- {
- /*
- * In essence, the seed is treated as part of the data being hashed,
- * but for simplicity, we pretend that it's padded with four bytes of
- * zeroes so that the seed constitutes a 12-byte chunk.
- */
- a += (uint32) (seed >> 32);
- b += (uint32) seed;
- mix(a, b, c);
- }
-
- /* If the source pointer is word-aligned, we use word-wide fetches */
- if (((uintptr_t) k & UINT32_ALIGN_MASK) == 0)
- {
- /* Code path for aligned source data */
- const uint32 *ka = (const uint32 *) k;
-
- /* handle most of the key */
- while (len >= 12)
- {
- a += ka[0];
- b += ka[1];
- c += ka[2];
- mix(a, b, c);
- ka += 3;
- len -= 12;
- }
-
- /* handle the last 11 bytes */
- k = (const unsigned char *) ka;
-#ifdef WORDS_BIGENDIAN
- switch (len)
- {
- case 11:
- c += ((uint32) k[10] << 8);
- /* fall through */
- case 10:
- c += ((uint32) k[9] << 16);
- /* fall through */
- case 9:
- c += ((uint32) k[8] << 24);
- /* fall through */
- case 8:
- /* the lowest byte of c is reserved for the length */
- b += ka[1];
- a += ka[0];
- break;
- case 7:
- b += ((uint32) k[6] << 8);
- /* fall through */
- case 6:
- b += ((uint32) k[5] << 16);
- /* fall through */
- case 5:
- b += ((uint32) k[4] << 24);
- /* fall through */
- case 4:
- a += ka[0];
- break;
- case 3:
- a += ((uint32) k[2] << 8);
- /* fall through */
- case 2:
- a += ((uint32) k[1] << 16);
- /* fall through */
- case 1:
- a += ((uint32) k[0] << 24);
- /* case 0: nothing left to add */
- }
-#else /* !WORDS_BIGENDIAN */
- switch (len)
- {
- case 11:
- c += ((uint32) k[10] << 24);
- /* fall through */
- case 10:
- c += ((uint32) k[9] << 16);
- /* fall through */
- case 9:
- c += ((uint32) k[8] << 8);
- /* fall through */
- case 8:
- /* the lowest byte of c is reserved for the length */
- b += ka[1];
- a += ka[0];
- break;
- case 7:
- b += ((uint32) k[6] << 16);
- /* fall through */
- case 6:
- b += ((uint32) k[5] << 8);
- /* fall through */
- case 5:
- b += k[4];
- /* fall through */
- case 4:
- a += ka[0];
- break;
- case 3:
- a += ((uint32) k[2] << 16);
- /* fall through */
- case 2:
- a += ((uint32) k[1] << 8);
- /* fall through */
- case 1:
- a += k[0];
- /* case 0: nothing left to add */
- }
-#endif /* WORDS_BIGENDIAN */
- }
- else
- {
- /* Code path for non-aligned source data */
-
- /* handle most of the key */
- while (len >= 12)
- {
-#ifdef WORDS_BIGENDIAN
- a += (k[3] + ((uint32) k[2] << 8) + ((uint32) k[1] << 16) + ((uint32) k[0] << 24));
- b += (k[7] + ((uint32) k[6] << 8) + ((uint32) k[5] << 16) + ((uint32) k[4] << 24));
- c += (k[11] + ((uint32) k[10] << 8) + ((uint32) k[9] << 16) + ((uint32) k[8] << 24));
-#else /* !WORDS_BIGENDIAN */
- a += (k[0] + ((uint32) k[1] << 8) + ((uint32) k[2] << 16) + ((uint32) k[3] << 24));
- b += (k[4] + ((uint32) k[5] << 8) + ((uint32) k[6] << 16) + ((uint32) k[7] << 24));
- c += (k[8] + ((uint32) k[9] << 8) + ((uint32) k[10] << 16) + ((uint32) k[11] << 24));
-#endif /* WORDS_BIGENDIAN */
- mix(a, b, c);
- k += 12;
- len -= 12;
- }
-
- /* handle the last 11 bytes */
-#ifdef WORDS_BIGENDIAN
- switch (len)
- {
- case 11:
- c += ((uint32) k[10] << 8);
- /* fall through */
- case 10:
- c += ((uint32) k[9] << 16);
- /* fall through */
- case 9:
- c += ((uint32) k[8] << 24);
- /* fall through */
- case 8:
- /* the lowest byte of c is reserved for the length */
- b += k[7];
- /* fall through */
- case 7:
- b += ((uint32) k[6] << 8);
- /* fall through */
- case 6:
- b += ((uint32) k[5] << 16);
- /* fall through */
- case 5:
- b += ((uint32) k[4] << 24);
- /* fall through */
- case 4:
- a += k[3];
- /* fall through */
- case 3:
- a += ((uint32) k[2] << 8);
- /* fall through */
- case 2:
- a += ((uint32) k[1] << 16);
- /* fall through */
- case 1:
- a += ((uint32) k[0] << 24);
- /* case 0: nothing left to add */
- }
-#else /* !WORDS_BIGENDIAN */
- switch (len)
- {
- case 11:
- c += ((uint32) k[10] << 24);
- /* fall through */
- case 10:
- c += ((uint32) k[9] << 16);
- /* fall through */
- case 9:
- c += ((uint32) k[8] << 8);
- /* fall through */
- case 8:
- /* the lowest byte of c is reserved for the length */
- b += ((uint32) k[7] << 24);
- /* fall through */
- case 7:
- b += ((uint32) k[6] << 16);
- /* fall through */
- case 6:
- b += ((uint32) k[5] << 8);
- /* fall through */
- case 5:
- b += k[4];
- /* fall through */
- case 4:
- a += ((uint32) k[3] << 24);
- /* fall through */
- case 3:
- a += ((uint32) k[2] << 16);
- /* fall through */
- case 2:
- a += ((uint32) k[1] << 8);
- /* fall through */
- case 1:
- a += k[0];
- /* case 0: nothing left to add */
- }
-#endif /* WORDS_BIGENDIAN */
- }
-
- final(a, b, c);
-
- /* report the result */
- return ((uint64) b << 32) | c;
-}
-
-/*
- * Both the seed and the magic number added at the end are from
- * https://stackoverflow.com/a/67189122
-*/
-
-static uint64
-hash_bytes_uint32_extended(uint32 k)
-{
- uint32 a,
- b,
- c;
- uint64 seed = HASH_PARTITION_SEED;
-
- a = b = c = 0x9e3779b9 + (uint32) sizeof(uint32) + 3923095;
-
- if (seed != 0)
- {
- a += (uint32) (seed >> 32);
- b += (uint32) seed;
- mix(a, b, c);
- }
-
- a += k;
-
- final(a, b, c);
-
- /* report the result */
- return ((uint64) b << 32) | c;
-}
-
-uint64 hashint8extended(int64 val)
-{
- /* Same approach as hashint8 */
- uint32 lohalf = (uint32) val;
- uint32 hihalf = (uint32) (val >> 32);
-
- lohalf ^= (val >= 0) ? hihalf : ~hihalf;
-
- return hash_bytes_uint32_extended(lohalf);
-}
diff --git a/plugins/pgdog-routing/src/comment.rs b/plugins/pgdog-routing/src/comment.rs
deleted file mode 100644
index 45bef8da..00000000
--- a/plugins/pgdog-routing/src/comment.rs
+++ /dev/null
@@ -1,46 +0,0 @@
-//! Parse shards/sharding keys from comments.
-
-use once_cell::sync::Lazy;
-use pg_query::{protobuf::Token, scan, Error};
-use regex::Regex;
-use uuid::Uuid;
-
-use crate::sharding_function;
-
-static SHARD: Lazy<Regex> = Lazy::new(|| Regex::new(r#"pgdog_shard: *([0-9]+)"#).unwrap());
-static SHARDING_KEY: Lazy<Regex> =
- Lazy::new(|| Regex::new(r#"pgdog_sharding_key: *([0-9a-zA-Z]+)"#).unwrap());
-
-/// Extract shard number from a comment.
-///
-/// Comment style uses the C-style comments (not SQL comments!)
-/// as to allow the comment to appear anywhere in the query.
-///
-/// See [`SHARD`] and [`SHARDING_KEY`] for the style of comment we expect.
-///
-pub fn shard(query: &str, shards: usize) -> Result<Option<usize>, Error> {
- let tokens = scan(query)?;
-
- for token in tokens.tokens.iter() {
- if token.token == Token::CComment as i32 {
- let comment = &query[token.start as usize..token.end as usize];
- if let Some(cap) = SHARDING_KEY.captures(comment) {
- if let Some(sharding_key) = cap.get(1) {
- if let Ok(value) = sharding_key.as_str().parse::<i64>() {
- return Ok(Some(sharding_function::bigint(value, shards)));
- }
- if let Ok(value) = sharding_key.as_str().parse::<Uuid>() {
- return Ok(Some(sharding_function::uuid(value, shards)));
- }
- }
- }
- if let Some(cap) = SHARD.captures(comment) {
- if let Some(shard) = cap.get(1) {
- return Ok(shard.as_str().parse::<usize>().ok());
- }
- }
- }
- }
-
- Ok(None)
-}
diff --git a/plugins/pgdog-routing/src/copy.rs b/plugins/pgdog-routing/src/copy.rs
deleted file mode 100644
index 86b7b467..00000000
--- a/plugins/pgdog-routing/src/copy.rs
+++ /dev/null
@@ -1,143 +0,0 @@
-//! Handle COPY.
-
-use csv::ReaderBuilder;
-use pg_query::{protobuf::CopyStmt, NodeEnum};
-use pgdog_plugin::bindings::*;
-
-use crate::sharding_function::bigint;
-
-/// Parse COPY statement.
-pub fn parse(stmt: &CopyStmt) -> Result<Copy, pg_query::Error> {
- if !stmt.is_from {
- return Ok(Copy::invalid());
- }
-
- if let Some(ref rel) = stmt.relation {
- let mut headers = false;
- let mut csv = false;
- let mut delimiter = ',';
-
- let mut columns = vec![];
-
- for column in &stmt.attlist {
- if let Some(NodeEnum::String(ref column)) = column.node {
- columns.push(column.sval.as_str());
- }
- }
-
- for option in &stmt.options {
- if let Some(NodeEnum::DefElem(ref elem)) = option.node {
- match elem.defname.to_lowercase().as_str() {
- "format" => {
- if let Some(ref arg) = elem.arg {
- if let Some(NodeEnum::String(ref string)) = arg.node {
- if string.sval.to_lowercase().as_str() == "csv" {
- csv = true;
- }
- }
- }
- }
-
- "delimiter" => {
- if let Some(ref arg) = elem.arg {
- if let Some(NodeEnum::String(ref string)) = arg.node {
- delimiter = string.sval.chars().next().unwrap_or(',');
- }
- }
- }
-
- "header" => {
- headers = true;
- }
-
- _ => (),
- }
- }
- }
-
- if csv {
- return Ok(Copy::new(&rel.relname, headers, delimiter, &columns));
- }
- }
-
- Ok(Copy::invalid())
-}
-
-/// Split copy data into individual rows
-/// and determine where each row should go.
-pub fn copy_data(input: CopyInput, shards: usize) -> Result<CopyOutput, csv::Error> {
- let data = input.data();
- let mut csv = ReaderBuilder::new()
- .has_headers(input.headers())
- .delimiter(input.delimiter() as u8)
- .from_reader(data);
-
- let mut rows = vec![];
-
- while let Some(record) = csv.records().next() {
- let record = record?;
- if let Some(position) = record.position() {
- let start = position.byte() as usize;
- let end = start + record.as_slice().len();
- // N.B.: includes \n character which indicates the end of a single CSV record.
- // If CSV is encoded using Windows \r\n, this will break.
- if let Some(row_data) = data.get(start..=end + 1) {
- let key = record.iter().nth(input.sharding_column());
- let shard = key
- .and_then(|k| k.parse::<i64>().ok().map(|k| bigint(k, shards) as i64))
- .unwrap_or(-1);
-
- let row = CopyRow::new(row_data, shard as i32);
- rows.push(row);
- }
- }
- }
-
- Ok(CopyOutput::new(&rows).with_header(if csv.has_headers() {
- csv.headers().ok().map(|s| {
- s.into_iter()
- .collect::<Vec<_>>()
- .join(input.delimiter().to_string().as_str())
- + "\n" // New line indicating the end of a CSV line.
- })
- } else {
- None
- }))
-}
-
-#[cfg(test)]
-mod test {
-
- use super::*;
-
- #[test]
- fn test_copy() {
- let stmt = "COPY test_table FROM 'some_file.csv' CSV HEADER DELIMITER ';'";
- let ast = pg_query::parse(stmt).unwrap();
- let copy = ast.protobuf.stmts.first().unwrap().stmt.clone().unwrap();
-
- let copy = match copy.node {
- Some(NodeEnum::CopyStmt(ref stmt)) => parse(stmt).unwrap(),
- _ => panic!("not COPY"),
- };
-
- assert_eq!(copy.copy_format, CopyFormat_CSV);
- assert_eq!(copy.delimiter(), ';');
- assert!(copy.has_headers());
- assert_eq!(copy.table_name(), "test_table");
-
- let data = "id;email\n1;test@test.com\n2;admin@test.com\n";
- let input = CopyInput::new(data.as_bytes(), 0, copy.has_headers(), ';');
- let output = copy_data(input, 4).unwrap();
-
- let mut rows = output.rows().iter();
- assert_eq!(rows.next().unwrap().shard, bigint(1, 4) as i32);
- assert_eq!(rows.next().unwrap().shard, bigint(2, 4) as i32);
- assert_eq!(output.header(), Some("id;email\n"));
-
- unsafe {
- copy.deallocate();
- output.deallocate();
- }
- }
-}
diff --git a/plugins/pgdog-routing/src/lib.rs b/plugins/pgdog-routing/src/lib.rs
deleted file mode 100644
index 131aca75..00000000
--- a/plugins/pgdog-routing/src/lib.rs
+++ /dev/null
@@ -1,137 +0,0 @@
-//! Parse queries using pg_query and route all SELECT queries
-//! to replicas. All other queries are routed to a primary.
-
-use once_cell::sync::Lazy;
-use pg_query::{parse, NodeEnum};
-use pgdog_plugin::bindings::{Config, Input, Output};
-use pgdog_plugin::Route;
-
-use tracing::{debug, level_filters::LevelFilter};
-use tracing::{error, trace};
-use tracing_subscriber::{fmt, prelude::*, EnvFilter};
-
-use std::io::IsTerminal;
-use std::sync::atomic::{AtomicUsize, Ordering};
-
-static SHARD_ROUND_ROBIN: Lazy<AtomicUsize> = Lazy::new(|| AtomicUsize::new(0));
-
-pub mod comment;
-pub mod copy;
-pub mod order_by;
-pub mod sharding_function;
-
-#[no_mangle]
-pub extern "C" fn pgdog_init() {
- let format = fmt::layer()
- .with_ansi(std::io::stderr().is_terminal())
- .with_file(false);
-
- let filter = EnvFilter::builder()
- .with_default_directive(LevelFilter::INFO.into())
- .from_env_lossy();
-
- tracing_subscriber::registry()
- .with(format)
- .with(filter)
- .init();
-
- // TODO: This is more for fun/demo, but in prod, we want
- // this logger to respect options passed to pgDog proper, e.g.
- // use JSON output.
- debug!("🐕 pgDog routing plugin v{}", env!("CARGO_PKG_VERSION"));
-}
-
-#[no_mangle]
-pub extern "C" fn pgdog_route_query(input: Input) -> Output {
- if let Some(query) = input.query() {
- match route_internal(query.query(), input.config) {
- Ok(output) => output,
- Err(_) => Output::new_forward(Route::unknown()),
- }
- } else if let Some(copy_input) = input.copy() {
- match copy::copy_data(copy_input, input.config.shards as usize) {
- Ok(output) => Output::new_copy_rows(output),
- Err(err) => {
- error!("{:?}", err);
- Output::skip()
- }
- }
- } else {
- Output::skip()
- }
-}
-
-fn route_internal(query: &str, config: Config) -> Result<Output, pg_query::Error> {
- let shards = config.shards;
- let databases = config.databases();
-
- // Shortcut for typical single shard replicas-only/primary-only deployments.
- if shards == 1 {
- let read_only = databases.iter().all(|d| d.replica());
- let write_only = databases.iter().all(|d| d.primary());
- if read_only {
- return Ok(Output::new_forward(Route::read(0)));
- }
- if write_only {
- return Ok(Output::new_forward(Route::read(0)));
- }
- }
-
- let ast = parse(query)?;
- trace!("{:#?}", ast);
-
- let shard = comment::shard(query, shards as usize)?;
-
- // For cases like SELECT NOW(), or SELECT 1, etc.
- let tables = ast.tables();
- if tables.is_empty() && shard.is_none() {
- // Better than random for load distribution.
- let shard_counter = SHARD_ROUND_ROBIN.fetch_add(1, Ordering::Relaxed);
- return Ok(Output::new_forward(Route::read(
- shard_counter % shards as usize,
- )));
- }
-
- if let Some(query) = ast.protobuf.stmts.first() {
- if let Some(ref node) = query.stmt {
- match node.node {
- Some(NodeEnum::SelectStmt(ref stmt)) => {
- let order_by = order_by::extract(stmt)?;
- let mut route = if let Some(shard) = shard {
- Route::read(shard)
- } else {
- Route::read_all()
- };
-
- if !order_by.is_empty() {
- route.order_by(&order_by);
- }
-
- return Ok(Output::new_forward(route));
- }
-
- Some(NodeEnum::CopyStmt(ref stmt)) => {
- return Ok(Output::new_copy(copy::parse(stmt)?))
- }
-
- Some(_) => (),
-
- None => (),
- }
- }
- }
-
- Ok(if let Some(shard) = shard {
- Output::new_forward(Route::write(shard))
- } else {
- Output::new_forward(Route::write_all())
- })
-}
-
-#[no_mangle]
-pub extern "C" fn pgdog_fini() {
- debug!(
- "🐕 pgDog routing plugin v{} shutting down",
- env!("CARGO_PKG_VERSION")
- );
-}
diff --git a/plugins/pgdog-routing/src/order_by.rs b/plugins/pgdog-routing/src/order_by.rs
deleted file mode 100644
index 2e1083aa..00000000
--- a/plugins/pgdog-routing/src/order_by.rs
+++ /dev/null
@@ -1,56 +0,0 @@
-//! Handle the ORDER BY clause.
-
-use pg_query::{
- protobuf::{a_const::*, *},
- Error, NodeEnum,
-};
-use pgdog_plugin::*;
-
-/// Extract sorting columns.
-///
-/// If a query spans multiple shards, this allows pgDog to apply
-/// sorting rules Postgres used and show the rows in the correct order.
-///
-pub fn extract(stmt: &SelectStmt) -> Result<Vec<OrderBy>, Error> {
- let mut order_by = vec![];
- for clause in &stmt.sort_clause {
- if let Some(NodeEnum::SortBy(ref sort_by)) = clause.node {
- let asc = matches!(sort_by.sortby_dir, 0..=2);
- if let Some(ref node) = sort_by.node {
- if let Some(ref node) = node.node {
- match node {
- NodeEnum::AConst(aconst) => {
- if let Some(Val::Ival(ref integer)) = aconst.val {
- order_by.push(OrderBy::column_index(
- integer.ival as usize,
- if asc {
- OrderByDirection_ASCENDING
- } else {
- OrderByDirection_DESCENDING
- },
- ));
- }
- }
-
- NodeEnum::ColumnRef(column_ref) => {
- if let Some(field) = column_ref.fields.first() {
- if let Some(NodeEnum::String(ref string)) = field.node {
- order_by.push(OrderBy::column_name(
- &string.sval,
- if asc {
- OrderByDirection_ASCENDING
- } else {
- OrderByDirection_DESCENDING
- },
- ));
- }
- }
- }
- _ => (),
- }
- }
- }
- }
- }
- Ok(order_by)
-}
diff --git a/plugins/pgdog-routing/src/sharding_function.rs b/plugins/pgdog-routing/src/sharding_function.rs
deleted file mode 100644
index 8bc8df10..00000000
--- a/plugins/pgdog-routing/src/sharding_function.rs
+++ /dev/null
@@ -1,143 +0,0 @@
-//! PostgreSQL hash functions.
-//!
-//! This module delegates most of the hashing work directly
-//! to PostgreSQL internal functions that we copied in `postgres_hash` C library.
-//!
-
-use uuid::Uuid;
-
-#[link(name = "postgres_hash")]
-extern "C" {
- /// Hash any size data using its bytes representation.
- fn hash_bytes_extended(k: *const u8, keylen: i64) -> u64;
- /// Special hashing function for BIGINT (i64).
- fn hashint8extended(k: i64) -> u64;
- /// Combine multiple hashes into one in the case of multi-column hashing keys.
- fn hash_combine64(a: u64, b: u64) -> u64;
-}
-
-/// Safe wrapper around `hash_bytes_extended`.
-fn hash_slice(k: &[u8]) -> u64 {
- unsafe { hash_bytes_extended(k.as_ptr(), k.len() as i64) }
-}
-
-/// Calculate shard for a BIGINT value.
-pub fn bigint(value: i64, shards: usize) -> usize {
- let hash = unsafe { hashint8extended(value) };
- let combined = unsafe { hash_combine64(0, hash as u64) };
-
- combined as usize % shards
-}
-
-/// Calculate shard for a UUID value.
-pub fn uuid(value: Uuid, shards: usize) -> usize {
- let hash = hash_slice(value.as_bytes().as_slice());
- let combined = unsafe { hash_combine64(0, hash) };
-
- combined as usize % shards
-}
-
-#[cfg(test)]
-mod test {
- use super::*;
- use postgres::{Client, NoTls};
- use rand::Rng;
-
- #[test]
- fn test_bigint() {
- let tables = r#"
- BEGIN;
-
- DROP TABLE IF EXISTS sharding_func CASCADE;
-
- CREATE TABLE sharding_func (id BIGINT)
- PARTITION BY HASH(id);
-
- CREATE TABLE sharding_func_0
- PARTITION OF sharding_func
- FOR VALUES WITH (modulus 4, remainder 0);
-
- CREATE TABLE sharding_func_1
- PARTITION OF sharding_func
- FOR VALUES WITH (modulus 4, remainder 1);
-
- CREATE TABLE sharding_func_2
- PARTITION OF sharding_func
- FOR VALUES WITH (modulus 4, remainder 2);
-
- CREATE TABLE sharding_func_3
- PARTITION OF sharding_func
- FOR VALUES WITH (modulus 4, remainder 3);
- "#;
-
- let mut client = Client::connect(
- "host=localhost user=pgdog password=pgdog dbname=pgdog",
- NoTls,
- )
- .expect("client to connect");
-
- client.batch_execute(tables).expect("create tables");
-
- for _ in 0..4096 {
- let v = rand::thread_rng().gen::<i64>();
- // Our hashing function.
- let shard = bigint(v as i64, 4);
-
- // Check that Postgres did the same thing.
- // Note: we are inserting directly into the subtable.
- let table = format!("sharding_func_{}", shard);
- client
- .query(&format!("INSERT INTO {} (id) VALUES ($1)", table), &[&v])
- .expect("insert");
- }
- }
-
- #[test]
- fn test_uuid() {
- let tables = r#"
- BEGIN;
-
- DROP TABLE IF EXISTS sharding_func_uuid CASCADE;
-
- CREATE TABLE sharding_func_uuid (id UUID)
- PARTITION BY HASH(id);
-
- CREATE TABLE sharding_func_uuid_0
- PARTITION OF sharding_func_uuid
- FOR VALUES WITH (modulus 4, remainder 0);
-
- CREATE TABLE sharding_func_uuid_1
- PARTITION OF sharding_func_uuid
- FOR VALUES WITH (modulus 4, remainder 1);
-
- CREATE TABLE sharding_func_uuid_2
- PARTITION OF sharding_func_uuid
- FOR VALUES WITH (modulus 4, remainder 2);
-
- CREATE TABLE sharding_func_uuid_3
- PARTITION OF sharding_func_uuid
- FOR VALUES WITH (modulus 4, remainder 3);
- "#;
-
- let mut client = Client::connect(
- "host=localhost user=pgdog password=pgdog dbname=pgdog",
- NoTls,
- )
- .expect("client to connect");
-
- client.batch_execute(tables).expect("create tables");
-
- for _ in 0..4096 {
- let v = Uuid::new_v4();
- // Our hashing function.
- let shard = uuid(v, 4);
-
- // Check that Postgres did the same thing.
- // Note: we are inserting directly into the subtable.
- let table = format!("sharding_func_uuid_{}", shard);
- client
- .query(&format!("INSERT INTO {} (id) VALUES ($1)", table), &[&v])
- .expect("insert");
- }
- }
-}
[parent: 2ce762d02087]