Перейти к содержанию

11. Асинхронные потоки и Channel<T>

О главе

Цель: понять, чем yield return отличается от await, как устроен асинхронный итератор (IAsyncEnumerable<T>, await foreach) и почему он не аллоцирует на элемент. Разобрать по коду Channel<T>: режимы переполнения, обратное давление, завершение, синхронные продолжения. Узнать, что LINQ для IAsyncEnumerable теперь входит в BCL (.NET 10).

Лабораторная: start/ — асинхронный итератор и ограниченный канал с тремя TODO (finally, DropOldest, LINQ). final/ — девять опытов: обычный итератор, await foreach (отмена, break, WithCancellation, AsyncLocal), аллокации, LINQ из BCL, режимы FullMode, завершение канала, несколько читателей, AllowSynchronousContinuations, стоимость элемента канала. Код — в конце главы.

Визуализация: Обратное давление — Играйте производителем и потребителем при каждом FullMode.

Статус: ✅ проверено на Ubuntu 26.04 (2 ядра), рантайм 10.0.12. Листинги CoreLib — декомпиляция System.Private.CoreLib 10.0.12 (.\tools\disasm.ps1 -Assembly corelib -Type …); Channel<T> — из System.Threading.Channels.dll того же рантайма; листинги итераторов — ilspycmd с -ds AsyncAwait=false -ds YieldReturn=false на final. Перепроверено на Windows 11 (8 ядер): вывод совпадает, отличаются только времена и шум замеров аллокаций.

11.1. Обычный итератор: yield return

Начнём с итераторов без async: это тот же приём, что и машина состояний async, только возобновляет её не уведомление, а вызывающий.

final/Sources.cs: обычный итератор
1
2
3
4
5
6
7
public static IEnumerable<int> Numbers(int max)
{
    Console.WriteLine("    [итератор] тело началось");
    for (int i = 1; i <= max; i++)
        yield return i;
    Console.WriteLine("    [итератор] тело закончилось");
}

yield return x значит «отдать x вызывающему и замереть на этой строке до следующего запроса». Компилятор превращает такой метод в класс (sealed class, не структуру), который реализует IEnumerable<int> и IEnumerator<int> сразу:

Numbers после компиляции (ILSpy без восстановления yield, Release)
private sealed class <Numbers>d__0 : IEnumerable<int>, IEnumerable, IEnumerator<int>, IEnumerator, IDisposable
{
    private int <>1__state;
    private int <>2__current;
    private int <>l__initialThreadId;
    private int max;
    private int <i>5__2;

    public <Numbers>d__0(int <>1__state)
    {
        this.<>1__state = <>1__state;
        <>l__initialThreadId = Environment.CurrentManagedThreadId;
    }

    private bool MoveNext()
    {
        switch (<>1__state)
        {
        default:
            return false;
        case 0:
            <>1__state = -1;
            Console.WriteLine("    [итератор] тело началось");
            <i>5__2 = 1;
            break;
        case 1:
            <>1__state = -1;
            <i>5__2++;
            break;
        }
        if (<i>5__2 <= max)
        {
            <>2__current = <i>5__2;
            <>1__state = 1;
            return true;
        }
        Console.WriteLine("    [итератор] тело закончилось");
        return false;
    }

    IEnumerator<int> IEnumerable<int>.GetEnumerator()
    {
        <Numbers>d__0 obj;
        if (<>1__state == -2 && <>l__initialThreadId == Environment.CurrentManagedThreadId)
        {
            <>1__state = 0;
            obj = this;
        }
        else
        {
            obj = new <Numbers>d__0(0);
        }
        ...
    }
}
  • Строки 3–4 и 7: состояние, текущий элемент (<>2__current) и поднятая локальная i. Те же приёмы, что в главе 2: номер состояния, поля вместо локальных.
  • Строки 15–39: MoveNext — обычный метод, возвращающий bool. Каждый вызов проходит тело до следующего yield return: кладёт значение в <>2__current, запоминает состояние (1) и возвращает true. Когда тело кончилось, возвращает false. Пока никто не вызвал MoveNext, ни одна строка тела не выполняется: итератор ленив.
  • Строки 41–54: GetEnumerator. Если последовательность ещё не начата (состояние -2) и вызов сделан на том же потоке, в котором её создали, перечислитель — это сам объект последовательности (строка 47: obj = this). Для второго вызова (или с другого потока) создаётся новый экземпляр. Поэтому на один обход приходится один объект.
final/Program.cs
Console.WriteLine("== 1. Обычный итератор ==");
IEnumerable<int> numbers = Sources.Numbers(2);
Console.WriteLine("  метод вызван, тело ещё не началось");
foreach (int n in numbers)
    Console.WriteLine($"  получено {n}");
IEnumerable<int> seq = Sources.Numbers(1);
IEnumerator<int> e1 = seq.GetEnumerator();
IEnumerator<int> e2 = seq.GetEnumerator();
Console.WriteLine($"  первый GetEnumerator вернул саму последовательность: {ReferenceEquals(seq, e1)}, второй: {ReferenceEquals(seq, e2)}");
Console.WriteLine($"  тип: {seq.GetType().Name.Replace("<", "‹").Replace(">", "›")}");
1
2
3
4
5
6
7
8
== 1. Обычный итератор ==
  метод вызван, тело ещё не началось
    [итератор] тело началось
  получено 1
  получено 2
    [итератор] тело закончилось
  первый GetEnumerator вернул саму последовательность: True, второй: False
  тип: ‹Numbers›d__0
  • Строка 2 вывода: Numbers(2) вызван (строка 13 кода), а «тело началось» напечатано только в первом foreach (строки 3–6).
  • Строка 7: первый GetEnumerator вернул саму последовательность, второй — новый объект.
yield return (итератор) await (async-метод)
Что значит пауза отдать значение вызывающему ждать результат чего-то
Кто возобновляет вызывающий, следующим MoveNext(), синхронно, в том же потоке завершение того, чего ждали (в любом потоке)
Возвращает вызывающему элементы по одному Task, один раз в конце

Асинхронный итератор объединяет обе возможности: внутри можно и await, и yield return.

11.2. IAsyncEnumerable<T> и await foreach

final/Sources.cs: асинхронный итератор
public static async IAsyncEnumerable<int> CountAsync(int max, [EnumeratorCancellation] CancellationToken ct = default)
{
    try
    {
        for (int i = 1; i <= max; i++)
        {
            await Task.Delay(100, ct);
            yield return i;
        }
    }
    finally
    {
        Console.WriteLine("    [итератор] finally");
    }
}

Потребитель:

1
2
3
4
await foreach (int n in CountAsync(10, ct))
{
    // тело цикла
}

Компилятор разворачивает его так (смысл; в машине состояний потребителя это try/catch с запоминанием исключения, потому что в сгенерированном коде await в finally недоступен):

IAsyncEnumerator<int> e = CountAsync(10, ct).GetAsyncEnumerator();
try
{
    while (await e.MoveNextAsync())      // ValueTask<bool>
    {
        int n = e.Current;
        // тело цикла
    }
}
finally
{
    await e.DisposeAsync();              // и при break, и при исключении
}
  • Строка 4: MoveNextAsync() возвращает ValueTask<bool>: true — есть следующий элемент (он в Current), false — конец.
  • Строки 10–12: DisposeAsync вызывается всегда, в том числе при break и исключении в теле цикла. В итераторе это выполнит блок finally.
final/Program.cs
Console.WriteLine("== 2. await foreach ==");
using (var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(450)))
{
    try
    {
        await foreach (int n in Sources.CountAsync(10, cts.Token))
            Console.WriteLine($"  получено {n}");
    }
    catch (OperationCanceledException e) { Console.WriteLine($"  отменено по таймауту: {e.GetType().Name}"); }
}
Console.WriteLine("  --- break на втором элементе:");
await foreach (int n in Sources.CountAsync(10))
{
    Console.WriteLine($"  получено {n}");
    if (n == 2) break;                          // DisposeAsync -> finally в итераторе
}
Console.WriteLine("  --- WithCancellation:");
using (var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(250)))
{
    try
    {
        await foreach (int n in Sources.CountAsync(10).WithCancellation(cts.Token))
            Console.WriteLine($"  получено {n}");
    }
    catch (OperationCanceledException) { Console.WriteLine("  отменено токеном из WithCancellation"); }
}
var asyncSeq = Sources.CountAsync(1);
IAsyncEnumerator<int> a1 = asyncSeq.GetAsyncEnumerator();
IAsyncEnumerator<int> a2 = asyncSeq.GetAsyncEnumerator();
Console.WriteLine($"  первый GetAsyncEnumerator вернул саму последовательность: {ReferenceEquals(asyncSeq, a1)}, второй: {ReferenceEquals(asyncSeq, a2)}");
await a1.DisposeAsync();
await a2.DisposeAsync();

Console.WriteLine("  --- AsyncLocal, записанный внутри итератора:");
foreach (bool pause in new[] { false, true })
{
    Sources.Ambient.Value = 1;
    Console.WriteLine($"  итератор {(pause ? "с await Task.Yield()" : "без await")}:");
    await foreach (int n in Sources.WritesAmbientAsync(pause))
        Console.WriteLine($"  потребитель, элемент {n}: Ambient = {Sources.Ambient.Value}");
}
== 2. await foreach ==
  получено 1
  получено 2
  получено 3
  получено 4
    [итератор] finally
  отменено по таймауту: TaskCanceledException
  --- break на втором элементе:
  получено 1
  получено 2
    [итератор] finally
  --- WithCancellation:
  получено 1
  получено 2
    [итератор] finally
  отменено токеном из WithCancellation
  первый GetAsyncEnumerator вернул саму последовательность: True, второй: False
  --- AsyncLocal, записанный внутри итератора:
  итератор без await:
    [итератор] до записи Ambient = 1
    [итератор] перед yield Ambient = 99
  потребитель, элемент 0: Ambient = 1
    [итератор] перед yield Ambient = 1
  потребитель, элемент 1: Ambient = 1
  итератор с await Task.Yield():
    [итератор] до записи Ambient = 1
    [итератор] перед yield Ambient = 99
  потребитель, элемент 0: Ambient = 1
    [итератор] перед yield Ambient = 1
  потребитель, элемент 1: Ambient = 1
  • Строки 2–7 вывода (строки 25–33 кода): таймаут 450 мс прервал цикл на пятом шаге: 4 элемента, затем Task.Delay внутри итератора бросил TaskCanceledException (глава 9, §9.2). Блок finally итератора выполнился.
  • Строки 8–11 (строки 34–39): break на втором элементе. Потребитель вышел из цикла, DisposeAsync разбудил итератор в режиме «закрытия», и finally выполнился. Если бы DisposeAsync не вызвали, finally не выполнился бы никогда.
  • Строки 12–15 (строки 40–49): WithCancellation(token). Токен пришёл не параметром метода, а при GetAsyncEnumerator(token). Он подставился в параметр с атрибутом [EnumeratorCancellation] (в нашем итераторе это ct). Без этого атрибута токен из WithCancellation пропал бы: параметр не получил бы его (предупреждение CS8425).
  • Строка 16 (строки 50–55): первый GetAsyncEnumerator вернул саму последовательность, как у обычного итератора.

Предскажите и задание (TODO 1)

В start/Program.cs у строк стоят комментарии PREDICT: сколько элементов придёт до отмены и какого типа будет исключение. Запустите и сверьтесь с таблицей выше:

cd chapters
dotnet run -c Release --project 11-async-streams\start

Затем добавьте в CountAsync блок try/finally с печатью и проверьте, что finally выполняется и при отмене, и при досрочном break (добавьте второй цикл с break). А что будет, если не вызвать DisposeAsync руками, пользуясь GetAsyncEnumerator напрямую?

AsyncLocal в итераторе

Строки 17–27 вывода (опыт «AsyncLocal», строки 57–64 кода): итератор записал Ambient = 99, увидел 99 перед первым yield return, но потребитель увидел 1, а перед вторым yield return внутри итератора снова 1. Запись не пережила даже один шаг. Дело в MoveNextAsync:

Сгенерированный MoveNextAsync и DisposeAsync (ILSpy, Release; фрагмент <CountAsync>d__1)
ValueTask<bool> IAsyncEnumerator<int>.MoveNextAsync()
{
    if (<>1__state == -2)
    {
        return default;
    }
    <>v__promiseOfValueOrEnd.Reset();
    <CountAsync>d__1 stateMachine = this;
    <>t__builder.MoveNext(ref stateMachine);
    short version = <>v__promiseOfValueOrEnd.Version;
    if (<>v__promiseOfValueOrEnd.GetStatus(version) == ValueTaskSourceStatus.Succeeded)
    {
        return new ValueTask<bool>(<>v__promiseOfValueOrEnd.GetResult(version));
    }
    return new ValueTask<bool>(this, version);
}

ValueTask IAsyncDisposable.DisposeAsync()
{
    if (<>1__state >= -1)
    {
        throw new NotSupportedException();
    }
    if (<>1__state == -2)
    {
        return default;
    }
    <>w__disposeMode = true;
    <>v__promiseOfValueOrEnd.Reset();
    <CountAsync>d__1 stateMachine = this;
    <>t__builder.MoveNext(ref stateMachine);
    return new ValueTask(this, <>v__promiseOfValueOrEnd.Version);
}
  • Строка 9: <>t__builder.MoveNext(ref stateMachine) — это AsyncMethodBuilderCore.Start (листинг ниже). Тот самый Start, который в главе 7, §7.3 возвращает вызывающему его ExecutionContext после синхронной части. Каждый MoveNextAsync — как вызов отдельного async-метода: что итератор записал в AsyncLocal, потребитель не увидит, и на следующем шаге итератор тоже не увидит своей записи.
  • Строки 10–14: если итератор дошёл до yield return без приостановки, статус уже Succeeded, и возвращается готовый ValueTask<bool> с результатом, без объекта (§11.3).
  • Строка 15: иначе в ValueTask кладётся сам итератор как источник (IValueTaskSource<bool>) с версией.
  • Строки 20–23: DisposeAsync в середине MoveNextAsync (состояние >= -1, то есть шаг идёт) бросает NotSupportedException: параллельно итерировать один перечислитель нельзя.
  • Строки 28–32: DisposeAsync включает режим закрытия (<>w__disposeMode) и пускает машину ещё раз: она выходит из тела через finally.

Факт: AsyncLocal, записанный в асинхронном итераторе, не переживает шаг

То же правило, что для обычного async-метода (глава 7, §7.3), но заметнее: там запись терялась на выходе из метода, а здесь теряется между соседними элементами. Если нужно протащить значение через итератор, передавайте его параметром или кладите в объект-держатель, а не в AsyncLocal. Читать AsyncLocal в итераторе можно: значение потребителя на каждом шаге видно.

11.3. Что генерирует компилятор для async + yield return

CountAsync после компиляции: класс и главные поля (ILSpy, Release)
private sealed class <CountAsync>d__1 : IAsyncEnumerable<int>, IAsyncEnumerator<int>, IAsyncDisposable, IValueTaskSource<bool>, IValueTaskSource, IAsyncStateMachine
{
    public int <>1__state;
    public AsyncIteratorMethodBuilder <>t__builder;
    public ManualResetValueTaskSourceCore<bool> <>v__promiseOfValueOrEnd;
    private int <>2__current;
    private bool <>w__disposeMode;
    private CancellationTokenSource <>x__combinedTokens;
    private int <>l__initialThreadId;
    private CancellationToken ct;
    public CancellationToken <>3__ct;
    private int max;
    private int <i>5__2;
    private TaskAwaiter <>u__1;
    ...
}

Это один класс на все роли:

  • Строка 1: он и последовательность (IAsyncEnumerable<int>), и перечислитель (IAsyncEnumerator<int>), и машина состояний (IAsyncStateMachine), и источник для ValueTask (IValueTaskSource<bool>). Один объект на обход.
  • Строки 4–5: AsyncIteratorMethodBuilder (builder для итераторов) и ManualResetValueTaskSourceCore<bool> — готовая структура из CoreLib, которая хранит результат MoveNextAsync и продолжение потребителя. Именно из неё выросла идея «ValueTask без Task».
  • Строка 6: текущий элемент.
  • Строка 7: режим закрытия для DisposeAsync.
  • Строка 8: если и параметр метода, и GetAsyncEnumerator(token) дали токены, GetAsyncEnumerator создаёт связанный источник (CreateLinkedTokenSource, глава 9, §9.5), а по концу итерации освобождает его.
  • Строки 10–11: <>3__ct — значение параметра из вызова метода, ct — рабочая копия для текущего перечислителя.

Состояния сгенерированной машины:

<>1__state Значение
-2 создан, ещё не использован (можно отдать самого себя как перечислитель) или завершён
-3 перечислитель выдан, итерация не начата
-1 тело выполняется
0, 1, … приостановлен на await номер N (как в главе 2)
-4 приостановлен на yield return: ждёт следующего MoveNextAsync

А AsyncIteratorMethodBuilder устроен просто:

CoreLib 10.0.12: AsyncIteratorMethodBuilder
public struct AsyncIteratorMethodBuilder
{
    private Task<VoidTaskResult> m_task;

    public static AsyncIteratorMethodBuilder Create()
    {
        return default;
    }

    public void MoveNext<TStateMachine>(ref TStateMachine stateMachine) where TStateMachine : IAsyncStateMachine
    {
        AsyncMethodBuilderCore.Start(ref stateMachine);
    }

    public void AwaitUnsafeOnCompleted<TAwaiter, TStateMachine>(ref TAwaiter awaiter, ref TStateMachine stateMachine) where TAwaiter : ICriticalNotifyCompletion where TStateMachine : IAsyncStateMachine
    {
        AsyncTaskMethodBuilder<VoidTaskResult>.AwaitUnsafeOnCompleted(ref awaiter, ref stateMachine, ref m_task);
    }

    public void Complete()
    {
        if (m_task == null)
        {
            m_task = Task.s_cachedCompleted;
            return;
        }
        AsyncTaskMethodBuilder<VoidTaskResult>.SetExistingTaskResult(m_task, default);
        ...
    }
}
  • Строки 10–13: MoveNext — это AsyncMethodBuilderCore.Start (глава 7): отсюда терялись AsyncLocal в §11.2.
  • Строки 15–18: приостановка на await внутри итератора — тот же AsyncTaskMethodBuilder<VoidTaskResult>.AwaitUnsafeOnCompleted, что в главе 3: бокс (глава 1) создаётся на первой приостановке. Бокс — отдельный объект, но создаётся один раз на весь обход, а не на элемент.
  • Строки 20–29: конец итерации.

Сам yield return в MoveNext итератора выглядит так (фрагмент того же класса):

CountAsync.MoveNext: yield return и конец (фрагмент)
IL_00b0:
awaiter.GetResult();
<>2__current = <i>5__2;                     // 1. положить элемент в Current
num = (<>1__state = -4);                    // 2. запомнить: ждём следующего MoveNextAsync
goto IL_01a6;
...
IL_01a6:
<>v__promiseOfValueOrEnd.SetResult(result: true);   // 3. завершить «обещание» MoveNextAsync значением true

// конец тела:
<>1__state = -2;
<>2__current = 0;
<>t__builder.Complete();
<>v__promiseOfValueOrEnd.SetResult(result: false);  // MoveNextAsync вернёт false
  • Строки 3–4 и 8: yield return — это запись в Current, смена состояния на -4 и SetResult(true) на ManualResetValueTaskSourceCore. Тот, кто ждёт MoveNextAsync, получает true.
  • Строки 11–14: конец тела или исключение (SetException) завершают то же «обещание» значением false или ошибкой.

11.4. Сколько стоит элемент

final/Program.cs
Console.WriteLine();
Console.WriteLine("== 3. Аллокации на элемент: 100 000 элементов ==");
await Measure("синхронный проход, yield return без await", async () => { await foreach (int _ in Sources.SyncAsync(100_000)) { } });
await Measure("с паузой await Task.Yield() на каждом элементе", async () => { await foreach (int _ in Sources.YieldingAsync(100_000)) { } });
await Measure("то же, ConfigureAwait(false)", async () => { await foreach (int _ in Sources.YieldingAsync(100_000).ConfigureAwait(false)) { } });
1
2
3
4
== 3. Аллокации на элемент: 100 000 элементов ==
  синхронный проход, yield return без await       :    0.00 байт на элемент
  с паузой await Task.Yield() на каждом элементе  :    0.00 байт на элемент
  то же, ConfigureAwait(false)                    :   -0.02 байт на элемент
  • Строка 2 вывода: итератор без await (все 100 000 MoveNextAsync завершаются синхронно): ноль байт на элемент. Из §11.2: готовый ValueTask<bool> без объекта.
  • Строка 3: итератор, который приостанавливается на каждом элементе, тоже ноль. Источник ValueTask — сам итератор (ManualResetValueTaskSourceCore, Reset() перед каждым шагом), а продолжение потребителя хранится в нём. Бокс для await Task.Yield() создаётся один раз (§11.3).
  • Строка 4: отрицательное число — шум измерения. Общая аллокация на 100 000 элементов — десятки байт на весь обход.

Сравните с вариантом, где каждый шаг возвращал бы Task<bool>: по объекту на элемент. Так выглядела бы цена асинхронного перечисления без ValueTask (глава 12).

Следствия для кода:

  • IAsyncEnumerable<T> дёшев на элемент. Но каждый обход создаёт объекты (итератор, возможно связанный источник токенов). Для «горячих» коротких последовательностей это заметно.
  • ConfigureAwait(false) для цикла: await foreach (var x in seq.ConfigureAwait(false)) применяется к MoveNextAsync и DisposeAsync потребителя (глава 5). Внутри самого итератора он нужен на его await, как в любом async-методе.
  • Не смешивайте с буферизацией. IAsyncEnumerable отдаёт по одному элементу и не читает вперёд: это стриминг из БД и HTTP (ASP.NET Core умеет отдавать IAsyncEnumerable<T> как поток JSON).

11.5. LINQ для IAsyncEnumerable из BCL (.NET 10)

Раньше LINQ-операторы для IAsyncEnumerable<T> давал отдельный пакет System.Linq.Async. В .NET 10 они входят в платформу: сборка System.Linq.AsyncEnumerable лежит рядом с System.Linq в Microsoft.NETCore.App, и наша лабораторная использует Where, Select, ToListAsync без единого PackageReference.

final/Program.cs
Console.WriteLine("== 4. LINQ над IAsyncEnumerable (BCL, .NET 10) ==");
List<int> evens = await Sources.CountAsync(6).Where(n => n % 2 == 0).Select(n => n * 10).ToListAsync();
Console.WriteLine($"  чётные ×10: {string.Join(", ", evens)}");
int count = await Sources.SyncAsync(1000).CountAsync(n => n % 7 == 0);
Console.WriteLine($"  CountAsync(кратные 7 из 1000): {count}");
int firstBig = await Sources.SyncAsync(1000).FirstAsync(n => n > 500);
Console.WriteLine($"  FirstAsync(> 500): {firstBig}");
List<int> asyncPredicate = await Sources.SyncAsync(10).Where(async (n, ct) => { await Task.Yield(); return n % 3 == 0; }).ToListAsync();
Console.WriteLine($"  Where с async-предикатом: {string.Join(", ", asyncPredicate)}");
List<int> taken = await Sources.CountAsync(10).Take(3).ToListAsync();
Console.WriteLine($"  Take(3): {string.Join(", ", taken)} (источник остановлен: см. finally выше)");
using (var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(250)))
{
    try { await Sources.CountAsync(10).ToListAsync(cts.Token); }
    catch (OperationCanceledException) { Console.WriteLine("  ToListAsync(ct) отменён"); }
}
== 4. LINQ над IAsyncEnumerable (BCL, .NET 10) ==
    [итератор] finally
  чётные ×10: 20, 40, 60
  CountAsync(кратные 7 из 1000): 143
  FirstAsync(> 500): 501
  Where с async-предикатом: 0, 3, 6, 9
    [итератор] finally
  Take(3): 1, 2, 3 (источник остановлен: см. finally выше)
    [итератор] finally
  ToListAsync(ct) отменён
  • Методы: в сборке есть привычный набор (Where, Select, SelectMany, Take, Skip, Concat, Zip, OrderBy, GroupBy, Join, Distinct, Chunk, Reverse и другие) и «терминальные» операторы с суффиксом Async: ToListAsync, ToArrayAsync, ToDictionaryAsync, CountAsync, FirstAsync, AnyAsync, AggregateAsync, SumAsync, MinAsync/MaxAsync и т. д. Они возвращают ValueTask<T> и принимают CancellationToken.
  • Асинхронные предикаты: у большинства операторов есть перегрузка с Func<T, CancellationToken, ValueTask<bool>> (строка 80 кода: Where(async (n, ct) => …)), а не только с обычным Func<T, bool>.
  • Конец потребления освобождает источник: Take(3) остановился после третьего элемента, и в выводе [итератор] finally (строка 7) — DisposeAsync источника вызван. Так же ToListAsync(ct), отменённый токеном (строки 9–10): итератор закрылся.
  • Ленивость: Where и Select ничего не делают, пока не вызван терминальный оператор или await foreach.

Для перевода обычной последовательности в асинхронную: IEnumerable<T>.ToAsyncEnumerable().

Задание (TODO 3)

В start/Program.cs оставьте из CountAsync(10) только нечётные элементы и умножьте их на 10, одним выражением LINQ над IAsyncEnumerable<int> и с ToListAsync. Не подключайте никаких пакетов.

.NET 11

Сравнил открытые статические методы System.Linq.AsyncEnumerable отражением: в .NET 10.0.12 их 95, в .NET 11.0 RC1 (11.0.100-rc.1.26425.128) — 101. Разница только в соединениях: появился FullJoin (две перегрузки), а у Join, GroupJoin, LeftJoin, RightJoin — перегрузки без компаратора (в 10.0.12 только с компаратором, LeftJoin и RightJoin там уже есть). Остальной набор операторов тот же. Сравнение по числу параметров, без разбора типов; для релиза .NET 11 повторить. В этой книге все опыты сняты на рантайме 10.0.12.

11.6. Channel<T>: асинхронная очередь «производитель — потребитель»

Channel<T> — стандартный способ делать фоновую обработку в ASP.NET Core (очередь заданий + BackgroundService), замена BlockingCollection<T> в async-коде. Две стороны: Writer (TryWrite, WriteAsync, Complete) и Reader (TryRead, ReadAsync, ReadAllAsync, Completion). Какие каналы создаёт Channel:

System.Threading.Channels 10.0.12: фабрика Channel (сокращено)
public static Channel<T> CreateUnbounded<T>()
{
    return new UnboundedChannel<T>(runContinuationsAsynchronously: true);
}

public static Channel<T> CreateUnbounded<T>(UnboundedChannelOptions options)
{
    if (options.SingleReader)
    {
        return new SingleConsumerUnboundedChannel<T>(!options.AllowSynchronousContinuations);
    }
    return new UnboundedChannel<T>(!options.AllowSynchronousContinuations);
}

public static Channel<T> CreateBounded<T>(BoundedChannelOptions options, Action<T>? itemDropped)
{
    if (options.Capacity <= 0)
    {
        return new RendezvousChannel<T>(options.FullMode, !options.AllowSynchronousContinuations, itemDropped);
    }
    return new BoundedChannel<T>(options.Capacity, options.FullMode, !options.AllowSynchronousContinuations, itemDropped);
}

public static Channel<T> CreateUnboundedPrioritized<T>(UnboundedPrioritizedChannelOptions<T> options)
{
    return new UnboundedPrioritizedChannel<T>(!options.AllowSynchronousContinuations, options.Comparer);
}
  • Строки 1–4, 10–12, 19–21: все продолжения читателей по умолчанию асинхронные: runContinuationsAsynchronously = !AllowSynchronousContinuations, а AllowSynchronousContinuations по умолчанию false (§11.9). Это тот же принцип, что RunContinuationsAsynchronously у TaskCompletionSource (глава 5).
  • Строки 8–11: SingleReader = true выбирает другую реализацию, оптимизированную под одного читателя.
  • Строки 17–19: ёмкость 0 (и меньше) даёт RendezvousChannel: канал без буфера, запись завершается, только когда её принял читатель. В опыте 5в: capacity 0 (рандеву): RendezvousChannel1, TryWrite без читателя → False.
  • Строки 24–27: UnboundedPrioritizedChannel — канал с приоритетами (порядок по Comparer).

Типичный сценарий: производитель пишет, потребитель читает через await foreach:

var channel = Channel.CreateBounded<int>(new BoundedChannelOptions(capacity: 3)
{
    FullMode = BoundedChannelFullMode.Wait,    // производитель будет ждать при переполнении
    SingleReader = true,
});

var producer = Task.Run(async () =>
{
    for (int i = 1; i <= 8; i++)
        await channel.Writer.WriteAsync(i);    // ждёт, если очередь полна (обратное давление)
    channel.Writer.Complete();                 // сигнал «больше не будет»
});

var consumer = Task.Run(async () =>
{
    await foreach (var item in channel.Reader.ReadAllAsync())
        await Task.Delay(200);                 // медленный потребитель
});

await Task.WhenAll(producer, consumer);

ReadAllAsync — это обычный асинхронный итератор над теми же WaitToReadAsync и TryRead:

System.Threading.Channels 10.0.12: ChannelReader<T>.ReadAllAsync
public virtual async IAsyncEnumerable<T> ReadAllAsync([EnumeratorCancellation] CancellationToken cancellationToken = default)
{
    while (await WaitToReadAsync(cancellationToken).ConfigureAwait(continueOnCapturedContext: false))
    {
        T item;
        while (TryRead(out item))
        {
            yield return item;
        }
    }
}
  • Строка 3: ждать, пока в канале что-то появится (или он закроется). ConfigureAwait(false) внутри.
  • Строки 5–8: выбрать из канала всё, что есть, синхронно и по одному отдавать потребителю. Пока в канале есть элементы, MoveNextAsync завершается синхронно, без приостановки (§11.4).

Интерактивная схема

Играйте производителем и потребителем при каждом FullMode. Открыть на весь экран.

11.7. Обратное давление и режимы переполнения

System.Threading.Channels 10.0.12: BoundedChannel.TryWrite (сокращено)
public override bool TryWrite(T item)
{
    ...
    lock (parent.SyncObj)
    {
        if (parent._doneWriting != null)
        {
            return false;
        }
        int count = parent._items.Count;
        if (count != 0)
        {
            if (count < parent._bufferedCapacity)
            {
                parent._items.EnqueueTail(item);
                return true;
            }
            if (parent._mode == BoundedChannelFullMode.Wait)
            {
                return false;
            }
            if (parent._mode == BoundedChannelFullMode.DropWrite)
            {
                parent._itemDropped?.Invoke(item);
                return true;
            }
            T obj = ((parent._mode == BoundedChannelFullMode.DropNewest) ? parent._items.DequeueTail() : parent._items.DequeueHead());
            parent._items.EnqueueTail(item);
            parent._itemDropped?.Invoke(obj);
            return true;
        }
        ...                                       // канал пуст: отдать элемент ждущему читателю или положить в очередь
    }
}
  • Строки 6–9: после Complete() писать нельзя: false.
  • Строки 13–17: есть место — положить в хвост.
  • Строки 18–21: Wait — места нет, TryWrite возвращает false. WriteAsync в этом режиме вернёт ValueTask, который завершится, когда читатель освободит место (обратное давление: писатель замедляется до темпа читателя).
  • Строки 22–26: DropWrite — выбросить записываемый элемент и вернуть true (!). Вызывающий не узнает, что элемент потерян, кроме как через itemDropped.
  • Строки 27–30: DropNewest убирает последний в буфере (хвост), DropOldest — первый (голову), потом записывает новый.
final/Program.cs
Console.WriteLine("== 5б. Режимы переполнения: 6 записей в канал на 3 элемента без читателя ==");
foreach (var mode in new[] { BoundedChannelFullMode.DropOldest, BoundedChannelFullMode.DropNewest, BoundedChannelFullMode.DropWrite })
{
    var dropped = new List<int>();
    var ch = Channel.CreateBounded<int>(new BoundedChannelOptions(3) { FullMode = mode }, dropped.Add);
    var results = new List<bool>();
    for (int i = 1; i <= 6; i++) results.Add(ch.Writer.TryWrite(i));
    ch.Writer.Complete();
    var kept = new List<int>();
    await foreach (int item in ch.Reader.ReadAllAsync()) kept.Add(item);
    Console.WriteLine($"  {mode,-10}: TryWrite → {string.Join("", results.Select(r => r ? "T" : "F"))}, выброшено [{string.Join(", ", dropped)}], осталось [{string.Join(", ", kept)}]");
}
{
    var ch = Channel.CreateBounded<int>(new BoundedChannelOptions(3) { FullMode = BoundedChannelFullMode.Wait });
    for (int i = 1; i <= 3; i++) ch.Writer.TryWrite(i);
    Console.WriteLine($"  Wait      : четвёртый TryWrite → {ch.Writer.TryWrite(4)}; WriteAsync вернул ValueTask, завершённый: {ch.Writer.WriteAsync(4).IsCompleted}");
}
1
2
3
4
5
== 5б. Режимы переполнения: 6 записей в канал на 3 элемента без читателя ==
  DropOldest: TryWrite → TTTTTT, выброшено [1, 2, 3], осталось [4, 5, 6]
  DropNewest: TryWrite → TTTTTT, выброшено [3, 4, 5], осталось [1, 2, 6]
  DropWrite : TryWrite → TTTTTT, выброшено [4, 5, 6], осталось [1, 2, 3]
  Wait      : четвёртый TryWrite → False; WriteAsync вернул ValueTask, завершённый: False
Режим Что выбрасывается при переполнении Что остаётся из 1…6 при ёмкости 3 TryWrite
Wait ничего: писатель ждёт — (запись откладывается) false при полном буфере
DropOldest самый старый из буфера 4, 5, 6 (последние) true
DropNewest самый новый из буфера 1, 2, 6 true
DropWrite пишущийся элемент 1, 2, 3 (первые) true (!)

Факт: в режимах Drop* TryWrite всегда возвращает true

Даже если элемент тут же выброшен (DropWrite) или вытеснил другой. Единственный способ узнать о потере — колбэк itemDropped в Channel.CreateBounded(options, itemDropped). Колбэк вызывается вне блокировки канала, на потоке писателя: делайте его быстрым.

Задание (TODO 2)

В start/Program.cs замените FullMode = Wait на DropOldest. Что теперь прочитает потребитель? Запустите программу несколько раз: одинаков ли результат? Объясните разницу с помощью itemDropped.

Опыт: темп потребителя и Wait против DropOldest

Wait, 8 записей, читатель берёт по 200 мс, ёмкость 3:
     5 мс  записано 1
     7 мс  записано 2
     7 мс  записано 3
    10 мс  прочитано 1
    11 мс  записано 4
   211 мс  прочитано 2
   211 мс  записано 5
   411 мс  прочитано 3
   412 мс  записано 6
   612 мс  прочитано 4
   612 мс  записано 7
   812 мс  прочитано 5
   812 мс  записано 8
  1012 мс  прочитано 6 …

(final/Program.cs, функция RunChannel; полный вывод — запуском.) Видно обратное давление: писатель быстро записал 1, 2, 3, после чего его запись 4 стала возможной, только когда читатель взял 1; дальше каждая запись приходит сразу после чтения (в темпе читателя, 200 мс). Все 8 элементов дошли.

С DropOldest писатель не ждёт никого: все 8 записей за миллисекунду, пять элементов вытеснены (1…5), читатель получает 6, 7, 8.

Факт: какие элементы переживут DropOldest, зависит от гонки

В 10 запусках сценария DropOldest читатель получил 6, 7, 8 в шести и 2, 6, 7, 8 в четырёх: в этих запусках он успел взять элемент 2 раньше, чем писатель вытеснил его. Это нормально: режим гарантирует «в канале всегда последние N», а не «читатель увидит именно эти». Если потребитель должен видеть конкретные элементы, DropOldest не подходит.

Как работает ожидание в режиме Wait:

System.Threading.Channels 10.0.12: BoundedChannel.WriteAsync, ветка Wait (сокращено)
if (parent._mode == BoundedChannelFullMode.Wait)
{
    if (!cancellationToken.CanBeCanceled)
    {
        BlockedWriteAsyncOperation<T> writerSingleton = _writerSingleton;
        if (writerSingleton.TryOwnAndReset())
        {
            writerSingleton.Item = item;
            ChannelUtilities.Enqueue(ref parent._blockedWritersHead, writerSingleton);
            return writerSingleton.ValueTask;
        }
    }
    BlockedWriteAsyncOperation<T> blockedWriteAsyncOperation = new BlockedWriteAsyncOperation<T>(runContinuationsAsynchronously: true, cancellationToken, pooled: false, _parent.CancellationCallbackDelegate)
    {
        Item = item
    };
    ChannelUtilities.Enqueue(ref parent._blockedWritersHead, blockedWriteAsyncOperation);
    return blockedWriteAsyncOperation.ValueTask;
}
  • Строки 3–12: если токена отмены нет, используется один заранее созданный и повторно используемый объект _writerSingleton (TryOwnAndReset: «возьми, если свободен»). Поэтому ожидающая запись не аллоцирует.
  • Строки 13–16: если свободного синглтона нет или есть токен отмены, создаётся новый объект (pooled: false).
  • Строки 9 и 17: блокированный писатель — это очередь операций (_blockedWritersHead), а возвращаемое значение — ValueTask поверх объекта-источника. Читатель, освободив место, завершит операцию первого блокированного писателя.

11.8. Завершение канала и несколько читателей

final/Program.cs
Console.WriteLine();
Console.WriteLine("== 5в. Завершение канала ==");
var closing = Channel.CreateUnbounded<int>();
closing.Writer.TryWrite(1);
closing.Writer.TryWrite(2);
closing.Writer.Complete();
Console.WriteLine($"  после Complete: TryWrite → {closing.Writer.TryWrite(3)}; Completion завершена: {closing.Reader.Completion.IsCompleted} (остались непрочитанные)");
try { await closing.Writer.WriteAsync(3); }
catch (ChannelClosedException) { Console.WriteLine("  WriteAsync в закрытый канал: ChannelClosedException"); }
await foreach (int item in closing.Reader.ReadAllAsync()) Console.WriteLine($"  дочитано {item}");
Console.WriteLine($"  после дочитывания Completion: {closing.Reader.Completion.Status}");
var failing = Channel.CreateUnbounded<int>();
failing.Writer.TryWrite(1);
failing.Writer.Complete(new InvalidOperationException("производитель упал"));
try { await foreach (int item in failing.Reader.ReadAllAsync()) Console.WriteLine($"  получено {item}"); }
catch (InvalidOperationException e) { Console.WriteLine($"  Complete(exception): читатель получил «{e.Message}» после дочитывания"); }
var rendezvous = Channel.CreateBounded<int>(0);
Console.WriteLine($"  capacity 0 (рандеву): {rendezvous.GetType().Name}, TryWrite без читателя → {rendezvous.Writer.TryWrite(1)}");

Console.WriteLine();
Console.WriteLine("== 5г. Несколько читателей ==");
var jobs = Channel.CreateUnbounded<int>(new UnboundedChannelOptions { SingleReader = false });
var handled = new int[3];
Task[] workers = Enumerable.Range(0, 3).Select(id => Task.Run(async () =>
{
    await foreach (int job in jobs.Reader.ReadAllAsync())
    {
        await Task.Delay(10);
        Interlocked.Increment(ref handled[id]);
    }
})).ToArray();
for (int i = 0; i < 30; i++) jobs.Writer.TryWrite(i);
jobs.Writer.Complete();
await Task.WhenAll(workers);
Console.WriteLine($"  30 заданий обработано читателями: {string.Join(" + ", handled)} = {handled.Sum()}; Completion: {jobs.Reader.Completion.Status}");
== 5в. Завершение канала ==
  после Complete: TryWrite → False; Completion завершена: False (остались непрочитанные)
  WriteAsync в закрытый канал: ChannelClosedException
  дочитано 1
  дочитано 2
  после дочитывания Completion: RanToCompletion
  получено 1
  Complete(exception): читатель получил «производитель упал» после дочитывания
  capacity 0 (рандеву): RendezvousChannel`1, TryWrite без читателя → False

== 5г. Несколько читателей ==
  30 заданий обработано читателями: 10 + 10 + 10 = 30; Completion: RanToCompletion
  • Строка 2 вывода: Complete() закрывает канал для записи (TryWrite → false), но не для чтения. Reader.Completion ещё не завершена, пока не прочитаны оставшиеся элементы.
  • Строка 3: WriteAsync в закрытый канал бросает ChannelClosedException.
  • Строки 4–6: читатель дочитывает остаток, и только потом Completion переходит в RanToCompletion. Так делается корректное завершение: производитель вызывает Complete(), потребитель выходит из await foreach сам.
  • Строки 7–8: Complete(exception) — так писатель сообщает об ошибке: читатель получает уже накопленные элементы и затем исключение из ReadAllAsync / ReadAsync.
  • Строка 12: три читателя поделили 30 заданий поровну (10 + 10 + 10), и все вышли из await foreach, когда писатель вызвал Complete(). Один канал на нескольких читателей работает как очередь заданий: каждое задание получает ровно один читатель (для этого SingleReader оставляют false).

11.9. AllowSynchronousContinuations: чужой код внутри вашего TryWrite

Когда читатель ждёт на пустом канале, а писатель делает TryWrite, нужно «разбудить» читателя. Как это делается, решает AsyncOperation.SignalCompletion:

System.Threading.Channels 10.0.12: AsyncOperation.SignalCompletion (сокращено)
private protected void SignalCompletion()
{
    ...
    object capturedContext = _capturedContext;
    if ((capturedContext == null || capturedContext is ExecutionContext) ? true : false)
    {
        if (RunContinuationsAsynchronously)
        {
            UnsafeQueueSetCompletionAndInvokeContinuation();      // поставить в пул
            return;
        }
    }
    else
    {
        SynchronizationContext synchronizationContext = ...;
        if (synchronizationContext != null)
        {
            if (RunContinuationsAsynchronously || synchronizationContext != SynchronizationContext.Current)
            {
                synchronizationContext.Post(…);                    // вернуть в контекст читателя
                return;
            }
        }
        else
        {
            TaskScheduler taskScheduler = ...;
            if (RunContinuationsAsynchronously || taskScheduler != TaskScheduler.Current)
            {
                Task.Factory.StartNew(…, taskScheduler);           // в планировщик читателя
                return;
            }
        }
    }
    SetCompletionAndInvokeContinuation();                          // иначе — выполнить продолжение прямо здесь
}

Это та же логика выбора места продолжения, что в главе 5, только написанная заново для каналов:

  • Строки 5–11: читатель без контекста: продолжение ставится в пул, если асинхронные продолжения включены (по умолчанию).
  • Строки 15–23: читатель в SynchronizationContext — Post в него.
  • Строки 25–31: читатель в нестандартном TaskScheduler — задача в нём.
  • Строка 34: иначе (синхронные продолжения разрешены) — выполнить продолжение читателя прямо здесь, внутри вызова писателя.
final/Program.cs
Console.WriteLine();
Console.WriteLine("== 6. AllowSynchronousContinuations ==");
foreach (bool allow in new[] { false, true })
{
    var ch = Channel.CreateUnbounded<int>(new UnboundedChannelOptions { SingleReader = true, AllowSynchronousContinuations = allow });
    Flag.InsideWrite = false;
    bool ranInsideTryWrite = false;
    Task reader = Task.Run(async () =>
    {
        await ch.Reader.ReadAsync();                    // читатель ждёт на пустом канале
        ranInsideTryWrite = Flag.InsideWrite;
        Thread.Sleep(200);                              // «обработка» в продолжении читателя
    });
    await Task.Delay(100);
    long tryWriteMs = 0;
    await Task.Run(() =>
    {
        var sw = Stopwatch.StartNew();
        Flag.InsideWrite = true;
        ch.Writer.TryWrite(1);
        Flag.InsideWrite = false;
        tryWriteMs = sw.ElapsedMilliseconds;
    });
    await reader;
    Console.WriteLine($"  AllowSynchronousContinuations={allow,-5}: TryWrite занял {tryWriteMs,3} мс, продолжение читателя выполнилось внутри TryWrite: {ranInsideTryWrite}");
}
1
2
3
== 6. AllowSynchronousContinuations ==
  AllowSynchronousContinuations=False: TryWrite занял   0 мс, продолжение читателя выполнилось внутри TryWrite: False
  AllowSynchronousContinuations=True : TryWrite занял 200 мс, продолжение читателя выполнилось внутри TryWrite: True

Читатель после ReadAsync «обрабатывает» 200 мс (строка 160 кода). С AllowSynchronousContinuations = true его обработка выполнилась внутри TryWrite писателя: TryWrite занял те же 200 мс, и писатель стоял всё это время (строка 3 вывода). По умолчанию (false) TryWrite вернулся сразу (строка 2), а читатель обработал элемент в пуле.

Включать AllowSynchronousContinuations можно, только если продолжение читателя гарантированно короткое и не блокирует (оптимизация: нет постановки в пул), и писатель не держит блокировок, которые продолжение может захотеть. В остальных случаях оставляйте false. Это то же правило, что с TaskCreationOptions.RunContinuationsAsynchronously (глава 5, §5.6) и SemaphoreSlim (глава 10, §10.6).

11.10. Сколько стоит элемент канала

final/Program.cs
Console.WriteLine();
Console.WriteLine("== 7. Стоимость элемента канала: 200 000 элементов, один писатель и один читатель ==");
foreach (var (name, options) in new (string, BoundedChannelOptions)[]
{
    ("Bounded(64)", new BoundedChannelOptions(64)),
    ("Bounded(64), SingleReader + SingleWriter", new BoundedChannelOptions(64) { SingleReader = true, SingleWriter = true }),
})
{
    var (ms, bytes) = await Pump(options);
    Console.WriteLine($"  {name,-42}: {ms,4} мс, {bytes:F1} байт на элемент");
}
1
2
3
== 7. Стоимость элемента канала: 200 000 элементов, один писатель и один читатель ==
  Bounded(64)                               :   97 мс, 0.1 байт на элемент
  Bounded(64), SingleReader + SingleWriter  :   72 мс, 0.0 байт на элемент

Байты считаются по всему процессу (GC.GetTotalAllocatedBytes), то есть в обоих потоках. Около нуля байт на элемент: ожидающие операции канала — повторно используемые объекты (§11.7), ReadAllAsync — итератор без аллокаций на элемент (§11.4), значения int хранятся в Deque<T>. SingleReader и SingleWriter включают более лёгкие пути (меньше блокировок): в двух прогонах на 25–30 % быстрее (76 → 53 мс, 97 → 72 мс).

Перед выбором настроек ответьте на вопросы:

Вопрос Настройка
Нужна ли защита от переполнения памяти? Bounded (Unbounded растёт без предела)
Что делать писателю при переполнении? FullMode: Wait (давление), Drop* (потеря)
Один читатель или несколько? SingleReader
Один писатель или несколько? SingleWriter
Могут ли продолжения читателя быть синхронными? AllowSynchronousContinuations (по умолчанию нет)

11.11. Итоги

  • yield return и await — разные «паузы»: первая отдаёт значение вызывающему и возобновляется им, вторая ждёт уведомления. Асинхронный итератор — класс, объединяющий последовательность, перечислитель, машину состояний и источник для ValueTask<bool>.
  • await foreach всегда вызывает DisposeAsync: при break и исключениях выполняются finally итератора. Токен из WithCancellation попадает в параметр с [EnumeratorCancellation].
  • Каждый MoveNextAsync выполняется через AsyncMethodBuilderCore.Start: AsyncLocal, записанный внутри итератора, не переживает шаг.
  • Элемент асинхронного итератора не аллоцирует (ни синхронный, ни приостановленный): ValueTask поверх ManualResetValueTaskSourceCore.
  • В .NET 10 LINQ для IAsyncEnumerable входит в BCL (System.Linq.AsyncEnumerable), в том числе операторы с асинхронными предикатами.
  • Channel<T>: Bounded/Unbounded, SingleReader/SingleWriter; FullMode определяет судьбу переполнения (Wait — давление; Drop* — потеря, а TryWrite возвращает true); Complete() закрывает запись, читатели дочитывают остаток; Complete(exception) передаёт ошибку.
  • Продолжения читателей асинхронные по умолчанию. AllowSynchronousContinuations = true выполняет код читателя внутри TryWrite писателя.
  • Канал почти не аллоцирует на элемент.

Код лабораторной

Запуск из папки главы: dotnet run -c Release --project start или --project final.

start/Program.cs
// Глава 11, заготовка. Асинхронный итератор и ограниченный канал.
// PREDICT: сколько элементов успеет прийти до отмены? Как будут чередоваться «записано» и «прочитано»?
// TODO 1 (§11.2): добавьте в CountAsync блок try/finally с печатью и убедитесь, что finally выполняется при отмене и при break.
// TODO 2 (§11.5): замените FullMode = Wait на DropOldest. Что прочитает потребитель? Одинаков ли результат от запуска к запуску?
// TODO 3 (§11.4): отфильтруйте нечётные и умножьте на 10 с помощью LINQ над IAsyncEnumerable (без System.Linq.Async: он есть в BCL).
using System.Runtime.CompilerServices;
using System.Threading.Channels;

Console.WriteLine("== IAsyncEnumerable ==");
using var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(450));
try
{
    await foreach (var n in CountAsync(10, cts.Token))
        Console.WriteLine($"получено {n}");                  // PREDICT: до какого n?
}
catch (OperationCanceledException e)
{
    Console.WriteLine($"отменено по таймауту: {e.GetType().Name}");   // PREDICT: какой тип?
}

Console.WriteLine();
Console.WriteLine("== Channel ==");
var channel = Channel.CreateBounded<int>(new BoundedChannelOptions(capacity: 3)
{
    FullMode = BoundedChannelFullMode.Wait,    // производитель будет ждать при переполнении
    SingleReader = true,
});

var producer = Task.Run(async () =>
{
    for (int i = 1; i <= 8; i++)
    {
        await channel.Writer.WriteAsync(i);
        Console.WriteLine($"  записано {i}");
    }
    channel.Writer.Complete();
});

var consumer = Task.Run(async () =>
{
    await foreach (var item in channel.Reader.ReadAllAsync())
    {
        Console.WriteLine($"прочитано {item}");
        await Task.Delay(200);                 // медленный потребитель
    }
});

await Task.WhenAll(producer, consumer);

static async IAsyncEnumerable<int> CountAsync(int max, [EnumeratorCancellation] CancellationToken ct = default)
{
    for (int i = 1; i <= max; i++)
    {
        await Task.Delay(100, ct);
        yield return i;
    }
}
final/Program.cs
// Глава 11, итог. Итераторы, асинхронные потоки и каналы:
//   1) обычный итератор: ленивость, один объект на последовательность и перечислитель (§11.1);
//   2) await foreach: отмена, break и finally, WithCancellation, AsyncLocal в итераторе (§11.2);
//   3) сколько аллокаций стоит элемент асинхронного итератора (§11.4);
//   4) LINQ над IAsyncEnumerable из BCL (.NET 10) (§11.5);
//   5) Channel<T>: Wait, DropOldest, DropWrite, DropNewest, завершение, несколько читателей (§11.7, §11.8);
//   6) AllowSynchronousContinuations: продолжение читателя внутри TryWrite (§11.9);
//   7) стоимость элемента канала (§11.10).
using System.Diagnostics;
using System.Threading.Channels;

Console.WriteLine("== 1. Обычный итератор ==");
IEnumerable<int> numbers = Sources.Numbers(2);
Console.WriteLine("  метод вызван, тело ещё не началось");
foreach (int n in numbers)
    Console.WriteLine($"  получено {n}");
IEnumerable<int> seq = Sources.Numbers(1);
IEnumerator<int> e1 = seq.GetEnumerator();
IEnumerator<int> e2 = seq.GetEnumerator();
Console.WriteLine($"  первый GetEnumerator вернул саму последовательность: {ReferenceEquals(seq, e1)}, второй: {ReferenceEquals(seq, e2)}");
Console.WriteLine($"  тип: {seq.GetType().Name.Replace("<", "‹").Replace(">", "›")}");

Console.WriteLine();
Console.WriteLine("== 2. await foreach ==");
using (var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(450)))
{
    try
    {
        await foreach (int n in Sources.CountAsync(10, cts.Token))
            Console.WriteLine($"  получено {n}");
    }
    catch (OperationCanceledException e) { Console.WriteLine($"  отменено по таймауту: {e.GetType().Name}"); }
}
Console.WriteLine("  --- break на втором элементе:");
await foreach (int n in Sources.CountAsync(10))
{
    Console.WriteLine($"  получено {n}");
    if (n == 2) break;                          // DisposeAsync -> finally в итераторе
}
Console.WriteLine("  --- WithCancellation:");
using (var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(250)))
{
    try
    {
        await foreach (int n in Sources.CountAsync(10).WithCancellation(cts.Token))
            Console.WriteLine($"  получено {n}");
    }
    catch (OperationCanceledException) { Console.WriteLine("  отменено токеном из WithCancellation"); }
}
var asyncSeq = Sources.CountAsync(1);
IAsyncEnumerator<int> a1 = asyncSeq.GetAsyncEnumerator();
IAsyncEnumerator<int> a2 = asyncSeq.GetAsyncEnumerator();
Console.WriteLine($"  первый GetAsyncEnumerator вернул саму последовательность: {ReferenceEquals(asyncSeq, a1)}, второй: {ReferenceEquals(asyncSeq, a2)}");
await a1.DisposeAsync();
await a2.DisposeAsync();

Console.WriteLine("  --- AsyncLocal, записанный внутри итератора:");
foreach (bool pause in new[] { false, true })
{
    Sources.Ambient.Value = 1;
    Console.WriteLine($"  итератор {(pause ? "с await Task.Yield()" : "без await")}:");
    await foreach (int n in Sources.WritesAmbientAsync(pause))
        Console.WriteLine($"  потребитель, элемент {n}: Ambient = {Sources.Ambient.Value}");
}

Console.WriteLine();
Console.WriteLine("== 3. Аллокации на элемент: 100 000 элементов ==");
await Measure("синхронный проход, yield return без await", async () => { await foreach (int _ in Sources.SyncAsync(100_000)) { } });
await Measure("с паузой await Task.Yield() на каждом элементе", async () => { await foreach (int _ in Sources.YieldingAsync(100_000)) { } });
await Measure("то же, ConfigureAwait(false)", async () => { await foreach (int _ in Sources.YieldingAsync(100_000).ConfigureAwait(false)) { } });

Console.WriteLine();
Console.WriteLine("== 4. LINQ над IAsyncEnumerable (BCL, .NET 10) ==");
List<int> evens = await Sources.CountAsync(6).Where(n => n % 2 == 0).Select(n => n * 10).ToListAsync();
Console.WriteLine($"  чётные ×10: {string.Join(", ", evens)}");
int count = await Sources.SyncAsync(1000).CountAsync(n => n % 7 == 0);
Console.WriteLine($"  CountAsync(кратные 7 из 1000): {count}");
int firstBig = await Sources.SyncAsync(1000).FirstAsync(n => n > 500);
Console.WriteLine($"  FirstAsync(> 500): {firstBig}");
List<int> asyncPredicate = await Sources.SyncAsync(10).Where(async (n, ct) => { await Task.Yield(); return n % 3 == 0; }).ToListAsync();
Console.WriteLine($"  Where с async-предикатом: {string.Join(", ", asyncPredicate)}");
List<int> taken = await Sources.CountAsync(10).Take(3).ToListAsync();
Console.WriteLine($"  Take(3): {string.Join(", ", taken)} (источник остановлен: см. finally выше)");
using (var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(250)))
{
    try { await Sources.CountAsync(10).ToListAsync(cts.Token); }
    catch (OperationCanceledException) { Console.WriteLine("  ToListAsync(ct) отменён"); }
}

Console.WriteLine();
await RunChannel(BoundedChannelFullMode.Wait);
Console.WriteLine();
await RunChannel(BoundedChannelFullMode.DropOldest);
Console.WriteLine();
Console.WriteLine("== 5б. Режимы переполнения: 6 записей в канал на 3 элемента без читателя ==");
foreach (var mode in new[] { BoundedChannelFullMode.DropOldest, BoundedChannelFullMode.DropNewest, BoundedChannelFullMode.DropWrite })
{
    var dropped = new List<int>();
    var ch = Channel.CreateBounded<int>(new BoundedChannelOptions(3) { FullMode = mode }, dropped.Add);
    var results = new List<bool>();
    for (int i = 1; i <= 6; i++) results.Add(ch.Writer.TryWrite(i));
    ch.Writer.Complete();
    var kept = new List<int>();
    await foreach (int item in ch.Reader.ReadAllAsync()) kept.Add(item);
    Console.WriteLine($"  {mode,-10}: TryWrite → {string.Join("", results.Select(r => r ? "T" : "F"))}, выброшено [{string.Join(", ", dropped)}], осталось [{string.Join(", ", kept)}]");
}
{
    var ch = Channel.CreateBounded<int>(new BoundedChannelOptions(3) { FullMode = BoundedChannelFullMode.Wait });
    for (int i = 1; i <= 3; i++) ch.Writer.TryWrite(i);
    Console.WriteLine($"  Wait      : четвёртый TryWrite → {ch.Writer.TryWrite(4)}; WriteAsync вернул ValueTask, завершённый: {ch.Writer.WriteAsync(4).IsCompleted}");
}

Console.WriteLine();
Console.WriteLine("== 5в. Завершение канала ==");
var closing = Channel.CreateUnbounded<int>();
closing.Writer.TryWrite(1);
closing.Writer.TryWrite(2);
closing.Writer.Complete();
Console.WriteLine($"  после Complete: TryWrite → {closing.Writer.TryWrite(3)}; Completion завершена: {closing.Reader.Completion.IsCompleted} (остались непрочитанные)");
try { await closing.Writer.WriteAsync(3); }
catch (ChannelClosedException) { Console.WriteLine("  WriteAsync в закрытый канал: ChannelClosedException"); }
await foreach (int item in closing.Reader.ReadAllAsync()) Console.WriteLine($"  дочитано {item}");
Console.WriteLine($"  после дочитывания Completion: {closing.Reader.Completion.Status}");
var failing = Channel.CreateUnbounded<int>();
failing.Writer.TryWrite(1);
failing.Writer.Complete(new InvalidOperationException("производитель упал"));
try { await foreach (int item in failing.Reader.ReadAllAsync()) Console.WriteLine($"  получено {item}"); }
catch (InvalidOperationException e) { Console.WriteLine($"  Complete(exception): читатель получил «{e.Message}» после дочитывания"); }
var rendezvous = Channel.CreateBounded<int>(0);
Console.WriteLine($"  capacity 0 (рандеву): {rendezvous.GetType().Name}, TryWrite без читателя → {rendezvous.Writer.TryWrite(1)}");

Console.WriteLine();
Console.WriteLine("== 5г. Несколько читателей ==");
var jobs = Channel.CreateUnbounded<int>(new UnboundedChannelOptions { SingleReader = false });
var handled = new int[3];
Task[] workers = Enumerable.Range(0, 3).Select(id => Task.Run(async () =>
{
    await foreach (int job in jobs.Reader.ReadAllAsync())
    {
        await Task.Delay(10);
        Interlocked.Increment(ref handled[id]);
    }
})).ToArray();
for (int i = 0; i < 30; i++) jobs.Writer.TryWrite(i);
jobs.Writer.Complete();
await Task.WhenAll(workers);
Console.WriteLine($"  30 заданий обработано читателями: {string.Join(" + ", handled)} = {handled.Sum()}; Completion: {jobs.Reader.Completion.Status}");

Console.WriteLine();
Console.WriteLine("== 6. AllowSynchronousContinuations ==");
foreach (bool allow in new[] { false, true })
{
    var ch = Channel.CreateUnbounded<int>(new UnboundedChannelOptions { SingleReader = true, AllowSynchronousContinuations = allow });
    Flag.InsideWrite = false;
    bool ranInsideTryWrite = false;
    Task reader = Task.Run(async () =>
    {
        await ch.Reader.ReadAsync();                    // читатель ждёт на пустом канале
        ranInsideTryWrite = Flag.InsideWrite;
        Thread.Sleep(200);                              // «обработка» в продолжении читателя
    });
    await Task.Delay(100);
    long tryWriteMs = 0;
    await Task.Run(() =>
    {
        var sw = Stopwatch.StartNew();
        Flag.InsideWrite = true;
        ch.Writer.TryWrite(1);
        Flag.InsideWrite = false;
        tryWriteMs = sw.ElapsedMilliseconds;
    });
    await reader;
    Console.WriteLine($"  AllowSynchronousContinuations={allow,-5}: TryWrite занял {tryWriteMs,3} мс, продолжение читателя выполнилось внутри TryWrite: {ranInsideTryWrite}");
}

Console.WriteLine();
Console.WriteLine("== 7. Стоимость элемента канала: 200 000 элементов, один писатель и один читатель ==");
foreach (var (name, options) in new (string, BoundedChannelOptions)[]
{
    ("Bounded(64)", new BoundedChannelOptions(64)),
    ("Bounded(64), SingleReader + SingleWriter", new BoundedChannelOptions(64) { SingleReader = true, SingleWriter = true }),
})
{
    var (ms, bytes) = await Pump(options);
    Console.WriteLine($"  {name,-42}: {ms,4} мс, {bytes:F1} байт на элемент");
}

static async Task RunChannel(BoundedChannelFullMode mode)
{
    Console.WriteLine($"== 5а. Channel на 3 элемента, медленный читатель, FullMode={mode} ==");
    var sw = Stopwatch.StartNew();
    var channel = Channel.CreateBounded<int>(
        new BoundedChannelOptions(capacity: 3) { FullMode = mode, SingleReader = true },
        itemDropped: dropped => Console.WriteLine($"  {sw.ElapsedMilliseconds,4} мс  выброшено {dropped}"));

    var producer = Task.Run(async () =>
    {
        for (int i = 1; i <= 8; i++)
        {
            await channel.Writer.WriteAsync(i);       // при Wait ждёт место (обратное давление), при DropOldest не ждёт
            Console.WriteLine($"  {sw.ElapsedMilliseconds,4} мс  записано {i}");
        }
        channel.Writer.Complete();
    });

    var consumer = Task.Run(async () =>
    {
        await foreach (var item in channel.Reader.ReadAllAsync())
        {
            Console.WriteLine($"  {sw.ElapsedMilliseconds,4} мс  прочитано {item}");
            await Task.Delay(200);
        }
    });

    await Task.WhenAll(producer, consumer);
}

static async Task Measure(string title, Func<Task> action)
{
    await action();                                   // прогрев JIT
    long before = GC.GetAllocatedBytesForCurrentThread();
    await action();
    long bytes = GC.GetAllocatedBytesForCurrentThread() - before;
    Console.WriteLine($"  {title,-48}: {bytes / 100_000.0,7:F2} байт на элемент");
}

// Гоняет 200 000 элементов через канал; байты считаются по всему процессу (оба потока).
static async Task<(long Ms, double BytesPerItem)> Pump(BoundedChannelOptions options)
{
    const int N = 200_000;
    var ch = Channel.CreateBounded<int>(options);
    long before = GC.GetTotalAllocatedBytes(precise: true);
    var sw = Stopwatch.StartNew();
    Task producer = Task.Run(async () => { for (int i = 0; i < N; i++) await ch.Writer.WriteAsync(i); ch.Writer.Complete(); });
    Task consumer = Task.Run(async () => { await foreach (int _ in ch.Reader.ReadAllAsync()) { } });
    await Task.WhenAll(producer, consumer);
    return (sw.ElapsedMilliseconds, (GC.GetTotalAllocatedBytes(precise: true) - before) / (double)N);
}

static class Flag
{
    public static volatile bool InsideWrite;
}
final/Sources.cs
using System.Runtime.CompilerServices;

// Итераторы для опытов главы 11.
static class Sources
{
    // Обычный (синхронный) итератор: yield return без async.
    public static IEnumerable<int> Numbers(int max)
    {
        Console.WriteLine("    [итератор] тело началось");
        for (int i = 1; i <= max; i++)
            yield return i;
        Console.WriteLine("    [итератор] тело закончилось");
    }

    // Асинхронный итератор с настоящей паузой перед каждым элементом.
    public static async IAsyncEnumerable<int> CountAsync(int max, [EnumeratorCancellation] CancellationToken ct = default)
    {
        try
        {
            for (int i = 1; i <= max; i++)
            {
                await Task.Delay(100, ct);
                yield return i;
            }
        }
        finally
        {
            Console.WriteLine("    [итератор] finally");
        }
    }

    // Асинхронный итератор, который ни разу не приостанавливается: yield return без await.
    public static async IAsyncEnumerable<int> SyncAsync(int max)
    {
        for (int i = 0; i < max; i++)
            yield return i;
        await Task.CompletedTask;
    }

    // Асинхронный итератор, который приостанавливается на каждом элементе.
    public static async IAsyncEnumerable<int> YieldingAsync(int max)
    {
        for (int i = 0; i < max; i++)
        {
            await Task.Yield();
            yield return i;
        }
    }

    public static readonly AsyncLocal<int> Ambient = new();

    // Итератор записывает AsyncLocal и смотрит на него перед каждым yield return.
    public static async IAsyncEnumerable<int> WritesAmbientAsync(bool pause)
    {
        Console.WriteLine($"    [итератор] до записи Ambient = {Ambient.Value}");
        Ambient.Value = 99;
        for (int i = 0; i < 2; i++)
        {
            if (pause) await Task.Yield();
            Console.WriteLine($"    [итератор] перед yield Ambient = {Ambient.Value}");
            yield return i;
        }
    }
}