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;
88 class ReadSelectable final :
public internal::Selectable {
90 explicit ReadSelectable(
Channel* channel) : channel_(channel) {}
92 bool Handle(internal::CaseInSelectClause* state,
bool enqueue)
override {
93 return channel_->HandleRead(state, enqueue);
96 void Unregister(internal::CaseInSelectClause* state)
override {
97 channel_->Unregister(&channel_->readers_, state);
104 class WriteSelectable final :
public internal::Selectable {
106 explicit WriteSelectable(
Channel* channel) : channel_(channel) {}
108 bool Handle(internal::CaseInSelectClause* state,
bool enqueue)
override {
109 return channel_->HandleWrite(state, enqueue);
112 void Unregister(internal::CaseInSelectClause* state)
override {
113 channel_->Unregister(&channel_->writers_, state);
122 : capacity_(capacity),
124 read_selectable_(this),
126 write_selectable_(this) {}
133 if (readers_ !=
nullptr || writers_ !=
nullptr) {
144 return queue_.size();
157 void Write(T&& item) {
Select({OnWrite(std::move(item))}); }
159 void Write(
const T& item)
requires std::is_copy_constructible_v<T> {
163 Case OnRead(T*
out,
bool* ok) {
return {&read_selectable_,
out, ok}; }
165 Case OnWrite(T&& item) {
return MakeWriteCase(&item, kMove); }
167 Case OnWrite(
const T& item)
requires std::is_copy_constructible_v<T> {
168 return MakeWriteCase(
const_cast<T*
>(&item), kCopy);
171 void Close() { CloseImpl(); }
173 using State = internal::CaseInSelectClause;
174 using Selector = internal::Selector;
175 using Wakes = absl::InlinedVector<std::shared_ptr<Selector>, 4>;
177 static T TransferItem(T* item, Transfer strategy) {
178 if (strategy == kMove) {
179 return std::move(*item);
181 if constexpr (std::is_copy_constructible_v<T>) {
187 static void Wake(Wakes* wakes) {
188 for (
const auto& selector : *wakes) {
189 selector->cv.Signal();
193 static void UnlinkIfQueued(State** head, State* state) {
194 if (state->prev !=
nullptr) {
199 static bool LockPair(State* first, State* second) {
200 Selector* left = first->selector.get();
201 Selector* right = second->selector.get();
205 if (std::less<Selector*>{}(right, left)) {
206 std::swap(left, right);
210 if (left->picked_case_index == Selector::kNonePicked &&
211 right->picked_case_index == Selector::kNonePicked) {
219 static void UnlockPair(State* first, State* second) {
220 first->selector->mu.Unlock();
221 second->selector->mu.Unlock();
224 static State* FindMatch(State* head, State* counterpart) {
225 if (head ==
nullptr) {
228 State* current = head;
230 if (LockPair(current, counterpart)) {
233 current = current->next;
234 }
while (current != head);
240 void AdmitWriter(Wakes* wakes) {
241 if (writers_ ==
nullptr || queue_.size() >= capacity_) {
244 State* current = writers_;
246 State* next = current->next;
247 MutexLock selector_lock(¤t->selector->mu);
248 if (current->selector->picked_case_index == Selector::kNonePicked) {
249 T* item = current->GetCase()->template GetArgPtr<T>(0);
251 *current->GetCase()->template GetArgPtr<Transfer>(1);
252 queue_.push_back(TransferItem(item, strategy));
254 UnlinkIfQueued(&writers_, current);
255 wakes->push_back(current->selector);
259 }
while (current != writers_ && writers_ !=
nullptr);
262 bool HandleRead(State*
reader,
bool enqueue) {
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;
272 MutexLock selector_lock(&
reader->selector->mu);
274 *
out = std::move(queue_.front());
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);
291 UnlinkIfQueued(&writers_,
writer);
292 wakes.push_back(
writer->selector);
296 MutexLock selector_lock(&
reader->selector->mu);
297 if (
reader->selector->picked_case_index != Selector::kNonePicked) {
299 }
else if (closed_) {
303 }
else if (enqueue) {
312 bool HandleWrite(State*
writer,
bool enqueue) {
316 MutexLock lock(&mu_);
320 State*
reader = closed_ ? nullptr : FindMatch(readers_,
writer);
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);
330 UnlinkIfQueued(&readers_,
reader);
331 wakes.push_back(
reader->selector);
335 MutexLock selector_lock(&
writer->selector->mu);
336 if (
writer->selector->picked_case_index != Selector::kNonePicked) {
338 }
else if (queue_.size() < capacity_) {
339 T* item =
writer->GetCase()->template GetArgPtr<T>(0);
341 *
writer->GetCase()->template GetArgPtr<Transfer>(1);
342 queue_.push_back(TransferItem(item, strategy));
345 }
else if (enqueue) {
354 void Unregister(State** head, State* state) {
355 MutexLock lock(&mu_);
356 UnlinkIfQueued(head, state);
359 Case MakeWriteCase(T* item, Transfer strategy) {
360 Case result(&write_selectable_);
362 result.AddArg(strategy == kCopy ? &kCopy : &kMove);
369 MutexLock lock(&mu_);
370 if (closed_ || writers_ !=
nullptr) {
374 while (readers_ !=
nullptr) {
375 State* state = readers_;
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);
389 const std::size_t capacity_;
391 std::deque<T> queue_;
392 bool closed_ =
false;
393 State* readers_ =
nullptr;
394 State* writers_ =
nullptr;
396 ReadSelectable read_selectable_;
398 WriteSelectable write_selectable_;