| 85 | } |
| 86 | |
| 87 | func whereFilter(filter expr.Node, task TaskRunner, cols map[string]int) MessageHandler { |
| 88 | out := task.MessageOut() |
| 89 | |
| 90 | //u.Debugf("prepare filter %s", filter) |
| 91 | return func(ctx *plan.Context, msg schema.Message) bool { |
| 92 | |
| 93 | var filterValue value.Value |
| 94 | var ok bool |
| 95 | //u.Debugf("WHERE: T:%T body%#v", msg, msg.Body()) |
| 96 | switch mt := msg.(type) { |
| 97 | case *datasource.SqlDriverMessage: |
| 98 | //u.Debugf("WHERE: T:%T vals:%#v", msg, mt.Vals) |
| 99 | //u.Debugf("cols: %#v", cols) |
| 100 | msgReader := mt.ToMsgMap(cols) |
| 101 | filterValue, ok = vm.Eval(msgReader, filter) |
| 102 | case *datasource.SqlDriverMessageMap: |
| 103 | filterValue, ok = vm.Eval(mt, filter) |
| 104 | if !ok { |
| 105 | u.Warnf("wtf %s %#v", filter, mt) |
| 106 | } |
| 107 | //u.Debugf("WHERE: result:%v T:%T \n\trow:%#v \n\tvals:%#v", filterValue, msg, mt, mt.Values()) |
| 108 | //u.Debugf("cols: %#v", cols) |
| 109 | default: |
| 110 | if msgReader, isContextReader := msg.(expr.ContextReader); isContextReader { |
| 111 | filterValue, ok = vm.Eval(msgReader, filter) |
| 112 | if !ok { |
| 113 | u.Warnf("wat? %v filterval:%#v expr: %s", filter.String(), filterValue, filter) |
| 114 | } |
| 115 | } else { |
| 116 | u.Errorf("could not convert to message reader: %T", msg) |
| 117 | } |
| 118 | } |
| 119 | //u.Debugf("msg: %#v", msgReader) |
| 120 | //u.Infof("evaluating: ok?%v result=%v filter expr: '%s'", ok, filterValue.ToString(), filter.String()) |
| 121 | if !ok { |
| 122 | u.Debugf("could not evaluate: %T %#v", msg, msg) |
| 123 | return false |
| 124 | } |
| 125 | switch valTyped := filterValue.(type) { |
| 126 | case value.BoolValue: |
| 127 | if valTyped.Val() == false { |
| 128 | //u.Debugf("Filtering out: T:%T v:%#v", valTyped, valTyped) |
| 129 | return true |
| 130 | } |
| 131 | case nil: |
| 132 | return false |
| 133 | default: |
| 134 | if valTyped.Nil() { |
| 135 | return false |
| 136 | } |
| 137 | } |
| 138 | |
| 139 | //u.Debugf("about to send from where to forward: %#v", msg) |
| 140 | select { |
| 141 | case out <- msg: |
| 142 | return true |
| 143 | case <-task.SigChan(): |
| 144 | return false |