diff --git a/vendor/github.com/weaveworks/tcptracer-bpf/pkg/tracer/tracer.go b/vendor/github.com/weaveworks/tcptracer-bpf/pkg/tracer/tracer.go index 8d725a36f..ad21afdbb 100644 --- a/vendor/github.com/weaveworks/tcptracer-bpf/pkg/tracer/tracer.go +++ b/vendor/github.com/weaveworks/tcptracer-bpf/pkg/tracer/tracer.go @@ -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) } } diff --git a/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/elf.go b/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/elf.go index 5f9926e9c..c5ce0438e 100644 --- a/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/elf.go +++ b/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/elf.go @@ -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), + } } } } diff --git a/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/module.go b/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/module.go index 4270806be..9b850d311 100644 --- a/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/module.go +++ b/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/module.go @@ -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 } diff --git a/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/perf.go b/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/perf.go index 75ba3f8e6..1ae543a16 100644 --- a/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/perf.go +++ b/vendor/github.com/weaveworks/tcptracer-bpf/vendor/github.com/iovisor/gobpf/elf/perf.go @@ -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 { diff --git a/vendor/manifest b/vendor/manifest index 63d1395ec..a595c8e7e 100644 --- a/vendor/manifest +++ b/vendor/manifest @@ -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 },