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}