From c951fa79cda2d6210636e6050d6790f1527da2a9 Mon Sep 17 00:00:00 2001 From: Marcel Meier Date: Wed, 12 Feb 2025 11:40:03 +0100 Subject: [PATCH 1/6] added changed from pr 915 --- postgres-appliance/major_upgrade/inplace_upgrade.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/postgres-appliance/major_upgrade/inplace_upgrade.py b/postgres-appliance/major_upgrade/inplace_upgrade.py index 0390cfbd4..f71b577be 100644 --- a/postgres-appliance/major_upgrade/inplace_upgrade.py +++ b/postgres-appliance/major_upgrade/inplace_upgrade.py @@ -172,13 +172,15 @@ def ensure_replicas_state(self, cluster): """ self.replica_connections = {} streaming = {a: l for a, l in self.postgresql.query( - ("SELECT client_addr, pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_current_{0}_{1}()," + ("SELECT client_addr, application_name, pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_current_{0}_{1}()," " COALESCE(replay_{1}, '0/0'))::bigint FROM pg_catalog.pg_stat_replication") .format(self.postgresql.wal_name, self.postgresql.lsn_name))} def ensure_replica_state(member): ip = member.conn_kwargs().get('host') lag = streaming.get(ip) + if os.getenv('USE_APPLICATION_NAME_IN_UPGRADE'): + lag = streaming.get(member.name) if lag is None: return logger.error('Member %s is not streaming from the primary', member.name) if lag > 16*1024*1024: From f340fd0f1b11aac9d94ee361f5fb5232a73f33f6 Mon Sep 17 00:00:00 2001 From: Marcel Meier Date: Wed, 12 Feb 2025 12:06:46 +0100 Subject: [PATCH 2/6] hopefully fixing comments from pr 915 --- postgres-appliance/major_upgrade/inplace_upgrade.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/postgres-appliance/major_upgrade/inplace_upgrade.py b/postgres-appliance/major_upgrade/inplace_upgrade.py index f71b577be..169f7216c 100644 --- a/postgres-appliance/major_upgrade/inplace_upgrade.py +++ b/postgres-appliance/major_upgrade/inplace_upgrade.py @@ -171,16 +171,17 @@ def ensure_replicas_state(self, cluster): to all of them and puts into the `self.replica_connections` dict for a future usage. """ self.replica_connections = {} - streaming = {a: l for a, l in self.postgresql.query( + streaming = {(addr, name): lag for addr, name, lag in self.postgresql.query( ("SELECT client_addr, application_name, pg_catalog.pg_{0}_{1}_diff(pg_catalog.pg_current_{0}_{1}()," " COALESCE(replay_{1}, '0/0'))::bigint FROM pg_catalog.pg_stat_replication") .format(self.postgresql.wal_name, self.postgresql.lsn_name))} def ensure_replica_state(member): ip = member.conn_kwargs().get('host') - lag = streaming.get(ip) - if os.getenv('USE_APPLICATION_NAME_IN_UPGRADE'): - lag = streaming.get(member.name) + lag = streaming.get((ip, member.name)) + if lag is None and os.getenv('USE_APPLICATION_NAME_IN_UPGRADE'): + # Try looking up by any IP address matching the member name + lag = next((lag for (_, app_name), lag in streaming.items() if app_name == member.name), None) if lag is None: return logger.error('Member %s is not streaming from the primary', member.name) if lag > 16*1024*1024: From 6a94f40c639a7beacb3d8577a35933e66dfe7173 Mon Sep 17 00:00:00 2001 From: Marcel Meier Date: Mon, 17 Feb 2025 16:31:52 +0100 Subject: [PATCH 3/6] added proxy addresses to rsync config --- .../major_upgrade/inplace_upgrade.py | 21 ++++++++++++++----- 1 file changed, 16 insertions(+), 5 deletions(-) diff --git a/postgres-appliance/major_upgrade/inplace_upgrade.py b/postgres-appliance/major_upgrade/inplace_upgrade.py index 169f7216c..f7884e32a 100644 --- a/postgres-appliance/major_upgrade/inplace_upgrade.py +++ b/postgres-appliance/major_upgrade/inplace_upgrade.py @@ -200,7 +200,12 @@ def ensure_replica_state(member): cur.execute('SELECT pg_catalog.pg_is_in_recovery()') if not cur.fetchone()[0]: return logger.error('Member %s is not running as replica!', member.name) - self.replica_connections[member.name] = (ip, cur) + + # Find the streaming address for this member + streaming_addr = next((addr for (addr, app_name), _ in streaming.items() + if app_name == member.name), None) + + self.replica_connections[member.name] = (ip, cur, streaming_addr) return True return all(ensure_replica_state(member) for member in cluster.members if member.name != self.postgresql.name) @@ -246,7 +251,7 @@ def wait_for_replicas(self, checkpoint_lsn): for _ in polling_loop(60): synced = True - for name, (_, cur) in self.replica_connections.items(): + for name, (_, cur, _) in self.replica_connections.items(): prev = status.get(name) if prev and prev >= checkpoint_lsn: continue @@ -279,8 +284,14 @@ def create_rsyncd_configs(self): secrets_file = os.path.join(self.rsyncd_conf_dir, 'rsyncd.secrets') auth_users = ','.join(self.replica_connections.keys()) - replica_ips = ','.join(str(v[0]) for v in self.replica_connections.values()) + # Collect both connection and streaming IPs from replica_connections + replica_ips = {str(v[0]) for v in self.replica_connections.values()} # Connection IPs + replica_ips.update(str(v[2]) for v in self.replica_connections.values() if v[2]) # Streaming IPs + + # Filter out None values and join IPs + replica_ips = ','.join(filter(None, replica_ips)) + with open(self.rsyncd_conf, 'w') as f: f.write("""port = {0} use chroot = false @@ -324,7 +335,7 @@ def stop_rsyncd(self): logger.error('Failed to remove %s: %r', self.rsyncd_conf_dir, e) def checkpoint(self, member): - name, (_, cur) = member + name, (_, cur, _) = member try: cur.execute('CHECKPOINT') return name, True @@ -338,7 +349,7 @@ def rsync_replicas(self, primary_ip): logger.info('Notifying replicas %s to start rsync', ','.join(self.replica_connections.keys())) ret = True status = {} - for name, (ip, cur) in self.replica_connections.items(): + for name, (ip, cur, _) in self.replica_connections.items(): try: cur.execute("SELECT pg_catalog.pg_backend_pid()") pid = cur.fetchone()[0] From 7f5f604712d97b7b9d66702bd3ed4ea0fe459075 Mon Sep 17 00:00:00 2001 From: Marcel Meier Date: Tue, 18 Feb 2025 10:01:56 +0100 Subject: [PATCH 4/6] updated comments --- postgres-appliance/major_upgrade/inplace_upgrade.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/postgres-appliance/major_upgrade/inplace_upgrade.py b/postgres-appliance/major_upgrade/inplace_upgrade.py index f7884e32a..3f4516627 100644 --- a/postgres-appliance/major_upgrade/inplace_upgrade.py +++ b/postgres-appliance/major_upgrade/inplace_upgrade.py @@ -201,11 +201,12 @@ def ensure_replica_state(member): if not cur.fetchone()[0]: return logger.error('Member %s is not running as replica!', member.name) - # Find the streaming address for this member - streaming_addr = next((addr for (addr, app_name), _ in streaming.items() + # determine the "client_ip" seen from leader + # differs from "ip" when using proxy sidecars (service mesh e.g. istio) + client_ip = next((addr for (addr, app_name), _ in streaming.items() if app_name == member.name), None) - self.replica_connections[member.name] = (ip, cur, streaming_addr) + self.replica_connections[member.name] = (ip, cur, client_ip) return True return all(ensure_replica_state(member) for member in cluster.members if member.name != self.postgresql.name) @@ -285,7 +286,7 @@ def create_rsyncd_configs(self): auth_users = ','.join(self.replica_connections.keys()) - # Collect both connection and streaming IPs from replica_connections + # Collect both host IP and the client IP in case of proxy sidecars replica_ips = {str(v[0]) for v in self.replica_connections.values()} # Connection IPs replica_ips.update(str(v[2]) for v in self.replica_connections.values() if v[2]) # Streaming IPs From 4a60d1de2d917a01cdc53db984b48121bb19a653 Mon Sep 17 00:00:00 2001 From: Marcel Meier Date: Wed, 19 Feb 2025 14:19:49 +0100 Subject: [PATCH 5/6] added env to ENVIRONMENT.rst --- ENVIRONMENT.rst | 1 + 1 file changed, 1 insertion(+) diff --git a/ENVIRONMENT.rst b/ENVIRONMENT.rst index ea2e3df81..04a0fcbf0 100644 --- a/ENVIRONMENT.rst +++ b/ENVIRONMENT.rst @@ -110,6 +110,7 @@ Environment Configuration Settings - **ENABLE_WAL_PATH_COMPAT**: old Spilo images were generating wal path in the backup store using the following template ``/spilo/{WAL_BUCKET_SCOPE_PREFIX}{SCOPE}{WAL_BUCKET_SCOPE_SUFFIX}/wal/``, while new images adding one additional directory (``{PGVERSION}``) to the end. In order to avoid (unlikely) issues with restoring WALs (from S3/GC/and so on) when switching to ``spilo-13`` please set the ``ENABLE_WAL_PATH_COMPAT=true`` when deploying old cluster with ``spilo-13`` for the first time. After that the environment variable could be removed. Change of the WAL path also mean that backups stored in the old location will not be cleaned up automatically. - **WALE_DISABLE_S3_SSE**, **WALG_DISABLE_S3_SSE**: by default wal-e/wal-g are configured to encrypt files uploaded to S3. In order to disable it you can set this environment variable to ``true``. - **USE_OLD_LOCALES**: whether to use old locales from Ubuntu 18.04 in the Ubuntu 22.04-based image. Default is false. +- **USE_APPLICATION_NAME_IN_UPGRADE**: whether to use the application name in the upgrade script. Default is false. Usable for usage with service meshs. wal-g ----- From bbd72dd289e0164073d9bd8a17631e28b56c301a Mon Sep 17 00:00:00 2001 From: Marcel Meier Date: Thu, 6 Mar 2025 13:52:02 +0100 Subject: [PATCH 6/6] adapted default behaviour of USE_APPLICATION_NAME_IN_UPGRADE --- postgres-appliance/major_upgrade/inplace_upgrade.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/postgres-appliance/major_upgrade/inplace_upgrade.py b/postgres-appliance/major_upgrade/inplace_upgrade.py index 3f4516627..adb4f0d19 100644 --- a/postgres-appliance/major_upgrade/inplace_upgrade.py +++ b/postgres-appliance/major_upgrade/inplace_upgrade.py @@ -179,7 +179,7 @@ def ensure_replicas_state(self, cluster): def ensure_replica_state(member): ip = member.conn_kwargs().get('host') lag = streaming.get((ip, member.name)) - if lag is None and os.getenv('USE_APPLICATION_NAME_IN_UPGRADE'): + if lag is None and os.getenv('USE_APPLICATION_NAME_IN_UPGRADE') == 'true': # Try looking up by any IP address matching the member name lag = next((lag for (_, app_name), lag in streaming.items() if app_name == member.name), None) if lag is None: