(hash string, dest *sql.DB, dataSource *config.DataSource)
| 198 | } |
| 199 | |
| 200 | func (api *API) fetchSingle(hash string, dest *sql.DB, dataSource *config.DataSource) error { |
| 201 | dataConnection, err := api.ReadDataConnection(hash, dataSource.DataConnection) |
| 202 | if err != nil { |
| 203 | return err |
| 204 | } |
| 205 | switch dataConnection.Type { |
| 206 | case "csv": |
| 207 | f, err := api.readFile(dataConnection.Path) |
| 208 | if err != nil { |
| 209 | return err |
| 210 | } |
| 211 | r := csv.NewReader(f) |
| 212 | firstLine, err := r.Read() |
| 213 | if err != nil { |
| 214 | return err |
| 215 | } |
| 216 | columns := firstLine |
| 217 | for i := range columns { |
| 218 | columns[i] = strconv.Quote(columns[i]) |
| 219 | } |
| 220 | |
| 221 | _, err = dest.Exec(fmt.Sprintf("CREATE TABLE %s (%s)", dataSource.ID, strings.Join(columns, ","))) |
| 222 | if err != nil { |
| 223 | return err |
| 224 | } |
| 225 | |
| 226 | params := []string{} |
| 227 | for i := range columns { |
| 228 | params = append(params, fmt.Sprintf("$%d", i+1)) |
| 229 | } |
| 230 | stmt, err := dest.Prepare(fmt.Sprintf("INSERT INTO %s VALUES (%s)", dataSource.ID, strings.Join(params, ","))) |
| 231 | if err != nil { |
| 232 | return err |
| 233 | } |
| 234 | defer stmt.Close() |
| 235 | |
| 236 | for { |
| 237 | record, err := r.Read() |
| 238 | if err != nil { |
| 239 | if err == io.EOF { |
| 240 | return nil |
| 241 | } |
| 242 | return err |
| 243 | } |
| 244 | if len(record) != len(columns) { |
| 245 | return errors.New("inconsistent number of fields") |
| 246 | } |
| 247 | values := make([]interface{}, len(record)) |
| 248 | for i := range record { |
| 249 | values[i] = record[i] |
| 250 | } |
| 251 | _, err = stmt.Exec(values...) |
| 252 | if err != nil { |
| 253 | return err |
| 254 | } |
| 255 | } |
| 256 | case "postgres": |
| 257 | db, err := sql.Open(dataConnection.Type, dataConnection.ConnectionString) |
no test coverage detected