From 8454bcb3907f4e444a9a06ee109acacbd0a430d7 Mon Sep 17 00:00:00 2001 From: Matthew Nibecker Date: Wed, 26 Aug 2026 11:40:22 -0700 Subject: [PATCH] optimizer: make optimizeSourcePaths work with sub queries Fixes #6206 --- compiler/dag/op.go | 31 +++++++++++++++++++++++++++++++ compiler/optimizer/optimizer.go | 7 +++---- 2 files changed, 34 insertions(+), 4 deletions(-) diff --git a/compiler/dag/op.go b/compiler/dag/op.go index f10915684d..648419a93d 100644 --- a/compiler/dag/op.go +++ b/compiler/dag/op.go @@ -410,3 +410,34 @@ func WalkT[T any](v reflect.Value, post func(T) T) { } } } + +func WalkTWithError[T any](v reflect.Value, post func(T) (T, error)) error { + switch v.Kind() { + case reflect.Array, reflect.Slice: + for i := range v.Len() { + if err := WalkTWithError(v.Index(i), post); err != nil { + return err + } + } + case reflect.Interface, reflect.Pointer: + if err := WalkTWithError(v.Elem(), post); err != nil { + return err + } + case reflect.Struct: + for _, field := range v.Fields() { + if err := WalkTWithError(field, post); err != nil { + return err + } + } + } + if v.CanSet() { + if t, ok := v.Interface().(T); ok { + r, err := post(t) + if err != nil { + return err + } + v.Set(reflect.ValueOf(r)) + } + } + return nil +} diff --git a/compiler/optimizer/optimizer.go b/compiler/optimizer/optimizer.go index 056bb19c92..7dc4471962 100644 --- a/compiler/optimizer/optimizer.go +++ b/compiler/optimizer/optimizer.go @@ -132,8 +132,7 @@ func (o *Optimizer) Optimize(main *dag.Main) error { seq = replaceSortAndHeadOrTailWithTop(seq) o.optimizeParallels(seq) seq = mergeFilters(seq) - seq, err := o.optimizeSourcePaths(seq) - if err != nil { + if err := o.optimizeSourcePaths(&seq); err != nil { return err } seq = removePassOps(seq) @@ -198,8 +197,8 @@ func (o *Optimizer) OptimizeDeleter(main *dag.Main, replicas int) error { return nil } -func (o *Optimizer) optimizeSourcePaths(seq dag.Seq) (dag.Seq, error) { - return walkEntries(seq, func(seq dag.Seq) (dag.Seq, error) { +func (o *Optimizer) optimizeSourcePaths(seq *dag.Seq) error { + return dag.WalkTWithError(reflect.ValueOf(seq), func(seq dag.Seq) (dag.Seq, error) { if len(seq) == 0 { return nil, errors.New("internal error: optimizer encountered empty sequential operator") }