Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
452fa9e
fix(cluster): preserve authority across transient quorum publication
Oct 7, 2026
e2de79c
fix(cluster): yield pending serving checks before cache requests
Oct 7, 2026
8cf0a7e
test(cluster): expose formation identity loss on qualification resamp…
Oct 7, 2026
10affa5
fix(cluster): preserve formation identity from one admission observation
Oct 7, 2026
bcc5755
fix(cluster): retain data frames through pending admission
Oct 7, 2026
a1a7203
test(cluster): complete pending admission fixture integration
Oct 7, 2026
224fc51
fix(cluster): wait for pending GES admission before reserving
Oct 7, 2026
5fa0cf3
test(cluster): reproduce member COMMIT carrier loss on stable builds
Oct 7, 2026
f218c7d
test(cluster): cover COMMIT request ordering and late peer receipts
Oct 7, 2026
5e9b322
fix(cluster): retain pending member COMMIT proof across observation gaps
Oct 7, 2026
b455b3f
fix(cluster): retain admitted operations across pending observations
Oct 7, 2026
f545e51
test(cluster): cover native checkpoint admission retries
Oct 7, 2026
0e7ed42
fix(cluster): retry pending control-root checkpoint reads
Oct 8, 2026
cd2089b
feat(buffer): keep read-only version anchors in native mapping
Oct 8, 2026
ca1f133
feat(snapshot): bind CR keys to retained snapshot lifetimes
Oct 8, 2026
c59a0e4
feat(buffer): define scoped read-only version metadata
Oct 8, 2026
cdec4c4
fix(buffer): retire read-only versions through native invalidation
Oct 8, 2026
c0ed437
feat(buffer): publish immutable CR pages through native buffer owners
Oct 8, 2026
04dafa6
feat(heap): reuse snapshot CR pages in native index fetch owners
Oct 8, 2026
c8f5ab6
fix(cr): reuse retained local terminal proofs and bypass private shar…
Oct 8, 2026
2069b48
Preserve native acquisition and visibility on CR cache misses
Oct 8, 2026
6d7006e
Register native lock requests from exact master grants
Oct 8, 2026
7967127
Refresh R11 source census after native request registration fix
Oct 8, 2026
6b4d0b8
test(cluster): reproduce compatible GRD owner capacity refusal
Oct 8, 2026
390b9f4
test(cluster): cover GRD capacity queues and exact cancellation
Oct 8, 2026
9274fe0
test: cover GRD BAST fanout and bounded pool exhaustion
Oct 8, 2026
d82647b
test: reproduce cold GRD attach cleanup callback ordering
Oct 8, 2026
94910bf
test: cover full LMON convert grant fanout on release
Oct 8, 2026
06c29dd
test: verify exact dedup keys for capacity release grants
Oct 8, 2026
e5d3dd5
test: cover full control drain batches and exact retirement
Oct 8, 2026
0ad57ec
Scale exact GRD owner capacity to the configured cohort
Oct 8, 2026
3aae334
test: preserve GES grant delivery across projection OOM
Oct 8, 2026
665df3d
Use current origin authority for shared transaction waits
Oct 8, 2026
2d857ab
Reuse retained terminal proofs during scratch visibility scans
Oct 8, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
268 changes: 263 additions & 5 deletions src/backend/access/heap/heapam.c
Original file line number Diff line number Diff line change
Expand Up @@ -4901,7 +4901,7 @@ static bool
heap_hot_r4_search_scratch(const BufferTag *tag,
const ItemPointerData *logical_root,
Relation relation, Snapshot snapshot,
HeapHotSearchResult *result)
HeapHotSearchResult *result, bool scoped_full)
{
ClusterR4HotScratchTestContext context;
Page page = (Page) result->scratch_page;
Expand Down Expand Up @@ -4950,7 +4950,16 @@ heap_hot_r4_search_scratch(const BufferTag *tag,

if (offnum < FirstOffsetNumber
|| offnum > PageGetMaxOffsetNumber(page))
{
/* A later index insertion may name a new root beyond this older
* snapshot-complete page. Only the continuous scan/snapshot owner
* proves that absence; a traversed HOT/redirect edge never does. */
if (scoped_full && at_chain_start
&& offnum > PageGetMaxOffsetNumber(page)
&& offnum <= MaxHeapTuplesPerPage)
return false;
heap_hot_r4_unknown("HOT offset is absent from the FULL page");
}
if (offnum > MaxHeapTuplesPerPage || visited[offnum])
heap_hot_r4_unknown("HOT chain contains a cycle");
if (++depth > MaxHeapTuplesPerPage)
Expand All @@ -4962,9 +4971,11 @@ heap_hot_r4_search_scratch(const BufferTag *tag,
{
/* The completed CR producer removes versions born after the
* statement SCN. An absent first index root is invisible, not a
* broken traversed edge. The outer caller still revalidates its
* current input and never marks this index entry all-dead. */
if (at_chain_start && !ItemIdIsUsed(lp)
* broken traversed edge. The outer caller still revalidates its
* current input or continuous scan scope, and never marks this
* index entry all-dead. */
if (at_chain_start
&& (!ItemIdIsUsed(lp) || (scoped_full && ItemIdIsDead(lp)))
&& ItemIdGetOffset(lp) == 0 && ItemIdGetLength(lp) == 0)
return false;
if (ItemIdIsRedirected(lp) && at_chain_start)
Expand Down Expand Up @@ -4999,6 +5010,18 @@ heap_hot_r4_search_scratch(const BufferTag *tag,
prev_xmax, HeapTupleHeaderGetXmin(result->tuple.t_data)))
heap_hot_r4_unknown("HOT predecessor identity does not match xmin");

/* This optional cache has no current authority for MultiXact I/O.
* Return to the original visibility owner before evaluating this
* known unsupported shape; actual UNKNOWN outcomes still throw. */
if (scoped_full
&& (result->tuple.t_data->t_infomask & HEAP_XMAX_IS_MULTI) != 0
&& (result->tuple.t_data->t_infomask & HEAP_XMAX_INVALID) == 0
&& !HEAP_XMAX_IS_LOCKED_ONLY(result->tuple.t_data->t_infomask))
{
result->cr_unsupported = true;
return false;
}

if (HeapTupleSatisfiesMVCCScratch(&result->tuple, snapshot, &context))
return true;

Expand All @@ -5008,7 +5031,14 @@ heap_hot_r4_search_scratch(const BufferTag *tag,
|| ItemPointerGetBlockNumber(&result->tuple.t_data->t_ctid) != blkno)
heap_hot_r4_unknown("HOT chain crosses the reconstructed block");
if ((result->tuple.t_data->t_infomask & HEAP_XMAX_IS_MULTI) != 0)
{
if (scoped_full)
{
result->cr_unsupported = true;
return false;
}
heap_hot_r4_unknown("scratch HOT predecessor requires MultiXact I/O");
}

offnum = ItemPointerGetOffsetNumber(&result->tuple.t_data->t_ctid);
at_chain_start = false;
Expand Down Expand Up @@ -5153,7 +5183,7 @@ heap_hot_r4_full_cycle(Buffer buffer, const BufferTag *tag,
&& build_reason == CLUSTER_CR_BUILD_NONE)
{
if (heap_hot_r4_search_scratch(tag, logical_root, relation,
snapshot, result))
snapshot, result, false))
result->kind = HEAP_HOT_SEARCH_OWNED_SCRATCH;
}
}
Expand Down Expand Up @@ -5190,6 +5220,232 @@ heap_hot_r4_full_failure(SCN read_scn, ClusterCrBuildResult build_result,
cluster_cr_build_reason_name(build_reason))));
pg_unreachable();
}

/* The original transaction AS and active Relation protect this scan's locator.
* This is deliberately not a cross-scan or catalog relation-identity proof. */
static bool
heap_index_cr_eligible(IndexFetchHeapData *hscan, Snapshot snapshot)
{
Relation relation = hscan->xs_base.rel;

return cluster_shared_config && cluster_shared_catalog
&& cluster_storage_mode_enabled()
&& relation != NULL && relation->rd_refcnt > 0
&& RelationIsPermanent(relation) && !relation->rd_rel->relisshared
&& !IsCatalogRelation(relation)
&& (relation->rd_rel->relkind == RELKIND_RELATION
|| relation->rd_rel->relkind == RELKIND_MATVIEW)
&& snapshot != NULL && snapshot->snapshot_type == SNAPSHOT_MVCC
&& snapshot->cluster_source == SNAPSHOT_SOURCE_CLUSTER
&& !TransactionIdIsValid(GetTopTransactionIdIfAny())
&& !IsolationIsSerializable()
&& CurrentResourceOwner != NULL
&& CheckRelationLockedByMe(relation, AccessShareLock, true);
}

static bool
heap_index_cr_scope_matches(IndexFetchHeapData *hscan, Snapshot snapshot,
uint64 snapshot_id)
{
const HeapReadOnlyCrScope *scope = &hscan->cr_scope;
Relation relation = hscan->xs_base.rel;

return scope->scan_id != 0 && scope->reserved == 0
&& scope->relation == relation && scope->owner == CurrentResourceOwner
&& scope->relation_oid == RelationGetRelid(relation)
&& RelFileLocatorEquals(scope->locator, relation->rd_locator)
&& scope->command_id == snapshot->curcid
&& scope->snapshot_id == snapshot_id
&& scope->read_scn == snapshot->read_scn
&& scope->read_epoch == snapshot->read_epoch;
}

static void
heap_index_cr_recheck(IndexFetchHeapData *hscan, Snapshot snapshot,
const ClusterSemanticAdmissionToken *admission)
{
uint64 snapshot_id;

if (!heap_index_cr_eligible(hscan, snapshot)
|| !cluster_snapshot_cr_identity_v1(snapshot, &snapshot_id)
|| !heap_index_cr_scope_matches(hscan, snapshot, snapshot_id)
|| !cluster_semantic_activation_recheck(admission))
heap_hot_r4_unknown("read-only scan or TARGET identity changed");
}

/*
* One native index-fetch owner can reuse its immutable FULL versions before
* taking current S. The key is not authority: actual retained snapshot and
* fresh TARGET admission cover every hit, the original scratch resolver and
* publication. No buffer/content/mapping lock is held during R4 or TT work.
*/
bool
heap_index_fetch_cr_result(IndexFetchHeapData *hscan, ItemPointer tid,
Snapshot snapshot, HeapHotSearchResult *result)
{
ClusterSemanticAdmissionToken admission;
ClusterSemanticAdmissionResult admission_result;
ClusterSnapshotReadScopeV1 read_scope;
BufferCrKey key;
ItemPointerData logical_root = *tid;
uint64 snapshot_id;
bool found;
bool handled = false;

if (!heap_index_cr_eligible(hscan, snapshot))
{
memset(&hscan->cr_scope, 0, sizeof(hscan->cr_scope));
return false;
}
admission_result = cluster_semantic_activation_enter(
CLUSTER_SEMANTIC_FEATURE_R4_SYNC_CR_V1,
CLUSTER_SEMANTIC_TARGET_SIDE, &admission);
if (admission_result != CLUSTER_SEMANTIC_ADMISSION_OK)
{
memset(&hscan->cr_scope, 0, sizeof(hscan->cr_scope));
if (admission_result == CLUSTER_SEMANTIC_ADMISSION_TARGET_DISABLED)
return false;
heap_hot_r4_unknown("read-only scan TARGET admission refused");
}
cluster_snapshot_read_enter_v1(&read_scope, snapshot);
PG_TRY();
{
PG_TRY();
{
if (!cluster_snapshot_cr_identity_v1(snapshot, &snapshot_id))
heap_hot_r4_unknown("read-only scan has no retained snapshot identity");
if (!heap_index_cr_scope_matches(hscan, snapshot, snapshot_id))
{
HeapReadOnlyCrScope *scope = &hscan->cr_scope;

memset(scope, 0, sizeof(*scope));
if (!BufTableNewCRScope(&scope->scan_id))
heap_hot_r4_unknown("read-only scan identity is exhausted");
scope->relation = hscan->xs_base.rel;
scope->owner = CurrentResourceOwner;
scope->locator = scope->relation->rd_locator;
scope->relation_oid = RelationGetRelid(scope->relation);
scope->command_id = snapshot->curcid;
scope->snapshot_id = snapshot_id;
scope->read_scn = snapshot->read_scn;
scope->read_epoch = snapshot->read_epoch;
}
memset(&key, 0, sizeof(key));
InitBufferTag(&key.tag, &hscan->cr_scope.locator, MAIN_FORKNUM,
ItemPointerGetBlockNumber(&logical_root));
key.scan_identity = hscan->cr_scope.scan_id;
key.snapshot_identity = snapshot_id;
key.read_scn = snapshot->read_scn;
key.read_epoch = snapshot->read_epoch;
memset(result, 0, sizeof(*result));
if (cluster_bufmgr_cr_copy_v1(&key, result->scratch_page))
{
heap_index_cr_recheck(hscan, snapshot, &admission);
found = heap_hot_r4_search_scratch(&key.tag, &logical_root,
hscan->xs_base.rel, snapshot, result, true);
heap_index_cr_recheck(hscan, snapshot, &admission);
handled = !result->cr_unsupported;
result->kind = found ? HEAP_HOT_SEARCH_OWNED_SCRATCH
: HEAP_HOT_SEARCH_NOT_FOUND;
if (found)
*tid = result->tuple.t_self;
}
/* A miss must acquire current S through the original owner.
* In particular, do not enter holder-moved retry with no holder,
* or force FULL for tuples the live path can already decide. */
}
PG_FINALLY();
{
cluster_snapshot_read_exit_v1(&read_scope);
cluster_semantic_activation_leave(&admission);
}
PG_END_TRY();
}
PG_CATCH();
{
memset(&hscan->cr_scope, 0, sizeof(hscan->cr_scope));
PG_RE_THROW();
}
PG_END_TRY();
return handled;
}
/* Inspect once per publication, not once per tuple hit. Unsupported pages
* stay with the original current/FULL consumer and are never installed. */
static bool
heap_index_cr_page_supported(Page page)
{
PageHeader header = (PageHeader) page;
OffsetNumber offnum;

if (!heap_hot_r4_scratch_page_valid(page))
return false;
for (offnum = FirstOffsetNumber; offnum <= PageGetMaxOffsetNumber(page); offnum++)
{
ItemId lp = PageGetItemId(page, offnum);
Size offset = ItemIdGetOffset(lp);
Size length = ItemIdGetLength(lp);
HeapTupleHeader tuple;

if (!ItemIdIsNormal(lp))
continue;
if (length < SizeofHeapTupleHeader || offset < header->pd_upper
|| offset > header->pd_special || length > header->pd_special - offset)
return false;
tuple = (HeapTupleHeader) ((char *) page + offset);
if (tuple->t_hoff < SizeofHeapTupleHeader || tuple->t_hoff > length)
return false;
if ((tuple->t_infomask & HEAP_XMAX_IS_MULTI) != 0
&& (tuple->t_infomask & HEAP_XMAX_INVALID) == 0
&& (!HEAP_XMAX_IS_LOCKED_ONLY(tuple->t_infomask)
|| (tuple->t_infomask2 & HEAP_HOT_UPDATED) != 0))
return false;
}
return true;
}

/* Called after current SHARE is released, only for an original, revalidated
* FULL result. Neither a selected current tuple nor a failed build qualifies. */
void
heap_index_publish_cr_result(IndexFetchHeapData *hscan, BlockNumber block,
Snapshot snapshot, const HeapHotSearchResult *result)
{
ClusterSemanticAdmissionToken admission;
ClusterSnapshotReadScopeV1 read_scope;
BufferCrKey key;
Buffer volatile reservation = InvalidBuffer;

if (!result->cr_full_page || hscan->cr_scope.scan_id == 0
|| !heap_index_cr_page_supported((Page) result->scratch_page))
return;
if (cluster_semantic_activation_enter(CLUSTER_SEMANTIC_FEATURE_R4_SYNC_CR_V1,
CLUSTER_SEMANTIC_TARGET_SIDE, &admission)
!= CLUSTER_SEMANTIC_ADMISSION_OK)
heap_hot_r4_unknown("read-only FULL publication TARGET admission refused");
cluster_snapshot_read_enter_v1(&read_scope, snapshot);
PG_TRY();
{
heap_index_cr_recheck(hscan, snapshot, &admission);
memset(&key, 0, sizeof(key));
InitBufferTag(&key.tag, &hscan->cr_scope.locator, MAIN_FORKNUM, block);
key.scan_identity = hscan->cr_scope.scan_id;
key.snapshot_identity = hscan->cr_scope.snapshot_id;
key.read_scn = hscan->cr_scope.read_scn;
key.read_epoch = hscan->cr_scope.read_epoch;
reservation = cluster_bufmgr_cr_reserve_v1();
heap_index_cr_recheck(hscan, snapshot, &admission);
(void) cluster_bufmgr_cr_publish_v1(reservation, &key, result->scratch_page);
heap_index_cr_recheck(hscan, snapshot, &admission);
}
PG_FINALLY();
{
if (BufferIsValid(reservation))
ReleaseBuffer(reservation);
cluster_snapshot_read_exit_v1(&read_scope);
cluster_semantic_activation_leave(&admission);
}
PG_END_TRY();
}

#endif

#ifdef USE_CLUSTER_UNIT
Expand Down Expand Up @@ -5342,13 +5598,15 @@ heap_hot_search_buffer_result(ItemPointer tid, Relation relation, Buffer buffer,
build_reason);
if (outcome == HEAP_HOT_R4_CYCLE_FOUND)
{
result->cr_full_page = true;
*tid = result->tuple.t_self;
if (all_dead)
*all_dead = false;
return HEAP_HOT_SEARCH_OWNED_SCRATCH;
}
if (outcome == HEAP_HOT_R4_CYCLE_NOT_FOUND)
{
result->cr_full_page = true;
if (all_dead)
*all_dead = false;
heap_hot_r4_log_miss(relation, buffer, snapshot, &logical_root,
Expand Down
Loading
Loading