ingestJobs now has its own lock (itself); ingestTasks still uses this for locking

This commit is contained in:
Samuel H. Kenyon
2014-04-23 16:10:49 -04:00
parent 8b9742ea91
commit dcfd2f276d
@@ -21,6 +21,7 @@ package org.sleuthkit.autopsy.ingest;
import java.beans.PropertyChangeListener;
import java.beans.PropertyChangeSupport;
import java.util.ArrayList;
import java.util.Collection;
import java.util.HashMap;
import java.util.List;
import java.util.concurrent.CancellationException;
@@ -61,8 +62,8 @@ public class IngestManager {
private final ExecutorService startIngestJobsExecutor = Executors.newSingleThreadExecutor();
private final ExecutorService dataSourceIngestTasksExecutor = Executors.newSingleThreadExecutor();
private final ExecutorService fileIngestTasksExecutor = Executors.newFixedThreadPool(MAX_NUMBER_OF_FILE_INGEST_THREADS);
private final HashMap<Long, IngestJob> ingestJobs = new HashMap<>(); // Maps job ids to jobs
private final HashMap<Long, Future<?>> ingestTasks = new HashMap<>(); // Maps task ids to task cancellation handles
private final HashMap<Long, IngestJob> ingestJobs = new HashMap<>(); // Maps job ids to jobs. Guarded by ingestJobs.
private final HashMap<Long, Future<?>> ingestTasks = new HashMap<>(); // Maps task ids to task cancellation handles. Guarded by this.
private AtomicLong ingestJobId = new AtomicLong(0L);
private AtomicLong ingestTaskId = new AtomicLong(0L);
private volatile IngestUI ingestMessageBox;
@@ -124,12 +125,17 @@ public class IngestManager {
*
* @return True if any ingest jobs are in progress, false otherwise
*/
public synchronized boolean isIngestRunning() {
return (ingestJobs.isEmpty() == false);
public boolean isIngestRunning() {
synchronized(ingestJobs) {
return (ingestJobs.isEmpty() == false);
}
}
synchronized void addFileToIngestJob(long ingestJobId, AbstractFile file) {
IngestJob job = ingestJobs.get(ingestJobId);
void addFileToIngestJob(long ingestJobId, AbstractFile file) {
IngestJob job;
synchronized(ingestJobs) {
job = ingestJobs.get(ingestJobId);
}
if (job != null) {
scheduler.getFileIngestScheduler().queueFile(job, file);
}
@@ -301,44 +307,54 @@ public class IngestManager {
}
}
private synchronized void stopIngestTasks() {
private void stopIngestTasks() {
// First mark all of the ingest jobs as cancelled. This way the
// ingest modules will know they are being shut down due to
// cancellation when the cancelled run ingest module tasks release
// their pipelines.
for (IngestJob job : ingestJobs.values()) {
Collection<IngestJob> jobList;
synchronized(ingestJobs) {
jobList = ingestJobs.values();
}
for (IngestJob job : jobList) {
job.cancel();
}
// Cancel the run ingest module tasks, setting the state of the threads
// running them to interrupted.
for (Future<?> task : ingestTasks.values()) {
task.cancel(true);
synchronized(this) {
for (Future<?> task : ingestTasks.values()) {
task.cancel(true);
}
}
// Jettision the remaining data source and file ingest tasks.
scheduler.getFileIngestScheduler().emptyQueues();
scheduler.getDataSourceIngestScheduler().emptyQueues();
}
synchronized void reportStartIngestJobsTaskDone(long taskId) {
private synchronized void reportStartIngestJobsTaskDone(long taskId) {
ingestTasks.remove(taskId);
}
synchronized void reportRunIngestModulesTaskDone(long taskId) {
ingestTasks.remove(taskId);
List<Long> completedJobs = new ArrayList<>();
for (IngestJob job : ingestJobs.values()) {
job.releaseIngestPipelinesForThread(taskId);
if (job.areIngestPipelinesShutDown() == true) {
completedJobs.add(job.getId());
}
private void reportRunIngestModulesTaskDone(long taskId) {
synchronized(this) {
ingestTasks.remove(taskId);
}
for (Long jobId : completedJobs) {
IngestJob job = ingestJobs.remove(jobId);
fireIngestJobEvent(job.isCancelled() ? IngestEvent.INGEST_JOB_CANCELLED.toString() : IngestEvent.INGEST_JOB_COMPLETED.toString(), jobId);
List<Long> completedJobs = new ArrayList<>();
synchronized(ingestJobs) {
for (IngestJob job : ingestJobs.values()) {
job.releaseIngestPipelinesForThread(taskId);
if (job.areIngestPipelinesShutDown() == true) {
completedJobs.add(job.getId());
}
}
for (Long jobId : completedJobs) {
IngestJob job = ingestJobs.remove(jobId);
fireIngestJobEvent(job.isCancelled() ? IngestEvent.INGEST_JOB_CANCELLED.toString() : IngestEvent.INGEST_JOB_COMPLETED.toString(), jobId);
}
}
}
@@ -384,7 +400,7 @@ public class IngestManager {
// Create an ingest job.
IngestJob ingestJob = new IngestJob(IngestManager.this.ingestJobId.incrementAndGet(), dataSource, moduleTemplates, processUnallocatedSpace);
synchronized (IngestManager.this) {
synchronized (ingestJobs) {
ingestJobs.put(ingestJob.getId(), ingestJob);
}
@@ -420,7 +436,7 @@ public class IngestManager {
"IngestManager.StartIngestJobsTask.run.startupErr.dlgTitle"), JOptionPane.ERROR_MESSAGE);
// Jettison the ingest job and move on to the next one.
synchronized (IngestManager.this) {
synchronized (ingestJobs) {
ingestJob.cancel();
ingestJobs.remove(ingestJob.getId());
}