mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-14 18:40:48 +02:00
* 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.
405 lines
16 KiB
Go
405 lines
16 KiB
Go
package iceberg
|
|
|
|
import (
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"net/http"
|
|
"path"
|
|
"strings"
|
|
|
|
"github.com/apache/iceberg-go/table"
|
|
"github.com/google/uuid"
|
|
"github.com/gorilla/mux"
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
|
"github.com/seaweedfs/seaweedfs/weed/s3api/s3tables"
|
|
)
|
|
|
|
// handleUpdateTable commits updates to a table.
|
|
// Implements the Iceberg REST Catalog commit protocol.
|
|
func (s *Server) handleUpdateTable(w http.ResponseWriter, r *http.Request) {
|
|
vars := mux.Vars(r)
|
|
namespace := parseNamespace(vars["namespace"])
|
|
tableName := vars["table"]
|
|
|
|
if len(namespace) == 0 || tableName == "" {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", "Namespace and table name are required")
|
|
return
|
|
}
|
|
|
|
bucketName := getBucketFromPrefix(r)
|
|
bucketARN := buildTableBucketARN(bucketName)
|
|
|
|
// Extract identity from context
|
|
identityName := s3_constants.GetIdentityNameFromContext(r)
|
|
|
|
// Parse commit request and keep statistics updates separate because iceberg-go v0.4.0
|
|
// does not decode set/remove-statistics update actions yet.
|
|
var raw struct {
|
|
Identifier *TableIdentifier `json:"identifier,omitempty"`
|
|
Requirements json.RawMessage `json:"requirements"`
|
|
Updates []json.RawMessage `json:"updates"`
|
|
}
|
|
if err := json.NewDecoder(r.Body).Decode(&raw); err != nil {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", "Invalid request body: "+err.Error())
|
|
return
|
|
}
|
|
|
|
var req CommitTableRequest
|
|
req.Identifier = raw.Identifier
|
|
var statisticsUpdates []statisticsUpdate
|
|
if len(raw.Requirements) > 0 {
|
|
normalized, err := normalizeRequirements(raw.Requirements)
|
|
if err != nil {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", "Invalid requirements: "+err.Error())
|
|
return
|
|
}
|
|
if err := json.Unmarshal(normalized, &req.Requirements); err != nil {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", "Invalid requirements: "+err.Error())
|
|
return
|
|
}
|
|
}
|
|
if len(raw.Updates) > 0 {
|
|
var err error
|
|
req.Updates, statisticsUpdates, err = parseCommitUpdates(raw.Updates)
|
|
if err != nil {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", "Invalid updates: "+err.Error())
|
|
return
|
|
}
|
|
}
|
|
// Manifest repair runs once, as soon as the table location is known; on
|
|
// commit retries the updates already reference the repaired files. Repair
|
|
// is best effort end to end: the originals parsed already, so a repair
|
|
// that fails to re-parse is discarded rather than failing the commit.
|
|
manifestsRepaired := false
|
|
repairManifests := func(location string) {
|
|
if manifestsRepaired {
|
|
return
|
|
}
|
|
manifestsRepaired = true
|
|
repaired, changed := s.repairAddSnapshotManifests(r.Context(), location, raw.Updates)
|
|
if !changed {
|
|
return
|
|
}
|
|
repairedUpdates, repairedStatistics, err := parseCommitUpdates(repaired)
|
|
if err != nil {
|
|
glog.Warningf("Iceberg: repaired updates failed to parse, keeping originals: %v", err)
|
|
return
|
|
}
|
|
raw.Updates = repaired
|
|
req.Updates = repairedUpdates
|
|
statisticsUpdates = repairedStatistics
|
|
}
|
|
|
|
maxCommitAttempts := 3
|
|
generatedLegacyUUID := uuid.New()
|
|
stageCreateEnabled := isStageCreateEnabled()
|
|
for attempt := 1; attempt <= maxCommitAttempts; attempt++ {
|
|
getReq := &s3tables.GetTableRequest{
|
|
TableBucketARN: bucketARN,
|
|
Namespace: namespace,
|
|
Name: tableName,
|
|
}
|
|
var getResp s3tables.GetTableResponse
|
|
|
|
err := s.filerClient.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error {
|
|
mgrClient := s3tables.NewManagerClient(client)
|
|
return s.tablesManager.Execute(r.Context(), mgrClient, "GetTable", getReq, &getResp, identityName)
|
|
})
|
|
if err != nil {
|
|
if isS3TablesNotFound(err) {
|
|
location := fmt.Sprintf("s3://%s/%s", bucketName, path.Join(flattenNamespacePath(namespace), tableName))
|
|
tableUUID := generatedLegacyUUID
|
|
baseMetadataVersion := 0
|
|
baseMetadataLocation := ""
|
|
var baseMetadata table.Metadata
|
|
|
|
var latestMarker *stageCreateMarker
|
|
if stageCreateEnabled {
|
|
var markerErr error
|
|
latestMarker, markerErr = s.loadLatestStageCreateMarker(r.Context(), bucketName, namespace, tableName)
|
|
if markerErr != nil {
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to load stage-create marker: "+markerErr.Error())
|
|
return
|
|
}
|
|
}
|
|
if latestMarker != nil {
|
|
if latestMarker.Location != "" {
|
|
location = strings.TrimSuffix(latestMarker.Location, "/")
|
|
}
|
|
if latestMarker.TableUUID != "" {
|
|
if parsedUUID, parseErr := uuid.Parse(latestMarker.TableUUID); parseErr == nil {
|
|
tableUUID = parsedUUID
|
|
}
|
|
}
|
|
|
|
stagedMetadataLocation := latestMarker.StagedMetadataLocation
|
|
if stagedMetadataLocation == "" {
|
|
stagedMetadataLocation = fmt.Sprintf("%s/metadata/v1.metadata.json", strings.TrimSuffix(location, "/"))
|
|
}
|
|
stagedLocation := tableLocationFromMetadataLocation(stagedMetadataLocation)
|
|
stagedFileName := path.Base(stagedMetadataLocation)
|
|
stagedBucket, stagedPath, parseLocationErr := parseS3Location(stagedLocation)
|
|
if parseLocationErr != nil {
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Invalid staged metadata location: "+parseLocationErr.Error())
|
|
return
|
|
}
|
|
if err := confineMetadataLocation(stagedBucket, stagedPath, bucketName); err != nil {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", err.Error())
|
|
return
|
|
}
|
|
stagedMetadataBytes, loadErr := s.loadMetadataFile(r.Context(), stagedBucket, stagedPath, stagedFileName)
|
|
if loadErr != nil {
|
|
if !errors.Is(loadErr, filer_pb.ErrNotFound) {
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to load staged metadata: "+loadErr.Error())
|
|
return
|
|
}
|
|
} else if len(stagedMetadataBytes) > 0 {
|
|
stagedMetadata, parseErr := table.ParseMetadataBytes(stagedMetadataBytes)
|
|
if parseErr != nil {
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to parse staged metadata: "+parseErr.Error())
|
|
return
|
|
}
|
|
// Staged metadata is only a template for table creation; commit starts from version 1.
|
|
baseMetadata = stagedMetadata
|
|
baseMetadataLocation = ""
|
|
baseMetadataVersion = 0
|
|
if stagedMetadata.TableUUID() != uuid.Nil {
|
|
tableUUID = stagedMetadata.TableUUID()
|
|
}
|
|
}
|
|
}
|
|
|
|
hasAssertCreate := hasAssertCreateRequirement(req.Requirements)
|
|
hasStagedTemplate := baseMetadata != nil
|
|
if !(stageCreateEnabled && (hasAssertCreate || hasStagedTemplate)) {
|
|
writeError(w, http.StatusNotFound, "NoSuchTableException", fmt.Sprintf("Table does not exist: %s", tableName))
|
|
return
|
|
}
|
|
// From here the commit creates the table, writing its metadata
|
|
// file before the create that authorizes it.
|
|
if authErr := s.authorizeCreateTable(r.Context(), bucketARN, namespace, tableName, identityName); authErr != nil {
|
|
writeManagerError(w, authErr)
|
|
return
|
|
}
|
|
|
|
for _, requirement := range req.Requirements {
|
|
validateAgainst := table.Metadata(nil)
|
|
if hasStagedTemplate && requirement.GetType() != requirementAssertCreate {
|
|
validateAgainst = baseMetadata
|
|
}
|
|
if requirementErr := requirement.Validate(validateAgainst); requirementErr != nil {
|
|
writeError(w, http.StatusConflict, "CommitFailedException", "Requirement failed: "+requirementErr.Error())
|
|
return
|
|
}
|
|
}
|
|
|
|
if baseMetadata == nil {
|
|
var buildErr error
|
|
if baseMetadata, buildErr = newTableMetadata(tableUUID, location, nil, nil, nil, nil); buildErr != nil {
|
|
glog.Errorf("Iceberg: CommitTable placeholder metadata for %s: %v", tableName, buildErr)
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to build current metadata")
|
|
return
|
|
}
|
|
}
|
|
|
|
createBucket, createPath, createLocErr := parseS3Location(location)
|
|
if createLocErr != nil {
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Invalid table location: "+createLocErr.Error())
|
|
return
|
|
}
|
|
if err := confineMetadataLocation(createBucket, createPath, bucketName); err != nil {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", err.Error())
|
|
return
|
|
}
|
|
|
|
repairManifests(location)
|
|
|
|
result, reqErr := s.finalizeCreateOnCommit(r.Context(), createOnCommitInput{
|
|
bucketARN: bucketARN,
|
|
markerBucket: bucketName,
|
|
namespace: namespace,
|
|
tableName: tableName,
|
|
identityName: identityName,
|
|
location: location,
|
|
tableUUID: tableUUID,
|
|
baseMetadata: baseMetadata,
|
|
baseMetadataLoc: baseMetadataLocation,
|
|
baseMetadataVer: baseMetadataVersion,
|
|
updates: req.Updates,
|
|
statisticsUpdates: statisticsUpdates,
|
|
})
|
|
if reqErr != nil {
|
|
writeError(w, reqErr.status, reqErr.errType, reqErr.message)
|
|
return
|
|
}
|
|
writeJSON(w, http.StatusOK, result)
|
|
return
|
|
}
|
|
glog.V(1).Infof("Iceberg: CommitTable GetTable error: %v", err)
|
|
writeManagerError(w, err)
|
|
return
|
|
}
|
|
|
|
location := tableLocationFromMetadataLocation(getResp.MetadataLocation)
|
|
if location == "" {
|
|
location = fmt.Sprintf("s3://%s/%s", bucketName, path.Join(flattenNamespacePath(namespace), tableName))
|
|
}
|
|
locBucket, locPath, locErr := parseS3Location(location)
|
|
if locErr != nil {
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Invalid table location: "+locErr.Error())
|
|
return
|
|
}
|
|
if err := confineMetadataLocation(locBucket, locPath, bucketName); err != nil {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", err.Error())
|
|
return
|
|
}
|
|
tableUUID := uuid.Nil
|
|
if getResp.Metadata != nil && getResp.Metadata.Iceberg != nil && getResp.Metadata.Iceberg.TableUUID != "" {
|
|
if parsed, parseErr := uuid.Parse(getResp.Metadata.Iceberg.TableUUID); parseErr == nil {
|
|
tableUUID = parsed
|
|
}
|
|
}
|
|
if tableUUID == uuid.Nil {
|
|
tableUUID = generatedLegacyUUID
|
|
}
|
|
|
|
var currentMetadata table.Metadata
|
|
if getResp.Metadata != nil && len(getResp.Metadata.FullMetadata) > 0 {
|
|
currentMetadata, err = table.ParseMetadataBytes(getResp.Metadata.FullMetadata)
|
|
if err != nil {
|
|
glog.Errorf("Iceberg: Failed to parse current metadata for %s: %v", tableName, err)
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to parse current metadata")
|
|
return
|
|
}
|
|
} else {
|
|
currentMetadata, err = newTableMetadata(tableUUID, location, nil, nil, nil, nil)
|
|
if err != nil {
|
|
glog.Errorf("Iceberg: CommitTable placeholder metadata for %s: %v", tableName, err)
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to build current metadata")
|
|
return
|
|
}
|
|
}
|
|
|
|
for _, requirement := range req.Requirements {
|
|
if err := requirement.Validate(currentMetadata); err != nil {
|
|
writeError(w, http.StatusConflict, "CommitFailedException", "Requirement failed: "+err.Error())
|
|
return
|
|
}
|
|
}
|
|
|
|
repairManifests(location)
|
|
|
|
builder, err := table.MetadataBuilderFromBase(currentMetadata, getResp.MetadataLocation)
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to create metadata builder: "+err.Error())
|
|
return
|
|
}
|
|
for _, update := range req.Updates {
|
|
if err := update.Apply(builder); err != nil {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", "Failed to apply update: "+err.Error())
|
|
return
|
|
}
|
|
}
|
|
|
|
newMetadata, err := builder.Build()
|
|
if err != nil {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", "Failed to build new metadata: "+err.Error())
|
|
return
|
|
}
|
|
|
|
metadataVersion := getResp.MetadataVersion + 1
|
|
metadataFileName := fmt.Sprintf("v%d.metadata.json", metadataVersion)
|
|
newMetadataLocation := fmt.Sprintf("%s/metadata/%s", strings.TrimSuffix(location, "/"), metadataFileName)
|
|
|
|
metadataBytes, err := json.Marshal(newMetadata)
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to serialize metadata: "+err.Error())
|
|
return
|
|
}
|
|
// iceberg-go does not currently support set/remove-statistics updates in MetadataBuilder.
|
|
// Patch the encoded metadata JSON and parse it back to keep the response object consistent.
|
|
metadataBytes, err = applyStatisticsUpdates(metadataBytes, statisticsUpdates)
|
|
if err != nil {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", "Failed to apply statistics updates: "+err.Error())
|
|
return
|
|
}
|
|
metadataBytes = refreshDefaultNameMapping(metadataBytes, newMetadata)
|
|
// Same spec-compliance fixup we apply on create-table; ensures
|
|
// v{N}.metadata.json files written during commit are also readable by
|
|
// strict Iceberg clients reading directly from S3, and that the
|
|
// FullMetadata persisted in S3Tables stays consistent.
|
|
metadataBytes = ensureMetadataSpecCompliance(metadataBytes)
|
|
newMetadata, err = table.ParseMetadataBytes(metadataBytes)
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to parse committed metadata: "+err.Error())
|
|
return
|
|
}
|
|
|
|
metadataBucket, metadataPath, err := parseS3Location(location)
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Invalid table location: "+err.Error())
|
|
return
|
|
}
|
|
if err := confineMetadataLocation(metadataBucket, metadataPath, bucketName); err != nil {
|
|
writeError(w, http.StatusBadRequest, "BadRequestException", err.Error())
|
|
return
|
|
}
|
|
metadataFileName, newMetadataLocation, err = s.stageCommitMetadata(r.Context(), metadataBucket, metadataPath, location, metadataFileName, metadataBytes)
|
|
if err != nil {
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to save metadata file: "+err.Error())
|
|
return
|
|
}
|
|
|
|
updateReq := &s3tables.UpdateTableRequest{
|
|
TableBucketARN: bucketARN,
|
|
Namespace: namespace,
|
|
Name: tableName,
|
|
VersionToken: getResp.VersionToken,
|
|
Metadata: &s3tables.TableMetadata{
|
|
Iceberg: &s3tables.IcebergMetadata{
|
|
TableUUID: tableUUID.String(),
|
|
},
|
|
FullMetadata: metadataBytes,
|
|
},
|
|
MetadataVersion: metadataVersion,
|
|
MetadataLocation: newMetadataLocation,
|
|
}
|
|
|
|
err = s.filerClient.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error {
|
|
mgrClient := s3tables.NewManagerClient(client)
|
|
return s.tablesManager.Execute(r.Context(), mgrClient, "UpdateTable", updateReq, nil, identityName)
|
|
})
|
|
if err == nil {
|
|
result := CommitTableResponse{
|
|
MetadataLocation: newMetadataLocation,
|
|
Metadata: newMetadata,
|
|
}
|
|
writeJSON(w, http.StatusOK, result)
|
|
return
|
|
}
|
|
|
|
if isS3TablesConflict(err) {
|
|
if cleanupErr := s.deleteMetadataFile(r.Context(), metadataBucket, metadataPath, metadataFileName); cleanupErr != nil {
|
|
glog.V(1).Infof("Iceberg: failed to cleanup metadata file %s on conflict: %v", newMetadataLocation, cleanupErr)
|
|
}
|
|
if attempt < maxCommitAttempts {
|
|
glog.V(1).Infof("Iceberg: CommitTable conflict for %s (attempt %d/%d), retrying", tableName, attempt, maxCommitAttempts)
|
|
sleepBeforeCommitRetry(attempt)
|
|
continue
|
|
}
|
|
writeError(w, http.StatusConflict, "CommitFailedException", "Version token mismatch")
|
|
return
|
|
}
|
|
|
|
if cleanupErr := s.deleteMetadataFile(r.Context(), metadataBucket, metadataPath, metadataFileName); cleanupErr != nil {
|
|
glog.V(1).Infof("Iceberg: failed to cleanup metadata file %s after update failure: %v", newMetadataLocation, cleanupErr)
|
|
}
|
|
glog.Errorf("Iceberg: CommitTable UpdateTable error: %v", err)
|
|
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to commit table update: "+err.Error())
|
|
return
|
|
}
|
|
}
|