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.
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#ifndef THREAD_FIBER_CHANNEL_H_
16#define THREAD_FIBER_CHANNEL_H_
17
18#include <deque>
19#include <string>
20#include <type_traits>
21#include <utility>
22
23#include "thread/boost_primitives.h"
24#include "thread/cases.h"
26#include "thread/fiber.h"
27#include "thread/select.h"
28
30enum class CopyOrMove {
31 Copy,
32 Move,
33};
34
35static constexpr auto kCopy = CopyOrMove::Copy;
36static constexpr auto kMove = CopyOrMove::Move;
37
38template <typename T>
39T CopyOrMoveOut(T* absl_nonnull item, CopyOrMove strategy) {
40 if (strategy == kCopy) {
41 if constexpr (std::is_copy_constructible_v<T>) {
42 return *static_cast<const T*>(item);
43 } else {
44 LOG(FATAL) << "Copy requested for a move-only channel value";
45 ABSL_ASSUME(false);
46 }
47 }
48
49 if (strategy == kMove) {
50 return std::move(*item);
51 }
52
53 LOG(FATAL) << "Invalid CopyOrMove strategy: " << static_cast<int>(strategy);
54 ABSL_ASSUME(false);
55}
56} // namespace thread::internal
57
58namespace thread {
59template <class T>
60requires std::is_move_assignable_v<T> class Channel;
61}
62
63namespace thread::internal {
64template <typename T>
65struct ReadSelectable final : Selectable {
66 explicit ReadSelectable(Channel<T>* absl_nonnull channel)
67 : channel(channel) {}
68
69 bool Handle(CaseInSelectClause* absl_nonnull reader, bool enqueue) override;
70 void Unregister(CaseInSelectClause* absl_nonnull c) override;
71
72 Channel<T>* absl_nonnull channel;
73};
74
75template <typename T>
76struct WriteSelectable final : Selectable {
77 explicit WriteSelectable(Channel<T>* absl_nonnull channel)
78 : channel(channel) {}
79
80 bool Handle(CaseInSelectClause* absl_nonnull writer, bool enqueue) override;
81 void Unregister(CaseInSelectClause* absl_nonnull c) override;
82
83 Channel<T>* absl_nonnull channel;
84};
85} // namespace thread::internal
86
87namespace thread {
94template <class T>
95class Reader {
96 public:
97 // This class is not copyable or movable.
98 Reader(const Reader&) = delete;
99 Reader& operator=(const Reader&) = delete;
100
111 bool Read(T* absl_nonnull item);
112
142 Case OnRead(T* absl_nonnull item, bool* absl_nonnull ok);
143
144 private:
145 friend class Channel<T>;
146 explicit Reader(Channel<T>* absl_nonnull channel);
147
148 Channel<T>* absl_nonnull channel_;
149};
150
160template <class T>
161class Writer {
162 public:
163 // This class is not copyable or movable.
164 Writer(const Writer&) = delete;
165 Writer& operator=(const Writer&) = delete;
166
174 void Write(const T& item);
175
183 void Write(T&& item);
184
189 void Close();
190
214 Case OnWrite(const T& item);
215
239 Case OnWrite(T&& item);
240
252 bool WriteUnlessCancelled(const T& item);
253
269 bool WriteUnlessCancelled(T&& item);
270
271 private:
272 friend class Channel<T>;
273 explicit Writer(Channel<T>* absl_nonnull channel_state);
274
275 Channel<T>* absl_nonnull channel_;
276};
277
287template <typename T>
288requires std::is_move_assignable_v<T> class Channel {
289 public:
297 explicit Channel(size_t capacity)
298 : capacity_(capacity), closed_(false), rd_(this), wr_(this) {
299 DCHECK(Invariants());
300 }
301
302 // This class is not copyable or movable.
303 Channel(const Channel&) = delete;
304 Channel& operator=(const Channel&) = delete;
305
307 // Ensure exclusive access (to e.g. prevent concurrent Close()).
308 thread::MutexLock lock(&mu_);
309 DCHECK(Invariants());
310 }
311
321 Reader<T>* absl_nonnull reader() { return &reader_; }
322
334 Writer<T>* absl_nonnull writer() { return &writer_; }
335
343 [[nodiscard]] size_t length() const { return Length(); }
344
345 private:
346 friend struct internal::ReadSelectable<T>;
347 friend struct internal::WriteSelectable<T>;
348 friend class Reader<T>;
349 friend class Writer<T>;
350
351 void Close();
352
353 bool Get(T* absl_nonnull dst);
354
355 size_t Length() const {
356 thread::MutexLock lock(&mu_);
357 return queue_.size();
358 }
359
360 Case OnRead(T* absl_nonnull dst, bool* absl_nonnull ok) {
361 return {&rd_, dst, ok};
362 }
363
364 Case OnWrite(const T& item) requires(std::is_copy_constructible_v<T>) {
365 return {&wr_, &item, &internal::kCopy};
366 }
367
368 Case OnWrite(T&& item) { return {&wr_, &item, &internal::kMove}; }
369
370 const size_t capacity_;
371
372 mutable thread::Mutex mu_;
373 std::deque<T> queue_ ABSL_GUARDED_BY(mu_);
374 bool closed_ ABSL_GUARDED_BY(mu_);
375
376 internal::ChannelWaiterState waiters_;
377
378 Reader<T> reader_{this};
379 internal::ReadSelectable<T> rd_;
380
381 Writer<T> writer_{this};
382 internal::WriteSelectable<T> wr_;
383
384 bool Invariants() const ABSL_EXCLUSIVE_LOCKS_REQUIRED
385
386 (mu_) {
387 CHECK_LE(queue_.size(), capacity_) << "Channel queue size exceeds capacity";
388 return true;
389 }
390};
391
392template <typename T>
393requires std::is_move_assignable_v<T> void Channel<T>::Close() {
394 thread::MutexLock lock(&mu_);
395 DCHECK(Invariants());
396 CHECK(!closed_) << "Calling Close() on closed channel";
397 CHECK(waiters_.writers == nullptr)
398 << "Calling Close() on a channel with blocked writers";
399 closed_ = true;
400 this->waiters_.CloseAndReleaseReaders();
401 DCHECK(Invariants());
402}
403
404template <typename T>
405requires std::is_move_assignable_v<T> bool Channel<T>::Get(
406 T* absl_nonnull dst) {
407 bool result;
408 Select({OnRead(dst, &result)});
409 return result;
410}
411
412template <typename T>
413Reader<T>::Reader(Channel<T>* absl_nonnull channel) : channel_(channel) {}
414
415template <typename T>
416bool Reader<T>::Read(T* absl_nonnull item) {
417 return channel_->Get(item);
418}
419
420template <typename T>
421Case Reader<T>::OnRead(T* absl_nonnull item, bool* absl_nonnull ok) {
422 return channel_->OnRead(item, ok);
423}
424
425template <typename T>
426Writer<T>::Writer(Channel<T>* absl_nonnull channel_state)
427 : channel_(channel_state) {}
428
429template <typename T>
430void Writer<T>::Write(const T& item) {
431 Select({OnWrite(item)});
432}
433
434template <typename T>
435void Writer<T>::Write(T&& item) {
436 Select({OnWrite(std::move(item))});
437}
438
439template <typename T>
441 channel_->Close();
442}
443
444template <typename T>
445Case Writer<T>::OnWrite(const T& item) {
446 return channel_->OnWrite(item);
447}
448
449template <typename T>
451 return channel_->OnWrite(std::move(item));
452}
453
454template <typename T>
456 return !Cancelled() && Select({OnCancel(), OnWrite(item)}) == 1;
457}
458
459template <typename T>
461 return !Cancelled() && Select({OnCancel(), OnWrite(std::move(item))}) == 1;
462}
463} // namespace thread
464
465namespace thread::internal {
466template <typename T>
467bool ReadSelectable<T>::Handle(CaseInSelectClause* absl_nonnull reader,
468
469 bool enqueue) {
470 thread::MutexLock lock(&channel->mu_);
471 DCHECK(channel->Invariants());
472
473 T* absl_nonnull dst_item = reader->GetCase()->GetArgPtr<T>(0);
474 bool* absl_nonnull dst_ok = reader->GetCase()->GetArgPtr<bool>(1);
475
476 // Is there a buffered item to read?
477 if (!channel->queue_.empty()) {
478 DVLOG(2) << "Get from buffer";
479 reader->selector->mu.Lock();
480 if (reader->selector->picked_case_index == Selector::kNonePicked) {
481 // Move out of the buffer.
482 *dst_item = std::move(channel->queue_.front());
483 channel->queue_.pop_front();
484 *dst_ok = true;
485 channel->waiters_.UnlockAndReleaseReader(reader);
486
487 // Potentially admit a waiting writer.
488 if (CaseInSelectClause* absl_nonnull unblocked_writer;
489 channel->waiters_.GetWaitingWriter(&unblocked_writer)) {
490 auto* absl_nonnull item = unblocked_writer->GetCase()->GetArgPtr<T>(0);
491 auto copy_or_move =
492 *unblocked_writer->GetCase()->GetArgPtr<CopyOrMove>(1);
493
494 channel->queue_.push_back(CopyOrMoveOut(item, copy_or_move));
495 channel->waiters_.UnlockAndReleaseWriter(unblocked_writer);
496 }
497 } else {
498 // While we weren't technically able to proceed, there's no point in
499 // Select() processing further cases, so we'll still return true below.
500 reader->selector->mu.Unlock();
501 }
502 DCHECK(channel->Invariants());
503 return true;
504 }
505
506 // Try to transfer directly from waiting writer to reader
507 if (CaseInSelectClause* absl_nonnull writer;
508 channel->waiters_.GetMatchingWriter(reader, &writer)) {
509 auto* absl_nonnull item = writer->GetCase()->GetArgPtr<T>(0);
510 auto copy_or_move = *writer->GetCase()->GetArgPtr<CopyOrMove>(1);
511
512 *dst_item = CopyOrMoveOut(item, copy_or_move);
513 *dst_ok = true;
514
515 channel->waiters_.UnlockAndReleaseReader(reader);
516 channel->waiters_.UnlockAndReleaseWriter(writer);
517 DCHECK(channel->Invariants());
518 return true;
519 }
520
521 reader->selector->mu.Lock();
522 // We must guarantee that this case is eligible to proceed before any
523 // side effects can occur.
524 if (reader->selector->picked_case_index != Selector::kNonePicked) {
525 reader->selector->mu.Unlock();
526 // Already handled item
527 DVLOG(2) << "Read cancelled since another selector case done";
528 DCHECK(channel->Invariants());
529 return true;
530 }
531
532 if (channel->closed_) {
533 DVLOG(2) << "Read failing because channel closed";
534 *dst_ok = false;
535 channel->waiters_.UnlockAndReleaseReader(reader);
536 return true;
537 }
538
539 if (enqueue) {
540 // Register with waiting readers
541 DVLOG(2) << "Read waiting";
542 internal::PushBack(&channel->waiters_.readers, reader);
543 }
544
545 reader->selector->mu.Unlock();
546 DCHECK(channel->Invariants());
547 return false;
548}
549
550template <typename T>
551void ReadSelectable<T>::Unregister(CaseInSelectClause* absl_nonnull c) {
552 thread::MutexLock lock(&channel->mu_);
553 internal::UnlinkFromList(&channel->waiters_.readers, c);
554}
555
556template <typename T>
557bool WriteSelectable<T>::Handle(CaseInSelectClause* absl_nonnull writer,
558
559 bool enqueue) {
560 thread::MutexLock lock(&channel->mu_);
561 DCHECK(channel->Invariants());
562 CHECK(!channel->closed_) << "Calling Write() on closed channel";
563
564 // First try to transfer directly from writer to a waiting reader
565 if (CaseInSelectClause* absl_nonnull reader;
566 channel->waiters_.GetMatchingReader(writer, &reader)) {
567 auto* absl_nonnull writer_item = writer->GetCase()->GetArgPtr<T>(0);
568 auto copy_or_move = *writer->GetCase()->GetArgPtr<CopyOrMove>(1);
569
570 auto* absl_nonnull reader_item = reader->GetCase()->GetArgPtr<T>(0);
571 bool* absl_nonnull reader_ok = reader->GetCase()->GetArgPtr<bool>(1);
572
573 *reader_item = CopyOrMoveOut(writer_item, copy_or_move);
574 *reader_ok = true;
575
576 channel->waiters_.UnlockAndReleaseReader(reader);
577 channel->waiters_.UnlockAndReleaseWriter(writer);
578
579 DCHECK(channel->Invariants());
580 return true;
581 }
582
583 writer->selector->mu.Lock();
584 // We must guarantee that this case is eligible to proceed before any
585 // side effects can occur.
586 if (writer->selector->picked_case_index != Selector::kNonePicked) {
587 writer->selector->mu.Unlock();
588 // Already handled item
589 DVLOG(2) << "Write cancelled since another selector case done";
590 DCHECK(channel->Invariants());
591 return true;
592 }
593
594 // Is there room to buffer item?
595 if (channel->queue_.size() < channel->capacity_) {
596 DVLOG(2) << "Add to buffer";
597
598 T* absl_nonnull item = writer->GetCase()->GetArgPtr<T>(0);
599 const CopyOrMove copy_or_move =
600 *writer->GetCase()->GetArgPtr<CopyOrMove>(1);
601
602 channel->queue_.push_back(CopyOrMoveOut(item, copy_or_move));
603 channel->waiters_.UnlockAndReleaseWriter(writer);
604
605 DCHECK(channel->Invariants());
606 return true;
607 }
608
609 if (enqueue) {
610 // Register with waiting writers
611 DVLOG(2) << "Write waiting";
612 internal::PushBack(&channel->waiters_.writers, writer);
613 }
614
615 writer->selector->mu.Unlock();
616 DCHECK(channel->Invariants());
617 return false;
618}
619
620template <typename T>
621void WriteSelectable<T>::Unregister(CaseInSelectClause* absl_nonnull c) {
622 thread::MutexLock lock(&channel->mu_);
623 internal::UnlinkFromList(&channel->waiters_.writers, c);
624}
625} // namespace thread::internal
626
627#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:306
friend class Reader< T >
Definition channel.h:348
Channel(size_t capacity)
Constructs a Channel with the given capacity.
Definition channel.h:297
Reader< T > *absl_nonnull reader()
Returns a Reader that can be used to read items from the channel.
Definition channel.h:321
Channel(const Channel &)=delete
size_t length() const
Returns the instantaneous length of the channel.
Definition channel.h:343
Definition boost_primitives.h:95
Definition boost_primitives.h:34
A Reader is used to read items from a Channel.
Definition channel.h:95
bool Read(T *absl_nonnull item)
Reads an item from the channel, blocking until it can be read.
Definition channel.h:416
Case OnRead(T *absl_nonnull item, bool *absl_nonnull ok)
Returns a Case that can be used to wait for and synchronise on reading an item from the channel.
Definition channel.h:421
Reader(const Reader &)=delete
Reader & operator=(const Reader &)=delete
A Writer is used to write items to a Channel.
Definition channel.h:161
Writer(const Writer &)=delete
void Close()
Marks the channel as closed, notifying any waiting readers.
Definition channel.h:440
Case OnWrite(const T &item)
Returns a Case that can be used to write the given item to the channel.
Definition channel.h:445
Writer & operator=(const Writer &)=delete
bool WriteUnlessCancelled(const T &item)
Returns false iff the calling fiber is cancelled before the value can be written.
Definition channel.h:455
void Write(const T &item)
Writes the given item to the channel, blocking until it can be written.
Definition channel.h:430
Definition channel.h:29
void PushBack(CaseInSelectClause **head, CaseInSelectClause *element)
Definition cases.h:133
void UnlinkFromList(CaseInSelectClause **head, CaseInSelectClause *element)
Definition cases.h:149
static constexpr auto kMove
Definition channel.h:36
T CopyOrMoveOut(T *absl_nonnull item, CopyOrMove strategy)
Definition channel.h:39
CopyOrMove
Definition channel.h:30
static constexpr auto kCopy
Definition channel.h:35
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