3030import java .sql .SQLException ;
3131import java .sql .Statement ;
3232import java .util .ArrayList ;
33+ import java .util .Arrays ;
3334import java .util .HashMap ;
35+ import java .util .HashSet ;
3436import java .util .List ;
3537import java .util .Map ;
3638import java .util .Properties ;
@@ -89,7 +91,10 @@ public void testPutGetDeleteSequence() throws Exception {
8991
9092 List <PartitionId > partitionIds = clusterMap .getWritablePartitionIds (MockClusterMap .DEFAULT_PARTITION_CLASS );
9193 List <Long > partitions = partitionIds .stream ().map (p -> p .getId ()).collect (Collectors .toList ());
92- int blobsPerContainer = 5 ;
94+ // Each container covers delete-only, TTL-only, mixed, and source-only partitions.
95+ RepairRequestRecord .OperationType [] operationTypes =
96+ {DeleteRequest , DeleteRequest , TtlUpdateRequest , TtlUpdateRequest , TtlUpdateRequest , DeleteRequest ,
97+ TtlUpdateRequest , DeleteRequest };
9398
9499 String hostName1 = "localhost1" ;
95100 String hostName2 = "localhost2" ;
@@ -102,23 +107,26 @@ public void testPutGetDeleteSequence() throws Exception {
102107 // Prepare RepairRequests and insert them to the DB.
103108 // Map<Partition ID, Map<BlobId, RepairRequestRecord>>
104109 Map <Integer , Map <String , RepairRequestRecord >> records = new HashMap <>();
110+ for (PartitionId partitionId : partitionIds ) {
111+ records .put ((int ) partitionId .getId (), new HashMap <>());
112+ }
105113 for (Account account : accountService .getAllAccounts ()) {
106114 for (Container container : account .getAllContainers ()) {
107- for (int i = 0 ; i < blobsPerContainer ; i ++) {
108- PartitionId partitionId = partitionIds .get (random . nextInt ( partitionIds . size ()) );
115+ for (int i = 0 ; i < operationTypes . length ; i ++) {
116+ PartitionId partitionId = partitionIds .get (i / 2 );
109117 String blobId = generateBlobId (account , container , partitionId );
110- RepairRequestRecord .OperationType operationType = i % 2 == 0 ? TtlUpdateRequest : DeleteRequest ;
111- long operationTime = System . currentTimeMillis () - random . nextInt ( 1000 );
118+ RepairRequestRecord .OperationType operationType = operationTypes [ i ] ;
119+ long operationTime = time . milliseconds ( );
112120 short lifeVersion = -1 ;
113121 long expirationTime =
114- i % 2 == 0 ? Utils .Infinite_Time : System .currentTimeMillis () + TimeUnit .HOURS .toMillis (1 );
115- String hostName = random .nextInt (2 ) == 0 ? hostName1 : hostName2 ;
116- int hostPort = random .nextInt (2 ) == 0 ? hostPort1 : hostPort2 ;
122+ operationType == TtlUpdateRequest ? Utils .Infinite_Time : operationTime + TimeUnit .HOURS .toMillis (1 );
123+ boolean sourceOnlyPartition = partitionId .equals (partitionIds .get (3 ));
124+ String hostName = !sourceOnlyPartition && i % 2 == 0 ? hostName2 : hostName1 ;
125+ int hostPort = !sourceOnlyPartition && i % 2 != 0 ? hostPort2 : hostPort1 ;
117126 RepairRequestRecord record =
118127 new RepairRequestRecord (blobId , (int ) partitionId .getId (), hostName , hostPort , operationType ,
119128 operationTime , lifeVersion , expirationTime );
120129 repairRequestsDb .putRepairRequests (record );
121- records .putIfAbsent ((int ) partitionId .getId (), new HashMap <>());
122130 records .get ((int ) partitionId .getId ()).put (record .getBlobId (), record );
123131 }
124132 }
@@ -127,34 +135,31 @@ public void testPutGetDeleteSequence() throws Exception {
127135 // on one node with name as thisNodeName and port as thisNodePort,
128136 // read the database but exclude all the records which has the source replica as this node.
129137 // we should run ODR to fix the requests on the nodes except the source replica
130- Set <Long > partitionsNeedRepair = repairRequestsDb .getPartitionsNeedRepair (thisNodeName , thisNodePort , partitions );
138+ Set <Long > expectedPartitionsNeedRepair = new HashSet <>(Arrays .asList (partitions .get (1 ), partitions .get (2 )));
139+ assertEquals (expectedPartitionsNeedRepair ,
140+ repairRequestsDb .getPartitionsNeedRepair (thisNodeName , thisNodePort , partitions ));
131141 for (PartitionId id : partitionIds ) {
132142 Pair <List <RepairRequestRecord >, Long > dbRecords =
133143 repairRequestsDb .getRepairRequestsExcludingHost ((int ) id .getId (), thisNodeName , thisNodePort , 0 );
134144 List <RepairRequestRecord > recordFromStore = dbRecords .getFirst ();
135- Set <Long > partitionsNeedRepairUpdated ;
136145 Map <String , RepairRequestRecord > orgRecords = records .get ((int ) id .getId ());
137- if (recordFromStore .size () > 0 ) {
138- for (RepairRequestRecord record : recordFromStore ) {
139- RepairRequestRecord org = orgRecords .get (record .getBlobId ());
140- assertEquals ("Record does not match expectation " , org , record );
141- assertTrue ("should exclude this node" ,
142- !thisNodeName .equals (record .getSourceHostName ()) || thisNodePort != record .getSourceHostPort ());
143- orgRecords .remove (record .getBlobId ());
144- repairRequestsDb .removeRepairRequests (record .getBlobId (), record .getOperationType ());
145- }
146- partitionsNeedRepairUpdated = repairRequestsDb .getPartitionsNeedRepair (thisNodeName , thisNodePort , partitions );
147- assertTrue (partitionsNeedRepair .contains (id .getId ()));
148- assertFalse (partitionsNeedRepairUpdated .contains (id .getId ()));
149- assertEquals (partitionsNeedRepair .size (), partitionsNeedRepairUpdated .size () + 1 );
150- } else {
151- partitionsNeedRepairUpdated = repairRequestsDb .getPartitionsNeedRepair (thisNodeName , thisNodePort , partitions );
152- assertEquals (partitionsNeedRepairUpdated , partitionsNeedRepair );
153- assertFalse (partitionsNeedRepairUpdated .contains (id .getId ()));
146+ long expectedRecordCount = orgRecords .values ().stream ()
147+ .filter (record -> !thisNodeName .equals (record .getSourceHostName ())
148+ || thisNodePort != record .getSourceHostPort ()).count ();
149+ assertEquals ("Eligible record count does not match" , expectedRecordCount , recordFromStore .size ());
150+ for (RepairRequestRecord record : recordFromStore ) {
151+ RepairRequestRecord org = orgRecords .get (record .getBlobId ());
152+ assertEquals ("Record does not match expectation " , org , record );
153+ assertTrue ("should exclude this node" ,
154+ !thisNodeName .equals (record .getSourceHostName ()) || thisNodePort != record .getSourceHostPort ());
155+ orgRecords .remove (record .getBlobId ());
156+ repairRequestsDb .removeRepairRequests (record .getBlobId (), record .getOperationType ());
154157 }
155- partitionsNeedRepair = partitionsNeedRepairUpdated ;
158+ expectedPartitionsNeedRepair .remove (id .getId ());
159+ assertEquals (expectedPartitionsNeedRepair ,
160+ repairRequestsDb .getPartitionsNeedRepair (thisNodeName , thisNodePort , partitions ));
156161 }
157- assertTrue (partitionsNeedRepair . size () == 0 );
162+ assertTrue (expectedPartitionsNeedRepair . isEmpty () );
158163
159164 // get the remaining records.
160165 for (PartitionId id : partitionIds ) {
0 commit comments