001/* 002 * Licensed to DuraSpace under one or more contributor license agreements. 003 * See the NOTICE file distributed with this work for additional information 004 * regarding copyright ownership. 005 * 006 * DuraSpace licenses this file to you under the Apache License, 007 * Version 2.0 (the "License"); you may not use this file except in 008 * compliance with the License. You may obtain a copy of the License at 009 * 010 * http://www.apache.org/licenses/LICENSE-2.0 011 * 012 * Unless required by applicable law or agreed to in writing, software 013 * distributed under the License is distributed on an "AS IS" BASIS, 014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. 015 * See the License for the specific language governing permissions and 016 * limitations under the License. 017 */ 018package org.fcrepo.integration.http.api; 019 020import static javax.ws.rs.core.HttpHeaders.CONTENT_TYPE; 021import static javax.ws.rs.core.Response.Status.CREATED; 022import static org.junit.Assert.assertEquals; 023 024import java.io.ByteArrayInputStream; 025import java.io.IOException; 026import java.util.ArrayList; 027import java.util.List; 028 029import org.apache.http.HttpResponse; 030import org.apache.http.client.HttpClient; 031import org.apache.http.client.methods.HttpGet; 032import org.apache.http.client.methods.HttpPatch; 033import org.apache.http.client.methods.HttpRequestBase; 034import org.apache.http.entity.BasicHttpEntity; 035import org.junit.Test; 036import org.springframework.test.context.TestExecutionListeners; 037 038/** 039 * This "test" is a utility for collecting the timing of concurrent operations operations. 040 * It takes roughly 2 minutes to complete and should only be run if the timing metrics are wanted. 041 * In order to activate this utility, the following System Property must be set: 042 * <p/> 043 * mvn -Dfcrepo.test.http.concurrent install 044 * 045 * @author lsitu 046 */ 047@TestExecutionListeners( 048 listeners = { TestIsolationExecutionListener.class }, 049 mergeMode = TestExecutionListeners.MergeMode.MERGE_WITH_DEFAULTS) 050public class FedoraCrudConcurrentIT extends AbstractResourceIT { 051 052 private static final String TEST_ACTIVATION_PROPERTY = "fcrepo.test.http.concurrent"; 053 054 @Test 055 public void testConcurrentIngest() throws Exception { 056 setLogger(); 057 058 if (System.getProperty(TEST_ACTIVATION_PROPERTY) == null) { 059 logger.info("Not running tests because system property not set: {}", TEST_ACTIVATION_PROPERTY); 060 return; 061 } 062 063 final int[] numThreadsToTest = {2, 4, 8, 16, 32}; 064 logger.info("# Start CRUD concurrent performance testing..."); 065 for (final int aNumThreadsToTest : numThreadsToTest) { 066 startCrudConcurrentPerformanceTest(aNumThreadsToTest); 067 } 068 } 069 070 /** 071 * Test CRUD concurrent access performance: 072 * create/update/delete object, create/update/delete content file 073 * 074 * @param numThreads to run 075 * @throws Exception on error 076 */ 077 private void startCrudConcurrentPerformanceTest(final int numThreads) throws Exception { 078 String pid = null; 079 080 // Tasks to run 081 final List<HttpRunner> tasks = new ArrayList<>(); 082 final List<String> pids = new ArrayList<>(); 083 084 // Create object 085 logger.info("# Starting " + numThreads + " concurrent threads to create object..."); 086 for (int i = 0; i < numThreads; i++) { 087 pid = getRandomUniqueId(); 088 pids.add(pid); 089 final String taskName = "Thread " + (i + 1) + " to create object " + pid; 090 final HttpRequestBase request = postObjMethod("/"); 091 request.addHeader("Slug", pid); 092 final HttpRunner task = new HttpRunner(request, taskName); 093 task.setExpectedStatusCode(CREATED.getStatusCode()); 094 tasks.add(task); 095 } 096 startThreads(tasks); 097 long totalResponseTime = getTotalResponseTime(numThreads, tasks); 098 logger.info("** Average response time for {} concurrent threads to CREATE object: {} ms", 099 numThreads, 100 totalResponseTime / numThreads); 101 102 tasks.clear(); 103 // Update objects 104 logger.info("# Starting " + numThreads + " concurrent threads to update object..."); 105 for (int i = 0; i < numThreads; i++) { 106 pid = pids.get(i); 107 final String taskName = "Thread " + (i + 1) + " to update object"; 108 final HttpPatch request = patchObjMethod(pid); 109 request.addHeader(CONTENT_TYPE, "application/sparql-update"); 110 final String subjectUri = request.getURI().toString(); 111 final BasicHttpEntity e = new BasicHttpEntity(); 112 e.setContent(new ByteArrayInputStream( 113 ("INSERT { <" + subjectUri + "> <http://purl.org/dc/elements/1.1/title> " 114 + "\"Title: " + taskName + pid + "\" } WHERE {}" 115 ).getBytes())); 116 request.setEntity(e); 117 final HttpRunner task = new HttpRunner(request, taskName); 118 task.setExpectedStatusCode(204); 119 tasks.add(task); 120 } 121 startThreads(tasks); 122 totalResponseTime = getTotalResponseTime(numThreads, tasks); 123 logger.info("** Average response time for {} concurrent threads to UPDATE object: {} ms", 124 numThreads, 125 totalResponseTime / numThreads); 126 127 128 tasks.clear(); 129 // Ingest new content 130 logger.info("# Starting " + numThreads + " concurrent threads to inget content..."); 131 for (int i = 0; i < numThreads; i++) { 132 pid = pids.get(i); 133 final String taskName = "Thread " + (i + 1) + " to ingest content file to object"; 134 final HttpRequestBase request = putDSMethod(pid, "ds", "This is a content file: " + taskName + pid); 135 final HttpRunner task = new HttpRunner(request, taskName); 136 task.setExpectedStatusCode(CREATED.getStatusCode()); 137 tasks.add(task); 138 } 139 startThreads(tasks); 140 totalResponseTime = getTotalResponseTime(numThreads, tasks); 141 logger.info("** Average response time for {} concurrent threads to INGEST content file: {} ms", 142 numThreads, 143 totalResponseTime / numThreads); 144 145 146 tasks.clear(); 147 // Update content 148 logger.info("# Starting " + numThreads + " concurrent threads to update content..."); 149 for (int i = 0; i < numThreads; i++) { 150 pid = pids.get(i); 151 final String taskName = "Thread " + (i + 1) + " to update content file in object"; 152 final HttpRequestBase request = putDSMethod(pid, 153 "ds", 154 "This is an updated content file: " + taskName + pid); 155 final HttpRunner task = new HttpRunner(request, taskName); 156 task.setExpectedStatusCode(204); 157 tasks.add(task); 158 } 159 startThreads(tasks); 160 totalResponseTime = getTotalResponseTime(numThreads, tasks); 161 logger.info("** Average response time for {} concurrent threads to UPDATE content file: {} ms", 162 numThreads, 163 totalResponseTime / numThreads); 164 165 166 tasks.clear(); 167 // Retrieve content 168 logger.info("# Starting " + numThreads + " concurrent threads to retrieve content..."); 169 for (int i = 0; i < numThreads; i++) { 170 pid = pids.get(i); 171 final String taskName = "Thread " + (i + 1) + " to retrieve content file in object"; 172 final HttpRequestBase request = getDSMethod(pid, "ds"); 173 final HttpRunner task = new HttpRunner(request, taskName); 174 task.setExpectedStatusCode(200); 175 tasks.add(task); 176 } 177 startThreads(tasks); 178 totalResponseTime = getTotalResponseTime(numThreads, tasks); 179 logger.info("** Average response time for {} concurrent threads to RETRIEVE content file: {} ms", 180 numThreads, 181 totalResponseTime / numThreads); 182 183 184 tasks.clear(); 185 // Delete content file 186 logger.info("# Starting " + numThreads + " concurrent threads to delete content file..."); 187 for (int i = 0; i < numThreads; i++) { 188 pid = pids.get(i); 189 final String taskName = "Thread " + (i + 1) + " to delete content file in object"; 190 final HttpRequestBase request = deleteObjMethod(pid + "/ds"); 191 final HttpRunner task = new HttpRunner(request, taskName); 192 task.setExpectedStatusCode(204); 193 tasks.add(task); 194 } 195 startThreads(tasks); 196 totalResponseTime = getTotalResponseTime(numThreads, tasks); 197 logger.info("** Average response time for {} concurrent threads to DELETE content file: {} ms", 198 numThreads, 199 totalResponseTime / numThreads); 200 201 202 tasks.clear(); 203 // Retrieve objects 204 logger.info("# Starting " + numThreads + " concurrent threads to retrieve object..."); 205 for (int i = 0; i < numThreads; i++) { 206 pid = pids.get(i); 207 final String taskName = "Thread " + (i + 1) + " to retrieve object"; 208 final HttpGet request = getObjMethod(pid); 209 final HttpRunner task = new HttpRunner(request, taskName); 210 task.setExpectedStatusCode(200); 211 tasks.add(task); 212 } 213 startThreads(tasks); 214 totalResponseTime = getTotalResponseTime(numThreads, tasks); 215 logger.info("** Average response time for {} concurrent threads to RETRIEVE object: {} ms", 216 numThreads, 217 totalResponseTime / numThreads); 218 219 220 tasks.clear(); 221 // Delete objects 222 logger.info("# Starting " + numThreads + " concurrent threads to delete object..."); 223 for (int i = 0; i < numThreads; i++) { 224 pid = pids.get(i); 225 final String taskName = "Thread " + (i + 1) + " to delete object"; 226 final HttpRequestBase request = deleteObjMethod(pid); 227 final HttpRunner task = new HttpRunner(request, taskName); 228 task.setExpectedStatusCode(204); 229 tasks.add(task); 230 } 231 startThreads(tasks); 232 totalResponseTime = getTotalResponseTime(numThreads, tasks); 233 logger.info("** Average response time for {} concurrent threads to DELETE object: {} ms", 234 numThreads, 235 totalResponseTime / numThreads); 236 237 } 238 239 private static long getTotalResponseTime(final int numThreads, 240 final List<HttpRunner> tasks) throws InterruptedException { 241 Thread.sleep(1000); 242 long totalResponseTime = 0; 243 for (int i = 0; i < numThreads; i++) { 244 totalResponseTime += tasks.get(i).responseTime; 245 } 246 return totalResponseTime; 247 } 248 249 private static void startThreads(final List<HttpRunner> tasks) throws InterruptedException { 250 for (final HttpRunner task : tasks) { 251 252 final Thread thread = new Thread(task); 253 thread.run(); 254 thread.join(); 255 } 256 } 257 258 /** 259 * Task to run http request for CRUD concurrent performance test. 260 * 261 * @author lsitu 262 */ 263 class HttpRunner implements Runnable { 264 265 private HttpClient httpClient = null; 266 267 private HttpResponse response = null; 268 269 private HttpRequestBase request = null; 270 271 private String taskName = null; 272 273 private long responseTime = 0; 274 275 private int statusCode = 0; 276 277 private int expectedStatusCode = 0; 278 279 public HttpRunner(final HttpRequestBase request, final String taskName) { 280 this.taskName = taskName; 281 this.request = request; 282 // Use its own HttpClient instance to make sure each performance test 283 // won't affected by a single HttpClient instance with multiple connections. 284 httpClient = createClient(); 285 } 286 287 @Override 288 public void run() { 289 try { 290 final long startTime = System.currentTimeMillis(); 291 response = httpClient.execute(request); 292 final long endTime = System.currentTimeMillis(); 293 responseTime = endTime - startTime; 294 statusCode = response.getStatusLine().getStatusCode(); 295 logger.info("{} {} with status {} in {} ms.", 296 taskName, request.getURI().toString(), 297 statusCode, String.valueOf(responseTime)); 298 assertEquals(taskName + " exited abnormally.", expectedStatusCode, statusCode); 299 } catch (final IOException e) { 300 logger.error("Error {} {} got IOException: {}", taskName, request.getURI().toString(), e.getMessage()); 301 } finally { 302 request.releaseConnection(); 303 } 304 } 305 306 public HttpResponse getResponse() { 307 return response; 308 } 309 310 311 public HttpRequestBase getRequest() { 312 return request; 313 } 314 315 316 public int getStatusCode() { 317 return statusCode; 318 } 319 320 321 public void setExpectedStatusCode(final int expectedStatusCode) { 322 this.expectedStatusCode = expectedStatusCode; 323 } 324 325 } 326}