mirror of
https://github.com/weaveworks/scope.git
synced 2026-07-28 09:41:57 +00:00
plugins/traffic-control: added status caching
To avoid to run `tc` every tic to get the qdisc status, we cache it. The cache is updated every time a value is changed via the control interface. NOTE: This will keep track only of the changes by done the plugins interface, if the qdisc parameters are changed by another applcation or by the user via terminal, such changes will not be caught.
This commit is contained in:
@@ -50,6 +50,14 @@ type Plugin struct {
|
||||
clients []containerClient
|
||||
}
|
||||
|
||||
type trafficControlStatus struct {
|
||||
latency uint32
|
||||
pakLoss float32
|
||||
// timestmap?
|
||||
}
|
||||
|
||||
var trafficControlStatusCache map[string]trafficControlStatus
|
||||
|
||||
func main() {
|
||||
const socket = "/var/run/scope/plugins/traffic-control.sock"
|
||||
|
||||
@@ -68,6 +76,7 @@ func main() {
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to create a plugin: %v", err)
|
||||
}
|
||||
trafficControlStatusCache = make(map[string]trafficControlStatus)
|
||||
if err := plugin.Serve(listener); err != nil {
|
||||
log.Fatalf("failed to serve: %v", err)
|
||||
}
|
||||
|
||||
@@ -2,6 +2,7 @@ package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"os/exec"
|
||||
"strings"
|
||||
|
||||
@@ -43,6 +44,16 @@ func DoTrafficControl(pid int, delay uint32) error {
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to perform traffic control: %v", err)
|
||||
}
|
||||
// cache parameters
|
||||
if netNSID, err := getNSID(netNS); err != nil {
|
||||
log.Error(netNSID)
|
||||
return fmt.Errorf("failed to get network namespace ID: %v", err)
|
||||
} else {
|
||||
trafficControlStatusCache[netNSID] = trafficControlStatus{
|
||||
latency: delay,
|
||||
pakLoss: 0,
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -54,18 +65,21 @@ func getLatency(pid int) (string, error) {
|
||||
} else if output == "" {
|
||||
return "-", fmt.Errorf("Error: output is empty")
|
||||
}
|
||||
outputSplited := split(output)
|
||||
for i, s := range outputSplited {
|
||||
return parseLatency(output)
|
||||
}
|
||||
|
||||
func parseLatency(statusString string) (string, error) {
|
||||
statusStringSplited := split(statusString)
|
||||
for i, s := range statusStringSplited {
|
||||
if s == "delay" {
|
||||
if i < len(outputSplited)-1 {
|
||||
output = outputSplited[i+1]
|
||||
if i < len(statusStringSplited)-1 {
|
||||
return statusStringSplited[i+1], nil
|
||||
} else {
|
||||
output = "-"
|
||||
return "-", nil
|
||||
}
|
||||
return output, nil
|
||||
}
|
||||
}
|
||||
return output, fmt.Errorf("delay not found")
|
||||
return "N/A", fmt.Errorf("delay not found")
|
||||
}
|
||||
|
||||
func getPktLoss(pid int) (string, error) {
|
||||
@@ -90,11 +104,36 @@ func getPktLoss(pid int) (string, error) {
|
||||
return output, fmt.Errorf("delay not found")
|
||||
}
|
||||
|
||||
func parsePktLoss(statusString string) (string, error) {
|
||||
statusStringSplited := split(statusString)
|
||||
for i, s := range statusStringSplited {
|
||||
if s == "loss" {
|
||||
if i < len(statusStringSplited)-1 {
|
||||
return statusStringSplited[i+1], nil
|
||||
} else {
|
||||
return "-", nil
|
||||
}
|
||||
}
|
||||
}
|
||||
return "N/A", fmt.Errorf("delay not found")
|
||||
}
|
||||
|
||||
// TODO @alepuccetti: return the ful status structure
|
||||
func getStatus(pid int) (string, error) {
|
||||
cmd := split("tc qdisc show dev eth0")
|
||||
netNS := fmt.Sprintf("/proc/%d/ns/net", pid)
|
||||
netNSID, err := getNSID(netNS)
|
||||
if err != nil {
|
||||
log.Error(netNSID)
|
||||
return "-", fmt.Errorf("failed to get network namespace ID: %v", err)
|
||||
}
|
||||
if status, ok := trafficControlStatusCache[netNSID]; ok {
|
||||
log.Info("Happy caching")
|
||||
return string(status.latency), nil
|
||||
}
|
||||
log.Info("cache miss: execute tc")
|
||||
cmd := split("tc qdisc show dev eth0")
|
||||
var output string
|
||||
err := ns.WithNetNSPath(netNS, func(hostNS ns.NetNS) error {
|
||||
err = ns.WithNetNSPath(netNS, func(hostNS ns.NetNS) error {
|
||||
if cmdOut, err := exec.Command(cmd[0], cmd[1:]...).CombinedOutput(); err != nil {
|
||||
log.Error(string(cmdOut))
|
||||
output = "-"
|
||||
@@ -107,6 +146,15 @@ func getStatus(pid int) (string, error) {
|
||||
return output, err
|
||||
}
|
||||
|
||||
func getNSID(nsPath string) (string, error) {
|
||||
if nsID, err := os.Readlink(nsPath); err != nil {
|
||||
log.Error(nsID)
|
||||
return "", fmt.Errorf("failed to execute command: tc qdisc show dev eth0: %v", err)
|
||||
} else {
|
||||
return nsID, nil
|
||||
}
|
||||
}
|
||||
|
||||
func split(cmd string) []string {
|
||||
return strings.Split(cmd, " ")
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user