From e6d9bcc1cb99d0da7c21a21f432359afe859c890 Mon Sep 17 00:00:00 2001 From: Lili Cosic Date: Fri, 20 Jul 2018 07:31:02 +0200 Subject: [PATCH] Add CRI probe When the probe.cri is enabled the CRI probe will be used to gather the container information via the CRI API. For now only the basic information is included in the generated report, those that we can get via the CRI ListContainersRequest. --- probe/cri/registry.go | 72 ++++++++++++++++++++++++++++++++++++++++++ probe/cri/reporter.go | 73 +++++++++++++++++++++++++++++++++++++++++++ prog/main.go | 7 +++++ prog/probe.go | 10 ++++++ 4 files changed, 162 insertions(+) create mode 100644 probe/cri/registry.go create mode 100644 probe/cri/reporter.go diff --git a/probe/cri/registry.go b/probe/cri/registry.go new file mode 100644 index 000000000..c346ec59b --- /dev/null +++ b/probe/cri/registry.go @@ -0,0 +1,72 @@ +package cri + +import ( + "fmt" + "net" + "net/url" + "time" + + "google.golang.org/grpc" + + client "github.com/weaveworks/scope/cri/runtime" +) + +const unixProtocol = "unix" + +func dial(addr string, timeout time.Duration) (net.Conn, error) { + return net.DialTimeout(unixProtocol, addr, timeout) +} + +func getAddressAndDialer(endpoint string) (string, func(addr string, timeout time.Duration) (net.Conn, error), error) { + protocol, addr, err := parseEndpointWithFallbackProtocol(endpoint, unixProtocol) + if err != nil { + return "", nil, err + } + if protocol != unixProtocol { + return "", nil, fmt.Errorf("endpoint was not unix socket %v", protocol) + } + + return addr, dial, nil +} + +func parseEndpointWithFallbackProtocol(endpoint string, fallbackProtocol string) (protocol string, addr string, err error) { + if protocol, addr, err = parseEndpoint(endpoint); err != nil && protocol == "" { + fallbackEndpoint := fallbackProtocol + "://" + endpoint + protocol, addr, err = parseEndpoint(fallbackEndpoint) + if err != nil { + return "", "", err + } + } + return +} + +func parseEndpoint(endpoint string) (string, string, error) { + u, err := url.Parse(endpoint) + if err != nil { + return "", "", err + } + + if u.Scheme == "tcp" { + return "tcp", u.Host, nil + } else if u.Scheme == "unix" { + return "unix", u.Path, nil + } else if u.Scheme == "" { + return "", "", fmt.Errorf("Using %q as endpoint is deprecated, please consider using full url format", endpoint) + } else { + return u.Scheme, "", fmt.Errorf("protocol %q not supported", u.Scheme) + } +} + +// NewCRIClient creates client to CRI. +func NewCRIClient(endpoint string) (client.RuntimeServiceClient, error) { + addr, dailer, err := getAddressAndDialer(endpoint) + if err != nil { + return nil, err + } + conn, err := grpc.Dial(addr, grpc.WithInsecure(), grpc.WithDialer(dailer)) + if err != nil { + return nil, err + } + + return client.NewRuntimeServiceClient(conn), nil +} diff --git a/probe/cri/reporter.go b/probe/cri/reporter.go new file mode 100644 index 000000000..dec1e2985 --- /dev/null +++ b/probe/cri/reporter.go @@ -0,0 +1,73 @@ +package cri + +import ( + "context" + "fmt" + + client "github.com/weaveworks/scope/cri/runtime" + "github.com/weaveworks/scope/probe/docker" + "github.com/weaveworks/scope/report" +) + +// Reporter generate Reports containing Container and ContainerImage topologies +type Reporter struct { + cri client.RuntimeServiceClient +} + +// NewReporter makes a new Reporter +func NewReporter(cri client.RuntimeServiceClient) *Reporter { + reporter := &Reporter{ + cri: cri, + } + + return reporter +} + +// Name of this reporter, for metrics gathering +func (Reporter) Name() string { return "CRI" } + +// Report generates a Report containing Container topologies +func (r *Reporter) Report() (report.Report, error) { + result := report.MakeReport() + containerTopol, err := r.containerTopology() + if err != nil { + return report.MakeReport(), err + } + + result.Container = result.Container.Merge(containerTopol) + return result, nil +} + +func (r *Reporter) containerTopology() (report.Topology, error) { + result := report.MakeTopology(). + WithMetadataTemplates(docker.ContainerImageMetadataTemplates). + WithTableTemplates(docker.ContainerImageTableTemplates) + + ctx := context.Background() + resp, err := r.cri.ListContainers(ctx, &client.ListContainersRequest{}) + if err != nil { + return result, err + } + + for _, c := range resp.Containers { + result.AddNode(getNode(c)) + } + + return result, nil +} + +func getNode(c *client.Container) report.Node { + result := report.MakeNodeWith(report.MakeContainerNodeID(c.Id), map[string]string{ + docker.ContainerName: c.Metadata.Name, + docker.ContainerID: c.Id, + docker.ContainerState: fmt.Sprintf("%v", c.State), + docker.ContainerRestartCount: fmt.Sprintf("%v", c.Metadata.Attempt), + docker.ImageID: c.ImageRef, + docker.ImageName: c.Image.Image, + }).WithParents(report.MakeSets(). + Add(report.ContainerImage, report.MakeStringSet(report.MakeContainerImageNodeID(c.ImageRef))), + ) + result = result.AddPrefixPropertyList(docker.LabelPrefix, c.Labels) + + return result +} diff --git a/prog/main.go b/prog/main.go index cb1f7ad00..cb6cb68f3 100644 --- a/prog/main.go +++ b/prog/main.go @@ -120,6 +120,9 @@ type probeFlags struct { dockerInterval time.Duration dockerBridge string + criEnabled bool + criEndpoint string + kubernetesEnabled bool kubernetesNodeName string kubernetesClientConfig kubernetes.ClientConfig @@ -304,6 +307,10 @@ func setupFlags(flags *flags) { flag.DurationVar(&flags.probe.dockerInterval, "probe.docker.interval", 10*time.Second, "how often to update Docker attributes") flag.StringVar(&flags.probe.dockerBridge, "probe.docker.bridge", "docker0", "the docker bridge name") + // CRI + flag.BoolVar(&flags.probe.criEnabled, "probe.cri", true, "collect CRI-related attributes for processes") + flag.StringVar(&flags.probe.criEndpoint, "probe.cri.endpoint", "unix///var/run/dockershim.sock", "The endpoint to connect to the CRI") + // K8s flag.BoolVar(&flags.probe.kubernetesEnabled, "probe.kubernetes", false, "collect kubernetes-related attributes for containers") flag.StringVar(&flags.probe.kubernetesClientConfig.Server, "probe.kubernetes.api", "", "The address and port of the Kubernetes API server (deprecated in favor of equivalent probe.kubernetes.server)") diff --git a/prog/probe.go b/prog/probe.go index fb989f760..cf8683943 100644 --- a/prog/probe.go +++ b/prog/probe.go @@ -26,6 +26,7 @@ import ( "github.com/weaveworks/scope/probe/appclient" "github.com/weaveworks/scope/probe/awsecs" "github.com/weaveworks/scope/probe/controls" + "github.com/weaveworks/scope/probe/cri" "github.com/weaveworks/scope/probe/docker" "github.com/weaveworks/scope/probe/endpoint" "github.com/weaveworks/scope/probe/host" @@ -253,6 +254,15 @@ func probeMain(flags probeFlags, targets []appclient.Target) { } } + if flags.criEnabled { + client, err := cri.NewCRIClient(flags.criEndpoint) + if err != nil { + log.Errorf("CRI: failed to start registry: %v", err) + } else { + p.AddReporter(cri.NewReporter(client)) + } + } + if flags.kubernetesEnabled { if client, err := kubernetes.NewClient(flags.kubernetesClientConfig); err == nil { defer client.Stop()