aboutsummaryrefslogtreecommitdiffhomepage
path: root/tdactor/td/actor/MultiPromise.h
blob: 73b24d5d1ced99d1e6b2dfd9b911583bff12b18c (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
//
// Copyright Aliaksei Levin (levlam@telegram.org), Arseny Smirnov (arseny30@gmail.com) 2014-2022
//
// 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/actor/PromiseFuture.h"

#include "td/utils/common.h"
#include "td/utils/Promise.h"
#include "td/utils/Status.h"

namespace td {

class MultiPromiseInterface {
 public:
  virtual void add_promise(Promise<> &&promise) = 0;
  virtual Promise<> get_promise() = 0;

  virtual size_t promise_count() const = 0;
  virtual void set_ignore_errors(bool ignore_errors) = 0;

  MultiPromiseInterface() = default;
  MultiPromiseInterface(const MultiPromiseInterface &) = delete;
  MultiPromiseInterface &operator=(const MultiPromiseInterface &) = delete;
  MultiPromiseInterface(MultiPromiseInterface &&) = default;
  MultiPromiseInterface &operator=(MultiPromiseInterface &&) = default;
  virtual ~MultiPromiseInterface() = default;
};

class MultiPromise final : public MultiPromiseInterface {
 public:
  void add_promise(Promise<> &&promise) final {
    impl_->add_promise(std::move(promise));
  }
  Promise<> get_promise() final {
    return impl_->get_promise();
  }

  size_t promise_count() const final {
    return impl_->promise_count();
  }
  void set_ignore_errors(bool ignore_errors) final {
    impl_->set_ignore_errors(ignore_errors);
  }

  MultiPromise() = default;
  explicit MultiPromise(unique_ptr<MultiPromiseInterface> impl) : impl_(std::move(impl)) {
  }

 private:
  unique_ptr<MultiPromiseInterface> impl_;
};

class MultiPromiseActor final
    : public Actor
    , public MultiPromiseInterface {
 public:
  explicit MultiPromiseActor(string name) : name_(std::move(name)) {
  }

  void add_promise(Promise<Unit> &&promise) final;

  Promise<Unit> get_promise() final;

  void set_ignore_errors(bool ignore_errors) final;

  size_t promise_count() const final;

 private:
  void set_result(Result<Unit> &&result);

  string name_;
  vector<Promise<Unit>> promises_;     // promises waiting for result
  vector<FutureActor<Unit>> futures_;  // futures waiting for result of the queries
  size_t received_results_ = 0;
  bool ignore_errors_ = false;
  Result<Unit> result_;

  void raw_event(const Event::Raw &event) final;

  void tear_down() final;

  void on_start_migrate(int32) final {
    UNREACHABLE();
  }
  void on_finish_migrate() final {
    UNREACHABLE();
  }
};

template <>
class ActorTraits<MultiPromiseActor> {
 public:
  static constexpr bool need_context = false;
  static constexpr bool need_start_up = true;
};

class MultiPromiseActorSafe final : public MultiPromiseInterface {
 public:
  void add_promise(Promise<Unit> &&promise) final;
  Promise<Unit> get_promise() final;
  void set_ignore_errors(bool ignore_errors) final;
  size_t promise_count() const final;
  explicit MultiPromiseActorSafe(string name) : multi_promise_(td::make_unique<MultiPromiseActor>(std::move(name))) {
  }
  MultiPromiseActorSafe(const MultiPromiseActorSafe &other) = delete;
  MultiPromiseActorSafe &operator=(const MultiPromiseActorSafe &other) = delete;
  MultiPromiseActorSafe(MultiPromiseActorSafe &&other) = delete;
  MultiPromiseActorSafe &operator=(MultiPromiseActorSafe &&other) = delete;
  ~MultiPromiseActorSafe() final;

 private:
  unique_ptr<MultiPromiseActor> multi_promise_;
};

}  // namespace td