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}