NOTE

2.12 channel

CSP concurrency model, channel usage and closing patterns, happens-before, hchan internals, creation, receive, send, close, and references.

GoCreated Updated 2 min readhistorical

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 channel synchronization primitive. At the same time, it also provides shared-memory synchronization primitives such as sync.* and atomic.*.

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 sync package; 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

concurrent.md

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

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
  • 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

6. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub