From dbfecb626b5985273bee2a3b5e1a05dd0ae3059b Mon Sep 17 00:00:00 2001 From: Richard Cordovano Date: Mon, 7 Jun 2021 12:04:35 -0400 Subject: [PATCH] 7332 artifact pipeline work --- .../sleuthkit/autopsy/ingest/IngestJob.java | 2 +- .../autopsy/ingest/IngestJobPipeline.java | 550 +++++++++--------- 2 files changed, 277 insertions(+), 275 deletions(-) diff --git a/Core/src/org/sleuthkit/autopsy/ingest/IngestJob.java b/Core/src/org/sleuthkit/autopsy/ingest/IngestJob.java index c5704e873e..321d38c627 100644 --- a/Core/src/org/sleuthkit/autopsy/ingest/IngestJob.java +++ b/Core/src/org/sleuthkit/autopsy/ingest/IngestJob.java @@ -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(); } /** diff --git a/Core/src/org/sleuthkit/autopsy/ingest/IngestJobPipeline.java b/Core/src/org/sleuthkit/autopsy/ingest/IngestJobPipeline.java index f960aaf575..850ff26294 100644 --- a/Core/src/org/sleuthkit/autopsy/ingest/IngestJobPipeline.java +++ b/Core/src/org/sleuthkit/autopsy/ingest/IngestJobPipeline.java @@ -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 sortedModules, final Map javaModules, final Map jythonModules) { + private static void addToIngestPipelineTemplate(final List sortedModules, final Map javaModules, final Map jythonModules) { final List autopsyModules = new ArrayList<>(); final List 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 mapping, Map jythonMapping, IngestModuleTemplate template) { + private static void addModuleTemplateToSortingMap(Map mapping, Map 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 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 javaDataSourceModuleTemplates = new LinkedHashMap<>(); Map jythonDataSourceModuleTemplates = new LinkedHashMap<>(); @@ -351,86 +350,79 @@ final class IngestJobPipeline { Map 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 firstStageDataSourceModuleTemplates = addConfiguredIngestModuleTemplates(javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates, pipelineConfig.getStageOneDataSourceIngestPipelineConfig()); - List secondStageDataSourceModuleTemplates = addConfiguredIngestModuleTemplates(javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates, pipelineConfig.getStageTwoDataSourceIngestPipelineConfig()); - List fileIngestModuleTemplates = addConfiguredIngestModuleTemplates(javaFileModuleTemplates, jythonFileModuleTemplates, pipelineConfig.getFileIngestPipelineConfig()); - List artifactModuleTemplates = new ArrayList<>(); + List firstStageDataSourcePipelineTemplate = createIngestPipelineTemplate(javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates, pipelineConfig.getStageOneDataSourceIngestPipelineConfig()); + List secondStageDataSourcePipelineTemplate = createIngestPipelineTemplate(javaDataSourceModuleTemplates, jythonDataSourceModuleTemplates, pipelineConfig.getStageTwoDataSourceIngestPipelineConfig()); + List filePipelineTemplate = createIngestPipelineTemplate(javaFileModuleTemplates, jythonFileModuleTemplates, pipelineConfig.getFileIngestPipelineConfig()); + List 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 addConfiguredIngestModuleTemplates(Map javaIngestModuleTemplates, Map jythonIngestModuleTemplates, List pipelineConfig) { - List templates = new ArrayList<>(); + private static List createIngestPipelineTemplate(Map javaIngestModuleTemplates, Map jythonIngestModuleTemplates, List pipelineConfig) { + List 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 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 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 artifacts) { List 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(),