psycopg2/tests/test_replication.py

201 lines
6.9 KiB
Python
Raw Normal View History

2015-10-15 19:01:43 +03:00
#!/usr/bin/env python
# test_replication.py - unit test for replication protocol
#
# Copyright (C) 2015 Daniele Varrazzo <daniele.varrazzo@gmail.com>
#
# psycopg2 is free software: you can redistribute it and/or modify it
# under the terms of the GNU Lesser General Public License as published
# by the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
#
# In addition, as a special exception, the copyright holders give
# permission to link this program with the OpenSSL library (or with
# modified versions of OpenSSL that use the same license as OpenSSL),
# and distribute linked combinations including the two.
#
# You must obey the GNU Lesser General Public License in all respects for
# all of the code used other than OpenSSL.
#
# psycopg2 is distributed in the hope that it will be useful, but WITHOUT
# ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or
# FITNESS FOR A PARTICULAR PURPOSE. See the GNU Lesser General Public
# License for more details.
import psycopg2
import psycopg2.extensions
from psycopg2.extras import PhysicalReplicationConnection, LogicalReplicationConnection
from psycopg2.extras import StopReplication
2015-10-15 19:01:43 +03:00
2015-10-19 18:02:18 +03:00
import testconfig
2015-10-15 19:01:43 +03:00
from testutils import unittest
from testutils import skip_before_postgres
from testutils import ConnectingTestCase
class ReplicationTestCase(ConnectingTestCase):
def setUp(self):
2015-10-19 18:02:18 +03:00
if not testconfig.repl_dsn:
2015-10-16 17:36:03 +03:00
self.skipTest("replication tests disabled by default")
2015-10-19 18:02:18 +03:00
2015-10-15 19:01:43 +03:00
super(ReplicationTestCase, self).setUp()
2015-10-19 18:02:18 +03:00
self.slot = testconfig.repl_slot
2015-10-15 19:01:43 +03:00
self._slots = []
def tearDown(self):
# first close all connections, as they might keep the slot(s) active
super(ReplicationTestCase, self).tearDown()
import time
time.sleep(0.025) # sometimes the slot is still active, wait a little
2015-10-15 19:01:43 +03:00
if self._slots:
kill_conn = self.connect()
2015-10-15 19:01:43 +03:00
if kill_conn:
kill_cur = kill_conn.cursor()
for slot in self._slots:
kill_cur.execute("SELECT pg_drop_replication_slot(%s)", (slot,))
kill_conn.commit()
2015-10-15 19:01:43 +03:00
kill_conn.close()
2015-10-19 18:02:18 +03:00
def create_replication_slot(self, cur, slot_name=testconfig.repl_slot, **kwargs):
2015-10-15 19:01:43 +03:00
cur.create_replication_slot(slot_name, **kwargs)
self._slots.append(slot_name)
2015-10-19 18:02:18 +03:00
def drop_replication_slot(self, cur, slot_name=testconfig.repl_slot):
2015-10-15 19:01:43 +03:00
cur.drop_replication_slot(slot_name)
self._slots.remove(slot_name)
2015-10-19 18:02:18 +03:00
# generate some events for our replication stream
def make_replication_events(self):
conn = self.connect()
if conn is None: return
cur = conn.cursor()
try:
cur.execute("DROP TABLE dummy1")
except psycopg2.ProgrammingError:
conn.rollback()
cur.execute("CREATE TABLE dummy1 AS SELECT * FROM generate_series(1, 5) AS id")
conn.commit()
2015-10-15 19:01:43 +03:00
class ReplicationTest(ReplicationTestCase):
@skip_before_postgres(9, 0)
def test_physical_replication_connection(self):
conn = self.repl_connect(connection_factory=PhysicalReplicationConnection)
if conn is None: return
cur = conn.cursor()
cur.execute("IDENTIFY_SYSTEM")
cur.fetchall()
@skip_before_postgres(9, 4)
def test_logical_replication_connection(self):
conn = self.repl_connect(connection_factory=LogicalReplicationConnection)
if conn is None: return
cur = conn.cursor()
cur.execute("IDENTIFY_SYSTEM")
cur.fetchall()
@skip_before_postgres(9, 4) # slots require 9.4
def test_create_replication_slot(self):
conn = self.repl_connect(connection_factory=PhysicalReplicationConnection)
if conn is None: return
cur = conn.cursor()
2015-10-19 18:02:18 +03:00
self.create_replication_slot(cur)
self.assertRaises(psycopg2.ProgrammingError, self.create_replication_slot, cur)
2015-10-15 19:01:43 +03:00
2015-10-16 17:36:03 +03:00
@skip_before_postgres(9, 4) # slots require 9.4
def test_start_on_missing_replication_slot(self):
conn = self.repl_connect(connection_factory=PhysicalReplicationConnection)
if conn is None: return
cur = conn.cursor()
2015-10-19 18:02:18 +03:00
self.assertRaises(psycopg2.ProgrammingError, cur.start_replication, self.slot)
2015-10-16 17:36:03 +03:00
2015-10-19 18:02:18 +03:00
self.create_replication_slot(cur)
cur.start_replication(self.slot)
2015-10-16 17:36:03 +03:00
@skip_before_postgres(9, 4) # slots require 9.4
def test_start_and_recover_from_error(self):
conn = self.repl_connect(connection_factory=LogicalReplicationConnection)
if conn is None: return
cur = conn.cursor()
self.create_replication_slot(cur, output_plugin='test_decoding')
# try with invalid options
cur.start_replication(slot_name=self.slot, options={'invalid_param': 'value'})
def consume(msg):
pass
# we don't see the error from the server before we try to read the data
self.assertRaises(psycopg2.DataError, cur.consume_stream, consume)
# try with correct command
cur.start_replication(slot_name=self.slot)
@skip_before_postgres(9, 4) # slots require 9.4
def test_stop_replication(self):
conn = self.repl_connect(connection_factory=LogicalReplicationConnection)
if conn is None: return
cur = conn.cursor()
2015-10-19 18:02:18 +03:00
self.create_replication_slot(cur, output_plugin='test_decoding')
2015-10-19 18:02:18 +03:00
self.make_replication_events()
2015-10-19 18:02:18 +03:00
cur.start_replication(self.slot)
def consume(msg):
raise StopReplication()
self.assertRaises(StopReplication, cur.consume_stream, consume)
2015-10-15 19:01:43 +03:00
class AsyncReplicationTest(ReplicationTestCase):
2015-10-19 18:02:18 +03:00
@skip_before_postgres(9, 4) # slots require 9.4
2015-10-15 19:01:43 +03:00
def test_async_replication(self):
conn = self.repl_connect(connection_factory=LogicalReplicationConnection, async=1)
if conn is None: return
self.wait(conn)
cur = conn.cursor()
2015-10-19 18:02:18 +03:00
self.create_replication_slot(cur, output_plugin='test_decoding')
2015-10-15 19:01:43 +03:00
self.wait(cur)
2015-10-19 18:02:18 +03:00
cur.start_replication(self.slot)
2015-10-15 19:01:43 +03:00
self.wait(cur)
2015-10-19 18:02:18 +03:00
self.make_replication_events()
self.msg_count = 0
def consume(msg):
# just check the methods
log = "%s: %s" % (cur.io_timestamp, repr(msg))
2015-10-19 18:02:18 +03:00
self.msg_count += 1
if self.msg_count > 3:
cur.send_feedback(reply=True)
2015-10-19 18:02:18 +03:00
raise StopReplication()
cur.send_feedback(flush_lsn=msg.data_start)
# cannot be used in asynchronous mode
self.assertRaises(psycopg2.ProgrammingError, cur.consume_stream, consume)
2015-10-19 18:02:18 +03:00
def process_stream():
from select import select
while True:
msg = cur.read_message()
2015-10-19 18:02:18 +03:00
if msg:
consume(msg)
else:
select([cur], [], [])
self.assertRaises(StopReplication, process_stream)
2015-10-15 19:01:43 +03:00
def test_suite():
return unittest.TestLoader().loadTestsFromName(__name__)
if __name__ == "__main__":
unittest.main()