mirror of
https://github.com/weaveworks/scope.git
synced 2026-08-18 03:46:45 +00:00
vendor: update tcptracer-bpf and gobpf
This includes: - https://github.com/iovisor/gobpf/pull/70 perf: close go channels idiomatically - https://github.com/iovisor/gobpf/pull/70 close channels on the sender side & fix closing race - https://github.com/weaveworks/tcptracer-bpf/pull/50 vendor: update gobpf
This commit is contained in:
+19
-4
@@ -79,10 +79,19 @@ func NewTracer(cb Callback) (*Tracer, error) {
|
||||
for {
|
||||
select {
|
||||
case <-stopChan:
|
||||
// On stop, stopChan will be closed but the other channels will
|
||||
// also be closed shortly after. The select{} has no priorities,
|
||||
// therefore, the "ok" value must be checked below.
|
||||
return
|
||||
case data := <-channelV4:
|
||||
case data, ok := <-channelV4:
|
||||
if !ok {
|
||||
return // see explanation above
|
||||
}
|
||||
cb.TCPEventV4(tcpV4ToGo(&data))
|
||||
case lost := <-lostChanV4:
|
||||
case lost, ok := <-lostChanV4:
|
||||
if !ok {
|
||||
return // see explanation above
|
||||
}
|
||||
cb.LostV4(lost)
|
||||
}
|
||||
}
|
||||
@@ -93,9 +102,15 @@ func NewTracer(cb Callback) (*Tracer, error) {
|
||||
select {
|
||||
case <-stopChan:
|
||||
return
|
||||
case data := <-channelV6:
|
||||
case data, ok := <-channelV6:
|
||||
if !ok {
|
||||
return // see explanation above
|
||||
}
|
||||
cb.TCPEventV6(tcpV6ToGo(&data))
|
||||
case lost := <-lostChanV6:
|
||||
case lost, ok := <-lostChanV6:
|
||||
if !ok {
|
||||
return // see explanation above
|
||||
}
|
||||
cb.LostV6(lost)
|
||||
}
|
||||
}
|
||||
|
||||
Generated
Vendored
+20
-2
@@ -515,6 +515,7 @@ func (b *Module) Load(parameters map[string]SectionParams) error {
|
||||
isCgroupSkb := strings.HasPrefix(secName, "cgroup/skb")
|
||||
isCgroupSock := strings.HasPrefix(secName, "cgroup/sock")
|
||||
isSocketFilter := strings.HasPrefix(secName, "socket")
|
||||
isTracepoint := strings.HasPrefix(secName, "tracepoint/")
|
||||
|
||||
var progType uint32
|
||||
switch {
|
||||
@@ -528,9 +529,11 @@ func (b *Module) Load(parameters map[string]SectionParams) error {
|
||||
progType = uint32(C.BPF_PROG_TYPE_CGROUP_SOCK)
|
||||
case isSocketFilter:
|
||||
progType = uint32(C.BPF_PROG_TYPE_SOCKET_FILTER)
|
||||
case isTracepoint:
|
||||
progType = uint32(C.BPF_PROG_TYPE_TRACEPOINT)
|
||||
}
|
||||
|
||||
if isKprobe || isKretprobe || isCgroupSkb || isCgroupSock || isSocketFilter {
|
||||
if isKprobe || isKretprobe || isCgroupSkb || isCgroupSock || isSocketFilter || isTracepoint {
|
||||
rdata, err := rsection.Data()
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -579,6 +582,12 @@ func (b *Module) Load(parameters map[string]SectionParams) error {
|
||||
insns: insns,
|
||||
fd: int(progFd),
|
||||
}
|
||||
case isTracepoint:
|
||||
b.tracepointPrograms[secName] = &TracepointProgram{
|
||||
Name: secName,
|
||||
insns: insns,
|
||||
fd: int(progFd),
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -596,6 +605,7 @@ func (b *Module) Load(parameters map[string]SectionParams) error {
|
||||
isCgroupSkb := strings.HasPrefix(secName, "cgroup/skb")
|
||||
isCgroupSock := strings.HasPrefix(secName, "cgroup/sock")
|
||||
isSocketFilter := strings.HasPrefix(secName, "socket")
|
||||
isTracepoint := strings.HasPrefix(secName, "tracepoint/")
|
||||
|
||||
var progType uint32
|
||||
switch {
|
||||
@@ -609,9 +619,11 @@ func (b *Module) Load(parameters map[string]SectionParams) error {
|
||||
progType = uint32(C.BPF_PROG_TYPE_CGROUP_SOCK)
|
||||
case isSocketFilter:
|
||||
progType = uint32(C.BPF_PROG_TYPE_SOCKET_FILTER)
|
||||
case isTracepoint:
|
||||
progType = uint32(C.BPF_PROG_TYPE_TRACEPOINT)
|
||||
}
|
||||
|
||||
if isKprobe || isKretprobe || isCgroupSkb || isCgroupSock || isSocketFilter {
|
||||
if isKprobe || isKretprobe || isCgroupSkb || isCgroupSock || isSocketFilter || isTracepoint {
|
||||
data, err := section.Data()
|
||||
if err != nil {
|
||||
return err
|
||||
@@ -655,6 +667,12 @@ func (b *Module) Load(parameters map[string]SectionParams) error {
|
||||
insns: insns,
|
||||
fd: int(progFd),
|
||||
}
|
||||
case isTracepoint:
|
||||
b.tracepointPrograms[secName] = &TracepointProgram{
|
||||
Name: secName,
|
||||
insns: insns,
|
||||
fd: int(progFd),
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Generated
Vendored
+97
-21
@@ -96,11 +96,12 @@ type Module struct {
|
||||
fileReader io.ReaderAt
|
||||
file *elf.File
|
||||
|
||||
log []byte
|
||||
maps map[string]*Map
|
||||
probes map[string]*Kprobe
|
||||
cgroupPrograms map[string]*CgroupProgram
|
||||
socketFilters map[string]*SocketFilter
|
||||
log []byte
|
||||
maps map[string]*Map
|
||||
probes map[string]*Kprobe
|
||||
cgroupPrograms map[string]*CgroupProgram
|
||||
socketFilters map[string]*SocketFilter
|
||||
tracepointPrograms map[string]*TracepointProgram
|
||||
}
|
||||
|
||||
// Kprobe represents a kprobe or kretprobe and has to be declared
|
||||
@@ -134,13 +135,22 @@ type SocketFilter struct {
|
||||
fd int
|
||||
}
|
||||
|
||||
// TracepointProgram represents a tracepoint program
|
||||
type TracepointProgram struct {
|
||||
Name string
|
||||
insns *C.struct_bpf_insn
|
||||
fd int
|
||||
efd int
|
||||
}
|
||||
|
||||
func NewModule(fileName string) *Module {
|
||||
return &Module{
|
||||
fileName: fileName,
|
||||
probes: make(map[string]*Kprobe),
|
||||
cgroupPrograms: make(map[string]*CgroupProgram),
|
||||
socketFilters: make(map[string]*SocketFilter),
|
||||
log: make([]byte, 65536),
|
||||
fileName: fileName,
|
||||
probes: make(map[string]*Kprobe),
|
||||
cgroupPrograms: make(map[string]*CgroupProgram),
|
||||
socketFilters: make(map[string]*SocketFilter),
|
||||
tracepointPrograms: make(map[string]*TracepointProgram),
|
||||
log: make([]byte, 65536),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -186,6 +196,22 @@ func writeKprobeEvent(probeType, eventName, funcName, maxactiveStr string) (int,
|
||||
return kprobeId, nil
|
||||
}
|
||||
|
||||
func perfEventOpenTracepoint(id int, progFd int) (int, error) {
|
||||
efd, err := C.perf_event_open_tracepoint(C.int(id), -1 /* pid */, 0 /* cpu */, -1 /* group_fd */, C.PERF_FLAG_FD_CLOEXEC)
|
||||
if efd < 0 {
|
||||
return -1, fmt.Errorf("perf_event_open error: %v", err)
|
||||
}
|
||||
|
||||
if _, _, err := syscall.Syscall(syscall.SYS_IOCTL, uintptr(efd), C.PERF_EVENT_IOC_ENABLE, 0); err != 0 {
|
||||
return -1, fmt.Errorf("error enabling perf event: %v", err)
|
||||
}
|
||||
|
||||
if _, _, err := syscall.Syscall(syscall.SYS_IOCTL, uintptr(efd), C.PERF_EVENT_IOC_SET_BPF, uintptr(progFd)); err != 0 {
|
||||
return -1, fmt.Errorf("error attaching bpf program to perf event: %v", err)
|
||||
}
|
||||
return int(efd), nil
|
||||
}
|
||||
|
||||
// EnableKprobe enables a kprobe/kretprobe identified by secName.
|
||||
// For kretprobes, you can configure the maximum number of instances
|
||||
// of the function that can be probed simultaneously with maxactive.
|
||||
@@ -222,22 +248,43 @@ func (b *Module) EnableKprobe(secName string, maxactive int) error {
|
||||
return err
|
||||
}
|
||||
|
||||
efd := C.perf_event_open_tracepoint(C.int(kprobeId), -1 /* pid */, 0 /* cpu */, -1 /* group_fd */, C.PERF_FLAG_FD_CLOEXEC)
|
||||
if efd < 0 {
|
||||
return fmt.Errorf("perf_event_open for kprobe error")
|
||||
probe.efd, err = perfEventOpenTracepoint(kprobeId, progFd)
|
||||
return err
|
||||
}
|
||||
|
||||
func writeTracepointEvent(category, name string) (int, error) {
|
||||
tracepointIdFile := fmt.Sprintf("/sys/kernel/debug/tracing/events/%s/%s/id", category, name)
|
||||
tracepointIdBytes, err := ioutil.ReadFile(tracepointIdFile)
|
||||
if err != nil {
|
||||
return -1, fmt.Errorf("cannot read tracepoint id %q: %v", tracepointIdFile, err)
|
||||
}
|
||||
|
||||
_, _, err2 := syscall.Syscall(syscall.SYS_IOCTL, uintptr(efd), C.PERF_EVENT_IOC_ENABLE, 0)
|
||||
if err2 != 0 {
|
||||
return fmt.Errorf("error enabling perf event: %v", err2)
|
||||
tracepointId, err := strconv.Atoi(strings.TrimSpace(string(tracepointIdBytes)))
|
||||
if err != nil {
|
||||
return -1, fmt.Errorf("invalid tracepoint id: %v\n", err)
|
||||
}
|
||||
|
||||
_, _, err2 = syscall.Syscall(syscall.SYS_IOCTL, uintptr(efd), C.PERF_EVENT_IOC_SET_BPF, uintptr(progFd))
|
||||
if err2 != 0 {
|
||||
return fmt.Errorf("error enabling perf event: %v", err2)
|
||||
return tracepointId, nil
|
||||
}
|
||||
|
||||
func (b *Module) EnableTracepoint(secName string) error {
|
||||
prog, ok := b.tracepointPrograms[secName]
|
||||
if !ok {
|
||||
return fmt.Errorf("no such tracepoint program %q", secName)
|
||||
}
|
||||
probe.efd = int(efd)
|
||||
return nil
|
||||
progFd := prog.fd
|
||||
|
||||
tracepointGroup := strings.SplitN(secName, "/", 3)
|
||||
category := tracepointGroup[1]
|
||||
name := tracepointGroup[2]
|
||||
|
||||
tracepointId, err := writeTracepointEvent(category, name)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
prog.efd, err = perfEventOpenTracepoint(tracepointId, progFd)
|
||||
return err
|
||||
}
|
||||
|
||||
// IterKprobes returns a channel that emits the kprobes that included in the
|
||||
@@ -277,6 +324,17 @@ func (b *Module) IterCgroupProgram() <-chan *CgroupProgram {
|
||||
return ch
|
||||
}
|
||||
|
||||
func (b *Module) IterTracepointProgram() <-chan *TracepointProgram {
|
||||
ch := make(chan *TracepointProgram)
|
||||
go func() {
|
||||
for name := range b.tracepointPrograms {
|
||||
ch <- b.tracepointPrograms[name]
|
||||
}
|
||||
close(ch)
|
||||
}()
|
||||
return ch
|
||||
}
|
||||
|
||||
func (b *Module) CgroupProgram(name string) *CgroupProgram {
|
||||
return b.cgroupPrograms[name]
|
||||
}
|
||||
@@ -412,6 +470,21 @@ func (b *Module) closeProbes() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b *Module) closeTracepointPrograms() error {
|
||||
for _, program := range b.tracepointPrograms {
|
||||
if program.efd != -1 {
|
||||
if err := syscall.Close(program.efd); err != nil {
|
||||
return fmt.Errorf("error closing perf event fd: %v", err)
|
||||
}
|
||||
program.efd = -1
|
||||
}
|
||||
if err := syscall.Close(program.fd); err != nil {
|
||||
return fmt.Errorf("error closing tracepoint program fd: %v", err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b *Module) closeCgroupPrograms() error {
|
||||
for _, program := range b.cgroupPrograms {
|
||||
if err := syscall.Close(program.fd); err != nil {
|
||||
@@ -490,6 +563,9 @@ func (b *Module) CloseExt(options map[string]CloseOptions) error {
|
||||
if err := b.closeCgroupPrograms(); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := b.closeTracepointPrograms(); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := b.closeSocketFilters(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
Generated
Vendored
+26
-7
@@ -109,7 +109,7 @@ type PerfMap struct {
|
||||
pageCount int
|
||||
receiverChan chan []byte
|
||||
lostChan chan uint64
|
||||
pollStop chan bool
|
||||
pollStop chan struct{}
|
||||
timestamp func(*[]byte) uint64
|
||||
}
|
||||
|
||||
@@ -125,6 +125,9 @@ func InitPerfMap(b *Module, mapName string, receiverChan chan []byte, lostChan c
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("no map with name %s", mapName)
|
||||
}
|
||||
if receiverChan == nil {
|
||||
return nil, fmt.Errorf("receiverChan is nil")
|
||||
}
|
||||
// Maps are initialized in b.Load(), nothing to do here
|
||||
return &PerfMap{
|
||||
name: mapName,
|
||||
@@ -132,7 +135,7 @@ func InitPerfMap(b *Module, mapName string, receiverChan chan []byte, lostChan c
|
||||
pageCount: m.pageCount,
|
||||
receiverChan: receiverChan,
|
||||
lostChan: lostChan,
|
||||
pollStop: make(chan bool),
|
||||
pollStop: make(chan struct{}),
|
||||
}, nil
|
||||
}
|
||||
|
||||
@@ -163,6 +166,13 @@ func (pm *PerfMap) PollStart() {
|
||||
pageSize := os.Getpagesize()
|
||||
state := C.struct_read_state{}
|
||||
|
||||
defer func() {
|
||||
close(pm.receiverChan)
|
||||
if pm.lostChan != nil {
|
||||
close(pm.lostChan)
|
||||
}
|
||||
}()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-pm.pollStop:
|
||||
@@ -208,7 +218,11 @@ func (pm *PerfMap) PollStart() {
|
||||
}
|
||||
case C.PERF_RECORD_LOST:
|
||||
if pm.lostChan != nil {
|
||||
pm.lostChan <- lost.Lost
|
||||
select {
|
||||
case pm.lostChan <- lost.Lost:
|
||||
case <-pm.pollStop:
|
||||
return
|
||||
}
|
||||
}
|
||||
default:
|
||||
// ignore unknown events
|
||||
@@ -226,7 +240,11 @@ func (pm *PerfMap) PollStart() {
|
||||
// elements also must not be processed now.
|
||||
break harvestLoop
|
||||
}
|
||||
pm.receiverChan <- incoming.bytesArray[0]
|
||||
select {
|
||||
case pm.receiverChan <- incoming.bytesArray[0]:
|
||||
case <-pm.pollStop:
|
||||
return
|
||||
}
|
||||
// remove first element
|
||||
incoming.bytesArray = incoming.bytesArray[1:]
|
||||
}
|
||||
@@ -238,10 +256,11 @@ func (pm *PerfMap) PollStart() {
|
||||
}()
|
||||
}
|
||||
|
||||
// PollStop stops the goroutine that polls the perf event map. Make
|
||||
// sure to close the receiverChan only *after* calling PollStop.
|
||||
// PollStop stops the goroutine that polls the perf event map.
|
||||
// Callers must not close receiverChan or lostChan: they will be automatically
|
||||
// closed on the sender side.
|
||||
func (pm *PerfMap) PollStop() {
|
||||
pm.pollStop <- true
|
||||
close(pm.pollStop)
|
||||
}
|
||||
|
||||
func perfEventPoll(fds []C.int) error {
|
||||
|
||||
Vendored
+1
-1
@@ -1020,7 +1020,7 @@
|
||||
"importpath": "github.com/weaveworks/tcptracer-bpf",
|
||||
"repository": "https://github.com/weaveworks/tcptracer-bpf",
|
||||
"vcs": "git",
|
||||
"revision": "53a889d82fbd8a4fc4ef3f0013a9f88a3e9f4ccc",
|
||||
"revision": "e080bd747dc6b62d4ed3ed2b7f0be4801bef8faf",
|
||||
"branch": "master",
|
||||
"notests": true
|
||||
},
|
||||
|
||||
Reference in New Issue
Block a user