Files
truf-server/tests/test_postgres_runtime.py
2026-09-30 20:30:56 +03:00

1403 lines
70 KiB
Python

from concurrent.futures import Future
import ctypes
from datetime import datetime, timezone
import os
from pathlib import Path
import signal
import subprocess
from types import SimpleNamespace
import sys
import tempfile
import unittest
from unittest import mock
ROOT = Path(__file__).resolve().parents[1]
APP_DIR = ROOT / 'app'
sys.path.insert(0, str(APP_DIR))
import postgres_runtime
import process_identity
from paths import apply_path_config
from postgres_runtime import (
ClusterIdentityError,
OnlineUnavailable,
PostgresBackend,
PostgresController,
PostgresState,
ProbeKind,
ProbeResult,
StartResult,
StopResult,
)
from runtime_security import ensure_private_directory
class PostgresRuntimePathTests(unittest.TestCase):
def setUp(self):
self.enterContext(mock.patch.object(
postgres_runtime, 'require_trusted_native_executable', side_effect=lambda path: path,
))
@staticmethod
def identity(paths, data_directory):
executable_keys = ('postgres', 'pg_ctl', 'pg_isready', 'pg_controldata')
return {
'schema': postgres_runtime.CLUSTER_IDENTITY_SCHEMA,
'private_file_ready': True,
'data_directory': data_directory,
'pg_major': 16,
'executables': {key: paths[key] for key in executable_keys},
'executable_sha256': {key: 'fixture-hash' for key in executable_keys},
'database': 'truf',
'user': 'truf',
'port': 5432,
'system_identifier': '12345',
'created_at': '2026-01-01T00:00:00+00:00',
}
def test_default_data_directory_is_unchanged(self):
with tempfile.TemporaryDirectory() as temp_dir:
runtime_dir = os.path.join(temp_dir, 'runtime')
paths = postgres_runtime.postgres_runtime_paths({
'global': {'runtime_dir': runtime_dir},
})
self.assertEqual(
paths['data_dir'],
postgres_runtime.canonical_path(os.path.join(runtime_dir, 'postgres', 'data')),
)
def test_external_data_directory_expands_without_moving_runtime_assets(self):
with tempfile.TemporaryDirectory() as temp_dir:
config = apply_path_config({
'global': {
'root_dir': temp_dir,
'project_dir': os.path.join(temp_dir, 'app'),
'runtime_dir': '{root_dir}/runtime',
'postgres_data_dir': '{root_dir}/external/../external/data',
},
})
paths = postgres_runtime.postgres_runtime_paths(config)
postgres_dir = postgres_runtime.canonical_path(os.path.join(temp_dir, 'runtime', 'postgres'))
self.assertEqual(
paths['data_dir'],
postgres_runtime.canonical_path(os.path.join(temp_dir, 'external', 'data')),
)
self.assertEqual(paths['postgres_dir'], postgres_dir)
self.assertEqual(paths['log_dir'], postgres_runtime.canonical_path(os.path.join(postgres_dir, 'logs')))
self.assertEqual(
postgres_runtime.canonical_path(paths['identity_path']),
postgres_runtime.canonical_path(os.path.join(postgres_dir, 'cluster_identity.json')),
)
def test_external_data_identity_mismatch_remains_fail_closed(self):
with tempfile.TemporaryDirectory() as temp_dir:
runtime_dir = os.path.join(temp_dir, 'runtime')
external_data = os.path.join(temp_dir, 'external-data')
config = {'global': {
'runtime_dir': runtime_dir,
'postgres_data_dir': external_data,
}}
paths = postgres_runtime.postgres_runtime_paths(config)
stale_identity = self.identity(
paths,
postgres_runtime.canonical_path(os.path.join(runtime_dir, 'postgres', 'data')),
)
with mock.patch.object(postgres_runtime, 'read_private_json', return_value=stale_identity), \
mock.patch.object(postgres_runtime, 'configured_cluster_values', return_value={
'database': 'truf', 'user': 'truf', 'port': 5432,
}):
with self.assertRaisesRegex(ClusterIdentityError, 'data_directory mismatch'):
postgres_runtime.verify_cluster_identity(config)
def test_external_data_identity_verifies_only_when_exactly_bound(self):
with tempfile.TemporaryDirectory() as temp_dir:
config = {'global': {
'runtime_dir': os.path.join(temp_dir, 'runtime'),
'postgres_data_dir': os.path.join(temp_dir, 'external-data'),
}}
paths = postgres_runtime.postgres_runtime_paths(config)
identity = self.identity(paths, paths['data_dir'])
with mock.patch.object(postgres_runtime, 'read_private_json', return_value=identity), \
mock.patch.object(postgres_runtime, 'configured_cluster_values', return_value={
'database': 'truf', 'user': 'truf', 'port': 5432,
}), mock.patch.object(postgres_runtime.os.path, 'isfile', return_value=True), \
mock.patch.object(postgres_runtime, 'sha256_file', return_value='fixture-hash'), \
mock.patch.object(postgres_runtime, '_data_directory_major', return_value=16), \
mock.patch.object(postgres_runtime, '_system_identifier', return_value='12345'):
verified = postgres_runtime.verify_cluster_identity(config)
self.assertEqual(verified['data_directory'], paths['data_dir'])
def test_bootstrap_major_executes_and_matches_bundled_postgres(self):
with tempfile.TemporaryDirectory() as temp_dir:
data_dir = os.path.join(temp_dir, 'data')
os.makedirs(data_dir)
Path(os.path.join(data_dir, 'PG_VERSION')).write_text('16\n', encoding='ascii')
paths = {'data_dir': data_dir, 'postgres': os.path.join(temp_dir, 'postgres.exe')}
completed = SimpleNamespace(returncode=0, stdout='postgres (PostgreSQL) 16.4')
with mock.patch.object(postgres_runtime, '_run_bounded', return_value=completed) as run:
self.assertEqual(postgres_runtime._bootstrap_postgres_major(paths), 16)
run.assert_called_once_with([paths['postgres'], '--version'], timeout=10)
completed.stdout = 'postgres (PostgreSQL) 17.0'
with mock.patch.object(postgres_runtime, '_run_bounded', return_value=completed), \
self.assertRaisesRegex(ClusterIdentityError, 'does not match'):
postgres_runtime._bootstrap_postgres_major(paths)
def test_bootstrap_identity_uses_executable_version_validation(self):
paths = {
'data_dir': 'data', 'identity_path': 'identity',
'postgres': 'postgres', 'pg_ctl': 'pg_ctl',
'pg_isready': 'pg_isready', 'pg_controldata': 'pg_controldata',
}
with mock.patch.object(postgres_runtime, 'postgres_runtime_paths', return_value=paths), \
mock.patch.object(postgres_runtime.os.path, 'isdir', return_value=True), \
mock.patch.object(postgres_runtime.os.path, 'isfile', return_value=True), \
mock.patch.object(postgres_runtime, '_require_bootstrap_offline'), \
mock.patch.object(postgres_runtime, '_bootstrap_postgres_major', return_value=16) as major, \
mock.patch.object(postgres_runtime, 'sha256_file', return_value='fixture-hash'), \
mock.patch.object(postgres_runtime, '_system_identifier', return_value='12345'), \
mock.patch.object(postgres_runtime, 'write_private_json_exclusive'), \
mock.patch.object(postgres_runtime, 'verify_cluster_identity', return_value={'verified': True}):
self.assertEqual(postgres_runtime.bootstrap_cluster_identity({}), {'verified': True})
major.assert_called_once_with(paths)
def test_periodic_identity_verification_never_probes_postgres_version(self):
with tempfile.TemporaryDirectory() as temp_dir:
config = {'global': {'runtime_dir': os.path.join(temp_dir, 'runtime')}}
paths = postgres_runtime.postgres_runtime_paths(config)
identity = self.identity(paths, paths['data_dir'])
with mock.patch.object(postgres_runtime, 'read_private_json', return_value=identity), \
mock.patch.object(postgres_runtime, 'configured_cluster_values', return_value={
'database': 'truf', 'user': 'truf', 'port': 5432,
}), mock.patch.object(postgres_runtime.os.path, 'isfile', return_value=True), \
mock.patch.object(postgres_runtime, 'sha256_file', return_value='fixture-hash'), \
mock.patch.object(postgres_runtime, '_data_directory_major', return_value=16), \
mock.patch.object(postgres_runtime, '_system_identifier', return_value='12345'), \
mock.patch.object(postgres_runtime, '_bundled_postgres_major') as executable_major:
postgres_runtime.verify_cluster_identity(config)
executable_major.assert_not_called()
def test_periodic_identity_verification_still_rejects_executable_hash_drift(self):
with tempfile.TemporaryDirectory() as temp_dir:
config = {'global': {'runtime_dir': os.path.join(temp_dir, 'runtime')}}
paths = postgres_runtime.postgres_runtime_paths(config)
identity = self.identity(paths, paths['data_dir'])
with mock.patch.object(postgres_runtime, 'read_private_json', return_value=identity), \
mock.patch.object(postgres_runtime, 'configured_cluster_values', return_value={
'database': 'truf', 'user': 'truf', 'port': 5432,
}), mock.patch.object(postgres_runtime.os.path, 'isfile', return_value=True), \
mock.patch.object(postgres_runtime, 'sha256_file', return_value='drifted-hash'):
with self.assertRaisesRegex(ClusterIdentityError, 'executable hash mismatch'):
postgres_runtime.verify_cluster_identity(config)
def test_periodic_identity_verification_still_rejects_pg_version_drift(self):
with tempfile.TemporaryDirectory() as temp_dir:
config = {'global': {'runtime_dir': os.path.join(temp_dir, 'runtime')}}
paths = postgres_runtime.postgres_runtime_paths(config)
identity = self.identity(paths, paths['data_dir'])
with mock.patch.object(postgres_runtime, 'read_private_json', return_value=identity), \
mock.patch.object(postgres_runtime, 'configured_cluster_values', return_value={
'database': 'truf', 'user': 'truf', 'port': 5432,
}), mock.patch.object(postgres_runtime.os.path, 'isfile', return_value=True), \
mock.patch.object(postgres_runtime, 'sha256_file', return_value='fixture-hash'), \
mock.patch.object(postgres_runtime, '_data_directory_major', return_value=17):
with self.assertRaisesRegex(ClusterIdentityError, 'major version changed'):
postgres_runtime.verify_cluster_identity(config)
def test_periodic_identity_verification_still_rejects_system_identifier_drift(self):
with tempfile.TemporaryDirectory() as temp_dir:
config = {'global': {'runtime_dir': os.path.join(temp_dir, 'runtime')}}
paths = postgres_runtime.postgres_runtime_paths(config)
identity = self.identity(paths, paths['data_dir'])
with mock.patch.object(postgres_runtime, 'read_private_json', return_value=identity), \
mock.patch.object(postgres_runtime, 'configured_cluster_values', return_value={
'database': 'truf', 'user': 'truf', 'port': 5432,
}), mock.patch.object(postgres_runtime.os.path, 'isfile', return_value=True), \
mock.patch.object(postgres_runtime, 'sha256_file', return_value='fixture-hash'), \
mock.patch.object(postgres_runtime, '_data_directory_major', return_value=16), \
mock.patch.object(postgres_runtime, '_system_identifier', return_value='99999') as system_id:
with self.assertRaisesRegex(ClusterIdentityError, 'system identifier changed'):
postgres_runtime.verify_cluster_identity(config)
system_id.assert_called_once_with(paths)
def test_live_probe_uses_online_system_identifier_without_pg_controldata(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.config = {}
backend._verified_identity = None
identity = {'port': 5432, 'data_directory': 'data'}
with mock.patch.object(backend, '_root_job_error', return_value=None), \
mock.patch.object(backend, '_prune_logging'), \
mock.patch.object(backend, '_online_query', return_value=(False, 'epoch')), \
mock.patch.object(postgres_runtime, 'verify_cluster_identity', return_value=identity) as verify:
result = backend.probe()
self.assertEqual(result.kind, ProbeKind.READY)
verify.assert_called_once_with({}, verify_offline_system_identifier=False)
def test_offline_probe_verifies_system_identifier_before_reporting_stopped(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.config = {}
backend._verified_identity = None
backend._expected_process = None
backend._start_requested_wall_time = None
backend._accepted_start_at_monotonic = None
identity = {'port': 5432, 'data_directory': 'data'}
with mock.patch.object(backend, '_root_job_error', return_value=None), \
mock.patch.object(backend, '_prune_logging'), \
mock.patch.object(backend, '_online_query', side_effect=OnlineUnavailable('offline')), \
mock.patch.object(backend, '_expected_is_live', return_value=False), \
mock.patch.object(backend, '_hint_running', return_value=False), \
mock.patch.object(postgres_runtime, '_parse_postmaster_pid', return_value=None), \
mock.patch.object(postgres_runtime, '_listener_present', return_value=False), \
mock.patch.object(postgres_runtime, 'verify_cluster_identity', return_value=identity) as verify:
result = backend.probe()
self.assertEqual(result.kind, ProbeKind.STOPPED)
self.assertEqual(verify.call_args_list, [
mock.call({}, verify_offline_system_identifier=False),
mock.call({}),
])
def test_offline_probe_accepts_stale_pid_for_definitively_absent_process(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.config = {}
backend._verified_identity = None
backend._expected_process = None
backend._start_requested_wall_time = None
backend._accepted_start_at_monotonic = None
identity = {'port': 5432, 'data_directory': 'data'}
with mock.patch.object(backend, '_root_job_error', return_value=None), \
mock.patch.object(backend, '_prune_logging'), \
mock.patch.object(backend, '_online_query', side_effect=OnlineUnavailable('offline')), \
mock.patch.object(backend, '_expected_is_live', return_value=False), \
mock.patch.object(backend, '_open_postmaster', side_effect=postgres_runtime.PostmasterProcessAbsent('absent')), \
mock.patch.object(backend, '_hint_running', return_value=False), \
mock.patch.object(postgres_runtime, '_parse_postmaster_pid', return_value=20892), \
mock.patch.object(postgres_runtime, '_listener_present', return_value=False), \
mock.patch.object(postgres_runtime, 'verify_cluster_identity', return_value=identity) as verify:
result = backend.probe()
self.assertEqual(result.kind, ProbeKind.STOPPED)
self.assertEqual(verify.call_args_list, [
mock.call({}, verify_offline_system_identifier=False),
mock.call({}),
])
@unittest.skipUnless(os.name == 'nt', 'Windows retained process semantics')
def test_open_process_rejects_definitively_exited_retained_handle(self):
def exited(_handle, code_pointer):
code_pointer._obj.value = 0
return True
with mock.patch.object(process_identity, '_OPEN_PROCESS', return_value=1234), \
mock.patch.object(process_identity, '_WAIT_FOR_SINGLE_OBJECT', return_value=0), \
mock.patch.object(process_identity, '_GET_EXIT_CODE_PROCESS', side_effect=exited), \
mock.patch.object(process_identity, '_CLOSE_HANDLE') as close_handle, \
mock.patch.object(process_identity, '_windows_identity') as inspect_identity:
with self.assertRaisesRegex(process_identity.ProcessExitedError, 'already exited'):
process_identity.open_process(20892)
inspect_identity.assert_not_called()
close_handle.assert_called_once_with(1234)
@unittest.skipUnless(os.name == 'nt', 'Windows retained process semantics')
def test_live_uninspectable_process_remains_fail_closed(self):
def still_active(_handle, code_pointer):
code_pointer._obj.value = 259
return True
native_error = ctypes.WinError(31)
with mock.patch.object(process_identity, '_OPEN_PROCESS', return_value=1234), \
mock.patch.object(process_identity, '_WAIT_FOR_SINGLE_OBJECT', return_value=258), \
mock.patch.object(process_identity, '_GET_EXIT_CODE_PROCESS', side_effect=still_active), \
mock.patch.object(process_identity, '_CLOSE_HANDLE') as close_handle, \
mock.patch.object(process_identity, '_windows_identity', side_effect=native_error):
with self.assertRaises(process_identity.ProcessIdentityError) as caught:
process_identity.open_process(20892)
self.assertNotIsInstance(caught.exception, process_identity.ProcessExitedError)
self.assertIs(caught.exception.__cause__, native_error)
close_handle.assert_called_once_with(1234)
@unittest.skipUnless(os.name == 'nt', 'Windows retained process semantics')
def test_retained_windows_liveness_uses_wait_handle_not_still_active_code(self):
identity = process_identity.ProcessIdentity(
pid=20892, creation_time='windows-filetime:1',
creation_time_unix=1.0, executable=os.path.abspath(sys.executable),
in_job=False,
)
retained = process_identity.RetainedProcess(identity, handle=1234)
with mock.patch.object(
process_identity, '_WAIT_FOR_SINGLE_OBJECT', return_value=258,
), mock.patch.object(
process_identity, '_GET_EXIT_CODE_PROCESS',
side_effect=AssertionError('STILL_ACTIVE must not decide liveness'),
):
self.assertTrue(retained.is_running())
with mock.patch.object(
process_identity, '_WAIT_FOR_SINGLE_OBJECT', return_value=0,
):
self.assertFalse(retained.is_running())
retained._closed = True
def test_retained_posix_termination_prefers_exact_pidfd_signal(self):
identity = process_identity.ProcessIdentity(
pid=20892, creation_time='proc-start-ticks:1',
creation_time_unix=1.0, executable='/usr/bin/python3',
in_job=False,
)
retained = process_identity.RetainedProcess(identity, pidfd=77)
with mock.patch.object(process_identity.os, 'name', 'posix'), \
mock.patch.object(
process_identity.signal, 'pidfd_send_signal', create=True,
) as send, mock.patch.object(
process_identity.os, 'kill',
side_effect=AssertionError('numeric PID signal is forbidden with pidfd'),
):
retained.terminate()
send.assert_called_once_with(77, signal.SIGTERM, None, 0)
retained._closed = True
def test_posix_identity_binding_opens_pidfd_before_proc_identity(self):
events = []
identity = process_identity.ProcessIdentity(
pid=20892, creation_time='proc-start-ticks:1',
creation_time_unix=1.0, executable='/usr/bin/python3',
in_job=False,
)
def pin(_pid, _flags):
events.append('pidfd_open')
return 77
def inspect(_pid):
events.append('identity')
return identity
with mock.patch.object(process_identity.os, 'name', 'posix'), \
mock.patch.object(
process_identity.os, 'pidfd_open', side_effect=pin, create=True,
), mock.patch.object(
process_identity, '_pidfd_live', return_value=True,
), mock.patch.object(
process_identity, '_posix_identity', side_effect=inspect,
):
retained = process_identity.open_process(20892)
self.assertEqual(events, ['pidfd_open', 'identity', 'identity'])
retained._closed = True
def test_posix_pidfd_exit_during_identity_binding_fails_closed(self):
identity = process_identity.ProcessIdentity(
pid=20892, creation_time='proc-start-ticks:1',
creation_time_unix=1.0, executable='/usr/bin/python3',
in_job=False,
)
with mock.patch.object(process_identity.os, 'name', 'posix'), \
mock.patch.object(
process_identity.os, 'pidfd_open', return_value=77, create=True,
), mock.patch.object(
process_identity, '_pidfd_live', side_effect=(True, False),
), mock.patch.object(
process_identity, '_posix_identity', return_value=identity,
), mock.patch.object(process_identity.os, 'close') as close:
with self.assertRaises(process_identity.ProcessExitedError):
process_identity.open_process(20892)
close.assert_called_once_with(77)
def test_postmaster_treats_only_typed_exited_process_as_absent(self):
with tempfile.TemporaryDirectory() as temp_dir:
data_dir = os.path.join(temp_dir, 'data')
os.makedirs(data_dir)
Path(os.path.join(data_dir, 'postmaster.pid')).write_text(
f'20892\n{data_dir}\n1700000000\n',
encoding='ascii',
)
backend = PostgresBackend.__new__(PostgresBackend)
identity = {
'data_directory': postgres_runtime.canonical_path(data_dir),
'executables': {'postgres': postgres_runtime.canonical_path('postgres.exe')},
}
with mock.patch.object(
postgres_runtime,
'open_process',
side_effect=process_identity.ProcessExitedError('exited'),
):
with self.assertRaisesRegex(postgres_runtime.PostmasterProcessAbsent, 'exited process'):
backend._open_postmaster(identity)
unknown = process_identity.ProcessIdentityError('uninspectable')
unknown.__cause__ = ctypes.WinError(31) if os.name == 'nt' else OSError(5, 'uninspectable')
with mock.patch.object(postgres_runtime, 'open_process', side_effect=unknown):
with self.assertRaises(ClusterIdentityError) as caught:
backend._open_postmaster(identity)
self.assertNotIsInstance(caught.exception, postgres_runtime.PostmasterProcessAbsent)
class ImmediateExecutor:
def submit(self, function):
future = Future()
try:
future.set_result(function())
except BaseException as exc:
future.set_exception(exc)
return future
class ControlledExecutor:
def __init__(self):
self.submissions = []
def submit(self, function):
future = Future()
self.submissions.append((function.__name__, future))
return future
class FakeBackend:
def __init__(self, probes, starts=None, stops=None):
self.probes = list(probes)
self.starts = list(starts or [StartResult(True, 'accepted')])
self.stops = list(stops or [StopResult(True, True, 'stopped')])
self.probe_calls = 0
self.start_calls = 0
self.stop_calls = 0
self.close_calls = 0
def probe(self):
self.probe_calls += 1
return self.probes.pop(0)
def start(self):
self.start_calls += 1
return self.starts.pop(0)
def stop(self):
self.stop_calls += 1
result = self.stops.pop(0)
if isinstance(result, BaseException):
raise result
return result
def close(self):
self.close_calls += 1
def immediate_controller(backend, **options):
return PostgresController(
backend,
executor=ImmediateExecutor(),
health_interval_sec=options.pop('health_interval_sec', 5),
stable_ready_interval_sec=options.pop('stable_ready_interval_sec', 10),
**options,
)
class PostgresControllerTests(unittest.TestCase):
def test_maintenance_start_waits_for_authenticated_ready(self):
backend = FakeBackend([
ProbeResult(ProbeKind.RECOVERING, 'starting'),
ProbeResult(ProbeKind.READY, 'authenticated'),
])
probe = postgres_runtime.maintenance_start({}, backend=backend)
self.assertEqual(probe.kind, ProbeKind.READY)
self.assertEqual(backend.start_calls, 1)
self.assertEqual(backend.stop_calls, 0)
self.assertEqual(backend.close_calls, 0)
def test_maintenance_stop_refuses_foreign_cluster(self):
backend = FakeBackend([], stops=[StopResult(False, False, 'foreign')])
with self.assertRaisesRegex(ClusterIdentityError, 'foreign'):
postgres_runtime.maintenance_stop({}, backend=backend)
self.assertEqual(backend.stop_calls, 1)
self.assertEqual(backend.close_calls, 0)
def test_failed_maintenance_start_leaves_supplied_backend_with_its_caller(self):
backend = FakeBackend([ProbeResult(ProbeKind.OWNED_START_UNCERTAIN, 'unconfirmed start')])
with self.assertRaisesRegex(ClusterIdentityError, 'unconfirmed start'):
postgres_runtime.maintenance_start({}, backend=backend)
self.assertEqual(backend.start_calls, 1)
self.assertEqual(backend.stop_calls, 0)
self.assertEqual(backend.close_calls, 0)
def test_bundled_utility_subprocess_uses_machine_readable_locale(self):
completed = SimpleNamespace(returncode=0, stdout='')
with mock.patch.object(postgres_runtime.subprocess, 'run', return_value=completed) as run, \
mock.patch.object(postgres_runtime, 'require_trusted_native_executable', side_effect=lambda path: path):
self.assertIs(postgres_runtime._run_bounded(['pg_controldata'], 1), completed)
self.assertEqual(run.call_args.kwargs['env']['LC_ALL'], 'C')
self.assertEqual(run.call_args.kwargs['env']['LANG'], 'C')
def test_start_accepted_then_early_death_enters_30_second_backoff(self):
backend = FakeBackend([
ProbeResult(ProbeKind.STOPPED, 'offline'),
ProbeResult(ProbeKind.STOPPED, 'early death'),
])
controller = immediate_controller(backend)
controller.tick(0)
controller.tick(0)
controller.tick(0)
self.assertEqual(controller.state, PostgresState.STARTING)
self.assertEqual(controller.failures, 0)
controller.tick(0)
self.assertEqual(controller.state, PostgresState.BACKOFF)
self.assertEqual(controller.failures, 1)
self.assertEqual(controller.next_action_at, 30)
self.assertEqual(backend.start_calls, 1)
def test_live_recovery_is_only_reprobed_and_never_restarted(self):
backend = FakeBackend([
ProbeResult(ProbeKind.RECOVERING, 'expected recovery'),
ProbeResult(ProbeKind.RECOVERING, 'still recovering'),
])
controller = immediate_controller(backend)
controller.tick(0)
controller.tick(0)
self.assertEqual(controller.state, PostgresState.RECOVERING)
controller.tick(5)
controller.tick(5)
self.assertEqual(controller.state, PostgresState.RECOVERING)
self.assertEqual(backend.start_calls, 0)
self.assertEqual(backend.stop_calls, 0)
def test_failure_streak_resets_only_after_stable_ready_interval(self):
backend = FakeBackend([
ProbeResult(ProbeKind.READY, 'ready once'),
ProbeResult(ProbeKind.READY, 'ready and stable'),
])
controller = immediate_controller(backend, stable_ready_interval_sec=10)
controller.failures = 3
controller.tick(0)
controller.tick(0)
self.assertEqual(controller.state, PostgresState.STABILIZING)
self.assertEqual(controller.failures, 3)
controller.tick(5)
controller.tick(10)
self.assertEqual(controller.state, PostgresState.READY)
self.assertEqual(controller.failures, 0)
def test_periodic_ready_probe_does_not_reenter_stabilizing(self):
backend = FakeBackend([
ProbeResult(ProbeKind.READY, 'initial'),
ProbeResult(ProbeKind.READY, 'stable'),
ProbeResult(ProbeKind.READY, 'periodic'),
])
controller = immediate_controller(backend, stable_ready_interval_sec=5)
controller.tick(0)
controller.tick(0)
controller.tick(5)
controller.tick(5)
self.assertEqual(controller.state, PostgresState.READY)
controller.tick(10)
controller.tick(10)
self.assertEqual(controller.state, PostgresState.READY)
def test_same_postmaster_transient_timeout_holds_ready_during_grace(self):
backend = FakeBackend([
ProbeResult(ProbeKind.READY, 'initial', 'epoch-a'),
ProbeResult(ProbeKind.RECOVERING, 'connection timeout', 'epoch-a'),
ProbeResult(ProbeKind.RECOVERING, 'connection timeout', 'epoch-a'),
])
controller = immediate_controller(
backend,
stable_ready_interval_sec=0,
health_interval_sec=15,
ready_loss_grace_sec=45,
)
controller.tick(0)
controller.tick(0)
self.assertEqual(controller.state, PostgresState.READY)
controller.tick(15)
controller.tick(15)
self.assertEqual(controller.state, PostgresState.READY)
self.assertIn('transient readiness loss', controller.detail)
controller.tick(30)
controller.tick(30)
self.assertEqual(controller.state, PostgresState.READY)
def test_same_postmaster_timeout_exceeding_grace_enters_recovery(self):
backend = FakeBackend([
ProbeResult(ProbeKind.READY, 'initial', 'epoch-a'),
ProbeResult(ProbeKind.RECOVERING, 'connection timeout', 'epoch-a'),
ProbeResult(ProbeKind.RECOVERING, 'connection timeout', 'epoch-a'),
])
controller = immediate_controller(
backend,
stable_ready_interval_sec=0,
health_interval_sec=45,
ready_loss_grace_sec=45,
)
controller.tick(0)
controller.tick(0)
controller.tick(45)
controller.tick(45)
self.assertEqual(controller.state, PostgresState.READY)
controller.tick(90)
controller.tick(90)
self.assertEqual(controller.state, PostgresState.RECOVERING)
def test_different_postmaster_recovery_bypasses_ready_loss_grace(self):
backend = FakeBackend([
ProbeResult(ProbeKind.READY, 'initial', 'epoch-a'),
ProbeResult(ProbeKind.RECOVERING, 'replacement unavailable', 'epoch-b'),
])
controller = immediate_controller(
backend,
stable_ready_interval_sec=0,
health_interval_sec=15,
ready_loss_grace_sec=45,
)
controller.tick(0)
controller.tick(0)
controller.tick(15)
controller.tick(15)
self.assertEqual(controller.state, PostgresState.RECOVERING)
def test_replaced_postmaster_must_stabilize_again(self):
backend = FakeBackend([
ProbeResult(ProbeKind.READY, 'first', 'epoch-a'),
ProbeResult(ProbeKind.READY, 'stable', 'epoch-a'),
ProbeResult(ProbeKind.READY, 'replacement', 'epoch-b'),
])
controller = immediate_controller(backend, stable_ready_interval_sec=5)
controller.tick(0)
controller.tick(0)
controller.tick(5)
controller.tick(5)
self.assertEqual(controller.state, PostgresState.READY)
controller.tick(10)
controller.tick(10)
self.assertEqual(controller.state, PostgresState.STABILIZING)
def test_foreign_identity_never_receives_lifecycle_action(self):
backend = FakeBackend([
ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, 'foreign listener'),
])
controller = immediate_controller(backend)
controller.tick(0)
controller.tick(0)
controller.tick(1000)
self.assertEqual(controller.state, PostgresState.FOREIGN_OR_CONFIG_ERROR)
self.assertEqual(backend.start_calls, 0)
self.assertEqual(backend.stop_calls, 0)
def test_late_probe_result_after_shutdown_is_ignored(self):
backend = FakeBackend([])
executor = ControlledExecutor()
controller = PostgresController(backend, executor=executor)
controller.tick(0)
self.assertEqual(executor.submissions[0][0], 'probe')
controller.request_stop()
executor.submissions[0][1].set_result(ProbeResult(ProbeKind.STOPPED, 'late'))
controller.tick(0)
self.assertEqual(controller.state, PostgresState.STOPPING)
self.assertEqual([name for name, _ in executor.submissions], ['probe', 'stop'])
self.assertEqual(backend.start_calls, 0)
executor.submissions[1][1].set_result(StopResult(True, True, 'stopped'))
controller.tick(0)
self.assertEqual(controller.state, PostgresState.STOPPED)
def test_shutdown_before_first_tick_performs_identity_safe_stop(self):
backend = FakeBackend([])
controller = immediate_controller(backend)
controller.request_stop()
controller.tick(0)
self.assertEqual(controller.state, PostgresState.STOPPING)
controller.tick(0)
self.assertEqual(controller.state, PostgresState.STOPPED)
self.assertEqual(backend.probe_calls, 0)
self.assertEqual(backend.start_calls, 0)
self.assertEqual(backend.stop_calls, 1)
def test_stop_refusal_and_exception_are_stop_failed(self):
refused = immediate_controller(
FakeBackend([], stops=[StopResult(False, False, 'recovery left running')]),
)
refused.request_stop()
refused.tick(0)
refused.tick(0)
self.assertEqual(refused.state, PostgresState.STOP_FAILED)
self.assertFalse(refused.stop_succeeded)
failed = immediate_controller(FakeBackend([], stops=[RuntimeError('worker failed')]))
failed.request_stop()
failed.tick(0)
failed.tick(0)
self.assertEqual(failed.state, PostgresState.STOP_FAILED)
self.assertIn('worker failed', failed.detail)
def test_foreign_state_is_revalidated_without_automatic_start(self):
backend = FakeBackend([
ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, 'foreign listener'),
ProbeResult(ProbeKind.STOPPED, 'listener disappeared'),
])
controller = immediate_controller(backend, health_interval_sec=5)
controller.tick(0)
controller.tick(0)
controller.tick(5)
controller.tick(5)
self.assertEqual(controller.state, PostgresState.FOREIGN_OR_CONFIG_ERROR)
self.assertEqual(backend.probe_calls, 2)
self.assertEqual(backend.start_calls, 0)
def test_foreign_start_result_latches_inert_across_later_stopped_probe(self):
backend = FakeBackend(
[
ProbeResult(ProbeKind.STOPPED, 'offline'),
ProbeResult(ProbeKind.STOPPED, 'still offline'),
],
starts=[StartResult(False, 'identity changed', foreign_or_config_error=True)],
)
controller = immediate_controller(backend, health_interval_sec=5)
controller.tick(0)
controller.tick(0)
controller.tick(0)
self.assertEqual(controller.state, PostgresState.FOREIGN_OR_CONFIG_ERROR)
self.assertTrue(controller.snapshot()['lifecycle_inert'])
controller.tick(5)
controller.tick(5)
self.assertEqual(controller.state, PostgresState.FOREIGN_OR_CONFIG_ERROR)
self.assertEqual(backend.start_calls, 1)
def test_recovery_adoption_remains_inert_until_online_identity(self):
backend = FakeBackend([
ProbeResult(ProbeKind.RECOVERING, 'adopted recovery'),
ProbeResult(ProbeKind.STOPPED, 'recovery disappeared'),
])
controller = immediate_controller(backend, health_interval_sec=5)
controller.tick(0)
controller.tick(0)
self.assertTrue(controller.snapshot()['lifecycle_inert'])
controller.tick(5)
controller.tick(5)
self.assertEqual(controller.state, PostgresState.FOREIGN_OR_CONFIG_ERROR)
self.assertEqual(backend.start_calls, 0)
def test_incomplete_stop_remains_nonterminal_until_result_arrives(self):
backend = FakeBackend([])
executor = ControlledExecutor()
controller = PostgresController(backend, executor=executor)
controller.request_stop()
controller.tick(0)
self.assertFalse(controller.terminal)
self.assertEqual(controller.state, PostgresState.STOPPING)
def test_late_accepted_start_after_authority_drift_schedules_compensating_stop(self):
backend = FakeBackend([])
executor = ControlledExecutor()
controller = PostgresController(backend, executor=executor)
controller.tick(0)
executor.submissions[0][1].set_result(ProbeResult(ProbeKind.STOPPED, 'offline'))
controller.tick(0)
self.assertEqual([name for name, _ in executor.submissions], ['probe', 'start'])
controller.inhibit_lifecycle('authority drift')
executor.submissions[1][1].set_result(StartResult(True, 'late accepted'))
controller.tick(0)
self.assertEqual([name for name, _ in executor.submissions], ['probe', 'start', 'stop'])
self.assertEqual(controller.state, PostgresState.STOPPING)
self.assertFalse(controller.stop_succeeded)
controller.request_stop()
executor.submissions[2][1].set_result(StopResult(True, True, 'compensated'))
controller.tick(0)
self.assertEqual(controller.state, PostgresState.STOPPED)
class FakeRetainedPostmaster:
def __init__(self, pid=321, creation_time='epoch'):
self.pid = pid
self.identity = SimpleNamespace(creation_time=creation_time, creation_time_unix=1700000000.0)
self.closed = False
def is_running(self):
return not self.closed
def close(self):
self.closed = True
class StoppableRetainedPostmaster(FakeRetainedPostmaster):
def __init__(self, pid=321, creation_time='epoch'):
super().__init__(pid, creation_time)
self.running = True
def is_running(self):
return self.running and not self.closed
class PostgresBackendTests(unittest.TestCase):
def test_authenticated_handoff_can_stop_after_sql_readiness_is_lost(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.config = {}
backend.paths = {}
backend.connect_timeout_sec = 1
backend.query_timeout_ms = 1000
backend.stop_timeout_sec = 5
backend._expected_process = None
backend._start_requested_wall_time = None
backend._accepted_start_at_monotonic = None
backend._started_postmaster_observed = False
backend._authenticated_postmaster = None
process = StoppableRetainedPostmaster()
identity = {
'data_directory': postgres_runtime.canonical_path('fixture-data'),
'database': 'truf', 'user': 'truf', 'port': 55432, 'pg_major': 16,
'system_identifier': '123456789',
'executables': {'postgres': 'fixture-postgres', 'pg_ctl': 'fixture-pg_ctl'},
}
connection = mock.Mock()
connection.execute.return_value.fetchone.side_effect = [
{'database': 'truf', 'user_name': 'truf', 'port': 55432,
'data_directory': identity['data_directory'], 'in_recovery': False,
'postmaster_start_time': datetime.fromtimestamp(1700000000, timezone.utc)},
{'present': True, 'permitted': True}, {'system_identifier': '123456789'},
]
def stopped(*_args, **_kwargs):
process.running = False
return SimpleNamespace(returncode=0, stdout='')
with mock.patch.object(backend, '_root_job_error', return_value=None), \
mock.patch.object(postgres_runtime, 'verify_cluster_identity', return_value=identity), \
mock.patch.object(postgres_runtime, 'canonical_postgres_url', return_value='fixture-dsn'), \
mock.patch.object(postgres_runtime, 'database_url_from_env', return_value='fixture-dsn'), \
mock.patch.object(postgres_runtime, 'connect_postgres', side_effect=[
connection, RuntimeError('synthetic schema unavailable')]), \
mock.patch.object(backend, '_open_postmaster', return_value=process), \
mock.patch.object(postgres_runtime, '_parse_postmaster_pid', return_value=process.pid), \
mock.patch.object(postgres_runtime, '_run_bounded', side_effect=stopped) as run:
self.assertEqual(backend.probe().kind, ProbeKind.READY)
self.assertFalse(backend.owns_start, 'authentication is not launch ownership')
self.assertEqual(backend._authenticated_postmaster, (identity, process))
result = backend.stop()
self.assertTrue(result.completed and result.stopped)
self.assertTrue(process.closed)
self.assertIsNone(backend._authenticated_postmaster)
run.assert_called_once()
def test_missing_system_identity_authentication_does_not_mint_stop_proof(self):
for availability in ({'present': False, 'permitted': False},
{'present': True, 'permitted': False}):
with self.subTest(availability=availability):
backend = PostgresBackend.__new__(PostgresBackend)
backend.connect_timeout_sec = 1
backend.query_timeout_ms = 1000
backend._expected_process = None
backend._authenticated_postmaster = None
process = FakeRetainedPostmaster()
identity = {
'data_directory': postgres_runtime.canonical_path('fixture-data'),
'database': 'truf', 'user': 'truf', 'port': 55432, 'pg_major': 16,
'system_identifier': '123456789',
}
connection = mock.Mock()
connection.execute.return_value.fetchone.side_effect = [
{'database': 'truf', 'user_name': 'truf', 'port': 55432,
'data_directory': identity['data_directory'], 'in_recovery': False,
'postmaster_start_time': datetime.fromtimestamp(1700000000, timezone.utc)},
availability,
]
with mock.patch.object(postgres_runtime, 'database_url_from_env', return_value='fixture-dsn'), \
mock.patch.object(postgres_runtime, 'canonical_postgres_url', return_value='fixture-dsn'), \
mock.patch.object(postgres_runtime, 'connect_postgres', return_value=connection), \
mock.patch.object(backend, '_open_postmaster', return_value=process):
self.assertFalse(backend._online_query(identity)[0])
self.assertIsNone(backend._authenticated_postmaster)
backend._verified_identity = identity
backend._accepted_start_at_monotonic = None
backend._started_postmaster_observed = False
with mock.patch.object(backend, 'probe', return_value=ProbeResult(ProbeKind.RECOVERING)), \
mock.patch.object(backend, '_stop_retained_postmaster') as stop:
result = backend.stop()
self.assertFalse(result.completed or result.stopped)
stop.assert_not_called()
def test_recovery_stop_requires_the_same_authenticated_binding_and_handle(self):
for proof in ('none', 'changed-binding', 'different-process'):
with self.subTest(proof=proof):
backend = PostgresBackend.__new__(PostgresBackend)
process = FakeRetainedPostmaster()
identity = {'system_identifier': '123', 'data_directory': 'fixture-data'}
backend._expected_process = process
backend._verified_identity = identity
backend._accepted_start_at_monotonic = None
backend._started_postmaster_observed = False
backend._authenticated_postmaster = {
'none': None,
'changed-binding': (dict(identity, system_identifier='456'), process),
'different-process': (identity, FakeRetainedPostmaster()),
}[proof]
with mock.patch.object(backend, 'probe', return_value=ProbeResult(ProbeKind.RECOVERING)), \
mock.patch.object(postgres_runtime, '_run_bounded') as run:
result = backend.stop()
self.assertFalse(result.completed or result.stopped)
self.assertFalse(process.closed)
run.assert_not_called()
def test_foreign_probe_cannot_stop_a_retained_accepted_start(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend._expected_process = FakeRetainedPostmaster()
backend._accepted_start_at_monotonic = postgres_runtime.time.monotonic()
backend._verified_identity = {'data_directory': 'fixture-data'}
backend._authenticated_postmaster = (backend._verified_identity, backend._expected_process)
with mock.patch.object(backend, 'probe', return_value=ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR)), \
mock.patch.object(backend, '_stop_retained_postmaster') as stop:
result = backend.stop()
self.assertFalse(result.completed or result.stopped)
self.assertFalse(backend._expected_process.closed)
stop.assert_not_called()
def test_interrupted_close_does_not_poison_a_freshly_verified_handle(self):
backend = PostgresBackend.__new__(PostgresBackend)
old = FakeRetainedPostmaster()
replacement = FakeRetainedPostmaster()
backend._expected_process = old
backend._authenticated_postmaster = ({'system_identifier': '123'}, old)
def interrupted_close():
old.closed = True
raise KeyboardInterrupt()
with mock.patch.object(old, 'close', side_effect=interrupted_close):
with self.assertRaises(KeyboardInterrupt):
backend.close()
self.assertIs(backend._expected_process, old)
self.assertIs(backend._remember_process(replacement), replacement)
self.assertFalse(replacement.closed)
self.assertIsNot(backend._authenticated_postmaster[1], replacement)
def test_interrupted_or_failed_launcher_retains_the_start_promise_for_compensation(self):
outcomes = (
KeyboardInterrupt(), SystemExit(0), OSError('launcher outcome unavailable'),
SimpleNamespace(returncode=1, stdout='launcher failed after possible fork'),
)
for outcome in outcomes:
with self.subTest(outcome=type(outcome).__name__):
backend = PostgresBackend.__new__(PostgresBackend)
backend.paths = {'log_path': 'fixture-postgres.log'}
backend._expected_process = None
backend._start_requested_wall_time = None
backend._accepted_start_at_monotonic = None
backend._started_postmaster_observed = False
backend._verified_identity = {
'data_directory': 'fixture-data', 'port': 55432,
'executables': {'pg_ctl': 'fixture-pg_ctl'},
}
self.assertFalse(backend.owns_start)
with mock.patch.object(backend, 'probe', return_value=ProbeResult(ProbeKind.STOPPED)), \
mock.patch.object(postgres_runtime, 'IS_WINDOWS', False), \
mock.patch.object(postgres_runtime, '_run_bounded') as run:
if isinstance(outcome, BaseException):
run.side_effect = outcome
else:
run.return_value = outcome
if isinstance(outcome, (KeyboardInterrupt, SystemExit)):
with self.assertRaises(type(outcome)):
backend.start()
else:
result = backend.start()
self.assertTrue(result.accepted)
self.assertTrue(result.uncertain)
self.assertTrue(backend.owns_start)
self.assertIsNotNone(backend._accepted_start_at_monotonic)
self.assertIsNotNone(backend._start_requested_wall_time)
def test_foreign_identity_stop_never_signals_or_closes_retained_process(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend._expected_process = FakeRetainedPostmaster()
backend._accepted_start_at_monotonic = None
with mock.patch.object(backend, 'probe', return_value=ProbeResult(ProbeKind.FOREIGN_OR_CONFIG_ERROR, 'foreign')), \
mock.patch.object(postgres_runtime, '_run_bounded') as run, \
mock.patch.object(postgres_runtime, 'open_process') as lookup:
result = backend.stop()
self.assertFalse(result.completed)
self.assertFalse(result.stopped)
self.assertIn('foreign', result.detail)
self.assertFalse(backend._expected_process.closed)
run.assert_not_called()
lookup.assert_not_called()
def test_pg_ctl_start_timeout_retains_uncertain_start_ownership(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.paths = {'log_path': r'D:\logs\postgres.log'}
backend._start_requested_wall_time = None
backend._accepted_start_at_monotonic = None
backend._started_postmaster_observed = False
backend._verified_identity = {
'data_directory': r'D:\bound-data',
'executables': {'pg_ctl': r'D:\bin\pg_ctl.exe'},
'port': 5544,
}
with mock.patch.object(backend, 'probe', return_value=ProbeResult(ProbeKind.STOPPED, 'offline')), \
mock.patch.object(postgres_runtime, 'IS_WINDOWS', False), \
mock.patch('postgres_runtime._run_bounded', side_effect=subprocess.TimeoutExpired('pg_ctl', 30)):
result = backend.start()
self.assertTrue(result.accepted)
self.assertTrue(result.uncertain)
self.assertIsNotNone(backend._accepted_start_at_monotonic)
self.assertIsNotNone(backend._start_requested_wall_time)
def test_shutdown_waits_for_late_postmaster_after_start_timeout_and_identity_stops_it(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.stop_timeout_sec = 5
backend.start_settle_timeout_sec = 0.1
backend._accepted_start_at_monotonic = postgres_runtime.time.monotonic()
backend._start_requested_wall_time = postgres_runtime.time.time()
backend._started_postmaster_observed = False
backend._expected_process = None
process = StoppableRetainedPostmaster()
identity = {
'data_directory': r'D:\bound-data',
'executables': {'pg_ctl': r'D:\bin\pg_ctl.exe'},
}
backend._verified_identity = identity
probes = {'count': 0}
def late_probe():
probes['count'] += 1
if probes['count'] >= 2:
backend._expected_process = process
return ProbeResult(ProbeKind.RECOVERING, 'uncertain start')
def stopped_after_command(*_args, **_kwargs):
process.running = False
return SimpleNamespace(returncode=0, stdout='')
with mock.patch.object(backend, 'probe', side_effect=late_probe), \
mock.patch('postgres_runtime._parse_postmaster_pid', return_value=process.pid), \
mock.patch('postgres_runtime._run_bounded', side_effect=stopped_after_command):
result = backend.stop()
self.assertTrue(result.completed)
self.assertTrue(result.stopped)
self.assertGreaterEqual(probes['count'], 2)
def test_stop_reprobes_recent_accepted_start_then_identity_stops_late_postmaster(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.stop_timeout_sec = 5
backend.start_settle_timeout_sec = 1
backend._accepted_start_at_monotonic = postgres_runtime.time.monotonic()
backend._start_requested_wall_time = postgres_runtime.time.time()
process = StoppableRetainedPostmaster()
backend._expected_process = process
backend._verified_identity = {
'data_directory': r'D:\bound-data',
'executables': {'pg_ctl': r'D:\bin\pg_ctl.exe'},
}
def stopped_after_command(*_args, **_kwargs):
process.running = False
return SimpleNamespace(returncode=0, stdout='')
with mock.patch.object(backend, 'probe', side_effect=[
ProbeResult(ProbeKind.STOPPED, 'not published yet'),
ProbeResult(ProbeKind.READY, 'late postmaster ready'),
]), mock.patch('postgres_runtime._parse_postmaster_pid', return_value=process.pid), \
mock.patch('postgres_runtime._run_bounded', side_effect=stopped_after_command):
result = backend.stop()
self.assertTrue(result.completed)
self.assertTrue(result.stopped)
self.assertTrue(process.closed)
def test_stop_reports_uncertain_when_recent_accepted_start_never_publishes(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.start_settle_timeout_sec = 0.001
backend._accepted_start_at_monotonic = postgres_runtime.time.monotonic()
backend._start_requested_wall_time = postgres_runtime.time.time()
backend._expected_process = None
backend._verified_identity = None
with mock.patch.object(backend, 'probe', return_value=ProbeResult(ProbeKind.STOPPED, 'offline')):
result = backend.stop()
self.assertFalse(result.completed)
self.assertFalse(result.stopped)
self.assertIn('uncertain', result.detail)
def test_restart_during_local_recovery_is_adopted_as_recovering(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.config = {}
backend.paths = {}
backend._expected_process = None
backend._start_requested_wall_time = None
backend._verified_identity = None
identity = {
'data_directory': r'D:\bound-data',
'executables': {'postgres': r'D:\bin\postgres.exe'},
'port': 5432,
}
process = FakeRetainedPostmaster()
with mock.patch.object(backend, '_root_job_error', return_value=None), \
mock.patch('postgres_runtime.verify_cluster_identity', return_value=identity), \
mock.patch.object(backend, '_online_query', side_effect=OnlineUnavailable('not query-ready')), \
mock.patch('postgres_runtime._parse_postmaster_pid', return_value=process.pid), \
mock.patch.object(backend, '_open_postmaster', return_value=process), \
mock.patch('postgres_runtime._listener_present', return_value=True):
result = backend.probe()
self.assertEqual(result.kind, ProbeKind.RECOVERING)
self.assertIs(backend._expected_process, process)
self.assertIn('awaiting authenticated identity', result.detail)
def test_postmaster_adoption_requires_pid_file_data_directory_and_start_time(self):
with tempfile.TemporaryDirectory() as temp_dir:
data_dir = os.path.join(temp_dir, 'data')
os.makedirs(data_dir)
start_time = 1700000000
Path(os.path.join(data_dir, 'postmaster.pid')).write_text(
f'321\n{data_dir}\n{start_time}\n',
encoding='ascii',
)
backend = PostgresBackend.__new__(PostgresBackend)
backend._start_requested_wall_time = None
process = mock.Mock()
process.identity = SimpleNamespace(
executable=postgres_runtime.canonical_path(os.path.join(temp_dir, 'postgres.exe')),
in_job=False,
creation_time_unix=float(start_time),
)
identity = {
'data_directory': postgres_runtime.canonical_path(data_dir),
'executables': {'postgres': process.identity.executable},
}
with mock.patch('postgres_runtime.open_process', return_value=process):
self.assertIs(backend._open_postmaster(identity), process)
Path(os.path.join(data_dir, 'postmaster.pid')).write_text(
f'321\n{os.path.join(temp_dir, "other")}\n{start_time}\n',
encoding='ascii',
)
with self.assertRaisesRegex(ClusterIdentityError, 'data directory'):
backend._open_postmaster(identity)
def test_linux_postmaster_adoption_ignores_wall_clock_skew(self):
with tempfile.TemporaryDirectory() as temp_dir:
data_dir = os.path.join(temp_dir, 'data')
os.makedirs(data_dir)
start_time = 1700000000
Path(os.path.join(data_dir, 'postmaster.pid')).write_text(
f'321\n{data_dir}\n{start_time}\n',
encoding='ascii',
)
backend = PostgresBackend.__new__(PostgresBackend)
backend._start_requested_wall_time = None
process = mock.Mock()
process.identity = SimpleNamespace(
executable=postgres_runtime.canonical_path(os.path.join(temp_dir, 'postgres')),
in_job=False,
creation_time_unix=float(start_time) - 22.0,
)
identity = {
'data_directory': postgres_runtime.canonical_path(data_dir),
'executables': {'postgres': process.identity.executable},
}
with mock.patch.object(postgres_runtime, 'IS_WINDOWS', False), \
mock.patch('postgres_runtime.open_process', return_value=process):
self.assertIs(backend._open_postmaster(identity), process)
def test_windows_postmaster_adoption_rejects_wall_clock_skew(self):
with tempfile.TemporaryDirectory() as temp_dir:
data_dir = os.path.join(temp_dir, 'data')
os.makedirs(data_dir)
start_time = 1700000000
Path(os.path.join(data_dir, 'postmaster.pid')).write_text(
f'321\n{data_dir}\n{start_time}\n',
encoding='ascii',
)
backend = PostgresBackend.__new__(PostgresBackend)
backend._start_requested_wall_time = None
process = mock.Mock()
process.identity = SimpleNamespace(
executable=postgres_runtime.canonical_path(os.path.join(temp_dir, 'postgres.exe')),
in_job=False,
creation_time_unix=float(start_time) - 22.0,
)
identity = {
'data_directory': postgres_runtime.canonical_path(data_dir),
'executables': {'postgres': process.identity.executable},
}
with mock.patch.object(postgres_runtime, 'IS_WINDOWS', True), \
mock.patch('postgres_runtime.open_process', return_value=process):
with self.assertRaisesRegex(ClusterIdentityError, 'start time does not match'):
backend._open_postmaster(identity)
def test_linux_online_query_ignores_postmaster_wall_clock_skew(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.connect_timeout_sec = 1
backend.query_timeout_ms = 1000
backend._expected_process = None
backend._authenticated_postmaster = None
process = FakeRetainedPostmaster()
process.identity.creation_time_unix = 1700000000.0 - 22.0
identity = {
'data_directory': postgres_runtime.canonical_path('fixture-data'),
'database': 'truf', 'user': 'truf', 'port': 55432, 'pg_major': 16,
'system_identifier': '123456789',
}
connection = mock.Mock()
connection.execute.return_value.fetchone.side_effect = [
{'database': 'truf', 'user_name': 'truf', 'port': 55432,
'data_directory': identity['data_directory'], 'in_recovery': False,
'postmaster_start_time': datetime.fromtimestamp(1700000000, timezone.utc)},
{'present': True, 'permitted': True}, {'system_identifier': '123456789'},
]
with mock.patch.object(postgres_runtime, 'IS_WINDOWS', False), \
mock.patch.object(postgres_runtime, 'database_url_from_env', return_value='fixture-dsn'), \
mock.patch.object(postgres_runtime, 'canonical_postgres_url', return_value='fixture-dsn'), \
mock.patch.object(postgres_runtime, 'connect_postgres', return_value=connection), \
mock.patch.object(backend, '_open_postmaster', return_value=process):
in_recovery, _epoch = backend._online_query(identity)
self.assertFalse(in_recovery)
self.assertEqual(backend._authenticated_postmaster, (identity, process))
def test_windows_online_query_rejects_postmaster_wall_clock_skew(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.connect_timeout_sec = 1
backend.query_timeout_ms = 1000
backend._expected_process = None
backend._authenticated_postmaster = None
process = FakeRetainedPostmaster()
process.identity.creation_time_unix = 1700000000.0 - 22.0
identity = {
'data_directory': postgres_runtime.canonical_path('fixture-data'),
'database': 'truf', 'user': 'truf', 'port': 55432, 'pg_major': 16,
'system_identifier': '123456789',
}
connection = mock.Mock()
connection.execute.return_value.fetchone.side_effect = [
{'database': 'truf', 'user_name': 'truf', 'port': 55432,
'data_directory': identity['data_directory'], 'in_recovery': False,
'postmaster_start_time': datetime.fromtimestamp(1700000000, timezone.utc)},
{'present': True, 'permitted': True}, {'system_identifier': '123456789'},
]
with mock.patch.object(postgres_runtime, 'IS_WINDOWS', True), \
mock.patch.object(postgres_runtime, 'database_url_from_env', return_value='fixture-dsn'), \
mock.patch.object(postgres_runtime, 'canonical_postgres_url', return_value='fixture-dsn'), \
mock.patch.object(postgres_runtime, 'connect_postgres', return_value=connection), \
mock.patch.object(backend, '_open_postmaster', return_value=process):
with self.assertRaisesRegex(ClusterIdentityError, 'creation time mismatch'):
backend._online_query(identity)
self.assertIsNone(backend._authenticated_postmaster)
def test_pg_ctl_start_forces_bound_data_directory(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.paths = {'log_path': r'D:\logs\postgres.log'}
backend._start_requested_wall_time = None
identity = {
'data_directory': r'D:\bound data',
'executables': {'pg_ctl': r'D:\bin\pg_ctl.exe'},
'port': 5544,
}
backend._verified_identity = identity
with mock.patch.object(backend, 'probe', return_value=ProbeResult(ProbeKind.STOPPED, 'offline')), \
mock.patch.object(postgres_runtime, 'IS_WINDOWS', False), \
mock.patch('postgres_runtime._run_bounded', return_value=SimpleNamespace(returncode=0, stdout='')) as run:
result = backend.start()
self.assertTrue(result.accepted)
self.assertFalse(run.call_args.kwargs['capture_output'])
command = run.call_args.args[0]
options = command[command.index('-o') + 1]
self.assertIn('data_directory=', options)
self.assertIn(r'D:\bound data', options)
def test_pg_ctl_start_forces_collector_rotation_options(self):
with tempfile.TemporaryDirectory() as temp_dir:
postgres_dir = os.path.join(temp_dir, 'postgres')
log_dir = os.path.join(postgres_dir, 'logs')
os.makedirs(postgres_dir)
ensure_private_directory(temp_dir, reject_reparse=True)
ensure_private_directory(postgres_dir, reject_reparse=True)
backend = PostgresBackend.__new__(PostgresBackend)
backend.paths = {
'log_path': os.path.join(postgres_dir, 'postgres.log'),
'log_dir': log_dir,
}
backend.log_max_mb = 7
backend.log_keep = 3
backend._start_requested_wall_time = None
backend._accepted_start_at_monotonic = None
backend._started_postmaster_observed = False
identity = {
'data_directory': os.path.join(temp_dir, 'data'),
'executables': {'pg_ctl': os.path.join(temp_dir, 'pg_ctl.exe')},
'port': 5544,
}
backend._verified_identity = identity
with mock.patch.object(backend, 'probe', return_value=ProbeResult(ProbeKind.STOPPED, 'offline')), \
mock.patch.object(postgres_runtime, 'IS_WINDOWS', False), \
mock.patch('postgres_runtime._run_bounded', return_value=SimpleNamespace(returncode=0, stdout='')) as run:
result = backend.start()
self.assertTrue(result.accepted)
options = run.call_args.args[0][run.call_args.args[0].index('-o') + 1]
self.assertIn('logging_collector=on', options)
self.assertIn('log_rotation_size=7MB', options)
self.assertIn('log_truncate_on_rotation=on', options)
def test_postgres_collector_retention_prunes_old_private_segments(self):
with tempfile.TemporaryDirectory() as temp_dir:
ensure_private_directory(temp_dir, reject_reparse=True)
paths = []
for index in range(5):
path = os.path.join(temp_dir, f'postgresql-20260719-00000{index}.log')
Path(path).write_text('x', encoding='ascii')
postgres_runtime.harden_private_file(path)
os.utime(path, (index + 1, index + 1))
paths.append(path)
removed = postgres_runtime.prune_postgres_collector_logs(temp_dir, 2)
self.assertEqual(removed, 3)
self.assertEqual(sum(os.path.exists(path) for path in paths), 2)
def test_pg_ctl_success_is_not_stopped_while_retained_process_is_live(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.stop_timeout_sec = 5
backend._verified_identity = {
'data_directory': r'D:\bound-data',
'executables': {'pg_ctl': r'D:\bin\pg_ctl.exe'},
}
backend._expected_process = FakeRetainedPostmaster(pid=321)
with mock.patch.object(backend, 'probe', return_value=ProbeResult(ProbeKind.READY, 'ready')), \
mock.patch('postgres_runtime._parse_postmaster_pid', return_value=321), \
mock.patch('postgres_runtime._run_bounded', return_value=SimpleNamespace(returncode=0, stdout='')):
result = backend.stop()
self.assertFalse(result.stopped)
self.assertIn('still running', result.detail)
def test_pg_ctl_timeout_is_reported_as_not_stopped(self):
backend = PostgresBackend.__new__(PostgresBackend)
backend.stop_timeout_sec = 5
backend._verified_identity = {
'data_directory': r'D:\bound-data',
'executables': {'pg_ctl': r'D:\bin\pg_ctl.exe'},
}
backend._expected_process = FakeRetainedPostmaster(pid=321)
with mock.patch.object(backend, 'probe', return_value=ProbeResult(ProbeKind.READY, 'ready')), \
mock.patch('postgres_runtime._parse_postmaster_pid', return_value=321), \
mock.patch('postgres_runtime._run_bounded', side_effect=TimeoutError('timed out')):
result = backend.stop()
self.assertFalse(result.completed)
self.assertFalse(result.stopped)
self.assertIn('timed out', result.detail)
class PostgresEntrypointTests(unittest.TestCase):
def test_preflight_failure_precedes_environment_and_cluster_tools(self):
config = {'global': {'runtime_dir': r'D:\fixture\runtime'}}
with mock.patch.object(sys, 'argv', ['postgres_runtime.py', 'verify', '--config', 'config.yaml']), \
mock.patch.object(postgres_runtime, '_load_config', return_value=config), \
mock.patch.object(postgres_runtime, 'preflight_lifecycle_paths', side_effect=OSError('unsafe config')) as preflight, \
mock.patch.object(postgres_runtime, 'load_postgres_environment') as load_environment, \
mock.patch.object(postgres_runtime, 'verify_cluster_identity') as verify:
with self.assertRaisesRegex(SystemExit, 'unsafe config'):
postgres_runtime.main()
load_environment.assert_not_called()
verify.assert_not_called()
preflight.assert_called_once_with(
os.path.abspath('config.yaml'), config, authority_profile='server',
)
if __name__ == '__main__':
unittest.main()