()
| 52 | } |
| 53 | |
| 54 | func (m *JoinKey) Run() error { |
| 55 | defer m.Ctx.Recover() |
| 56 | defer close(m.msgOutCh) |
| 57 | |
| 58 | outCh := m.MessageOut() |
| 59 | inCh := m.MessageIn() |
| 60 | joinNodes := m.p.Source.Stmt.JoinNodes() |
| 61 | |
| 62 | for { |
| 63 | |
| 64 | select { |
| 65 | case <-m.SigChan(): |
| 66 | //u.Debugf("got signal quit") |
| 67 | return nil |
| 68 | case msg, ok := <-inCh: |
| 69 | if !ok { |
| 70 | //u.Debugf("NICE, got msg shutdown") |
| 71 | return nil |
| 72 | } |
| 73 | |
| 74 | //u.Infof("In joinkey msg %#v", msg) |
| 75 | msgTypeSwitch: |
| 76 | switch mt := msg.(type) { |
| 77 | case *datasource.SqlDriverMessageMap: |
| 78 | vals := make([]string, len(joinNodes)) |
| 79 | for i, node := range joinNodes { |
| 80 | joinVal, ok := vm.Eval(mt, node) |
| 81 | //u.Debugf("evaluating: ok?%v T:%T result=%v node '%v'", ok, joinVal, joinVal.ToString(), node.String()) |
| 82 | if !ok { |
| 83 | u.Errorf("could not evaluate: %T %#v %v", joinVal, joinVal, msg) |
| 84 | break msgTypeSwitch |
| 85 | } |
| 86 | vals[i] = joinVal.ToString() |
| 87 | } |
| 88 | //u.Infof("joinkey: %v row:%v", vals, mt) |
| 89 | key := strings.Join(vals, string(byte(0))) |
| 90 | mt.SetKeyHashed(key) |
| 91 | outCh <- mt |
| 92 | default: |
| 93 | return fmt.Errorf("To use JoinKey must use SqlDriverMessageMap but got %T", msg) |
| 94 | } |
| 95 | |
| 96 | } |
| 97 | } |
| 98 | } |
| 99 | |
| 100 | // Scans 2 source tasks for rows, evaluate keys, use for join |
| 101 | // |
nothing calls this directly
no test coverage detected