postjevsql.git / tests / support / postgres.rs
1//! A throwaway Postgres per test, run from the devshell's binaries of
2//! the major buck built for (`POSTJEVSQL_POSTGRES_BIN`,
3//! `POSTJEVSQL_PG_MAJOR`): `initdb` into a temp dir, `postgres` as a
4//! child process, the built extension served in place. Same shape as a
5//! testcontainers module (builder, `start`, teardown on drop), without
6//! docker.
7//!
8//! The library is found through `dynamic_library_path` on every major.
9//! The control and install script are found through
10//! `extension_control_path` from PG18. PG17 has no such setting and
11//! reads them only from its own sharedir, which it locates relative to
12//! the `postgres` it was started as: nixpkgs' relative-to-symlinks
13//! patch keeps it from resolving the link. So on PG17 the instance runs
14//! from a prefix of symlinks in its temp dir, whose
15//! `share/postgresql/extension` also links the built tree's files
16//! (measured on 17, 2026-09-26).
17//!
18//! Teardown cannot leak: the postmaster gets `PR_SET_PDEATHSIG`, so the
19//! kernel sends it SIGQUIT (immediate shutdown) when the test thread
20//! that spawned it exits, even if the test process is SIGKILLed.
21//!
22//! At most one instance per CPU runs at once, machine-wide: each holds
23//! an `flock` on one of `available_parallelism()` slot files in
24//! `/tmp/postjevsql-test-slots-<uid>` for its life. buck runs several
25//! test binaries at once, and each runs one test per core, each with its
26//! own server, so without the cap a full suite on 3 cores reached a load
27//! average of 28 (2026-09-26). The path is fixed rather than
28//! `temp_dir()` so every worktree's daemon and every test sandbox share
29//! the one set of slots.
30//!
31//! A test that needs several instances at once takes their slots together
32//! with `TestPostgres::start_all`, all or none. Taking them one at a time
33//! is hold-and-wait: on 3 cores, three `sidecar_cli` tests each held a
34//! target's slot and waited for a sidecar's until buck killed them at 10
35//! minutes (2026-09-27). So a second `start` on a thread that already
36//! holds a slot panics, naming `start_all`, instead of risking that. At
37//! least 2 slots, so a pair always fits.
38//!
39//! It listens only on a Unix socket in its own temp dir, never on TCP. A
40//! port picked by binding 0 and releasing it can be taken by any other
41//! process before the postmaster binds it, and then the test either
42//! waits for a server that never starts or talks to another test's.
43
44use std::fs::File;
45use std::os::fd::AsRawFd;
46use std::os::unix::process::CommandExt;
47use std::collections::HashMap;
48use std::path::{Path, PathBuf};
49use std::process::{Child, Command, Stdio};
50use std::sync::Mutex;
51use std::thread::ThreadId;
52use std::time::{Duration, Instant};
53
54use tempfile::TempDir;
55
56pub struct TestPostgres {
57    ext: Option<PathBuf>,
58    settings: Vec<(String, String)>,
59    env: Vec<(String, String)>,
60}
61
62pub struct Instance {
63    child: Child,
64    dir: TempDir,
65    // Released after `Drop` has stopped the server.
66    _slot: Slot,
67}
68
69impl Default for TestPostgres {
70    fn default() -> Self {
71        Self::new()
72    }
73}
74
75impl TestPostgres {
76    pub fn new() -> Self {
77        TestPostgres {
78            ext: None,
79            settings: vec![("fsync".into(), "off".into())],
80            env: Vec::new(),
81        }
82    }
83
84    /// Serves the built extension tree without installing it anywhere.
85    pub fn with_extension(mut self, ext_dir: &Path) -> Self {
86        let ext = ext_dir.canonicalize().expect("extension tree exists");
87        let path = ext.to_str().expect("utf-8 path").to_owned();
88        let controls = has_extension_control_path();
89        self.ext = Some(ext);
90        let this = self.setting("dynamic_library_path", &format!("{path}/lib:$libdir"));
91        if controls {
92            this.setting("extension_control_path", &format!("{path}/share:$system"))
93        } else {
94            this
95        }
96    }
97
98    pub fn setting(mut self, name: &str, value: &str) -> Self {
99        self.settings.push((name.into(), value.into()));
100        self
101    }
102
103    pub fn env_var(mut self, name: &str, value: &str) -> Self {
104        self.env.push((name.into(), value.into()));
105        self
106    }
107
108    /// One instance. A test that needs two at once uses `start_all`.
109    pub async fn start(self) -> Instance {
110        let [instance] = Self::start_all([self]).await;
111        instance
112    }
113
114    /// Several instances whose slots are taken together, all or none.
115    pub async fn start_all<const N: usize>(configs: [TestPostgres; N]) -> [Instance; N] {
116        let mut slots = acquire_slots(N).await.into_iter();
117        let mut instances = Vec::with_capacity(N);
118        for config in configs {
119            instances.push(config.launch(slots.next().expect("one slot per instance")).await);
120        }
121        let Ok(instances) = instances.try_into() else { unreachable!("one instance per config") };
122        instances
123    }
124
125    async fn launch(self, slot: Slot) -> Instance {
126        let dir = tempfile::Builder::new().prefix("postjevsql-pg").tempdir().expect("temp dir");
127        let bin = match &self.ext {
128            Some(ext) if !has_extension_control_path() => linked_prefix(dir.path(), ext),
129            _ => PathBuf::from(env("POSTJEVSQL_POSTGRES_BIN")),
130        };
131        let data = dir.path().join("data");
132
133        let status = Command::new(bin.join("initdb"))
134            .args(["--username=postgres", "--auth=trust", "--encoding=UTF8", "--locale=C", "--no-sync"])
135            .arg("-D")
136            .arg(&data)
137            .stdout(Stdio::null())
138            .status()
139            .expect("initdb runs");
140        assert!(status.success(), "initdb failed");
141
142        let mut postgres = Command::new(bin.join("postgres"));
143        postgres
144            .arg("-D")
145            .arg(&data)
146            .args(["-c", "listen_addresses="])
147            .arg("-c")
148            .arg(format!("unix_socket_directories={}", dir.path().display()))
149            // Only what the test sets: a developer's real key or libpq
150            // settings must never reach the instance under test.
151            .env_clear()
152            .env("PATH", std::env::var_os("PATH").unwrap_or_default())
153            .envs(self.env)
154            .stdout(Stdio::null());
155        for (name, value) in &self.settings {
156            postgres.arg("-c").arg(format!("{name}={value}"));
157        }
158        // SAFETY: prctl is async-signal-safe and touches no Rust state.
159        unsafe {
160            postgres.pre_exec(|| {
161                if libc::prctl(libc::PR_SET_PDEATHSIG, libc::SIGQUIT) != 0 {
162                    return Err(std::io::Error::last_os_error());
163                }
164                Ok(())
165            });
166        }
167        let child = postgres.spawn().expect("postgres starts");
168        let instance = Instance { child, dir, _slot: slot };
169        instance.wait_ready().await;
170        instance
171    }
172}
173
174impl Instance {
175    pub fn conn_str(&self) -> String {
176        format!("host={} user=postgres dbname=postgres", self.dir.path().display())
177    }
178
179    async fn wait_ready(&self) {
180        let deadline = Instant::now() + Duration::from_secs(30);
181        loop {
182            match tokio_postgres::connect(&self.conn_str(), tokio_postgres::NoTls).await {
183                Ok(_) => return,
184                Err(e) if Instant::now() > deadline => panic!("postgres never became ready: {e}"),
185                Err(_) => tokio::time::sleep(Duration::from_millis(20)).await,
186            }
187        }
188    }
189}
190
191impl Drop for Instance {
192    fn drop(&mut self) {
193        // SIGQUIT: immediate shutdown, which also stops the backends.
194        // SAFETY: the pid is our own un-reaped child, so it is not reused.
195        unsafe { libc::kill(self.child.id() as libc::pid_t, libc::SIGQUIT) };
196        let _ = self.child.wait();
197    }
198}
199
200/// A locked slot file, counted against the thread that took it. Closing
201/// the file releases the lock, and the kernel closes it if the process
202/// dies. The file is close-on-exec, so the server never inherits it.
203struct Slot {
204    _file: File,
205    thread: ThreadId,
206}
207
208/// Slots held per thread, so a second acquisition on one thread is caught.
209static HELD: Mutex<Option<HashMap<ThreadId, usize>>> = Mutex::new(None);
210
211fn held(thread: ThreadId, delta: isize) -> usize {
212    let mut map = HELD.lock().unwrap_or_else(|e| e.into_inner());
213    let count = map.get_or_insert_with(HashMap::new).entry(thread).or_default();
214    *count = count.checked_add_signed(delta).expect("slot count");
215    *count
216}
217
218impl Drop for Slot {
219    fn drop(&mut self) {
220        held(self.thread, -1);
221    }
222}
223
224/// Waits until `n` slots are free at once and returns them locked. Never
225/// holds some while waiting for the rest.
226async fn acquire_slots(n: usize) -> Vec<Slot> {
227    let thread = std::thread::current().id();
228    assert!(
229        held(thread, 0) == 0,
230        "this thread already holds a Postgres slot; start every instance a test needs \
231         together with TestPostgres::start_all, or waiting for a second can deadlock"
232    );
233    // SAFETY: getuid cannot fail and touches no Rust state.
234    let dir = PathBuf::from(format!("/tmp/postjevsql-test-slots-{}", unsafe { libc::getuid() }));
235    std::fs::create_dir_all(&dir).expect("slot dir");
236    let slots = std::thread::available_parallelism().map_or(2, |n| n.get().max(2));
237    assert!(n <= slots, "{n} instances at once, but only {slots} slots");
238    loop {
239        let mut files = Vec::with_capacity(n);
240        for i in 0..slots {
241            let file = File::options()
242                .create(true)
243                .truncate(false)
244                .write(true)
245                .open(dir.join(i.to_string()))
246                .expect("slot file");
247            // SAFETY: the fd is open and owned by `file`.
248            if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) } == 0 {
249                files.push(file);
250                if files.len() == n {
251                    held(thread, n as isize);
252                    return files.into_iter().map(|_file| Slot { _file, thread }).collect();
253                }
254            }
255        }
256        // Too few free: release what was taken before waiting.
257        drop(files);
258        tokio::time::sleep(Duration::from_millis(50)).await;
259    }
260}
261
262/// `extension_control_path` arrived in PG18.
263fn has_extension_control_path() -> bool {
264    let major: u32 = env("POSTJEVSQL_PG_MAJOR").parse().expect("POSTJEVSQL_PG_MAJOR is a number");
265    major >= 18
266}
267
268/// A prefix of symlinks to the server's installation, with the built
269/// tree's control and SQL files added to its extension directory.
270/// Returns its `bin`.
271fn linked_prefix(dir: &Path, ext: &Path) -> PathBuf {
272    let real = PathBuf::from(env("POSTJEVSQL_POSTGRES_BIN"));
273    let real = real.parent().expect("bin has a parent");
274    let prefix = dir.join("prefix");
275    link_tree(real, &prefix, &[Path::new("share/postgresql/extension")]);
276    let extension = prefix.join("share/postgresql/extension");
277    for entry in std::fs::read_dir(ext.join("share/extension")).expect("built extension share") {
278        let entry = entry.expect("dir entry");
279        std::os::unix::fs::symlink(entry.path(), extension.join(entry.file_name())).expect("link extension file");
280    }
281    prefix.join("bin")
282}
283
284/// Mirrors `from` at `to`: a symlink per entry, except that each of
285/// `open` (relative to `from`) and its ancestors is a real directory
286/// whose entries are linked in turn.
287fn link_tree(from: &Path, to: &Path, open: &[&Path]) {
288    std::fs::create_dir_all(to).expect("prefix dir");
289    for entry in std::fs::read_dir(from).expect("installation dir") {
290        let entry = entry.expect("dir entry");
291        let name = PathBuf::from(entry.file_name());
292        let inner: Vec<&Path> = open.iter().filter_map(|p| p.strip_prefix(&name).ok()).collect();
293        if inner.is_empty() {
294            std::os::unix::fs::symlink(entry.path(), to.join(&name)).expect("link installation entry");
295        } else {
296            let inner: Vec<&Path> = inner.into_iter().filter(|p| !p.as_os_str().is_empty()).collect();
297            link_tree(&entry.path(), &to.join(&name), &inner);
298        }
299    }
300}
301
302fn env(name: &str) -> String {
303    std::env::var(name).unwrap_or_else(|_| panic!("{name} is unset; run under buck2 test"))
304}