Skip to content

Commit 423066f

Browse files
committed
Add interconnect regression tests
1 parent 6c230b3 commit 423066f

11 files changed

Lines changed: 496 additions & 0 deletions

File tree

.github/workflows/build-cloudberry.yml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -298,6 +298,7 @@ jobs:
298298
"contrib/formatter_fixedwidth:installcheck",
299299
"contrib/hstore:installcheck",
300300
"contrib/indexscan:installcheck",
301+
"contrib/interconnect:installcheck",
301302
"contrib/pg_trgm:installcheck",
302303
"contrib/indexscan:installcheck",
303304
"contrib/pgcrypto:installcheck",

contrib/interconnect/Makefile

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,10 @@ include $(top_builddir)/contrib/interconnect/Makefile.interconnect
77
MODULE_big = interconnect
88
PGFILEDESC = "interconnect - inter connection module"
99

10+
EXTENSION = interconnect
11+
EXTENSION_VERSION = 1.0
12+
DATA = interconnect--$(EXTENSION_VERSION).sql
13+
1014
OBJS = \
1115
$(WIN32RES) \
1216
ic_common.o \
@@ -33,6 +37,8 @@ OBJS += proxy/ic_proxy_iobuf.o
3337
SHLIB_LINK += $(filter -luv, $(LIBS))
3438
endif # enable_ic_proxy
3539

40+
REGRESS = interconnect
41+
3642
ifdef USE_PGXS
3743
PG_CONFIG = pg_config
3844
PGXS := $(shell $(PG_CONFIG) --pgxs)

contrib/interconnect/README.md

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -271,3 +271,35 @@ udpifc result:
271271

272272
Notice that: Lower TPS does not mean the protocol is slower, might means that the cpu time taken by the protocol is low. For the udpifc, it satisfies the highest tps required by `cbdb`. at the same time it occupies a lower cpu than other types of interconnect.
273273

274+
# interconnect statistics
275+
276+
This extension provides cumulative interconnect statistics for Cloudberry Database, including queue sizes, buffer usage, retransmits, packet errors, and other UDPIFC‑related metrics.
277+
278+
It exposes three views with statistics at different aggregation levels:
279+
gp_interconnect_stats — total cluster‑wide stats;
280+
gp_interconnect_stats_per_segment — stats per segment;
281+
gp_interconnect_stats_per_host — stats grouped by host.
282+
283+
## How to create the extension
284+
285+
Add interconnect to shared_preload_libraries and restart the cluster.
286+
287+
```
288+
gpconfig -c shared_preload_libraries -v \
289+
"$(psql -At -c \
290+
"SELECT array_to_string( \
291+
array_append( \
292+
string_to_array( \
293+
current_setting('shared_preload_libraries'), \
294+
','), \
295+
'interconnect'), \
296+
',')" \
297+
postgres)"
298+
gpstop -ra
299+
```
300+
301+
Create the extension in your database.
302+
303+
```
304+
CREATE EXTENSION interconnect;
305+
```
Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
1+
-- Capture current interconnect stats as baseline for future comparisons
2+
SELECT * FROM gp_interconnect_stats \gset prev_
3+
-- Verify that all baseline interconnect statistics are >= 0 (no negative values)
4+
SELECT
5+
:prev_total_recv_queue_size >= 0,
6+
:prev_recv_queue_conting_time >= 0,
7+
:prev_total_capacity >= 0,
8+
:prev_capacity_counting_time >= 0,
9+
:prev_total_buffers >= 0,
10+
:prev_buffer_counting_time >= 0,
11+
:prev_retransmits >= 0,
12+
:prev_startup_cached_pkts >= 0,
13+
:prev_mismatches >= 0,
14+
:prev_crs_errors >= 0,
15+
:prev_snd_pkt_num >= 0,
16+
:prev_recv_pkt_num >= 0,
17+
:prev_disordered_pkt_num >= 0,
18+
:prev_duplicate_pkt_num >= 0,
19+
:prev_recv_ack_num >= 0,
20+
:prev_status_query_msg_num >= 0;
21+
?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column?
22+
----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------
23+
t | t | t | t | t | t | t | t | t | t | t | t | t | t | t | t
24+
(1 row)
25+
26+
-- Create test table to generate interconnect traffic
27+
CREATE TABLE test_ic_data
28+
AS SELECT generate_series(1, 1000) AS id
29+
DISTRIBUTED RANDOMLY;
30+
-- Re-capture current state: overwrite prev with latest values
31+
SELECT * FROM gp_interconnect_stats \gset prev2_
32+
-- Check if current statistics are >= baseline values after first data insertion
33+
SELECT
34+
:prev2_total_recv_queue_size >= :prev_total_recv_queue_size,
35+
:prev2_recv_queue_conting_time >= :prev_recv_queue_conting_time,
36+
:prev2_total_capacity >= :prev_total_capacity,
37+
:prev2_capacity_counting_time >= :prev_capacity_counting_time,
38+
:prev2_total_buffers >= :prev_total_buffers,
39+
:prev2_buffer_counting_time >= :prev_buffer_counting_time,
40+
:prev2_retransmits >= :prev_retransmits,
41+
:prev2_startup_cached_pkts >= :prev_startup_cached_pkts,
42+
:prev2_mismatches >= :prev_mismatches,
43+
:prev2_crs_errors >= :prev_crs_errors,
44+
:prev2_snd_pkt_num >= :prev_snd_pkt_num,
45+
:prev2_recv_pkt_num >= :prev_recv_pkt_num,
46+
:prev2_disordered_pkt_num >= :prev_disordered_pkt_num,
47+
:prev2_duplicate_pkt_num >= :prev_duplicate_pkt_num,
48+
:prev2_recv_ack_num >= :prev_recv_ack_num,
49+
:prev2_status_query_msg_num >= :prev_status_query_msg_num;
50+
?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column?
51+
----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------
52+
t | t | t | t | t | t | t | t | t | t | t | t | t | t | t | t
53+
(1 row)
54+
55+
-- Insert additional data to further test interconnect statistics changes under load
56+
INSERT INTO test_ic_data SELECT generate_series(1001, 2000);
57+
-- Re‑check if current statistics remain >= baseline after second data insertion
58+
SELECT
59+
total_recv_queue_size >= :prev2_total_recv_queue_size,
60+
recv_queue_conting_time >= :prev2_recv_queue_conting_time,
61+
total_capacity >= :prev2_total_capacity,
62+
capacity_counting_time >= :prev2_capacity_counting_time,
63+
total_buffers >= :prev2_total_buffers,
64+
buffer_counting_time >= :prev2_buffer_counting_time,
65+
retransmits >= :prev2_retransmits,
66+
startup_cached_pkts >= :prev2_startup_cached_pkts,
67+
mismatches >= :prev2_mismatches,
68+
crs_errors >= :prev2_crs_errors,
69+
snd_pkt_num >= :prev2_snd_pkt_num,
70+
recv_pkt_num >= :prev2_recv_pkt_num,
71+
disordered_pkt_num >= :prev2_disordered_pkt_num,
72+
duplicate_pkt_num >= :prev2_duplicate_pkt_num,
73+
recv_ack_num >= :prev2_recv_ack_num,
74+
status_query_msg_num >= :prev2_status_query_msg_num
75+
FROM gp_interconnect_stats;
76+
?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column? | ?column?
77+
----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------+----------
78+
t | t | t | t | t | t | t | t | t | t | t | t | t | t | t | t
79+
(1 row)
80+
81+
DROP TABLE test_ic_data;
82+
DROP EXTENSION interconnect;

contrib/interconnect/ic_modules.c

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,13 +16,16 @@
1616
#include "ic_common.h"
1717
#include "tcp/ic_tcp.h"
1818
#include "udp/ic_udpifc.h"
19+
#include "storage/ipc.h"
1920

2021
#ifdef ENABLE_IC_PROXY
2122
#include "proxy/ic_proxy_server.h"
2223
#endif
2324

2425
PG_MODULE_MAGIC;
2526

27+
shmem_startup_hook_type prev_shmem_startup_hook = NULL;
28+
2629
MotionIPCLayer tcp_ipc_layer = {
2730
.ic_type = INTERCONNECT_TYPE_TCP,
2831
.type_name = "tcp",
@@ -141,6 +144,16 @@ MotionIPCLayer udpifc_ipc_layer = {
141144
.GetMotionSentRecordTypmod = GetMotionSentRecordTypmod,
142145
};
143146

147+
static void
148+
InterconnectShmemInit(void)
149+
{
150+
if (prev_shmem_startup_hook)
151+
prev_shmem_startup_hook();
152+
153+
if (Gp_interconnect_type == INTERCONNECT_TYPE_UDPIFC)
154+
InterconnectShmemInitUDPIFC();
155+
}
156+
144157
void
145158
_PG_init(void)
146159
{
@@ -153,4 +166,16 @@ _PG_init(void)
153166
RegisterIPCLayerImpl(&tcp_ipc_layer);
154167
RegisterIPCLayerImpl(&udpifc_ipc_layer);
155168
RegisterIPCLayerImpl(&proxy_ipc_layer);
169+
170+
if (Gp_interconnect_type == INTERCONNECT_TYPE_UDPIFC)
171+
RequestAddinShmemSpace(sizeof(ICStatisticsShmem));
172+
173+
prev_shmem_startup_hook = shmem_startup_hook;
174+
shmem_startup_hook = InterconnectShmemInit;
175+
}
176+
177+
void
178+
_PG_fini(void)
179+
{
180+
shmem_startup_hook = prev_shmem_startup_hook;
156181
}

contrib/interconnect/ic_modules.h

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,5 +18,6 @@ extern MotionIPCLayer proxy_ipc_layer;
1818
extern MotionIPCLayer udpifc_ipc_layer;
1919

2020
extern void _PG_init(void);
21+
extern void _PG_fini(void);
2122

2223
#endif // INTER_CONNECT_H
Lines changed: 110 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
1+
/* contrib/interconnect/interconnect--1.0.sql */
2+
3+
-- complain if script is sourced in psql, rather than via CREATE EXTENSION
4+
\echo Use "CREATE EXTENSION interconnect" to load this file. \quit
5+
6+
CREATE FUNCTION __gp_interconnect_get_stats_f_on_master(
7+
OUT gp_segment_id smallint,
8+
OUT total_recv_queue_size bigint,
9+
OUT recv_queue_conting_time bigint,
10+
OUT total_capacity bigint,
11+
OUT capacity_counting_time bigint,
12+
OUT total_buffers bigint,
13+
OUT buffer_counting_time bigint,
14+
OUT retransmits bigint,
15+
OUT startup_cached_pkts bigint,
16+
OUT mismatches bigint,
17+
OUT crs_errors bigint,
18+
OUT snd_pkt_num bigint,
19+
OUT recv_pkt_num bigint,
20+
OUT disordered_pkt_num bigint,
21+
OUT duplicate_pkt_num bigint,
22+
OUT recv_ack_num bigint,
23+
OUT status_query_msg_num bigint
24+
)
25+
RETURNS SETOF record
26+
LANGUAGE C VOLATILE EXECUTE ON MASTER
27+
AS '$libdir/interconnect', 'gp_interconnect_get_stats';
28+
29+
CREATE FUNCTION __gp_interconnect_get_stats_f_on_segments(
30+
OUT gp_segment_id smallint,
31+
OUT total_recv_queue_size bigint,
32+
OUT recv_queue_conting_time bigint,
33+
OUT total_capacity bigint,
34+
OUT capacity_counting_time bigint,
35+
OUT total_buffers bigint,
36+
OUT buffer_counting_time bigint,
37+
OUT retransmits bigint,
38+
OUT startup_cached_pkts bigint,
39+
OUT mismatches bigint,
40+
OUT crs_errors bigint,
41+
OUT snd_pkt_num bigint,
42+
OUT recv_pkt_num bigint,
43+
OUT disordered_pkt_num bigint,
44+
OUT duplicate_pkt_num bigint,
45+
OUT recv_ack_num bigint,
46+
OUT status_query_msg_num bigint
47+
)
48+
RETURNS SETOF record LANGUAGE C VOLATILE EXECUTE ON ALL SEGMENTS
49+
AS '$libdir/interconnect', 'gp_interconnect_get_stats';
50+
51+
52+
-- Cummulative interconnect statistics per segment
53+
CREATE VIEW gp_interconnect_stats_per_segment AS
54+
SELECT c.hostname, s.* FROM (
55+
SELECT * FROM __gp_interconnect_get_stats_f_on_master()
56+
UNION ALL
57+
SELECT * FROM __gp_interconnect_get_stats_f_on_segments()
58+
) s
59+
JOIN pg_catalog.gp_segment_configuration AS c
60+
ON s.gp_segment_id = c.content AND c.role = 'p';
61+
62+
GRANT SELECT ON gp_interconnect_stats_per_segment TO public;
63+
64+
-- Cummulative interconnect statistics
65+
CREATE VIEW gp_interconnect_stats AS
66+
SELECT
67+
sum(total_recv_queue_size) as total_recv_queue_size
68+
, sum(recv_queue_conting_time) as recv_queue_conting_time
69+
, sum(total_capacity) as total_capacity
70+
, sum(capacity_counting_time) as capacity_counting_time
71+
, sum(total_buffers) as total_buffers
72+
, sum(buffer_counting_time) as buffer_counting_time
73+
, sum(retransmits) as retransmits
74+
, sum(startup_cached_pkts) as startup_cached_pkts
75+
, sum(mismatches) as mismatches
76+
, sum(crs_errors) as crs_errors
77+
, sum(snd_pkt_num) as snd_pkt_num
78+
, sum(recv_pkt_num) as recv_pkt_num
79+
, sum(disordered_pkt_num) as disordered_pkt_num
80+
, sum(duplicate_pkt_num) as duplicate_pkt_num
81+
, sum(recv_ack_num) as recv_ack_num
82+
, sum(status_query_msg_num) as status_query_msg_num
83+
FROM gp_interconnect_stats_per_segment;
84+
85+
GRANT SELECT ON gp_interconnect_stats TO public;
86+
87+
-- Cummulative interconnect statistics grouped by host
88+
CREATE VIEW gp_interconnect_stats_per_host AS
89+
SELECT
90+
hostname
91+
, sum(total_recv_queue_size) as total_recv_queue_size
92+
, sum(recv_queue_conting_time) as recv_queue_conting_time
93+
, sum(total_capacity) as total_capacity
94+
, sum(capacity_counting_time) as capacity_counting_time
95+
, sum(total_buffers) as total_buffers
96+
, sum(buffer_counting_time) as buffer_counting_time
97+
, sum(retransmits) as retransmits
98+
, sum(startup_cached_pkts) as startup_cached_pkts
99+
, sum(mismatches) as mismatches
100+
, sum(crs_errors) as crs_errors
101+
, sum(snd_pkt_num) as snd_pkt_num
102+
, sum(recv_pkt_num) as recv_pkt_num
103+
, sum(disordered_pkt_num) as disordered_pkt_num
104+
, sum(duplicate_pkt_num) as duplicate_pkt_num
105+
, sum(recv_ack_num) as recv_ack_num
106+
, sum(status_query_msg_num) as status_query_msg_num
107+
FROM gp_interconnect_stats_per_segment
108+
GROUP BY hostname;
109+
110+
GRANT SELECT ON gp_interconnect_stats_per_host TO public;
Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,4 @@
1+
comment = 'Cummulative statistics from UDPIFC interconnect protocol'
2+
default_version = '1.0'
3+
relocatable = false
4+
schema = public

0 commit comments

Comments
 (0)