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 -[![Documentation](https://img.shields.io/badge/documentation-blue?style=flat)](https://pgdog.dev) +[![Documentation](https://img.shields.io/badge/documentation-blue?style=flat)](https://docsrs.pgdog.dev/pgdog_plugin/index.html) [![Latest crate](https://img.shields.io/crates/v/pgdog-plugin.svg)](https://crates.io/crates/pgdog-plugin) -[![Reference docs](https://img.shields.io/docsrs/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]