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

Valery Kharseko
2 days ago 80481f756d71bd58b4bda627758dcd774e9d5dcd
opendj-server-legacy/src/test/java/org/opends/server/replication/UpdateOperationTest.java
@@ -95,6 +95,7 @@
import org.opends.server.types.Operation;
import org.opends.server.types.OperationType;
import org.opends.server.types.RawModification;
import org.opends.server.util.StaticUtils;
import org.opends.server.util.TestTimer;
import org.opends.server.util.TestTimer.CallableVoid;
import org.opends.server.util.TimeThread;
@@ -3583,6 +3584,508 @@
    }
  }
  /**
   * Test case for [Issue 954]: a change parked as waiting for another one is given back
   * when the replay which parked it is unwound, and the session is restarted for it.
   * <p>
   * A change which waits for another one is parked and stays owned by the replay thread
   * which parked it, while that thread goes on to the changes which follow: it is handed
   * out again by {@code getNextUpdate()}, which every replay loop of this domain runs once
   * it is done, so it is replayed by whichever thread clears the change it was waiting for.
   * A replay which is unwound leaves the thread which parked it without that road - it
   * takes the next delivery off the shared queue instead, and never comes back to the
   * change it parked - and every redelivery of a change a replay thread owns is refused as
   * a duplicate. On a domain which then goes quiet that change is where this replica's
   * ServerState, and every change behind it from every master, stops.
   * <p>
   * The replay which is unwound here is one which applied its change: the change it was
   * replaying is committed and owns nothing anymore by the time the give-back on the way
   * out of {@code replay()} runs, so the restart that give-back asks for is the one thing
   * which has the parked change delivered again - a replay which failed would have asked
   * for the same restart on the road of its own change. The parked change travels the
   * replication server, and nothing but a new session brings it back.
   */
  @Test
  public void aChangeParkedByAnUnwoundReplayIsDeliveredAgain() throws Exception
  {
    testSetUp("aChangeParkedByAnUnwoundReplayIsDeliveredAgain");
    logger.error(LocalizableMessage.raw(
        "Starting replication test : aChangeParkedByAnUnwoundReplayIsDeliveredAgain"));
    final DN waitedOn = addEntryForChange("user.954.1");
    final String waitedOnUUID = getEntry(waitedOn, 1, true).parseAttribute("entryuuid").asString();
    final DN other = addEntryForChange("user.954.2");
    final String otherUUID = getEntry(other, 1, true).parseAttribute("entryuuid").asString();
    final LDAPReplicationDomain domain = MultimasterReplication.findDomain(baseDN, null);
    domain.resetUnreplayedChangeAlertThrottle();
    final long inProgress = getMonitorAttrValue(baseDN, "changes-in-progress-size");
    assertEquals(getMonitorAttrValue(baseDN, "dependent-changes-size"), 0,
        "no change of this domain is waiting for another one when this test starts");
    // The replica the changes come from: one of its own, so that the CSNs of this test are
    // never covered by the ServerState another one left behind.
    final int serverId = 30;
    final CSNGenerator gen = new CSNGenerator(serverId, TimeThread.getTime());
    final CSN failing = gen.newCSN();
    final CSN parked = gen.newCSN();
    final CSN unwound = gen.newCSN();
    final List<Modification> failingMods = generatemods("description", "the replay of this change fails");
    final String parkedDescription = "the change which was parked as a dependency";
    final List<Modification> parkedMods = generatemods("description", parkedDescription);
    final List<Modification> unwoundMods =
        generatemods("description", "the replay of this change is unwound once it is applied");
    /*
     * The change whose replay fails is the barrier the parked change waits behind, so its
     * budget must not be spent while this test is setting up: it is shortened once the
     * change which was parked is back, and put back in the finally below.
     */
    setReplayGiveUpDelay("unlimited");
    try
    {
      /*
       * One replay thread, so that the change which is parked and the replay which is
       * unwound after it are the same thread's: a parked change is left owned by the thread
       * which parked it, and this is about a thread which does not come back to it.
       */
      setNumUpdateReplayThreads(1);
      try
      {
        /*
         * The change which is parked is published to the replication server rather than
         * handed to the domain: a change which travelled a session is one the replication
         * server sends again over the next session of this domain, and over nothing else -
         * which is what the restart the give-back asks for has to be pinned against.
         */
        final ReplicationBroker broker =
            openReplicationSession(baseDN, serverId, 100, replServerPort, 1000);
        try
        {
          /*
           * A change whose replay failed stays listed and uncommitted - it is what holds this
           * domain's ServerState back - and stays among the changes the newer ones are checked
           * against, so a change which follows it on the same entry has to wait for it. It has
           * to be a change which is asked for again rather than stepped over: one whose
           * operation is built and then refused, since #928 has a modify whose entry DN does
           * not parse reported once and recorded as replayed.
           */
          deliverUntilMonitorReaches(domain, "changes-in-progress-size", inProgress + 1,
              () -> new ModifyMsgWhoseOperationRefusesAControl(failing, waitedOn, failingMods, waitedOnUUID),
              "the change whose replay fails must be listed as one which is not in the data");
          // The change which is parked as waiting for it by the replay thread it was given to.
          broker.publish(new ModifyMsg(parked, waitedOn, parkedMods, waitedOnUUID));
          assertMonitorAttrValueEventually(baseDN, "dependent-changes-size", 1,
              "a change which waits for one that is not in the data must be parked");
          /*
           * How many restarts in a row the session has been through, the barrier's among
           * them: the restart the give-back asks for on the road out of a JVM which has run
           * out of memory is owed no backoff, and a restart which waits the backoff out is
           * the one road which moves this count. It can not move otherwise between here and
           * the reading below: a replay which made it puts the count back to zero only once
           * nothing is failing anymore, and the barrier keeps failing until it is given up.
           */
          final int restarts = domain.getConsecutiveSessionRestarts();
          /*
           * The replay which is unwound while that same thread still holds the parked change.
           * It applies its change, and the ack of its delivery is where the JVM runs out of
           * memory: that is the one throw from the ack which is not caught, past every catch
           * the replay itself has, and it leaves replay() with the change this thread was
           * replaying committed and owned by nobody. What the give-back on the way out finds
           * to give back is the parked change alone, and the restart it asks for is the one
           * road which has that change delivered again: the road a failed replay takes to
           * ask for its own change is not run. An Error met replaying a change does not get
           * here, and neither does any other throw from the ack: both are reported and the
           * replay carries on to the road the change itself decided (issue #922).
           *
           * It is made on another entry, so that it is replayed rather than parked in its
           * turn, and it is delivered until the parked change has been given back: a delivery
           * of a change which is being replayed, or was applied, is refused as the duplicate
           * it is.
           */
          deliverUntilMonitorReaches(domain, "dependent-changes-size", 0,
              () -> new ModifyMsgWhoseAckRunsOutOfMemoryOnceApplied(unwound, other, unwoundMods, otherUUID),
              "a change parked by a replay which was unwound must be given back");
          /*
           * The thread which met the error is gone - an OutOfMemoryError ends the replay
           * thread which met it (issue #923) - and it was the whole pool, so the delivery the
           * restarted session brings waits in the replay queue: the pool is brought back for
           * it, at any number now that the thread it had to be the same as is over. The change
           * is parked again by the thread which takes it, the barrier being still there: that
           * it is parked at all is what says a new session delivered it.
           */
          setNumUpdateReplayThreads(2);
          assertMonitorAttrValueEventually(baseDN, "dependent-changes-size", 1,
              "the change which was given back must be delivered again by the session which"
                  + " was restarted for it");
          assertEquals(domain.getConsecutiveSessionRestarts(), restarts,
              "a thread an OutOfMemoryError is ending must not wait the backoff out before"
                  + " it restarts the session for the changes it gave back");
          /*
           * The barrier is lifted, which lets the ServerState past it and hands the parked
           * change to the thread which cleared it.
           */
          giveUpOn(domain, failing, waitedOn, failingMods, waitedOnUUID);
          checkEntryHasAttributeValue(waitedOn, "description", parkedDescription, 30,
              "the change which was parked must be applied by the delivery which took it over");
        }
        finally
        {
          broker.stop();
        }
        assertEquals(resetNumUpdateReplayThreads(), 0,
            "the number of replay threads could not be put back");
      }
      finally
      {
        /*
         * Best-effort on the way out of a red, with no assertion to replace the failure
         * which is being reported: on the normal road the number is put back, and checked,
         * above.
         */
        resetNumUpdateReplayThreads();
      }
    }
    finally
    {
      try
      {
        /*
         * A red before the barrier was given up would leave it listed for the rest of this
         * class: it never travelled the replication server, so no session restart brings a
         * delivery which could give it up, and a change which is failing keeps every later
         * session restart at its backoff and the ServerState behind it. The failure which
         * is being reported is the one that matters, so this does not replace it.
         */
        if (!domain.getServerState().cover(failing))
        {
          giveUpOn(domain, failing, waitedOn, failingMods, waitedOnUUID);
        }
      }
      catch (Throwable cleanupFailure)
      {
        logger.error(LocalizableMessage.raw(
            "the barrier of aChangeParkedByAnUnwoundReplayIsDeliveredAgain could not be given"
                + " up on the way out: %s", StaticUtils.stackTraceToSingleLineString(cleanupFailure)));
      }
      finally
      {
        resetReplayGiveUpDelay();
      }
    }
  }
  /**
   * Gives up on the change no delivery can replay: the budget is shortened to nothing, so
   * that the change is given up on as soon as one more delivery of it fails. Nothing sends
   * that change again - it never travelled a session - so its deliveries are made here.
   */
  private void giveUpOn(final LDAPReplicationDomain domain, final CSN failing, final DN dn,
      final List<Modification> mods, final String entryUUID) throws Exception
  {
    setReplayGiveUpDelay("0ms");
    TestTimer timer = new TestTimer.Builder()
      .maxSleep(120, SECONDS)
      .sleepTimes(500, MILLISECONDS)
      .toTimer();
    timer.repeatUntilSuccess(new CallableVoid()
    {
      @Override
      public void call() throws Exception
      {
        if (!domain.getServerState().cover(failing))
        {
          domain.processUpdate(new ModifyMsgWhoseOperationRefusesAControl(failing, dn, mods, entryUUID));
        }
        assertTrue(domain.getServerState().cover(failing),
            "the change no delivery can replay must be given up on");
      }
    });
  }
  /**
   * Test case for [Issue 986]: a change parked as waiting for another one is given back
   * when the replay thread which parked it is stopped with the pool.
   * <p>
   * A parked change stays owned by the replay thread which parked it while that thread
   * goes back to the pool and takes the changes which follow: {@code getNextUpdate()} is
   * what hands it out again, to whichever replay thread clears the change it was waiting
   * for. Changing the number of replay threads stops the whole pool and creates another
   * one, so a thread which parked a change and went back to the queue is joined while it
   * is idle, and it would end still recorded as the owner of that change - a thread which
   * does not exist anymore, while every redelivery of a change a replay thread owns is
   * refused as a duplicate. On a domain which then goes quiet that change is where this
   * replica's ServerState, and every change behind it from every master, stops.
   */
  @Test
  public void aChangeParkedByAThreadThePoolStoppedIsDeliveredAgain() throws Exception
  {
    testSetUp("aChangeParkedByAThreadThePoolStoppedIsDeliveredAgain");
    logger.error(LocalizableMessage.raw(
        "Starting replication test : aChangeParkedByAThreadThePoolStoppedIsDeliveredAgain"));
    final DN waitedOn = addEntryForChange("user.986.1");
    final String waitedOnUUID = getEntry(waitedOn, 1, true).parseAttribute("entryuuid").asString();
    final DN other = addEntryForChange("user.986.2");
    final String otherUUID = getEntry(other, 1, true).parseAttribute("entryuuid").asString();
    final LDAPReplicationDomain domain = MultimasterReplication.findDomain(baseDN, null);
    domain.resetUnreplayedChangeAlertThrottle();
    final long inProgress = getMonitorAttrValue(baseDN, "changes-in-progress-size");
    assertEquals(getMonitorAttrValue(baseDN, "dependent-changes-size"), 0,
        "no change of this domain is waiting for another one when this test starts");
    final CSNGenerator gen = new CSNGenerator(26, TimeThread.getTime());
    final CSN failing = gen.newCSN();
    final CSN parked = gen.newCSN();
    final CSN applied = gen.newCSN();
    final List<Modification> failingMods = generatemods("description", "the replay of this change fails");
    final String parkedDescription = "the change which was parked by a thread the pool stopped";
    final List<Modification> parkedMods = generatemods("description", parkedDescription);
    final String appliedDescription = "the change which is applied before the pool is stopped";
    final List<Modification> appliedMods = generatemods("description", appliedDescription);
    /*
     * One replay thread, so that the change which is parked is parked by the thread the
     * pool then stops: a parked change is left owned by the thread which parked it, and
     * this is about a thread which is not there anymore to be given it back.
     */
    setNumUpdateReplayThreads(1);
    try
    {
      /*
       * A change whose replay failed stays listed and uncommitted - it is what holds this
       * domain's ServerState back - and stays among the changes the newer ones are checked
       * against, so a change which follows it on the same entry has to wait for it.
       */
      deliverUntilMonitorReaches(domain, "changes-in-progress-size", inProgress + 1,
          () -> new ModifyMsgWhoseOperationRefusesAControl(failing, waitedOn, failingMods, waitedOnUUID),
          "the change whose replay fails must stay listed as one which is not in the data");
      /*
       * The change which is parked as waiting for it. The replay thread it was given to
       * parks it and goes back to the pool: nothing is waiting for it there, so it is idle
       * and it still owns the change it parked.
       */
      deliverUntilMonitorReaches(domain, "dependent-changes-size", 1,
          () -> new ModifyMsg(parked, waitedOn, parkedMods, waitedOnUUID),
          "a change which waits for one that is not in the data must be parked");
      /*
       * A change which is applied, on another entry so that it is replayed rather than
       * parked in its turn. It is there for the count of the deliveries this session took
       * off it, replayed-updates: a session which is started counts from zero, so that
       * count going back to zero below is what says the session was restarted - and it is
       * not left to the give-back alone to put the count above zero before the pool is
       * stopped. The restart the change whose replay failed asked for is over by now: the
       * thread which ran it is the one which parked the change above.
       */
      domain.processUpdate(new ModifyMsg(applied, other, appliedMods, otherUUID));
      checkEntryHasAttributeValue(other, "description", appliedDescription, 30,
          "the change made on the other entry must be applied");
      // Counted once its ack is out, which is after the change is in the data.
      final TestTimer counted = new TestTimer.Builder()
        .maxSleep(30, SECONDS)
        .sleepTimes(200, MILLISECONDS)
        .toTimer();
      counted.repeatUntilSuccess(new CallableVoid()
      {
        @Override
        public void call() throws Exception
        {
          assertTrue(getMonitorAttrValue(baseDN, "replayed-updates") > 0,
              "the session must have counted the deliveries it took off before the pool is stopped");
        }
      });
      /*
       * How many restarts in a row the session has been through, the barrier's among them:
       * the restart the give-back of a stopped thread asks for is owed no backoff - what
       * went away is a replay thread, not the backend - and a restart which waits the
       * backoff out is the one road which moves this count. It can not move otherwise
       * between here and the reading below: a replay which made it puts the count back to
       * zero only once nothing is failing anymore, and the barrier keeps failing until it
       * is given up, below. The restart the barrier asked for is over: the thread which
       * ran it is the one which parked the change and applied the other, above.
       */
      final int restarts = domain.getConsecutiveSessionRestarts();
      /*
       * The pool is stopped and created again, the way an administrator changing the
       * number of replay threads has it: the thread which parked the change is joined
       * where it waits for the next delivery, and the change it owns is nobody's.
       */
      setNumUpdateReplayThreads(2);
      assertMonitorAttrValueEventually(baseDN, "dependent-changes-size", 0,
          "a change parked by a replay thread the pool stopped must be given back");
      /*
       * The give-back alone leaves the count where it was, plus the change it handed back:
       * only a session which is started counts from zero, and the changes of this test
       * never travelled a session, so nothing is delivered over the one which is started
       * before the deliveries made below.
       */
      assertMonitorAttrValueEventually(baseDN, "replayed-updates", 0,
          "the session must be restarted for the changes which were given back");
      assertEquals(domain.getConsecutiveSessionRestarts(), restarts,
          "the restart run for the changes a stopped thread gave back must not wait the"
              + " backoff out: what went away is a replay thread, not the backend");
      /*
       * Nothing sends these changes again - they never travelled a session - so the
       * deliveries which take over from the ones the stopped pool left behind are made
       * here. The change no delivery can replay is given up on, which lets the ServerState
       * past it, and the change which was parked behind it is applied.
       */
      setReplayGiveUpDelay(TEST_GIVE_UP_DELAY);
      TestTimer timer = new TestTimer.Builder()
        .maxSleep(120, SECONDS)
        .sleepTimes(500, MILLISECONDS)
        .toTimer();
      timer.repeatUntilSuccess(new CallableVoid()
      {
        @Override
        public void call() throws Exception
        {
          final ServerState state = domain.getServerState();
          if (!state.cover(failing))
          {
            domain.processUpdate(
                new ModifyMsgWhoseOperationRefusesAControl(failing, waitedOn, failingMods, waitedOnUUID));
          }
          if (!state.cover(parked))
          {
            domain.processUpdate(new ModifyMsg(parked, waitedOn, parkedMods, waitedOnUUID));
          }
          assertTrue(state.cover(parked),
              "the change which was given back must be replayed by the delivery which takes it over");
        }
      });
      checkEntryHasAttributeValue(waitedOn, "description", parkedDescription, 30,
          "the change which was parked must be applied by the delivery which took it over");
      assertEquals(resetNumUpdateReplayThreads(), 0,
          "the number of replay threads could not be put back");
    }
    finally
    {
      try
      {
        /*
         * A red before the barrier was given up would leave it listed for the rest of this
         * class: it never travelled the replication server, so no session restart brings a
         * delivery which could give it up, and a change which is failing keeps every later
         * session restart at its backoff and the ServerState behind it. The failure which
         * is being reported is the one that matters, so this does not replace it.
         */
        if (!domain.getServerState().cover(failing))
        {
          giveUpOn(domain, failing, waitedOn, failingMods, waitedOnUUID);
        }
      }
      catch (Throwable cleanupFailure)
      {
        logger.error(LocalizableMessage.raw(
            "the barrier of aChangeParkedByAThreadThePoolStoppedIsDeliveredAgain could not be"
                + " given up on the way out: %s", StaticUtils.stackTraceToSingleLineString(cleanupFailure)));
      }
      finally
      {
        try
        {
          resetReplayGiveUpDelay();
        }
        finally
        {
          /*
           * Best-effort on the way out of a red, with no assertion to replace the failure
           * which is being reported: on the normal road the number is put back, and checked,
           * above.
           */
          resetNumUpdateReplayThreads();
        }
      }
    }
  }
  /**
   * Delivers a change until a monitor attribute of the domain reaches the expected value.
   * <p>
   * A delivery is dropped rather than queued while the listener thread is down, which it
   * is for as long as a recovery is restarting the session, so a change which has to reach
   * a replay thread is delivered until it does. A delivery of a change a replay thread
   * owns is refused as the duplicate it is, so the deliveries which follow the one that
   * was taken cost nothing.
   */
  private void deliverUntilMonitorReaches(final LDAPReplicationDomain domain,
      final String attributeName, final long expected,
      final Supplier<? extends LDAPUpdateMsg> delivery, final String message) throws Exception
  {
    TestTimer timer = new TestTimer.Builder()
      .maxSleep(20, SECONDS)
      .sleepTimes(500, MILLISECONDS)
      .toTimer();
    timer.repeatUntilSuccess(new CallableVoid()
    {
      @Override
      public void call() throws Exception
      {
        if (getMonitorAttrValue(baseDN, attributeName) != expected)
        {
          domain.processUpdate(delivery.get());
        }
        assertEquals(getMonitorAttrValue(baseDN, attributeName), expected, message);
      }
    });
  }
  /**
   * Sets how many replay threads this server runs, the way an administrator would: the
   * pool is stopped and created again with that number.
   */
  private static void setNumUpdateReplayThreads(int threads) throws Exception
  {
    assertEquals(TestCaseUtils.applyModifications(true,
        "dn: " + SYNCHRO_PLUGIN_DN,
        "changetype: modify",
        "replace: ds-cfg-num-update-replay-threads",
        "ds-cfg-num-update-replay-threads: " + threads), 0,
        "the number of replay threads could not be changed");
  }
  /**
   * Puts the number of replay threads back to what this server computes for itself, which
   * is what it runs with when the configuration carries no number of its own.
   *
   * @return the result code of the change, so that a finally can call this without an
   *         assertion which would replace the failure it is on the way out of
   */
  private static int resetNumUpdateReplayThreads() throws Exception
  {
    return TestCaseUtils.applyModifications(true,
        "dn: " + SYNCHRO_PLUGIN_DN,
        "changetype: modify",
        "delete: ds-cfg-num-update-replay-threads");
  }
  /** Adds the entry a change of these tests is made on. */
  private DN addEntryForChange(String uid) throws Exception
  {
    return TestCaseUtils.addEntry(
        "dn: uid=" + uid + "," + baseDN,
        "objectClass: top",
        "objectClass: person",
        "objectClass: organizationalPerson",
        "objectClass: inetOrgPerson",
        "uid: " + uid,
        "cn: Aaccf Amar",
        "sn: Amar").getName();
  }
  /** A delivery of a change whose replay does not run to its end. */
  private interface FailingDelivery
  {
@@ -3603,16 +4106,7 @@
  private void assertChangeIsDeliveredAgainAfter(
      FailingDelivery delivery, int serverId, String uid, final String description) throws Exception
  {
    Entry tmp = TestCaseUtils.addEntry(
        "dn: uid=" + uid + "," + baseDN,
        "objectClass: top",
        "objectClass: person",
        "objectClass: organizationalPerson",
        "objectClass: inetOrgPerson",
        "uid: " + uid,
        "cn: Aaccf Amar",
        "sn: Amar");
    final DN dn = tmp.getName();
    final DN dn = addEntryForChange(uid);
    final String uuid = getEntry(dn, 1, true).parseAttribute("entryuuid").asString();
    final LDAPReplicationDomain domain = MultimasterReplication.findDomain(baseDN, null);
@@ -3842,6 +4336,36 @@
  }
  /**
   * A ModifyMsg whose ack runs out of memory on the way out of a replay which applied it.
   * <p>
   * The change is committed before the ack of its delivery is published, and the ack is
   * where the JVM runs out of memory: the one throw from there which is not caught, so the
   * replay is unwound with the change it was replaying in the data and owned by nobody -
   * commit() cleared the owner - and the thread ends on the error. What the give-back on
   * the way out of {@code replay()} has left to give back is the changes this thread parked
   * as waiting for another one, and the restart it asks for them is the one which is run:
   * the road a failed replay takes to ask for its own change again is not on the way.
   * <p>
   * Nothing on the way in reads what throws here: a message handed to the domain rather
   * than published is not one this server acknowledges to anybody.
   */
  private static final class ModifyMsgWhoseAckRunsOutOfMemoryOnceApplied extends ModifyMsg
  {
    private ModifyMsgWhoseAckRunsOutOfMemoryOnceApplied(
        CSN csn, DN dn, List<Modification> mods, String entryUUID)
    {
      super(csn, dn, mods, entryUUID);
    }
    @Override
    public boolean isAssured()
    {
      // Read first thing by processUpdateDone(), which is what publishes the ack.
      throw new OutOfMemoryError("the ack of this applied delivery runs out of memory");
    }
  }
  /**
   * An AddMsg whose ack throws on the way out of a replay which applied it.
   * <p>
   * The change is committed before the ack of its delivery is published, so this is the