From f922ea19c82088878ac9ebe340fd238dab04b0e9 Mon Sep 17 00:00:00 2001 From: Alfonso Acosta Date: Thu, 4 Feb 2016 12:33:04 +0000 Subject: [PATCH] Rate-limit reading proc files Use a reader in the background, dynamically rate-limited, reading the required files in a loop --- .../procspy/background_reader_linux.go | 110 ++++++++++++++++++ probe/endpoint/procspy/proc_linux.go | 4 +- probe/endpoint/procspy/spy_linux.go | 7 +- probe/endpoint/reporter.go | 1 + 4 files changed, 118 insertions(+), 4 deletions(-) create mode 100644 probe/endpoint/procspy/background_reader_linux.go diff --git a/probe/endpoint/procspy/background_reader_linux.go b/probe/endpoint/procspy/background_reader_linux.go new file mode 100644 index 000000000..3c5ef0f2f --- /dev/null +++ b/probe/endpoint/procspy/background_reader_linux.go @@ -0,0 +1,110 @@ +package procspy + +import ( + "bytes" + "fmt" + "log" + "math" + "sync" + "time" + + "github.com/weaveworks/scope/probe/process" +) + +const ( + initialRateLimit = 100 * time.Millisecond // read 10 namespaces per second + maxRateLimit = 250 * time.Millisecond // lead at least 4 namespaces per second + targetWalkTime = 10 * time.Second + + maxRateLimitF = float64(maxRateLimit) + targetWalkTimeF = float64(targetWalkTime) +) + +type backgroundReader struct { + walker process.Walker + mtx sync.Mutex + walkingBuf *bytes.Buffer + readyBuf *bytes.Buffer + readySockets map[uint64]*Proc +} + +// HACK: Pretty ugly singleton interface (particularly the part of passing +// the walker to StartBackgroundReader() and ignoring it in in Connections() ) +// experimenting with this for now. +var singleton *backgroundReader + +func getBackgroundReader() (*backgroundReader, error) { + var err error + if singleton == nil { + err = fmt.Errorf("background reader hasn't yet been started") + } + return singleton, err +} + +// StartBackgroundReader starts a ratelimited background goroutine to +// read the expensive files from proc. +func StartBackgroundReader(walker process.Walker) { + if singleton != nil { + return + } + singleton = &backgroundReader{ + walker: walker, + walkingBuf: bytes.NewBuffer(make([]byte, 0, 5000)), + readyBuf: bytes.NewBuffer(make([]byte, 0, 5000)), + } + go singleton.loop() +} + +func (br *backgroundReader) loop() { + rateLimit := initialRateLimit + + namespaceTicker := time.Tick(rateLimit) + + for { + start := time.Now() + sockets, err := walkProcPid(br.walkingBuf, br.walker, namespaceTicker) + if err != nil { + log.Printf("background reader: error walking /proc: %s\n", err) + continue + } + walkTime := time.Now().Sub(start) + walkTimeF := float64(walkTime) + + log.Printf("debug: background reader: full pass took %s\n", walkTime) + if walkTimeF/targetWalkTimeF > 1.5 { + log.Printf( + "warn: background reader: full pass took %s: 50%% more than expected (%s)\n", + walkTime, + targetWalkTime, + ) + } + + // Adjust rate limit to more-accurately meet the target walk time in next iteration + scaledRateLimit := targetWalkTimeF / walkTimeF * float64(rateLimit) + rateLimit = time.Duration(math.Min(scaledRateLimit, maxRateLimitF)) + log.Printf("debug: background reader: new rate limit %s\n", rateLimit) + + namespaceTicker = time.Tick(rateLimit) + + // Swap buffers + br.mtx.Lock() + br.readyBuf, br.walkingBuf = br.walkingBuf, br.readyBuf + br.readySockets = sockets + br.mtx.Unlock() + + br.walkingBuf.Reset() + + // Sleep during spare time + time.Sleep(targetWalkTime - walkTime) + } +} + +func (br *backgroundReader) getWalkedProcPid(buf *bytes.Buffer) map[uint64]*Proc { + br.mtx.Lock() + defer br.mtx.Unlock() + + reader := bytes.NewReader(br.readyBuf.Bytes()) + buf.ReadFrom(reader) + + return br.readySockets +} diff --git a/probe/endpoint/procspy/proc_linux.go b/probe/endpoint/procspy/proc_linux.go index 87ac4bd43..9545415f5 100644 --- a/probe/endpoint/procspy/proc_linux.go +++ b/probe/endpoint/procspy/proc_linux.go @@ -7,6 +7,7 @@ import ( "path/filepath" "strconv" "syscall" + "time" log "github.com/Sirupsen/logrus" "github.com/armon/go-metrics" @@ -143,7 +144,7 @@ func walkNamespacePid(buf *bytes.Buffer, sockets map[uint64]*Proc, namespaceProc // /proc/PID/net/tcp{,6} for each namespace and sees if the ./fd/* files of each // process in that namespace are symlinks to sockets. Returns a map from socket // ID (inode) to PID. -func walkProcPid(buf *bytes.Buffer, walker process.Walker) (map[uint64]*Proc, error) { +func walkProcPid(buf *bytes.Buffer, walker process.Walker, namespaceTicker <-chan time.Time) (map[uint64]*Proc, error) { var ( sockets = map[uint64]*Proc{} // map socket inode -> process namespaces = map[uint64][]*process.Process{} // map network namespace id -> processes @@ -171,6 +172,7 @@ func walkProcPid(buf *bytes.Buffer, walker process.Walker) (map[uint64]*Proc, er }) for _, procs := range namespaces { + <-namespaceTicker walkNamespacePid(buf, sockets, procs) } diff --git a/probe/endpoint/procspy/spy_linux.go b/probe/endpoint/procspy/spy_linux.go index 68b0d7c0b..4429aa277 100644 --- a/probe/endpoint/procspy/spy_linux.go +++ b/probe/endpoint/procspy/spy_linux.go @@ -33,17 +33,18 @@ func (c *pnConnIter) Next() *Connection { } // cbConnections sets Connections() -var cbConnections = func(processes bool, walker process.Walker) (ConnIter, error) { +var cbConnections = func(processes bool, _ process.Walker) (ConnIter, error) { // buffer for contents of /proc//net/tcp buf := bufPool.Get().(*bytes.Buffer) buf.Reset() var procs map[uint64]*Proc if processes { - var err error - if procs, err = walkProcPid(buf, walker); err != nil { + br, err := getBackgroundReader() + if err != nil { return nil, err } + procs = br.getWalkedProcPid(buf) } if buf.Len() == 0 { diff --git a/probe/endpoint/reporter.go b/probe/endpoint/reporter.go index c2f9c83f0..2adf1c9e2 100644 --- a/probe/endpoint/reporter.go +++ b/probe/endpoint/reporter.go @@ -49,6 +49,7 @@ var SpyDuration = prometheus.NewSummaryVec( // is stored in the Endpoint topology. It optionally enriches that topology // with process (PID) information. func NewReporter(hostID, hostName string, includeProcesses bool, useConntrack bool, procWalker process.Walker) *Reporter { + procspy.StartBackgroundReader(procWalker) return &Reporter{ hostID: hostID, hostName: hostName,