Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
67f6669
Move appendable object creation outside txn
jiangxinmeng1 Jul 8, 2026
1184e0d
fix(tae): avoid updating aobject create ts on append commit
jiangxinmeng1 Jul 8, 2026
06e4754
fix(tae): preserve appendable stats in apply table data
jiangxinmeng1 Jul 8, 2026
eff17a6
Merge branch 'main' into move-aobject-create-outside-txn
mergify[bot] Jul 8, 2026
e732b32
fix(tae): update replay aobject create ts for delete entries
jiangxinmeng1 Jul 8, 2026
3e815e3
fix tae merge notifier race
jiangxinmeng1 Jul 14, 2026
dcdab9c
fix empty appendable object flush
jiangxinmeng1 Jul 14, 2026
157f999
set merge notifier only in write mode
jiangxinmeng1 Jul 14, 2026
80b1047
fix tae flush tombstone transfer scan
jiangxinmeng1 Jul 15, 2026
5df3fc2
fix tae tombstone range scan duplicate versions
jiangxinmeng1 Jul 15, 2026
bd08a99
add tae appendable object create coverage
jiangxinmeng1 Jul 16, 2026
0927bb6
add flush transfer range diagnostics
jiangxinmeng1 Jul 17, 2026
62c7397
adjust flush transfer diagnostics
jiangxinmeng1 Jul 17, 2026
e05bd7e
fix tae tombstone scan during active drop
jiangxinmeng1 Jul 20, 2026
968624f
fix tae empty appendable object replay
jiangxinmeng1 Jul 20, 2026
0372eae
fix tae dropped appendable object dedup
jiangxinmeng1 Jul 20, 2026
c15aecc
fix tae committed timestamp monotonicity
jiangxinmeng1 Jul 20, 2026
ff8dcd5
fix tae catalog publication and tombstone scan
jiangxinmeng1 Jul 21, 2026
c6c5409
fix tombstone range scan regression
jiangxinmeng1 Jul 22, 2026
4f242df
fix gofmt check
jiangxinmeng1 Jul 22, 2026
24cebe7
Merge branch 'main' into move-aobject-create-outside-txn
XuPeng-SH Jul 22, 2026
99c07ae
fix(tae): preserve object apply and mode switch lifecycle
jiangxinmeng1 Jul 23, 2026
72354ce
fix(tae): reject readonly txn admission cleanly
jiangxinmeng1 Jul 23, 2026
a887723
Merge branch 'main' into move-aobject-create-outside-txn
jiangxinmeng1 Jul 23, 2026
624a273
test(tae): close mode switch txn servers
jiangxinmeng1 Jul 23, 2026
97839bc
fix(tae): bound replay drain and snapshot dump
jiangxinmeng1 Jul 23, 2026
9874009
fix(tae): use moerr for scheduler stop
jiangxinmeng1 Jul 23, 2026
36dbad3
fix(tae): reject stopped scheduler generation IO
jiangxinmeng1 Jul 23, 2026
9da7616
Merge commit 'd8a464eb65ea33fd151ae185ae73cc87a4c68289' into move-aob…
jiangxinmeng1 Jul 24, 2026
84f0581
fix(tae): isolate merge scheduler generations
jiangxinmeng1 Jul 24, 2026
dcce800
fix(tae): preserve merge completion accounting
jiangxinmeng1 Jul 24, 2026
5a99ab3
Merge remote-tracking branch 'upstream/main' into move-aobject-create…
jiangxinmeng1 Jul 24, 2026
6b216d2
Merge remote-tracking branch 'upstream/main' into move-aobject-create…
jiangxinmeng1 Jul 24, 2026
91c4c5b
Merge branch 'main' into move-aobject-create-outside-txn
mergify[bot] Jul 24, 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
8 changes: 8 additions & 0 deletions pkg/objectio/object_stats.go
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,14 @@ func (des *ObjectStats) GetAppendable() bool {
return des[reservedOffset]&ObjectFlag_Appendable != 0
}

func SetObjectStatsAppendable(des *ObjectStats, appendable bool) {
if appendable {
des[reservedOffset] |= ObjectFlag_Appendable
} else {
des[reservedOffset] &^= ObjectFlag_Appendable
}
}

func (des *ObjectStats) GetSorted() bool {
return des[reservedOffset]&ObjectFlag_Sorted != 0
}
Expand Down
2 changes: 1 addition & 1 deletion pkg/vm/engine/disttae/change_handle.go
Original file line number Diff line number Diff line change
Expand Up @@ -434,7 +434,7 @@ func (h *PartitionChangesHandle) getNextChangeHandle(ctx context.Context) (end b
if err != nil {
return
}
h.currentChangeHandle, err = logtailreplay.NewChangesHandlerWithCheckpointEntries(
h.currentChangeHandle, err = logtailreplay.NewChangesHandlerWithCheckpointRangeRecovery(
ctx,
h.tbl.tableId,
h.tbl.proc.Load().GetService(),
Expand Down
32 changes: 32 additions & 0 deletions pkg/vm/engine/disttae/logtailreplay/change_handle.go
Original file line number Diff line number Diff line change
Expand Up @@ -1605,6 +1605,38 @@ func NewChangesHandlerWithCheckpointRange(
)
}

// NewChangesHandlerWithCheckpointRangeRecovery rebuilds CollectChanges(start,
// end) from checkpoint metadata using range-aware object selection while
// preserving CDC/checkpoint recovery merge semantics.
func NewChangesHandlerWithCheckpointRangeRecovery(
ctx context.Context,
tid uint64,
sid string,
checkpoints []*checkpoint.CheckpointEntry,
start, end types.TS,
skipDeletes bool,
maxRow uint32,
primarySeqnum int,
mp *mpool.MPool,
fs fileservice.FileService,
) (changeHandle *ChangeHandler, err error) {
return newChangesHandlerWithCheckpointEntries(
ctx,
tid,
sid,
checkpoints,
start,
end,
skipDeletes,
maxRow,
primarySeqnum,
mp,
fs,
checkpointObjectSelectionRange,
true,
)
}

// NewChangesHandlerWithPartitionStateRange rebuilds CollectChanges(start, end)
// from the partition state visible at the range end snapshot.
//
Expand Down
21 changes: 21 additions & 0 deletions pkg/vm/engine/tae/catalog/catalogreplay.go
Original file line number Diff line number Diff line change
Expand Up @@ -239,6 +239,27 @@ func (catalog *Catalog) onReplayUpdateObject(
}
if cmd.mvccNode.DeletedAt.Equal(&txnif.UncommitTS) {
cobj, err := rel.GetObjectByID(cmd.ID.ObjectID(), cmd.node.IsTombstone)
if moerr.IsMoErrCode(err, moerr.OkExpectedEOB) && cmd.mvccNode.BaseNode.ObjectStats.GetAppendable() {
rel.Lock()
cobj, err = rel.GetObjectByID(cmd.ID.ObjectID(), cmd.node.IsTombstone)
if moerr.IsMoErrCode(err, moerr.OkExpectedEOB) {
var factory ObjectDataFactory
if catalog.DataFactory != nil {
factory = catalog.DataFactory.MakeObjectFactory()
}
cobj = NewCommittedObjectEntry(
rel,
cmd.mvccNode.CreatedAt,
cmd.mvccNode.BaseNode.ObjectStats,
factory,
cmd.node.IsTombstone,
)
cobj.ObjectNode = *cmd.node
rel.AddEntryLocked(cobj)
err = nil
}
rel.Unlock()
}
if err != nil {
panic(fmt.Sprintf("obj %v not existed, table:\n%v", cmd.ID.String(), rel.StringWithLevel(3)))
}
Expand Down
28 changes: 28 additions & 0 deletions pkg/vm/engine/tae/catalog/object.go
Original file line number Diff line number Diff line change
Expand Up @@ -386,6 +386,34 @@ func NewObjectEntry(
return e
}

func NewCommittedObjectEntry(
table *TableEntry,
ts types.TS,
stats objectio.ObjectStats,
dataFactory ObjectDataFactory,
isTombstone bool,
) *ObjectEntry {
e := &ObjectEntry{
table: table,
ObjectNode: ObjectNode{
SortHint: table.GetDB().catalog.NextObject(),
IsTombstone: isTombstone,
},
EntryMVCCNode: EntryMVCCNode{
CreatedAt: ts,
},
CreateNode: txnbase.NewTxnMVCCNodeWithTS(ts),
ObjectState: ObjectState_Create_ApplyCommit,
ObjectMVCCNode: ObjectMVCCNode{
ObjectStats: stats,
},
}
if dataFactory != nil {
e.objData = dataFactory(e)
}
return e
}

func NewReplayObjectEntry() *ObjectEntry {
e := &ObjectEntry{}
return e
Expand Down
45 changes: 45 additions & 0 deletions pkg/vm/engine/tae/catalog/object_list.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import (
"github.com/matrixorigin/matrixone/pkg/objectio"
"github.com/matrixorigin/matrixone/pkg/vm/engine/tae/common"
"github.com/matrixorigin/matrixone/pkg/vm/engine/tae/iface/txnif"
"github.com/matrixorigin/matrixone/pkg/vm/engine/tae/txn/txnbase"
"github.com/tidwall/btree"
"go.uber.org/zap"
)
Expand Down Expand Up @@ -314,6 +315,50 @@ func (l *ObjectList) DeleteAllEntries(id *objectio.ObjectId) error {
return nil
}

func (l *ObjectList) UpdateCreateTS(id *objectio.ObjectId, ts types.TS) (*ObjectEntry, error) {
l.Lock()
defer l.Unlock()
oldTS, ok := l.maxTs_objectID[*id]
if !ok {
return nil, moerr.GetOkExpectedEOB()
}
oldTree := l.tree.Load()
newTree := oldTree.Copy()
nodes := l.getNodesSnap(newTree, oldTS, id, true)
if len(nodes) == 0 {
return nil, moerr.GetOkExpectedEOB()
}
oldNode := nodes[0]
newNode := oldNode.Clone()
if oldNode.IsDEntry() {
newPrev := oldNode.prevVersion.Clone()
newPrev.CreatedAt = ts
newPrev.CreateNode = txnbase.NewTxnMVCCNodeWithTS(ts)
newPrev.nextVersion = newNode
newNode.CreatedAt = ts
newNode.CreateNode = txnbase.NewTxnMVCCNodeWithTS(ts)
newNode.prevVersion = newPrev
newTree.Delete(oldNode)
newTree.Delete(oldNode.prevVersion)
newTree.Set(newNode)
newTree.Set(newPrev)
} else {
newNode.CreatedAt = ts
newNode.CreateNode = txnbase.NewTxnMVCCNodeWithTS(ts)
newTree.Delete(oldNode)
newTree.Set(newNode)
}
maxTS := newNode.CreatedAt
if maxTS.LT(&newNode.DeletedAt) {
maxTS = newNode.DeletedAt
}
l.maxTs_objectID[*id] = maxTS
if !l.tree.CompareAndSwap(oldTree, newTree) {
panic("concurrent mutation")
}
return newNode, nil
}

// WaitUntilCommitted waits for entries in txn active zone with prepareTS > ts to move to committed zone.
// As CreateObject will be called in txn queue, when WaitUntilCommitted returns, all creating objects in txn active zone are invisible to ts, because they are created after ts.
func (l *ObjectList) WaitUntilCommitted(ts types.TS) {
Expand Down
40 changes: 36 additions & 4 deletions pkg/vm/engine/tae/catalog/table.go
Original file line number Diff line number Diff line change
Expand Up @@ -268,6 +268,17 @@ func (entry *TableEntry) getTombstoneObjectByID(id *types.Objectid) (obj *Object
return entry.tombstoneObjects.GetObjectByID(id)
}

func (entry *TableEntry) UpdateObjectCreateTS(id *types.Objectid, isTombstone bool, ts types.TS) error {
obj, err := entry.getObjectList(isTombstone).UpdateCreateTS(id, ts)
if err != nil {
return err
}
if obj.GetObjectData() != nil {
obj.GetObjectData().UpdateMeta(obj)
}
return nil
}

func (entry *TableEntry) MakeTombstoneObjectIt() btree.IterG[*ObjectEntry] {
return entry.tombstoneObjects.tree.Load().Iter()
}
Expand Down Expand Up @@ -346,8 +357,27 @@ func (entry *TableEntry) CreateObject(
entry.Lock()
defer entry.Unlock()
created = NewObjectEntry(entry, txn, *opts.Stats, dataFactory, opts.IsTombstone)
if entry.GetCatalog().mergeNotifier != nil && !opts.Stats.GetAppendable() {
entry.GetCatalog().mergeNotifier.OnCreateNonAppendObject(ToMergeTable(entry))
if !opts.Stats.GetAppendable() {
if notifier := entry.GetCatalog().getMergeNotifier(); notifier != nil {
notifier.OnCreateNonAppendObject(ToMergeTable(entry))
}
}
entry.AddEntryLocked(created)
return
}

func (entry *TableEntry) CreateCommittedObject(
ts types.TS,
opts *objectio.CreateObjOpt,
dataFactory ObjectDataFactory,
) (created *ObjectEntry, err error) {
entry.Lock()
defer entry.Unlock()
created = NewCommittedObjectEntry(entry, ts, *opts.Stats, dataFactory, opts.IsTombstone)
if !opts.Stats.GetAppendable() {
if notifier := entry.GetCatalog().getMergeNotifier(); notifier != nil {
notifier.OnCreateNonAppendObject(ToMergeTable(entry))
}
}
entry.AddEntryLocked(created)
return
Expand Down Expand Up @@ -815,8 +845,10 @@ func (entry *TableEntry) ApplyCommit(id string) (err error) {
entry.TableNode.schema.Store(schema)

// create table commit
if lastestNode.DeletedAt.IsEmpty() && entry.GetCatalog().mergeNotifier != nil {
entry.GetCatalog().mergeNotifier.OnCreateTableCommit(ToMergeTable(entry))
if lastestNode.DeletedAt.IsEmpty() {
if notifier := entry.GetCatalog().getMergeNotifier(); notifier != nil {
notifier.OnCreateTableCommit(ToMergeTable(entry))
}
}

return
Expand Down
4 changes: 4 additions & 0 deletions pkg/vm/engine/tae/catalog/tableForMerge.go
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,10 @@ func (catalog *Catalog) SetMergeNotifier(notifier MergeNotifierOnCatalog) {
catalog.mergeNotifier = notifier
}

func (catalog *Catalog) getMergeNotifier() MergeNotifierOnCatalog {
return catalog.mergeNotifier
}

func (catalog *Catalog) InitSource() iter.Seq[MergeTable] {
return func(yield func(MergeTable) bool) {
p := new(LoopProcessor)
Expand Down
41 changes: 41 additions & 0 deletions pkg/vm/engine/tae/catalog/table_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,47 @@ func TestObjectList(t *testing.T) {
t.Log(ll.getNodes(entry1.ID(), false))
}

func TestObjectListUpdateCreateTSWithDeleteEntry(t *testing.T) {
ll := NewObjectList(false)
nobjid := objectio.NewObjectid()
createTS := types.BuildTS(10, 0)
deleteTS := types.BuildTS(20, 0)
updatedCreateTS := types.BuildTS(5, 0)
createEntry := &ObjectEntry{
ObjectNode: ObjectNode{SortHint: 1},
EntryMVCCNode: EntryMVCCNode{
CreatedAt: createTS,
},
ObjectMVCCNode: ObjectMVCCNode{ObjectStats: *objectio.NewObjectStatsWithObjectID(&nobjid, true, false, false)},
CreateNode: txnbase.NewTxnMVCCNodeWithTS(createTS),
ObjectState: ObjectState_Create_ApplyCommit,
}
deleteEntry := createEntry.Clone()
deleteEntry.DeletedAt = deleteTS
deleteEntry.DeleteNode = txnbase.NewTxnMVCCNodeWithTS(deleteTS)
deleteEntry.ObjectState = ObjectState_Delete_ApplyCommit
updatedCreateEntry := createEntry.Clone()
updatedCreateEntry.nextVersion = deleteEntry
deleteEntry.prevVersion = updatedCreateEntry

ll.modify(nil, deleteEntry, updatedCreateEntry)
updated, err := ll.UpdateCreateTS(createEntry.ID(), updatedCreateTS)
require.NoError(t, err)
require.True(t, updated.IsDEntry())

nodes := ll.GetAllNodes(createEntry.ID())
require.Len(t, nodes, 2)
require.Equal(t, updatedCreateTS, nodes[0].CreatedAt)
require.Equal(t, updatedCreateTS, nodes[0].CreateNode.GetPrepare())
require.Equal(t, updatedCreateTS, nodes[1].CreatedAt)
require.Equal(t, updatedCreateTS, nodes[1].CreateNode.GetPrepare())
require.Same(t, nodes[0].prevVersion, nodes[1])
require.Same(t, nodes[1].nextVersion, nodes[0])
require.Equal(t, 2, ll.tree.Load().Len())
require.NoError(t, ll.DeleteAllEntries(createEntry.ID()))
require.Zero(t, ll.tree.Load().Len())
}

func TestGetSoftdeleteObjects(t *testing.T) {
db := MockDBEntryWithAccInfo(0, 0)
tbl := MockTableEntryWithDB(db, 1)
Expand Down
Loading
Loading