From 5eca11c27462f4ae2064638d4f2cd8e8ab96c582 Mon Sep 17 00:00:00 2001 From: jiangxinmeng Date: Mon, 10 Aug 2026 10:33:56 +0800 Subject: [PATCH 1/5] fix: support legacy tombstone abort columns --- pkg/objectio/ioutil/funcs.go | 4 +- pkg/objectio/ioutil/funcs_test.go | 72 +++++++++++++++++++++++++++++++ 2 files changed, 74 insertions(+), 2 deletions(-) create mode 100644 pkg/objectio/ioutil/funcs_test.go diff --git a/pkg/objectio/ioutil/funcs.go b/pkg/objectio/ioutil/funcs.go index 68504f5dd8017..23cccd869af50 100644 --- a/pkg/objectio/ioutil/funcs.go +++ b/pkg/objectio/ioutil/funcs.go @@ -351,7 +351,7 @@ func IsRowDeletedByLocation( tss := vector.MustFixedColNoTypeCheck[types.TS](&data[1]) abortVec := &data[2] var aborts []bool - if !abortVec.IsConstNull() { + if !abortVec.IsConstNull() && abortVec.Length() == len(rowids) { aborts = vector.MustFixedColNoTypeCheck[bool](abortVec) } for i := idx; i < len(rowids); i++ { @@ -469,7 +469,7 @@ func EvalDeleteMaskFromDNCreatedTombstones( noTSCheck := false var aborts []bool - if abortVec != nil && !abortVec.IsConstNull() { + if abortVec != nil && !abortVec.IsConstNull() && abortVec.Length() == len(rowids) { aborts = vector.MustFixedColWithTypeCheck[bool](abortVec) } if end-start > 10 && aborts == nil { diff --git a/pkg/objectio/ioutil/funcs_test.go b/pkg/objectio/ioutil/funcs_test.go new file mode 100644 index 0000000000000..898512f52b1bc --- /dev/null +++ b/pkg/objectio/ioutil/funcs_test.go @@ -0,0 +1,72 @@ +// 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.NewVec(types.T_bool.ToType()) + defer abortColumn.Free(mp) + + rows := EvalDeleteMaskFromDNCreatedTombstones( + rowIDs, commitTS, abortColumn, objectio.BlockObject{}, types.BuildTSForTest(2, 0), &blockID, + ) + 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 := EvalDeleteMaskFromDNCreatedTombstones( + rowIDs, commitTS, abortColumn, objectio.BlockObject{}, types.BuildTSForTest(2, 0), &blockID, + ) + 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() + }) +} From 6993c2f0000f235ae7c3ea8604c02d8c83ce3b8a Mon Sep 17 00:00:00 2001 From: jiangxinmeng Date: Mon, 10 Aug 2026 14:12:54 +0800 Subject: [PATCH 2/5] fix: validate tombstone abort metadata --- pkg/objectio/ioutil/funcs.go | 101 +++++++++++++++--- pkg/objectio/ioutil/funcs_test.go | 48 ++++++++- pkg/objectio/ioutil/loadfuncs_test.go | 87 +++++++++++++++ .../disttae/local_disttae_datasource.go | 10 +- .../disttae/logtailreplay/partition_state.go | 41 +++---- pkg/vm/engine/disttae/txn_table.go | 19 +--- 6 files changed, 252 insertions(+), 54 deletions(-) diff --git a/pkg/objectio/ioutil/funcs.go b/pkg/objectio/ioutil/funcs.go index 23cccd869af50..d01c9c8fa802b 100644 --- a/pkg/objectio/ioutil/funcs.go +++ b/pkg/objectio/ioutil/funcs.go @@ -349,16 +349,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() && abortVec.Length() == len(rowids) { - 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 } @@ -402,7 +401,7 @@ func FillBlockDeleteMask( if createdByCN { deleteMask = EvalDeleteMaskFromCNCreatedTombstones(blockId, &persistedDeletes[0]) } else { - deleteMask = EvalDeleteMaskFromDNCreatedTombstones( + deleteMask, err = EvalDeleteMaskFromDNCreatedTombstones( &persistedDeletes[0], &persistedDeletes[1], &persistedDeletes[2], @@ -445,9 +444,83 @@ 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 column with length zero is the explicitly supported legacy representation +// for an object written before abort metadata was persisted. +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 a zero-length column 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.Length() == 0 { + return TombstoneAbortColumn{}, nil + } + if abortVec.Length() != expectedRows { + return TombstoneAbortColumn{}, moerr.NewInvalidInputNoCtxf( + "tombstone abort column has %d rows, expected %d", + abortVec.Length(), expectedRows, + ) + } + if abortVec.IsConstNull() { + return TombstoneAbortColumn{}, moerr.NewInvalidInputNoCtx("tombstone abort column contains null values") + } + 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( @@ -457,22 +530,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 nil, err + } start, end := FindStartEndOfBlockFromSortedRowids(rowids, blockid) if start >= end { return } noTSCheck := false - var aborts []bool - if abortVec != nil && !abortVec.IsConstNull() && abortVec.Length() == len(rowids) { - 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) @@ -489,7 +562,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() diff --git a/pkg/objectio/ioutil/funcs_test.go b/pkg/objectio/ioutil/funcs_test.go index 898512f52b1bc..5539a2df0c4a5 100644 --- a/pkg/objectio/ioutil/funcs_test.go +++ b/pkg/objectio/ioutil/funcs_test.go @@ -43,9 +43,10 @@ func TestEvalDeleteMaskFromDNCreatedTombstonesAbortColumnCompatibility(t *testin abortColumn := vector.NewVec(types.T_bool.ToType()) defer abortColumn.Free(mp) - rows := EvalDeleteMaskFromDNCreatedTombstones( + 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)) @@ -60,9 +61,10 @@ func TestEvalDeleteMaskFromDNCreatedTombstonesAbortColumnCompatibility(t *testin } defer abortColumn.Free(mp) - rows := EvalDeleteMaskFromDNCreatedTombstones( + 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)) @@ -70,3 +72,45 @@ func TestEvalDeleteMaskFromDNCreatedTombstonesAbortColumnCompatibility(t *testin rows.Release() }) } + +func TestValidateTombstoneAbortColumn(t *testing.T) { + mp := mpool.MustNewZero() + defer mpool.DeleteMPool(mp) + + legacy := vector.NewVec(types.T_bool.ToType()) + 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) + _, err = ValidateTombstoneAbortColumn(3, constNull) + require.Error(t, err) + constNull.Free(mp) +} diff --git a/pkg/objectio/ioutil/loadfuncs_test.go b/pkg/objectio/ioutil/loadfuncs_test.go index 7df6be4c9128c..770284d274ca5 100644 --- a/pkg/objectio/ioutil/loadfuncs_test.go +++ b/pkg/objectio/ioutil/loadfuncs_test.go @@ -151,6 +151,93 @@ 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(2) + input.Vecs[0] = vector.NewVec(types.T_Rowid.ToType()) + input.Vecs[1] = 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], 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, 0, cacheVectors[2].Length()) + 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 diff --git a/pkg/vm/engine/disttae/local_disttae_datasource.go b/pkg/vm/engine/disttae/local_disttae_datasource.go index 7060211240b46..64aea7f2c9d70 100644 --- a/pkg/vm/engine/disttae/local_disttae_datasource.go +++ b/pkg/vm/engine/disttae/local_disttae_datasource.go @@ -1631,13 +1631,15 @@ func (ls *LocalDisttaeDataSource) batchApplyTombstoneObjects( var deletedRowIds []objectio.Rowid var commit []types.TS - var aborts []bool + var abortColumn ioutil.TombstoneAbortColumn deletedRowIds = vector.MustFixedColWithTypeCheck[objectio.Rowid](&cacheVectors[0]) if !obj.GetCNCreated() { commit = vector.MustFixedColWithTypeCheck[types.TS](&cacheVectors[1]) - if !cacheVectors[2].IsConstNull() { - aborts = vector.MustFixedColWithTypeCheck[bool](&cacheVectors[2]) + var abortErr error + abortColumn, abortErr = ioutil.ValidateTombstoneAbortColumn(len(deletedRowIds), &cacheVectors[2]) + if abortErr != nil { + return abortErr } } @@ -1647,7 +1649,7 @@ func (ls *LocalDisttaeDataSource) batchApplyTombstoneObjects( for j := s; j < e; j++ { if rowIds[i].EQ(&deletedRowIds[j]) && - (aborts == nil || !aborts[j]) && + (!abortColumn.IsPresent() || !abortColumn.IsAborted(j)) && (commit == nil || commit[j].LE(&ls.snapshotTS)) { deletedMask.Add(uint64(i)) break diff --git a/pkg/vm/engine/disttae/logtailreplay/partition_state.go b/pkg/vm/engine/disttae/logtailreplay/partition_state.go index 645639592258c..3cb7e5c477525 100644 --- a/pkg/vm/engine/disttae/logtailreplay/partition_state.go +++ b/pkg/vm/engine/disttae/logtailreplay/partition_state.go @@ -1253,13 +1253,13 @@ func (p *PartitionState) countVisibleRowsInAppendableObject( return true } commitTSCol := vector.MustFixedColWithTypeCheck[types.TS](&cacheVectors[0]) - abortVec := &cacheVectors[1] - var aborts []bool - if !abortVec.IsConstNull() { - aborts = vector.MustFixedColWithTypeCheck[bool](abortVec) + abortColumn, err := ioutil.ValidateTombstoneAbortColumn(len(commitTSCol), &cacheVectors[1]) + if err != nil { + loadErr = err + return false } for row, ts := range commitTSCol { - if (aborts == nil || !aborts[row]) && ts.LE(&snapshot) { + if (!abortColumn.IsPresent() || !abortColumn.IsAborted(row)) && ts.LE(&snapshot) { count++ } } @@ -1607,13 +1607,14 @@ func (p *PartitionState) countTombstoneStatsLinear( rowIds := vector.MustFixedColNoTypeCheck[types.Rowid](&persistedDeletes[0]) var commitTSs []types.TS - var aborts []bool + var abortColumn ioutil.TombstoneAbortColumn // When cnCreated=false (TN created), ReadDeletes reads [Rowid, CommitTS] at indices [0, 1] // When cnCreated=true (CN created), ReadDeletes only reads [Rowid] at index [0], no CommitTS if needCheckCommitTs && len(persistedDeletes) > 2 { commitTSs = vector.MustFixedColNoTypeCheck[types.TS](&persistedDeletes[1]) - if !persistedDeletes[2].IsConstNull() { - aborts = vector.MustFixedColNoTypeCheck[bool](&persistedDeletes[2]) + abortColumn, readErr = ioutil.ValidateTombstoneAbortColumn(len(rowIds), &persistedDeletes[2]) + if readErr != nil { + return false } } @@ -1627,7 +1628,7 @@ func (p *PartitionState) countTombstoneStatsLinear( continue } - if (aborts != nil && aborts[j]) || + if (abortColumn.IsPresent() && abortColumn.IsAborted(j)) || (needCheckCommitTs && len(commitTSs) > 0 && commitTSs[j].GT(&snapshot)) { continue } @@ -1711,13 +1712,14 @@ func (p *PartitionState) countTombstoneStatsWithMap( rowIds := vector.MustFixedColNoTypeCheck[types.Rowid](&persistedDeletes[0]) var commitTSs []types.TS - var aborts []bool + var abortColumn ioutil.TombstoneAbortColumn // When cnCreated=false (TN created), ReadDeletes reads [Rowid, CommitTS] at indices [0, 1] // When cnCreated=true (CN created), ReadDeletes only reads [Rowid] at index [0], no CommitTS if needCheckCommitTs && len(persistedDeletes) > 2 { commitTSs = vector.MustFixedColNoTypeCheck[types.TS](&persistedDeletes[1]) - if !persistedDeletes[2].IsConstNull() { - aborts = vector.MustFixedColNoTypeCheck[bool](&persistedDeletes[2]) + abortColumn, readErr = ioutil.ValidateTombstoneAbortColumn(len(rowIds), &persistedDeletes[2]) + if readErr != nil { + return false } } @@ -1730,7 +1732,7 @@ func (p *PartitionState) countTombstoneStatsWithMap( continue } - if (aborts != nil && aborts[j]) || + if (abortColumn.IsPresent() && abortColumn.IsAborted(j)) || (needCheckCommitTs && len(commitTSs) > 0 && commitTSs[j].GT(&snapshot)) { continue } @@ -1810,7 +1812,7 @@ type tombstoneBlockIterator struct { blockIdx int rowIds []types.Rowid commitTSs []types.TS - aborts []bool + abortColumn ioutil.TombstoneAbortColumn rowIdx int needCheckTS bool snapshot types.TS @@ -1850,14 +1852,13 @@ func (it *tombstoneBlockIterator) loadNextBlock() bool { if it.needCheckTS && len(it.persistedDel) > 2 { it.commitTSs = vector.MustFixedColNoTypeCheck[types.TS](&it.persistedDel[1]) - if !it.persistedDel[2].IsConstNull() { - it.aborts = vector.MustFixedColNoTypeCheck[bool](&it.persistedDel[2]) - } else { - it.aborts = nil + it.abortColumn, it.err = ioutil.ValidateTombstoneAbortColumn(len(it.rowIds), &it.persistedDel[2]) + if it.err != nil { + return false } } else { it.commitTSs = nil - it.aborts = nil + it.abortColumn = ioutil.TombstoneAbortColumn{} } it.rowIdx = 0 @@ -1889,7 +1890,7 @@ func (it *tombstoneBlockIterator) next() bool { // Persisted tombstone iterator for { for it.rowIdx < len(it.rowIds) { - if (it.aborts != nil && it.aborts[it.rowIdx]) || + if (it.abortColumn.IsPresent() && it.abortColumn.IsAborted(it.rowIdx)) || (it.needCheckTS && len(it.commitTSs) > 0 && it.commitTSs[it.rowIdx].GT(&it.snapshot)) { it.rowIdx++ continue diff --git a/pkg/vm/engine/disttae/txn_table.go b/pkg/vm/engine/disttae/txn_table.go index f6f32921d264e..4615041e8296f 100644 --- a/pkg/vm/engine/disttae/txn_table.go +++ b/pkg/vm/engine/disttae/txn_table.go @@ -2616,15 +2616,9 @@ func pkCommitTSMatchedInRange( return false, false } timestamps := vector.MustFixedColWithTypeCheck[types.TS](commitTSVec) - var aborts []bool - if abortVec != nil && !abortVec.IsConstNull() { - if abortVec.GetType().Oid != types.T_bool { - return false, false - } - aborts = vector.MustFixedColWithTypeCheck[bool](abortVec) - if len(aborts) != len(timestamps) { - return false, false - } + abortColumn, err := ioutil.ValidateTombstoneAbortColumn(len(timestamps), abortVec) + if err != nil { + return false, false } for _, sel := range sels { if sel < 0 || int(sel) >= len(timestamps) { @@ -2633,11 +2627,8 @@ func pkCommitTSMatchedInRange( if commitTSVec.IsNull(uint64(sel)) { return false, false } - if aborts != nil { - if abortVec.IsNull(uint64(sel)) { - return false, false - } - if aborts[sel] { + if abortColumn.IsPresent() { + if abortColumn.IsAborted(int(sel)) { continue } } From d1caf4d88f6d6957b93b14218398d890f858dcdb Mon Sep 17 00:00:00 2001 From: jiangxinmeng Date: Mon, 10 Aug 2026 14:40:28 +0800 Subject: [PATCH 3/5] fix: declare ReadDeletes return values --- pkg/objectio/ioutil/funcs.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pkg/objectio/ioutil/funcs.go b/pkg/objectio/ioutil/funcs.go index d01c9c8fa802b..974678aacd44e 100644 --- a/pkg/objectio/ioutil/funcs.go +++ b/pkg/objectio/ioutil/funcs.go @@ -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" @@ -422,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 From 42572a7eeb566deb4bc42f2c2dbe990c9822a5c8 Mon Sep 17 00:00:00 2001 From: jiangxinmeng Date: Mon, 10 Aug 2026 14:58:29 +0800 Subject: [PATCH 4/5] fix: accept legacy const-null abort columns --- pkg/objectio/ioutil/funcs.go | 29 ++++++++++++++++++--------- pkg/objectio/ioutil/funcs_test.go | 14 +++++++++---- pkg/objectio/ioutil/loadfuncs_test.go | 11 ++++++---- 3 files changed, 36 insertions(+), 18 deletions(-) diff --git a/pkg/objectio/ioutil/funcs.go b/pkg/objectio/ioutil/funcs.go index 974678aacd44e..24c6d386082f7 100644 --- a/pkg/objectio/ioutil/funcs.go +++ b/pkg/objectio/ioutil/funcs.go @@ -473,8 +473,10 @@ func ReadDeletes( } // TombstoneAbortColumn is the validated abort metadata for a tombstone block. -// A column with length zero is the explicitly supported legacy representation -// for an object written before abort metadata was persisted. +// 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 } @@ -488,9 +490,10 @@ func (c TombstoneAbortColumn) IsAborted(row int) bool { } // ValidateTombstoneAbortColumn validates the abort column against the row -// count shared by all tombstone consumers. Only a zero-length column is -// treated as the legacy no-abort representation. All other malformed layouts -// are returned as errors instead of being silently treated as live rows. +// 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, @@ -504,8 +507,17 @@ func ValidateTombstoneAbortColumn( 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{}, nil + return TombstoneAbortColumn{}, moerr.NewInvalidInputNoCtx("tombstone abort column is empty") } if abortVec.Length() != expectedRows { return TombstoneAbortColumn{}, moerr.NewInvalidInputNoCtxf( @@ -513,9 +525,6 @@ func ValidateTombstoneAbortColumn( abortVec.Length(), expectedRows, ) } - if abortVec.IsConstNull() { - return TombstoneAbortColumn{}, moerr.NewInvalidInputNoCtx("tombstone abort column contains null values") - } for i := 0; i < expectedRows; i++ { if abortVec.IsNull(uint64(i)) { return TombstoneAbortColumn{}, moerr.NewInvalidInputNoCtxf("tombstone abort row %d is null", i) @@ -538,7 +547,7 @@ func EvalDeleteMaskFromDNCreatedTombstones( rowids := vector.MustFixedColWithTypeCheck[types.Rowid](deletedRows) aborts, err := ValidateTombstoneAbortColumn(len(rowids), abortVec) if err != nil { - return nil, err + return objectio.NullBitmap, err } start, end := FindStartEndOfBlockFromSortedRowids(rowids, blockid) if start >= end { diff --git a/pkg/objectio/ioutil/funcs_test.go b/pkg/objectio/ioutil/funcs_test.go index 5539a2df0c4a5..6981834245240 100644 --- a/pkg/objectio/ioutil/funcs_test.go +++ b/pkg/objectio/ioutil/funcs_test.go @@ -40,7 +40,7 @@ func TestEvalDeleteMaskFromDNCreatedTombstonesAbortColumnCompatibility(t *testin defer commitTS.Free(mp) t.Run("legacy checkpoint without abort column", func(t *testing.T) { - abortColumn := vector.NewVec(types.T_bool.ToType()) + abortColumn := vector.NewConstNull(types.T_bool.ToType(), 3, mp) defer abortColumn.Free(mp) rows, err := EvalDeleteMaskFromDNCreatedTombstones( @@ -77,7 +77,7 @@ func TestValidateTombstoneAbortColumn(t *testing.T) { mp := mpool.MustNewZero() defer mpool.DeleteMPool(mp) - legacy := vector.NewVec(types.T_bool.ToType()) + legacy := vector.NewConstNull(types.T_bool.ToType(), 3, mp) column, err := ValidateTombstoneAbortColumn(3, legacy) require.NoError(t, err) require.False(t, column.IsPresent()) @@ -110,7 +110,13 @@ func TestValidateTombstoneAbortColumn(t *testing.T) { constAbort.Free(mp) constNull := vector.NewConstNull(types.T_bool.ToType(), 3, mp) - _, err = ValidateTombstoneAbortColumn(3, constNull) - require.Error(t, err) + 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) } diff --git a/pkg/objectio/ioutil/loadfuncs_test.go b/pkg/objectio/ioutil/loadfuncs_test.go index 770284d274ca5..0d87101db5bdf 100644 --- a/pkg/objectio/ioutil/loadfuncs_test.go +++ b/pkg/objectio/ioutil/loadfuncs_test.go @@ -157,14 +157,16 @@ func TestReadDeletesSupportsLegacyTombstoneWithoutAbortColumn(t *testing.T) { mp := mpool.MustNewZero() defer mpool.DeleteMPool(mp) - input := batch.NewWithSize(2) + input := batch.NewWithSize(3) input.Vecs[0] = vector.NewVec(types.T_Rowid.ToType()) - input.Vecs[1] = vector.NewVec(types.T_TS.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], types.BuildTS(1, 0), 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) @@ -187,7 +189,8 @@ func TestReadDeletesSupportsLegacyTombstoneWithoutAbortColumn(t *testing.T) { defer release() require.Equal(t, 2, cacheVectors[0].Length()) require.Equal(t, 2, cacheVectors[1].Length()) - require.Equal(t, 0, cacheVectors[2].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()) From e863c805f5db11b656e6f0ca6c09ebbb008c272e Mon Sep 17 00:00:00 2001 From: jiangxinmeng Date: Mon, 10 Aug 2026 15:18:34 +0800 Subject: [PATCH 5/5] fix: propagate initial tombstone merge errors --- .../disttae/logtailreplay/partition_state.go | 25 ++++++++--- .../logtailreplay/partition_state_test.go | 42 +++++++++++++++++++ 2 files changed, 62 insertions(+), 5 deletions(-) diff --git a/pkg/vm/engine/disttae/logtailreplay/partition_state.go b/pkg/vm/engine/disttae/logtailreplay/partition_state.go index 3cb7e5c477525..5bf2e0220f09d 100644 --- a/pkg/vm/engine/disttae/logtailreplay/partition_state.go +++ b/pkg/vm/engine/disttae/logtailreplay/partition_state.go @@ -1958,6 +1958,14 @@ func (p *PartitionState) countTombstoneStatsWithMerge( stats TombstoneStats, ) (TombstoneStats, error) { iterators := make([]*tombstoneBlockIterator, 0, len(objects)) + releaseIterators := func() { + for _, it := range iterators { + if it.release != nil { + it.release() + it.release = nil + } + } + } for _, obj := range objects { cnCreated := obj.GetCNCreated() @@ -1996,9 +2004,18 @@ func (p *PartitionState) countTombstoneStatsWithMerge( p: p, } - if it.next() { - iterators = append(iterators, it) + if !it.next() { + if it.release != nil { + it.release() + it.release = nil + } + if it.err != nil { + releaseIterators() + return stats, it.err + } + continue } + iterators = append(iterators, it) } // Add in-memory tombstones as an iterator @@ -2039,10 +2056,8 @@ func (p *PartitionState) countTombstoneStatsWithMerge( } } + releaseIterators() for _, it := range iterators { - if it.release != nil { - it.release() - } if it.err != nil { return stats, it.err } diff --git a/pkg/vm/engine/disttae/logtailreplay/partition_state_test.go b/pkg/vm/engine/disttae/logtailreplay/partition_state_test.go index 3fa0fb9d894c4..135cb0a6185d4 100644 --- a/pkg/vm/engine/disttae/logtailreplay/partition_state_test.go +++ b/pkg/vm/engine/disttae/logtailreplay/partition_state_test.go @@ -1973,6 +1973,48 @@ func TestCountTombstoneRowsReadError(t *testing.T) { } } +func TestCountTombstoneStatsWithMergePropagatesInitialReadError(t *testing.T) { + ctx := context.Background() + fs := testutil.NewSharedFS() + mp := mpool.MustNewZero() + defer mpool.DeleteMPool(mp) + + blockID := types.Blockid{} + bat := batch.NewWithSize(3) + bat.Vecs[0] = vector.NewVec(types.T_Rowid.ToType()) + bat.Vecs[1] = vector.NewVec(types.T_TS.ToType()) + bat.Vecs[2] = vector.NewVec(types.T_bool.ToType()) + for offset := uint32(1); offset <= 3; offset++ { + require.NoError(t, vector.AppendFixed(bat.Vecs[0], types.NewRowid(&blockID, offset), false, mp)) + require.NoError(t, vector.AppendFixed(bat.Vecs[1], types.BuildTS(1, 0), false, mp)) + } + // Persist a malformed abort column shorter than the rowid/commit columns. + require.NoError(t, vector.AppendFixed(bat.Vecs[2], false, false, mp)) + bat.SetRowCount(3) + + writer := ioutil.ConstructWriter( + 0, + []uint16{0, objectio.SEQNUM_COMMITTS, objectio.SEQNUM_ABORT}, + -1, + false, + false, + fs, + ) + _, err := writer.WriteBatch(bat) + require.NoError(t, err) + _, _, err = writer.Sync(ctx) + require.NoError(t, err) + + objects := []objectio.ObjectEntry{{ + ObjectStats: writer.GetObjectStats(), + CreateTime: types.BuildTS(1, 0), + }} + _, err = NewPartitionState("", false, 42, false).countTombstoneStatsWithMerge( + ctx, types.BuildTS(10, 0), fs, objects, TombstoneStats{}, + ) + require.Error(t, err) +} + func TestCountTombstoneRowsEdgeCases(t *testing.T) { // Test various edge cases for CountTombstoneRows ctx := context.Background()