Skip to content

Commit 8b24221

Browse files
lmanganiqxip
andauthored
query passthrough (#18)
* passthrough * flight passthrough * Secure pass-through --------- Co-authored-by: qxip <qxip@mini-ams.local>
1 parent a368e77 commit 8b24221

2 files changed

Lines changed: 48 additions & 18 deletions

File tree

querier/flightsql.go

Lines changed: 2 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -146,23 +146,8 @@ func (s *FlightSQLServer) GetFlightInfo(ctx context.Context, desc *flight.Flight
146146
}
147147
}
148148

149-
// Parse the query to extract time range
150-
parsed, err := s.queryClient.ParseQuery(query, dbName)
151-
if err != nil {
152-
log.Printf("Failed to parse query: %v", err)
153-
return nil, fmt.Errorf("failed to parse query: %w", err)
154-
}
155-
156-
// Find relevant files based on the parsed query
157-
files, err := s.queryClient.FindRelevantFiles(ctx, parsed.DbName, parsed.Measurement, parsed.TimeRange)
158-
if err != nil {
159-
log.Printf("Failed to find relevant files: %v", err)
160-
return nil, fmt.Errorf("failed to find relevant files: %w", err)
161-
}
162-
log.Printf("Found %d relevant files for query", len(files))
163-
164-
// Execute the query using our existing QueryClient
165-
results, err := s.queryClient.Query(ctx, query, parsed.DbName) // Use the parsed database name
149+
// Use QueryClient.Query which now handles all fallback logic
150+
results, err := s.queryClient.Query(ctx, query, dbName)
166151
if err != nil {
167152
log.Printf("Query execution failed: %v", err)
168153
return nil, fmt.Errorf("failed to execute query: %w", err)

querier/queryClient.go

Lines changed: 46 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -670,7 +670,52 @@ func (c *QueryClient) Query(ctx context.Context, query, dbName string) ([]map[st
670670
// Parse the query
671671
parsed, err := c.ParseQuery(query, dbName)
672672
if err != nil {
673-
return nil, err
673+
// Fallback: Directly execute the query in DuckDB if ParseQuery fails
674+
// Set DuckDB safety option to disable external access
675+
if _, errSet := c.DB.Exec("SET enable_external_access = false;"); errSet != nil {
676+
return nil, fmt.Errorf("failed to disable external access: %v", errSet)
677+
}
678+
if _, errSet := c.DB.Exec("SET lock_configuration = true;"); errSet != nil {
679+
return nil, fmt.Errorf("failed to lock configuration: %v", errSet)
680+
}
681+
682+
stmt, err2 := c.DB.Prepare(query)
683+
if err2 != nil {
684+
return nil, fmt.Errorf("failed to prepare query: %v", err2)
685+
}
686+
defer stmt.Close()
687+
688+
rows, err2 := stmt.Query()
689+
if err2 != nil {
690+
return nil, fmt.Errorf("query execution failed: %v", err2)
691+
}
692+
defer rows.Close()
693+
694+
columns, err2 := rows.Columns()
695+
if err2 != nil {
696+
return nil, fmt.Errorf("failed to get columns: %v", err2)
697+
}
698+
699+
var result []map[string]interface{}
700+
for rows.Next() {
701+
values := make([]interface{}, len(columns))
702+
valuePtrs := make([]interface{}, len(columns))
703+
for i := range values {
704+
valuePtrs[i] = &values[i]
705+
}
706+
if err := rows.Scan(valuePtrs...); err != nil {
707+
return nil, fmt.Errorf("error scanning row: %v", err)
708+
}
709+
row := make(map[string]interface{})
710+
for i, col := range columns {
711+
row[col] = values[i]
712+
}
713+
result = append(result, row)
714+
}
715+
if err := rows.Err(); err != nil {
716+
return nil, fmt.Errorf("error iterating rows: %v", err)
717+
}
718+
return result, nil
674719
}
675720

676721
// Find relevant files

0 commit comments

Comments
 (0)