MCPcopy Create free account
hub / github.com/crossjoin-io/crossjoin / fetchSingle

Method fetchSingle

api/datasets.go:200–318  ·  view source on GitHub ↗
(hash string, dest *sql.DB, dataSource *config.DataSource)

Source from the content-addressed store, hash-verified

198}
199
200func (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)

Callers 1

createDatasetMethod · 0.95

Calls 2

ReadDataConnectionMethod · 0.95
readFileMethod · 0.95

Tested by

no test coverage detected