Go under the hood
Go: Under the Hood

10.3 Отправка, получение и прямая передача

10.2 очертил скелет hchan: мьютекс, кольцевой буфер buf, а также очередь отправителей sendq и очередь получателей recvq. Настоящий раздел оживляет этот скелет и отвечает на вопрос, что именно происходит внутри рантайма при одиночном ch <- v и одиночном v := <-ch. Разобравшись в этом пути отправки/получения, оба наиболее часто задаваемых вопроса о каналах — почему небуферизованный канал является единственной точкой рандеву и почему для него получение происходит прежде, чем завершается соответствующая отправка (11.9), — сводятся к одному механизму: прямой отправке / прямому получению.

Проектирование отправки и получения должно удовлетворять трём ограничениям одновременно. Первое — корректность: данные не должны теряться, а закрытый канал не должен поглощать новые значения. Второе — скорость: при отсутствии конкуренции одиночная отправка или получение должны сводиться не более чем к «захватить блокировку, скопировать один раз, отпустить блокировку», а на горячем пути по возможности и вовсе обходить её. Третье — справедливость: когда несколько отправителей или получателей блокируются на одном канале, порядок их пробуждения должен быть предсказуемым (10.3.5). Ниже мы сначала рассмотрим отправку, затем симметрично набросаем получение и наконец придём к оптимизации, объединяющей оба пути.

10.3.1 Трёхсторонний выбор в chansend

Начнём с интерактивного рисунка для формирования интуитивного понимания: буферизованный канал представляет собой очередь фиксированной длины; отправка помещает значение в буфер, получение извлекает его; когда буфер заполнен, отправитель блокируется, когда буфер пуст — блокируется получатель. Вы можете менять cap, а также вручную выполнять отправку и получение.

Компилятор транслирует ch <- v в chansend1, которая переадресует вызов к более общей функции chansend. Третий аргумент chansend, block, отличает блокирующую отправку/получение от неблокирующей ветви внутри select (10.5). Если отвлечься от детектора гонок, «пузыря» synctest и кода сбора статистики, основная логика представляет собой чёткий трёхсторонний выбор:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
// chansend: отправляет значение, на которое указывает ep, в канал (сокращённая схема)
func chansend(c *hchan, ep unsafe.Pointer, block bool) bool {
    if c == nil {                  // отправка в nil-канал: блокироваться вечно
        if !block { return false }
        gopark(nil, nil, waitReasonChanSendNilChan, ...) // никогда не возвращается
        throw("unreachable")
    }

    // неблокирующий быстрый путь: отказ, определяемый без захвата блокировки (см. 10.3.4)
    if !block && c.closed == 0 && full(c) {
        return false
    }

    lock(&c.lock)

    if c.closed != 0 {             // отправка в закрытый канал: паника
        unlock(&c.lock)
        panic(plainError("send on closed channel"))
    }

    // ветвь первая: в recvq есть ожидающий получатель -> прямая передача, минуя buf
    if sg := c.recvq.dequeue(); sg != nil {
        send(c, sg, ep, func() { unlock(&c.lock) })
        return true
    }

    // ветвь вторая: в буфере ещё есть место -> скопировать значение в кольцевой buf
    if c.qcount < c.dataqsiz {
        qp := chanbuf(c, c.sendx)
        typedmemmove(c.elemtype, qp, ep) // скопировать в слот buf[sendx]
        c.sendx++
        if c.sendx == c.dataqsiz { c.sendx = 0 } // перенос на начало кольца
        c.qcount++
        unlock(&c.lock)
        return true
    }

    // ветвь третья: нет получателя, буфер тоже заполнен -> поместить себя в sendq и уступить CPU
    // ... см. 10.3.3
}

Приоритет среди трёх ветвей сам по себе является проектным решением: при наличии ожидающего получателя значение всегда передаётся ему напрямую, даже если это буферизованный канал с незаполненным буфером. Интуитивно может казаться, что следует «сначала заполнять буфер», однако пока recvq непуст, буфер в данный момент обязательно пуст (иначе получатель давно взял бы значение из буфера и не был бы заблокирован), поэтому обход буфера с прямой доставкой не только допустим, но и экономит одно копирование. Именно этот приоритет лежит в основе оптимизации, описанной в следующем разделе.

Вторая ветвь является типичным случаем для буферизованного канала: когда буфер не заполнен, отправка сводится к «скопировать значение в buf[sendx], продвинуть sendx». Совместно с recvx на стороне получения sendx использует массив buf фиксированной длины как кольцевую очередь: когда sendx == dataqsiz, он возвращается к 0 — это и есть весь смысл слова «кольцо».

10.3.2 Прямая передача: пропуск копирования «внутрь» и «наружу»

Функция send, вызываемая в первой ветви, — наиболее содержательная часть реализации канала, достойная пристального изучения. Её ядром является sendDirect:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
func send(c *hchan, sg *sudog, ep unsafe.Pointer, unlockf func()) {
    if sg.elem != nil {
        sendDirect(c.elemtype, sg, ep) // скопировать прямо в слот на стеке получателя
        sg.elem = nil
    }
    gp := sg.g
    unlockf()                          // сначала снять блокировку
    gp.param = unsafe.Pointer(sg)
    sg.success = true
    goready(gp, ...)                   // затем разбудить получателя (см. 9.4)
}

func sendDirect(t *_type, sg *sudog, src unsafe.Pointer) {
    // src находится на «моём» стеке, dst — слот на стеке другой горутины
    dst := sg.elem
    typeBitsBulkBarrier(t, uintptr(dst), uintptr(src), t.Size_) // барьер записи
    memmove(dst, src, t.Size_)         // один memmove — прямо между стеками
}

Получатель, заблокированный в recvq, ранее записал в свой sudog указание «положи значение по этому адресу» (sg.elem, указывающий на слот переменной-получателя на его стеке). Поэтому отправитель не обращается к buf: одним memmove он переносит данные непосредственно из собственного слота стека в слот стека получателя. Это и есть прямая передача.

Что она даёт? Сравним с «через буфер»: отправитель сначала копирует значение в buf (in), а после пробуждения получатель копирует его из buf в собственную переменную (out) — два копирования туда и обратно плюс занятие и освобождение слота буфера. Прямая передача сливает оба копирования в один межстековый memmove. Цена этого — необходимость выполнять его, удерживая блокировку канала, когда получатель находится в состоянии _Gwaiting (9.3) и ещё не запущен. Именно потому, что получатель в этот момент не выполняется, никакой пользовательский код не конкурирует с межстековой записью, и запись прямо «в чужой стек» безопасна. typeBitsBulkBarrier здесь обязателен: межстековая запись значения, содержащего указатели, должна быть видна сборщику мусора (13).

Есть и легко упускаемый порядок операций: send сначала вызывает unlockf() для снятия блокировки, и только после этого вызывает goready для пробуждения получателя. Данные уже скопированы до снятия блокировки, поэтому в момент пробуждения значение уже находится в переменной получателя; ему не нужно снова обращаться к каналу — можно сразу возвращаться. Обратите внимание: goready лишь помечает получателя как готового к выполнению и возвращает его в очередь на выполнение (9.4); переключение не происходит немедленно.

10.3.3 Блокирующий путь: постановка в очередь и делегирование завершения

Если нет ни ожидающего получателя, ни места в буфере (небуферизованный канал «всегда заполнен», см. full в следующем разделе), отправитель может только заблокироваться. Он оборачивает себя в sudog, помещает его в sendq, а затем вызывает gopark для уступки CPU:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
    // chansend, ветвь третья: блокировка
    gp := getg()
    mysg := acquireSudog()
    mysg.elem.set(ep)        // зафиксировать «значение для отправки находится по этому адресу»
    mysg.g = gp
    mysg.c.set(c)
    gp.waiting = mysg
    c.sendq.enqueue(mysg)    // поместить в очередь ожидания отправки
    gp.parkingOnChan.Store(true)
    // уступить CPU, состояние становится _Gwaiting; chanparkcommit снимает блокировку после парковки
    gopark(chanparkcommit, unsafe.Pointer(&c.lock), waitReasonChanSend, ...)

    // == после пробуждения каким-либо получателем, продолжение отсюда ==
    KeepAlive(ep)            // удерживать значение живым, пока получатель не скопирует его
    closed := !mysg.success  // если success == false при пробуждении, нас разбудил close
    gp.waiting = nil
    releaseSudog(mysg)
    if closed {
        panic(plainError("send on closed channel"))
    }
    return true

Обратите внимание на двойственность этого пути: заблокированный отправитель оставляет в mysg.elem указание «где находится значение», а будущий получатель, попав в свою первую ветвь, вызовет recv для извлечения значения из этого sudog (recvDirect), а затем goready для пробуждения отправителя. park и goready строго попарны: отправитель паркуется в sendq, поскольку буфер заполнен, и освобождается через goready получателя; получатель паркуется в recvq, поскольку буфер пуст, и освобождается через goready отправителя. Любая из сторон может прийти первой; та, что приходит второй, несёт ответственность за завершение всей транзакции и пробуждение пришедшей первой. Это в точности механизм park/ready из 9.4, применённый конкретно к каналам.

chanparkcommit, передаваемый в gopark, — ключевая деталь. Горутина не может парковаться, удерживая блокировку (иначе блокировка никогда не будет снята), но и снимать блокировку до парковки нельзя (иначе горутина может быть разбужена в окне после снятия блокировки, но до фактического перехода в _Gwaiting, что приведёт к повреждению состояния). Решение — отложить снятие блокировки до момента «парковка успешно выполнена», делегировав это chanparkcommit; это соглашение об обратном вызове unlockf функции gopark (9.4). Булево значение mysg.success — секретный сигнал между отправителем и тем, кто его разбудит: true при нормальном завершении получателем, false при пробуждении через close; отправитель использует его, чтобы решить — вернуться нормально или паниковать.

Объединяя блокирующий путь с двумя предшествующими ветвями, полная картина chansend выглядит следующим образом:

flowchart TD
    S["ch &lt;- v, i.e. chansend(c, ep, block)"] --> NIL{"c == nil?"}
    NIL -->|yes| PARK0["парковаться вечно (nil-канал)"]
    NIL -->|no| FAST{"!block и не закрыт и заполнен?"}
    FAST -->|yes| RF["return false (select пропускает)"]
    FAST -->|no| LOCK["lock(c.lock)"]
    LOCK --> CL{"c.closed?"}
    CL -->|yes| PANIC["паника: send on closed channel"]
    CL -->|no| RECVQ{"в recvq есть получатель?"}
    RECVQ -->|yes| DIRECT["send / sendDirect:<br/>прямая межстековая передача, goready получателю"]
    RECVQ -->|no| BUF{"qcount &lt; dataqsiz?<br/>в буфере есть место?"}
    BUF -->|yes| COPY["typedmemmove в buf[sendx], продвинуть sendx"]
    BUF -->|no| BLK{"block?"}
    BLK -->|no| RF2["return false"]
    BLK -->|yes| ENQ["поместить в sendq, gopark (_Gwaiting)<br/>ждать goready от получателя"]
    DIRECT --> OK["return true"]
    COPY --> OK

10.3.4 chanrecv и рандеву небуферизованного канала

Получение v := <-ch (компилируется в chanrecv1) и v, ok := <-ch (chanrecv2) оба перенаправляются в chanrecv, которая почти зеркально отражает трёхсторонний выбор chansend, добавляя лишь одну ветвь для случая «канал закрыт и данных нет — вернуть нулевое значение»:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
// chanrecv: получить одно значение из канала (сокращённая схема)
func chanrecv(c *hchan, ep unsafe.Pointer, block bool) (selected, received bool) {
    if c == nil { /* то же, что при отправке: nil-канал блокируется навсегда */ }

    // неблокирующий быстрый путь: канал не готов и не закрыт — сразу промах
    if !block && empty(c) {
        if atomic.Load(&c.closed) == 0 { return }
        if empty(c) { /* закрыт и пуст: вернуть нулевое значение */ }
    }

    lock(&c.lock)
    if c.closed != 0 && c.qcount == 0 {        // закрыт и данных нет
        unlock(&c.lock)
        if ep != nil { typedmemclr(c.elemtype, ep) } // записать нулевое значение
        return true, false                     // received == false
    }
    if sg := c.sendq.dequeue(); sg != nil {    // ветвь первая: есть ожидающий отправитель
        recv(c, sg, ep, func() { unlock(&c.lock) })
        return true, true
    }
    if c.qcount > 0 {                          // ветвь вторая: данные в буфере
        qp := chanbuf(c, c.recvx)
        typedmemmove(c.elemtype, ep, qp)       // скопировать из buf[recvx]
        typedmemclr(c.elemtype, qp)            // очистить слот для GC
        c.recvx++
        if c.recvx == c.dataqsiz { c.recvx = 0 }
        c.qcount--
        unlock(&c.lock)
        return true, true
    }
    // ветвь третья: данных для получения нет -> поместить в recvq и gopark (симметрично блокирующему пути отправки)
}

Функция recv на стороне получения завершает симметрию прямой передачи. Она должна различать случаи наличия и отсутствия буфера:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
20
func recv(c *hchan, sg *sudog, ep unsafe.Pointer, unlockf func()) {
    if c.dataqsiz == 0 {
        // небуферизованный канал: скопировать прямо со стека отправителя в стек получателя
        if ep != nil { recvDirect(c.elemtype, sg, ep) }
    } else {
        // буферизованный, буфер заполнен: голова очереди уходит получателю, значение отправителя занимает освободившийся хвост
        qp := chanbuf(c, c.recvx)
        if ep != nil { typedmemmove(c.elemtype, ep, qp) } // buf -> получатель
        typedmemmove(c.elemtype, qp, sg.elem.get())       // отправитель -> buf
        c.recvx++
        if c.recvx == c.dataqsiz { c.recvx = 0 }
        c.sendx = c.recvx
    }
    sg.elem.set(nil)
    gp := sg.g
    unlockf()
    gp.param = unsafe.Pointer(sg)
    sg.success = true
    goready(gp, ...)   // разбудить заблокированного отправителя
}

Небуферизованный канал (dataqsiz == 0) использует recvDirect — одиночный memmove из слота стека отправителя к получателю. Техника та же, что и sendDirect на стороне отправки, только в двух направлениях в зависимости от того, кто из двух участников пришёл первым и кто заблокировался в очереди. Буферизованный канал при «заполненном буфере и поставленном в очередь отправителе» выполняет изящный приём: голова очереди передаётся получателю, а освободившийся слот используется для принятия значения от отправителя, стоящего в очереди, — одна операция извлечения в паре с одной операцией вставки, и кольцевая очередь продвигается ровно на одну позицию без холостого вращения.

Суть небуферизованного канала теперь очевидна: его buf имеет ёмкость ноль, и обе функции full и empty деградируют до «есть ли кто-то в противоположной очереди». Поэтому успешная отправка/получение требуют присутствия обеих сторон одновременно: либо отправитель встречает ожидающего получателя (send/sendDirect), либо получатель встречает ожидающего отправителя (recv/recvDirect), и тот, кто приходит первым, всегда паркуется и ждёт. Это и есть рандеву: небуферизованный канал ничего не хранит; он завершает единственную межстековую передачу значения лишь в тот момент, когда встречаются обе стороны.

Это также напрямую объясняет правило модели памяти (11.9), поначалу кажущееся загадочным: для небуферизованного канала получение происходит прежде, чем завершается соответствующая отправка. Причина прямо в коде: когда получатель выполняет recvDirect, или когда отправитель выполняет send/sendDirect, и копирование значения, и goready происходят до того, как «сторона, пришедшая первой, будет разбужена и сможет вернуться из chansend / chanrecv». Иными словами, заблокированный отправитель может продолжить выполнение только после того, как получатель забрал значение и вызвал goready. Этот порядок в коде и является источником отношения happens-before в модели памяти.

10.3.5 Неблокирующий быстрый путь и тонкость упорядочивания памяти

Ветвь default в select и неблокирующее использование с ok входят в отправку/получение с block == false. Для них предусмотрен ранний выход без блокировки: на стороне отправки — !block && c.closed == 0 && full(c), на стороне получения — !block && empty(c). Обе функции, full и empty, читают лишь одно-два слова:

1
2
3
4
5
6
func full(c *hchan) bool {
    if c.dataqsiz == 0 {
        return c.recvq.first == nil   // небуферизованный: нет ожидающего получателя означает «заполнен»
    }
    return c.qcount == c.dataqsiz     // буферизованный: заполненный буфер означает «полон»
}

Этот быстрый путь содержит тонкость в части упорядочивания памяти (11.9), на которую явно указывают комментарии в исходном коде. На стороне отправки сначала читается c.closed, затем full(c) — сначала подтверждается, что канал не закрыт, а потом — что он не готов к приёму. Ключевой аргумент: закрытый канал никогда не может перейти обратно из «отправка невозможна» в «отправка возможна». Поэтому даже если канал окажется закрытым между двумя этими чтениями, между ними обязательно был момент, когда канал одновременно не был закрыт и не был готов к отправке, и рантайм трактует это как наблюдение канала в тот момент, сообщая «отправка не может быть выполнена». Именно это свойство монотонности является страховкой, и поэтому два обычных чтения, даже переупорядоченные процессором или компилятором, дают верный результат: атомарные операции здесь не нужны, что экономит затраты на горячем пути. Продвижение вперёд опирается не на эти два чтения, а на побочные эффекты, которые chanrecv и closechan производят при снятии блокировки, обновляя видимость c.closed и full для текущего потока. Такая техника — «заменить одну монотонную гарантию на атомарную операцию» — весьма характерна для lock-free быстрых путей (ср. обсуждение упорядочивания памяти в 11.9).

10.3.6 Справедливость FIFO и урок истории

Очереди отправки/получения sendq/recvq являются FIFO: заблокированные горутины добавляются в хвост и извлекаются из головы, поэтому несколько ожидающих пробуждаются в порядке прихода. Это не случайная деталь реализации. Команда Go однажды обсуждала в issue #11506, «должен ли канал гарантировать пробуждение в порядке FIFO», и вывод был таков: рантайм действительно поддерживает FIFO, хотя спецификация языка этого явно не обещает. Для программ, зависящих от порядка пробуждения, важно понимать эту границу: наблюдаемое поведение FIFO — это гарантия текущей реализации, но не гарантия спецификации.

Рассматривая обе стороны отправки и получения вместе, рантайм каналов предстаёт весьма компактным и симметричным решением: единственная блокировка обслуживает трёхсторонний выбор; на горячем пути выполняется либо прямая передача, либо запись в кольцевой буфер; лишь на холодном пути производится постановка в очередь с парковкой. Прямая передача использует предпосылку «получатель не выполняется» для пропуска копирования туда и обратно; небуферизованный канал есть не что иное, как вырожденный буфер с ёмкостью ноль, а семантика рандеву и отношение happens-before в модели памяти следуют непосредственно из этого. Цена за производительность никогда не бесплатна: быстрота прямой передачи опирается на то, что механизм шедулера park/ready (9.4) уже обеспечил «пусть пришедшая первой сторона ждёт, а пришедшая второй завершит транзакцию». 10.4 объясняет, как close будит всех ожидающих, а select из 10.5 расширяет эту одноканальную отправку/получение до более сложного уровня — «наблюдения за несколькими каналами одновременно».

Дополнительные материалы

  1. The Go Authors. runtime/chan.go (chansend, chanrecv, send, recv, sendDirect, recvDirect, full, empty), Go 1.26. https://github.com/golang/go/blob/master/src/runtime/chan.go
  2. The Go Authors. runtime/proc.go (gopark, goready, ready), Go 1.26. https://github.com/golang/go/blob/master/src/runtime/proc.go
  3. The Go Authors. The Go Memory Model (раздел о коммуникации через каналы), редакция от 6 июня 2022 г. https://go.dev/ref/mem
  4. Go issue #11506. runtime: make channel FIFO ordering explicit / guaranteed? https://github.com/golang/go/issues/11506
  5. C. A. R. Hoare. «Communicating Sequential Processes.» Communications of the ACM, 21(8), 1978. https://doi.org/10.1145/359576.359585
  6. Эта книга: 10.2 hchan: внутренняя структура канала, 10.4 Семантика close, 10.5 Реализация select, 9.4 Цикл шедулера, 11.9 Модель согласованности памяти.