| 182 | } |
| 183 | |
| 184 | func (m *Upsert) insertRows(rows [][]*rel.ValueColumn) (int64, error) { |
| 185 | for i, row := range rows { |
| 186 | select { |
| 187 | case <-m.SigChan(): |
| 188 | if i == 0 { |
| 189 | return 0, nil |
| 190 | } |
| 191 | return int64(i) - 1, nil |
| 192 | default: |
| 193 | vals := make([]driver.Value, len(row)) |
| 194 | for x, val := range row { |
| 195 | if val.Expr != nil { |
| 196 | exprVal, ok := vm.Eval(nil, val.Expr) |
| 197 | if !ok { |
| 198 | u.Errorf("Could not evaluate: %v", val.Expr) |
| 199 | return 0, fmt.Errorf("Could not evaluate expression: %v", val.Expr) |
| 200 | } |
| 201 | vals[x] = exprVal.Value() |
| 202 | } else { |
| 203 | vals[x] = val.Value.Value() |
| 204 | } |
| 205 | } |
| 206 | |
| 207 | if _, err := m.db.Put(m.Ctx.Context, nil, vals); err != nil { |
| 208 | u.Errorf("Could not put values: fordb T:%T %v", m.db, err) |
| 209 | return 0, err |
| 210 | } |
| 211 | } |
| 212 | } |
| 213 | return int64(len(rows)), nil |
| 214 | } |
| 215 | |
| 216 | func (m *DeletionTask) Close() error { |
| 217 | m.Lock() |