diff --git a/.github/workflows/s3-tables-tests.yml b/.github/workflows/s3-tables-tests.yml index c215e32af..f146a3c29 100644 --- a/.github/workflows/s3-tables-tests.yml +++ b/.github/workflows/s3-tables-tests.yml @@ -1013,6 +1013,81 @@ jobs: path: test/s3tables/lifecycle/test-output.log retention-days: 3 + duckdb-lance-tests: + name: DuckDB Lance Integration Tests + runs-on: ubuntu-22.04 + timeout-minutes: 30 + + steps: + - name: Check out code + uses: actions/checkout@v7 + with: + # The job uploads a test log on failure; nothing here needs to push, + # so do not leave a token in the checkout for it to pick up. + persist-credentials: false + + - name: Set up Go + uses: actions/setup-go@v7 + with: + go-version-file: 'go.mod' + id: go + + - name: Configure Docker Hub mirror + run: | + echo '{"registry-mirrors": ["https://mirror.gcr.io"]}' | sudo tee /etc/docker/daemon.json + sudo systemctl restart docker + + - name: Pre-pull images + run: | + pull() { for i in 1 2 3; do docker pull "$1" && return 0; sleep 15; done; return 1; } + pull duckdb/duckdb:latest + pull python:3.11-slim + + - name: Run go mod tidy + run: go mod tidy + + - name: Build SeaweedFS + run: | + cd weed && go build -buildvcs=false . + + - name: Run DuckDB Lance Integration Tests + timeout-minutes: 25 + working-directory: test/s3tables/catalog_duckdb_lance + run: | + set -x + set -o pipefail + echo "=== System Information ===" + uname -a + free -h + df -h + docker info + echo "=== Starting DuckDB Lance Tests ===" + + go test -v -timeout 20m . 2>&1 | tee test-output.log || { + echo "DuckDB Lance integration tests failed" + exit 1 + } + + - name: Show test output on failure + if: failure() + working-directory: test/s3tables/catalog_duckdb_lance + run: | + echo "=== Test Output ===" + if [ -f test-output.log ]; then + tail -200 test-output.log + fi + + echo "=== Process information ===" + ps aux | grep -E "(weed|test|docker|duckdb)" || true + + - name: Upload test logs on failure + if: failure() + uses: actions/upload-artifact@v7 + with: + name: duckdb-lance-test-logs + path: test/s3tables/catalog_duckdb_lance/test-output.log + retention-days: 3 + s3-tables-build-verification: name: S3 Tables Build Verification runs-on: ubuntu-22.04 diff --git a/test/s3tables/catalog_duckdb_lance/Dockerfile.seed b/test/s3tables/catalog_duckdb_lance/Dockerfile.seed new file mode 100644 index 000000000..40b9202b1 --- /dev/null +++ b/test/s3tables/catalog_duckdb_lance/Dockerfile.seed @@ -0,0 +1,15 @@ +# Seeds a Lance dataset so DuckDB has something to read. The catalog records +# where a table lives and does not carry its data, so the writing half is +# pylance; DuckDB then reads it over S3 with no catalog involved. +FROM python:3.11-slim + +WORKDIR /app + +RUN pip install --no-cache-dir \ + "lance-namespace==0.8.6" \ + "pylance==10.0.0" \ + "pyarrow==25.0.1" + +COPY seed_table.py /app/ + +CMD ["python3", "/app/seed_table.py", "--help"] diff --git a/test/s3tables/catalog_duckdb_lance/README.md b/test/s3tables/catalog_duckdb_lance/README.md new file mode 100644 index 000000000..bf1656c04 --- /dev/null +++ b/test/s3tables/catalog_duckdb_lance/README.md @@ -0,0 +1,66 @@ +# DuckDB Lance Integration Test + +Reads SeaweedFS Lance tables from DuckDB's `lance` core extension, the +counterpart of the DuckDB Iceberg tests in `../catalog/`. + +## What this one proves that the others do not + +The LanceDB and Spark suites go through the catalog. DuckDB does not: it reaches +the data over S3 with no namespace involved. That exercises the other half of the +design — a table bucket's layout is a valid Lance dataset directory, so a table +stays readable when the catalog is not in the path. + +## What it does + +`TestDuckDBLance`: + +1. Starts a `weed mini` cluster with S3 and the Lance Namespace enabled. +2. Creates a table bucket declared `LANCE`, then declares a table through the + namespace and writes a dataset into it with pylance. +3. Runs `duckdb_lance_ops.sql` in `duckdb/duckdb:latest`. + +| Step | What it proves | +| --- | --- | +| `__lance_scan(s3://…)` | DuckDB reads a table this catalog created | +| `DESCRIBE` | the schema survived, vector column included | +| `WHERE id < 5` | the filter path | +| `lance_vector_search` | vector search over data behind SeaweedFS | +| a `.lance` path | the replacement scan works on a suffixed path | +| a suffix-less path | and does **not** see one without the suffix | + +## The `.lance` suffix + +DuckDB's replacement scan — `SELECT * FROM 's3://…'` — recognises a Lance dataset +by a `.lance` path suffix. Tables created through this catalog deliberately have +none: the catalog entry *is* the dataset directory, table names may not contain a +dot, and a suffix would leak into ARNs and policies. + +So from DuckDB, a table in a SeaweedFS table bucket is read with +`__lance_scan('s3://bucket/namespace/table')` rather than the bare `FROM 's3://…'` +form. The test asserts both halves, so if the extension ever recognises a +suffix-less directory the test fails and tells us to update the documentation. + +## Credentials + +```sql +CREATE SECRET seaweedfs ( + TYPE lance, + ACCESS_KEY_ID '…', SECRET_ACCESS_KEY '…', + REGION 'us-east-1', + ENDPOINT 'http://seaweed:8333', + ALLOW_HTTP true, + VIRTUAL_HOSTED_STYLE_REQUEST false +); +``` + +Those are object_store's key names, the same as everywhere else Lance touches +storage. + +## Running it + + cd test/s3tables/catalog_duckdb_lance + (cd ../../../weed && go build .) # the harness runs this binary + go test -run TestDuckDBLance -v -timeout 30m . + +Skipped without Docker, in `-short` mode, and if the DuckDB image cannot load the +extension. diff --git a/test/s3tables/catalog_duckdb_lance/duckdb_lance_ops.sql b/test/s3tables/catalog_duckdb_lance/duckdb_lance_ops.sql new file mode 100644 index 000000000..2a64fd187 --- /dev/null +++ b/test/s3tables/catalog_duckdb_lance/duckdb_lance_ops.sql @@ -0,0 +1,48 @@ +-- Read a SeaweedFS Lance table from DuckDB. +-- +-- The lance extension reaches the data over S3 rather than through the +-- namespace, so what is exercised here is the layout on the S3 door and the +-- credentials, not the catalog protocol. That is the point: this is the client +-- that proves a Lance table stays readable without a catalog at all. +-- +-- Placeholders are substituted by the Go harness: __ENDPOINT__, __KEY__, +-- __SECRET__, __TABLE__ (s3://bucket/ns/table), __SUFFIXED__ (the same data at a +-- path ending in .lance). + +INSTALL lance; +LOAD lance; + +CREATE SECRET seaweedfs ( + TYPE lance, + ACCESS_KEY_ID '__KEY__', + SECRET_ACCESS_KEY '__SECRET__', + REGION 'us-east-1', + ENDPOINT '__ENDPOINT__', + ALLOW_HTTP true, + VIRTUAL_HOSTED_STYLE_REQUEST false +); + +-- 1. The table the catalog created, read by URI. A table bucket's layout is a +-- valid Lance dataset directory, which is what makes this possible. +SELECT 'scan_rows=' || count(*) FROM __lance_scan('__TABLE__'); + +-- 2. Its schema survived, vector column included. Ordered, because string_agg +-- over an unordered relation is free to return the names in any order. +SELECT 'scan_columns=' || string_agg(column_name, ',' ORDER BY column_name) +FROM (DESCRIBE SELECT * FROM __lance_scan('__TABLE__')); + +-- 3. A filter, so it is not only a full scan. +SELECT 'filtered_rows=' || count(*) FROM __lance_scan('__TABLE__') WHERE id < 5; + +-- 4. Vector search, which is what the format is for. No index is built here, so +-- this is a brute-force search; the ids nearest the query are what matters. +SELECT 'nearest=' || string_agg(id::VARCHAR, ',' ORDER BY _distance ASC, id ASC) +FROM lance_vector_search('__TABLE__', 'vector', + [1.0, 2.0, 3.0, 4.0, 5.0, 6.0, 7.0, 8.0], k := 3); + +-- 5. DuckDB's replacement scan recognises a path by its .lance suffix. Tables +-- this catalog creates deliberately have no suffix - the name is the table's +-- name, and a suffix would leak into ARNs and policies - so the bare +-- SELECT ... FROM 's3://…' form does not see them, and __lance_scan is the +-- way in. This asserts both halves so that a change upstream is noticed. +SELECT 'suffixed_rows=' || count(*) FROM '__SUFFIXED__'; diff --git a/test/s3tables/catalog_duckdb_lance/duckdb_lance_test.go b/test/s3tables/catalog_duckdb_lance/duckdb_lance_test.go new file mode 100644 index 000000000..997cab260 --- /dev/null +++ b/test/s3tables/catalog_duckdb_lance/duckdb_lance_test.go @@ -0,0 +1,390 @@ +// Package duckdblance reads SeaweedFS Lance tables from DuckDB's lance +// extension, the counterpart of the DuckDB Iceberg tests next door. +// +// DuckDB reaches the data over S3 rather than through the namespace, so what +// this proves is the other half of the design: a table bucket's layout is a +// valid Lance dataset directory, and a table stays readable with no catalog in +// the path at all. It also pins the one place that costs us - DuckDB's +// replacement scan recognises a dataset by its .lance suffix, which tables +// created through this catalog deliberately do not have. +package duckdblance + +import ( + "context" + "fmt" + "math/rand" + "net/http" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/test/testutil" +) + +const ( + // Tracks whatever DuckDB users are actually running; the lance extension is + // a core extension, so there is nothing to pin beyond the image. + duckDBImage = "duckdb/duckdb:latest" + // Seeds a dataset for DuckDB to read. + seedImage = "seaweedfs-lance-seed" + startupTimeout = 60 * time.Second + clientTimeout = 15 * time.Minute + seededRows = 128 +) + +// TestDuckDBLance reads a table this catalog created, from DuckDB. +func TestDuckDBLance(t *testing.T) { + if testing.Short() { + t.Skip("Skipping integration test in short mode") + } + if !hasDocker() { + t.Skip("Docker not available, skipping DuckDB Lance integration test") + } + + env := newEnvironment(t) + defer env.cleanup() + + env.start(t) + + bucket := "duckdb-" + randomSuffix() + env.createTableBucket(t, bucket) + + requireDuckDBLance(t) + buildSeedImage(t) + env.seedTable(t, bucket) + env.runDuckDB(t, bucket) +} + +type environment struct { + weedBinary string + dataDir string + bindIP string + + masterPort int + masterGrpcPort int + volumePort int + volumeGrpcPort int + filerPort int + filerGrpcPort int + s3Port int + s3GrpcPort int + lancePort int + + accessKey string + secretKey string + + weedCancel context.CancelFunc + weedCmd *exec.Cmd +} + +func newEnvironment(t *testing.T) *environment { + t.Helper() + + wd, err := os.Getwd() + if err != nil { + t.Fatalf("get working directory: %v", err) + } + seaweedDir := wd + for i := 0; i < 5; i++ { + if _, err := os.Stat(filepath.Join(seaweedDir, "go.mod")); err == nil { + break + } + seaweedDir = filepath.Dir(seaweedDir) + } + + weedBinary := filepath.Join(seaweedDir, "weed", "weed") + if info, statErr := os.Stat(weedBinary); statErr == nil && !info.IsDir() { + // `make test` builds first; a plain `go test` will otherwise drive a + // binary from days ago and report a pass for code it never ran. + t.Logf("using %s, built %s", weedBinary, info.ModTime().Format(time.RFC3339)) + } else { + weedBinary = "weed" + if _, err := exec.LookPath(weedBinary); err != nil { + t.Skip("weed binary not found, skipping integration test") + } + } + + dataDir, err := os.MkdirTemp("", "seaweed-duckdb-lance-test-*") + if err != nil { + t.Fatalf("create temp dir: %v", err) + } + + ports := testutil.MustAllocatePorts(t, 9) + return &environment{ + weedBinary: weedBinary, + dataDir: dataDir, + bindIP: testutil.FindBindIP(), + masterPort: ports[0], + masterGrpcPort: ports[1], + volumePort: ports[2], + volumeGrpcPort: ports[3], + filerPort: ports[4], + filerGrpcPort: ports[5], + s3Port: ports[6], + s3GrpcPort: ports[7], + lancePort: ports[8], + accessKey: "AKIAIOSFODNN7EXAMPLE", + secretKey: "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY", + } +} + +func (env *environment) start(t *testing.T) { + t.Helper() + + iamConfigPath, err := testutil.WriteIAMConfig(env.dataDir, env.accessKey, env.secretKey) + if err != nil { + t.Fatalf("write IAM config: %v", err) + } + + ctx, cancel := context.WithCancel(context.Background()) + env.weedCancel = cancel + + cmd := exec.CommandContext(ctx, env.weedBinary, "mini", + "-master.port", fmt.Sprintf("%d", env.masterPort), + "-master.port.grpc", fmt.Sprintf("%d", env.masterGrpcPort), + "-volume.port", fmt.Sprintf("%d", env.volumePort), + "-volume.port.grpc", fmt.Sprintf("%d", env.volumeGrpcPort), + "-filer.port", fmt.Sprintf("%d", env.filerPort), + "-filer.port.grpc", fmt.Sprintf("%d", env.filerGrpcPort), + "-s3.port", fmt.Sprintf("%d", env.s3Port), + "-s3.port.grpc", fmt.Sprintf("%d", env.s3GrpcPort), + "-s3.port.lance", fmt.Sprintf("%d", env.lancePort), + "-s3.config", iamConfigPath, + "-ip", env.bindIP, + "-ip.bind", "0.0.0.0", + "-dir", env.dataDir, + ) + cmd.Dir = env.dataDir + cmd.Stdout = os.Stdout + cmd.Stderr = os.Stderr + cmd.Env = append(os.Environ(), + "AWS_ACCESS_KEY_ID="+env.accessKey, + "AWS_SECRET_ACCESS_KEY="+env.secretKey, + ) + + if err := cmd.Start(); err != nil { + t.Fatalf("start SeaweedFS: %v", err) + } + env.weedCmd = cmd + + // The namespace answers /v1/table once it is serving, which is a cheaper + // readiness check than waiting on a bucket that does not exist yet. + url := fmt.Sprintf("http://%s:%d/v1/table", env.bindIP, env.lancePort) + if !waitForHTTP(url, startupTimeout) { + t.Fatalf("the Lance namespace did not become ready at %s", url) + } +} + +// waitForHTTP polls until the URL answers at all. An auth refusal counts: it +// means the server is up, which is the only thing being waited on. +func waitForHTTP(url string, timeout time.Duration) bool { + client := &http.Client{Timeout: 2 * time.Second} + deadline := time.Now().Add(timeout) + for time.Now().Before(deadline) { + resp, err := client.Get(url) + if err != nil { + time.Sleep(500 * time.Millisecond) + continue + } + status := resp.StatusCode + resp.Body.Close() + if status < 500 { + return true + } + time.Sleep(500 * time.Millisecond) + } + return false +} + +func (env *environment) cleanup() { + if env.weedCancel != nil { + env.weedCancel() + } + if env.weedCmd != nil { + _ = env.weedCmd.Wait() + } + if env.dataDir != "" { + _ = os.RemoveAll(env.dataDir) + } +} + +// createTableBucket makes the bucket LanceDB will read through, declared LANCE +// so the catalog refuses anything of another format in it. +func (env *environment) createTableBucket(t *testing.T, bucket string) { + t.Helper() + + ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) + defer cancel() + + cmd := exec.CommandContext(ctx, env.weedBinary, "shell", + fmt.Sprintf("-master=%s:%d.%d", env.bindIP, env.masterPort, env.masterGrpcPort), + ) + cmd.Stdin = strings.NewReader(fmt.Sprintf( + "s3tables.bucket -create -name %s -format LANCE -account 000000000000\nexit\n", bucket)) + out, err := cmd.CombinedOutput() + if err != nil { + t.Fatalf("create table bucket %s: %v\n%s", bucket, err, out) + } + // weed shell reports a command's own failure on stdout and still exits 0, so + // the exit code alone would let a missing bucket through and turn a setup + // failure into a confusing engine failure later. + if !env.tableBucketExists(t, bucket) { + t.Fatalf("table bucket %s was not created:\n%s", bucket, out) + } + t.Logf("created LANCE table bucket %s", bucket) +} + +// tableBucketExists asks the namespace, which lists table buckets at its root. +func (env *environment) tableBucketExists(t *testing.T, bucket string) bool { + t.Helper() + + url := fmt.Sprintf("http://%s:%d/v1/namespace/%s/exists", env.bindIP, env.lancePort, bucket) + resp, err := http.Post(url, "application/json", strings.NewReader("{}")) + if err != nil { + t.Fatalf("ask the namespace whether %s exists: %v", bucket, err) + } + defer resp.Body.Close() + return resp.StatusCode == http.StatusOK +} + +// runSpark runs the SQL driver inside the stock Spark image. The connector is +// pulled from Maven at submit time, the way the Iceberg Spark suite pulls its +// runtime, so nothing has to be built here. +// requireDuckDBLance skips rather than fails when the image cannot load the +// extension, so a DuckDB build without it reads as "not applicable" rather than +// as a broken catalog. +func requireDuckDBLance(t *testing.T) { + t.Helper() + + const ready = "lance extension ready" + cmd := exec.Command("docker", "run", "--rm", "--entrypoint", "duckdb", duckDBImage, + "-c", fmt.Sprintf("INSTALL lance; LOAD lance; SELECT '%s' AS marker;", ready)) + out, err := cmd.CombinedOutput() + if err != nil || !strings.Contains(string(out), ready) { + t.Skipf("DuckDB image cannot load the lance extension: %v\n%s", err, out) + } +} + +func buildSeedImage(t *testing.T) { + t.Helper() + + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute) + defer cancel() + + cmd := exec.CommandContext(ctx, "docker", "build", "-t", seedImage, "-f", "Dockerfile.seed", ".") + if out, err := cmd.CombinedOutput(); err != nil { + t.Fatalf("build the seed image: %v\n%s", err, out) + } +} + +// seedTable declares a table through the namespace and writes a dataset into +// it, which is the split the catalog serves: it records where a table lives and +// does not carry its data. +func (env *environment) seedTable(t *testing.T, bucket string) { + t.Helper() + + ctx, cancel := context.WithTimeout(context.Background(), clientTimeout) + defer cancel() + + cmd := exec.CommandContext(ctx, "docker", "run", "--rm", + "--add-host", "host.docker.internal:host-gateway", + "-e", "AWS_ACCESS_KEY_ID="+env.accessKey, + "-e", "AWS_SECRET_ACCESS_KEY="+env.secretKey, + "-e", "AWS_REGION=us-east-1", + seedImage, + "python3", "/app/seed_table.py", + "--namespace-url", fmt.Sprintf("http://host.docker.internal:%d", env.lancePort), + "--s3-endpoint", fmt.Sprintf("http://host.docker.internal:%d", env.s3Port), + "--bucket", bucket, + "--rows", fmt.Sprintf("%d", seededRows), + "--access-key", env.accessKey, + "--secret-key", env.secretKey, + ) + out, err := cmd.CombinedOutput() + t.Logf("seeder:\n%s", out) + if err != nil { + t.Fatalf("seeding the table failed: %v", err) + } +} + +// runDuckDB substitutes the endpoint and paths into the SQL and runs it. +func (env *environment) runDuckDB(t *testing.T, bucket string) { + t.Helper() + + sqlBytes, err := os.ReadFile("duckdb_lance_ops.sql") + if err != nil { + t.Fatalf("read the SQL: %v", err) + } + table := fmt.Sprintf("s3://%s/ml/embeddings", bucket) + sql := strings.NewReplacer( + "__ENDPOINT__", fmt.Sprintf("http://host.docker.internal:%d", env.s3Port), + "__KEY__", env.accessKey, + "__SECRET__", env.secretKey, + "__TABLE__", table, + "__SUFFIXED__", table+"-direct.lance", + ).Replace(string(sqlBytes)) + + ctx, cancel := context.WithTimeout(context.Background(), clientTimeout) + defer cancel() + + cmd := exec.CommandContext(ctx, "docker", "run", "--rm", "-i", + "--add-host", "host.docker.internal:host-gateway", + "--entrypoint", "duckdb", duckDBImage, "-c", sql) + out, err := cmd.CombinedOutput() + t.Logf("DuckDB output:\n%s", out) + if err != nil { + t.Fatalf("the DuckDB query failed: %v", err) + } + + output := string(out) + for _, want := range []string{ + fmt.Sprintf("scan_rows=%d", seededRows), + "scan_columns=id,title,vector", + "filtered_rows=5", + "nearest=1,0,2", + fmt.Sprintf("suffixed_rows=%d", seededRows), + } { + if !strings.Contains(output, want) { + t.Fatalf("DuckDB did not report %q", want) + } + } + + // The other half of the suffix rule: a table this catalog created has no + // .lance suffix, so DuckDB's replacement scan does not see it and + // __lance_scan is the way in. If this ever starts working, the docs saying + // otherwise are wrong. + bare := exec.CommandContext(ctx, "docker", "run", "--rm", "-i", + "--add-host", "host.docker.internal:host-gateway", + "--entrypoint", "duckdb", duckDBImage, + "-c", fmt.Sprintf("INSTALL lance; LOAD lance; SELECT count(*) FROM '%s';", table)) + bareOut, bareErr := bare.CombinedOutput() + // DuckDB exits nonzero for any error, so the exit status alone does not tell + // "the replacement scan refused the path" from "the query never ran": require + // the catalog error either way. + if !strings.Contains(string(bareOut), "does not exist") { + t.Fatalf("a suffix-less path did not fail the way the docs say it does (%v); "+ + "if the replacement scan now reads it, update them:\n%s", bareErr, bareOut) + } + t.Logf("a suffix-less path is not seen by the replacement scan, as expected") +} + +// hasDocker reports whether a Docker daemon answers. Bounded, because an +// unhealthy daemon makes `docker version` hang, and this runs before the test +// has a timeout of its own: better to skip than to eat the whole budget. +func hasDocker() bool { + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + return exec.CommandContext(ctx, "docker", "version").Run() == nil +} + +func randomSuffix() string { + const charset = "abcdefghijklmnopqrstuvwxyz0123456789" + suffix := make([]byte, 8) + for i := range suffix { + suffix[i] = charset[rand.Intn(len(charset))] + } + return string(suffix) +} diff --git a/test/s3tables/catalog_duckdb_lance/seed_table.py b/test/s3tables/catalog_duckdb_lance/seed_table.py new file mode 100644 index 000000000..051f2f519 --- /dev/null +++ b/test/s3tables/catalog_duckdb_lance/seed_table.py @@ -0,0 +1,79 @@ +#!/usr/bin/env python3 +"""Declare a table through the namespace and write a dataset into it. + +Also writes the same rows to a second path whose name ends in `.lance`, because +DuckDB's replacement scan recognises a dataset by that suffix and tables created +through this catalog do not have one. The test asserts both behaviours so a +change upstream is noticed rather than silently making the docs wrong. +""" + +import argparse +import sys +import warnings + +warnings.filterwarnings("ignore") + +import lance +import lance_namespace as ln +import pyarrow as pa + +DIM = 8 + + +def rows(count): + return pa.table( + { + "id": pa.array(list(range(count)), type=pa.int64()), + "title": pa.array([f"row-{i}" for i in range(count)]), + "vector": pa.array( + [[float(i) + d for d in range(DIM)] for i in range(count)], + type=pa.list_(pa.float32(), DIM), + ), + } + ) + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--namespace-url", required=True) + parser.add_argument("--s3-endpoint", required=True) + parser.add_argument("--bucket", required=True) + parser.add_argument("--namespace", default="ml") + parser.add_argument("--table", default="embeddings") + parser.add_argument("--rows", type=int, default=128) + parser.add_argument("--access-key", default="any") + parser.add_argument("--secret-key", default="any") + args = parser.parse_args() + + storage = { + "aws_endpoint": args.s3_endpoint, + "allow_http": "true", + "aws_access_key_id": args.access_key, + "aws_secret_access_key": args.secret_key, + "aws_region": "us-east-1", + } + + ns = ln.connect("rest", {"uri": args.namespace_url}) + ns.create_namespace(ln.CreateNamespaceRequest(id=[args.bucket], mode="EXIST_OK")) + ns.create_namespace( + ln.CreateNamespaceRequest(id=[args.bucket, args.namespace], mode="EXIST_OK") + ) + declared = ns.declare_table( + ln.DeclareTableRequest(id=[args.bucket, args.namespace, args.table]) + ) + lance.write_dataset(rows(args.rows), declared.location, storage_options=storage, + mode="overwrite") + print(f"seeded {args.rows} rows at {declared.location}") + + # The same data under a name DuckDB's replacement scan recognises. Written + # directly rather than declared, because a table name containing a dot is + # not a valid catalog name. + suffixed = f"s3://{args.bucket}/{args.namespace}/{args.table}-direct.lance" + lance.write_dataset(rows(args.rows), suffixed, storage_options=storage, + mode="overwrite") + print(f"seeded {args.rows} rows at {suffixed}") + return 0 + + +if __name__ == "__main__": + sys.exit(main())