| 158 | } |
| 159 | |
| 160 | void OutputBufferManager::removeTask(const std::string& taskId) { |
| 161 | auto buffer = |
| 162 | buffers_.withLock([&](auto& buffers) -> std::shared_ptr<OutputBuffer> { |
| 163 | auto it = buffers.find(taskId); |
| 164 | if (it == buffers.end()) { |
| 165 | // Already removed. |
| 166 | return nullptr; |
| 167 | } |
| 168 | auto taskBuffer = it->second; |
| 169 | buffers.erase(taskId); |
| 170 | return taskBuffer; |
| 171 | }); |
| 172 | if (buffer != nullptr) { |
| 173 | buffer->terminate(); |
| 174 | } |
| 175 | } |
| 176 | |
| 177 | std::string OutputBufferManager::toString() { |
| 178 | return buffers_.withLock([](const auto& buffers) { |