@@ -100,6 +100,19 @@ public final class SegmentRing implements QuietCloseable {
100100 // Logical head of sealedSegments. Head removal nulls one entry and advances
101101 // this index; occasional compaction bounds unused prefix slots.
102102 private int sealedHead ;
103+ // Frontier index into sealedSegments: every sealed segment in
104+ // [sealedHead, firstNonDurableSealed) has been proven durable by an earlier
105+ // periodic pass and never needs re-scanning -- publishedCursor is frozen at
106+ // seal and durableCursor only advances, so durability never regresses. The
107+ // periodic sync (copyPendingSyncSegments) skips this proven-durable prefix,
108+ // keeping the steady-state copy-under-monitor O(1) instead of O(live-sealed)
109+ // work that would otherwise grow with a producer-outpaces-drain backlog.
110+ // Rotation seals only already-durable predecessors (the
111+ // requestSyncBeforeRotation gate), so the sole source of non-durable sealed
112+ // segments is a crash-recovery resume; those are covered because the
113+ // frontier starts at 0. Maintained entirely under this monitor, in the same
114+ // coordinate space as sealedHead (shifted by compaction, reset on clear).
115+ private int firstNonDurableSealed ;
103116 // High-water byte offset within the active segment at which we proactively
104117 // ask the segment manager to provision a spare (if one isn't already
105118 // installed). Computed once as 3/4 of segment capacity -- leaves the manager
@@ -668,29 +681,23 @@ private static RetainedSegmentMembership newDefaultMembership(
668681 }
669682 if (DEFAULT_MEMBERSHIP_MODE == RetainedSegmentMembershipMode .LINEAR ) {
670683 if (observer == null ) {
671- return new RetainedSegmentMembership () {
672- @ Override
673- public boolean contains (MmapSegment segment ) {
674- for (int i = 0 , n = chain .size (); i < n ; i ++) {
675- if (chain .get (i ) == segment ) {
676- return true ;
677- }
678- }
679- return false ;
680- }
681- };
682- }
683- return new RetainedSegmentMembership () {
684- @ Override
685- public boolean contains (MmapSegment segment ) {
684+ return segment -> {
686685 for (int i = 0 , n = chain .size (); i < n ; i ++) {
687- observer .onMembershipOperation ();
688686 if (chain .get (i ) == segment ) {
689687 return true ;
690688 }
691689 }
692690 return false ;
691+ };
692+ }
693+ return segment -> {
694+ for (int i = 0 , n = chain .size (); i < n ; i ++) {
695+ observer .onMembershipOperation ();
696+ if (chain .get (i ) == segment ) {
697+ return true ;
698+ }
693699 }
700+ return false ;
694701 };
695702 }
696703
@@ -699,19 +706,11 @@ public boolean contains(MmapSegment segment) {
699706 retained .put (chain .get (i ), Boolean .TRUE );
700707 }
701708 if (observer == null ) {
702- return new RetainedSegmentMembership () {
703- @ Override
704- public boolean contains (MmapSegment segment ) {
705- return retained .containsKey (segment );
706- }
707- };
709+ return retained ::containsKey ;
708710 }
709- return new RetainedSegmentMembership () {
710- @ Override
711- public boolean contains (MmapSegment segment ) {
712- observer .onMembershipOperation ();
713- return retained .containsKey (segment );
714- }
711+ return segment -> {
712+ observer .onMembershipOperation ();
713+ return retained .containsKey (segment );
715714 };
716715 }
717716
@@ -1070,6 +1069,46 @@ synchronized void copyLiveSegmentsForSync(ObjList<MmapSegment> target) {
10701069 }
10711070 }
10721071
1072+ /**
1073+ * Copies the live segments that may still need a durability barrier: every
1074+ * sealed segment from the {@link #firstNonDurableSealed} frontier onward,
1075+ * plus the active segment. First advances the frontier past any sealed
1076+ * segments an earlier pass (or rotation's pre-seal barrier) has since made
1077+ * durable. Used by the periodic sync path in place of
1078+ * {@link #copyLiveSegmentsForSync}: the proven-durable prefix would
1079+ * otherwise be re-copied under this monitor and re-scanned every tick as
1080+ * no-op {@link MmapSegment#syncPublished()} early-returns -- O(live-sealed)
1081+ * work that grows without bound under a producer-outpaces-drain backlog.
1082+ * The frontier is a conservative lower bound (it only ever advances past
1083+ * segments observed durable, and durability never regresses), so this can
1084+ * never skip a segment that still needs a barrier.
1085+ */
1086+ synchronized void copyPendingSyncSegments (ObjList <MmapSegment > target ) {
1087+ target .clear ();
1088+ // Invariant maintained by every mutation site (rotation append, trim's
1089+ // removeSealedHead, compaction shift, close). A frontier that drifted
1090+ // above size would silently skip un-fsynced segments, so guard it in
1091+ // tests; the clamp below keeps production safe if it is ever violated.
1092+ assert firstNonDurableSealed >= sealedHead && firstNonDurableSealed <= sealedSegments .size ()
1093+ : "durability frontier out of range: firstNonDurableSealed=" + firstNonDurableSealed
1094+ + " sealedHead=" + sealedHead + " size=" + sealedSegments .size ();
1095+ int i = firstNonDurableSealed ;
1096+ if (i < sealedHead ) {
1097+ i = sealedHead ;
1098+ }
1099+ int n = sealedSegments .size ();
1100+ while (i < n && sealedSegments .get (i ).isPublishedDurable ()) {
1101+ i ++;
1102+ }
1103+ firstNonDurableSealed = i ;
1104+ for (; i < n ; i ++) {
1105+ target .add (sealedSegments .get (i ));
1106+ }
1107+ if (active != null ) {
1108+ target .add (active );
1109+ }
1110+ }
1111+
10731112 void enablePeriodicSync () {
10741113 periodicSyncEnabled = true ;
10751114 syncRequested = true ;
@@ -1115,6 +1154,7 @@ public synchronized void close() {
11151154 }
11161155 sealedSegments .clear ();
11171156 sealedHead = 0 ;
1157+ firstNonDurableSealed = 0 ;
11181158 for (int i = 0 , n = pendingTrims .size (); i < n ; i ++) {
11191159 pendingTrims .getQuick (i ).close ();
11201160 }
@@ -1211,7 +1251,6 @@ public synchronized MmapSegment firstTrimmable() {
12111251 return lastSeq <= ackedFsn ? segment : null ;
12121252 }
12131253
1214- /** Active segment -- exposed for the I/O thread's "send next batch" path. */
12151254 /**
12161255 * Walks every published frame in the ring (sealed segments plus the active
12171256 * segment) and returns the FSN of the LAST frame whose payload does NOT
@@ -1419,6 +1458,20 @@ public synchronized int getPendingTrimCount() {
14191458 return pendingTrims .size ();
14201459 }
14211460
1461+ /**
1462+ * Number of live segments the periodic path would barrier this tick,
1463+ * advancing the durability frontier exactly as a real tick does. In the
1464+ * steady state (every sealed segment proven durable) this collapses to 1
1465+ * -- the active segment -- proving the proven-durable sealed prefix is no
1466+ * longer copied/scanned under the monitor.
1467+ */
1468+ @ TestOnly
1469+ public synchronized int pendingSyncSegmentCountForTest () {
1470+ ObjList <MmapSegment > scratch = new ObjList <>();
1471+ copyPendingSyncSegments (scratch );
1472+ return scratch .size ();
1473+ }
1474+
14221475 @ TestOnly
14231476 public synchronized MmapSegment pinSegmentContainingForTest (long fsn ) {
14241477 return pinSegmentContaining (fsn );
@@ -1618,16 +1671,30 @@ private void compactSealedSegments() {
16181671 int liveCount = sealedSegments .size () - sealedHead ;
16191672 trimMovedReferences += liveCount ;
16201673 sealedSegments .remove (0 , sealedHead - 1 );
1674+ // The durable-prefix frontier lives in the same index space as the
1675+ // entries we just shifted down by sealedHead; move it with them.
1676+ // The invariant firstNonDurableSealed >= sealedHead keeps this >= 0.
1677+ firstNonDurableSealed -= sealedHead ;
1678+ if (firstNonDurableSealed < 0 ) {
1679+ firstNonDurableSealed = 0 ;
1680+ }
16211681 sealedHead = 0 ;
16221682 }
16231683 }
16241684
16251685 private void removeSealedHead () {
16261686 sealedSegments .setQuick (sealedHead ++, null );
1687+ // If the removed head WAS the frontier (a non-durable but already-ACKed
1688+ // recovery-resumed segment can be trimmed before its first barrier),
1689+ // keep the frontier at or ahead of the live head.
1690+ if (firstNonDurableSealed < sealedHead ) {
1691+ firstNonDurableSealed = sealedHead ;
1692+ }
16271693 int size = sealedSegments .size ();
16281694 if (sealedHead == size ) {
16291695 sealedSegments .clear ();
16301696 sealedHead = 0 ;
1697+ firstNonDurableSealed = 0 ;
16311698 } else if (sealedHead >= 64 && sealedHead >= size - sealedHead ) {
16321699 compactSealedSegments ();
16331700 }
0 commit comments