Improve probe docker code quality & test coverage.

- Move docker probe code into it's own module
- Put PIDTree behind and interface for mocking
- Disaggregate dockerTagger into a registry, tagger and reporter
- Similarly disaggregate tests
- Add mocks for docker container and registry
- Add test for docker events & stats
This commit is contained in:
Tom Wilkie
2015-06-18 17:09:33 +00:00
parent a6ef295bd8
commit 314af5ca89
20 changed files with 1062 additions and 596 deletions
+218
View File
@@ -0,0 +1,218 @@
package docker
import (
"bufio"
"encoding/json"
"fmt"
"io"
"log"
"net"
"net/http"
"net/http/httputil"
"net/url"
"strconv"
"strings"
"sync"
docker "github.com/fsouza/go-dockerclient"
"github.com/weaveworks/scope/report"
)
// These constants are keys used in node metadata
// TODO: use these constants in report/{mapping.go, detailed_node.go} - pending some circular references
const (
NetworkRxDropped = "network_rx_dropped"
NetworkRxBytes = "network_rx_bytes"
NetworkRxErrors = "network_rx_errors"
NetworkTxPackets = "network_tx_packets"
NetworkTxDropped = "network_tx_dropped"
NetworkRxPackets = "network_rx_packets"
NetworkTxErrors = "network_tx_errors"
NetworkTxBytes = "network_tx_bytes"
MemoryMaxUsage = "memory_max_usage"
MemoryUsage = "memory_usage"
MemoryFailcnt = "memory_failcnt"
MemoryLimit = "memory_limit"
CPUPercpuUsage = "cpu_per_cpu_usage"
CPUUsageInUsermode = "cpu_usage_in_usermode"
CPUTotalUsage = "cpu_total_usage"
CPUUsageInKernelmode = "cpu_usage_in_kernelmode"
CPUSystemCPUUsage = "cpu_system_cpu_usage"
)
// Exported for testing
var (
DialStub = net.Dial
NewClientConnStub = newClientConn
)
func newClientConn(c net.Conn, r *bufio.Reader) ClientConn {
return httputil.NewClientConn(c, r)
}
// ClientConn is exported for testing
type ClientConn interface {
Do(req *http.Request) (resp *http.Response, err error)
Close() error
}
// Container represents a docker container
type Container interface {
ID() string
Image() string
PID() int
GetNodeMetadata() report.NodeMetadata
StartGatheringStats() error
StopGatheringStats()
}
type container struct {
sync.RWMutex
container *docker.Container
statsConn ClientConn
latestStats *docker.Stats
}
// NewContainer creates a new Container
func NewContainer(c *docker.Container) Container {
return &container{container: c}
}
func (c *container) ID() string {
return c.container.ID
}
func (c *container) Image() string {
return c.container.Image
}
func (c *container) PID() int {
return c.container.State.Pid
}
func (c *container) StartGatheringStats() error {
c.Lock()
defer c.Unlock()
if c.statsConn != nil {
return fmt.Errorf("already gather stats for container %s", c.container.ID)
}
go func() {
log.Printf("docker container: collecting stats for %s", c.container.ID)
req, err := http.NewRequest("GET", fmt.Sprintf("/containers/%s/stats", c.container.ID), nil)
if err != nil {
log.Printf("docker container: %v", err)
return
}
req.Header.Set("User-Agent", "weavescope")
url, err := url.Parse(endpoint)
if err != nil {
log.Printf("docker container: %v", err)
return
}
dial, err := net.Dial(url.Scheme, url.Path)
if err != nil {
log.Printf("docker container: %v", err)
return
}
conn := NewClientConnStub(dial, nil)
resp, err := conn.Do(req)
if err != nil {
log.Printf("docker container: %v", err)
return
}
c.Lock()
c.statsConn = conn
c.Unlock()
defer func() {
c.Lock()
defer c.Unlock()
log.Printf("docker container: stopped collecting stats for %s", c.container.ID)
c.statsConn = nil
c.latestStats = nil
}()
stats := &docker.Stats{}
decoder := json.NewDecoder(resp.Body)
for err := decoder.Decode(&stats); err != io.EOF; err = decoder.Decode(&stats) {
if err != nil {
log.Printf("docker container: error reading event %v", err)
return
}
c.Lock()
c.latestStats = stats
c.Unlock()
stats = &docker.Stats{}
}
}()
return nil
}
// called whilst holding t.Lock()
func (c *container) StopGatheringStats() {
c.Lock()
defer c.Unlock()
if c.statsConn == nil {
return
}
c.statsConn.Close()
c.statsConn = nil
c.latestStats = nil
return
}
// called whilst holding t.RLock()
func (c *container) GetNodeMetadata() report.NodeMetadata {
c.RLock()
defer c.RUnlock()
result := report.NodeMetadata{
ContainerID: c.ID(),
ContainerName: strings.TrimPrefix(c.container.Name, "/"),
ImageID: c.container.Image,
}
if c.latestStats == nil {
return result
}
result.Merge(report.NodeMetadata{
NetworkRxDropped: strconv.FormatUint(c.latestStats.Network.RxDropped, 10),
NetworkRxBytes: strconv.FormatUint(c.latestStats.Network.RxBytes, 10),
NetworkRxErrors: strconv.FormatUint(c.latestStats.Network.RxErrors, 10),
NetworkTxPackets: strconv.FormatUint(c.latestStats.Network.TxPackets, 10),
NetworkTxDropped: strconv.FormatUint(c.latestStats.Network.TxDropped, 10),
NetworkRxPackets: strconv.FormatUint(c.latestStats.Network.RxPackets, 10),
NetworkTxErrors: strconv.FormatUint(c.latestStats.Network.TxErrors, 10),
NetworkTxBytes: strconv.FormatUint(c.latestStats.Network.TxBytes, 10),
MemoryMaxUsage: strconv.FormatUint(c.latestStats.MemoryStats.MaxUsage, 10),
MemoryUsage: strconv.FormatUint(c.latestStats.MemoryStats.Usage, 10),
MemoryFailcnt: strconv.FormatUint(c.latestStats.MemoryStats.Failcnt, 10),
MemoryLimit: strconv.FormatUint(c.latestStats.MemoryStats.Limit, 10),
// CPUPercpuUsage: strconv.FormatUint(stats.CPUStats.CPUUsage.PercpuUsage, 10),
CPUUsageInUsermode: strconv.FormatUint(c.latestStats.CPUStats.CPUUsage.UsageInUsermode, 10),
CPUTotalUsage: strconv.FormatUint(c.latestStats.CPUStats.CPUUsage.TotalUsage, 10),
CPUUsageInKernelmode: strconv.FormatUint(c.latestStats.CPUStats.CPUUsage.UsageInKernelmode, 10),
CPUSystemCPUUsage: strconv.FormatUint(c.latestStats.CPUStats.SystemCPUUsage, 10),
})
return result
}
+67
View File
@@ -0,0 +1,67 @@
package docker_test
import (
"bufio"
"encoding/json"
"io"
"net"
"net/http"
"runtime"
"testing"
client "github.com/fsouza/go-dockerclient"
"github.com/weaveworks/scope/probe/docker"
)
type mockConnection struct {
reader *io.PipeReader
}
func (c *mockConnection) Do(req *http.Request) (resp *http.Response, err error) {
return &http.Response{
Body: c.reader,
}, nil
}
func (c *mockConnection) Close() error {
return c.reader.Close()
}
func TestContainer(t *testing.T) {
oldDialStub, oldNewClientConnStub := docker.DialStub, docker.NewClientConnStub
defer func() { docker.DialStub, docker.NewClientConnStub = oldDialStub, oldNewClientConnStub }()
docker.DialStub = func(network, address string) (net.Conn, error) {
return nil, nil
}
reader, writer := io.Pipe()
connection := &mockConnection{reader}
docker.NewClientConnStub = func(c net.Conn, r *bufio.Reader) docker.ClientConn {
return connection
}
c := docker.NewContainer(container1)
err := c.StartGatheringStats()
if err != nil {
t.Errorf("%v", err)
}
defer c.StopGatheringStats()
runtime.Gosched() // wait for StartGatheringStats goroutine to call connection.Do
// Send some stats to the docker container
stats := &client.Stats{}
stats.MemoryStats.Usage = 12345
err = json.NewEncoder(writer).Encode(&stats)
if err != nil {
t.Errorf("%v", err)
}
runtime.Gosched() // wait for StartGatheringStats goroutine to receive the stats
// Now see if we go them
nmd := c.GetNodeMetadata()
if nmd[docker.MemoryUsage] != "12345" {
t.Errorf("want 12345, got %s", nmd[docker.MemoryUsage])
}
}
+288
View File
@@ -0,0 +1,288 @@
package docker
import (
"log"
"sync"
"time"
docker_client "github.com/fsouza/go-dockerclient"
)
// Consts exported for testing.
const (
StartEvent = "start"
DieEvent = "die"
endpoint = "unix:///var/run/docker.sock"
)
// Vars exported for testing.
var (
NewDockerClientStub = newDockerClient
NewContainerStub = NewContainer
)
// Registry keeps track of running docker containers and their images
type Registry interface {
Stop()
LockedPIDLookup(f func(func(int) Container))
WalkContainers(f func(Container))
WalkImages(f func(*docker_client.APIImages))
}
type registry struct {
sync.RWMutex
quit chan chan struct{}
interval time.Duration
client Client
containers map[string]Container
containersByPID map[int]Container
images map[string]*docker_client.APIImages
}
// Client interface for mocking.
type Client interface {
ListContainers(docker_client.ListContainersOptions) ([]docker_client.APIContainers, error)
InspectContainer(string) (*docker_client.Container, error)
ListImages(docker_client.ListImagesOptions) ([]docker_client.APIImages, error)
AddEventListener(chan<- *docker_client.APIEvents) error
RemoveEventListener(chan *docker_client.APIEvents) error
}
func newDockerClient(endpoint string) (Client, error) {
return docker_client.NewClient(endpoint)
}
// NewRegistry returns a usable Registry. Don't forget to Stop it.
func NewRegistry(interval time.Duration) (Registry, error) {
client, err := NewDockerClientStub(endpoint)
if err != nil {
return nil, err
}
r := &registry{
containers: map[string]Container{},
containersByPID: map[int]Container{},
images: map[string]*docker_client.APIImages{},
client: client,
interval: interval,
quit: make(chan chan struct{}),
}
go r.loop()
return r, nil
}
// Stop stops the Docker registry's event subscriber.
func (r *registry) Stop() {
ch := make(chan struct{})
r.quit <- ch
<-ch
}
func (r *registry) loop() {
for {
// NB listenForEvents blocks.
// Returning false means we should exit.
if !r.listenForEvents() {
return
}
// Sleep here so we don't hammer the
// logs if docker is down
time.Sleep(r.interval)
}
}
func (r *registry) listenForEvents() bool {
// First we empty the store lists.
// This ensure any containers that went away inbetween calls to
// listenForEvents don't hang around.
r.reset()
// Next, start listening for events. We do this before fetching
// the list of containers so we don't miss containers created
// after listing but before listening for events.
events := make(chan *docker_client.APIEvents)
if err := r.client.AddEventListener(events); err != nil {
log.Printf("docker registry: %s", err)
return true
}
defer func() {
if err := r.client.RemoveEventListener(events); err != nil {
log.Printf("docker registry: %s", err)
}
}()
if err := r.updateContainers(); err != nil {
log.Printf("docker registry: %s", err)
return true
}
if err := r.updateImages(); err != nil {
log.Printf("docker registry: %s", err)
return true
}
otherUpdates := time.Tick(r.interval)
for {
select {
case event := <-events:
r.handleEvent(event)
case <-otherUpdates:
if err := r.updateImages(); err != nil {
log.Printf("docker registry: %s", err)
return true
}
case ch := <-r.quit:
r.Lock()
defer r.Unlock()
for _, c := range r.containers {
c.StopGatheringStats()
}
close(ch)
return false
}
}
}
func (r *registry) reset() {
r.Lock()
defer r.Unlock()
for _, c := range r.containers {
c.StopGatheringStats()
}
r.containers = map[string]Container{}
r.containersByPID = map[int]Container{}
r.images = map[string]*docker_client.APIImages{}
}
func (r *registry) updateContainers() error {
apiContainers, err := r.client.ListContainers(docker_client.ListContainersOptions{All: true})
if err != nil {
return err
}
for _, apiContainer := range apiContainers {
if err := r.addContainer(apiContainer.ID); err != nil {
return err
}
}
return nil
}
func (r *registry) updateImages() error {
images, err := r.client.ListImages(docker_client.ListImagesOptions{})
if err != nil {
return err
}
r.Lock()
defer r.Unlock()
for i := range images {
image := &images[i]
r.images[image.ID] = image
}
return nil
}
func (r *registry) handleEvent(event *docker_client.APIEvents) {
switch event.Status {
case DieEvent:
containerID := event.ID
r.removeContainer(containerID)
case StartEvent:
containerID := event.ID
if err := r.addContainer(containerID); err != nil {
log.Printf("docker registry: %s", err)
}
}
}
func (r *registry) addContainer(containerID string) error {
dockerContainer, err := r.client.InspectContainer(containerID)
if err != nil {
// Don't spam the logs if the container was short lived
if _, ok := err.(*docker_client.NoSuchContainer); ok {
return nil
}
return err
}
if !dockerContainer.State.Running {
// We get events late, and the containers sometimes have already
// stopped. Not an error, so don't return it.
return nil
}
r.Lock()
defer r.Unlock()
c := NewContainerStub(dockerContainer)
r.containers[containerID] = c
r.containersByPID[dockerContainer.State.Pid] = c
return c.StartGatheringStats()
}
func (r *registry) removeContainer(containerID string) {
r.Lock()
defer r.Unlock()
container, ok := r.containers[containerID]
if !ok {
return
}
delete(r.containers, containerID)
delete(r.containersByPID, container.PID())
container.StopGatheringStats()
}
// LockedPIDLookup runs f under a read lock, and gives f a function for
// use doing pid->container lookups.
func (r *registry) LockedPIDLookup(f func(func(int) Container)) {
r.RLock()
defer r.RUnlock()
lookup := func(pid int) Container {
return r.containersByPID[pid]
}
f(lookup)
}
// WalkContainers runs f on every running containers the registry knows of.
func (r *registry) WalkContainers(f func(Container)) {
r.RLock()
defer r.RUnlock()
for _, container := range r.containers {
f(container)
}
}
// WalkImages runs f on every image of running containers the registry
// knows of. f may be run on the same image more than once.
func (r *registry) WalkImages(f func(*docker_client.APIImages)) {
r.RLock()
defer r.RUnlock()
// Loop over containers so we only emit images for running containers.
for _, container := range r.containers {
image, ok := r.images[container.Image()]
if ok {
f(image)
}
}
}
+241
View File
@@ -0,0 +1,241 @@
package docker_test
import (
"reflect"
"runtime"
"sort"
"sync"
"testing"
"time"
client "github.com/fsouza/go-dockerclient"
"github.com/weaveworks/scope/probe/docker"
"github.com/weaveworks/scope/report"
"github.com/weaveworks/scope/test"
)
type mockContainer struct {
c *client.Container
}
func (c *mockContainer) ID() string {
return c.c.ID
}
func (c *mockContainer) PID() int {
return c.c.State.Pid
}
func (c *mockContainer) Image() string {
return c.c.Image
}
func (c *mockContainer) StartGatheringStats() error {
return nil
}
func (c *mockContainer) StopGatheringStats() {}
func (c *mockContainer) GetNodeMetadata() report.NodeMetadata {
return report.NodeMetadata{
docker.ContainerID: c.c.ID,
docker.ContainerName: c.c.Name,
docker.ImageID: c.c.Image,
}
}
type mockDockerClient struct {
sync.Mutex
apiContainers []client.APIContainers
containers map[string]*client.Container
apiImages []client.APIImages
events []chan<- *client.APIEvents
}
func (m *mockDockerClient) ListContainers(client.ListContainersOptions) ([]client.APIContainers, error) {
return m.apiContainers, nil
}
func (m *mockDockerClient) InspectContainer(id string) (*client.Container, error) {
return m.containers[id], nil
}
func (m *mockDockerClient) ListImages(client.ListImagesOptions) ([]client.APIImages, error) {
m.Lock()
defer m.Unlock()
return m.apiImages, nil
}
func (m *mockDockerClient) AddEventListener(events chan<- *client.APIEvents) error {
m.Lock()
defer m.Unlock()
m.events = append(m.events, events)
return nil
}
func (m *mockDockerClient) RemoveEventListener(events chan *client.APIEvents) error {
m.Lock()
defer m.Unlock()
for i, c := range m.events {
if c == events {
m.events = append(m.events[:i], m.events[i+1:]...)
}
}
return nil
}
func (m *mockDockerClient) send(event *client.APIEvents) {
m.Lock()
defer m.Unlock()
for _, c := range m.events {
c <- event
}
}
var (
container1 = &client.Container{
ID: "ping",
Name: "pong",
Image: "baz",
State: client.State{Pid: 1, Running: true},
}
container2 = &client.Container{
ID: "wiff",
Name: "waff",
Image: "baz",
State: client.State{Pid: 1, Running: true},
}
apiContainer1 = client.APIContainers{ID: "ping"}
apiImage1 = client.APIImages{ID: "baz", RepoTags: []string{"bang", "not-chosen"}}
mockClient = mockDockerClient{
apiContainers: []client.APIContainers{apiContainer1},
containers: map[string]*client.Container{"ping": container1},
apiImages: []client.APIImages{apiImage1},
}
)
func setupStubs(mdc *mockDockerClient, f func()) {
oldDockerClient, oldNewContainer := docker.NewDockerClientStub, docker.NewContainerStub
defer func() { docker.NewDockerClientStub, docker.NewContainerStub = oldDockerClient, oldNewContainer }()
docker.NewDockerClientStub = func(endpoint string) (docker.Client, error) {
return mdc, nil
}
docker.NewContainerStub = func(c *client.Container) docker.Container {
return &mockContainer{c}
}
f()
}
type containers []docker.Container
func (c containers) Len() int { return len(c) }
func (c containers) Swap(i, j int) { c[i], c[j] = c[j], c[i] }
func (c containers) Less(i, j int) bool { return c[i].ID() < c[j].ID() }
func allContainers(r docker.Registry) []docker.Container {
result := []docker.Container{}
r.WalkContainers(func(c docker.Container) {
result = append(result, c)
})
sort.Sort(containers(result))
return result
}
func allImages(r docker.Registry) []*client.APIImages {
result := []*client.APIImages{}
r.WalkImages(func(i *client.APIImages) {
result = append(result, i)
})
return result
}
func TestRegistry(t *testing.T) {
mdc := mockClient // take a copy
setupStubs(&mdc, func() {
registry, _ := docker.NewRegistry(10 * time.Second)
defer registry.Stop()
runtime.Gosched()
{
have := allContainers(registry)
want := []docker.Container{&mockContainer{container1}}
if !reflect.DeepEqual(want, have) {
t.Errorf("%s", test.Diff(want, have))
}
}
{
have := allImages(registry)
want := []*client.APIImages{&apiImage1}
if !reflect.DeepEqual(want, have) {
t.Errorf("%s", test.Diff(want, have))
}
}
})
}
func TestRegistryEvents(t *testing.T) {
mdc := mockClient // take a copy
setupStubs(&mdc, func() {
registry, _ := docker.NewRegistry(10 * time.Second)
defer registry.Stop()
runtime.Gosched()
{
mdc.Lock()
mdc.containers["wiff"] = container2
mdc.Unlock()
mdc.send(&client.APIEvents{Status: docker.StartEvent, ID: "wiff"})
runtime.Gosched()
have := allContainers(registry)
want := []docker.Container{&mockContainer{container1}, &mockContainer{container2}}
if !reflect.DeepEqual(want, have) {
t.Errorf("%s", test.Diff(want, have))
}
}
{
mdc.Lock()
delete(mdc.containers, "wiff")
mdc.Unlock()
mdc.send(&client.APIEvents{Status: docker.DieEvent, ID: "wiff"})
runtime.Gosched()
have := allContainers(registry)
want := []docker.Container{&mockContainer{container1}}
if !reflect.DeepEqual(want, have) {
t.Errorf("%s", test.Diff(want, have))
}
}
{
mdc.Lock()
delete(mdc.containers, "ping")
mdc.Unlock()
mdc.send(&client.APIEvents{Status: docker.DieEvent, ID: "ping"})
runtime.Gosched()
have := allContainers(registry)
want := []docker.Container{}
if !reflect.DeepEqual(want, have) {
t.Errorf("%s", test.Diff(want, have))
}
}
{
mdc.send(&client.APIEvents{Status: docker.DieEvent, ID: "doesntexist"})
runtime.Gosched()
have := allContainers(registry)
want := []docker.Container{}
if !reflect.DeepEqual(want, have) {
t.Errorf("%s", test.Diff(want, have))
}
}
})
}
+66
View File
@@ -0,0 +1,66 @@
package docker
import (
docker_client "github.com/fsouza/go-dockerclient"
"github.com/weaveworks/scope/report"
)
// Keys for use in NodeMetadata
const (
ContainerName = "docker_container_name"
ImageID = "docker_image_id"
ImageName = "docker_image_name"
)
// Reporter generate Reports containing Container and ContainerImage topologies
type Reporter struct {
registry Registry
scope string
}
// NewReporter makes a new Reporter
func NewReporter(registry Registry, scope string) *Reporter {
return &Reporter{
registry: registry,
scope: scope,
}
}
// Report generates a Report containing Container and ContainerImage topologies
func (r *Reporter) Report() report.Report {
result := report.MakeReport()
result.Container.Merge(r.containerTopology())
result.ContainerImage.Merge(r.containerImageTopology())
return result
}
func (r *Reporter) containerTopology() report.Topology {
result := report.NewTopology()
r.registry.WalkContainers(func(c Container) {
nodeID := report.MakeContainerNodeID(r.scope, c.ID())
result.NodeMetadatas[nodeID] = c.GetNodeMetadata()
})
return result
}
func (r *Reporter) containerImageTopology() report.Topology {
result := report.NewTopology()
r.registry.WalkImages(func(image *docker_client.APIImages) {
nmd := report.NodeMetadata{
ImageID: image.ID,
}
if len(image.RepoTags) > 0 {
nmd[ImageName] = image.RepoTags[0]
}
nodeID := report.MakeContainerNodeID(r.scope, image.ID)
result.NodeMetadatas[nodeID] = nmd
})
return result
}
+79
View File
@@ -0,0 +1,79 @@
package docker_test
import (
"reflect"
"testing"
client "github.com/fsouza/go-dockerclient"
"github.com/weaveworks/scope/probe/docker"
"github.com/weaveworks/scope/report"
"github.com/weaveworks/scope/test"
)
type mockRegistry struct {
containersByPID map[int]docker.Container
images map[string]*client.APIImages
}
func (r *mockRegistry) Stop() {}
func (r *mockRegistry) LockedPIDLookup(f func(func(int) docker.Container)) {
f(func(pid int) docker.Container {
return r.containersByPID[pid]
})
}
func (r *mockRegistry) WalkContainers(f func(docker.Container)) {
for _, c := range r.containersByPID {
f(c)
}
}
func (r *mockRegistry) WalkImages(f func(*client.APIImages)) {
for _, i := range r.images {
f(i)
}
}
var (
mockRegistryInstance = &mockRegistry{
containersByPID: map[int]docker.Container{
1: &mockContainer{container1},
},
images: map[string]*client.APIImages{
"baz": &apiImage1,
},
}
)
func TestReporter(t *testing.T) {
want := report.MakeReport()
want.Container = report.Topology{
Adjacency: report.Adjacency{},
EdgeMetadatas: report.EdgeMetadatas{},
NodeMetadatas: report.NodeMetadatas{
report.MakeContainerNodeID("", "ping"): report.NodeMetadata{
docker.ContainerID: "ping",
docker.ContainerName: "pong",
docker.ImageID: "baz",
},
},
}
want.ContainerImage = report.Topology{
Adjacency: report.Adjacency{},
EdgeMetadatas: report.EdgeMetadatas{},
NodeMetadatas: report.NodeMetadatas{
report.MakeContainerNodeID("", "baz"): report.NodeMetadata{
docker.ImageID: "baz",
docker.ImageName: "bang",
},
},
}
reporter := docker.NewReporter(mockRegistryInstance, "")
have := reporter.Report()
if !reflect.DeepEqual(want, have) {
t.Errorf("%s", test.Diff(want, have))
}
}
+87
View File
@@ -0,0 +1,87 @@
package docker
import (
"strconv"
"github.com/weaveworks/scope/probe/tag"
"github.com/weaveworks/scope/report"
)
// These constants are keys used in node metadata
// TODO: use these constants in report/{mapping.go, detailed_node.go} - pending some circular references
const (
ContainerID = "docker_container_id"
)
// These vars are exported for testing.
var (
NewPIDTreeStub = tag.NewPIDTree
)
// Tagger is a tagger that tags Docker container information to process
// nodes that have a PID.
type Tagger struct {
procRoot string
registry Registry
}
// NewTagger returns a usable Tagger.
func NewTagger(registry Registry, procRoot string) *Tagger {
return &Tagger{
registry: registry,
procRoot: procRoot,
}
}
// Tag implements Tagger.
func (t *Tagger) Tag(r report.Report) (report.Report, error) {
pidTree, err := NewPIDTreeStub(t.procRoot)
if err != nil {
return report.MakeReport(), err
}
t.tag(pidTree, &r.Process)
return r, nil
}
func (t *Tagger) tag(pidTree tag.PIDTree, topology *report.Topology) {
for nodeID, nodeMetadata := range topology.NodeMetadatas {
pidStr, ok := nodeMetadata["pid"]
if !ok {
continue
}
pid, err := strconv.ParseUint(pidStr, 10, 64)
if err != nil {
continue
}
var (
c Container
candidate = int(pid)
)
t.registry.LockedPIDLookup(func(lookup func(int) Container) {
for {
c = lookup(candidate)
if c != nil {
break
}
candidate, err = pidTree.GetParent(candidate)
if err != nil {
break
}
}
})
if c == nil {
continue
}
md := report.NodeMetadata{
ContainerID: c.ID(),
}
topology.NodeMetadatas[nodeID].Merge(md)
}
}
+60
View File
@@ -0,0 +1,60 @@
package docker_test
import (
"fmt"
"reflect"
"testing"
"github.com/weaveworks/scope/probe/docker"
"github.com/weaveworks/scope/probe/tag"
"github.com/weaveworks/scope/report"
"github.com/weaveworks/scope/test"
)
type mockPIDTree struct {
parents map[int]int
}
func (m *mockPIDTree) GetParent(pid int) (int, error) {
parent, ok := m.parents[pid]
if !ok {
return -1, fmt.Errorf("Not found %d", pid)
}
return parent, nil
}
func (m *mockPIDTree) ProcessTopology(hostID string) report.Topology {
panic("")
}
func TestTagger(t *testing.T) {
oldPIDTree := docker.NewPIDTreeStub
defer func() { docker.NewPIDTreeStub = oldPIDTree }()
docker.NewPIDTreeStub = func(procRoot string) (tag.PIDTree, error) {
return &mockPIDTree{map[int]int{2: 1}}, nil
}
var (
pid1NodeID = report.MakeProcessNodeID("somehost.com", "1")
pid2NodeID = report.MakeProcessNodeID("somehost.com", "2")
wantNodeMetadata = report.NodeMetadata{docker.ContainerID: "ping"}
)
input := report.MakeReport()
input.Process.NodeMetadatas[pid1NodeID] = report.NodeMetadata{"pid": "1"}
input.Process.NodeMetadatas[pid2NodeID] = report.NodeMetadata{"pid": "2"}
want := report.MakeReport()
want.Process.NodeMetadatas[pid1NodeID] = report.NodeMetadata{"pid": "1"}.Merge(wantNodeMetadata)
want.Process.NodeMetadatas[pid2NodeID] = report.NodeMetadata{"pid": "2"}.Merge(wantNodeMetadata)
tagger := docker.NewTagger(mockRegistryInstance, "/irrelevant")
have, err := tagger.Tag(input)
if err != nil {
t.Errorf("%v", err)
}
if !reflect.DeepEqual(want, have) {
t.Errorf("%s", test.Diff(want, have))
}
}