Compare commits

...
17 Commits
Author SHA1 Message Date
Alexander KukushkinandGitHub 6909ce0c7a Release 1.5.6 (#1020)
* Update release notes
* Bump version
2019-04-03 14:44:31 +02:00
Alexander KukushkinandGitHub cd9d9ca0c3 Make sure we don't enforce ssl_version (#1010)
Fixes https://github.com/zalando/patroni/issues/1009
2019-04-02 16:49:32 +02:00
Alexander KukushkinandGitHub a0a2da238e Couple of minor improvements (#1019)
1. Fix race condition on shutdown. It is very annoying when you cancel behave tests but postgres remains running.
2. Dump pg_controldata output to logs when "recovering" stopped postgres. It will help to investigate some annoying issues.
2019-04-02 16:49:21 +02:00
Pavlo GolubandAlexander Kukushkin b53a29c022 Fix unit-tests for Windows (#1014)
Closes #1013
2019-04-02 13:58:17 +02:00
Alexander KukushkinandGitHub e38fe78b56 Fix callbacks behavior (mostly for standby cluster) (#998)
First of all, this patch changes the behavior of `on_start`/`on_restart` callbacks, they will be called only when postgres is started or restarted without role changes. In case if the member is promoted or demoted only the `on_role_change` callback will be executed. `on_role_change` was never called for standby leader, only `on_start`/`on_restart` and with a wrong role argument.
Before that `on_role_change` was never called for standby leader, only `on_start`/`on_restart` and with a wrong role argument.

In addition to that, the REST API will return standby_leader role for the leader of the standby cluster.

Closes https://github.com/zalando/patroni/issues/988
2019-03-29 10:28:07 +01:00
Lukas VogelandAlexander Kukushkin e059e30560 Ingore VSCode files (#1007)
Fixes #1008
2019-03-22 16:26:37 -04:00
Alexander KukushkinandGitHub 680444ae13 Reduce lock time taken by dcs.get_cluster() (#989)
`dcs.cluster` and `dcs.get_cluster()` are using the same lock resource and therefore when get_cluster call is slow due to the slowness of DCS it was also affecting the `dcs.cluster` call, which in return was making health-check requests slow.
2019-03-12 22:37:11 +01:00
Alexander KukushkinandGitHub 92720882aa Reset is_leader flag for every removal of leader key (#990)
This is the next improvement of #777
2019-03-12 22:10:46 +01:00
Alexander KukushkinandGitHub 13c88e8b7a Replace self-execute with multiprocessing.Process (#994)
In addition to that transfer postmaster pid to Patroni process with the help of multiprocessing.Pipe instead of using stdin-stdout pipes.

Closes https://github.com/zalando/patroni/issues/992
2019-03-12 10:40:37 +01:00
Alexander KukushkinandGitHub 4a4258fc3f Mock external resources (#995)
unit tests should not accidentally hit running Postgres, DCS or filesystem unless we want it explicitly.
2019-03-12 10:39:42 +01:00
Ants AasmaandAlexander Kukushkin 2204454094 Remove unnecessary usage of relpath (#1002)
`os.path.relpath` depends on being able to resolve the working directory.
This will fail if Patroni is started in a directory that is later unlinked from the filesystem, creating an unnecessary exception when loading from DCS.
2019-03-11 12:31:06 +01:00
Alexander KukushkinandGitHub 9e19b43869 Rename cluster name to demo (#1000)
* Assign hostname to haproxy container
* Tune vim config
2019-03-11 10:57:26 +01:00
Julien TachoiresandAlexander Kukushkin 13f7aede61 Wait for callback end if it could not be killed (#985)
that could happen if the script is running under a sudo
2019-03-11 10:35:36 +01:00
Julien TachoiresandAlexander Kukushkin 8a7ba57457 Clean target directory when pg_wal|pg_xlog is a symlink (#997) 2019-03-11 10:33:46 +01:00
Alexander KukushkinandGitHub c64d51f79c Better support for static etcd cluster (#986)
if the `etcd.use_proxies` is set to true, Patroni will stick to the list of hosts specified in the `etcd.hosts` and avoid doing topology discovery. Such mode might be useful when you know that you connect to the etcd cluster via the set of proxies or when th etcd cluster has static topology.
2019-03-07 11:36:36 +01:00
Alexander KukushkinandGitHub f0990532dc Update docker-compose demo cluster (#980)
1. Multi-stage build with an extensive cleanup of useless files and optional image compression
2. Start three-node etcd cluster
3. Start three-node Patroni cluster
4. One container with haproxy
5. All container names are prefixed with "demo-" and don't have suffixes
6. Decommission dev_patroni_cluster.sh script, docker-compose is now standard de-facto.
7. Provide more examples in the docker/README.md
2019-03-07 11:32:03 +01:00
Alexander KukushkinandGitHub e670122c80 Use create_replica_methods from standby_cluster for replica bootstrap (#981)
It might happen that the standby cluster is configured to be created and replay WAL files from the different source than when it is running not in standby mode.  This is necessary to avoid writing WAL files and backups into the old place after promotion.

The easiest way to achieve such behavior is passing RemoteMember object to `Postgresql.clone` method instead of the usual Member object.
2019-02-21 11:37:50 +01:00
43 changed files with 828 additions and 553 deletions
+3
View File
@@ -54,3 +54,6 @@ docs/source/_templates/
# Pycharm IDE
.idea/
#VSCode IDE
.vscode/
+6 -2
View File
@@ -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
View File
@@ -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
View File
@@ -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
View File
@@ -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)
-104
View File
@@ -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
View File
@@ -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
+5
View File
@@ -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
+1
View File
@@ -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
View File
@@ -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
+43
View File
@@ -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
-------------
+2 -2
View File
@@ -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
+17 -4
View File
@@ -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
+4 -3
View File
@@ -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)
+16 -5
View File
@@ -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
View File
@@ -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()
+3
View File
@@ -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:
+7 -2
View File
@@ -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
View File
@@ -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
View File
@@ -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
View File
@@ -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.
+9 -5
View File
@@ -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
View File
@@ -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):
+15 -8
View File
@@ -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"""
+10 -4
View File
@@ -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
View File
@@ -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
View File
@@ -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
View File
@@ -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
View File
@@ -1 +1 @@
__version__ = '1.5.5'
__version__ = '1.5.6'
+3 -2
View File
@@ -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:
+3
View File
@@ -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])
+1 -1
View File
@@ -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'",
+1 -1
View File
@@ -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
View File
@@ -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
View File
@@ -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
View File
@@ -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
+7 -5
View File
@@ -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)
+6 -9
View File
@@ -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
+20 -4
View File
@@ -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())
+18 -1
View File
@@ -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)):
+3
View File
@@ -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):
+2 -1
View File
@@ -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):