001/* 002 * Licensed to the Apache Software Foundation (ASF) under one 003 * or more contributor license agreements. See the NOTICE file 004 * distributed with this work for additional information 005 * regarding copyright ownership. The ASF licenses this file 006 * to you under the Apache License, Version 2.0 (the 007 * "License"); you may not use this file except in compliance 008 * with the License. You may obtain a copy of the License at 009 * 010 * http://www.apache.org/licenses/LICENSE-2.0 011 * 012 * Unless required by applicable law or agreed to in writing, software 013 * distributed under the License is distributed on an "AS IS" BASIS, 014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. 015 * See the License for the specific language governing permissions and 016 * limitations under the License. 017 */ 018package org.apache.hadoop.hbase.master; 019 020import static org.junit.jupiter.api.Assertions.assertEquals; 021import static org.junit.jupiter.api.Assertions.assertTrue; 022import static org.mockito.Mockito.mock; 023import static org.mockito.Mockito.when; 024 025import java.util.Collections; 026import java.util.concurrent.CyclicBarrier; 027import java.util.concurrent.ExecutorService; 028import java.util.concurrent.Executors; 029import java.util.concurrent.Future; 030import java.util.concurrent.TimeUnit; 031import org.apache.hadoop.conf.Configuration; 032import org.apache.hadoop.hbase.HBaseConfiguration; 033import org.apache.hadoop.hbase.HConstants; 034import org.apache.hadoop.hbase.RegionMetricsBuilder; 035import org.apache.hadoop.hbase.ServerMetrics; 036import org.apache.hadoop.hbase.ServerMetricsBuilder; 037import org.apache.hadoop.hbase.ServerName; 038import org.apache.hadoop.hbase.TableName; 039import org.apache.hadoop.hbase.client.RegionInfo; 040import org.apache.hadoop.hbase.client.RegionInfoBuilder; 041import org.apache.hadoop.hbase.master.assignment.AssignmentManager; 042import org.apache.hadoop.hbase.master.assignment.RegionStates; 043import org.apache.hadoop.hbase.testclassification.MasterTests; 044import org.apache.hadoop.hbase.testclassification.SmallTests; 045import org.apache.hadoop.hbase.util.Bytes; 046import org.junit.jupiter.api.BeforeEach; 047import org.junit.jupiter.api.Tag; 048import org.junit.jupiter.api.Test; 049 050@Tag(MasterTests.TAG) 051@Tag(SmallTests.TAG) 052public class TestServerManager { 053 054 private static final class DummyMasterServices extends MockNoopMasterServices { 055 private final AssignmentManager am; 056 057 DummyMasterServices(Configuration conf) { 058 super(conf); 059 am = mock(AssignmentManager.class); 060 RegionStates rss = mock(RegionStates.class); 061 when(am.getRegionStates()).thenReturn(rss); 062 } 063 064 @Override 065 public AssignmentManager getAssignmentManager() { 066 return am; 067 } 068 } 069 070 private ServerManager sm; 071 private RegionInfo region; 072 073 @BeforeEach 074 public void setUp() { 075 Configuration conf = HBaseConfiguration.create(); 076 sm = new ServerManager(new DummyMasterServices(conf), new DummyRegionServerList()); 077 region = RegionInfoBuilder.newBuilder(TableName.valueOf("t")).build(); 078 } 079 080 private long lastFlushed(RegionInfo ri) { 081 return sm.getLastFlushedSequenceId(ri.getEncodedNameAsBytes()).getLastFlushedSequenceId(); 082 } 083 084 @Test 085 public void testReportRegionOpenSeedsFlushedSequenceId() { 086 assertEquals(HConstants.NO_SEQNUM, lastFlushed(region)); 087 sm.reportRegionOpen(region, 42L); 088 assertEquals(42L, lastFlushed(region)); 089 } 090 091 @Test 092 public void testReportRegionOpenDoesNotRegressExistingValue() { 093 sm.reportRegionOpen(region, 100L); 094 // A later OPEN carrying a smaller openSeqNum (e.g. after a restart replayed less) must not 095 // clobber a higher watermark already seeded here or supplied by a heartbeat. 096 sm.reportRegionOpen(region, 50L); 097 assertEquals(100L, lastFlushed(region)); 098 } 099 100 @Test 101 public void testReportRegionOpenIgnoresNoSeqNum() { 102 sm.reportRegionOpen(region, HConstants.NO_SEQNUM); 103 assertEquals(HConstants.NO_SEQNUM, lastFlushed(region)); 104 } 105 106 @Test 107 public void testReportRegionOpenIgnoresNegativeSeqNum() { 108 sm.reportRegionOpen(region, -5L); 109 assertEquals(HConstants.NO_SEQNUM, lastFlushed(region)); 110 } 111 112 /** 113 * HBASE-30335: the OPEN-time seed ({@link ServerManager#reportRegionOpen}, which uses 114 * {@code merge(Math::max)}) and the heartbeat handler ({@link ServerManager#regionServerReport}) 115 * both write {@code flushedSequenceIdByRegion}. A stale in-flight heartbeat from the 116 * soon-to-be-dead source RS carries a lower {@code completedSequenceId}. Whatever the 117 * interleaving, the watermark must never regress below the seed: if the heartbeat lands first the 118 * seed lifts it to {@code openSeqNum}; if it lands after, the heartbeat's read-modify-write must 119 * refuse to lower it. Only a non-atomic check-then-put in the heartbeat path (the pre-fix bug) 120 * could let the stale value clobber the seed. This drives both writers concurrently over many 121 * rounds to catch that race. 122 */ 123 @Test 124 public void testConcurrentStaleHeartbeatDoesNotClobberOpenSeed() throws Exception { 125 final long seedSeqId = 200L; 126 final long staleSeqId = 100L; 127 // Register the server so regionServerReport takes the heartbeat (updateLastFlushedSequenceIds) 128 // path instead of the new-server-registration path. 129 ServerName sn = ServerName.valueOf("rs.example.org", 16020, 1L); 130 sm.recordNewServerWithLock(sn, ServerMetricsBuilder.of(sn)); 131 132 ExecutorService pool = Executors.newFixedThreadPool(2); 133 try { 134 for (int i = 0; i < 500; i++) { 135 // Fresh region per round so no round is masked by a prior round's watermark. 136 final RegionInfo ri = RegionInfoBuilder.newBuilder(TableName.valueOf("concurrentSeed")) 137 .setStartKey(Bytes.toBytes(i)).setEndKey(Bytes.toBytes(i + 1)).build(); 138 final ServerMetrics staleReport = ServerMetricsBuilder.newBuilder(sn) 139 .setRegionMetrics(Collections.singletonList(RegionMetricsBuilder 140 .newBuilder(ri.getRegionName()).setCompletedSequenceId(staleSeqId).build())) 141 .build(); 142 final CyclicBarrier barrier = new CyclicBarrier(2); 143 Future<?> seedTask = pool.submit(() -> { 144 await(barrier); 145 sm.reportRegionOpen(ri, seedSeqId); 146 }); 147 Future<?> heartbeatTask = pool.submit(() -> { 148 await(barrier); 149 try { 150 sm.regionServerReport(sn, staleReport); 151 } catch (Exception e) { 152 throw new RuntimeException(e); 153 } 154 }); 155 seedTask.get(30, TimeUnit.SECONDS); 156 heartbeatTask.get(30, TimeUnit.SECONDS); 157 assertTrue(lastFlushed(ri) >= seedSeqId, "round " + i + ": watermark regressed to " 158 + lastFlushed(ri) + ", stale heartbeat clobbered the openSeqNum seed"); 159 } 160 } finally { 161 pool.shutdownNow(); 162 } 163 } 164 165 private static void await(CyclicBarrier barrier) { 166 try { 167 barrier.await(30, TimeUnit.SECONDS); 168 } catch (Exception e) { 169 throw new RuntimeException(e); 170 } 171 } 172}