| 28 | } |
| 29 | |
| 30 | void TaskIdGenerator::TryUpdateTaskIndex(const HashMap<int64_t, uint32_t>& task_index_state) { |
| 31 | for (auto& pair : stream_id2task_index_counter_) { |
| 32 | const int64_t i64_stream_id = EncodeStreamIdToInt64(pair.first); |
| 33 | uint32_t initial_task_index = 0; |
| 34 | if (task_index_state.count(i64_stream_id) != 0) { |
| 35 | initial_task_index = task_index_state.at(i64_stream_id); |
| 36 | } |
| 37 | pair.second = std::max(pair.second, initial_task_index); |
| 38 | } |
| 39 | |
| 40 | // try update the task_index_init_state |
| 41 | for (const auto& pair : task_index_state) { |
| 42 | const auto& key = pair.first; |
| 43 | const auto& val = pair.second; |
| 44 | if (task_index_init_state_.count(key) != 0) { |
| 45 | task_index_init_state_[key] = std::max(task_index_init_state_.at(key), val); |
| 46 | } else { |
| 47 | task_index_init_state_[key] = val; |
| 48 | } |
| 49 | } |
| 50 | } |
| 51 | |
| 52 | TaskId TaskIdGenerator::Generate(const StreamId& stream_id) { |
| 53 | std::unique_lock<std::mutex> lock(mutex_); |
no test coverage detected