diff --git a/probe/endpoint/ebpf.go b/probe/endpoint/ebpf.go new file mode 100644 index 000000000..08847e3c3 --- /dev/null +++ b/probe/endpoint/ebpf.go @@ -0,0 +1,148 @@ +package endpoint + +import ( + "bufio" + "fmt" + "net" + "os" + "os/exec" + "strconv" + "strings" + + log "github.com/Sirupsen/logrus" +) + +type event int + +const ( + Connect event = iota + Accept + Close +) + +type connectionEvent struct { + Type event + Pid int + Command string + SourceAddress net.IP + DestAddress net.IP + SourcePort int + DestPort int + Netns int +} + +type EbpfTracker struct { + Cmd *exec.Cmd +} + +func NewEbpfTracker(bccProgramPath string) *EbpfTracker { + cmd := exec.Command(bccProgramPath) + env := os.Environ() + cmd.Env = append(env, "PYTHONUNBUFFERED=1") + + stderr, err := cmd.StderrPipe() + if err != nil { + log.Errorf("bcc error: %v", err) + return nil + } + go logPipe("bcc stderr:", stderr) + + tracker := &EbpfTracker{ + Cmd: cmd, + } + go tracker.run() + return tracker +} + +func (t *EbpfTracker) run() { + stdout, err := t.Cmd.StdoutPipe() + if err != nil { + log.Errorf("conntrack error: %v", err) + return + } + + if err := t.Cmd.Start(); err != nil { + log.Errorf("bcc error: %v", err) + return + } + + defer func() { + if err := t.Cmd.Wait(); err != nil { + log.Errorf("bcc error: %v", err) + } + }() + + reader := bufio.NewReader(stdout) + // skip fist line + if _, err := reader.ReadString('\n'); err != nil { + log.Errorf("bcc error: %v", err) + return + } + + defer log.Infof("bcc exiting") + + scn := bufio.NewScanner(reader) + for scn.Scan() { + txt := scn.Text() + line := strings.Fields(txt) + + pid, err := strconv.Atoi(line[1]) + if err != nil { + log.Errorf("error parsing pid %q: %v", line[1], err) + continue + } + + sAddr := net.ParseIP(line[3]) + if sAddr == nil { + log.Errorf("error parsing sAddr %q: %v", line[3], err) + continue + } + + dAddr := net.ParseIP(line[4]) + if sAddr == nil { + log.Errorf("error parsing dAddr %q: %v", line[4], err) + continue + } + + sPort, err := strconv.Atoi(line[5]) + if err != nil { + log.Errorf("error parsing sPort %q: %v", line[5], err) + continue + } + + dPort, err := strconv.Atoi(line[6]) + if err != nil { + log.Errorf("error parsing dPort %q: %v", line[6], err) + continue + } + + netns, err := strconv.Atoi(line[6]) + if err != nil { + log.Errorf("error parsing netns %q: %v", line[7], err) + continue + } + + var evt event + switch line[0] { + case "connect": + evt = Connect + case "accept": + evt = Accept + case "close": + evt = Close + } + + e := connectionEvent{ + Type: evt, + Pid: pid, + Command: line[2], + SourceAddress: sAddr, + DestAddress: dAddr, + SourcePort: sPort, + DestPort: dPort, + Netns: netns, + } + + fmt.Println(e) + } +} diff --git a/probe/endpoint/ebpf/main.go b/probe/endpoint/ebpf/main.go new file mode 100644 index 000000000..d7cdabe30 --- /dev/null +++ b/probe/endpoint/ebpf/main.go @@ -0,0 +1,22 @@ +package main + +import ( + "fmt" + "os" + "time" + + "github.com/weaveworks/scope/probe/endpoint" +) + +func main() { + tr := endpoint.NewEbpfTracker("/usr/local/share/bcc/examples/tracing/tcpv4tracer.py") + + if tr == nil { + fmt.Fprintf(os.Stderr, "error creating tracker\n") + os.Exit(1) + } + + time.Sleep(100 * time.Second) + + fmt.Println("done") +}