From 651c93fecbee7f04547e28625e9daa19d0ad2dec Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Toma=C5=BE=20Jerman?= Date: Tue, 17 Nov 2020 15:56:11 +0100 Subject: [PATCH] Cleanup, rename, restructure core envoy bits --- pkg/envoy/builder.go | 36 +++++++++++++++++++++++++++--------- pkg/envoy/encoder.go | 37 +++++++++++++++++++++++-------------- pkg/envoy/graph.go | 12 ++++++------ pkg/envoy/node.go | 2 +- 4 files changed, 57 insertions(+), 30 deletions(-) diff --git a/pkg/envoy/builder.go b/pkg/envoy/builder.go index f652679ee..3ca8cfaa9 100644 --- a/pkg/envoy/builder.go +++ b/pkg/envoy/builder.go @@ -2,22 +2,39 @@ package envoy import ( "context" + "sync" "github.com/cortezaproject/corteza-server/pkg/envoy/resource" ) type ( graphBuilder struct { - pp []Processor + pp []Preparer } - // Some special bits - Processor interface { - Process(ctx context.Context, state *ExecState) error + // Rc is an alias for the ResourceState channel + Rc chan *ResourceState + // Ec is an alias for the error channel + Ec chan error + + Preparer interface { + Prepare(ctx context.Context, ee ...*ResourceState) error + } + + // Encoder encodes all resources provided by the Rc until nil is passed + // + // Encoding errors are passed via Ec. + Encoder interface { + Encode(ctx context.Context, wg *sync.WaitGroup, sChan Rc, eChan Ec) + } + + PrepareEncoder interface { + Preparer + Encoder } ) -func NewGraphBuilder(pp ...Processor) *graphBuilder { +func NewGraphBuilder(pp ...Preparer) *graphBuilder { return &graphBuilder{ pp: pp, } @@ -67,8 +84,9 @@ func (b *graphBuilder) Build(ctx context.Context, rr ...resource.Interface) (*gr g.invert() // Do any dep. related preprocessing - var state *ExecState + var state *ResourceState var err error + err = func() error { for { state, err = g.Next(ctx) @@ -83,7 +101,7 @@ func (b *graphBuilder) Build(ctx context.Context, rr ...resource.Interface) (*gr nState.MissingDeps = mMap[state.Res] for _, p := range b.pp { - err = p.Process(ctx, nState) + err = p.Prepare(ctx, nState) if err != nil { return err } @@ -91,12 +109,12 @@ func (b *graphBuilder) Build(ctx context.Context, rr ...resource.Interface) (*gr } }() - g.Relink() - if err != nil { return nil, err } + g.Relink() g.reset() + return g, nil } diff --git a/pkg/envoy/encoder.go b/pkg/envoy/encoder.go index 0db72f6b0..126ae2ee4 100644 --- a/pkg/envoy/encoder.go +++ b/pkg/envoy/encoder.go @@ -2,33 +2,42 @@ package envoy import ( "context" + "sync" ) type ( - Encoder interface { - Encode(ctx context.Context, ss ...*ExecState) error - } - Provider interface { - Next(ctx context.Context) (*ExecState, error) + Next(ctx context.Context) (*ResourceState, error) } ) +// @todo errors! func Encode(ctx context.Context, p Provider, ee ...Encoder) error { + ec := make(Ec) + rcc := make([]Rc, len(ee)) + + var wg sync.WaitGroup + wg.Add(len(ee)) + + for i, e := range ee { + rcc[i] = make(Rc) + go e.Encode(ctx, &wg, rcc[i], ec) + } + for { - state, err := p.Next(ctx) + rs, err := p.Next(ctx) if err != nil { return err - } else if state == nil { - return nil } - for _, e := range ee { - // Any dep conflicts should be handled by the Encoder - err = e.Encode(ctx, state) - if err != nil { - return err - } + for _, rc := range rcc { + rc <- rs + } + if rs == nil { + break } } + + wg.Wait() + return nil } diff --git a/pkg/envoy/graph.go b/pkg/envoy/graph.go index 45eabaae1..3c26df1a1 100644 --- a/pkg/envoy/graph.go +++ b/pkg/envoy/graph.go @@ -22,7 +22,7 @@ type ( conflicting nodeMap } - ExecState struct { + ResourceState struct { Res resource.Interface MissingDeps resource.RefSet Conflicting bool @@ -112,15 +112,15 @@ func (g *graph) Relink() { } } -func (g *graph) Next(ctx context.Context) (s *ExecState, err error) { +func (g *graph) Next(ctx context.Context) (s *ResourceState, err error) { return g.next(ctx, nil) } -func (g *graph) NextOf(ctx context.Context, resTrype string) (s *ExecState, err error) { +func (g *graph) NextOf(ctx context.Context, resTrype string) (s *ResourceState, err error) { return g.next(ctx, &iterFilter{resourceType: resTrype}) } -func (g *graph) next(ctx context.Context, f *iterFilter) (s *ExecState, err error) { +func (g *graph) next(ctx context.Context, f *iterFilter) (s *ResourceState, err error) { upNN := g.removeProcessed(g.nodes()) if f != nil { @@ -280,8 +280,8 @@ func (g *graph) nodeResource(nn ...*node) []resource.Interface { return rr } -func (g *graph) prepExecState(n *node, conflicting bool) *ExecState { - return &ExecState{ +func (g *graph) prepExecState(n *node, conflicting bool) *ResourceState { + return &ResourceState{ Res: n.res, DepResources: g.nodeResource(g.childNodes(n)...), ParentResources: g.nodeResource(g.parentNodes(n)...), diff --git a/pkg/envoy/node.go b/pkg/envoy/node.go index 3e84fb402..a746494dc 100644 --- a/pkg/envoy/node.go +++ b/pkg/envoy/node.go @@ -40,7 +40,7 @@ func (nn nodeSet) filter(f func(n *node) bool) nodeSet { func (nn nodeSet) findByRef(ref *resource.Ref) *node { for _, n := range nn { if n.res.ResourceType() == ref.ResourceType { - if n.res.Identifiers().HasAny(ref.Identifiers.StringSlice()...) { + if n.res.Identifiers().HasAny(ref.Identifiers) { return n } }