mirror of
https://github.com/outbackdingo/patroni.git
synced 2026-08-26 15:40:21 +00:00
Compare commits
@@ -54,3 +54,6 @@ docs/source/_templates/
|
||||
|
||||
# Pycharm IDE
|
||||
.idea/
|
||||
|
||||
#VSCode IDE
|
||||
.vscode/
|
||||
|
||||
+6
-2
@@ -1,6 +1,10 @@
|
||||
sudo: true
|
||||
dist: trusty
|
||||
language: python
|
||||
addons:
|
||||
apt:
|
||||
packages:
|
||||
- expect-dev # for unbuffer
|
||||
env:
|
||||
global:
|
||||
- ETCDVERSION=3.0.17 ZKVERSION=3.4.11 CONSULVERSION=0.7.4
|
||||
@@ -131,11 +135,11 @@ script:
|
||||
|
||||
if [[ $TEST_SUITE != "behave" ]]; then
|
||||
echo Running unit tests using python${pv}
|
||||
$TEST_SUITE test
|
||||
unbuffer $TEST_SUITE test
|
||||
$TEST_SUITE flake8
|
||||
elif [[ $pv != $EXCLUDE_BEHAVE ]]; then
|
||||
echo Running acceptance tests using python${pv}
|
||||
if ! PATH=.:/usr/lib/postgresql/9.6/bin:$PATH $TEST_SUITE; then
|
||||
if ! PATH=.:/usr/lib/postgresql/9.6/bin:$PATH unbuffer $TEST_SUITE; then
|
||||
# output all log files when tests are failing
|
||||
grep . features/output/*_failed/*postgres?.*
|
||||
exit 1
|
||||
|
||||
+141
-30
@@ -1,45 +1,156 @@
|
||||
## This Dockerfile is meant to aid in the building and debugging patroni whilst developing on your local machine
|
||||
## It has all the necessary components to play/debug with a single node appliance, running etcd
|
||||
FROM postgres:10
|
||||
MAINTAINER Alexander Kukushkin <[email protected]>
|
||||
ARG PG_MAJOR=10
|
||||
ARG COMPRESS=false
|
||||
ARG PGHOME=/home/postgres
|
||||
ARG PGDATA=$PGHOME/data
|
||||
ARG LC_ALL=C.UTF-8
|
||||
ARG LANG=C.UTF-8
|
||||
|
||||
RUN export DEBIAN_FRONTEND=noninteractive \
|
||||
FROM postgres:$PG_MAJOR as builder
|
||||
|
||||
ARG PGHOME
|
||||
ARG PGDATA
|
||||
ARG LC_ALL
|
||||
ARG LANG
|
||||
|
||||
ENV ETCDVERSION=2.3.8 CONFDVERSION=0.16.0
|
||||
|
||||
RUN set -ex \
|
||||
&& export DEBIAN_FRONTEND=noninteractive \
|
||||
&& echo 'APT::Install-Recommends "0";\nAPT::Install-Suggests "0";' > /etc/apt/apt.conf.d/01norecommend \
|
||||
&& apt-get update -y \
|
||||
&& apt-get upgrade -y \
|
||||
# postgres:10 is based on debian, which has patroni package. We will install all required dependencies
|
||||
&& apt-get install -s patroni | sed -n -e '/^Inst patroni /d' -e 's/^Inst \([^ ]\+\) .*$/\1/p' \
|
||||
| xargs apt-get install -y curl jq haproxy locales python3-etcd python3-kazoo \
|
||||
|
||||
## Make sure we have a en_US.UTF-8 locale available
|
||||
# postgres:10 is based on debian, which has the patroni package. We will install all required dependencies
|
||||
&& apt-cache depends patroni | sed -n -e 's/.*Depends: \(python3-.\+\)$/\1/p' \
|
||||
| grep -Ev '^python3-(sphinx|etcd|consul|kazoo|kubernetes)' \
|
||||
| xargs apt-get install -y vim curl less jq locales haproxy sudo \
|
||||
python3-etcd python3-kazoo python3-pip busybox \
|
||||
&& pip3 install dumb-init \
|
||||
\
|
||||
# Cleanup all locales but en_US.UTF-8
|
||||
&& find /usr/share/i18n/charmaps/ -type f ! -name UTF-8.gz -delete \
|
||||
&& find /usr/share/i18n/locales/ -type f ! -name en_US ! -name en_GB ! -name i18n ! -name iso14651_t1 ! -name iso14651_t1_common ! -name 'translit_*' -delete \
|
||||
&& echo 'en_US.UTF-8 UTF-8' > /usr/share/i18n/SUPPORTED \
|
||||
\
|
||||
# Make sure we have a en_US.UTF-8 locale available
|
||||
&& localedef -i en_US -c -f UTF-8 -A /usr/share/locale/locale.alias en_US.UTF-8 \
|
||||
|
||||
&& mkdir -p /home/postgres \
|
||||
&& chown postgres:postgres /home/postgres \
|
||||
|
||||
# Clean up
|
||||
&& apt-get purge -y libpython2.7-stdlib libpython2.7-minimal \
|
||||
\
|
||||
# haproxy dummy config
|
||||
&& echo 'global\n stats socket /run/haproxy/admin.sock mode 660 level admin' > /etc/haproxy/haproxy.cfg \
|
||||
\
|
||||
# vim config
|
||||
&& echo 'syntax on\nfiletype plugin indent on\nset mouse-=a\nautocmd FileType yaml setlocal ts=2 sts=2 sw=2 expandtab' > /etc/vim/vimrc.local \
|
||||
\
|
||||
# Prepare postgres/patroni/haproxy environment
|
||||
&& mkdir -p $PGHOME/.config/patroni /patroni /run/haproxy \
|
||||
&& ln -s ../../patroni.yaml $PGHOME/.config/patroni/patronictl.yaml \
|
||||
&& ln -s /patronictl.py /usr/local/bin/patronictl \
|
||||
&& sed -i "s|/var/lib/postgresql.*|$PGHOME:/bin/bash|" /etc/passwd \
|
||||
&& chown -R postgres:postgres /var/log \
|
||||
\
|
||||
# Download etcd
|
||||
&& curl -sL https://github.com/coreos/etcd/releases/download/v${ETCDVERSION}/etcd-v${ETCDVERSION}-linux-amd64.tar.gz \
|
||||
| tar xz -C /usr/local/bin --strip=1 --wildcards --no-anchored etcd etcdctl \
|
||||
\
|
||||
# Download confd
|
||||
&& curl -sL https://github.com/kelseyhightower/confd/releases/download/v${CONFDVERSION}/confd-${CONFDVERSION}-linux-amd64 \
|
||||
> /usr/local/bin/confd && chmod +x /usr/local/bin/confd \
|
||||
\
|
||||
# Clean up all useless packages and some files
|
||||
&& apt-get purge -y --allow-remove-essential python3-pip gzip bzip2 util-linux e2fsprogs \
|
||||
libmagic1 bsdmainutils login ncurses-bin libmagic-mgc e2fslibs bsdutils \
|
||||
exim4-config gnupg-agent dirmngr libpython2.7-stdlib libpython2.7-minimal \
|
||||
&& apt-get autoremove -y \
|
||||
&& apt-get clean -y \
|
||||
&& rm -rf /var/lib/apt/lists/* /root/.cache
|
||||
&& rm -rf /var/lib/apt/lists/* \
|
||||
/root/.cache \
|
||||
/var/cache/debconf/* \
|
||||
/etc/rc?.d \
|
||||
/etc/systemd \
|
||||
/docker-entrypoint* \
|
||||
/sbin/pam* \
|
||||
/sbin/swap* \
|
||||
/sbin/unix* \
|
||||
/usr/local/bin/gosu \
|
||||
/usr/sbin/[acgipr]* \
|
||||
/usr/sbin/*user* \
|
||||
/usr/share/doc* \
|
||||
/usr/share/man \
|
||||
/usr/share/info \
|
||||
/usr/share/i18n/locales/translit_hangul \
|
||||
/usr/share/locale/?? \
|
||||
/usr/share/locale/??_?? \
|
||||
/usr/share/postgresql/*/man \
|
||||
/usr/share/postgresql-common/pg_wrapper \
|
||||
/usr/share/vim/vim80/doc \
|
||||
/usr/share/vim/vim80/lang \
|
||||
/usr/share/vim/vim80/tutor \
|
||||
# /var/lib/dpkg/info/* \
|
||||
&& find /usr/bin -xtype l -delete \
|
||||
&& find /var/log -type f -exec truncate --size 0 {} \; \
|
||||
&& find /usr/lib/python3/dist-packages -name '*test*' | xargs rm -fr \
|
||||
&& find /lib/x86_64-linux-gnu/security -type f ! -name pam_env.so ! -name pam_permit.so ! -name pam_unix.so -delete
|
||||
|
||||
ENV ETCDVERSION 3.2.23
|
||||
RUN curl -L https://github.com/coreos/etcd/releases/download/v${ETCDVERSION}/etcd-v${ETCDVERSION}-linux-amd64.tar.gz \
|
||||
| tar xz -C /usr/local/bin --strip=1 --wildcards --no-anchored etcd etcdctl
|
||||
# perform compression if it is necessary
|
||||
ARG COMPRESS
|
||||
RUN if [ "$COMPRESS" = "true" ]; then \
|
||||
set -ex \
|
||||
# Allow certain sudo commands from postgres
|
||||
&& echo 'postgres ALL=(ALL) NOPASSWD: /bin/tar xpJf /a.tar.xz -C /, /bin/rm /a.tar.xz, /bin/ln -snf dash /bin/sh' >> /etc/sudoers \
|
||||
&& ln -snf busybox /bin/sh \
|
||||
&& files="/bin/sh /usr/bin/sudo /usr/lib/sudo/sudoers.so /lib/x86_64-linux-gnu/security/pam_*.so" \
|
||||
&& libs="$(ldd $files | awk '{print $3;}' | grep '^/' | sort -u) /lib/x86_64-linux-gnu/ld-linux-x86-64.so.* /lib/x86_64-linux-gnu/libnsl.so.* /lib/x86_64-linux-gnu/libnss_compat.so.*" \
|
||||
&& (echo /var/run $files $libs | tr ' ' '\n' && realpath $files $libs) | sort -u | sed 's/^\///' > /exclude \
|
||||
&& find /etc/alternatives -xtype l -delete \
|
||||
&& save_dirs="usr lib var bin sbin etc/ssl etc/init.d etc/alternatives etc/apt" \
|
||||
&& XZ_OPT=-e9v tar -X /exclude -cpJf a.tar.xz $save_dirs \
|
||||
# we call "cat /exclude" to avoid including files from the $save_dirs that are also among
|
||||
# the exceptions listed in the /exclude, as "uniq -u" eliminates all non-unique lines.
|
||||
# By calling "cat /exclude" a second time we guarantee that there will be at least two lines
|
||||
# for each exception and therefore they will be excluded from the output passed to 'rm'.
|
||||
&& /bin/busybox sh -c "(find $save_dirs -not -type d && cat /exclude /exclude && echo exclude) | sort | uniq -u | xargs /bin/busybox rm" \
|
||||
&& /bin/busybox --install -s \
|
||||
&& /bin/busybox sh -c "find $save_dirs -type d -depth -exec rmdir -p {} \; 2> /dev/null"; \
|
||||
fi
|
||||
|
||||
ENV CONFDVERSION 0.16.0
|
||||
RUN curl -L https://github.com/kelseyhightower/confd/releases/download/v${CONFDVERSION}/confd-${CONFDVERSION}-linux-amd64 > /usr/local/bin/confd \
|
||||
&& chmod +x /usr/local/bin/confd
|
||||
FROM scratch
|
||||
COPY --from=builder / /
|
||||
|
||||
ADD patronictl.py patroni.py docker/entrypoint.sh /
|
||||
ADD patroni /patroni/
|
||||
ADD extras/confd /etc/confd
|
||||
LABEL maintainer="Alexander Kukushkin <[email protected]>"
|
||||
|
||||
RUN sed -i 's/env python/&3/' patroni*.py && ln -s /patronictl.py /usr/local/bin/patronictl && mkdir /data/ /run/haproxy \
|
||||
&& touch /pgpass /patroni.yml && chown postgres:postgres -R /patroni/ /data/ /pgpass /patroni.yml /etc/haproxy /var/run/ /var/lib/ /var/log/
|
||||
ARG PG_MAJOR
|
||||
ARG COMPRESS
|
||||
ARG PGHOME
|
||||
ARG PGDATA
|
||||
ARG LC_ALL
|
||||
ARG LANG
|
||||
|
||||
EXPOSE 2379 5432 8008
|
||||
ARG PGBIN=/usr/lib/postgresql/$PG_MAJOR/bin
|
||||
|
||||
ENV LC_ALL=$LC_ALL LANG=$LANG EDITOR=/usr/bin/editor
|
||||
ENV PGDATA=$PGDATA PATH=$PATH:$PGBIN
|
||||
|
||||
COPY patroni /patroni/
|
||||
COPY extras/confd/conf.d/haproxy.toml /etc/confd/conf.d/
|
||||
COPY extras/confd/templates/haproxy.tmpl /etc/confd/templates/
|
||||
COPY patroni*.py docker/entrypoint.sh /
|
||||
COPY postgres?.yml $PGHOME/
|
||||
|
||||
WORKDIR $PGHOME
|
||||
|
||||
RUN sed -i 's/env python/&3/' /patroni*.py \
|
||||
# "fix" patroni configs
|
||||
&& sed -i 's/^\( connect_address:\| - host\)/#&/' postgres?.yml \
|
||||
&& sed -i 's/^ listen: 127.0.0.1/ listen: 0.0.0.0/' postgres?.yml \
|
||||
&& sed -i "s|^\( data_dir: \).*|\1$PGDATA|" postgres?.yml \
|
||||
&& sed -i "s|^#\( bin_dir: \).*|\1$PGBIN|" postgres?.yml \
|
||||
&& sed -i 's/^ - encoding: UTF8/ - locale: en_US.UTF-8\n&/' postgres?.yml \
|
||||
&& sed -i 's/^\(scope\|name\|etcd\| host\| authentication\| pg_hba\| parameters\):/#&/' postgres?.yml \
|
||||
&& sed -i 's/^ \(replication\|superuser\|unix_socket_directories\|\(\( \)\{0,1\}\(username\|password\)\)\):/#&/' postgres?.yml \
|
||||
&& sed -i 's/^ parameters:/ pg_hba:\n - local all all trust\n - host replication all all md5\n - host all all all md5\n&\n max_connections: 100/' postgres?.yml \
|
||||
&& if [ "$COMPRESS" = "true" ]; then chmod u+s /usr/bin/sudo; fi \
|
||||
&& chown -R postgres:postgres $PGHOME /run /etc/haproxy
|
||||
|
||||
ENV LC_ALL=en_US.UTF-8 LANG=en_US.UTF-8
|
||||
ENTRYPOINT ["/bin/bash", "/entrypoint.sh"]
|
||||
USER postgres
|
||||
|
||||
ENTRYPOINT ["/bin/sh", "/entrypoint.sh"]
|
||||
|
||||
+63
-52
@@ -1,58 +1,69 @@
|
||||
# docker compose file for running a 3-node PostgreSQL cluster
|
||||
# with etcd as the SIS
|
||||
# with 3-node etcd cluster as the DCS and one haproxy node
|
||||
version: "2"
|
||||
|
||||
patroni_etcd:
|
||||
container_name: patroni_etcd
|
||||
image: patroni
|
||||
command: --etcd
|
||||
networks:
|
||||
demo:
|
||||
|
||||
dbnode1:
|
||||
image: patroni
|
||||
hostname: dbnode1
|
||||
links:
|
||||
- patroni_etcd:patroni_etcd
|
||||
volumes:
|
||||
- ./patroni:/patroni
|
||||
env_file: docker/patroni-secrets.env
|
||||
environment:
|
||||
PATRONI_ETCD_URL: http://patroni_etcd:2379
|
||||
PATRONI_NAME: dbnode1
|
||||
PATRONI_SCOPE: testcluster
|
||||
services:
|
||||
etcd1:
|
||||
image: patroni
|
||||
networks: [ demo ]
|
||||
env_file: docker/etcd.env
|
||||
container_name: demo-etcd1
|
||||
hostname: etcd1
|
||||
command: etcd -name etcd1 -initial-advertise-peer-urls http://etcd1:2380
|
||||
|
||||
dbnode2:
|
||||
image: patroni
|
||||
hostname: dbnode2
|
||||
links:
|
||||
- patroni_etcd:patroni_etcd
|
||||
volumes:
|
||||
- ./patroni:/patroni
|
||||
env_file: docker/patroni-secrets.env
|
||||
environment:
|
||||
PATRONI_ETCD_URL: http://patroni_etcd:2379
|
||||
PATRONI_NAME: dbnode2
|
||||
PATRONI_SCOPE: testcluster
|
||||
etcd2:
|
||||
image: patroni
|
||||
networks: [ demo ]
|
||||
env_file: docker/etcd.env
|
||||
container_name: demo-etcd2
|
||||
hostname: etcd2
|
||||
command: etcd -name etcd2 -initial-advertise-peer-urls http://etcd2:2380
|
||||
|
||||
dbnode3:
|
||||
image: patroni
|
||||
hostname: dbnode3
|
||||
links:
|
||||
- patroni_etcd:patroni_etcd
|
||||
volumes:
|
||||
- ./patroni:/patroni
|
||||
env_file: docker/patroni-secrets.env
|
||||
environment:
|
||||
PATRONI_ETCD_URL: http://patroni_etcd:2379
|
||||
PATRONI_NAME: dbnode3
|
||||
PATRONI_SCOPE: testcluster
|
||||
etcd3:
|
||||
image: patroni
|
||||
networks: [ demo ]
|
||||
env_file: docker/etcd.env
|
||||
container_name: demo-etcd3
|
||||
hostname: etcd3
|
||||
command: etcd -name etcd3 -initial-advertise-peer-urls http://etcd3:2380
|
||||
|
||||
haproxy:
|
||||
image: patroni
|
||||
links:
|
||||
- patroni_etcd:patroni_etcd
|
||||
ports:
|
||||
- "5000:5000"
|
||||
- "5001:5001"
|
||||
environment:
|
||||
PATRONI_ETCD_URL: http://patroni_etcd:2379
|
||||
PATRONI_SCOPE: testcluster
|
||||
command: --confd
|
||||
patroni1:
|
||||
image: patroni
|
||||
networks: [ demo ]
|
||||
env_file: docker/patroni.env
|
||||
hostname: patroni1
|
||||
container_name: demo-patroni1
|
||||
environment:
|
||||
PATRONI_NAME: patroni1
|
||||
|
||||
patroni2:
|
||||
image: patroni
|
||||
networks: [ demo ]
|
||||
env_file: docker/patroni.env
|
||||
hostname: patroni2
|
||||
container_name: demo-patroni2
|
||||
environment:
|
||||
PATRONI_NAME: patroni2
|
||||
|
||||
patroni3:
|
||||
image: patroni
|
||||
networks: [ demo ]
|
||||
env_file: docker/patroni.env
|
||||
hostname: patroni3
|
||||
container_name: demo-patroni3
|
||||
environment:
|
||||
PATRONI_NAME: patroni3
|
||||
|
||||
haproxy:
|
||||
image: patroni
|
||||
networks: [ demo ]
|
||||
env_file: docker/patroni.env
|
||||
hostname: haproxy
|
||||
container_name: demo-haproxy
|
||||
ports:
|
||||
- "5000:5000"
|
||||
- "5001:5001"
|
||||
command: haproxy
|
||||
|
||||
+103
-33
@@ -1,47 +1,117 @@
|
||||
# Patroni Dockerfile
|
||||
You can run Patroni in a docker container using this Dockerfile, or by using one of the Docker image at
|
||||
|
||||
https://registry.opensource.zalan.do/v1/repositories/acid/patroni/tags
|
||||
You can run Patroni in a docker container using this Dockerfile
|
||||
|
||||
This Dockerfile is meant in aiding development of Patroni and quick testing of features. It is not a production-worthy
|
||||
Dockerfile
|
||||
|
||||
docker build -t patroni .
|
||||
|
||||
# Examples
|
||||
|
||||
## Standalone Patroni
|
||||
|
||||
docker run -d registry.opensource.zalan.do/acid/patroni:1.0-SNAPSHOT
|
||||
docker run -d patroni
|
||||
|
||||
## Multiple Patroni's communicating with a standalone etcd inside Docker
|
||||
|
||||
Basically what you would do would be:
|
||||
|
||||
* Run 1 container which provides etcd
|
||||
|
||||
docker run -d <IMAGE> --etcd-only
|
||||
|
||||
* Run n containers running Patroni, passing the `--etcd` option to the `docker run` command
|
||||
|
||||
docker run -d <IMAGE> --etcd=<IP FROM etcd CONTAINER:PORT>
|
||||
|
||||
To automate this you can run the following script:
|
||||
|
||||
dev_patroni_cluster.sh [OPTIONS]
|
||||
|
||||
Options:
|
||||
|
||||
--image IMAGE The Docker image to use for the cluster
|
||||
--members INT The number of members for the cluster
|
||||
--name NAME The name of the new cluster
|
||||
## Three-node Patroni cluster with three-node etcd cluster and one haproxy container using docker-compose
|
||||
|
||||
Example session:
|
||||
|
||||
$ ./dev_patroni_cluster.sh --image registry.opensource.zalan.do/acid/patroni:1.0-SNAPSHOT --members=2 --name=bravo
|
||||
The etcd container is 6be871a11cb373406ca5ea1c6b39e1.0-SNAPSHOTfdde9fb1d6177212d6ad0c0d1bd9b563, ip=172.17.1.24
|
||||
Started Patroni container 67e611f2eca7c40f9e6e0e24a4a8f2cba7e3e56d22a420e15ab9240a37a9d7a4, ip=172.17.1.25
|
||||
Started Patroni container 47dd12ae635ab83b039f5889e250048b606ed5e48e3650b69e365e7e1d4acbcf, ip=172.17.1.26
|
||||
$ docker-compose up -d
|
||||
Creating demo-haproxy ...
|
||||
Creating demo-patroni2 ...
|
||||
Creating demo-patroni1 ...
|
||||
Creating demo-patroni3 ...
|
||||
Creating demo-etcd2 ...
|
||||
Creating demo-etcd1 ...
|
||||
Creating demo-etcd3 ...
|
||||
Creating demo-haproxy
|
||||
Creating demo-patroni2
|
||||
Creating demo-patroni1
|
||||
Creating demo-patroni3
|
||||
Creating demo-etcd1
|
||||
Creating demo-etcd2
|
||||
Creating demo-etcd2 ... done
|
||||
|
||||
$ docker ps
|
||||
CONTAINER ID IMAGE COMMAND CREATED STATUS PORTS NAMES
|
||||
47dd12ae635a registry.opensource.zalan.do/acid/patroni:1.0-SNAPSHOT "/bin/bash /entrypoi 10 seconds ago Up 8 seconds 4001/tcp, 5432/tcp, 2380/tcp bravo_OR64g8bx
|
||||
67e611f2eca7 registry.opensource.zalan.do/acid/patroni:1.0-SNAPSHOT "/bin/bash /entrypoi 11 seconds ago Up 10 seconds 2380/tcp, 4001/tcp, 5432/tcp bravo_si9no8iz
|
||||
6be871a11cb3 registry.opensource.zalan.do/acid/patroni:1.0-SNAPSHOT "/bin/bash /entrypoi 12 seconds ago Up 10 seconds 4001/tcp, 5432/tcp, 2380/tcp bravo_etcd
|
||||
CONTAINER ID IMAGE COMMAND CREATED STATUS PORTS NAMES
|
||||
5b7a90b4cfbf patroni "/bin/sh /entrypoint…" 29 seconds ago Up 27 seconds demo-etcd2
|
||||
e30eea5222f2 patroni "/bin/sh /entrypoint…" 29 seconds ago Up 27 seconds demo-etcd1
|
||||
83bcf3cb208f patroni "/bin/sh /entrypoint…" 29 seconds ago Up 27 seconds demo-etcd3
|
||||
922532c56e7d patroni "/bin/sh /entrypoint…" 29 seconds ago Up 28 seconds demo-patroni3
|
||||
14f875e445f3 patroni "/bin/sh /entrypoint…" 29 seconds ago Up 28 seconds demo-patroni2
|
||||
110d1073b383 patroni "/bin/sh /entrypoint…" 29 seconds ago Up 28 seconds demo-patroni1
|
||||
5af5e6e36028 patroni "/bin/sh /entrypoint…" 29 seconds ago Up 28 seconds 0.0.0.0:5000-5001->5000-5001/tcp demo-haproxy
|
||||
|
||||
$ docker logs demo-patroni1
|
||||
2019-02-20 08:19:32,714 INFO: Failed to import patroni.dcs.consul
|
||||
2019-02-20 08:19:32,737 INFO: Selected new etcd server http://etcd3:2379
|
||||
2019-02-20 08:19:35,140 INFO: Lock owner: None; I am patroni1
|
||||
2019-02-20 08:19:35,174 INFO: trying to bootstrap a new cluster
|
||||
...
|
||||
2019-02-20 08:19:39,310 INFO: postmaster pid=37
|
||||
2019-02-20 08:19:39.314 UTC [37] LOG: listening on IPv4 address "0.0.0.0", port 5432
|
||||
2019-02-20 08:19:39.321 UTC [37] LOG: listening on Unix socket "/var/run/postgresql/.s.PGSQL.5432"
|
||||
2019-02-20 08:19:39.353 UTC [39] LOG: database system was shut down at 2019-02-20 08:19:36 UTC
|
||||
2019-02-20 08:19:39.354 UTC [40] FATAL: the database system is starting up
|
||||
localhost:5432 - rejecting connections
|
||||
2019-02-20 08:19:39.369 UTC [37] LOG: database system is ready to accept connections
|
||||
localhost:5432 - accepting connections
|
||||
2019-02-20 08:19:39,383 INFO: establishing a new patroni connection to the postgres cluster
|
||||
2019-02-20 08:19:39,408 INFO: running post_bootstrap
|
||||
2019-02-20 08:19:39,432 WARNING: Could not activate Linux watchdog device: "Can't open watchdog device: [Errno 2] No such file or directory: '/dev/watchdog'"
|
||||
2019-02-20 08:19:39,515 INFO: initialized a new cluster
|
||||
2019-02-20 08:19:49,424 INFO: Lock owner: patroni1; I am patroni1
|
||||
2019-02-20 08:19:49,447 INFO: Lock owner: patroni1; I am patroni1
|
||||
2019-02-20 08:19:49,480 INFO: no action. i am the leader with the lock
|
||||
2019-02-20 08:19:59,422 INFO: Lock owner: patroni1; I am patroni1
|
||||
|
||||
$ docker exec -ti demo-patroni1 bash
|
||||
postgres@patroni1:~$ patronictl list
|
||||
+-------------+----------+------------+--------+---------+----+-----------+
|
||||
| Cluster | Member | Host | Role | State | TL | Lag in MB |
|
||||
+-------------+----------+------------+--------+---------+----+-----------+
|
||||
| testcluster | patroni1 | 172.21.0.3 | Leader | running | 1 | 0 |
|
||||
| testcluster | patroni2 | 172.21.0.4 | | running | 1 | 0 |
|
||||
| testcluster | patroni3 | 172.21.0.5 | | running | 1 | 0 |
|
||||
+-------------+----------+------------+--------+---------+----+-----------+
|
||||
|
||||
postgres@patroni1:~$ etcdctl ls --recursive --sort -p /service/testcluster
|
||||
/service/testcluster/config
|
||||
/service/testcluster/initialize
|
||||
/service/testcluster/leader
|
||||
/service/testcluster/members/
|
||||
/service/testcluster/members/patroni1
|
||||
/service/testcluster/members/patroni2
|
||||
/service/testcluster/members/patroni3
|
||||
/service/testcluster/optime/
|
||||
/service/testcluster/optime/leader
|
||||
|
||||
postgres@patroni1:~$ etcdctl member list
|
||||
1bab629f01fa9065: name=etcd3 peerURLs=http://etcd3:2380 clientURLs=http://etcd3:2379 isLeader=false
|
||||
8ecb6af518d241cc: name=etcd2 peerURLs=http://etcd2:2380 clientURLs=http://etcd2:2379 isLeader=true
|
||||
b2e169fcb8a34028: name=etcd1 peerURLs=http://etcd1:2380 clientURLs=http://etcd1:2379 isLeader=false
|
||||
postgres@patroni1:~$ exit
|
||||
|
||||
$ psql -h localhost -p 5000 -U postgres -W
|
||||
Password: postgres
|
||||
psql (11.2 (Ubuntu 11.2-1.pgdg18.04+1), server 10.7 (Debian 10.7-1.pgdg90+1))
|
||||
Type "help" for help.
|
||||
|
||||
localhost/postgres=# select pg_is_in_recovery();
|
||||
pg_is_in_recovery
|
||||
───────────────────
|
||||
f
|
||||
(1 row)
|
||||
|
||||
localhost/postgres=# \q
|
||||
|
||||
$ psql -h localhost -p 5001 -U postgres -W
|
||||
Password: postgres
|
||||
psql (11.2 (Ubuntu 11.2-1.pgdg18.04+1), server 10.7 (Debian 10.7-1.pgdg90+1))
|
||||
Type "help" for help.
|
||||
|
||||
localhost/postgres=# select pg_is_in_recovery();
|
||||
pg_is_in_recovery
|
||||
───────────────────
|
||||
t
|
||||
(1 row)
|
||||
|
||||
@@ -1,104 +0,0 @@
|
||||
#!/bin/bash
|
||||
|
||||
DOCKER_IMAGE="registry.opensource.zalan.do/acid/patroni:1.0-SNAPSHOT"
|
||||
MEMBERS=3
|
||||
|
||||
|
||||
function usage()
|
||||
{
|
||||
cat <<__EOF__
|
||||
Usage: $0
|
||||
|
||||
Options:
|
||||
|
||||
--image IMAGE The Docker image to use for the cluster
|
||||
--members INT The number of members for the cluster
|
||||
--name NAME The name of the new cluster
|
||||
|
||||
Examples:
|
||||
|
||||
$0 --image ${DOCKER_IMAGE}
|
||||
$0
|
||||
$0 --image ${DOCKER_IMAGE} --members=2
|
||||
__EOF__
|
||||
}
|
||||
|
||||
|
||||
optspec=":-:"
|
||||
while getopts "$optspec" optchar; do
|
||||
case "${optchar}" in
|
||||
-)
|
||||
case "${OPTARG}" in
|
||||
help)
|
||||
usage
|
||||
exit 0
|
||||
;;
|
||||
name)
|
||||
PATRONI_SCOPE="${!OPTIND}"; OPTIND=$(( $OPTIND + 1 ))
|
||||
;;
|
||||
name=*)
|
||||
PATRONI_SCOPE="${OPTARG#*=}"
|
||||
;;
|
||||
image)
|
||||
DOCKER_IMAGE="${!OPTIND}"; OPTIND=$(( $OPTIND + 1 ))
|
||||
;;
|
||||
image=*)
|
||||
DOCKER_IMAGE="${OPTARG#*=}"
|
||||
;;
|
||||
members)
|
||||
MEMBERS="${!OPTIND}"; OPTIND=$(( $OPTIND + 1 ))
|
||||
;;
|
||||
members=*)
|
||||
MEMBERS="${OPTARG#*=}"
|
||||
;;
|
||||
*)
|
||||
if [ "$OPTERR" = 1 ] && [ "${optspec:0:1}" != ":" ]; then
|
||||
echo "Unknown option --${OPTARG}" >&2
|
||||
fi
|
||||
;;
|
||||
esac;;
|
||||
*)
|
||||
if [ "$OPTERR" != 1 ] || [ "${optspec:0:1}" = ":" ]; then
|
||||
echo "Non-option argument: '-${OPTARG}'" >&2
|
||||
usage
|
||||
exit 1
|
||||
fi
|
||||
;;
|
||||
esac
|
||||
done
|
||||
|
||||
if [ -z ${PATRONI_SCOPE} ]; then
|
||||
PATRONI_SCOPE=$(cat /dev/urandom | LC_ALL=C tr -dc 'a-z0-9' | head -c 8)
|
||||
fi
|
||||
|
||||
function docker_run()
|
||||
{
|
||||
local name=$1
|
||||
shift
|
||||
container=$(docker run -d --name=$name $*)
|
||||
container_ip=$(docker inspect --format '{{ .NetworkSettings.IPAddress }}' ${container})
|
||||
echo "Started container ${name}, ip=${container_ip}"
|
||||
}
|
||||
|
||||
|
||||
ETCD_CONTAINER="${PATRONI_SCOPE}_etcd"
|
||||
docker_run ${ETCD_CONTAINER} ${DOCKER_IMAGE} --etcd
|
||||
|
||||
DOCKER_ARGS="--link=${ETCD_CONTAINER}:${ETCD_CONTAINER} -e PATRONI_SCOPE=${PATRONI_SCOPE} -e PATRONI_ETCD_URL=http://${ETCD_CONTAINER}:2379"
|
||||
PATRONI_ENV=$(sed 's/#.*//g' docker/patroni-secrets.env | sed -n 's/^PATRONI_.*$/-e &/p' | tr '\n' ' ')
|
||||
PATRONI_VOLUME="-v $(dirname $(dirname $(realpath $0)))/patroni:/patroni"
|
||||
|
||||
for i in $(seq 1 "${MEMBERS}"); do
|
||||
container_name=postgres${i}
|
||||
docker_run "${PATRONI_SCOPE}_${container_name}" \
|
||||
$PATRONI_VOLUME \
|
||||
$DOCKER_ARGS \
|
||||
$PATRONI_ENV \
|
||||
-e PATRONI_NAME=${container_name} \
|
||||
${DOCKER_IMAGE}
|
||||
done
|
||||
|
||||
docker_run "${PATRONI_SCOPE}_haproxy" \
|
||||
-p=5000 -p=5001 \
|
||||
$DOCKER_ARGS \
|
||||
${DOCKER_IMAGE} --confd
|
||||
+44
-99
@@ -1,115 +1,60 @@
|
||||
#!/bin/bash
|
||||
#!/bin/sh
|
||||
|
||||
function usage()
|
||||
{
|
||||
cat <<__EOF__
|
||||
Usage: $0
|
||||
if [ -f /a.tar.xz ]; then
|
||||
echo "decompressing image..."
|
||||
sudo tar xpJf /a.tar.xz -C / > /dev/null 2>&1
|
||||
sudo rm /a.tar.xz
|
||||
sudo ln -snf dash /bin/sh
|
||||
fi
|
||||
|
||||
Options:
|
||||
readonly PATRONI_SCOPE=${PATRONI_SCOPE:-batman}
|
||||
PATRONI_NAMESPACE=${PATRONI_NAMESPACE:-/service}
|
||||
readonly PATRONI_NAMESPACE=${PATRONI_NAMESPACE%/}
|
||||
readonly DOCKER_IP=$(hostname --ip-address)
|
||||
|
||||
--etcd Do not run Patroni, run a standalone etcd
|
||||
--confd Do not run Patroni, run a standalone confd
|
||||
--zookeeper Do not run Patroni, run a standalone zookeeper
|
||||
|
||||
Examples:
|
||||
|
||||
$0 --etcd
|
||||
$0 --confd
|
||||
$0 --zookeeper
|
||||
$0
|
||||
__EOF__
|
||||
}
|
||||
|
||||
DOCKER_IP=$(hostname --ip-address)
|
||||
PATRONI_SCOPE=${PATRONI_SCOPE:-batman}
|
||||
ETCD_ARGS="--data-dir /tmp/etcd.data -advertise-client-urls=http://${DOCKER_IP}:2379 -listen-client-urls=http://0.0.0.0:2379"
|
||||
|
||||
optspec=":vh-:"
|
||||
while getopts "$optspec" optchar; do
|
||||
case "${optchar}" in
|
||||
-)
|
||||
case "${OPTARG}" in
|
||||
confd)
|
||||
haproxy -f /etc/haproxy/haproxy.cfg -p /var/run/haproxy.pid -D
|
||||
CONFD="confd -prefix=${PATRONI_NAMESPACE:-/service}/$PATRONI_SCOPE -interval=10 -backend"
|
||||
if [ ! -z ${PATRONI_ZOOKEEPER_HOSTS} ]; then
|
||||
while ! /usr/share/zookeeper/bin/zkCli.sh -server ${PATRONI_ZOOKEEPER_HOSTS} ls /; do
|
||||
sleep 1
|
||||
done
|
||||
exec $CONFD zookeeper -node ${PATRONI_ZOOKEEPER_HOSTS}
|
||||
else
|
||||
while ! curl -s ${PATRONI_ETCD_URL}/v2/members | jq -r '.members[0].clientURLs[0]' | grep -q http; do
|
||||
sleep 1
|
||||
done
|
||||
exec $CONFD etcd -node $PATRONI_ETCD_URL
|
||||
fi
|
||||
;;
|
||||
etcd)
|
||||
exec etcd $ETCD_ARGS
|
||||
;;
|
||||
zookeeper)
|
||||
exec /usr/share/zookeeper/bin/zkServer.sh start-foreground
|
||||
;;
|
||||
cheat)
|
||||
CHEAT=1
|
||||
;;
|
||||
help)
|
||||
usage
|
||||
exit 0
|
||||
;;
|
||||
*)
|
||||
if [ "$OPTERR" = 1 ] && [ "${optspec:0:1}" != ":" ]; then
|
||||
echo "Unknown option --${OPTARG}" >&2
|
||||
fi
|
||||
;;
|
||||
esac;;
|
||||
*)
|
||||
if [ "$OPTERR" != 1 ] || [ "${optspec:0:1}" = ":" ]; then
|
||||
echo "Non-option argument: '-${OPTARG}'" >&2
|
||||
usage
|
||||
exit 1
|
||||
fi
|
||||
;;
|
||||
esac
|
||||
done
|
||||
case "$1" in
|
||||
haproxy)
|
||||
haproxy -f /etc/haproxy/haproxy.cfg -p /var/run/haproxy.pid -D
|
||||
CONFD="confd -prefix=$PATRONI_NAMESPACE/$PATRONI_SCOPE -interval=10 -backend"
|
||||
if [ ! -z "$PATRONI_ZOOKEEPER_HOSTS" ]; then
|
||||
while ! /usr/share/zookeeper/bin/zkCli.sh -server $PATRONI_ZOOKEEPER_HOSTS ls /; do
|
||||
sleep 1
|
||||
done
|
||||
exec dumb-init $CONFD zookeeper -node $PATRONI_ZOOKEEPER_HOSTS
|
||||
else
|
||||
while ! etcdctl cluster-health 2> /dev/null; do
|
||||
sleep 1
|
||||
done
|
||||
exec dumb-init $CONFD etcd -node $(echo $ETCDCTL_ENDPOINTS | sed 's/,/ -node /g')
|
||||
fi
|
||||
;;
|
||||
etcd)
|
||||
exec "$@" -advertise-client-urls http://$DOCKER_IP:2379
|
||||
;;
|
||||
zookeeper)
|
||||
exec /usr/share/zookeeper/bin/zkServer.sh start-foreground
|
||||
;;
|
||||
esac
|
||||
|
||||
## We start an etcd
|
||||
if [[ -z ${PATRONI_ETCD_URL} && -z ${PATRONI_ZOOKEEPER_HOSTS} ]]; then
|
||||
etcd $ETCD_ARGS > /var/log/etcd.log 2> /var/log/etcd.err &
|
||||
if [ -z "$PATRONI_ETCD_HOSTS" ] && [ -z "$PATRONI_ZOOKEEPER_HOSTS" ]; then
|
||||
export PATRONI_ETCD_URL="http://127.0.0.1:2379"
|
||||
etcd --data-dir /tmp/etcd.data -advertise-client-urls=$PATRONI_ETCD_URL -listen-client-urls=http://0.0.0.0:2379 > /var/log/etcd.log 2> /var/log/etcd.err &
|
||||
fi
|
||||
|
||||
export PATRONI_SCOPE
|
||||
export PATRONI_NAME="${PATRONI_NAME:-${HOSTNAME}}"
|
||||
export PATRONI_RESTAPI_CONNECT_ADDRESS="${DOCKER_IP}:8008"
|
||||
export PATRONI_NAMESPACE
|
||||
export PATRONI_NAME="${PATRONI_NAME:-$(hostname)}"
|
||||
export PATRONI_RESTAPI_CONNECT_ADDRESS="$DOCKER_IP:8008"
|
||||
export PATRONI_RESTAPI_LISTEN="0.0.0.0:8008"
|
||||
export PATRONI_admin_PASSWORD="${PATRONI_admin_PASSWORD:=admin}"
|
||||
export PATRONI_admin_PASSWORD="${PATRONI_admin_PASSWORD:-admin}"
|
||||
export PATRONI_admin_OPTIONS="${PATRONI_admin_OPTIONS:-createdb, createrole}"
|
||||
export PATRONI_POSTGRESQL_CONNECT_ADDRESS="${DOCKER_IP}:5432"
|
||||
export PATRONI_POSTGRESQL_CONNECT_ADDRESS="$DOCKER_IP:5432"
|
||||
export PATRONI_POSTGRESQL_LISTEN="0.0.0.0:5432"
|
||||
export PATRONI_POSTGRESQL_DATA_DIR="data/${PATRONI_SCOPE}"
|
||||
export PATRONI_POSTGRESQL_DATA_DIR="${PATRONI_POSTGRESQL_DATA_DIR:-$PGDATA}"
|
||||
export PATRONI_REPLICATION_USERNAME="${PATRONI_REPLICATION_USERNAME:-replicator}"
|
||||
export PATRONI_REPLICATION_PASSWORD="${PATRONI_REPLICATION_PASSWORD:-abcd}"
|
||||
export PATRONI_REPLICATION_PASSWORD="${PATRONI_REPLICATION_PASSWORD:-replicate}"
|
||||
export PATRONI_SUPERUSER_USERNAME="${PATRONI_SUPERUSER_USERNAME:-postgres}"
|
||||
export PATRONI_SUPERUSER_PASSWORD="${PATRONI_SUPERUSER_PASSWORD:-postgres}"
|
||||
export PATRONI_POSTGRESQL_PGPASS="$HOME/.pgpass"
|
||||
|
||||
cat > /patroni.yml <<__EOF__
|
||||
bootstrap:
|
||||
dcs:
|
||||
postgresql:
|
||||
use_pg_rewind: true
|
||||
|
||||
pg_hba:
|
||||
- host all all 0.0.0.0/0 md5
|
||||
- host replication replicator ${DOCKER_IP}/16 md5
|
||||
__EOF__
|
||||
|
||||
mkdir -p "$HOME/.config/patroni"
|
||||
[ -h "$HOME/.config/patroni/patronictl.yaml" ] || ln -s /patroni.yml "$HOME/.config/patroni/patronictl.yaml"
|
||||
|
||||
[ -z $CHEAT ] && exec python3 /patroni.py /patroni.yml
|
||||
|
||||
while true; do
|
||||
sleep 60
|
||||
done
|
||||
exec python3 /patroni.py postgres0.yml
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
ETCD_LISTEN_PEER_URLS=http://0.0.0.0:2380
|
||||
ETCD_LISTEN_CLIENT_URLS=http://0.0.0.0:2379
|
||||
ETCD_INITIAL_CLUSTER=etcd1=http://etcd1:2380,etcd2=http://etcd2:2380,etcd3=http://etcd3:2380
|
||||
ETCD_INITIAL_CLUSTER_STATE=new
|
||||
ETCD_INITIAL_CLUSTER_TOKEN=tutorial
|
||||
@@ -1,3 +1,6 @@
|
||||
PATRONI_SCOPE=demo
|
||||
PATRONI_ETCD_HOSTS='etcd1:2379','etcd2:2379','etcd3:2379'
|
||||
|
||||
PATRONI_RESTAPI_USERNAME=admin
|
||||
PATRONI_RESTAPI_PASSWORD=admin
|
||||
PATRONI_SUPERUSER_USERNAME=postgres
|
||||
@@ -6,3 +9,6 @@ PATRONI_REPLICATION_USERNAME=replicator
|
||||
PATRONI_REPLICATION_PASSWORD=replicate
|
||||
PATRONI_admin_PASSWORD=admin
|
||||
PATRONI_admin_OPTIONS=createdb,createrole
|
||||
|
||||
# for etcdctl
|
||||
ETCDCTL_ENDPOINTS=http://etcd1:2379,http://etcd2:2379,http://etcd3:2379
|
||||
@@ -50,6 +50,7 @@ Etcd
|
||||
- **PATRONI\_ETCD\_PROXY**: proxy url for the etcd. If you are connecting to the etcd using proxy, use this parameter instead of **PATRONI\_ETCD\_URL**
|
||||
- **PATRONI\_ETCD\_URL**: url for the etcd, in format: http(s)://(username:password@)host:port
|
||||
- **PATRONI\_ETCD\_HOSTS**: list of etcd endpoints in format 'host1:port1','host2:port2',etc...
|
||||
- **PATRONI\_ETCD\_USE\_PROXIES**: If this parameter is set to true, Patroni will consider **hosts** as a list of proxies and will not perform a topology discovery of etcd cluster but stick to a fixed list of **hosts**.
|
||||
- **PATRONI\_ETCD\_PROTOCOL**: http or https, if not specified http is used. If the **url** or **proxy** is specified - will take protocol from them.
|
||||
- **PATRONI\_ETCD\_HOST**: the host:port for the etcd endpoint.
|
||||
- **PATRONI\_ETCD\_SRV**: Domain to search the SRV record(s) for cluster autodiscovery.
|
||||
|
||||
+7
-6
@@ -50,8 +50,8 @@ Bootstrap configuration
|
||||
- **slots**: define permanent replication slots. These slots will be preserved during switchover/failover. Patroni will try to create slots before opening connections to the cluster.
|
||||
- **my_slot_name**: the name of replication slot. It is the responsibility of the operator to make sure that there are no clashes in names between replication slots automatically created by Patroni for members and permanent replication slots.
|
||||
- **type**: slot type. Could be ``physical`` or ``logical``. If the slot is logical, you have to additionally define ``database`` and ``plugin``.
|
||||
**database**: the database name where logical slots should be created.
|
||||
**plugin**: the plugin name for the logical slot.
|
||||
- **database**: the database name where logical slots should be created.
|
||||
- **plugin**: the plugin name for the logical slot.
|
||||
- **method**: custom script to use for bootstrapping this cluster.
|
||||
See :ref:`custom bootstrap methods documentation <custom_bootstrap>` for details.
|
||||
When ``initdb`` is specified revert to the default ``initdb`` command. ``initdb`` is also triggered when no ``method``
|
||||
@@ -97,6 +97,7 @@ Most of the parameters are optional, but you have to specify one of the **host**
|
||||
|
||||
- **host**: the host:port for the etcd endpoint.
|
||||
- **hosts**: list of etcd endpoint in format host1:port1,host2:port2,etc... Could be a comma separated string or an actual yaml list.
|
||||
- **use\_proxies**: If this parameter is set to true, Patroni will consider **hosts** as a list of proxies and will not perform a topology discovery of etcd cluster.
|
||||
- **url**: url for the etcd
|
||||
- **proxy**: proxy url for the etcd. If you are connecting to the etcd using proxy, use this parameter instead of **url**
|
||||
- **srv**: Domain to search the SRV record(s) for cluster autodiscovery.
|
||||
@@ -138,10 +139,10 @@ PostgreSQL
|
||||
- **password**: replication password; the user will be created during initialization.
|
||||
- **callbacks**: callback scripts to run on certain actions. Patroni will pass the action, role and cluster name. (See scripts/aws.py as an example of how to write them.)
|
||||
- **on\_reload**: run this script when configuration reload is triggered.
|
||||
- **on\_restart**: run this script when the cluster restarts.
|
||||
- **on\_role\_change**: run this script when the cluster is being promoted or demoted.
|
||||
- **on\_start**: run this script when the cluster starts.
|
||||
- **on\_stop**: run this script when the cluster stops.
|
||||
- **on\_restart**: run this script when the postgres restarts (without changing role).
|
||||
- **on\_role\_change**: run this script when the postgres is being promoted or demoted.
|
||||
- **on\_start**: run this script when the postgres starts.
|
||||
- **on\_stop**: run this script when the postgres stops.
|
||||
- **connect\_address**: IP address + port through which Postgres is accessible from other nodes and applications.
|
||||
- **create\_replica\_methods**: an ordered list of the create methods for turning a Patroni node into a new replica.
|
||||
"basebackup" is the default method; other methods are assumed to refer to scripts, each of which is configured as its
|
||||
|
||||
@@ -3,6 +3,49 @@
|
||||
Release notes
|
||||
=============
|
||||
|
||||
Version 1.5.6
|
||||
-------------
|
||||
|
||||
**New features**
|
||||
|
||||
- Support work with etcd cluster via set of proxies (Alexander Kukushkin)
|
||||
|
||||
It might happen that etcd cluster is not accessible directly but via set of proxies. In this case Patroni will not perform etcd topology discovery but just round-robin via proxy hosts. Behavior is controlled by `etcd.use_proxies`.
|
||||
|
||||
- Changed callbacks behavior when role on the node is changed (Alexander)
|
||||
|
||||
If the role was changed from `master` or `standby_leader` to `replica` or from `replica` to `standby_leader`, `on_restart` callback will not be called anymore in favor of `on_role_change` callback.
|
||||
|
||||
- Change the way how we start postgres (Alexander)
|
||||
|
||||
Use `multiprocessing.Process` instead of executing itself and `multiprocessing.Pipe` to transmit the postmaster pid to the Patroni process. Before that we were using pipes, what was leaving postmaster process with stdin closed.
|
||||
|
||||
**Bug fixes**
|
||||
|
||||
- Fix role returned by REST API for the standby leader (Alexander)
|
||||
|
||||
It was incorrectly returning `replica` instead of `standby_leader`
|
||||
|
||||
- Wait for callback end if it could not be killed (Julien Tachoires)
|
||||
|
||||
Patroni doesn't have enough privileges to terminate the callback script running under `sudo` what was cancelling the new callback. If the running script could not be killed, Patroni will wait until it finishes and then run the next callback.
|
||||
|
||||
- Reduce lock time taken by dcs.get_cluster method (Alexander)
|
||||
|
||||
Due to the lock being held DCS slowness was affecting the REST API health checks causing false positives.
|
||||
|
||||
- Improve cleaning of PGDATA when `pg_wal`/`pg_xlog` is a symlink (Julien)
|
||||
|
||||
In this case Patroni will explicitly remove files from the target directory.
|
||||
|
||||
- Remove unnecessary usage of os.path.relpath (Ants Aasma)
|
||||
|
||||
It depends on being able to resolve the working directory, what will fail if Patroni is started in a directory that is later unlinked from the filesystem.
|
||||
|
||||
- Do not enforce ssl version when communicating with Etcd (Alexander)
|
||||
|
||||
For some unknown reason python3-etcd on debian and ubuntu are not based on the latest version of the package and therefore it enforces TLSv1 which is not supported by Etcd v3. We solved this problem on Patroni side.
|
||||
|
||||
Version 1.5.5
|
||||
-------------
|
||||
|
||||
|
||||
@@ -111,9 +111,9 @@ class PatroniController(AbstractController):
|
||||
with open(os.path.join(self._data_dir, 'label'), 'w') as f:
|
||||
f.write(content)
|
||||
|
||||
def read_label(self):
|
||||
def read_label(self, label):
|
||||
try:
|
||||
with open(os.path.join(self._data_dir, 'label'), 'r') as f:
|
||||
with open(os.path.join(self._data_dir, label), 'r') as f:
|
||||
return f.read().strip()
|
||||
except IOError:
|
||||
return None
|
||||
|
||||
@@ -2,11 +2,11 @@ Feature: standby cluster
|
||||
Scenario: check permanent logical slots are preserved on failover/switchover
|
||||
Given I start postgres1
|
||||
Then postgres1 is a leader after 10 seconds
|
||||
And I sleep for 2 seconds
|
||||
And I sleep for 3 seconds
|
||||
When I issue a PATCH request to http://127.0.0.1:8009/config with {"loop_wait": 2, "slots": {"pm_1": {"type": "physical"}}, "postgresql": {"parameters": {"wal_level": "logical"}}}
|
||||
Then I receive a response code 200
|
||||
And Response on GET http://127.0.0.1:8009/config contains slots after 10 seconds
|
||||
And I sleep for 2 seconds
|
||||
And I sleep for 3 seconds
|
||||
When I issue a PATCH request to http://127.0.0.1:8009/config with {"slots": {"test_logical": {"type": "logical", "database": "postgres", "plugin": "test_decoding"}}}
|
||||
Then I receive a response code 200
|
||||
When I start postgres0 with callback configured
|
||||
@@ -14,7 +14,7 @@ Feature: standby cluster
|
||||
And replication works from postgres1 to postgres0 after 15 seconds
|
||||
When I shut down postgres1
|
||||
Then postgres0 is a leader after 10 seconds
|
||||
And I sleep for 2 seconds
|
||||
And "members/postgres0" key in DCS has role=master after 3 seconds
|
||||
When I issue a GET request to http://127.0.0.1:8008/
|
||||
Then I receive a response code 200
|
||||
And there is a label with "test_logical" in postgres0 data directory
|
||||
@@ -24,6 +24,12 @@ Feature: standby cluster
|
||||
Then postgres1 is a leader of batman1 after 10 seconds
|
||||
When I add the table foo to postgres0
|
||||
Then table foo is present on postgres1 after 20 seconds
|
||||
When I issue a GET request to http://127.0.0.1:8009/master
|
||||
Then I receive a response code 200
|
||||
When I issue a GET request to http://127.0.0.1:8009/standby_leader
|
||||
Then I receive a response code 200
|
||||
And I receive a response role standby_leader
|
||||
And there is a postgres1_cb.log with "on_start replica batman1\non_role_change standby_leader batman1" in postgres1 data directory
|
||||
When I start postgres2 in a cluster batman1
|
||||
Then postgres2 role is the replica after 24 seconds
|
||||
And table foo is present on postgres2 after 20 seconds
|
||||
@@ -31,4 +37,11 @@ Feature: standby cluster
|
||||
Scenario: check failover
|
||||
When I kill postgres1
|
||||
And I kill postmaster on postgres1
|
||||
Then postgres2 is replicating from postgres0 after 20 seconds
|
||||
Then postgres2 is replicating from postgres0 after 32 seconds
|
||||
When I issue a GET request to http://127.0.0.1:8010/master
|
||||
Then I receive a response code 200
|
||||
When I issue a GET request to http://127.0.0.1:8010/standby_leader
|
||||
Then I receive a response code 200
|
||||
And I receive a response role standby_leader
|
||||
And replication works from postgres0 to postgres2 after 15 seconds
|
||||
And there is a postgres2_cb.log with "on_start replica batman1\non_role_change standby_leader batman1" in postgres2 data directory
|
||||
|
||||
@@ -9,9 +9,10 @@ def start_patroni_with_a_name_value_tag(context, name, tag_name, tag_value):
|
||||
return context.pctl.start(name, custom_config={'tags': {tag_name: tag_value}})
|
||||
|
||||
|
||||
@then('There is a label with "{content:w}" in {name:w} data directory')
|
||||
def check_label(context, content, name):
|
||||
label = context.pctl.read_label(name)
|
||||
@then('There is a {label} with "{content}" in {name:w} data directory')
|
||||
def check_label(context, label, content, name):
|
||||
label = context.pctl.read_label(name, label)
|
||||
label = label.replace('\n', '\\n')
|
||||
assert label == content, "{0} is not equal to {1}".format(label, content)
|
||||
|
||||
|
||||
|
||||
@@ -9,6 +9,8 @@ SELECT * FROM pg_catalog.pg_stat_replication
|
||||
WHERE application_name = '{0}'
|
||||
"""
|
||||
|
||||
callback = "bash -c 'echo \"${*: -3:1} ${*: -2:1} ${*: -1:1}\" >> data/$1/$1_cb.log' -- "
|
||||
|
||||
|
||||
@step('I start {name:w} with callback configured')
|
||||
def start_patroni_with_callbacks(context, name):
|
||||
@@ -24,12 +26,15 @@ def start_patroni_with_callbacks(context, name):
|
||||
@step('I start {name:w} in a cluster {cluster_name:w}')
|
||||
def start_patroni(context, name, cluster_name):
|
||||
return context.pctl.start(name, custom_config={
|
||||
"scope": cluster_name
|
||||
"scope": cluster_name,
|
||||
"postgresql": {
|
||||
"callbacks": {c: callback + name for c in ('on_start', 'on_stop', 'on_restart', 'on_role_change')}
|
||||
}
|
||||
})
|
||||
|
||||
|
||||
@step('I start {name:w} in a standby cluster {cluster_name:w} as a clone of {name2:w}')
|
||||
def start_patroni_stanby_cluster(context, name, cluster_name, name2):
|
||||
def start_patroni_standby_cluster(context, name, cluster_name, name2):
|
||||
# we need to remove patroni.dynamic.json in order to "bootstrap" standby cluster with existing PGDATA
|
||||
os.unlink(os.path.join(context.pctl._processes[name]._data_dir, 'patroni.dynamic.json'))
|
||||
port = context.pctl._processes[name2]._connkwargs.get('port')
|
||||
@@ -37,12 +42,18 @@ def start_patroni_stanby_cluster(context, name, cluster_name, name2):
|
||||
"scope": cluster_name,
|
||||
"bootstrap": {
|
||||
"dcs": {
|
||||
"ttl": 20,
|
||||
"loop_wait": 2,
|
||||
"retry_timeout": 5,
|
||||
"standby_cluster": {
|
||||
"host": "localhost",
|
||||
"port": port,
|
||||
"primary_slot_name": "pm_1",
|
||||
}
|
||||
}
|
||||
},
|
||||
"postgresql": {
|
||||
"callbacks": {c: callback + name for c in ('on_start', 'on_stop', 'on_restart', 'on_role_change')}
|
||||
}
|
||||
})
|
||||
return context.pctl.start(name)
|
||||
@@ -60,8 +71,8 @@ def check_replication_status(context, pg_name1, pg_name2, timeout):
|
||||
)
|
||||
|
||||
if cur and len(cur.fetchall()) != 0:
|
||||
return True
|
||||
break
|
||||
|
||||
time.sleep(1)
|
||||
|
||||
return False
|
||||
else:
|
||||
assert False, "{0} is not replicating from {1} after {2} seconds".format(pg_name1, pg_name2, timeout)
|
||||
|
||||
+23
-33
@@ -85,9 +85,10 @@ class Patroni(object):
|
||||
self._received_sighup = True
|
||||
|
||||
def sigterm_handler(self, *args):
|
||||
if not self._received_sigterm:
|
||||
self._received_sigterm = True
|
||||
sys.exit()
|
||||
with self._sigterm_lock:
|
||||
if not self._received_sigterm:
|
||||
self._received_sigterm = True
|
||||
sys.exit()
|
||||
|
||||
@property
|
||||
def noloadbalance(self):
|
||||
@@ -106,11 +107,16 @@ class Patroni(object):
|
||||
elif self.ha.watch(nap_time):
|
||||
self.next_run = time.time()
|
||||
|
||||
@property
|
||||
def received_sigterm(self):
|
||||
with self._sigterm_lock:
|
||||
return self._received_sigterm
|
||||
|
||||
def run(self):
|
||||
self.api.start()
|
||||
self.next_run = time.time()
|
||||
|
||||
while not self._received_sigterm:
|
||||
while not self.received_sigterm:
|
||||
if self._received_sighup:
|
||||
self._received_sighup = False
|
||||
if self.config.reload_local_configuration():
|
||||
@@ -128,13 +134,18 @@ class Patroni(object):
|
||||
self.schedule_next_run()
|
||||
|
||||
def setup_signal_handlers(self):
|
||||
from threading import Lock
|
||||
|
||||
self._received_sighup = False
|
||||
self._sigterm_lock = Lock()
|
||||
self._received_sigterm = False
|
||||
if os.name != 'nt':
|
||||
signal.signal(signal.SIGHUP, self.sighup_handler)
|
||||
signal.signal(signal.SIGTERM, self.sigterm_handler)
|
||||
|
||||
def shutdown(self):
|
||||
with self._sigterm_lock:
|
||||
self._received_sigterm = True
|
||||
try:
|
||||
self.api.shutdown()
|
||||
except Exception:
|
||||
@@ -153,34 +164,11 @@ def patroni_main():
|
||||
logging.shutdown()
|
||||
|
||||
|
||||
def pg_ctl_start(args):
|
||||
import subprocess
|
||||
if os.name != 'nt':
|
||||
os.setsid()
|
||||
postmaster = subprocess.Popen(args)
|
||||
print(postmaster.pid)
|
||||
|
||||
|
||||
def call_self(args, **kwargs):
|
||||
"""This function executes Patroni once again with provided arguments.
|
||||
|
||||
:args: list of arguments to call Patroni with.
|
||||
:returns: `Popen` object"""
|
||||
|
||||
exe = [sys.executable]
|
||||
if not getattr(sys, 'frozen', False): # Binary distribution?
|
||||
exe.append(sys.argv[0])
|
||||
|
||||
import subprocess
|
||||
return subprocess.Popen(exe + args, **kwargs)
|
||||
|
||||
|
||||
def main():
|
||||
if os.getpid() != 1:
|
||||
if len(sys.argv) > 5 and sys.argv[1] == 'pg_ctl_start':
|
||||
return pg_ctl_start(sys.argv[2:])
|
||||
return patroni_main()
|
||||
|
||||
# Patroni started with PID=1, it looks like we are in the container
|
||||
pid = 0
|
||||
|
||||
# Looks like we are in a docker, so we will act like init
|
||||
@@ -199,16 +187,18 @@ def main():
|
||||
if pid:
|
||||
os.kill(pid, signo)
|
||||
|
||||
signal.signal(signal.SIGCHLD, sigchld_handler)
|
||||
if os.name != 'nt':
|
||||
signal.signal(signal.SIGCHLD, sigchld_handler)
|
||||
signal.signal(signal.SIGHUP, passtochild)
|
||||
signal.signal(signal.SIGQUIT, passtochild)
|
||||
signal.signal(signal.SIGUSR1, passtochild)
|
||||
signal.signal(signal.SIGUSR2, passtochild)
|
||||
signal.signal(signal.SIGINT, passtochild)
|
||||
signal.signal(signal.SIGUSR1, passtochild)
|
||||
signal.signal(signal.SIGUSR2, passtochild)
|
||||
signal.signal(signal.SIGABRT, passtochild)
|
||||
signal.signal(signal.SIGTERM, passtochild)
|
||||
|
||||
patroni = call_self(sys.argv[1:])
|
||||
import multiprocessing
|
||||
patroni = multiprocessing.Process(target=patroni_main)
|
||||
patroni.start()
|
||||
pid = patroni.pid
|
||||
patroni.wait()
|
||||
patroni.join()
|
||||
|
||||
@@ -446,6 +446,9 @@ class RestApiHandler(BaseHTTPRequestHandler):
|
||||
})
|
||||
}
|
||||
|
||||
if result['role'] == 'replica' and self.server.patroni.ha.is_standby_cluster():
|
||||
result['role'] = self.server.patroni.postgresql.role
|
||||
|
||||
if row[1] > 0:
|
||||
result['timeline'] = row[1]
|
||||
else:
|
||||
|
||||
@@ -19,8 +19,13 @@ class CallbackExecutor(Thread):
|
||||
def call(self, cmd):
|
||||
with self._lock:
|
||||
if self._process and self._process.poll() is None:
|
||||
self._process.kill()
|
||||
logger.warning('Killed the old callback process because it was still running: %s', self._cmd)
|
||||
try:
|
||||
self._process.kill()
|
||||
logger.warning('Killed the old callback process because it was still running: %s', self._cmd)
|
||||
except OSError:
|
||||
logger.exception('Failed to kill the old callback')
|
||||
logger.warning('Wait until callback end')
|
||||
self._process.wait()
|
||||
self._cmd = cmd
|
||||
self._callback_event.set()
|
||||
|
||||
|
||||
+6
-6
@@ -233,7 +233,7 @@ class Config(object):
|
||||
ret[section][param] = value
|
||||
|
||||
_set_section_values('restapi', ['listen', 'connect_address', 'certfile', 'keyfile'])
|
||||
_set_section_values('postgresql', ['listen', 'connect_address', 'data_dir', 'pgpass', 'bin_dir'])
|
||||
_set_section_values('postgresql', ['listen', 'connect_address', 'config_dir', 'data_dir', 'pgpass', 'bin_dir'])
|
||||
_set_section_values('log', ['level', 'format', 'dateformat', 'dir', 'file_size', 'file_num', 'loggers'])
|
||||
|
||||
def _parse_dict(value):
|
||||
@@ -285,10 +285,10 @@ class Config(object):
|
||||
if param.startswith(Config.PATRONI_ENV_PREFIX):
|
||||
# PATRONI_(ETCD|CONSUL|ZOOKEEPER|EXHIBITOR|...)_(HOSTS?|PORT|..)
|
||||
name, suffix = (param[8:].split('_', 1) + [''])[:2]
|
||||
if suffix in ('HOST', 'HOSTS', 'PORT', 'PROTOCOL', 'SRV', 'URL', 'PROXY', 'CACERT', 'CERT', 'KEY',
|
||||
'VERIFY', 'TOKEN', 'CHECKS', 'DC', 'REGISTER_SERVICE', 'SERVICE_CHECK_INTERVAL',
|
||||
'NAMESPACE', 'CONTEXT', 'USE_ENDPOINTS', 'SCOPE_LABEL', 'ROLE_LABEL', 'POD_IP',
|
||||
'PORTS', 'LABELS') and name:
|
||||
if suffix in ('HOST', 'HOSTS', 'PORT', 'USE_PROXIES', 'PROTOCOL', 'SRV', 'URL', 'PROXY',
|
||||
'CACERT', 'CERT', 'KEY', 'VERIFY', 'TOKEN', 'CHECKS', 'DC', 'REGISTER_SERVICE',
|
||||
'SERVICE_CHECK_INTERVAL', 'NAMESPACE', 'CONTEXT', 'USE_ENDPOINTS', 'SCOPE_LABEL',
|
||||
'ROLE_LABEL', 'POD_IP', 'PORTS', 'LABELS') and name:
|
||||
value = os.environ.pop(param)
|
||||
if suffix == 'PORT':
|
||||
value = value and parse_int(value)
|
||||
@@ -296,7 +296,7 @@ class Config(object):
|
||||
value = value and _parse_list(value)
|
||||
elif suffix == 'LABELS':
|
||||
value = _parse_dict(value)
|
||||
elif suffix == 'REGISTER_SERVICE':
|
||||
elif suffix in ('USE_PROXIES', 'REGISTER_SERVICE'):
|
||||
value = parse_bool(value)
|
||||
if value:
|
||||
ret[name.lower()][suffix.lower()] = value
|
||||
|
||||
+2
-2
@@ -63,8 +63,8 @@ def parse_dcs(dcs):
|
||||
elif scheme not in DCS_DEFAULTS:
|
||||
raise PatroniCtlException('Unknown dcs scheme: {}'.format(scheme))
|
||||
|
||||
dcs_info = DCS_DEFAULTS[scheme]
|
||||
return yaml.load(dcs_info['template'].format(host=parsed.hostname or 'localhost', port=port or dcs_info['port']))
|
||||
default = DCS_DEFAULTS[scheme]
|
||||
return yaml.safe_load(default['template'].format(host=parsed.hostname or 'localhost', port=port or default['port']))
|
||||
|
||||
|
||||
def load_config(path, dcs):
|
||||
|
||||
+18
-8
@@ -9,6 +9,7 @@ import pkgutil
|
||||
import re
|
||||
import six
|
||||
import sys
|
||||
import time
|
||||
|
||||
from collections import defaultdict, namedtuple
|
||||
from copy import deepcopy
|
||||
@@ -529,6 +530,7 @@ class AbstractDCS(object):
|
||||
|
||||
self._ctl = bool(config.get('patronictl', False))
|
||||
self._cluster = None
|
||||
self._cluster_valid_till = 0
|
||||
self._cluster_thread_lock = Lock()
|
||||
self._last_leader_operation = ''
|
||||
self.event = Event()
|
||||
@@ -576,6 +578,10 @@ class AbstractDCS(object):
|
||||
def set_ttl(self, ttl):
|
||||
"""Set the new ttl value for leader key"""
|
||||
|
||||
@abc.abstractmethod
|
||||
def ttl(self):
|
||||
"""Get new ttl value"""
|
||||
|
||||
@abc.abstractmethod
|
||||
def set_retry_timeout(self, retry_timeout):
|
||||
"""Set the new value for retry_timeout"""
|
||||
@@ -603,22 +609,26 @@ class AbstractDCS(object):
|
||||
instance would be demoted."""
|
||||
|
||||
def get_cluster(self):
|
||||
try:
|
||||
cluster = self._load_cluster()
|
||||
except Exception:
|
||||
self.reset_cluster()
|
||||
raise
|
||||
|
||||
with self._cluster_thread_lock:
|
||||
try:
|
||||
self._load_cluster()
|
||||
except Exception:
|
||||
self._cluster = None
|
||||
raise
|
||||
return self._cluster
|
||||
self._cluster = cluster
|
||||
self._cluster_valid_till = time.time() + self.ttl
|
||||
return cluster
|
||||
|
||||
@property
|
||||
def cluster(self):
|
||||
with self._cluster_thread_lock:
|
||||
return self._cluster
|
||||
return self._cluster if self._cluster_valid_till > time.time() else None
|
||||
|
||||
def reset_cluster(self):
|
||||
with self._cluster_thread_lock:
|
||||
self._cluster = None
|
||||
self._cluster_valid_till = 0
|
||||
|
||||
@abc.abstractmethod
|
||||
def _write_leader_optime(self, last_operation):
|
||||
@@ -684,7 +694,7 @@ class AbstractDCS(object):
|
||||
"""Create or update `/config` key"""
|
||||
|
||||
@abc.abstractmethod
|
||||
def touch_member(self, data, ttl=None, permanent=False):
|
||||
def touch_member(self, data, permanent=False):
|
||||
"""Update member key in DCS.
|
||||
This method should create or update key with the name = '/members/' + `~self._name`
|
||||
and value = data in a given DCS.
|
||||
|
||||
@@ -237,6 +237,10 @@ class Consul(AbstractDCS):
|
||||
self._session = None
|
||||
self.__do_not_watch = True
|
||||
|
||||
@property
|
||||
def ttl(self):
|
||||
return self._client.http.ttl
|
||||
|
||||
def set_retry_timeout(self, retry_timeout):
|
||||
self._retry.deadline = retry_timeout
|
||||
self._client.http.set_read_timeout(retry_timeout)
|
||||
@@ -299,7 +303,7 @@ class Consul(AbstractDCS):
|
||||
nodes = {}
|
||||
for node in results:
|
||||
node['Value'] = (node['Value'] or b'').decode('utf-8')
|
||||
nodes[os.path.relpath(node['Key'], path).replace('\\', '/')] = node
|
||||
nodes[node['Key'][len(path):].lstrip('/')] = node
|
||||
|
||||
# get initialize flag
|
||||
initialize = nodes.get(self._INITIALIZE)
|
||||
@@ -342,15 +346,15 @@ class Consul(AbstractDCS):
|
||||
sync = nodes.get(self._SYNC)
|
||||
sync = SyncState.from_node(sync and sync['ModifyIndex'], sync and sync['Value'])
|
||||
|
||||
self._cluster = Cluster(initialize, config, leader, last_leader_operation, members, failover, sync, history)
|
||||
return Cluster(initialize, config, leader, last_leader_operation, members, failover, sync, history)
|
||||
except NotFound:
|
||||
self._cluster = Cluster(None, None, None, None, [], None, None, None)
|
||||
return Cluster(None, None, None, None, [], None, None, None)
|
||||
except Exception:
|
||||
logger.exception('get_cluster')
|
||||
raise ConsulError('Consul is not responding properly')
|
||||
|
||||
@catch_consul_errors
|
||||
def touch_member(self, data, ttl=None, permanent=False):
|
||||
def touch_member(self, data, permanent=False):
|
||||
cluster = self.cluster
|
||||
member = cluster and cluster.get_member(self._name, fallback_to_leader=False)
|
||||
create_member = not permanent and self.refresh_session()
|
||||
@@ -503,6 +507,7 @@ class Consul(AbstractDCS):
|
||||
return self.retry(self._client.kv.delete, self.sync_path, cas=index)
|
||||
|
||||
def watch(self, leader_index, timeout):
|
||||
self._last_session_refresh = 0
|
||||
if self.__do_not_watch:
|
||||
self.__do_not_watch = False
|
||||
return True
|
||||
@@ -521,5 +526,4 @@ class Consul(AbstractDCS):
|
||||
try:
|
||||
return super(Consul, self).watch(None, timeout)
|
||||
finally:
|
||||
self._last_session_refresh = 0
|
||||
self.event.clear()
|
||||
|
||||
+47
-32
@@ -87,9 +87,12 @@ class Client(etcd.Client):
|
||||
args = {p: config.get(p) for p in ('host', 'port', 'protocol', 'use_proxies', 'username', 'password',
|
||||
'cert', 'ca_cert') if config.get(p)}
|
||||
super(Client, self).__init__(read_timeout=config['retry_timeout'], **args)
|
||||
# For some reason python3-etcd on debian and ubuntu are not based on the latest version
|
||||
# Workaround for the case when https://github.com/jplana/python-etcd/pull/196 is not applied
|
||||
self.http.connection_pool_kw.pop('ssl_version', None)
|
||||
self._config = config
|
||||
self._load_machines_cache()
|
||||
self._allow_reconnect = not self._use_proxies
|
||||
self._allow_reconnect = True
|
||||
|
||||
def _build_request_parameters(self):
|
||||
kwargs = {'headers': self._get_headers(), 'redirect': self.allow_redirect}
|
||||
@@ -128,8 +131,11 @@ class Client(etcd.Client):
|
||||
while True:
|
||||
try:
|
||||
response = self.http.request(self._MGET, self._base_uri + self.version_prefix + '/machines', **kwargs)
|
||||
machines = [n.strip() for n in self._handle_server_response(response).data.decode('utf-8').split(',')]
|
||||
data = self._handle_server_response(response).data.decode('utf-8')
|
||||
machines = [m.strip() for m in data.split(',') if m.strip()]
|
||||
logger.debug("Retrieved list of machines: %s", machines)
|
||||
if not machines:
|
||||
raise etcd.EtcdException
|
||||
random.shuffle(machines)
|
||||
for url in machines:
|
||||
r = urlparse(url)
|
||||
@@ -186,10 +192,8 @@ class Client(etcd.Client):
|
||||
# Update machines_cache if previous attempt of update has failed
|
||||
if self._update_machines_cache:
|
||||
self._load_machines_cache()
|
||||
elif time.time() - self._machines_cache_updated > self._machines_cache_ttl:
|
||||
self._machines_cache = self.machines
|
||||
if self._base_uri in self._machines_cache:
|
||||
self._machines_cache.remove(self._base_uri)
|
||||
elif not self._use_proxies and time.time() - self._machines_cache_updated > self._machines_cache_ttl:
|
||||
self._refresh_machines_cache()
|
||||
self._machines_cache_updated = time.time()
|
||||
|
||||
kwargs.update(self._build_request_parameters())
|
||||
@@ -206,10 +210,8 @@ class Client(etcd.Client):
|
||||
|
||||
if response is False:
|
||||
some_request_failed = True
|
||||
if some_request_failed and not self._use_proxies:
|
||||
self._machines_cache = self.machines
|
||||
if self._base_uri in self._machines_cache:
|
||||
self._machines_cache.remove(self._base_uri)
|
||||
if some_request_failed:
|
||||
self._refresh_machines_cache()
|
||||
except etcd.EtcdConnectionFailed as e:
|
||||
if isinstance(e, etcd.EtcdWatchTimedOut) and self._machines_cache:
|
||||
self._base_uri = self._next_server()
|
||||
@@ -268,6 +270,21 @@ class Client(etcd.Client):
|
||||
return list(set(ret))
|
||||
return [uri(self.protocol, host, port)]
|
||||
|
||||
def _get_machines_cache_from_config(self):
|
||||
if 'proxy' in self._config:
|
||||
return [uri(self.protocol, self._config['host'], self._config['port'])]
|
||||
|
||||
machines_cache = []
|
||||
if 'srv' in self._config:
|
||||
machines_cache = self._get_machines_cache_from_srv(self._config['srv'])
|
||||
|
||||
if not machines_cache and 'hosts' in self._config:
|
||||
machines_cache = list(self._config['hosts'])
|
||||
|
||||
if not machines_cache and 'host' in self._config:
|
||||
machines_cache = self._get_machines_cache_from_dns(self._config['host'], self._config['port'])
|
||||
return machines_cache
|
||||
|
||||
def _load_machines_cache(self):
|
||||
"""This method should fill up `_machines_cache` from scratch.
|
||||
It could happen only in two cases:
|
||||
@@ -279,19 +296,7 @@ class Client(etcd.Client):
|
||||
if 'srv' not in self._config and 'host' not in self._config and 'hosts' not in self._config:
|
||||
raise Exception('Neither srv, hosts, host nor url are defined in etcd section of config')
|
||||
|
||||
if self._use_proxies:
|
||||
self._machines_cache = [uri(self.protocol, self._config['host'], self._config['port'])]
|
||||
else:
|
||||
self._machines_cache = []
|
||||
|
||||
if 'srv' in self._config:
|
||||
self._machines_cache = self._get_machines_cache_from_srv(self._config['srv'])
|
||||
|
||||
if not self._machines_cache and 'hosts' in self._config:
|
||||
self._machines_cache = list(self._config['hosts'])
|
||||
|
||||
if not self._machines_cache and 'host' in self._config:
|
||||
self._machines_cache = self._get_machines_cache_from_dns(self._config['host'], self._config['port'])
|
||||
self._machines_cache = self._get_machines_cache_from_config()
|
||||
|
||||
# Can not bootstrap list of etcd-cluster members, giving up
|
||||
if not self._machines_cache:
|
||||
@@ -299,14 +304,16 @@ class Client(etcd.Client):
|
||||
|
||||
# After filling up initial list of machines_cache we should ask etcd-cluster about actual list
|
||||
self._base_uri = self._next_server()
|
||||
self._machines_cache = self.machines
|
||||
|
||||
if self._base_uri in self._machines_cache:
|
||||
self._machines_cache.remove(self._base_uri)
|
||||
self._refresh_machines_cache()
|
||||
|
||||
self._update_machines_cache = False
|
||||
self._machines_cache_updated = time.time()
|
||||
|
||||
def _refresh_machines_cache(self):
|
||||
self._machines_cache = self._get_machines_cache_from_config() if self._use_proxies else self.machines
|
||||
if self._base_uri in self._machines_cache:
|
||||
self._machines_cache.remove(self._base_uri)
|
||||
|
||||
|
||||
class Etcd(AbstractDCS):
|
||||
|
||||
@@ -426,6 +433,8 @@ class Etcd(AbstractDCS):
|
||||
while not client:
|
||||
try:
|
||||
client = Client(config, dns_resolver)
|
||||
if 'use_proxies' in config and not client.machines:
|
||||
raise etcd.EtcdException
|
||||
except etcd.EtcdException:
|
||||
logger.info('waiting on etcd')
|
||||
time.sleep(5)
|
||||
@@ -437,6 +446,10 @@ class Etcd(AbstractDCS):
|
||||
self._ttl = ttl
|
||||
self._client.set_machines_cache_ttl(ttl*10)
|
||||
|
||||
@property
|
||||
def ttl(self):
|
||||
return self._ttl
|
||||
|
||||
def set_retry_timeout(self, retry_timeout):
|
||||
self._retry.deadline = retry_timeout
|
||||
self._client.set_read_timeout(retry_timeout)
|
||||
@@ -446,9 +459,10 @@ class Etcd(AbstractDCS):
|
||||
return Member.from_node(node.modifiedIndex, os.path.basename(node.key), node.ttl, node.value)
|
||||
|
||||
def _load_cluster(self):
|
||||
cluster = None
|
||||
try:
|
||||
result = self.retry(self._client.read, self.client_path(''), recursive=True)
|
||||
nodes = {os.path.relpath(node.key, result.key).replace('\\', '/'): node for node in result.leaves}
|
||||
nodes = {node.key[len(result.key):].lstrip('/'): node for node in result.leaves}
|
||||
|
||||
# get initialize flag
|
||||
initialize = nodes.get(self._INITIALIZE)
|
||||
@@ -486,17 +500,18 @@ class Etcd(AbstractDCS):
|
||||
sync = nodes.get(self._SYNC)
|
||||
sync = SyncState.from_node(sync and sync.modifiedIndex, sync and sync.value)
|
||||
|
||||
self._cluster = Cluster(initialize, config, leader, last_leader_operation, members, failover, sync, history)
|
||||
cluster = Cluster(initialize, config, leader, last_leader_operation, members, failover, sync, history)
|
||||
except etcd.EtcdKeyNotFound:
|
||||
self._cluster = Cluster(None, None, None, None, [], None, None, None)
|
||||
cluster = Cluster(None, None, None, None, [], None, None, None)
|
||||
except Exception as e:
|
||||
self._handle_exception(e, 'get_cluster', raise_ex=EtcdError('Etcd is not responding properly'))
|
||||
self._has_failed = False
|
||||
return cluster
|
||||
|
||||
@catch_etcd_errors
|
||||
def touch_member(self, data, ttl=None, permanent=False):
|
||||
def touch_member(self, data, permanent=False):
|
||||
data = json.dumps(data, separators=(',', ':'))
|
||||
return self._client.set(self.member_path, data, None if permanent else ttl or self._ttl)
|
||||
return self._client.set(self.member_path, data, None if permanent else self._ttl)
|
||||
|
||||
@catch_etcd_errors
|
||||
def take_leader(self):
|
||||
|
||||
@@ -107,6 +107,7 @@ class Kubernetes(AbstractDCS):
|
||||
self._leader_observed_time = None
|
||||
self._leader_resource_version = None
|
||||
self._leader_observed_subsets = []
|
||||
self._config_resource_version = None
|
||||
self.__do_not_watch = False
|
||||
|
||||
def retry(self, *args, **kwargs):
|
||||
@@ -124,6 +125,10 @@ class Kubernetes(AbstractDCS):
|
||||
self.__do_not_watch = self._ttl != ttl
|
||||
self._ttl = ttl
|
||||
|
||||
@property
|
||||
def ttl(self):
|
||||
return self._ttl
|
||||
|
||||
def set_retry_timeout(self, retry_timeout):
|
||||
self._retry.deadline = retry_timeout
|
||||
self._api.set_timeout(retry_timeout)
|
||||
@@ -146,6 +151,7 @@ class Kubernetes(AbstractDCS):
|
||||
|
||||
config = nodes.get(self.config_path)
|
||||
metadata = config and config.metadata
|
||||
self._config_resource_version = metadata.resource_version if metadata else None
|
||||
annotations = metadata and metadata.annotations or {}
|
||||
|
||||
# get initialize flag
|
||||
@@ -201,7 +207,7 @@ class Kubernetes(AbstractDCS):
|
||||
metadata = sync and sync.metadata
|
||||
sync = SyncState.from_node(metadata and metadata.resource_version, metadata and metadata.annotations)
|
||||
|
||||
self._cluster = Cluster(initialize, config, leader, last_leader_operation, members, failover, sync, history)
|
||||
return Cluster(initialize, config, leader, last_leader_operation, members, failover, sync, history)
|
||||
except Exception:
|
||||
logger.exception('get_cluster')
|
||||
raise KubernetesError('Kubernetes API is not responding properly')
|
||||
@@ -280,7 +286,10 @@ class Kubernetes(AbstractDCS):
|
||||
if self.__subsets and not patch and not resource_version:
|
||||
self._should_create_config_service = True
|
||||
self._create_config_service()
|
||||
return self.patch_or_create(self.config_path, annotations, resource_version, patch, retry)
|
||||
ret = self.patch_or_create(self.config_path, annotations, resource_version, patch, retry)
|
||||
if ret:
|
||||
self._config_resource_version = ret.metadata.resource_version
|
||||
return ret
|
||||
|
||||
def _create_config_service(self):
|
||||
metadata = k8s_client.V1ObjectMeta(namespace=self._namespace, name=self.config_path, labels=self._labels)
|
||||
@@ -349,11 +358,10 @@ class Kubernetes(AbstractDCS):
|
||||
return self.patch_or_create(self.failover_path, annotations, index, bool(index or patch), False)
|
||||
|
||||
def set_config_value(self, value, index=None):
|
||||
patch = bool(index or self.cluster and self.cluster.config and self.cluster.config.index)
|
||||
return self.patch_or_create_config({self._CONFIG: value}, index, patch, False)
|
||||
return self.patch_or_create_config({self._CONFIG: value}, index, bool(self._config_resource_version), False)
|
||||
|
||||
@catch_kubernetes_errors
|
||||
def touch_member(self, data, ttl=None, permanent=False):
|
||||
def touch_member(self, data, permanent=False):
|
||||
cluster = self.cluster
|
||||
if cluster and cluster.leader and cluster.leader.name == self._name:
|
||||
role = 'promoted' if data['role'] in ('replica', 'promoted') else 'master'
|
||||
@@ -386,15 +394,14 @@ class Kubernetes(AbstractDCS):
|
||||
self.reset_cluster()
|
||||
|
||||
def cancel_initialization(self):
|
||||
self.patch_or_create_config({self._INITIALIZE: None}, self.cluster.config.index, True)
|
||||
self.patch_or_create_config({self._INITIALIZE: None}, self._config_resource_version, True)
|
||||
|
||||
@catch_kubernetes_errors
|
||||
def delete_cluster(self):
|
||||
self.retry(self._api.delete_collection_namespaced_kind, self._namespace, label_selector=self._label_selector)
|
||||
|
||||
def set_history_value(self, value):
|
||||
patch = bool(self.cluster and self.cluster.config and self.cluster.config.index)
|
||||
return self.patch_or_create_config({self._HISTORY: value}, None, patch, False)
|
||||
return self.patch_or_create_config({self._HISTORY: value}, None, bool(self._config_resource_version), False)
|
||||
|
||||
def set_sync_state_value(self, value, index=None):
|
||||
"""Unused"""
|
||||
|
||||
@@ -126,6 +126,10 @@ class ZooKeeper(AbstractDCS):
|
||||
self._client.restart()
|
||||
return True
|
||||
|
||||
@property
|
||||
def ttl(self):
|
||||
return self._client._session_timeout
|
||||
|
||||
def set_retry_timeout(self, retry_timeout):
|
||||
retry = self._client.retry if isinstance(self._client.retry, KazooRetry) else self._client._retry
|
||||
retry.deadline = retry_timeout
|
||||
@@ -206,16 +210,18 @@ class ZooKeeper(AbstractDCS):
|
||||
failover = self.get_node(self.failover_path, watch=self.cluster_watcher) if self._FAILOVER in nodes else None
|
||||
failover = failover and Failover.from_node(failover[1].version, failover[0])
|
||||
|
||||
self._cluster = Cluster(initialize, config, leader, last_leader_operation, members, failover, sync, history)
|
||||
return Cluster(initialize, config, leader, last_leader_operation, members, failover, sync, history)
|
||||
|
||||
def _load_cluster(self):
|
||||
if self._fetch_cluster or self._cluster is None:
|
||||
cluster = self.cluster
|
||||
if self._fetch_cluster or cluster is None:
|
||||
try:
|
||||
self._client.retry(self._inner_load_cluster)
|
||||
cluster = self._client.retry(self._inner_load_cluster)
|
||||
except Exception:
|
||||
logger.exception('get_cluster')
|
||||
self.cluster_watcher(None)
|
||||
raise ZooKeeperError('ZooKeeper in not responding properly')
|
||||
return cluster
|
||||
|
||||
def _create(self, path, value, retry=False, ephemeral=False):
|
||||
try:
|
||||
@@ -264,7 +270,7 @@ class ZooKeeper(AbstractDCS):
|
||||
return self._create(self.initialize_path, sysid, retry=True) if create_new \
|
||||
else self._client.retry(self._client.set, self.initialize_path, sysid)
|
||||
|
||||
def touch_member(self, data, ttl=None, permanent=False):
|
||||
def touch_member(self, data, permanent=False):
|
||||
cluster = self.cluster
|
||||
member = cluster and cluster.get_member(self._name, fallback_to_leader=False)
|
||||
encoded_data = json.dumps(data, separators=(',', ':')).encode('utf-8')
|
||||
|
||||
+60
-37
@@ -112,11 +112,11 @@ class Ha(object):
|
||||
|
||||
def is_leader(self):
|
||||
with self._is_leader_lock:
|
||||
return self._is_leader and not self._leader_access_is_restricted
|
||||
return self._is_leader > time.time() and not self._leader_access_is_restricted
|
||||
|
||||
def set_is_leader(self, value):
|
||||
with self._is_leader_lock:
|
||||
self._is_leader = value
|
||||
self._is_leader = time.time() + self.dcs.ttl if value else 0
|
||||
|
||||
def set_leader_access_is_restricted(self, value):
|
||||
with self._is_leader_lock:
|
||||
@@ -130,7 +130,7 @@ class Ha(object):
|
||||
self.old_cluster = cluster
|
||||
self.cluster = cluster
|
||||
|
||||
if self.cluster.is_unlocked() or self.cluster.leader.name != self.state_handler.name:
|
||||
if not self.has_lock(False):
|
||||
self.set_is_leader(False)
|
||||
|
||||
self._leader_timeline = None if cluster.is_unlocked() else cluster.leader.timeline
|
||||
@@ -154,9 +154,10 @@ class Ha(object):
|
||||
self.watchdog.keepalive()
|
||||
return ret
|
||||
|
||||
def has_lock(self):
|
||||
def has_lock(self, info=True):
|
||||
lock_owner = self.cluster.leader and self.cluster.leader.name
|
||||
logger.info('Lock owner: %s; I am %s', lock_owner, self.state_handler.name)
|
||||
if info:
|
||||
logger.info('Lock owner: %s; I am %s', lock_owner, self.state_handler.name)
|
||||
return lock_owner == self.state_handler.name
|
||||
|
||||
def get_effective_tags(self):
|
||||
@@ -206,6 +207,9 @@ class Ha(object):
|
||||
return self.dcs.touch_member(data)
|
||||
|
||||
def clone(self, clone_member=None, msg='(without leader)'):
|
||||
if self.is_standby_cluster() and not isinstance(clone_member, RemoteMember):
|
||||
clone_member = self.get_remote_member(clone_member)
|
||||
|
||||
if self.state_handler.clone(clone_member):
|
||||
logger.info('bootstrapped %s', msg)
|
||||
cluster = self.dcs.get_cluster()
|
||||
@@ -301,6 +305,7 @@ class Ha(object):
|
||||
timeout = None
|
||||
|
||||
data = self.state_handler.controldata()
|
||||
logger.info('pg_controldata:\n%s\n', '\n'.join(' {0}: {1}'.format(k, v) for k, v in data.items()))
|
||||
if data.get('Database cluster state') in ('in production', 'shutting down', 'in crash recovery') and \
|
||||
not self._crash_recovery_executed and (self.cluster.is_unlocked() or self.state_handler.can_rewind):
|
||||
self._crash_recovery_executed = True
|
||||
@@ -310,6 +315,7 @@ class Ha(object):
|
||||
|
||||
self.load_cluster_from_dcs()
|
||||
|
||||
role = 'replica'
|
||||
if self.is_standby_cluster() or not self.has_lock():
|
||||
if not self.state_handler.rewind_executed:
|
||||
self.state_handler.trigger_check_diverged_lsn()
|
||||
@@ -318,6 +324,7 @@ class Ha(object):
|
||||
|
||||
if self.has_lock(): # in standby cluster
|
||||
msg = "starting as a standby leader because i had the session lock"
|
||||
role = 'standby_leader'
|
||||
node_to_follow = self._get_node_to_follow(self.cluster)
|
||||
elif self.is_standby_cluster() and self.cluster.is_unlocked():
|
||||
msg = "trying to follow a remote master because standby cluster is unhealthy"
|
||||
@@ -332,15 +339,13 @@ class Ha(object):
|
||||
self.recovering = True
|
||||
|
||||
self._async_executor.schedule('restarting after failure')
|
||||
self._async_executor.run_async(self.state_handler.follow, (node_to_follow, timeout))
|
||||
self._async_executor.run_async(self.state_handler.follow, (node_to_follow, role, timeout))
|
||||
return msg
|
||||
|
||||
def _get_node_to_follow(self, cluster):
|
||||
# determine the node to follow. If replicatefrom tag is set,
|
||||
# try to follow the node mentioned there, otherwise, follow the leader.
|
||||
is_leader = self.cluster.leader and self.state_handler.name == self.cluster.leader.name
|
||||
|
||||
if self.is_standby_cluster() and (is_leader or self.cluster.is_unlocked()):
|
||||
if self.is_standby_cluster() and (self.cluster.is_unlocked() or self.has_lock(False)):
|
||||
node_to_follow = self.get_remote_master()
|
||||
elif self.patroni.replicatefrom and self.patroni.replicatefrom != self.state_handler.name:
|
||||
node_to_follow = cluster.get_member(self.patroni.replicatefrom)
|
||||
@@ -372,9 +377,17 @@ class Ha(object):
|
||||
if self._handle_rewind_or_reinitialize():
|
||||
return self._async_executor.scheduled_action
|
||||
|
||||
if not self.state_handler.check_recovery_conf(node_to_follow):
|
||||
role = 'standby_leader' if isinstance(node_to_follow, RemoteMember) and self.has_lock(False) else 'replica'
|
||||
# It might happen that leader key in the standby cluster references non-exiting member.
|
||||
# In this case it is safe to continue running without changing recovery.conf
|
||||
if self.is_standby_cluster() and role == 'replica' and not (node_to_follow and node_to_follow.conn_url):
|
||||
return 'continue following the old known standby leader'
|
||||
elif not self.state_handler.check_recovery_conf(node_to_follow):
|
||||
self._async_executor.schedule('changing primary_conninfo and restarting')
|
||||
self._async_executor.run_async(self.state_handler.follow, (node_to_follow,))
|
||||
self._async_executor.run_async(self.state_handler.follow, (node_to_follow, role))
|
||||
elif role == 'standby_leader' and self.state_handler.role != role:
|
||||
self.state_handler.set_role(role)
|
||||
self.state_handler.call_nowait(ACTION_ON_ROLE_CHANGE)
|
||||
|
||||
return follow_reason
|
||||
|
||||
@@ -416,7 +429,10 @@ class Ha(object):
|
||||
time.sleep(2)
|
||||
picked, allow_promote = self.state_handler.pick_synchronous_standby(self.cluster)
|
||||
if allow_promote:
|
||||
cluster = self.dcs.get_cluster()
|
||||
try:
|
||||
cluster = self.dcs.get_cluster()
|
||||
except DCSError:
|
||||
return logger.warning("Could not get cluster state from DCS during process_sync_replication()")
|
||||
if cluster.sync.leader and cluster.sync.leader != self.state_handler.name:
|
||||
logger.info("Synchronous replication key updated by someone else")
|
||||
return
|
||||
@@ -489,9 +505,7 @@ class Ha(object):
|
||||
self.dcs.set_history_value(json.dumps(history, separators=(',', ':')))
|
||||
|
||||
def enforce_follow_remote_master(self, message):
|
||||
self.state_handler.set_role('standby_leader')
|
||||
demote_reason = 'cannot be a real master in standby cluster'
|
||||
|
||||
return self.follow(demote_reason, message)
|
||||
|
||||
def enforce_master_role(self, message, promote_message):
|
||||
@@ -693,10 +707,14 @@ class Ha(object):
|
||||
|
||||
return self._is_healthiest_node(members.values())
|
||||
|
||||
def release_leader_key_voluntarily(self):
|
||||
def _delete_leader(self):
|
||||
self.set_is_leader(False)
|
||||
self.dcs.delete_leader()
|
||||
self.touch_member()
|
||||
self.dcs.reset_cluster()
|
||||
|
||||
def release_leader_key_voluntarily(self):
|
||||
self._delete_leader()
|
||||
self.touch_member()
|
||||
logger.info("Leader key released")
|
||||
|
||||
def demote(self, mode):
|
||||
@@ -724,7 +742,8 @@ class Ha(object):
|
||||
self.set_is_leader(False)
|
||||
|
||||
if mode_control['release']:
|
||||
self.release_leader_key_voluntarily()
|
||||
with self._async_executor:
|
||||
self.release_leader_key_voluntarily()
|
||||
time.sleep(2) # Give a time to somebody to take the leader lock
|
||||
if mode_control['offline']:
|
||||
node_to_follow, leader = None, None
|
||||
@@ -843,7 +862,7 @@ class Ha(object):
|
||||
# standby leader disappeared, and this is a healthiest
|
||||
# replica, so it should become a new standby leader.
|
||||
# This imply that we need to start following a remote master
|
||||
msg = 'promoted self to a standby leader because i had the session lock'
|
||||
msg = 'promoted self to a standby leader by acquiring session lock'
|
||||
return self.enforce_follow_remote_master(msg)
|
||||
else:
|
||||
return self.enforce_master_role(
|
||||
@@ -872,8 +891,7 @@ class Ha(object):
|
||||
if self.cluster.failover and self.cluster.failover.candidate == self.state_handler.name:
|
||||
return 'waiting to become master after promote...'
|
||||
|
||||
self.dcs.delete_leader()
|
||||
self.dcs.reset_cluster()
|
||||
self._delete_leader()
|
||||
return 'removed leader lock because postgres is not running as master'
|
||||
|
||||
if self.state_handler.is_leader() and self._leader_access_is_restricted:
|
||||
@@ -890,7 +908,9 @@ class Ha(object):
|
||||
# in case of standby cluster we don't really need to
|
||||
# enforce anything, since the leader is not a master.
|
||||
# So just remind the role.
|
||||
msg = 'no action. i am the standby leader with the lock'
|
||||
msg = 'no action. i am the standby leader with the lock' \
|
||||
if self.state_handler.role == 'standby_leader' else \
|
||||
'promoted self to a standby leader because i had the session lock'
|
||||
return self.enforce_follow_remote_master(msg)
|
||||
else:
|
||||
return self.enforce_master_role(
|
||||
@@ -909,6 +929,9 @@ class Ha(object):
|
||||
return 'not promoting because failed to update leader lock in DCS'
|
||||
else:
|
||||
logger.info('does not have lock')
|
||||
if self.is_standby_cluster():
|
||||
return self.follow('cannot be a real master in standby cluster',
|
||||
'no action. i am a secondary and i am following a standby leader', refresh=False)
|
||||
return self.follow('demoting self because i do not have the lock and i was a leader',
|
||||
'no action. i am a secondary and i am following a leader', refresh=False)
|
||||
|
||||
@@ -1040,7 +1063,7 @@ class Ha(object):
|
||||
if self.cluster.is_unlocked():
|
||||
return 'Cluster has no leader, can not reinitialize'
|
||||
|
||||
if self.cluster.leader.name == self.state_handler.name:
|
||||
if self.has_lock(False):
|
||||
return 'I am the leader, can not reinitialize'
|
||||
|
||||
if force:
|
||||
@@ -1086,8 +1109,7 @@ class Ha(object):
|
||||
self.watchdog.disable()
|
||||
if self.has_lock():
|
||||
self.state_handler.set_role('demoted')
|
||||
self.dcs.delete_leader()
|
||||
self.dcs.reset_cluster()
|
||||
self._delete_leader()
|
||||
return 'removed leader key after trying and failing to start postgres'
|
||||
return 'failed to start postgres'
|
||||
self._crash_recovery_executed = False
|
||||
@@ -1253,16 +1275,15 @@ class Ha(object):
|
||||
if not self.state_handler.is_healthy():
|
||||
if self.is_paused():
|
||||
if self.has_lock():
|
||||
self.dcs.delete_leader()
|
||||
self.dcs.reset_cluster()
|
||||
self._delete_leader()
|
||||
return 'removed leader lock because postgres is not running'
|
||||
# Normally we don't start Postgres in a paused state. We make an exception for the demoted primary
|
||||
# that needs to be started after it had been stopped by demote. When there is no need to call rewind
|
||||
# the demote code follows through to starting Postgres right away, however, in the rewind case
|
||||
# it returns from demote and reaches this point to start PostgreSQL again after rewind. In that
|
||||
# case it makes no sense to continue to recover() unless rewind has finished successfully.
|
||||
elif self.state_handler.rewind_failed or not self.state_handler.need_rewind \
|
||||
or not self.state_handler.can_rewind_or_reinitialize_allowed:
|
||||
elif self.state_handler.rewind_failed or not self.state_handler.rewind_executed and not \
|
||||
(self.state_handler.need_rewind and self.state_handler.can_rewind_or_reinitialize_allowed):
|
||||
return 'postgres is not running'
|
||||
|
||||
# try to start dead postgres
|
||||
@@ -1325,13 +1346,11 @@ class Ha(object):
|
||||
(" Leaving watchdog running." if self.watchdog.is_running else ""))
|
||||
|
||||
def watch(self, timeout):
|
||||
cluster = self.cluster
|
||||
# watch on leader key changes if the postgres is running and leader is known and current node is not lock owner
|
||||
if not self._async_executor.busy and cluster and cluster.leader \
|
||||
and cluster.leader.name != self.state_handler.name:
|
||||
leader_index = cluster.leader.index
|
||||
else:
|
||||
if self._async_executor.busy or self.cluster.is_unlocked() or self.has_lock(False):
|
||||
leader_index = None
|
||||
else:
|
||||
leader_index = self.cluster.leader.index
|
||||
|
||||
return self.dcs.watch(leader_index, timeout)
|
||||
|
||||
@@ -1341,7 +1360,7 @@ class Ha(object):
|
||||
This usually happens on the master or if the node is running async action"""
|
||||
self.dcs.event.set()
|
||||
|
||||
def get_remote_master(self):
|
||||
def get_remote_member(self, member=None):
|
||||
""" In case of standby cluster this will tel us from which remote
|
||||
master to stream. Config can be both patroni config or
|
||||
cluster.config.data
|
||||
@@ -1349,12 +1368,16 @@ class Ha(object):
|
||||
cluster_params = self.get_standby_cluster_config()
|
||||
|
||||
if cluster_params:
|
||||
unique_name = 'remote_master:{}'.format(uuid.uuid1())
|
||||
name = member.name if member else 'remote_master:{}'.format(uuid.uuid1())
|
||||
|
||||
data = {k: v for k, v in cluster_params.items() if k in RemoteMember.allowed_keys()}
|
||||
data['no_replication_slot'] = 'primary_slot_name' not in cluster_params
|
||||
conn_kwargs = {k: cluster_params[k] for k in ('host', 'port') if k in cluster_params}
|
||||
conn_kwargs = member.conn_kwargs() if member else \
|
||||
{k: cluster_params[k] for k in ('host', 'port') if k in cluster_params}
|
||||
if conn_kwargs:
|
||||
data['conn_kwargs'] = conn_kwargs
|
||||
|
||||
return RemoteMember(unique_name, data)
|
||||
return RemoteMember(name, data)
|
||||
|
||||
def get_remote_master(self):
|
||||
return self.get_remote_member()
|
||||
|
||||
+35
-15
@@ -30,6 +30,7 @@ ACTION_ON_STOP = "on_stop"
|
||||
ACTION_ON_RESTART = "on_restart"
|
||||
ACTION_ON_RELOAD = "on_reload"
|
||||
ACTION_ON_ROLE_CHANGE = "on_role_change"
|
||||
ACTION_NOOP = "noop"
|
||||
|
||||
STATE_RUNNING = 'running'
|
||||
STATE_REJECT = 'rejecting connections'
|
||||
@@ -697,10 +698,11 @@ class Postgresql(object):
|
||||
# if basebackup succeeds, exit with success
|
||||
break
|
||||
else:
|
||||
if not self.data_directory_empty() and not self.config.get(replica_method, {}).get('keep_data', False):
|
||||
self.remove_data_directory()
|
||||
else:
|
||||
logger.info('Leaving data directory uncleaned')
|
||||
if not self.data_directory_empty():
|
||||
if self.config.get(replica_method, {}).get('keep_data', False):
|
||||
logger.info('Leaving data directory uncleaned')
|
||||
else:
|
||||
self.remove_data_directory()
|
||||
|
||||
cmd = replica_method
|
||||
method_config = {}
|
||||
@@ -870,7 +872,7 @@ class Postgresql(object):
|
||||
self._pending_restart = True
|
||||
return effective_configuration
|
||||
|
||||
def start(self, timeout=None, block_callbacks=False, task=None):
|
||||
def start(self, timeout=None, task=None, block_callbacks=False, role=None):
|
||||
"""Start PostgreSQL
|
||||
|
||||
Waits for postmaster to open ports or terminate so pg_isready can be used to check startup completion
|
||||
@@ -890,7 +892,7 @@ class Postgresql(object):
|
||||
if not block_callbacks:
|
||||
self.__cb_pending = ACTION_ON_START
|
||||
|
||||
self.set_role(self.get_postgres_role_from_data_directory())
|
||||
self.set_role(role or self.get_postgres_role_from_data_directory())
|
||||
|
||||
self.set_state('starting')
|
||||
self._pending_restart = False
|
||||
@@ -1091,7 +1093,7 @@ class Postgresql(object):
|
||||
|
||||
return self.state == 'running'
|
||||
|
||||
def restart(self, timeout=None, task=None):
|
||||
def restart(self, timeout=None, task=None, block_callbacks=False, role=None):
|
||||
"""Restarts PostgreSQL.
|
||||
|
||||
When timeout parameter is set the call will block either until PostgreSQL has started, failed to start or
|
||||
@@ -1100,8 +1102,9 @@ class Postgresql(object):
|
||||
:returns: True when restart was successful and timeout did not expire when waiting.
|
||||
"""
|
||||
self.set_state('restarting')
|
||||
self.__cb_pending = ACTION_ON_RESTART
|
||||
ret = self.stop(block_callbacks=True) and self.start(timeout=timeout, block_callbacks=True, task=task)
|
||||
if not block_callbacks:
|
||||
self.__cb_pending = ACTION_ON_RESTART
|
||||
ret = self.stop(block_callbacks=True) and self.start(timeout, task, True, role)
|
||||
if not ret and not self.is_starting():
|
||||
self.set_state('restart failed ({0})'.format(self.state))
|
||||
return ret
|
||||
@@ -1436,7 +1439,7 @@ class Postgresql(object):
|
||||
def rewind_failed(self):
|
||||
return self._rewind_state == REWIND_STATUS.FAILED
|
||||
|
||||
def follow(self, member, timeout=None):
|
||||
def follow(self, member, role='replica', timeout=None):
|
||||
is_remote_master = isinstance(member, RemoteMember)
|
||||
no_replication_slot = is_remote_master and member.no_replication_slot
|
||||
restore_command = is_remote_master and member.restore_command
|
||||
@@ -1444,7 +1447,8 @@ class Postgresql(object):
|
||||
archive_cleanup = is_remote_master and member.archive_cleanup_command
|
||||
|
||||
primary_conninfo = self.primary_conninfo(member)
|
||||
change_role = self.role in ('master', 'demoted')
|
||||
change_role = self.cb_called and (self.role in ('master', 'demoted') or
|
||||
not {'standby_leader', 'replica'} - {self.role, role})
|
||||
|
||||
recovery_params = self.config.get('recovery_conf', {}).copy()
|
||||
recovery_params.update({'standby_mode': 'on', 'recovery_target_timeline': 'latest'})
|
||||
@@ -1463,11 +1467,17 @@ class Postgresql(object):
|
||||
|
||||
self.write_recovery_conf(recovery_params)
|
||||
|
||||
# When we demoting the master or standby_leader to replica or promoting replica to a standby_leader
|
||||
# and we know for sure that postgres was already running before, we will only execute on_role_change
|
||||
# callback and prevent execution of on_restart/on_start callback.
|
||||
# If the role remains the same (replica or standby_leader), we will execute on_start or on_restart
|
||||
if change_role:
|
||||
self.__cb_pending = ACTION_NOOP
|
||||
|
||||
if self.is_running():
|
||||
self.restart()
|
||||
self.restart(block_callbacks=change_role, role=role)
|
||||
else:
|
||||
self.start(timeout=timeout)
|
||||
self.set_role('replica')
|
||||
self.start(timeout=timeout, block_callbacks=change_role, role=role)
|
||||
|
||||
if change_role:
|
||||
# TODO: postpone this until start completes, or maybe do even earlier
|
||||
@@ -1506,7 +1516,7 @@ class Postgresql(object):
|
||||
logger.exception('unable to restore configuration files from backup')
|
||||
|
||||
def _wait_promote(self, wait_seconds):
|
||||
for _ in polling_loop(wait_seconds - 1):
|
||||
for _ in polling_loop(wait_seconds):
|
||||
data = self.controldata()
|
||||
if data.get('Database cluster state') == 'in production':
|
||||
return True
|
||||
@@ -1732,6 +1742,16 @@ $$""".format(name, ' '.join(options)), name, password, password)
|
||||
elif os.path.isfile(self._data_dir):
|
||||
os.remove(self._data_dir)
|
||||
elif os.path.isdir(self._data_dir):
|
||||
|
||||
# let's see if pg_xlog|pg_wal is a symlink, in this case we
|
||||
# should clean the target
|
||||
for pg_wal_dir in ('pg_xlog', 'pg_wal'):
|
||||
pg_wal_path = os.path.join(self._data_dir, pg_wal_dir)
|
||||
if os.path.exists(pg_wal_path) and os.path.islink(pg_wal_path):
|
||||
pg_wal_realpath = os.path.realpath(pg_wal_path)
|
||||
logger.info('Removing WAL directory: %s', pg_wal_realpath)
|
||||
shutil.rmtree(pg_wal_realpath)
|
||||
|
||||
shutil.rmtree(self._data_dir)
|
||||
except (IOError, OSError):
|
||||
logger.exception('Could not remove data directory %s', self._data_dir)
|
||||
|
||||
+20
-6
@@ -1,12 +1,11 @@
|
||||
import logging
|
||||
import multiprocessing
|
||||
import os
|
||||
import psutil
|
||||
import re
|
||||
import signal
|
||||
import subprocess
|
||||
|
||||
from patroni import call_self
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
STOP_SIGNALS = {
|
||||
@@ -16,6 +15,18 @@ STOP_SIGNALS = {
|
||||
}
|
||||
|
||||
|
||||
def pg_ctl_start(conn, cmdline, env):
|
||||
if os.name != 'nt':
|
||||
os.setsid()
|
||||
try:
|
||||
postmaster = subprocess.Popen(cmdline, close_fds=True, env=env)
|
||||
conn.send(postmaster.pid)
|
||||
except Exception:
|
||||
logger.exception('Failed to execute %s', cmdline)
|
||||
conn.send(None)
|
||||
conn.close()
|
||||
|
||||
|
||||
class PostmasterProcess(psutil.Process):
|
||||
|
||||
def __init__(self, pid):
|
||||
@@ -159,10 +170,13 @@ class PostmasterProcess(psutil.Process):
|
||||
pass
|
||||
cmdline = [pgcommand, '-D', data_dir, '--config-file={}'.format(conf)] + options
|
||||
logger.debug("Starting postgres: %s", " ".join(cmdline))
|
||||
proc = call_self(['pg_ctl_start'] + cmdline, close_fds=(os.name != 'nt'),
|
||||
stdout=subprocess.PIPE, env=env)
|
||||
pid = int(proc.stdout.readline().strip())
|
||||
proc.wait()
|
||||
parent_conn, child_conn = multiprocessing.Pipe(False)
|
||||
proc = multiprocessing.Process(target=pg_ctl_start, args=(child_conn, cmdline, env))
|
||||
proc.start()
|
||||
pid = parent_conn.recv()
|
||||
proc.join()
|
||||
if pid is None:
|
||||
return
|
||||
logger.info('postmaster pid=%s', pid)
|
||||
|
||||
# TODO: In an extremely unlikely case, the process could have exited and the pid reassigned. The start
|
||||
|
||||
+1
-1
@@ -1 +1 @@
|
||||
__version__ = '1.5.5'
|
||||
__version__ = '1.5.6'
|
||||
|
||||
@@ -1,6 +1,5 @@
|
||||
import collections
|
||||
import ctypes
|
||||
import fcntl
|
||||
import os
|
||||
import platform
|
||||
from patroni.watchdog.base import WatchdogBase, WatchdogError
|
||||
@@ -162,7 +161,9 @@ class LinuxWatchdogDevice(WatchdogBase):
|
||||
Raises OSError or IOError (Python 2) when the ioctl fails."""
|
||||
if self._fd is None:
|
||||
raise WatchdogError("Watchdog device is closed")
|
||||
fcntl.ioctl(self._fd, func, arg, True)
|
||||
if os.name != 'nt':
|
||||
import fcntl
|
||||
fcntl.ioctl(self._fd, func, arg, True)
|
||||
|
||||
def get_support(self):
|
||||
if self._support_cache is None:
|
||||
|
||||
@@ -17,6 +17,9 @@ class TestCallbackExecutor(unittest.TestCase):
|
||||
|
||||
self.assertIsNone(ce.call([]))
|
||||
|
||||
mock_popen.return_value.kill.side_effect = OSError
|
||||
self.assertRaises(Exception, ce.call, [])
|
||||
|
||||
mock_popen.side_effect = Exception
|
||||
ce = CallbackExecutor()
|
||||
ce._callback_event.wait = Mock(side_effect=[None, Exception])
|
||||
|
||||
@@ -52,7 +52,7 @@ class TestConfig(unittest.TestCase):
|
||||
'PATRONI_ETCD_KEY': '/key',
|
||||
'PATRONI_CONSUL_HOST': '127.0.0.1:8500',
|
||||
'PATRONI_CONSUL_REGISTER_SERVICE': 'on',
|
||||
'PATRONI_KUBERNETES_LABELS': 'a:b:c',
|
||||
'PATRONI_KUBERNETES_LABELS': 'a: b: c',
|
||||
'PATRONI_KUBERNETES_SCOPE_LABEL': 'a',
|
||||
'PATRONI_KUBERNETES_PORTS': '[{"name": "postgresql"}]',
|
||||
'PATRONI_ZOOKEEPER_HOSTS': "'host1:2181','host2:2181'",
|
||||
|
||||
@@ -83,7 +83,7 @@ class TestConsul(unittest.TestCase):
|
||||
self.c = Consul({'ttl': 30, 'scope': 'test', 'name': 'postgresql1', 'host': 'localhost:1', 'retry_timeout': 10,
|
||||
'register_service': True})
|
||||
self.c._base_path = '/service/good'
|
||||
self.c._load_cluster()
|
||||
self.c.get_cluster()
|
||||
|
||||
@patch('time.sleep', Mock(side_effect=SleepException))
|
||||
@patch.object(consul.Consul.Session, 'create', Mock(side_effect=ConsulException))
|
||||
|
||||
+4
-2
@@ -533,7 +533,8 @@ class TestCtl(unittest.TestCase):
|
||||
show_diff(b"foo:\n bar: \xc3\xb6\xc3\xb6\n".decode('utf-8'),
|
||||
b"foo:\n bar: \xc3\xbc\xc3\xbc\n".decode('utf-8'))
|
||||
|
||||
def test_invoke_editor(self):
|
||||
@patch('subprocess.call', return_value=1)
|
||||
def test_invoke_editor(self, mock_subprocess_call):
|
||||
for e in ('', 'false'):
|
||||
os.environ['EDITOR'] = e
|
||||
self.assertRaises(PatroniCtlException, invoke_editor, 'foo: bar\n', 'test')
|
||||
@@ -548,13 +549,14 @@ class TestCtl(unittest.TestCase):
|
||||
def test_edit_config(self, mock_get_dcs):
|
||||
mock_get_dcs.return_value = self.e
|
||||
mock_get_dcs.return_value.get_cluster = get_cluster_initialized_with_leader
|
||||
mock_get_dcs.return_value.set_config_value = Mock(return_value=False)
|
||||
os.environ['EDITOR'] = 'true'
|
||||
self.runner.invoke(ctl, ['edit-config', 'dummy'])
|
||||
self.runner.invoke(ctl, ['edit-config', 'dummy', '-s', 'foo=bar'])
|
||||
self.runner.invoke(ctl, ['edit-config', 'dummy', '--replace', 'postgres0.yml'])
|
||||
self.runner.invoke(ctl, ['edit-config', 'dummy', '--apply', '-'], input='foo: bar')
|
||||
self.runner.invoke(ctl, ['edit-config', 'dummy', '--force', '--apply', '-'], input='foo: bar')
|
||||
mock_get_dcs.return_value.set_config_value = Mock(return_value=True)
|
||||
mock_get_dcs.return_value.set_config_value.return_value = True
|
||||
self.runner.invoke(ctl, ['edit-config', 'dummy', '--force', '--apply', '-'], input='foo: bar')
|
||||
|
||||
@patch('patroni.ctl.get_dcs')
|
||||
|
||||
+18
-16
@@ -148,13 +148,14 @@ def socket_getaddrinfo(*args):
|
||||
def http_request(method, url, **kwargs):
|
||||
if url == 'http://localhost:2379/timeout':
|
||||
raise ReadTimeoutError(None, None, None)
|
||||
ret = MockResponse()
|
||||
if url == 'http://localhost:2379/v2/machines':
|
||||
ret = MockResponse()
|
||||
ret.content = 'http://localhost:2379,http://localhost:4001'
|
||||
return ret
|
||||
if url == 'http://localhost:2379/':
|
||||
return MockResponse()
|
||||
raise socket.error
|
||||
elif url == 'http://localhost:4001/v2/machines':
|
||||
ret.content = ''
|
||||
elif url != 'http://localhost:2379/':
|
||||
raise socket.error
|
||||
return ret
|
||||
|
||||
|
||||
class TestDnsCachingResolver(unittest.TestCase):
|
||||
@@ -209,8 +210,8 @@ class TestClient(unittest.TestCase):
|
||||
self.client.api_execute('/', 'POST', timeout=0)
|
||||
self.client._machines_cache = [self.client._base_uri]
|
||||
self.assertRaises(etcd.EtcdWatchTimedOut, self.client.api_execute, '/timeout', 'POST', params={'wait': 'true'})
|
||||
self.assertRaises(etcd.EtcdWatchTimedOut, self.client.api_execute, '/timeout', 'POST', params={'wait': 'true'})
|
||||
self.assertRaises(etcd.EtcdException, self.client.api_execute, '/', '')
|
||||
self.client._update_machines_cache = True
|
||||
with patch.object(Client, '_load_machines_cache', Mock(side_effect=etcd.EtcdException)):
|
||||
self.assertRaises(etcd.EtcdException, self.client.api_execute, '/', 'GET')
|
||||
|
||||
@@ -263,17 +264,18 @@ class TestEtcd(unittest.TestCase):
|
||||
|
||||
@patch('dns.resolver.query', dns_query)
|
||||
def test_get_etcd_client(self):
|
||||
with patch.object(Client, 'machines') as mock_machines:
|
||||
with patch('time.sleep', Mock(side_effect=SleepException)),\
|
||||
patch.object(Client, 'machines') as mock_machines:
|
||||
mock_machines.__get__ = Mock(side_effect=etcd.EtcdException)
|
||||
with patch('time.sleep', Mock(side_effect=SleepException)):
|
||||
self.assertRaises(SleepException, self.etcd.get_etcd_client,
|
||||
{'discovery_srv': 'test', 'retry_timeout': 10, 'cacert': '1', 'key': '1', 'cert': 1})
|
||||
self.assertRaises(SleepException, self.etcd.get_etcd_client,
|
||||
{'url': 'https://test:2379', 'retry_timeout': 10})
|
||||
self.assertRaises(SleepException, self.etcd.get_etcd_client,
|
||||
{'proxy': 'https://user:password@test:2379', 'retry_timeout': 10})
|
||||
self.assertRaises(SleepException, self.etcd.get_etcd_client,
|
||||
{'hosts': 'foo:4001,bar', 'retry_timeout': 10})
|
||||
self.assertRaises(SleepException, self.etcd.get_etcd_client,
|
||||
{'discovery_srv': 'test', 'retry_timeout': 10, 'cacert': '1', 'key': '1', 'cert': 1})
|
||||
self.assertRaises(SleepException, self.etcd.get_etcd_client,
|
||||
{'url': 'https://test:2379', 'retry_timeout': 10})
|
||||
self.assertRaises(SleepException, self.etcd.get_etcd_client,
|
||||
{'hosts': 'foo:4001,bar', 'retry_timeout': 10})
|
||||
mock_machines.__get__ = Mock(return_value=[])
|
||||
self.assertRaises(SleepException, self.etcd.get_etcd_client,
|
||||
{'proxy': 'https://user:password@test:2379', 'retry_timeout': 10})
|
||||
|
||||
def test_get_cluster(self):
|
||||
cluster = self.etcd.get_cluster()
|
||||
|
||||
+25
-12
@@ -152,6 +152,8 @@ def run_async(self, func, args=()):
|
||||
@patch.object(Postgresql, 'checkpoint', Mock())
|
||||
@patch.object(Postgresql, 'cancellable_subprocess_call', Mock(return_value=0))
|
||||
@patch.object(Postgresql, '_get_local_timeline_lsn_from_replication_connection', Mock(return_value=[2, 10]))
|
||||
@patch.object(Postgresql, 'get_master_timeline', Mock(return_value=2))
|
||||
@patch.object(Postgresql, 'restore_configuration_files', Mock())
|
||||
@patch.object(etcd.Client, 'write', etcd_write)
|
||||
@patch.object(etcd.Client, 'read', etcd_read)
|
||||
@patch.object(etcd.Client, 'delete', Mock(side_effect=etcd.EtcdException))
|
||||
@@ -211,7 +213,8 @@ class TestHa(unittest.TestCase):
|
||||
self.assertEqual(self.ha.run_cycle(), 'trying to bootstrap a new standby leader')
|
||||
|
||||
@patch.object(Cluster, 'get_clone_member',
|
||||
Mock(return_value=Member(0, 'test', 1, {'api_url': 'http://127.0.0.1:8011/patroni'})))
|
||||
Mock(return_value=Member(0, 'test', 1, {'api_url': 'http://127.0.0.1:8011/patroni',
|
||||
'conn_url': 'postgres://127.0.0.1:5432/postgres'})))
|
||||
@patch.object(Postgresql, 'create_replica', Mock(return_value=0))
|
||||
def test_start_as_cascade_replica_in_standby_cluster(self):
|
||||
self.p.data_directory_empty = true
|
||||
@@ -321,7 +324,7 @@ class TestHa(unittest.TestCase):
|
||||
self.assertEqual(self.ha.run_cycle(), 'Not promoting self because watchdog could not be activated')
|
||||
|
||||
def test_leader_with_lock(self):
|
||||
self.ha.cluster = get_cluster_not_initialized_without_leader()
|
||||
self.ha.cluster = get_cluster_initialized_with_leader()
|
||||
self.ha.cluster.is_unlocked = false
|
||||
self.ha.has_lock = true
|
||||
self.assertEqual(self.ha.run_cycle(), 'no action. i am the leader with the lock')
|
||||
@@ -576,6 +579,7 @@ class TestHa(unittest.TestCase):
|
||||
self.ha.is_paused = true
|
||||
self.assertFalse(self.ha.is_healthiest_node())
|
||||
|
||||
@patch('requests.get', requests_get)
|
||||
def test__is_healthiest_node(self):
|
||||
self.assertTrue(self.ha._is_healthiest_node(self.ha.old_cluster.members))
|
||||
self.p.is_leader = false
|
||||
@@ -670,16 +674,19 @@ class TestHa(unittest.TestCase):
|
||||
self.p.is_leader = false
|
||||
self.p.name = 'leader'
|
||||
self.ha.cluster = get_standby_cluster_initialized_with_only_leader()
|
||||
msg = 'no action. i am the standby leader with the lock'
|
||||
self.assertEqual(self.ha.run_cycle(), msg)
|
||||
self.p.check_recovery_conf = true
|
||||
self.assertEqual(self.ha.run_cycle(), 'promoted self to a standby leader because i had the session lock')
|
||||
self.assertEqual(self.ha.run_cycle(), 'no action. i am the standby leader with the lock')
|
||||
|
||||
def test_process_healthy_standby_cluster_as_cascade_replica(self):
|
||||
self.p.is_leader = false
|
||||
self.p.name = 'replica'
|
||||
self.ha.cluster = get_standby_cluster_initialized_with_only_leader()
|
||||
msg = 'no action. i am a secondary and i am following a leader'
|
||||
self.assertEqual(self.ha.run_cycle(), msg)
|
||||
self.assertEqual(self.ha.run_cycle(), 'no action. i am a secondary and i am following a standby leader')
|
||||
with patch.object(Leader, 'conn_url', PropertyMock(return_value='')):
|
||||
self.assertEqual(self.ha.run_cycle(), 'continue following the old known standby leader')
|
||||
|
||||
@patch('requests.get', requests_get)
|
||||
def test_process_unhealthy_standby_cluster_as_standby_leader(self):
|
||||
self.p.is_leader = false
|
||||
self.p.name = 'leader'
|
||||
@@ -687,8 +694,7 @@ class TestHa(unittest.TestCase):
|
||||
self.ha.cluster.is_unlocked = true
|
||||
self.ha.sysid_valid = true
|
||||
self.p._sysid = True
|
||||
msg = 'promoted self to a standby leader because i had the session lock'
|
||||
self.assertEqual(self.ha.run_cycle(), msg)
|
||||
self.assertEqual(self.ha.run_cycle(), 'promoted self to a standby leader by acquiring session lock')
|
||||
|
||||
@patch.object(Postgresql, 'rewind_or_reinitialize_needed_and_possible', Mock(return_value=True))
|
||||
@patch.object(Postgresql, 'can_rewind', PropertyMock(return_value=True))
|
||||
@@ -859,11 +865,15 @@ class TestHa(unittest.TestCase):
|
||||
self.ha.run_cycle()
|
||||
self.assertEqual(self.ha.dcs.write_sync_state.call_count, 2)
|
||||
|
||||
# Test changing sync standby failed due to race
|
||||
# Test updating sync standby key failed due to DCS being not accessible
|
||||
self.ha.dcs.write_sync_state = Mock(return_value=True)
|
||||
self.ha.dcs.get_cluster = Mock(side_effect=DCSError('foo'))
|
||||
self.ha.run_cycle()
|
||||
|
||||
# Test changing sync standby failed due to race
|
||||
self.ha.dcs.get_cluster = Mock(return_value=get_cluster_initialized_with_leader(sync=('somebodyelse', None)))
|
||||
self.ha.run_cycle()
|
||||
self.assertEqual(self.ha.dcs.write_sync_state.call_count, 1)
|
||||
self.assertEqual(self.ha.dcs.write_sync_state.call_count, 2)
|
||||
|
||||
# Test sync set to '*' when synchronous_mode_strict is enabled
|
||||
mock_set_sync.reset_mock()
|
||||
@@ -998,13 +1008,16 @@ class TestHa(unittest.TestCase):
|
||||
# will not say bootstrap from leader as replica can't self elect
|
||||
self.assertEqual(self.ha.run_cycle(), "trying to bootstrap from replica 'other'")
|
||||
|
||||
@patch('psycopg2.connect', psycopg2_connect)
|
||||
def test_update_cluster_history(self):
|
||||
self.p.get_master_timeline = Mock(return_value=1)
|
||||
self.ha.has_lock = true
|
||||
self.ha.cluster.is_unlocked = false
|
||||
self.assertEqual(self.ha.run_cycle(), 'no action. i am the leader with the lock')
|
||||
for tl in (1, 3):
|
||||
self.p.get_master_timeline = Mock(return_value=tl)
|
||||
self.assertEqual(self.ha.run_cycle(), 'no action. i am the leader with the lock')
|
||||
|
||||
@patch('sys.exit', return_value=1)
|
||||
@patch('requests.get', requests_get)
|
||||
def test_abort_join(self, exit_mock):
|
||||
self.ha.cluster = get_cluster_not_initialized_without_leader()
|
||||
self.p.is_leader = false
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import time
|
||||
import unittest
|
||||
|
||||
from mock import Mock, patch
|
||||
@@ -33,13 +34,14 @@ class TestKubernetes(unittest.TestCase):
|
||||
@patch.object(k8s_client.CoreV1Api, 'list_namespaced_pod', mock_list_namespaced_pod)
|
||||
def setUp(self):
|
||||
self.k = Kubernetes({'ttl': 30, 'scope': 'test', 'name': 'p-0', 'retry_timeout': 10, 'labels': {'f': 'b'}})
|
||||
with patch('time.time', Mock(return_value=1)):
|
||||
self.k.get_cluster()
|
||||
|
||||
def test_get_cluster(self):
|
||||
with patch.object(k8s_client.CoreV1Api, 'list_namespaced_config_map', mock_list_namespaced_config_map), \
|
||||
patch.object(k8s_client.CoreV1Api, 'list_namespaced_pod', mock_list_namespaced_pod), \
|
||||
patch('time.time', Mock(return_value=time.time() + 31)):
|
||||
self.k.get_cluster()
|
||||
|
||||
@patch.object(k8s_client.CoreV1Api, 'list_namespaced_config_map', mock_list_namespaced_config_map)
|
||||
@patch.object(k8s_client.CoreV1Api, 'list_namespaced_pod', mock_list_namespaced_pod)
|
||||
def test_get_cluster(self):
|
||||
self.k.get_cluster()
|
||||
with patch.object(k8s_client.CoreV1Api, 'list_namespaced_pod', Mock(side_effect=Exception)):
|
||||
self.assertRaises(KubernetesError, self.k.get_cluster)
|
||||
|
||||
|
||||
@@ -67,16 +67,12 @@ class TestPatroni(unittest.TestCase):
|
||||
patroni_main()
|
||||
|
||||
@patch('os.getpid')
|
||||
@patch('subprocess.Popen', )
|
||||
@patch('multiprocessing.Process')
|
||||
@patch('patroni.patroni_main', Mock())
|
||||
def test_patroni_main(self, mock_popen, mock_getpid):
|
||||
def test_patroni_main(self, mock_process, mock_getpid):
|
||||
mock_getpid.return_value = 2
|
||||
_main()
|
||||
|
||||
with patch('sys.frozen', Mock(return_value=True), create=True), patch('os.setsid', Mock()):
|
||||
sys.argv = ['/patroni', 'pg_ctl_start', 'postgres', '-D', '/data', '--max_connections=100']
|
||||
_main()
|
||||
|
||||
mock_getpid.return_value = 1
|
||||
|
||||
def mock_signal(signo, handler):
|
||||
@@ -91,13 +87,13 @@ class TestPatroni(unittest.TestCase):
|
||||
ref = {'passtochild': lambda signo, stack_frame: 0}
|
||||
|
||||
def mock_sighup(signo, handler):
|
||||
if signo == signal.SIGHUP:
|
||||
if hasattr(signal, 'SIGHUP') and signo == signal.SIGHUP:
|
||||
ref['passtochild'] = handler
|
||||
|
||||
def mock_wait():
|
||||
def mock_join():
|
||||
ref['passtochild'](0, None)
|
||||
|
||||
mock_popen.return_value.wait = mock_wait
|
||||
mock_process.return_value.join = mock_join
|
||||
with patch('signal.signal', mock_sighup), patch('os.kill', Mock()):
|
||||
self.assertIsNone(_main())
|
||||
|
||||
@@ -122,6 +118,7 @@ class TestPatroni(unittest.TestCase):
|
||||
self.assertRaises(SystemExit, self.p.sigterm_handler)
|
||||
|
||||
def test_schedule_next_run(self):
|
||||
self.p.ha.cluster = Mock()
|
||||
self.p.ha.dcs.watch = Mock(return_value=True)
|
||||
self.p.schedule_next_run()
|
||||
self.p.next_run = time.time() - self.p.dcs.loop_wait - 1
|
||||
|
||||
@@ -15,6 +15,7 @@ from patroni.postmaster import PostmasterProcess
|
||||
from patroni.utils import RetryFailedError
|
||||
from six.moves import builtins
|
||||
from threading import Thread, current_thread
|
||||
from tempfile import gettempdir
|
||||
|
||||
|
||||
class MockCursor(object):
|
||||
@@ -189,7 +190,8 @@ class TestPostgresql(unittest.TestCase):
|
||||
if not os.path.exists(self.data_dir):
|
||||
os.makedirs(self.data_dir)
|
||||
self.p = Postgresql({'name': 'test0', 'scope': 'batman', 'data_dir': self.data_dir,
|
||||
'config_dir': self.config_dir, 'retry_timeout': 10, 'pgpass': '/tmp/pgpass0',
|
||||
'config_dir': self.config_dir, 'retry_timeout': 10,
|
||||
'pgpass': os.path.join(gettempdir(), 'pgpass0'),
|
||||
'listen': '127.0.0.2, 127.0.0.3:5432', 'connect_address': '127.0.0.2:5432',
|
||||
'authentication': {'superuser': {'username': 'test', 'password': 'test'},
|
||||
'replication': {'username': 'replicator', 'password': 'rep-pass'}},
|
||||
@@ -403,6 +405,7 @@ class TestPostgresql(unittest.TestCase):
|
||||
@patch.object(Postgresql, 'is_running', Mock(return_value=False))
|
||||
@patch.object(Postgresql, 'start', Mock())
|
||||
def test_follow(self):
|
||||
self.p.call_nowait('on_start')
|
||||
m = RemoteMember('1', {'restore_command': '2', 'recovery_min_apply_delay': 3, 'archive_cleanup_command': '4'})
|
||||
self.p.follow(m)
|
||||
|
||||
@@ -624,6 +627,7 @@ class TestPostgresql(unittest.TestCase):
|
||||
self.assertTrue('host replication replicator 127.0.0.1/32 md5\n' in lines)
|
||||
|
||||
@patch.object(Postgresql, 'cancellable_subprocess_call')
|
||||
@patch.object(Postgresql, 'get_major_version', Mock(return_value=90600))
|
||||
def test_custom_bootstrap(self, mock_cancellable_subprocess_call):
|
||||
self.p.config.pop('pg_hba')
|
||||
config = {'method': 'foo', 'foo': {'command': 'bar'}}
|
||||
@@ -632,7 +636,7 @@ class TestPostgresql(unittest.TestCase):
|
||||
self.assertFalse(self.p.bootstrap(config))
|
||||
|
||||
mock_cancellable_subprocess_call.return_value = 0
|
||||
with patch('subprocess.Popen', Mock(side_effect=Exception("42"))),\
|
||||
with patch('multiprocessing.Process', Mock(side_effect=Exception("42"))),\
|
||||
patch('os.path.isfile', Mock(return_value=True)),\
|
||||
patch('os.unlink', Mock()),\
|
||||
patch.object(Postgresql, 'save_configuration_files', Mock()),\
|
||||
@@ -653,10 +657,12 @@ class TestPostgresql(unittest.TestCase):
|
||||
|
||||
@patch('time.sleep', Mock())
|
||||
@patch('os.unlink', Mock())
|
||||
@patch('shutil.copy', Mock())
|
||||
@patch('os.path.isfile', Mock(return_value=True))
|
||||
@patch.object(Postgresql, 'run_bootstrap_post_init', Mock(return_value=True))
|
||||
@patch.object(Postgresql, '_custom_bootstrap', Mock(return_value=True))
|
||||
@patch.object(Postgresql, 'start', Mock(return_value=True))
|
||||
@patch.object(Postgresql, 'get_major_version', Mock(return_value=90600))
|
||||
def test_post_bootstrap(self):
|
||||
config = {'method': 'foo', 'foo': {'command': 'bar'}}
|
||||
self.p.bootstrap(config)
|
||||
@@ -680,7 +686,7 @@ class TestPostgresql(unittest.TestCase):
|
||||
self.p.set_state('stopped')
|
||||
self.p.reload_config({'authentication': {'superuser': {'username': 'p', 'password': 'p'},
|
||||
'replication': {'username': 'r', 'password': 'r'}},
|
||||
'listen': '*', 'retry_timeout': 10, 'parameters': {'hba_file': 'foo'}})
|
||||
'listen': '*', 'retry_timeout': 10, 'parameters': {'wal_level': '', 'hba_file': 'foo'}})
|
||||
with patch.object(Postgresql, 'restart', Mock()) as mock_restart:
|
||||
self.p.post_bootstrap({}, task)
|
||||
mock_restart.assert_called_once()
|
||||
@@ -717,10 +723,18 @@ class TestPostgresql(unittest.TestCase):
|
||||
self.assertEqual(self.p.get_postgres_role_from_data_directory(), 'replica')
|
||||
|
||||
def test_remove_data_directory(self):
|
||||
def _symlink(src, dst):
|
||||
try:
|
||||
os.symlink(src, dst)
|
||||
except OSError:
|
||||
if os.name == 'nt': # os.symlink under Windows needs admin rights skip it
|
||||
pass
|
||||
os.makedirs(os.path.join(self.data_dir, 'foo'))
|
||||
_symlink('foo', os.path.join(self.data_dir, 'pg_wal'))
|
||||
self.p.remove_data_directory()
|
||||
open(self.data_dir, 'w').close()
|
||||
self.p.remove_data_directory()
|
||||
os.symlink('unexisting', self.data_dir)
|
||||
_symlink('unexisting', self.data_dir)
|
||||
with patch('os.unlink', Mock(side_effect=OSError)):
|
||||
self.p.remove_data_directory()
|
||||
self.p.remove_data_directory()
|
||||
@@ -990,7 +1004,9 @@ class TestPostgresql(unittest.TestCase):
|
||||
self.p.cleanup_archive_status()
|
||||
|
||||
@patch('os.unlink', Mock())
|
||||
@patch('os.listdir', Mock(return_value=[]))
|
||||
@patch('os.path.isfile', Mock(return_value=True))
|
||||
@patch.object(Postgresql, 'read_postmaster_opts', Mock(return_value={}))
|
||||
@patch.object(Postgresql, 'single_user_mode', Mock(return_value=0))
|
||||
def test_fix_cluster_state(self):
|
||||
self.assertTrue(self.p.fix_cluster_state())
|
||||
|
||||
@@ -6,6 +6,18 @@ from patroni.postmaster import PostmasterProcess
|
||||
from six.moves import builtins
|
||||
|
||||
|
||||
class MockProcess(object):
|
||||
def __init__(self, target, args):
|
||||
self.target = target
|
||||
self.args = args
|
||||
|
||||
def start(self):
|
||||
self.target(*self.args)
|
||||
|
||||
def join(self):
|
||||
pass
|
||||
|
||||
|
||||
class TestPostmasterProcess(unittest.TestCase):
|
||||
@patch('psutil.Process.__init__', Mock())
|
||||
def test_init(self):
|
||||
@@ -82,18 +94,23 @@ class TestPostmasterProcess(unittest.TestCase):
|
||||
self.assertIsNone(proc.wait_for_user_backends_to_close())
|
||||
|
||||
@patch('subprocess.Popen')
|
||||
@patch('os.setsid', Mock(), create=True)
|
||||
@patch('multiprocessing.Process', MockProcess)
|
||||
@patch.object(PostmasterProcess, 'from_pid')
|
||||
@patch.object(PostmasterProcess, '_from_pidfile')
|
||||
def test_start(self, mock_frompidfile, mock_frompid, mock_popen):
|
||||
mock_frompidfile.return_value._is_postmaster_process.return_value = False
|
||||
mock_frompid.return_value = "proc 123"
|
||||
mock_popen.return_value.stdout.readline.return_value = '123'
|
||||
mock_popen.return_value.pid = 123
|
||||
self.assertEqual(PostmasterProcess.start('true', '/tmp', '/tmp/test.conf', []), "proc 123")
|
||||
mock_frompid.assert_called_with(123)
|
||||
|
||||
mock_frompidfile.side_effect = psutil.NoSuchProcess(123)
|
||||
self.assertEqual(PostmasterProcess.start('true', '/tmp', '/tmp/test.conf', []), "proc 123")
|
||||
|
||||
mock_popen.side_effect = Exception
|
||||
self.assertIsNone(PostmasterProcess.start('true', '/tmp', '/tmp/test.conf', []))
|
||||
|
||||
@patch('psutil.Process.__init__', Mock(side_effect=psutil.NoSuchProcess(123)))
|
||||
def test_read_postmaster_pidfile(self):
|
||||
with patch.object(builtins, 'open', Mock(side_effect=IOError)):
|
||||
|
||||
@@ -2,6 +2,7 @@ import ctypes
|
||||
import patroni.watchdog.linux as linuxwd
|
||||
import sys
|
||||
import unittest
|
||||
import os
|
||||
|
||||
from mock import patch, Mock, PropertyMock
|
||||
from patroni.watchdog import Watchdog, WatchdogError
|
||||
@@ -61,6 +62,7 @@ def mock_close(fd):
|
||||
mock_devices[fd].open = False
|
||||
|
||||
|
||||
@unittest.skipIf(os.name == 'nt', "Windows not supported")
|
||||
@patch('os.open', mock_open)
|
||||
@patch('os.write', mock_write)
|
||||
@patch('os.close', mock_close)
|
||||
@@ -174,6 +176,7 @@ class TestNullWatchdog(unittest.TestCase):
|
||||
self.assertIsInstance(NullWatchdog.from_config({}), NullWatchdog)
|
||||
|
||||
|
||||
@unittest.skipIf(os.name == 'nt', "Windows not supported")
|
||||
class TestLinuxWatchdogDevice(unittest.TestCase):
|
||||
|
||||
def setUp(self):
|
||||
|
||||
@@ -17,6 +17,7 @@ class MockKazooClient(Mock):
|
||||
|
||||
def __init__(self, *args, **kwargs):
|
||||
super(MockKazooClient, self).__init__()
|
||||
self._session_timeout = 30
|
||||
|
||||
@property
|
||||
def client_id(self):
|
||||
@@ -24,7 +25,7 @@ class MockKazooClient(Mock):
|
||||
|
||||
@staticmethod
|
||||
def retry(func, *args, **kwargs):
|
||||
func(*args, **kwargs)
|
||||
return func(*args, **kwargs)
|
||||
|
||||
def get(self, path, watch=None):
|
||||
if not isinstance(path, six.string_types):
|
||||
|
||||
Reference in New Issue
Block a user