Skip to content
113 changes: 98 additions & 15 deletions pkg/objectio/ioutil/funcs.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import (
"sort"

"github.com/matrixorigin/matrixone/pkg/common/bitmap"
"github.com/matrixorigin/matrixone/pkg/common/moerr"
"github.com/matrixorigin/matrixone/pkg/container/types"
"github.com/matrixorigin/matrixone/pkg/container/vector"
"github.com/matrixorigin/matrixone/pkg/fileservice"
Expand Down Expand Up @@ -349,16 +350,15 @@ func IsRowDeletedByLocation(
deleted = (idx < len(rowids)) && (rowids[idx].EQ(row))
} else {
tss := vector.MustFixedColNoTypeCheck[types.TS](&data[1])
abortVec := &data[2]
var aborts []bool
if !abortVec.IsConstNull() {
aborts = vector.MustFixedColNoTypeCheck[bool](abortVec)
aborts, err := ValidateTombstoneAbortColumn(len(rowids), &data[2])
if err != nil {
return false, err
}
for i := idx; i < len(rowids); i++ {
if !rowids[i].EQ(row) {
break
}
if (aborts == nil || !aborts[i]) && tss[i].LE(snapshotTS) {
if (!aborts.IsPresent() || !aborts.IsAborted(i)) && tss[i].LE(snapshotTS) {
deleted = true
break
}
Expand Down Expand Up @@ -402,7 +402,7 @@ func FillBlockDeleteMask(
if createdByCN {
deleteMask = EvalDeleteMaskFromCNCreatedTombstones(blockId, &persistedDeletes[0])
} else {
deleteMask = EvalDeleteMaskFromDNCreatedTombstones(
deleteMask, err = EvalDeleteMaskFromDNCreatedTombstones(
&persistedDeletes[0],
&persistedDeletes[1],
&persistedDeletes[2],
Expand All @@ -423,7 +423,7 @@ func ReadDeletes(
isPersistedByCN bool,
cacheVectors containers.Vectors,
pkType *types.Type,
) (objectio.ObjectDataMeta, func(), error) {
) (meta objectio.ObjectDataMeta, release func(), err error) {

var cols []uint16
var typs []types.Type
Expand All @@ -445,9 +445,92 @@ func ReadDeletes(
typs = append(typs[:1], append([]types.Type{*pkType}, typs[1:]...)...)
}

return LoadTombstoneColumns(
meta, release, err = LoadTombstoneColumns(
ctx, cols, typs, fs, deltaLoc, cacheVectors, nil, fileservice.Policy(0),
)
if err != nil || isPersistedByCN {
return
}
if len(cacheVectors) < 3 {
err = moerr.NewInvalidInputNoCtxf(
"tombstone column count %d is missing commit-ts or abort metadata",
len(cacheVectors),
)
if release != nil {
release()
}
release = nil
return
}
rowCount := cacheVectors[0].Length()
if _, err = ValidateTombstoneAbortColumn(rowCount, &cacheVectors[len(cacheVectors)-1]); err != nil {
if release != nil {
release()
}
release = nil
}
return
}

// TombstoneAbortColumn is the validated abort metadata for a tombstone block.
// A const-null bool column with the block row count is the explicitly supported
// legacy representation for an object written before abort metadata was
// persisted. The object reader synthesizes this marker when the physical
// column is absent.
type TombstoneAbortColumn struct {
vec *vector.Vector
}

func (c TombstoneAbortColumn) IsPresent() bool {
return c.vec != nil
}

func (c TombstoneAbortColumn) IsAborted(row int) bool {
return vector.GetFixedAtNoTypeCheck[bool](c.vec, row)
}

// ValidateTombstoneAbortColumn validates the abort column against the row
// count shared by all tombstone consumers. Only the reader's const-null
// sentinel is treated as the legacy no-abort representation. All other
// malformed layouts are returned as errors instead of being silently treated
// as live rows.
func ValidateTombstoneAbortColumn(
expectedRows int,
abortVec *vector.Vector,
) (TombstoneAbortColumn, error) {
if abortVec == nil {
return TombstoneAbortColumn{}, nil
}
if abortVec.GetType().Oid != types.T_bool {
return TombstoneAbortColumn{}, moerr.NewInvalidInputNoCtxf(
"tombstone abort column has type %s, expected bool",
abortVec.GetType().String(),
)
}
if abortVec.IsConstNull() {
if abortVec.Length() == expectedRows {
return TombstoneAbortColumn{}, nil
}
return TombstoneAbortColumn{}, moerr.NewInvalidInputNoCtxf(
"tombstone abort const-null column has %d rows, expected %d",
abortVec.Length(), expectedRows,
)
}
if abortVec.Length() == 0 {
return TombstoneAbortColumn{}, moerr.NewInvalidInputNoCtx("tombstone abort column is empty")
}
if abortVec.Length() != expectedRows {
return TombstoneAbortColumn{}, moerr.NewInvalidInputNoCtxf(
"tombstone abort column has %d rows, expected %d",
abortVec.Length(), expectedRows,
)
}
for i := 0; i < expectedRows; i++ {
if abortVec.IsNull(uint64(i)) {
return TombstoneAbortColumn{}, moerr.NewInvalidInputNoCtxf("tombstone abort row %d is null", i)
}
}
return TombstoneAbortColumn{vec: abortVec}, nil
}

func EvalDeleteMaskFromDNCreatedTombstones(
Expand All @@ -457,22 +540,22 @@ func EvalDeleteMaskFromDNCreatedTombstones(
meta objectio.BlockObject,
ts *types.TS,
blockid *types.Blockid,
) (rows objectio.Bitmap) {
) (rows objectio.Bitmap, err error) {
if deletedRows == nil {
return
}
rowids := vector.MustFixedColWithTypeCheck[types.Rowid](deletedRows)
aborts, err := ValidateTombstoneAbortColumn(len(rowids), abortVec)
if err != nil {
return objectio.NullBitmap, err
}
start, end := FindStartEndOfBlockFromSortedRowids(rowids, blockid)
if start >= end {
return
}

noTSCheck := false
var aborts []bool
if abortVec != nil && !abortVec.IsConstNull() {
aborts = vector.MustFixedColWithTypeCheck[bool](abortVec)
}
if end-start > 10 && aborts == nil {
if end-start > 10 && !aborts.IsPresent() {
// fast path is true if the maxTS is less than the snapshotTS
// this means that all the rows between start and end are visible
layout := objectio.ResolveSpecialColumnLayout(meta)
Expand All @@ -489,7 +572,7 @@ func EvalDeleteMaskFromDNCreatedTombstones(
} else {
tss := vector.MustFixedColWithTypeCheck[types.TS](commitTSVec)
for i := end - 1; i >= start; i-- {
if (aborts != nil && aborts[i]) || tss[i].GT(ts) {
if (aborts.IsPresent() && aborts.IsAborted(i)) || tss[i].GT(ts) {
continue
}
row := rowids[i].GetRowOffset()
Expand Down
122 changes: 122 additions & 0 deletions pkg/objectio/ioutil/funcs_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,122 @@
// Copyright 2021 Matrix Origin
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

package ioutil

import (
"testing"

"github.com/stretchr/testify/require"

"github.com/matrixorigin/matrixone/pkg/common/mpool"
"github.com/matrixorigin/matrixone/pkg/container/types"
"github.com/matrixorigin/matrixone/pkg/container/vector"
"github.com/matrixorigin/matrixone/pkg/objectio"
)

func TestEvalDeleteMaskFromDNCreatedTombstonesAbortColumnCompatibility(t *testing.T) {
mp := mpool.MustNewZero()
defer mpool.DeleteMPool(mp)

blockID := types.Blockid{}
rowIDs := vector.NewVec(objectio.RowidType)
commitTS := vector.NewVec(objectio.TSType)
for _, offset := range []uint32{1, 2, 3} {
require.NoError(t, vector.AppendFixed(rowIDs, types.NewRowid(&blockID, offset), false, mp))
require.NoError(t, vector.AppendFixed(commitTS, types.BuildTS(1, 0), false, mp))
}
defer rowIDs.Free(mp)
defer commitTS.Free(mp)

t.Run("legacy checkpoint without abort column", func(t *testing.T) {
abortColumn := vector.NewConstNull(types.T_bool.ToType(), 3, mp)
defer abortColumn.Free(mp)

rows, err := EvalDeleteMaskFromDNCreatedTombstones(
rowIDs, commitTS, abortColumn, objectio.BlockObject{}, types.BuildTSForTest(2, 0), &blockID,
)
require.NoError(t, err)
require.True(t, rows.IsValid())
require.True(t, rows.Contains(1))
require.True(t, rows.Contains(2))
require.True(t, rows.Contains(3))
rows.Release()
})

t.Run("current checkpoint with abort column", func(t *testing.T) {
abortColumn := vector.NewVec(types.T_bool.ToType())
for _, aborted := range []bool{false, true, false} {
require.NoError(t, vector.AppendFixed(abortColumn, aborted, false, mp))
}
defer abortColumn.Free(mp)

rows, err := EvalDeleteMaskFromDNCreatedTombstones(
rowIDs, commitTS, abortColumn, objectio.BlockObject{}, types.BuildTSForTest(2, 0), &blockID,
)
require.NoError(t, err)
require.True(t, rows.IsValid())
require.True(t, rows.Contains(1))
require.False(t, rows.Contains(2))
require.True(t, rows.Contains(3))
rows.Release()
})
}

func TestValidateTombstoneAbortColumn(t *testing.T) {
mp := mpool.MustNewZero()
defer mpool.DeleteMPool(mp)

legacy := vector.NewConstNull(types.T_bool.ToType(), 3, mp)
column, err := ValidateTombstoneAbortColumn(3, legacy)
require.NoError(t, err)
require.False(t, column.IsPresent())
legacy.Free(mp)

partial := vector.NewVec(types.T_bool.ToType())
require.NoError(t, vector.AppendFixed(partial, true, false, mp))
_, err = ValidateTombstoneAbortColumn(3, partial)
require.Error(t, err)
partial.Free(mp)

wrongType := vector.NewVec(types.T_int8.ToType())
_, err = ValidateTombstoneAbortColumn(3, wrongType)
require.Error(t, err)
wrongType.Free(mp)

nullAbort := vector.NewVec(types.T_bool.ToType())
require.NoError(t, vector.AppendFixed(nullAbort, false, true, mp))
_, err = ValidateTombstoneAbortColumn(1, nullAbort)
require.Error(t, err)
nullAbort.Free(mp)

constAbort, err := vector.NewConstFixed(types.T_bool.ToType(), true, 3, mp)
require.NoError(t, err)
column, err = ValidateTombstoneAbortColumn(3, constAbort)
require.NoError(t, err)
require.True(t, column.IsPresent())
require.True(t, column.IsAborted(0))
require.True(t, column.IsAborted(2))
constAbort.Free(mp)

constNull := vector.NewConstNull(types.T_bool.ToType(), 3, mp)
column, err = ValidateTombstoneAbortColumn(3, constNull)
require.NoError(t, err)
require.False(t, column.IsPresent())
constNull.Free(mp)

empty := vector.NewVec(types.T_bool.ToType())
_, err = ValidateTombstoneAbortColumn(3, empty)
require.Error(t, err)
empty.Free(mp)
}
90 changes: 90 additions & 0 deletions pkg/objectio/ioutil/loadfuncs_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,96 @@ func TestAppendableVisibilityFiltersAbortFromMaterializeAndSearch(t *testing.T)
require.Equal(t, []int64{0}, sels, "cached search must return only live visible rows")
}

func TestReadDeletesSupportsLegacyTombstoneWithoutAbortColumn(t *testing.T) {
ctx := context.Background()
fs := testutil.NewSharedFS()
mp := mpool.MustNewZero()
defer mpool.DeleteMPool(mp)

input := batch.NewWithSize(3)
input.Vecs[0] = vector.NewVec(types.T_Rowid.ToType())
input.Vecs[1] = vector.NewVec(types.T_int32.ToType())
input.Vecs[2] = vector.NewVec(types.T_TS.ToType())
defer input.Clean(mp)
blockID := types.Blockid{}
for offset := uint32(1); offset <= 2; offset++ {
require.NoError(t, vector.AppendFixed(input.Vecs[0], types.NewRowid(&blockID, offset), false, mp))
require.NoError(t, vector.AppendFixed(input.Vecs[1], int32(offset), false, mp))
require.NoError(t, vector.AppendFixed(input.Vecs[2], types.BuildTS(1, 0), false, mp))
}
input.SetRowCount(2)

writer := ConstructTombstoneWriter(objectio.HiddenColumnSelection_CommitTS, fs)
_, err := writer.WriteBatch(input)
require.NoError(t, err)
blocks, _, err := writer.Sync(ctx)
require.NoError(t, err)
require.Len(t, blocks, 1)

location := objectio.BuildLocation(
writer.GetName(),
blocks[0].GetExtent(),
uint32(input.RowCount()),
blocks[0].GetID(),
)
cacheVectors := containers.NewVectors(3)
_, release, err := ReadDeletes(ctx, location, fs, false, cacheVectors, nil)
require.NoError(t, err)
defer release()
require.Equal(t, 2, cacheVectors[0].Length())
require.Equal(t, 2, cacheVectors[1].Length())
require.Equal(t, 2, cacheVectors[2].Length())
require.True(t, cacheVectors[2].IsConstNull())
abortColumn, err := ValidateTombstoneAbortColumn(2, &cacheVectors[2])
require.NoError(t, err)
require.False(t, abortColumn.IsPresent())
}

func TestReadDeletesRejectsMalformedAbortColumn(t *testing.T) {
ctx := context.Background()
fs := testutil.NewSharedFS()
mp := mpool.MustNewZero()
defer mpool.DeleteMPool(mp)

input := batch.NewWithSize(3)
input.Vecs[0] = vector.NewVec(types.T_Rowid.ToType())
input.Vecs[1] = vector.NewVec(types.T_TS.ToType())
input.Vecs[2] = vector.NewVec(types.T_bool.ToType())
defer input.Clean(mp)
blockID := types.Blockid{}
for offset := uint32(1); offset <= 3; offset++ {
require.NoError(t, vector.AppendFixed(input.Vecs[0], types.NewRowid(&blockID, offset), false, mp))
require.NoError(t, vector.AppendFixed(input.Vecs[1], types.BuildTS(1, 0), false, mp))
}
require.NoError(t, vector.AppendFixed(input.Vecs[2], false, false, mp))
input.SetRowCount(3)

writer := ConstructWriter(
0,
[]uint16{0, objectio.SEQNUM_COMMITTS, objectio.SEQNUM_ABORT},
-1,
false,
false,
fs,
)
_, err := writer.WriteBatch(input)
require.NoError(t, err)
blocks, _, err := writer.Sync(ctx)
require.NoError(t, err)
require.Len(t, blocks, 1)

location := objectio.BuildLocation(
writer.GetName(),
blocks[0].GetExtent(),
uint32(input.RowCount()),
blocks[0].GetID(),
)
cacheVectors := containers.NewVectors(3)
_, release, err := ReadDeletes(ctx, location, fs, false, cacheVectors, nil)
require.Error(t, err)
require.Nil(t, release)
}

func (d *releaseTrackingData) Slice(length int) fscache.Data {
d.Data = d.Data.Slice(length)
return d
Expand Down
Loading
Loading