• Home
  • Line#
  • Scopes#
  • Navigate#
  • Raw
  • Download
1 //
2 //
3 // Copyright 2017 gRPC authors.
4 //
5 // Licensed under the Apache License, Version 2.0 (the "License");
6 // you may not use this file except in compliance with the License.
7 // You may obtain a copy of the License at
8 //
9 //     http://www.apache.org/licenses/LICENSE-2.0
10 //
11 // Unless required by applicable law or agreed to in writing, software
12 // distributed under the License is distributed on an "AS IS" BASIS,
13 // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14 // See the License for the specific language governing permissions and
15 // limitations under the License.
16 //
17 //
18 
19 #include "src/core/client_channel/backup_poller.h"
20 
21 #include <grpc/support/alloc.h>
22 #include <grpc/support/port_platform.h>
23 #include <grpc/support/sync.h>
24 #include <inttypes.h>
25 
26 #include "absl/log/log.h"
27 #include "absl/status/status.h"
28 #include "src/core/config/config_vars.h"
29 #include "src/core/lib/iomgr/closure.h"
30 #include "src/core/lib/iomgr/error.h"
31 #include "src/core/lib/iomgr/iomgr.h"
32 #include "src/core/lib/iomgr/pollset.h"
33 #include "src/core/lib/iomgr/pollset_set.h"
34 #include "src/core/lib/iomgr/timer.h"
35 #include "src/core/util/memory.h"
36 #include "src/core/util/time.h"
37 
38 #define DEFAULT_POLL_INTERVAL_MS 5000
39 
40 namespace {
41 struct backup_poller {
42   grpc_timer polling_timer;
43   grpc_closure run_poller_closure;
44   grpc_closure shutdown_closure;
45   gpr_mu* pollset_mu;
46   grpc_pollset* pollset;  // guarded by pollset_mu
47   bool shutting_down;     // guarded by pollset_mu
48   gpr_refcount refs;
49   gpr_refcount shutdown_refs;
50 };
51 }  // namespace
52 
53 static gpr_mu g_poller_mu;
54 static backup_poller* g_poller = nullptr;  // guarded by g_poller_mu
55 // g_poll_interval_ms is set only once at the first time
56 // grpc_client_channel_start_backup_polling() is called, after that it is
57 // treated as const.
58 static grpc_core::Duration g_poll_interval =
59     grpc_core::Duration::Milliseconds(DEFAULT_POLL_INTERVAL_MS);
60 // TODO(hork): delete the backup poller when EventEngine is rolled out
61 // everywhere.
62 static bool g_backup_polling_disabled;
63 
grpc_client_channel_global_init_backup_polling()64 void grpc_client_channel_global_init_backup_polling() {
65   // Disable backup polling if EventEngine is used everywhere.
66   g_backup_polling_disabled = grpc_core::IsEventEngineClientEnabled() &&
67                               grpc_core::IsEventEngineListenerEnabled() &&
68                               grpc_core::IsEventEngineDnsEnabled();
69   if (g_backup_polling_disabled) {
70     return;
71   }
72 
73   gpr_mu_init(&g_poller_mu);
74   int32_t poll_interval_ms =
75       grpc_core::ConfigVars::Get().ClientChannelBackupPollIntervalMs();
76   if (poll_interval_ms < 0) {
77     LOG(ERROR) << "Invalid GRPC_CLIENT_CHANNEL_BACKUP_POLL_INTERVAL_MS: "
78                << poll_interval_ms << ", default value "
79                << g_poll_interval.millis() << " will be used.";
80   } else {
81     g_poll_interval = grpc_core::Duration::Milliseconds(poll_interval_ms);
82   }
83 }
84 
backup_poller_shutdown_unref(backup_poller * p)85 static void backup_poller_shutdown_unref(backup_poller* p) {
86   if (gpr_unref(&p->shutdown_refs)) {
87     grpc_pollset_destroy(p->pollset);
88     gpr_free(p->pollset);
89     gpr_free(p);
90   }
91 }
92 
done_poller(void * arg,grpc_error_handle)93 static void done_poller(void* arg, grpc_error_handle /*error*/) {
94   backup_poller_shutdown_unref(static_cast<backup_poller*>(arg));
95 }
96 
g_poller_unref()97 static void g_poller_unref() {
98   gpr_mu_lock(&g_poller_mu);
99   if (gpr_unref(&g_poller->refs)) {
100     backup_poller* p = g_poller;
101     g_poller = nullptr;
102     gpr_mu_unlock(&g_poller_mu);
103     gpr_mu_lock(p->pollset_mu);
104     p->shutting_down = true;
105     grpc_pollset_shutdown(
106         p->pollset, GRPC_CLOSURE_INIT(&p->shutdown_closure, done_poller, p,
107                                       grpc_schedule_on_exec_ctx));
108     gpr_mu_unlock(p->pollset_mu);
109     grpc_timer_cancel(&p->polling_timer);
110     backup_poller_shutdown_unref(p);
111   } else {
112     gpr_mu_unlock(&g_poller_mu);
113   }
114 }
115 
run_poller(void * arg,grpc_error_handle error)116 static void run_poller(void* arg, grpc_error_handle error) {
117   backup_poller* p = static_cast<backup_poller*>(arg);
118   if (!error.ok()) {
119     if (error != absl::CancelledError()) {
120       GRPC_LOG_IF_ERROR("run_poller", error);
121     }
122     backup_poller_shutdown_unref(p);
123     return;
124   }
125   gpr_mu_lock(p->pollset_mu);
126   if (p->shutting_down) {
127     gpr_mu_unlock(p->pollset_mu);
128     backup_poller_shutdown_unref(p);
129     return;
130   }
131   grpc_error_handle err =
132       grpc_pollset_work(p->pollset, nullptr, grpc_core::Timestamp::Now());
133   gpr_mu_unlock(p->pollset_mu);
134   GRPC_LOG_IF_ERROR("Run client channel backup poller", err);
135   grpc_timer_init(&p->polling_timer,
136                   grpc_core::Timestamp::Now() + g_poll_interval,
137                   &p->run_poller_closure);
138 }
139 
g_poller_init_locked()140 static void g_poller_init_locked() {
141   if (g_poller == nullptr) {
142     g_poller = grpc_core::Zalloc<backup_poller>();
143     g_poller->pollset =
144         static_cast<grpc_pollset*>(gpr_zalloc(grpc_pollset_size()));
145     g_poller->shutting_down = false;
146     grpc_pollset_init(g_poller->pollset, &g_poller->pollset_mu);
147     gpr_ref_init(&g_poller->refs, 0);
148     // one for timer cancellation, one for pollset shutdown, one for g_poller
149     gpr_ref_init(&g_poller->shutdown_refs, 3);
150     GRPC_CLOSURE_INIT(&g_poller->run_poller_closure, run_poller, g_poller,
151                       grpc_schedule_on_exec_ctx);
152     grpc_timer_init(&g_poller->polling_timer,
153                     grpc_core::Timestamp::Now() + g_poll_interval,
154                     &g_poller->run_poller_closure);
155   }
156 }
157 
grpc_client_channel_start_backup_polling(grpc_pollset_set * interested_parties)158 void grpc_client_channel_start_backup_polling(
159     grpc_pollset_set* interested_parties) {
160   if (g_backup_polling_disabled ||
161       g_poll_interval == grpc_core::Duration::Zero() ||
162       grpc_iomgr_run_in_background()) {
163     return;
164   }
165   gpr_mu_lock(&g_poller_mu);
166   g_poller_init_locked();
167   gpr_ref(&g_poller->refs);
168   // Get a reference to g_poller->pollset before releasing g_poller_mu to make
169   // TSAN happy. Otherwise, reading from g_poller (i.e g_poller->pollset) after
170   // releasing the lock and setting g_poller to NULL in g_poller_unref() is
171   // being flagged as a data-race by TSAN
172   grpc_pollset* pollset = g_poller->pollset;
173   gpr_mu_unlock(&g_poller_mu);
174 
175   grpc_pollset_set_add_pollset(interested_parties, pollset);
176 }
177 
grpc_client_channel_stop_backup_polling(grpc_pollset_set * interested_parties)178 void grpc_client_channel_stop_backup_polling(
179     grpc_pollset_set* interested_parties) {
180   if (g_backup_polling_disabled ||
181       g_poll_interval == grpc_core::Duration::Zero() ||
182       grpc_iomgr_run_in_background()) {
183     return;
184   }
185   grpc_pollset_set_del_pollset(interested_parties, g_poller->pollset);
186   g_poller_unref();
187 }
188