From 650e2449042763ccac6bd923f4549db9a5b31253 Mon Sep 17 00:00:00 2001 From: Alexander Kukushkin Date: Fri, 4 Sep 2015 16:06:44 +0200 Subject: [PATCH] Refactor directory structure in preparation for building pypi-package --- patroni.py | 120 +---------------------- patroni/__init__.py | 119 ++++++++++++++++++++++ patroni/__main__.py | 5 + {helpers => patroni}/api.py | 0 {helpers => patroni}/dcs.py | 2 +- {helpers => patroni}/etcd.py | 4 +- {helpers => patroni}/ha.py | 2 +- {helpers => patroni}/postgresql.py | 2 +- {helpers => patroni/scripts}/__init__.py | 0 {scripts => patroni/scripts}/aws.py | 0 {scripts => patroni/scripts}/restore.py | 0 {helpers => patroni}/utils.py | 0 patroni/version.py | 1 + {helpers => patroni}/zookeeper.py | 4 +- scripts/__init__.py | 0 setup.py | 18 ++-- tests/test_api.py | 2 +- tests/test_aws.py | 2 +- tests/test_etcd.py | 4 +- tests/test_ha.py | 6 +- tests/test_patroni.py | 12 +-- tests/test_postgresql.py | 4 +- tests/test_restore.py | 2 +- tests/test_utils.py | 2 +- tests/test_zookeeper.py | 8 +- 25 files changed, 166 insertions(+), 153 deletions(-) create mode 100644 patroni/__init__.py create mode 100644 patroni/__main__.py rename {helpers => patroni}/api.py (100%) rename {helpers => patroni}/dcs.py (99%) rename {helpers => patroni}/etcd.py (98%) rename {helpers => patroni}/ha.py (99%) rename {helpers => patroni}/postgresql.py (99%) rename {helpers => patroni/scripts}/__init__.py (100%) rename {scripts => patroni/scripts}/aws.py (100%) rename {scripts => patroni/scripts}/restore.py (100%) rename {helpers => patroni}/utils.py (100%) create mode 100644 patroni/version.py rename {helpers => patroni}/zookeeper.py (98%) delete mode 100644 scripts/__init__.py diff --git a/patroni.py b/patroni.py index 3be2aa57..f34ef133 100755 --- a/patroni.py +++ b/patroni.py @@ -1,123 +1,5 @@ #!/usr/bin/env python -import logging -import os -import sys -import time -import yaml - -from helpers.api import RestApiServer -from helpers.etcd import Etcd -from helpers.ha import Ha -from helpers.postgresql import Postgresql -from helpers.utils import setup_signal_handlers, sleep, reap_children -from helpers.zookeeper import ZooKeeper - -logger = logging.getLogger(__name__) - - -class Patroni: - - def __init__(self, config): - self.nap_time = config['loop_wait'] - self.postgresql = Postgresql(config['postgresql']) - self.ha = Ha(self.postgresql, self.get_dcs(self.postgresql.name, config)) - host, port = config['restapi']['listen'].split(':') - self.api = RestApiServer(self, config['restapi']) - self.next_run = time.time() - self.shutdown_member_ttl = 300 - - @staticmethod - def get_dcs(name, config): - if 'etcd' in config: - return Etcd(name, config['etcd']) - if 'zookeeper' in config: - return ZooKeeper(name, config['zookeeper']) - raise Exception('Can not find sutable configuration of distributed configuration store') - - def touch_member(self, ttl=None): - connection_string = self.postgresql.connection_string + '?application_name=' + self.api.connection_string - if self.ha.cluster: - for m in self.ha.cluster.members: - # Do not update member TTL when it is far from being expired - if m.name == self.postgresql.name and m.real_ttl() > self.shutdown_member_ttl: - return True - return self.ha.dcs.touch_member(connection_string, ttl) - - def initialize(self): - # wait for etcd to be available - while not self.touch_member(): - logger.info('waiting on DCS') - sleep(5) - - # is data directory empty? - if self.postgresql.data_directory_empty(): - # racing to initialize - if self.ha.dcs.race('/initialize'): - self.postgresql.initialize() - self.ha.dcs.take_leader() - self.postgresql.start() - self.postgresql.create_replication_user() - self.postgresql.create_connection_users() - else: - while True: - leader = self.ha.dcs.current_leader() - if leader and self.postgresql.sync_from_leader(leader): - self.postgresql.write_recovery_conf(leader) - self.postgresql.start() - break - sleep(5) - elif self.postgresql.is_running(): - self.postgresql.load_replication_slots() - - def schedule_next_run(self): - self.next_run += self.nap_time - current_time = time.time() - nap_time = self.next_run - current_time - if nap_time <= 0: - self.next_run = current_time - else: - self.ha.dcs.sleep(nap_time) - - def run(self): - self.api.start() - self.next_run = time.time() - - while True: - self.touch_member() - logger.info(self.ha.run_cycle()) - try: - if self.ha.state_handler.is_leader(): - self.ha.cluster and self.ha.state_handler.create_replication_slots(self.ha.cluster) - else: - self.ha.state_handler.drop_replication_slots() - except: - logger.exception('Exception when changing replication slots') - reap_children() - self.schedule_next_run() - - -def main(): - logging.basicConfig(format='%(asctime)s %(levelname)s: %(message)s', level=logging.INFO) - logging.getLogger('requests').setLevel(logging.WARNING) - setup_signal_handlers() - - if len(sys.argv) < 2 or not os.path.isfile(sys.argv[1]): - print('Usage: {} config.yml'.format(sys.argv[0])) - return - - with open(sys.argv[1], 'r') as f: - config = yaml.load(f) - - patroni = Patroni(config) - try: - patroni.initialize() - patroni.run() - except KeyboardInterrupt: - pass - finally: - patroni.touch_member(patroni.shutdown_member_ttl) # schedule member removal - patroni.postgresql.stop() - patroni.ha.dcs.delete_leader() +from patroni import main if __name__ == '__main__': diff --git a/patroni/__init__.py b/patroni/__init__.py new file mode 100644 index 00000000..572f4263 --- /dev/null +++ b/patroni/__init__.py @@ -0,0 +1,119 @@ +import logging +import os +import sys +import time +import yaml + +from patroni.api import RestApiServer +from patroni.etcd import Etcd +from patroni.ha import Ha +from patroni.postgresql import Postgresql +from patroni.utils import setup_signal_handlers, sleep, reap_children +from patroni.zookeeper import ZooKeeper + +logger = logging.getLogger(__name__) + + +class Patroni: + + def __init__(self, config): + self.nap_time = config['loop_wait'] + self.postgresql = Postgresql(config['postgresql']) + self.ha = Ha(self.postgresql, self.get_dcs(self.postgresql.name, config)) + host, port = config['restapi']['listen'].split(':') + self.api = RestApiServer(self, config['restapi']) + self.next_run = time.time() + self.shutdown_member_ttl = 300 + + @staticmethod + def get_dcs(name, config): + if 'etcd' in config: + return Etcd(name, config['etcd']) + if 'zookeeper' in config: + return ZooKeeper(name, config['zookeeper']) + raise Exception('Can not find sutable configuration of distributed configuration store') + + def touch_member(self, ttl=None): + connection_string = self.postgresql.connection_string + '?application_name=' + self.api.connection_string + if self.ha.cluster: + for m in self.ha.cluster.members: + # Do not update member TTL when it is far from being expired + if m.name == self.postgresql.name and m.real_ttl() > self.shutdown_member_ttl: + return True + return self.ha.dcs.touch_member(connection_string, ttl) + + def initialize(self): + # wait for etcd to be available + while not self.touch_member(): + logger.info('waiting on DCS') + sleep(5) + + # is data directory empty? + if self.postgresql.data_directory_empty(): + # racing to initialize + if self.ha.dcs.race('/initialize'): + self.postgresql.initialize() + self.ha.dcs.take_leader() + self.postgresql.start() + self.postgresql.create_replication_user() + self.postgresql.create_connection_users() + else: + while True: + leader = self.ha.dcs.current_leader() + if leader and self.postgresql.sync_from_leader(leader): + self.postgresql.write_recovery_conf(leader) + self.postgresql.start() + break + sleep(5) + elif self.postgresql.is_running(): + self.postgresql.load_replication_slots() + + def schedule_next_run(self): + self.next_run += self.nap_time + current_time = time.time() + nap_time = self.next_run - current_time + if nap_time <= 0: + self.next_run = current_time + else: + self.ha.dcs.sleep(nap_time) + + def run(self): + self.api.start() + self.next_run = time.time() + + while True: + self.touch_member() + logger.info(self.ha.run_cycle()) + try: + if self.ha.state_handler.is_leader(): + self.ha.cluster and self.ha.state_handler.create_replication_slots(self.ha.cluster) + else: + self.ha.state_handler.drop_replication_slots() + except: + logger.exception('Exception when changing replication slots') + reap_children() + self.schedule_next_run() + + +def main(): + logging.basicConfig(format='%(asctime)s %(levelname)s: %(message)s', level=logging.INFO) + logging.getLogger('requests').setLevel(logging.WARNING) + setup_signal_handlers() + + if len(sys.argv) < 2 or not os.path.isfile(sys.argv[1]): + print('Usage: {} config.yml'.format(sys.argv[0])) + return + + with open(sys.argv[1], 'r') as f: + config = yaml.load(f) + + patroni = Patroni(config) + try: + patroni.initialize() + patroni.run() + except KeyboardInterrupt: + pass + finally: + patroni.touch_member(patroni.shutdown_member_ttl) # schedule member removal + patroni.postgresql.stop() + patroni.ha.dcs.delete_leader() diff --git a/patroni/__main__.py b/patroni/__main__.py new file mode 100644 index 00000000..3abcbfc3 --- /dev/null +++ b/patroni/__main__.py @@ -0,0 +1,5 @@ +from patroni import main + + +if __name__ == '__main__': + main() diff --git a/helpers/api.py b/patroni/api.py similarity index 100% rename from helpers/api.py rename to patroni/api.py diff --git a/helpers/dcs.py b/patroni/dcs.py similarity index 99% rename from helpers/dcs.py rename to patroni/dcs.py index c7140c22..f5eb4d2c 100644 --- a/helpers/dcs.py +++ b/patroni/dcs.py @@ -1,7 +1,7 @@ import abc from collections import namedtuple -from helpers.utils import calculate_ttl, sleep +from patroni.utils import calculate_ttl, sleep from six.moves.urllib_parse import urlparse, urlunparse, parse_qsl diff --git a/helpers/etcd.py b/patroni/etcd.py similarity index 98% rename from helpers/etcd.py rename to patroni/etcd.py index add736c5..0ebb99f4 100644 --- a/helpers/etcd.py +++ b/patroni/etcd.py @@ -8,8 +8,8 @@ import socket from dns.exception import DNSException from dns import resolver -from helpers.dcs import AbstractDCS, Cluster, DCSError, Member, parse_connection_string -from helpers.utils import sleep +from patroni.dcs import AbstractDCS, Cluster, DCSError, Member, parse_connection_string +from patroni.utils import sleep from requests.exceptions import RequestException logger = logging.getLogger(__name__) diff --git a/helpers/ha.py b/patroni/ha.py similarity index 99% rename from helpers/ha.py rename to patroni/ha.py index 8e283c55..f31b2b26 100644 --- a/helpers/ha.py +++ b/patroni/ha.py @@ -1,6 +1,6 @@ import logging -from helpers.dcs import DCSError +from patroni.dcs import DCSError from psycopg2 import InterfaceError, OperationalError logger = logging.getLogger(__name__) diff --git a/helpers/postgresql.py b/patroni/postgresql.py similarity index 99% rename from helpers/postgresql.py rename to patroni/postgresql.py index 3fc7983a..2c86800d 100644 --- a/helpers/postgresql.py +++ b/patroni/postgresql.py @@ -5,7 +5,7 @@ import shlex import shutil import subprocess -from helpers.utils import sleep +from patroni.utils import sleep from six.moves.urllib_parse import urlparse logger = logging.getLogger(__name__) diff --git a/helpers/__init__.py b/patroni/scripts/__init__.py similarity index 100% rename from helpers/__init__.py rename to patroni/scripts/__init__.py diff --git a/scripts/aws.py b/patroni/scripts/aws.py similarity index 100% rename from scripts/aws.py rename to patroni/scripts/aws.py diff --git a/scripts/restore.py b/patroni/scripts/restore.py similarity index 100% rename from scripts/restore.py rename to patroni/scripts/restore.py diff --git a/helpers/utils.py b/patroni/utils.py similarity index 100% rename from helpers/utils.py rename to patroni/utils.py diff --git a/patroni/version.py b/patroni/version.py new file mode 100644 index 00000000..11d27f8c --- /dev/null +++ b/patroni/version.py @@ -0,0 +1 @@ +__version__ = '0.1' diff --git a/helpers/zookeeper.py b/patroni/zookeeper.py similarity index 98% rename from helpers/zookeeper.py rename to patroni/zookeeper.py index cb2918cd..bbc05002 100644 --- a/helpers/zookeeper.py +++ b/patroni/zookeeper.py @@ -3,10 +3,10 @@ import random import requests import time -from helpers.dcs import AbstractDCS, Cluster, DCSError, Member, parse_connection_string -from helpers.utils import sleep from kazoo.client import KazooClient, KazooState from kazoo.exceptions import NoNodeError, NodeExistsError +from patroni.dcs import AbstractDCS, Cluster, DCSError, Member, parse_connection_string +from patroni.utils import sleep from requests.exceptions import RequestException logger = logging.getLogger(__name__) diff --git a/scripts/__init__.py b/scripts/__init__.py deleted file mode 100644 index e69de29b..00000000 diff --git a/setup.py b/setup.py index 62d05a69..4ed8ca37 100644 --- a/setup.py +++ b/setup.py @@ -19,13 +19,20 @@ if sys.version_info < (2, 7, 0): __location__ = os.path.join(os.getcwd(), os.path.dirname(inspect.getfile(inspect.currentframe()))) +def read_version(package): + data = {} + with open(os.path.join(package, 'version.py'), 'r') as fd: + exec(fd.read(), data) + return data['__version__'] + + NAME = 'patroni' -MAIN_PACKAGE = 'patroni.py' -HELPERS = 'helpers' +MAIN_PACKAGE = NAME SCRIPTS = 'scripts' +VERSION = read_version(MAIN_PACKAGE) VERSION = '0.1' DESCRIPTION = 'A Template for PostgreSQL HA with etcd' -LICENSE = 'The MIT License' +LICENSE = 'MIT License' COVERAGE_XML = True COVERAGE_HTML = False @@ -62,8 +69,7 @@ class PyTest(TestCommand): def finalize_options(self): TestCommand.finalize_options(self) if self.cov_xml or self.cov_html: - self.cov = ['--cov', MAIN_PACKAGE, '--cov', HELPERS, '--cov', SCRIPTS, '--cov-report', - 'term-missing'] + self.cov = ['--cov', MAIN_PACKAGE, '--cov', MAIN_PACKAGE, '--cov-report', 'term-missing'] if self.cov_xml: self.cov.extend(['--cov-report', 'xml']) if self.cov_html: @@ -82,7 +88,7 @@ class PyTest(TestCommand): params['plugins'] = ['cov'] if self.junitxml: params['args'] += self.junitxml - params['args'] += ['--doctest-modules', HELPERS, '--doctest-modules', SCRIPTS, '-s'] + params['args'] += ['--doctest-modules', MAIN_PACKAGE, '-s', '-vv'] errno = pytest.main(**params) sys.exit(errno) diff --git a/tests/test_api.py b/tests/test_api.py index 91b36943..aecf3df6 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -1,7 +1,7 @@ import psycopg2 import unittest -from helpers.api import RestApiHandler, RestApiServer +from patroni.api import RestApiHandler, RestApiServer from six import BytesIO as IO from test_postgresql import psycopg2_connect diff --git a/tests/test_aws.py b/tests/test_aws.py index 84d495fe..09c357f4 100644 --- a/tests/test_aws.py +++ b/tests/test_aws.py @@ -2,7 +2,7 @@ import unittest import requests import boto.ec2 from collections import namedtuple -from scripts.aws import AWSConnection +from patroni.scripts.aws import AWSConnection from requests.exceptions import RequestException diff --git a/tests/test_etcd.py b/tests/test_etcd.py index 38692c52..bc88e6a7 100644 --- a/tests/test_etcd.py +++ b/tests/test_etcd.py @@ -8,9 +8,9 @@ import time import unittest from dns.exception import DNSException -from helpers.dcs import Cluster, DCSError, Member -from helpers.etcd import Client, Etcd from mock import Mock, patch +from patroni.dcs import Cluster, DCSError, Member +from patroni.etcd import Client, Etcd class MockResponse: diff --git a/tests/test_ha.py b/tests/test_ha.py index abfc4ac8..46aa9a0a 100644 --- a/tests/test_ha.py +++ b/tests/test_ha.py @@ -1,9 +1,9 @@ import unittest -from helpers.dcs import Cluster, DCSError -from helpers.etcd import Client, Etcd -from helpers.ha import Ha from mock import Mock, patch +from patroni.dcs import Cluster, DCSError +from patroni.etcd import Client, Etcd +from patroni.ha import Ha from test_etcd import etcd_read, etcd_write diff --git a/tests/test_patroni.py b/tests/test_patroni.py index 67ac1ffd..979e603c 100644 --- a/tests/test_patroni.py +++ b/tests/test_patroni.py @@ -1,5 +1,5 @@ import datetime -import helpers.zookeeper +import patroni.zookeeper import psycopg2 import subprocess import sys @@ -7,12 +7,12 @@ import time import unittest import yaml -from helpers.api import RestApiServer -from helpers.dcs import Cluster, Member -from helpers.etcd import Etcd -from helpers.zookeeper import ZooKeeper from mock import Mock, patch +from patroni.api import RestApiServer +from patroni.dcs import Cluster, Member +from patroni.etcd import Etcd from patroni import Patroni, main +from patroni.zookeeper import ZooKeeper from six.moves import BaseHTTPServer from test_etcd import Client, etcd_read, etcd_write from test_ha import true, false @@ -70,7 +70,7 @@ class TestPatroni(unittest.TestCase): Postgresql.write_recovery_conf = self.write_recovery_conf def test_get_dcs(self): - helpers.zookeeper.KazooClient = MockKazooClient + patroni.zookeeper.KazooClient = MockKazooClient self.assertIsInstance(self.p.get_dcs('', {'zookeeper': {'scope': '', 'hosts': ''}}), ZooKeeper) self.assertRaises(Exception, self.p.get_dcs, '', {}) diff --git a/tests/test_postgresql.py b/tests/test_postgresql.py index a63c240b..23fd9d74 100644 --- a/tests/test_postgresql.py +++ b/tests/test_postgresql.py @@ -4,8 +4,8 @@ import shutil import subprocess import unittest -from helpers.dcs import Cluster, Member -from helpers.postgresql import Postgresql +from patroni.dcs import Cluster, Member +from patroni.postgresql import Postgresql def nop(*args, **kwargs): diff --git a/tests/test_restore.py b/tests/test_restore.py index 38b34c6a..2ffd8a58 100644 --- a/tests/test_restore.py +++ b/tests/test_restore.py @@ -1,7 +1,7 @@ import unittest from mock import MagicMock, patch import os -from scripts.restore import Restore, WALERestore +from patroni.scripts.restore import Restore, WALERestore def fake_cursor_fetchone(*args, **kwargs): diff --git a/tests/test_utils.py b/tests/test_utils.py index 312277b6..9a383ca0 100644 --- a/tests/test_utils.py +++ b/tests/test_utils.py @@ -2,7 +2,7 @@ import os import time import unittest -from helpers.utils import reap_children, sigchld_handler, sigterm_handler, sleep +from patroni.utils import reap_children, sigchld_handler, sigterm_handler, sleep def nop(*args, **kwargs): diff --git a/tests/test_zookeeper.py b/tests/test_zookeeper.py index 53f5ab82..bf06fd19 100644 --- a/tests/test_zookeeper.py +++ b/tests/test_zookeeper.py @@ -1,8 +1,8 @@ -import helpers.zookeeper +import patroni.zookeeper import requests import unittest -from helpers.zookeeper import ExhibitorEnsembleProvider, ZooKeeper, ZooKeeperError +from patroni.zookeeper import ExhibitorEnsembleProvider, ZooKeeper, ZooKeeperError from kazoo.client import KazooState from kazoo.exceptions import NoNodeError, NodeExistsError from kazoo.protocol.states import ZnodeStat @@ -105,7 +105,7 @@ class TestExhibitorEnsembleProvider(unittest.TestCase): def set_up(self): requests.get = requests_get - helpers.zookeeper.sleep = exhibitor_sleep + patroni.zookeeper.sleep = exhibitor_sleep def test_init(self): self.assertRaises(Exception, ExhibitorEnsembleProvider, ['localhost'], 8181) @@ -119,7 +119,7 @@ class TestZooKeeper(unittest.TestCase): def set_up(self): requests.get = requests_get - helpers.zookeeper.KazooClient = MockKazooClient + patroni.zookeeper.KazooClient = MockKazooClient self.zk = ZooKeeper('foo', {'exhibitor': {'hosts': ['localhost', 'exhibitor'], 'port': 8181}, 'scope': 'test'}) def test_session_listener(self):