Можно ли безопасно вызывать continuation AsyncStream из нескольких конкурентных задач без собственной блокировки?
Да, AsyncStream.Continuation предназначена для безопасных вызовов из разных конкурентных задач, поэтому дополнительная блокировка только для защиты самого вызова yield не требуется. Однако потокобезопасность не устанавливает порядок значений: при одновременных вызовах порядок выдачи может быть непредсказуемым, а политика буферизации может привести к отбрасыванию элементов.
AsyncStream появился как стандартный способ адаптировать callback- и событийные API к модели async/await. До этого разработчику приходилось самостоятельно синхронизировать доступ к буферу событий, отслеживать завершение и корректно обрабатывать отмену.
Continuation отделяет производителя событий от потребителя асинхронной последовательности. Благодаря этому callback может вызываться из произвольного потока или задачи, а потребитель получает значения через for await.
Несколько источников могут одновременно отправлять события в одну последовательность: например, сетевой callback, таймер и системное уведомление. Наивная реализация с общей коллекцией или собственным счетчиком может привести к гонкам данных.
У AsyncStream.Continuation внутренние операции публикации синхронизированы. Но это не делает атомарной внешнюю логику: проверка состояния, вычисление значения и последующий yield остаются отдельными действиями. Нельзя также рассчитывать на конкретный порядок конкурентных публикаций.
Вызов yield безопасен из разных задач, потому что continuation самостоятельно координирует доступ к внутреннему состоянию последовательности. Значение либо помещается в буфер и становится доступным потребителю, либо результат публикации сообщает, что оно было отброшено согласно политике буфера или последовательность уже завершена.
Завершение следует выполнять только после того, как все производители закончили публикацию. Если вызвать finish слишком рано, последующие значения не будут доставлены. Отмена потребителя также может завершить поток, поэтому производителям долгих операций нужно проверять отмену или реагировать на onTermination.
В примере две дочерние задачи конкурентно вызывают yield, а finish выполняется после завершения группы. Безопасность публикации обеспечивается continuation, но порядок 1, 2 не гарантируется.
Если требуется строгий порядок, его нужно организовать отдельно: например, присвоить событиям последовательные номера и упорядочить их до публикации либо отправлять их через один сериализованный владелец состояния. Actor может защитить такое состояние, но не отменяет необходимость определить правила упорядочивания.
Сервис получает события от нескольких callback-источников и передаёт их в UI через AsyncStream. Рассматривались три варианта: общий массив с блокировкой, отдельный actor и непосредственный вызов continuation.yield из каждого callback.
Общий массив с блокировкой даёт контроль над порядком, но усложняет ожидание новых событий и обработку отмены. Actor хорошо защищает дополнительное состояние, но добавляет асинхронные переходы и не нужен, если требуется только публикация событий.
Был выбран непосредственный yield из callback-ов. Это минимальное решение для безопасной публикации; при этом команда явно отказалась от гарантии порядка и обработала результат публикации, чтобы учитывать переполнение буфера и завершение потребителя.
Нет, конкурентные вызовы yield не дают прикладной гарантии порядка. Даже если одна задача логически начала публикацию раньше другой, планировщик и внутренняя синхронизация могут привести к другой последовательности выдачи. Если порядок существенен, его нужно моделировать явно, например через номера событий или единственного сериализованного производителя.
При выбранной политике ограниченного буфера новые элементы могут отбрасываться, когда потребитель не успевает их обрабатывать. Результат yield позволяет отличить успешную постановку от отбрасывания или публикации после завершения последовательности. Поэтому yield нельзя безусловно трактовать как гарантию доставки каждого события.
Нет. Она защищает внутреннее состояние AsyncStream, но не внешние переменные, коллекции, счётчики или объекты, к которым обращаются callback-и. Такое состояние нужно отдельно изолировать с помощью actor, блокировки или другой подходящей модели синхронизации. Иначе гонка данных может возникнуть до вызова yield или после него, несмотря на безопасность самой continuation.