Symbian platform (C++)
SDK native APIs, runtime and tooling
Loading...
Searching...
No Matches
bounded_channel.h
Go to the documentation of this file.
1// Copyright 2026 The Action Engine Authors and the Symbian SDK Authors.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7// http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14//
15// SDK-owned fallible bounded queue, separate from A11
16// cpp/thread/thread/channel.h at fcccb8cb6e1e67d7ba0822ac14cece9ee4c7091b.
17// This queue serves nonblocking SDK mailboxes. It does not change the
18// signatures or closed-write contract of thread::Channel.
19
20#ifndef SYMBIAN_CONCURRENCY_BOUNDED_CHANNEL_H_
21#define SYMBIAN_CONCURRENCY_BOUNDED_CHANNEL_H_
22
23#include <cstddef>
24#include <deque>
25#include <type_traits>
26#include <utility>
27
28#include "absl/status/status.h"
29#include "absl/status/statusor.h"
30#include "thread/boost_primitives.h"
31
33
34template <typename T>
35class BoundedChannel;
36
37template <typename T>
39 public:
40 explicit BoundedReader(BoundedChannel<T>* channel) : channel_(channel) {}
41
42 bool Read(T* out) { return channel_->Read(out); }
43
44 absl::StatusOr<bool> TryRead(T* out) { return channel_->TryRead(out); }
45
46 private:
47 BoundedChannel<T>* channel_;
48};
49
50template <typename T>
52 public:
53 explicit BoundedWriter(BoundedChannel<T>* channel) : channel_(channel) {}
54
55 absl::Status Write(T&& item) { return channel_->Write(std::move(item)); }
56
57 absl::Status TryWrite(T&& item) {
58 return channel_->TryWrite(std::move(item));
59 }
60
61 absl::Status Write(const T& item) requires std::is_copy_constructible_v<T> {
62 return channel_->Write(item);
63 }
64
65 absl::Status TryWrite(
66 const T& item) requires std::is_copy_constructible_v<T> {
67 return channel_->TryWrite(item);
68 }
69
70 void Close() { channel_->Close(); }
71
72 private:
73 BoundedChannel<T>* channel_;
74};
75
76// A bounded multi-producer, multi-consumer FIFO. Capacity must be positive.
77// Blocking calls park according to the selected thread primitive backend:
78// Boost fibers on hosts, guest ARM fibers on the event thread, or an OS thread
79// outside a fiber. A failed write
80// preserves an rvalue source. Closed writes return Status under the default
81// no-exceptions profile.
82template <typename T>
84 static_assert(std::is_move_assignable_v<T>);
85
86 public:
87 explicit BoundedChannel(std::size_t capacity)
88 : capacity_(capacity), reader_(this), writer_(this) {}
89
92
93 BoundedReader<T>* reader() { return &reader_; }
94
95 BoundedWriter<T>* writer() { return &writer_; }
96
97 absl::Status TryWrite(T&& item) { return TryWriteMoved(&item); }
98
99 absl::Status TryWrite(
100 const T& item) requires std::is_copy_constructible_v<T> {
101 {
102 thread::MutexLock lock(&mu_);
103 absl::Status ready = CheckWritable();
104 if (!ready.ok()) {
105 return ready;
106 }
107 queue_.push_back(item);
108 }
109 readers_.Signal();
110 return absl::OkStatus();
111 }
112
113 absl::Status Write(T&& item) { return WriteMoved(&item); }
114
115 absl::Status Write(const T& item) requires std::is_copy_constructible_v<T> {
116 {
117 thread::MutexLock lock(&mu_);
118 if (capacity_ == 0) {
119 return absl::InvalidArgumentError(
120 "BoundedChannel capacity must be positive");
121 }
122 while (!closed_ && queue_.size() >= capacity_) {
123 writers_.Wait(&mu_);
124 }
125 if (closed_) {
126 return absl::FailedPreconditionError("BoundedChannel is closed");
127 }
128 queue_.push_back(item);
129 }
130 readers_.Signal();
131 return absl::OkStatus();
132 }
133
134 // True means one item was read; false means closed and drained. An open,
135 // empty channel returns Unavailable so pollers can distinguish the states.
136 absl::StatusOr<bool> TryRead(T* out) {
137 {
138 thread::MutexLock lock(&mu_);
139 if (queue_.empty()) {
140 if (closed_) {
141 return false;
142 }
143 return absl::UnavailableError("BoundedChannel is empty");
144 }
145 *out = std::move(queue_.front());
146 queue_.pop_front();
147 }
148 writers_.Signal();
149 return true;
150 }
151
152 bool Read(T* out) {
153 {
154 thread::MutexLock lock(&mu_);
155 while (!closed_ && queue_.empty()) {
156 readers_.Wait(&mu_);
157 }
158 if (queue_.empty()) {
159 return false;
160 }
161 *out = std::move(queue_.front());
162 queue_.pop_front();
163 }
164 writers_.Signal();
165 return true;
166 }
167
168 // Close rejects writes, wakes blocked readers/writers, and permits queued
169 // items to drain. Repeating Close is harmless for owner shutdown.
170 void Close() {
171 {
172 thread::MutexLock lock(&mu_);
173 closed_ = true;
174 }
175 readers_.SignalAll();
176 writers_.SignalAll();
177 }
178
179 // Discard closes and releases queued values outside the queue lock.
180 void Discard() {
181 std::deque<T> retired;
182 {
183 thread::MutexLock lock(&mu_);
184 closed_ = true;
185 retired.swap(queue_);
186 }
187 readers_.SignalAll();
188 writers_.SignalAll();
189 }
190
191 std::size_t Size() const {
192 thread::MutexLock lock(&mu_);
193 return queue_.size();
194 }
195
196 private:
197 absl::Status CheckWritable() const {
198 if (closed_) {
199 return absl::FailedPreconditionError("BoundedChannel is closed");
200 }
201 if (capacity_ == 0) {
202 return absl::InvalidArgumentError(
203 "BoundedChannel capacity must be positive");
204 }
205 if (queue_.size() >= capacity_) {
206 return absl::ResourceExhaustedError("BoundedChannel is full");
207 }
208 return absl::OkStatus();
209 }
210
211 absl::Status TryWriteMoved(T* item) {
212 {
213 thread::MutexLock lock(&mu_);
214 absl::Status ready = CheckWritable();
215 if (!ready.ok()) {
216 return ready;
217 }
218 queue_.push_back(std::move(*item));
219 }
220 readers_.Signal();
221 return absl::OkStatus();
222 }
223
224 absl::Status WriteMoved(T* item) {
225 {
226 thread::MutexLock lock(&mu_);
227 if (capacity_ == 0) {
228 return absl::InvalidArgumentError(
229 "BoundedChannel capacity must be positive");
230 }
231 while (!closed_ && queue_.size() >= capacity_) {
232 writers_.Wait(&mu_);
233 }
234 if (closed_) {
235 return absl::FailedPreconditionError("BoundedChannel is closed");
236 }
237 queue_.push_back(std::move(*item));
238 }
239 readers_.Signal();
240 return absl::OkStatus();
241 }
242
243 const std::size_t capacity_;
244 mutable thread::Mutex mu_;
245 thread::CondVar readers_;
246 thread::CondVar writers_;
247 std::deque<T> queue_;
248 bool closed_ = false;
249 BoundedReader<T> reader_;
250 BoundedWriter<T> writer_;
251};
252
253} // namespace symbian::concurrency
254
255#endif // SYMBIAN_CONCURRENCY_BOUNDED_CHANNEL_H_
Definition bounded_channel.h:83
BoundedChannel(std::size_t capacity)
Definition bounded_channel.h:87
BoundedChannel(const BoundedChannel &)=delete
BoundedWriter< T > * writer()
Definition bounded_channel.h:95
void Close()
Definition bounded_channel.h:170
absl::StatusOr< bool > TryRead(T *out)
Definition bounded_channel.h:136
absl::Status Write(const T &item)
Definition bounded_channel.h:115
BoundedChannel & operator=(const BoundedChannel &)=delete
absl::Status TryWrite(T &&item)
Definition bounded_channel.h:97
bool Read(T *out)
Definition bounded_channel.h:152
BoundedReader< T > * reader()
Definition bounded_channel.h:93
std::size_t Size() const
Definition bounded_channel.h:191
absl::Status Write(T &&item)
Definition bounded_channel.h:113
void Discard()
Definition bounded_channel.h:180
absl::Status TryWrite(const T &item)
Definition bounded_channel.h:99
Definition bounded_channel.h:38
bool Read(T *out)
Definition bounded_channel.h:42
BoundedReader(BoundedChannel< T > *channel)
Definition bounded_channel.h:40
absl::StatusOr< bool > TryRead(T *out)
Definition bounded_channel.h:44
Definition bounded_channel.h:51
BoundedWriter(BoundedChannel< T > *channel)
Definition bounded_channel.h:53
absl::Status TryWrite(T &&item)
Definition bounded_channel.h:57
absl::Status Write(T &&item)
Definition bounded_channel.h:55
absl::Status TryWrite(const T &item)
Definition bounded_channel.h:65
void Close()
Definition bounded_channel.h:70
absl::Status Write(const T &item)
Definition bounded_channel.h:61
Definition boost_primitives.h:110
void Signal() noexcept
Definition boost_primitives.h:210
void Wait(Mutex *mu) noexcept
Definition boost_primitives.h:116
void SignalAll() noexcept
Definition boost_primitives.h:230
Definition boost_primitives.h:95
Definition boost_primitives.h:34
Definition bounded_channel.h:32
uint8_t out
Definition usb.cc:194