Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions cmd/arc/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -1047,6 +1047,10 @@ func main() {
queryHandler.SetTieringManager(tieringManager)
log.Info().Msg("Tiering manager wired to query handler for multi-tier queries")

// Wire tiering manager to databases handler for cold-tier database/measurement listing
databasesHandler.SetTieringManager(tieringManager)
log.Info().Msg("Tiering manager wired to databases handler for cold-tier listing")

// Wire tiering manager to arrow buffer for automatic file registration
arrowBuffer.SetTieringManager(tieringManager)
log.Info().Msg("Tiering manager wired to arrow buffer for auto-registration")
Expand Down
119 changes: 96 additions & 23 deletions internal/api/databases.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,15 +9,17 @@ import (

"github.com/basekick-labs/arc/internal/config"
"github.com/basekick-labs/arc/internal/storage"
"github.com/basekick-labs/arc/internal/tiering"
"github.com/gofiber/fiber/v2"
"github.com/rs/zerolog"
)

// DatabasesHandler handles database management API endpoints
type DatabasesHandler struct {
storage storage.Backend
deleteConfig *config.DeleteConfig
logger zerolog.Logger
storage storage.Backend
deleteConfig *config.DeleteConfig
tieringManager *tiering.Manager
logger zerolog.Logger
}

// CreateDatabaseRequest represents a request to create a new database
Expand Down Expand Up @@ -70,6 +72,12 @@ func NewDatabasesHandler(storage storage.Backend, deleteConfig *config.DeleteCon
}
}

// SetTieringManager sets the tiering manager for multi-tier database/measurement listing.
// This is called after initialization when tiering is enabled and licensed.
func (h *DatabasesHandler) SetTieringManager(tm *tiering.Manager) {
h.tieringManager = tm
}

// RegisterRoutes registers the database management routes
func (h *DatabasesHandler) RegisterRoutes(app *fiber.App) {
app.Get("/api/v1/databases", h.handleList)
Expand Down Expand Up @@ -428,32 +436,52 @@ func (h *DatabasesHandler) listDatabases(ctx context.Context) ([]string, error)
// listDatabasesWithMeasurementCounts returns all databases with their measurement counts
// using minimal storage calls instead of N+1 calls.
// Strategy: Get list of databases first (1 call), then get all files (1 call) to count measurements.
// Also checks tiering metadata for databases that may only have data in cold storage.
func (h *DatabasesHandler) listDatabasesWithMeasurementCounts(ctx context.Context) ([]DatabaseInfo, error) {
// First, get list of all databases (includes empty databases with just .arc-database marker)
// First, get list of all databases from hot tier (includes empty databases with just .arc-database marker)
databases, err := h.listDatabases(ctx)
if err != nil {
return nil, err
}

if len(databases) == 0 {
// Build a map of database -> set of measurements from files
dbMeasurements := make(map[string]map[string]bool)

// Initialize all known databases from hot tier (some may be empty)
for _, db := range databases {
dbMeasurements[db] = make(map[string]bool)
}

// Also check tiering metadata for cold-only databases
if h.tieringManager != nil {
metadata := h.tieringManager.GetMetadata()
if metadata != nil {
coldDatabases, err := metadata.GetAllDatabases(ctx)
if err != nil {
h.logger.Warn().Err(err).Msg("Failed to get databases from tiering metadata")
} else {
// Add any databases that only exist in cold tier
for _, db := range coldDatabases {
if dbMeasurements[db] == nil {
dbMeasurements[db] = make(map[string]bool)
}
}
}
}
}

// If no databases at all, return empty
if len(dbMeasurements) == 0 {
return []DatabaseInfo{}, nil
}

// Single storage call to get all files for measurement counting
// Single storage call to get all files for measurement counting (hot tier)
files, err := h.storage.List(ctx, "")
if err != nil {
return nil, err
}

// Build a map of database -> set of measurements from files
dbMeasurements := make(map[string]map[string]bool)

// Initialize all known databases (some may be empty)
for _, db := range databases {
dbMeasurements[db] = make(map[string]bool)
}

// Count measurements from files
// Count measurements from hot tier files
for _, file := range files {
parts := strings.SplitN(file, "/", 3)
if len(parts) < 2 {
Expand All @@ -476,9 +504,34 @@ func (h *DatabasesHandler) listDatabasesWithMeasurementCounts(ctx context.Contex
dbMeasurements[db][measurement] = true
}

// Also get measurements from tiering metadata (for cold-only measurements)
if h.tieringManager != nil {
metadata := h.tieringManager.GetMetadata()
if metadata != nil {
for db := range dbMeasurements {
coldMeasurements, err := metadata.GetMeasurementsByDatabase(ctx, db)
if err != nil {
h.logger.Warn().Err(err).Str("database", db).Msg("Failed to get measurements from tiering metadata")
continue
}
for _, m := range coldMeasurements {
if !strings.HasPrefix(m, ".") && !strings.HasPrefix(m, "_") {
dbMeasurements[db][m] = true
}
}
}
}
}

// Convert to sorted list of DatabaseInfo
result := make([]DatabaseInfo, 0, len(databases))
for _, db := range databases {
sortedDatabases := make([]string, 0, len(dbMeasurements))
for db := range dbMeasurements {
sortedDatabases = append(sortedDatabases, db)
}
sort.Strings(sortedDatabases)

result := make([]DatabaseInfo, 0, len(sortedDatabases))
for _, db := range sortedDatabases {
result = append(result, DatabaseInfo{
Name: db,
MeasurementCount: len(dbMeasurements[db]),
Expand All @@ -489,27 +542,47 @@ func (h *DatabasesHandler) listDatabasesWithMeasurementCounts(ctx context.Contex
}

func (h *DatabasesHandler) listMeasurements(ctx context.Context, database string) ([]string, error) {
var measurements []string
// Use a set to deduplicate measurements from hot and cold tiers
measurementSet := make(map[string]bool)

// Use DirectoryLister if available
// Get measurements from hot tier
if lister, ok := h.storage.(storage.DirectoryLister); ok {
dirs, err := lister.ListDirectories(ctx, database+"/")
if err != nil {
return nil, err
}
measurements = dirs
for _, m := range dirs {
measurementSet[m] = true
}
} else {
// Fall back to List and extract unique subdirectories
files, err := h.storage.List(ctx, database+"/")
if err != nil {
return nil, err
}
measurements = extractSubdirectories(files, database)
for _, m := range extractSubdirectories(files, database) {
measurementSet[m] = true
}
}

// Also get measurements from tiering metadata (for cold-only measurements)
if h.tieringManager != nil {
metadata := h.tieringManager.GetMetadata()
if metadata != nil {
coldMeasurements, err := metadata.GetMeasurementsByDatabase(ctx, database)
if err != nil {
h.logger.Warn().Err(err).Str("database", database).Msg("Failed to get measurements from tiering metadata")
} else {
for _, m := range coldMeasurements {
measurementSet[m] = true
}
}
}
}

// Filter out hidden directories and sort
filtered := make([]string, 0, len(measurements))
for _, m := range measurements {
filtered := make([]string, 0, len(measurementSet))
for m := range measurementSet {
if !strings.HasPrefix(m, ".") && !strings.HasPrefix(m, "_") {
filtered = append(filtered, m)
}
Expand Down
47 changes: 43 additions & 4 deletions internal/api/query.go
Original file line number Diff line number Diff line change
Expand Up @@ -1496,6 +1496,33 @@ func (h *QueryHandler) buildReadParquetExpr(path, originalSQL, keyword string) s
return keyword + " read_parquet('" + path + "', " + options + ")"
}

// buildReadParquetExprForMeasurement builds a read_parquet expression for a database/measurement pair.
// This is the tiering-aware version used by the fast path that takes database and measurement
// separately instead of a pre-constructed path, allowing proper tiering metadata lookup.
func (h *QueryHandler) buildReadParquetExprForMeasurement(database, measurement, originalSQL, keyword string) string {
// Check if tiering is enabled and cold tier is configured
if h.tieringManager != nil {
router := h.tieringManager.GetRouter()
if router != nil {
// Get glob paths for both tiers
tieredPaths := router.GetGlobPathsForQuery(database, measurement)

// If cold tier is configured and enabled, build multi-tier query
if _, hasCold := tieredPaths[tiering.TierCold]; hasCold {
h.logger.Debug().
Str("database", database).
Str("measurement", measurement).
Msg("Tiering enabled: building multi-tier query (fast path)")
return h.buildMultiTierReadParquet(database, measurement, tieredPaths, keyword)
}
}
}

// Fall back to single-tier behavior (hot tier only)
path := h.getStoragePath(database, measurement)
return h.buildReadParquetExpr(path, originalSQL, keyword)
}

// buildReadParquetExprForParallel builds a read_parquet expression and returns
// parallel execution info if the query can benefit from parallel partition scanning.
// Returns (sql_expression, parallel_info) where parallel_info is non-nil if parallel is recommended.
Expand Down Expand Up @@ -1740,9 +1767,8 @@ func (h *QueryHandler) convertSingleTableQuery(sql, sqlLower, database string) s
return sql
}

// Build replacement
path := h.getStoragePath(database, tableName)
replacement := h.buildReadParquetExpr(path, sql, "FROM")
// Build replacement - use tiering-aware method that checks both hot and cold tiers
replacement := h.buildReadParquetExprForMeasurement(database, tableName, sql, "FROM")

return sql[:idx] + replacement + sql[end:]
}
Expand Down Expand Up @@ -1780,7 +1806,20 @@ func (h *QueryHandler) convertSingleTableQueryForParallel(sql, sqlLower, databas
return sql, nil
}

// Build replacement with parallel info
// Check tiering first - if cold tier exists, use tiering-aware method (no parallel for multi-tier)
if h.tieringManager != nil {
router := h.tieringManager.GetRouter()
if router != nil {
tieredPaths := router.GetGlobPathsForQuery(database, tableName)
if _, hasCold := tieredPaths[tiering.TierCold]; hasCold {
// Use tiering-aware method - parallel not supported for multi-tier queries
replacement := h.buildMultiTierReadParquet(database, tableName, tieredPaths, "FROM")
return sql[:idx] + replacement + sql[end:], nil
}
}
}

// Build replacement with parallel info (hot tier only)
path := h.getStoragePath(database, tableName)
replacement, parallelInfo := h.buildReadParquetExprForParallel(path, sql, "FROM")

Expand Down
60 changes: 60 additions & 0 deletions internal/tiering/metadata.go
Original file line number Diff line number Diff line change
Expand Up @@ -219,6 +219,66 @@ func (s *MetadataStore) GetFilesByDatabase(ctx context.Context, database string)
return s.scanFiles(rows)
}

// GetAllDatabases returns all unique database names from the tier metadata.
// This includes databases that may only have data in cold storage.
func (s *MetadataStore) GetAllDatabases(ctx context.Context) ([]string, error) {
s.mu.RLock()
defer s.mu.RUnlock()

query := `SELECT DISTINCT database FROM tier_files ORDER BY database`

rows, err := s.db.QueryContext(ctx, query)
if err != nil {
return nil, fmt.Errorf("failed to query databases: %w", err)
}
defer rows.Close()

var databases []string
for rows.Next() {
var db string
if err := rows.Scan(&db); err != nil {
return nil, fmt.Errorf("failed to scan database: %w", err)
}
databases = append(databases, db)
}

if err := rows.Err(); err != nil {
return nil, fmt.Errorf("error iterating databases: %w", err)
}

return databases, nil
}

// GetMeasurementsByDatabase returns all unique measurements for a database from the tier metadata.
// This includes measurements that may only have data in cold storage.
func (s *MetadataStore) GetMeasurementsByDatabase(ctx context.Context, database string) ([]string, error) {
s.mu.RLock()
defer s.mu.RUnlock()

query := `SELECT DISTINCT measurement FROM tier_files WHERE database = ? ORDER BY measurement`

rows, err := s.db.QueryContext(ctx, query, database)
if err != nil {
return nil, fmt.Errorf("failed to query measurements: %w", err)
}
defer rows.Close()

var measurements []string
for rows.Next() {
var m string
if err := rows.Scan(&m); err != nil {
return nil, fmt.Errorf("failed to scan measurement: %w", err)
}
measurements = append(measurements, m)
}

if err := rows.Err(); err != nil {
return nil, fmt.Errorf("error iterating measurements: %w", err)
}

return measurements, nil
}

// UpdateTier updates the tier for a file
func (s *MetadataStore) UpdateTier(ctx context.Context, path string, newTier Tier) error {
s.mu.Lock()
Expand Down