diff --git a/internal/ducklake/writer.go b/internal/ducklake/writer.go index 165b376..28123da 100644 --- a/internal/ducklake/writer.go +++ b/internal/ducklake/writer.go @@ -27,6 +27,7 @@ type Config struct { S3Endpoint string S3Region string CatalogName string + ReadOnly bool } type Writer struct { @@ -143,10 +144,16 @@ func Configure(ctx context.Context, db *sql.DB, cfg Config) error { return fmt.Errorf("exec %q: %w", stmt, err) } } - attach := fmt.Sprintf("ATTACH '%s' AS %s (DATA_PATH '%s')", + attachOptions := []string{ + fmt.Sprintf("DATA_PATH '%s'", strings.ReplaceAll(cfg.DataPath, "'", "''")), + } + if cfg.ReadOnly { + attachOptions = append(attachOptions, "READ_ONLY") + } + attach := fmt.Sprintf("ATTACH '%s' AS %s (%s)", duckLakeAttachPath(cfg.CatalogPath, catalogStore), quoteIdent(catalogName(cfg)), - strings.ReplaceAll(cfg.DataPath, "'", "''"), + strings.Join(attachOptions, ", "), ) if _, err := db.ExecContext(ctx, attach); err != nil { return fmt.Errorf("attach ducklake catalog: %w", err) diff --git a/internal/server/ducklake_test.go b/internal/server/ducklake_test.go index 36341b7..b0dd276 100644 --- a/internal/server/ducklake_test.go +++ b/internal/server/ducklake_test.go @@ -2,6 +2,7 @@ package server import ( "context" + "fmt" "log/slog" "os" "path/filepath" @@ -12,34 +13,96 @@ import ( "github.com/viggy28/streambed/internal/wal" ) -func TestDuckLakeServerRegistersViews(t *testing.T) { +func TestDuckLakeServerUsesAttachedCatalogDirectly(t *testing.T) { ctx := context.Background() dir := t.TempDir() - catalogPath := filepath.Join(dir, "catalog.sqlite") + catalogPath := filepath.Join(dir, "catalog.ducklake") dataPath := filepath.Join(dir, "data") + "/" logger := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelError})) writer, err := ducklake.NewWriter(ctx, ducklake.Config{ - CatalogPath: catalogPath, - DataPath: dataPath, - }, nil, 10, time.Second, logger) + CatalogPath: catalogPath, + CatalogStore: "duckdb", + DataPath: dataPath, + }, nil, 100, time.Hour, logger) if err != nil { t.Fatalf("ducklake writer: %v", err) } - _, err = writer.HandleEvent(ctx, wal.RowEvent{ - Schema: "public", - Table: "orders", - Columns: []wal.Column{{Name: "id", OID: 23, IsKey: true}, {Name: "name", OID: 25}}, - KeyColumns: []int{0}, - Op: wal.OpInsert, - Values: []wal.ColumnValue{ - {Name: "id", OID: 23, Value: []byte("1")}, - {Name: "name", OID: 25, Value: []byte("alice")}, - }, - }) + defer writer.Close() + + writeDuckLakeRow(t, writer, "orders", 1, "alice") + if err := writer.FlushAll(ctx); err != nil { + t.Fatalf("initial flush: %v", err) + } + + srv, err := NewServer(ServerConfig{ + ListenAddr: ":0", + TargetFormat: "ducklake", + DuckLakeCatalog: catalogPath, + DuckLakeCatalogStore: "duckdb", + DuckLakeDataPath: dataPath, + }, nil, logger) if err != nil { - t.Fatalf("handle event: %v", err) + t.Fatalf("NewServer: %v", err) } + defer srv.Close() + + assertDuckLakeCount(t, srv, `streambed.public.orders`, 1) // fully qualified + assertDuckLakeCount(t, srv, `public.orders`, 1) // current catalog + assertDuckLakeCount(t, srv, `orders`, 1) // current catalog and schema + + if _, err := srv.duckDB.Exec(`INSERT INTO streambed.public.orders VALUES (99, 'mallory')`); err == nil { + t.Fatal("query-side DuckLake attachment accepted a write") + } + // Mutating the query session must not disturb the writer's attachment. The + // next query also gets a fresh session, so this DETACH is self-contained. + if _, err := srv.duckDB.Exec(`USE memory; DETACH streambed`); err != nil { + t.Fatalf("detach query-side catalog: %v", err) + } + + // Existing tables expose newly committed rows without recreating a view. + writeDuckLakeRow(t, writer, "orders", 2, "bob") + if err := writer.FlushAll(ctx); err != nil { + t.Fatalf("existing-table flush: %v", err) + } + assertDuckLakeCount(t, srv, `public.orders`, 2) + + // Newly committed tables are visible without waiting for catalog discovery. + writeDuckLakeRow(t, writer, "customers", 1, "carol") + if err := writer.FlushAll(ctx); err != nil { + t.Fatalf("new-table flush: %v", err) + } + assertDuckLakeCount(t, srv, `public.customers`, 1) + + var viewCount int + if err := srv.duckDB.QueryRow(` + SELECT count(*) + FROM information_schema.tables + WHERE table_type = 'VIEW' AND table_name IN ('orders', 'public_orders', 'customers', 'public_customers') + `).Scan(&viewCount); err != nil { + t.Fatalf("list compatibility views: %v", err) + } + if viewCount != 0 { + t.Fatalf("got %d DuckLake compatibility views, want none", viewCount) + } +} + +func TestDuckLakeServerAttachesCatalogInStandaloneMode(t *testing.T) { + ctx := context.Background() + dir := t.TempDir() + catalogPath := filepath.Join(dir, "catalog.ducklake") + dataPath := filepath.Join(dir, "data") + "/" + logger := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelError})) + + writer, err := ducklake.NewWriter(ctx, ducklake.Config{ + CatalogPath: catalogPath, + CatalogStore: "duckdb", + DataPath: dataPath, + }, nil, 100, time.Hour, logger) + if err != nil { + t.Fatalf("ducklake writer: %v", err) + } + writeDuckLakeRow(t, writer, "orders", 1, "alice") if err := writer.FlushAll(ctx); err != nil { t.Fatalf("flush: %v", err) } @@ -48,29 +111,52 @@ func TestDuckLakeServerRegistersViews(t *testing.T) { } srv, err := NewServer(ServerConfig{ - ListenAddr: ":0", - TargetFormat: "ducklake", - DuckLakeCatalog: catalogPath, - DuckLakeDataPath: dataPath, + ListenAddr: ":0", + TargetFormat: "ducklake", + DuckLakeCatalog: catalogPath, + DuckLakeCatalogStore: "duckdb", + DuckLakeDataPath: dataPath, }, nil, logger) if err != nil { t.Fatalf("NewServer: %v", err) } defer srv.Close() - if err := srv.refreshAndRegister(ctx); err != nil { - t.Fatalf("refreshAndRegister: %v", err) - } - var count int - if err := srv.duckDB.QueryRow(`SELECT count(*) FROM orders`).Scan(&count); err != nil { - t.Fatalf("query unqualified view: %v", err) + + assertDuckLakeCount(t, srv, `streambed.public.orders`, 1) + assertDuckLakeCount(t, srv, `public.orders`, 1) + assertDuckLakeCount(t, srv, `orders`, 1) +} + +func writeDuckLakeRow(t *testing.T, writer *ducklake.Writer, table string, id int, name string) { + t.Helper() + _, err := writer.HandleEvent(context.Background(), wal.RowEvent{ + Schema: "public", + Table: table, + Columns: []wal.Column{{Name: "id", OID: 23, IsKey: true}, {Name: "name", OID: 25}}, + KeyColumns: []int{0}, + Op: wal.OpInsert, + Values: []wal.ColumnValue{ + {Name: "id", OID: 23, Value: []byte(fmt.Sprint(id))}, + {Name: "name", OID: 25, Value: []byte(name)}, + }, + }) + if err != nil { + t.Fatalf("write %s row: %v", table, err) } - if count != 1 { - t.Fatalf("got count=%d, want 1", count) +} + +func assertDuckLakeCount(t *testing.T, srv *Server, table string, want int) { + t.Helper() + srv.duckDBMu.Lock() + defer srv.duckDBMu.Unlock() + if err := srv.resetDuckLakeQueryDB(context.Background()); err != nil { + t.Fatalf("reset DuckLake query DB: %v", err) } - if err := srv.duckDB.QueryRow(`SELECT count(*) FROM public_orders`).Scan(&count); err != nil { - t.Fatalf("query qualified view: %v", err) + var got int + if err := srv.duckDB.QueryRow(`SELECT count(*) FROM ` + table).Scan(&got); err != nil { + t.Fatalf("query %s: %v", table, err) } - if count != 1 { - t.Fatalf("got qualified count=%d, want 1", count) + if got != want { + t.Fatalf("query %s: got count=%d, want %d", table, got, want) } } diff --git a/internal/server/server.go b/internal/server/server.go index 97cbb0d..c195c67 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -7,6 +7,7 @@ import ( "log/slog" "os" "strings" + "sync" "time" duckdb "github.com/duckdb/duckdb-go/v2" @@ -30,16 +31,17 @@ type ServerConfig struct { } // Server implements a Postgres-wire-compatible query interface backed by DuckDB. -// It reads Iceberg tables from S3 and serves them to any Postgres client. +// It serves Iceberg views or tables from an attached DuckLake catalog. type Server struct { - cfg ServerConfig - catalog *TableCatalog - duckDB *sql.DB - logger *slog.Logger + cfg ServerConfig + catalog *TableCatalog + duckDB *sql.DB + duckDBMu sync.Mutex + logger *slog.Logger } -// NewServer creates a query server. Initializes an embedded DuckDB instance -// configured with S3 credentials and loads the Iceberg + httpfs extensions. +// NewServer creates a query server and initializes DuckDB for the selected +// lakehouse target. func NewServer(cfg ServerConfig, s3Client storage.ObjectStorage, logger *slog.Logger) (*Server, error) { if cfg.TargetFormat == "" { cfg.TargetFormat = "iceberg" @@ -50,16 +52,24 @@ func NewServer(cfg ServerConfig, s3Client storage.ObjectStorage, logger *slog.Lo } var catalog *TableCatalog if cfg.TargetFormat == "ducklake" { + // The query server owns an isolated, read-only DuckDB session so client + // SQL cannot mutate the writer's DuckLake attachment. + db.SetMaxOpenConns(1) if err := ducklake.Configure(context.Background(), db, ducklake.Config{ CatalogPath: cfg.DuckLakeCatalog, CatalogStore: cfg.DuckLakeCatalogStore, DataPath: cfg.DuckLakeDataPath, S3Endpoint: cfg.S3Endpoint, S3Region: cfg.S3Region, + ReadOnly: true, }); err != nil { db.Close() return nil, fmt.Errorf("configure ducklake: %w", err) } + if err := configureDuckLakeQuerySession(db); err != nil { + db.Close() + return nil, err + } } else { conn, _ := db.Conn(context.Background()) // configure DuckDB on a specific connection @@ -79,6 +89,44 @@ func NewServer(cfg ServerConfig, s3Client storage.ObjectStorage, logger *slog.Lo }, nil } +func (s *Server) resetDuckLakeQueryDB(ctx context.Context) error { + db, err := sql.Open("duckdb", "") + if err != nil { + return fmt.Errorf("open fresh duckdb query session: %w", err) + } + db.SetMaxOpenConns(1) + if err := ducklake.Configure(ctx, db, ducklake.Config{ + CatalogPath: s.cfg.DuckLakeCatalog, + CatalogStore: s.cfg.DuckLakeCatalogStore, + DataPath: s.cfg.DuckLakeDataPath, + S3Endpoint: s.cfg.S3Endpoint, + S3Region: s.cfg.S3Region, + ReadOnly: true, + }); err != nil { + db.Close() + return fmt.Errorf("configure fresh ducklake query session: %w", err) + } + if err := configureDuckLakeQuerySession(db); err != nil { + db.Close() + return err + } + oldDB := s.duckDB + s.duckDB = db + return oldDB.Close() +} + +func configureDuckLakeQuerySession(db *sql.DB) error { + if _, err := db.Exec("USE streambed"); err != nil { + return fmt.Errorf("use ducklake catalog: %w", err) + } + // Match PostgreSQL's usual default schema while retaining DuckLake's main + // schema as a fallback. Missing schemas in the search path are allowed. + if _, err := db.Exec("SET search_path = 'streambed.public,streambed.main'"); err != nil { + return fmt.Errorf("configure ducklake search path: %w", err) + } + return nil +} + // configureDuckDBPerConn configures DuckDB on a specific connection -- delete it func configureDuckDBPerConn(ctx context.Context, con *sql.Conn, cfg ServerConfig) error { stmts := []string{ @@ -189,17 +237,18 @@ func configureDuckDB(db *sql.DB, cfg ServerConfig) error { } // Start begins listening for Postgres client connections and serving queries. -// It performs an initial catalog refresh, starts a background refresh goroutine, -// and blocks until ctx is cancelled. +// Iceberg mode refreshes discovered views in the background; DuckLake mode uses +// its attached catalog directly. Start blocks until ctx is cancelled. func (s *Server) Start(ctx context.Context) error { - // Initial catalog refresh and view registration - if err := s.refreshAndRegister(ctx); err != nil { - s.logger.Warn("initial catalog refresh failed (will retry)", "error", err) + // Iceberg tables must be discovered in S3 and exposed as views. DuckLake + // tables are already visible through the directly attached catalog. + if s.cfg.TargetFormat != "ducklake" { + if err := s.refreshAndRegister(ctx); err != nil { + s.logger.Warn("initial catalog refresh failed (will retry)", "error", err) + } + go s.refreshLoop(ctx) } - // Periodic catalog refresh in background - go s.refreshLoop(ctx) - // Create psql-wire server srv, err := wire.NewServer(s.handleParse, wire.Logger(s.logger), @@ -241,6 +290,17 @@ func (s *Server) handleParse(ctx context.Context, query string) (wire.PreparedSt s.logger.Debug("query received", "query", query) + s.duckDBMu.Lock() + defer s.duckDBMu.Unlock() + if s.cfg.TargetFormat == "ducklake" { + // Give every client query an isolated, read-only attachment. Besides + // containing session mutations such as DETACH, reopening refreshes the + // snapshot cached by DuckDB-backed metadata catalogs. + if err := s.resetDuckLakeQueryDB(ctx); err != nil { + return nil, err + } + } + // Execute query against DuckDB rows, err := s.duckDB.QueryContext(ctx, query) if err != nil { @@ -311,9 +371,8 @@ func (s *Server) handleParse(ctx context.Context, query string) (wire.PreparedSt // If DuckDB's engine was invalidated by a FATAL error (e.g., corrupt Iceberg // table), it re-opens the DuckDB instance and retries view registration. func (s *Server) refreshAndRegister(ctx context.Context) error { - if s.cfg.TargetFormat == "ducklake" { - return s.registerDuckLakeViews(ctx) - } + s.duckDBMu.Lock() + defer s.duckDBMu.Unlock() if err := s.catalog.Refresh(ctx); err != nil { return err } @@ -340,49 +399,6 @@ func (s *Server) refreshAndRegister(ctx context.Context) error { return s.catalog.RegisterViews(s.duckDB) } -func (s *Server) registerDuckLakeViews(ctx context.Context) error { - const catalogName = "streambed" - rows, err := s.duckDB.QueryContext(ctx, - "SELECT table_schema, table_name FROM information_schema.tables WHERE table_catalog = ? AND table_type = 'BASE TABLE'", - catalogName, - ) - if err != nil { - return fmt.Errorf("list ducklake tables: %w", err) - } - defer rows.Close() - type tableInfo struct { - schema string - table string - } - var tables []tableInfo - nameCount := map[string]int{} - for rows.Next() { - var info tableInfo - if err := rows.Scan(&info.schema, &info.table); err != nil { - return err - } - tables = append(tables, info) - nameCount[info.table]++ - } - if err := rows.Err(); err != nil { - return err - } - for _, info := range tables { - src := fmt.Sprintf("%s.%s.%s", quoteIdent(catalogName), quoteIdent(info.schema), quoteIdent(info.table)) - qualified := quoteIdent(info.schema + "_" + info.table) - if _, err := s.duckDB.ExecContext(ctx, fmt.Sprintf("CREATE OR REPLACE VIEW %s AS SELECT * FROM %s", qualified, src)); err != nil { - return fmt.Errorf("register ducklake view %s_%s: %w", info.schema, info.table, err) - } - if nameCount[info.table] == 1 { - if _, err := s.duckDB.ExecContext(ctx, fmt.Sprintf("CREATE OR REPLACE VIEW %s AS SELECT * FROM %s", quoteIdent(info.table), src)); err != nil { - return fmt.Errorf("register ducklake view %s: %w", info.table, err) - } - } - } - s.logger.Info("ducklake views registered", "count", len(tables)) - return nil -} - // refreshLoop periodically refreshes the catalog and re-registers views. func (s *Server) refreshLoop(ctx context.Context) { ticker := time.NewTicker(30 * time.Second) @@ -402,6 +418,8 @@ func (s *Server) refreshLoop(ctx context.Context) { // Close shuts down the DuckDB connection. func (s *Server) Close() error { + s.duckDBMu.Lock() + defer s.duckDBMu.Unlock() return s.duckDB.Close() } @@ -456,7 +474,3 @@ func normalizeValue(v any) any { return v } } - -func quoteIdent(s string) string { - return `"` + strings.ReplaceAll(s, `"`, `""`) + `"` -}