22 Commits

Author SHA1 Message Date
75c5d76047 Add a metric for tracking the number of granted locks 2025-09-23 01:14:50 -04:00
43cd162313 Refactor to make pylint happy 2025-09-23 01:12:49 -04:00
29bfd07dad Run code through black 2025-07-14 01:58:08 -04:00
60589c2058 Update undiscovered item behavior 2025-07-14 01:39:57 -04:00
ea3aca3455 Filter out initial logical replication sync workers
* Slots are crated for each table during the initial sync, which only
  live for the duration of the copy for that table.

* The initial sync backend workers' names are also based on the table
  being copied.

* Filter out the above in Zabbix discovery based on application_name and
  slot_name.
2025-07-07 13:15:03 -04:00
cc71547f5f Include application_name in replication discovery data 2025-07-07 13:07:11 -04:00
107d5056d6 Merge branch 'main' into develop 2025-07-06 03:37:37 -04:00
e930178e9c Merge branch 'rel/1.0.4' 2025-07-06 03:35:39 -04:00
5afc940df8 Bump release version to 1.0.4 2025-07-06 03:35:00 -04:00
3961aa3448 Add Gentoo ebuild for 1.0.4-rc1 2025-07-06 03:32:27 -04:00
8fd57032e7 Update Makefile to support rc versions, bump version
* Support packaging RC versions for deb and rpm packages

* Bump version to 1.0.4-rc1
2025-07-06 03:28:36 -04:00
5ede7dea07 Fix package-all make target
* Fix the package-all target in the Makefile

* Remove storage of raw sequence stats
2025-07-05 12:28:36 -04:00
5ea007c3f6 Merge branch 'dev/sequences' into develop 2025-07-05 01:18:01 -04:00
7cb0f7ad40 Fix typo in activity query 2025-07-05 01:17:35 -04:00
83fa12ec54 Correct types for slot LSN lag metrics 2025-07-04 02:46:25 -04:00
98b74d9aed Add sequence metrics to Zabbix template 2025-07-04 02:45:43 -04:00
45953848e2 Remove some type casts from working around Decimal types 2025-07-03 10:16:20 -04:00
6116f4f885 Fix missing json default 2025-07-03 02:08:45 -04:00
24d1214855 Teach json how to serialize decimals 2025-07-03 01:47:06 -04:00
3c39d8aa97 Remove CIDR prefix from replication id 2025-07-03 01:13:58 -04:00
ebb084aa9d Add initial query for sequence usage 2025-07-01 02:29:56 -04:00
86d5e8917b Merge branch 'dev/io_stats' into develop 2025-07-01 01:30:21 -04:00
9 changed files with 1110 additions and 466 deletions

74
GENTOO/pgmon-1.0.4.ebuild Normal file
View File

@@ -0,0 +1,74 @@
# Copyright 2024 Gentoo Authors
# Distributed under the terms of the GNU General Public License v2
EAPI=8
PYTHON_COMPAT=( python3_{6..13} )
inherit python-r1 systemd
DESCRIPTION="PostgreSQL monitoring bridge"
HOMEPAGE="None"
LICENSE="BSD"
SLOT="0"
KEYWORDS="amd64"
SRC_URI="https://code2.shh-dot-com.org/james/${PN}/releases/download/v${PV}/${P}.tar.bz2"
IUSE="-systemd"
DEPEND="
${PYTHON_DEPS}
dev-python/psycopg:2
dev-python/pyyaml
dev-python/requests
app-admin/logrotate
"
RDEPEND="${DEPEND}"
BDEPEND=""
#RESTRICT="fetch"
#S="${WORKDIR}/${PN}"
#pkg_nofetch() {
# einfo "Please download"
# einfo " - ${P}.tar.bz2"
# einfo "from ${HOMEPAGE} and place it in your DISTDIR directory."
# einfo "The file should be owned by portage:portage."
#}
src_compile() {
true
}
src_install() {
# Install init script
if ! use systemd ; then
newinitd "openrc/pgmon.initd" pgmon
newconfd "openrc/pgmon.confd" pgmon
fi
# Install systemd unit
if use systemd ; then
systemd_dounit "systemd/pgmon.service"
fi
# Install script
exeinto /usr/bin
newexe "src/pgmon.py" pgmon
# Install default config
diropts -o root -g root -m 0755
insinto /etc/pgmon
doins "sample-config/pgmon.yml"
doins "sample-config/pgmon-metrics.yml"
# Install logrotate config
insinto /etc/logrotate.d
newins "logrotate/pgmon.logrotate" pgmon
# Install man page
doman manpages/pgmon.1
}

View File

@@ -3,7 +3,22 @@ PACKAGE_NAME := pgmon
SCRIPT := src/$(PACKAGE_NAME).py SCRIPT := src/$(PACKAGE_NAME).py
VERSION := $(shell grep -m 1 '^VERSION = ' "$(SCRIPT)" | sed -ne 's/.*"\(.*\)".*/\1/p') # Figure out the version components
# Note: The release is for RPM packages, where prerelease releases are written as 0.<release>
FULL_VERSION := $(shell grep -m 1 '^VERSION = ' "$(SCRIPT)" | sed -ne 's/.*"\(.*\)".*/\1/p')
VERSION := $(shell echo $(FULL_VERSION) | sed -n 's/\(.*\)\(-rc.*\|$$\)/\1/p')
RELEASE := $(shell echo $(FULL_VERSION) | sed -n 's/.*-rc\([0-9]\+\)$$/\1/p')
ifeq ($(RELEASE),)
RPM_RELEASE := 1
RPM_VERSION := $(VERSION)-$(RPM_RELEASE)
DEB_VERSION := $(VERSION)
else
RPM_RELEASE := 0.$(RELEASE)
RPM_VERSION := $(VERSION)-$(RPM_RELEASE)
DEB_VERSION := $(VERSION)~rc$(RELEASE)
endif
# Where packages are built # Where packages are built
BUILD_DIR := build BUILD_DIR := build
@@ -23,18 +38,22 @@ SUPPORTED := ubuntu-20.04 \
# These targets are the main ones to use for most things. # These targets are the main ones to use for most things.
## ##
.PHONY: all clean tgz test query-tests install-common install-openrc install-systemd .PHONY: all clean tgz lint format test query-tests install-common install-openrc install-systemd
all: package-all
version:
@echo "full version=$(FULL_VERSION) version=$(VERSION) rel=$(RELEASE) rpm=$(RPM_VERSION) deb=$(DEB_VERSION)"
# Build all packages # Build all packages
.PHONY: package-all .PHONY: package-all
all: $(foreach distro_release, $(SUPPORTED), package-$(distro_release)) package-all: $(foreach distro_release, $(SUPPORTED), package-$(distro_release))
# Gentoo package (tar.gz) creation # Gentoo package (tar.gz) creation
.PHONY: package-gentoo .PHONY: package-gentoo
package-gentoo: package-gentoo:
mkdir -p $(BUILD_DIR)/gentoo mkdir -p $(BUILD_DIR)/gentoo
tar --transform "s,^,$(PACKAGE_NAME)-$(VERSION)/," -acjf $(BUILD_DIR)/gentoo/$(PACKAGE_NAME)-$(VERSION).tar.bz2 --exclude .gitignore $(shell git ls-tree --full-tree --name-only -r HEAD) tar --transform "s,^,$(PACKAGE_NAME)-$(FULL_VERSION)/," -acjf $(BUILD_DIR)/gentoo/$(PACKAGE_NAME)-$(FULL_VERSION).tar.bz2 --exclude .gitignore $(shell git ls-tree --full-tree --name-only -r HEAD)
# Create a deb package # Create a deb package
@@ -55,12 +74,24 @@ tgz:
rm -rf $(BUILD_DIR)/tgz/root rm -rf $(BUILD_DIR)/tgz/root
mkdir -p $(BUILD_DIR)/tgz/root mkdir -p $(BUILD_DIR)/tgz/root
$(MAKE) install-openrc DESTDIR=$(BUILD_DIR)/tgz/root $(MAKE) install-openrc DESTDIR=$(BUILD_DIR)/tgz/root
tar -cz -f $(BUILD_DIR)/tgz/$(PACKAGE_NAME)-$(VERSION).tgz -C $(BUILD_DIR)/tgz/root . tar -cz -f $(BUILD_DIR)/tgz/$(PACKAGE_NAME)-$(FULL_VERSION).tgz -C $(BUILD_DIR)/tgz/root .
# Clean up the build directory # Clean up the build directory
clean: clean:
rm -rf $(BUILD_DIR) rm -rf $(BUILD_DIR)
# Check for lint
lint:
pylint src/pgmon.py
pylint src/test_pgmon.py
black --check --diff src/pgmon.py
black --check --diff src/test_pgmon.py
# Format the code using black
format:
black src/pgmon.py
black src/test_pylint.py
# Run unit tests for the script # Run unit tests for the script
test: test:
cd src ; python3 -m unittest cd src ; python3 -m unittest
@@ -129,28 +160,28 @@ debian-%-install-test:
docker run --rm \ docker run --rm \
-v ./$(BUILD_DIR):/output \ -v ./$(BUILD_DIR):/output \
debian:$* \ debian:$* \
bash -c 'apt-get update && apt-get install -y /output/$(PACKAGE_NAME)-$(VERSION)-debian-$*.deb' bash -c 'apt-get update && apt-get install -y /output/$(PACKAGE_NAME)-$(DEB_VERSION)-debian-$*.deb'
# Run a RedHat install test # Run a RedHat install test
rockylinux-%-install-test: rockylinux-%-install-test:
docker run --rm \ docker run --rm \
-v ./$(BUILD_DIR):/output \ -v ./$(BUILD_DIR):/output \
rockylinux:$* \ rockylinux:$* \
bash -c 'dnf makecache && dnf install -y /output/$(PACKAGE_NAME)-$(VERSION)-1.el$*.noarch.rpm' bash -c 'dnf makecache && dnf install -y /output/$(PACKAGE_NAME)-$(RPM_VERSION).el$*.noarch.rpm'
# Run an Ubuntu install test # Run an Ubuntu install test
ubuntu-%-install-test: ubuntu-%-install-test:
docker run --rm \ docker run --rm \
-v ./$(BUILD_DIR):/output \ -v ./$(BUILD_DIR):/output \
ubuntu:$* \ ubuntu:$* \
bash -c 'apt-get update && apt-get install -y /output/$(PACKAGE_NAME)-$(VERSION)-ubuntu-$*.deb' bash -c 'apt-get update && apt-get install -y /output/$(PACKAGE_NAME)-$(DEB_VERSION)-ubuntu-$*.deb'
# Run an OracleLinux install test (this is for EL7 since CentOS7 images no longer exist) # Run an OracleLinux install test (this is for EL7 since CentOS7 images no longer exist)
oraclelinux-%-install-test: oraclelinux-%-install-test:
docker run --rm \ docker run --rm \
-v ./$(BUILD_DIR):/output \ -v ./$(BUILD_DIR):/output \
oraclelinux:7 \ oraclelinux:7 \
bash -c 'yum makecache && yum install -y /output/$(PACKAGE_NAME)-$(VERSION)-1.el7.noarch.rpm' bash -c 'yum makecache && yum install -y /output/$(PACKAGE_NAME)-$(RPM_VERSION).el7.noarch.rpm'
# Run a Gentoo install test # Run a Gentoo install test
gentoo-install-test: gentoo-install-test:
@@ -192,28 +223,28 @@ package-image-%:
actually-package-debian-%: actually-package-debian-%:
$(MAKE) install-systemd DESTDIR=/output/debian-$* $(MAKE) install-systemd DESTDIR=/output/debian-$*
cp -r --preserve=mode DEBIAN /output/debian-$*/ cp -r --preserve=mode DEBIAN /output/debian-$*/
dpkg-deb -Zgzip --build /output/debian-$* "/output/$(PACKAGE_NAME)-$(VERSION)-debian-$*.deb" dpkg-deb -Zgzip --build /output/debian-$* "/output/$(PACKAGE_NAME)-$(DEB_VERSION)-debian-$*.deb"
# RedHat package creation # RedHat package creation
actually-package-rockylinux-%: actually-package-rockylinux-%:
mkdir -p /output/rockylinux-$*/{BUILD,RPMS,SOURCES,SPECS,SRPMS} mkdir -p /output/rockylinux-$*/{BUILD,RPMS,SOURCES,SPECS,SRPMS}
sed -e "s/@@VERSION@@/$(VERSION)/g" RPM/$(PACKAGE_NAME).spec > /output/rockylinux-$*/SPECS/$(PACKAGE_NAME).spec sed -e "s/@@VERSION@@/$(VERSION)/g" -e "s/@@RELEASE@@/$(RPM_RELEASE)/g" RPM/$(PACKAGE_NAME).spec > /output/rockylinux-$*/SPECS/$(PACKAGE_NAME).spec
rpmbuild --define '_topdir /output/rockylinux-$*' \ rpmbuild --define '_topdir /output/rockylinux-$*' \
--define 'version $(VERSION)' \ --define 'version $(RPM_VERSION)' \
-bb /output/rockylinux-$*/SPECS/$(PACKAGE_NAME).spec -bb /output/rockylinux-$*/SPECS/$(PACKAGE_NAME).spec
cp /output/rockylinux-$*/RPMS/noarch/$(PACKAGE_NAME)-$(VERSION)-1.el$*.noarch.rpm /output/ cp /output/rockylinux-$*/RPMS/noarch/$(PACKAGE_NAME)-$(RPM_VERSION).el$*.noarch.rpm /output/
# Ubuntu package creation # Ubuntu package creation
actually-package-ubuntu-%: actually-package-ubuntu-%:
$(MAKE) install-systemd DESTDIR=/output/ubuntu-$* $(MAKE) install-systemd DESTDIR=/output/ubuntu-$*
cp -r --preserve=mode DEBIAN /output/ubuntu-$*/ cp -r --preserve=mode DEBIAN /output/ubuntu-$*/
dpkg-deb -Zgzip --build /output/ubuntu-$* "/output/$(PACKAGE_NAME)-$(VERSION)-ubuntu-$*.deb" dpkg-deb -Zgzip --build /output/ubuntu-$* "/output/$(PACKAGE_NAME)-$(DEB_VERSION)-ubuntu-$*.deb"
# OracleLinux package creation # OracleLinux package creation
actually-package-oraclelinux-%: actually-package-oraclelinux-%:
mkdir -p /output/oraclelinux-$*/{BUILD,RPMS,SOURCES,SPECS,SRPMS} mkdir -p /output/oraclelinux-$*/{BUILD,RPMS,SOURCES,SPECS,SRPMS}
sed -e "s/@@VERSION@@/$(VERSION)/g" RPM/$(PACKAGE_NAME)-el7.spec > /output/oraclelinux-$*/SPECS/$(PACKAGE_NAME).spec sed -e "s/@@VERSION@@/$(VERSION)/g" -e "s/@@RELEASE@@/$(RPM_RELEASE)/g" RPM/$(PACKAGE_NAME)-el7.spec > /output/oraclelinux-$*/SPECS/$(PACKAGE_NAME).spec
rpmbuild --define '_topdir /output/oraclelinux-$*' \ rpmbuild --define '_topdir /output/oraclelinux-$*' \
--define 'version $(VERSION)' \ --define 'version $(RPM_VERSION)' \
-bb /output/oraclelinux-$*/SPECS/$(PACKAGE_NAME).spec -bb /output/oraclelinux-$*/SPECS/$(PACKAGE_NAME).spec
cp /output/oraclelinux-$*/RPMS/noarch/$(PACKAGE_NAME)-$(VERSION)-1.el$*.noarch.rpm /output/ cp /output/oraclelinux-$*/RPMS/noarch/$(PACKAGE_NAME)-$(RPM_VERSION).el$*.noarch.rpm /output/

View File

@@ -1,6 +1,6 @@
Name: pgmon Name: pgmon
Version: @@VERSION@@ Version: @@VERSION@@
Release: 1%{?dist} Release: @@RELEASE@@%{?dist}
Summary: A bridge to sit between monitoring tools and PostgreSQL Summary: A bridge to sit between monitoring tools and PostgreSQL
License: MIT License: MIT

View File

@@ -1,6 +1,6 @@
Name: pgmon Name: pgmon
Version: @@VERSION@@ Version: @@VERSION@@
Release: 1%{?dist} Release: @@RELEASE@@%{?dist}
Summary: A bridge to sit between monitoring tools and PostgreSQL Summary: A bridge to sit between monitoring tools and PostgreSQL
License: MIT License: MIT

4
pylintrc Normal file
View File

@@ -0,0 +1,4 @@
[MASTER]
py-version=3.5
disable=fixme

View File

@@ -8,14 +8,22 @@ metrics:
0: > 0: >
SELECT datname AS dbname SELECT datname AS dbname
FROM pg_database FROM pg_database
# Note: If the user lacks sufficient privileges, these fields will be NULL.
# The WHERE clause is intended to prevent Zabbix from discovering a
# connection it cannot monitor. Ideally this would generate an error
# instead.
discover_rep: discover_rep:
type: set type: set
query: query:
0: > 0: >
SELECT client_addr || '_' || regexp_replace(application_name, '[ ,]', '_', 'g') AS repid, SELECT host(client_addr) || '_' || regexp_replace(application_name, '[ ,]', '_', 'g') AS repid,
application_name,
client_addr, client_addr,
state state
FROM pg_stat_replication FROM pg_stat_replication
WHERE state IS NOT NULL
discover_slots: discover_slots:
type: set type: set
query: query:
@@ -36,6 +44,7 @@ metrics:
active active
FROM pg_replication_slots FROM pg_replication_slots
## ##
# cluster-wide metrics # cluster-wide metrics
## ##
@@ -85,29 +94,29 @@ metrics:
FROM pg_stat_bgwriter bg FROM pg_stat_bgwriter bg
CROSS JOIN pg_stat_checkpointer cp CROSS JOIN pg_stat_checkpointer cp
io_per_backend: io_per_backend:
type: set type: set
query: query:
160000: > 160000: >
SELECT backend_type, SELECT backend_type,
COALESCE(SUM(reads * op_bytes), 0)::bigint AS reads, COALESCE(SUM(reads * op_bytes), 0) AS reads,
COALESCE(SUM(read_time), 0)::bigint AS read_time, COALESCE(SUM(read_time), 0) AS read_time,
COALESCE(SUM(writes * op_bytes), 0)::bigint AS writes, COALESCE(SUM(writes * op_bytes), 0) AS writes,
COALESCE(SUM(write_time), 0)::bigint AS write_time, COALESCE(SUM(write_time), 0) AS write_time,
COALESCE(SUM(writebacks * op_bytes), 0)::bigint AS writebacks, COALESCE(SUM(writebacks * op_bytes), 0) AS writebacks,
COALESCE(SUM(writeback_time), 0)::bigint AS writeback_time, COALESCE(SUM(writeback_time), 0) AS writeback_time,
COALESCE(SUM(extends * op_bytes), 0)::bigint AS extends, COALESCE(SUM(extends * op_bytes), 0) AS extends,
COALESCE(SUM(extend_time), 0)::bigint AS extend_time, COALESCE(SUM(extend_time), 0) AS extend_time,
COALESCE(SUM(op_bytes), 0)::bigint AS op_bytes, COALESCE(SUM(op_bytes), 0) AS op_bytes,
COALESCE(SUM(hits), 0)::bigint AS hits, COALESCE(SUM(hits), 0) AS hits,
COALESCE(SUM(evictions), 0)::bigint AS evictions, COALESCE(SUM(evictions), 0) AS evictions,
COALESCE(SUM(reuses), 0)::bigint AS reuses, COALESCE(SUM(reuses), 0) AS reuses,
COALESCE(SUM(fsyncs), 0)::bigint AS fsyncs, COALESCE(SUM(fsyncs), 0) AS fsyncs,
COALESCE(SUM(fsync_time), 0)::bigint AS fsync_time COALESCE(SUM(fsync_time), 0) AS fsync_time
FROM pg_stat_io FROM pg_stat_io
GROUP BY backend_type GROUP BY backend_type
## ##
# Per-database metrics # Per-database metrics
## ##
@@ -139,7 +148,7 @@ metrics:
NULL AS sessions_abandoned, NULL AS sessions_abandoned,
NULL AS sessions_fatal, NULL AS sessions_fatal,
NULL AS sessions_killed, NULL AS sessions_killed,
extract('epoch' from stats_reset)::float AS stats_reset extract('epoch' from stats_reset) AS stats_reset
FROM pg_stat_database WHERE datname = %(dbname)s FROM pg_stat_database WHERE datname = %(dbname)s
140000: > 140000: >
SELECT numbackends, SELECT numbackends,
@@ -166,7 +175,7 @@ metrics:
sessions_abandoned, sessions_abandoned,
sessions_fatal, sessions_fatal,
sessions_killed, sessions_killed,
extract('epoch' from stats_reset)::float AS stats_reset extract('epoch' from stats_reset) AS stats_reset
FROM pg_stat_database WHERE datname = %(dbname)s FROM pg_stat_database WHERE datname = %(dbname)s
test_args: test_args:
dbname: postgres dbname: postgres
@@ -189,13 +198,52 @@ metrics:
0: > 0: >
SELECT state, SELECT state,
count(*) AS backend_count, count(*) AS backend_count,
COALESCE(EXTRACT(EPOCH FROM max(now() - state_change))::float, 0) AS max_state_time COALESCE(EXTRACT(EPOCH FROM max(now() - state_change)), 0) AS max_state_time
FROM pg_stat_activity FROM pg_stat_activity
WHERE datname = %(dbname)s WHERE datname = %(dbname)s
GROUP BY state GROUP BY state
test_args: test_args:
dbname: postgres dbname: postgres
sequence_usage:
type: value
query:
# 9.2 lacks lateral joins, the pg_sequence_last_value function, and the pg_sequences view
# 0: >
# SELECT COALESCE(MAX(pg_sequence_last_value(c.oid)::float / (pg_sequence_parameters(oid)).maximum_value), 0) AS max_usage
# FROM pg_class c
# WHERE c.relkind = 'S'
# 9.3 - 9.6 lacks the pg_sequence_last_value function, and pg_sequences view
# 90300: >
# SELECT COALESCE(MAX(pg_sequence_last_value(c.oid)::float / s.maximum_value), 0) AS max_usage
# FROM pg_class c
# CROSS JOIN LATERAL pg_sequence_parameters(c.oid) AS s
# WHERE c.relkind = 'S'
100000: SELECT COALESCE(MAX(last_value::float / max_value), 0) AS max_usage FROM pg_sequences;
test_args:
dbname: postgres
sequence_visibility:
type: row
query:
100000: >
SELECT COUNT(*) FILTER (WHERE has_sequence_privilege(c.oid, 'SELECT,USAGE')) AS visible_sequences,
COUNT(*) AS total_sequences
FROM pg_class AS c
WHERE relkind = 'S'
locks:
type: row
query:
0:
SELECT COUNT(*) AS total,
SUM(CASE WHEN granted THEN 1 ELSE 0 END) AS granted
FROM pg_locks
90400: >
SELECT COUNT(*) AS total,
COUNT(*) FILTER (WHERE granted) AS granted
FROM pg_locks
## ##
# Per-replication metrics # Per-replication metrics
@@ -205,7 +253,7 @@ metrics:
query: query:
90400: > 90400: >
SELECT pid, usename, SELECT pid, usename,
EXTRACT(EPOCH FROM backend_start)::integer AS backend_start, EXTRACT(EPOCH FROM backend_start) AS backend_start,
state, state,
pg_xlog_location_diff(pg_current_xlog_location(), sent_location) AS sent_lsn, pg_xlog_location_diff(pg_current_xlog_location(), sent_location) AS sent_lsn,
pg_xlog_location_diff(pg_current_xlog_location(), write_location) AS write_lsn, pg_xlog_location_diff(pg_current_xlog_location(), write_location) AS write_lsn,
@@ -216,20 +264,21 @@ metrics:
NULL AS replay_lag, NULL AS replay_lag,
sync_state sync_state
FROM pg_stat_replication FROM pg_stat_replication
WHERE client_addr || '_' || regexp_replace(application_name, '[ ,]', '_', 'g') = %(repid)s WHERE host(client_addr) || '_' || regexp_replace(application_name, '[ ,]', '_', 'g') = %(repid)s
100000: > 100000: >
SELECT pid, usename, SELECT pid, usename,
EXTRACT(EPOCH FROM backend_start)::integer AS backend_start, EXTRACT(EPOCH FROM backend_start) AS backend_start,
state, state,
pg_wal_lsn_diff(pg_current_wal_lsn(), sent_lsn) AS sent_lsn, pg_wal_lsn_diff(pg_current_wal_lsn(), sent_lsn) AS sent_lsn,
pg_wal_lsn_diff(pg_current_wal_lsn(), write_lsn) AS write_lsn, pg_wal_lsn_diff(pg_current_wal_lsn(), write_lsn) AS write_lsn,
pg_wal_lsn_diff(pg_current_wal_lsn(), flush_lsn) AS flush_lsn, pg_wal_lsn_diff(pg_current_wal_lsn(), flush_lsn) AS flush_lsn,
pg_wal_lsn_diff(pg_current_wal_lsn(), replay_lsn) AS replay_lsn, pg_wal_lsn_diff(pg_current_wal_lsn(), replay_lsn) AS replay_lsn,
COALESCE(EXTRACT(EPOCH FROM write_lag), 0)::integer AS write_lag, COALESCE(EXTRACT(EPOCH FROM write_lag), 0) AS write_lag,
COALESCE(EXTRACT(EPOCH FROM flush_lag), 0)::integer AS flush_lag, COALESCE(EXTRACT(EPOCH FROM flush_lag), 0) AS flush_lag,
COALESCE(EXTRACT(EPOCH FROM replay_lag), 0)::integer AS replay_lag, COALESCE(EXTRACT(EPOCH FROM replay_lag), 0) AS replay_lag,
sync_state sync_state
FROM pg_stat_replication WHERE client_addr || '_' || regexp_replace(application_name, '[ ,]', '_', 'g') = %(repid)s FROM pg_stat_replication
WHERE host(client_addr) || '_' || regexp_replace(application_name, '[ ,]', '_', 'g') = %(repid)s
test_args: test_args:
repid: 127.0.0.1_test_rep repid: 127.0.0.1_test_rep
@@ -261,6 +310,7 @@ metrics:
test_args: test_args:
slot: test_slot slot: test_slot
## ##
# Debugging # Debugging
## ##

View File

@@ -1,103 +1,141 @@
#!/usr/bin/env python3 #!/usr/bin/env python3
"""
pgmon is a monitoring intermediary that sits between a PostgreSQL cluster and a monitoring systen
that is capable of parsing JSON responses over an HTTP connection.
"""
# pylint: disable=too-few-public-methods
import yaml
import json import json
import time import time
import os import os
import sys import sys
import signal
import argparse import argparse
import logging import logging
import re
from decimal import Decimal
from urllib.parse import urlparse, parse_qs
from contextlib import contextmanager
from datetime import datetime, timedelta from datetime import datetime, timedelta
from http.server import BaseHTTPRequestHandler
from http.server import ThreadingHTTPServer
from threading import Lock
import yaml
import psycopg2 import psycopg2
from psycopg2.extras import RealDictCursor from psycopg2.extras import RealDictCursor
from psycopg2.pool import ThreadedConnectionPool from psycopg2.pool import ThreadedConnectionPool
from contextlib import contextmanager
import signal
from threading import Thread, Lock, Semaphore
from http.server import BaseHTTPRequestHandler, HTTPServer
from http.server import ThreadingHTTPServer
from urllib.parse import urlparse, parse_qs
import requests import requests
import re
VERSION = "1.0.3"
# Configuration VERSION = "1.0.4"
config = {}
# Dictionary of current PostgreSQL connection pools
connections_lock = Lock()
connections = {}
# Dictionary of unhappy databases. Keys are database names, value is the time class Context:
# the database was determined to be unhappy plus the cooldown setting. So, """
# basically it's the time when we should try to connect to the database again. The global context for connections, config, version, nad IPC
unhappy_cooldown = {} """
# Version information # Configuration
cluster_version = None config = {}
cluster_version_next_check = None
cluster_version_lock = Lock()
# PostgreSQL latest version information # Dictionary of current PostgreSQL connection pools
latest_version = None connections_lock = Lock()
latest_version_next_check = None connections = {}
latest_version_lock = Lock()
release_supported = None
# Running state (used to gracefully shut down) # Dictionary of unhappy databases. Keys are database names, value is the time
running = True # the database was determined to be unhappy plus the cooldown setting. So,
# basically it's the time when we should try to connect to the database again.
unhappy_cooldown = {}
# The http server object # Version information
httpd = None cluster_version = None
cluster_version_next_check = None
cluster_version_lock = Lock()
# Where the config file lives # PostgreSQL latest version information
config_file = None latest_version = None
latest_version_next_check = None
latest_version_lock = Lock()
release_supported = None
# Configure logging # Running state (used to gracefully shut down)
log = logging.getLogger(__name__) running = True
formatter = logging.Formatter(
"%(asctime)s - %(levelname)s - %(filename)s: %(funcName)s() line %(lineno)d: %(message)s" # The http server object
) httpd = None
console_log_handler = logging.StreamHandler()
console_log_handler.setFormatter(formatter) # Where the config file lives
log.addHandler(console_log_handler) config_file = None
# Configure logging
log = logging.getLogger(__name__)
@classmethod
def init_logging(cls):
"""
Actually initialize the logging framework. Since we don't ever instantiate the Context
class, this provides a way to make a few modifications to the log handler.
"""
formatter = logging.Formatter(
"%(asctime)s - %(levelname)s - %(filename)s: "
"%(funcName)s() line %(lineno)d: %(message)s"
)
console_log_handler = logging.StreamHandler()
console_log_handler.setFormatter(formatter)
cls.log.addHandler(console_log_handler)
# Error types # Error types
class ConfigError(Exception): class ConfigError(Exception):
pass """
Error type for all config related errors.
"""
class DisconnectedError(Exception): class DisconnectedError(Exception):
pass """
Error indicating a previously active connection to the database has been disconnected.
"""
class UnhappyDBError(Exception): class UnhappyDBError(Exception):
pass """
Error indicating that a database the code has been asked to connect to is on the unhappy list.
"""
class UnknownMetricError(Exception): class UnknownMetricError(Exception):
pass """
Error indicating that an undefined metric was requested.
"""
class MetricVersionError(Exception): class MetricVersionError(Exception):
pass """
Error indicating that there is no suitable query for a metric that was requested for the
version of PostgreSQL being monitored.
"""
class LatestVersionCheckError(Exception): class LatestVersionCheckError(Exception):
pass """
Error indicating that there was a problem retrieving or parsing the latest version information.
"""
# Default config settings # Default config settings
default_config = { DEFAULT_CONFIG = {
# The address the agent binds to # The address the agent binds to
"address": "127.0.0.1", "address": "127.0.0.1",
# The port the agent listens on for requests # The port the agent listens on for requests
@@ -175,27 +213,16 @@ def update_deep(d1, d2):
return d1 return d1
def read_config(path, included=False): def validate_metric(path, name, metric):
""" """
Read a config file. Validate a metric definition from a given file. If any query definitions come from external
files, the metric dict will be updated with the actual query.
params: Params:
path: path to the file to read path: path to the file which contains this definition
included: is this file included by another file? name: name of the metric
metric: the dictionary containing the metric definition
""" """
# Read config file
log.info("Reading log file: {}".format(path))
with open(path, "r") as f:
try:
cfg = yaml.safe_load(f)
except yaml.parser.ParserError as e:
raise ConfigError("Inavlid config file: {}: {}".format(path, e))
# Since we use it a few places, get the base directory from the config
config_base = os.path.dirname(path)
# Read any external queries and validate metric definitions
for name, metric in cfg.get("metrics", {}).items():
# Validate return types # Validate return types
try: try:
if metric["type"] not in ["value", "row", "column", "set"]: if metric["type"] not in ["value", "row", "column", "set"]:
@@ -204,14 +231,14 @@ def read_config(path, included=False):
metric["type"], name, path metric["type"], name, path
) )
) )
except KeyError: except KeyError as e:
raise ConfigError( raise ConfigError(
"No type specified for metric {} in {}".format(name, path) "No type specified for metric {} in {}".format(name, path)
) ) from e
# Ensure queries exist # Ensure queries exist
query_dict = metric.get("query", {}) query_dict = metric.get("query", {})
if type(query_dict) is not dict: if not isinstance(query_dict, dict):
raise ConfigError( raise ConfigError(
"Query definition should be a dictionary, got: {} for metric {} in {}".format( "Query definition should be a dictionary, got: {} for metric {} in {}".format(
query_dict, name, path query_dict, name, path
@@ -222,22 +249,47 @@ def read_config(path, included=False):
raise ConfigError("Missing queries for metric {} in {}".format(name, path)) raise ConfigError("Missing queries for metric {} in {}".format(name, path))
# Read external sql files and validate version keys # Read external sql files and validate version keys
config_base = os.path.dirname(path)
for vers, query in metric["query"].items(): for vers, query in metric["query"].items():
try: try:
int(vers) int(vers)
except: except Exception as e:
raise ConfigError( raise ConfigError(
"Invalid version: {} for metric {} in {}".format(vers, name, path) "Invalid version: {} for metric {} in {}".format(vers, name, path)
) ) from e
# Read in the external query and update the definition in the metricdictionary
if query.startswith("file:"): if query.startswith("file:"):
query_path = query[5:] query_path = query[5:]
if not query_path.startswith("/"): if not query_path.startswith("/"):
query_path = os.path.join(config_base, query_path) query_path = os.path.join(config_base, query_path)
with open(query_path, "r") as f: with open(query_path, "r", encoding="utf-8") as f:
metric["query"][vers] = f.read() metric["query"][vers] = f.read()
def read_config(path, included=False):
"""
Read a config file.
params:
path: path to the file to read
included: is this file included by another file?
"""
# Read config file
Context.log.info("Reading log file: %s", path)
with open(path, "r", encoding="utf-8") as f:
try:
cfg = yaml.safe_load(f)
except yaml.parser.ParserError as e:
raise ConfigError("Inavlid config file: {}: {}".format(path, e)) from e
# Read any external queries and validate metric definitions
for name, metric in cfg.get("metrics", {}).items():
validate_metric(path, name, metric)
# Read any included config files # Read any included config files
config_base = os.path.dirname(path)
for inc in cfg.get("include", []): for inc in cfg.get("include", []):
# Prefix relative paths with the directory from the current config # Prefix relative paths with the directory from the current config
if not inc.startswith("/"): if not inc.startswith("/"):
@@ -248,14 +300,14 @@ def read_config(path, included=False):
# config # config
if included: if included:
return cfg return cfg
else:
new_config = {} new_config = {}
update_deep(new_config, default_config) update_deep(new_config, DEFAULT_CONFIG)
update_deep(new_config, cfg) update_deep(new_config, cfg)
# Minor sanity checks # Minor sanity checks
if len(new_config["metrics"]) == 0: if len(new_config["metrics"]) == 0:
log.error("No metrics are defined") Context.log.error("No metrics are defined")
raise ConfigError("No metrics defined") raise ConfigError("No metrics defined")
# Validate the new log level before changing the config # Validate the new log level before changing the config
@@ -268,14 +320,17 @@ def read_config(path, included=False):
]: ]:
raise ConfigError("Invalid log level: {}".format(new_config["log_level"])) raise ConfigError("Invalid log level: {}".format(new_config["log_level"]))
global config Context.config = new_config
config = new_config
# Apply changes to log level # Apply changes to log level
log.setLevel(logging.getLevelName(config["log_level"].upper())) Context.log.setLevel(logging.getLevelName(Context.config["log_level"].upper()))
# Return the config (mostly to make pylint happy, but also in case I opt to remove the side
# effect and make this more functional.
return Context.config
def signal_handler(sig, frame): def signal_handler(sig, frame): # pylint: disable=unused-argument
""" """
Function for handling signals Function for handling signals
@@ -286,19 +341,22 @@ def signal_handler(sig, frame):
# Signal everything to shut down # Signal everything to shut down
if sig in [signal.SIGINT, signal.SIGTERM, signal.SIGQUIT]: if sig in [signal.SIGINT, signal.SIGTERM, signal.SIGQUIT]:
log.info("Shutting down ...") Context.log.info("Shutting down ...")
global running Context.running = False
running = False if Context.httpd is not None:
if httpd is not None: Context.httpd.socket.close()
httpd.socket.close()
# Signal a reload # Signal a reload
if sig == signal.SIGHUP: if sig == signal.SIGHUP:
log.warning("Received config reload signal") Context.log.warning("Received config reload signal")
read_config(config_file) read_config(Context.config_file)
class ConnectionPool(ThreadedConnectionPool): class ConnectionPool(ThreadedConnectionPool):
"""
Threaded connection pool that has a context manager.
"""
def __init__(self, dbname, minconn, maxconn, *args, **kwargs): def __init__(self, dbname, minconn, maxconn, *args, **kwargs):
# Make sure dbname isn't different in the kwargs # Make sure dbname isn't different in the kwargs
kwargs["dbname"] = dbname kwargs["dbname"] = dbname
@@ -307,7 +365,14 @@ class ConnectionPool(ThreadedConnectionPool):
self.name = dbname self.name = dbname
@contextmanager @contextmanager
def connection(self, timeout=None): def connection(self, timeout):
"""
Connection context manager for our connection pool. This will attempt to retrieve a
connection until the timeout is reached.
Params:
timeout: how long to keep trying to get a connection bedore giving up
"""
conn = None conn = None
timeout_time = datetime.now() + timedelta(timeout) timeout_time = datetime.now() + timedelta(timeout)
# We will continue to try to get a connection slot until we time out # We will continue to try to get a connection slot until we time out
@@ -331,34 +396,37 @@ class ConnectionPool(ThreadedConnectionPool):
def get_pool(dbname): def get_pool(dbname):
""" """
Get a database connection pool. Get a database connection pool.
Params:
dbname: the name of the database for which a connection pool should be returned.
""" """
# Check if the db is unhappy and wants to be left alone # Check if the db is unhappy and wants to be left alone
if dbname in unhappy_cooldown: if dbname in Context.unhappy_cooldown:
if unhappy_cooldown[dbname] > datetime.now(): if Context.unhappy_cooldown[dbname] > datetime.now():
raise UnhappyDBError() raise UnhappyDBError()
# Create a connection pool if it doesn't already exist # Create a connection pool if it doesn't already exist
if dbname not in connections: if dbname not in Context.connections:
with connections_lock: with Context.connections_lock:
# Make sure nobody created the pool while we were waiting on the # Make sure nobody created the pool while we were waiting on the
# lock # lock
if dbname not in connections: if dbname not in Context.connections:
log.info("Creating connection pool for: {}".format(dbname)) Context.log.info("Creating connection pool for: %s", dbname)
# Actually create the connection pool # Actually create the connection pool
connections[dbname] = ConnectionPool( Context.connections[dbname] = ConnectionPool(
dbname, dbname,
int(config["min_pool_size"]), int(Context.config["min_pool_size"]),
int(config["max_pool_size"]), int(Context.config["max_pool_size"]),
application_name="pgmon", application_name="pgmon",
host=config["dbhost"], host=Context.config["dbhost"],
port=config["dbport"], port=Context.config["dbport"],
user=config["dbuser"], user=Context.config["dbuser"],
connect_timeout=int(config["connect_timeout"]), connect_timeout=int(Context.config["connect_timeout"]),
sslmode=config["ssl_mode"], sslmode=Context.config["ssl_mode"],
) )
# Clear the unhappy indicator if present # Clear the unhappy indicator if present
unhappy_cooldown.pop(dbname, None) Context.unhappy_cooldown.pop(dbname, None)
return connections[dbname] return Context.connections[dbname]
def handle_connect_failure(pool): def handle_connect_failure(pool):
@@ -366,8 +434,8 @@ def handle_connect_failure(pool):
Mark the database as being unhappy so we can leave it alone for a while Mark the database as being unhappy so we can leave it alone for a while
""" """
dbname = pool.name dbname = pool.name
unhappy_cooldown[dbname] = datetime.now() + timedelta( Context.unhappy_cooldown[dbname] = datetime.now() + timedelta(
seconds=int(config["reconnect_cooldown"]) seconds=int(Context.config["reconnect_cooldown"])
) )
@@ -391,37 +459,54 @@ def get_query(metric, version):
raise MetricVersionError("Missing metric query for PostgreSQL {}".format(version)) raise MetricVersionError("Missing metric query for PostgreSQL {}".format(version))
def json_encode_special(obj):
"""
Encoder function to handle types the standard JSON package doesn't know what
to do with
"""
if isinstance(obj, Decimal):
return float(obj)
raise TypeError("Cannot serialize object of {}".format(type(obj)))
def run_query_no_retry(pool, return_type, query, args): def run_query_no_retry(pool, return_type, query, args):
""" """
Run the query with no explicit retry code Run the query with no explicit retry code
""" """
with pool.connection(float(config["connect_timeout"])) as conn: with pool.connection(float(Context.config["connect_timeout"])) as conn:
try: try:
with conn.cursor(cursor_factory=RealDictCursor) as curs: with conn.cursor(cursor_factory=RealDictCursor) as curs:
output = None
curs.execute(query, args) curs.execute(query, args)
res = curs.fetchall() res = curs.fetchall()
if return_type == "value": if return_type == "value":
if len(res) == 0: if len(res) == 0:
return "" output = ""
return str(list(res[0].values())[0]) output = str(list(res[0].values())[0])
elif return_type == "row": elif return_type == "row":
if len(res) == 0: # if len(res) == 0:
return "[]" # return "[]"
return json.dumps(res[0]) output = json.dumps(res[0], default=json_encode_special)
elif return_type == "column": elif return_type == "column":
if len(res) == 0: # if len(res) == 0:
return "[]" # return "[]"
return json.dumps([list(r.values())[0] for r in res]) output = json.dumps(
[list(r.values())[0] for r in res], default=json_encode_special
)
elif return_type == "set": elif return_type == "set":
return json.dumps(res) output = json.dumps(res, default=json_encode_special)
except:
dbname = pool.name
if dbname in unhappy_cooldown:
raise UnhappyDBError()
elif conn.closed != 0:
raise DisconnectedError()
else: else:
raise ConfigError(
"Invalid query return type: {}".format(return_type)
)
return output
except Exception as e:
dbname = pool.name
if dbname in Context.unhappy_cooldown:
raise UnhappyDBError() from e
if conn.closed != 0:
raise DisconnectedError() from e
raise raise
@@ -443,7 +528,7 @@ def run_query(pool, return_type, query, args):
try: try:
return run_query_no_retry(pool, return_type, query, args) return run_query_no_retry(pool, return_type, query, args)
except DisconnectedError: except DisconnectedError:
log.warning("Stale PostgreSQL connection found ... trying again") Context.log.warning("Stale PostgreSQL connection found ... trying again")
# This sleep is an annoying hack to give the pool workers time to # This sleep is an annoying hack to give the pool workers time to
# actually mark the connection, otherwise it can be given back in the # actually mark the connection, otherwise it can be given back in the
# next connection() call # next connection() call
@@ -451,9 +536,9 @@ def run_query(pool, return_type, query, args):
time.sleep(1) time.sleep(1)
try: try:
return run_query_no_retry(pool, return_type, query, args) return run_query_no_retry(pool, return_type, query, args)
except: except Exception as e:
handle_connect_failure(pool) handle_connect_failure(pool)
raise UnhappyDBError() raise UnhappyDBError() from e
def get_cluster_version(): def get_cluster_version():
@@ -461,40 +546,39 @@ def get_cluster_version():
Get the PostgreSQL version if we don't already know it, or if it's been Get the PostgreSQL version if we don't already know it, or if it's been
too long sice the last time it was checked. too long sice the last time it was checked.
""" """
global cluster_version
global cluster_version_next_check
# If we don't know the version or it's past the recheck time, get the # If we don't know the version or it's past the recheck time, get the
# version from the database. Only one thread needs to do this, so they all # version from the database. Only one thread needs to do this, so they all
# try to grab the lock, and then make sure nobody else beat them to it. # try to grab the lock, and then make sure nobody else beat them to it.
if ( if (
cluster_version is None Context.cluster_version is None
or cluster_version_next_check is None or Context.cluster_version_next_check is None
or cluster_version_next_check < datetime.now() or Context.cluster_version_next_check < datetime.now()
): ):
with cluster_version_lock: with Context.cluster_version_lock:
# Only check if nobody already got the version before us # Only check if nobody already got the version before us
if ( if (
cluster_version is None Context.cluster_version is None
or cluster_version_next_check is None or Context.cluster_version_next_check is None
or cluster_version_next_check < datetime.now() or Context.cluster_version_next_check < datetime.now()
): ):
log.info("Checking PostgreSQL cluster version") Context.log.info("Checking PostgreSQL cluster version")
pool = get_pool(config["dbname"]) pool = get_pool(Context.config["dbname"])
cluster_version = int( Context.cluster_version = int(
run_query(pool, "value", "SHOW server_version_num", None) run_query(pool, "value", "SHOW server_version_num", None)
) )
cluster_version_next_check = datetime.now() + timedelta( Context.cluster_version_next_check = datetime.now() + timedelta(
seconds=int(config["version_check_period"]) seconds=int(Context.config["version_check_period"])
) )
log.info("Got PostgreSQL cluster version: {}".format(cluster_version)) Context.log.info(
log.debug( "Got PostgreSQL cluster version: %s", Context.cluster_version
"Next PostgreSQL cluster version check will be after: {}".format(
cluster_version_next_check
) )
Context.log.debug(
"Next PostgreSQL cluster version check will be after: %s",
Context.cluster_version_next_check,
) )
return cluster_version return Context.cluster_version
def version_num_to_release(version_num): def version_num_to_release(version_num):
@@ -507,7 +591,6 @@ def version_num_to_release(version_num):
""" """
if version_num // 10000 < 10: if version_num // 10000 < 10:
return version_num // 10000 + (version_num % 10000 // 100 / 10) return version_num // 10000 + (version_num % 10000 // 100 / 10)
else:
return version_num // 10000 return version_num // 10000
@@ -516,7 +599,7 @@ def parse_version_rss(raw_rss, release):
Parse the raw RSS from the versions.rss feed to extract the latest version of Parse the raw RSS from the versions.rss feed to extract the latest version of
PostgreSQL that's availabe for the cluster being monitored. PostgreSQL that's availabe for the cluster being monitored.
This sets these global variables: This sets these Context variables:
latest_version latest_version
release_supported release_supported
@@ -526,8 +609,6 @@ def parse_version_rss(raw_rss, release):
raw_rss: The raw rss text from versions.rss raw_rss: The raw rss text from versions.rss
release: The PostgreSQL release we care about (ex: 9.2, 14) release: The PostgreSQL release we care about (ex: 9.2, 14)
""" """
global latest_version
global release_supported
# Regular expressions for parsing the RSS document # Regular expressions for parsing the RSS document
version_line = re.compile( version_line = re.compile(
@@ -547,75 +628,75 @@ def parse_version_rss(raw_rss, release):
version = m.group(1) version = m.group(1)
parts = list(map(int, version.split("."))) parts = list(map(int, version.split(".")))
if parts[0] < 10: if parts[0] < 10:
latest_version = int( Context.latest_version = int(
"{}{:02}{:02}".format(parts[0], parts[1], parts[2]) "{}{:02}{:02}".format(parts[0], parts[1], parts[2])
) )
else: else:
latest_version = int("{}00{:02}".format(parts[0], parts[1])) Context.latest_version = int("{}00{:02}".format(parts[0], parts[1]))
elif release_found: elif release_found:
# The next line after the version tells if the version is supported # The next line after the version tells if the version is supported
if unsupported_line.match(line): if unsupported_line.match(line):
release_supported = False Context.release_supported = False
else: else:
release_supported = True Context.release_supported = True
break break
# Make sure we actually found it # Make sure we actually found it
if not release_found: if not release_found:
raise LatestVersionCheckError("Current release ({}) not found".format(release)) raise LatestVersionCheckError("Current release ({}) not found".format(release))
log.info( Context.log.info(
"Got latest PostgreSQL version: {} supported={}".format( "Got latest PostgreSQL version: %s supported=%s",
latest_version, release_supported Context.latest_version,
) Context.release_supported,
)
log.debug(
"Next latest PostgreSQL version check will be after: {}".format(
latest_version_next_check
) )
Context.log.debug(
"Next latest PostgreSQL version check will be after: %s",
Context.latest_version_next_check,
) )
def get_latest_version(): def get_latest_version():
""" """
Get the latest supported version of the major PostgreSQL release running on the server being monitored. Get the latest supported version of the major PostgreSQL release running on the server being
monitored.
""" """
global latest_version_next_check
# If we don't know the latest version or it's past the recheck time, get the # If we don't know the latest version or it's past the recheck time, get the
# version from the PostgreSQL RSS feed. Only one thread needs to do this, so # version from the PostgreSQL RSS feed. Only one thread needs to do this, so
# they all try to grab the lock, and then make sure nobody else beat them to it. # they all try to grab the lock, and then make sure nobody else beat them to it.
if ( if (
latest_version is None Context.latest_version is None
or latest_version_next_check is None or Context.latest_version_next_check is None
or latest_version_next_check < datetime.now() or Context.latest_version_next_check < datetime.now()
): ):
# Note: we get the cluster version here before grabbing the latest_version_lock # Note: we get the cluster version here before grabbing the latest_version_lock
# lock so it's not held while trying to talk with the DB. # lock so it's not held while trying to talk with the DB.
release = version_num_to_release(get_cluster_version()) release = version_num_to_release(get_cluster_version())
with latest_version_lock: with Context.latest_version_lock:
# Only check if nobody already got the version before us # Only check if nobody already got the version before us
if ( if (
latest_version is None Context.latest_version is None
or latest_version_next_check is None or Context.latest_version_next_check is None
or latest_version_next_check < datetime.now() or Context.latest_version_next_check < datetime.now()
): ):
log.info("Checking latest PostgreSQL version") Context.log.info("Checking latest PostgreSQL version")
latest_version_next_check = datetime.now() + timedelta( Context.latest_version_next_check = datetime.now() + timedelta(
seconds=int(config["latest_version_check_period"]) seconds=int(Context.config["latest_version_check_period"])
) )
# Grab the RSS feed # Grab the RSS feed
raw_rss = requests.get("https://www.postgresql.org/versions.rss") raw_rss = requests.get(
"https://www.postgresql.org/versions.rss", timeout=30
)
if raw_rss.status_code != 200: if raw_rss.status_code != 200:
raise LatestVersionCheckError("code={}".format(r.status_code)) raise LatestVersionCheckError("code={}".format(raw_rss.status_code))
# Parse the RSS body and set global variables # Parse the RSS body and set Context variables
parse_version_rss(raw_rss.text, release) parse_version_rss(raw_rss.text, release)
return latest_version return Context.latest_version
def sample_metric(dbname, metric_name, args, retry=True): def sample_metric(dbname, metric_name, args, retry=True):
@@ -624,9 +705,9 @@ def sample_metric(dbname, metric_name, args, retry=True):
""" """
# Get the metric definition # Get the metric definition
try: try:
metric = config["metrics"][metric_name] metric = Context.config["metrics"][metric_name]
except KeyError: except KeyError as e:
raise UnknownMetricError("Unknown metric: {}".format(metric_name)) raise UnknownMetricError("Unknown metric: {}".format(metric_name)) from e
# Get the connection pool for the database, or create one if it doesn't # Get the connection pool for the database, or create one if it doesn't
# already exist. # already exist.
@@ -641,7 +722,6 @@ def sample_metric(dbname, metric_name, args, retry=True):
# Execute the quert # Execute the quert
if retry: if retry:
return run_query(pool, metric["type"], query, args) return run_query(pool, metric["type"], query, args)
else:
return run_query_no_retry(pool, metric["type"], query, args) return run_query_no_retry(pool, metric["type"], query, args)
@@ -650,12 +730,17 @@ def test_queries():
Run all of the metric queries against a database and check the results Run all of the metric queries against a database and check the results
""" """
# We just use the default db for tests # We just use the default db for tests
dbname = config["dbname"] dbname = Context.config["dbname"]
# Loop through all defined metrics. # Loop through all defined metrics.
for name, metric in config["metrics"].items(): for name, metric in Context.config["metrics"].items():
# If the metric has arguments to use while testing, grab those # If the metric has arguments to use while testing, grab those
args = metric.get("test_args", {}) args = metric.get("test_args", {})
print("Testing {} [{}]".format(name, ", ".join(["{}={}".format(key, value) for key, value in args.items()]))) print(
"Testing {} [{}]".format(
name,
", ".join(["{}={}".format(key, value) for key, value in args.items()]),
)
)
# When testing against a docker container, we may end up connecting # When testing against a docker container, we may end up connecting
# before the service is truly up (it restarts during the initialization # before the service is truly up (it restarts during the initialization
# phase). To cope with this, we'll allow a few connection failures. # phase). To cope with this, we'll allow a few connection failures.
@@ -692,9 +777,8 @@ class SimpleHTTPRequestHandler(BaseHTTPRequestHandler):
""" """
Override to suppress standard request logging Override to suppress standard request logging
""" """
pass
def do_GET(self): def do_GET(self): # pylint: disable=invalid-name
""" """
Handle a request. This is just a wrapper around the actual handler Handle a request. This is just a wrapper around the actual handler
code to keep things more readable. code to keep things more readable.
@@ -702,7 +786,7 @@ class SimpleHTTPRequestHandler(BaseHTTPRequestHandler):
try: try:
self._handle_request() self._handle_request()
except BrokenPipeError: except BrokenPipeError:
log.error("Client disconnected, exiting handler") Context.log.error("Client disconnected, exiting handler")
def _handle_request(self): def _handle_request(self):
""" """
@@ -715,7 +799,6 @@ class SimpleHTTPRequestHandler(BaseHTTPRequestHandler):
if metric_name == "agent_version": if metric_name == "agent_version":
self._reply(200, VERSION) self._reply(200, VERSION)
return
elif metric_name == "latest_version_info": elif metric_name == "latest_version_info":
try: try:
get_latest_version() get_latest_version()
@@ -723,46 +806,40 @@ class SimpleHTTPRequestHandler(BaseHTTPRequestHandler):
200, 200,
json.dumps( json.dumps(
{ {
"latest": latest_version, "latest": Context.latest_version,
"supported": 1 if release_supported else 0, "supported": 1 if Context.release_supported else 0,
} }
), ),
) )
except LatestVersionCheckError as e: except LatestVersionCheckError as e:
log.error("Failed to retrieve latest version information: {}".format(e)) Context.log.error(
"Failed to retrieve latest version information: %s", e
)
self._reply(503, "Failed to retrieve latest version info") self._reply(503, "Failed to retrieve latest version info")
return else:
# Note: parse_qs returns the values as a list. Since we always expect # Note: parse_qs returns the values as a list. Since we always expect
# single values, just grab the first from each. # single values, just grab the first from each.
args = {key: values[0] for key, values in parsed_query.items()} args = {key: values[0] for key, values in parsed_query.items()}
# Get the dbname. If none was provided, use the default from the # Get the dbname. If none was provided, use the default from the
# config. # config.
dbname = args.get("dbname", config["dbname"]) dbname = args.get("dbname", Context.config["dbname"])
# Sample the metric # Sample the metric
try: try:
self._reply(200, sample_metric(dbname, metric_name, args)) self._reply(200, sample_metric(dbname, metric_name, args))
return except UnknownMetricError:
except UnknownMetricError as e: Context.log.error("Unknown metric: %s", metric_name)
log.error("Unknown metric: {}".format(metric_name))
self._reply(404, "Unknown metric") self._reply(404, "Unknown metric")
return except MetricVersionError:
except MetricVersionError as e: Context.log.error("Failed to find an query version for %s", metric_name)
log.error(
"Failed to find a version of {} for {}".format(metric_name, version)
)
self._reply(404, "Unsupported version") self._reply(404, "Unsupported version")
return except UnhappyDBError:
except UnhappyDBError as e: Context.log.info("Database %s is unhappy, please be patient", dbname)
log.info("Database {} is unhappy, please be patient".format(dbname))
self._reply(503, "Database unavailable") self._reply(503, "Database unavailable")
return except Exception as e: # pylint: disable=broad-exception-caught
except Exception as e: Context.log.error("Error running query: %s", e)
log.error("Error running query: {}".format(e))
self._reply(500, "Unexpected error: {}".format(e)) self._reply(500, "Unexpected error: {}".format(e))
return
def _reply(self, code, content): def _reply(self, code, content):
""" """
@@ -775,7 +852,14 @@ class SimpleHTTPRequestHandler(BaseHTTPRequestHandler):
self.wfile.write(bytes(content, "utf-8")) self.wfile.write(bytes(content, "utf-8"))
if __name__ == "__main__": def main():
"""
Main application routine
"""
# Initialize the logging framework
Context.init_logging()
# Handle cli args # Handle cli args
parser = argparse.ArgumentParser( parser = argparse.ArgumentParser(
prog="pgmon", description="A PostgreSQL monitoring agent" prog="pgmon", description="A PostgreSQL monitoring agent"
@@ -796,33 +880,35 @@ if __name__ == "__main__":
args = parser.parse_args() args = parser.parse_args()
# Set the config file path # Set the config file path
config_file = args.config_file Context.config_file = args.config_file
# Read the config file # Read the config file
read_config(config_file) read_config(Context.config_file)
# Run query tests and exit if test mode is enabled # Run query tests and exit if test mode is enabled
if args.test: if args.test:
errors = test_queries() if test_queries() > 0:
if errors > 0:
sys.exit(1) sys.exit(1)
else:
sys.exit(0) sys.exit(0)
# Set up the http server to receive requests # Set up the http server to receive requests
server_address = (config["address"], config["port"]) server_address = (Context.config["address"], Context.config["port"])
httpd = ThreadingHTTPServer(server_address, SimpleHTTPRequestHandler) Context.httpd = ThreadingHTTPServer(server_address, SimpleHTTPRequestHandler)
# Set up the signal handler # Set up the signal handler
signal.signal(signal.SIGINT, signal_handler) signal.signal(signal.SIGINT, signal_handler)
signal.signal(signal.SIGHUP, signal_handler) signal.signal(signal.SIGHUP, signal_handler)
# Handle requests. # Handle requests.
log.info("Listening on port {}...".format(config["port"])) Context.log.info("Listening on port %s...", Context.config["port"])
while running: while Context.running:
httpd.handle_request() Context.httpd.handle_request()
# Clean up PostgreSQL connections # Clean up PostgreSQL connections
# TODO: Improve this ... not sure it actually closes all the connections cleanly # TODO: Improve this ... not sure it actually closes all the connections cleanly
for pool in connections.values(): for pool in Context.connections.values():
pool.close() pool.close()
if __name__ == "__main__":
main()

View File

@@ -1,16 +1,26 @@
"""
Unit tests for pgmon
"""
# pylint: disable=too-many-lines
import unittest import unittest
import os
from datetime import datetime, timedelta from datetime import datetime, timedelta
import tempfile import tempfile
import logging import logging
from decimal import Decimal
import json
import pgmon import pgmon
# Silence most logging output # Silence most logging output
logging.disable(logging.CRITICAL) logging.disable(logging.CRITICAL)
versions_rss = """ VERSIONS_RSS = """
<?xml version="1.0" encoding="utf-8"?> <?xml version="1.0" encoding="utf-8"?>
<rss version="2.0" xmlns:atom="http://www.w3.org/2005/Atom"><channel><title>PostgreSQL latest versions</title><link>https://www.postgresql.org/</link><description>PostgreSQL latest versions</description><atom:link href="https://www.postgresql.org/versions.rss" rel="self"/><language>en-us</language><lastBuildDate>Thu, 08 May 2025 00:00:00 +0000</lastBuildDate><item><title>17.5 <rss version="2.0" xmlns:atom="http://www.w3.org/2005/Atom"><channel><title>PostgreSQL latest versions</title><link>https://www.postgresql.org/</link><description>PostgreSQL latest versions</description><atom:link href="https://www.postgresql.org/versions.rss" rel="self"/><language>en-us</language><lastBuildDate>Thu, 08 May 2025 00:00:00 +0000</lastBuildDate><item><title>17.5
</title><link>https://www.postgresql.org/docs/17/release-17-5.html</link><description>17.5 is the latest release in the 17 series. </title><link>https://www.postgresql.org/docs/17/release-17-5.html</link><description>17.5 is the latest release in the 17 series.
@@ -100,12 +110,18 @@ This version is unsupported!
""" """
class TestPgmonMethods(unittest.TestCase): class TestPgmonMethods(unittest.TestCase): # pylint: disable=too-many-public-methods
"""
Unit test class for pgmon
"""
## ##
# update_deep # update_deep
## ##
def test_update_deep__empty_cases(self): def test_update_deep__empty_cases(self):
# Test empty dict cases """
Test various empty dictionary permutations
"""
d1 = {} d1 = {}
d2 = {} d2 = {}
pgmon.update_deep(d1, d2) pgmon.update_deep(d1, d2)
@@ -125,7 +141,9 @@ class TestPgmonMethods(unittest.TestCase):
self.assertEqual(d2, d1) self.assertEqual(d2, d1)
def test_update_deep__scalars(self): def test_update_deep__scalars(self):
# Test adding/updating scalar values """
Test adding/updating scalar values
"""
d1 = {"foo": 1, "bar": "text", "hello": "world"} d1 = {"foo": 1, "bar": "text", "hello": "world"}
d2 = {"foo": 2, "baz": "blah"} d2 = {"foo": 2, "baz": "blah"}
pgmon.update_deep(d1, d2) pgmon.update_deep(d1, d2)
@@ -133,7 +151,9 @@ class TestPgmonMethods(unittest.TestCase):
self.assertEqual(d2, {"foo": 2, "baz": "blah"}) self.assertEqual(d2, {"foo": 2, "baz": "blah"})
def test_update_deep__lists(self): def test_update_deep__lists(self):
# Test adding to lists """
Test adding to lists
"""
d1 = {"lst1": []} d1 = {"lst1": []}
d2 = {"lst1": [1, 2]} d2 = {"lst1": [1, 2]}
pgmon.update_deep(d1, d2) pgmon.update_deep(d1, d2)
@@ -169,7 +189,9 @@ class TestPgmonMethods(unittest.TestCase):
self.assertEqual(d2, {"obj1": {"l1": [3, 4]}}) self.assertEqual(d2, {"obj1": {"l1": [3, 4]}})
def test_update_deep__dicts(self): def test_update_deep__dicts(self):
# Test adding to lists """
Test adding to dictionaries
"""
d1 = {"obj1": {}} d1 = {"obj1": {}}
d2 = {"obj1": {"a": 1, "b": 2}} d2 = {"obj1": {"a": 1, "b": 2}}
pgmon.update_deep(d1, d2) pgmon.update_deep(d1, d2)
@@ -196,7 +218,9 @@ class TestPgmonMethods(unittest.TestCase):
self.assertEqual(d2, {"obj1": {"d1": {"a": 5, "c": 12}}}) self.assertEqual(d2, {"obj1": {"d1": {"a": 5, "c": 12}}})
def test_update_deep__types(self): def test_update_deep__types(self):
# Test mismatched types """
Test mismatched types
"""
d1 = {"foo": 5} d1 = {"foo": 5}
d2 = None d2 = None
self.assertRaises(TypeError, pgmon.update_deep, d1, d2) self.assertRaises(TypeError, pgmon.update_deep, d1, d2)
@@ -215,15 +239,19 @@ class TestPgmonMethods(unittest.TestCase):
## ##
def test_get_pool__simple(self): def test_get_pool__simple(self):
# Just get a pool in a normal case """
pgmon.config.update(pgmon.default_config) Test getting a pool in a normal case
"""
pgmon.Context.config.update(pgmon.DEFAULT_CONFIG)
pool = pgmon.get_pool("postgres") pool = pgmon.get_pool("postgres")
self.assertIsNotNone(pool) self.assertIsNotNone(pool)
def test_get_pool__unhappy(self): def test_get_pool__unhappy(self):
# Test getting an unhappy database pool """
pgmon.config.update(pgmon.default_config) Test getting an unhappy database pool
pgmon.unhappy_cooldown["postgres"] = datetime.now() + timedelta(60) """
pgmon.Context.config.update(pgmon.DEFAULT_CONFIG)
pgmon.Context.unhappy_cooldown["postgres"] = datetime.now() + timedelta(60)
self.assertRaises(pgmon.UnhappyDBError, pgmon.get_pool, "postgres") self.assertRaises(pgmon.UnhappyDBError, pgmon.get_pool, "postgres")
# Test getting a different database when there's an unhappy one # Test getting a different database when there's an unhappy one
@@ -235,30 +263,37 @@ class TestPgmonMethods(unittest.TestCase):
## ##
def test_handle_connect_failure__simple(self): def test_handle_connect_failure__simple(self):
# Test adding to an empty unhappy list """
pgmon.config.update(pgmon.default_config) Test adding to an empty unhappy list
pgmon.unhappy_cooldown = {} """
pgmon.Context.config.update(pgmon.DEFAULT_CONFIG)
pgmon.Context.unhappy_cooldown = {}
pool = pgmon.get_pool("postgres") pool = pgmon.get_pool("postgres")
pgmon.handle_connect_failure(pool) pgmon.handle_connect_failure(pool)
self.assertGreater(pgmon.unhappy_cooldown["postgres"], datetime.now()) self.assertGreater(pgmon.Context.unhappy_cooldown["postgres"], datetime.now())
# Test adding another database # Test adding another database
pool = pgmon.get_pool("template0") pool = pgmon.get_pool("template0")
pgmon.handle_connect_failure(pool) pgmon.handle_connect_failure(pool)
self.assertGreater(pgmon.unhappy_cooldown["postgres"], datetime.now()) self.assertGreater(pgmon.Context.unhappy_cooldown["postgres"], datetime.now())
self.assertGreater(pgmon.unhappy_cooldown["template0"], datetime.now()) self.assertGreater(pgmon.Context.unhappy_cooldown["template0"], datetime.now())
self.assertEqual(len(pgmon.unhappy_cooldown), 2) self.assertEqual(len(pgmon.Context.unhappy_cooldown), 2)
## ##
# get_query # get_query
## ##
def test_get_query__basic(self): def test_get_query__basic(self):
# Test getting a query with one version """
Test getting a query with just a default version.
"""
metric = {"type": "value", "query": {0: "DEFAULT"}} metric = {"type": "value", "query": {0: "DEFAULT"}}
self.assertEqual(pgmon.get_query(metric, 100000), "DEFAULT") self.assertEqual(pgmon.get_query(metric, 100000), "DEFAULT")
def test_get_query__versions(self): def test_get_query__versions(self):
"""
Test getting queries when multiple versions are present.
"""
metric = {"type": "value", "query": {0: "DEFAULT", 110000: "NEW"}} metric = {"type": "value", "query": {0: "DEFAULT", 110000: "NEW"}}
# Test getting the default version of a query with no lower bound and a newer # Test getting the default version of a query with no lower bound and a newer
@@ -278,6 +313,9 @@ class TestPgmonMethods(unittest.TestCase):
self.assertEqual(pgmon.get_query(metric, 100000), "OLD") self.assertEqual(pgmon.get_query(metric, 100000), "OLD")
def test_get_query__missing_version(self): def test_get_query__missing_version(self):
"""
Test trying to get a query that is not defined for the requested version.
"""
metric = {"type": "value", "query": {96000: "OLD", 110000: "NEW", 150000: ""}} metric = {"type": "value", "query": {96000: "OLD", 110000: "NEW", 150000: ""}}
# Test getting a metric that only exists for newer versions # Test getting a metric that only exists for newer versions
@@ -291,11 +329,16 @@ class TestPgmonMethods(unittest.TestCase):
## ##
def test_read_config__simple(self): def test_read_config__simple(self):
pgmon.config = {} """
Test reading a simple config.
"""
pgmon.Context.config = {}
# Test reading just a metric and using the defaults for everything else # Test reading just a metric and using the defaults for everything else
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
# This is a comment! # This is a comment!
@@ -307,18 +350,20 @@ metrics:
""" """
) )
pgmon.read_config(f"{tmpdirname}/config.yml") pgmon.read_config(os.path.join(tmpdirname, "config.yml"))
self.assertEqual( self.assertEqual(
pgmon.config["max_pool_size"], pgmon.default_config["max_pool_size"] pgmon.Context.config["max_pool_size"], pgmon.DEFAULT_CONFIG["max_pool_size"]
) )
self.assertEqual(pgmon.config["dbuser"], pgmon.default_config["dbuser"]) self.assertEqual(pgmon.Context.config["dbuser"], pgmon.DEFAULT_CONFIG["dbuser"])
pgmon.config = {} pgmon.Context.config = {}
# Test reading a basic config # Test reading a basic config
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
# This is a comment! # This is a comment!
@@ -354,22 +399,27 @@ metrics:
""" """
) )
pgmon.read_config(f"{tmpdirname}/config.yml") pgmon.read_config(os.path.join(tmpdirname, "config.yml"))
self.assertEqual(pgmon.config["dbuser"], "someone") self.assertEqual(pgmon.Context.config["dbuser"], "someone")
self.assertEqual(pgmon.config["metrics"]["test1"]["type"], "value") self.assertEqual(pgmon.Context.config["metrics"]["test1"]["type"], "value")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") self.assertEqual(pgmon.Context.config["metrics"]["test1"]["query"][0], "TEST1")
self.assertEqual(pgmon.config["metrics"]["test2"]["query"][0], "TEST2") self.assertEqual(pgmon.Context.config["metrics"]["test2"]["query"][0], "TEST2")
def test_read_config__include(self): def test_read_config__include(self):
pgmon.config = {} """
Test including one config from another.
"""
pgmon.Context.config = {}
# Test reading a config that includes other files (absolute and relative paths, # Test reading a config that includes other files (absolute and relative paths,
# multiple levels) # multiple levels)
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
f"""--- """---
# This is a comment! # This is a comment!
min_pool_size: 1 min_pool_size: 1
max_pool_size: 2 max_pool_size: 2
@@ -381,13 +431,17 @@ reconnect_cooldown: 15
version_check_period: 3600 version_check_period: 3600
include: include:
- dbsettings.yml - dbsettings.yml
- {tmpdirname}/metrics.yml - {}/metrics.yml
""" """.format(
tmpdirname
)
) )
with open(f"{tmpdirname}/dbsettings.yml", "w") as f: with open(
os.path.join(tmpdirname, "dbsettings.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
f"""--- """---
dbuser: someone dbuser: someone
dbhost: localhost dbhost: localhost
dbport: 5555 dbport: 5555
@@ -395,9 +449,11 @@ dbname: template0
""" """
) )
with open(f"{tmpdirname}/metrics.yml", "w") as f: with open(
os.path.join(tmpdirname, "metrics.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
f"""--- """---
metrics: metrics:
test1: test1:
type: value type: value
@@ -412,9 +468,11 @@ include:
""" """
) )
with open(f"{tmpdirname}/more_metrics.yml", "w") as f: with open(
os.path.join(tmpdirname, "more_metrics.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
f"""--- """---
metrics: metrics:
test3: test3:
type: value type: value
@@ -422,20 +480,25 @@ metrics:
0: TEST3 0: TEST3
""" """
) )
pgmon.read_config(f"{tmpdirname}/config.yml") pgmon.read_config(os.path.join(tmpdirname, "config.yml"))
self.assertEqual(pgmon.config["max_idle_time"], 10) self.assertEqual(pgmon.Context.config["max_idle_time"], 10)
self.assertEqual(pgmon.config["dbuser"], "someone") self.assertEqual(pgmon.Context.config["dbuser"], "someone")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") self.assertEqual(pgmon.Context.config["metrics"]["test1"]["query"][0], "TEST1")
self.assertEqual(pgmon.config["metrics"]["test2"]["query"][0], "TEST2") self.assertEqual(pgmon.Context.config["metrics"]["test2"]["query"][0], "TEST2")
self.assertEqual(pgmon.config["metrics"]["test3"]["query"][0], "TEST3") self.assertEqual(pgmon.Context.config["metrics"]["test3"]["query"][0], "TEST3")
def test_read_config__reload(self): def test_read_config__reload(self):
pgmon.config = {} """
Test reloading a config.
"""
pgmon.Context.config = {}
# Test rereading a config to update an existing config # Test rereading a config to update an existing config
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
# This is a comment! # This is a comment!
@@ -463,12 +526,14 @@ metrics:
""" """
) )
pgmon.read_config(f"{tmpdirname}/config.yml") pgmon.read_config(os.path.join(tmpdirname, "config.yml"))
# Just make sure the first config was read # Just make sure the first config was read
self.assertEqual(len(pgmon.config["metrics"]), 2) self.assertEqual(len(pgmon.Context.config["metrics"]), 2)
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
# This is a comment! # This is a comment!
@@ -481,18 +546,23 @@ metrics:
""" """
) )
pgmon.read_config(f"{tmpdirname}/config.yml") pgmon.read_config(os.path.join(tmpdirname, "config.yml"))
self.assertEqual(pgmon.config["min_pool_size"], 7) self.assertEqual(pgmon.Context.config["min_pool_size"], 7)
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "NEW1") self.assertEqual(pgmon.Context.config["metrics"]["test1"]["query"][0], "NEW1")
self.assertEqual(len(pgmon.config["metrics"]), 1) self.assertEqual(len(pgmon.Context.config["metrics"]), 1)
def test_read_config__query_file(self): def test_read_config__query_file(self):
pgmon.config = {} """
Test reading a query definition from a separate file
"""
pgmon.Context.config = {}
# Read a config file that reads a query from a file # Read a config file that reads a query from a file
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
metrics: metrics:
@@ -503,22 +573,30 @@ metrics:
""" """
) )
with open(f"{tmpdirname}/some_query.sql", "w") as f: with open(
os.path.join(tmpdirname, "some_query.sql"), "w", encoding="utf-8"
) as f:
f.write("This is a query") f.write("This is a query")
pgmon.read_config(f"{tmpdirname}/config.yml") pgmon.read_config(os.path.join(tmpdirname, "config.yml"))
self.assertEqual( self.assertEqual(
pgmon.config["metrics"]["test1"]["query"][0], "This is a query" pgmon.Context.config["metrics"]["test1"]["query"][0], "This is a query"
) )
def test_read_config__invalid(self): def init_invalid_config_test(self):
pgmon.config = {} """
Initialize an invalid config read test. Basically just set up a simple valid config in
order to confirm that an invalid read does not modify the live config.
"""
pgmon.Context.config = {}
# For all of these tests, we start with a valid config and also ensure that # For all of these tests, we start with a valid config and also ensure that
# it is not modified when a new config read fails # it is not modified when a new config read fails
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
metrics: metrics:
@@ -529,20 +607,48 @@ metrics:
""" """
) )
pgmon.read_config(f"{tmpdirname}/config.yml") pgmon.read_config(os.path.join(tmpdirname, "config.yml"))
# Just make sure the config was read # Just make sure the config was read
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") self.assertEqual(pgmon.Context.config["metrics"]["test1"]["query"][0], "TEST1")
def verify_invalid_config_test(self):
"""
Verify that an invalid read did not modify the live config.
"""
self.assertEqual(pgmon.Context.config["dbuser"], "postgres")
self.assertEqual(pgmon.Context.config["metrics"]["test1"]["query"][0], "TEST1")
def test_read_config__missing(self):
"""
Test reading a nonexistant config file.
"""
# Set up the test
self.init_invalid_config_test()
# Test reading a nonexistant config file # Test reading a nonexistant config file
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
self.assertRaises( self.assertRaises(
FileNotFoundError, pgmon.read_config, f"{tmpdirname}/missing.yml" FileNotFoundError,
pgmon.read_config,
os.path.join(tmpdirname, "missing.yml"),
) )
# Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__invalid(self):
"""
Test reading an invalid config file.
"""
# Set up the test
self.init_invalid_config_test()
# Test reading an invalid config file # Test reading an invalid config file
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""[default] """[default]
This looks a lot like an ini file to me This looks a lot like an ini file to me
@@ -551,12 +657,26 @@ Or maybe a TOML?
""" """
) )
self.assertRaises( self.assertRaises(
pgmon.ConfigError, pgmon.read_config, f"{tmpdirname}/config.yml" pgmon.ConfigError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
# Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__invalid_include(self):
"""
Test reading an invalid config file.
"""
# Set up the test
self.init_invalid_config_test()
# Test reading a config that includes an invalid file # Test reading a config that includes an invalid file
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
dbuser: evil dbuser: evil
@@ -570,14 +690,26 @@ include:
""" """
) )
self.assertRaises( self.assertRaises(
FileNotFoundError, pgmon.read_config, f"{tmpdirname}/config.yml" FileNotFoundError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
self.assertEqual(pgmon.config["dbuser"], "postgres")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") # Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__invalid_log_level(self):
"""
Test reading an invalid log level from a config file.
"""
# Set up the test
self.init_invalid_config_test()
# Test invalid log level # Test invalid log level
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
log_level: noisy log_level: noisy
@@ -590,14 +722,26 @@ metrics:
""" """
) )
self.assertRaises( self.assertRaises(
pgmon.ConfigError, pgmon.read_config, f"{tmpdirname}/config.yml" pgmon.ConfigError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
self.assertEqual(pgmon.config["dbuser"], "postgres")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") # Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__invalid_type(self):
"""
Test reading an invalid query result type form a config file.
"""
# Set up the test
self.init_invalid_config_test()
# Test invalid query return type # Test invalid query return type
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
dbuser: evil dbuser: evil
@@ -609,32 +753,57 @@ metrics:
""" """
) )
self.assertRaises( self.assertRaises(
pgmon.ConfigError, pgmon.read_config, f"{tmpdirname}/config.yml" pgmon.ConfigError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
self.assertEqual(pgmon.config["dbuser"], "postgres")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") # Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__invalid_query_dict(self):
"""
Test reading an invalid query definition structure type form a config file. In other words
what's supposed to be a dictionary of the form version => query, we give it something else.
"""
# Set up the test
self.init_invalid_config_test()
# Test invalid query dict type # Test invalid query dict type
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
dbuser: evil dbuser: evil
metrics: metrics:
test1: test1:
type: lots_of_data type: row
query: EVIL1 query: EVIL1
""" """
) )
self.assertRaises( self.assertRaises(
pgmon.ConfigError, pgmon.read_config, f"{tmpdirname}/config.yml" pgmon.ConfigError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
self.assertEqual(pgmon.config["dbuser"], "postgres")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") # Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__missing_type(self):
"""
Test reading a metric with a missing result type from a config file.
"""
# Set up the test
self.init_invalid_config_test()
# Test incomplete metric: missing type # Test incomplete metric: missing type
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
dbuser: evil dbuser: evil
@@ -645,14 +814,26 @@ metrics:
""" """
) )
self.assertRaises( self.assertRaises(
pgmon.ConfigError, pgmon.read_config, f"{tmpdirname}/config.yml" pgmon.ConfigError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
self.assertEqual(pgmon.config["dbuser"], "postgres")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") # Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__missing_queries(self):
"""
Test reading a metric with no queries from a config file.
"""
# Set up the test
self.init_invalid_config_test()
# Test incomplete metric: missing queries # Test incomplete metric: missing queries
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
dbuser: evil dbuser: evil
@@ -662,14 +843,26 @@ metrics:
""" """
) )
self.assertRaises( self.assertRaises(
pgmon.ConfigError, pgmon.read_config, f"{tmpdirname}/config.yml" pgmon.ConfigError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
self.assertEqual(pgmon.config["dbuser"], "postgres")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") # Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__empty_query_dict(self):
"""
Test reading a fetric with an empty query dict from a config file.
"""
# Set up the test
self.init_invalid_config_test()
# Test incomplete metric: empty queries # Test incomplete metric: empty queries
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
dbuser: evil dbuser: evil
@@ -680,14 +873,26 @@ metrics:
""" """
) )
self.assertRaises( self.assertRaises(
pgmon.ConfigError, pgmon.read_config, f"{tmpdirname}/config.yml" pgmon.ConfigError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
self.assertEqual(pgmon.config["dbuser"], "postgres")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") # Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__none_query_dict(self):
"""
Test reading a metric where the query dict is None from a config file.
"""
# Set up the test
self.init_invalid_config_test()
# Test incomplete metric: query dict is None # Test incomplete metric: query dict is None
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
dbuser: evil dbuser: evil
@@ -698,28 +903,53 @@ metrics:
""" """
) )
self.assertRaises( self.assertRaises(
pgmon.ConfigError, pgmon.read_config, f"{tmpdirname}/config.yml" pgmon.ConfigError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
self.assertEqual(pgmon.config["dbuser"], "postgres")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") # Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__missing_metrics(self):
"""
Test reading a config file with no metrics.
"""
# Set up the test
self.init_invalid_config_test()
# Test reading a config with no metrics # Test reading a config with no metrics
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
dbuser: evil dbuser: evil
""" """
) )
self.assertRaises( self.assertRaises(
pgmon.ConfigError, pgmon.read_config, f"{tmpdirname}/config.yml" pgmon.ConfigError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
self.assertEqual(pgmon.config["dbuser"], "postgres")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") # Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__missing_query_file(self):
"""
Test reading a metric from a config file where the query definition cones from a missing
file.
"""
# Set up the test
self.init_invalid_config_test()
# Test reading a query defined in a file but the file is missing # Test reading a query defined in a file but the file is missing
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
dbuser: evil dbuser: evil
@@ -731,14 +961,26 @@ metrics:
""" """
) )
self.assertRaises( self.assertRaises(
FileNotFoundError, pgmon.read_config, f"{tmpdirname}/config.yml" FileNotFoundError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
self.assertEqual(pgmon.config["dbuser"], "postgres")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") # Confirm nothing changed
self.verify_invalid_config_test()
def test_read_config__invalid_version(self):
"""
Test reading a metric with an invalid PostgreSQL version from a config file.
"""
# Set up the test
self.init_invalid_config_test()
# Test invalid query versions # Test invalid query versions
with tempfile.TemporaryDirectory() as tmpdirname: with tempfile.TemporaryDirectory() as tmpdirname:
with open(f"{tmpdirname}/config.yml", "w") as f: with open(
os.path.join(tmpdirname, "config.yml"), "w", encoding="utf-8"
) as f:
f.write( f.write(
"""--- """---
dbuser: evil dbuser: evil
@@ -750,42 +992,96 @@ metrics:
""" """
) )
self.assertRaises( self.assertRaises(
pgmon.ConfigError, pgmon.read_config, f"{tmpdirname}/config.yml" pgmon.ConfigError,
pgmon.read_config,
os.path.join(tmpdirname, "config.yml"),
) )
self.assertEqual(pgmon.config["dbuser"], "postgres")
self.assertEqual(pgmon.config["metrics"]["test1"]["query"][0], "TEST1") # Confirm nothing changed
self.verify_invalid_config_test()
##
# version_num
##
def test_version_num_to_release__valid(self): def test_version_num_to_release__valid(self):
"""
Test converting PostgreSQL versions before and after 10 when the numbering scheme changed.
"""
self.assertEqual(pgmon.version_num_to_release(90602), 9.6) self.assertEqual(pgmon.version_num_to_release(90602), 9.6)
self.assertEqual(pgmon.version_num_to_release(130002), 13) self.assertEqual(pgmon.version_num_to_release(130002), 13)
def test_parse_version_rss__simple(self): ##
pgmon.parse_version_rss(versions_rss, 13) # parse_version_rss
self.assertEqual(pgmon.latest_version, 130021) ##
self.assertTrue(pgmon.release_supported)
pgmon.parse_version_rss(versions_rss, 9.6) def test_parse_version_rss__supported(self):
self.assertEqual(pgmon.latest_version, 90624) """
self.assertFalse(pgmon.release_supported) Test parsing a supported version from the RSS feed
"""
pgmon.parse_version_rss(VERSIONS_RSS, 13)
self.assertEqual(pgmon.Context.latest_version, 130021)
self.assertTrue(pgmon.Context.release_supported)
def test_parse_version_rss__unsupported(self):
"""
Test parsing an unsupported version from the RSS feed
"""
pgmon.parse_version_rss(VERSIONS_RSS, 9.6)
self.assertEqual(pgmon.Context.latest_version, 90624)
self.assertFalse(pgmon.Context.release_supported)
def test_parse_version_rss__missing(self): def test_parse_version_rss__missing(self):
# Test asking about versions that don't exist """
Test asking about versions that don't exist in the RSS feed
"""
self.assertRaises( self.assertRaises(
pgmon.LatestVersionCheckError, pgmon.parse_version_rss, versions_rss, 9.7 pgmon.LatestVersionCheckError, pgmon.parse_version_rss, VERSIONS_RSS, 9.7
) )
self.assertRaises( self.assertRaises(
pgmon.LatestVersionCheckError, pgmon.parse_version_rss, versions_rss, 99 pgmon.LatestVersionCheckError, pgmon.parse_version_rss, VERSIONS_RSS, 99
) )
##
# get_latest_version
##
def test_get_latest_version(self): def test_get_latest_version(self):
"""
Test getting the latest version from the actual RSS feed
"""
# Define a cluster version here so the test doesn't need a database # Define a cluster version here so the test doesn't need a database
pgmon.cluster_version_next_check = datetime.now() + timedelta(hours=1) pgmon.Context.cluster_version_next_check = datetime.now() + timedelta(hours=1)
pgmon.cluster_version = 90623 pgmon.Context.cluster_version = 90623
# Set up a default config # Set up a default config
pgmon.update_deep(pgmon.config, pgmon.default_config) pgmon.update_deep(pgmon.Context.config, pgmon.DEFAULT_CONFIG)
# Make sure we can pull the RSS file (we assume the 9.6 series won't be getting # Make sure we can pull the RSS file (we assume the 9.6 series won't be getting
# any more updates) # any more updates)
self.assertEqual(pgmon.get_latest_version(), 90624) self.assertEqual(pgmon.get_latest_version(), 90624)
##
# json_encode_special
##
def test_json_encode_special(self):
"""
Test encoding Decimal types as JSON
"""
# Confirm that we're getting the right type
self.assertFalse(isinstance(Decimal("0.5"), float))
self.assertTrue(isinstance(pgmon.json_encode_special(Decimal("0.5")), float))
# Make sure we get sane values
self.assertEqual(pgmon.json_encode_special(Decimal("0.5")), 0.5)
self.assertEqual(pgmon.json_encode_special(Decimal("12")), 12.0)
# Make sure we can still fail for other types
self.assertRaises(TypeError, pgmon.json_encode_special, object)
# Make sure we can actually serialize a Decimal
self.assertEqual(
json.dumps(Decimal("2.5"), default=pgmon.json_encode_special), "2.5"
)

View File

@@ -167,7 +167,8 @@ zabbix_export:
operator: NOT_MATCHES_REGEX operator: NOT_MATCHES_REGEX
formulaid: A formulaid: A
lifetime: 30d lifetime: 30d
enabled_lifetime_type: DISABLE_NEVER enabled_lifetime_type: DISABLE_AFTER
enabled_lifetime: 1d
item_prototypes: item_prototypes:
- uuid: a30babe4a6f4440bba2a3ee46eff7ce2 - uuid: a30babe4a6f4440bba2a3ee46eff7ce2
name: 'Time spent executing statements on {#DBNAME}' name: 'Time spent executing statements on {#DBNAME}'
@@ -785,6 +786,77 @@ zabbix_export:
value: PostgreSQL value: PostgreSQL
- tag: Database - tag: Database
value: '{#DBNAME}' value: '{#DBNAME}'
- uuid: 5960120dd01c4926b0fc1fbe9c011507
name: 'Database max sequence usage in {#DBNAME}'
type: HTTP_AGENT
key: 'pgmon_db_max_sequence[{#DBNAME}]'
delay: 5m
value_type: FLOAT
units: '%'
description: 'The percent of the currently configured value range for the most utilized sequence.'
url: 'http://localhost:{$AGENT_PORT}/sequence_usage'
query_fields:
- name: dbname
value: '{#DBNAME}'
tags:
- tag: Application
value: PostgreSQL
- tag: Database
value: '{#DBNAME}'
- uuid: 48b9cc80ac4d4aee9e9f3a5d6f7d4a95
name: 'Total number of sequences on {#DBNAME}'
type: DEPENDENT
key: 'pgmon_db_sequences[total,{#DBNAME}]'
delay: '0'
description: 'Total number of sequences in the database.'
preprocessing:
- type: JSONPATH
parameters:
- $.total_sequences
master_item:
key: 'pgmon_db_sequence_visibility[{#DBNAME}]'
tags:
- tag: Application
value: PostgreSQL
- tag: Database
value: '{#DBNAME}'
- uuid: 6521a9bab2ac47bf85429832d289bbac
name: 'Visible sequences on {#DBNAME}'
type: DEPENDENT
key: 'pgmon_db_sequences[visible,{#DBNAME}]'
delay: '0'
description: 'Number of sequences in the database for which Zabbix can see stats.'
preprocessing:
- type: JSONPATH
parameters:
- $.visible_sequences
master_item:
key: 'pgmon_db_sequence_visibility[{#DBNAME}]'
tags:
- tag: Application
value: PostgreSQL
- tag: Database
value: '{#DBNAME}'
- uuid: 00f2da3eb99940839410a6ecd5df153f
name: 'Database sequence visibility in {#DBNAME}'
type: HTTP_AGENT
key: 'pgmon_db_sequence_visibility[{#DBNAME}]'
delay: 30m
history: '0'
value_type: TEXT
trends: '0'
description: 'Statistics about the number of sequences that exist and the number Zabbix can actually see stats for.'
url: 'http://localhost:{$AGENT_PORT}/sequence_visibility'
query_fields:
- name: dbname
value: '{#DBNAME}'
tags:
- tag: Application
value: PostgreSQL
- tag: Database
value: '{#DBNAME}'
- tag: Type
value: Raw
- uuid: 492b3cac15f348c2b85f97b69c114d1b - uuid: 492b3cac15f348c2b85f97b69c114d1b
name: 'Database Stats for {#DBNAME}' name: 'Database Stats for {#DBNAME}'
type: HTTP_AGENT type: HTTP_AGENT
@@ -803,6 +875,17 @@ zabbix_export:
value: '{#DBNAME}' value: '{#DBNAME}'
- tag: Type - tag: Type
value: Raw value: Raw
trigger_prototypes:
- uuid: d29d0fd9d9d34b5ebd649592b0829ce5
expression: 'last(/PostgreSQL by pgmon/pgmon_db_sequences[total,{#DBNAME}]) <> last(/PostgreSQL by pgmon/pgmon_db_sequences[visible,{#DBNAME}])'
name: 'Sequences not visible to Zabbix on {#DBNAME}'
priority: WARNING
description: 'There are sequences for which Zabbix cannot see usage statistics'
tags:
- tag: Application
value: PostgreSQL
- tag: Component
value: Sequence
graph_prototypes: graph_prototypes:
- uuid: 1f7de43b77714f819e61c31273712b70 - uuid: 1f7de43b77714f819e61c31273712b70
name: 'DML Totals for {#DBNAME}' name: 'DML Totals for {#DBNAME}'
@@ -900,6 +983,9 @@ zabbix_export:
type: DEPENDENT type: DEPENDENT
key: pgmon_discover_io_backend_types key: pgmon_discover_io_backend_types
delay: '0' delay: '0'
lifetime: 30d
enabled_lifetime_type: DISABLE_AFTER
enabled_lifetime: 1h
item_prototypes: item_prototypes:
- uuid: b1ac2e56b30f4812bf33ce973ef16b10 - uuid: b1ac2e56b30f4812bf33ce973ef16b10
name: 'I/O Evictions by {#BACKEND_TYPE}' name: 'I/O Evictions by {#BACKEND_TYPE}'
@@ -1490,8 +1576,15 @@ zabbix_export:
type: HTTP_AGENT type: HTTP_AGENT
key: pgmon_discover_rep key: pgmon_discover_rep
delay: 10m delay: 10m
filter:
conditions:
- macro: '{#APPLICATION_NAME}'
value: '^pg_[0-9]+_sync_[0-9]+_[0-9]+$'
operator: NOT_MATCHES_REGEX
formulaid: A
lifetime: 30d lifetime: 30d
enabled_lifetime_type: DISABLE_NEVER enabled_lifetime_type: DISABLE_AFTER
enabled_lifetime: 7d
item_prototypes: item_prototypes:
- uuid: 3a5a60620e6a4db694e47251148d82f5 - uuid: 3a5a60620e6a4db694e47251148d82f5
name: 'Flush lag for {#REPID}' name: 'Flush lag for {#REPID}'
@@ -1500,6 +1593,7 @@ zabbix_export:
delay: '0' delay: '0'
history: 90d history: 90d
value_type: FLOAT value_type: FLOAT
units: s
description: 'Time elapsed between flushing recent WAL locally and receiving notification that this standby server has written and flushed it (but not yet applied it). This can be used to gauge the delay that synchronous_commit level on incurred while committing if this server was configured as a synchronous standby.' description: 'Time elapsed between flushing recent WAL locally and receiving notification that this standby server has written and flushed it (but not yet applied it). This can be used to gauge the delay that synchronous_commit level on incurred while committing if this server was configured as a synchronous standby.'
preprocessing: preprocessing:
- type: JSONPATH - type: JSONPATH
@@ -1515,13 +1609,12 @@ zabbix_export:
- tag: Database - tag: Database
value: '{#DBNAME}' value: '{#DBNAME}'
- uuid: 624f8f085a3642c9a10a03361c17763d - uuid: 624f8f085a3642c9a10a03361c17763d
name: 'Last flush LSN for {#REPID}' name: 'Last flush LSN lag for {#REPID}'
type: DEPENDENT type: DEPENDENT
key: 'pgmon_rep[flush_lsn,repid={#REPID}]' key: 'pgmon_rep[flush_lsn,repid={#REPID}]'
delay: '0' delay: '0'
history: 90d history: 90d
value_type: TEXT units: B
trends: '0'
description: 'Last write-ahead log location flushed to disk by this standby server' description: 'Last write-ahead log location flushed to disk by this standby server'
preprocessing: preprocessing:
- type: JSONPATH - type: JSONPATH
@@ -1543,6 +1636,7 @@ zabbix_export:
delay: '0' delay: '0'
history: 90d history: 90d
value_type: FLOAT value_type: FLOAT
units: s
description: 'Time elapsed between flushing recent WAL locally and receiving notification that this standby server has written, flushed and applied it. This can be used to gauge the delay that synchronous_commit level remote_apply incurred while committing if this server was configured as a synchronous standby.' description: 'Time elapsed between flushing recent WAL locally and receiving notification that this standby server has written, flushed and applied it. This can be used to gauge the delay that synchronous_commit level remote_apply incurred while committing if this server was configured as a synchronous standby.'
preprocessing: preprocessing:
- type: JSONPATH - type: JSONPATH
@@ -1558,13 +1652,12 @@ zabbix_export:
- tag: Database - tag: Database
value: '{#DBNAME}' value: '{#DBNAME}'
- uuid: fe1bed51845d4694bae8f53deed4846d - uuid: fe1bed51845d4694bae8f53deed4846d
name: 'Last replay LSN for {#REPID}' name: 'Last replay LSN lag for {#REPID}'
type: DEPENDENT type: DEPENDENT
key: 'pgmon_rep[replay_lsn,repid={#REPID}]' key: 'pgmon_rep[replay_lsn,repid={#REPID}]'
delay: '0' delay: '0'
history: 90d history: 90d
value_type: TEXT units: B
trends: '0'
description: 'Last write-ahead log location replayed into the database on this standby server' description: 'Last write-ahead log location replayed into the database on this standby server'
preprocessing: preprocessing:
- type: JSONPATH - type: JSONPATH
@@ -1580,13 +1673,12 @@ zabbix_export:
- tag: Database - tag: Database
value: '{#DBNAME}' value: '{#DBNAME}'
- uuid: 68c179d0e33f45f9bf82d2d4125763f0 - uuid: 68c179d0e33f45f9bf82d2d4125763f0
name: 'Last sent LSN for {#REPID}' name: 'Last sent LSN lag for {#REPID}'
type: DEPENDENT type: DEPENDENT
key: 'pgmon_rep[sent_lsn,repid={#REPID}]' key: 'pgmon_rep[sent_lsn,repid={#REPID}]'
delay: '0' delay: '0'
history: 90d history: 90d
value_type: TEXT units: B
trends: '0'
description: 'Last write-ahead log location sent on this connection' description: 'Last write-ahead log location sent on this connection'
preprocessing: preprocessing:
- type: JSONPATH - type: JSONPATH
@@ -1636,6 +1728,7 @@ zabbix_export:
delay: '0' delay: '0'
history: 90d history: 90d
value_type: FLOAT value_type: FLOAT
units: s
description: 'Time elapsed between flushing recent WAL locally and receiving notification that this standby server has written it (but not yet flushed it or applied it). This can be used to gauge the delay that synchronous_commit level remote_write incurred while committing if this server was configured as a synchronous standby.' description: 'Time elapsed between flushing recent WAL locally and receiving notification that this standby server has written it (but not yet flushed it or applied it). This can be used to gauge the delay that synchronous_commit level remote_write incurred while committing if this server was configured as a synchronous standby.'
preprocessing: preprocessing:
- type: JSONPATH - type: JSONPATH
@@ -1651,13 +1744,12 @@ zabbix_export:
- tag: Database - tag: Database
value: '{#DBNAME}' value: '{#DBNAME}'
- uuid: 57fb03cf63af4b0a91d8e36d6ff64d30 - uuid: 57fb03cf63af4b0a91d8e36d6ff64d30
name: 'Last write LSN for {#REPID}' name: 'Last write LSN lag for {#REPID}'
type: DEPENDENT type: DEPENDENT
key: 'pgmon_rep[write_lsn,repid={#REPID}]' key: 'pgmon_rep[write_lsn,repid={#REPID}]'
delay: '0' delay: '0'
history: 90d history: 90d
value_type: TEXT units: B
trends: '0'
description: 'Time elapsed between flushing recent WAL locally and receiving notification that this standby server has written it (but not yet flushed it or applied it). This can be used to gauge the delay that synchronous_commit level remote_write incurred while committing if this server was configured as a synchronous standby.' description: 'Time elapsed between flushing recent WAL locally and receiving notification that this standby server has written it (but not yet flushed it or applied it). This can be used to gauge the delay that synchronous_commit level remote_write incurred while committing if this server was configured as a synchronous standby.'
preprocessing: preprocessing:
- type: JSONPATH - type: JSONPATH
@@ -1694,6 +1786,8 @@ zabbix_export:
value: Raw value: Raw
url: 'http://localhost:{$AGENT_PORT}/discover_rep' url: 'http://localhost:{$AGENT_PORT}/discover_rep'
lld_macro_paths: lld_macro_paths:
- lld_macro: '{#APPLICATION_NAME}'
path: $.application_name
- lld_macro: '{#CLIENT_ADDR}' - lld_macro: '{#CLIENT_ADDR}'
path: $.client_addr path: $.client_addr
- lld_macro: '{#REPID}' - lld_macro: '{#REPID}'
@@ -1705,6 +1799,15 @@ zabbix_export:
type: HTTP_AGENT type: HTTP_AGENT
key: pgmon_discover_slots key: pgmon_discover_slots
delay: 10m delay: 10m
filter:
conditions:
- macro: '{#SLOT_NAME}'
value: '^pg_[0-9]+_sync_[0-9]+_[0-9]+$'
operator: NOT_MATCHES_REGEX
formulaid: A
lifetime: 30d
enabled_lifetime_type: DISABLE_AFTER
enabled_lifetime: 7d
item_prototypes: item_prototypes:
- uuid: 536c5f82e3074ddfbfd842b3a2e8d46c - uuid: 536c5f82e3074ddfbfd842b3a2e8d46c
name: 'Slot {#SLOT_NAME} - Confirmed Flushed Bytes Lag' name: 'Slot {#SLOT_NAME} - Confirmed Flushed Bytes Lag'