diff --git a/.github/workflows/build-cloudberry-rocky8.yml b/.github/workflows/build-cloudberry-rocky8.yml index d48e436265b..8bc0e12e0e0 100644 --- a/.github/workflows/build-cloudberry-rocky8.yml +++ b/.github/workflows/build-cloudberry-rocky8.yml @@ -320,6 +320,11 @@ jobs: "gpcontrib/gp_sparse_vector:installcheck", "gpcontrib/gp_toolkit:installcheck"] }, + {"test":"gpcontrib-gp-relaccess-stats", + "make_configs":["gpcontrib/gp_relaccess_stats:installcheck"], + "extension":"gp_relaccess_stats", + "shared_preload_libraries":"gp_relaccess_stats" + }, {"test":"gpcontrib-gp-stats-collector", "make_configs":["gpcontrib/gp_stats_collector:installcheck"], "extension":"gp_stats_collector" diff --git a/.github/workflows/build-cloudberry.yml b/.github/workflows/build-cloudberry.yml index 6289785bb14..f43e3510815 100644 --- a/.github/workflows/build-cloudberry.yml +++ b/.github/workflows/build-cloudberry.yml @@ -271,6 +271,11 @@ jobs: }, "enable_core_check":false }, + {"test":"gpcontrib-gp-relaccess-stats", + "make_configs":["gpcontrib/gp_relaccess_stats:installcheck"], + "extension":"gp_relaccess_stats", + "shared_preload_libraries":"gp_relaccess_stats" + }, {"test":"gpcontrib-gp-stats-collector", "make_configs":["gpcontrib/gp_stats_collector:installcheck"], "extension":"gp_stats_collector" diff --git a/gpcontrib/Makefile b/gpcontrib/Makefile old mode 100644 new mode 100755 index 32c134c95e6..a03c881a92f --- a/gpcontrib/Makefile +++ b/gpcontrib/Makefile @@ -23,6 +23,7 @@ ifeq "$(enable_debug_extensions)" "yes" gp_exttable_fdw \ gp_legacy_string_agg \ gp_relsizes_stats \ + gp_relaccess_stats \ gp_replica_check \ gp_toolkit \ pg_hint_plan \ @@ -33,6 +34,7 @@ else gp_internal_tools \ gp_legacy_string_agg \ gp_relsizes_stats \ + gp_relaccess_stats \ gp_exttable_fdw \ gp_toolkit \ pg_hint_plan diff --git a/gpcontrib/gp_relaccess_stats/.gitignore b/gpcontrib/gp_relaccess_stats/.gitignore new file mode 100644 index 00000000000..8031bbbf6eb --- /dev/null +++ b/gpcontrib/gp_relaccess_stats/.gitignore @@ -0,0 +1,5 @@ +*.o +*.so +.vscode +compile_commands.json +results diff --git a/gpcontrib/gp_relaccess_stats/Makefile b/gpcontrib/gp_relaccess_stats/Makefile new file mode 100755 index 00000000000..eb973c33c34 --- /dev/null +++ b/gpcontrib/gp_relaccess_stats/Makefile @@ -0,0 +1,19 @@ +MODULE_big = gp_relaccess_stats +OBJS = ./src/gp_relaccess_stats.o +EXTENSION = gp_relaccess_stats +EXTVERSION = 1.0 +DATA = $(wildcard sql/*--*.sql) +REGRESS = gp_relaccess_stats +REGRESS_OPTS = --inputdir=test/ +PGFILEDESC = "gp_relaccess_stats - facility to track how and when tables, partitions or views were accessed" +PG_CXXFLAGS += $(COMMON_CPP_FLAGS) + +ifdef USE_PGXS + PG_CONFIG = pg_config + PGXS := $(shell $(PG_CONFIG) --pgxs) + include $(PGXS) +else + top_builddir = ../.. + include $(top_builddir)/src/Makefile.global + include $(top_srcdir)/contrib/contrib-global.mk +endif diff --git a/gpcontrib/gp_relaccess_stats/README.md b/gpcontrib/gp_relaccess_stats/README.md new file mode 100644 index 00000000000..93614978c71 --- /dev/null +++ b/gpcontrib/gp_relaccess_stats/README.md @@ -0,0 +1,105 @@ + + +# gp_relaccess_stats: Table access monitoring tool for Greenplum + +## Features +gp_relaccess_stats is an extension that records access statistics for Greenplum tables and views. Allowing users to see what objects were used, when and by whom. For example, this allows DBAs to find objects that are not used anymore or objects that are being misused. + +Features include: +* support of both tables (regular, external or partitioned) and views +* separate tracking of select, insert, update and delete queries +* separate tracking of last read and write timestamps +* tracking of the last user who accessed the object +* per-database configuration +* in-memory stats survive server restarts (but not crashes) + +### Supported versions and platforms +For now it is being tested only for GP6 and Linux. Though, there are no apparent reasons why it should not be working on newer GP versions (or even PG with slight code modification) or other OSes. + +### Installation +Install from source: +```bash +# get the source code somewhere +git clone git@github.com:Smyatkin-Maxim/gp_relaccess_stats.git +cd gp_relaccess_stats +# Build it. Building would require GP installed nearby and sourcing greenplum_path.sh +source /greenplum_path.sh +make && make install +``` + +### Configuration +As this extension does extensive usage of hooks and shared memory, you need to load gp_relaccess_stats.so on start-up: +``` +gpconfig -c shared_preload_libraries -v 'gp_relaccess_stats' && gpstop -ra +``` +gp_relaccess_stats configuration parameters: +| **Parameter** | **Type** | **Default** | **Default** | +| ---------------- | --------------- | ------------ | ------------ | +| `gp_relaccess_stats.enabled` | bool | false | Using `gp_relaccess_stats.enabled` you can enable/disable stats collection either globally or for each database separately. The second option is preferred.| +| `gp_relaccess_stats.max_tables` | integer | 65536 | `gp_relaccess_stats.max_tables` is a hard limit on how many tables can be cached in shared memory. Feel free to make this number higher if necessary, as the overhead is only about 160 bytes per table. Note, that stats cache for a specific table is evicted from memory any time you execute `relaccess_stats_update()` or `relaccess_stats_dump()` and new tables can be recorded. If you call these functions often enough, there is no need for high gp_relaccess_stats.max_tables| +| `gp_relaccess_stats.dump_on_overflow` | bool | false | This parameter configures what happens in case `gp_relaccess_stats.max_tables` was not enough. If set to `true`, `relaccess_stats_dump()` will be called implicitly and stats cache will be freed. Otherwice, you will get a WARNING saying that there is no room for new stats. Is this case, stats for some tables will be lost.| + +### Usage +The first thing you need to do after `CREATE EXTENSION` and configuring - execute `SELECT relaccess_stats_init();` in a specific database. This function will fill `relaccess_stats` table with empty stats for each table and partition in this database. This is optional, but will come handy when you try to find tables that haven't been used recently, for example. + +Then, either manually or with a cron job start executing `select relaccess_stats_update()`. This function takes all stats cached in shared memory and all stats stored in pg_stat dir (e.g, dumps after restarts, or when `max_tables` was exceeded) and upserts them into `relaccess_stats` table. + +The `relaccess_stats` table itself looks like this: +| **Column** | **Description** | +| ---------------- | --------------- | +| relid | OID of the relation | +| relname | Name of the relation at last access | +| last_reader_id | OID of user who read the table last | +| last_writer_id | OID of user who wrote the table last | +| last_read | Timestamp of the most recent select | +| last_write | Timestamp of the most recent insert/delete/update/truncate | +| n_select_queries | | +| n_select_queries | | +| n_select_queries | | +| n_select_queries | | +| n_truncate_queries | | + +**NOTE**: n_*_queries columns count the number of queries executed, not the number of rows read, inserted, deleted or updated. + +This table has a view associated with it: `relaccess_stats_root_tables_aggregated`. This view has exactly same columns, however it only shows partitioned tables. To be more specific, it shows aggregated stats for each partitioned table. +For example, if we have 1 insert into `tbl1_prt_1` and 3 inserts into `tbl1_prt_2`, then `select * from relaccess_stats_root_tables_aggregated where relname = 'tbl1'` will show us only root table with n_insert_queries = 4. This view, however, has some limitations. See the next section for more detail. + +Another useful function is `relaccess_stats_dump()`, which simply moves cached stats from shared memory to temporary files in pg_stat directory. This function is cheaper than `relaccess_stats_update` but will evict stats cache if needed. Though, stats in temporary files can also get lost. Hence, it is recommended to stick with frequent `select relaccess_stats_update()` calls. + +To better understand when it's time to dump or update the stats one might check `select relaccess.relaccess_stats_fillfactor();`. It will show current usage of stats hash table in percents. For example if shared memory for our relaccess hash table is 70% full we will get relaccess_stats_fillfactor=70. It would be a good idea to dump or update when fillfactor is around 70%. + +### Limitations and gotchas +There is a number of interesting edge-cases in this simple extension: +* `relaccess_stats_root_tables_aggregated` shows info only about tables that exist **now**. We simply can`t get information about inheritance relationship for deleted tables. +* Stats don't rollback on savepoint rollback. We will see n_select_queries incremented by 2 in the following case: +```sql +BEGIN; +SELECT * FROM tbl; +SAVEPOINT sp; +SELECT * FROM tbl; +ROLLBACK TO SAVEPOINT sp; +COMMIT; +``` +There is no technical reason for this limitation. It cat be fixed when there will be need for that. +* Update stats often! Otherwise, data can be lost if any of it happens: 1) there was a crash, 2) `max_tables` exceeded w/o `dump_on_overflow`, 3) temporary pg_stat dir got cleaned. +* no `truncate only` support. There is a TODO in code in case it is ever needed. +* Updates and Deletes also increment n_select_queries. Every update and delete also read the table. That is, n_select_queries get incremented as well. If you need **only** selects, query like this `SELECT n_select_queries - (n_update_queries + n_delete_queries) ... FROM relaccess_stats ...;`. For this same reason last_read and last_reader_id change on update and delete queries. +* view = view + tables. It looks like whenever you select from view, n_select_queries get incremented for both the view and tables it references. +* obviously, we don't know any timestamps before we started tracking. So, the first timestamps are initialized with 0 (something around year 2000), which means those tables haven't been accessed since gp_relaccess_stats was enabled. diff --git a/gpcontrib/gp_relaccess_stats/gp_relaccess_stats.control b/gpcontrib/gp_relaccess_stats/gp_relaccess_stats.control new file mode 100644 index 00000000000..e1c5f5189b9 --- /dev/null +++ b/gpcontrib/gp_relaccess_stats/gp_relaccess_stats.control @@ -0,0 +1,6 @@ +# gp_relaccess_stats extension +comment = 'gp_relaccess_stats - facility to track how and when tables, partitions or views were accesseds' +default_version = '1.1' +module_pathname = '$libdir/gp_relaccess_stats' +relocatable = true +trusted = true diff --git a/gpcontrib/gp_relaccess_stats/sql/gp_relaccess_stats--1.0--1.1.sql b/gpcontrib/gp_relaccess_stats/sql/gp_relaccess_stats--1.0--1.1.sql new file mode 100644 index 00000000000..799474e6f2f --- /dev/null +++ b/gpcontrib/gp_relaccess_stats/sql/gp_relaccess_stats--1.0--1.1.sql @@ -0,0 +1,52 @@ +/* gp_relaccess_stats--1.0--1.1.sql */ + +-- complain if script is sourced in psql, rather than via CREATE EXTENSION +\echo Use "ALTER EXTENSION gp_relaccess_stats UPDATE TO '1.1'" to load this file. \quit + +DROP VIEW relaccess.relaccess_stats_root_tables_aggregated; + +ALTER TABLE relaccess.relaccess_stats ALTER COLUMN n_select_queries TYPE int8; +ALTER TABLE relaccess.relaccess_stats ALTER COLUMN n_insert_queries TYPE int8; +ALTER TABLE relaccess.relaccess_stats ALTER COLUMN n_update_queries TYPE int8; +ALTER TABLE relaccess.relaccess_stats ALTER COLUMN n_delete_queries TYPE int8; +ALTER TABLE relaccess.relaccess_stats ALTER COLUMN n_truncate_queries TYPE int8; + +-- This utility view shows **ONLY** stats on **EXISTING** partitioned tables in aggregated form +CREATE VIEW relaccess.relaccess_stats_root_tables_aggregated AS ( + WITH RECURSIVE parents AS ( + SELECT inhrelid AS child, inhparent AS parent FROM pg_inherits + UNION ALL + SELECT prev.child, next.inhparent AS parent FROM parents AS prev JOIN pg_inherits AS next ON prev.parent = next.inhrelid + ), part_to_root_mapping AS ( + SELECT DISTINCT child AS partid, min(parent) OVER (partition BY child) AS rootid FROM parents + ), parts_including_roots AS ( + SELECT rootid as partid, rootid FROM (SELECT DISTINCT rootid FROM part_to_root_mapping) AS p + UNION + SELECT * FROM part_to_root_mapping + ), with_root_id AS ( + SELECT part_tbl.rootid, stats.* FROM relaccess.relaccess_stats stats JOIN parts_including_roots part_tbl ON (stats.relid = part_tbl.partid) + ), without_last_user AS ( + SELECT rootid AS relid, + rootid::regclass::text AS relname, + max(last_read) AS last_read, + max(last_write) AS last_write, + sum(n_select_queries) AS n_select_queries, + sum(n_insert_queries) AS n_insert_queries, + sum(n_update_queries) AS n_update_queries, + sum(n_delete_queries) AS n_delete_queries, + sum(n_truncate_queries) AS n_truncate_queries + FROM with_root_id outer_tbl GROUP BY rootid + ) + SELECT relid, + relname, + (SELECT last_reader_id FROM with_root_id w WHERE w.rootid = wo.relid AND wo.last_read = w.last_read LIMIT 1) AS last_reader_id, + (SELECT last_writer_id FROM with_root_id w WHERE w.rootid = wo.relid AND wo.last_write = w.last_write LIMIT 1) AS last_writer_id, + last_read, + last_write, + n_select_queries, + n_insert_queries, + n_update_queries, + n_delete_queries, + n_truncate_queries + FROM without_last_user wo +); diff --git a/gpcontrib/gp_relaccess_stats/sql/gp_relaccess_stats--1.1.sql b/gpcontrib/gp_relaccess_stats/sql/gp_relaccess_stats--1.1.sql new file mode 100755 index 00000000000..6d25b63132d --- /dev/null +++ b/gpcontrib/gp_relaccess_stats/sql/gp_relaccess_stats--1.1.sql @@ -0,0 +1,143 @@ +/* gp_relaccess_stats--1.1.sql */ + +-- complain if script is sourced in psql, rather than via CREATE EXTENSION +\echo Use "CREATE EXTENSION gp_relaccess_stats" to load this file. \quit + +CREATE SCHEMA IF NOT EXISTS relaccess; + +CREATE TABLE relaccess.relaccess_stats ( + relid Oid, + relname Name, + last_reader_id Oid, + last_writer_id Oid, + last_read timestamptz, + last_write timestamptz, + n_select_queries int8, + n_insert_queries int8, + n_update_queries int8, + n_delete_queries int8, + n_truncate_queries int8 +) DISTRIBUTED BY (relid); + +CREATE FUNCTION relaccess.relaccess_stats_dump() +RETURNS SETOF void +AS 'MODULE_PATHNAME', 'relaccess_stats_dump' +LANGUAGE C VOLATILE EXECUTE ON MASTER; + +CREATE FUNCTION relaccess.relaccess_stats_update() +RETURNS SETOF void +AS 'MODULE_PATHNAME', 'relaccess_stats_update' +LANGUAGE C VOLATILE EXECUTE ON MASTER; + +CREATE FUNCTION relaccess.relaccess_stats_fillfactor() +RETURNS SETOF INT2 +AS 'MODULE_PATHNAME', 'relaccess_stats_fillfactor' +LANGUAGE C VOLATILE EXECUTE ON MASTER; + +CREATE FUNCTION relaccess.__get_db_stats_from_dump() +RETURNS SETOF relaccess.relaccess_stats +AS 'MODULE_PATHNAME', 'relaccess_stats_from_dump' +LANGUAGE C VOLATILE EXECUTE ON MASTER; + +CREATE FUNCTION relaccess.__relaccess_upsert_from_dump_file() RETURNS VOID +LANGUAGE plpgsql VOLATILE AS +$func$ +BEGIN + EXECUTE 'DROP TABLE IF EXISTS relaccess_stats_tmp'; + EXECUTE 'CREATE TEMP TABLE relaccess_stats_tmp (LIKE relaccess.relaccess_stats) distributed by (relid)'; + EXECUTE 'DROP TABLE IF EXISTS relaccess_stats_tmp_aggregated'; + EXECUTE 'CREATE TEMP TABLE relaccess_stats_tmp_aggregated (LIKE relaccess.relaccess_stats) distributed by (relid)'; + EXECUTE 'INSERT INTO relaccess_stats_tmp SELECT * FROM relaccess.__get_db_stats_from_dump()'; + EXECUTE 'WITH aggregated_wo_relname_and_user AS ( + SELECT relid, max(last_read) AS last_read, max(last_write) AS last_write, sum(n_select_queries) AS n_select_queries, + sum(n_insert_queries) AS n_insert_queries, sum(n_update_queries) AS n_update_queries, sum(n_delete_queries) AS n_delete_queries, sum(n_truncate_queries) AS n_truncate_queries + FROM relaccess_stats_tmp GROUP BY relid + ) + INSERT INTO relaccess_stats_tmp_aggregated + SELECT relid, + (SELECT relname FROM relaccess_stats_tmp w WHERE w.relid = wo.relid AND greatest(wo.last_read, wo.last_write) IN (w.last_read, w.last_write) LIMIT 1) AS relname, + (SELECT last_reader_id FROM relaccess_stats_tmp w WHERE w.relid = wo.relid AND wo.last_read = w.last_read LIMIT 1) AS last_reader_id, + (SELECT last_writer_id FROM relaccess_stats_tmp w WHERE w.relid = wo.relid AND wo.last_write = w.last_write LIMIT 1) AS last_writer_id, + last_read, + last_write, + n_select_queries, + n_insert_queries, + n_update_queries, + n_delete_queries, + n_truncate_queries FROM aggregated_wo_relname_and_user AS wo'; + EXECUTE 'DROP TABLE IF EXISTS relaccess_stats_tmp'; + EXECUTE 'INSERT INTO relaccess.relaccess_stats + SELECT relid, relname, last_reader_id, last_writer_id, last_read, last_write, 0, 0, 0, 0, 0 + FROM relaccess_stats_tmp_aggregated stage + WHERE NOT EXISTS ( + SELECT 1 FROM relaccess.relaccess_stats orig WHERE orig.relid = stage.relid)'; + EXECUTE 'UPDATE relaccess.relaccess_stats orig SET + relname = stage.relname, + n_select_queries = orig.n_select_queries + stage.n_select_queries, + n_insert_queries = orig.n_insert_queries + stage.n_insert_queries, + n_update_queries = orig.n_update_queries + stage.n_update_queries, + n_delete_queries = orig.n_delete_queries + stage.n_delete_queries, + n_truncate_queries = orig.n_truncate_queries + stage.n_truncate_queries + FROM relaccess_stats_tmp_aggregated stage + WHERE orig.relid = stage.relid'; + EXECUTE 'UPDATE relaccess.relaccess_stats orig SET + last_reader_id = stage.last_reader_id, last_read = stage.last_read + FROM relaccess_stats_tmp_aggregated stage + WHERE orig.relid = stage.relid AND orig.last_read < stage.last_read'; + EXECUTE 'UPDATE relaccess.relaccess_stats orig SET + last_writer_id = stage.last_writer_id, last_write = stage.last_write + FROM relaccess_stats_tmp_aggregated stage + WHERE orig.relid = stage.relid AND orig.last_write < stage.last_write'; + EXECUTE 'DROP TABLE IF EXISTS relaccess_stats_tmp_aggregated'; +END +$func$; + +CREATE FUNCTION relaccess.relaccess_stats_init() RETURNS VOID AS +$$ + WITH relations AS ( + SELECT oid as relid, relname, relowner FROM pg_catalog.pg_class WHERE relkind in ('r', 'v', 'm', 'f', 'p') + ) + INSERT INTO relaccess.relaccess_stats + SELECT relid, relname, relowner, relowner, '2000-01-01 03:00:00', '2000-01-01 03:00:00', 0, 0, 0, 0, 0 + FROM relations AS all_rels WHERE NOT EXISTS(SELECT 1 FROM relaccess.relaccess_stats orig WHERE orig.relid = all_rels.relid); +$$ LANGUAGE SQL VOLATILE; + +-- This utility view shows **ONLY** stats on **EXISTING** partitioned tables in aggregated form +CREATE VIEW relaccess.relaccess_stats_root_tables_aggregated AS ( + WITH RECURSIVE parents AS ( + SELECT inhrelid AS child, inhparent AS parent FROM pg_inherits + UNION ALL + SELECT prev.child, next.inhparent AS parent FROM parents AS prev JOIN pg_inherits AS next ON prev.parent = next.inhrelid + ), part_to_root_mapping AS ( + SELECT DISTINCT child AS partid, min(parent) OVER (partition BY child) AS rootid FROM parents + ), parts_including_roots AS ( + SELECT rootid as partid, rootid FROM (SELECT DISTINCT rootid FROM part_to_root_mapping) AS p + UNION + SELECT * FROM part_to_root_mapping + ), with_root_id AS ( + SELECT part_tbl.rootid, stats.* FROM relaccess.relaccess_stats stats JOIN parts_including_roots part_tbl ON (stats.relid = part_tbl.partid) + ), without_last_user AS ( + SELECT rootid AS relid, + rootid::regclass::text AS relname, + max(last_read) AS last_read, + max(last_write) AS last_write, + sum(n_select_queries) AS n_select_queries, + sum(n_insert_queries) AS n_insert_queries, + sum(n_update_queries) AS n_update_queries, + sum(n_delete_queries) AS n_delete_queries, + sum(n_truncate_queries) AS n_truncate_queries + FROM with_root_id outer_tbl GROUP BY rootid + ) + SELECT relid, + relname, + (SELECT last_reader_id FROM with_root_id w WHERE w.rootid = wo.relid AND wo.last_read = w.last_read LIMIT 1) AS last_reader_id, + (SELECT last_writer_id FROM with_root_id w WHERE w.rootid = wo.relid AND wo.last_write = w.last_write LIMIT 1) AS last_writer_id, + last_read, + last_write, + n_select_queries, + n_insert_queries, + n_update_queries, + n_delete_queries, + n_truncate_queries + FROM without_last_user wo +); diff --git a/gpcontrib/gp_relaccess_stats/src/gp_relaccess_stats.c b/gpcontrib/gp_relaccess_stats/src/gp_relaccess_stats.c new file mode 100755 index 00000000000..fdb05f60459 --- /dev/null +++ b/gpcontrib/gp_relaccess_stats/src/gp_relaccess_stats.c @@ -0,0 +1,817 @@ +#include "postgres.h" +#include "access/table.h" +#include "access/xact.h" +#include "access/hash.h" +#include "catalog/objectaccess.h" +#include "catalog/pg_database.h" +#include "cdb/cdbvars.h" +#include "commands/dbcommands.h" +#include "executor/executor.h" +#include "executor/spi.h" +#include "funcapi.h" +#include "miscadmin.h" +#include "pg_config_ext.h" +#include "pgstat.h" +#include "storage/ipc.h" +#include "storage/lwlock.h" +#include "storage/shmem.h" +#include "storage/spin.h" +#include "utils/builtins.h" +#include "utils/datetime.h" +#include "utils/lsyscache.h" +#include "utils/memutils.h" +#include "utils/timestamp.h" +#include "tcop/utility.h" + +#include +#include +#include + +/** + * gp_relaccess_stats collects runtime access stats on db objects: relations and + * views. Stats include last read and write timestamps, last user, last known + * relname and number of select, insert, update, delete or truncate queries. + * Only committed actions are recorded. + * + * To track those actions we use: + * - ExecutorCheckPerms hook for select, insert, update and delete statements + * - ProcessUtility hook for truncate statements + * + * Intermediate data is stored in three hash tables. + * One lives in shared memory and is cleaned only when dumped to disc: + * - relaccesses - represents all recorded accesses since last dump to disc. + * And two live in coordinator`s local memory and are cleaned on every commit + * or rollback: + * - local_access_entries - represent all record accesses in for this + * transaction only + * - relname_cache - maps relid to relname for relations used in this + * transaction only + * + * Ultimately all recorded stats should end up in relaccess_stats table when a + * user executes relaccess_stats_update(). But any intermediate stats will be + * dumped to disc. This might happen for either or those reasons: + * - shmem is exceeded + * - server is restarted + * - manual execution of relaccess_stats_dump() + * In this case stats are offloaded to disc into pg_stat directory into separate + * file per each tracked database: pg_stat/relaccess_stats_dump_.csv Those + * files are upserted into relaccess_stats when relaccess_stats_update() is + * called + */ + +PG_MODULE_MAGIC; + +void _PG_init(void); +void _PG_fini(void); +PG_FUNCTION_INFO_V1(relaccess_stats_update); +PG_FUNCTION_INFO_V1(relaccess_stats_dump); +PG_FUNCTION_INFO_V1(relaccess_stats_fillfactor); +PG_FUNCTION_INFO_V1(relaccess_stats_from_dump); + +static void relaccess_stats_update_internal(void); +static void relaccess_dump_to_files(bool only_this_db); +static void relaccess_dump_to_files_internal(HTAB *files); +static void relaccess_upsert_from_file(void); +static void relaccess_shmem_startup(void); +static void relaccess_shmem_shutdown(int code, Datum arg); +static uint32 relaccess_hash_fn(const void *key, Size keysize); +static int relaccess_match_fn(const void *key1, const void *key2, Size keysize); +static uint32 local_relaccess_hash_fn(const void *key, Size keysize); +static int local_relaccess_match_fn(const void *key1, const void *key2, + Size keysize); +static bool collect_relaccess_hook(List *rangeTable, bool ereport_on_violation); +static void relaccess_xact_callback(XactEvent event, void *arg); +static void collect_truncate_hook(PlannedStmt *pstmt, const char *queryString, + bool readOnlyTree, + ProcessUtilityContext context, + ParamListInfo params, + QueryEnvironment *queryEnv, + DestReceiver *dest, QueryCompletion *qc); +static void relaccess_executor_end_hook(QueryDesc *query_desc); +static void relaccess_drop_hook(ObjectAccessType access, Oid classId, + Oid objectId, int subId, void *arg); +static void memorize_local_access_entry(Oid relid, AclMode perms); +static void update_relname_cache(Oid relid, char *relname); +static StringInfoData get_dump_filename(Oid dbid); + +static shmem_startup_hook_type prev_shmem_startup_hook = NULL; +static ExecutorCheckPerms_hook_type prev_check_perms_hook = NULL; +static ProcessUtility_hook_type next_ProcessUtility_hook = NULL; +static ExecutorEnd_hook_type prev_ExecutorEnd_hook = NULL; +static object_access_hook_type prev_object_access_hook = NULL; + +typedef struct relaccessHashKey { + Oid dbid; + Oid relid; +} relaccessHashKey; + +typedef struct relaccessEntry { + relaccessHashKey key; + char relname[NAMEDATALEN]; + Oid last_reader_id; + Oid last_writer_id; + TimestampTz last_read; + TimestampTz last_write; + int64 n_select; + int64 n_insert; + int64 n_update; + int64 n_delete; + int64 n_truncate; +} relaccessEntry; + +typedef struct relaccessGlobalData { + LWLock *relaccess_ht_lock; + LWLock *relaccess_file_lock; +} relaccessGlobalData; + +typedef struct localAccessKey { + Oid relid; + int stmt_cnt; +} localAccessKey; + +typedef struct localAccessEntry { + localAccessKey key; + Oid last_reader_id, last_writer_id; + Timestamp last_read, last_write; + AclMode perms; +} localAccessEntry; + +typedef struct relnameCacheEntry { + Oid relid; + char relname[NAMEDATALEN]; +} relnameCacheEntry; + +typedef struct fileDumpEntry { + Oid dbid; + char *filename; + FILE *file; +} fileDumpEntry; + +static int32 relaccess_size; +static bool dump_on_overflow; +static bool is_enabled; +static relaccessGlobalData *data; +static HTAB *relaccesses; +static HTAB *local_access_entries = NULL; +static const int32 LOCAL_HTAB_SZ = 128; +static HTAB *relname_cache = NULL; +static const int32 RELCACHE_SZ = 16; +static const int32 FILE_CACHE_SZ = 16; +static int stmt_counter = 0; +static bool had_ht_overflow = false; + +#define IS_POSTGRES_DB \ + (strcmp("postgres", get_database_name(MyDatabaseId)) == 0) + +#define is_write(perms) \ + (((perms) & (ACL_INSERT | ACL_UPDATE | ACL_DELETE | ACL_TRUNCATE)) != 0) + +#define is_read(perms) (!is_write(perms) && ((perms) & ACL_SELECT) != 0) + +static void relaccess_shmem_startup() { + bool found; + HASHCTL info; + + if (prev_shmem_startup_hook) + prev_shmem_startup_hook(); + + LWLockAcquire(AddinShmemInitLock, LW_EXCLUSIVE); + + data = (relaccessGlobalData *)(ShmemInitStruct( + "relaccess_stats", sizeof(relaccessGlobalData), &found)); + if (!found) { + LWLockPadded *locks = GetNamedLWLockTranche("gp_relaccess_stats"); + data->relaccess_ht_lock = &locks[0].lock; + data->relaccess_file_lock = &locks[1].lock; + } + + memset(&info, 0, sizeof(info)); + info.keysize = sizeof(relaccessHashKey); + info.entrysize = sizeof(relaccessEntry); + info.hash = relaccess_hash_fn; + info.match = relaccess_match_fn; + relaccesses = ShmemInitHash( + "relaccess_stats hash", relaccess_size, relaccess_size, &info, + (HASH_ELEM | HASH_FUNCTION | HASH_COMPARE | HASH_FIXED_SIZE)); + + LWLockRelease(AddinShmemInitLock); + + if (!IsUnderPostmaster) { + on_shmem_exit(relaccess_shmem_shutdown, (Datum)0); + } +} + +static void relaccess_shmem_shutdown(int code, Datum arg) { + if (code || !data || !relaccesses) { + return; + } + LWLockAcquire(data->relaccess_ht_lock, LW_EXCLUSIVE); + relaccess_dump_to_files(false); + LWLockRelease(data->relaccess_ht_lock); +} + +static uint32 relaccess_hash_fn(const void *key, Size keysize) { + const relaccessHashKey *k = (const relaccessHashKey *)key; + return hash_uint32((uint32)k->dbid) ^ hash_uint32((uint32)k->relid); +} + +static int relaccess_match_fn(const void *key1, const void *key2, + Size keysize) { + const relaccessHashKey *k1 = (const relaccessHashKey *)key1; + const relaccessHashKey *k2 = (const relaccessHashKey *)key2; + return (k1->dbid == k2->dbid && k1->relid == k2->relid) ? 0 : 1; +} + +static uint32 local_relaccess_hash_fn(const void *key, Size keysize) { + const localAccessKey *k = (const localAccessKey *)key; + return hash_uint32((uint32)k->stmt_cnt) ^ hash_uint32((uint32)k->relid); +} + +static int local_relaccess_match_fn(const void *key1, const void *key2, + Size keysize) { + const localAccessKey *k1 = (const localAccessKey *)key1; + const localAccessKey *k2 = (const localAccessKey *)key2; + return (k1->stmt_cnt == k2->stmt_cnt && k1->relid == k2->relid) ? 0 : 1; +} + +void _PG_init(void) { + Size size; + if (Gp_role != GP_ROLE_DISPATCH) { + return; + } + if (!process_shared_preload_libraries_in_progress) { + return; + } + + DefineCustomIntVariable( + "gp_relaccess_stats.max_tables", + "Sets the maximum number of tables cached by gp_relaccess_stats.", NULL, + &relaccess_size, 65536, 128, INT_MAX, PGC_POSTMASTER, 0, NULL, NULL, + NULL); + + DefineCustomBoolVariable("gp_relaccess_stats.dump_on_overflow", + "Selects whether we should dump to .csv in case " + "gp_relaccess_stats.max_tables is exceeded.", + NULL, &dump_on_overflow, false, PGC_SIGHUP, 0, NULL, + NULL, NULL); + + DefineCustomBoolVariable( + "gp_relaccess_stats.enabled", + "Collect table access stats globally or for a specific database. " + "Note that shared memory is initialized indepemdent of this argument.", + NULL, &is_enabled, false, PGC_SUSET, 0, NULL, NULL, NULL); + + prev_shmem_startup_hook = shmem_startup_hook; + shmem_startup_hook = relaccess_shmem_startup; + prev_check_perms_hook = ExecutorCheckPerms_hook; + ExecutorCheckPerms_hook = collect_relaccess_hook; + next_ProcessUtility_hook = ProcessUtility_hook; + ProcessUtility_hook = collect_truncate_hook; + prev_ExecutorEnd_hook = ExecutorEnd_hook; + ExecutorEnd_hook = relaccess_executor_end_hook; + prev_object_access_hook = object_access_hook; + object_access_hook = relaccess_drop_hook; + RequestNamedLWLockTranche("gp_relaccess_stats", 2); + size = MAXALIGN(sizeof(relaccessGlobalData)); + size = add_size(size, + hash_estimate_size(relaccess_size, sizeof(relaccessEntry))); + RequestAddinShmemSpace(size); + RegisterXactCallback(relaccess_xact_callback, NULL); + HASHCTL ctl; + MemSet(&ctl, 0, sizeof(ctl)); + ctl.keysize = sizeof(localAccessKey); + ctl.entrysize = sizeof(localAccessEntry); + ctl.hash = local_relaccess_hash_fn; + ctl.match = local_relaccess_match_fn; + local_access_entries = + hash_create("Transaction-wide relaccess entries", LOCAL_HTAB_SZ, &ctl, + HASH_ELEM | HASH_FUNCTION | HASH_COMPARE); + MemSet(&ctl, 0, sizeof(ctl)); + ctl.keysize = sizeof(Oid); + ctl.entrysize = sizeof(relnameCacheEntry); + ctl.hash = oid_hash; + relname_cache = hash_create("Transaction-wide relation name cache", + RELCACHE_SZ, &ctl, HASH_ELEM | HASH_FUNCTION); +} + +void _PG_fini(void) { + if (Gp_role != GP_ROLE_DISPATCH) { + return; + } + shmem_startup_hook = prev_shmem_startup_hook; + ExecutorCheckPerms_hook = prev_check_perms_hook; + ProcessUtility_hook = next_ProcessUtility_hook; + ExecutorEnd_hook = prev_ExecutorEnd_hook; + object_access_hook = prev_object_access_hook; +} + +static bool collect_relaccess_hook(List *rangeTable, + bool ereport_on_violation) { + if (prev_check_perms_hook && + !prev_check_perms_hook(rangeTable, ereport_on_violation)) { + return false; + } + if (Gp_role == GP_ROLE_DISPATCH && is_enabled) { + ListCell *l; + foreach (l, rangeTable) { + RangeTblEntry *rte = (RangeTblEntry *)lfirst(l); + if (rte->rtekind != RTE_RELATION) { + continue; + } + Oid relid = rte->relid; + AclMode requiredPerms = rte->requiredPerms; + if (is_read(requiredPerms) || is_write(requiredPerms)) { + memorize_local_access_entry(relid, requiredPerms); + update_relname_cache(relid, NULL); + } + } + } + return true; +} + +static void collect_truncate_hook(PlannedStmt *pstmt, const char *queryString, + bool readOnlyTree, + ProcessUtilityContext context, + ParamListInfo params, + QueryEnvironment *queryEnv, + DestReceiver *dest, QueryCompletion *qc) { + Node *parsetree = pstmt->utilityStmt; + if (nodeTag(parsetree) == T_TruncateStmt && is_enabled && + Gp_role == GP_ROLE_DISPATCH) { + TruncateStmt *stmt = (TruncateStmt *)parsetree; + ListCell *cell; + /** + * TODO: TRUNCATE may be called with ONLY option which limits it only to + *the root partition. Otherwise it will truncate all child partitions. We + *might wish to track the difference by explicitly adding records for each + *truncated partition in the future if it proves useful + **/ + foreach (cell, stmt->relations) { + RangeVar *rv = lfirst(cell); + Relation rel = table_openrv(rv, AccessExclusiveLock); + Oid relid = rel->rd_id; + table_close(rel, NoLock); + memorize_local_access_entry(relid, ACL_TRUNCATE); + update_relname_cache(relid, rv->relname); + } + } + if (next_ProcessUtility_hook) { + (*next_ProcessUtility_hook)(pstmt, queryString, readOnlyTree, context, + params, queryEnv, dest, qc); + } else { + standard_ProcessUtility(pstmt, queryString, readOnlyTree, context, params, + queryEnv, dest, qc); + } +} + +#define UPDATE_STAT(lowercase, uppercase) \ + dst_entry->n_##lowercase += (src_entry->perms & ACL_##uppercase ? 1 : 0) + +// if there is a better way to cleanup a postgres hashtable +// w/o recreating it, I didn't find it +#define CLEAR_HTAB(entryType, hmap, key_name) \ + { \ + HASH_SEQ_STATUS hash_seq; \ + entryType *src_entry; \ + hash_seq_init(&hash_seq, hmap); \ + while ((src_entry = hash_seq_search(&hash_seq)) != NULL) { \ + bool found; \ + hash_search(hmap, &src_entry->key_name, HASH_REMOVE, &found); \ + Assert(found); \ + } \ + } + +static void relaccess_xact_callback(XactEvent event, void *arg) { + if (Gp_role != GP_ROLE_DISPATCH || !is_enabled) { + return; + } + // TODO: add support for savepoint rollbacks + Assert(GetCurrentTransactionNestLevel() == 1); + if (event == XACT_EVENT_COMMIT) { + HASH_SEQ_STATUS hash_seq; + localAccessEntry *src_entry; + hash_seq_init(&hash_seq, local_access_entries); + LWLockAcquire(data->relaccess_ht_lock, LW_EXCLUSIVE); + while ((src_entry = hash_seq_search(&hash_seq)) != NULL) { + bool found; + relaccessHashKey key; + key.dbid = MyDatabaseId; + key.relid = src_entry->key.relid; + long n_access_records = hash_get_num_entries(relaccesses); + relaccessEntry *dst_entry = NULL; + Assert(n_access_records <= relaccess_size); + if (n_access_records == relaccess_size) { + // no room for new entries. Perhaps this relid is already being tracked? + dst_entry = + (relaccessEntry *)hash_search(relaccesses, &key, HASH_FIND, &found); + } else { + dst_entry = (relaccessEntry *)hash_search(relaccesses, &key, + HASH_ENTER_NULL, &found); + } + if (dst_entry || dump_on_overflow) { + if (!dst_entry) { + // we are out of shared memory and need to dump + relaccess_dump_to_files(false); + // we MUST have enough space now, unless we were unable to dump + dst_entry = (relaccessEntry *)hash_search(relaccesses, &key, + HASH_ENTER_NULL, &found); + if (!dst_entry) { + // still no memory left + if (!had_ht_overflow) { + elog(WARNING, ("gp_relaccess_stats.max_tables is exceeded and we " + "are unable to dump hashtables to disk. " + "Will start loosing some relaccess stats")); + had_ht_overflow = true; + } + continue; + } else { + had_ht_overflow = false; + } + } + if (!found) { + dst_entry->last_reader_id = InvalidOid; + dst_entry->last_writer_id = InvalidOid; + dst_entry->last_read = 0; + dst_entry->last_write = 0; + dst_entry->n_select = 0; + dst_entry->n_insert = 0; + dst_entry->n_update = 0; + dst_entry->n_delete = 0; + dst_entry->n_truncate = 0; + } + UPDATE_STAT(select, SELECT); + UPDATE_STAT(insert, INSERT); + UPDATE_STAT(update, UPDATE); + UPDATE_STAT(delete, DELETE); + UPDATE_STAT(truncate, TRUNCATE); + if (src_entry->last_read > dst_entry->last_read) { + dst_entry->last_read = src_entry->last_read; + dst_entry->last_reader_id = src_entry->last_reader_id; + } + if (src_entry->last_write > dst_entry->last_write) { + dst_entry->last_write = src_entry->last_write; + dst_entry->last_writer_id = src_entry->last_writer_id; + } + relnameCacheEntry *namecache_entry = (relnameCacheEntry *)hash_search( + relname_cache, &key.relid, HASH_ENTER, &found); + Assert(namecache_entry); + strlcpy(dst_entry->relname, namecache_entry->relname, + sizeof(dst_entry->relname)); + } else { + if (!had_ht_overflow) { + elog(WARNING, "gp_relaccess_stats.max_tables is exceeded! New table " + "events will be lost. " + "Please execute relaccess_stats_update() and consider " + "setting a hihger value"); + } + had_ht_overflow = true; + } + } + LWLockRelease(data->relaccess_ht_lock); + CLEAR_HTAB(localAccessEntry, local_access_entries, key); + CLEAR_HTAB(relnameCacheEntry, relname_cache, relid); + } else if (event == XACT_EVENT_ABORT) { + CLEAR_HTAB(localAccessEntry, local_access_entries, key); + CLEAR_HTAB(relnameCacheEntry, relname_cache, relid); + } +} + +Datum relaccess_stats_update(PG_FUNCTION_ARGS) { + FuncCallContext *funcctx; + + if (SRF_IS_FIRSTCALL()) { + funcctx = SRF_FIRSTCALL_INIT(); + funcctx->max_calls = 1; + relaccess_stats_update_internal(); + } + + funcctx = SRF_PERCALL_SETUP(); + if (funcctx->call_cntr < funcctx->max_calls) { + SRF_RETURN_NEXT(funcctx, (Datum)0); + } + SRF_RETURN_DONE(funcctx); +} + +Datum relaccess_stats_dump(PG_FUNCTION_ARGS) { + FuncCallContext *funcctx; + + if (SRF_IS_FIRSTCALL()) { + funcctx = SRF_FIRSTCALL_INIT(); + funcctx->max_calls = 1; + LWLockAcquire(data->relaccess_ht_lock, LW_EXCLUSIVE); + relaccess_dump_to_files(true); + LWLockRelease(data->relaccess_ht_lock); + } + + funcctx = SRF_PERCALL_SETUP(); + if (funcctx->call_cntr < funcctx->max_calls) { + SRF_RETURN_NEXT(funcctx, (Datum)0); + } + SRF_RETURN_DONE(funcctx); +} + +Datum relaccess_stats_fillfactor(PG_FUNCTION_ARGS) { + FuncCallContext *funcctx; + + if (SRF_IS_FIRSTCALL()) { + int16_t fillfactor; + + funcctx = SRF_FIRSTCALL_INIT(); + funcctx->max_calls = 1; + LWLockAcquire(data->relaccess_ht_lock, LW_SHARED); + fillfactor = hash_get_num_entries(relaccesses) * 100 / relaccess_size; + LWLockRelease(data->relaccess_ht_lock); + funcctx->user_fctx = (void *)(intptr_t)fillfactor; + } + + funcctx = SRF_PERCALL_SETUP(); + if (funcctx->call_cntr < funcctx->max_calls) { + SRF_RETURN_NEXT(funcctx, + Int16GetDatum((int16_t)(intptr_t)funcctx->user_fctx)); + } + SRF_RETURN_DONE(funcctx); +} + +Datum relaccess_stats_from_dump(PG_FUNCTION_ARGS) { + FuncCallContext *funcctx; + List *stats_entries = NIL; + + if (SRF_IS_FIRSTCALL()) { + funcctx = SRF_FIRSTCALL_INIT(); + MemoryContext oldcontext = + MemoryContextSwitchTo(funcctx->multi_call_memory_ctx); + TupleDesc tupdesc = CreateTemplateTupleDesc(11); + TupleDescInitEntry(tupdesc, (AttrNumber)1, "relid", OIDOID, -1 /* typmod */, + 0 /* attdim */); + TupleDescInitEntry(tupdesc, (AttrNumber)2, "relname", NAMEOID, + -1 /* typmod */, 0 /* attdim */); + TupleDescInitEntry(tupdesc, (AttrNumber)3, "last_reader_id", OIDOID, + -1 /* typmod */, 0 /* attdim */); + TupleDescInitEntry(tupdesc, (AttrNumber)4, "last_writer_id", OIDOID, + -1 /* typmod */, 0 /* attdim */); + TupleDescInitEntry(tupdesc, (AttrNumber)5, "last_read", TIMESTAMPTZOID, + -1 /* typmod */, 0 /* attdim */); + TupleDescInitEntry(tupdesc, (AttrNumber)6, "last_write", TIMESTAMPTZOID, + -1 /* typmod */, 0 /* attdim */); + TupleDescInitEntry(tupdesc, (AttrNumber)7, "n_select_queries", INT8OID, + -1 /* typmod */, 0 /* attdim */); + TupleDescInitEntry(tupdesc, (AttrNumber)8, "n_insert_queries", INT8OID, + -1 /* typmod */, 0 /* attdim */); + TupleDescInitEntry(tupdesc, (AttrNumber)9, "n_update_queries", INT8OID, + -1 /* typmod */, 0 /* attdim */); + TupleDescInitEntry(tupdesc, (AttrNumber)10, "n_delete_queries", INT8OID, + -1 /* typmod */, 0 /* attdim */); + TupleDescInitEntry(tupdesc, (AttrNumber)11, "n_truncate_queries", INT8OID, + -1 /* typmod */, 0 /* attdim */); + funcctx->tuple_desc = BlessTupleDesc(tupdesc); + StringInfoData dump_file = get_dump_filename(MyDatabaseId); + FILE *dump = AllocateFile(dump_file.data, "rb"); + pfree(dump_file.data); + if (dump) { + while (true) { + relaccessEntry *entry = palloc(sizeof(relaccessEntry)); + if (fread(entry, sizeof(relaccessEntry), 1, dump) != 1) { + pfree(entry); + break; + } + stats_entries = lappend(stats_entries, entry); + } + FreeFile(dump); + } + funcctx->user_fctx = stats_entries; + MemoryContextSwitchTo(oldcontext); + } + + funcctx = SRF_PERCALL_SETUP(); + stats_entries = (List *)funcctx->user_fctx; + + while (true) { + if (stats_entries == NIL) { + SRF_RETURN_DONE(funcctx); + } + relaccessEntry *entry = linitial(stats_entries); + stats_entries = list_delete_first(stats_entries); + Datum values[11]; + bool nulls[11]; + MemSet(nulls, 0, sizeof(nulls)); + values[0] = ObjectIdGetDatum(entry->key.relid); + values[1] = CStringGetDatum(entry->relname); + values[2] = ObjectIdGetDatum(entry->last_reader_id); + values[3] = ObjectIdGetDatum(entry->last_writer_id); + values[4] = TimestampTzGetDatum(entry->last_read); + values[5] = TimestampTzGetDatum(entry->last_write); + values[6] = Int64GetDatum(entry->n_select); + values[7] = Int64GetDatum(entry->n_insert); + values[8] = Int64GetDatum(entry->n_update); + values[9] = Int64GetDatum(entry->n_delete); + values[10] = Int64GetDatum(entry->n_truncate); + HeapTuple tuple = heap_form_tuple(funcctx->tuple_desc, values, nulls); + Datum result = HeapTupleGetDatum(tuple); + funcctx->user_fctx = stats_entries; + /** NOTE: Cannot delete entry from this iteration right now. + * For now let's rely on multi_call_memory_ctx until there is a proven + * memory problem with this codepath + */ + // pfree(entry); + SRF_RETURN_NEXT(funcctx, result); + } +} + +static void relaccess_stats_update_internal() { + LWLockAcquire(data->relaccess_ht_lock, LW_EXCLUSIVE); + relaccess_dump_to_files(true); + LWLockRelease(data->relaccess_ht_lock); + relaccess_upsert_from_file(); +} + +static void add_file_dump_entry(Oid dbid, HTAB *ht) { + bool found; + fileDumpEntry *file_entry = hash_search(ht, &dbid, HASH_ENTER, &found); + if (!found) { + file_entry->dbid = dbid; + StringInfoData filename = get_dump_filename(file_entry->dbid); + file_entry->filename = filename.data; + file_entry->file = AllocateFile(file_entry->filename, "ab"); + } +} + +static void relaccess_dump_to_files(bool only_this_db) { + HTAB *file_mapping; + HASHCTL ctl; + MemSet(&ctl, 0, sizeof(ctl)); + ctl.keysize = sizeof(Oid); + ctl.entrysize = sizeof(fileDumpEntry); + ctl.hash = oid_hash; + file_mapping = hash_create("Relaccess dump files", FILE_CACHE_SZ, &ctl, + HASH_ELEM | HASH_FUNCTION); + LWLockAcquire(data->relaccess_file_lock, LW_EXCLUSIVE); + if (only_this_db) { + add_file_dump_entry(MyDatabaseId, file_mapping); + } else { + HASH_SEQ_STATUS hash_seq; + relaccessEntry *access_entry; + hash_seq_init(&hash_seq, relaccesses); + while ((access_entry = hash_seq_search(&hash_seq)) != NULL) { + add_file_dump_entry(access_entry->key.dbid, file_mapping); + } + } + relaccess_dump_to_files_internal(file_mapping); + HASH_SEQ_STATUS hash_seq; + hash_seq_init(&hash_seq, file_mapping); + fileDumpEntry *entry; + while ((entry = hash_seq_search(&hash_seq)) != NULL) { + FreeFile(entry->file); + pfree(entry->filename); + } + LWLockRelease(data->relaccess_file_lock); + hash_destroy(file_mapping); +} + +static void relaccess_dump_to_files_internal(HTAB *files) { + HASH_SEQ_STATUS hash_seq; + relaccessEntry *entry; + hash_seq_init(&hash_seq, relaccesses); + while ((entry = hash_seq_search(&hash_seq)) != NULL) { + bool found; + fileDumpEntry *dumpfile = + hash_search(files, &entry->key.dbid, HASH_FIND, &found); + if (!found) { + // we don't want to dump events from this DB + continue; + } + if (fwrite(entry, sizeof(relaccessEntry), 1, dumpfile->file) != 1) { + hash_seq_term(&hash_seq); + ereport(WARNING, + (errcode_for_file_access(), + errmsg("could not write gp_relaccess_stats file \"%s\": %m", + dumpfile->filename))); + break; + } + hash_search(relaccesses, &entry->key, HASH_REMOVE, &found); + had_ht_overflow = false; + } +} + +static void relaccess_upsert_from_file() { + int ret; + if ((ret = SPI_connect()) < 0) { + elog(ERROR, "SPI connect failure - returned %d", ret); + } + LWLockAcquire(data->relaccess_file_lock, LW_EXCLUSIVE); + StringInfoData filename = get_dump_filename(MyDatabaseId); + StringInfoData query; + initStringInfo(&query); + appendStringInfo(&query, + "SELECT relaccess.__relaccess_upsert_from_dump_file()"); + ret = SPI_execute(query.data, false, 1); + unlink(filename.data); + LWLockRelease(data->relaccess_file_lock); + SPI_finish(); + if (ret < 0) { + elog(ERROR, "SPI execute failure - returned %d", ret); + } +} + +static void update_relname_cache(Oid relid, char *relname) { + bool found; + relnameCacheEntry *relname_entry = (relnameCacheEntry *)hash_search( + relname_cache, &relid, HASH_ENTER, &found); + if (!found) { + relname_entry->relid = relid; + if (!relname) { + strlcpy(relname_entry->relname, get_rel_name(relid), + sizeof(relname_entry->relname)); + } else { + strlcpy(relname_entry->relname, relname, sizeof(relname_entry->relname)); + } + } else { + /** + * NOTE: as we don't handle the 'else' clause here, there will be cases when + * we write outdated table names, like below: + * BEGIN; + * INSERT INTO tbl VALUES (1); + * ALTER TABLE tbl RENAME TO new_tbl; + * SELECT * FROM new_tbl; + * COMMIT; + * In this case both INSERT and SELECT stmts would be counted with the + * old'tbl' name, as we don't update our cache for already known relids in + * the same transaction. This is a deliberate decision for performance + * reasons. + */ + } +} + +static void memorize_local_access_entry(Oid relid, AclMode perms) { + bool found; + localAccessKey key; + key.stmt_cnt = stmt_counter; + key.relid = relid; + localAccessEntry *entry = (localAccessEntry *)hash_search( + local_access_entries, &key, HASH_ENTER, &found); + if (!found) { + entry->last_read = entry->last_write = InvalidOid; + entry->perms = perms; + entry->last_read = 0; + entry->last_write = 0; + } else { + entry->perms |= perms; + } + TimestampTz curts = GetCurrentTimestamp(); + if (is_read(perms)) { + entry->last_reader_id = GetUserId(); + entry->last_read = curts; + } + if (is_write(perms)) { + entry->last_writer_id = GetUserId(); + entry->last_write = curts; + } +} + +static void relaccess_executor_end_hook(QueryDesc *query_desc) { + if (prev_ExecutorEnd_hook) { + prev_ExecutorEnd_hook(query_desc); + } else { + standard_ExecutorEnd(query_desc); + } + // Unfortunately, we cannot safely rely on gp_command_counter as + // it is being incremented more than once for many statements. + // So we have to maintain our own statement counter. + stmt_counter++; +} + +static StringInfoData get_dump_filename(Oid dbid) { + StringInfoData filename; + initStringInfoOfSize(&filename, 256); + appendStringInfo(&filename, "%s/relaccess_stats_dump_%d.csv", + PGSTAT_STAT_PERMANENT_DIRECTORY, dbid); + return filename; +} + +static void relaccess_drop_hook(ObjectAccessType access, Oid classId, + Oid objectId, int subId, void *arg) { + if (prev_object_access_hook) { + prev_object_access_hook(access, classId, objectId, subId, arg); + } + // we don't want shared memory and .csv files hanging around forever + // for databases that we've dropped. + // This function cleans up both files and shmem + if (classId == DatabaseRelationId && access == OAT_DROP) { + LWLockAcquire(data->relaccess_ht_lock, LW_EXCLUSIVE); + HASH_SEQ_STATUS hash_seq; + relaccessEntry *entry; + hash_seq_init(&hash_seq, relaccesses); + while ((entry = hash_seq_search(&hash_seq)) != NULL) { + if (entry->key.dbid == objectId) { + bool found; + hash_search(relaccesses, &entry->key, HASH_REMOVE, &found); + had_ht_overflow = false; + } + } + LWLockRelease(data->relaccess_ht_lock); + LWLockAcquire(data->relaccess_file_lock, LW_EXCLUSIVE); + StringInfoData filename = get_dump_filename(objectId); + unlink(filename.data); + pfree(filename.data); + LWLockRelease(data->relaccess_file_lock); + } +} diff --git a/gpcontrib/gp_relaccess_stats/test/expected/gp_relaccess_stats.out b/gpcontrib/gp_relaccess_stats/test/expected/gp_relaccess_stats.out new file mode 100644 index 00000000000..566f313c915 --- /dev/null +++ b/gpcontrib/gp_relaccess_stats/test/expected/gp_relaccess_stats.out @@ -0,0 +1,406 @@ + GP_IGNORE: formatted by atmsort.pm +CREATE EXTENSION gp_relaccess_stats; +-- get rid of NOTICEs +SET client_min_messages TO WARNING; +SET search_path TO relaccess; +DROP TABLE IF EXISTS tbl1 CASCADE; +DROP TABLE IF EXISTS tbl2 CASCADE; +DROP TABLE IF EXISTS tbl3 CASCADE; +DROP TABLE IF EXISTS tbl4 CASCADE; +DROP TABLE IF EXISTS new_tbl1 CASCADE; +DROP TABLE IF EXISTS p3_sales CASCADE; +DROP TABLE IF EXISTS public.last_usr_checks CASCADE; +DROP USER IF EXISTS select_usr; +DROP USER IF EXISTS update_usr; +DROP USER IF EXISTS insert_usr; +DROP USER IF EXISTS delete_usr; +DROP USER IF EXISTS truncate_usr; +-- make sure tracking is ON +SET gp_relaccess_stats.enabled TO 'on'; +SELECT relaccess_stats_init(); + relaccess_stats_init +---------------------- + +(1 row) + +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +TRUNCATE relaccess_stats; +-- test simple actions one by one in separate transactions +CREATE TABLE tbl1 (a INTEGER); +INSERT INTO tbl1 VALUES(1); +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1'; + n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +------------------+------------------+------------------+------------------+-------------------- + 0 | 1 | 0 | 0 | 0 +(1 row) + +SELECT * FROM tbl1; + a +--- + 1 +(1 row) + +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1'; + n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +------------------+------------------+------------------+------------------+-------------------- + 1 | 1 | 0 | 0 | 0 +(1 row) + +UPDATE tbl1 SET a = -a; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1'; + n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +------------------+------------------+------------------+------------------+-------------------- + 2 | 1 | 1 | 0 | 0 +(1 row) + +DELETE FROM tbl1 WHERE a < 0; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1'; + n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +------------------+------------------+------------------+------------------+-------------------- + 3 | 1 | 1 | 1 | 0 +(1 row) + +TRUNCATE tbl1; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1'; + n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +------------------+------------------+------------------+------------------+-------------------- + 3 | 1 | 1 | 1 | 1 +(1 row) + +-- verify that rename table works +ALTER TABLE tbl1 RENAME TO new_tbl1; +INSERT INTO new_tbl1 VALUES(1); +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relname = 'tbl1'; + n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +------------------+------------------+------------------+------------------+-------------------- +(0 rows) + +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'new_tbl1'::regclass::oid AND relname = 'new_tbl1'; + n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +------------------+------------------+------------------+------------------+-------------------- + 3 | 2 | 1 | 1 | 1 +(1 row) + +TRUNCATE relaccess_stats; +-- multitable truncate +CREATE TABLE tbl1 (a integer); +CREATE TABLE tbl2 (a integer); +TRUNCATE tbl1, tbl2; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT relname, n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats + WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1' OR relid = 'tbl2'::regclass::oid AND relname = 'tbl2' ORDER BY relname; + relname | n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +---------+------------------+------------------+------------------+------------------+-------------------- + tbl1 | 0 | 0 | 0 | 0 | 1 + tbl2 | 0 | 0 | 0 | 0 | 1 +(2 rows) + +TRUNCATE relaccess_stats; +-- test a more complicated statement +CREATE TABLE tbl3 (a integer); +CREATE TABLE tbl4 (a integer); +BEGIN; +-- should give +1 insert for tbl1 and +1 select for other tables +INSERT INTO tbl1 SELECT * FROM tbl2 UNION SELECT * FROM tbl3 UNION SELECT * FROM tbl4; +-- nothing in there before we commit +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT COUNT(*) FROM relaccess_stats WHERE relname LIKE ('tbl_'); + count +------- + 0 +(1 row) + +COMMIT; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT relname, n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries + FROM relaccess_stats WHERE relname LIKE ('tbl_') AND relname::regclass::oid = relid ORDER BY relname; + relname | n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +---------+------------------+------------------+------------------+------------------+-------------------- + tbl1 | 0 | 1 | 0 | 0 | 0 + tbl2 | 1 | 0 | 0 | 0 | 0 + tbl3 | 1 | 0 | 0 | 0 | 0 + tbl4 | 1 | 0 | 0 | 0 | 0 +(4 rows) + +TRUNCATE relaccess_stats; +-- test views +CREATE VIEW v1_2_3 AS (SELECT * FROM tbl2 UNION SELECT * FROM tbl3 UNION SELECT * FROM tbl4); +INSERT INTO tbl1 SELECT * FROM v1_2_3; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT relname, n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries + FROM relaccess_stats WHERE relname = 'v1_2_3' OR relname LIKE ('tbl_') ORDER BY relname; + relname | n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +---------+------------------+------------------+------------------+------------------+-------------------- + tbl1 | 0 | 1 | 0 | 0 | 0 + tbl2 | 1 | 0 | 0 | 0 | 0 + tbl3 | 1 | 0 | 0 | 0 | 0 + tbl4 | 1 | 0 | 0 | 0 | 0 + v1_2_3 | 1 | 0 | 0 | 0 | 0 +(5 rows) + +TRUNCATE relaccess_stats; +-- test timestamps difference +BEGIN; +INSERT INTO tbl1 VALUES (1); +SELECT pg_sleep(1); + pg_sleep +---------- + +(1 row) + +SELECT COUNT(*) FROM tbl1; + count +------- + 1 +(1 row) + +COMMIT; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT EXTRACT(EPOCH FROM (last_read - last_write)) >= 1 FROM relaccess_stats WHERE relname = 'tbl1' AND relid = 'tbl1'::regclass::oid; + ?column? +---------- + t +(1 row) + +TRUNCATE relaccess_stats; +-- test nested partitions lookup +BEGIN; +CREATE TABLE p3_sales (id int, year int, month int, day int, + region text) +DISTRIBUTED BY (id) +PARTITION BY RANGE (year) + SUBPARTITION BY RANGE (month) + SUBPARTITION TEMPLATE ( + START (1) END (13) EVERY (1), + DEFAULT SUBPARTITION other_months ) + SUBPARTITION BY LIST (region) + SUBPARTITION TEMPLATE ( + SUBPARTITION usa VALUES ('usa'), + SUBPARTITION europe VALUES ('europe'), + SUBPARTITION asia VALUES ('asia'), + DEFAULT SUBPARTITION other_regions ) +( START (2002) END (2012) EVERY (1), + DEFAULT PARTITION outlying_years ); +-- 3 inserts into p3_sales root table +INSERT INTO p3_sales SELECT i, i%43+1980, i%12, i%25, 'asia' FROM generate_series(1, 100)i; +INSERT INTO p3_sales SELECT i, i%43+1980, i%12, i%25, 'europe' FROM generate_series(1, 100)i; +INSERT INTO p3_sales SELECT i, i%43+1980, i%12, i%25, 'usa' FROM generate_series(1, 100)i; +-- insert and select to/from specific leaf level partition +INSERT INTO p3_sales_1_prt_11_2_prt_12_3_prt_usa SELECT * FROM p3_sales_1_prt_11_2_prt_12_3_prt_usa; +COMMIT; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT relname, n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries +FROM relaccess_stats WHERE relname LIKE 'p3_sales%' ORDER BY relname; + relname | n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +--------------------------------------+------------------+------------------+------------------+------------------+-------------------- + p3_sales | 0 | 3 | 0 | 0 | 0 + p3_sales_1_prt_11_2_prt_12_3_prt_usa | 1 | 1 | 0 | 0 | 0 +(2 rows) + +SELECT relname, n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries +FROM relaccess_stats_root_tables_aggregated WHERE relname LIKE 'p3_sales%' ORDER BY relname; + relname | n_select_queries | n_insert_queries | n_update_queries | n_delete_queries | n_truncate_queries +----------+------------------+------------------+------------------+------------------+-------------------- + p3_sales | 1 | 4 | 0 | 0 | 0 +(1 row) + +-- test last_reader and last_writer +CREATE USER select_usr; +CREATE USER update_usr; +CREATE USER insert_usr; +CREATE USER delete_usr; +CREATE USER truncate_usr; +CREATE TABLE public.last_usr_checks(a integer); +GRANT ALL ON TABLE public.last_usr_checks TO select_usr, update_usr, insert_usr, delete_usr, truncate_usr; +SET ROLE select_usr; +SELECT COUNT(*) FROM public.last_usr_checks; + count +------- + 0 +(1 row) + +RESET ROLE; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT (SELECT last_reader_id FROM relaccess_stats WHERE RELNAME = 'last_usr_checks') = (SELECT oid FROM pg_roles WHERE rolname = 'select_usr'); + ?column? +---------- + t +(1 row) + +SET ROLE insert_usr; +INSERT INTO public.last_usr_checks VALUES (-1), (0), (1); +RESET ROLE; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT (SELECT last_writer_id FROM relaccess_stats WHERE RELNAME = 'last_usr_checks') = (SELECT oid FROM pg_roles WHERE rolname = 'insert_usr'); + ?column? +---------- + t +(1 row) + +SET ROLE update_usr; +UPDATE public.last_usr_checks SET a = a*10 WHERE a < 0; +RESET ROLE; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT (SELECT last_writer_id FROM relaccess_stats WHERE RELNAME = 'last_usr_checks') = (SELECT oid FROM pg_roles WHERE rolname = 'update_usr'); + ?column? +---------- + t +(1 row) + +SET ROLE delete_usr; +DELETE FROM public.last_usr_checks WHERE a >= 0; +RESET ROLE; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT (SELECT last_writer_id FROM relaccess_stats WHERE RELNAME = 'last_usr_checks') = (SELECT oid FROM pg_roles WHERE rolname = 'delete_usr'); + ?column? +---------- + t +(1 row) + +SET ROLE truncate_usr; +TRUNCATE public.last_usr_checks; +RESET ROLE; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT (SELECT last_writer_id FROM relaccess_stats WHERE RELNAME = 'last_usr_checks') = (SELECT oid FROM pg_roles WHERE rolname = 'truncate_usr'); + ?column? +---------- + t +(1 row) + +RESET ROLE; +-- make sure we can turn it OFF +SET gp_relaccess_stats.enabled TO 'off'; +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +TRUNCATE relaccess_stats; +SELECT * FROM tbl1; + a +--- + 1 +(1 row) + +SELECT relaccess_stats_update(); + relaccess_stats_update +------------------------ + +(1 row) + +SELECT count(*) FROM relaccess_stats; + count +------- + 0 +(1 row) + +RESET gp_relaccess_stats.enabled; +DROP TABLE tbl1 CASCADE; +DROP TABLE tbl2 CASCADE; +DROP TABLE tbl3 CASCADE; +DROP TABLE tbl4 CASCADE; +DROP TABLE new_tbl1 CASCADE; +DROP TABLE p3_sales CASCADE; +DROP TABLE public.last_usr_checks CASCADE; +DROP USER select_usr; +DROP USER update_usr; +DROP USER insert_usr; +DROP USER delete_usr; +DROP USER truncate_usr; diff --git a/gpcontrib/gp_relaccess_stats/test/sql/gp_relaccess_stats.sql b/gpcontrib/gp_relaccess_stats/test/sql/gp_relaccess_stats.sql new file mode 100644 index 00000000000..cb96d8eb9fa --- /dev/null +++ b/gpcontrib/gp_relaccess_stats/test/sql/gp_relaccess_stats.sql @@ -0,0 +1,186 @@ +CREATE EXTENSION gp_relaccess_stats; + +-- get rid of NOTICEs +SET client_min_messages TO WARNING; +SET search_path TO relaccess; +DROP TABLE IF EXISTS tbl1 CASCADE; +DROP TABLE IF EXISTS tbl2 CASCADE; +DROP TABLE IF EXISTS tbl3 CASCADE; +DROP TABLE IF EXISTS tbl4 CASCADE; +DROP TABLE IF EXISTS new_tbl1 CASCADE; +DROP TABLE IF EXISTS p3_sales CASCADE; +DROP TABLE IF EXISTS public.last_usr_checks CASCADE; +DROP USER IF EXISTS select_usr; +DROP USER IF EXISTS update_usr; +DROP USER IF EXISTS insert_usr; +DROP USER IF EXISTS delete_usr; +DROP USER IF EXISTS truncate_usr; + +-- make sure tracking is ON +SET gp_relaccess_stats.enabled TO 'on'; +SELECT relaccess_stats_init(); +SELECT relaccess_stats_update(); +TRUNCATE relaccess_stats; + +-- test simple actions one by one in separate transactions +CREATE TABLE tbl1 (a INTEGER); + +INSERT INTO tbl1 VALUES(1); +SELECT relaccess_stats_update(); +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1'; + +SELECT * FROM tbl1; +SELECT relaccess_stats_update(); +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1'; + +UPDATE tbl1 SET a = -a; +SELECT relaccess_stats_update(); +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1'; + +DELETE FROM tbl1 WHERE a < 0; +SELECT relaccess_stats_update(); +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1'; + +TRUNCATE tbl1; +SELECT relaccess_stats_update(); +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1'; + +-- verify that rename table works +ALTER TABLE tbl1 RENAME TO new_tbl1; +INSERT INTO new_tbl1 VALUES(1); +SELECT relaccess_stats_update(); +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relname = 'tbl1'; +SELECT n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats WHERE relid = 'new_tbl1'::regclass::oid AND relname = 'new_tbl1'; + +TRUNCATE relaccess_stats; +-- multitable truncate +CREATE TABLE tbl1 (a integer); +CREATE TABLE tbl2 (a integer); +TRUNCATE tbl1, tbl2; +SELECT relaccess_stats_update(); +SELECT relname, n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries FROM relaccess_stats + WHERE relid = 'tbl1'::regclass::oid AND relname = 'tbl1' OR relid = 'tbl2'::regclass::oid AND relname = 'tbl2' ORDER BY relname; + +TRUNCATE relaccess_stats; +-- test a more complicated statement +CREATE TABLE tbl3 (a integer); +CREATE TABLE tbl4 (a integer); + +BEGIN; +-- should give +1 insert for tbl1 and +1 select for other tables +INSERT INTO tbl1 SELECT * FROM tbl2 UNION SELECT * FROM tbl3 UNION SELECT * FROM tbl4; +-- nothing in there before we commit +SELECT relaccess_stats_update(); +SELECT COUNT(*) FROM relaccess_stats WHERE relname LIKE ('tbl_'); +COMMIT; +SELECT relaccess_stats_update(); +SELECT relname, n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries + FROM relaccess_stats WHERE relname LIKE ('tbl_') AND relname::regclass::oid = relid ORDER BY relname; + +TRUNCATE relaccess_stats; +-- test views +CREATE VIEW v1_2_3 AS (SELECT * FROM tbl2 UNION SELECT * FROM tbl3 UNION SELECT * FROM tbl4); +INSERT INTO tbl1 SELECT * FROM v1_2_3; +SELECT relaccess_stats_update(); +SELECT relname, n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries + FROM relaccess_stats WHERE relname = 'v1_2_3' OR relname LIKE ('tbl_') ORDER BY relname; + +TRUNCATE relaccess_stats; +-- test timestamps difference +BEGIN; +INSERT INTO tbl1 VALUES (1); +SELECT pg_sleep(1); +SELECT COUNT(*) FROM tbl1; +COMMIT; +SELECT relaccess_stats_update(); +SELECT EXTRACT(EPOCH FROM (last_read - last_write)) >= 1 FROM relaccess_stats WHERE relname = 'tbl1' AND relid = 'tbl1'::regclass::oid; +TRUNCATE relaccess_stats; + +-- test nested partitions lookup +BEGIN; +CREATE TABLE p3_sales (id int, year int, month int, day int, + region text) +DISTRIBUTED BY (id) +PARTITION BY RANGE (year) + SUBPARTITION BY RANGE (month) + SUBPARTITION TEMPLATE ( + START (1) END (13) EVERY (1), + DEFAULT SUBPARTITION other_months ) + SUBPARTITION BY LIST (region) + SUBPARTITION TEMPLATE ( + SUBPARTITION usa VALUES ('usa'), + SUBPARTITION europe VALUES ('europe'), + SUBPARTITION asia VALUES ('asia'), + DEFAULT SUBPARTITION other_regions ) +( START (2002) END (2012) EVERY (1), + DEFAULT PARTITION outlying_years ); +-- 3 inserts into p3_sales root table +INSERT INTO p3_sales SELECT i, i%43+1980, i%12, i%25, 'asia' FROM generate_series(1, 100)i; +INSERT INTO p3_sales SELECT i, i%43+1980, i%12, i%25, 'europe' FROM generate_series(1, 100)i; +INSERT INTO p3_sales SELECT i, i%43+1980, i%12, i%25, 'usa' FROM generate_series(1, 100)i; +-- insert and select to/from specific leaf level partition +INSERT INTO p3_sales_1_prt_11_2_prt_12_3_prt_usa SELECT * FROM p3_sales_1_prt_11_2_prt_12_3_prt_usa; +COMMIT; +SELECT relaccess_stats_update(); +SELECT relname, n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries +FROM relaccess_stats WHERE relname LIKE 'p3_sales%' ORDER BY relname; +SELECT relname, n_select_queries, n_insert_queries, n_update_queries, n_delete_queries, n_truncate_queries +FROM relaccess_stats_root_tables_aggregated WHERE relname LIKE 'p3_sales%' ORDER BY relname; + +-- test last_reader and last_writer +CREATE USER select_usr; +CREATE USER update_usr; +CREATE USER insert_usr; +CREATE USER delete_usr; +CREATE USER truncate_usr; +CREATE TABLE public.last_usr_checks(a integer); +GRANT ALL ON TABLE public.last_usr_checks TO select_usr, update_usr, insert_usr, delete_usr, truncate_usr; +SET ROLE select_usr; +SELECT COUNT(*) FROM public.last_usr_checks; +RESET ROLE; +SELECT relaccess_stats_update(); +SELECT (SELECT last_reader_id FROM relaccess_stats WHERE RELNAME = 'last_usr_checks') = (SELECT oid FROM pg_roles WHERE rolname = 'select_usr'); +SET ROLE insert_usr; +INSERT INTO public.last_usr_checks VALUES (-1), (0), (1); +RESET ROLE; +SELECT relaccess_stats_update(); +SELECT (SELECT last_writer_id FROM relaccess_stats WHERE RELNAME = 'last_usr_checks') = (SELECT oid FROM pg_roles WHERE rolname = 'insert_usr'); +SET ROLE update_usr; +UPDATE public.last_usr_checks SET a = a*10 WHERE a < 0; +RESET ROLE; +SELECT relaccess_stats_update(); +SELECT (SELECT last_writer_id FROM relaccess_stats WHERE RELNAME = 'last_usr_checks') = (SELECT oid FROM pg_roles WHERE rolname = 'update_usr'); +SET ROLE delete_usr; +DELETE FROM public.last_usr_checks WHERE a >= 0; +RESET ROLE; +SELECT relaccess_stats_update(); +SELECT (SELECT last_writer_id FROM relaccess_stats WHERE RELNAME = 'last_usr_checks') = (SELECT oid FROM pg_roles WHERE rolname = 'delete_usr'); +SET ROLE truncate_usr; +TRUNCATE public.last_usr_checks; +RESET ROLE; +SELECT relaccess_stats_update(); +SELECT (SELECT last_writer_id FROM relaccess_stats WHERE RELNAME = 'last_usr_checks') = (SELECT oid FROM pg_roles WHERE rolname = 'truncate_usr'); +RESET ROLE; + +-- make sure we can turn it OFF +SET gp_relaccess_stats.enabled TO 'off'; +SELECT relaccess_stats_update(); +TRUNCATE relaccess_stats; +SELECT * FROM tbl1; +SELECT relaccess_stats_update(); +SELECT count(*) FROM relaccess_stats; +RESET gp_relaccess_stats.enabled; + +DROP TABLE tbl1 CASCADE; +DROP TABLE tbl2 CASCADE; +DROP TABLE tbl3 CASCADE; +DROP TABLE tbl4 CASCADE; +DROP TABLE new_tbl1 CASCADE; +DROP TABLE p3_sales CASCADE; +DROP TABLE public.last_usr_checks CASCADE; +DROP USER select_usr; +DROP USER update_usr; +DROP USER insert_usr; +DROP USER delete_usr; +DROP USER truncate_usr; + diff --git a/pom.xml b/pom.xml index cb2c25c20c4..ba89d1b19ab 100644 --- a/pom.xml +++ b/pom.xml @@ -1281,6 +1281,10 @@ code or new licensing patterns. gpcontrib/gp_relsizes_stats/gp_relsizes_stats.control gpcontrib/gp_relsizes_stats/test/postgresql.conf.add + gpcontrib/gp_relaccess_stats/Makefile + gpcontrib/gp_relaccess_stats/src/gp_relaccess_stats.c + gpcontrib/gp_relaccess_stats/gp_relaccess_stats.control + gpcontrib/reject_partition_fullscan/Makefile gpcontrib/reject_partition_fullscan/reject_partition_fullscan.control