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

Valery Kharseko
yesterday 70d9a179cdd8d975b44e1815c249e20d9f91097f
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
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
/*
 * 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 2014-2016 ForgeRock AS.
 * Portions Copyright 2026 3A Systems, LLC.
 */
package org.opends.server.replication.service;
 
import static java.util.concurrent.TimeUnit.MILLISECONDS;
import static java.util.concurrent.TimeUnit.NANOSECONDS;
 
import java.util.Collection;
import java.util.Map.Entry;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
 
import org.forgerock.opendj.ldap.DN;
import org.opends.server.replication.common.CSN;
 
/**
 * Class useful for the case where DS/RS instances are collocated inside the
 * same JVM. It synchronizes the shutdown of the DS and RS sides.
 * <p>
 * More specifically, it ensures a ReplicaOfflineMsg sent by the DS is
 * relayed/forwarded by the collocated RS to the other RSs in the topology
 * before the whole process shuts down.
 * <p>
 * The state is kept per domain and per instance: the collocated DS and RS
 * sides coordinate through the single instance MultimasterReplication hands
 * to both of them.
 *
 * @since OPENDJ-1453
 */
public class DSRSShutdownSync
{
  /**
   * How long a ReplicaOfflineMsg may hold back the shutdown of the collocated
   * RS, in milliseconds, counted from the moment the message was announced.
   */
  public static final long REPLICA_OFFLINE_GRACE_PERIOD = 5000;
 
  private final long gracePeriod;
 
  /**
   * The ReplicaOfflineMsg which is still owed a forward, per domain and per
   * replica of that domain.
   * <p>
   * An entry lives until every replication server the message was queued for
   * has forwarded it, so it legitimately holds a message some of them have
   * already sent: what is pending is the forward, not the message.
   * <p>
   * It is kept per domain because a domain sends this message whenever its
   * replication service is disabled - an online import, a restore, a
   * configuration change - and not only when the process shuts down. A single
   * entry for the whole process would be the one of the first such message and
   * would leave no grace period at all to the shutdown this class exists for.
   * <p>
   * It is kept per replica because the collocated RS relays the message of
   * every replica connected to it, and the forward of another replica's
   * message says nothing about this one.
   * <p>
   * Each entry knows the replication servers its message was queued for, because each of them is
   * served by its own writer: the forward of one of them says nothing about the others, whose
   * queue the shutdown is about to clear.
   */
  private final ConcurrentMap<DN, ConcurrentMap<Integer, PendingOfflineMsg>> replicaOfflineMsgs =
      new ConcurrentHashMap<>();
  /** Monitor notified whenever a ReplicaOfflineMsg has been forwarded. */
  private final Object forwardedMonitor = new Object();
 
  /** Creates a synchronization object using the default grace period. */
  public DSRSShutdownSync()
  {
    this(REPLICA_OFFLINE_GRACE_PERIOD);
  }
 
  /**
   * Creates a synchronization object using the provided grace period.
   *
   * @param gracePeriod
   *          how long a ReplicaOfflineMsg may hold back the shutdown, in milliseconds
   */
  DSRSShutdownSync(long gracePeriod)
  {
    this.gracePeriod = gracePeriod;
  }
 
  /**
   * Message is about to be sent.
   * <p>
   * The announcement comes before the message is published rather than after: a collocated
   * replication server can forward the message as soon as it is on the wire, and a forward which
   * finds nothing announced has nothing to clear. The announcement of a message the broker then
   * refuses is taken back by {@link #replicaOfflineMsgNotSent(DN, CSN)}.
   * <p>
   * A replica announces itself offline on every disableService(), so this may take the place of
   * an earlier announcement of the same replica which is still owed its forward. The earlier one
   * is kept behind the new one: a forward of the newer message, which the replication server
   * queued behind the earlier one, covers both, and a withdrawal of the newer one gives the
   * earlier one its wait back. It is kept only while its own grace period runs: past it, the
   * announcement holds nothing back any more, and keeping it would chain every announcement of
   * a replica whose message nobody in this process forwards - a directory server without a
   * collocated replication server, or connected to a remote one - for the life of the process.
   *
   * @param baseDN
   *          the domain for which the message is being sent
   * @param offlineCSN
   *          the CSN of the message, which identifies both the replica which announces itself
   *          offline and the announcement being waited for
   */
  public void replicaOfflineMsgSent(DN baseDN, CSN offlineCSN)
  {
    final long announcedAt = System.nanoTime();
    replicaOfflineMsgs
        .computeIfAbsent(baseDN, dn -> new ConcurrentHashMap<Integer, PendingOfflineMsg>())
        .compute(offlineCSN.getServerId(), (serverId, displaced) ->
            new PendingOfflineMsg(offlineCSN, announcedAt,
                displaced != null && gracePeriodLeft(displaced, announcedAt) > 0 ? displaced : null));
  }
 
  /**
   * The message which was announced was not sent after all: the broker had no session to write
   * it to, or was stopped before it could.
   * <p>
   * The announcement is made before the message is published, since a collocated replication
   * server can forward it as soon as it is on the wire, so the announcement of a message the
   * broker then refused has to be taken back: nobody will forward it, and the shutdown would
   * spend the whole grace period waiting for that forward. Only the announcement carrying that
   * CSN is withdrawn, and the announcement it displaced - an earlier message of the same replica
   * which did go out and is still owed its forward - takes its place again.
   * <p>
   * Whatever is reported about that earlier message while the announcement of the refused one
   * stands in its place is not seen by it. A forward, or the loss of a peer it was queued for,
   * is lost, and the shutdown then waits out what is left of the earlier message's own grace
   * period; the peers it is queued for, if they are recorded in that window, are lost too, with
   * the opposite effect - the first forward ends its wait, as for a message no peer was recorded
   * for. That window is the one publish the broker refuses: at once on a connection error or a
   * pending recovery, the broker's retry loop up to the reconnect when it has no session. The
   * wait it can cost is bounded by a grace period which is already running.
   *
   * @param baseDN
   *          the domain for which the message was announced
   * @param offlineCSN
   *          the CSN of the message which was not sent
   */
  public void replicaOfflineMsgNotSent(DN baseDN, CSN offlineCSN)
  {
    final ConcurrentMap<Integer, PendingOfflineMsg> msgs = replicaOfflineMsgs.get(baseDN);
    if (msgs != null)
    {
      final int serverId = offlineCSN.getServerId();
      final PendingOfflineMsg pending = msgs.get(serverId);
      if (pending != null && pending.csn.equals(offlineCSN))
      {
        /*
         * The displaced announcement may owe nothing any more: the forward which released it
         * can have been reported while this announcement was being made, so that its remove(),
         * which matches the entry it read, found this one in its place. Given its place back,
         * such an announcement would hold the shutdown for the rest of its grace period, since
         * nobody will report that forward again.
         */
        if (pending.displaced != null && !pending.displaced.isFullyForwarded())
        {
          msgs.replace(serverId, pending, pending.displaced);
        }
        else
        {
          msgs.remove(serverId, pending);
        }
      }
    }
    notifyForwarded();
  }
 
  /**
   * Message has been queued for the replication servers which must forward it.
   * <p>
   * This must be called before the message is queued for any of them: a replication server can
   * forward it as soon as it is in its queue, and a forward which finds no recipient recorded
   * ends the wait at once.
   *
   * @param baseDN
   *          the domain for which the message has been sent
   * @param offlineCSN
   *          the CSN of the message which is being queued
   * @param replicationServerIds
   *          the server ids of the replication servers the message is being queued for
   */
  public void replicaOfflineMsgDispatched(
      DN baseDN, CSN offlineCSN, Collection<Integer> replicationServerIds)
  {
    final ConcurrentMap<Integer, PendingOfflineMsg> msgs = replicaOfflineMsgs.get(baseDN);
    if (msgs == null)
    {
      return;
    }
    final int serverId = offlineCSN.getServerId();
    final PendingOfflineMsg pending = msgs.get(serverId);
    /*
     * The message being queued may be an older announcement of the same replica - one which was
     * queued behind a backlog since an earlier import. The replication servers it goes to say
     * nothing about the announcement the shutdown is waiting for.
     */
    if (pending != null && pending.csn.equals(offlineCSN)
        && pending.awaitForwardsFrom(replicationServerIds))
    {
      // queued for nobody: there is nothing to wait for
      msgs.remove(serverId, pending);
      notifyForwarded();
    }
  }
 
  /**
   * Message has been forwarded to one of the replication servers it was queued for.
   *
   * @param baseDN
   *          the domain for which the message has been sent
   * @param forwardedCSN
   *          the CSN of the forwarded message
   * @param replicationServerId
   *          the server id of the replication server the message has been forwarded to
   */
  public void replicaOfflineMsgForwarded(DN baseDN, CSN forwardedCSN, int replicationServerId)
  {
    final ConcurrentMap<Integer, PendingOfflineMsg> msgs = replicaOfflineMsgs.get(baseDN);
    if (msgs != null)
    {
      final int serverId = forwardedCSN.getServerId();
      final PendingOfflineMsg pending = msgs.get(serverId);
      /*
       * A replica announces itself offline on every disableService(), so the message which is
       * forwarded now may be an older one - queued behind a backlog since an earlier import, or
       * synthesized from the offline CSN of the changelog for a server which is catching up.
       * Such a forward says nothing about the announcement the shutdown is waiting for, and must
       * not consume its grace period.
       */
      if (pending != null && pending.csn.isOlderThanOrEqualTo(forwardedCSN)
          && pending.forwardedBy(replicationServerId))
      {
        msgs.remove(serverId, pending);
      }
    }
    notifyForwarded();
  }
 
  /**
   * A replication server the message may have been queued for will not forward it: it is gone,
   * or the message was dropped on its way out.
   * <p>
   * Whatever it was given can no longer reach it, so the shutdown must not spend the rest of its
   * grace period waiting for it.
   *
   * @param baseDN
   *          the domain the replication server is connected to
   * @param replicationServerId
   *          the server id of the replication server which will not forward the message
   */
  public void replicaOfflineMsgNotForwarded(DN baseDN, int replicationServerId)
  {
    final ConcurrentMap<Integer, PendingOfflineMsg> msgs = replicaOfflineMsgs.get(baseDN);
    if (msgs != null)
    {
      for (Entry<Integer, PendingOfflineMsg> entry : msgs.entrySet())
      {
        final PendingOfflineMsg pending = entry.getValue();
        if (pending.giveUpOn(replicationServerId))
        {
          msgs.remove(entry.getKey(), pending);
        }
      }
    }
    notifyForwarded();
  }
 
  /** Wakes up the shutdown, which re-reads what is left to wait for. */
  private void notifyForwarded()
  {
    synchronized (forwardedMonitor)
    {
      forwardedMonitor.notifyAll();
    }
  }
 
  /**
   * Whether the shutdown of a domain can proceed, i.e. its ReplicaOfflineMsg
   * has been forwarded by every replication server it was queued for, or its
   * grace period has expired.
   * <p>
   * The shutdown itself blocks on {@link #awaitReplicaOfflineMsgsForwarded(Collection, long)}
   * rather than polling this; it is the same state, observable without waiting for it.
   *
   * @param baseDN
   *          the baseDN of the domain being shut down
   * @return true if the shutdown of this domain need not wait any longer, i.e. its message was
   *         forwarded or its grace period has expired, false otherwise
   */
  public boolean canShutdown(DN baseDN)
  {
    return remainingGracePeriod(baseDN) <= 0;
  }
 
  /**
   * Returns the time by which every wait of one shutdown must be over.
   * <p>
   * A process shuts its domains down one after the other and each of them may have a message
   * pending, so a deadline computed once and shared by all of them keeps the whole shutdown
   * bounded by one grace period instead of one per domain.
   *
   * @return the point in time, on the {@link System#nanoTime()} clock, by which the waits must
   *         be over
   */
  public long newShutdownDeadline()
  {
    return System.nanoTime() + MILLISECONDS.toNanos(gracePeriod);
  }
 
  /**
   * Waits for the ReplicaOfflineMsg of every provided domain to be forwarded, or for their grace
   * periods or the provided deadline to expire.
   * <p>
   * This must be called before the server handlers of those domains are stopped: stopping them
   * deactivates their consumer, clears their message queue and closes their session, after which
   * the message can no longer be forwarded.
   * <p>
   * All the domains of one shutdown wait together rather than one after the other, so that the
   * shutdown is bounded by one grace period without the wait of one domain spending the grace
   * period of the next.
   *
   * @param baseDNs
   *          the baseDNs of the domains whose messages must be forwarded
   * @param deadline
   *          the point in time, on the {@link System#nanoTime()} clock, by which this wait must
   *          be over whatever the domains announce in the meantime - see
   *          {@link #newShutdownDeadline()}. A deadline which is not in the future returns
   *          without waiting at all, for a caller which has nothing to wait for.
   */
  public void awaitReplicaOfflineMsgsForwarded(Collection<DN> baseDNs, long deadline)
  {
    if (deadline - System.nanoTime() <= 0)
    {
      return;
    }
    synchronized (forwardedMonitor)
    {
      while (true)
      {
        final long timeout = Math.min(remainingGracePeriod(baseDNs),
            NANOSECONDS.toMillis(deadline - System.nanoTime()));
        if (timeout <= 0)
        {
          return;
        }
        try
        {
          forwardedMonitor.wait(timeout);
        }
        catch (InterruptedException e)
        {
          /*
           * Give up waiting. The interrupt is deliberately not restored: what follows this call is
           * the rest of the shutdown - joining the reader and writer thread of every handler, then
           * closing the changelog DB - and an interrupt flag would make all of it give up too.
           */
          return;
        }
      }
    }
  }
 
  /**
   * Returns the time left, in milliseconds, to forward the ReplicaOfflineMsg of the replica of
   * the provided domains which has the longest to wait, zero or less if none of them has a
   * message pending.
   */
  private long remainingGracePeriod(Collection<DN> baseDNs)
  {
    long remaining = 0;
    for (DN baseDN : baseDNs)
    {
      remaining = Math.max(remaining, remainingGracePeriod(baseDN));
    }
    return remaining;
  }
 
  /**
   * Returns the time left, in milliseconds, to forward the ReplicaOfflineMsg of the replica of
   * this domain which has the longest to wait, zero or less if no message of this domain is
   * pending.
   */
  private long remainingGracePeriod(DN baseDN)
  {
    final ConcurrentMap<Integer, PendingOfflineMsg> msgs = replicaOfflineMsgs.get(baseDN);
    if (msgs == null)
    {
      return 0;
    }
    final long now = System.nanoTime();
    long remaining = 0;
    for (PendingOfflineMsg pending : msgs.values())
    {
      remaining = Math.max(remaining, gracePeriodLeft(pending, now));
    }
    return remaining;
  }
 
  /**
   * Returns the time left, in milliseconds, of the grace period of one announcement, zero or
   * less once it has expired.
   */
  private long gracePeriodLeft(PendingOfflineMsg pending, long now)
  {
    return gracePeriod - NANOSECONDS.toMillis(now - pending.sentTime);
  }
 
  /**
   * A ReplicaOfflineMsg a replica announced and which has not been forwarded yet.
   * <p>
   * This deliberately does not override {@code equals}: the two-argument
   * {@link ConcurrentMap#remove(Object, Object)} of the forward guard must match the very
   * announcement it read, not another one which happens to carry the same values.
   */
  private static final class PendingOfflineMsg
  {
    /** The CSN of the message, so that the forward of an older one is not taken for this one. */
    private final CSN csn;
    /** When the message was announced, on the {@link System#nanoTime()} clock. */
    private final long sentTime;
    /**
     * The announcement of the same replica this one took the place of and which is still owed its
     * forward, null when there was none or when its grace period had already expired. It is
     * given its place back if this message is withdrawn.
     */
    private final PendingOfflineMsg displaced;
    /**
     * The replication servers the message was queued for and which have not forwarded it yet,
     * null as long as it has not been queued for anybody.
     */
    private volatile Set<Integer> awaitedForwarders;
 
    private PendingOfflineMsg(CSN csn, long sentTime, PendingOfflineMsg displaced)
    {
      this.csn = csn;
      this.sentTime = sentTime;
      this.displaced = displaced;
    }
 
    /**
     * Records the replication servers the message is being queued for, and returns whether there
     * is none of them, i.e. nothing left to wait for.
     */
    private boolean awaitForwardsFrom(Collection<Integer> replicationServerIds)
    {
      final Set<Integer> awaited = ConcurrentHashMap.newKeySet();
      awaited.addAll(replicationServerIds);
      awaitedForwarders = awaited;
      return awaited.isEmpty();
    }
 
    /**
     * Records the forward of one replication server, and returns whether nothing is left to wait
     * for.
     */
    private boolean forwardedBy(int replicationServerId)
    {
      final Set<Integer> awaited = awaitedForwarders;
      if (awaited == null)
      {
        /*
         * The message never went through the collocated RS - a replica which picked a remote one
         * is announcing itself offline. Nobody is known to owe a forward, so keep the behaviour
         * the wait had before the recipients were tracked: the first forward ends it.
         */
        return true;
      }
      awaited.remove(replicationServerId);
      return awaited.isEmpty();
    }
 
    /**
     * Gives up on the forward of one replication server, and returns whether nothing is left to
     * wait for. Unlike a forward, this releases nothing while no recipient is known: a peer going
     * away says nothing about a message it was never given.
     */
    private boolean giveUpOn(int replicationServerId)
    {
      final Set<Integer> awaited = awaitedForwarders;
      return awaited != null && awaited.remove(replicationServerId) && awaited.isEmpty();
    }
 
    /**
     * Returns whether every replication server the message was queued for has forwarded it, or
     * has been given up on: nothing is left to wait for. False while no recipient is known.
     */
    private boolean isFullyForwarded()
    {
      final Set<Integer> awaited = awaitedForwarders;
      return awaited != null && awaited.isEmpty();
    }
 
    @Override
    public String toString()
    {
      final Set<Integer> awaited = awaitedForwarders;
      return "PendingOfflineMsg(" + csn
          + (awaited != null ? ", awaiting the forward of " + awaited : "") + ")";
    }
  }
}