Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
b025343
fix(tae): persist and filter append abort metadata
jiangxinmeng1 Jul 29, 2026
301a570
Merge remote-tracking branch 'upstream/main' into codex/abort-column
jiangxinmeng1 Jul 29, 2026
70e953f
fix(logtail): use moerr for snapshot metadata error
jiangxinmeng1 Jul 29, 2026
9524124
fix(tae): remove redundant loop variable copy
jiangxinmeng1 Jul 29, 2026
cd122db
fix tae abort column CI regressions
jiangxinmeng1 Jul 29, 2026
91462cb
Merge remote-tracking branch 'upstream/main' into codex/abort-column
jiangxinmeng1 Aug 3, 2026
63c9a89
fix tae abort visibility regressions
jiangxinmeng1 Aug 3, 2026
83211ad
Merge remote-tracking branch 'upstream/main' into codex/abort-column
jiangxinmeng1 Aug 3, 2026
269c20b
fix mixed-version aobject abort layout
jiangxinmeng1 Aug 3, 2026
f3b5785
Merge remote-tracking branch 'upstream/main' into codex/abort-column
jiangxinmeng1 Aug 4, 2026
d891917
fix abort row offset remapping
jiangxinmeng1 Aug 4, 2026
30f350d
fix(disttae): keep rowid leading in persisted changes
jiangxinmeng1 Aug 4, 2026
c8a2acf
Merge remote-tracking branch 'upstream/main' into codex/abort-column
jiangxinmeng1 Aug 4, 2026
10de16b
fix(tae): preserve v9 aobject row coordinates
jiangxinmeng1 Aug 4, 2026
04a43fd
Merge branch 'main' into codex/abort-column
jiangxinmeng1 Aug 5, 2026
a718fc1
fix: filter rollback rows in backup reads
jiangxinmeng1 Aug 5, 2026
5f70c3f
Merge remote-tracking branch 'upstream/main' into codex/abort-column
jiangxinmeng1 Aug 5, 2026
3e95f76
Merge remote-tracking branch 'upstream/main' into codex/abort-column
jiangxinmeng1 Aug 5, 2026
a898420
fix: filter v9 rollback sentinels
jiangxinmeng1 Aug 5, 2026
3afb38a
Merge branch 'main' into codex/abort-column
jiangxinmeng1 Aug 6, 2026
f58c452
Merge branch 'main' into codex/abort-column
mergify[bot] Aug 6, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 9 additions & 8 deletions pkg/defines/const.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,14 +38,15 @@ const (
MORPCMinVersion int64 = math.MinInt64
MORPCVersion1 int64 = 1
MORPCVersion2 int64 = 2
MORPCVersion3 int64 = 3 // start from 1.3.0
MORPCVersion4 int64 = 4 // start from 2.0.1
MORPCVersion5 int64 = 5 // assignment-aware CHAR/VARCHAR casts
MORPCVersion6 int64 = 6 // ordered aggregate pipeline configuration
MORPCVersion7 int64 = 7 // structured CHECK constraint metadata and enforcement
MORPCVersion8 int64 = 8 // versioned exact runtime-filter key contract
MORPCVersion9 int64 = 9 // AUTO_INCREMENT epoch-fenced commit
MORPCLatestVersion = MORPCVersion9
MORPCVersion3 int64 = 3 // start from 1.3.0
MORPCVersion4 int64 = 4 // start from 2.0.1
MORPCVersion5 int64 = 5 // assignment-aware CHAR/VARCHAR casts
MORPCVersion6 int64 = 6 // ordered aggregate pipeline configuration
MORPCVersion7 int64 = 7 // structured CHECK constraint metadata and enforcement
MORPCVersion8 int64 = 8 // versioned exact runtime-filter key contract
MORPCVersion9 int64 = 9 // AUTO_INCREMENT epoch-fenced commit
MORPCVersion10 int64 = 10 // persisted appendable-object abort metadata
MORPCLatestVersion = MORPCVersion10
)

// DefaultLockWaitTimeoutSeconds is shared by the frontend default and by
Expand Down
153 changes: 120 additions & 33 deletions pkg/objectio/const.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ package objectio
import (
"fmt"
"math"
"slices"

"github.com/matrixorigin/matrixone/pkg/container/types"
"github.com/matrixorigin/matrixone/pkg/vm/engine/tae/index"
Expand Down Expand Up @@ -74,6 +75,75 @@ const (

const HiddenColumnSelection_None HiddenColumnSelection = 0

const InvalidSpecialColumnPosition = math.MaxUint16

// SpecialColumnLayout describes the physical metadata positions of the
// appendable-only columns. Object writers map special seqnums, in declaration
// order, immediately after MaxSeqnum. PhysicalAddr may therefore appear before
// or after CommitTS, depending on the writer. The positions are stable across
// sparse user schemas and do not depend on the total column count.
//
// Old appendable objects contain only CommitTS. New appendable objects contain
// CommitTS followed by Abort.
type SpecialColumnLayout struct {
PhysicalAddr uint16
CommitTS uint16
Abort uint16
}

func (layout SpecialColumnLayout) Resolve(seqnum uint16) (uint16, bool) {
switch seqnum {
case SEQNUM_COMMITTS:
return layout.CommitTS, layout.CommitTS != InvalidSpecialColumnPosition
case SEQNUM_ABORT:
return layout.Abort, layout.Abort != InvalidSpecialColumnPosition
default:
return InvalidSpecialColumnPosition, false
}
}

// ResolveSpecialColumnLayout resolves appendable special columns by their
// format-defined positions and validates their types. This keeps old
// commitTS-only objects readable while preventing a user TS/bool column from
// being mistaken for a hidden column.
func ResolveSpecialColumnLayout(block BlockObject) SpecialColumnLayout {
layout := SpecialColumnLayout{
PhysicalAddr: InvalidSpecialColumnPosition,
CommitTS: InvalidSpecialColumnPosition,
Abort: InvalidSpecialColumnPosition,
}
metaColumnCount := block.GetMetaColumnCount()
if metaColumnCount == 0 {
return layout
}

pos := block.GetMaxSeqnum() + 1
if pos < metaColumnCount &&
block.ColumnMeta(pos).DataType() == uint8(types.T_Rowid) {
layout.PhysicalAddr = pos
pos++
}

commitPos := pos
if commitPos >= metaColumnCount ||
block.ColumnMeta(commitPos).DataType() != uint8(types.T_TS) {
return layout
}
layout.CommitTS = commitPos

abortPos := commitPos + 1
if abortPos < metaColumnCount &&
block.ColumnMeta(abortPos).DataType() == uint8(types.T_Rowid) {
layout.PhysicalAddr = abortPos
abortPos++
}
if abortPos < metaColumnCount &&
block.ColumnMeta(abortPos).DataType() == uint8(types.T_bool) {
layout.Abort = abortPos
}
return layout
}

var (
TombstoneSeqnums_CN_Created = []uint16{0, 1}
TombstoneSeqnums_CN_Created_PhyAddr = []uint16{0, 1, SEQNUM_ROWID}
Expand All @@ -97,41 +167,53 @@ func IsPhysicalAddr(attr string) bool {
return attr == PhysicalAddr_Attr
}

func normalizeTombstoneHiddenColumns(hidden HiddenColumnSelection) HiddenColumnSelection {
// Abort is part of appendable MVCC metadata and is never meaningful without
// the commit timestamp that defines the row's visibility interval.
if hidden&HiddenColumnSelection_Abort != 0 {
hidden |= HiddenColumnSelection_CommitTS
}
return hidden
}

func GetTombstoneAttrs(hidden HiddenColumnSelection) []string {
hidden = normalizeTombstoneHiddenColumns(hidden)
var attrs []string
if hidden&HiddenColumnSelection_PhysicalAddr != 0 &&
hidden&HiddenColumnSelection_CommitTS != 0 {
return TombstoneAttrs_TN_Created_PhyAddr
}
if hidden&HiddenColumnSelection_PhysicalAddr != 0 {
return TombstoneAttrs_CN_Created_PhyAddr
attrs = TombstoneAttrs_TN_Created_PhyAddr
} else if hidden&HiddenColumnSelection_PhysicalAddr != 0 {
attrs = TombstoneAttrs_CN_Created_PhyAddr
} else if hidden&HiddenColumnSelection_CommitTS != 0 {
attrs = TombstoneAttrs_TN_Created
} else {
attrs = TombstoneAttrs_CN_Created
}
if hidden&HiddenColumnSelection_CommitTS != 0 {
return TombstoneAttrs_TN_Created
attrs = slices.Clone(attrs)
if hidden&HiddenColumnSelection_Abort != 0 {
attrs = append(attrs, TombstoneAttr_Abort_Attr)
}
return TombstoneAttrs_CN_Created
}

func GetTombstoneCommitTSAttrIdx(columnCnt uint16) uint16 {
if columnCnt == 3 {
return TombstoneAttr_NA_CommitTs_Idx
} else if columnCnt == 4 {
return TombstoneAttr_A_CommitTs_Idx
}
panic(fmt.Sprintf("invalid tombstone column count %d", columnCnt))
return attrs
}

func GetTombstoneSeqnums(hidden HiddenColumnSelection) []uint16 {
hidden = normalizeTombstoneHiddenColumns(hidden)
var seqnums []uint16
if hidden&HiddenColumnSelection_PhysicalAddr != 0 &&
hidden&HiddenColumnSelection_CommitTS != 0 {
return TombstoneSeqnums_DN_Created_PhyAddr
}
if hidden&HiddenColumnSelection_PhysicalAddr != 0 {
return TombstoneSeqnums_CN_Created_PhyAddr
seqnums = TombstoneSeqnums_DN_Created_PhyAddr
} else if hidden&HiddenColumnSelection_PhysicalAddr != 0 {
seqnums = TombstoneSeqnums_CN_Created_PhyAddr
} else if hidden&HiddenColumnSelection_CommitTS != 0 {
seqnums = TombstoneSeqnums_DN_Created
} else {
seqnums = TombstoneSeqnums_CN_Created
}
if hidden&HiddenColumnSelection_CommitTS != 0 {
return TombstoneSeqnums_DN_Created
seqnums = slices.Clone(seqnums)
if hidden&HiddenColumnSelection_Abort != 0 {
seqnums = append(seqnums, SEQNUM_ABORT)
}
return TombstoneSeqnums_CN_Created
return seqnums
}

func GetTombstoneSchema(
Expand All @@ -144,33 +226,38 @@ func GetTombstoneSchema(
func GetTombstoneTypes(
pk types.Type, hidden HiddenColumnSelection,
) []types.Type {
hidden = normalizeTombstoneHiddenColumns(hidden)
var typs []types.Type
if hidden&HiddenColumnSelection_PhysicalAddr != 0 &&
hidden&HiddenColumnSelection_CommitTS != 0 {
return []types.Type{
typs = []types.Type{
RowidType,
pk,
TSType,
RowidType,
}
}
if hidden&HiddenColumnSelection_PhysicalAddr != 0 {
return []types.Type{
} else if hidden&HiddenColumnSelection_PhysicalAddr != 0 {
typs = []types.Type{
RowidType,
pk,
RowidType,
}
}
if hidden&HiddenColumnSelection_CommitTS != 0 {
return []types.Type{
} else if hidden&HiddenColumnSelection_CommitTS != 0 {
typs = []types.Type{
RowidType,
pk,
TSType,
}
} else {
typs = []types.Type{
RowidType,
pk,
}
}
return []types.Type{
RowidType,
pk,
if hidden&HiddenColumnSelection_Abort != 0 {
typs = append(typs, types.T_bool.ToType())
}
return typs
}

func MustGetPhysicalColumnPosition(seqnums []uint16, colTypes []types.Type) int {
Expand Down
65 changes: 65 additions & 0 deletions pkg/objectio/constructors.go
Original file line number Diff line number Diff line change
Expand Up @@ -594,6 +594,17 @@ func FilterCachedRowsByCommitTS(
data fscache.Data,
sels []int64,
snapshot types.TS,
) ([]int64, error) {
return FilterCachedRowsByCommitTSAndAbort(data, nil, sels, snapshot)
}

// FilterCachedRowsByCommitTSAndAbort removes rows newer than snapshot and rows
// marked aborted. A nil or const-null abort vector is the legacy object format.
func FilterCachedRowsByCommitTSAndAbort(
data fscache.Data,
abortData fscache.Data,
sels []int64,
snapshot types.TS,
) ([]int64, error) {
var commits vector.Vector
if err := bindCachedVectorForScope(&commits, data); err != nil {
Expand All @@ -603,6 +614,19 @@ func FilterCachedRowsByCommitTS(
if commits.GetType().Oid != types.T_TS || commits.IsConstNull() {
return nil, moerr.NewInvalidInputNoCtx("object commit-ts column is unavailable")
}
var aborts vector.Vector
hasAborts := abortData != nil
if hasAborts {
if err := bindCachedVectorForScope(&aborts, abortData); err != nil {
return nil, err
}
defer aborts.Free(nil)
if aborts.IsConstNull() {
hasAborts = false
} else if aborts.GetType().Oid != types.T_bool || aborts.Length() != commits.Length() {
return nil, moerr.NewInvalidInputNoCtx("object abort column is unavailable")
}
}

filtered := sels[:0]
for _, sel := range sels {
Expand All @@ -616,6 +640,14 @@ func FilterCachedRowsByCommitTS(
if commits.IsNull(uint64(sel)) {
return nil, moerr.NewInvalidInputNoCtxf("object commit-ts row %d is null", sel)
}
if hasAborts {
if aborts.IsNull(uint64(sel)) {
return nil, moerr.NewInvalidInputNoCtxf("object abort row %d is null", sel)
}
if vector.GetFixedAtNoTypeCheck[bool](&aborts, int(sel)) {
continue
}
}
commit := vector.GetFixedAtNoTypeCheck[types.TS](&commits, int(sel))
if !commit.GT(&snapshot) {
filtered = append(filtered, sel)
Expand All @@ -631,6 +663,18 @@ func AnyCachedTSInRange(
data fscache.Data,
sels []int64,
from, to types.TS,
) (matched bool, usable bool, err error) {
return AnyCachedTSInRangeWithAbort(data, nil, sels, from, to)
}

// AnyCachedTSInRangeWithAbort checks selected commit timestamps while ignoring
// rows marked aborted. A nil or const-null abort vector represents the legacy
// commitTS-only object format.
func AnyCachedTSInRangeWithAbort(
data fscache.Data,
abortData fscache.Data,
sels []int64,
from, to types.TS,
) (matched bool, usable bool, err error) {
var commits vector.Vector
if err = bindCachedVectorForScope(&commits, data); err != nil {
Expand All @@ -640,10 +684,31 @@ func AnyCachedTSInRange(
if commits.GetType().Oid != types.T_TS || commits.IsConstNull() {
return false, false, nil
}
var aborts vector.Vector
hasAborts := abortData != nil
if hasAborts {
if err = bindCachedVectorForScope(&aborts, abortData); err != nil {
return
}
defer aborts.Free(nil)
if aborts.IsConstNull() {
hasAborts = false
} else if aborts.GetType().Oid != types.T_bool || aborts.Length() != commits.Length() {
return false, false, nil
}
}
for _, sel := range sels {
if sel < 0 || sel >= int64(commits.Length()) || commits.IsNull(uint64(sel)) {
return false, false, nil
}
if hasAborts {
if aborts.IsNull(uint64(sel)) {
return false, false, nil
}
if vector.GetFixedAtNoTypeCheck[bool](&aborts, int(sel)) {
continue
}
}
commit := vector.GetFixedAtNoTypeCheck[types.TS](&commits, int(sel))
if commit.GT(&from) && commit.LE(&to) {
return true, true, nil
Expand Down
25 changes: 6 additions & 19 deletions pkg/objectio/funcs.go
Original file line number Diff line number Diff line change
Expand Up @@ -168,32 +168,19 @@ func ReadOneBlockWithMeta(

blkmeta := meta.GetBlockMeta(uint32(blk))
maxSeqnum := blkmeta.GetMaxSeqnum()
specialLayout := ResolveSpecialColumnLayout(blkmeta)
for i, seqnum := range seqnums {
// special columns
if seqnum >= SEQNUM_UPPER {
metaColCnt := blkmeta.GetMetaColumnCount()
switch seqnum {
case SEQNUM_COMMITTS:
if metaColCnt == 0 {
putFillHolder(i, 0)
continue
}
seqnum = metaColCnt - 1
case SEQNUM_ABORT:
panic("not support")
default:
var ok bool
if seqnum != SEQNUM_COMMITTS && seqnum != SEQNUM_ABORT {
panic(fmt.Sprintf("bad path to read special column %d", seqnum))
}
// Type alone is insufficient: the last user column may itself be
// T_TS. A hidden commit-TS column must sit beyond MaxSeqnum.
// If the last column is not commits, do not read it:
// 1. created by cn
// 2. old version tn nonappendable block
col := blkmeta.ColumnMeta(seqnum)
hasHiddenColumn := metaColCnt > maxSeqnum+1
if !hasHiddenColumn || col.DataType() != uint8(types.T_TS) {
seqnum, ok = specialLayout.Resolve(seqnum)
if !ok {
putFillHolder(i, seqnum)
} else {
col := blkmeta.ColumnMeta(seqnum)
ext := col.Location()
ioVec.Entries = append(ioVec.Entries, newColumnIOEntry(ext, factory))
}
Expand Down
Loading
Loading