| 345 | } |
| 346 | |
| 347 | void Pipe::addSource(ProcessorPtr source) |
| 348 | { |
| 349 | checkSource(*source); |
| 350 | const auto & source_header_ptr = source->getOutputs().front().getSharedHeader(); |
| 351 | |
| 352 | if (output_ports.empty()) |
| 353 | header = source_header_ptr; |
| 354 | else |
| 355 | assertBlocksHaveEqualStructure(*header, *source_header_ptr, "Pipes"); |
| 356 | |
| 357 | if (collected_processors) |
| 358 | collected_processors->emplace_back(source); |
| 359 | |
| 360 | output_ports.push_back(&source->getOutputs().front()); |
| 361 | processors->emplace_back(std::move(source)); |
| 362 | |
| 363 | max_parallel_streams = std::max<size_t>(max_parallel_streams, output_ports.size()); |
| 364 | } |
| 365 | |
| 366 | void Pipe::addTotalsSource(ProcessorPtr source) |
| 367 | { |
no test coverage detected