Coverage Report

Created: 2026-09-01 06:31

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/brpc/src/bthread/work_stealing_queue.h
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 10 17:40:58 CST 2012
21
22
#ifndef BTHREAD_WORK_STEALING_QUEUE_H
23
#define BTHREAD_WORK_STEALING_QUEUE_H
24
25
#include "butil/macros.h"
26
#include "butil/atomicops.h"
27
#include "butil/logging.h"
28
29
namespace bthread {
30
31
template <typename T>
32
class WorkStealingQueue {
33
public:
34
    WorkStealingQueue()
35
10
        : _bottom(1)
36
10
        , _capacity(0)
37
10
        , _buffer(nullptr)
38
10
        , _top(1) {
39
10
    }
40
41
0
    ~WorkStealingQueue() {
42
0
        delete [] _buffer;
43
0
        _buffer = nullptr;
44
0
    }
45
46
9
    int init(size_t capacity) {
47
9
        if (_capacity != 0) {
48
0
            LOG(ERROR) << "Already initialized";
49
0
            return -1;
50
0
        }
51
9
        if (capacity == 0) {
52
0
            LOG(ERROR) << "Invalid capacity=" << capacity;
53
0
            return -1;
54
0
        }
55
9
        if (capacity & (capacity - 1)) {
56
0
            LOG(ERROR) << "Invalid capacity=" << capacity
57
0
                       << " which must be power of 2";
58
0
            return -1;
59
0
        }
60
9
        _buffer = new T[capacity];
61
9
        _capacity = capacity;
62
9
        return 0;
63
9
    }
64
65
    // Push an item into the queue.
66
    // Returns true on pushed.
67
    // May run in parallel with steal().
68
    // Never run in parallel with pop() or another push().
69
0
    bool push(const T& x) {
70
0
        const size_t b = _bottom.load(butil::memory_order_relaxed);
71
0
        const size_t t = _top.load(butil::memory_order_acquire);
72
0
        if (b >= t + _capacity) { // Full queue.
73
0
            return false;
74
0
        }
75
0
        _buffer[b & (_capacity - 1)] = x;
76
0
        _bottom.store(b + 1, butil::memory_order_release);
77
0
        return true;
78
0
    }
79
80
    // Pop an item from the queue.
81
    // Returns true on popped and the item is written to `val'.
82
    // May run in parallel with steal().
83
    // Never run in parallel with push() or another pop().
84
1
    bool pop(T* val) {
85
1
        const size_t b = _bottom.load(butil::memory_order_relaxed);
86
1
        size_t t = _top.load(butil::memory_order_relaxed);
87
1
        if (t >= b) {
88
            // fast check since we call pop() in each sched.
89
            // Stale _top which is smaller should not enter this branch.
90
1
            return false;
91
1
        }
92
0
        const size_t newb = b - 1;
93
0
        _bottom.store(newb, butil::memory_order_relaxed);
94
0
        butil::atomic_thread_fence(butil::memory_order_seq_cst);
95
0
        t = _top.load(butil::memory_order_relaxed);
96
0
        if (t > newb) {
97
0
            _bottom.store(b, butil::memory_order_relaxed);
98
0
            return false;
99
0
        }
100
0
        *val = _buffer[newb & (_capacity - 1)];
101
0
        if (t != newb) {
102
0
            return true;
103
0
        }
104
        // Single last element, compete with steal()
105
0
        const bool popped = _top.compare_exchange_strong(
106
0
            t, t + 1, butil::memory_order_seq_cst, butil::memory_order_relaxed);
107
0
        _bottom.store(b, butil::memory_order_relaxed);
108
0
        return popped;
109
0
    }
110
111
    // Steal one item from the queue.
112
    // Returns true on stolen.
113
    // May run in parallel with push() pop() or another steal().
114
16
    bool steal(T* val) {
115
16
        size_t t = _top.load(butil::memory_order_acquire);
116
16
        size_t b = _bottom.load(butil::memory_order_acquire);
117
16
        if (t >= b) {
118
            // Permit false negative for performance considerations.
119
16
            return false;
120
16
        }
121
0
        do {
122
0
            butil::atomic_thread_fence(butil::memory_order_seq_cst);
123
0
            b = _bottom.load(butil::memory_order_acquire);
124
0
            if (t >= b) {
125
0
                return false;
126
0
            }
127
0
            *val = _buffer[t & (_capacity - 1)];
128
0
        } while (!_top.compare_exchange_weak(t, t + 1,
129
0
                                               butil::memory_order_seq_cst,
130
0
                                               butil::memory_order_relaxed));
131
0
        return true;
132
0
    }
133
134
0
    size_t volatile_size() const {
135
0
        const size_t b = _bottom.load(butil::memory_order_relaxed);
136
0
        const size_t t = _top.load(butil::memory_order_relaxed);
137
0
        return (b <= t ? 0 : (b - t));
138
0
    }
139
140
0
    size_t capacity() const { return _capacity; }
141
142
private:
143
    // Copying a concurrent structure makes no sense.
144
    DISALLOW_COPY_AND_ASSIGN(WorkStealingQueue);
145
146
    butil::atomic<size_t> _bottom;
147
    size_t _capacity;
148
    T* _buffer;
149
    BAIDU_CACHELINE_ALIGNMENT butil::atomic<size_t> _top;
150
};
151
152
}  // namespace bthread
153
154
#endif  // BTHREAD_WORK_STEALING_QUEUE_H