diff options
| author | levlam <levlam@telegram.org> | 2023-01-13 13:09:38 +0300 |
|---|---|---|
| committer | levlam <levlam@telegram.org> | 2023-01-13 13:09:38 +0300 |
| commit | b1883d357c8b5d00bbcfb4aa157d94f9aadef6d2 (patch) | |
| tree | 25f814e562a1e9b9708f28f3da27922556528d9e /td/telegram/QueryMerger.h | |
| parent | 75bdc6292b488c195645a01b9242c7dab9d01fd1 (diff) | |
Add QueryMerger.
Diffstat (limited to 'td/telegram/QueryMerger.h')
| -rw-r--r-- | td/telegram/QueryMerger.h | 54 |
1 files changed, 54 insertions, 0 deletions
diff --git a/td/telegram/QueryMerger.h b/td/telegram/QueryMerger.h new file mode 100644 index 000000000..e87f81666 --- /dev/null +++ b/td/telegram/QueryMerger.h @@ -0,0 +1,54 @@ +// +// Copyright Aliaksei Levin (levlam@telegram.org), Arseny Smirnov (arseny30@gmail.com) 2014-2023 +// +// Distributed under the Boost Software License, Version 1.0. (See accompanying +// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt) +// +#pragma once + +#include "td/actor/actor.h" + +#include "td/utils/common.h" +#include "td/utils/FlatHashMap.h" +#include "td/utils/Promise.h" +#include "td/utils/Slice.h" +#include "td/utils/Status.h" + +#include <functional> +#include <queue> + +namespace td { + +// merges queries into a single request +class QueryMerger final : public Actor { + public: + QueryMerger(Slice name, size_t max_concurrent_query_count, size_t max_merged_query_count); + + using MergeFunction = std::function<void(vector<int64> query_ids, Promise<Unit> &&promise)>; + void set_merge_function(MergeFunction merge_function) { + merge_function_ = std::move(merge_function); + } + + void add_query(int64 query_id, Promise<Unit> &&promise); + + private: + struct QueryInfo { + vector<Promise<Unit>> promises_; + }; + + size_t query_count_ = 0; + size_t max_concurrent_query_count_; + size_t max_merged_query_count_; + + MergeFunction merge_function_; + std::queue<int64> pending_queries_; + FlatHashMap<int64, QueryInfo> queries_; + + void send_query(vector<int64> query_ids); + + void on_get_query_result(vector<int64> query_ids, Result<Unit> &&result); + + void loop() final; +}; + +} // namespace td |
