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
11 changes: 9 additions & 2 deletions internal/ducklake/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ type Config struct {
S3Endpoint string
S3Region string
CatalogName string
ReadOnly bool
}

type Writer struct {
Expand Down Expand Up @@ -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)
Expand Down
152 changes: 119 additions & 33 deletions internal/server/ducklake_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package server

import (
"context"
"fmt"
"log/slog"
"os"
"path/filepath"
Expand All @@ -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)
}
Expand All @@ -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)
}
}
Loading
Loading