Skip to content

Commit 350a219

Browse files
authored
Merge pull request #28 from gigapi/feat/s3
Feat: S3 support
2 parents 830bef7 + b31b30c commit 350a219

4 files changed

Lines changed: 137 additions & 17 deletions

File tree

go.mod

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ go 1.24.2
44

55
require (
66
github.com/apache/arrow/go/v14 v14.0.2
7-
github.com/gigapi/gigapi-config v0.0.8
7+
github.com/gigapi/gigapi-config v0.0.9
88
github.com/gigapi/gigapi/v2 v2.0.13
99
github.com/gigapi/metadata v0.0.4
1010
github.com/marcboeker/go-duckdb/v2 v2.2.0

go.sum

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -38,8 +38,8 @@ github.com/frankban/quicktest v1.14.6 h1:7Xjx+VpznH+oBnejlPUj8oUpdxnVs4f8XU8WnHk
3838
github.com/frankban/quicktest v1.14.6/go.mod h1:4ptaffx2x8+WTWXmUCuVU6aPUX1/Mz7zb5vbUoiM6w0=
3939
github.com/fsnotify/fsnotify v1.7.0 h1:8JEhPFa5W2WU7YfeZzPNqzMP6Lwt7L2715Ggo0nosvA=
4040
github.com/fsnotify/fsnotify v1.7.0/go.mod h1:40Bi/Hjc2AVfZrqy+aj+yEI+/bRxZnMJyTJwOpGvigM=
41-
github.com/gigapi/gigapi-config v0.0.8 h1:fyOwS3by+qMlHLthmIgX1gXEchH2PUH2l4x/y9sRyOg=
42-
github.com/gigapi/gigapi-config v0.0.8/go.mod h1:/hD+d1odWyElSP9++ZPLULXthWHXvx0shNypdbZDYA8=
41+
github.com/gigapi/gigapi-config v0.0.9 h1:7VWmMJb2dKGDReasOgu45bzOmrfj/WCiBcO9g4MNgXw=
42+
github.com/gigapi/gigapi-config v0.0.9/go.mod h1:/hD+d1odWyElSP9++ZPLULXthWHXvx0shNypdbZDYA8=
4343
github.com/gigapi/gigapi/v2 v2.0.13 h1:m4w2Sx+AQZl+SPS+5oCtCe7G1888/2nF6m+d1vsiHlQ=
4444
github.com/gigapi/gigapi/v2 v2.0.13/go.mod h1:xoh5uz+EQCK43cI08wMJlyhAZ1368OOPlJjmOr1TGis=
4545
github.com/gigapi/metadata v0.0.4 h1:kQdreAUyeguIrwSlPukzZ36Y9HlKVhcaBgrg8crK1tY=

querier/layerDesc.go

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,67 @@
1+
package querier
2+
3+
import (
4+
"fmt"
5+
"github.com/gigapi/gigapi-config/config"
6+
"net/url"
7+
"strings"
8+
)
9+
10+
type querierLayerDesc struct {
11+
layer config.LayersConfiguration
12+
path string
13+
hostname string
14+
key string
15+
secret string
16+
bucket string
17+
secure bool
18+
urlStyle string
19+
}
20+
21+
func getQuerierLayerDesc(layer config.LayersConfiguration) (querierLayerDesc, error) {
22+
switch layer.Type {
23+
case "fs":
24+
return getQuerierLayerDescFs(layer)
25+
case "s3":
26+
return getQuerierLayerDescS3(layer)
27+
}
28+
return querierLayerDesc{}, fmt.Errorf("Unsupported layer type: %s", layer.Type)
29+
}
30+
31+
func getQuerierLayerDescFs(layer config.LayersConfiguration) (querierLayerDesc, error) {
32+
return querierLayerDesc{
33+
layer: layer,
34+
path: strings.TrimPrefix(layer.URL, "file://"),
35+
}, nil
36+
}
37+
38+
func getQuerierLayerDescS3(layer config.LayersConfiguration) (querierLayerDesc, error) {
39+
s3Url, err := url.Parse(layer.URL)
40+
if err != nil {
41+
return querierLayerDesc{}, err
42+
}
43+
if s3Url.Scheme != "s3" {
44+
return querierLayerDesc{}, fmt.Errorf("Invalid S3 URL: %s", layer.URL)
45+
}
46+
47+
pathParts := strings.SplitN(s3Url.Path, "/", 2)
48+
49+
res := querierLayerDesc{
50+
layer: layer,
51+
hostname: s3Url.Host,
52+
bucket: pathParts[0],
53+
path: pathParts[1],
54+
secure: s3Url.Query().Get("secure") != "false",
55+
key: layer.Auth.Key,
56+
secret: layer.Auth.Secret,
57+
}
58+
if s3Url.User != nil {
59+
res.key = s3Url.User.Username()
60+
res.secret, _ = s3Url.User.Password()
61+
}
62+
res.urlStyle = "vhost"
63+
if s3Url.Query().Get("url-style") == "path" {
64+
res.urlStyle = "path"
65+
}
66+
return res, nil
67+
}

querier/queryClient.go

Lines changed: 67 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,8 @@ type QueryClient struct {
3030
DataDir string
3131
DB *sql.DB
3232
DefaultTimeRange int64 // 10 minutes in nanoseconds
33+
34+
layers map[string]querierLayerDesc
3335
}
3436

3537
// NewQueryClient creates a new QueryClient
@@ -48,6 +50,16 @@ func (q *QueryClient) Initialize() error {
4850
return fmt.Errorf("failed to initialize DuckDB: %v", err)
4951
}
5052
q.DB = db
53+
q.layers = map[string]querierLayerDesc{}
54+
55+
for _, l := range config.Config.Gigapi.Layers {
56+
desc, err := getQuerierLayerDesc(l)
57+
if err != nil {
58+
return fmt.Errorf("failed to initialize layer %s: %v", l.Name, err)
59+
}
60+
q.layers[l.Name] = desc
61+
}
62+
5163
return nil
5264
}
5365

@@ -422,7 +434,7 @@ func getTableIndex(table *shared.Table) (metadata.TableIndex, error) {
422434

423435
// Find relevant parquet files based on time range
424436
func (q *QueryClient) FindRelevantFiles(ctx context.Context, dbName, measurement string,
425-
timeRange TimeRange) ([]string, error) {
437+
timeRange TimeRange) ([]*metadata.IndexEntry, error) {
426438
// If no time range specified, get all files
427439
var relevantFiles []string
428440
// log.Printf("Getting relevant files for %s.%s within time range %v to %v", dbName, measurement,
@@ -451,19 +463,13 @@ func (q *QueryClient) FindRelevantFiles(ctx context.Context, dbName, measurement
451463
}
452464

453465
ies, err := idx.GetQuerier().Query(opts)
454-
layersMap := make(map[string]string)
455-
for _, l := range config.Config.Gigapi.Layers {
456-
layersMap[l.Name] = strings.TrimPrefix(l.URL, "file://")
457-
}
458-
for _, ie := range ies {
459-
relevantFiles = append(relevantFiles, filepath.Join(layersMap[ie.Layer], dbName, measurement, "data", ie.Path))
466+
if err != nil {
467+
return nil, err
460468
}
461-
462-
if len(relevantFiles) == 0 {
469+
if len(ies) == 0 {
463470
log.Printf("No files found in any directory for %s.%s", dbName, measurement)
464471
}
465-
466-
return relevantFiles, nil
472+
return ies, err
467473
}
468474

469475
// Find all files for a measurement
@@ -660,6 +666,49 @@ func (c *QueryClient) getDBIndex(ctx context.Context) (metadata.DBIndex, error)
660666
return nil, fmt.Errorf("unsupported metadata type: %s", config.Config.Gigapi.Metadata.Type)
661667
}
662668

669+
func (c *QueryClient) buildFilesList(db string, table string, files []*metadata.IndexEntry) ([]string, error) {
670+
var fileList []string
671+
s3Layers := map[string]bool{}
672+
673+
for _, f := range files {
674+
l, ok := c.layers[f.Layer]
675+
if !ok {
676+
return nil, fmt.Errorf("layer %s not found", f.Layer)
677+
}
678+
switch l.layer.Type {
679+
case "fs":
680+
fileList = append(fileList, filepath.Join(l.path, db, table, "data", f.Path))
681+
case "s3":
682+
_path := l.path
683+
if _path != "" {
684+
_path += "/"
685+
}
686+
path := fmt.Sprintf("s3://%s%s/%s/%s", _path, db, table, f.Path)
687+
fileList = append(fileList, path)
688+
s3Layers[l.layer.Name] = true
689+
}
690+
}
691+
for lName, _ := range s3Layers {
692+
sanitizedLName := lName
693+
sanitizedLName = regexp.MustCompile("[^a-zA-Z0-9_]").ReplaceAllString(sanitizedLName, "_")
694+
secretName := fmt.Sprintf("secret_%s", sanitizedLName)
695+
l := c.layers[lName]
696+
_, err := c.DB.Exec(fmt.Sprintf(`CREATE OR REPLACE SECRET %s (
697+
TYPE s3,
698+
USE_SSL %t,
699+
KEY_ID %s,
700+
SECRET %s,
701+
ENDPOINT '%s',
702+
SCOPE 's3://%s',
703+
URL_STYLE '%s'
704+
);`, secretName, l.secure, l.key, l.secret, l.hostname, l.bucket, l.urlStyle))
705+
if err != nil {
706+
return nil, fmt.Errorf("failed to create secret: %v", err)
707+
}
708+
}
709+
return fileList, nil
710+
}
711+
663712
// Query executes a query against the database
664713
func (c *QueryClient) Query(ctx context.Context, query, dbName string) ([]map[string]interface{}, error) {
665714
// Ensure we have a context
@@ -766,9 +815,13 @@ func (c *QueryClient) Query(ctx context.Context, query, dbName string) ([]map[st
766815
}
767816

768817
// Find relevant files
769-
files, err := c.FindRelevantFiles(ctx, parsed.DbName, parsed.Measurement, parsed.TimeRange)
770-
if err != nil || len(files) == 0 {
771-
return nil, fmt.Errorf("no relevant files found for query")
818+
entries, err := c.FindRelevantFiles(ctx, parsed.DbName, parsed.Measurement, parsed.TimeRange)
819+
if err != nil {
820+
return nil, err
821+
}
822+
files, err := c.buildFilesList(parsed.DbName, parsed.Measurement, entries)
823+
if err != nil {
824+
return nil, err
772825
}
773826

774827
start := time.Now()

0 commit comments

Comments
 (0)