Symbian platform (C++)
SDK native APIs, runtime and tooling
Loading...
Searching...
No Matches
channel.h
Go to the documentation of this file.
1// Copyright 2026 The Action Engine Authors and the Symbian SDK Authors.
2// Licensed under the Apache License, Version 2.0.
3// Guest adaptation of A11 cpp/thread/thread/channel.h at
4// fcccb8cb6e1e67d7ba0822ac14cece9ee4c7091b.
5
6#ifndef THREAD_FIBER_CHANNEL_H_
7#define THREAD_FIBER_CHANNEL_H_
8
9#include <cstddef>
10#include <cstdlib>
11#include <deque>
12#include <functional>
13#include <memory>
14#include <type_traits>
15#include <utility>
16
17#include "absl/container/inlined_vector.h"
18#include "thread/boost_primitives.h"
19#include "thread/fiber.h"
20#include "thread/select.h"
21
22namespace thread {
23
24template <typename T>
25requires std::is_move_assignable_v<T> class Channel;
26
27template <typename T>
28class Reader {
29 public:
30 Reader(const Reader&) = delete;
31 Reader& operator=(const Reader&) = delete;
32
33 bool Read(T* out) { return channel_->Read(out); }
34
35 Case OnRead(T* out, bool* ok) { return channel_->OnRead(out, ok); }
36
37 private:
38 friend class Channel<T>;
39
40 explicit Reader(Channel<T>* channel) : channel_(channel) {}
41
42 Channel<T>* channel_;
43};
44
45template <typename T>
46class Writer {
47 public:
48 Writer(const Writer&) = delete;
49 Writer& operator=(const Writer&) = delete;
50
51 void Write(T&& item) { channel_->Write(std::move(item)); }
52
53 void Write(const T& item) { channel_->Write(item); }
54
55 Case OnWrite(T&& item) { return channel_->OnWrite(std::move(item)); }
56
57 Case OnWrite(const T& item) { return channel_->OnWrite(item); }
58
59 bool WriteUnlessCancelled(T&& item) {
60 return !thread::Cancelled() &&
61 Select({thread::OnCancel(), OnWrite(std::move(item))}) == 1;
62 }
63
64 bool WriteUnlessCancelled(const T& item) {
65 return !thread::Cancelled() &&
66 Select({thread::OnCancel(), OnWrite(item)}) == 1;
67 }
68
69 void Close() { channel_->Close(); }
70
71 private:
72 friend class Channel<T>;
73
74 explicit Writer(Channel<T>* channel) : channel_(channel) {}
75
76 Channel<T>* channel_;
77};
78
79// A11-compatible selectable buffered/rendezvous channel. Every transfer and
80// selection decision is made under the channel and selector locks. Wakes run
81// after those locks are released, so a ready fiber can reenter the channel.
82template <typename T>
83requires std::is_move_assignable_v<T> class Channel {
84 enum class Transfer { kCopy, kMove };
85 inline static constexpr Transfer kCopy = Transfer::kCopy;
86 inline static constexpr Transfer kMove = Transfer::kMove;
87
88 class ReadSelectable final : public internal::Selectable {
89 public:
90 explicit ReadSelectable(Channel* channel) : channel_(channel) {}
91
92 bool Handle(internal::CaseInSelectClause* state, bool enqueue) override {
93 return channel_->HandleRead(state, enqueue);
94 }
95
96 void Unregister(internal::CaseInSelectClause* state) override {
97 channel_->Unregister(&channel_->readers_, state);
98 }
99
100 private:
101 Channel* channel_;
102 };
103
104 class WriteSelectable final : public internal::Selectable {
105 public:
106 explicit WriteSelectable(Channel* channel) : channel_(channel) {}
107
108 bool Handle(internal::CaseInSelectClause* state, bool enqueue) override {
109 return channel_->HandleWrite(state, enqueue);
110 }
111
112 void Unregister(internal::CaseInSelectClause* state) override {
113 channel_->Unregister(&channel_->writers_, state);
114 }
115
116 private:
117 Channel* channel_;
118 };
119
120 public:
121 explicit Channel(std::size_t capacity)
122 : capacity_(capacity),
123 reader_(this),
124 read_selectable_(this),
125 writer_(this),
126 write_selectable_(this) {}
127
128 Channel(const Channel&) = delete;
129 Channel& operator=(const Channel&) = delete;
130
132 MutexLock lock(&mu_);
133 if (readers_ != nullptr || writers_ != nullptr) {
134 std::abort();
135 }
136 }
137
138 Reader<T>* reader() { return &reader_; }
139
140 Writer<T>* writer() { return &writer_; }
141
142 std::size_t length() const {
143 MutexLock lock(&mu_);
144 return queue_.size();
145 }
146
147 private:
148 friend class Reader<T>;
149 friend class Writer<T>;
150
151 bool Read(T* out) {
152 bool ok = false;
153 Select({OnRead(out, &ok)});
154 return ok;
155 }
156
157 void Write(T&& item) { Select({OnWrite(std::move(item))}); }
158
159 void Write(const T& item) requires std::is_copy_constructible_v<T> {
160 Select({OnWrite(item)});
161 }
162
163 Case OnRead(T* out, bool* ok) { return {&read_selectable_, out, ok}; }
164
165 Case OnWrite(T&& item) { return MakeWriteCase(&item, kMove); }
166
167 Case OnWrite(const T& item) requires std::is_copy_constructible_v<T> {
168 return MakeWriteCase(const_cast<T*>(&item), kCopy);
169 }
170
171 void Close() { CloseImpl(); }
172
173 using State = internal::CaseInSelectClause;
174 using Selector = internal::Selector;
175 using Wakes = absl::InlinedVector<std::shared_ptr<Selector>, 4>;
176
177 static T TransferItem(T* item, Transfer strategy) {
178 if (strategy == kMove) {
179 return std::move(*item);
180 }
181 if constexpr (std::is_copy_constructible_v<T>) {
182 return *item;
183 }
184 std::abort();
185 }
186
187 static void Wake(Wakes* wakes) {
188 for (const auto& selector : *wakes) {
189 selector->cv.Signal();
190 }
191 }
192
193 static void UnlinkIfQueued(State** head, State* state) {
194 if (state->prev != nullptr) {
195 internal::UnlinkFromList(head, state);
196 }
197 }
198
199 static bool LockPair(State* first, State* second) {
200 Selector* left = first->selector.get();
201 Selector* right = second->selector.get();
202 if (left == right) {
203 return false;
204 }
205 if (std::less<Selector*>{}(right, left)) {
206 std::swap(left, right);
207 }
208 left->mu.Lock();
209 right->mu.Lock();
210 if (left->picked_case_index == Selector::kNonePicked &&
211 right->picked_case_index == Selector::kNonePicked) {
212 return true;
213 }
214 right->mu.Unlock();
215 left->mu.Unlock();
216 return false;
217 }
218
219 static void UnlockPair(State* first, State* second) {
220 first->selector->mu.Unlock();
221 second->selector->mu.Unlock();
222 }
223
224 static State* FindMatch(State* head, State* counterpart) {
225 if (head == nullptr) {
226 return nullptr;
227 }
228 State* current = head;
229 do {
230 if (LockPair(current, counterpart)) {
231 return current;
232 }
233 current = current->next;
234 } while (current != head);
235 return nullptr;
236 }
237
238 // Called with channel mu held. A queued writer can refill one freed buffer
239 // slot, including after the reader's immediate-case transfer.
240 void AdmitWriter(Wakes* wakes) {
241 if (writers_ == nullptr || queue_.size() >= capacity_) {
242 return;
243 }
244 State* current = writers_;
245 do {
246 State* next = current->next;
247 MutexLock selector_lock(&current->selector->mu);
248 if (current->selector->picked_case_index == Selector::kNonePicked) {
249 T* item = current->GetCase()->template GetArgPtr<T>(0);
250 Transfer strategy =
251 *current->GetCase()->template GetArgPtr<Transfer>(1);
252 queue_.push_back(TransferItem(item, strategy));
253 current->TryPick();
254 UnlinkIfQueued(&writers_, current);
255 wakes->push_back(current->selector);
256 return;
257 }
258 current = next;
259 } while (current != writers_ && writers_ != nullptr);
260 }
261
262 bool HandleRead(State* reader, bool enqueue) {
263 Wakes wakes;
264 bool ready = false;
265 {
266 MutexLock lock(&mu_);
267 T* out = reader->GetCase()->template GetArgPtr<T>(0);
268 bool* ok = reader->GetCase()->template GetArgPtr<bool>(1);
269 if (!queue_.empty()) {
270 bool consumed = false;
271 {
272 MutexLock selector_lock(&reader->selector->mu);
273 if (reader->TryPick()) {
274 *out = std::move(queue_.front());
275 queue_.pop_front();
276 *ok = true;
277 consumed = true;
278 }
279 }
280 if (consumed) {
281 AdmitWriter(&wakes);
282 }
283 ready = true;
284 } else if (State* writer = FindMatch(writers_, reader)) {
285 T* item = writer->GetCase()->template GetArgPtr<T>(0);
286 Transfer strategy = *writer->GetCase()->template GetArgPtr<Transfer>(1);
287 *out = TransferItem(item, strategy);
288 *ok = true;
289 reader->TryPick();
290 writer->TryPick();
291 UnlinkIfQueued(&writers_, writer);
292 wakes.push_back(writer->selector);
293 UnlockPair(writer, reader);
294 ready = true;
295 } else {
296 MutexLock selector_lock(&reader->selector->mu);
297 if (reader->selector->picked_case_index != Selector::kNonePicked) {
298 ready = true;
299 } else if (closed_) {
300 *ok = false;
301 reader->TryPick();
302 ready = true;
303 } else if (enqueue) {
304 internal::PushBack(&readers_, reader);
305 }
306 }
307 }
308 Wake(&wakes);
309 return ready;
310 }
311
312 bool HandleWrite(State* writer, bool enqueue) {
313 Wakes wakes;
314 bool ready = false;
315 {
316 MutexLock lock(&mu_);
317 if (closed_) {
318 std::abort(); // A11's OnWrite on a closed channel is fatal.
319 }
320 State* reader = closed_ ? nullptr : FindMatch(readers_, writer);
321 if (reader != nullptr) {
322 T* item = writer->GetCase()->template GetArgPtr<T>(0);
323 Transfer strategy = *writer->GetCase()->template GetArgPtr<Transfer>(1);
324 T* out = reader->GetCase()->template GetArgPtr<T>(0);
325 bool* ok = reader->GetCase()->template GetArgPtr<bool>(1);
326 *out = TransferItem(item, strategy);
327 *ok = true;
328 reader->TryPick();
329 writer->TryPick();
330 UnlinkIfQueued(&readers_, reader);
331 wakes.push_back(reader->selector);
332 UnlockPair(reader, writer);
333 ready = true;
334 } else {
335 MutexLock selector_lock(&writer->selector->mu);
336 if (writer->selector->picked_case_index != Selector::kNonePicked) {
337 ready = true;
338 } else if (queue_.size() < capacity_) {
339 T* item = writer->GetCase()->template GetArgPtr<T>(0);
340 Transfer strategy =
341 *writer->GetCase()->template GetArgPtr<Transfer>(1);
342 queue_.push_back(TransferItem(item, strategy));
343 writer->TryPick();
344 ready = true;
345 } else if (enqueue) {
346 internal::PushBack(&writers_, writer);
347 }
348 }
349 }
350 Wake(&wakes);
351 return ready;
352 }
353
354 void Unregister(State** head, State* state) {
355 MutexLock lock(&mu_);
356 UnlinkIfQueued(head, state);
357 }
358
359 Case MakeWriteCase(T* item, Transfer strategy) {
360 Case result(&write_selectable_);
361 result.AddArg(item);
362 result.AddArg(strategy == kCopy ? &kCopy : &kMove);
363 return result;
364 }
365
366 void CloseImpl() {
367 Wakes wakes;
368 {
369 MutexLock lock(&mu_);
370 if (closed_ || writers_ != nullptr) {
371 std::abort();
372 }
373 closed_ = true;
374 while (readers_ != nullptr) {
375 State* state = readers_;
376 {
377 MutexLock selector_lock(&state->selector->mu);
378 if (state->TryPick()) {
379 *state->GetCase()->template GetArgPtr<bool>(1) = false;
380 wakes.push_back(state->selector);
381 }
382 internal::UnlinkFromList(&readers_, state);
383 }
384 }
385 }
386 Wake(&wakes);
387 }
388
389 const std::size_t capacity_;
390 mutable Mutex mu_;
391 std::deque<T> queue_;
392 bool closed_ = false;
393 State* readers_ = nullptr;
394 State* writers_ = nullptr;
395 Reader<T> reader_;
396 ReadSelectable read_selectable_;
397 Writer<T> writer_;
398 WriteSelectable write_selectable_;
399};
400
401} // namespace thread
402
403#endif // THREAD_FIBER_CHANNEL_H_
A Channel is a bounded FIFO queue that allows multiple readers and writers to communicate with each o...
Definition channel.h:83
Writer< T > *absl_nonnull writer()
Returns a Writer that can be used to write items to the channel.
Definition channel.h:334
friend class Writer< T >
Definition channel.h:349
Channel & operator=(const Channel &)=delete
~Channel()
Definition channel.h:131
friend class Reader< T >
Definition channel.h:348
Writer< T > * writer()
Definition channel.h:140
Channel(std::size_t capacity)
Definition channel.h:121
Reader< T > *absl_nonnull reader()
Returns a Reader that can be used to read items from the channel.
Definition channel.h:321
std::size_t length() const
Definition channel.h:142
Channel(const Channel &)=delete
Reader< T > * reader()
Definition channel.h:138
Definition boost_primitives.h:95
A Reader is used to read items from a Channel.
Definition channel.h:95
friend class Channel< T >
Definition channel.h:145
bool Read(T *out)
Definition channel.h:33
Case OnRead(T *out, bool *ok)
Definition channel.h:35
Reader(const Reader &)=delete
Reader & operator=(const Reader &)=delete
A Writer is used to write items to a Channel.
Definition channel.h:161
Case OnWrite(T &&item)
Definition channel.h:55
friend class Channel< T >
Definition channel.h:272
Writer(const Writer &)=delete
void Close()
Definition channel.h:69
Case OnWrite(const T &item)
Definition channel.h:57
Writer & operator=(const Writer &)=delete
bool WriteUnlessCancelled(const T &item)
Definition channel.h:64
void Write(T &&item)
Definition channel.h:51
void Write(const T &item)
Definition channel.h:53
bool WriteUnlessCancelled(T &&item)
Definition channel.h:59
void PushBack(CaseInSelectClause **head, CaseInSelectClause *element)
Definition cases.h:133
void UnlinkFromList(CaseInSelectClause **head, CaseInSelectClause *element)
Definition cases.h:149
Definition channel.h:29
bool Cancelled()
Definition fiber.cc:289
Case OnCancel()
Definition fiber.cc:294
int Select(const CaseArray &cases)
Returns the index of the first case that is ready, blocking until one is.
Definition select.h:18
A Case represents a selectable case in a Select statement.
Definition cases.h:55
uint8_t out
Definition usb.cc:194