Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
## Unreleased

* [FEATURE] testutil: Add GatherAndFormat to encode a subset of metrics from a Gatherer. #2091
* [FEATURE] prometheus: `NewProcessCollector` on Darwin now reports `process_network_receive_bytes_total` and `process_network_transmit_bytes_total`, read via the `com.apple.network.statistics` kernel control socket (no cgo, no third-party dependency). Closes #1590. #2083

## 1.24.1 / 2026-07-23

Expand Down
9 changes: 7 additions & 2 deletions prometheus/process_collector_darwin.go
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,11 @@ func (c *processCollector) processCollect(ch chan<- Metric) {
c.reportError(ch, c.maxVsize, err)
}

// TODO: socket(PF_SYSTEM) to fetch "com.apple.network.statistics" might
// be able to get the per-process network send/receive counts.
if rxBytes, txBytes, err := getNetworkBytes(); err == nil {
ch <- MustNewConstMetric(c.inBytes, CounterValue, float64(rxBytes))
ch <- MustNewConstMetric(c.outBytes, CounterValue, float64(txBytes))
} else {
c.reportError(ch, c.inBytes, err)
c.reportError(ch, c.outBytes, err)
}
}
4 changes: 4 additions & 0 deletions prometheus/process_collector_darwin_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,11 +56,15 @@ func TestDarwinProcessCollector(t *testing.T) {
regexp.MustCompile("\nprocess_open_fds [1-9]"),
regexp.MustCompile("\nprocess_virtual_memory_max_bytes (-1|[1-9])"),
regexp.MustCompile("\nprocess_start_time_seconds [0-9]"),
regexp.MustCompile("\nprocess_network_receive_bytes_total [0-9]"),
regexp.MustCompile("\nprocess_network_transmit_bytes_total [0-9]"),
regexp.MustCompile("\nfoobar_process_cpu_seconds_total [0-9]"),
regexp.MustCompile("\nfoobar_process_max_fds [1-9]"),
regexp.MustCompile("\nfoobar_process_open_fds [1-9]"),
regexp.MustCompile("\nfoobar_process_virtual_memory_max_bytes (-1|[1-9])"),
regexp.MustCompile("\nfoobar_process_start_time_seconds [0-9]"),
regexp.MustCompile("\nfoobar_process_network_receive_bytes_total [0-9]"),
regexp.MustCompile("\nfoobar_process_network_transmit_bytes_total [0-9]"),
} {
if !re.Match(buf.Bytes()) {
t.Errorf("want body to match %s\n%s", re, buf.String())
Expand Down
3 changes: 0 additions & 3 deletions prometheus/process_collector_mem_cgo_darwin.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,9 +43,6 @@ func (c *processCollector) describe(ch chan<- *Desc) {
ch <- c.startTime
ch <- c.rss
ch <- c.vsize

/* the process could be collected but not implemented yet
ch <- c.inBytes
ch <- c.outBytes
*/
}
4 changes: 2 additions & 2 deletions prometheus/process_collector_mem_nocgo_darwin.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,11 +29,11 @@ func (c *processCollector) describe(ch chan<- *Desc) {
ch <- c.maxFDs
ch <- c.maxVsize
ch <- c.startTime
ch <- c.inBytes
ch <- c.outBytes

/* the process could be collected but not implemented yet
ch <- c.rss
ch <- c.vsize
ch <- c.inBytes
ch <- c.outBytes
*/
}
304 changes: 304 additions & 0 deletions prometheus/process_collector_netstat_darwin.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,304 @@
// Copyright The Prometheus Authors
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

//go:build darwin && !ios

package prometheus

import (
"bytes"
"encoding/binary"
"fmt"
"os"
"time"

"golang.org/x/sys/unix"
)

// This file implements a client for the undocumented "com.apple.network.statistics"
// kernel control socket, which is what Apple's own `nettop`/`netstat` tools use to
// report per-process network byte counters. There is no public Apple API for this
// and no cgo or third-party dependency is used: the protocol is implemented directly
// on top of golang.org/x/sys/unix, using the same PF_SYSTEM/SYSPROTO_CONTROL socket
// mechanism as the utun driver. Struct layouts below mirror bsd/net/ntstat.h from
// Apple's XNU source (https://github.com/apple/darwin-xnu/blob/main/bsd/net/ntstat.h).

const (
netStatControlName = "com.apple.network.statistics"

nstatProviderTCPUserland uint32 = 3
nstatProviderUDPUserland uint32 = 5

nstatMsgTypeAddAllSrcs uint32 = 1002
nstatMsgTypeRemSrc uint32 = 1003
nstatMsgTypeQuerySrc uint32 = 1004

nstatMsgTypeSuccess uint32 = 0
nstatMsgTypeError uint32 = 1
nstatMsgTypeSrcAdded uint32 = 10001
nstatMsgTypeSrcCounts uint32 = 10004

// Restrict subscription to sources owned by a single pid, rather than
// system-wide (which would require elevated privileges anyway).
nstatFilterSpecificUserByPid uint64 = 0x01000000

nstatReadTimeout = 2 * time.Second
)

type nstatMsgHdr struct {
Context uint64
Type uint32
Length uint16
Flags uint16
}

type nstatMsgAddAllSrcs struct {
Hdr nstatMsgHdr
Filter uint64
Events uint64
Provider uint32
TargetPid int32
TargetUUID [16]byte
}

// nstatMsgSrcAdded mirrors the leading fields of struct nstat_msg_src_added;
// trailing Provider/Reserved fields are present on the wire but unused here,
// so binary.Read simply leaves them unread.
type nstatMsgSrcAdded struct {
Hdr nstatMsgHdr
SrcRef uint64
}

type nstatMsgQuerySrcReq struct {
Hdr nstatMsgHdr
SrcRef uint64
}

// nstatCounts mirrors the leading fields of struct nstat_counts. The real
// struct has more trailing fields (retransmits, RTT estimates, etc.); we only
// need the byte counters, and binary.Read leaves the rest of the message
// unread, so trimming here avoids coupling to the full XNU layout.
type nstatCounts struct {
RxPackets uint64
RxBytes uint64
TxPackets uint64
TxBytes uint64
}

type nstatMsgSrcCounts struct {
Hdr nstatMsgHdr
SrcRef uint64
EventFlags uint64
Counts nstatCounts
}

// nstatMsgErr mirrors the leading fields of struct nstat_msg_error; the
// trailing Reserved field is unused here.
type nstatMsgErr struct {
Hdr nstatMsgHdr
Error uint32
}

// getNetworkBytes returns the total bytes received and sent over the network
// by the current process, summed across its TCP and UDP sockets.
func getNetworkBytes() (rxBytes, txBytes uint64, err error) {
fd, err := openNstatSocket()
if err != nil {
return 0, 0, err
}
defer unix.Close(fd)

pid := int32(os.Getpid())
for _, provider := range []uint32{nstatProviderTCPUserland, nstatProviderUDPUserland} {
refs, err := nstatCollectSrcRefs(fd, provider, pid)
if err != nil {
return 0, 0, fmt.Errorf("nstat: enumerating sources for provider %d: %w", provider, err)
}
for _, ref := range refs {
rx, tx, err := nstatQueryCounts(fd, ref)
if err != nil {
// The source may have been torn down between enumeration and
// query (e.g. a connection just closed); skip it rather than
// failing the whole collection.
continue
}
rxBytes += rx
txBytes += tx
}
}

return rxBytes, txBytes, nil
}

func openNstatSocket() (int, error) {
fd, err := unix.Socket(unix.AF_SYSTEM, unix.SOCK_DGRAM, 2 /* SYSPROTO_CONTROL */)
if err != nil {
return -1, fmt.Errorf("nstat: socket: %w", err)
}

ctlInfo := &unix.CtlInfo{}
copy(ctlInfo.Name[:], netStatControlName)
if err := unix.IoctlCtlInfo(fd, ctlInfo); err != nil {
unix.Close(fd)
return -1, fmt.Errorf("nstat: IoctlCtlInfo: %w", err)
}

if err := unix.Connect(fd, &unix.SockaddrCtl{ID: ctlInfo.Id}); err != nil {
unix.Close(fd)
return -1, fmt.Errorf("nstat: connect: %w", err)
}

tv := unix.NsecToTimeval(nstatReadTimeout.Nanoseconds())
if err := unix.SetsockoptTimeval(fd, unix.SOL_SOCKET, unix.SO_RCVTIMEO, &tv); err != nil {
unix.Close(fd)
return -1, fmt.Errorf("nstat: SetsockoptTimeval: %w", err)
}

return fd, nil
}

// nstatCollectSrcRefs subscribes to all sources of the given provider owned by
// pid, and returns the srcrefs the kernel reports. The kernel replies with zero
// or more SRC_ADDED messages followed by a SUCCESS message carrying the same
// context, which marks the end of enumeration.
func nstatCollectSrcRefs(fd int, provider uint32, pid int32) ([]uint64, error) {
const ctx = 1

req := nstatMsgAddAllSrcs{
Hdr: nstatMsgHdr{Context: ctx, Type: nstatMsgTypeAddAllSrcs},
Filter: nstatFilterSpecificUserByPid,
Provider: provider,
TargetPid: pid,
}
req.Hdr.Length = uint16(binary.Size(req))
if err := nstatSend(fd, req); err != nil {
return nil, err
}

var refs []uint64
buf := make([]byte, 4096)
for {
n, err := nstatRead(fd, buf)
if err != nil {
return nil, err
}

hdr, err := nstatReadHdr(buf[:n])
if err != nil {
return nil, err
}

switch hdr.Type {
case nstatMsgTypeSrcAdded:
var m nstatMsgSrcAdded
if err := binary.Read(bytes.NewReader(buf[:n]), binary.LittleEndian, &m); err != nil {
return nil, fmt.Errorf("nstat: decoding SRC_ADDED: %w", err)
}
refs = append(refs, m.SrcRef)
case nstatMsgTypeSuccess:
if hdr.Context == ctx {
return refs, nil
}
case nstatMsgTypeError:
nErr, err := nstatReadErr(buf[:n])
if err != nil {
return nil, err
}
return nil, fmt.Errorf("kernel returned errno %d", nErr)
}
}
}

func nstatQueryCounts(fd int, srcref uint64) (rxBytes, txBytes uint64, err error) {
const ctx = 2

req := nstatMsgQuerySrcReq{
Hdr: nstatMsgHdr{Context: ctx, Type: nstatMsgTypeQuerySrc},
SrcRef: srcref,
}
req.Hdr.Length = uint16(binary.Size(req))
if err := nstatSend(fd, req); err != nil {
return 0, 0, err
}

buf := make([]byte, 4096)
for {
n, err := nstatRead(fd, buf)
if err != nil {
return 0, 0, err
}

hdr, err := nstatReadHdr(buf[:n])
if err != nil {
return 0, 0, err
}

switch hdr.Type {
case nstatMsgTypeSrcCounts:
var m nstatMsgSrcCounts
if err := binary.Read(bytes.NewReader(buf[:n]), binary.LittleEndian, &m); err != nil {
return 0, 0, fmt.Errorf("nstat: decoding SRC_COUNTS: %w", err)
}
return m.Counts.RxBytes, m.Counts.TxBytes, nil
case nstatMsgTypeError:
nErr, err := nstatReadErr(buf[:n])
if err != nil {
return 0, 0, err
}
return 0, 0, fmt.Errorf("kernel returned errno %d", nErr)
}
}
}

func nstatSend(fd int, msg any) error {
buf := &bytes.Buffer{}
if err := binary.Write(buf, binary.LittleEndian, msg); err != nil {
return fmt.Errorf("nstat: encoding request: %w", err)
}
if _, err := unix.Write(fd, buf.Bytes()); err != nil {
return fmt.Errorf("nstat: write: %w", err)
}
return nil
}

// nstatRead reads one datagram from the control socket. SO_RCVTIMEO is set
// once on the socket in openNstatSocket, so a stalled kernel response can't
// hang collection forever.
func nstatRead(fd int, buf []byte) (int, error) {
n, err := unix.Read(fd, buf)
if err != nil {
return 0, fmt.Errorf("nstat: read: %w", err)
}
if n < nstatMsgHdrSize {
return 0, fmt.Errorf("nstat: short read (%d bytes)", n)
}
return n, nil
}

const nstatMsgHdrSize = 16

func nstatReadHdr(buf []byte) (nstatMsgHdr, error) {
var hdr nstatMsgHdr
if err := binary.Read(bytes.NewReader(buf), binary.LittleEndian, &hdr); err != nil {
return hdr, fmt.Errorf("nstat: decoding header: %w", err)
}
return hdr, nil
}

func nstatReadErr(buf []byte) (uint32, error) {
var m nstatMsgErr
if err := binary.Read(bytes.NewReader(buf), binary.LittleEndian, &m); err != nil {
return 0, fmt.Errorf("nstat: decoding ERROR: %w", err)
}
return m.Error, nil
}
Loading