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

Valery Kharseko
yesterday 1aa253d7f6f530b9c738ebaf1baf42471dcf01db
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
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
/*
 * 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.backends.jeb;
 
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.fail;
import static org.assertj.core.api.Assertions.failBecauseExceptionWasNotThrown;
import static org.forgerock.opendj.config.ConfigurationMock.mockCfg;
import static org.forgerock.opendj.ldap.ByteString.valueOfUtf8;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.opends.messages.BackendMessages.NOTE_CONFIG_DB_CACHE_REQUIRES_RESTART;
import static org.opends.server.util.CollectionUtils.newTreeSet;
import static org.opends.server.util.StaticUtils.MB;
 
import java.io.File;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.CyclicBarrier;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
 
import org.forgerock.i18n.LocalizableMessage;
import org.forgerock.opendj.config.server.ConfigChangeResult;
import org.forgerock.opendj.config.server.ConfigException;
import org.forgerock.opendj.ldap.ByteString;
import org.forgerock.opendj.ldap.DN;
import org.forgerock.opendj.ldap.ResultCode;
import org.forgerock.opendj.server.config.server.JEBackendCfg;
import org.opends.server.DirectoryServerTestCase;
import org.opends.server.TestCaseUtils;
import org.opends.server.backends.pluggable.spi.AccessMode;
import org.opends.server.backends.pluggable.spi.ReadOperation;
import org.opends.server.backends.pluggable.spi.ReadableTransaction;
import org.opends.server.backends.pluggable.spi.StorageRuntimeException;
import org.opends.server.backends.pluggable.spi.TreeName;
import org.opends.server.backends.pluggable.spi.WriteOperation;
import org.opends.server.backends.pluggable.spi.WriteableTransaction;
import org.opends.server.core.MemoryQuota;
import org.opends.server.core.ServerContext;
import org.opends.server.extensions.DiskSpaceMonitor;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeClass;
import org.testng.annotations.BeforeMethod;
import org.testng.annotations.Test;
 
import com.sleepycat.je.LockConflictException;
import com.sleepycat.je.LockTimeoutException;
 
/**
 * Tests what a {@link JEStorage} takes as it opens and gives back when the open fails - the twin
 * of the same cases on {@code PDBStorageTest} - and the replay of a {@link JEStorage#write} whose
 * transaction JE ends with a {@link LockConflictException}.
 * <p>
 * The conflicts are the engine's own. A deadlock is made by two writers locking two records in opposite order,
 * which JE resolves by throwing at a random victim; a conflict on every attempt is made by a transaction which
 * keeps a record locked while the storage runs with a {@code je.lock.timeout} - the shipped configuration sets
 * none, so a writer waits for a lock rather than times out, but an operator can set one through
 * {@code ds-cfg-je-property}, and JE then reports plain contention as a {@link LockTimeoutException}, which it
 * documents as "abort and retry" just like the deadlock. A {@code DeadlockException} cannot be built by a test:
 * its constructor needs the internal locker it invalidates.
 */
@SuppressWarnings("javadoc")
public class JEStorageTest extends DirectoryServerTestCase
{
  private static final String BACKEND_ID = "JEStorageTest";
  /**
   * A parent directory under which the storage's own directory is a regular file, so that the open
   * fails once its configuration is built - the memory reserved - and before the environment is:
   * what a backend whose directory the server cannot use meets.
   */
  private static final String BLOCKED_DB_DIRECTORY = BACKEND_ID + "-blocked";
  /** A window no run of replays can spend, so that a test of the attempt cap is only ever ended by the cap. */
  private static final long UNREACHABLE_RETRY_WINDOW_NANOS = 300L * 1000L * 1000L * 1000L; //5 min
  /** A window a single attempt outlasts, so that a test of the window reaches it without seconds of build time. */
  private static final long SHORT_RETRY_WINDOW_NANOS = 200L * 1000L * 1000L; //200 ms
  /** A lock wait shorter than any bound, so that a conflict is reported promptly on every attempt. */
  private static final String SHORT_LOCK_TIMEOUT = "20 ms";
  /** A lock wait longer than {@link #SHORT_RETRY_WINDOW_NANOS}, so that one attempt outlasts the window on its own. */
  private static final String LOCK_TIMEOUT_LONGER_THAN_SHORT_WINDOW = "300 ms";
  /** How long a test waits for a thread it started, in seconds; well past any bound the tests configure. */
  private static final long WAIT_SECONDS = 60;
  /** A cache size the quota of the test JVM grants several times over, in bytes. */
  private static final long SMALL_CACHE = 64L * MB;
 
  private final TreeName treeName = new TreeName("dc=test", "test");
  private ServerContext serverContext;
  private JEStorage storage;
 
  @BeforeClass
  public static void startServer() throws Exception
  {
    TestCaseUtils.startServer();
  }
 
  @BeforeMethod
  public void setUp() throws Exception
  {
    serverContext = mock(ServerContext.class);
    when(serverContext.getMemoryQuota()).thenReturn(new MemoryQuota());
    when(serverContext.getDiskSpaceMonitor()).thenReturn(mock(DiskSpaceMonitor.class));
 
    storage = new JEStorage(createBackendCfg(), serverContext);
    // the environment is removed on the way in as well as on the way out: a build whose JVM died never ran
    // tearDown(), and this class shares a fixed db-directory across methods and across builds
    storage.removeStorageFiles();
    storage.open(AccessMode.READ_WRITE);
    createTreeWithTwoRecords();
  }
 
  @AfterMethod
  public void tearDown()
  {
    closeAndRemove(storage);
  }
 
  /**
   * Closes the storage and removes its environment, keeping whichever of the two failed first. Removing it
   * from a finally would let a removal failure replace the close() failure (JLS 14.20.2) - and a close() that
   * throws is exactly the case the removal is here for.
   */
  private static void closeAndRemove(JEStorage storage)
  {
    RuntimeException failure = null;
    try
    {
      storage.close();
    }
    catch (RuntimeException e)
    {
      failure = e;
    }
    try
    {
      storage.removeStorageFiles();
    }
    catch (RuntimeException e)
    {
      if (failure == null)
      {
        failure = e;
      }
      else
      {
        failure.addSuppressed(e);
      }
    }
    if (failure != null)
    {
      throw failure;
    }
  }
 
  /**
   * An open which fails gives back what it took before it failed: the memory it reserved for the
   * cache, and the listener the constructor registered on the backend configuration. Nothing else
   * will - a root container does not close a storage which did not open - and a backend whose
   * directory the server cannot use is enabled again and again, each attempt draining one cache
   * size.
   */
  @Test
  public void aStorageWhoseOpenFailedGivesBackWhatItTook() throws Exception
  {
    final JEBackendCfg cfg = createBackendCfg();
    final JEStorage second = blockedStorage(cfg);
    final MemoryQuota quota = serverContext.getMemoryQuota();
    final long availableBefore = quota.getAvailableMemory();
    try
    {
      second.open(AccessMode.READ_WRITE);
      fail("the storage was expected not to open over a directory which is a file");
    }
    catch (ConfigException expected)
    {
      // What a backend directory the server cannot use does.
    }
    finally
    {
      unblock(second);
    }
 
    assertThat(quota.getAvailableMemory()).isEqualTo(availableBefore);
    verify(cfg).removeJEChangeListener(second);
  }
 
  /**
   * A storage whose open failed has given everything back already, so closing it afterwards takes
   * nothing more - {@code BackendImpl.importLDIF} closes the storage of its root container however
   * the import ended - and does not fail on what the open never got to.
   */
  @Test
  public void closingAStorageWhoseOpenFailedTakesNothingMore() throws Exception
  {
    final JEStorage second = blockedStorage(createBackendCfg());
    final MemoryQuota quota = serverContext.getMemoryQuota();
    final long availableBefore = quota.getAvailableMemory();
    try
    {
      second.open(AccessMode.READ_WRITE);
      fail("the storage was expected not to open over a directory which is a file");
    }
    catch (ConfigException expected)
    {
      // What a backend directory the server cannot use does.
    }
    finally
    {
      unblock(second);
    }
 
    second.close();
 
    assertThat(quota.getAvailableMemory()).isEqualTo(availableBefore);
  }
 
  /**
   * A storage which is open refuses to open again before it takes anything, and what it holds is
   * left as it is: the refusal is a guard against a programming error, not a failed open with
   * something to give back.
   */
  @Test
  public void openingAnOpenStorageIsRefusedAndTakesNothing() throws Exception
  {
    final MemoryQuota quota = serverContext.getMemoryQuota();
    final long availableBefore = quota.getAvailableMemory();
    try
    {
      storage.open(AccessMode.READ_WRITE);
      fail("a storage which is open was expected to refuse to open again");
    }
    catch (IllegalStateException expected)
    {
      // The guard against a double open.
    }
 
    assertThat(quota.getAvailableMemory()).isEqualTo(availableBefore);
    // Still open: a read reaches the environment.
    assertThat(read("missing")).isNull();
  }
 
  /**
   * A cache size changed while the storage is open is given back as it was taken: the close
   * releases what the open reserved, not what the configuration says by then. Read from the
   * configuration at both ends, a change in between drifts the quota by the difference for the
   * life of the JVM - the open which follows reserves the new size and pays nothing back.
   */
  @Test
  public void aCacheGrownWhileOpenIsGivenBackAsItWasTaken() throws Exception
  {
    final MemoryQuota quota = serverContext.getMemoryQuota();
    closeAndRemove(storage);
    final long availableBefore = quota.getAvailableMemory();
    storage = new JEStorage(createBackendCfg(SMALL_CACHE), serverContext);
    storage.open(AccessMode.READ_WRITE);
    assertThat(quota.getAvailableMemory()).isEqualTo(availableBefore - SMALL_CACHE);
 
    storage.applyConfigurationChange(createBackendCfg(2 * SMALL_CACHE));
    storage.close();
 
    assertThat(quota.getAvailableMemory()).isEqualTo(availableBefore);
  }
 
  /** The shrink is the same drift the other way: the difference stays reserved by nobody. */
  @Test
  public void aCacheShrunkWhileOpenIsGivenBackAsItWasTaken() throws Exception
  {
    final MemoryQuota quota = serverContext.getMemoryQuota();
    closeAndRemove(storage);
    final long availableBefore = quota.getAvailableMemory();
    storage = new JEStorage(createBackendCfg(2 * SMALL_CACHE), serverContext);
    storage.open(AccessMode.READ_WRITE);
 
    final ConfigChangeResult ccr = storage.applyConfigurationChange(createBackendCfg(SMALL_CACHE));
    // A shrink asks for the restart as a growth does: the cache keeps the size it was opened with.
    assertThat(ccr.adminActionRequired()).isTrue();
    assertThat(ccr.getMessages()).hasSize(1);
    assertThat(ccr.getMessages().get(0).toString()).isEqualTo(
        NOTE_CONFIG_DB_CACHE_REQUIRES_RESTART.get(BACKEND_ID, 2 * SMALL_CACHE, SMALL_CACHE).toString());
    storage.close();
 
    assertThat(quota.getAvailableMemory()).isEqualTo(availableBefore);
  }
 
  /**
   * The cache is sized when the environment opens and this storage never resizes it, so a change
   * of the cache size is applied by the next open of the backend - and the operator is told so,
   * rather than that the change applied.
   */
  @Test
  public void aCacheSizeChangedWhileOpenAsksForARestart() throws Exception
  {
    closeAndRemove(storage);
    storage = new JEStorage(createBackendCfg(SMALL_CACHE), serverContext);
    storage.open(AccessMode.READ_WRITE);
 
    final ConfigChangeResult ccr = storage.applyConfigurationChange(createBackendCfg(2 * SMALL_CACHE));
 
    assertThat(ccr.getResultCode()).isEqualTo(ResultCode.SUCCESS);
    assertThat(ccr.adminActionRequired()).isTrue();
    assertThat(ccr.getMessages()).hasSize(1);
    assertThat(ccr.getMessages().get(0).ordinal()).isEqualTo(NOTE_CONFIG_DB_CACHE_REQUIRES_RESTART.ordinal());
    assertThat(ccr.getMessages().get(0).toString()).isEqualTo(
        NOTE_CONFIG_DB_CACHE_REQUIRES_RESTART.get(BACKEND_ID, SMALL_CACHE, 2 * SMALL_CACHE).toString());
 
    // The cache still runs at the size it was opened with, whatever the change before said: back to
    // that size, there is nothing left to restart for.
    final ConfigChangeResult back = storage.applyConfigurationChange(createBackendCfg(SMALL_CACHE));
    assertThat(back.adminActionRequired()).isFalse();
    assertThat(back.getMessages()).isEmpty();
  }
 
  /**
   * The default cache is sized by db-cache-percent, db-cache-size left at 0: the restart is asked
   * for by the size the percentage comes to, not by db-cache-size, which does not move.
   */
  @Test
  public void aCacheSizedByPercentAsksForARestartOnlyWhenThePercentChanges() throws Exception
  {
    final MemoryQuota quota = serverContext.getMemoryQuota();
    closeAndRemove(storage);
    storage = new JEStorage(createBackendCfg(0L, 10), serverContext);
    storage.open(AccessMode.READ_WRITE);
    final JEBackendCfg unchangedCache = createBackendCfg(0L, 10);
    when(unchangedCache.isDBTxnNoSync()).thenReturn(true);
 
    final ConfigChangeResult unchanged = storage.applyConfigurationChange(unchangedCache);
    assertThat(unchanged.adminActionRequired()).isFalse();
    assertThat(unchanged.getMessages()).isEmpty();
 
    final ConfigChangeResult ccr = storage.applyConfigurationChange(createBackendCfg(0L, 20));
    assertThat(ccr.adminActionRequired()).isTrue();
    assertThat(ccr.getMessages()).hasSize(1);
    assertThat(ccr.getMessages().get(0).toString()).isEqualTo(NOTE_CONFIG_DB_CACHE_REQUIRES_RESTART.get(
        BACKEND_ID, quota.memPercentToBytes(10), quota.memPercentToBytes(20)).toString());
  }
 
  /**
   * A storage which has not opened runs no cache to restart, and a change of the cache size asks it
   * for none. The listener is registered by the constructor already.
   */
  @Test
  public void aStorageWhichIsNotOpenAsksForNoRestart() throws Exception
  {
    final JEStorage unopened = new JEStorage(createBackendCfg(SMALL_CACHE), serverContext);
    try
    {
      final ConfigChangeResult ccr = unopened.applyConfigurationChange(createBackendCfg(2 * SMALL_CACHE));
 
      assertThat(ccr.adminActionRequired()).isFalse();
      for (LocalizableMessage message : ccr.getMessages())
      {
        assertThat(message.ordinal()).isNotEqualTo(NOTE_CONFIG_DB_CACHE_REQUIRES_RESTART.ordinal());
      }
    }
    finally
    {
      unopened.close();
    }
  }
 
  /** A change which leaves the cache size alone asks for nothing, as before. */
  @Test
  public void aChangeWhichLeavesTheCacheSizeAloneAsksForNothing() throws Exception
  {
    closeAndRemove(storage);
    storage = new JEStorage(createBackendCfg(SMALL_CACHE), serverContext);
    storage.open(AccessMode.READ_WRITE);
    final JEBackendCfg unchangedCache = createBackendCfg(SMALL_CACHE);
    when(unchangedCache.isDBTxnNoSync()).thenReturn(true);
 
    final ConfigChangeResult ccr = storage.applyConfigurationChange(unchangedCache);
 
    assertThat(ccr.getResultCode()).isEqualTo(ResultCode.SUCCESS);
    assertThat(ccr.adminActionRequired()).isFalse();
    assertThat(ccr.getMessages()).isEmpty();
  }
 
  /**
   * A change of the cache size is admitted against what the storage holds of the quota, which is
   * what the next open has to add to. Once a change has been admitted but not applied, the
   * configuration says the new size while the reservation is still the old one, and a check
   * against the configuration would admit a second change the server has no memory for.
   */
  @Test
  public void aCacheSizeChangeIsAdmittedAgainstWhatTheStorageHolds() throws Exception
  {
    final MemoryQuota quota = serverContext.getMemoryQuota();
    closeAndRemove(storage);
    storage = new JEStorage(createBackendCfg(SMALL_CACHE), serverContext);
    storage.open(AccessMode.READ_WRITE);
    storage.applyConfigurationChange(createBackendCfg(2 * SMALL_CACHE));
    // Room for two caches and a bit: the difference to the configured size, not to the reserved one.
    assertThat(quota.acquireMemory(quota.getAvailableMemory() - 2 * SMALL_CACHE - MB)).isTrue();
 
    final List<LocalizableMessage> reasons = new ArrayList<>();
    assertThat(storage.isConfigurationChangeAcceptable(createBackendCfg(4 * SMALL_CACHE), reasons))
        .as("four caches, with one reserved and two and a bit free").isFalse();
    assertThat(storage.isConfigurationChangeAcceptable(createBackendCfg(3 * SMALL_CACHE), reasons))
        .as("three caches, with one reserved and two and a bit free").isTrue();
  }
 
  /**
   * A reservation the quota refused is not given back on close. The open goes ahead without it -
   * the quota is a budget, not a lock - but a close which released what was never taken would
   * hand the quota memory the server does not have.
   */
  @Test
  public void aReservationTheQuotaRefusedIsNotGivenBackOnClose() throws Exception
  {
    final MemoryQuota quota = serverContext.getMemoryQuota();
    closeAndRemove(storage);
    // Half a cache left in the quota: the reservation of a whole one is refused.
    assertThat(quota.acquireMemory(quota.getAvailableMemory() - SMALL_CACHE / 2)).isTrue();
    final long availableBefore = quota.getAvailableMemory();
    storage = new JEStorage(createBackendCfg(SMALL_CACHE), serverContext);
    storage.open(AccessMode.READ_WRITE);
    assertThat(quota.getAvailableMemory()).isEqualTo(availableBefore);
 
    storage.close();
 
    assertThat(quota.getAvailableMemory()).isEqualTo(availableBefore);
  }
 
  /**
   * After an open the quota refused, the storage holds nothing of the quota, and a change which
   * leaves the cache size alone - any other property, the disable an online import makes - still
   * asks the quota for nothing: every change of the backend entry is put to this storage.
   */
  @Test
  public void aChangeWhichLeavesTheCacheSizeAloneIsAdmittedAfterARefusedReservation() throws Exception
  {
    openWithTheReservationRefused();
    final JEBackendCfg unchangedCache = createBackendCfg(SMALL_CACHE);
    when(unchangedCache.isDBTxnNoSync()).thenReturn(true);
 
    assertThat(storage.isConfigurationChangeAcceptable(unchangedCache, new ArrayList<LocalizableMessage>()))
        .isTrue();
  }
 
  /**
   * A growth after an open the quota refused is measured against what the storage holds, which is
   * nothing: a quarter of a cache more than configured is a cache and a quarter more than held.
   */
  @Test
  public void aGrowthAfterARefusedReservationIsMeasuredAgainstNothingHeld() throws Exception
  {
    openWithTheReservationRefused();
 
    assertThat(storage.isConfigurationChangeAcceptable(
        createBackendCfg(SMALL_CACHE + SMALL_CACHE / 4), new ArrayList<LocalizableMessage>()))
        .as("a cache and a quarter, with nothing held and half a cache free").isFalse();
  }
 
  /** A shrink asks the quota for nothing, even with none of it left. */
  @Test
  public void aShrinkIsAdmittedWithTheQuotaExhausted() throws Exception
  {
    final MemoryQuota quota = serverContext.getMemoryQuota();
    closeAndRemove(storage);
    storage = new JEStorage(createBackendCfg(SMALL_CACHE), serverContext);
    storage.open(AccessMode.READ_WRITE);
    assertThat(quota.acquireMemory(quota.getAvailableMemory())).isTrue();
 
    assertThat(storage.isConfigurationChangeAcceptable(
        createBackendCfg(SMALL_CACHE / 2), new ArrayList<LocalizableMessage>())).isTrue();
  }
 
  /**
   * After a shrink while open, the storage still holds the cache it was opened with, and a growth
   * back within that asks the quota for nothing, even with none of it left: measured against the
   * configuration alone, it would ask the quota for the negative difference to what is held.
   */
  @Test
  public void aGrowthWithinWhatIsHeldAfterAShrinkAsksTheQuotaForNothing() throws Exception
  {
    final MemoryQuota quota = serverContext.getMemoryQuota();
    closeAndRemove(storage);
    storage = new JEStorage(createBackendCfg(2 * SMALL_CACHE), serverContext);
    storage.open(AccessMode.READ_WRITE);
    storage.applyConfigurationChange(createBackendCfg(SMALL_CACHE));
    assertThat(quota.acquireMemory(quota.getAvailableMemory())).isTrue();
 
    assertThat(storage.isConfigurationChangeAcceptable(
        createBackendCfg(SMALL_CACHE + SMALL_CACHE / 2), new ArrayList<LocalizableMessage>())).isTrue();
  }
 
  /**
   * A storage which is not open yet - its listener is registered by the constructor, the open comes
   * later - admits a change of a cache sized by percent: the size is counted by the quota of the
   * server context, not by the one the open keeps, which is not there yet.
   */
  @Test
  public void aStorageWhichIsNotOpenAdmitsAChangeOfItsCachePercent() throws Exception
  {
    final JEStorage unopened = new JEStorage(createBackendCfg(0L, 10), serverContext);
    try
    {
      assertThat(unopened.isConfigurationChangeAcceptable(
          createBackendCfg(0L, 20), new ArrayList<LocalizableMessage>())).isTrue();
    }
    finally
    {
      unopened.close();
    }
  }
 
  /** Opens a storage of one cache with half a cache left in the quota, so that its reservation is refused. */
  private void openWithTheReservationRefused() throws Exception
  {
    final MemoryQuota quota = serverContext.getMemoryQuota();
    closeAndRemove(storage);
    assertThat(quota.acquireMemory(quota.getAvailableMemory() - SMALL_CACHE / 2)).isTrue();
    final long availableBefore = quota.getAvailableMemory();
    storage = new JEStorage(createBackendCfg(SMALL_CACHE), serverContext);
    storage.open(AccessMode.READ_WRITE);
    assertThat(quota.getAvailableMemory()).isEqualTo(availableBefore);
  }
 
  /** A storage whose directory is a regular file, which no open of it can use. */
  private JEStorage blockedStorage(JEBackendCfg cfg) throws Exception
  {
    when(cfg.getDBDirectory()).thenReturn(BLOCKED_DB_DIRECTORY);
    final JEStorage blocked = new JEStorage(cfg, serverContext);
    final File directory = blocked.getDirectory();
    directory.getParentFile().mkdirs();
    if (!directory.isFile())
    {
      assertThat(directory.createNewFile()).as("the file in the way of %s", directory).isTrue();
    }
    return blocked;
  }
 
  /** Removes the file in the way of the given storage's directory, and the directory it was made in. */
  private static void unblock(JEStorage blocked)
  {
    final File directory = blocked.getDirectory();
    directory.delete();
    directory.getParentFile().delete();
  }
 
  /**
   * Replaces the storage under test with one bounded by the given values and, when a lock timeout is given, one
   * whose lock waits end in a {@link LockTimeoutException} after that long, the way an operator's
   * {@code ds-cfg-je-property: je.lock.timeout=...} makes them end. Bounding the loop stops the two bounds racing
   * each other: with the shipped values a run of replays spends a random share of the window on backoff alone,
   * so an attempt cap test can be ended by the ten second window instead, and a window test has to make every
   * attempt outlast seconds of that window to reach it.
   */
  private void reopenWithReplayBounds(int maxRetries, long retryWindowNanos, String lockTimeout) throws Exception
  {
    closeAndRemove(storage);
    storage = new JEStorage(createBackendCfg(lockTimeout), serverContext, maxRetries, retryWindowNanos);
    storage.open(AccessMode.READ_WRITE);
    createTreeWithTwoRecords();
  }
 
  /**
   * Two writers which lock the two records in opposite order, each holding its first record's write lock before
   * asking for the other's. JE detects the cycle and ends one of the two transactions - chosen at random - with a
   * {@code DeadlockException}; the other is granted its lock once the victim has aborted.
   */
  private final class OpposedWriter implements Runnable
  {
    private final String name;
    private final String first;
    private final String second;
    private final CyclicBarrier bothHoldTheirFirstLock;
    private final boolean swallowTheConflict;
    final AtomicInteger attempts = new AtomicInteger();
    final AtomicReference<Throwable> failure = new AtomicReference<>();
    final Thread thread;
 
    OpposedWriter(String name, String first, String second, CyclicBarrier bothHoldTheirFirstLock,
        boolean swallowTheConflict)
    {
      this.name = name;
      this.first = first;
      this.second = second;
      this.bothHoldTheirFirstLock = bothHoldTheirFirstLock;
      this.swallowTheConflict = swallowTheConflict;
      this.thread = new Thread(this, name);
    }
 
    @Override
    public void run()
    {
      try
      {
        storage.write(new WriteOperation()
        {
          @Override
          public void run(WriteableTransaction txn) throws Exception
          {
            final int attempt = attempts.incrementAndGet();
            txn.put(treeName, valueOfUtf8(first), valueOfUtf8(name + attempt));
            if (attempt == 1)
            {
              // only the first attempt meets the other writer there: a replay would wait on the barrier forever
              bothHoldTheirFirstLock.await(WAIT_SECONDS, TimeUnit.SECONDS);
            }
            try
            {
              txn.put(treeName, valueOfUtf8(second), valueOfUtf8(name + attempt));
            }
            catch (StorageRuntimeException e)
            {
              if (!swallowTheConflict)
              {
                throw e;
              }
              // an operation which catches what the transaction raised, the way DN2URI.targetEntryReferrals does
            }
          }
        });
      }
      catch (Throwable e)
      {
        failure.set(e);
      }
    }
 
    void startAndJoin(OpposedWriter other) throws InterruptedException
    {
      thread.start();
      other.thread.start();
      thread.join(TimeUnit.SECONDS.toMillis(WAIT_SECONDS));
      other.thread.join(TimeUnit.SECONDS.toMillis(WAIT_SECONDS));
      assertThat(thread.isAlive()).as(name + " finished").isFalse();
      assertThat(other.thread.isAlive()).as(other.name + " finished").isFalse();
    }
  }
 
  /**
   * Both writers commit, and the victim's transaction was rolled back whole and replayed. How many times is
   * JE's to decide: the abort of the victim hands the survivor the lock it waited for, but on the record version
   * the abort undoes, so the survivor has to lock the version put back, and a replay which wins that race forms
   * the deadlock again with the roles drawn afresh - the backoff makes that rare, not impossible.
   */
  private void assertDeadlockVictimReplayedAndBothCommitted(OpposedWriter ab, OpposedWriter ba) throws Exception
  {
    assertThat(ab.failure.get()).as("ab").isNull();
    assertThat(ba.failure.get()).as("ba").isNull();
    assertThat(ab.attempts.get() + ba.attempts.get()).as("at least one of the two was replayed")
        .isGreaterThanOrEqualTo(3);
    assertThat(Math.max(ab.attempts.get(), ba.attempts.get())).isLessThanOrEqualTo(JEStorage.MAX_RETRIES);
    // whichever committed last wrote both records with the number of the attempt which committed, so the
    // records agree - which they would not, had a victim's first record survived its abort
    final ByteString a = read("a");
    assertThat(a).isEqualTo(read("b"));
    final OpposedWriter last = a.toString().startsWith("ab") ? ab : ba;
    assertThat(a).isEqualTo(valueOfUtf8(last.name + last.attempts.get()));
  }
 
  @Test
  public void testDeadlockVictimIsReplayed() throws Exception
  {
    final CyclicBarrier barrier = new CyclicBarrier(2);
    final OpposedWriter ab = new OpposedWriter("ab", "a", "b", barrier, false);
    final OpposedWriter ba = new OpposedWriter("ba", "b", "a", barrier, false);
 
    ab.startAndJoin(ba);
 
    assertDeadlockVictimReplayedAndBothCommitted(ab, ba);
  }
 
  /**
   * JE raises a lock conflict from inside the operation, and an operation may catch it there. The transaction is
   * then abort-only, and it is {@code commit()} which raises the conflict again - so the loop has to treat a
   * conflict raised by the commit as one to replay, not only one raised by the operation. Unlike PersistIt, an
   * attempt which swallowed its conflict therefore commits nothing.
   */
  @Test
  public void testConflictSwallowedInsideTheOperationIsRaisedAgainByCommitAndReplayed() throws Exception
  {
    final CyclicBarrier barrier = new CyclicBarrier(2);
    final OpposedWriter ab = new OpposedWriter("ab", "a", "b", barrier, true);
    final OpposedWriter ba = new OpposedWriter("ba", "b", "a", barrier, true);
 
    ab.startAndJoin(ba);
 
    assertDeadlockVictimReplayedAndBothCommitted(ab, ba);
  }
 
  /**
   * A transaction of its own which keeps a record write-locked until released, so that every attempt of a write
   * asking for that record ends the way the storage's lock timeout ends it.
   */
  private final class LockHolder implements Runnable
  {
    private final String key;
    private final CountDownLatch held = new CountDownLatch(1);
    private final CountDownLatch release = new CountDownLatch(1);
    private final AtomicReference<Throwable> failure = new AtomicReference<>();
    private final Thread thread = new Thread(this, "lock holder");
 
    LockHolder(String key)
    {
      this.key = key;
    }
 
    @Override
    public void run()
    {
      try
      {
        storage.write(new WriteOperation()
        {
          @Override
          public void run(WriteableTransaction txn) throws Exception
          {
            txn.put(treeName, valueOfUtf8(key), valueOfUtf8("held"));
            held.countDown();
            release.await(WAIT_SECONDS, TimeUnit.SECONDS);
          }
        });
      }
      catch (Throwable e)
      {
        failure.set(e);
      }
    }
 
    LockHolder start() throws InterruptedException
    {
      thread.start();
      assertThat(held.await(WAIT_SECONDS, TimeUnit.SECONDS)).as("the holder took the lock").isTrue();
      return this;
    }
 
    void releaseAndJoin() throws InterruptedException
    {
      release.countDown();
      thread.join(TimeUnit.SECONDS.toMillis(WAIT_SECONDS));
      assertThat(thread.isAlive()).as("the holder finished").isFalse();
      assertThat(failure.get()).as("the holder's own write").isNull();
    }
  }
 
  @Test
  public void testWriteGivesUpAfterTheAttemptCap() throws Exception
  {
    // the message is the same at any cap, so this one is spent in two backoffs rather than in the shipped ladder
    final int maxRetries = 3;
    reopenWithReplayBounds(maxRetries, UNREACHABLE_RETRY_WINDOW_NANOS, SHORT_LOCK_TIMEOUT);
    final LockHolder holder = new LockHolder("a").start();
    try
    {
      final AtomicInteger attempts = new AtomicInteger();
      try
      {
        storage.write(new WriteOperation()
        {
          @Override
          public void run(WriteableTransaction txn) throws Exception
          {
            attempts.incrementAndGet();
            txn.put(treeName, valueOfUtf8("a"), valueOfUtf8("abandoned"));
          }
        });
        failBecauseExceptionWasNotThrown(StorageRuntimeException.class);
      }
      catch (StorageRuntimeException e)
      {
        assertThat(e.getMessage()).contains("JEStorageTest").contains(maxRetries + " attempts");
        // and which of the two bounds ran out, since the attempt count alone does not say
        assertThat(e.getMessage()).contains("attempt cap");
        // write() unwraps a StorageRuntimeException that carries a cause, which would replace this message with
        // the bare conflict, and it is the message the config change paths report
        assertThat(e.getCause()).isNull();
        assertThat(e.getSuppressed()).hasSize(1);
        assertThat(e.getSuppressed()[0]).isInstanceOf(LockTimeoutException.class);
      }
      assertThat(attempts.get()).isEqualTo(maxRetries);
    }
    finally
    {
      holder.releaseAndJoin();
    }
    assertThat(read("a")).isEqualTo(valueOfUtf8("held"));
  }
 
  @Test
  public void testWriteIsReplayedUntilTheConflictClears() throws Exception
  {
    reopenWithReplayBounds(JEStorage.MAX_RETRIES, UNREACHABLE_RETRY_WINDOW_NANOS, SHORT_LOCK_TIMEOUT);
    final LockHolder holder = new LockHolder("a").start();
    try
    {
      final AtomicInteger attempts = new AtomicInteger();
      storage.write(new WriteOperation()
      {
        @Override
        public void run(WriteableTransaction txn) throws Exception
        {
          if (attempts.incrementAndGet() == 4)
          {
            // released, and committed, before the lock is asked for, so that this attempt is the one which
            // gets it rather than the one which times out on the holder's commit
            holder.releaseAndJoin();
          }
          txn.put(treeName, valueOfUtf8("a"), valueOfUtf8("applied"));
        }
      });
      assertThat(attempts.get()).isEqualTo(4);
    }
    finally
    {
      holder.releaseAndJoin();
    }
    assertThat(read("a")).isEqualTo(valueOfUtf8("applied"));
  }
 
  /**
   * With a lock timeout longer than the window, a single attempt outlasts the whole window. Giving up on the
   * window alone would then replay nothing, in the very case where the replay is likeliest to succeed: the
   * transaction that was blocking this one has just finished.
   */
  @Test
  public void testWriteIsReplayedOnceWhenTheFirstAttemptOutlastsTheWindow() throws Exception
  {
    reopenWithReplayBounds(JEStorage.MAX_RETRIES, SHORT_RETRY_WINDOW_NANOS, LOCK_TIMEOUT_LONGER_THAN_SHORT_WINDOW);
    final LockHolder holder = new LockHolder("a").start();
    try
    {
      final AtomicInteger attempts = new AtomicInteger();
      storage.write(new WriteOperation()
      {
        @Override
        public void run(WriteableTransaction txn) throws Exception
        {
          if (attempts.incrementAndGet() == 2)
          {
            holder.releaseAndJoin();
          }
          txn.put(treeName, valueOfUtf8("a"), valueOfUtf8("outlasted"));
        }
      });
      assertThat(attempts.get()).isEqualTo(2);
    }
    finally
    {
      holder.releaseAndJoin();
    }
    assertThat(read("a")).isEqualTo(valueOfUtf8("outlasted"));
  }
 
  @Test
  public void testWriteGivesUpOnTheWindowWhenAttemptsAreSlow() throws Exception
  {
    reopenWithReplayBounds(JEStorage.MAX_RETRIES, SHORT_RETRY_WINDOW_NANOS, LOCK_TIMEOUT_LONGER_THAN_SHORT_WINDOW);
    final LockHolder holder = new LockHolder("a").start();
    try
    {
      final AtomicInteger attempts = new AtomicInteger();
      try
      {
        storage.write(new WriteOperation()
        {
          @Override
          public void run(WriteableTransaction txn) throws Exception
          {
            attempts.incrementAndGet();
            // a conflict this slow to report spends the wall clock window long before the attempt cap
            txn.put(treeName, valueOfUtf8("a"), valueOfUtf8("abandoned"));
          }
        });
        failBecauseExceptionWasNotThrown(StorageRuntimeException.class);
      }
      catch (StorageRuntimeException e)
      {
        // the window is what ended it, and it says so: an assertion on the attempt count alone would also pass
        // for a give up on attempt 1, which is the regression the attempt > 1 exemption exists to prevent
        assertThat(e.getMessage()).contains("retry window");
      }
      // one attempt beyond the first: the first spends the window, the exemption grants the replay, and the
      // check after that replay is the one that gives up
      assertThat(attempts.get()).isEqualTo(2);
    }
    finally
    {
      holder.releaseAndJoin();
    }
    assertThat(read("a")).isEqualTo(valueOfUtf8("held"));
  }
 
  /**
   * An interrupt reaches the loop in its backoff sleep alone, and it must not reach JE afterwards: a thread
   * carrying the interrupt flag invalidates the whole environment on its next call - the transaction registry's
   * latch is acquired interruptibly - which is why the loop leaves the flag as the sleep cleared it, and reports
   * the interrupt with the conflict instead. Delivering the interrupt to the sleep alone takes a write which makes
   * no call of JE at all while the flag is set: this one runs on the storage's import environment, whose writes
   * open no transaction, with an operation which raises a conflict JE raised earlier rather than asking JE for a
   * new one - a transaction aborted with the flag set would take the environment down before the sleep is reached.
   */
  @Test
  public void testInterruptedWriteReportsTheConflictItWasReplaying() throws Exception
  {
    final LockConflictException conflict = aConflictOfJEsOwn();
    // the storage that raised it gave up at its first attempt; this one is bounded as shipped
    closeAndRemove(storage);
    storage = new JEStorage(createBackendCfg(), serverContext);
    storage.startImport();
 
    final AtomicInteger attempts = new AtomicInteger();
    final boolean interruptedAfterwards;
    try
    {
      storage.write(new WriteOperation()
      {
        @Override
        public void run(WriteableTransaction txn) throws Exception
        {
          attempts.incrementAndGet();
          Thread.currentThread().interrupt();
          throw new StorageRuntimeException(conflict);
        }
      });
      failBecauseExceptionWasNotThrown(StorageRuntimeException.class);
      return;
    }
    catch (StorageRuntimeException e)
    {
      interruptedAfterwards = Thread.interrupted();
      // the conflict, not the interrupt, is what the caller is told about - but through the same shape the
      // exhausted loop uses, since a bare conflict reaches every caller as its own class name
      assertThat(e.getMessage()).contains("JEStorageTest").contains("interrupted");
      assertThat(e.getSuppressed()).contains(conflict).hasAtLeastOneElementOfType(InterruptedException.class);
      assertThat(e.getCause()).isNull();
    }
    finally
    {
      Thread.interrupted();
    }
    // the flag the sleep cleared stays clear: restored, it would end the environment on the caller's next call
    assertThat(interruptedAfterwards).as("interrupt flag after the write").isFalse();
    // one attempt even though the first backoff is a random 0-49 ms and so is sometimes 0: Thread.sleep() checks
    // the interrupt flag before it checks for a zero duration, so the replay is never reached
    assertThat(attempts.get()).isEqualTo(1);
  }
 
  /** A conflict raised by JE itself: the lock timeout a write gave up on at its first attempt. */
  private LockConflictException aConflictOfJEsOwn() throws Exception
  {
    reopenWithReplayBounds(1, UNREACHABLE_RETRY_WINDOW_NANOS, SHORT_LOCK_TIMEOUT);
    final LockHolder holder = new LockHolder("a").start();
    try
    {
      storage.write(new WriteOperation()
      {
        @Override
        public void run(WriteableTransaction txn) throws Exception
        {
          txn.put(treeName, valueOfUtf8("a"), valueOfUtf8("abandoned"));
        }
      });
      throw new AssertionError("the write was applied although its record was held");
    }
    catch (StorageRuntimeException e)
    {
      assertThat(e.getSuppressed()).hasSize(1);
      return (LockConflictException) e.getSuppressed()[0];
    }
    finally
    {
      holder.releaseAndJoin();
    }
  }
 
  /**
   * The delay grows with the attempt and stays under the cap, so that a contention the first delays did not
   * outlast still has a chance to clear without the replays overrunning the window on sleep alone.
   */
  @Test
  public void testRetryDelayGrowsAndStaysBounded()
  {
    long previousBound = 0;
    for (int attempt = 1; attempt <= JEStorage.MAX_RETRIES; attempt++)
    {
      long bound = 0;
      for (int i = 0; i < 100; i++)
      {
        final long delay = JEStorage.retryDelayMillis(attempt);
        assertThat(delay).as("attempt %d", attempt).isGreaterThanOrEqualTo(0).isLessThan(1000);
        bound = Math.max(bound, delay);
      }
      if (attempt == 1)
      {
        assertThat(bound).as("attempt 1 delays past the first tier").isLessThan(50);
      }
      assertThat(bound).as("attempt %d did not grow past attempt %d", attempt, attempt - 1)
          .isGreaterThanOrEqualTo(previousBound / 2);
      previousBound = bound;
    }
    // and the growth is real rather than a delay that never leaves the first tier
    long grown = 0;
    for (int i = 0; i < 100; i++)
    {
      grown = Math.max(grown, JEStorage.retryDelayMillis(JEStorage.MAX_RETRIES));
    }
    assertThat(grown).as("the last attempts still sleep within the first attempt's bound").isGreaterThan(500);
  }
 
  private void createTreeWithTwoRecords() throws Exception
  {
    storage.write(new WriteOperation()
    {
      @Override
      public void run(WriteableTransaction txn) throws Exception
      {
        txn.openTree(treeName, true);
        txn.put(treeName, valueOfUtf8("a"), valueOfUtf8("0"));
        txn.put(treeName, valueOfUtf8("b"), valueOfUtf8("0"));
      }
    });
  }
 
  private ByteString read(final String key) throws Exception
  {
    return storage.read(new ReadOperation<ByteString>()
    {
      @Override
      public ByteString run(ReadableTransaction txn) throws Exception
      {
        return txn.read(treeName, valueOfUtf8(key));
      }
    });
  }
 
  private static JEBackendCfg createBackendCfg()
  {
    return createBackendCfg(0L);
  }
 
  /** A configuration whose cache is the given size in bytes, or a fifth of the quota when it is zero. */
  private static JEBackendCfg createBackendCfg(long cacheSize)
  {
    return createBackendCfg(cacheSize, 20);
  }
 
  /** A configuration whose cache is the given size in bytes, or the given percent of the quota when it is zero. */
  private static JEBackendCfg createBackendCfg(long cacheSize, int cachePercent)
  {
    final JEBackendCfg backendCfg = mockCfg(JEBackendCfg.class);
    when(backendCfg.dn()).thenReturn(DN.valueOf("ds-cfg-backend-id=" + BACKEND_ID + ",cn=Backends,cn=config"));
    when(backendCfg.getBackendId()).thenReturn(BACKEND_ID);
    when(backendCfg.getDBDirectory()).thenReturn(BACKEND_ID);
    when(backendCfg.getDBDirectoryPermissions()).thenReturn("755");
    when(backendCfg.getDBCacheSize()).thenReturn(cacheSize);
    when(backendCfg.getDBCachePercent()).thenReturn(cachePercent);
    when(backendCfg.getDBNumCleanerThreads()).thenReturn(2);
    when(backendCfg.getDBNumLockTables()).thenReturn(63);
    return backendCfg;
  }
 
  /**
   * The configuration of the storage under test, with the lock timeout an operator would set through
   * {@code ds-cfg-je-property}, or the shipped one - none - when null.
   */
  private static JEBackendCfg createBackendCfg(String lockTimeout)
  {
    final JEBackendCfg backendCfg = createBackendCfg();
    if (lockTimeout != null)
    {
      when(backendCfg.getJEProperty()).thenReturn(newTreeSet("je.lock.timeout=" + lockTimeout));
    }
    return backendCfg;
  }
}