aboutsummaryrefslogtreecommitdiffhomepage
path: root/td/telegram/QueryMerger.h
diff options
context:
space:
mode:
authorlevlam <levlam@telegram.org>2023-01-13 13:09:38 +0300
committerlevlam <levlam@telegram.org>2023-01-13 13:09:38 +0300
commitb1883d357c8b5d00bbcfb4aa157d94f9aadef6d2 (patch)
tree25f814e562a1e9b9708f28f3da27922556528d9e /td/telegram/QueryMerger.h
parent75bdc6292b488c195645a01b9242c7dab9d01fd1 (diff)
Add QueryMerger.
Diffstat (limited to 'td/telegram/QueryMerger.h')
-rw-r--r--td/telegram/QueryMerger.h54
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