Source file src/runtime/mcleanup.go

     1  // Copyright 2024 The Go Authors. All rights reserved.
     2  // Use of this source code is governed by a BSD-style
     3  // license that can be found in the LICENSE file.
     4  
     5  package runtime
     6  
     7  import (
     8  	"internal/abi"
     9  	"internal/cpu"
    10  	"internal/goarch"
    11  	"internal/runtime/atomic"
    12  	"internal/runtime/math"
    13  	"internal/runtime/sys"
    14  	"unsafe"
    15  )
    16  
    17  // AddCleanup attaches a cleanup function to ptr. Some time after ptr is no longer
    18  // reachable, the runtime will call cleanup(arg) in a separate goroutine.
    19  //
    20  // A typical use is that ptr is an object wrapping an underlying resource (e.g.,
    21  // a File object wrapping an OS file descriptor), arg is the underlying resource
    22  // (e.g., the OS file descriptor), and the cleanup function releases the underlying
    23  // resource (e.g., by calling the close system call).
    24  //
    25  // There are few constraints on ptr. In particular, multiple cleanups may be
    26  // attached to the same pointer, or to different pointers within the same
    27  // allocation.
    28  //
    29  // If ptr is reachable from cleanup or arg, ptr will never be collected
    30  // and the cleanup will never run. As a protection against simple cases of this,
    31  // AddCleanup panics if arg is equal to ptr.
    32  //
    33  // There is no specified order in which cleanups will run.
    34  // In particular, if several objects point to each other and all become
    35  // unreachable at the same time, their cleanups all become eligible to run
    36  // and can run in any order. This is true even if the objects form a cycle.
    37  //
    38  // Cleanups run concurrently with any user-created goroutines.
    39  // Cleanups may also run concurrently with one another (unlike finalizers).
    40  // If a cleanup function must run for a long time, it should create a new goroutine
    41  // to avoid blocking the execution of other cleanups.
    42  //
    43  // If ptr has both a cleanup and a finalizer, the cleanup will only run once
    44  // it has been finalized and becomes unreachable without an associated finalizer.
    45  //
    46  // The cleanup(arg) call is not always guaranteed to run; in particular it is not
    47  // guaranteed to run before program exit.
    48  //
    49  // Cleanups are not guaranteed to run if the size of T is zero bytes, because
    50  // it may share same address with other zero-size objects in memory. See
    51  // https://go.dev/ref/spec#Size_and_alignment_guarantees.
    52  //
    53  // It is not guaranteed that a cleanup will run for objects allocated
    54  // in initializers for package-level variables. Such objects may be
    55  // linker-allocated, not heap-allocated.
    56  //
    57  // Note that because cleanups may execute arbitrarily far into the future
    58  // after an object is no longer referenced, the runtime is allowed to perform
    59  // a space-saving optimization that batches objects together in a single
    60  // allocation slot. The cleanup for an unreferenced object in such an
    61  // allocation may never run if it always exists in the same batch as a
    62  // referenced object. Typically, this batching only happens for tiny
    63  // (on the order of 16 bytes or less) and pointer-free objects.
    64  //
    65  // A cleanup may run as soon as an object becomes unreachable.
    66  // In order to use cleanups correctly, the program must ensure that
    67  // the object is reachable until it is safe to run its cleanup.
    68  // Objects stored in global variables, or that can be found by tracing
    69  // pointers from a global variable, are reachable. A function argument or
    70  // receiver may become unreachable at the last point where the function
    71  // mentions it. To ensure a cleanup does not get called prematurely,
    72  // pass the object to the [KeepAlive] function after the last point
    73  // where the object must remain reachable.
    74  //
    75  //go:nocheckptr
    76  func AddCleanup[T, S any](ptr *T, cleanup func(S), arg S) Cleanup {
    77  	// This is marked nocheckptr because checkptr doesn't understand the
    78  	// pointer manipulation done when looking at closure pointers.
    79  	// Similar code in mbitmap.go works because the functions are
    80  	// go:nosplit, which implies go:nocheckptr (CL 202158).
    81  
    82  	// Explicitly force ptr and cleanup to escape to the heap.
    83  	ptr = abi.Escape(ptr)
    84  	cleanup = abi.Escape(cleanup)
    85  
    86  	// The pointer to the object must be valid.
    87  	if ptr == nil {
    88  		panic("runtime.AddCleanup: ptr is nil")
    89  	}
    90  	usptr := uintptr(unsafe.Pointer(ptr))
    91  
    92  	// Check that arg is not equal to ptr.
    93  	// Use the static type of arg, since S may itself be an interface.
    94  	argType := abi.TypeFor[S]()
    95  	kind := argType.Kind()
    96  	if kind == abi.Pointer || kind == abi.UnsafePointer {
    97  		if unsafe.Pointer(ptr) == *((*unsafe.Pointer)(unsafe.Pointer(&arg))) {
    98  			panic("runtime.AddCleanup: ptr is equal to arg, cleanup will never run")
    99  		}
   100  	}
   101  	if inUserArenaChunk(usptr) {
   102  		// Arena-allocated objects are not eligible for cleanup.
   103  		panic("runtime.AddCleanup: ptr is arena-allocated")
   104  	}
   105  	if debug.sbrk != 0 {
   106  		// debug.sbrk never frees memory, so no cleanup will ever run
   107  		// (and we don't have the data structures to record them).
   108  		// Return a noop cleanup.
   109  		return Cleanup{}
   110  	}
   111  
   112  	// Create new storage for the argument.
   113  	var argv *S
   114  	if size := unsafe.Sizeof(arg); size < maxTinySize && argType.PtrBytes == 0 {
   115  		// Side-step the tiny allocator to avoid liveness issues, since this box
   116  		// will be treated like a root by the GC. We model the box as an array of
   117  		// uintptrs to guarantee maximum allocator alignment.
   118  		//
   119  		// TODO(mknyszek): Consider just making space in cleanupFn for this. The
   120  		// unfortunate part of this is it would grow specialCleanup by 16 bytes, so
   121  		// while there wouldn't be an allocation, *every* cleanup would take the
   122  		// memory overhead hit.
   123  		box := new([maxTinySize / goarch.PtrSize]uintptr)
   124  		argv = (*S)(unsafe.Pointer(box))
   125  	} else {
   126  		argv = new(S)
   127  	}
   128  	*argv = arg
   129  
   130  	// Find the containing object.
   131  	base, span, _ := findObject(usptr, 0, 0)
   132  	if base == 0 {
   133  		if isGoPointerWithoutSpan(unsafe.Pointer(ptr)) {
   134  			// Cleanup is a noop.
   135  			return Cleanup{}
   136  		}
   137  		panic("runtime.AddCleanup: ptr not in allocated block")
   138  	}
   139  
   140  	// Check that arg is not within ptr.
   141  	if kind == abi.Pointer || kind == abi.UnsafePointer {
   142  		argPtr := uintptr(*(*unsafe.Pointer)(unsafe.Pointer(&arg)))
   143  		if argPtr >= base && argPtr < base+span.elemsize {
   144  			// It's possible that both pointers are separate
   145  			// parts of a tiny allocation, which is OK.
   146  			// We side-stepped the tiny allocator above for
   147  			// the allocation for the cleanup,
   148  			// but the argument itself can still overlap
   149  			// with the value to which the cleanup is attached.
   150  			if span.spanclass != tinySpanClass {
   151  				panic("runtime.AddCleanup: ptr is within arg, cleanup will never run")
   152  			}
   153  		}
   154  	}
   155  
   156  	// Check that the cleanup function doesn't close over the pointer.
   157  	cleanupFV := unsafe.Pointer(*(**funcval)(unsafe.Pointer(&cleanup)))
   158  	cBase, cSpan, _ := findObject(uintptr(cleanupFV), 0, 0)
   159  	if cBase != 0 {
   160  		tp := cSpan.typePointersOfUnchecked(cBase)
   161  		for {
   162  			var addr uintptr
   163  			if tp, addr = tp.next(cBase + cSpan.elemsize); addr == 0 {
   164  				break
   165  			}
   166  			ptr := *(*uintptr)(unsafe.Pointer(addr))
   167  			if ptr >= base && ptr < base+span.elemsize {
   168  				panic("runtime.AddCleanup: cleanup function closes over ptr, cleanup will never run")
   169  			}
   170  		}
   171  	}
   172  
   173  	// Create another G if necessary.
   174  	if gcCleanups.needG() {
   175  		gcCleanups.createGs()
   176  	}
   177  
   178  	id := addCleanup(unsafe.Pointer(ptr), cleanupFn{
   179  		// Instantiate a caller function to call the cleanup, that is cleanup(*argv).
   180  		//
   181  		// TODO(mknyszek): This allocates because the generic dictionary argument
   182  		// gets closed over, but callCleanup doesn't even use the dictionary argument,
   183  		// so theoretically that could be removed, eliminating an allocation.
   184  		call: callCleanup[S],
   185  		fn:   *(**funcval)(unsafe.Pointer(&cleanup)),
   186  		arg:  unsafe.Pointer(argv),
   187  	})
   188  	if debug.checkfinalizers != 0 {
   189  		cleanupFn := *(**funcval)(unsafe.Pointer(&cleanup))
   190  		setCleanupContext(unsafe.Pointer(ptr), abi.TypeFor[T](), sys.GetCallerPC(), cleanupFn.fn, id)
   191  	}
   192  	return Cleanup{
   193  		id:  id,
   194  		ptr: usptr,
   195  	}
   196  }
   197  
   198  // callCleanup is a helper for calling cleanups in a polymorphic way.
   199  //
   200  // In practice, all it does is call fn(*arg). arg must be a *T.
   201  //
   202  //go:noinline
   203  func callCleanup[T any](fn *funcval, arg unsafe.Pointer) {
   204  	cleanup := *(*func(T))(unsafe.Pointer(&fn))
   205  	cleanup(*(*T)(arg))
   206  }
   207  
   208  // Cleanup is a handle to a cleanup call for a specific object.
   209  type Cleanup struct {
   210  	// id is the unique identifier for the cleanup within the arena.
   211  	id uint64
   212  	// ptr contains the pointer to the object.
   213  	ptr uintptr
   214  }
   215  
   216  // Stop cancels the cleanup call. Stop will have no effect if the cleanup call
   217  // has already been queued for execution (because ptr became unreachable).
   218  // To guarantee that Stop removes the cleanup function, the caller must ensure
   219  // that the pointer that was passed to AddCleanup is reachable across the call to Stop.
   220  func (c Cleanup) Stop() {
   221  	if c.id == 0 {
   222  		// id is set to zero when the cleanup is a noop.
   223  		return
   224  	}
   225  
   226  	// The following block removes the Special record of type cleanup for the object c.ptr.
   227  	span := spanOfHeap(c.ptr)
   228  	if span == nil {
   229  		return
   230  	}
   231  	// Ensure that the span is swept.
   232  	// Sweeping accesses the specials list w/o locks, so we have
   233  	// to synchronize with it. And it's just much safer.
   234  	mp := acquirem()
   235  	span.ensureSwept()
   236  
   237  	offset := c.ptr - span.base()
   238  
   239  	var found *special
   240  	lock(&span.speciallock)
   241  
   242  	iter, exists := span.specialFindSplicePoint(offset, _KindSpecialCleanup)
   243  	if exists {
   244  		for {
   245  			s := *iter
   246  			if s == nil {
   247  				// Reached the end of the linked list. Stop searching at this point.
   248  				break
   249  			}
   250  			if offset == s.offset && _KindSpecialCleanup == s.kind &&
   251  				(*specialCleanup)(unsafe.Pointer(s)).id == c.id {
   252  				// The special is a cleanup and contains a matching cleanup id.
   253  				*iter = s.next
   254  				found = s
   255  				break
   256  			}
   257  			if offset < s.offset || (offset == s.offset && _KindSpecialCleanup < s.kind) {
   258  				// The special is outside the region specified for that kind of
   259  				// special. The specials are sorted by kind.
   260  				break
   261  			}
   262  			// Try the next special.
   263  			iter = &s.next
   264  		}
   265  	}
   266  	if span.specials == nil {
   267  		spanHasNoSpecials(span)
   268  	}
   269  	unlock(&span.speciallock)
   270  	releasem(mp)
   271  
   272  	if found == nil {
   273  		return
   274  	}
   275  	lock(&mheap_.speciallock)
   276  	mheap_.specialCleanupAlloc.free(unsafe.Pointer(found))
   277  	unlock(&mheap_.speciallock)
   278  
   279  	if debug.checkfinalizers != 0 {
   280  		clearCleanupContext(c.ptr, c.id)
   281  	}
   282  }
   283  
   284  const cleanupBlockSize = 512
   285  
   286  // cleanupBlock is an block of cleanups to be executed.
   287  //
   288  // cleanupBlock is allocated from non-GC'd memory, so any heap pointers
   289  // must be specially handled. The GC and cleanup queue currently assume
   290  // that the cleanup queue does not grow during marking (but it can shrink).
   291  type cleanupBlock struct {
   292  	cleanupBlockHeader
   293  	cleanups [(cleanupBlockSize - unsafe.Sizeof(cleanupBlockHeader{})) / unsafe.Sizeof(cleanupFn{})]cleanupFn
   294  }
   295  
   296  var cleanupFnPtrMask = [...]uint8{0b111}
   297  
   298  // cleanupFn represents a cleanup function with it's argument, yet to be called.
   299  type cleanupFn struct {
   300  	// call is an adapter function that understands how to safely call fn(*arg).
   301  	call func(*funcval, unsafe.Pointer)
   302  	fn   *funcval       // cleanup function passed to AddCleanup.
   303  	arg  unsafe.Pointer // pointer to argument to pass to cleanup function.
   304  }
   305  
   306  var cleanupBlockPtrMask [cleanupBlockSize / goarch.PtrSize / 8]byte
   307  
   308  type cleanupBlockHeader struct {
   309  	_ sys.NotInHeap
   310  	lfnode
   311  	alllink *cleanupBlock
   312  
   313  	// n is sometimes accessed atomically.
   314  	//
   315  	// The invariant depends on what phase the garbage collector is in.
   316  	// During the sweep phase (gcphase == _GCoff), each block has exactly
   317  	// one owner, so it's always safe to update this without atomics.
   318  	// But if this *could* be updated during the mark phase, it must be
   319  	// updated atomically to synchronize with the garbage collector
   320  	// scanning the block as a root.
   321  	n uint32
   322  }
   323  
   324  // enqueue pushes a single cleanup function into the block.
   325  //
   326  // Returns if this enqueue call filled the block. This is odd,
   327  // but we want to flush full blocks eagerly to get cleanups
   328  // running as soon as possible.
   329  //
   330  // Must only be called if the GC is in the sweep phase (gcphase == _GCoff),
   331  // because it does not synchronize with the garbage collector.
   332  func (b *cleanupBlock) enqueue(c cleanupFn) bool {
   333  	b.cleanups[b.n] = c
   334  	b.n++
   335  	return b.full()
   336  }
   337  
   338  // full returns true if the cleanup block is full.
   339  func (b *cleanupBlock) full() bool {
   340  	return b.n == uint32(len(b.cleanups))
   341  }
   342  
   343  // empty returns true if the cleanup block is empty.
   344  func (b *cleanupBlock) empty() bool {
   345  	return b.n == 0
   346  }
   347  
   348  // take moves as many cleanups as possible from b into a.
   349  func (a *cleanupBlock) take(b *cleanupBlock) {
   350  	dst := a.cleanups[a.n:]
   351  	if uint32(len(dst)) >= b.n {
   352  		// Take all.
   353  		copy(dst, b.cleanups[:])
   354  		a.n += b.n
   355  		b.n = 0
   356  	} else {
   357  		// Partial take. Copy from the tail to avoid having
   358  		// to move more memory around.
   359  		copy(dst, b.cleanups[b.n-uint32(len(dst)):b.n])
   360  		a.n = uint32(len(a.cleanups))
   361  		b.n -= uint32(len(dst))
   362  	}
   363  }
   364  
   365  // cleanupQueue is a queue of ready-to-run cleanup functions.
   366  type cleanupQueue struct {
   367  	// Stack of full cleanup blocks.
   368  	full      lfstack
   369  	workUnits atomic.Uint64 // length of full; decrement before pop from full, increment after push to full
   370  	_         [cpu.CacheLinePadSize - unsafe.Sizeof(lfstack(0)) - unsafe.Sizeof(atomic.Uint64{})]byte
   371  
   372  	// Stack of free cleanup blocks.
   373  	free lfstack
   374  
   375  	// flushed indicates whether all local cleanupBlocks have been
   376  	// flushed, and we're in a period of time where this condition is
   377  	// stable (after the last sweeper, before the next sweep phase
   378  	// begins).
   379  	flushed atomic.Bool // Next to free because frequently accessed together.
   380  
   381  	_ [cpu.CacheLinePadSize - unsafe.Sizeof(lfstack(0)) - 1]byte
   382  
   383  	// Linked list of all cleanup blocks.
   384  	all atomic.UnsafePointer // *cleanupBlock
   385  	_   [cpu.CacheLinePadSize - unsafe.Sizeof(atomic.UnsafePointer{})]byte
   386  
   387  	// Goroutine block state.
   388  	lock mutex
   389  
   390  	// sleeping is the list of sleeping cleanup goroutines.
   391  	//
   392  	// Protected by lock.
   393  	sleeping gList
   394  
   395  	// asleep is the number of cleanup goroutines sleeping.
   396  	//
   397  	// Read without lock, written only with the lock held.
   398  	// When the lock is held, the lock holder may only observe
   399  	// asleep.Load() == sleeping.n.
   400  	//
   401  	// To make reading without the lock safe as a signal to wake up
   402  	// a goroutine and handle new work, it must always be greater
   403  	// than or equal to sleeping.n. In the periods of time that it
   404  	// is strictly greater, it may cause spurious calls to wake.
   405  	asleep atomic.Uint32
   406  
   407  	// running indicates the number of cleanup goroutines actively
   408  	// executing user cleanup functions at any point in time.
   409  	//
   410  	// Read and written to without lock.
   411  	running atomic.Uint32
   412  
   413  	// ng is the number of cleanup goroutines.
   414  	//
   415  	// Read without lock, written only with lock held.
   416  	ng atomic.Uint32
   417  
   418  	// needg is the number of new cleanup goroutines that
   419  	// need to be created.
   420  	//
   421  	// Read without lock, written only with lock held.
   422  	needg atomic.Uint32
   423  
   424  	// Cleanup queue stats.
   425  
   426  	// queued represents a monotonic count of queued cleanups. This is sharded across
   427  	// Ps via the field cleanupsQueued in each p, so reading just this value is insufficient.
   428  	// In practice, this value only includes the queued count of dead Ps.
   429  	//
   430  	// Writes are protected by STW.
   431  	queued uint64
   432  
   433  	// executed is a monotonic count of executed cleanups.
   434  	//
   435  	// Read and updated atomically.
   436  	executed atomic.Uint64
   437  }
   438  
   439  // addWork indicates that n units of parallelizable work have been added to the queue.
   440  func (q *cleanupQueue) addWork(n int) {
   441  	q.workUnits.Add(int64(n))
   442  }
   443  
   444  // tryTakeWork is an attempt to dequeue some work by a cleanup goroutine.
   445  // This might fail if there's no work to do.
   446  func (q *cleanupQueue) tryTakeWork() bool {
   447  	for {
   448  		wu := q.workUnits.Load()
   449  		if wu == 0 {
   450  			return false
   451  		}
   452  		// CAS to prevent us from going negative.
   453  		if q.workUnits.CompareAndSwap(wu, wu-1) {
   454  			return true
   455  		}
   456  	}
   457  }
   458  
   459  // enqueue queues a single cleanup for execution.
   460  //
   461  // Called by the sweeper, and only the sweeper.
   462  func (q *cleanupQueue) enqueue(c cleanupFn) {
   463  	mp := acquirem()
   464  	pp := mp.p.ptr()
   465  	b := pp.cleanups
   466  	if b == nil {
   467  		if q.flushed.Load() {
   468  			q.flushed.Store(false)
   469  		}
   470  		b = (*cleanupBlock)(q.free.pop())
   471  		if b == nil {
   472  			b = (*cleanupBlock)(persistentalloc(cleanupBlockSize, tagAlign, &memstats.gcMiscSys))
   473  			for {
   474  				next := (*cleanupBlock)(q.all.Load())
   475  				b.alllink = next
   476  				if q.all.CompareAndSwap(unsafe.Pointer(next), unsafe.Pointer(b)) {
   477  					break
   478  				}
   479  			}
   480  		}
   481  		pp.cleanups = b
   482  	}
   483  	if full := b.enqueue(c); full {
   484  		q.full.push(&b.lfnode)
   485  		pp.cleanups = nil
   486  		q.addWork(1)
   487  	}
   488  	pp.cleanupsQueued++
   489  	releasem(mp)
   490  }
   491  
   492  // dequeue pops a block of cleanups from the queue. Blocks until one is available
   493  // and never returns nil.
   494  func (q *cleanupQueue) dequeue() *cleanupBlock {
   495  	for {
   496  		if q.tryTakeWork() {
   497  			// Guaranteed to be non-nil.
   498  			return (*cleanupBlock)(q.full.pop())
   499  		}
   500  		lock(&q.lock)
   501  		// Increment asleep first. We may have to undo this if we abort the sleep.
   502  		// We must update asleep first because the scheduler might not try to wake
   503  		// us up when work comes in between the last check of workUnits and when we
   504  		// go to sleep. (It may see asleep as 0.) By incrementing it here, we guarantee
   505  		// after this point that if new work comes in, someone will try to grab the
   506  		// lock and wake us. However, this also means that if we back out, we may cause
   507  		// someone to spuriously grab the lock and try to wake us up, only to fail.
   508  		// This should be very rare because the window here is incredibly small: the
   509  		// window between now and when we decrement q.asleep below.
   510  		q.asleep.Add(1)
   511  
   512  		// Re-check workUnits under the lock and with asleep updated. If it's still zero,
   513  		// then no new work came in, and it's safe for us to go to sleep. If new work
   514  		// comes in after this point, then the scheduler will notice that we're sleeping
   515  		// and wake us up.
   516  		if q.workUnits.Load() > 0 {
   517  			// Undo the q.asleep update and try to take work again.
   518  			q.asleep.Add(-1)
   519  			unlock(&q.lock)
   520  			continue
   521  		}
   522  		q.sleeping.push(getg())
   523  		goparkunlock(&q.lock, waitReasonCleanupWait, traceBlockSystemGoroutine, 1)
   524  	}
   525  }
   526  
   527  // flush pushes all active cleanup blocks to the full list and wakes up cleanup
   528  // goroutines to handle them.
   529  //
   530  // Must only be called at a point when we can guarantee that no more cleanups
   531  // are being queued, such as after the final sweeper for the cycle is done
   532  // but before the next mark phase.
   533  func (q *cleanupQueue) flush() {
   534  	mp := acquirem()
   535  	flushed := 0
   536  	emptied := 0
   537  	missing := 0
   538  
   539  	// Coalesce the partially-filled blocks to present a more accurate picture of demand.
   540  	// We use the number of coalesced blocks to process as a signal for demand to create
   541  	// new cleanup goroutines.
   542  	var cb *cleanupBlock
   543  	for _, pp := range allp {
   544  		if pp == nil {
   545  			// This function is reachable via mallocgc in the
   546  			// middle of procresize, when allp has been resized,
   547  			// but the new Ps not allocated yet.
   548  			missing++
   549  			continue
   550  		}
   551  		b := pp.cleanups
   552  		if b == nil {
   553  			missing++
   554  			continue
   555  		}
   556  		pp.cleanups = nil
   557  		if cb == nil {
   558  			cb = b
   559  			continue
   560  		}
   561  		// N.B. After take, either cb is full, b is empty, or both.
   562  		cb.take(b)
   563  		if cb.full() {
   564  			q.full.push(&cb.lfnode)
   565  			flushed++
   566  			cb = b
   567  			b = nil
   568  		}
   569  		if b != nil && b.empty() {
   570  			q.free.push(&b.lfnode)
   571  			emptied++
   572  		}
   573  	}
   574  	if cb != nil {
   575  		q.full.push(&cb.lfnode)
   576  		flushed++
   577  	}
   578  	if flushed != 0 {
   579  		q.addWork(flushed)
   580  	}
   581  	if flushed+emptied+missing != len(allp) {
   582  		throw("failed to correctly flush all P-owned cleanup blocks")
   583  	}
   584  	q.flushed.Store(true)
   585  	releasem(mp)
   586  }
   587  
   588  // needsWake returns true if cleanup goroutines may need to be awoken or created to handle cleanup load.
   589  func (q *cleanupQueue) needsWake() bool {
   590  	return q.workUnits.Load() > 0 && (q.asleep.Load() > 0 || q.ng.Load() < maxCleanupGs())
   591  }
   592  
   593  // wake wakes up one or more goroutines to process the cleanup queue. If there aren't
   594  // enough sleeping goroutines to handle the demand, wake will arrange for new goroutines
   595  // to be created.
   596  func (q *cleanupQueue) wake() {
   597  	lock(&q.lock)
   598  
   599  	// Figure out how many goroutines to wake, and how many extra goroutines to create.
   600  	// Wake one goroutine for each work unit.
   601  	var wake, extra uint32
   602  	work := q.workUnits.Load()
   603  	asleep := uint64(q.asleep.Load())
   604  	if work > asleep {
   605  		wake = uint32(asleep)
   606  		if work > uint64(math.MaxUint32) {
   607  			// Protect against overflow.
   608  			extra = math.MaxUint32
   609  		} else {
   610  			extra = uint32(work - asleep)
   611  		}
   612  	} else {
   613  		wake = uint32(work)
   614  		extra = 0
   615  	}
   616  	if extra != 0 {
   617  		// Signal that we should create new goroutines, one for each extra work unit,
   618  		// up to maxCleanupGs.
   619  		newg := min(extra, maxCleanupGs()-q.ng.Load())
   620  		if newg > 0 {
   621  			q.needg.Add(int32(newg))
   622  		}
   623  	}
   624  	if wake == 0 {
   625  		// Nothing to do.
   626  		unlock(&q.lock)
   627  		return
   628  	}
   629  
   630  	// Take ownership of waking 'wake' goroutines.
   631  	//
   632  	// Nobody else will wake up these goroutines, so they're guaranteed
   633  	// to be sitting on q.sleeping, waiting for us to wake them.
   634  	q.asleep.Add(-int32(wake))
   635  
   636  	// Collect them and schedule them.
   637  	var list gList
   638  	for range wake {
   639  		list.push(q.sleeping.pop())
   640  	}
   641  	unlock(&q.lock)
   642  
   643  	injectglist(&list)
   644  	return
   645  }
   646  
   647  func (q *cleanupQueue) needG() bool {
   648  	have := q.ng.Load()
   649  	if have >= maxCleanupGs() {
   650  		return false
   651  	}
   652  	if have == 0 {
   653  		// Make sure we have at least one.
   654  		return true
   655  	}
   656  	return q.needg.Load() > 0
   657  }
   658  
   659  func (q *cleanupQueue) createGs() {
   660  	lock(&q.lock)
   661  	have := q.ng.Load()
   662  	need := min(q.needg.Swap(0), maxCleanupGs()-have)
   663  	if have == 0 && need == 0 {
   664  		// Make sure we have at least one.
   665  		need = 1
   666  	}
   667  	if need > 0 {
   668  		q.ng.Add(int32(need))
   669  	}
   670  	unlock(&q.lock)
   671  
   672  	for range need {
   673  		go runCleanups()
   674  	}
   675  }
   676  
   677  func (q *cleanupQueue) beginRunningCleanups() {
   678  	// Update runningCleanups and running atomically with respect
   679  	// to goroutine profiles by disabling preemption.
   680  	mp := acquirem()
   681  	getg().runningCleanups.Store(true)
   682  	q.running.Add(1)
   683  	releasem(mp)
   684  }
   685  
   686  func (q *cleanupQueue) endRunningCleanups() {
   687  	// Update runningCleanups and running atomically with respect
   688  	// to goroutine profiles by disabling preemption.
   689  	mp := acquirem()
   690  	getg().runningCleanups.Store(false)
   691  	q.running.Add(-1)
   692  	releasem(mp)
   693  }
   694  
   695  func (q *cleanupQueue) readQueueStats() (queued, executed uint64) {
   696  	executed = q.executed.Load()
   697  	queued = q.queued
   698  
   699  	// N.B. This is inconsistent, but that's intentional. It's just an estimate.
   700  	// Read this _after_ reading executed to decrease the chance that we observe
   701  	// an inconsistency in the statistics (executed > queued).
   702  	for _, pp := range allp {
   703  		queued += pp.cleanupsQueued
   704  	}
   705  	return
   706  }
   707  
   708  func maxCleanupGs() uint32 {
   709  	// N.B. Left as a function to make changing the policy easier.
   710  	return uint32(max(gomaxprocs/4, 1))
   711  }
   712  
   713  // gcCleanups is the global cleanup queue.
   714  var gcCleanups cleanupQueue
   715  
   716  // runCleanups is the entrypoint for all cleanup-running goroutines.
   717  func runCleanups() {
   718  	for {
   719  		b := gcCleanups.dequeue()
   720  		if raceenabled {
   721  			// Approximately: adds a happens-before edge between the cleanup
   722  			// argument being mutated and the call to the cleanup below.
   723  			racefingo()
   724  		}
   725  
   726  		gcCleanups.beginRunningCleanups()
   727  		for i := 0; i < int(b.n); i++ {
   728  			c := b.cleanups[i]
   729  			b.cleanups[i] = cleanupFn{}
   730  
   731  			var racectx uintptr
   732  			if raceenabled {
   733  				// Enter a new race context so the race detector can catch
   734  				// potential races between cleanups, even if they execute on
   735  				// the same goroutine.
   736  				//
   737  				// Synchronize on fn. This would fail to find races on the
   738  				// closed-over values in fn (suppose arg is passed to multiple
   739  				// AddCleanup calls) if arg was not unique, but it is.
   740  				racerelease(unsafe.Pointer(c.arg))
   741  				racectx = raceEnterNewCtx()
   742  				raceacquire(unsafe.Pointer(c.arg))
   743  			}
   744  
   745  			// Execute the next cleanup.
   746  			c.call(c.fn, c.arg)
   747  
   748  			if raceenabled {
   749  				// Restore the old context.
   750  				raceRestoreCtx(racectx)
   751  			}
   752  		}
   753  		gcCleanups.endRunningCleanups()
   754  		gcCleanups.executed.Add(int64(b.n))
   755  
   756  		atomic.Store(&b.n, 0) // Synchronize with markroot. See comment in cleanupBlockHeader.
   757  		gcCleanups.free.push(&b.lfnode)
   758  	}
   759  }
   760  
   761  // blockUntilEmpty blocks until either the cleanup queue is emptied
   762  // and the cleanups have been executed, or the timeout is reached.
   763  // Returns true if the cleanup queue was emptied.
   764  // This is used by the sync and unique tests.
   765  func (q *cleanupQueue) blockUntilEmpty(timeout int64) bool {
   766  	start := nanotime()
   767  	for nanotime()-start < timeout {
   768  		lock(&q.lock)
   769  		// The queue is empty when there's no work left to do *and* all the cleanup goroutines
   770  		// are asleep. If they're not asleep, they may be actively working on a block.
   771  		if q.flushed.Load() && q.full.empty() && uint32(q.sleeping.size) == q.ng.Load() {
   772  			unlock(&q.lock)
   773  			return true
   774  		}
   775  		unlock(&q.lock)
   776  		Gosched()
   777  	}
   778  	return false
   779  }
   780  
   781  //go:linkname unique_runtime_blockUntilEmptyCleanupQueue unique.runtime_blockUntilEmptyCleanupQueue
   782  func unique_runtime_blockUntilEmptyCleanupQueue(timeout int64) bool {
   783  	return gcCleanups.blockUntilEmpty(timeout)
   784  }
   785  
   786  //go:linkname sync_test_runtime_blockUntilEmptyCleanupQueue sync_test.runtime_blockUntilEmptyCleanupQueue
   787  func sync_test_runtime_blockUntilEmptyCleanupQueue(timeout int64) bool {
   788  	return gcCleanups.blockUntilEmpty(timeout)
   789  }
   790  
   791  // raceEnterNewCtx creates a new racectx and switches the current
   792  // goroutine to it. Returns the old racectx.
   793  //
   794  // Must be running on a user goroutine. nosplit to match other race
   795  // instrumentation.
   796  //
   797  //go:nosplit
   798  func raceEnterNewCtx() uintptr {
   799  	// We use the existing ctx as the spawn context, but gp.gopc
   800  	// as the spawn PC to make the error output a little nicer
   801  	// (pointing to AddCleanup, where the goroutines are created).
   802  	//
   803  	// We also need to carefully indicate to the race detector
   804  	// that the goroutine stack will only be accessed by the new
   805  	// race context, to avoid false positives on stack locations.
   806  	// We do this by marking the stack as free in the first context
   807  	// and then re-marking it as allocated in the second. Crucially,
   808  	// there must be (1) no race operations and (2) no stack changes
   809  	// in between. (1) is easy to avoid because we're in the runtime
   810  	// so there's no implicit race instrumentation. To avoid (2) we
   811  	// defensively become non-preemptible so the GC can't stop us,
   812  	// and rely on the fact that racemalloc, racefreem, and racectx
   813  	// are nosplit.
   814  	mp := acquirem()
   815  	gp := getg()
   816  	ctx := getg().racectx
   817  	racefree(unsafe.Pointer(gp.stack.lo), gp.stack.hi-gp.stack.lo)
   818  	getg().racectx = racectxstart(gp.gopc, ctx)
   819  	racemalloc(unsafe.Pointer(gp.stack.lo), gp.stack.hi-gp.stack.lo)
   820  	releasem(mp)
   821  	return ctx
   822  }
   823  
   824  // raceRestoreCtx restores ctx on the goroutine. It is the inverse of
   825  // raceenternewctx and must be called with its result.
   826  //
   827  // Must be running on a user goroutine. nosplit to match other race
   828  // instrumentation.
   829  //
   830  //go:nosplit
   831  func raceRestoreCtx(ctx uintptr) {
   832  	mp := acquirem()
   833  	gp := getg()
   834  	racefree(unsafe.Pointer(gp.stack.lo), gp.stack.hi-gp.stack.lo)
   835  	racectxend(getg().racectx)
   836  	racemalloc(unsafe.Pointer(gp.stack.lo), gp.stack.hi-gp.stack.lo)
   837  	getg().racectx = ctx
   838  	releasem(mp)
   839  }
   840  

View as plain text