| 279 | } |
| 280 | |
| 281 | bool TestPlan::runCurrentProcessor(std::function<void(const std::shared_ptr<core::ProcessContext>, const std::shared_ptr<core::ProcessSession>)> verify) { |
| 282 | if (!finalized) { |
| 283 | finalize(); |
| 284 | } |
| 285 | logger_->log_info("Rerunning current processor %d, processor_queue_.size %d, processor_contexts_.size %d", location, processor_queue_.size(), processor_contexts_.size()); |
| 286 | std::lock_guard<std::recursive_mutex> guard(mutex); |
| 287 | |
| 288 | std::shared_ptr<core::Processor> processor = processor_queue_.at(location); |
| 289 | std::shared_ptr<core::ProcessContext> context = processor_contexts_.at(location); |
| 290 | std::shared_ptr<core::ProcessSession> current_session = std::make_shared<core::ProcessSession>(context); |
| 291 | process_sessions_.push_back(current_session); |
| 292 | current_flowfile_ = nullptr; |
| 293 | processor->incrementActiveTasks(); |
| 294 | processor->setScheduledState(core::ScheduledState::RUNNING); |
| 295 | if (verify != nullptr) { |
| 296 | verify(context, current_session); |
| 297 | } else { |
| 298 | logger_->log_info("Running %s", processor->getName()); |
| 299 | processor->onTrigger(context, current_session); |
| 300 | } |
| 301 | current_session->commit(); |
| 302 | return gsl::narrow<size_t>(location + 1) < processor_queue_.size(); |
| 303 | } |
| 304 | |
| 305 | std::set<std::shared_ptr<provenance::ProvenanceEventRecord>> TestPlan::getProvenanceRecords() { |
| 306 | return process_sessions_.at(location)->getProvenanceReporter()->getEvents(); |
no test coverage detected