mirror of https://github.com/OpenIdentityPlatform/OpenDJ.git

Valery Kharseko
15 hours ago 600df926522919b63bfca2816ef9588b6f1c6e34
opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/PendingChangesTest.java
@@ -16,12 +16,14 @@
package org.opends.server.replication.plugin;
import static org.mockito.Matchers.any;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.testng.Assert.*;
import org.forgerock.opendj.ldap.DN;
import org.mockito.ArgumentCaptor;
import org.opends.server.DirectoryServerTestCase;
import org.opends.server.replication.common.CSN;
@@ -29,23 +31,36 @@
import org.opends.server.replication.protocol.LDAPUpdateMsg;
import org.opends.server.replication.protocol.ReplicaOfflineMsg;
import org.opends.server.replication.protocol.UpdateMsg;
import org.opends.server.replication.service.DSRSShutdownSync;
import org.opends.server.replication.service.ReplicationDomain;
import org.opends.server.types.operation.PluginOperation;
import org.testng.annotations.BeforeClass;
import org.testng.annotations.Test;
/**
 * Tests the bookkeeping a replica does on its own changes: they are published in the order of
 * their CSNs, and the announcement that the replica goes offline is only reported as sent when
 * it really was published.
 * their CSNs, and the announcement that the replica goes offline is made before the message is
 * published, since the shutdown of a collocated replication server waits for that message to be
 * forwarded - and stands only for a message the broker reports as written.
 * <p>
 * These tests need no server: the changes are built by a CSNGenerator, which reads the time
 * service, and the time service is up as soon as its class is loaded.
 * These tests need no server: the changes are built by a CSNGenerator, which needs nothing but
 * a server id.
 */
@SuppressWarnings("javadoc")
@Test(groups = { "precommit", "replication" }, sequential = true)
public class PendingChangesTest extends DirectoryServerTestCase
{
  private static final int SERVER_ID = 42;
  /** A peer replication server the collocated one relays the message to. */
  private static final int RS_ID = 11;
  private static DN baseDN;
  @BeforeClass
  public static void classSetup() throws Exception
  {
    baseDN = DN.valueOf("dc=example,dc=com");
  }
  @Test
  public void replicaOfflineMsgTheBrokerPublishedIsReportedAsSent() throws Exception
@@ -62,6 +77,47 @@
  }
  /**
   * The collocated replication server forwards the message as soon as it is on the wire, so an
   * announcement made after the publish is one the forward found nothing to clear: nothing will
   * ever remove it, and the shutdown waits out its whole grace period for a message which the
   * topology already has.
   * <p>
   * The forward is reported from inside publish(), which is where the message reaches the
   * session, so the race is reproduced rather than waited for.
   */
  @Test
  public void theReplicaOfflineMsgIsAnnouncedBeforeItIsPublished() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    final ReplicationDomain domain = mock(ReplicationDomain.class);
    forwardWhilePublishing(domain, shutdownSync);
    final PendingChanges pendingChanges = newPendingChanges(domain, shutdownSync);
    pendingChanges.putReplicaOfflineMsg();
    assertTrue(shutdownSync.canShutdown(baseDN),
        "the message was forwarded, so nothing must hold the shutdown back any longer");
  }
  /**
   * The announcement of a message the broker took stands until a peer forwards it: the
   * withdrawal is for the message the broker refused, and an announcement taken back after a
   * publish which succeeded would leave the shutdown nothing to wait for.
   */
  @Test
  public void theAnnouncementOfAPublishedMessageStandsUntilItIsForwarded() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    final PendingChanges pendingChanges =
        newPendingChanges(domainWhichPublishes(true), shutdownSync);
    pendingChanges.putReplicaOfflineMsg();
    assertFalse(shutdownSync.canShutdown(baseDN),
        "the message went out and nobody has forwarded it yet, so the shutdown must wait for it");
  }
  /**
   * The broker writes nothing when it has no usable session, when the changes which come before
   * this one still have to be republished by the recovery, or when it is stopped in between - and
   * what was not written must not be reported as sent: the shutdown of a collocated replication
@@ -80,6 +136,27 @@
  }
  /**
   * The announcement is made before the message is published, and the broker may refuse it
   * once it is: an announcement which stayed would be one nobody will ever forward, and the
   * shutdown would wait out its whole grace period for a message which never left. It is
   * therefore withdrawn - and it is a withdrawal, not an announcement which was never made: the
   * shutdown is held back while the broker holds the message.
   */
  @Test
  public void theAnnouncementOfAReplicaOfflineMsgTheBrokerRefusedIsWithdrawn() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    final ReplicationDomain domain = mock(ReplicationDomain.class);
    refuseWhilePublishing(domain, shutdownSync);
    final PendingChanges pendingChanges = newPendingChanges(domain, shutdownSync);
    assertNull(pendingChanges.putReplicaOfflineMsg(), "the broker refused the message");
    assertTrue(shutdownSync.canShutdown(baseDN),
        "the message never reached the wire, so nothing must hold the shutdown back");
  }
  /**
   * The message carries the newest CSN of the replica, so a change which is still in flight
   * holds it back, and the broker is never even asked to publish it.
   */
@@ -117,6 +194,32 @@
  }
  /**
   * The announcement follows the publication rather than the queueing, so a message which a
   * change in flight holds back is not announced: neither while it waits, nor when the change
   * which held it back completes and the message is given up on. Announcing it either time
   * would leave the shutdown waiting out its whole grace period for a forward which cannot
   * happen.
   */
  @Test
  public void theReplicaOfflineMsgHeldBackByAChangeInFlightIsNeverAnnounced() throws Exception
  {
    final DSRSShutdownSync shutdownSync = new DSRSShutdownSync();
    final ReplicationDomain domain = domainWhichPublishes(true);
    final PendingChanges pendingChanges = newPendingChanges(domain, shutdownSync);
    final CSN inFlight = pendingChanges.putLocalOperation(newLocalOperation());
    pendingChanges.putReplicaOfflineMsg();
    assertTrue(shutdownSync.canShutdown(baseDN),
        "the message is still queued behind a change in flight, and nothing was announced");
    pendingChanges.commitAndPushCommittedChanges(inFlight, mock(LDAPUpdateMsg.class));
    assertTrue(shutdownSync.canShutdown(baseDN),
        "the message was given up on with the change which held it back, and never announced");
  }
  /**
   * A change the broker refused leaves the pending changes all the same: the replica has done
   * it, its ServerState says so, and it is by finding that state ahead of the one its
   * replication server reports that the next session republishes the change from the historical
@@ -137,9 +240,55 @@
    assertEquals(pendingChanges.size(), 0, "and is not queued for a second attempt");
  }
  /**
   * Reports the forward of the message from within the publish which puts it on the wire, after
   * checking from there that the announcement is already in place, and publishes it. The forward
   * is the one of a peer the message was never recorded as queued for, which is what a forward
   * racing the announcement looks like.
   */
  private void forwardWhilePublishing(
      final ReplicationDomain domain, final DSRSShutdownSync shutdownSync)
  {
    doAnswer(invocation -> {
      final UpdateMsg msg = (UpdateMsg) invocation.getArguments()[0];
      if (msg instanceof ReplicaOfflineMsg)
      {
        assertFalse(shutdownSync.canShutdown(baseDN),
            "the message must be announced before it is published");
        shutdownSync.replicaOfflineMsgForwarded(baseDN, msg.getCSN(), RS_ID);
      }
      return true;
    }).when(domain).publish(any(UpdateMsg.class));
  }
  /**
   * Refuses to publish the message, the way a broker with no usable session does, after checking
   * from within the publish that the announcement is already in place.
   */
  private void refuseWhilePublishing(
      final ReplicationDomain domain, final DSRSShutdownSync shutdownSync)
  {
    doAnswer(invocation -> {
      assertFalse(shutdownSync.canShutdown(baseDN),
          "the message must be announced before it is published");
      return false;
    }).when(domain).publish(any(UpdateMsg.class));
  }
  private PendingChanges newPendingChanges(ReplicationDomain domain)
  {
    return new PendingChanges(new CSNGenerator(SERVER_ID, 0), domain);
    return newPendingChanges(domain, new DSRSShutdownSync());
  }
  /**
   * The pending changes of a replica whose domain announces itself through the shutdown sync,
   * with the announcer LDAPReplicationDomain hands its own pending changes.
   */
  private PendingChanges newPendingChanges(
      final ReplicationDomain domain, final DSRSShutdownSync shutdownSync)
  {
    return new PendingChanges(new CSNGenerator(SERVER_ID, 0), domain,
        new ShutdownSyncAnnouncer(shutdownSync, baseDN));
  }
  /** A domain whose broker accepts, or refuses, whatever it is given to publish. */