#include "pipeline.h" #include "logger.h" #include "module.h" void Pipeline::run(std::string input_file, std::string output_directory, std::map parameters, std::string input_level, bool ui, std::shared_ptr>> uiCallList, std::shared_ptr uiCallListMutex) { logger->debug("Starting " + name); std::vector lastFiles; int stepC = 0; bool foundLevel; for (PipelineStep &step : steps) { if (!foundLevel) { foundLevel = step.level_name == input_level; logger->warn("Data is already at level " + step.level_name + ", skipping"); continue; } logger->warn("Processing data to level " + step.level_name); std::vector files; for (PipelineModule &modStep : step.modules) { std::map final_parameters = modStep.parameters; for (const std::pair ¶m : parameters) if (final_parameters.count(param.first) > 0) final_parameters[param.first] = param.second; else final_parameters.emplace(param.first, param.second); logger->debug("Parameters :"); for (const std::pair ¶m : final_parameters) logger->debug(" - " + param.first + " : " + param.second); std::shared_ptr module = modules_registry[modStep.module_name](stepC == 0 ? input_file : lastFiles[0], output_directory + "/" + name, final_parameters); if (ui) { uiCallListMutex->lock(); uiCallList->push_back(module); uiCallListMutex->unlock(); } module->process(); if (ui) { uiCallListMutex->lock(); uiCallList->clear(); uiCallListMutex->unlock(); } std::vector newfiles = module->getOutputs(); files.insert(files.end(), newfiles.begin(), newfiles.end()); } lastFiles = files; stepC++; } }