7332 artifact pipeline work

This commit is contained in:
Richard Cordovano
2021-06-07 12:04:35 -04:00
parent 9944c0c0eb
commit dbfecb626b
2 changed files with 277 additions and 275 deletions
@@ -174,7 +174,7 @@ public final class IngestJob {
}
// Streaming ingest jobs will only have one data source
IngestJobPipeline streamingIngestPipeline = ingestJobPipelines.values().iterator().next();
streamingIngestPipeline.notifyStreamedDataSourceReady();
streamingIngestPipeline.addStreamedDataSource();
}
/**
@@ -100,18 +100,21 @@ final class IngestJobPipeline {
*/
INITIALIZATION,
/*
* The pipeline is running file ingest modules on files streamed to it
* by a data source processor. If configured to have data artifact
* ingest modules, it is running them on artifacts generated by the
* analysis of the streamed files.
* This stage is unique to a streaming mode ingest job. The pipeline is
* running file ingest modules on files streamed to it via
* addStreamedFiles(). If configured to have data artifact ingest
* modules, the pipeline is also running them on data artifacts
* generated by the analysis of the streamed files. This stage ends when
* the data source is streamed to the pipeline via
* addStreamedDataSource().
*/
FIRST_STAGE_FILE_STREAMING,
FIRST_STAGE_STREAMING,
/*
* The pipeline is running the following three types of ingest modules:
* higher priority data source level ingest modules, file ingest
* modules, and data artifact ingest modules.
*/
FIRST_STAGE_ALL_TASKS,
FIRST_STAGE,
/**
* The pipeline is running lower priority, usually long-running, data
* source level ingest modules and data artifact ingest modules.
@@ -124,6 +127,13 @@ final class IngestJobPipeline {
};
private volatile Stages stage = IngestJobPipeline.Stages.INITIALIZATION;
/*
* The stage field is volatile to allow it to be read by multiple threads.
* This lock is used not to guard the stage field, but to make stage
* transitions atomic.
*/
private final Object stageTransitionLock = new Object();
/*
* An ingest pipeline has separate data source level ingest module pipelines
* for the first and second stages. Longer running, lower priority modules
@@ -260,7 +270,7 @@ final class IngestJobPipeline {
* @param jythonModules The ingest module templates for modules implemented
* using Jython.
*/
private static void sortModuleTemplates(final List<IngestModuleTemplate> sortedModules, final Map<String, IngestModuleTemplate> javaModules, final Map<String, IngestModuleTemplate> jythonModules) {
private static void addToIngestPipelineTemplate(final List<IngestModuleTemplate> sortedModules, final Map<String, IngestModuleTemplate> javaModules, final Map<String, IngestModuleTemplate> jythonModules) {
final List<IngestModuleTemplate> autopsyModules = new ArrayList<>();
final List<IngestModuleTemplate> thirdPartyModules = new ArrayList<>();
Stream.concat(javaModules.entrySet().stream(), jythonModules.entrySet().stream()).forEach((templateEntry) -> {
@@ -304,7 +314,7 @@ final class IngestJobPipeline {
* @param jythonMapping Mapping for Jython ingest module templates.
* @param template The ingest module template.
*/
private static void addModuleTemplateToImplLangMap(Map<String, IngestModuleTemplate> mapping, Map<String, IngestModuleTemplate> jythonMapping, IngestModuleTemplate template) {
private static void addModuleTemplateToSortingMap(Map<String, IngestModuleTemplate> mapping, Map<String, IngestModuleTemplate> jythonMapping, IngestModuleTemplate template) {
String className = template.getModuleFactory().getClass().getCanonicalName();
String jythonName = getModuleNameFromJythonClassName(className);
if (jythonName != null) {
@@ -323,25 +333,14 @@ final class IngestJobPipeline {
private void createIngestModulePipelines() throws InterruptedException {
/*
* Get the enabled ingest module templates from the ingest job settings.
* An ingest module template combines an ingest module factory with job
* level ingest module settings to support the creation of any number of
* fully configured instances of a given ingest module. An ingest module
* factory may be able to create multiple types of ingest modules.
*/
List<IngestModuleTemplate> enabledTemplates = settings.getEnabledIngestModuleTemplates();
/**
* Sort the ingest module templates into buckets based on the module
* types the ingest module factory can create. A template may go into
* more than one bucket. The buckets are actually maps of ingest module
* factory class names to ingest module templates. These maps are used
* to go from an ingest module factory class name read from the pipeline
* configuration file to the corresponding ingest module template.
*
* There are actually two maps for each module type bucket. One map is
* for Java modules and the other one is for Jython modules. The
* templates are separated this way so that Java modules that are not in
* the pipeline config file can be placed before the Jython modules.
* types the template can be used to create. A template may go into more
* than one bucket. Each bucket actually consists of two collections:
* one for Java modules and one for Jython modules.
*/
Map<String, IngestModuleTemplate> javaDataSourceModuleTemplates = new LinkedHashMap<>();
Map<String, IngestModuleTemplate> jythonDataSourceModuleTemplates = new LinkedHashMap<>();
@@ -351,86 +350,79 @@ final class IngestJobPipeline {
Map<String, IngestModuleTemplate> jythonArtifactModuleTemplates = new LinkedHashMap<>();
for (IngestModuleTemplate template : enabledTemplates) {
if (template.isDataSourceIngestModuleTemplate()) {
addModuleTemplateToImplLangMap(javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates, template);
addModuleTemplateToSortingMap(javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates, template);
}
if (template.isFileIngestModuleTemplate()) {
addModuleTemplateToImplLangMap(javaFileModuleTemplates, jythonFileModuleTemplates, template);
addModuleTemplateToSortingMap(javaFileModuleTemplates, jythonFileModuleTemplates, template);
}
if (template.isDataArtifactIngestModuleTemplate()) {
addModuleTemplateToImplLangMap(javaArtifactModuleTemplates, jythonArtifactModuleTemplates, template);
addModuleTemplateToSortingMap(javaArtifactModuleTemplates, jythonArtifactModuleTemplates, template);
}
}
/**
* Take the module templates that have pipeline configuration file
* entries out of the buckets and put them in lists representing ingest
* task pipelines, in the order prescribed by the file. Note that the
* pipeline configuration file currently only supports specifying data
* source level and file ingest module pipeline layouts.
* Take the module templates that have pipeline configuration entries
* out of the buckets and add them to ingest module pipeline templates
* in the order prescribed by the pipeline configuration.
*/
IngestPipelinesConfiguration pipelineConfig = IngestPipelinesConfiguration.getInstance();
List<IngestModuleTemplate> firstStageDataSourceModuleTemplates = addConfiguredIngestModuleTemplates(javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates, pipelineConfig.getStageOneDataSourceIngestPipelineConfig());
List<IngestModuleTemplate> secondStageDataSourceModuleTemplates = addConfiguredIngestModuleTemplates(javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates, pipelineConfig.getStageTwoDataSourceIngestPipelineConfig());
List<IngestModuleTemplate> fileIngestModuleTemplates = addConfiguredIngestModuleTemplates(javaFileModuleTemplates, jythonFileModuleTemplates, pipelineConfig.getFileIngestPipelineConfig());
List<IngestModuleTemplate> artifactModuleTemplates = new ArrayList<>();
List<IngestModuleTemplate> firstStageDataSourcePipelineTemplate = createIngestPipelineTemplate(javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates, pipelineConfig.getStageOneDataSourceIngestPipelineConfig());
List<IngestModuleTemplate> secondStageDataSourcePipelineTemplate = createIngestPipelineTemplate(javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates, pipelineConfig.getStageTwoDataSourceIngestPipelineConfig());
List<IngestModuleTemplate> filePipelineTemplate = createIngestPipelineTemplate(javaFileModuleTemplates, jythonFileModuleTemplates, pipelineConfig.getFileIngestPipelineConfig());
List<IngestModuleTemplate> artifactPipelineTemplate = new ArrayList<>();
/**
* Add any module templates remaining in the buckets to the appropriate
* ingest task pipeline. Note that any data source level ingest modules
* that were not listed in the configuration file are added to the first
* stage data source pipeline, Java modules are added before Jython
* modules, and Core Autopsy modules are added before third party
* modules.
* Add any ingest module templates remaining in the buckets to the
* appropriate ingest module pipeline templates. Data source level
* ingest modules templates that were not listed in the pipeline
* configuration are added to the first stage data source pipeline
* template, Java modules are added before Jython modules and Core
* Autopsy modules are added before third party modules.
*/
sortModuleTemplates(firstStageDataSourceModuleTemplates, javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates);
sortModuleTemplates(fileIngestModuleTemplates, javaFileModuleTemplates, jythonFileModuleTemplates);
sortModuleTemplates(artifactModuleTemplates, javaArtifactModuleTemplates, jythonArtifactModuleTemplates);
addToIngestPipelineTemplate(firstStageDataSourcePipelineTemplate, javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates);
addToIngestPipelineTemplate(filePipelineTemplate, javaFileModuleTemplates, jythonFileModuleTemplates);
addToIngestPipelineTemplate(artifactPipelineTemplate, javaArtifactModuleTemplates, jythonArtifactModuleTemplates);
/**
* Construct the actual ingest task pipelines from the ordered lists.
* Construct the ingest module pipelines from the ingest module pipeline
* templates.
*/
firstStageDataSourceIngestPipeline = new DataSourceIngestPipeline(this, firstStageDataSourceModuleTemplates);
secondStageDataSourceIngestPipeline = new DataSourceIngestPipeline(this, secondStageDataSourceModuleTemplates);
firstStageDataSourceIngestPipeline = new DataSourceIngestPipeline(this, firstStageDataSourcePipelineTemplate);
secondStageDataSourceIngestPipeline = new DataSourceIngestPipeline(this, secondStageDataSourcePipelineTemplate);
int numberOfFileIngestThreads = IngestManager.getInstance().getNumberOfFileIngestThreads();
for (int i = 0; i < numberOfFileIngestThreads; ++i) {
FileIngestPipeline pipeline = new FileIngestPipeline(this, fileIngestModuleTemplates);
FileIngestPipeline pipeline = new FileIngestPipeline(this, filePipelineTemplate);
fileIngestPipelinesQueue.put(pipeline);
fileIngestPipelines.add(pipeline);
}
artifactIngestPipeline = new DataArtifactIngestPipeline(this, artifactModuleTemplates);
artifactIngestPipeline = new DataArtifactIngestPipeline(this, artifactPipelineTemplate);
}
/**
* Uses an input collection of ingest module templates and a pipeline
* configuration, i.e., an ordered list of ingest module factory class
* names, to create an ordered output list of ingest module templates for an
* ingest task pipeline. The ingest module templates are removed from the
* input collection as they are added to the output collection.
* Creates an ingest module pipeline template that can be used to construct
* an ingest module pipeline.
*
* @param javaIngestModuleTemplates A mapping of Java ingest module
* factory class names to ingest module
* templates.
* @param jythonIngestModuleTemplates A mapping of Jython ingest module
* factory proxy class names to ingest
* module templates.
* @param pipelineConfig An ordered list of ingest module
* factory class names representing an
* ingest pipeline, read from the
* pipeline configuration file.
* @param javaIngestModuleTemplates Ingest module templates for ingest
* modules implemented using Java.
* @param jythonIngestModuleTemplates Ingest module templates for ingest
* modules implemented using Jython.
* @param pipelineConfig An ordered list of the ingest modules
* that belong in the ingest pipeline for
* which the template is being created.
*
* @return An ordered list of ingest module templates, i.e., an
* uninstantiated pipeline.
* @return An ordered list of ingest module templates, i.e., a template for
* creating ingest module pipelines.
*/
private static List<IngestModuleTemplate> addConfiguredIngestModuleTemplates(Map<String, IngestModuleTemplate> javaIngestModuleTemplates, Map<String, IngestModuleTemplate> jythonIngestModuleTemplates, List<String> pipelineConfig) {
List<IngestModuleTemplate> templates = new ArrayList<>();
private static List<IngestModuleTemplate> createIngestPipelineTemplate(Map<String, IngestModuleTemplate> javaIngestModuleTemplates, Map<String, IngestModuleTemplate> jythonIngestModuleTemplates, List<String> pipelineConfig) {
List<IngestModuleTemplate> pipelineTemplate = new ArrayList<>();
for (String moduleClassName : pipelineConfig) {
if (javaIngestModuleTemplates.containsKey(moduleClassName)) {
templates.add(javaIngestModuleTemplates.remove(moduleClassName));
pipelineTemplate.add(javaIngestModuleTemplates.remove(moduleClassName));
} else if (jythonIngestModuleTemplates.containsKey(moduleClassName)) {
templates.add(jythonIngestModuleTemplates.remove(moduleClassName));
pipelineTemplate.add(jythonIngestModuleTemplates.remove(moduleClassName));
}
}
return templates;
return pipelineTemplate;
}
/**
@@ -533,13 +525,7 @@ final class IngestJobPipeline {
* @return True or false.
*/
boolean hasFileIngestModules() {
if (!fileIngestPipelines.isEmpty()) {
/*
* Note that the file ingest task pipelines are identical.
*/
return !fileIngestPipelines.get(0).isEmpty();
}
return false;
return (fileIngestPipelines.isEmpty() == false);
}
/**
@@ -563,9 +549,9 @@ final class IngestJobPipeline {
recordIngestJobStartUpInfo();
if (hasFirstStageDataSourceIngestModules() || hasFileIngestModules() || hasDataArtifactIngestModules()) {
if (parentJob.getIngestMode() == IngestJob.Mode.STREAMING) {
startFileStreaming();
startFirstStageInStreamingMode();
} else {
startFirstStage();
startFirstStageInBatchMode();
}
} else if (hasSecondStageDataSourceIngestModules()) {
startSecondStage();
@@ -680,58 +666,62 @@ final class IngestJobPipeline {
* of the files in the data source (excepting carved and derived files) have
* already been added to the case database by the data source processor.
*/
private void startFirstStage() {
logInfoMessage("Starting first stage analysis in batch mode"); //NON-NLS
stage = Stages.FIRST_STAGE_ALL_TASKS;
private void startFirstStageInBatchMode() {
synchronized (stageTransitionLock) {
logInfoMessage("Starting first stage analysis in batch mode"); //NON-NLS
stage = Stages.FIRST_STAGE;
/*
* Do a count of the files the data source processor has added to the
* case database. This estimate will be used for ingest progress
* snapshots and for the file ingest progress bar if running with a GUI.
*/
if (hasFileIngestModules()) {
long filesToProcess = dataSource.accept(new GetFilesCountVisitor());;
synchronized (fileIngestProgressLock) {
estimatedFilesToProcess = filesToProcess;
}
}
/*
* If running with a GUI, start ingest progress bars in the lower right
* hand corner of the main application window.
*/
if (doUI) {
/*
* Do a count of the files the data source processor has added to
* the case database. This estimate will be used for ingest progress
* snapshots and for the file ingest progress bar if running with a
* GUI.
*/
if (hasFileIngestModules()) {
startFileIngestProgressBar();
long filesToProcess = dataSource.accept(new GetFilesCountVisitor());;
synchronized (fileIngestProgressLock) {
estimatedFilesToProcess = filesToProcess;
}
}
if (hasFirstStageDataSourceIngestModules()) {
startDataSourceIngestProgressBar();
}
if (hasDataArtifactIngestModules()) {
startArtifactIngestProgressBar();
}
}
/*
* Make the first stage data source level ingest pipeline the current
* data source level pipeline.
*/
currentDataSourceIngestPipeline = firstStageDataSourceIngestPipeline;
/*
* If running with a GUI, start ingest progress bars in the lower
* right hand corner of the main application window.
*/
if (doUI) {
if (hasFileIngestModules()) {
startFileIngestProgressBar();
}
if (hasFirstStageDataSourceIngestModules()) {
startDataSourceIngestProgressBar();
}
if (hasDataArtifactIngestModules()) {
startArtifactIngestProgressBar();
}
}
/*
* Schedule the first stage ingest tasks and then immediately check for
* stage completion. This is necessary because it is possible that zero
* tasks will actually make it to task execution due to the file filter
* or other ingest job settings. In that case, there will never be a
* stage completion check in an ingest thread executing an ingest task,
* so such a job would run forever without a check here.
*/
if (!files.isEmpty() && hasFileIngestModules()) {
taskScheduler.scheduleFileIngestTasks(this, files);
} else if (hasFirstStageDataSourceIngestModules() || hasFileIngestModules() || hasDataArtifactIngestModules()) {
taskScheduler.scheduleIngestTasks(this);
/*
* Make the first stage data source level ingest pipeline the
* current data source level pipeline.
*/
currentDataSourceIngestPipeline = firstStageDataSourceIngestPipeline;
/*
* Schedule the first stage ingest tasks and then immediately check
* for stage completion. This is necessary because it is possible
* that zero tasks will actually make it to task execution due to
* the file filter or other ingest job settings. In that case, there
* will never be a stage completion check in an ingest thread
* executing an ingest task, so such a job would run forever without
* a check here.
*/
if (!files.isEmpty() && hasFileIngestModules()) {
taskScheduler.scheduleFileIngestTasks(this, files);
} else if (hasFirstStageDataSourceIngestModules() || hasFileIngestModules() || hasDataArtifactIngestModules()) {
taskScheduler.scheduleIngestTasks(this);
}
checkForStageCompleted();
}
checkForStageCompleted();
}
/**
@@ -740,38 +730,40 @@ final class IngestJobPipeline {
* adds them to the case database and file level analysis can begin before
* data source level analysis.
*/
private void startFileStreaming() {
logInfoMessage("Starting first stage analysis in streaming mode"); //NON-NLS
stage = Stages.FIRST_STAGE_FILE_STREAMING;
private void startFirstStageInStreamingMode() {
synchronized (stageTransitionLock) {
logInfoMessage("Starting first stage analysis in streaming mode"); //NON-NLS
stage = Stages.FIRST_STAGE_STREAMING;
if (doUI) {
/*
* If running with a GUI, start ingest progress bars in the lower
* right hand corner of the main application window.
*/
if (hasFileIngestModules()) {
if (doUI) {
/*
* Note that because estimated files remaining to process still
* has its initial value of zero, the progress bar will start in
* the "indeterminate" state. An estimate of the files to
* process can be computed in
* If running with a GUI, start ingest progress bars in the
* lower right hand corner of the main application window.
*/
startFileIngestProgressBar();
if (hasFileIngestModules()) {
/*
* Note that because estimated files remaining to process
* still has its initial value of zero, the progress bar
* will start in the "indeterminate" state. An estimate of
* the files to process can be computed in
*/
startFileIngestProgressBar();
}
if (hasDataArtifactIngestModules()) {
startArtifactIngestProgressBar();
}
}
if (hasDataArtifactIngestModules()) {
startArtifactIngestProgressBar();
}
}
if (hasDataArtifactIngestModules()) {
/*
* Schedule artifact ingest tasks for any artifacts currently in the
* case database. This needs to be done before any files or the data
* source are streamed in to avoid analyzing data artifacts added to
* the case database by the data source level or file level ingest
* tasks.
*/
taskScheduler.scheduleDataArtifactIngestTasks(this);
if (hasDataArtifactIngestModules()) {
/*
* Schedule artifact ingest tasks for any artifacts currently in
* the case database. This needs to be done before any files or
* the data source are streamed in to avoid analyzing data
* artifacts added to the case database by the data source level
* or file level ingest tasks.
*/
taskScheduler.scheduleDataArtifactIngestTasks(this);
}
}
}
@@ -779,46 +771,49 @@ final class IngestJobPipeline {
* Notifies the ingest pipeline running in streaming mode that the data
* source is now ready for analysis.
*/
void notifyStreamedDataSourceReady() {
logInfoMessage("Starting full first stage analysis in streaming mode"); //NON-NLS
stage = IngestJobPipeline.Stages.FIRST_STAGE_ALL_TASKS;
currentDataSourceIngestPipeline = firstStageDataSourceIngestPipeline;
void addStreamedDataSource() {
synchronized (stageTransitionLock) {
logInfoMessage("Starting full first stage analysis in streaming mode"); //NON-NLS
stage = IngestJobPipeline.Stages.FIRST_STAGE;
currentDataSourceIngestPipeline = firstStageDataSourceIngestPipeline;
/*
* Do a count of the files the data source processor has added to the
* case database. This estimate will be used for ingest progress
* snapshots and for the file ingest progress bar if running with a GUI.
*/
long filesToProcess = dataSource.accept(new GetFilesCountVisitor()) - processedFiles;
synchronized (fileIngestProgressLock) {
if (processedFiles <= filesToProcess) {
filesToProcess -= processedFiles;
}
estimatedFilesToProcess = filesToProcess;
if (doUI && fileIngestProgressBar != null) {
fileIngestProgressBar.switchToDeterminate((int) estimatedFilesToProcess);
}
}
if (doUI) {
if (hasFirstStageDataSourceIngestModules()) {
startDataSourceIngestProgressBar();
}
}
currentDataSourceIngestPipeline = firstStageDataSourceIngestPipeline;
if (hasFirstStageDataSourceIngestModules()) {
IngestJobPipeline.taskScheduler.scheduleDataSourceIngestTask(this);
} else {
/*
* If no data source level ingest task is scheduled at this time and
* all of the file level and artifact ingest tasks scheduled during
* the initial file streaming stage have already executed, there
* will never be a stage completion check in an ingest thread
* executing an ingest task, so such a job would run forever without
* a check here.
* Do a count of the files the data source processor has added to
* the case database. This estimate will be used for ingest progress
* snapshots and for the file ingest progress bar if running with a
* GUI.
*/
checkForStageCompleted();
long filesToProcess = dataSource.accept(new GetFilesCountVisitor()) - processedFiles;
synchronized (fileIngestProgressLock) {
if (processedFiles <= filesToProcess) {
filesToProcess -= processedFiles;
}
estimatedFilesToProcess = filesToProcess;
if (doUI && fileIngestProgressBar != null) {
fileIngestProgressBar.switchToDeterminate((int) estimatedFilesToProcess);
}
}
if (doUI) {
if (hasFirstStageDataSourceIngestModules()) {
startDataSourceIngestProgressBar();
}
}
currentDataSourceIngestPipeline = firstStageDataSourceIngestPipeline;
if (hasFirstStageDataSourceIngestModules()) {
IngestJobPipeline.taskScheduler.scheduleDataSourceIngestTask(this);
} else {
/*
* If no data source level ingest task is scheduled at this time
* and all of the file level and artifact ingest tasks scheduled
* during the initial file streaming stage have already
* executed, there will never be a stage completion check in an
* ingest thread executing an ingest task, so such a job would
* run forever without a check here.
*/
checkForStageCompleted();
}
}
}
@@ -826,16 +821,18 @@ final class IngestJobPipeline {
* Starts the second stage ingest task pipelines.
*/
private void startSecondStage() {
if (!cancelled && hasSecondStageDataSourceIngestModules()) {
logInfoMessage(String.format("Starting second stage ingest task pipelines for %s (objID=%d, jobID=%d)", dataSource.getName(), dataSource.getId(), parentJob.getId())); //NON-NLS
stage = IngestJobPipeline.Stages.SECOND_STAGE;
synchronized (stageTransitionLock) {
if (hasSecondStageDataSourceIngestModules()) {
logInfoMessage(String.format("Starting second stage ingest task pipelines for %s (objID=%d, jobID=%d)", dataSource.getName(), dataSource.getId(), parentJob.getId())); //NON-NLS
stage = IngestJobPipeline.Stages.SECOND_STAGE;
if (doUI) {
startDataSourceIngestProgressBar();
if (doUI) {
startDataSourceIngestProgressBar();
}
currentDataSourceIngestPipeline = secondStageDataSourceIngestPipeline;
taskScheduler.scheduleDataSourceIngestTask(this);
}
currentDataSourceIngestPipeline = secondStageDataSourceIngestPipeline;
taskScheduler.scheduleDataSourceIngestTask(this);
}
}
@@ -925,17 +922,19 @@ final class IngestJobPipeline {
* completed and does a stage transition if they are.
*/
private void checkForStageCompleted() {
if (stage == Stages.FIRST_STAGE_FILE_STREAMING) {
return;
}
if (taskScheduler.currentTasksAreCompleted(this)) {
switch (stage) {
case FIRST_STAGE_ALL_TASKS:
finishFirstStage();
break;
case SECOND_STAGE:
shutDown();
break;
synchronized (stageTransitionLock) {
if (stage == Stages.FIRST_STAGE_STREAMING) {
return;
}
if (taskScheduler.currentTasksAreCompleted(this)) {
switch (stage) {
case FIRST_STAGE:
finishFirstStage();
break;
case SECOND_STAGE:
shutDown();
break;
}
}
}
}
@@ -945,88 +944,92 @@ final class IngestJobPipeline {
* job and starts the second stage, if appropriate.
*/
private void finishFirstStage() {
logInfoMessage("Finished first stage analysis"); //NON-NLS
synchronized (stageTransitionLock) {
logInfoMessage("Finished first stage analysis"); //NON-NLS
shutDownIngestModulePipeline(currentDataSourceIngestPipeline);
while (!fileIngestPipelinesQueue.isEmpty()) {
FileIngestPipeline pipeline = fileIngestPipelinesQueue.poll();
shutDownIngestModulePipeline(pipeline);
}
shutDownIngestModulePipeline(currentDataSourceIngestPipeline);
while (!fileIngestPipelinesQueue.isEmpty()) {
FileIngestPipeline pipeline = fileIngestPipelinesQueue.poll();
shutDownIngestModulePipeline(pipeline);
}
if (doUI) {
synchronized (dataSourceIngestProgressLock) {
if (dataSourceIngestProgressBar != null) {
dataSourceIngestProgressBar.finish();
dataSourceIngestProgressBar = null;
if (doUI) {
synchronized (dataSourceIngestProgressLock) {
if (dataSourceIngestProgressBar != null) {
dataSourceIngestProgressBar.finish();
dataSourceIngestProgressBar = null;
}
}
synchronized (fileIngestProgressLock) {
if (fileIngestProgressBar != null) {
fileIngestProgressBar.finish();
fileIngestProgressBar = null;
}
}
}
synchronized (fileIngestProgressLock) {
if (fileIngestProgressBar != null) {
fileIngestProgressBar.finish();
fileIngestProgressBar = null;
}
if (!cancelled && hasSecondStageDataSourceIngestModules()) {
startSecondStage();
} else {
shutDown();
}
}
if (!cancelled && hasSecondStageDataSourceIngestModules()) {
startSecondStage();
} else {
shutDown();
}
}
/**
* Shuts down the ingest module pipelines and progress bars for this job.
*/
private void shutDown() {
logInfoMessage("Finished all tasks"); //NON-NLS
stage = IngestJobPipeline.Stages.FINALIZATION;
synchronized (stageTransitionLock) {
logInfoMessage("Finished all tasks"); //NON-NLS
stage = IngestJobPipeline.Stages.FINALIZATION;
shutDownIngestModulePipeline(currentDataSourceIngestPipeline);
shutDownIngestModulePipeline(artifactIngestPipeline);
shutDownIngestModulePipeline(currentDataSourceIngestPipeline);
shutDownIngestModulePipeline(artifactIngestPipeline);
if (doUI) {
synchronized (dataSourceIngestProgressLock) {
if (dataSourceIngestProgressBar != null) {
dataSourceIngestProgressBar.finish();
dataSourceIngestProgressBar = null;
if (doUI) {
synchronized (dataSourceIngestProgressLock) {
if (dataSourceIngestProgressBar != null) {
dataSourceIngestProgressBar.finish();
dataSourceIngestProgressBar = null;
}
}
synchronized (fileIngestProgressLock) {
if (fileIngestProgressBar != null) {
fileIngestProgressBar.finish();
fileIngestProgressBar = null;
}
}
synchronized (artifactIngestProgressLock) {
if (artifactIngestProgressBar != null) {
artifactIngestProgressBar.finish();
artifactIngestProgressBar = null;
}
}
}
synchronized (fileIngestProgressLock) {
if (fileIngestProgressBar != null) {
fileIngestProgressBar.finish();
fileIngestProgressBar = null;
if (ingestJobInfo != null) {
if (cancelled) {
try {
ingestJobInfo.setIngestJobStatus(IngestJobStatusType.CANCELLED);
} catch (TskCoreException ex) {
logErrorMessage(Level.WARNING, "Failed to update ingest job status in case database", ex);
}
} else {
try {
ingestJobInfo.setIngestJobStatus(IngestJobStatusType.COMPLETED);
} catch (TskCoreException ex) {
logErrorMessage(Level.WARNING, "Failed to update ingest job status in case database", ex);
}
}
}
synchronized (artifactIngestProgressLock) {
if (artifactIngestProgressBar != null) {
artifactIngestProgressBar.finish();
artifactIngestProgressBar = null;
}
}
}
if (ingestJobInfo != null) {
if (cancelled) {
try {
ingestJobInfo.setIngestJobStatus(IngestJobStatusType.CANCELLED);
ingestJobInfo.setEndDateTime(new Date());
} catch (TskCoreException ex) {
logErrorMessage(Level.WARNING, "Failed to update ingest job status in case database", ex);
logErrorMessage(Level.WARNING, "Failed to set job end date in case database", ex);
}
} else {
try {
ingestJobInfo.setIngestJobStatus(IngestJobStatusType.COMPLETED);
} catch (TskCoreException ex) {
logErrorMessage(Level.WARNING, "Failed to update ingest job status in case database", ex);
}
}
try {
ingestJobInfo.setEndDateTime(new Date());
} catch (TskCoreException ex) {
logErrorMessage(Level.WARNING, "Failed to set job end date in case database", ex);
}
}
@@ -1173,7 +1176,7 @@ final class IngestJobPipeline {
*/
void addStreamedFiles(List<Long> fileObjIds) {
if (hasFileIngestModules()) {
if (stage.equals(Stages.FIRST_STAGE_FILE_STREAMING)) {
if (stage.equals(Stages.FIRST_STAGE_STREAMING)) {
IngestJobPipeline.taskScheduler.scheduleStreamedFileIngestTasks(this, fileObjIds);
} else {
logErrorMessage(Level.SEVERE, "Adding streaming files to job during stage " + stage.toString() + " not supported");
@@ -1189,8 +1192,8 @@ final class IngestJobPipeline {
* @param files A list of the files to add.
*/
void addFiles(List<AbstractFile> files) {
if (stage.equals(Stages.FIRST_STAGE_FILE_STREAMING)
|| stage.equals(Stages.FIRST_STAGE_ALL_TASKS)) {
if (stage.equals(Stages.FIRST_STAGE_STREAMING)
|| stage.equals(Stages.FIRST_STAGE)) {
taskScheduler.fastTrackFileIngestTasks(this, files);
} else {
logErrorMessage(Level.SEVERE, "Adding streaming files to job during stage " + stage.toString() + " not supported");
@@ -1214,8 +1217,8 @@ final class IngestJobPipeline {
*/
void addDataArtifacts(List<DataArtifact> artifacts) {
List<DataArtifact> artifactsToAnalyze = new ArrayList<>(artifacts);
if (stage.equals(Stages.FIRST_STAGE_FILE_STREAMING)
|| stage.equals(Stages.FIRST_STAGE_ALL_TASKS)
if (stage.equals(Stages.FIRST_STAGE_STREAMING)
|| stage.equals(Stages.FIRST_STAGE)
|| stage.equals(Stages.SECOND_STAGE)) {
taskScheduler.scheduleDataArtifactIngestTasks(this, artifactsToAnalyze);
} else {
@@ -1525,7 +1528,6 @@ final class IngestJobPipeline {
snapShotTime = new Date().getTime();
}
tasksSnapshot = taskScheduler.getTasksSnapshotForJob(pipelineId);
}
return new Snapshot(dataSource.getName(),