iceberg: support multi-table transaction commit (#10066)
* iceberg: support multi-table transaction commit Add handleCommitTransaction for POST /v1/transactions/commit. Validation is atomic across all table-changes (resolve, load, evaluate every requirement before any write); metadata writes and pointer flips are best-effort with rollback, so this is not crash-atomic. * iceberg: route transactions/commit with and without prefix * iceberg: test transaction commit request decoding * iceberg: restore full prior table state on transaction rollback * iceberg: test transaction rollback restores full prior table state * iceberg: only clean up metadata for rolled-back tables
This commit is contained in:
parent
628ce57625
commit
1ca628d3e9
348
weed/s3api/iceberg/handlers_transaction.go
Normal file
348
weed/s3api/iceberg/handlers_transaction.go
Normal file
@ -0,0 +1,348 @@
|
||||
package iceberg
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"path"
|
||||
"strings"
|
||||
|
||||
"github.com/apache/iceberg-go/table"
|
||||
"github.com/google/uuid"
|
||||
"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"
|
||||
)
|
||||
|
||||
// CommitTransactionRequest is sent to POST /v1/transactions/commit.
|
||||
type CommitTransactionRequest struct {
|
||||
TableChanges []tableChangeRequest `json:"table-changes"`
|
||||
}
|
||||
|
||||
type tableChangeRequest struct {
|
||||
Identifier *TableIdentifier `json:"identifier"`
|
||||
Requirements json.RawMessage `json:"requirements"`
|
||||
Updates []json.RawMessage `json:"updates"`
|
||||
}
|
||||
|
||||
// preparedTableCommit holds the per-table work resolved during validation so
|
||||
// pointer flips (and their rollback) can run after every table is validated.
|
||||
type preparedTableCommit struct {
|
||||
namespace []string
|
||||
tableName string
|
||||
tableUUID uuid.UUID
|
||||
versionToken string
|
||||
metadataBucket string
|
||||
metadataPath string
|
||||
metadataFileName string
|
||||
metadataBytes []byte
|
||||
metadataVersion int
|
||||
newMetadataLoc string
|
||||
|
||||
// prior table state captured before the flip, so a rollback can revert
|
||||
// every field handleUpdateTable would have mutated (not just the location).
|
||||
prevMetadataLoc string
|
||||
prevMetadataVersion int
|
||||
prevMetadata *s3tables.TableMetadata
|
||||
}
|
||||
|
||||
// handleCommitTransaction commits changes to multiple tables in one request.
|
||||
// Validation is atomic (all requirements evaluated before any write); pointer
|
||||
// flips are best-effort with rollback, so this is not crash-atomic.
|
||||
func (s *Server) handleCommitTransaction(w http.ResponseWriter, r *http.Request) {
|
||||
bucketName := getBucketFromPrefix(r)
|
||||
bucketARN := buildTableBucketARN(bucketName)
|
||||
identityName := s3_constants.GetIdentityNameFromContext(r)
|
||||
|
||||
var req CommitTransactionRequest
|
||||
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
|
||||
writeError(w, http.StatusBadRequest, "BadRequestException", "Invalid request body: "+err.Error())
|
||||
return
|
||||
}
|
||||
if len(req.TableChanges) == 0 {
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
return
|
||||
}
|
||||
|
||||
// Phase 1: resolve, load, and validate every table; build new metadata bytes.
|
||||
prepared := make([]preparedTableCommit, 0, len(req.TableChanges))
|
||||
for _, change := range req.TableChanges {
|
||||
if change.Identifier == nil || len(change.Identifier.Namespace) == 0 || change.Identifier.Name == "" {
|
||||
writeError(w, http.StatusBadRequest, "BadRequestException", "Each table change requires identifier namespace and name")
|
||||
return
|
||||
}
|
||||
pc, reqErr := s.prepareTableCommit(r.Context(), bucketName, bucketARN, identityName, change)
|
||||
if reqErr != nil {
|
||||
writeError(w, reqErr.status, reqErr.errType, reqErr.message)
|
||||
return
|
||||
}
|
||||
prepared = append(prepared, *pc)
|
||||
}
|
||||
|
||||
// Phase 2: write each new metadata.json object.
|
||||
for i := range prepared {
|
||||
pc := &prepared[i]
|
||||
if err := s.saveMetadataFile(r.Context(), pc.metadataBucket, pc.metadataPath, pc.metadataFileName, pc.metadataBytes); err != nil {
|
||||
// No pointer flipped yet, so every written file is safe to delete.
|
||||
s.cleanupPreparedMetadata(r.Context(), prepared[:i+1], nil)
|
||||
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to save metadata file: "+err.Error())
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// Phase 3: flip each table's pointer xattr; on failure roll back prior flips.
|
||||
for i := range prepared {
|
||||
pc := &prepared[i]
|
||||
if err := s.flipTablePointer(r.Context(), bucketARN, identityName, pc); err != nil {
|
||||
rolledBack := s.rollbackTablePointers(r.Context(), bucketARN, identityName, prepared[:i])
|
||||
s.cleanupPreparedMetadata(r.Context(), prepared, cleanupSafeMetadata(prepared, i, rolledBack))
|
||||
if isS3TablesConflict(err) {
|
||||
writeError(w, http.StatusConflict, "CommitFailedException", "Version token mismatch")
|
||||
return
|
||||
}
|
||||
glog.Errorf("Iceberg: CommitTransaction UpdateTable error: %v", err)
|
||||
writeError(w, http.StatusInternalServerError, "InternalServerError", "Failed to commit table update: "+err.Error())
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
w.WriteHeader(http.StatusNoContent)
|
||||
}
|
||||
|
||||
func (s *Server) prepareTableCommit(ctx context.Context, bucketName, bucketARN, identityName string, change tableChangeRequest) (*preparedTableCommit, *icebergRequestError) {
|
||||
namespace := []string(change.Identifier.Namespace)
|
||||
tableName := change.Identifier.Name
|
||||
|
||||
requirements, updates, statisticsUpdates, reqErr := parseTableChange(change)
|
||||
if reqErr != nil {
|
||||
return nil, reqErr
|
||||
}
|
||||
|
||||
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(ctx, mgrClient, "GetTable", getReq, &getResp, identityName)
|
||||
})
|
||||
if err != nil {
|
||||
if isS3TablesNotFound(err) {
|
||||
return nil, &icebergRequestError{http.StatusNotFound, "NoSuchTableException", fmt.Sprintf("Table does not exist: %s", tableName)}
|
||||
}
|
||||
glog.V(1).Infof("Iceberg: CommitTransaction GetTable error: %v", err)
|
||||
return nil, &icebergRequestError{http.StatusInternalServerError, "InternalServerError", err.Error()}
|
||||
}
|
||||
|
||||
location := tableLocationFromMetadataLocation(getResp.MetadataLocation)
|
||||
if location == "" {
|
||||
location = fmt.Sprintf("s3://%s/%s", bucketName, path.Join(flattenNamespacePath(namespace), tableName))
|
||||
}
|
||||
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 = uuid.New()
|
||||
}
|
||||
|
||||
var currentMetadata table.Metadata
|
||||
if getResp.Metadata != nil && len(getResp.Metadata.FullMetadata) > 0 {
|
||||
currentMetadata, err = table.ParseMetadataBytes(getResp.Metadata.FullMetadata)
|
||||
if err != nil {
|
||||
return nil, &icebergRequestError{http.StatusInternalServerError, "InternalServerError", "Failed to parse current metadata"}
|
||||
}
|
||||
} else {
|
||||
currentMetadata = newTableMetadata(tableUUID, location, nil, nil, nil, nil)
|
||||
}
|
||||
if currentMetadata == nil {
|
||||
return nil, &icebergRequestError{http.StatusInternalServerError, "InternalServerError", "Failed to build current metadata"}
|
||||
}
|
||||
|
||||
for _, requirement := range requirements {
|
||||
if err := requirement.Validate(currentMetadata); err != nil {
|
||||
return nil, &icebergRequestError{http.StatusConflict, "CommitFailedException", "Requirement failed: " + err.Error()}
|
||||
}
|
||||
}
|
||||
|
||||
builder, err := table.MetadataBuilderFromBase(currentMetadata, getResp.MetadataLocation)
|
||||
if err != nil {
|
||||
return nil, &icebergRequestError{http.StatusInternalServerError, "InternalServerError", "Failed to create metadata builder: " + err.Error()}
|
||||
}
|
||||
for _, update := range updates {
|
||||
if err := update.Apply(builder); err != nil {
|
||||
return nil, &icebergRequestError{http.StatusBadRequest, "BadRequestException", "Failed to apply update: " + err.Error()}
|
||||
}
|
||||
}
|
||||
newMetadata, err := builder.Build()
|
||||
if err != nil {
|
||||
return nil, &icebergRequestError{http.StatusBadRequest, "BadRequestException", "Failed to build new metadata: " + err.Error()}
|
||||
}
|
||||
|
||||
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 {
|
||||
return nil, &icebergRequestError{http.StatusInternalServerError, "InternalServerError", "Failed to serialize metadata: " + err.Error()}
|
||||
}
|
||||
metadataBytes, err = applyStatisticsUpdates(metadataBytes, statisticsUpdates)
|
||||
if err != nil {
|
||||
return nil, &icebergRequestError{http.StatusBadRequest, "BadRequestException", "Failed to apply statistics updates: " + err.Error()}
|
||||
}
|
||||
metadataBytes = ensureMetadataSpecCompliance(metadataBytes)
|
||||
|
||||
metadataBucket, metadataPath, err := parseS3Location(location)
|
||||
if err != nil {
|
||||
return nil, &icebergRequestError{http.StatusInternalServerError, "InternalServerError", "Invalid table location: " + err.Error()}
|
||||
}
|
||||
|
||||
return &preparedTableCommit{
|
||||
namespace: namespace,
|
||||
tableName: tableName,
|
||||
tableUUID: tableUUID,
|
||||
versionToken: getResp.VersionToken,
|
||||
metadataBucket: metadataBucket,
|
||||
metadataPath: metadataPath,
|
||||
metadataFileName: metadataFileName,
|
||||
metadataBytes: metadataBytes,
|
||||
metadataVersion: metadataVersion,
|
||||
newMetadataLoc: newMetadataLocation,
|
||||
prevMetadataLoc: getResp.MetadataLocation,
|
||||
prevMetadataVersion: getResp.MetadataVersion,
|
||||
prevMetadata: cloneTableMetadata(getResp.Metadata),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// cloneTableMetadata deep-copies the prior table metadata so a later rollback
|
||||
// restores the exact pre-transaction bytes without aliasing the get response.
|
||||
func cloneTableMetadata(m *s3tables.TableMetadata) *s3tables.TableMetadata {
|
||||
if m == nil {
|
||||
return nil
|
||||
}
|
||||
clone := &s3tables.TableMetadata{}
|
||||
if m.Iceberg != nil {
|
||||
iceberg := *m.Iceberg
|
||||
clone.Iceberg = &iceberg
|
||||
}
|
||||
if len(m.FullMetadata) > 0 {
|
||||
clone.FullMetadata = append(json.RawMessage(nil), m.FullMetadata...)
|
||||
}
|
||||
return clone
|
||||
}
|
||||
|
||||
func parseTableChange(change tableChangeRequest) (table.Requirements, table.Updates, []statisticsUpdate, *icebergRequestError) {
|
||||
var requirements table.Requirements
|
||||
if len(change.Requirements) > 0 {
|
||||
if err := json.Unmarshal(change.Requirements, &requirements); err != nil {
|
||||
return nil, nil, nil, &icebergRequestError{http.StatusBadRequest, "BadRequestException", "Invalid requirements: " + err.Error()}
|
||||
}
|
||||
}
|
||||
var updates table.Updates
|
||||
var statisticsUpdates []statisticsUpdate
|
||||
if len(change.Updates) > 0 {
|
||||
var err error
|
||||
updates, statisticsUpdates, err = parseCommitUpdates(change.Updates)
|
||||
if err != nil {
|
||||
return nil, nil, nil, &icebergRequestError{http.StatusBadRequest, "BadRequestException", "Invalid updates: " + err.Error()}
|
||||
}
|
||||
}
|
||||
return requirements, updates, statisticsUpdates, nil
|
||||
}
|
||||
|
||||
func (s *Server) flipTablePointer(ctx context.Context, bucketARN, identityName string, pc *preparedTableCommit) error {
|
||||
updateReq := &s3tables.UpdateTableRequest{
|
||||
TableBucketARN: bucketARN,
|
||||
Namespace: pc.namespace,
|
||||
Name: pc.tableName,
|
||||
VersionToken: pc.versionToken,
|
||||
Metadata: &s3tables.TableMetadata{
|
||||
Iceberg: &s3tables.IcebergMetadata{TableUUID: pc.tableUUID.String()},
|
||||
FullMetadata: pc.metadataBytes,
|
||||
},
|
||||
MetadataVersion: pc.metadataVersion,
|
||||
MetadataLocation: pc.newMetadataLoc,
|
||||
}
|
||||
return s.filerClient.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error {
|
||||
mgrClient := s3tables.NewManagerClient(client)
|
||||
return s.tablesManager.Execute(ctx, mgrClient, "UpdateTable", updateReq, nil, identityName)
|
||||
})
|
||||
}
|
||||
|
||||
// rollbackTablePointers restores already-flipped tables to their full prior
|
||||
// state. handleUpdateTable applies partial field updates, so the restore must
|
||||
// re-send every field the flip changed (location, version, and full metadata),
|
||||
// not just the location. Best-effort: a failed rollback is logged, not surfaced.
|
||||
// The returned slice flags, per flipped table, whether the restore succeeded so
|
||||
// the caller can avoid deleting metadata a still-flipped pointer references.
|
||||
func (s *Server) rollbackTablePointers(ctx context.Context, bucketARN, identityName string, flipped []preparedTableCommit) []bool {
|
||||
restored := make([]bool, len(flipped))
|
||||
for i := range flipped {
|
||||
pc := &flipped[i]
|
||||
updateReq := buildTableRestoreRequest(bucketARN, pc)
|
||||
err := s.filerClient.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error {
|
||||
mgrClient := s3tables.NewManagerClient(client)
|
||||
return s.tablesManager.Execute(ctx, mgrClient, "UpdateTable", updateReq, nil, identityName)
|
||||
})
|
||||
if err != nil {
|
||||
glog.Errorf("Iceberg: CommitTransaction rollback of %s failed: %v", pc.tableName, err)
|
||||
continue
|
||||
}
|
||||
restored[i] = true
|
||||
}
|
||||
return restored
|
||||
}
|
||||
|
||||
// buildTableRestoreRequest reconstructs the UpdateTableRequest that reverts a
|
||||
// flipped table to its captured prior state. Every field handleUpdateTable
|
||||
// would have mutated on the flip is re-supplied here so the partial update
|
||||
// fully restores it (ModifiedAt and VersionToken are always regenerated by the
|
||||
// handler and cannot be pinned back).
|
||||
func buildTableRestoreRequest(bucketARN string, pc *preparedTableCommit) *s3tables.UpdateTableRequest {
|
||||
return &s3tables.UpdateTableRequest{
|
||||
TableBucketARN: bucketARN,
|
||||
Namespace: pc.namespace,
|
||||
Name: pc.tableName,
|
||||
Metadata: cloneTableMetadata(pc.prevMetadata),
|
||||
MetadataVersion: pc.prevMetadataVersion,
|
||||
MetadataLocation: pc.prevMetadataLoc,
|
||||
}
|
||||
}
|
||||
|
||||
// cleanupSafeMetadata reports, per prepared table, whether deleting its newly
|
||||
// written metadata file is safe after a flip failed at failedIndex. A file is
|
||||
// safe to delete only when its table's pointer no longer references it: tables
|
||||
// at or past failedIndex were never flipped, and earlier tables are safe only
|
||||
// if their rollback succeeded. A table whose rollback failed still points at the
|
||||
// new metadata, so deleting it would strand the pointer on a missing file.
|
||||
func cleanupSafeMetadata(prepared []preparedTableCommit, failedIndex int, rolledBack []bool) []bool {
|
||||
safe := make([]bool, len(prepared))
|
||||
for i := range prepared {
|
||||
switch {
|
||||
case i >= failedIndex:
|
||||
safe[i] = true
|
||||
case i < len(rolledBack):
|
||||
safe[i] = rolledBack[i]
|
||||
}
|
||||
}
|
||||
return safe
|
||||
}
|
||||
|
||||
// cleanupPreparedMetadata deletes the newly written metadata file for each
|
||||
// prepared table flagged safe; unflagged tables are left in place because a
|
||||
// still-flipped pointer references them.
|
||||
func (s *Server) cleanupPreparedMetadata(ctx context.Context, prepared []preparedTableCommit, safe []bool) {
|
||||
for i := range prepared {
|
||||
pc := &prepared[i]
|
||||
if i < len(safe) && !safe[i] {
|
||||
glog.Errorf("Iceberg: CommitTransaction keeping metadata %s; %s rollback failed and still references it", pc.newMetadataLoc, pc.tableName)
|
||||
continue
|
||||
}
|
||||
if cleanupErr := s.deleteMetadataFile(ctx, pc.metadataBucket, pc.metadataPath, pc.metadataFileName); cleanupErr != nil {
|
||||
glog.V(1).Infof("Iceberg: failed to cleanup metadata file %s: %v", pc.newMetadataLoc, cleanupErr)
|
||||
}
|
||||
}
|
||||
}
|
||||
254
weed/s3api/iceberg/iceberg_transaction_test.go
Normal file
254
weed/s3api/iceberg/iceberg_transaction_test.go
Normal file
@ -0,0 +1,254 @@
|
||||
package iceberg
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3tables"
|
||||
)
|
||||
|
||||
func TestDecodeCommitTransactionRequest(t *testing.T) {
|
||||
body := `{"table-changes":[
|
||||
{"identifier":{"namespace":["ns"],"name":"t1"},"requirements":[{"type":"assert-table-uuid","uuid":"00000000-0000-0000-0000-000000000001"}],"updates":[{"action":"set-properties","updates":{"k":"v"}}]},
|
||||
{"identifier":{"namespace":["ns","sub"],"name":"t2"},"requirements":[],"updates":[]}
|
||||
]}`
|
||||
|
||||
var req CommitTransactionRequest
|
||||
if err := json.Unmarshal([]byte(body), &req); err != nil {
|
||||
t.Fatalf("Unmarshal() error = %v", err)
|
||||
}
|
||||
if len(req.TableChanges) != 2 {
|
||||
t.Fatalf("table-changes = %d, want 2", len(req.TableChanges))
|
||||
}
|
||||
if req.TableChanges[0].Identifier.Name != "t1" {
|
||||
t.Fatalf("change[0] name = %q, want t1", req.TableChanges[0].Identifier.Name)
|
||||
}
|
||||
if got := []string(req.TableChanges[1].Identifier.Namespace); len(got) != 2 || got[1] != "sub" {
|
||||
t.Fatalf("change[1] namespace = %v, want [ns sub]", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseTableChangeSeparatesStatistics(t *testing.T) {
|
||||
change := tableChangeRequest{
|
||||
Requirements: json.RawMessage(`[{"type":"assert-table-uuid","uuid":"00000000-0000-0000-0000-000000000001"}]`),
|
||||
Updates: []json.RawMessage{
|
||||
json.RawMessage(`{"action":"set-properties","updates":{"k":"v"}}`),
|
||||
json.RawMessage(`{"action":"set-statistics","snapshot-id":10,"statistics-path":"s3://bucket/table/metadata/stats.puffin","file-size-in-bytes":100,"file-footer-size-in-bytes":20,"blob-metadata":[]}`),
|
||||
},
|
||||
}
|
||||
|
||||
requirements, updates, statistics, reqErr := parseTableChange(change)
|
||||
if reqErr != nil {
|
||||
t.Fatalf("parseTableChange() error = %+v", reqErr)
|
||||
}
|
||||
if len(requirements) != 1 {
|
||||
t.Fatalf("requirements = %d, want 1", len(requirements))
|
||||
}
|
||||
if len(updates) != 1 {
|
||||
t.Fatalf("updates = %d, want 1", len(updates))
|
||||
}
|
||||
if len(statistics) != 1 || statistics[0].set == nil || statistics[0].set.SnapshotID != 10 {
|
||||
t.Fatalf("unexpected statistics updates: %#v", statistics)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseTableChangeRejectsInvalidRequirements(t *testing.T) {
|
||||
change := tableChangeRequest{Requirements: json.RawMessage(`{not json}`)}
|
||||
_, _, _, reqErr := parseTableChange(change)
|
||||
if reqErr == nil {
|
||||
t.Fatalf("parseTableChange() expected error for invalid requirements")
|
||||
}
|
||||
if reqErr.status != 400 {
|
||||
t.Fatalf("status = %d, want 400", reqErr.status)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCloneTableMetadataIsIndependent(t *testing.T) {
|
||||
orig := &s3tables.TableMetadata{
|
||||
Iceberg: &s3tables.IcebergMetadata{TableUUID: "uuid-1"},
|
||||
FullMetadata: json.RawMessage(`{"v":1}`),
|
||||
}
|
||||
clone := cloneTableMetadata(orig)
|
||||
if clone == orig || clone.Iceberg == orig.Iceberg {
|
||||
t.Fatalf("cloneTableMetadata returned aliased pointers")
|
||||
}
|
||||
|
||||
orig.Iceberg.TableUUID = "mutated"
|
||||
orig.FullMetadata[0] = 'X'
|
||||
if clone.Iceberg.TableUUID != "uuid-1" {
|
||||
t.Fatalf("clone Iceberg.TableUUID = %q, want uuid-1", clone.Iceberg.TableUUID)
|
||||
}
|
||||
if string(clone.FullMetadata) != `{"v":1}` {
|
||||
t.Fatalf("clone FullMetadata = %q, mutated with source", clone.FullMetadata)
|
||||
}
|
||||
|
||||
if cloneTableMetadata(nil) != nil {
|
||||
t.Fatalf("cloneTableMetadata(nil) = non-nil")
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildTableRestoreRequestCarriesPriorState(t *testing.T) {
|
||||
pc := &preparedTableCommit{
|
||||
namespace: []string{"ns", "sub"},
|
||||
tableName: "t1",
|
||||
prevMetadataLoc: "s3://bucket/ns/t1/metadata/v3.metadata.json",
|
||||
prevMetadataVersion: 3,
|
||||
prevMetadata: &s3tables.TableMetadata{
|
||||
Iceberg: &s3tables.IcebergMetadata{TableUUID: uuid.NewString()},
|
||||
FullMetadata: json.RawMessage(`{"format-version":2,"last-updated-ms":3}`),
|
||||
},
|
||||
}
|
||||
|
||||
req := buildTableRestoreRequest("arn:bucket", pc)
|
||||
|
||||
if req.TableBucketARN != "arn:bucket" || req.Name != "t1" {
|
||||
t.Fatalf("restore target = %s/%s, want arn:bucket/t1", req.TableBucketARN, req.Name)
|
||||
}
|
||||
if req.MetadataLocation != pc.prevMetadataLoc {
|
||||
t.Fatalf("restore MetadataLocation = %q, want %q", req.MetadataLocation, pc.prevMetadataLoc)
|
||||
}
|
||||
if req.MetadataVersion != pc.prevMetadataVersion {
|
||||
t.Fatalf("restore MetadataVersion = %d, want %d", req.MetadataVersion, pc.prevMetadataVersion)
|
||||
}
|
||||
if req.Metadata == nil || !bytes.Equal(req.Metadata.FullMetadata, pc.prevMetadata.FullMetadata) {
|
||||
t.Fatalf("restore FullMetadata = %q, want %q", req.Metadata.FullMetadata, pc.prevMetadata.FullMetadata)
|
||||
}
|
||||
if req.Metadata.Iceberg == nil || req.Metadata.Iceberg.TableUUID != pc.prevMetadata.Iceberg.TableUUID {
|
||||
t.Fatalf("restore TableUUID mismatch")
|
||||
}
|
||||
}
|
||||
|
||||
// TestRollbackRevertsEveryFlippedField mirrors handleUpdateTable's partial-field
|
||||
// update semantics to prove that after a flip + rollback the table is back to its
|
||||
// pre-transaction state. With the old location-only rollback the table kept the
|
||||
// new full metadata and version; the restore request must revert all of them.
|
||||
func TestRollbackRevertsEveryFlippedField(t *testing.T) {
|
||||
// prior on-disk state of an already-committed table.
|
||||
prior := tableEntry{
|
||||
metadataLocation: "s3://bucket/ns/t1/metadata/v3.metadata.json",
|
||||
metadataVersion: 3,
|
||||
metadata: &s3tables.TableMetadata{
|
||||
Iceberg: &s3tables.IcebergMetadata{TableUUID: "uuid-1"},
|
||||
FullMetadata: json.RawMessage(`{"format-version":2,"v":3}`),
|
||||
},
|
||||
}
|
||||
|
||||
pc := &preparedTableCommit{
|
||||
namespace: []string{"ns"},
|
||||
tableName: "t1",
|
||||
tableUUID: uuid.MustParse("00000000-0000-0000-0000-000000000001"),
|
||||
metadataBytes: json.RawMessage(`{"format-version":2,"v":4}`),
|
||||
metadataVersion: 4,
|
||||
newMetadataLoc: "s3://bucket/ns/t1/metadata/v4.metadata.json",
|
||||
prevMetadataLoc: prior.metadataLocation,
|
||||
prevMetadataVersion: prior.metadataVersion,
|
||||
prevMetadata: cloneTableMetadata(prior.metadata),
|
||||
}
|
||||
|
||||
// flip: apply the new commit, then roll it back via the restore request.
|
||||
state := prior.clone()
|
||||
state.applyUpdate(&s3tables.UpdateTableRequest{
|
||||
Metadata: &s3tables.TableMetadata{Iceberg: &s3tables.IcebergMetadata{TableUUID: pc.tableUUID.String()}, FullMetadata: pc.metadataBytes},
|
||||
MetadataVersion: pc.metadataVersion,
|
||||
MetadataLocation: pc.newMetadataLoc,
|
||||
})
|
||||
state.applyUpdate(buildTableRestoreRequest("arn:bucket", pc))
|
||||
|
||||
if state.metadataLocation != prior.metadataLocation {
|
||||
t.Fatalf("after rollback location = %q, want %q", state.metadataLocation, prior.metadataLocation)
|
||||
}
|
||||
if state.metadataVersion != prior.metadataVersion {
|
||||
t.Fatalf("after rollback version = %d, want %d", state.metadataVersion, prior.metadataVersion)
|
||||
}
|
||||
if !bytes.Equal(state.metadata.FullMetadata, prior.metadata.FullMetadata) {
|
||||
t.Fatalf("after rollback FullMetadata = %q, want %q", state.metadata.FullMetadata, prior.metadata.FullMetadata)
|
||||
}
|
||||
if state.metadata.Iceberg.TableUUID != prior.metadata.Iceberg.TableUUID {
|
||||
t.Fatalf("after rollback TableUUID = %q, want %q", state.metadata.Iceberg.TableUUID, prior.metadata.Iceberg.TableUUID)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCleanupSafeMetadataKeepsFailedRollback proves a flip failure at index 2
|
||||
// where table 0's rollback also fails leaves table 0's new metadata in place
|
||||
// (its pointer still references it), while the rolled-back table 1 and the
|
||||
// never-flipped tables 2 and 3 are safe to delete.
|
||||
func TestCleanupSafeMetadataKeepsFailedRollback(t *testing.T) {
|
||||
prepared := make([]preparedTableCommit, 4)
|
||||
rolledBack := []bool{false, true} // tables 0 and 1 were flipped; 0's rollback failed
|
||||
|
||||
safe := cleanupSafeMetadata(prepared, 2, rolledBack)
|
||||
|
||||
want := []bool{false, true, true, true}
|
||||
for i := range want {
|
||||
if safe[i] != want[i] {
|
||||
t.Fatalf("safe[%d] = %v, want %v (full %v)", i, safe[i], want[i], safe)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestCleanupSafeMetadataFirstFlipFails covers a flip failure at the first table:
|
||||
// nothing was flipped, so every written file is safe to delete.
|
||||
func TestCleanupSafeMetadataFirstFlipFails(t *testing.T) {
|
||||
prepared := make([]preparedTableCommit, 3)
|
||||
safe := cleanupSafeMetadata(prepared, 0, nil)
|
||||
for i, ok := range safe {
|
||||
if !ok {
|
||||
t.Fatalf("safe[%d] = false, want true (full %v)", i, safe)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestCleanupSafeMetadataAllRolledBack covers the common case where every prior
|
||||
// flip rolls back cleanly, so all metadata files are safe to delete.
|
||||
func TestCleanupSafeMetadataAllRolledBack(t *testing.T) {
|
||||
prepared := make([]preparedTableCommit, 3)
|
||||
safe := cleanupSafeMetadata(prepared, 2, []bool{true, true})
|
||||
for i, ok := range safe {
|
||||
if !ok {
|
||||
t.Fatalf("safe[%d] = false, want true (full %v)", i, safe)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// tableEntry models the persisted table fields and the partial-update rules
|
||||
// handleUpdateTable applies, so rollback can be exercised without a live filer.
|
||||
type tableEntry struct {
|
||||
metadataLocation string
|
||||
metadataVersion int
|
||||
metadata *s3tables.TableMetadata
|
||||
}
|
||||
|
||||
func (e tableEntry) clone() *tableEntry {
|
||||
return &tableEntry{
|
||||
metadataLocation: e.metadataLocation,
|
||||
metadataVersion: e.metadataVersion,
|
||||
metadata: cloneTableMetadata(e.metadata),
|
||||
}
|
||||
}
|
||||
|
||||
func (e *tableEntry) applyUpdate(req *s3tables.UpdateTableRequest) {
|
||||
if req.Metadata != nil {
|
||||
if e.metadata == nil {
|
||||
e.metadata = &s3tables.TableMetadata{}
|
||||
}
|
||||
if req.Metadata.Iceberg != nil {
|
||||
if e.metadata.Iceberg == nil {
|
||||
e.metadata.Iceberg = &s3tables.IcebergMetadata{}
|
||||
}
|
||||
if req.Metadata.Iceberg.TableUUID != "" {
|
||||
e.metadata.Iceberg.TableUUID = req.Metadata.Iceberg.TableUUID
|
||||
}
|
||||
}
|
||||
if len(req.Metadata.FullMetadata) > 0 {
|
||||
e.metadata.FullMetadata = req.Metadata.FullMetadata
|
||||
}
|
||||
}
|
||||
if req.MetadataLocation != "" {
|
||||
e.metadataLocation = req.MetadataLocation
|
||||
}
|
||||
if req.MetadataVersion > 0 {
|
||||
e.metadataVersion = req.MetadataVersion
|
||||
}
|
||||
}
|
||||
@ -108,6 +108,9 @@ func (s *Server) RegisterRoutes(router *mux.Router) {
|
||||
router.HandleFunc("/v1/namespaces/{namespace}/tables/{table}", s.Auth(s.handleDropTable)).Methods(http.MethodDelete)
|
||||
router.HandleFunc("/v1/namespaces/{namespace}/tables/{table}", s.Auth(s.handleUpdateTable)).Methods(http.MethodPost)
|
||||
|
||||
// Multi-table transaction commit - wrapped with Auth middleware
|
||||
router.HandleFunc("/v1/transactions/commit", s.Auth(s.handleCommitTransaction)).Methods(http.MethodPost)
|
||||
|
||||
// With prefix support - wrapped with Auth middleware
|
||||
router.HandleFunc("/v1/{prefix}/namespaces", s.Auth(s.handleListNamespaces)).Methods(http.MethodGet)
|
||||
router.HandleFunc("/v1/{prefix}/namespaces", s.Auth(s.handleCreateNamespace)).Methods(http.MethodPost)
|
||||
@ -122,6 +125,7 @@ func (s *Server) RegisterRoutes(router *mux.Router) {
|
||||
router.HandleFunc("/v1/{prefix}/namespaces/{namespace}/tables/{table}", s.Auth(s.handleTableExists)).Methods(http.MethodHead)
|
||||
router.HandleFunc("/v1/{prefix}/namespaces/{namespace}/tables/{table}", s.Auth(s.handleDropTable)).Methods(http.MethodDelete)
|
||||
router.HandleFunc("/v1/{prefix}/namespaces/{namespace}/tables/{table}", s.Auth(s.handleUpdateTable)).Methods(http.MethodPost)
|
||||
router.HandleFunc("/v1/{prefix}/transactions/commit", s.Auth(s.handleCommitTransaction)).Methods(http.MethodPost)
|
||||
|
||||
// Catch-all for debugging
|
||||
router.PathPrefix("/").HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
Loading…
Reference in New Issue
Block a user