Files
seaweedfs/weed/s3api/iceberg/handlers_commit.go
T
Chris LuandDevin 9a365549c3 catalog: refused writes answer 403 to callers that can read the entry (#11650)
* s3tables: answer a denied write with 403 when the caller can read the entry

Update/Delete/Rename answered every authorization refusal as not-found so
the denial leaked no existence signal. For a caller allowed to GetTable
(or GetView) the same entry, the veil hides nothing it could not load —
yet a refused write was still answered 404, so clients saw a table they
just loaded reported as missing.

When the write check fails, re-check read permission on the same entry:
readable entries get 403 AccessDenied; invisible ones keep the
not-found answer.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* iceberg: map UpdateTable errors through writeManagerError on commit paths

CommitTable and CommitTransaction answered any non-conflict UpdateTable
failure, including AccessDenied and NoSuchTable, with a bare 500. A 500
reads as outcome-unknown to clients (PyIceberg raises
CommitStateUnknownException) where a refused commit is a plain
ForbiddenException, matching what the drop and rename handlers already
emit through the same mapper.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* s3tables: keep the not-found veil over views and match the real read check on renames

UpdateTable and DeleteTable read shared metadata without checking the
entry kind, so a denied write on a view reported 403 to a GetTable-only
caller where a missing name reports 404. The rename visibility check
also fed resource tags into GetView evaluation that the real GetView
path never supplies.

* s3tables: gofmt handler_table.go

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-10-09 09:59:32 +08:00

405 lines
15 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)
writeManagerError(w, err)
return
}
}