diff --git a/patroni/dcs/__init__.py b/patroni/dcs/__init__.py index 193b7fe1..473cbe4a 100644 --- a/patroni/dcs/__init__.py +++ b/patroni/dcs/__init__.py @@ -758,10 +758,19 @@ class AbstractDCS(object): otherwise it should return `!False`""" @abc.abstractmethod - def delete_leader(self): - """Voluntarily remove leader key from DCS + def _delete_leader(self): + """Remove leader key from DCS. This method should remove leader key if current instance is the leader""" + def delete_leader(self, last_operation=None): + """Update optime/leader and voluntarily remove leader key from DCS. + This method should remove leader key if current instance is the leader. + :param last_operation: latest checkpoint location in bytes""" + + if last_operation: + self.write_leader_optime(last_operation) + return self._delete_leader() + @abc.abstractmethod def cancel_initialization(self): """ Removes the initialize key for a cluster """ diff --git a/patroni/dcs/consul.py b/patroni/dcs/consul.py index efa81ea1..c22d0f8e 100644 --- a/patroni/dcs/consul.py +++ b/patroni/dcs/consul.py @@ -503,7 +503,7 @@ class Consul(AbstractDCS): return self._client.kv.put(self.history_path, value) @catch_consul_errors - def delete_leader(self): + def _delete_leader(self): cluster = self.cluster if cluster and isinstance(cluster.leader, Leader) and cluster.leader.name == self._name: return self._client.kv.delete(self.leader_path, cas=cluster.leader.index) diff --git a/patroni/dcs/etcd.py b/patroni/dcs/etcd.py index 844948ea..cb13f4e7 100644 --- a/patroni/dcs/etcd.py +++ b/patroni/dcs/etcd.py @@ -625,7 +625,7 @@ class Etcd(AbstractDCS): return self.retry(self._client.write, self.initialize_path, sysid, prevExist=(not create_new)) @catch_etcd_errors - def delete_leader(self): + def _delete_leader(self): return self._client.delete(self.leader_path, prevValue=self._name) @catch_etcd_errors diff --git a/patroni/dcs/kubernetes.py b/patroni/dcs/kubernetes.py index 9244ca04..15f4463d 100644 --- a/patroni/dcs/kubernetes.py +++ b/patroni/dcs/kubernetes.py @@ -546,9 +546,15 @@ class Kubernetes(AbstractDCS): resource_version = cluster.config.index if cluster and cluster.config and cluster.config.index else None return self.patch_or_create_config({self._INITIALIZE: sysid}, resource_version) - def delete_leader(self): + def _delete_leader(self): + """Unused""" + + def delete_leader(self, last_operation=None): if self.cluster and isinstance(self.cluster.leader, Leader) and self.cluster.leader.name == self._name: - self.patch_or_create(self.leader_path, {self._LEADER: None}, self._leader_resource_version, True, False, []) + annotations = {self._LEADER: None} + if last_operation: + annotations[self._OPTIME] = last_operation + self.patch_or_create(self.leader_path, annotations, self._leader_resource_version, True, False, []) self.reset_cluster() def cancel_initialization(self): diff --git a/patroni/dcs/zookeeper.py b/patroni/dcs/zookeeper.py index 8e9f5637..336fca70 100644 --- a/patroni/dcs/zookeeper.py +++ b/patroni/dcs/zookeeper.py @@ -315,7 +315,7 @@ class ZooKeeper(AbstractDCS): def _update_leader(self): return True - def delete_leader(self): + def _delete_leader(self): self._client.restart() return True diff --git a/patroni/ha.py b/patroni/ha.py index 07a8122e..4b62159f 100644 --- a/patroni/ha.py +++ b/patroni/ha.py @@ -748,13 +748,13 @@ class Ha(object): return self._is_healthiest_node(members.values()) - def _delete_leader(self): + def _delete_leader(self, last_operation=None): self.set_is_leader(False) - self.dcs.delete_leader() + self.dcs.delete_leader(last_operation) self.dcs.reset_cluster() - def release_leader_key_voluntarily(self): - self._delete_leader() + def release_leader_key_voluntarily(self, last_operation=None): + self._delete_leader(last_operation) self.touch_member() logger.info("Leader key released") @@ -784,8 +784,9 @@ class Ha(object): self.set_is_leader(False) if mode_control['release']: + checkpoint_location = self.state_handler.latest_checkpoint_location() if mode == 'graceful' else None with self._async_executor: - self.release_leader_key_voluntarily() + self.release_leader_key_voluntarily(checkpoint_location) time.sleep(2) # Give a time to somebody to take the leader lock if mode_control['offline']: node_to_follow, leader = None, None @@ -1384,7 +1385,8 @@ class Ha(object): stop_timeout=self.master_stop_timeout())) if not self.state_handler.is_running(): if self.has_lock(): - self.dcs.delete_leader() + checkpoint_location = self.state_handler.latest_checkpoint_location() + self.dcs.delete_leader(checkpoint_location) self.touch_member() else: # XXX: what about when Patroni is started as the wrong user that has access to the watchdog device diff --git a/patroni/postgresql/__init__.py b/patroni/postgresql/__init__.py index 248c3a59..6bed69e9 100644 --- a/patroni/postgresql/__init__.py +++ b/patroni/postgresql/__init__.py @@ -13,7 +13,7 @@ from patroni.postgresql.bootstrap import Bootstrap from patroni.postgresql.cancellable import CancellableSubprocess from patroni.postgresql.config import ConfigHandler from patroni.postgresql.connection import Connection, get_connection_cursor -from patroni.postgresql.misc import parse_history, postgres_major_version_to_int +from patroni.postgresql.misc import parse_history, parse_lsn, postgres_major_version_to_int from patroni.postgresql.postmaster import PostmasterProcess from patroni.postgresql.slots import SlotsHandler from patroni.exceptions import PostgresConnectionException @@ -310,6 +310,17 @@ class Postgresql(object): except (TypeError, ValueError): logger.exception('Failed to parse timeline from pg_controldata output') + def latest_checkpoint_location(self): + """Returns checkpoint location for the cleanly shut down primary""" + + data = self.controldata() + lsn = data.get('Latest checkpoint location') + if data.get('Database cluster state') == 'shut down' and lsn: + try: + return str(parse_lsn(lsn)) + except (IndexError, ValueError) as e: + logger.error('Exception when parsing lsn %s: %r', lsn, e) + def is_running(self): """Returns PostmasterProcess if one is running on the data directory or None. If most recently seen process is running updates the cached process based on pid file.""" diff --git a/tests/test_ha.py b/tests/test_ha.py index 50a1bc75..80550ec8 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -154,7 +154,10 @@ def run_async(self, func, args=()): @patch.object(Postgresql, '_cluster_info_state_get', Mock(return_value=3)) @patch.object(Postgresql, 'call_nowait', Mock(return_value=True)) @patch.object(Postgresql, 'data_directory_empty', Mock(return_value=False)) -@patch.object(Postgresql, 'controldata', Mock(return_value={'Database system identifier': SYSID})) +@patch.object(Postgresql, 'controldata', Mock(return_value={ + 'Database system identifier': SYSID, + 'Database cluster state': 'shut down', + 'Latest checkpoint location': '0/12345678'})) @patch.object(SlotsHandler, 'sync_replication_slots', Mock()) @patch.object(ConfigHandler, 'append_pg_hba', Mock()) @patch.object(ConfigHandler, 'write_pgpass', Mock(return_value={})) diff --git a/tests/test_kubernetes.py b/tests/test_kubernetes.py index 5fb89c23..e82870f4 100644 --- a/tests/test_kubernetes.py +++ b/tests/test_kubernetes.py @@ -113,7 +113,7 @@ class TestKubernetes(unittest.TestCase): self.k.initialize() def test_delete_leader(self): - self.k.delete_leader() + self.k.delete_leader(1) def test_cancel_initialization(self): self.k.cancel_initialization() diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index 78a9eeb5..858b6f2a 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -338,6 +338,11 @@ class TestPostgresql(BaseTestPostgresql): with patch.object(Postgresql, '_query', Mock(side_effect=RetryFailedError(''))): self.assertRaises(PostgresConnectionException, self.p.is_leader) + @patch.object(Postgresql, 'controldata', + Mock(return_value={'Database cluster state': 'shut down', 'Latest checkpoint location': 'X/678'})) + def test_latest_checkpoint_location(self): + self.assertIsNone(self.p.latest_checkpoint_location()) + def test_reload(self): self.assertTrue(self.p.reload())