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

Valery Kharseko
21 hours ago 13d57e063cbd3e884f787744a9eaedbfb98b5a51
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
/*
 * The contents of this file are subject to the terms of the Common Development and
 * Distribution License (the License). You may not use this file except in compliance with the
 * License.
 *
 * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the
 * specific language governing permission and limitations under the License.
 *
 * When distributing Covered Software, include this CDDL Header Notice in each file and include
 * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL
 * Header, with the fields enclosed by brackets [] replaced by your own identifying
 * information: "Portions copyright [year] [name of copyright owner]".
 *
 * Copyright 2026 3A Systems, LLC.
 */
package org.opends.server.replication.plugin;
 
import static org.mockito.Matchers.any;
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.mockito.ArgumentCaptor;
import org.opends.server.DirectoryServerTestCase;
import org.opends.server.replication.common.CSN;
import org.opends.server.replication.common.CSNGenerator;
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.ReplicationDomain;
import org.opends.server.types.operation.PluginOperation;
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.
 * <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.
 */
@SuppressWarnings("javadoc")
@Test(groups = { "precommit", "replication" }, sequential = true)
public class PendingChangesTest extends DirectoryServerTestCase
{
  private static final int SERVER_ID = 42;
 
  @Test
  public void replicaOfflineMsgTheBrokerPublishedIsReportedAsSent() throws Exception
  {
    final ReplicationDomain domain = domainWhichPublishes(true);
    final PendingChanges pendingChanges = newPendingChanges(domain);
 
    final CSN offlineCSN = pendingChanges.putReplicaOfflineMsg();
 
    assertNotNull(offlineCSN, "the message was published and must be reported as sent");
    final UpdateMsg published = onlyMsgPublishedBy(domain);
    assertTrue(published instanceof ReplicaOfflineMsg, "published " + published);
    assertEquals(published.getCSN(), offlineCSN);
  }
 
  /**
   * 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
   * server waits out the whole grace period of a message it was told about and which never
   * reached the wire.
   */
  @Test
  public void replicaOfflineMsgTheBrokerRefusedIsNotReportedAsSent() throws Exception
  {
    final ReplicationDomain domain = domainWhichPublishes(false);
    final PendingChanges pendingChanges = newPendingChanges(domain);
 
    assertNull(pendingChanges.putReplicaOfflineMsg(), "the broker refused the message");
 
    assertTrue(onlyMsgPublishedBy(domain) instanceof ReplicaOfflineMsg, "it was attempted");
  }
 
  /**
   * 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.
   */
  @Test
  public void replicaOfflineMsgQueuedBehindAnUncommittedChangeIsNotReportedAsSent() throws Exception
  {
    final ReplicationDomain domain = domainWhichPublishes(true);
    final PendingChanges pendingChanges = newPendingChanges(domain);
    pendingChanges.putLocalOperation(newLocalOperation());
 
    assertNull(pendingChanges.putReplicaOfflineMsg(), "nothing was published");
 
    verify(domain, never()).publish(any(UpdateMsg.class));
  }
 
  /**
   * A message which could not be published is given up on rather than left queued: the replica
   * which could not announce itself offline is either shutting down, and the message dies with
   * the process, or it is being disabled for an import or a configuration change - and once it
   * comes back, announcing it offline on the session which follows would be a lie.
   */
  @Test
  public void replicaOfflineMsgWhichCouldNotBeSentIsNotPublishedLater() throws Exception
  {
    final ReplicationDomain domain = domainWhichPublishes(true);
    final PendingChanges pendingChanges = newPendingChanges(domain);
    final CSN changeCSN = pendingChanges.putLocalOperation(newLocalOperation());
    assertNull(pendingChanges.putReplicaOfflineMsg(), "nothing was published");
 
    // The change which held the message back completes.
    pendingChanges.commitAndPushCommittedChanges(changeCSN, mock(LDAPUpdateMsg.class));
 
    final UpdateMsg published = onlyMsgPublishedBy(domain);
    assertTrue(published instanceof LDAPUpdateMsg, "published " + published);
  }
 
  /**
   * 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
   * information of its entry. Only the offline announcement, which is stored nowhere, needs the
   * answer of the broker.
   */
  @Test
  public void changeTheBrokerRefusedStillLeavesThePendingChanges() throws Exception
  {
    final ReplicationDomain domain = domainWhichPublishes(false);
    final PendingChanges pendingChanges = newPendingChanges(domain);
    final CSN changeCSN = pendingChanges.putLocalOperation(newLocalOperation());
    assertEquals(pendingChanges.size(), 1);
 
    pendingChanges.commitAndPushCommittedChanges(changeCSN, mock(LDAPUpdateMsg.class));
 
    assertTrue(onlyMsgPublishedBy(domain) instanceof LDAPUpdateMsg, "the change was published");
    assertEquals(pendingChanges.size(), 0, "and is not queued for a second attempt");
  }
 
  private PendingChanges newPendingChanges(ReplicationDomain domain)
  {
    return new PendingChanges(new CSNGenerator(SERVER_ID, 0), domain);
  }
 
  /** A domain whose broker accepts, or refuses, whatever it is given to publish. */
  private ReplicationDomain domainWhichPublishes(boolean accepted)
  {
    final ReplicationDomain domain = mock(ReplicationDomain.class);
    when(domain.publish(any(UpdateMsg.class))).thenReturn(accepted);
    return domain;
  }
 
  /** A local operation, i.e. one this replica must publish to the other replicas. */
  private PluginOperation newLocalOperation()
  {
    final PluginOperation operation = mock(PluginOperation.class);
    when(operation.isSynchronizationOperation()).thenReturn(false);
    return operation;
  }
 
  /**
   * The single message the domain was asked to publish, failing the test if it published
   * anything else: a message which is not sent and a message which is sent twice are both the
   * kind of mistake these tests are about.
   */
  private UpdateMsg onlyMsgPublishedBy(ReplicationDomain domain)
  {
    final ArgumentCaptor<UpdateMsg> published = ArgumentCaptor.forClass(UpdateMsg.class);
    verify(domain).publish(published.capture());
    return published.getValue();
  }
}