From 692c0f3fb87c31e1b995412e37ab8368a48b5a13 Mon Sep 17 00:00:00 2001 From: Konstantin Knizhnik Date: Mon, 28 Apr 2025 16:24:18 +0300 Subject: [PATCH] Prepare to prewarm support (#11740) ## Problem See (original prewarm implementation) https://github.com/neondatabase/neon/pull/9197 (functions for storing/restoring LFC state) https://github.com/neondatabase/neon/pull/9587 (store prefetch results in LFC) https://github.com/neondatabase/neon/pull/10442 ## Summary of changes Preparation for prewarm implementation. --------- Co-authored-by: Konstantin Knizhnik --- pgxn/neon/communicator.c | 32 +++++--- pgxn/neon/communicator.h | 2 + pgxn/neon/file_cache.c | 142 +++++++++++++++++++++++++---------- pgxn/neon/libpagestore.c | 5 +- pgxn/neon/pagestore_client.h | 3 +- 5 files changed, 133 insertions(+), 51 deletions(-) diff --git a/pgxn/neon/communicator.c b/pgxn/neon/communicator.c index db3e053321..61bb3206e7 100644 --- a/pgxn/neon/communicator.c +++ b/pgxn/neon/communicator.c @@ -88,9 +88,6 @@ typedef PGAlignedBlock PGIOAlignedBlock; page_server_api *page_server; -static uint32 local_request_counter; -#define GENERATE_REQUEST_ID() (((NeonRequestId)MyProcPid << 32) | ++local_request_counter) - /* * Various settings related to prompt (fast) handling of PageStream responses * at any CHECK_FOR_INTERRUPTS point. @@ -788,6 +785,27 @@ prefetch_read(PrefetchRequest *slot) } } + +/* + * Wait completion of previosly registered prefetch request. + * Prefetch result should be placed in LFC by prefetch_wait_for. + */ +bool +communicator_prefetch_receive(BufferTag tag) +{ + PrfHashEntry *entry; + PrefetchRequest hashkey; + + hashkey.buftag = tag; + entry = prfh_lookup(MyPState->prf_hash, &hashkey); + if (entry != NULL && prefetch_wait_for(entry->slot->my_ring_index)) + { + prefetch_set_unused(entry->slot->my_ring_index); + return true; + } + return false; +} + /* * Disconnect hook - drop prefetches when the connection drops * @@ -906,7 +924,6 @@ prefetch_do_request(PrefetchRequest *slot, neon_request_lsns *force_request_lsns NeonGetPageRequest request = { .hdr.tag = T_NeonGetPageRequest, - .hdr.reqid = GENERATE_REQUEST_ID(), /* lsn and not_modified_since are filled in below */ .rinfo = BufTagGetNRelFileInfo(slot->buftag), .forknum = slot->buftag.forkNum, @@ -915,8 +932,6 @@ prefetch_do_request(PrefetchRequest *slot, neon_request_lsns *force_request_lsns Assert(mySlotNo == MyPState->ring_unused); - slot->reqid = request.hdr.reqid; - if (force_request_lsns) slot->request_lsns = *force_request_lsns; else @@ -934,6 +949,7 @@ prefetch_do_request(PrefetchRequest *slot, neon_request_lsns *force_request_lsns Assert(mySlotNo == MyPState->ring_unused); /* loop */ } + slot->reqid = request.hdr.reqid; /* update prefetch state */ MyPState->n_requests_inflight += 1; @@ -1937,7 +1953,6 @@ communicator_exists(NRelFileInfo rinfo, ForkNumber forkNum, neon_request_lsns *r { NeonExistsRequest request = { .hdr.tag = T_NeonExistsRequest, - .hdr.reqid = GENERATE_REQUEST_ID(), .hdr.lsn = request_lsns->request_lsn, .hdr.not_modified_since = request_lsns->not_modified_since, .rinfo = rinfo, @@ -2212,7 +2227,6 @@ communicator_nblocks(NRelFileInfo rinfo, ForkNumber forknum, neon_request_lsns * { NeonNblocksRequest request = { .hdr.tag = T_NeonNblocksRequest, - .hdr.reqid = GENERATE_REQUEST_ID(), .hdr.lsn = request_lsns->request_lsn, .hdr.not_modified_since = request_lsns->not_modified_since, .rinfo = rinfo, @@ -2285,7 +2299,6 @@ communicator_dbsize(Oid dbNode, neon_request_lsns *request_lsns) { NeonDbSizeRequest request = { .hdr.tag = T_NeonDbSizeRequest, - .hdr.reqid = GENERATE_REQUEST_ID(), .hdr.lsn = request_lsns->request_lsn, .hdr.not_modified_since = request_lsns->not_modified_since, .dbNode = dbNode, @@ -2353,7 +2366,6 @@ communicator_read_slru_segment(SlruKind kind, int64 segno, neon_request_lsns *re request = (NeonGetSlruSegmentRequest) { .hdr.tag = T_NeonGetSlruSegmentRequest, - .hdr.reqid = GENERATE_REQUEST_ID(), .hdr.lsn = request_lsns->request_lsn, .hdr.not_modified_since = request_lsns->not_modified_since, .kind = kind, diff --git a/pgxn/neon/communicator.h b/pgxn/neon/communicator.h index 72cba526c1..f55c4b10f1 100644 --- a/pgxn/neon/communicator.h +++ b/pgxn/neon/communicator.h @@ -37,6 +37,8 @@ extern int communicator_prefetch_lookupv(NRelFileInfo rinfo, ForkNumber forknum, BlockNumber nblocks, void **buffers, bits8 *mask); extern void communicator_prefetch_register_bufferv(BufferTag tag, neon_request_lsns *frlsns, BlockNumber nblocks, const bits8 *mask); +extern bool communicator_prefetch_receive(BufferTag tag); + extern int communicator_read_slru_segment(SlruKind kind, int64 segno, neon_request_lsns *request_lsns, void *buffer); diff --git a/pgxn/neon/file_cache.c b/pgxn/neon/file_cache.c index 8c2990e57a..e2c1f7682f 100644 --- a/pgxn/neon/file_cache.c +++ b/pgxn/neon/file_cache.c @@ -25,6 +25,7 @@ #include "pgstat.h" #include "port/pg_iovec.h" #include "postmaster/bgworker.h" +#include "postmaster/interrupt.h" #include RELFILEINFO_HDR #include "storage/buf_internals.h" #include "storage/fd.h" @@ -32,6 +33,8 @@ #include "storage/latch.h" #include "storage/lwlock.h" #include "storage/pg_shmem.h" +#include "storage/procsignal.h" +#include "tcop/tcopprot.h" #include "utils/builtins.h" #include "utils/dynahash.h" #include "utils/guc.h" @@ -46,6 +49,8 @@ #include "neon.h" #include "neon_lwlsncache.h" #include "neon_perf_counters.h" +#include "pagestore_client.h" +#include "communicator.h" #define CriticalAssert(cond) do if (!(cond)) elog(PANIC, "LFC: assertion %s failed at %s:%d: ", #cond, __FILE__, __LINE__); while (0) @@ -87,14 +92,14 @@ * 1Mb chunks can reduce hash map size to 320Mb. * 2. Improve access locality, subsequent pages will be allocated together improving seqscan speed */ -#define BLOCKS_PER_CHUNK 128 /* 1Mb chunk */ -/* - * Smaller chunk seems to be better for OLTP workload - */ -// #define BLOCKS_PER_CHUNK 8 /* 64kb chunk */ +#define MAX_BLOCKS_PER_CHUNK_LOG 7 /* 1Mb chunk */ +#define MAX_BLOCKS_PER_CHUNK (1 << MAX_BLOCKS_PER_CHUNK_LOG) + #define MB ((uint64)1024*1024) -#define SIZE_MB_TO_CHUNKS(size) ((uint32)((size) * MB / BLCKSZ / BLOCKS_PER_CHUNK)) +#define SIZE_MB_TO_CHUNKS(size) ((uint32)((size) * MB / BLCKSZ >> lfc_chunk_size_log)) + +#define BLOCK_TO_CHUNK_OFF(blkno) ((blkno) & (lfc_blocks_per_chunk-1)) /* * Blocks are read or written to LFC file outside LFC critical section. @@ -119,10 +124,11 @@ typedef struct FileCacheEntry uint32 hash; uint32 offset; uint32 access_count; - uint32 state[(BLOCKS_PER_CHUNK + 31) / 32 * 2]; /* two bits per block */ dlist_node list_node; /* LRU/holes list node */ + uint32 state[FLEXIBLE_ARRAY_MEMBER]; /* two bits per block */ } FileCacheEntry; +#define FILE_CACHE_ENRTY_SIZE MAXALIGN(offsetof(FileCacheEntry, state) + (lfc_blocks_per_chunk*2+31)/32*4) #define GET_STATE(entry, i) (((entry)->state[(i) / 16] >> ((i) % 16 * 2)) & 3) #define SET_STATE(entry, i, new_state) (entry)->state[(i) / 16] = ((entry)->state[(i) / 16] & ~(3 << ((i) % 16 * 2))) | ((new_state) << ((i) % 16 * 2)) @@ -136,6 +142,7 @@ typedef struct FileCacheControl uint32 size; /* size of cache file in chunks */ uint32 used; /* number of used chunks */ uint32 used_pages; /* number of used pages */ + uint32 pinned; /* number of pinned chunks */ uint32 limit; /* shared copy of lfc_size_limit */ uint64 hits; uint64 misses; @@ -158,6 +165,8 @@ static int lfc_desc = -1; static LWLockId lfc_lock; static int lfc_max_size; static int lfc_size_limit; +static int lfc_chunk_size_log = MAX_BLOCKS_PER_CHUNK_LOG; +static int lfc_blocks_per_chunk = MAX_BLOCKS_PER_CHUNK; static char *lfc_path; static uint64 lfc_generation; static FileCacheControl *lfc_ctl; @@ -206,7 +215,9 @@ lfc_switch_off(void) } lfc_ctl->generation += 1; lfc_ctl->size = 0; + lfc_ctl->pinned = 0; lfc_ctl->used = 0; + lfc_ctl->used_pages = 0; lfc_ctl->limit = 0; dlist_init(&lfc_ctl->lru); dlist_init(&lfc_ctl->holes); @@ -296,7 +307,7 @@ lfc_shmem_startup(void) lfc_lock = (LWLockId) GetNamedLWLockTranche("lfc_lock"); info.keysize = sizeof(BufferTag); - info.entrysize = sizeof(FileCacheEntry); + info.entrysize = FILE_CACHE_ENRTY_SIZE; /* * n_chunks+1 because we add new element to hash table before eviction @@ -342,7 +353,7 @@ lfc_shmem_request(void) prev_shmem_request_hook(); #endif - RequestAddinShmemSpace(sizeof(FileCacheControl) + hash_estimate_size(SIZE_MB_TO_CHUNKS(lfc_max_size) + 1, sizeof(FileCacheEntry))); + RequestAddinShmemSpace(sizeof(FileCacheControl) + hash_estimate_size(SIZE_MB_TO_CHUNKS(lfc_max_size) + 1, FILE_CACHE_ENRTY_SIZE)); RequestNamedLWLockTranche("lfc_lock", 1); } @@ -359,6 +370,24 @@ is_normal_backend(void) return lfc_ctl && MyProc && UsedShmemSegAddr && !IsParallelWorker(); } +static bool +lfc_check_chunk_size(int *newval, void **extra, GucSource source) +{ + if (*newval & (*newval - 1)) + { + elog(ERROR, "LFC chunk size should be power of two"); + return false; + } + return true; +} + +static void +lfc_change_chunk_size(int newval, void* extra) +{ + lfc_chunk_size_log = pg_ceil_log2_32(newval); +} + + static bool lfc_check_limit_hook(int *newval, void **extra, GucSource source) { @@ -415,11 +444,11 @@ lfc_change_limit_hook(int newval, void *extra) CriticalAssert(victim->access_count == 0); #ifdef FALLOC_FL_PUNCH_HOLE - if (fallocate(lfc_desc, FALLOC_FL_PUNCH_HOLE | FALLOC_FL_KEEP_SIZE, (off_t) victim->offset * BLOCKS_PER_CHUNK * BLCKSZ, BLOCKS_PER_CHUNK * BLCKSZ) < 0) + if (fallocate(lfc_desc, FALLOC_FL_PUNCH_HOLE | FALLOC_FL_KEEP_SIZE, (off_t) victim->offset * lfc_blocks_per_chunk * BLCKSZ, lfc_blocks_per_chunk * BLCKSZ) < 0) neon_log(LOG, "Failed to punch hole in file: %m"); #endif /* We remove the old entry, and re-enter a hole to the hash table */ - for (int i = 0; i < BLOCKS_PER_CHUNK; i++) + for (int i = 0; i < lfc_blocks_per_chunk; i++) { bool is_page_cached = GET_STATE(victim, i) == AVAILABLE; lfc_ctl->used_pages -= is_page_cached; @@ -508,6 +537,19 @@ lfc_init(void) NULL, NULL); + DefineCustomIntVariable("neon.file_cache_chunk_size", + "LFC chunk size in blocks (should be power of two)", + NULL, + &lfc_blocks_per_chunk, + MAX_BLOCKS_PER_CHUNK, + 1, + MAX_BLOCKS_PER_CHUNK, + PGC_POSTMASTER, + GUC_UNIT_BLOCKS, + lfc_check_chunk_size, + lfc_change_chunk_size, + NULL); + if (lfc_max_size == 0) return; @@ -530,7 +572,7 @@ lfc_cache_contains(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno) { BufferTag tag; FileCacheEntry *entry; - int chunk_offs = blkno & (BLOCKS_PER_CHUNK - 1); + int chunk_offs = BLOCK_TO_CHUNK_OFF(blkno); bool found = false; uint32 hash; @@ -539,7 +581,7 @@ lfc_cache_contains(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno) CopyNRelFileInfoToBufTag(tag, rinfo); tag.forkNum = forkNum; - tag.blockNum = blkno & ~(BLOCKS_PER_CHUNK - 1); + tag.blockNum = blkno - chunk_offs; CriticalAssert(BufTagGetRelNumber(&tag) != InvalidRelFileNumber); hash = get_hash_value(lfc_hash, &tag); @@ -577,9 +619,9 @@ lfc_cache_containsv(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, CriticalAssert(BufTagGetRelNumber(&tag) != InvalidRelFileNumber); - tag.blockNum = blkno & ~(BLOCKS_PER_CHUNK - 1); + chunk_offs = BLOCK_TO_CHUNK_OFF(blkno); + tag.blockNum = blkno - chunk_offs; hash = get_hash_value(lfc_hash, &tag); - chunk_offs = blkno & (BLOCKS_PER_CHUNK - 1); LWLockAcquire(lfc_lock, LW_SHARED); @@ -590,12 +632,12 @@ lfc_cache_containsv(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, } while (true) { - int this_chunk = Min(nblocks - i, BLOCKS_PER_CHUNK - chunk_offs); + int this_chunk = Min(nblocks - i, lfc_blocks_per_chunk - chunk_offs); entry = hash_search_with_hash_value(lfc_hash, &tag, hash, HASH_FIND, NULL); if (entry != NULL) { - for (; chunk_offs < BLOCKS_PER_CHUNK && i < nblocks; chunk_offs++, i++) + for (; chunk_offs < lfc_blocks_per_chunk && i < nblocks; chunk_offs++, i++) { if (GET_STATE(entry, chunk_offs) != UNAVAILABLE) { @@ -619,9 +661,9 @@ lfc_cache_containsv(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, * Prepare for the next iteration. We don't unlock here, as that'd * probably be more expensive than the gains it'd get us. */ - tag.blockNum = (blkno + i) & ~(BLOCKS_PER_CHUNK - 1); + chunk_offs = BLOCK_TO_CHUNK_OFF(blkno + i); + tag.blockNum = (blkno + i) - chunk_offs; hash = get_hash_value(lfc_hash, &tag); - chunk_offs = (blkno + i) & (BLOCKS_PER_CHUNK - 1); } LWLockRelease(lfc_lock); @@ -696,9 +738,9 @@ lfc_readv_select(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, while (nblocks > 0) { struct iovec iov[PG_IOV_MAX]; - int8 chunk_mask[BLOCKS_PER_CHUNK / 8] = {0}; - int chunk_offs = (blkno & (BLOCKS_PER_CHUNK - 1)); - int blocks_in_chunk = Min(nblocks, BLOCKS_PER_CHUNK - (blkno % BLOCKS_PER_CHUNK)); + uint8 chunk_mask[MAX_BLOCKS_PER_CHUNK / 8] = {0}; + int chunk_offs = BLOCK_TO_CHUNK_OFF(blkno); + int blocks_in_chunk = Min(nblocks, lfc_blocks_per_chunk - chunk_offs); int iteration_hits = 0; int iteration_misses = 0; uint64 io_time_us = 0; @@ -786,8 +828,10 @@ lfc_readv_select(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, /* Unlink entry from LRU list to pin it for the duration of IO operation */ if (entry->access_count++ == 0) + { + lfc_ctl->pinned += 1; dlist_delete(&entry->list_node); - + } generation = lfc_ctl->generation; entry_offset = entry->offset; @@ -836,7 +880,7 @@ lfc_readv_select(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, if (iteration_hits != 0) { /* chunk offset (# of pages) into the LFC file */ - off_t first_read_offset = (off_t) entry_offset * BLOCKS_PER_CHUNK; + off_t first_read_offset = (off_t) entry_offset * lfc_blocks_per_chunk; int nwrite = iov_last_used - first_block_in_chunk_read; /* offset of first IOV */ first_read_offset += chunk_offs + first_block_in_chunk_read; @@ -884,7 +928,10 @@ lfc_readv_select(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, CriticalAssert(entry->access_count > 0); if (--entry->access_count == 0) + { + lfc_ctl->pinned -= 1; dlist_push_tail(&lfc_ctl->lru, &entry->list_node); + } } else { @@ -961,7 +1008,7 @@ lfc_init_new_entry(FileCacheEntry* entry, uint32 hash) FileCacheEntry *victim = dlist_container(FileCacheEntry, list_node, dlist_pop_head_node(&lfc_ctl->lru)); - for (int i = 0; i < BLOCKS_PER_CHUNK; i++) + for (int i = 0; i < lfc_blocks_per_chunk; i++) { bool is_page_cached = GET_STATE(victim, i) == AVAILABLE; lfc_ctl->used_pages -= is_page_cached; @@ -979,14 +1026,14 @@ lfc_init_new_entry(FileCacheEntry* entry, uint32 hash) /* Can't add this chunk - we don't have the space for it */ hash_search_with_hash_value(lfc_hash, &entry->key, hash, HASH_REMOVE, NULL); - return false; } entry->access_count = 1; entry->hash = hash; + lfc_ctl->pinned += 1; - for (int i = 0; i < BLOCKS_PER_CHUNK; i++) + for (int i = 0; i < lfc_blocks_per_chunk; i++) SET_STATE(entry, i, UNAVAILABLE); return true; @@ -1031,7 +1078,7 @@ lfc_prefetch(NRelFileInfo rinfo, ForkNumber forknum, BlockNumber blkno, FileCacheBlockState state; XLogRecPtr lwlsn; - int chunk_offs = blkno & (BLOCKS_PER_CHUNK - 1); + int chunk_offs = BLOCK_TO_CHUNK_OFF(blkno); if (lfc_maybe_disabled()) /* fast exit if file cache is disabled */ return false; @@ -1041,7 +1088,7 @@ lfc_prefetch(NRelFileInfo rinfo, ForkNumber forknum, BlockNumber blkno, CriticalAssert(BufTagGetRelNumber(&tag) != InvalidRelFileNumber); - tag.blockNum = blkno & ~(BLOCKS_PER_CHUNK - 1); + tag.blockNum = blkno - chunk_offs; hash = get_hash_value(lfc_hash, &tag); cv = &lfc_ctl->cv[hash % N_COND_VARS]; @@ -1052,7 +1099,7 @@ lfc_prefetch(NRelFileInfo rinfo, ForkNumber forknum, BlockNumber blkno, LWLockRelease(lfc_lock); return false; } - + lwlsn = neon_get_lwlsn(rinfo, forknum, blkno); if (lwlsn > lsn) @@ -1081,7 +1128,10 @@ lfc_prefetch(NRelFileInfo rinfo, ForkNumber forknum, BlockNumber blkno, * operation */ if (entry->access_count++ == 0) + { + lfc_ctl->pinned += 1; dlist_delete(&entry->list_node); + } } else { @@ -1106,7 +1156,7 @@ lfc_prefetch(NRelFileInfo rinfo, ForkNumber forknum, BlockNumber blkno, pgstat_report_wait_start(WAIT_EVENT_NEON_LFC_WRITE); INSTR_TIME_SET_CURRENT(io_start); rc = pwrite(lfc_desc, buffer, BLCKSZ, - ((off_t) entry_offset * BLOCKS_PER_CHUNK + chunk_offs) * BLCKSZ); + ((off_t) entry_offset * lfc_blocks_per_chunk + chunk_offs) * BLCKSZ); INSTR_TIME_SET_CURRENT(io_end); pgstat_report_wait_end(); @@ -1132,7 +1182,10 @@ lfc_prefetch(NRelFileInfo rinfo, ForkNumber forknum, BlockNumber blkno, inc_page_cache_write_wait(time_spent_us); if (--entry->access_count == 0) + { + lfc_ctl->pinned -= 1; dlist_push_tail(&lfc_ctl->lru, &entry->list_node); + } state = GET_STATE(entry, chunk_offs); if (state == REQUESTED) { @@ -1199,8 +1252,8 @@ lfc_writev(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, while (nblocks > 0) { struct iovec iov[PG_IOV_MAX]; - int chunk_offs = blkno & (BLOCKS_PER_CHUNK - 1); - int blocks_in_chunk = Min(nblocks, BLOCKS_PER_CHUNK - (blkno % BLOCKS_PER_CHUNK)); + int chunk_offs = BLOCK_TO_CHUNK_OFF(blkno); + int blocks_in_chunk = Min(nblocks, lfc_blocks_per_chunk - chunk_offs); instr_time io_start, io_end; ConditionVariable* cv; @@ -1212,7 +1265,7 @@ lfc_writev(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, iov[i].iov_len = BLCKSZ; } - tag.blockNum = blkno & ~(BLOCKS_PER_CHUNK - 1); + tag.blockNum = blkno - chunk_offs; hash = get_hash_value(lfc_hash, &tag); cv = &lfc_ctl->cv[hash % N_COND_VARS]; @@ -1232,7 +1285,10 @@ lfc_writev(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, * operation */ if (entry->access_count++ == 0) + { + lfc_ctl->pinned += 1; dlist_delete(&entry->list_node); + } } else { @@ -1285,7 +1341,7 @@ lfc_writev(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, pgstat_report_wait_start(WAIT_EVENT_NEON_LFC_WRITE); INSTR_TIME_SET_CURRENT(io_start); rc = pwritev(lfc_desc, iov, blocks_in_chunk, - ((off_t) entry_offset * BLOCKS_PER_CHUNK + chunk_offs) * BLCKSZ); + ((off_t) entry_offset * lfc_blocks_per_chunk + chunk_offs) * BLCKSZ); INSTR_TIME_SET_CURRENT(io_end); pgstat_report_wait_end(); @@ -1312,7 +1368,10 @@ lfc_writev(NRelFileInfo rinfo, ForkNumber forkNum, BlockNumber blkno, inc_page_cache_write_wait(time_spent_us); if (--entry->access_count == 0) + { + lfc_ctl->pinned -= 1; dlist_push_tail(&lfc_ctl->lru, &entry->list_node); + } for (int i = 0; i < blocks_in_chunk; i++) { @@ -1438,7 +1497,12 @@ neon_get_lfc_stats(PG_FUNCTION_ARGS) break; case 8: key = "file_cache_chunk_size_pages"; - value = BLOCKS_PER_CHUNK; + value = lfc_blocks_per_chunk; + break; + case 9: + key = "file_cache_chunks_pinned"; + if (lfc_ctl) + value = lfc_ctl->pinned; break; default: SRF_RETURN_DONE(funcctx); @@ -1566,7 +1630,7 @@ local_cache_pages(PG_FUNCTION_ARGS) /* Skip hole tags */ if (NInfoGetRelNumber(BufTagGetNRelFileInfo(entry->key)) != 0) { - for (int i = 0; i < BLOCKS_PER_CHUNK; i++) + for (int i = 0; i < lfc_blocks_per_chunk; i++) n_pages += GET_STATE(entry, i) == AVAILABLE; } } @@ -1594,13 +1658,13 @@ local_cache_pages(PG_FUNCTION_ARGS) hash_seq_init(&status, lfc_hash); while ((entry = hash_seq_search(&status)) != NULL) { - for (int i = 0; i < BLOCKS_PER_CHUNK; i++) + for (int i = 0; i < lfc_blocks_per_chunk; i++) { if (NInfoGetRelNumber(BufTagGetNRelFileInfo(entry->key)) != 0) { if (GET_STATE(entry, i) == AVAILABLE) { - fctx->record[n].pageoffs = entry->offset * BLOCKS_PER_CHUNK + i; + fctx->record[n].pageoffs = entry->offset * lfc_blocks_per_chunk + i; fctx->record[n].relfilenode = NInfoGetRelNumber(BufTagGetNRelFileInfo(entry->key)); fctx->record[n].reltablespace = NInfoGetSpcOid(BufTagGetNRelFileInfo(entry->key)); fctx->record[n].reldatabase = NInfoGetDbOid(BufTagGetNRelFileInfo(entry->key)); diff --git a/pgxn/neon/libpagestore.c b/pgxn/neon/libpagestore.c index 64d38e7913..ccb072d6f9 100644 --- a/pgxn/neon/libpagestore.c +++ b/pgxn/neon/libpagestore.c @@ -48,7 +48,6 @@ #define MIN_RECONNECT_INTERVAL_USEC 1000 #define MAX_RECONNECT_INTERVAL_USEC 1000000 - enum NeonComputeMode { CP_MODE_PRIMARY = 0, CP_MODE_REPLICA, @@ -167,6 +166,9 @@ typedef struct WaitEventSet *wes_read; } PageServer; +static uint32 local_request_counter; +#define GENERATE_REQUEST_ID() (((NeonRequestId)MyProcPid << 32) | ++local_request_counter) + static PageServer page_servers[MAX_SHARDS]; static bool pageserver_flush(shardno_t shard_no); @@ -994,6 +996,7 @@ pageserver_send(shardno_t shard_no, NeonRequest *request) pageserver_conn = NULL; } + request->reqid = GENERATE_REQUEST_ID(); req_buff = nm_pack_request(request); /* diff --git a/pgxn/neon/pagestore_client.h b/pgxn/neon/pagestore_client.h index 0ab539fe56..9df202290d 100644 --- a/pgxn/neon/pagestore_client.h +++ b/pgxn/neon/pagestore_client.h @@ -65,7 +65,6 @@ typedef enum { SLRU_MULTIXACT_OFFSETS } SlruKind; - /*-- * supertype of all the Neon*Request structs below. * @@ -129,6 +128,7 @@ typedef struct int segno; } NeonGetSlruSegmentRequest; + /* supertype of all the Neon*Response structs below */ typedef NeonMessage NeonResponse; @@ -187,6 +187,7 @@ typedef struct { /* * Send this request to the PageServer associated with this shard. + * This function assigns request_id to the request which can be extracted by caller from request struct. */ bool (*send) (shardno_t shard_no, NeonRequest * request); /*