Merge pull request #3850 from weaveworks/remove-unmerged-nodes

Remove partially merged nodes from deltas
This commit is contained in:
Bryan Boreham
2021-04-23 12:05:52 +01:00
committed by GitHub
6 changed files with 40 additions and 2 deletions

View File

@@ -479,6 +479,7 @@ func (r *Registry) makeTopologyList(rep Reporter) CtxHandlerFunc {
respondWith(ctx, w, http.StatusInternalServerError, err)
return
}
report.UnsafeRemovePartMergedNodes(ctx)
respondWith(ctx, w, http.StatusOK, r.renderTopologies(ctx, report, req))
}
}
@@ -579,6 +580,7 @@ func (r *Registry) captureRenderer(rep Reporter, f rendererHandler) CtxHandlerFu
respondWith(ctx, w, http.StatusInternalServerError, err)
return
}
rpt.UnsafeRemovePartMergedNodes(ctx)
req.ParseForm()
renderer, filter, err := r.RendererForTopology(topologyID, req.Form, rpt)
if err != nil {

View File

@@ -182,6 +182,7 @@ func (wc *websocketState) update(ctx context.Context) error {
if err != nil {
return errors.Wrap(err, "Error generating report")
}
re.UnsafeRemovePartMergedNodes(ctx)
renderer, filter, err := topologyRegistry.RendererForTopology(wc.topologyID, wc.values, re)
if err != nil {
return errors.Wrap(err, "Error generating report")

View File

@@ -87,6 +87,7 @@ func getReport(b *testing.B) report.Report {
func benchmarkRender(b *testing.B, f func(report.Report)) {
r := getReport(b)
r.UnsafeRemovePartMergedNodes(context.Background())
b.ResetTimer()
for i := 0; i < b.N; i++ {
b.StopTimer()

View File

@@ -23,10 +23,10 @@ func TestCollector(t *testing.T) {
defer mtime.NowReset()
r1 := report.MakeReport()
r1.Endpoint.AddNode(report.MakeNode("foo"))
r1.Endpoint.AddNode(report.MakeNode("foo").WithTopology("bar"))
r2 := report.MakeReport()
r2.Endpoint.AddNode(report.MakeNode("foo"))
r2.Endpoint.AddNode(report.MakeNode("foo").WithTopology("bar"))
have, err := c.Report(ctx, mtime.Now())
if err != nil {
@@ -52,6 +52,7 @@ func TestCollector(t *testing.T) {
merged := report.MakeReport()
merged.UnsafeMerge(r1)
merged.UnsafeMerge(r2)
merged.UnsafeRemovePartMergedNodes(context.Background())
have, err = c.Report(ctx, mtime.Now())
if err != nil {
t.Error(err)

View File

@@ -236,3 +236,10 @@ func (n *Node) UnsafeUnMerge(other Node) bool {
// metrics don't overlap so just check if we have any
return remove && len(n.Metrics) == 0
}
// If a node is removed from source between two full reports, then we
// might only have a delta of its last state. Detect that from a blank topology,
// which cannot arise in a properly-formed Node.
func (n *Node) isPartMerged() bool {
return n.Topology == ""
}

View File

@@ -1,11 +1,14 @@
package report
import (
"context"
"fmt"
"math/rand"
"strings"
"time"
opentracing "github.com/opentracing/opentracing-go"
"github.com/weaveworks/scope/common/xfer"
)
@@ -363,6 +366,29 @@ func (r *Report) UnsafeUnMerge(other Report) {
})
}
// UnsafeRemovePartMergedNodes removes nodes that have not fully re-merged.
// E.g. if a node is removed from source between two full reports, then we
// might only have a delta of its last state. Remove that from the set.
// The original is modified.
func (r *Report) UnsafeRemovePartMergedNodes(ctx context.Context) {
dropped := map[string]int{}
r.WalkNamedTopologies(func(name string, t *Topology) {
for k, v := range t.Nodes {
if v.isPartMerged() {
delete(t.Nodes, k)
dropped[name]++
}
}
})
if span := opentracing.SpanFromContext(ctx); span != nil && len(dropped) > 0 {
msg := ""
for name, count := range dropped {
msg += fmt.Sprintf("%s: %d, ", name, count)
}
span.LogKV("dropped-part-merged", msg)
}
}
// WalkTopologies iterates through the Topologies of the report,
// potentially modifying them
func (r *Report) WalkTopologies(f func(*Topology)) {