15#ifndef THREAD_FIBER_CHANNEL_H_
16#define THREAD_FIBER_CHANNEL_H_
23#include "thread/boost_primitives.h"
24#include "thread/cases.h"
26#include "thread/fiber.h"
27#include "thread/select.h"
40 if (strategy ==
kCopy) {
41 if constexpr (std::is_copy_constructible_v<T>) {
42 return *
static_cast<const T*
>(item);
44 LOG(FATAL) <<
"Copy requested for a move-only channel value";
49 if (strategy ==
kMove) {
50 return std::move(*item);
53 LOG(FATAL) <<
"Invalid CopyOrMove strategy: " <<
static_cast<int>(strategy);
60requires std::is_move_assignable_v<T>
class Channel;
65struct ReadSelectable final : Selectable {
66 explicit ReadSelectable(Channel<T>* absl_nonnull channel)
69 bool Handle(CaseInSelectClause* absl_nonnull reader,
bool enqueue)
override;
70 void Unregister(CaseInSelectClause* absl_nonnull c)
override;
72 Channel<T>* absl_nonnull channel;
76struct WriteSelectable final : Selectable {
77 explicit WriteSelectable(Channel<T>* absl_nonnull channel)
80 bool Handle(CaseInSelectClause* absl_nonnull writer,
bool enqueue)
override;
81 void Unregister(CaseInSelectClause* absl_nonnull c)
override;
83 Channel<T>* absl_nonnull channel;
111 bool Read(T* absl_nonnull item);
174 void Write(
const T& item);
183 void Write(T&& item);
288requires std::is_move_assignable_v<T>
class Channel {
298 : capacity_(capacity), closed_(false), rd_(this), wr_(this) {
299 DCHECK(Invariants());
309 DCHECK(Invariants());
343 [[nodiscard]]
size_t length()
const {
return Length(); }
346 friend struct internal::ReadSelectable<T>;
347 friend struct internal::WriteSelectable<T>;
353 bool Get(T* absl_nonnull dst);
355 size_t Length()
const {
357 return queue_.size();
360 Case OnRead(T* absl_nonnull dst,
bool* absl_nonnull ok) {
361 return {&rd_, dst, ok};
364 Case OnWrite(
const T& item)
requires(std::is_copy_constructible_v<T>) {
370 const size_t capacity_;
373 std::deque<T> queue_ ABSL_GUARDED_BY(mu_);
374 bool closed_ ABSL_GUARDED_BY(mu_);
376 internal::ChannelWaiterState waiters_;
379 internal::ReadSelectable<T> rd_;
382 internal::WriteSelectable<T> wr_;
384 bool Invariants() const ABSL_EXCLUSIVE_LOCKS_REQUIRED
387 CHECK_LE(queue_.size(), capacity_) <<
"Channel queue size exceeds capacity";
393requires std::is_move_assignable_v<T>
void Channel<T>::Close() {
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";
400 this->waiters_.CloseAndReleaseReaders();
401 DCHECK(Invariants());
405requires std::is_move_assignable_v<T>
bool Channel<T>::Get(
406 T* absl_nonnull dst) {
408 Select({OnRead(dst, &result)});
417 return channel_->Get(item);
422 return channel_->OnRead(item, ok);
427 : channel_(channel_state) {}
436 Select({OnWrite(std::move(item))});
446 return channel_->OnWrite(item);
451 return channel_->OnWrite(std::move(item));
467bool ReadSelectable<T>::Handle(CaseInSelectClause* absl_nonnull reader,
471 DCHECK(channel->Invariants());
473 T* absl_nonnull dst_item = reader->GetCase()->GetArgPtr<T>(0);
474 bool* absl_nonnull dst_ok = reader->GetCase()->GetArgPtr<
bool>(1);
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) {
482 *dst_item = std::move(channel->queue_.front());
483 channel->queue_.pop_front();
485 channel->waiters_.UnlockAndReleaseReader(reader);
488 if (CaseInSelectClause* absl_nonnull unblocked_writer;
489 channel->waiters_.GetWaitingWriter(&unblocked_writer)) {
490 auto* absl_nonnull item = unblocked_writer->GetCase()->GetArgPtr<T>(0);
492 *unblocked_writer->GetCase()->GetArgPtr<
CopyOrMove>(1);
494 channel->queue_.push_back(
CopyOrMoveOut(item, copy_or_move));
495 channel->waiters_.UnlockAndReleaseWriter(unblocked_writer);
500 reader->selector->mu.Unlock();
502 DCHECK(channel->Invariants());
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);
515 channel->waiters_.UnlockAndReleaseReader(reader);
516 channel->waiters_.UnlockAndReleaseWriter(writer);
517 DCHECK(channel->Invariants());
521 reader->selector->mu.Lock();
524 if (reader->selector->picked_case_index != Selector::kNonePicked) {
525 reader->selector->mu.Unlock();
527 DVLOG(2) <<
"Read cancelled since another selector case done";
528 DCHECK(channel->Invariants());
532 if (channel->closed_) {
533 DVLOG(2) <<
"Read failing because channel closed";
535 channel->waiters_.UnlockAndReleaseReader(reader);
541 DVLOG(2) <<
"Read waiting";
545 reader->selector->mu.Unlock();
546 DCHECK(channel->Invariants());
551void ReadSelectable<T>::Unregister(CaseInSelectClause* absl_nonnull c) {
557bool WriteSelectable<T>::Handle(CaseInSelectClause* absl_nonnull writer,
561 DCHECK(channel->Invariants());
562 CHECK(!channel->closed_) <<
"Calling Write() on closed channel";
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);
570 auto* absl_nonnull reader_item = reader->GetCase()->GetArgPtr<T>(0);
571 bool* absl_nonnull reader_ok = reader->GetCase()->GetArgPtr<
bool>(1);
576 channel->waiters_.UnlockAndReleaseReader(reader);
577 channel->waiters_.UnlockAndReleaseWriter(writer);
579 DCHECK(channel->Invariants());
583 writer->selector->mu.Lock();
586 if (writer->selector->picked_case_index != Selector::kNonePicked) {
587 writer->selector->mu.Unlock();
589 DVLOG(2) <<
"Write cancelled since another selector case done";
590 DCHECK(channel->Invariants());
595 if (channel->queue_.size() < channel->capacity_) {
596 DVLOG(2) <<
"Add to buffer";
598 T* absl_nonnull item = writer->GetCase()->GetArgPtr<T>(0);
602 channel->queue_.push_back(
CopyOrMoveOut(item, copy_or_move));
603 channel->waiters_.UnlockAndReleaseWriter(writer);
605 DCHECK(channel->Invariants());
611 DVLOG(2) <<
"Write waiting";
615 writer->selector->mu.Unlock();
616 DCHECK(channel->Invariants());
621void WriteSelectable<T>::Unregister(CaseInSelectClause* absl_nonnull c) {
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
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
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