Commit 0f458aca authored by Dominik Charousset's avatar Dominik Charousset

Call force_emit_batches once per stream manager

parent c5999142
......@@ -1170,9 +1170,15 @@ scheduled_actor::advance_streams(actor_clock::time_point now) {
auto bitmask = stream_ticks_.timeouts(now, {max_batch_delay_ticks_,
credit_round_ticks_});
// Force batches on all output paths.
if ((bitmask & 0x01) != 0) {
if ((bitmask & 0x01) != 0 && !stream_managers_.empty()) {
std::vector<stream_manager*> managers;
managers.reserve(stream_managers_.size());
for (auto& kvp : stream_managers_)
kvp.second->out().force_emit_batches();
managers.emplace_back(kvp.second.get());
std::sort(managers.begin(), managers.end());
auto e = std::unique(managers.begin(), managers.end());
for (auto i = managers.begin(); i != e; ++i)
(*i)->out().force_emit_batches();
}
// Fill up credit on each input path.
if ((bitmask & 0x02) != 0) {
......
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment