(hash string, dataset config.Dataset)
| 141 | } |
| 142 | |
| 143 | func (api *API) createDataset(hash string, dataset config.Dataset) error { |
| 144 | filename := filepath.Join(api.dataDir, dataset.ID+".db") |
| 145 | |
| 146 | // Does the file exist? If so, remove it. |
| 147 | _, err := os.Stat(filename) |
| 148 | if err == nil { |
| 149 | log.Printf("`%s` already exists; removing", filename) |
| 150 | os.Remove(filename) |
| 151 | } |
| 152 | |
| 153 | db, err := sql.Open("sqlite3", filename) |
| 154 | if err != nil { |
| 155 | return err |
| 156 | } |
| 157 | defer db.Close() |
| 158 | _, err = db.Exec("PRAGMA synchronous = OFF") |
| 159 | if err != nil { |
| 160 | return err |
| 161 | } |
| 162 | _, err = db.Exec("PRAGMA journal_mode = MEMORY") |
| 163 | if err != nil { |
| 164 | return err |
| 165 | } |
| 166 | _, err = db.Exec("PRAGMA cache_size = -2000000") |
| 167 | if err != nil { |
| 168 | return err |
| 169 | } |
| 170 | |
| 171 | log.Printf("querying `%s`", dataset.DataSource.ID) |
| 172 | err = api.fetchSingle(hash, db, dataset.DataSource) |
| 173 | if err != nil { |
| 174 | return fmt.Errorf("fetch single: %w", err) |
| 175 | } |
| 176 | |
| 177 | for _, join := range dataset.Joins { |
| 178 | log.Printf("querying `%s`", join.DataSource.ID) |
| 179 | err = api.fetchSingle(hash, db, join.DataSource) |
| 180 | if err != nil { |
| 181 | return fmt.Errorf("fetch single as part of join: %w", err) |
| 182 | } |
| 183 | } |
| 184 | |
| 185 | joinClauses := "" |
| 186 | for _, join := range dataset.Joins { |
| 187 | joinColumns := []string{} |
| 188 | for _, cols := range join.Columns { |
| 189 | joinColumns = append(joinColumns, fmt.Sprintf(`%s."%s" = %s."%s"`, dataset.DataSource.ID, cols.LeftColumn, join.DataSource.ID, cols.RightColumn)) |
| 190 | } |
| 191 | joinClauses += fmt.Sprintf(" %s %s ON %s", join.Type, join.DataSource.ID, strings.Join(joinColumns, " AND ")) |
| 192 | } |
| 193 | |
| 194 | log.Println("joining data") |
| 195 | joinQuery := fmt.Sprintf("CREATE TABLE %s AS SELECT * FROM %s %s", dataset.ID, dataset.DataSource.ID, joinClauses) |
| 196 | _, err = db.Exec(joinQuery) |
| 197 | return err |
| 198 | } |
| 199 | |
| 200 | func (api *API) fetchSingle(hash string, dest *sql.DB, dataSource *config.DataSource) error { |
no test coverage detected