NOTE
2.12 channel
CSP concurrency model, channel usage and closing patterns, happens-before, hchan internals, creation, receive, send, close, and references.
This is a historical learning note and may contain outdated or incomplete understanding.
1. CSP Concurrency Model
- From the perspective of memory, there are only two kinds of parallel computation: shared memory and message communication.
- Shared-memory concurrency models usually provide mutexes as synchronization primitives.
- CSP stands for “Communicating Sequential Processes”. It is a message-communication-based concurrency model proposed by Tony Hoare in 1977.
- Golang implements the CSP concurrency model through the explicit
channelsynchronization primitive. At the same time, it also provides shared-memory synchronization primitives such assync.*andatomic.*.
2. What Is a channel?
Goroutines are used to execute concurrent tasks, while channels are used for synchronization and communication between goroutines.
2.1. channel and mutex
Do not communicate by sharing memory; instead, share memory by communicating. The first half refers to concurrent programming through components in the
syncpackage; the second half says that Go recommends using channels for concurrent programming. Essentially, the underlying implementation of a channel still uses a mutex to control concurrency. A channel is simply a higher-level concurrency primitive that encapsulates more functionality.
3. Using channels
3.1. Declaration
// chan T // declare a bidirectional channel
// chan<- T // declare a send-only channel
// <-chan T // declare a receive-only channel
func TestChannel4(t *testing.T) {
receiveChannel := make(<-chan int)
sendChannel := make(chan<- int)
channel := make(chan int)
fmt.Println(receiveChannel, sendChannel, channel)
}
// Output
0xc00004a2a0 0xc00004a300 0xc00004a360
// You can see that make returns a channel reference
3.2. Creation
func TestChannel1(t *testing.T) {
// Unbuffered
// Can be regarded as synchronous mode. The sender and receiver must be paired for the operation to succeed; otherwise it blocks.
ints := make(chan int)
// Buffered
// Can be regarded as asynchronous mode. The buffer must have remaining capacity for the operation to succeed; otherwise it also blocks.
ints2 := make(chan int, 10)
fmt.Println(ints, ints2)
}
// Output
0xc00004a2a0 0xc0000b6000
3.3. Send and Receive
func goroutineA(a <-chan int) {
val := <-a
fmt.Println("G1 received data: ", val)
return
}
func goroutineB(b <-chan int) {
val := <-b
fmt.Println("G2 received data: ", val)
return
}
func TestChannel2(t *testing.T) {
ch := make(chan int)
go goroutineA(ch)
go goroutineB(ch)
ch <- 3
time.Sleep(time.Second)
}
// Output
G2 received data: 3
3.4. Closing
func TestChannel3(t *testing.T) {
ch := make(chan int)
go goroutineC(ch)
time.Sleep(time.Second)
close(ch)
time.Sleep(time.Second)
}
func goroutineC(ch chan int) {
data, ok := <-ch
if !ok {
fmt.Println("channel closed, data:", data)
return
}
fmt.Println(data)
}
// Output
channel closed, data: 0
3.4.1. How to Gracefully Close a channel
Principle: do not close a channel from the receiver side, and do not close a channel when there are multiple senders. Depending on the number of senders and receivers, there are the following cases:
| sender | receiver | Handling | |
|---|---|---|---|
| Case 1 | 1 | 1 | Close from the sender side |
| Case 2 | 1 | N | Close from the sender side |
| Case 3 | N | 1 | Add a channel that carries the close signal. The receiver sends the instruction to close the data channel through the signal channel. After the senders observe the close signal, they stop sending data |
| Case 4 | N | M | Add a moderator. The M receivers all send it requests to close dataCh; after the moderator receives the first request, it directly issues the instruction to close dataCh |
- 1 : 1
func TestClose1(t *testing.T) {
rand.Seed(time.Now().UnixNano())
const Max = 100000
dataCh := make(chan int, 100)
stopCh := make(chan struct{})
// the sender
go func() {
for {
value := rand.Intn(Max)
if value == Max-1 {
fmt.Println("send stop signal to receiver.")
close(stopCh)
return
}
dataCh <- value
}
}()
// the receiver
go func() {
for {
select {
case value := <-dataCh:
fmt.Println(value)
case <-stopCh:
return
}
}
}()
select {
case <-time.After(time.Hour):
}
}
- 1 : N
func TestClose2(t *testing.T) {
rand.Seed(time.Now().UnixNano())
const Max = 100000
const NumReceivers = 10
dataCh := make(chan int, 100)
stopCh := make(chan struct{})
// the sender
go func() {
for {
value := rand.Intn(Max)
if value == Max-1 {
fmt.Println("send stop signal to receiver.")
close(stopCh)
return
}
dataCh <- value
}
}()
// the receivers
for i := 0; i < NumReceivers; i++ {
go func(id string) {
for {
select {
case value := <-dataCh:
fmt.Println(value)
case <-stopCh:
fmt.Println("receiver ", id , " return.")
return
}
}
}(strconv.Itoa(i))
}
select {
case <-time.After(time.Hour):
}
}
- N : 1
func TestClose3(t *testing.T) {
rand.Seed(time.Now().UnixNano())
const Max = 100000
const NumSenders = 1000
dataCh := make(chan int, 100)
stopCh := make(chan struct{})
// senders
for i := 0; i < NumSenders; i++ {
go func() {
for {
select {
case <-stopCh:
return
case dataCh <- rand.Intn(Max):
}
}
}()
}
// the receiver
go func() {
for value := range dataCh {
if value == Max-1 {
fmt.Println("send stop signal to senders.")
close(stopCh)
return
}
fmt.Println(value)
}
}()
select {
case <-time.After(time.Hour):
}
}
- N : M
func TestClose4(t *testing.T) {
rand.Seed(time.Now().UnixNano())
const Max = 100000
const NumReceivers = 10
const NumSenders = 1000
dataCh := make(chan int, 100)
stopCh := make(chan struct{})
// It must be a buffered channel.
toStop := make(chan string, 1)
var stoppedBy string
// moderator
go func() {
stoppedBy = <-toStop
fmt.Println(stoppedBy)
close(stopCh)
}()
// senders
for i := 0; i < NumSenders; i++ {
go func(id string) {
for {
value := rand.Intn(Max)
if value == 0 {
select {
case toStop <- "sender#" + id:
default:
}
return
}
select {
case <-stopCh:
return
case dataCh <- value:
}
}
}(strconv.Itoa(i))
}
// receivers
for i := 0; i < NumReceivers; i++ {
go func(id string) {
for {
select {
case <-stopCh:
return
case value := <-dataCh:
if value == Max-1 {
select {
case toStop <- "receiver#" + id:
default:
}
return
}
fmt.Println(value)
}
}
}(strconv.Itoa(i))
}
select {
case <-time.After(time.Hour):
}
}
3.5. Report When the Buffer Is Full or on a Timer
type Data struct {
topic int
}
func (d *Data) String() string {
return fmt.Sprintf("%v", d.topic)
}
const bufferSize = 2
func TestChannel3(t *testing.T) {
Queue := make(chan *Data, 1000)
go func() {
buffers := make([]*Data, 0, bufferSize)
for {
select {
case data := <-Queue:
buffers = append(buffers, data)
if len(buffers) == bufferSize {
fmt.Println("==========buffer is full============")
batchDealData(&buffers)
}
case <-time.After(time.Second * 5):
fmt.Println("=========periodic report=============")
batchDealData(&buffers)
}
}
}()
for i := 0; i < 100; i++ {
Queue <- &Data{topic: i}
}
time.Sleep(time.Hour)
}
func batchDealData(buffers *[]*Data) {
fmt.Println("process data:", len(*buffers), buffers)
*buffers = nil
}
4. happens before
5. Internals
5.1. Data Structure
You can see that a channel is implemented using a circular array + doubly linked lists + a lock.
type hchan struct {
// Number of elements in chan
qcount uint
// Length of the circular array underlying chan
dataqsiz uint
// Implemented using a circular array; pointer to that circular array
// Only for buffered channels
buf unsafe.Pointer
// Size of an element in chan
elemsize uint16
// Flag indicating whether chan is closed
closed uint32
// Element type in chan
elemtype *_type // element type
// Index of sent elements in the circular array
sendx uint // send index
// Index of received elements in the circular array
recvx uint // receive index
// Queue of goroutines waiting to receive, i.e. (<-chan)
recvq waitq // list of recv waiters
// Queue of goroutines waiting to send, i.e. (chan<-)
sendq waitq // list of send waiters
// Protect all fields in hchan
lock mutex
}
// Doubly linked list of sudogs
type waitq struct {
first *sudog // sudog is a wrapper used for channel waiting
last *sudog
}
The data structure for a channel with capacity 6 and int elements is as follows:

5.2. Creating a channel
func TestChannel12(t *testing.T) {
ints := make(chan int)
fmt.Println(ints)
}
go tool compile -S
0x0028 00040 (channel2_test.go:9) PCDATA $0, $1
0x0028 00040 (channel2_test.go:9) PCDATA $1, $0
0x0028 00040 (channel2_test.go:9) LEAQ type.chan int(SB), AX
0x002f 00047 (channel2_test.go:9) PCDATA $0, $0
0x002f 00047 (channel2_test.go:9) MOVQ AX, (SP)
0x0033 00051 (channel2_test.go:9) MOVQ $0, 8(SP)
0x003c 00060 (channel2_test.go:9) CALL runtime.makechan(SB)
0x0041 00065 (channel2_test.go:9) PCDATA $0, $1
0x0041 00065 (channel2_test.go:9) MOVQ 16(SP), AX
It calls runtime.makechan(SB).
const hchanSize = unsafe.Sizeof(hchan{}) + uintptr(-int(unsafe.Sizeof(hchan{}))&(maxAlign-1))
// Returns an hchan pointer
func makechan(t *chantype, size int64) *hchan {
elem := t.elem
// Code that checks channel size and alignment is omitted
// ……
var c *hchan
// If the element type contains no pointers or size is 0 (unbuffered)
// Only one memory allocation is performed
if elem.kind&kindNoPointers != 0 || size == 0 {
// If the hchan structure contains no pointers, GC will not scan the elements in chan
// Allocate only "hchan structure size + element size * count" bytes
c = (*hchan)(mallocgc(hchanSize+uintptr(size)*elem.size, nil, true))
// If it is a buffered channel and the element size is not 0 (an element type with size 0 is struct{})
if size > 0 && elem.size != 0 {
c.buf = add(unsafe.Pointer(c), hchanSize)
} else {
// race detector uses this location for synchronization
// Also prevents us from pointing beyond the allocation (see issue 9401).
// 1. For an unbuffered channel, buf is unused and points directly to the beginning of chan
// 2. For a buffered channel reaching this branch, the element has no pointers and the element type is struct{}, so it also has no effect
// Only the receive and send cursors are used; nothing is actually copied to c.buf (which would overwrite chan contents)
c.buf = unsafe.Pointer(c)
}
} else {
// Perform two memory allocations
c = new(hchan)
c.buf = newarray(elem, int(size))
}
c.elemsize = uint16(elem.size)
c.elemtype = elem
// Circular array length
c.dataqsiz = uint(size)
// Return the hchan pointer
return c
}

-
Creating a buffered channel

-
Creating an unbuffered channel

5.3. Receiving from a channel
func TestChannel13(t *testing.T) {
ints := make(chan int)
go func() {
data := <-ints
fmt.Println(data)
}()
go func() {
data, ok := <-ints
fmt.Println(ok, data)
}()
ints <- 10
time.Sleep(time.Hour)
}
go tool compile -S
0x0028 00040 (channel2_test.go:12) PCDATA $0, $0
0x0028 00040 (channel2_test.go:12) PCDATA $1, $0
0x0028 00040 (channel2_test.go:12) MOVQ $0, ""..autotmp_8+64(SP)
0x0031 00049 (channel2_test.go:12) PCDATA $0, $1
0x0031 00049 (channel2_test.go:12) PCDATA $1, $1
0x0031 00049 (channel2_test.go:12) MOVQ "".ints+104(SP), AX
0x0036 00054 (channel2_test.go:12) PCDATA $0, $0
0x0036 00054 (channel2_test.go:12) MOVQ AX, (SP)
0x003a 00058 (channel2_test.go:12) PCDATA $0, $1
0x003a 00058 (channel2_test.go:12) LEAQ ""..autotmp_8+64(SP), AX
0x003f 00063 (channel2_test.go:12) PCDATA $0, $0
0x003f 00063 (channel2_test.go:12) MOVQ AX, 8(SP)
0x0044 00068 (channel2_test.go:12) CALL runtime.chanrecv1(SB)
0x0049 00073 (channel2_test.go:12) MOVQ ""..autotmp_8+64(SP), AX
0x0028 00040 (channel2_test.go:16) PCDATA $0, $1
0x0028 00040 (channel2_test.go:16) PCDATA $1, $1
0x0028 00040 (channel2_test.go:16) MOVQ "".ints+128(SP), AX
0x0030 00048 (channel2_test.go:16) PCDATA $0, $0
0x0030 00048 (channel2_test.go:16) MOVQ AX, (SP)
0x0034 00052 (channel2_test.go:16) PCDATA $0, $1
0x0034 00052 (channel2_test.go:16) LEAQ ""..autotmp_10+72(SP), AX
0x0039 00057 (channel2_test.go:16) PCDATA $0, $0
0x0039 00057 (channel2_test.go:16) MOVQ AX, 8(SP)
0x003e 00062 (channel2_test.go:16) CALL runtime.chanrecv2(SB)
0x0043 00067 (channel2_test.go:16) MOVQ ""..autotmp_10+72(SP), AX
0x0048 00072 (channel2_test.go:16) MOVBLZX 16(SP), CX
0x004d 00077 (channel2_test.go:16) MOVQ CX, ""..autotmp_26+64(SP)
It calls runtime.chanrecv1(SB) and runtime.chanrecv2(SB).
// Handle the case without "ok"
func chanrecv1(c *hchan, elem unsafe.Pointer) {
chanrecv(c, elem, true)
}
// Handle the case with "ok". The returned "received" field indicates whether the channel is closed
func chanrecv2(c *hchan, elem unsafe.Pointer) (received bool) {
_, received = chanrecv(c, elem, true)
return
}
Eventually it calls chanrecv.
// Located in src/runtime/chan.go
// chanrecv receives an element from channel c and writes it to the memory address pointed to by ep.
// If ep is nil, the received value is ignored.
// If block == false, i.e. non-blocking receive, and no data can be received, return (false, false).
// Otherwise, if c is closed, clear the address pointed to by ep and return (true, false).
// Otherwise, fill the memory address pointed to by ep with the return value. Return (true, true).
// If ep is non-nil, it should point to the heap or the caller's stack.
func chanrecv(c *hchan, ep unsafe.Pointer, block bool) (selected, received bool) {
// Debug content omitted …………
// If this is a nil channel
if c == nil {
// If it should not block, directly return (false, false)
if !block {
return
}
// Otherwise, receiving from a nil channel parks the goroutine
gopark(nil, nil, "chan receive (nil chan)", traceEvGoStop, 2)
// Execution never reaches here
throw("unreachable")
}
// In non-blocking mode, quickly detect failure without acquiring the lock and return immediately
// When we observe that the channel is not ready to receive:
// 1. unbuffered: no goroutine is waiting in sendq
// 2. buffered: no elements in buf
// Then we observe closed == 0, meaning the channel is not closed.
// Because a channel cannot be reopened, it was also not closed at the time of the previous observation,
// so receive can be declared failed and return (false, false).
if !block && (c.dataqsiz == 0 && c.sendq.first == nil ||
c.dataqsiz > 0 && atomic.Loaduint(&c.qcount) == 0) &&
atomic.Load(&c.closed) == 0 {
return
}
var t0 int64
if blockprofilerate > 0 {
t0 = cputicks()
}
// Lock
lock(&c.lock)
// The channel is closed and the circular-array buf has no elements
// This handles an unbuffered closed channel and a buffered closed channel whose buf is empty
// That is, even when closed, a buffered channel can still receive elements if buf has elements
if c.closed != 0 && c.qcount == 0 {
if raceenabled {
raceacquire(unsafe.Pointer(c))
}
// Unlock
unlock(&c.lock)
if ep != nil {
// Receiving from a closed channel without ignoring the return value
// returns the zero value of the corresponding type
// typedmemclr clears memory at the corresponding address according to the type
typedmemclr(c.elemtype, ep)
}
// Receiving from a closed channel returns selected == true
return true, false
}
// A goroutine in the waiting-send queue means buf is full
// This may be:
// 1. an unbuffered channel
// 2. a buffered channel whose buf is full
// For 1, copy memory directly (sender goroutine -> receiver goroutine)
// For 2, receive the element at the head of the circular array and place the sender's element at the tail
if sg := c.sendq.dequeue(); sg != nil {
// Found a waiting sender. If buffer is size 0, receive value
// directly from sender. Otherwise, receive from head of queue
// and add sender's value to the tail of the queue (both map to
// the same buffer slot because the queue is full).
recv(c, sg, ep, func() { unlock(&c.lock) }, 3)
return true, true
}
// Buffered and buf contains elements, so receive normally
if c.qcount > 0 {
// Find the element to receive directly from the circular array
qp := chanbuf(c, c.recvx)
// …………
// The code does not ignore the received value. It is not "<- ch" but "val <- ch", and ep points to val
if ep != nil {
typedmemmove(c.elemtype, ep, qp)
}
// Clear the value at the corresponding position in the circular array
typedmemclr(c.elemtype, qp)
// Advance the receive cursor
c.recvx++
// Reset the receive cursor
if c.recvx == c.dataqsiz {
c.recvx = 0
}
// Decrease the number of elements in buf by 1
c.qcount--
// Unlock
unlock(&c.lock)
return true, true
}
if !block {
// Non-blocking receive: unlock. selected returns false because no value was received
unlock(&c.lock)
return false, false
}
// Next is the blocking case
// Construct a sudog
gp := getg()
mysg := acquireSudog()
mysg.releasetime = 0
if t0 != 0 {
mysg.releasetime = -1
}
// Save the address where data is to be received
mysg.elem = ep
mysg.waitlink = nil
gp.waiting = mysg
mysg.g = gp
mysg.selectdone = nil
mysg.c = c
gp.param = nil
// Enter the channel's receive-wait queue
c.recvq.enqueue(mysg)
// Park the current goroutine
goparkunlock(&c.lock, "chan receive", traceEvGoBlockRecv, 3)
// Woken up; continue here to do cleanup work
if mysg != gp.waiting {
throw("G waiting list is corrupted")
}
gp.waiting = nil
if mysg.releasetime > 0 {
blockevent(mysg.releasetime-t0, 2)
}
closed := gp.param == nil
gp.param = nil
mysg.c = nil
releaseSudog(mysg)
return true, !closed
}
- Receiving data from a buffered channel
- If the channel is empty

- Then a new sender arrives

- If the channel is empty
5.4. Sending to a channel
func TestChannel13(t *testing.T) {
ints := make(chan int)
go func() {
data := <-ints
fmt.Println(data)
}()
ints <- 10
time.Sleep(time.Hour)
}
go tool compile -S
0x0063 00099 (channel2_test.go:15) PCDATA $0, $1
0x0063 00099 (channel2_test.go:15) PCDATA $1, $0
0x0063 00099 (channel2_test.go:15) MOVQ "".ints+24(SP), AX
0x0068 00104 (channel2_test.go:15) PCDATA $0, $0
0x0068 00104 (channel2_test.go:15) MOVQ AX, (SP)
0x006c 00108 (channel2_test.go:15) PCDATA $0, $1
0x006c 00108 (channel2_test.go:15) LEAQ ""..stmp_0(SB), AX
0x0073 00115 (channel2_test.go:15) PCDATA $0, $0
0x0073 00115 (channel2_test.go:15) MOVQ AX, 8(SP)
0x0078 00120 (channel2_test.go:15) CALL runtime.chansend1(SB)
It calls runtime.chansend1(SB).
// Located in src/runtime/chan.go
func chansend(c *hchan, ep unsafe.Pointer, block bool, callerpc uintptr) bool {
// If channel is nil
if c == nil {
// If it cannot block, directly return false, indicating the send failed
if !block {
return false
}
// Park the current goroutine
gopark(nil, nil, "chan send (nil chan)", traceEvGoStop, 2)
throw("unreachable")
}
// Debug-related content omitted……
// For a non-blocking send, quickly detect failure
// If the channel is not closed and there is no extra buffer space. This may be:
// 1. an unbuffered channel with no goroutine in the receive-wait queue
// 2. a buffered channel whose circular array is full
if !block && c.closed == 0 && ((c.dataqsiz == 0 && c.recvq.first == nil) ||
(c.dataqsiz > 0 && c.qcount == c.dataqsiz)) {
return false
}
var t0 int64
if blockprofilerate > 0 {
t0 = cputicks()
}
// Lock the channel for concurrency safety
lock(&c.lock)
// If the channel is closed
if c.closed != 0 {
// Unlock
unlock(&c.lock)
// Panic directly
panic(plainError("send on closed channel"))
}
// If there is a goroutine in the receive queue, copy the data to be sent directly to the receiving goroutine
if sg := c.recvq.dequeue(); sg != nil {
send(c, sg, ep, func() { unlock(&c.lock) }, 3)
return true
}
// For a buffered channel, if buffer space remains
if c.qcount < c.dataqsiz {
// qp points to the sendx position of buf
qp := chanbuf(c, c.sendx)
// ……
// Copy data from ep to qp
typedmemmove(c.elemtype, qp, ep)
// Increment the send cursor
c.sendx++
// If the send cursor equals the capacity, reset it to 0
if c.sendx == c.dataqsiz {
c.sendx = 0
}
// Increment the number of elements in the buffer
c.qcount++
// Unlock
unlock(&c.lock)
return true
}
// If blocking is not required, return failure directly
if !block {
unlock(&c.lock)
return false
}
// The channel is full and the sender will block. Next, construct a sudog
// Get the pointer to the current goroutine
gp := getg()
mysg := acquireSudog()
mysg.releasetime = 0
if t0 != 0 {
mysg.releasetime = -1
}
mysg.elem = ep
mysg.waitlink = nil
mysg.g = gp
mysg.selectdone = nil
mysg.c = c
gp.waiting = mysg
gp.param = nil
// The current goroutine enters the send-wait queue
c.sendq.enqueue(mysg)
// Park the current goroutine
goparkunlock(&c.lock, "chan send", traceEvGoBlockSend, 3)
// Woken up from here (the channel now has an opportunity to send)
if mysg != gp.waiting {
throw("G waiting list is corrupted")
}
gp.waiting = nil
if gp.param == nil {
if c.closed == 0 {
throw("chansend: spurious wakeup")
}
// After wake-up, the channel is closed. Panic
panic(plainError("send on closed channel"))
}
gp.param = nil
if mysg.releasetime > 0 {
blockevent(mysg.releasetime-t0, 2)
}
// Remove the channel bound to mysg
mysg.c = nil
releaseSudog(mysg)
return true
}
- Sending data to a buffered channel
- If the channel is already full

- Then a receiver arrives

- If the channel is already full
- Sending data to an unbuffered channel
5.5. Closing a channel
func TestChannel13(t *testing.T) {
ints := make(chan int)
go func() {
data := <-ints
fmt.Println(data)
}()
time.Sleep(time.Second)
close(ints)
}
go tool compile -S
0x006c 00108 (channel2_test.go:16) PCDATA $0, $1
0x006c 00108 (channel2_test.go:16) PCDATA $1, $0
0x006c 00108 (channel2_test.go:16) MOVQ "".ints+24(SP), AX
0x0071 00113 (channel2_test.go:16) PCDATA $0, $0
0x0071 00113 (channel2_test.go:16) MOVQ AX, (SP)
0x0075 00117 (channel2_test.go:16) CALL runtime.closechan(SB)
It calls runtime.closechan(SB).
func closechan(c *hchan) {
// Closing a nil channel panics
if c == nil {
panic(plainError("close of nil channel"))
}
// Lock
lock(&c.lock)
// If the channel is already closed
if c.closed != 0 {
unlock(&c.lock)
// Panic
panic(plainError("close of closed channel"))
}
// …………
// Change the closed state
c.closed = 1
var glist *g
// Release all sudogs in the channel's receive-wait queue
for {
// Dequeue one sudog from the receive queue
sg := c.recvq.dequeue()
// Queue drained; break
if sg == nil {
break
}
// If elem is not nil, this receiver did not ignore the received data
// Assign the corresponding type's zero value
if sg.elem != nil {
typedmemclr(c.elemtype, sg.elem)
sg.elem = nil
}
if sg.releasetime != 0 {
sg.releasetime = cputicks()
}
// Get the goroutine
gp := sg.g
gp.param = nil
if raceenabled {
raceacquireg(gp, unsafe.Pointer(c))
}
// Link into a list
gp.schedlink.set(glist)
glist = gp
}
// Release sudogs in the channel's send-wait queue
// If present, these goroutines will panic
for {
// Dequeue one sudog from the send queue
sg := c.sendq.dequeue()
if sg == nil {
break
}
// The sender will panic
sg.elem = nil
if sg.releasetime != 0 {
sg.releasetime = cputicks()
}
gp := sg.g
gp.param = nil
if raceenabled {
raceacquireg(gp, unsafe.Pointer(c))
}
// Form a linked list
gp.schedlink.set(glist)
glist = gp
}
// Unlock
unlock(&c.lock)
// Ready all Gs now that we've dropped the channel lock.
// Traverse the linked list
for glist != nil {
// Take the last one
gp := glist
// Move one step forward to the next g to wake
glist = glist.schedlink.ptr()
gp.schedlink = 0
// Wake the corresponding goroutine
goready(gp, 3)
}
}
- Closing a channel
- Reading from a closed channel

- Reading from a closed channel
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub