From 8a532cc0cffd5d17ac586bd861878334e5de2b1f Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Thu, 20 Aug 2026 23:46:32 -0700 Subject: [PATCH] mini: state the format of a -tableBucket, do not infer it (#10851) * mini: state the format of a -tableBucket, do not infer it A table bucket holds one format and that format decides which catalog can serve it, but the flag only took names. The format came from miniTableBucketFormat(): Iceberg whenever its port was up, Lance only when it was not. So -tableBucket=vectors on a default mini quietly made an ICEBERG bucket that the Lance namespace then refused every table in, and the only way to get a Lance one was -s3.port.iceberg=0, which buys it by deleting the other catalog. One flag, two meanings, decided by an unrelated port. Each entry is now name[:FORMAT], unsuffixed meaning ICEBERG as before: weed mini -tableBucket=warehouse,vectors:LANCE Both catalogs stay up and both buckets are reachable. A name whose format has no endpoint here is skipped with a warning rather than created out of reach, and the Iceberg-only S3_TABLE_BUCKET default-routing hint gets the Iceberg names alone, without their suffixes. * mini: do not reuse a table bucket that holds another format CreateTableBucket answers BucketAlreadyExists on the name alone, so -tableBucket=vectors against a bucket created as LANCE logged "already exists" and moved on, and the Iceberg default-warehouse hint then pointed at it. Every table create against that catalog fails with "table bucket vectors holds LANCE tables", far from the flag that chose it. ensureMiniTableBuckets now reads the format of a bucket it did not create, warns when it is not the one asked for, and returns only the buckets that hold what was requested. S3_TABLE_BUCKET is seeded from that list, so an unprefixed Iceberg request falls back to its own default rather than committing into a Lance bucket. A bucket predating declared formats reports an empty one and still accepts either. * mini: normalize S3_TABLE_BUCKET whichever way the spec arrived The rewrite that keeps Lance names out of the Iceberg default warehouse only ran when the flag supplied the spec. Set the variable directly, as the docker quickstart does, and it reached the catalog untouched: S3_TABLE_BUCKET= vectors:LANCE,warehouse made the unprefixed default the literal string "vectors:LANCE", a bucket no lookup finds, while warehouse sat behind it. The variable is both mini's input and the catalog's routing hint, so it is now always rewritten from the buckets that came back holding Iceberg tables, and unset when there are none rather than left pointing somewhere stale. * mini: reuse a table bucket only when its format reads back An ordinary S3 bucket wearing the name answers CreateTableBucket with the same BucketAlreadyExists as a table bucket does, and the format lookup that follows returned "" for a failed read exactly as it does for a bucket predating declared formats. So -bucket=data -tableBucket=data reported nothing and published data as the Iceberg default warehouse, where every unprefixed request 404s on a bucket that is not a catalog. The lookup now returns its error, and only a bucket that reads back as the format asked for is reused. Anything else is left alone with a warning naming why, rather than routed to and discovered later. --- README.md | 2 +- weed/command/mini.go | 174 ++++++++++++++++++++++++------- weed/command/mini_bucket_test.go | 93 +++++++++++++++++ 3 files changed, 229 insertions(+), 40 deletions(-) diff --git a/README.md b/README.md index dfa7f8030..374121a8b 100644 --- a/README.md +++ b/README.md @@ -90,7 +90,7 @@ S3_BUCKET=my-bucket \ ./weed mini -dir=/data ``` -That's it — the S3 endpoint is at http://localhost:8333, `my-bucket` already exists, and `admin`/`secret` are valid credentials. `S3_BUCKET` accepts a comma-separated list (e.g. `raw,processed`); use `S3_TABLE_BUCKET` for S3 Tables (Iceberg) buckets. Drop any of the env vars to skip that piece (no AWS keys → S3 runs in unauthenticated "Allow All" mode for development). +That's it — the S3 endpoint is at http://localhost:8333, `my-bucket` already exists, and `admin`/`secret` are valid credentials. `S3_BUCKET` accepts a comma-separated list (e.g. `raw,processed`); use `S3_TABLE_BUCKET` for S3 Tables buckets, each `name` or `name:FORMAT` where the format is `ICEBERG` (the default) or `LANCE`. Drop any of the env vars to skip that piece (no AWS keys → S3 runs in unauthenticated "Allow All" mode for development). The same command starts everything else too: - **S3 Endpoint**: http://localhost:8333 diff --git a/weed/command/mini.go b/weed/command/mini.go index 9b8020c67..0e9d8f45f 100644 --- a/weed/command/mini.go +++ b/weed/command/mini.go @@ -394,7 +394,7 @@ var ( miniS3AllowDeleteBucketNotEmpty = cmdMini.Flag.Bool("s3.allowDeleteBucketNotEmpty", true, "allow recursive deleting all entries along with bucket") miniS3AutoCreateBucket = cmdMini.Flag.Bool("s3.autoCreateBucket", true, "create the bucket on upload if it does not exist, for admin identities only") miniBucket = cmdMini.Flag.String("bucket", "", "comma-separated S3 bucket names to create on startup if they do not already exist; leave empty to skip. Falls back to S3_BUCKET env var.") - miniTableBucket = cmdMini.Flag.String("tableBucket", "", "comma-separated S3 Tables bucket names to create on startup if they do not already exist; leave empty to skip. Falls back to S3_TABLE_BUCKET env var.") + miniTableBucket = cmdMini.Flag.String("tableBucket", "", "comma-separated S3 Tables buckets to create on startup if they do not already exist, each name[:FORMAT] with FORMAT one of ICEBERG (default) or LANCE, e.g. warehouse,vectors:LANCE; leave empty to skip. Falls back to S3_TABLE_BUCKET env var.") ) // getBindIp determines the bind IP address based on miniIp and miniBindIp flags @@ -1367,14 +1367,20 @@ func runMini(cmd *Command, args []string) bool { tableBucketSpec := *miniTableBucket if tableBucketSpec == "" { tableBucketSpec = os.Getenv("S3_TABLE_BUCKET") - } else if os.Getenv("S3_TABLE_BUCKET") == "" { - // The catalog routes unprefixed requests to the first S3_TABLE_BUCKET - // entry; let the -tableBucket flag mean the same thing. - os.Setenv("S3_TABLE_BUCKET", tableBucketSpec) } - if err := ensureMiniTableBuckets(tableBucketSpec); err != nil { + ready, err := ensureMiniTableBuckets(parseTableBucketList(tableBucketSpec)) + if err != nil { glog.Warningf("failed to ensure table buckets %q: %v", tableBucketSpec, err) } + // The Iceberg catalog routes unprefixed requests to the first + // S3_TABLE_BUCKET entry and looks the name up as given, so rewrite the + // variable whichever way the spec arrived: a Lance entry or a :FORMAT + // suffix left in place names a bucket the catalog cannot find. + if names := icebergRoutingNames(ready); len(names) > 0 { + os.Setenv("S3_TABLE_BUCKET", strings.Join(names, ",")) + } else { + os.Unsetenv("S3_TABLE_BUCKET") + } // Print welcome message after all services are running printWelcomeMessage() @@ -2012,67 +2018,157 @@ func ensureMiniBuckets(bucketSpec string) error { }) } -// ensureMiniTableBuckets creates each named S3 Tables bucket on the embedded -// filer if it does not already exist. bucketSpec is comma-separated; whitespace -// is trimmed and duplicates are dropped. Per-bucket failures are logged so one -// bad name does not block the rest. Buckets are owned by s3tables.DefaultAccountID -// since mini does not yet model multi-account ownership. -func ensureMiniTableBuckets(bucketSpec string) error { - names := parseBucketList(bucketSpec) - if len(names) == 0 { - return nil - } +// tableBucketEntry is one -tableBucket entry: a bucket name and the table +// format that bucket will hold. +type tableBucketEntry struct { + name string + format string +} - // A bucket holds one format, and the format decides which catalog serves it. - // Creating one in a format this mini does not serve leaves a bucket no - // client can reach, so take the format from the endpoint that is running. - format := miniTableBucketFormat() - if format == "" { - glog.Warningf("not creating table buckets %q: neither the Iceberg nor the Lance endpoint is enabled, so nothing could reach them", bucketSpec) - return nil +// ensureMiniTableBuckets creates each named S3 Tables bucket on the embedded +// filer if it does not already exist, and returns the ones that now hold the +// format that was asked for. Per-bucket failures are logged so one bad name +// does not block the rest. Buckets are owned by s3tables.DefaultAccountID +// since mini does not yet model multi-account ownership. +func ensureMiniTableBuckets(buckets []tableBucketEntry) ([]tableBucketEntry, error) { + // A bucket holds one format, and the format decides which catalog serves + // it, so one created in a format this mini does not serve is a bucket no + // client can reach. + var servable []tableBucketEntry + for _, bucket := range buckets { + if !miniServesTableFormat(bucket.format) { + // An unsuffixed name means Iceberg, so on a Lance-only mini say how to ask. + hint := "" + if bucket.format == s3tables.FormatIceberg && miniServesTableFormat(s3tables.FormatLance) { + hint = fmt.Sprintf("; name it %s:%s for the Lance namespace", bucket.name, s3tables.FormatLance) + } + glog.Warningf("not creating table bucket %s: no %s endpoint is enabled, so nothing could reach it%s", bucket.name, bucket.format, hint) + continue + } + servable = append(servable, bucket) + } + if len(servable) == 0 { + return nil, nil } filerAddress := pb.NewServerAddress(*miniIp, *miniFilerOptions.port, *miniFilerOptions.portGrpc) grpcDialOption := security.LoadClientTLS(util.GetViper(), "grpc.client") - return pb.WithGrpcFilerClient(false, 0, filerAddress, grpcDialOption, func(client filer_pb.SeaweedFilerClient) error { + var ready []tableBucketEntry + err := pb.WithGrpcFilerClient(false, 0, filerAddress, grpcDialOption, func(client filer_pb.SeaweedFilerClient) error { manager := s3tables.NewManager() mgrClient := s3tables.NewManagerClient(client) - for _, name := range names { + for _, bucket := range servable { ctx, cancel := context.WithTimeout(miniClientsCtx(), 5*time.Second) - req := &s3tables.CreateTableBucketRequest{Name: name, Format: format} + req := &s3tables.CreateTableBucketRequest{Name: bucket.name, Format: bucket.format} var resp s3tables.CreateTableBucketResponse err := manager.Execute(ctx, mgrClient, "CreateTableBucket", req, &resp, s3tables.DefaultAccountID) cancel() if err == nil { - glog.V(0).Infof("created %s table bucket %s", format, name) + glog.V(0).Infof("created %s table bucket %s", bucket.format, bucket.name) + ready = append(ready, bucket) continue } var s3Err *s3tables.S3TablesError if errors.As(err, &s3Err) && s3Err.Type == s3tables.ErrCodeBucketAlreadyExists { - glog.V(0).Infof("table bucket %s already exists", name) + // The name being taken says nothing about what holds it: it + // may be the other format, or an ordinary S3 bucket, which + // answers the same conflict. Only reuse what reads back as + // the format asked for, since the rest fails later at the + // client instead of here. + existing, readErr := existingTableBucketFormat(manager, mgrClient, bucket.name) + switch { + case readErr != nil: + glog.Warningf("%s already exists but does not read back as a table bucket: %v; leaving it alone", bucket.name, readErr) + case existing != "" && existing != bucket.format: + glog.Warningf("table bucket %s already exists holding %s tables, not %s; leaving it alone", bucket.name, existing, bucket.format) + default: + glog.V(0).Infof("table bucket %s already exists", bucket.name) + ready = append(ready, bucket) + } continue } - glog.Warningf("create table bucket %s: %v", name, err) + glog.Warningf("create table bucket %s: %v", bucket.name, err) } return nil }) + return ready, err } -// miniTableBucketFormat is the format a pre-created table bucket should hold: -// Iceberg when its catalog is running, else Lance, else none because neither -// server is up. -func miniTableBucketFormat() string { +// icebergRoutingNames is the bare names of the buckets holding Iceberg tables, +// which are the only ones the Iceberg catalog can take as its default warehouse. +func icebergRoutingNames(buckets []tableBucketEntry) []string { + var names []string + for _, bucket := range buckets { + if bucket.format == s3tables.FormatIceberg { + names = append(names, bucket.name) + } + } + return names +} + +// existingTableBucketFormat is the format the named table bucket holds. An empty +// format with no error is a bucket predating declared formats, which accepts +// either; an error means nothing was confirmed, an ordinary S3 bucket wearing +// the name included. +func existingTableBucketFormat(manager *s3tables.Manager, client *s3tables.ManagerClient, name string) (string, error) { + arn, err := s3tables.BuildBucketARN(s3tables.DefaultRegion, s3tables.DefaultAccountID, name) + if err != nil { + return "", err + } + ctx, cancel := context.WithTimeout(miniClientsCtx(), 5*time.Second) + defer cancel() + var resp s3tables.GetTableBucketResponse + if err := manager.Execute(ctx, client, "GetTableBucket", &s3tables.GetTableBucketRequest{TableBucketARN: arn}, &resp, s3tables.DefaultAccountID); err != nil { + return "", err + } + return resp.Format, nil +} + +// miniServesTableFormat reports whether the endpoint that serves format is +// running here: the Iceberg catalog for ICEBERG, the Lance namespace for LANCE. +func miniServesTableFormat(format string) bool { if miniEnableS3 == nil || !*miniEnableS3 { - return "" + return false } - if miniS3Options.portIceberg != nil && *miniS3Options.portIceberg > 0 { - return s3tables.FormatIceberg + switch format { + case s3tables.FormatIceberg: + return miniS3Options.portIceberg != nil && *miniS3Options.portIceberg > 0 + case s3tables.FormatLance: + return miniS3Options.portLance != nil && *miniS3Options.portLance > 0 } - if miniS3Options.portLance != nil && *miniS3Options.portLance > 0 { - return s3tables.FormatLance + return false +} + +// parseTableBucketList splits a comma-separated table bucket spec into +// deduplicated name[:FORMAT] entries, in the order given. The format is stated +// rather than read off whichever catalog happens to be listening, so one spec +// means one thing on every mini. +func parseTableBucketList(spec string) []tableBucketEntry { + if spec == "" { + return nil } - return "" + seen := make(map[string]bool) + var buckets []tableBucketEntry + for _, raw := range strings.Split(spec, ",") { + name, rawFormat, _ := strings.Cut(raw, ":") + name = strings.TrimSpace(name) + if name == "" || seen[name] { + continue + } + format := s3tables.FormatIceberg + if rawFormat = strings.TrimSpace(rawFormat); rawFormat != "" { + normalized, ok := s3tables.NormalizeFormat(rawFormat) + if !ok { + glog.Warningf("not creating table bucket %s: unsupported format %q", name, rawFormat) + continue + } + format = normalized + } + seen[name] = true + buckets = append(buckets, tableBucketEntry{name: name, format: format}) + } + return buckets } // parseBucketList splits a comma-separated bucket spec into a deduplicated list diff --git a/weed/command/mini_bucket_test.go b/weed/command/mini_bucket_test.go index 4a2ac66db..78cd48aaf 100644 --- a/weed/command/mini_bucket_test.go +++ b/weed/command/mini_bucket_test.go @@ -3,6 +3,8 @@ package command import ( "reflect" "testing" + + "github.com/seaweedfs/seaweedfs/weed/s3api/s3tables" ) func TestParseBucketList(t *testing.T) { @@ -28,3 +30,94 @@ func TestParseBucketList(t *testing.T) { }) } } + +func TestParseTableBucketList(t *testing.T) { + tests := []struct { + name string + in string + want []tableBucketEntry + }{ + {"empty", "", nil}, + {"defaults to iceberg", "warehouse", []tableBucketEntry{{"warehouse", s3tables.FormatIceberg}}}, + {"explicit format", "vectors:LANCE", []tableBucketEntry{{"vectors", s3tables.FormatLance}}}, + {"format is case insensitive", "vectors:lance", []tableBucketEntry{{"vectors", s3tables.FormatLance}}}, + {"mixed formats", "warehouse,vectors:LANCE", []tableBucketEntry{{"warehouse", s3tables.FormatIceberg}, {"vectors", s3tables.FormatLance}}}, + {"trims whitespace", " warehouse , vectors : LANCE ", []tableBucketEntry{{"warehouse", s3tables.FormatIceberg}, {"vectors", s3tables.FormatLance}}}, + {"drops unsupported format", "warehouse,vectors:DELTA", []tableBucketEntry{{"warehouse", s3tables.FormatIceberg}}}, + {"empty format is the default", "warehouse:", []tableBucketEntry{{"warehouse", s3tables.FormatIceberg}}}, + {"dedupes by name", "vectors:LANCE,vectors:ICEBERG", []tableBucketEntry{{"vectors", s3tables.FormatLance}}}, + {"drops empty entries", "one,,two:LANCE,", []tableBucketEntry{{"one", s3tables.FormatIceberg}, {"two", s3tables.FormatLance}}}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got := parseTableBucketList(tt.in) + if !reflect.DeepEqual(got, tt.want) { + t.Errorf("parseTableBucketList(%q) = %v, want %v", tt.in, got, tt.want) + } + }) + } +} + +// A bucket's format decides which endpoint has to be up for it, so turning one +// catalog off must not change what the other format means. +func TestMiniServesTableFormat(t *testing.T) { + icebergPort, lancePort := *miniS3Options.portIceberg, *miniS3Options.portLance + t.Cleanup(func() { + *miniS3Options.portIceberg, *miniS3Options.portLance = icebergPort, lancePort + }) + + *miniS3Options.portIceberg, *miniS3Options.portLance = 8181, 9101 + if !miniServesTableFormat(s3tables.FormatIceberg) || !miniServesTableFormat(s3tables.FormatLance) { + t.Errorf("both endpoints up: want both formats served") + } + + *miniS3Options.portIceberg = 0 + if miniServesTableFormat(s3tables.FormatIceberg) { + t.Errorf("iceberg port 0: want ICEBERG unserved") + } + if !miniServesTableFormat(s3tables.FormatLance) { + t.Errorf("iceberg port 0: want LANCE still served") + } + + *miniS3Options.portLance = 0 + if miniServesTableFormat(s3tables.FormatLance) { + t.Errorf("lance port 0: want LANCE unserved") + } + + if miniServesTableFormat("DELTA") { + t.Errorf("unknown format: want unserved") + } +} + +// The Iceberg catalog looks its default warehouse up by name, so a Lance entry +// or a leftover :FORMAT suffix in S3_TABLE_BUCKET names a bucket it cannot find. +func TestIcebergRoutingNames(t *testing.T) { + tests := []struct { + name string + in []tableBucketEntry + want []string + }{ + {"none", nil, nil}, + {"iceberg only", []tableBucketEntry{{"warehouse", s3tables.FormatIceberg}}, []string{"warehouse"}}, + {"lance is not routable", []tableBucketEntry{{"vectors", s3tables.FormatLance}}, nil}, + {"lance first", []tableBucketEntry{{"vectors", s3tables.FormatLance}, {"warehouse", s3tables.FormatIceberg}}, []string{"warehouse"}}, + {"order preserved", []tableBucketEntry{{"raw", s3tables.FormatIceberg}, {"vectors", s3tables.FormatLance}, {"curated", s3tables.FormatIceberg}}, []string{"raw", "curated"}}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got := icebergRoutingNames(tt.in) + if !reflect.DeepEqual(got, tt.want) { + t.Errorf("icebergRoutingNames(%v) = %v, want %v", tt.in, got, tt.want) + } + }) + } +} + +// A spec arriving through S3_TABLE_BUCKET carries suffixes too, and the catalog +// reads that same variable, so parsing and routing have to agree on the name. +func TestIcebergRoutingNamesFromEnvSpec(t *testing.T) { + got := icebergRoutingNames(parseTableBucketList("vectors:LANCE,warehouse")) + if want := []string{"warehouse"}; !reflect.DeepEqual(got, want) { + t.Errorf("routing names for env spec = %v, want %v", got, want) + } +}