Skip to content

Commit 20bbe92

Browse files
authored
Merge pull request #4920 from wazir-ahmed/purge-stats-query-digest
feat(stats): Purge query digest based on last seen
2 parents 3020ee4 + c0d39eb commit 20bbe92

14 files changed

Lines changed: 833 additions & 109 deletions

include/gen_utils.h

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -516,6 +516,9 @@ bool mywildcmp(const char *p, const char *str);
516516
std::string trim(const std::string& s);
517517
char* escape_string_single_quotes_and_backslashes(char* input, bool free_it);
518518
const char* escape_string_backslash_spaces(const char* input);
519+
time_t monotonic_time_to_realtime(time_t mt);
520+
time_t realtime_to_monotonic_time(time_t rt);
521+
519522
std::string strip_schema_from_query(const char* query, const char* schema,
520523
const std::vector<std::string>& tables = {}, bool ansi_quotes = false);
521524
/**

include/query_processor.h

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,7 @@ class QP_query_digest_stats {
8181
unsigned long long t, unsigned long long n, unsigned long long ra, unsigned long long rs,
8282
unsigned long long cnt = 1
8383
);
84+
void merge(const QP_query_digest_stats *other);
8485
~QP_query_digest_stats();
8586
char *get_digest_text(const umap_query_digest_text *digest_text_umap) const;
8687
char **get_row(umap_query_digest_text *digest_text_umap, query_digest_stats_pointers_t *qdsp);
@@ -345,7 +346,7 @@ class Query_Processor {
345346
uint32_t max_window
346347
);
347348
void get_query_digests_reset(umap_query_digest* uqd, umap_query_digest_text* uqdt);
348-
unsigned long long purge_query_digests(bool async_purge, bool parallel, char** msg);
349+
unsigned long long purge_query_digests(bool async_purge, bool parallel, time_t last_seen = 0);
349350

350351
void save_query_rules(SQLite3_result* resultset);
351352

@@ -457,8 +458,8 @@ class Query_Processor {
457458
DEFINE_HAS_METHOD_STRUCT(query_parser_first_comment_extended);
458459
DEFINE_HAS_METHOD_STRUCT(process_query_extended);
459460

460-
unsigned long long purge_query_digests_async(char** msg);
461-
unsigned long long purge_query_digests_sync(bool parallel);
461+
unsigned long long purge_query_digests_async(time_t last_seen = 0);
462+
unsigned long long purge_query_digests_sync(bool parallel, time_t last_seen = 0);
462463

463464
/**
464465
* @brief Searches for a matching rule in the supplied map, returning the destination hostgroup.

lib/Admin_Bootstrap.cpp

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -143,9 +143,6 @@ extern struct MHD_Daemon *Admin_HTTP_Server;
143143

144144
extern ProxySQL_Statistics *GloProxyStats;
145145

146-
template<enum SERVER_TYPE>
147-
int ProxySQL_Test___PurgeDigestTable(bool async_purge, bool parallel, char **msg);
148-
149146
extern char *ssl_key_fp;
150147
extern char *ssl_cert_fp;
151148
extern char *ssl_ca_fp;

lib/Admin_FlushVariables.cpp

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -118,9 +118,6 @@ extern struct MHD_Daemon *Admin_HTTP_Server;
118118

119119
extern ProxySQL_Statistics *GloProxyStats;
120120

121-
template<enum SERVER_TYPE>
122-
int ProxySQL_Test___PurgeDigestTable(bool async_purge, bool parallel, char **msg);
123-
124121
extern char *ssl_key_fp;
125122
extern char *ssl_cert_fp;
126123
extern char *ssl_ca_fp;

lib/Admin_Handler.cpp

Lines changed: 89 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -120,9 +120,6 @@ extern struct MHD_Daemon *Admin_HTTP_Server;
120120

121121
extern ProxySQL_Statistics *GloProxyStats;
122122

123-
template<enum SERVER_TYPE>
124-
int ProxySQL_Test___PurgeDigestTable(bool async_purge, bool parallel, char **msg);
125-
126123
extern char *ssl_key_fp;
127124
extern char *ssl_cert_fp;
128125
extern char *ssl_ca_fp;
@@ -327,6 +324,17 @@ const std::vector<std::string> LOAD_COREDUMP_FROM_MEMORY = {
327324
"LOAD COREDUMP TO RUNTIME" ,
328325
"LOAD COREDUMP TO RUN" };
329326

327+
const std::vector<std::string> CMD_PREFIX_PURGE_QUERY_DIGESTS = {
328+
"PURGE TABLE stats.stats_mysql_query_digest TO ",
329+
"PURGE TABLE stats_mysql_query_digest TO ",
330+
"PURGE stats.stats_mysql_query_digest TO ",
331+
"PURGE stats_mysql_query_digest TO ",
332+
"PURGE TABLE stats.stats_pgsql_query_digest TO ",
333+
"PURGE TABLE stats_pgsql_query_digest TO ",
334+
"PURGE stats.stats_pgsql_query_digest TO ",
335+
"PURGE stats_pgsql_query_digest TO ",
336+
};
337+
330338
extern unordered_map<string,std::tuple<string, vector<string>, vector<string>>> load_save_disk_commands;
331339

332340
// Helper function: Extract pattern from c.relname OPERATOR or LIKE clause
@@ -518,6 +526,17 @@ bool is_admin_command_or_alias(const std::vector<std::string>& cmds, char *query
518526
return false;
519527
}
520528

529+
const char * match_command_prefix(const std::vector<std::string>& cmd_prefix, char *query, int query_len) {
530+
for (auto &prefix : cmd_prefix) {
531+
if ((unsigned int) query_len >= prefix.length()
532+
&& !strncasecmp(prefix.c_str(), query, prefix.length()))
533+
{
534+
return prefix.c_str();
535+
}
536+
}
537+
538+
return nullptr;
539+
}
521540

522541

523542
template <typename S>
@@ -553,6 +572,38 @@ bool FlushCommandWrapper(S* sess, const string& modname, char *query_no_space, i
553572
return false;
554573
}
555574

575+
std::tuple<bool, enum SERVER_TYPE, time_t> parse_command_purge_query_digests(char *query, int query_len) {
576+
bool match = false;
577+
enum SERVER_TYPE server_type = SERVER_TYPE_MYSQL;
578+
time_t last_seen = 0;
579+
580+
const char *prefix = match_command_prefix(CMD_PREFIX_PURGE_QUERY_DIGESTS, query, query_len);
581+
if (prefix) {
582+
match = true;
583+
584+
if (strstr(prefix, "_pgsql_") != nullptr) {
585+
server_type = SERVER_TYPE_PGSQL;
586+
}
587+
588+
// parse timestamp
589+
mf_unique_ptr<char> ts_str(strdup(query + strlen(prefix)));
590+
char *ts_end = nullptr;
591+
long long ts = strtoll(trim_spaces_in_place(ts_str.get()), &ts_end, 10);
592+
593+
// ts_str should only contain digits and respresent a valid timestamp
594+
if ((*ts_end == 0) && (ts > 0)) {
595+
last_seen = realtime_to_monotonic_time(ts);
596+
if (last_seen <= 0) {
597+
// valid timestamp, but older than the monotonic clock epoch (system
598+
// boot): no entry can match, make the command a successful no-op
599+
last_seen = 1;
600+
}
601+
}
602+
}
603+
604+
return std::make_tuple(match, server_type, last_seen);
605+
}
606+
556607
template <typename S>
557608
bool admin_handler_command_kill_pgsql_connection(uint32_t session_thd_id, S* sess, ProxySQL_Admin* pa) {
558609
proxy_debug(PROXY_DEBUG_ADMIN, 4, "Trying to kill pgsql session %u\n", session_thd_id);
@@ -3442,10 +3493,8 @@ void admin_session_handler(S* sess, void *_pa, PtrSize_t *pkt) {
34423493
SPA->admindb->execute("DELETE FROM stats.stats_mysql_query_digest_reset");
34433494
SPA->vacuum_stats(true);
34443495
// purge the digest map, asynchronously, in single thread
3445-
char *msg = NULL;
3446-
int r1 = ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_MYSQL>(true, false, &msg);
3447-
SPA->send_ok_msg_to_client(sess, msg, r1, query_no_space);
3448-
free(msg);
3496+
int r1 = GloMyQPro->purge_query_digests(true, false);
3497+
SPA->send_ok_msg_to_client(sess, NULL, r1, query_no_space);
34493498
run_query=false;
34503499
goto __run_query;
34513500
}
@@ -3478,16 +3527,45 @@ void admin_session_handler(S* sess, void *_pa, PtrSize_t *pkt) {
34783527
SPA->admindb->execute("DELETE FROM stats.stats_pgsql_query_digest_reset");
34793528
SPA->vacuum_stats(true);
34803529
// purge the digest map, asynchronously, in single thread
3481-
char* msg = NULL;
3482-
int r1 = ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_PGSQL>(true, false, &msg);
3483-
SPA->send_ok_msg_to_client(sess, msg, r1, query_no_space);
3484-
free(msg);
3530+
int r1 = GloPgQPro->purge_query_digests(true, false);
3531+
SPA->send_ok_msg_to_client(sess, NULL, r1, query_no_space);
34853532
run_query = false;
34863533
goto __run_query;
34873534
}
34883535
}
34893536
}
34903537
}
3538+
3539+
// handles 'PURGE stats_mysql_query_digest TO <value>'.
3540+
// any entry in stats_mysql_query_digest where last_seen is less than <value> will be deleted.
3541+
if (!strncasecmp("PURGE ", query_no_space, strlen("PURGE "))
3542+
&& sess->session_type == PROXYSQL_SESSION_ADMIN
3543+
) {
3544+
auto result = parse_command_purge_query_digests(query_no_space, query_no_space_length);
3545+
bool match = std::get<0>(result);
3546+
3547+
if (match == true) {
3548+
int ret = 0;
3549+
enum SERVER_TYPE type = std::get<1>(result);
3550+
time_t last_seen = std::get<2>(result);
3551+
3552+
if (last_seen > 0) {
3553+
if (type == SERVER_TYPE_MYSQL) {
3554+
ret = GloMyQPro->purge_query_digests(true, false, last_seen);
3555+
} else if (type == SERVER_TYPE_PGSQL) {
3556+
ret = GloPgQPro->purge_query_digests(true, false, last_seen);
3557+
}
3558+
3559+
pa->send_ok_msg_to_client(sess, NULL, ret, query_no_space);
3560+
} else {
3561+
pa->send_error_msg_to_client(sess, "Invalid timestamp");
3562+
}
3563+
3564+
run_query = false;
3565+
goto __run_query;
3566+
}
3567+
}
3568+
34913569
#ifdef DEBUG
34923570
/**
34933571
* @brief Handles the 'PROXYSQL_SIMULATOR' command. Performing the operation specified in the payload

lib/ProxySQL_Admin.cpp

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -123,9 +123,6 @@ struct MHD_Daemon *Admin_HTTP_Server;
123123

124124
extern ProxySQL_Statistics *GloProxyStats;
125125

126-
template<enum SERVER_TYPE>
127-
int ProxySQL_Test___PurgeDigestTable(bool async_purge, bool parallel, char **msg);
128-
129126
extern char *ssl_key_fp;
130127
extern char *ssl_cert_fp;
131128
extern char *ssl_ca_fp;

lib/ProxySQL_Admin_Tests.cpp

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
#include <ctime>
12
#include <iostream> // std::cout
23
#include <sstream> // std::stringstream
34
#include <fstream>
@@ -56,12 +57,12 @@ bool ProxySQL_Test___Refresh_MySQL_Variables(unsigned int cnt) {
5657
}
5758

5859
template <enum SERVER_TYPE ST>
59-
int ProxySQL_Test___PurgeDigestTable(bool async_purge, bool parallel, char **msg) {
60+
int ProxySQL_Test___PurgeDigestTable(bool async_purge, bool parallel, time_t last_seen) {
6061
int r = 0;
6162
if constexpr (ST == SERVER_TYPE_MYSQL) {
62-
r = GloMyQPro->purge_query_digests(async_purge, parallel, msg);
63+
r = GloMyQPro->purge_query_digests(async_purge, parallel, last_seen);
6364
} else if constexpr (ST == SERVER_TYPE_PGSQL) {
64-
r = GloPgQPro->purge_query_digests(async_purge, parallel, msg);
65+
r = GloPgQPro->purge_query_digests(async_purge, parallel, last_seen);
6566
}
6667
return r;
6768
}
@@ -147,5 +148,5 @@ int ProxySQL_Test___GenerateRandomQueryInDigestTable(int n) {
147148
return n*1000;
148149
}
149150

150-
template int ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_MYSQL>(bool async_purge, bool parallel, char** msg);
151-
template int ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_PGSQL>(bool async_purge, bool parallel, char** msg);
151+
template int ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_MYSQL>(bool async_purge, bool parallel, time_t last_seen);
152+
template int ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_PGSQL>(bool async_purge, bool parallel, time_t last_seen);

lib/ProxySQL_Admin_Tests2.cpp

Lines changed: 5 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,7 @@ static void init_rand_del() {
4646
int ProxySQL_Test___GetDigestTable(bool reset, bool use_swap);
4747
bool ProxySQL_Test___Refresh_MySQL_Variables(unsigned int cnt);
4848
template<enum SERVER_TYPE>
49-
int ProxySQL_Test___PurgeDigestTable(bool async_purge, bool parallel, char **msg);
49+
int ProxySQL_Test___PurgeDigestTable(bool async_purge, bool parallel, time_t last_seen);
5050
int ProxySQL_Test___GenerateRandomQueryInDigestTable(int n);
5151

5252
void ProxySQL_Admin::map_test_mysql_firewall_whitelist_rules_cleanup() {
@@ -1179,21 +1179,20 @@ void ProxySQL_Admin::ProxySQL_Test_Handler(ProxySQL_Admin *SPA, S* sess, char *q
11791179
break;
11801180
case 4:
11811181
// purge the digest map, synchronously, in single thread
1182-
r1 = ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_MYSQL>(false, false, NULL);
1182+
r1 = ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_MYSQL>(false, false, test_arg1);
11831183
SPA->send_ok_msg_to_client(sess, NULL, r1, query_no_space);
11841184
run_query=false;
11851185
break;
11861186
case 5:
11871187
// purge the digest map, synchronously, in multiple threads
1188-
r1 = ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_MYSQL>(false, true, NULL);
1188+
r1 = ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_MYSQL>(false, true, test_arg1);
11891189
SPA->send_ok_msg_to_client(sess, NULL, r1, query_no_space);
11901190
run_query=false;
11911191
break;
11921192
case 6:
11931193
// purge the digest map, asynchronously, in single thread
1194-
r1 = ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_MYSQL>(true, false, &msg);
1195-
SPA->send_ok_msg_to_client(sess, msg, r1, query_no_space);
1196-
free(msg);
1194+
r1 = ProxySQL_Test___PurgeDigestTable<SERVER_TYPE_MYSQL>(true, false, test_arg1);
1195+
SPA->send_ok_msg_to_client(sess, NULL, r1, query_no_space);
11971196
run_query=false;
11981197
break;
11991198
case 7:

lib/QP_query_digest_stats.cpp

Lines changed: 23 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
#include "gen_utils.h"
12
#include "query_processor.h"
23

34
// reverse: reverse string s in place
@@ -79,6 +80,26 @@ void QP_query_digest_stats::add_time(
7980
}
8081
last_seen=n;
8182
}
83+
// Merges the counters of 'other' into this entry. Used when reconciling stats
84+
// collected while a purge operation was running (see purge_query_digests_async()).
85+
void QP_query_digest_stats::merge(const QP_query_digest_stats *other) {
86+
count_star += other->count_star;
87+
sum_time += other->sum_time;
88+
rows_affected += other->rows_affected;
89+
rows_sent += other->rows_sent;
90+
if (other->min_time && (min_time == 0 || other->min_time < min_time)) {
91+
min_time = other->min_time;
92+
}
93+
if (other->max_time > max_time) {
94+
max_time = other->max_time;
95+
}
96+
if (other->first_seen && (first_seen == 0 || other->first_seen < first_seen)) {
97+
first_seen = other->first_seen;
98+
}
99+
if (other->last_seen > last_seen) {
100+
last_seen = other->last_seen;
101+
}
102+
}
82103
QP_query_digest_stats::~QP_query_digest_stats() {
83104
if (digest_text) {
84105
free(digest_text);
@@ -151,16 +172,13 @@ char **QP_query_digest_stats::get_row(umap_query_digest_text *digest_text_umap,
151172
my_itoa(qdsp->count_star, count_star);
152173
pta[5]=qdsp->count_star;
153174

154-
time_t __now;
155-
time(&__now);
156-
unsigned long long curtime=monotonic_time();
157175
time_t seen_time;
158-
seen_time= __now - curtime/1000000 + first_seen/1000000;
176+
seen_time=monotonic_time_to_realtime(first_seen);
159177
//sprintf(qdsp->first_seen,"%ld", seen_time);
160178
my_itoa(qdsp->first_seen, seen_time);
161179
pta[6]=qdsp->first_seen;
162180

163-
seen_time= __now - curtime/1000000 + last_seen/1000000;
181+
seen_time=monotonic_time_to_realtime(last_seen);
164182
//sprintf(qdsp->last_seen,"%ld", seen_time);
165183
my_itoa(qdsp->last_seen, seen_time);
166184
pta[7]=qdsp->last_seen;

0 commit comments

Comments
 (0)