diff --git a/patroni/__main__.py b/patroni/__main__.py index e9dcd28a..49542f0c 100644 --- a/patroni/__main__.py +++ b/patroni/__main__.py @@ -1,3 +1,8 @@ +"""Patroni main entry point. + +Implement ``patroni`` main daemon and expose its entry point. +""" + import logging import os import signal @@ -16,8 +21,33 @@ logger = logging.getLogger(__name__) class Patroni(AbstractPatroniDaemon): + """Implement ``patroni`` command daemon. + + :ivar version: Patroni version. + :ivar dcs: DCS object. + :ivar watchdog: watchdog handler, if configured to use watchdog. + :ivar postgresql: managed Postgres instance. + :ivar api: REST API server instance of this node. + :ivar request: wrapper for performing HTTP requests. + :ivar ha: HA handler. + :ivar tags: cache of custom tags configured for this node. + :ivar next_run: time when to run the next HA loop cycle. + :ivar scheduled_restart: when a restart has been scheduled to occur, if any. In that case, should contain two keys: + * ``schedule``: timestamp when restart should occur; + * ``postmaster_start_time``: timestamp when Postgres was last started. + """ def __init__(self, config: 'Config') -> None: + """Create a :class:`Patroni` instance with the given *config*. + + Get a connection to the DCS, configure watchdog (if required), set up Patroni interface with Postgres, configure + the HA loop and bring the REST API up. + + .. note:: + Expected to be instantiated and run through :func:`~patroni.daemon.abstract_main`. + + :param config: Patroni configuration. + """ from patroni.api import RestApiServer from patroni.dcs import get_dcs from patroni.ha import Ha @@ -46,6 +76,17 @@ class Patroni(AbstractPatroniDaemon): self.scheduled_restart: Dict[str, Any] = {} def load_dynamic_configuration(self) -> None: + """Load Patroni dynamic configuration. + + Load dynamic configuration from the DCS, if `/config` key is available in the DCS, otherwise fall back to + ``bootstrap.dcs`` section from the configuration file. + + If the DCS connection fails returning the exception :class:`~patroni.exceptions.DCSError` an attempt will be + remade every 5 seconds. + + .. note:: + This method is called only once, at the time when Patroni is started. + """ from patroni.exceptions import DCSError while True: try: @@ -81,18 +122,47 @@ class Patroni(AbstractPatroniDaemon): return def get_tags(self) -> Dict[str, Any]: + """Get tags configured for this node, if any. + + Handle both predefined Patroni tags and custom defined tags. + + .. note:: + A custom tag is any tag added to the configuration ``tags`` section that is not one of ``clonefrom``, + ``nofailover``, ``noloadbalance`` or ``nosync``. + + For the Patroni predefined tags, the returning object will only contain them if they are enabled as they + all are boolean values that default to disabled. + + :returns: a dictionary of tags set for this node. The key is the tag name, and the value is the corresponding + tag value. + """ return {tag: value for tag, value in self.config.get('tags', {}).items() if tag not in ('clonefrom', 'nofailover', 'noloadbalance', 'nosync') or value} @property def nofailover(self) -> bool: + """``True`` if ``tags.nofailover`` configuration is enabled for this node, else ``False``.""" return bool(self.tags.get('nofailover', False)) @property def nosync(self) -> bool: + """``True`` if ``tags.nosync`` configuration is enabled for this node, else ``False``.""" return bool(self.tags.get('nosync', False)) def reload_config(self, sighup: bool = False, local: Optional[bool] = False) -> None: + """Apply new configuration values for ``patroni`` daemon. + + Reload: + * Cached tags; + * Request wrapper configuration; + * REST API configuration; + * Watchdog configuration; + * Postgres configuration; + * DCS configuration. + + :param sighup: if it is related to a SIGHUP signal. + :param local: if there has been changes to the local configuration file. + """ try: super(Patroni, self).reload_config(sighup, local) if local: @@ -107,14 +177,21 @@ class Patroni(AbstractPatroniDaemon): logger.exception('Failed to reload config_file=%s', self.config.config_file) @property - def replicatefrom(self): + def replicatefrom(self) -> Optional[str]: + """Value of ``tags.replicatefrom`` configuration, if any.""" return self.tags.get('replicatefrom') @property - def noloadbalance(self): + def noloadbalance(self) -> bool: + """``True`` if ``tags.noloadbalance`` configuration is enabled for this node, else ``False``.""" return bool(self.tags.get('noloadbalance', False)) def schedule_next_run(self) -> None: + """Schedule the next run of the ``patroni`` daemon main loop. + + Next run is scheduled based on previous run plus value of ``loop_wait`` configuration from DCS. If that has + already been exceeded, run the next cycle immediately. + """ self.next_run += self.dcs.loop_wait current_time = time.time() nap_time = self.next_run - current_time @@ -128,11 +205,21 @@ class Patroni(AbstractPatroniDaemon): self.next_run = time.time() def run(self) -> None: + """Run ``patroni`` daemon process main loop. + + Start the REST API and keep running HA cycles every ``loop_wait`` seconds. + """ self.api.start() self.next_run = time.time() super(Patroni, self).run() def _run_cycle(self) -> None: + """Run a cycle of the ``patroni`` daemon main loop. + + Run an HA cycle and schedule the next cycle run. If any dynamic configuration change request is detected, apply + the change and cache the new dynamic configuration values in ``patroni.dynamic.json`` file under Postgres data + directory. + """ logger.info(self.ha.run_cycle()) if self.dcs.cluster and self.dcs.cluster.config and self.dcs.cluster.config.data \ @@ -145,6 +232,10 @@ class Patroni(AbstractPatroniDaemon): self.schedule_next_run() def _shutdown(self) -> None: + """Perform shutdown of ``patroni`` daemon process. + + Shut down the REST API and the HA handler. + """ try: self.api.shutdown() except Exception: @@ -156,13 +247,34 @@ class Patroni(AbstractPatroniDaemon): def patroni_main(configfile: str) -> None: + """Configure and start ``patroni`` main daemon process. + + :param configfile: path to Patroni configuration file. + """ from multiprocessing import freeze_support + # Windows executables created by PyInstaller are frozen, thus we need to enable frozen support for + # :mod:`multiprocessing` to avoid :class:`RuntimeError` exceptions. freeze_support() abstract_main(Patroni, configfile) def process_arguments() -> Namespace: + """Process command-line arguments. + + Create a basic command-line parser through :func:`~patroni.daemon.get_base_arg_parser`, extend its capabilities by + adding these flags and parse command-line arguments.: + + * ``--validate-config`` -- used to validate the Patroni configuration file + * ``--generate-config`` -- used to generate Patroni configuration from a running PostgreSQL instance + * ``--generate-sample-config`` -- used to generate a sample Patroni configuration + + .. note:: + If running with ``--generate-config``, ``--generate-sample-config`` or ``--validate-flag`` will exit + after generating or validating configuration. + + :returns: parsed arguments, if not running with ``--validate-config`` flag. + """ from patroni.config_generator import generate_config parser = get_base_arg_parser() @@ -196,6 +308,16 @@ def process_arguments() -> Namespace: def main() -> None: + """Main entrypoint of :mod:`patroni.__main__`. + + Process command-line arguments, ensure :mod:`psycopg2` (or :mod:`psycopg`) attendee the pre-requisites and start + ``patroni`` daemon process. + + .. note:: + If running through a Docker container, make the main process take care of init process duties and run + ``patroni`` daemon as another process. In that case relevant signals received by the main process and forwarded + to ``patroni`` daemon process. + """ from patroni import check_psycopg args = process_arguments() @@ -211,7 +333,13 @@ def main() -> None: # Looks like we are in a docker, so we will act like init def sigchld_handler(signo: int, stack_frame: Optional[FrameType]) -> None: + """Handle ``SIGCHLD`` received by main process from ``patroni`` daemon when the daemon terminates. + + :param signo: signal number. + :param stack_frame: current stack frame. + """ try: + # log exit code of all children processes, and break loop when there is none left while True: ret = os.waitpid(-1, os.WNOHANG) if ret == (0, 0): @@ -221,7 +349,12 @@ def main() -> None: except OSError: pass - def passtochild(signo: int, stack_frame: Optional[FrameType]): + def passtochild(signo: int, stack_frame: Optional[FrameType]) -> None: + """Forward a signal *signo* from main process to child process. + + :param signo: signal number. + :param stack_frame: current stack frame. + """ if pid: os.kill(pid, signo) diff --git a/patroni/api.py b/patroni/api.py index 52d0d21a..ce0c1d37 100644 --- a/patroni/api.py +++ b/patroni/api.py @@ -460,7 +460,7 @@ class RestApiHandler(BaseHTTPRequestHandler): * ``patroni_postmaster_start_time``: epoch timestamp since Postmaster was started; * ``patroni_master``: ``1`` if this node holds the leader lock, else ``0``; * ``patroni_primary``: same as ``patroni_master``; - * ``patroni_xlog_location``: ``pg_wal_lsn_diff(pg_current_wal_lsn(), '0/0')`` if leader, else ``0``; + * ``patroni_xlog_location``: ``pg_wal_lsn_diff(pg_current_wal_flush_lsn(), '0/0')`` if leader, else ``0``; * ``patroni_standby_leader``: ``1`` if standby leader node, else ``0``; * ``patroni_replica``: ``1`` if a replica, else ``0``; * ``patroni_sync_standby``: ``1`` if a sync replica, else ``0``; @@ -1159,7 +1159,7 @@ class RestApiHandler(BaseHTTPRequestHandler): * ``server_version``: Postgres version without periods, e.g. ``150002`` for Postgres ``15.2``; * ``xlog``: dictionary. Its structure depends on ``role``: * If ``master``: - * ``location``: ``pg_current_wal_lsn()`` + * ``location``: ``pg_current_wal_flush_lsn()`` * If ``replica``: * ``received_location``: ``pg_wal_lsn_diff(pg_last_wal_receive_lsn(), '0/0')``; * ``replayed_location``: ``pg_wal_lsn_diff(pg_last_wal_replay_lsn(), '0/0)``; @@ -1197,8 +1197,8 @@ class RestApiHandler(BaseHTTPRequestHandler): " application_name, client_addr, w.state, sync_state, sync_priority" " FROM pg_catalog.pg_stat_get_wal_senders() w, pg_catalog.pg_stat_get_activity(pid)) AS ri") - row = self.query(stmt.format(postgresql.wal_name, postgresql.lsn_name), retry=retry)[0] - + row = self.query(stmt.format(postgresql.wal_name, postgresql.lsn_name, + postgresql.wal_flush), retry=retry)[0] result = { 'state': postgresql.state, 'postmaster_start_time': row[0], diff --git a/patroni/postgresql/__init__.py b/patroni/postgresql/__init__.py index dfec86e1..376ed071 100644 --- a/patroni/postgresql/__init__.py +++ b/patroni/postgresql/__init__.py @@ -57,8 +57,8 @@ class Postgresql(object): TL_LSN = ("CASE WHEN pg_catalog.pg_is_in_recovery() THEN 0 " "ELSE ('x' || pg_catalog.substr(pg_catalog.pg_{0}file_name(" "pg_catalog.pg_current_{0}_{1}()), 1, 8))::bit(32)::int END, " # primary timeline - "CASE WHEN pg_catalog.pg_is_in_recovery() THEN 0 " - "ELSE pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_current_{0}_{1}(), '0/0')::bigint END, " # write_lsn + "CASE WHEN pg_catalog.pg_is_in_recovery() THEN 0 ELSE " + "pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_current_{0}{2}_{1}(), '0/0')::bigint END, " # wal(_flush)?_lsn "pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_last_{0}_replay_{1}(), '0/0')::bigint, " "pg_catalog.pg_{0}_{1}_diff(COALESCE(pg_catalog.pg_last_{0}_receive_{1}(), '0/0'), '0/0')::bigint, " "pg_catalog.pg_is_in_recovery() AND pg_catalog.pg_is_{0}_replay_paused()") @@ -159,6 +159,11 @@ class Postgresql(object): def wal_name(self) -> str: return 'wal' if self._major_version >= 100000 else 'xlog' + @property + def wal_flush(self) -> str: + """For PostgreSQL 9.6 onwards we want to use pg_current_wal_flush_lsn()/pg_current_xlog_flush_location().""" + return '_flush' if self._major_version >= 90600 else '' + @property def lsn_name(self) -> str: return 'lsn' if self._major_version >= 100000 else 'location' @@ -216,7 +221,7 @@ class Postgresql(object): else: extra = "0, NULL, NULL, NULL, NULL, NULL, NULL" + extra - return ("SELECT " + self.TL_LSN + ", {2}").format(self.wal_name, self.lsn_name, extra) + return ("SELECT " + self.TL_LSN + ", {3}").format(self.wal_name, self.lsn_name, self.wal_flush, extra) @property def available_gucs(self) -> CaseInsensitiveSet: diff --git a/patroni/postgresql/sync.py b/patroni/postgresql/sync.py index 815afe56..e5814221 100644 --- a/patroni/postgresql/sync.py +++ b/patroni/postgresql/sync.py @@ -202,7 +202,7 @@ class _ReplicaList(List[_Replica]): swapping, but only if lag on this member is exceeding a threshold (``maximum_lag_on_syncnode``). :ivar max_lsn: maximum value of ``_Replica.lsn`` among all values. In case if there is just one - element in the list we take value of ``pg_current_wal_lsn()``. + element in the list we take value of ``pg_current_wal_flush_lsn()``. """ def __init__(self, postgresql: 'Postgresql', cluster: Cluster) -> None: diff --git a/tests/test_api.py b/tests/test_api.py index 44b5c7d9..c1db7618 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -36,6 +36,7 @@ class MockPostgresql(object): pending_restart = True wal_name = 'wal' lsn_name = 'lsn' + wal_flush = '_flush' POSTMASTER_START_TIME = 'pg_catalog.pg_postmaster_start_time()' TL_LSN = 'CASE WHEN pg_catalog.pg_is_in_recovery()' citus_handler = Mock() diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index a703f235..4f1e36da 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -955,3 +955,10 @@ class TestPostgresql2(BaseTestPostgresql): gucs = self.p.available_gucs self.assertIsInstance(gucs, CaseInsensitiveSet) self.assertEqual(gucs, mock_available_gucs.return_value) + + def test_cluster_info_query(self): + self.assertIn('diff(pg_catalog.pg_current_wal_flush_lsn(', self.p.cluster_info_query) + self.p._major_version = 90600 + self.assertIn('diff(pg_catalog.pg_current_xlog_flush_location(', self.p.cluster_info_query) + self.p._major_version = 90500 + self.assertIn('diff(pg_catalog.pg_current_xlog_location(', self.p.cluster_info_query)