Files
Chris Lu bea10e269f iceberg/s3tables: confine stored metadataLocation to the authorized table bucket (#11292)
* iceberg: confine commit/transaction/view-update write paths to authorized bucket

The create, register, and createView handlers already confine the client-
supplied metadata location to the caller table bucket and reject ".."
segments. The commit, create-on-commit, transaction, and view-update paths
read the stored metadataLocation back from the catalog and skipped the same
guard, so a location poisoned via the raw S3Tables UpdateTable API (which
persists metadataLocation verbatim) could escape the caller bucket through
a ".." segment that path.Join collapses in saveMetadataBlob.

Add confineMetadataLocation and apply it after parseS3Location on every
commit/update/transaction/view write path, mirroring the create/register/
createView check. Reject with 400 so a poisoned stored location fails the
commit instead of writing into another tenant bucket tree.

* s3tables: validate metadataLocation at the store layer

The raw S3Tables API (CreateTable, RegisterTable, UpdateTable, CreateView,
UpdateView) persisted the client-supplied metadataLocation verbatim with no
bucket-confinement or traversal check, so a caller could store a location
pointing outside its own bucket. The Iceberg REST gateway commit paths then
read that stored value back and wrote through it.

Add ValidateMetadataLocation and call it in every s3tables store handler
that accepts a metadataLocation, rejecting locations whose bucket differs
from the caller table bucket or whose path contains traversal segments. This
prevents a poisoned location from ever being persisted, complementing the
per-write-path guard added to the Iceberg commit handlers.

* iceberg/s3tables: validate location before repair and after idempotency check

Address review feedback:
- Move the commit-path confinement check ahead of repairManifests so a
  poisoned stored location cannot reach manifest repair I/O before the
  commit is rejected.
- Move ValidateMetadataLocation in CreateTable/CreateView to after the
  existing-resource check so idempotent retries that do not consume the
  requested location are not rejected for an unused bad location.
- Assert HTTP 400 in the cross-tenant reproduction tests so an unrelated
  failure cannot satisfy them.

* iceberg: confine staged metadata location before load in create-on-commit

The create-on-commit path parsed the staged metadata location from the
stage-create marker and called loadMetadataFile before validating that the
staged bucket/path stay within the authorized bucket. Add the same
confineMetadataLocation guard before the read so a tampered marker cannot
direct a cross-tenant metadata read.

* iceberg/s3tables: reject bucket-only metadata locations

ValidateMetadataLocation and confineMetadataLocation accepted s3://bucket
with an empty table path. metadataDirPath then maps every such table to the
shared <TablesPath>/<bucket>/metadata directory, so tables could overwrite
or read each other's metadata files. Require a non-empty table path in both
validators; the empty-location case (where the catalog derives one) is
unaffected.

* iceberg/s3tables: reject slash-only table paths in location validation

s3://bkt/// parses to tablePath="/" which passed the empty-string check
but path.Join cleans it away, mapping to the bucket-level metadata
directory shared across tables. Update isValidTablePath to require at
least one non-empty segment and mirror the same check in
ValidateMetadataLocation, closing the gap in all callers.
2026-09-13 13:48:13 -07:00

536 lines
17 KiB
Go

package s3tables
import (
"crypto/rand"
"encoding/hex"
"fmt"
"net/url"
"path"
"regexp"
"strings"
"time"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
)
const (
bucketNamePatternStr = `[a-z0-9-]+`
tableNamespacePatternStr = `[a-z0-9_.-]+`
tableNamePatternStr = `[a-z0-9_-]+`
)
const (
tableObjectRootDirName = ".objects"
)
// ARNPartitionPatternStr matches any AWS partition, not just the commercial
// one: aws-cn and aws-us-gov ARNs are valid and must reach the handlers.
const ARNPartitionPatternStr = `aws[-a-z0-9]*`
// ARNPrefixPatternStr is the leading, partition-tolerant part of every S3
// Tables ARN. The HTTP router shares it so routes and parsing agree.
const ARNPrefixPatternStr = `arn:` + ARNPartitionPatternStr + `:s3tables`
var (
bucketARNPattern = regexp.MustCompile(`^` + ARNPrefixPatternStr + `:[^:]*:[^:]*:bucket/(` + bucketNamePatternStr + `)$`)
tableARNPattern = regexp.MustCompile(`^` + ARNPrefixPatternStr + `:[^:]*:[^:]*:bucket/(` + bucketNamePatternStr + `)/table/(` + tableNamespacePatternStr + `)/(` + tableNamePatternStr + `)$`)
tagPattern = regexp.MustCompile(`^([\p{L}\p{Z}\p{N}_.:/=+\-@]*)$`)
)
// ARN parsing functions
// parseBucketNameFromARN extracts bucket name from table bucket ARN
// ARN format: arn:aws:s3tables:{region}:{account}:bucket/{bucket-name}
func parseBucketNameFromARN(arn string) (string, error) {
matches := bucketARNPattern.FindStringSubmatch(arn)
if len(matches) != 2 {
return "", fmt.Errorf("invalid bucket ARN: %s", arn)
}
bucketName := matches[1]
if !isValidBucketName(bucketName) {
return "", fmt.Errorf("invalid bucket name in ARN: %s", bucketName)
}
return bucketName, nil
}
// ParseBucketNameFromARN is a wrapper to validate bucket ARN for other packages.
func ParseBucketNameFromARN(arn string) (string, error) {
return parseBucketNameFromARN(arn)
}
// IsValidBucketName is a wrapper to validate a table bucket name for other packages.
func IsValidBucketName(name string) bool {
return isValidBucketName(name)
}
// parseTableFromARN extracts bucket name, namespace, and table name from ARN
// ARN format: arn:aws:s3tables:{region}:{account}:bucket/{bucket-name}/table/{namespace}/{table-name}
func parseTableFromARN(arn string) (bucketName, namespace, tableName string, err error) {
matches := tableARNPattern.FindStringSubmatch(arn)
if len(matches) != 4 {
return "", "", "", fmt.Errorf("invalid table ARN: %s", arn)
}
// Validate bucket name
bucketName = matches[1]
if err := validateBucketName(bucketName); err != nil {
return "", "", "", fmt.Errorf("invalid bucket name in ARN: %v", err)
}
namespace, err = validateNamespace([]string{matches[2]})
if err != nil {
return "", "", "", fmt.Errorf("invalid namespace in ARN: %v", err)
}
// URL decode and validate the table name from the ARN path component
tableNameUnescaped, err := url.PathUnescape(matches[3])
if err != nil {
return "", "", "", fmt.Errorf("invalid table name encoding in ARN: %v", err)
}
if _, err := validateTableName(tableNameUnescaped); err != nil {
return "", "", "", fmt.Errorf("invalid table name in ARN: %v", err)
}
return bucketName, namespace, tableNameUnescaped, nil
}
// Path helpers
// GetTableBucketPath returns the filer path for a table bucket
func GetTableBucketPath(bucketName string) string {
return path.Join(TablesPath, bucketName)
}
// GetNamespacePath returns the filer path for a namespace
func GetNamespacePath(bucketName, namespace string) string {
return path.Join(TablesPath, bucketName, namespace)
}
// GetTablePath returns the filer path for a table
func GetTablePath(bucketName, namespace, tableName string) string {
return path.Join(TablesPath, bucketName, namespace, tableName)
}
// TableDataDirFromMetadataLocation maps a table's s3:// metadata location to the
// filer directory holding its data. A renamed table is catalog-only, so its data
// stays at the original location while its catalog entry moves; this lets a drop
// purge the real data instead of the now-empty catalog path.
func TableDataDirFromMetadataLocation(metadataLocation string) string {
loc := strings.TrimSuffix(metadataLocation, "/")
if idx := strings.LastIndex(loc, "/metadata/"); idx != -1 {
loc = loc[:idx]
}
loc = strings.TrimPrefix(loc, "s3://")
if loc == "" {
return ""
}
return path.Join(TablesPath, loc)
}
// ValidateMetadataLocation checks that an s3:// metadata location stays within
// the authorized table bucket and rejects traversal segments that path.Join
// would collapse to escape the bucket directory. Empty locations are allowed
// (the catalog derives one). A non-empty location must include a table path so
// its metadata directory is table-specific, not shared at the bucket level.
func ValidateMetadataLocation(metadataLocation, bucketName string) error {
if metadataLocation == "" {
return nil
}
bucket, tablePath, err := parseS3Location(metadataLocation)
if err != nil {
return err
}
if bucket != bucketName {
return fmt.Errorf("metadata location must be within bucket %s", bucketName)
}
hasSegment := false
for _, segment := range strings.Split(tablePath, "/") {
if segment == "" {
continue
}
if segment == "." || segment == ".." || strings.ContainsAny(segment, "\\\x00") {
return fmt.Errorf("invalid metadata location path")
}
hasSegment = true
}
if !hasSegment {
return fmt.Errorf("metadata location must include a table path")
}
return nil
}
func parseS3Location(location string) (bucket, tablePath string, err error) {
if !strings.HasPrefix(location, "s3://") {
return "", "", fmt.Errorf("unsupported location: %s", location)
}
trimmed := strings.TrimPrefix(location, "s3://")
trimmed = strings.TrimSuffix(trimmed, "/")
if trimmed == "" {
return "", "", fmt.Errorf("invalid location: %s", location)
}
parts := strings.SplitN(trimmed, "/", 2)
bucket = parts[0]
if bucket == "" {
return "", "", fmt.Errorf("invalid location bucket: %s", location)
}
if len(parts) == 2 {
tablePath = parts[1]
}
return bucket, tablePath, nil
}
// GetTableObjectRootDir returns the root path for table bucket object storage
func GetTableObjectRootDir() string {
return path.Join(TablesPath, tableObjectRootDirName)
}
// GetTableObjectBucketPath returns the filer path for table bucket object storage
func GetTableObjectBucketPath(bucketName string) string {
return path.Join(GetTableObjectRootDir(), bucketName)
}
// Metadata structures
type tableBucketMetadata struct {
Name string `json:"name"`
CreatedAt time.Time `json:"createdAt"`
OwnerAccountID string `json:"ownerAccountId"`
Format string `json:"format,omitempty"`
}
// namespaceMetadata stores metadata for a namespace
type namespaceMetadata struct {
Namespace []string `json:"namespace"`
CreatedAt time.Time `json:"createdAt"`
OwnerAccountID string `json:"ownerAccountId"`
Properties map[string]string `json:"properties,omitempty"`
}
// tableMetadataInternal stores metadata for a table
type tableMetadataInternal struct {
Name string `json:"name"`
Namespace string `json:"namespace"`
Format string `json:"format"`
CreatedAt time.Time `json:"createdAt"`
ModifiedAt time.Time `json:"modifiedAt"`
OwnerAccountID string `json:"ownerAccountId"`
VersionToken string `json:"versionToken"`
MetadataVersion int `json:"metadataVersion"`
MetadataLocation string `json:"metadataLocation,omitempty"`
Metadata *TableMetadata `json:"metadata,omitempty"`
}
// IsTableBucketEntry returns true when the entry is marked as a table bucket.
func IsTableBucketEntry(entry *filer_pb.Entry) bool {
if entry == nil || entry.Extended == nil {
return false
}
_, ok := entry.Extended[ExtendedKeyTableBucket]
return ok
}
// EntryType returns the entry-type marker for a catalog entry. Tables and views
// share the same on-disk layout; the marker distinguishes them. An absent marker
// means table for back-compat.
func EntryType(extended map[string][]byte) string {
if extended == nil {
return EntryTypeTable
}
if v, ok := extended[ExtendedKeyEntryType]; ok && len(v) > 0 {
return string(v)
}
return EntryTypeTable
}
// Utility functions
// validateBucketName validates bucket name and returns an error if invalid.
// Bucket names must contain only lowercase letters, numbers, and hyphens.
// Length must be between 3 and 63 characters.
// Must start and end with a letter or digit.
// Reserved prefixes/suffixes are rejected.
func validateBucketName(name string) error {
if name == "" {
return fmt.Errorf("bucket name is required")
}
if len(name) < 3 || len(name) > 63 {
return fmt.Errorf("bucket name must be between 3 and 63 characters")
}
// Must start and end with a letter or digit
start := name[0]
end := name[len(name)-1]
if !((start >= 'a' && start <= 'z') || (start >= '0' && start <= '9')) {
return fmt.Errorf("bucket name must start with a letter or digit")
}
if !((end >= 'a' && end <= 'z') || (end >= '0' && end <= '9')) {
return fmt.Errorf("bucket name must end with a letter or digit")
}
// Allowed characters: a-z, 0-9, -
for i := 0; i < len(name); i++ {
ch := name[i]
if (ch >= 'a' && ch <= 'z') || (ch >= '0' && ch <= '9') || ch == '-' {
continue
}
return fmt.Errorf("bucket name can only contain lowercase letters, numbers, and hyphens")
}
// Reserved prefixes
reservedPrefixes := []string{"xn--", "sthree-", "amzn-s3-demo-", "aws"}
for _, p := range reservedPrefixes {
if strings.HasPrefix(name, p) {
return fmt.Errorf("bucket name cannot start with reserved prefix: %s", p)
}
}
// Reserved suffixes
reservedSuffixes := []string{"-s3alias", "--ol-s3", "--x-s3", "--table-s3"}
for _, s := range reservedSuffixes {
if strings.HasSuffix(name, s) {
return fmt.Errorf("bucket name cannot end with reserved suffix: %s", s)
}
}
return nil
}
// BuildBucketARN builds a bucket ARN with the provided region and account ID.
// If region is empty, the ARN will omit the region field.
func BuildBucketARN(region, accountID, bucketName string) (string, error) {
if bucketName == "" {
return "", fmt.Errorf("bucket name is required")
}
if err := validateBucketName(bucketName); err != nil {
return "", err
}
if accountID == "" {
accountID = DefaultAccountID
}
return buildARN(region, accountID, fmt.Sprintf("bucket/%s", bucketName)), nil
}
// BuildTableARN builds a table ARN with the provided region and account ID.
func BuildTableARN(region, accountID, bucketName, namespace, tableName string) (string, error) {
if bucketName == "" {
return "", fmt.Errorf("bucket name is required")
}
if err := validateBucketName(bucketName); err != nil {
return "", err
}
if namespace == "" {
return "", fmt.Errorf("namespace is required")
}
normalizedNamespace, err := validateNamespace([]string{namespace})
if err != nil {
return "", err
}
if tableName == "" {
return "", fmt.Errorf("table name is required")
}
normalizedTable, err := validateTableName(tableName)
if err != nil {
return "", err
}
if accountID == "" {
accountID = DefaultAccountID
}
return buildARN(region, accountID, fmt.Sprintf("bucket/%s/table/%s/%s", bucketName, normalizedNamespace, normalizedTable)), nil
}
func buildARN(region, accountID, resourcePath string) string {
return fmt.Sprintf("arn:%s:s3tables:%s:%s:%s", arnPartitionForRegion(region), region, accountID, resourcePath)
}
// arnPartitionForRegion returns the ARN partition a region belongs to, so an
// ARN this handler emits round-trips through a client in that partition.
func arnPartitionForRegion(region string) string {
switch {
case strings.HasPrefix(region, "cn-"):
return "aws-cn"
case strings.HasPrefix(region, "us-gov-"):
return "aws-us-gov"
case strings.HasPrefix(region, "us-iso-"):
return "aws-iso"
case strings.HasPrefix(region, "us-isob-"):
return "aws-iso-b"
case strings.HasPrefix(region, "eu-isoe-"):
return "aws-iso-e"
case strings.HasPrefix(region, "us-isof-"):
return "aws-iso-f"
case strings.HasPrefix(region, "eusc-"):
return "aws-eusc"
default:
return "aws"
}
}
// ValidateTags validates tags for S3 Tables.
func ValidateTags(tags map[string]string) error {
if len(tags) > 10 {
return fmt.Errorf("validate tags: %d tags more than 10", len(tags))
}
for k, v := range tags {
if len(k) > 128 {
return fmt.Errorf("validate tags: tag key longer than 128")
}
if !tagPattern.MatchString(k) {
return fmt.Errorf("validate tags key %s error, incorrect key", k)
}
if len(v) > 256 {
return fmt.Errorf("validate tags: tag value longer than 256")
}
if !tagPattern.MatchString(v) {
return fmt.Errorf("validate tags value %s error, incorrect value", v)
}
}
return nil
}
// isValidBucketName validates bucket name characters (kept for compatibility)
// Deprecated: use validateBucketName instead
func isValidBucketName(name string) bool {
return validateBucketName(name) == nil
}
// generateVersionToken generates a unique, unpredictable version token
func generateVersionToken() string {
b := make([]byte, 16)
if _, err := rand.Read(b); err != nil {
// Fallback to timestamp if crypto/rand fails
return fmt.Sprintf("%x", time.Now().UnixNano())
}
return hex.EncodeToString(b)
}
// splitPath splits a path into directory and name components using stdlib
func splitPath(p string) (dir, name string) {
dir = path.Dir(p)
name = path.Base(p)
return
}
func validateNamespacePart(name string) error {
if len(name) < 1 || len(name) > 255 {
return fmt.Errorf("namespace name must be between 1 and 255 characters")
}
// Prevent path traversal and multi-segment paths
if name == "." || name == ".." {
return fmt.Errorf("namespace name cannot be '.' or '..'")
}
if strings.Contains(name, "/") {
return fmt.Errorf("namespace name cannot contain '/'")
}
// Must start and end with a letter or digit
start := name[0]
end := name[len(name)-1]
if !((start >= 'a' && start <= 'z') || (start >= '0' && start <= '9')) {
return fmt.Errorf("namespace name must start with a letter or digit")
}
if !((end >= 'a' && end <= 'z') || (end >= '0' && end <= '9')) {
return fmt.Errorf("namespace name must end with a letter or digit")
}
// Allowed characters: a-z, 0-9, _, - (hyphen interior; start/end checked above)
for _, ch := range name {
if (ch >= 'a' && ch <= 'z') || (ch >= '0' && ch <= '9') || ch == '_' || ch == '-' {
continue
}
return fmt.Errorf("invalid namespace name: only 'a-z', '0-9', '_', and '-' are allowed")
}
// Reserved prefix
if strings.HasPrefix(name, "aws") {
return fmt.Errorf("namespace name cannot start with reserved prefix 'aws'")
}
return nil
}
func normalizeNamespace(namespace []string) ([]string, error) {
if len(namespace) == 0 {
return nil, fmt.Errorf("namespace is required")
}
parts := namespace
if len(namespace) == 1 {
parts = strings.Split(namespace[0], ".")
}
normalized := make([]string, 0, len(parts))
for _, part := range parts {
if err := validateNamespacePart(part); err != nil {
return nil, err
}
normalized = append(normalized, part)
}
return normalized, nil
}
// validateNamespace validates namespace identifiers and returns an internal namespace key.
// A single dotted namespace value is interpreted as multi-level namespace for compatibility
// with path-style APIs, for example "analytics.daily" => ["analytics", "daily"].
func validateNamespace(namespace []string) (string, error) {
parts, err := normalizeNamespace(namespace)
if err != nil {
return "", err
}
return flattenNamespace(parts), nil
}
// ParseNamespace parses a namespace string into namespace parts.
func ParseNamespace(namespace string) ([]string, error) {
return normalizeNamespace([]string{namespace})
}
// validateTableName validates a table name
func validateTableName(name string) (string, error) {
if len(name) < 1 || len(name) > 255 {
return "", fmt.Errorf("table name must be between 1 and 255 characters")
}
if name == "." || name == ".." || strings.Contains(name, "/") {
return "", fmt.Errorf("invalid table name: cannot be '.', '..' or contain '/'")
}
// First character must be a letter or digit
start := name[0]
if !((start >= 'a' && start <= 'z') || (start >= '0' && start <= '9')) {
return "", fmt.Errorf("table name must start with a letter or digit")
}
// Allowed characters: a-z, 0-9, _, - (start checked above)
for _, ch := range name {
if (ch >= 'a' && ch <= 'z') || (ch >= '0' && ch <= '9') || ch == '_' || ch == '-' {
continue
}
return "", fmt.Errorf("invalid table name: only 'a-z', '0-9', '_', and '-' are allowed")
}
return name, nil
}
// ValidateTableName is a wrapper to validate table name for other packages.
func ValidateTableName(name string) (string, error) {
return validateTableName(name)
}
// flattenNamespace joins namespace elements into a single string (using dots as per AWS S3 Tables)
func flattenNamespace(namespace []string) string {
if len(namespace) == 0 {
return ""
}
return strings.Join(namespace, ".")
}
func expandNamespace(namespace string) []string {
if namespace == "" {
return nil
}
parts, err := ParseNamespace(namespace)
if err != nil {
return []string{namespace}
}
return parts
}