Coverage Report

Created: 2026-09-01 06:30

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/brpc/src/bthread/butex.cpp
Line
Count
Source
1
// Licensed to the Apache Software Foundation (ASF) under one
2
// or more contributor license agreements.  See the NOTICE file
3
// distributed with this work for additional information
4
// regarding copyright ownership.  The ASF licenses this file
5
// to you under the Apache License, Version 2.0 (the
6
// "License"); you may not use this file except in compliance
7
// with the License.  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,
12
// software distributed under the License is distributed on an
13
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14
// KIND, either express or implied.  See the License for the
15
// specific language governing permissions and limitations
16
// under the License.
17
18
// bthread - An M:N threading library to make applications more concurrent.
19
20
// Date: Tue Jul 22 17:30:12 CST 2014
21
22
#include "butil/atomicops.h"                // butil::atomic
23
#include "butil/scoped_lock.h"              // BAIDU_SCOPED_LOCK
24
#include "butil/macros.h"
25
#include "butil/containers/flat_map.h"
26
#include "butil/containers/linked_list.h"   // LinkNode
27
#ifdef SHOW_BTHREAD_BUTEX_WAITER_COUNT_IN_VARS
28
#include "butil/memory/singleton_on_pthread_once.h"
29
#endif
30
#include "butil/logging.h"
31
#include "butil/object_pool.h"
32
#include "bthread/errno.h"                 // EWOULDBLOCK
33
#include "bthread/sys_futex.h"             // futex_*
34
#include "bthread/processor.h"             // cpu_relax
35
#include "bthread/task_control.h"          // TaskControl
36
#include "bthread/task_group.h"            // TaskGroup
37
#include "bthread/timer_thread.h"
38
#include "bthread/butex.h"
39
#include "bthread/mutex.h"
40
41
// This file implements butex.h
42
// Provides futex-like semantics which is sequenced wait and wake operations
43
// and guaranteed visibilities.
44
//
45
// If wait is sequenced before wake:
46
//    [thread1]             [thread2]
47
//    wait()                value = new_value
48
//                          wake()
49
// wait() sees unmatched value(fail to wait), or wake() sees the waiter.
50
//
51
// If wait is sequenced after wake:
52
//    [thread1]             [thread2]
53
//                          value = new_value
54
//                          wake()
55
//    wait()
56
// wake() must provide some sort of memory fence to prevent assignment
57
// of value to be reordered after it. Thus the value is visible to wait()
58
// as well.
59
60
namespace bthread {
61
62
#ifdef SHOW_BTHREAD_BUTEX_WAITER_COUNT_IN_VARS
63
struct ButexWaiterCount : public bvar::Adder<int64_t> {
64
    ButexWaiterCount() : bvar::Adder<int64_t>("bthread_butex_waiter_count") {}
65
};
66
inline bvar::Adder<int64_t>& butex_waiter_count() {
67
    return *butil::get_leaky_singleton<ButexWaiterCount>();
68
}
69
#endif
70
71
enum WaiterState {
72
    WAITER_STATE_NONE,
73
    WAITER_STATE_READY,
74
    WAITER_STATE_TIMEDOUT,
75
    WAITER_STATE_UNMATCHEDVALUE,
76
    WAITER_STATE_INTERRUPTED,
77
};
78
79
struct Butex;
80
81
struct ButexWaiter : public butil::LinkNode<ButexWaiter> {
82
    // tids of pthreads are 0
83
    bthread_t tid;
84
85
    // Erasing node from middle of LinkedList is thread-unsafe, we need
86
    // to hold its container's lock.
87
    butil::atomic<Butex*> container;
88
};
89
90
// non_pthread_task allocates this structure on stack and queue it in
91
// Butex::waiters.
92
struct ButexBthreadWaiter : public ButexWaiter {
93
    TaskMeta* task_meta;
94
    TimerThread::TaskId sleep_id;
95
    WaiterState waiter_state;
96
    int expected_value;
97
    Butex* initial_butex;
98
    TaskControl* control;
99
    const timespec* abstime;
100
};
101
102
// pthread_task or main_task allocates this structure on stack and queue it
103
// in Butex::waiters.
104
struct ButexPthreadWaiter : public ButexWaiter {
105
    butil::atomic<int> sig;
106
};
107
108
typedef butil::LinkedList<ButexWaiter> ButexWaiterList;
109
110
enum ButexPthreadSignal { PTHREAD_NOT_SIGNALLED, PTHREAD_SIGNALLED };
111
112
struct BAIDU_CACHELINE_ALIGNMENT Butex {
113
3
    Butex() {}
114
0
    ~Butex() {}
115
116
    butil::atomic<int> value;
117
    ButexWaiterList waiters;
118
    FastPthreadMutex waiter_lock;
119
};
120
121
BAIDU_CASSERT(offsetof(Butex, value) == 0, offsetof_value_must_0);
122
BAIDU_CASSERT(sizeof(Butex) == BAIDU_CACHELINE_SIZE, butex_fits_in_one_cacheline);
123
124
} // namespace bthread
125
126
namespace butil {
127
// Butex object returned to the ObjectPool<Butex> may be accessed,
128
// so ObjectPool<Butex> can not poison the memory region of Butex.
129
template <>
130
struct ObjectPoolWithASanPoison<bthread::Butex> : false_type {};
131
} // namespace butil
132
133
namespace bthread {
134
135
0
static void wakeup_pthread(ButexPthreadWaiter* pw) {
136
    // release fence makes wait_pthread see changes before wakeup.
137
0
    pw->sig.store(PTHREAD_SIGNALLED, butil::memory_order_release);
138
    // At this point, wait_pthread() possibly has woken up and destroyed `pw'.
139
    // In which case, futex_wake_private() should return EFAULT.
140
    // If crash happens in future, `pw' can be made TLS and never destroyed
141
    // to solve the issue.
142
0
    futex_wake_private(&pw->sig, 1);
143
0
}
144
145
bool erase_from_butex(ButexWaiter*, bool, WaiterState);
146
147
0
int wait_pthread(ButexPthreadWaiter& pw, const timespec* abstime) {
148
0
    timespec* ptimeout = nullptr;
149
0
    timespec timeout;
150
0
    int64_t timeout_us = 0;
151
0
    int rc;
152
153
0
    while (true) {
154
0
        if (abstime != nullptr) {
155
0
            timeout_us = butil::timespec_to_microseconds(*abstime) - butil::gettimeofday_us();
156
0
            timeout = butil::microseconds_to_timespec(timeout_us);
157
0
            ptimeout = &timeout;
158
0
        }
159
0
        if (timeout_us > MIN_SLEEP_US || abstime == nullptr) {
160
0
            rc = futex_wait_private(&pw.sig, PTHREAD_NOT_SIGNALLED, ptimeout);
161
0
            if (PTHREAD_NOT_SIGNALLED != pw.sig.load(butil::memory_order_acquire)) {
162
                // If `sig' is changed, wakeup_pthread() must be called and `pw'
163
                // is already removed from the butex.
164
                // Acquire fence makes this thread sees changes before wakeup.
165
0
                return rc;
166
0
            }
167
0
        } else {
168
0
            errno = ETIMEDOUT;
169
0
            rc = -1;
170
0
        }
171
        // Handle ETIMEDOUT when abstime is valid.
172
        // If futex_wait_private return EINTR, just continue the loop.
173
0
        if (rc != 0 && errno == ETIMEDOUT) {
174
            // wait futex timeout, `pw' is still in the queue, remove it.
175
0
            if (!erase_from_butex(&pw, false, WAITER_STATE_TIMEDOUT)) {
176
                // Another thread is erasing `pw' as well, wait for the signal.
177
                // Acquire fence makes this thread sees changes before wakeup.
178
0
                if (pw.sig.load(butil::memory_order_acquire) == PTHREAD_NOT_SIGNALLED) {
179
                    // already timedout, abstime and ptimeout are expired.
180
0
                    abstime = nullptr;
181
0
                    ptimeout = nullptr;
182
0
                    continue;
183
0
                }
184
0
            }
185
0
            return rc;
186
0
        }
187
0
    }
188
0
}
189
190
EXTERN_BAIDU_VOLATILE_THREAD_LOCAL(TaskGroup*, tls_task_group);
191
192
// Returns 0 when no need to unschedule or successfully unscheduled,
193
// -1 otherwise.
194
inline int unsleep_if_necessary(ButexBthreadWaiter* w,
195
0
                                TimerThread* timer_thread) {
196
0
    if (!w->sleep_id) {
197
0
        return 0;
198
0
    }
199
0
    if (timer_thread->unschedule(w->sleep_id) > 0) {
200
        // the callback is running.
201
0
        return -1;
202
0
    }
203
0
    w->sleep_id = 0;
204
0
    return 0;
205
0
}
206
207
// Use ObjectPool(which never frees memory) to solve the race between
208
// butex_wake() and butex_destroy(). The race is as follows:
209
//
210
//   class Event {
211
//   public:
212
//     void wait() {
213
//       _mutex.lock();
214
//       if (!_done) {
215
//         _cond.wait(&_mutex);
216
//       }
217
//       _mutex.unlock();
218
//     }
219
//     void signal() {
220
//       _mutex.lock();
221
//       if (!_done) {
222
//         _done = true;
223
//         _cond.signal();
224
//       }
225
//       _mutex.unlock();  /*1*/
226
//     }
227
//   private:
228
//     bool _done = false;
229
//     Mutex _mutex;
230
//     Condition _cond;
231
//   };
232
//
233
//   [Thread1]                         [Thread2]
234
//   foo() {
235
//     Event event;
236
//     pass_to_thread2(&event);  --->  event.signal();
237
//     event.wait();
238
//   } <-- event destroyed
239
//   
240
// Summary: Thread1 passes a stateful condition to Thread2 and waits until
241
// the condition being signalled, which basically means the associated
242
// job is done and Thread1 can release related resources including the mutex
243
// and condition. The scenario is fine and the code is correct.
244
// The race needs a closer look. The unlock at /*1*/ may have different 
245
// implementations, but in which the last step is probably an atomic store
246
// and butex_wake(), like this:
247
//
248
//   locked->store(0);
249
//   butex_wake(locked);
250
//
251
// The `locked' represents the locking status of the mutex. The issue is that
252
// just after the store(), the mutex is already unlocked and the code in
253
// Event.wait() may successfully grab the lock and go through everything
254
// left and leave foo() function, destroying the mutex and butex, making
255
// the butex_wake(locked) crash.
256
// To solve this issue, one method is to add reference before store and
257
// release the reference after butex_wake. However reference countings need
258
// to be added in nearly every user scenario of butex_wake(), which is very
259
// error-prone. Another method is never freeing butex, with the side effect 
260
// that butex_wake() may wake up an unrelated butex(the one reuses the memory)
261
// and cause spurious wakeups. According to our observations, the race is 
262
// infrequent, even rare. The extra spurious wakeups should be acceptable.
263
264
3
void* butex_create() {
265
3
    Butex* b = butil::get_object<Butex>();
266
3
    if (b) {
267
3
        return &b->value;
268
3
    }
269
0
    return nullptr;
270
3
}
271
272
0
void butex_destroy(void* butex) {
273
0
    if (!butex) {
274
0
        return;
275
0
    }
276
0
    Butex* b = static_cast<Butex*>(
277
0
        container_of(static_cast<butil::atomic<int>*>(butex), Butex, value));
278
0
    butil::return_object(b);
279
0
}
280
281
// if TaskGroup tls_task_group is belong to tag
282
0
inline bool is_same_tag(bthread_tag_t tag) {
283
0
    auto g = BAIDU_GET_VOLATILE_THREAD_LOCAL(tls_task_group);
284
0
    return g && g->tag() == tag;
285
0
}
286
287
//  nosignal is true & tag is same can return true
288
0
inline bool check_nosignal(bool nosignal, bthread_tag_t tag) {
289
0
    return nosignal && is_same_tag(tag);
290
0
}
291
292
// if tag is same return tls_task_group else choose one group with tag
293
0
inline TaskGroup* get_task_group(TaskControl* c, bthread_tag_t tag) {
294
0
    return is_same_tag(tag) ? BAIDU_GET_VOLATILE_THREAD_LOCAL(tls_task_group)
295
0
                            : c->choose_one_group(tag);
296
0
}
297
298
0
inline void run_in_local_task_group(TaskGroup* g, TaskMeta* next_meta, bool nosignal) {
299
0
    if (!nosignal) {
300
0
        TaskGroup::exchange(&g, next_meta);
301
0
    } else {
302
0
        g->ready_to_run(next_meta, nosignal);
303
0
    }
304
0
}
305
306
0
int butex_wake(void* arg, bool nosignal) {
307
0
    Butex* b = container_of(static_cast<butil::atomic<int>*>(arg), Butex, value);
308
0
    ButexWaiter* front = nullptr;
309
0
    {
310
0
        BAIDU_SCOPED_LOCK(b->waiter_lock);
311
0
        if (b->waiters.empty()) {
312
0
            return 0;
313
0
        }
314
0
        front = b->waiters.head()->value();
315
0
        front->RemoveFromList();
316
0
        front->container.store(nullptr, butil::memory_order_relaxed);
317
0
    }
318
0
    if (front->tid == 0) {
319
0
        wakeup_pthread(static_cast<ButexPthreadWaiter*>(front));
320
0
        return 1;
321
0
    }
322
0
    ButexBthreadWaiter* bbw = static_cast<ButexBthreadWaiter*>(front);
323
0
    unsleep_if_necessary(bbw, get_global_timer_thread());
324
0
    TaskGroup* g = get_task_group(bbw->control, bbw->task_meta->attr.tag);
325
0
    if (g == BAIDU_GET_VOLATILE_THREAD_LOCAL(tls_task_group)) {
326
0
        run_in_local_task_group(g, bbw->task_meta, nosignal);
327
0
    } else {
328
0
        g->ready_to_run_remote(bbw->task_meta, check_nosignal(nosignal, g->tag()));
329
0
    }
330
0
    return 1;
331
0
}
332
333
0
int butex_wake_n(void* arg, size_t n, bool nosignal) {
334
0
    Butex* b = container_of(static_cast<butil::atomic<int>*>(arg), Butex, value);
335
336
0
    ButexWaiterList bthread_waiters;
337
0
    ButexWaiterList pthread_waiters;
338
0
    {
339
0
        BAIDU_SCOPED_LOCK(b->waiter_lock);
340
0
        for (size_t i = 0; (n == 0 || i < n) && !b->waiters.empty(); ++i) {
341
0
            ButexWaiter* bw = b->waiters.head()->value();
342
0
            bw->RemoveFromList();
343
0
            bw->container.store(nullptr, butil::memory_order_relaxed);
344
0
            if (bw->tid) {
345
0
                bthread_waiters.Append(bw);
346
0
            } else {
347
0
                pthread_waiters.Append(bw);
348
0
            }
349
0
        }
350
0
    }
351
352
0
    int nwakeup = 0;
353
0
    while (!pthread_waiters.empty()) {
354
0
        ButexPthreadWaiter* bw = static_cast<ButexPthreadWaiter*>(
355
0
            pthread_waiters.head()->value());
356
0
        bw->RemoveFromList();
357
0
        wakeup_pthread(bw);
358
0
        ++nwakeup;
359
0
    }
360
0
    if (bthread_waiters.empty()) {
361
0
        return nwakeup;
362
0
    }
363
0
    butil::FlatMap<bthread_tag_t, TaskGroup*> nwakeups;
364
0
    nwakeups.init(FLAGS_task_group_ntags);
365
    // We will exchange with first waiter in the end.
366
0
    ButexBthreadWaiter* next = static_cast<ButexBthreadWaiter*>(
367
0
        bthread_waiters.head()->value());
368
0
    next->RemoveFromList();
369
0
    unsleep_if_necessary(next, get_global_timer_thread());
370
0
    ++nwakeup;
371
0
    while (!bthread_waiters.empty()) {
372
        // pop reversely
373
0
        ButexBthreadWaiter* w = static_cast<ButexBthreadWaiter*>(
374
0
            bthread_waiters.tail()->value());
375
0
        w->RemoveFromList();
376
0
        unsleep_if_necessary(w, get_global_timer_thread());
377
0
        auto g = get_task_group(w->control, w->task_meta->attr.tag);
378
0
        g->ready_to_run_general(w->task_meta, true);
379
0
        nwakeups[g->tag()] = g;
380
0
        ++nwakeup;
381
0
    }
382
0
    for (auto it = nwakeups.begin(); it != nwakeups.end(); ++it) {
383
0
        auto g = it->second;
384
0
        if (!check_nosignal(nosignal, g->tag())) {
385
0
            g->flush_nosignal_tasks_general();
386
0
        }
387
0
    }
388
0
    auto g = get_task_group(next->control, next->task_meta->attr.tag);
389
0
    if (g == BAIDU_GET_VOLATILE_THREAD_LOCAL(tls_task_group)) {
390
0
        run_in_local_task_group(g, next->task_meta, nosignal);
391
0
    } else {
392
0
        g->ready_to_run_remote(next->task_meta, check_nosignal(nosignal, g->tag()));
393
0
    }
394
0
    return nwakeup;
395
0
}
396
397
0
int butex_wake_all(void* arg, bool nosignal) {
398
0
    return butex_wake_n(arg, 0, nosignal);
399
0
}
400
401
0
int butex_wake_except(void* arg, bthread_t excluded_bthread) {
402
0
    Butex* b = container_of(static_cast<butil::atomic<int>*>(arg), Butex, value);
403
404
0
    ButexWaiterList bthread_waiters;
405
0
    ButexWaiterList pthread_waiters;
406
0
    {
407
0
        ButexWaiter* excluded_waiter = nullptr;
408
0
        BAIDU_SCOPED_LOCK(b->waiter_lock);
409
0
        while (!b->waiters.empty()) {
410
0
            ButexWaiter* bw = b->waiters.head()->value();
411
0
            bw->RemoveFromList();
412
413
0
            if (bw->tid) {
414
0
                if (bw->tid != excluded_bthread) {
415
0
                    bthread_waiters.Append(bw);
416
0
                    bw->container.store(nullptr, butil::memory_order_relaxed);
417
0
                } else {
418
0
                    excluded_waiter = bw;
419
0
                }
420
0
            } else {
421
0
                bw->container.store(nullptr, butil::memory_order_relaxed);
422
0
                pthread_waiters.Append(bw);
423
0
            }
424
0
        }
425
426
0
        if (excluded_waiter) {
427
0
            b->waiters.Append(excluded_waiter);
428
0
        }
429
0
    }
430
431
0
    int nwakeup = 0;
432
0
    while (!pthread_waiters.empty()) {
433
0
        ButexPthreadWaiter* bw = static_cast<ButexPthreadWaiter*>(
434
0
            pthread_waiters.head()->value());
435
0
        bw->RemoveFromList();
436
0
        wakeup_pthread(bw);
437
0
        ++nwakeup;
438
0
    }
439
440
0
    if (bthread_waiters.empty()) {
441
0
        return nwakeup;
442
0
    }
443
0
    butil::FlatMap<bthread_tag_t, TaskGroup*> nwakeups;
444
0
    nwakeups.init(FLAGS_task_group_ntags);
445
0
    do {
446
        // pop reversely
447
0
        ButexBthreadWaiter* w = static_cast<ButexBthreadWaiter*>(bthread_waiters.tail()->value());
448
0
        w->RemoveFromList();
449
0
        unsleep_if_necessary(w, get_global_timer_thread());
450
0
        auto g = get_task_group(w->control, w->task_meta->attr.tag);
451
0
        g->ready_to_run_general(w->task_meta, true);
452
0
        nwakeups[g->tag()] = g;
453
0
        ++nwakeup;
454
0
    } while (!bthread_waiters.empty());
455
0
    for (auto it = nwakeups.begin(); it != nwakeups.end(); ++it) {
456
0
        auto g = it->second;
457
0
        g->flush_nosignal_tasks_general();
458
0
    }
459
0
    return nwakeup;
460
0
}
461
462
0
int butex_requeue(void* arg, void* arg2) {
463
0
    Butex* b = container_of(static_cast<butil::atomic<int>*>(arg), Butex, value);
464
0
    Butex* m = container_of(static_cast<butil::atomic<int>*>(arg2), Butex, value);
465
466
0
    ButexWaiter* front = nullptr;
467
0
    {
468
0
        std::unique_lock<FastPthreadMutex> lck1(b->waiter_lock, std::defer_lock);
469
0
        std::unique_lock<FastPthreadMutex> lck2(m->waiter_lock, std::defer_lock);
470
0
        butil::double_lock(lck1, lck2);
471
0
        if (b->waiters.empty()) {
472
0
            return 0;
473
0
        }
474
475
0
        front = b->waiters.head()->value();
476
0
        front->RemoveFromList();
477
0
        front->container.store(nullptr, butil::memory_order_relaxed);
478
479
0
        while (!b->waiters.empty()) {
480
0
            ButexWaiter* bw = b->waiters.head()->value();
481
0
            bw->RemoveFromList();
482
0
            m->waiters.Append(bw);
483
0
            bw->container.store(m, butil::memory_order_relaxed);
484
0
        }
485
0
    }
486
487
0
    if (front->tid == 0) {  // which is a pthread
488
0
        wakeup_pthread(static_cast<ButexPthreadWaiter*>(front));
489
0
        return 1;
490
0
    }
491
0
    ButexBthreadWaiter* bbw = static_cast<ButexBthreadWaiter*>(front);
492
0
    unsleep_if_necessary(bbw, get_global_timer_thread());
493
0
    auto g = is_same_tag(bbw->task_meta->attr.tag)
494
0
                 ? BAIDU_GET_VOLATILE_THREAD_LOCAL(tls_task_group)
495
0
                 : nullptr;
496
0
    if (g) {
497
0
        TaskGroup::exchange(&g, bbw->task_meta);
498
0
    } else {
499
0
        g = bbw->control->choose_one_group(bbw->task_meta->attr.tag);
500
0
        g->ready_to_run_remote(bbw->task_meta);
501
0
    }
502
0
    return 1;
503
0
}
504
505
// Callable from multiple threads, at most one thread may wake up the waiter.
506
0
static void erase_from_butex_and_wakeup(void* arg) {
507
0
    erase_from_butex(static_cast<ButexWaiter*>(arg), true, WAITER_STATE_TIMEDOUT);
508
0
}
509
510
// Used in task_group.cpp
511
0
bool erase_from_butex_because_of_interruption(ButexWaiter* bw) {
512
0
    return erase_from_butex(bw, true, WAITER_STATE_INTERRUPTED);
513
0
}
514
515
0
inline bool erase_from_butex(ButexWaiter* bw, bool wakeup, WaiterState state) {
516
    // `bw' is guaranteed to be valid inside this function because waiter
517
    // will wait until this function being cancelled or finished.
518
    // NOTE: This function must be no-op when bw->container is nullptr.
519
0
    bool erased = false;
520
0
    Butex* b;
521
0
    int saved_errno = errno;
522
0
    while ((b = bw->container.load(butil::memory_order_acquire))) {
523
        // b can be nullptr when the waiter is scheduled but queued.
524
0
        BAIDU_SCOPED_LOCK(b->waiter_lock);
525
0
        if (b == bw->container.load(butil::memory_order_relaxed)) {
526
0
            bw->RemoveFromList();
527
0
            bw->container.store(nullptr, butil::memory_order_relaxed);
528
0
            if (bw->tid) {
529
0
                static_cast<ButexBthreadWaiter*>(bw)->waiter_state = state;
530
0
            }
531
0
            erased = true;
532
0
            break;
533
0
        }
534
0
    }
535
0
    if (erased && wakeup) {
536
0
        if (bw->tid) {
537
0
            ButexBthreadWaiter* bbw = static_cast<ButexBthreadWaiter*>(bw);
538
0
            auto g = get_task_group(bbw->control, bbw->task_meta->attr.tag);
539
0
            g->ready_to_run_general(bbw->task_meta);
540
0
        } else {
541
0
            ButexPthreadWaiter* pw = static_cast<ButexPthreadWaiter*>(bw);
542
0
            wakeup_pthread(pw);
543
0
        }
544
0
    }
545
0
    errno = saved_errno;
546
0
    return erased;
547
0
}
548
549
struct WaitForButexArgs {
550
    ButexBthreadWaiter* bw;
551
    bool prepend;
552
};
553
554
0
void wait_for_butex(void* arg) {
555
0
    auto args = static_cast<WaitForButexArgs*>(arg);
556
0
    ButexBthreadWaiter* const bw = args->bw;
557
0
    Butex* const b = bw->initial_butex;
558
    // 1: waiter with timeout should have waiter_state == WAITER_STATE_READY
559
    //    before they're queued, otherwise the waiter is already timedout
560
    //    and removed by TimerThread, in which case we should stop queueing.
561
    //
562
    // Visibility of waiter_state:
563
    //    [bthread]                         [TimerThread]
564
    //    waiter_state = TIMED
565
    //    tt_lock { add task }
566
    //                                      tt_lock { get task }
567
    //                                      waiter_lock { waiter_state=TIMEDOUT }
568
    //    waiter_lock { use waiter_state }
569
    // tt_lock represents TimerThread::_mutex. Visibility of waiter_state is
570
    // sequenced by two locks, both threads are guaranteed to see the correct
571
    // value.
572
0
    {
573
0
        BAIDU_SCOPED_LOCK(b->waiter_lock);
574
0
        if (b->value.load(butil::memory_order_relaxed) != bw->expected_value) {
575
0
            bw->waiter_state = WAITER_STATE_UNMATCHEDVALUE;
576
0
        } else if (bw->waiter_state == WAITER_STATE_READY/*1*/ &&
577
0
                   !bw->task_meta->interrupted) {
578
0
            if (args->prepend) {
579
0
                b->waiters.Prepend(bw);
580
0
            } else {
581
0
                b->waiters.Append(bw);
582
0
            }
583
0
            bw->container.store(b, butil::memory_order_relaxed);
584
#ifdef BRPC_BTHREAD_TRACER
585
            bw->control->_task_tracer.set_status(TASK_STATUS_SUSPENDED, bw->task_meta);
586
#endif // BRPC_BTHREAD_TRACER
587
0
            if (bw->abstime != nullptr) {
588
0
                bw->sleep_id = get_global_timer_thread()->schedule(
589
0
                    erase_from_butex_and_wakeup, bw, *bw->abstime);
590
0
                if (!bw->sleep_id) {  // TimerThread stopped.
591
0
                    errno = ESTOP;
592
0
                    erase_from_butex_and_wakeup(bw);
593
0
                }
594
0
            }
595
0
            return;
596
0
        }
597
0
    }
598
    
599
    // b->container is nullptr which makes erase_from_butex_and_wakeup() and
600
    // TaskGroup::interrupt() no-op, there's no race between following code and
601
    // the two functions. The on-stack ButexBthreadWaiter is safe to use and
602
    // bw->waiter_state will not change again.
603
    // unsleep_if_necessary(bw, get_global_timer_thread());
604
0
    BAIDU_GET_VOLATILE_THREAD_LOCAL(tls_task_group)->ready_to_run(bw->task_meta);
605
    // FIXME: jump back to original thread is buggy.
606
    
607
    // // Value unmatched or waiter is already woken up by TimerThread, jump
608
    // // back to original bthread.
609
    // TaskGroup* g = tls_task_group;
610
    // ReadyToRunArgs args = { g->current_tid(), false };
611
    // g->set_remained(TaskGroup::ready_to_run_in_worker, &args);
612
    // // 2: Don't run remained because we're already in a remained function
613
    // //    otherwise stack may overflow.
614
    // TaskGroup::sched_to(&g, bw->tid, false/*2*/);
615
0
}
616
617
static int butex_wait_from_pthread(TaskGroup* g, Butex* b, int expected_value,
618
0
                                   const timespec* abstime, bool prepend) {
619
0
    TaskMeta* task = nullptr;
620
0
    ButexPthreadWaiter pw;
621
0
    pw.tid = 0;
622
0
    pw.sig.store(PTHREAD_NOT_SIGNALLED, butil::memory_order_relaxed);
623
0
    int rc = 0;
624
    
625
0
    if (g) {
626
0
        task = g->current_task();
627
0
        task->current_waiter.store(&pw, butil::memory_order_release);
628
0
    }
629
0
    b->waiter_lock.lock();
630
0
    if (b->value.load(butil::memory_order_relaxed) != expected_value) {
631
0
        b->waiter_lock.unlock();
632
0
        errno = EWOULDBLOCK;
633
0
        rc = -1;
634
0
    } else if (task != nullptr && task->interrupted) {
635
0
        b->waiter_lock.unlock();
636
        // Race with set and may consume multiple interruptions, which are OK.
637
0
        task->interrupted = false;
638
0
        errno = EINTR;
639
0
        rc = -1;
640
0
    } else {
641
0
        if (prepend) {
642
0
            b->waiters.Prepend(&pw);
643
0
        } else {
644
0
            b->waiters.Append(&pw);
645
0
        }
646
0
        pw.container.store(b, butil::memory_order_relaxed);
647
0
        b->waiter_lock.unlock();
648
649
#ifdef SHOW_BTHREAD_BUTEX_WAITER_COUNT_IN_VARS
650
        bvar::Adder<int64_t>& num_waiters = butex_waiter_count();
651
        num_waiters << 1;
652
#endif
653
0
        rc = wait_pthread(pw, abstime);
654
#ifdef SHOW_BTHREAD_BUTEX_WAITER_COUNT_IN_VARS
655
        num_waiters << -1;
656
#endif
657
0
    }
658
0
    if (task) {
659
        // If current_waiter is nullptr, TaskGroup::interrupt() is running and
660
        // using pw, spin until current_waiter != nullptr.
661
0
        BT_LOOP_WHEN(task->current_waiter.exchange(
662
0
                         nullptr, butil::memory_order_acquire) == nullptr,
663
0
                     30/*nops before sched_yield*/);
664
0
        if (task->interrupted) {
665
0
            task->interrupted = false;
666
0
            if (rc == 0) {
667
0
                errno = EINTR;
668
0
                return -1;
669
0
            }
670
0
        }
671
0
    }
672
0
    return rc;
673
0
}
674
675
0
int butex_wait(void* arg, int expected_value, const timespec* abstime, bool prepend) {
676
0
    Butex* b = container_of(static_cast<butil::atomic<int>*>(arg), Butex, value);
677
0
    if (b->value.load(butil::memory_order_relaxed) != expected_value) {
678
0
        errno = EWOULDBLOCK;
679
        // Sometimes we may take actions immediately after unmatched butex,
680
        // this fence makes sure that we see changes before changing butex.
681
0
        butil::atomic_thread_fence(butil::memory_order_acquire);
682
0
        return -1;
683
0
    }
684
0
    TaskGroup* g = BAIDU_GET_VOLATILE_THREAD_LOCAL(tls_task_group);
685
0
    if (nullptr == g || g->is_current_pthread_task()) {
686
0
        return butex_wait_from_pthread(g, b, expected_value, abstime, prepend);
687
0
    }
688
0
    ButexBthreadWaiter bbw;
689
    // tid is 0 iff the thread is non-bthread
690
0
    bbw.tid = g->current_tid();
691
0
    bbw.container.store(nullptr, butil::memory_order_relaxed);
692
0
    bbw.task_meta = g->current_task();
693
0
    bbw.sleep_id = 0;
694
0
    bbw.waiter_state = WAITER_STATE_READY;
695
0
    bbw.expected_value = expected_value;
696
0
    bbw.initial_butex = b;
697
0
    bbw.control = g->control();
698
0
    bbw.abstime = abstime;
699
700
0
    if (abstime != nullptr) {
701
        // Schedule timer before queueing. If the timer is triggered before
702
        // queueing, cancel queueing. This is a kind of optimistic locking.
703
0
        if (butil::timespec_to_microseconds(*abstime) <
704
0
            (butil::gettimeofday_us() + MIN_SLEEP_US)) {
705
            // Already timed out.
706
0
            errno = ETIMEDOUT;
707
0
            return -1;
708
0
        }
709
0
    }
710
#ifdef SHOW_BTHREAD_BUTEX_WAITER_COUNT_IN_VARS
711
    bvar::Adder<int64_t>& num_waiters = butex_waiter_count();
712
    num_waiters << 1;
713
#endif
714
715
    // release fence matches with acquire fence in interrupt_and_consume_waiters
716
    // in task_group.cpp to guarantee visibility of `interrupted'.
717
0
    bbw.task_meta->current_waiter.store(&bbw, butil::memory_order_release);
718
0
    WaitForButexArgs args{ &bbw, prepend };
719
0
    g->set_remained(wait_for_butex, &args);
720
0
    TaskGroup::sched(&g);
721
722
    // erase_from_butex_and_wakeup (called by TimerThread) is possibly still
723
    // running and using bbw. The chance is small, just spin until it's done.
724
0
    BT_LOOP_WHEN(unsleep_if_necessary(&bbw, get_global_timer_thread()) < 0,
725
0
                 30/*nops before sched_yield*/);
726
    
727
    // If current_waiter is nullptr, TaskGroup::interrupt() is running and using bbw.
728
    // Spin until current_waiter != nullptr.
729
0
    BT_LOOP_WHEN(bbw.task_meta->current_waiter.exchange(
730
0
                     nullptr, butil::memory_order_acquire) == nullptr,
731
0
                 30/*nops before sched_yield*/);
732
#ifdef SHOW_BTHREAD_BUTEX_WAITER_COUNT_IN_VARS
733
    num_waiters << -1;
734
#endif
735
736
0
    bool is_interrupted = false;
737
0
    if (bbw.task_meta->interrupted) {
738
        // Race with set and may consume multiple interruptions, which are OK.
739
0
        bbw.task_meta->interrupted = false;
740
0
        is_interrupted = true;
741
0
    }
742
    // If timed out as well as value unmatched, return ETIMEDOUT.
743
0
    if (WAITER_STATE_TIMEDOUT == bbw.waiter_state) {
744
0
        errno = ETIMEDOUT;
745
0
        return -1;
746
0
    } else if (WAITER_STATE_UNMATCHEDVALUE == bbw.waiter_state) {
747
0
        errno = EWOULDBLOCK;
748
0
        return -1;
749
0
    } else if (is_interrupted) {
750
0
        errno = EINTR;
751
0
        return -1;
752
0
    }
753
0
    return 0;
754
0
}
755
756
}  // namespace bthread
757
758
namespace butil {
759
template <> struct ObjectPoolBlockMaxItem<bthread::Butex> {
760
    static const size_t value = 128;
761
};
762
}