diff --git a/pkg/objectio/object_stats.go b/pkg/objectio/object_stats.go index 04fed88ec9f6e..df22748280756 100644 --- a/pkg/objectio/object_stats.go +++ b/pkg/objectio/object_stats.go @@ -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 } diff --git a/pkg/vm/engine/disttae/change_handle.go b/pkg/vm/engine/disttae/change_handle.go index 3f90f213adc2b..238de56d44812 100644 --- a/pkg/vm/engine/disttae/change_handle.go +++ b/pkg/vm/engine/disttae/change_handle.go @@ -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(), diff --git a/pkg/vm/engine/disttae/logtailreplay/change_handle.go b/pkg/vm/engine/disttae/logtailreplay/change_handle.go index 9706cf1a52b03..9eb69f9bbdb99 100755 --- a/pkg/vm/engine/disttae/logtailreplay/change_handle.go +++ b/pkg/vm/engine/disttae/logtailreplay/change_handle.go @@ -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. // diff --git a/pkg/vm/engine/tae/catalog/catalogreplay.go b/pkg/vm/engine/tae/catalog/catalogreplay.go index fe7f7308d5602..022404d986f23 100644 --- a/pkg/vm/engine/tae/catalog/catalogreplay.go +++ b/pkg/vm/engine/tae/catalog/catalogreplay.go @@ -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))) } diff --git a/pkg/vm/engine/tae/catalog/object.go b/pkg/vm/engine/tae/catalog/object.go index be8e033ffe114..18563c124e544 100644 --- a/pkg/vm/engine/tae/catalog/object.go +++ b/pkg/vm/engine/tae/catalog/object.go @@ -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 diff --git a/pkg/vm/engine/tae/catalog/object_list.go b/pkg/vm/engine/tae/catalog/object_list.go index 62f1d55e3df9c..7c0a4ba938248 100644 --- a/pkg/vm/engine/tae/catalog/object_list.go +++ b/pkg/vm/engine/tae/catalog/object_list.go @@ -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" ) @@ -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) { diff --git a/pkg/vm/engine/tae/catalog/table.go b/pkg/vm/engine/tae/catalog/table.go index 2aabc9cbd8c73..76bf094b70214 100644 --- a/pkg/vm/engine/tae/catalog/table.go +++ b/pkg/vm/engine/tae/catalog/table.go @@ -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() } @@ -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 @@ -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 diff --git a/pkg/vm/engine/tae/catalog/tableForMerge.go b/pkg/vm/engine/tae/catalog/tableForMerge.go index 67ca0165560a2..4f7b5459056c1 100644 --- a/pkg/vm/engine/tae/catalog/tableForMerge.go +++ b/pkg/vm/engine/tae/catalog/tableForMerge.go @@ -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) diff --git a/pkg/vm/engine/tae/catalog/table_test.go b/pkg/vm/engine/tae/catalog/table_test.go index 86a9fee1cce9e..f79550144c12b 100644 --- a/pkg/vm/engine/tae/catalog/table_test.go +++ b/pkg/vm/engine/tae/catalog/table_test.go @@ -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) diff --git a/pkg/vm/engine/tae/db/controller.go b/pkg/vm/engine/tae/db/controller.go index df0e99b77417a..db398b195f1ec 100644 --- a/pkg/vm/engine/tae/db/controller.go +++ b/pkg/vm/engine/tae/db/controller.go @@ -16,6 +16,7 @@ package db import ( "context" + "errors" "fmt" "sync" "sync/atomic" @@ -273,44 +274,53 @@ func (c *Controller) handleToReplayCmd(cmd *controlCmd) { ) defer func() { - err2 := err - if err2 != nil { - err = rollbackSteps.Apply("DB-SwitchToReplay-Rollback", true, 1) - } - if err2 != nil { + switchErr := err + var rollbackErr error + if switchErr != nil { + rollbackErr = rollbackSteps.Apply("DB-SwitchToReplay-Rollback", true, 1) + err = switchErr logger = logutil.Error } - if err != nil { + if rollbackErr != nil { + err = errors.Join(switchErr, rollbackErr) logger = logutil.Fatal } logger( "DB-SwitchToReplay-Done", zap.String("cmd", cmd.String()), zap.Duration("duration", time.Since(start)), - zap.Any("rollback-error", err2), + zap.Error(rollbackErr), zap.Error(err), ) cmd.setError(err) }() - // 1. stop the merge scheduler - c.db.MergeScheduler.Stop() - rollbackSteps.Add("stop merge scheduler", func() error { - c.db.MergeScheduler.Start() - return nil - }) - // TODO + // 1. switch the checkpoint|diskcleaner to replay mode - // 2. switch the checkpoint|diskcleaner to replay mode - - // 2.1 remove GC disk cron job. no new GC job will be issued from now on + // 1.1 remove GC disk cron job. no new GC job will be issued from now on + gcDiskRunning := c.db.CronJobs.GetJob(CronJobs_Name_GCDisk) != nil + gcCheckpointRunning := c.db.CronJobs.GetJob(CronJobs_Name_GCCheckpoint) != nil RemoveCronJob(c.db, CronJobs_Name_GCDisk) RemoveCronJob(c.db, CronJobs_Name_GCCheckpoint) + rollbackSteps.Add("restore write cron jobs", func() error { + if gcDiskRunning { + if err := AddCronJob(c.db, CronJobs_Name_GCDisk, true); err != nil { + return err + } + } + if gcCheckpointRunning { + return AddCronJob(c.db, CronJobs_Name_GCCheckpoint, true) + } + return nil + }) // RemoveCronJob(c.db, CronJobs_Name_Scanner) if err = c.db.DiskCleaner.SwitchToReplayMode(ctx); err != nil { // Rollback return } + rollbackSteps.Add("switch disk cleaner to write mode", func() error { + return c.db.DiskCleaner.SwitchToWriteMode(context.Background()) + }) // 2.x TODO: checkpoint runner flushCfg := c.db.BGFlusher.GetCfg() c.db.BGFlusher.Stop() @@ -336,12 +346,33 @@ func (c *Controller) handleToReplayCmd(cmd *controlCmd) { // 5. freeze the write requests consumer // TODO - // 6. switch the txn mode to readonly mode - if err = c.db.TxnMgr.SwitchToReadonly(cmd.ctx); err != nil { + // 6. Prevent new transactions while the merge scheduler is still draining + // catalog notifications from transactions that were already in flight. + if err = c.db.TxnMgr.SwitchToReadonly(ctx); err != nil { c.db.TxnMgr.ToWriteMode() - // TODO: recover the previous state return } + rollbackSteps.Add("switch txn manager to write mode", func() error { + c.db.TxnMgr.ToWriteMode() + return nil + }) + + // No transaction can publish another catalog notification after WaitEmpty. + // Detach the producer, process every event already queued, and only then stop + // the consumer. + c.db.Catalog.SetMergeNotifier(nil) + rollbackSteps.Add("restore merge scheduler notifier", func() error { + c.db.Catalog.SetMergeNotifier(c.db.MergeScheduler) + return nil + }) + if _, err = c.db.MergeScheduler.Query(ctx, nil); err != nil { + return + } + c.db.MergeScheduler.Stop() + rollbackSteps.Add("start merge scheduler", func() error { + c.db.MergeScheduler.Start() + return nil + }) // 7. wait the logtail push queue to be flushed // TODO @@ -428,17 +459,19 @@ func (c *Controller) handleToWriteCmd(cmd *controlCmd) { // 3. switch the txnmgr to write mode c.db.TxnMgr.ToWriteMode() - // 4. unfreeze the write requests - if err = c.db.TxnServer.SwitchTxnHandleStateTo(rpc2.TxnLocalHandle); err != nil { - return - } - + c.db.Catalog.SetMergeNotifier(c.db.MergeScheduler) c.db.MergeScheduler.Start() rollbackSteps.Add("stop merge scheduler", func() error { c.db.MergeScheduler.Stop() + c.db.Catalog.SetMergeNotifier(nil) return nil }) + // 4. unfreeze the write requests + if err = c.db.TxnServer.SwitchTxnHandleStateTo(rpc2.TxnLocalHandle); err != nil { + return + } + // 5. start merge scheduler|checkpoint|diskcleaner // 5.1 TODO: start the merger|checkpoint|flusher c.db.BGFlusher.Restart() // TODO: Restart with new config @@ -779,11 +812,15 @@ func (c *Controller) AssembleDB(ctx context.Context) (err error) { merge.NewTNMergeExecutor(db.Runtime), merge.NewStdClock(), ) - db.MergeScheduler.Start() - rollbackSteps.Add("stop merge scheduler", func() error { - db.MergeScheduler.Stop() - return nil - }) + if db.IsWriteMode() { + db.Catalog.SetMergeNotifier(db.MergeScheduler) + db.MergeScheduler.Start() + rollbackSteps.Add("stop merge scheduler", func() error { + db.MergeScheduler.Stop() + db.Catalog.SetMergeNotifier(nil) + return nil + }) + } // start flusher and disk cleaner db.BGFlusher.Start() diff --git a/pkg/vm/engine/tae/db/merge/scheduler.go b/pkg/vm/engine/tae/db/merge/scheduler.go index f29edc7ee9b6d..890ed6c3fc75c 100644 --- a/pkg/vm/engine/tae/db/merge/scheduler.go +++ b/pkg/vm/engine/tae/db/merge/scheduler.go @@ -21,9 +21,11 @@ import ( "fmt" "math/rand" "slices" + "sync" "sync/atomic" "time" + "github.com/matrixorigin/matrixone/pkg/common/moerr" "github.com/matrixorigin/matrixone/pkg/common/rscthrottler" "github.com/matrixorigin/matrixone/pkg/container/batch" "github.com/matrixorigin/matrixone/pkg/logutil" @@ -40,6 +42,8 @@ const ( objectOpsTriggerThreshold = 5 ) +var ErrMergeSchedulerStopped = moerr.NewInternalErrorNoCtx("merge scheduler stopped") + type mergeTask struct { objs []*objectio.ObjectStats kind taskHostKind @@ -51,6 +55,18 @@ type mergeTask struct { doneCB *taskObserver } +type mergeSchedulerGeneration struct { + stopCh chan struct{} + ioChan chan *MMsg +} + +func newMergeSchedulerGeneration() *mergeSchedulerGeneration { + return &mergeSchedulerGeneration{ + stopCh: make(chan struct{}), + ioChan: make(chan *MMsg, 256), + } +} + func (r *mergeTask) String() string { return fmt.Sprintf( "mergeTask{isTombstone: %v, level: %d, note: %s, objs: %v, oSize: %s}", @@ -65,13 +81,13 @@ type MergeScheduler struct { supps map[uint64]*todoSupporter // control flow - allPaused bool - stopCh atomic.Pointer[chan struct{}] - stopRecv chan struct{} - stopped atomic.Bool + allPaused bool + generation atomic.Pointer[mergeSchedulerGeneration] + stopRecv chan struct{} + stopped atomic.Bool - msgChan chan *MMsg - ioChan chan *MMsg + msgChan chan *MMsg + bootstrapMsg *MMsg pad *launchPad defaultTrigger *MMsgTaskTrigger @@ -97,7 +113,6 @@ func NewMergeScheduler( stopRecv: make(chan struct{}, 1), msgChan: make(chan *MMsg, 4096), - ioChan: make(chan *MMsg, 256), pad: newLaunchPad(clock), defaultTrigger: DefaultTrigger.Clone(), @@ -118,14 +133,13 @@ func NewMergeScheduler( sched.handleAddTable(table) } if fn := cata.GetMergeSettingsBatchFn(); fn != nil { - sched.ioChan <- &MMsg{ + sched.bootstrapMsg = &MMsg{ Kind: MMsgKindConfigBootstrap, Value: MMsgConfigBootstrap{ ReadSettingsBatch: fn, }, } } - cata.SetMergeNotifier(sched) return sched @@ -137,9 +151,9 @@ func (a *MergeScheduler) PatchTestRscController(rc rscthrottler.RSCThrottler) { func (a *MergeScheduler) Stop() { if a.stopped.CompareAndSwap(false, true) { - ch := a.stopCh.Load() - if ch != nil { - close(*ch) + generation := a.generation.Load() + if generation != nil { + close(generation.stopCh) } <-a.stopRecv } @@ -147,10 +161,17 @@ func (a *MergeScheduler) Stop() { func (a *MergeScheduler) Start() { if a.stopped.CompareAndSwap(true, false) { - ch := make(chan struct{}) - a.stopCh.Store(&ch) - go a.handleMainLoop() - go a.handleIOLoop() + generation := newMergeSchedulerGeneration() + a.generation.Store(generation) + if a.bootstrapMsg != nil { + generation.ioChan <- &MMsg{ + Kind: a.bootstrapMsg.Kind, + Value: a.bootstrapMsg.Value, + generation: generation, + } + } + go a.handleMainLoop(generation) + go a.handleIOLoop(generation) } } @@ -188,18 +209,49 @@ func (a *MergeScheduler) OnMergeDone(table catalog.MergeTable, esize int) { } func (a *MergeScheduler) taskObserverFactory( - t catalog.MergeTable, + supp *todoSupporter, size int, + rc rscthrottler.RSCThrottler, ) *taskObserver { - return &taskObserver{f: func() { a.OnMergeDone(t, size) }} + return &taskObserver{f: func() { + supp.DoneTask() + rc.Release(int64(size)) + }} } type taskObserver struct { - f func() + mu sync.Mutex + admitted bool + completed bool + f func() } func (o *taskObserver) OnExecDone(_ any) { - o.f() + o.mu.Lock() + if o.completed { + o.mu.Unlock() + return + } + o.completed = true + run := o.admitted + o.mu.Unlock() + if run { + o.f() + } +} + +func (o *taskObserver) Admit() { + o.mu.Lock() + if o.admitted { + o.mu.Unlock() + return + } + o.admitted = true + run := o.completed + o.mu.Unlock() + if run { + o.f() + } } func (a *MergeScheduler) CNActiveObjectsString() string { return "" } @@ -460,8 +512,9 @@ func (b *MMsgTaskTrigger) Merge(o *MMsgTaskTrigger) *MMsgTaskTrigger { } type MMsg struct { - Kind MMsgKind - Value any + Kind MMsgKind + Value any + generation *mergeSchedulerGeneration } type todoItem struct { @@ -470,13 +523,77 @@ type todoItem struct { table catalog.MergeTable } -func (a *MergeScheduler) Query(table catalog.MergeTable) *QueryAnswer { - answer := make(chan *QueryAnswer) - a.msgChan <- &MMsg{ - Kind: MMsgKindQuery, - Value: MMsgQuery{Table: table, Answer: answer}, +func (a *MergeScheduler) Query( + ctx context.Context, + table catalog.MergeTable, +) (*QueryAnswer, error) { + generation := a.generation.Load() + if a.stopped.Load() || generation == nil { + return nil, ErrMergeSchedulerStopped + } + answer := make(chan *QueryAnswer, 1) + msg := &MMsg{ + Kind: MMsgKindQuery, + Value: MMsgQuery{Table: table, Answer: answer}, + generation: generation, + } + select { + case a.msgChan <- msg: + case <-ctx.Done(): + return nil, ctx.Err() + case <-generation.stopCh: + return nil, ErrMergeSchedulerStopped + } + select { + case answer := <-answer: + return answer, nil + case <-ctx.Done(): + return nil, ctx.Err() + case <-generation.stopCh: + return nil, ErrMergeSchedulerStopped + } +} + +func (a *MergeScheduler) sendIOForGeneration( + generation *mergeSchedulerGeneration, + msg *MMsg, +) bool { + if generation == nil { + return false + } + msg.generation = generation + select { + case <-generation.stopCh: + return false + default: + } + select { + case generation.ioChan <- msg: + return true + case <-generation.stopCh: + return false + } +} + +func (a *MergeScheduler) sendMsgForGeneration( + generation *mergeSchedulerGeneration, + msg *MMsg, +) bool { + if generation == nil { + return false + } + msg.generation = generation + select { + case <-generation.stopCh: + return false + default: + } + select { + case a.msgChan <- msg: + return true + case <-generation.stopCh: + return false } - return <-answer } func (a *MergeScheduler) PauseAll() { @@ -583,7 +700,7 @@ func (pq *todoPQ) Update(item *todoItem, ready time.Time) { } type todoSupporter struct { - mergingTaskCnt int + mergingTaskCnt atomic.Int64 vaccumTrigCount int objectOperations int totalDataMergeCnt int @@ -601,17 +718,29 @@ type todoSupporter struct { } func (m *todoSupporter) DoneTask() { - m.mergingTaskCnt-- - if m.mergingTaskCnt < 0 { - logutil.Error("MergeExecutorEvent", - zap.String("event", "mergingTaskCnt < 0"), - zap.String("table", m.todo.table.GetNameDesc()), - ) - m.mergingTaskCnt = 0 + for { + count := m.mergingTaskCnt.Load() + if count <= 0 { + logutil.Error("MergeExecutorEvent", + zap.String("event", "mergingTaskCnt <= 0"), + zap.String("table", m.todo.table.GetNameDesc()), + ) + return + } + if m.mergingTaskCnt.CompareAndSwap(count, count-1) { + return + } } } -func (a *MergeScheduler) ioVacuumCheck(msg MMsgVacuumCheck) { +func (m *todoSupporter) AddTask() { + m.mergingTaskCnt.Add(1) +} + +func (a *MergeScheduler) ioVacuumCheck( + generation *mergeSchedulerGeneration, + msg MMsgVacuumCheck, +) { stats, err := CalculateVacuumStats(context.Background(), msg.Table, msg.opts, @@ -628,19 +757,24 @@ func (a *MergeScheduler) ioVacuumCheck(msg MMsgVacuumCheck) { compactTasks := GatherCompactTasks(context.Background(), stats) if len(compactTasks) > 0 { - a.SendTrigger( - NewMMsgTaskTrigger(msg.Table). + if !a.sendMsgForGeneration(generation, &MMsg{ + Kind: MMsgKindTrigger, + Value: NewMMsgTaskTrigger(msg.Table). WithAssignedTasks(compactTasks), - ) + }) { + return + } // if the compact tasks is equal to the hollow top k, // it means the table is full of hollow objects, // so we need to trigger the vacuum check if len(compactTasks) == msg.opts.HollowTopK { f := func() { - a.SendTrigger( - NewMMsgTaskTrigger(msg.Table).WithVacuumCheck(msg.opts), - ) + a.sendMsgForGeneration(generation, &MMsg{ + Kind: MMsgKindTrigger, + Value: NewMMsgTaskTrigger(msg.Table). + WithVacuumCheck(msg.opts), + }) } a.clock.AfterFunc(time.Second*10, f) } @@ -648,10 +782,13 @@ func (a *MergeScheduler) ioVacuumCheck(msg MMsgVacuumCheck) { if stats.DelVacuumPercent > 0.5 { oneshotOpts := DefaultTombstoneOpts.Clone().WithOneShot(true) - a.SendTrigger( - NewMMsgTaskTrigger(msg.Table). + if !a.sendMsgForGeneration(generation, &MMsg{ + Kind: MMsgKindTrigger, + Value: NewMMsgTaskTrigger(msg.Table). WithTombstone(oneshotOpts), - ) + }) { + return + } } logutil.Info( @@ -665,7 +802,10 @@ func (a *MergeScheduler) ioVacuumCheck(msg MMsgVacuumCheck) { ) } -func (a *MergeScheduler) ioConfigBootstrap(msg MMsgConfigBootstrap) { +func (a *MergeScheduler) ioConfigBootstrap( + generation *mergeSchedulerGeneration, + msg MMsgConfigBootstrap, +) { bat, release := msg.ReadSettingsBatch() if bat == nil { logutil.Error( @@ -676,7 +816,20 @@ func (a *MergeScheduler) ioConfigBootstrap(msg MMsgConfigBootstrap) { defer release() count := 0 DecodeMergeSettingsBatchAnd(bat, func(tid uint64, setting *MergeSettings) { - a.SendConfig(tid, setting) + var trigger *MMsgTaskTrigger + var err error + if setting != nil { + trigger, err = setting.ToMMsgTaskTrigger() + if err != nil { + return + } + } + if !a.sendMsgForGeneration(generation, &MMsg{ + Kind: MMsgKindConfig, + Value: MMsgConfig{ID: tid, Trigger: trigger}, + }) { + return + } count++ }) @@ -687,24 +840,38 @@ func (a *MergeScheduler) ioConfigBootstrap(msg MMsgConfigBootstrap) { ) } -func (a *MergeScheduler) handleIOLoop() { - stopCh := *a.stopCh.Load() +func (a *MergeScheduler) handleIOLoop(generation *mergeSchedulerGeneration) { for { select { - case <-stopCh: + case <-generation.stopCh: return - case msg := <-a.ioChan: + default: + } + select { + case <-generation.stopCh: + return + case msg := <-generation.ioChan: + select { + case <-generation.stopCh: + return + default: + } + if msg.generation != generation { + continue + } switch msg.Kind { case MMsgKindVacuumCheck: - a.ioVacuumCheck(msg.Value.(MMsgVacuumCheck)) + a.ioVacuumCheck(generation, msg.Value.(MMsgVacuumCheck)) case MMsgKindConfigBootstrap: - a.ioConfigBootstrap(msg.Value.(MMsgConfigBootstrap)) + a.ioConfigBootstrap(generation, msg.Value.(MMsgConfigBootstrap)) } } } } -func (a *MergeScheduler) fallbackSchedVacuumCheck() { +func (a *MergeScheduler) fallbackSchedVacuumCheck( + generation *mergeSchedulerGeneration, +) { for _, supp := range a.supps { size := 0 for stat := range supp.todo.table.IterTombstoneItem() { @@ -712,19 +879,21 @@ func (a *MergeScheduler) fallbackSchedVacuumCheck() { } if size > 2*common.DefaultMaxOsizeObjBytes { a.clock.AfterFunc(time.Duration(rand.Intn(10))*time.Minute, func() { - a.ioChan <- &MMsg{ + a.sendIOForGeneration(generation, &MMsg{ Kind: MMsgKindVacuumCheck, Value: MMsgVacuumCheck{ Table: supp.todo.table, opts: DefaultVacuumOpts, }, - } + }) }) } } } -func (a *MergeScheduler) handleMainLoop() { +func (a *MergeScheduler) handleMainLoop( + generation *mergeSchedulerGeneration, +) { var nextReadyAtTimer = a.clock.NewTimer(time.Hour * 24) never := make(<-chan time.Time) @@ -733,9 +902,7 @@ func (a *MergeScheduler) handleMainLoop() { vacuumCheckTicker := a.clock.NewTicker(time.Hour * 1) - stopCh := *a.stopCh.Load() - - a.fallbackSchedVacuumCheck() + a.fallbackSchedVacuumCheck(generation) for { @@ -770,7 +937,7 @@ func (a *MergeScheduler) handleMainLoop() { } // DO NOT pop the task from the priority queue, // because the task may be updated - a.doSched(todo) + a.doSched(generation, todo) } // set the timer for the next task @@ -787,7 +954,7 @@ func (a *MergeScheduler) handleMainLoop() { } select { - case <-stopCh: + case <-generation.stopCh: // stop the loop heartbeat.Stop() a.stopRecv <- struct{}{} @@ -798,14 +965,14 @@ func (a *MergeScheduler) handleMainLoop() { a.rc.Refresh() // continue the loop case <-vacuumCheckTicker.Chan(): - a.fallbackSchedVacuumCheck() + a.fallbackSchedVacuumCheck(generation) case msg := <-a.msgChan: - a.dispatchMsg(msg) + a.dispatchMsg(generation, msg) drained := false for !drained { select { case msg := <-a.msgChan: - a.dispatchMsg(msg) + a.dispatchMsg(generation, msg) default: drained = true } @@ -816,14 +983,20 @@ func (a *MergeScheduler) handleMainLoop() { // region: handle msg -func (a *MergeScheduler) dispatchMsg(msg *MMsg) { +func (a *MergeScheduler) dispatchMsg( + generation *mergeSchedulerGeneration, + msg *MMsg, +) { + if msg.generation != nil && msg.generation != generation { + return + } switch msg.Kind { case MMsgKindSwitch: a.handleSwitch(msg.Value.(MMsgSwitch)) case MMsgKindQuery: a.handleQuery(msg.Value.(MMsgQuery)) case MMsgKindTrigger: - a.handleTaskTrigger(msg.Value.(*MMsgTaskTrigger)) + a.handleTaskTrigger(generation, msg.Value.(*MMsgTaskTrigger)) case MMsgKindConfig: a.handleConfig(msg.Value.(MMsgConfig)) case MMsgKindTableChange: @@ -838,7 +1011,10 @@ func (a *MergeScheduler) dispatchMsg(msg *MMsg) { } } -func (a *MergeScheduler) handleTaskTrigger(msg *MMsgTaskTrigger) { +func (a *MergeScheduler) handleTaskTrigger( + generation *mergeSchedulerGeneration, + msg *MMsgTaskTrigger, +) { supp := a.supps[msg.table.ID()] if supp == nil { // this table has been dropped and removed from the priority queue @@ -846,12 +1022,14 @@ func (a *MergeScheduler) handleTaskTrigger(msg *MMsgTaskTrigger) { } if msg.vacuum != nil { - a.ioChan <- &MMsg{ + if !a.sendIOForGeneration(generation, &MMsg{ Kind: MMsgKindVacuumCheck, Value: MMsgVacuumCheck{ Table: msg.table, opts: msg.vacuum, }, + }) { + return } supp.lastVacuumCheckTime = a.clock.Now() } @@ -939,7 +1117,7 @@ func (a *MergeScheduler) handleQuery(msg MMsgQuery) { answer.NextCheckDue = a.clock.Until(supp.todo.readyAt) answer.DataMergeCnt = supp.totalDataMergeCnt answer.TombstoneMergeCnt = supp.totalTombstoneMergeCnt - answer.PendingMergeCnt = supp.mergingTaskCnt + answer.PendingMergeCnt = int(supp.mergingTaskCnt.Load()) answer.VaccumTrigCount = supp.vaccumTrigCount answer.LastVaccumCheck = a.clock.Since(supp.lastVacuumCheckTime) if len(supp.triggers) > 0 { @@ -1013,7 +1191,10 @@ func (a *MergeScheduler) handleMergeDone(table catalog.MergeTable, esz int) { // region: schedule -func (a *MergeScheduler) doSched(todo *todoItem) { +func (a *MergeScheduler) doSched( + generation *mergeSchedulerGeneration, + todo *todoItem, +) { // this table is dropped if todo.table.HasDropCommitted() { delete(a.supps, todo.table.ID()) @@ -1036,7 +1217,7 @@ func (a *MergeScheduler) doSched(todo *todoItem) { now := a.clock.Now() // this table is merging, postpone the task - if supp.mergingTaskCnt > 0 { + if supp.mergingTaskCnt.Load() > 0 { a.pq.Update(todo, now.Add(a.baseInterval/2)) return } @@ -1073,25 +1254,27 @@ func (a *MergeScheduler) doSched(todo *todoItem) { // Gather tasks + rc := a.rc tasks := a.pad.gatherByTrigger( context.Background(), trigger, supp.lastMergeTime, - a.rc, + rc, ) afterGather := a.clock.Now() // Schedule tasks for _, task := range tasks { - task.doneCB = a.taskObserverFactory(todo.table, task.eSize) + task.doneCB = a.taskObserverFactory(supp, task.eSize, rc) if a.executor.ExecuteFor(todo.table, task) { - a.rc.Acquire(int64(task.eSize)) + rc.Acquire(int64(task.eSize)) + supp.AddTask() + task.doneCB.Admit() if task.isTombstone { supp.totalTombstoneMergeCnt++ } else { supp.totalDataMergeCnt++ } - supp.mergingTaskCnt++ if !task.isTombstone && task.oSize > common.DefaultMaxOsizeObjBytes { supp.vaccumTrigCount++ } @@ -1108,12 +1291,14 @@ func (a *MergeScheduler) doSched(todo *todoItem) { if trigger.vacuum != nil { vacuumOpts = trigger.vacuum } - a.ioChan <- &MMsg{ + if !a.sendIOForGeneration(generation, &MMsg{ Kind: MMsgKindVacuumCheck, Value: MMsgVacuumCheck{ Table: todo.table, opts: vacuumOpts, }, + }) { + return } supp.lastVacuumCheckTime = afterGather supp.totalVacuumCheckCnt++ diff --git a/pkg/vm/engine/tae/db/merge/scheduler_extra_test.go b/pkg/vm/engine/tae/db/merge/scheduler_extra_test.go index a254f36105d3e..7f57314c45d08 100644 --- a/pkg/vm/engine/tae/db/merge/scheduler_extra_test.go +++ b/pkg/vm/engine/tae/db/merge/scheduler_extra_test.go @@ -381,6 +381,13 @@ func TestCoverage_taskObserver(t *testing.T) { f: func() { called = true }, } obs.OnExecDone(nil) + assert.False(t, called) + obs.OnExecDone(nil) + assert.False(t, called) + obs.Admit() + assert.True(t, called) + obs.Admit() + obs.OnExecDone(nil) assert.True(t, called) } diff --git a/pkg/vm/engine/tae/db/merge/scheduler_test.go b/pkg/vm/engine/tae/db/merge/scheduler_test.go index ba5adf1d0921d..6b8c984f1cefe 100644 --- a/pkg/vm/engine/tae/db/merge/scheduler_test.go +++ b/pkg/vm/engine/tae/db/merge/scheduler_test.go @@ -18,6 +18,8 @@ import ( "container/heap" "context" "iter" + "sync" + "sync/atomic" "testing" "time" @@ -42,6 +44,15 @@ func (e *dummyExecutor) ExecuteFor(table catalog.MergeTable, task mergeTask) boo return true } +type delayedCompletionExecutor struct { + tasks chan mergeTask +} + +func (e *delayedCompletionExecutor) ExecuteFor(_ catalog.MergeTable, task mergeTask) bool { + e.tasks <- task + return true +} + type dummyCatalogSource struct { settingsFn func() (*batch.Batch, func()) initTables []catalog.MergeTable @@ -88,6 +99,17 @@ func (c *dummyCatalogSource) GetMergeSettingsBatchFn() func() (*batch.Batch, fun return c.settingsFn } +func requireQuery( + t *testing.T, + sched *MergeScheduler, + table catalog.MergeTable, +) *QueryAnswer { + t.Helper() + answer, err := sched.Query(context.Background(), table) + require.NoError(t, err) + return answer +} + type droppedMergeTable struct { catalog.MergeTable } @@ -101,7 +123,6 @@ func TestHandleTaskTriggerNilPointerFixed(t *testing.T) { scheduler := &MergeScheduler{ supps: make(map[uint64]*todoSupporter), msgChan: make(chan *MMsg, 4096), - ioChan: make(chan *MMsg, 256), clock: NewStdClock(), } @@ -119,7 +140,7 @@ func TestHandleTaskTriggerNilPointerFixed(t *testing.T) { // After the fix: This should NOT panic // The early nil check should return gracefully require.NotPanics(t, func() { - scheduler.handleTaskTrigger(msg) + scheduler.handleTaskTrigger(nil, msg) }, "Should not panic after moving nil check before vacuum check") } @@ -134,14 +155,14 @@ func TestDoSchedNilSupporter(t *testing.T) { mockTable := catalog.ToMergeTable(table) require.NotPanics(t, func() { - scheduler.doSched(&todoItem{table: mockTable}) + scheduler.doSched(nil, &todoItem{table: mockTable}) }) todo := &todoItem{table: mockTable, readyAt: scheduler.clock.Now()} heap.Push(&scheduler.pq, todo) require.NotPanics(t, func() { - scheduler.doSched(todo) + scheduler.doSched(nil, todo) }) require.Equal(t, 0, scheduler.pq.Len()) } @@ -180,19 +201,19 @@ func TestScheduler(t *testing.T) { { // switch on/off sched.PauseTable(tables[0]) - answer := sched.Query(tables[0]) + answer := requireQuery(t, sched, tables[0]) require.Equal(t, answer.AutoMergeOn, false) sched.ResumeTable(tables[0]) - answer = sched.Query(tables[0]) + answer = requireQuery(t, sched, tables[0]) require.Equal(t, answer.AutoMergeOn, true) // next check due will be 1s later because of the resume require.Greater(t, answer.NextCheckDue, 900*time.Millisecond) sched.PauseAll() - answer = sched.Query(nil) + answer = requireQuery(t, sched, nil) require.Equal(t, answer.GlobalAutoMergeOn, false) sched.ResumeAll() - answer = sched.Query(nil) + answer = requireQuery(t, sched, nil) require.Equal(t, answer.GlobalAutoMergeOn, true) } @@ -201,7 +222,7 @@ func TestScheduler(t *testing.T) { for i := 0; i < 6; i++ { sched.OnCreateNonAppendObject(tables[0]) } - answer := sched.Query(tables[0]) + answer := requireQuery(t, sched, tables[0]) require.Less(t, answer.NextCheckDue, 500*time.Millisecond) } @@ -209,7 +230,7 @@ func TestScheduler(t *testing.T) { { // create new table sched.OnCreateTableCommit(t1004) - answer := sched.Query(t1004) + answer := requireQuery(t, sched, t1004) require.Equal(t, answer.AutoMergeOn, true) sched.PauseTable(t1004) @@ -293,7 +314,7 @@ func TestScheduler(t *testing.T) { trigger.WithExpire(time.Now().Add(50 * time.Millisecond)) sched.SendTrigger(trigger) - answer := sched.Query(tables[0]) + answer := requireQuery(t, sched, tables[0]) require.Contains(t, answer.Triggers, "L2C: 10") // merge existing patch @@ -303,7 +324,7 @@ func TestScheduler(t *testing.T) { WithTombstone(DefaultTombstoneOpts.Clone().WithL2Count(100)), ) - answer = sched.Query(tables[0]) + answer = requireQuery(t, sched, tables[0]) require.Contains(t, answer.Triggers, "L2C: 100") } @@ -312,7 +333,7 @@ func TestScheduler(t *testing.T) { var answer *QueryAnswer for i := 0; i < 100; i++ { - answer = sched.Query(t1004) + answer = requireQuery(t, sched, t1004) if answer.DataMergeCnt == 1 { break } @@ -321,7 +342,7 @@ func TestScheduler(t *testing.T) { require.Equal(t, answer.DataMergeCnt, 1) for i := 0; i < 100; i++ { - answer = sched.Query(tables[1]) + answer = requireQuery(t, sched, tables[1]) if answer.DataMergeCnt == t1002TaskCnt { break } @@ -331,7 +352,7 @@ func TestScheduler(t *testing.T) { require.Equal(t, answer.VaccumTrigCount, 1) for i := 0; i < 100; i++ { - answer = sched.Query(tables[0]) + answer = requireQuery(t, sched, tables[0]) if answer.DataMergeCnt == 1 { break } @@ -342,12 +363,283 @@ func TestScheduler(t *testing.T) { { // dropped table will be removed from scheduler - answer := sched.Query(tables[2]) + answer := requireQuery(t, sched, tables[2]) require.Equal(t, answer.NotExists, true) } } +type blockingMergeTable struct { + catalog.MergeTable + item catalog.MergeTombstoneItem +} + +func (t *blockingMergeTable) IterTombstoneItem() iter.Seq[catalog.MergeTombstoneItem] { + return func(yield func(catalog.MergeTombstoneItem) bool) { + yield(t.item) + } +} + +type blockingMergeTombstoneItem struct { + stats *objectio.ObjectStats + createdAt types.TS + enteredOnce sync.Once + entered chan struct{} + release chan struct{} +} + +func (i *blockingMergeTombstoneItem) GetCreatedAt() types.TS { + return i.createdAt +} + +func (i *blockingMergeTombstoneItem) GetObjectStats() *objectio.ObjectStats { + return i.stats +} + +func (i *blockingMergeTombstoneItem) ForeachRowid( + context.Context, + any, + func(types.Rowid, bool, int) error, +) error { + i.enteredOnce.Do(func() { + close(i.entered) + }) + <-i.release + return nil +} + +func (i *blockingMergeTombstoneItem) MakeBufferBatch() (any, func()) { + return struct{}{}, func() {} +} + +func TestQueryAndStopBoundedWhenIOQueueFull(t *testing.T) { + db := catalog.MockDBEntryWithAccInfo(1, 1001) + baseTable := catalog.ToMergeTable(catalog.MockTableEntryWithDB(db, 1001)) + item := &blockingMergeTombstoneItem{ + stats: newTestObjectStats( + t, + 1, + 2, + 2*common.DefaultMaxOsizeObjBytes, + 1, + 0, + nil, + 0, + ), + createdAt: types.BuildTS(time.Now().Add(-time.Hour).UnixNano(), 0), + entered: make(chan struct{}), + release: make(chan struct{}), + } + table := &blockingMergeTable{ + MergeTable: baseTable, + item: item, + } + source := &dummyCatalogSource{initTables: []catalog.MergeTable{table}} + sched := NewMergeScheduler( + time.Hour, + source, + &dummyExecutor{}, + NewStdClock(), + ) + sched.Start() + generation := sched.generation.Load() + + var releaseOnce sync.Once + releaseIO := func() { + releaseOnce.Do(func() { + close(item.release) + }) + } + t.Cleanup(func() { + releaseIO() + sched.Stop() + }) + + require.NoError(t, sched.SendTrigger( + NewMMsgTaskTrigger(table).WithVacuumCheck(DefaultVacuumOpts), + )) + select { + case <-item.entered: + case <-time.After(time.Second): + t.Fatal("vacuum I/O did not start") + } + + for i := 0; i <= cap(generation.ioChan); i++ { + require.NoError(t, sched.SendTrigger( + NewMMsgTaskTrigger(table).WithVacuumCheck(DefaultVacuumOpts), + )) + } + require.Eventually(t, func() bool { + return len(generation.ioChan) == cap(generation.ioChan) + }, time.Second, time.Millisecond) + + queryCtx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + _, err := sched.Query(queryCtx, nil) + require.ErrorIs(t, err, context.DeadlineExceeded) + + for len(sched.msgChan) < cap(sched.msgChan) { + sched.msgChan <- &MMsg{ + Kind: MMsgKindTrigger, + Value: NewMMsgTaskTrigger(table). + WithVacuumCheck(DefaultVacuumOpts), + } + } + sendCtx, cancelSend := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancelSend() + _, err = sched.Query(sendCtx, nil) + require.ErrorIs(t, err, context.DeadlineExceeded) + + stopDone := make(chan struct{}) + go func() { + sched.Stop() + close(stopDone) + }() + select { + case <-stopDone: + case <-time.After(2 * time.Second): + t.Fatal("scheduler stop blocked behind the full I/O queue") + } + releaseIO() +} + +func TestStoppedGenerationIOCannotCrossRestart(t *testing.T) { + source := &dummyCatalogSource{} + sched := NewMergeScheduler( + time.Hour, + source, + &dummyExecutor{}, + NewStdClock(), + ) + sched.Start() + stoppedGeneration := sched.generation.Load() + sched.Stop() + + sched.Start() + t.Cleanup(sched.Stop) + currentGeneration := sched.generation.Load() + require.NotSame(t, stoppedGeneration, currentGeneration) + require.NotEqual(t, stoppedGeneration.ioChan, currentGeneration.ioChan) + + staleMsg := &MMsg{Kind: MMsgKindVacuumCheck} + for range cap(stoppedGeneration.ioChan) * 4 { + require.False(t, sched.sendIOForGeneration(stoppedGeneration, staleMsg)) + } + require.Empty(t, stoppedGeneration.ioChan) + require.Empty(t, currentGeneration.ioChan) + + // Simulate the exact race where an old callback passed its stop check and + // completes the send only after Stop returned and a new generation started. + // Its message stays on the stopped generation's private queue. + var staleIOProcessed atomic.Bool + stoppedGeneration.ioChan <- &MMsg{ + Kind: MMsgKindConfigBootstrap, + Value: MMsgConfigBootstrap{ + ReadSettingsBatch: func() (*batch.Batch, func()) { + staleIOProcessed.Store(true) + return nil, func() {} + }, + }, + generation: stoppedGeneration, + } + + staleAnswer := make(chan *QueryAnswer, 1) + sched.msgChan <- &MMsg{ + Kind: MMsgKindQuery, + Value: MMsgQuery{Answer: staleAnswer}, + generation: stoppedGeneration, + } + _, err := sched.Query(context.Background(), nil) + require.NoError(t, err) + require.Empty(t, staleAnswer) + require.Never(t, staleIOProcessed.Load, 50*time.Millisecond, time.Millisecond) + require.Empty(t, currentGeneration.ioChan) +} + +func TestMergeCompletionAccountingSurvivesRestart(t *testing.T) { + db := catalog.MockDBEntryWithAccInfo(1, 1001) + table := catalog.ToMergeTable(catalog.MockTableEntryWithDB(db, 1001)) + source := &dummyCatalogSource{initTables: []catalog.MergeTable{table}} + executor := &delayedCompletionExecutor{tasks: make(chan mergeTask, 1)} + rc := newSimRscController(common.Const1GBytes) + sched := NewMergeScheduler(time.Hour, source, executor, NewStdClock()) + sched.PatchTestRscController(rc) + sched.Start() + t.Cleanup(sched.Stop) + + initialAvailable := rc.Available() + require.NoError(t, sched.SendTrigger( + NewMMsgTaskTrigger(table).WithAssignedTasks([]mergeTask{{ + objs: []*objectio.ObjectStats{ + newTestObjectStats( + t, + 1, + 2, + 8*common.Const1MBytes, + 1000, + 1, + nil, + 0, + ), + }, + note: "delayed completion across restart", + }}), + )) + + var admitted mergeTask + select { + case admitted = <-executor.tasks: + case <-time.After(time.Second): + t.Fatal("merge task was not admitted") + } + require.Positive(t, admitted.eSize) + answer, err := sched.Query(context.Background(), table) + require.NoError(t, err) + require.Equal(t, 1, answer.PendingMergeCnt) + require.Equal(t, initialAvailable-int64(admitted.eSize), rc.Available()) + + sched.Stop() + sched.Start() + + admitted.doneCB.OnExecDone(nil) + answer, err = sched.Query(context.Background(), table) + require.NoError(t, err) + require.Zero(t, answer.PendingMergeCnt) + require.Equal(t, initialAvailable, rc.Available()) + + // Completion observers are allowed to be notified only once. A duplicate + // notification must not underflow the task count or release memory twice. + admitted.doneCB.OnExecDone(nil) + answer, err = sched.Query(context.Background(), table) + require.NoError(t, err) + require.Zero(t, answer.PendingMergeCnt) + require.Equal(t, initialAvailable, rc.Available()) +} + +func TestTaskObserverAdmissionAndCompletionExactlyOnce(t *testing.T) { + var calls atomic.Int64 + observer := &taskObserver{f: func() { + calls.Add(1) + }} + + const workers = 100 + var wg sync.WaitGroup + for i := range workers { + wg.Add(1) + go func() { + defer wg.Done() + if i%2 == 0 { + observer.Admit() + } else { + observer.OnExecDone(nil) + } + }() + } + wg.Wait() + + require.Equal(t, int64(1), calls.Load()) +} + func TestLaunchPad(t *testing.T) { pad := newLaunchPad(NewStdClock()) cata := catalog.MockCatalog(nil) diff --git a/pkg/vm/engine/tae/db/merge/simulator.go b/pkg/vm/engine/tae/db/merge/simulator.go index acc4d08d8f454..9757dbfad003b 100644 --- a/pkg/vm/engine/tae/db/merge/simulator.go +++ b/pkg/vm/engine/tae/db/merge/simulator.go @@ -86,15 +86,13 @@ func (c *fakeClock) Until(t time.Time) time.Duration { // region: resource controller type simRscController struct { - sync.Mutex limit atomic.Int64 - reserved int64 + reserved atomic.Int64 } func newSimRscController(initLimit int64) *simRscController { c := &simRscController{ - limit: atomic.Int64{}, - reserved: 0, + limit: atomic.Int64{}, } c.setMemLimit(initLimit) return c @@ -110,22 +108,30 @@ func (c *simRscController) Refresh() {} func (c *simRscController) PrintUsage() {} func (c *simRscController) Acquire(estMem int64) (int64, bool) { - c.reserved += estMem - return c.Available() - estMem, true + c.reserved.Add(estMem) + return c.Available(), true } func (c *simRscController) Release(estMem int64) int64 { - c.reserved -= estMem - if c.reserved < 0 { - c.reserved = 0 - logutil.Warnf("simRscController: releaseResources: %d", estMem) + for { + reserved := c.reserved.Load() + next := reserved - estMem + if next < 0 { + next = 0 + } + if c.reserved.CompareAndSwap(reserved, next) { + if reserved < estMem { + logutil.Warnf("simRscController: releaseResources: %d", estMem) + } + break + } } return c.Available() } func (c *simRscController) Available() int64 { - avail := c.limit.Load() - c.reserved + avail := c.limit.Load() - c.reserved.Load() if avail < 0 { avail = 0 } @@ -470,7 +476,7 @@ func (e *SExecutor) ExecuteFor(target catalog.MergeTable, task mergeTask) bool { e.clock.AfterFunc(taskCost, func() { if task.doneCB != nil { - task.doneCB.f() + task.doneCB.OnExecDone(nil) } for range newObjCount { e.scatalog.mergeSched.OnCreateNonAppendObject(target) @@ -896,6 +902,7 @@ func NewSimPlayer() *SimPlayer { sexecutor, sclock, ) + scatalog.SetMergeNotifier(sched) sched.PatchTestRscController(srsc) return &SimPlayer{ diff --git a/pkg/vm/engine/tae/db/replay.go b/pkg/vm/engine/tae/db/replay.go index 8c32c70ea02c3..8bceb03a7f6ea 100644 --- a/pkg/vm/engine/tae/db/replay.go +++ b/pkg/vm/engine/tae/db/replay.go @@ -30,6 +30,7 @@ import ( "sync" "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/catalog" + "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/logstore/driver" "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/logstore/wal" @@ -144,6 +145,15 @@ type WalReplayer struct { maxLSN atomic.Uint64 lsn uint64 + + replayAObjectCreates map[replayAObjectCreateKey]types.TS +} + +type replayAObjectCreateKey struct { + dbID uint64 + tableID uint64 + objectID objectio.ObjectId + isTombstone bool } func newWalReplayer( @@ -152,9 +162,10 @@ func newWalReplayer( lsn uint64, ) *WalReplayer { replayer := &WalReplayer{ - db: db, - fromTS: fromTS, - lsn: lsn, + db: db, + fromTS: fromTS, + lsn: lsn, + replayAObjectCreates: make(map[replayAObjectCreateKey]types.TS), } replayer.OnTimeStamp(fromTS) return replayer @@ -184,6 +195,28 @@ func (replayer *WalReplayer) PreReplayWal() { } func (replayer *WalReplayer) postReplayWal() error { + for key, createTS := range replayer.replayAObjectCreates { + dbEntry, err := replayer.db.Catalog.GetDatabaseByID(key.dbID) + if err != nil { + if moerr.IsMoErrCode(err, moerr.OkExpectedEOB) { + continue + } + return err + } + tableEntry, err := dbEntry.GetTableEntryByID(key.tableID) + if err != nil { + if moerr.IsMoErrCode(err, moerr.OkExpectedEOB) { + continue + } + return err + } + if err = tableEntry.UpdateObjectCreateTS(&key.objectID, key.isTombstone, createTS); err != nil { + if moerr.IsMoErrCode(err, moerr.OkExpectedEOB) { + continue + } + return err + } + } processor := new(catalog.LoopProcessor) processor.ObjectFn = func(entry *catalog.ObjectEntry) (err error) { if skippedTbl[entry.GetTable().ID] { @@ -197,6 +230,19 @@ func (replayer *WalReplayer) postReplayWal() error { return replayer.db.Catalog.RecurLoop(processor) } +func (replayer *WalReplayer) RecordReplayAObjectCreate(id *common.ID, isTombstone bool, ts types.TS) { + key := replayAObjectCreateKey{ + dbID: id.DbID, + tableID: id.TableID, + objectID: *id.ObjectID(), + isTombstone: isTombstone, + } + old, ok := replayer.replayAObjectCreates[key] + if !ok || ts.LT(&old) { + replayer.replayAObjectCreates[key] = ts + } +} + func (replayer *WalReplayer) Schedule( ctx context.Context, mode driver.ReplayMode, diff --git a/pkg/vm/engine/tae/db/test/db_test.go b/pkg/vm/engine/tae/db/test/db_test.go index 04b23e408364e..2ae2b724aed3e 100644 --- a/pkg/vm/engine/tae/db/test/db_test.go +++ b/pkg/vm/engine/tae/db/test/db_test.go @@ -1544,13 +1544,13 @@ func TestRollback1(t *testing.T) { tableMeta := rel.GetMeta().(*catalog.TableEntry) err = tableMeta.RecurLoop(processor) assert.Nil(t, err) - assert.Equal(t, objCnt, 1) + assert.Equal(t, 1, objCnt) assert.Nil(t, txn.Rollback(context.Background())) objCnt = 0 err = tableMeta.RecurLoop(processor) assert.Nil(t, err) - assert.Equal(t, objCnt, 0) + assert.Equal(t, 1, objCnt) txn, rel = testutil.GetDefaultRelation(t, db, schema.Name) obj, err := rel.CreateObject(false) @@ -1560,7 +1560,7 @@ func TestRollback1(t *testing.T) { objCnt = 0 err = tableMeta.RecurLoop(processor) assert.Nil(t, err) - assert.Equal(t, objCnt, 1) + assert.Equal(t, 2, objCnt) txn, rel = testutil.GetDefaultRelation(t, db, schema.Name) _, err = rel.GetObject(objMeta.ID(), false) @@ -1576,6 +1576,99 @@ func TestRollback1(t *testing.T) { t.Log(db.Catalog.SimplePPString(common.PPL1)) } +func TestFlushEmptyAppendableObjectReplay(t *testing.T) { + for _, tc := range []struct { + name string + checkpoint bool + }{ + {name: "wal"}, + {name: "checkpoint-collect", checkpoint: true}, + } { + t.Run(tc.name, func(t *testing.T) { + defer testutils.AfterTest(t)() + testutils.EnsureNoLeak(t) + ctx := context.Background() + + opts := config.WithLongScanAndCKPOpts(nil) + tae := testutil.NewTestEngine(ctx, ModuleName, t, opts) + defer tae.Close() + schema := catalog.MockSchema(2, 0) + tae.BindSchema(schema) + setupTxn, err := tae.StartTxn(nil) + require.NoError(t, err) + setupDB, err := testutil.CreateDatabase2(ctx, setupTxn, testutil.DefaultTestDB) + require.NoError(t, err) + _, err = testutil.CreateRelation2(ctx, setupTxn, setupDB, schema) + require.NoError(t, err) + require.NoError(t, setupTxn.Commit(ctx)) + + txn, rel := tae.GetRelation() + obj, err := rel.CreateObject(false) + require.NoError(t, err) + meta := obj.GetMeta().(*catalog.ObjectEntry) + objectID := *meta.ID() + createTS := meta.GetCreatedAt() + require.NoError(t, obj.Close()) + require.NoError(t, txn.Commit(ctx)) + require.True(t, meta.GetObjectData().PrepareCompact()) + + beforeFlushTxn, beforeFlushRel := tae.GetRelation() + flushTxn, _ := tae.GetRelation() + require.Zero(t, tae.Runtime.TransferTable.Len()) + flushStart := flushTxn.GetStartTS() + require.Truef(t, flushStart.GE(&createTS), "flush txn %s starts before object create %s", flushStart.ToString(), createTS.ToString()) + task, err := jobs.NewFlushTableTailTask( + nil, flushTxn, []*catalog.ObjectEntry{meta}, nil, tae.Runtime, + ) + require.NoError(t, err) + require.NoError(t, task.OnExec(ctx)) + require.Nil(t, task.GetCreatedObjects()) + require.NoError(t, flushTxn.Commit(ctx)) + require.Zero(t, tae.Runtime.TransferTable.Len()) + + dedupBat := catalog.MockBatch(schema, 1) + beforeFlushIt := beforeFlushRel.MakeObjectIt(false) + require.True(t, beforeFlushIt.Next()) + require.False(t, beforeFlushIt.Next()) + beforeFlushIt.Close() + require.NoError(t, beforeFlushRel.Append(ctx, dedupBat)) + require.NoError(t, beforeFlushTxn.Rollback(ctx)) + afterFlushTxn, afterFlushRel := tae.GetRelation() + afterFlushIt := afterFlushRel.MakeObjectIt(false) + require.False(t, afterFlushIt.Next()) + afterFlushIt.Close() + require.NoError(t, afterFlushRel.Append(ctx, dedupBat)) + require.NoError(t, afterFlushTxn.Rollback(ctx)) + dedupBat.Close() + + if tc.checkpoint { + anchor := catalog.MockBatch(schema, 1) + txn, rel = tae.GetRelation() + require.NoError(t, rel.Append(ctx, anchor)) + require.NoError(t, txn.Commit(ctx)) + anchor.Close() + tae.CompactBlocks(true) + tae.ForceCheckpoint() + tae.Restart(ctx) + txn, rel = tae.GetRelation() + checkpointed, err := rel.GetMeta().(*catalog.TableEntry).GetObjectByID(&objectID, false) + require.NoError(t, err) + require.True(t, checkpointed.HasDropCommitted()) + require.Equal(t, createTS, checkpointed.GetCreatedAt()) + require.NoError(t, txn.Commit(ctx)) + return + } + tae.Restart(ctx) + txn, rel = tae.GetRelation() + replayed, err := rel.GetMeta().(*catalog.TableEntry).GetObjectByID(&objectID, false) + require.NoError(t, err) + require.True(t, replayed.HasDropCommitted()) + require.Equal(t, createTS, replayed.GetCreatedAt()) + require.NoError(t, txn.Commit(ctx)) + }) + } +} + func TestMVCC1(t *testing.T) { defer testutils.AfterTest(t)() testutils.EnsureNoLeak(t) @@ -10364,6 +10457,55 @@ func TestCollectDeletesInRange1(t *testing.T) { tae.CheckCollectTombstoneInRange() } +func TestCollectDeletesInRangeWithActiveTombstoneDrop(t *testing.T) { + defer testutils.AfterTest(t)() + testutils.EnsureNoLeak(t) + ctx := context.Background() + + opts := config.WithLongScanAndCKPOpts(nil) + tae := testutil.NewTestEngine(ctx, ModuleName, t, opts) + defer tae.Close() + schema := catalog.MockSchemaAll(2, 1) + tae.BindSchema(schema) + bat := catalog.MockBatch(schema, 2) + defer bat.Close() + + tae.CreateRelAndAppend(bat, true) + + txn, rel := tae.GetRelation() + dataObj := testutil.GetOneObject(rel) + dataObjectID := *dataObj.GetID() + filter := handle.NewEQFilter(bat.Vecs[schema.GetSingleSortKeyIdx()].Get(0)) + require.NoError(t, rel.DeleteByFilter(ctx, filter)) + require.NoError(t, txn.Commit(ctx)) + + dropTxn, dropRel := tae.GetRelation() + defer dropTxn.Rollback(ctx) + tombstone := testutil.GetOneTombstoneMeta(dropRel) + // Model a concurrent tombstone flush after it installs the D entry but + // before the flush transaction starts committing. + require.NoError(t, dropTxn.GetStore().SoftDeleteObject(true, tombstone.AsCommonID())) + dropped := tombstone.GetLatestNode() + require.True(t, dropped.IsDEntry()) + require.False(t, dropped.HasDropCommitted()) + require.Equal(t, txnif.UncommitTS, dropped.GetDeletedAt()) + + tableEntry := dropRel.GetMeta().(*catalog.TableEntry) + deletes, err := tables.TombstoneRangeScanByObject( + ctx, + tableEntry, + dataObjectID, + types.TS{}, + dropTxn.GetStartTS(), + common.DefaultAllocator, + tae.Runtime.VectorPool.Small, + ) + require.NoError(t, err) + require.NotNil(t, deletes) + defer deletes.Close() + require.Equal(t, 1, deletes.Length()) +} + func TestCollectDeletesInRange2(t *testing.T) { defer testutils.AfterTest(t)() ctx := context.Background() @@ -11871,19 +12013,26 @@ func TestRW2(t *testing.T) { assert.True(t, moerr.IsMoErrCode(err, moerr.ErrTxnRWConflict)) } +func newTestTxnServer(t *testing.T) rpc.TxnServer { + t.Helper() + server, err := rpc.NewTxnServer("127.0.0.1:0", runtime.ServiceRuntime("")) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, server.Close()) + }) + return server +} + func Test_BasicTxnModeSwitch(t *testing.T) { ctx := context.Background() opts := config.WithLongScanAndCKPOpts(nil) tae := testutil.NewTestEngine(ctx, ModuleName, t, opts) - - var err error - tae.TxnServer, err = rpc.NewTxnServer("localhost:12345", runtime.ServiceRuntime("")) - require.NoError(t, err) + tae.TxnServer = newTestTxnServer(t) defer tae.Close() assert.True(t, tae.IsWriteMode()) - err = tae.SwitchTxnMode(ctx, 1, "todo") + err := tae.SwitchTxnMode(ctx, 1, "todo") assert.NoError(t, err) assert.True(t, tae.IsReplayMode()) assert.True(t, tae.TxnMgr.IsReplayMode()) @@ -11895,6 +12044,115 @@ func Test_BasicTxnModeSwitch(t *testing.T) { assert.Error(t, db.CheckCronJobs(tae.DB, db.DBTxnMode_Replay)) } +func prepareTxnModeSwitchWithInflightTxn( + t *testing.T, +) (*testutil.TestEngine, txnif.AsyncTxn, handle.Relation, string) { + t.Helper() + ctx := context.Background() + opts := config.WithLongScanAndCKPOpts(nil) + tae := testutil.NewTestEngine(ctx, ModuleName, t, opts) + tae.TxnServer = newTestTxnServer(t) + + schema := catalog.MockSchema(1, -1) + txn, err := tae.StartTxn(nil) + require.NoError(t, err) + database, err := txn.CreateDatabase("mode-switch", "", "") + require.NoError(t, err) + _, err = database.CreateRelation(schema) + require.NoError(t, err) + require.NoError(t, txn.Commit(ctx)) + + txn, rel := testutil.GetRelation(t, 0, tae.DB, "mode-switch", schema.Name) + return tae, txn, rel, schema.Name +} + +func TestTxnModeSwitchDrainsInflightCatalogEvents(t *testing.T) { + ctx := context.Background() + tae, txn, rel, _ := prepareTxnModeSwitchWithInflightTxn(t) + defer tae.Close() + + switchCtx, cancel := context.WithTimeout(ctx, 30*time.Second) + defer cancel() + switchErr := make(chan error, 1) + go func() { + switchErr <- tae.SwitchTxnMode(switchCtx, 1, "todo") + }() + + require.Eventually(t, func() bool { + return !tae.TxnMgr.IsWriteMode() + }, 10*time.Second, time.Millisecond) + + commitErr := make(chan error, 1) + go func() { + for i := 0; i < 4097; i++ { + if _, err := rel.CreateNonAppendableObject(false, nil); err != nil { + commitErr <- err + return + } + } + commitErr <- txn.Commit(ctx) + }() + + select { + case err := <-commitErr: + require.NoError(t, err) + case <-switchCtx.Done(): + t.Fatal("in-flight transaction blocked while publishing catalog events") + } + select { + case err := <-switchErr: + require.NoError(t, err) + case <-switchCtx.Done(): + t.Fatal("write to replay mode switch did not finish") + } + require.True(t, tae.IsReplayMode()) + require.True(t, tae.TxnMgr.IsReplayMode()) +} + +func TestTxnModeSwitchWaitCancelRollsBack(t *testing.T) { + tae, txn, _, tableName := prepareTxnModeSwitchWithInflightTxn(t) + defer tae.Close() + + switchCtx, cancel := context.WithCancel(context.Background()) + switchErr := make(chan error, 1) + go func() { + switchErr <- tae.SwitchTxnMode(switchCtx, 1, "todo") + }() + + require.Eventually(t, func() bool { + return !tae.TxnMgr.IsWriteMode() + }, 10*time.Second, time.Millisecond) + cancel() + + select { + case err := <-switchErr: + require.ErrorIs(t, err, context.Canceled) + case <-time.After(10 * time.Second): + t.Fatal("canceled mode switch did not roll back") + } + require.True(t, tae.IsWriteMode()) + require.True(t, tae.TxnMgr.IsWriteMode()) + require.NoError(t, db.CheckCronJobs(tae.DB, db.DBTxnMode_Write)) + require.NoError(t, txn.Rollback(context.Background())) + + // The notifier and scheduler must remain usable after rollback. + txn, rel := testutil.GetRelation(t, 0, tae.DB, "mode-switch", tableName) + _, err := rel.CreateNonAppendableObject(false, nil) + require.NoError(t, err) + require.NoError(t, txn.Commit(context.Background())) + queryDone := make(chan error, 1) + go func() { + _, err := tae.MergeScheduler.Query(context.Background(), nil) + queryDone <- err + }() + select { + case err := <-queryDone: + require.NoError(t, err) + case <-time.After(10 * time.Second): + t.Fatal("merge scheduler was not restarted after mode switch rollback") + } +} + func Test_OpenReplayDB1(t *testing.T) { ctx := context.Background() opts := config.WithLongScanAndCKPOpts(nil) @@ -11933,6 +12191,9 @@ func TestRW3(t *testing.T) { ctx := context.Background() opts := config.WithLongScanAndCKPOpts(nil) tae := testutil.NewTestEngine(ctx, ModuleName, t, opts) + defer func() { + t.Log(tae.Catalog.SimplePPString(common.PPL3)) + }() objCount := 100 schema := catalog.MockSchemaAll(1, -1) @@ -12842,13 +13103,31 @@ func Test_ApplyTableData(t *testing.T) { err = applyArg.Run() assert.NoError(t, err) - txn, rel := testutil.GetRelation(t, 0, tae.DB, "db2", "table2") - assert.NoError(t, txn.Commit(ctx)) - for i := 0; i < colCount; i++ { - rows := testutil.GetColumnRowsByScan(t, rel, i, true) - assert.Equal(t, 2, rows) + checkAppliedTable := func() { + t.Helper() + txn, rel := testutil.GetRelation(t, 0, tae.DB, "db2", "table2") + it := rel.MakeObjectIt(false) + objectCount := 0 + for it.Next() { + objectCount++ + meta := it.GetObject().GetMeta().(*catalog.ObjectEntry) + require.False(t, meta.IsAppendable()) + require.True(t, meta.ObjectPersisted()) + require.False(t, meta.GetObjectData().IsAppendable()) + } + require.NoError(t, it.Close()) + require.Equal(t, 1, objectCount) + for i := 0; i < colCount; i++ { + rows := testutil.GetColumnRowsByScan(t, rel, i, true) + require.Equal(t, 2, rows) + } + require.NoError(t, txn.Commit(ctx)) } + checkAppliedTable() + tae.Restart(ctx) + checkAppliedTable() + t.Log(tae.Catalog.SimplePPString(3)) } diff --git a/pkg/vm/engine/tae/db/test/replay_test.go b/pkg/vm/engine/tae/db/test/replay_test.go index fb2030b345a41..ab6a032eb1abe 100644 --- a/pkg/vm/engine/tae/db/test/replay_test.go +++ b/pkg/vm/engine/tae/db/test/replay_test.go @@ -255,8 +255,9 @@ func TestReplayCatalog3(t *testing.T) { assert.Nil(t, err) rel, err = e.GetRelationByName(schema.Name) assert.Nil(t, err) - obj, err = rel.CreateObject(false) + obj, err = rel.CreateNonAppendableObject(false, nil) assert.Nil(t, err) + testutil.MockObjectStats(t, obj) assert.Nil(t, txn.Commit(context.Background())) txn, _ = tae.StartTxn(nil) diff --git a/pkg/vm/engine/tae/db/test/tables_test.go b/pkg/vm/engine/tae/db/test/tables_test.go index d905e873828ab..4c76b18ae7021 100644 --- a/pkg/vm/engine/tae/db/test/tables_test.go +++ b/pkg/vm/engine/tae/db/test/tables_test.go @@ -174,8 +174,10 @@ func TestTxn1(t *testing.T) { blkCnt += uint32(objIt.GetObject().BlkCnt()) } objIt.Close() - assert.Equal(t, expectObjCnt, objCnt) - assert.Equal(t, expectBlkCnt, blkCnt) + // Appendable objects are created committed outside the creating + // transaction, so the explicitly created empty object is visible here. + assert.Equal(t, expectObjCnt+1, objCnt) + assert.Equal(t, expectBlkCnt+1, blkCnt) } t.Log(db.Catalog.SimplePPString(common.PPL1)) } diff --git a/pkg/vm/engine/tae/iface/handle/relation.go b/pkg/vm/engine/tae/iface/handle/relation.go index 36fda9bad3420..50a735f20df3e 100644 --- a/pkg/vm/engine/tae/iface/handle/relation.go +++ b/pkg/vm/engine/tae/iface/handle/relation.go @@ -54,6 +54,7 @@ type Relation interface { GetMeta() any CreateObject(bool) (Object, error) + CreateObjectWithOpt(isTombstone bool, opt *objectio.CreateObjOpt) (Object, error) CreateNonAppendableObject(isTombstone bool, opt *objectio.CreateObjOpt) (Object, error) GetObject(id *types.Objectid, isTombstone bool) (Object, error) SoftDeleteObject(id *types.Objectid, isTombstone bool) (err error) diff --git a/pkg/vm/engine/tae/iface/txnif/types.go b/pkg/vm/engine/tae/iface/txnif/types.go index 235e030aad2af..36425c50c0c46 100644 --- a/pkg/vm/engine/tae/iface/txnif/types.go +++ b/pkg/vm/engine/tae/iface/txnif/types.go @@ -293,6 +293,7 @@ type TxnStore interface { GetObject(id *common.ID, isTombstone bool) (handle.Object, error) CreateObject(dbId, tid uint64, isTombstone bool) (handle.Object, error) + CreateObjectWithOpt(dbId, tid uint64, isTombstone bool, opt *objectio.CreateObjOpt) (handle.Object, error) CreateNonAppendableObject(dbId, tid uint64, isTombstone bool, opt *objectio.CreateObjOpt) (handle.Object, error) SoftDeleteObject(isTombstone bool, id *common.ID) error diff --git a/pkg/vm/engine/tae/rpc/apply_table_data.go b/pkg/vm/engine/tae/rpc/apply_table_data.go index a915f1b0782b1..526d9d5a368e5 100644 --- a/pkg/vm/engine/tae/rpc/apply_table_data.go +++ b/pkg/vm/engine/tae/rpc/apply_table_data.go @@ -211,18 +211,28 @@ func (a *ApplyTableDataArg) Run() (err error) { panic(fmt.Sprintf("invalid object type: %d", objTypes[i])) } stats := objectio.ObjectStats(idVec.GetBytesAt(i)) - var obj handle.Object - if obj, err = a.rel.CreateNonAppendableObject( - isTombstone, - &objectio.CreateObjOpt{ - Stats: &stats, - IsTombstone: isTombstone, - }, - ); err != nil { + var createLiveAppendable bool + if stats, createLiveAppendable, err = prepareObjectStatsForApply(stats, isPersisted[i]); err != nil { return } + opt := &objectio.CreateObjOpt{ + Stats: &stats, + IsTombstone: isTombstone, + } + tableEntry := a.rel.GetMeta().(*catalog.TableEntry) + var obj handle.Object + if createLiveAppendable { + if obj, err = a.rel.CreateObjectWithOpt(isTombstone, opt); err != nil { + return + } + } else { + if obj, err = a.rel.CreateNonAppendableObject(isTombstone, opt); err != nil { + return + } + } + meta := obj.GetMeta().(*catalog.ObjectEntry) - schema := a.rel.GetMeta().(*catalog.TableEntry).GetLastestSchema(isTombstone) + schema := tableEntry.GetLastestSchema(isTombstone) attrs := schema.AllNames() attrs = append(attrs, objectio.TombstoneAttr_CommitTs_Attr) @@ -235,7 +245,6 @@ func (a *ApplyTableDataArg) Run() (err error) { } defer objectRelease() tnBat := containers.ToTNBatch(bat, a.mp) - meta := obj.GetMeta().(*catalog.ObjectEntry) var anodes []txnif.TxnEntry if anodes, err = meta.GetObjectData().ApplyDebugBatch(tnBat, a.txn); err != nil { return @@ -270,6 +279,32 @@ func (a *ApplyTableDataArg) Run() (err error) { } +func prepareObjectStatsForApply( + stats objectio.ObjectStats, + isPersisted bool, +) (prepared objectio.ObjectStats, createLiveAppendable bool, err error) { + prepared = stats + if !stats.GetAppendable() { + return + } + if isPersisted { + // A flushed appendable object is immutable file-backed history. Restore + // it as a non-appendable catalog object so it uses a persisted node and + // remains transactional for rollback and WAL replay. + objectio.SetObjectStatsAppendable(&prepared, false) + return prepared, false, nil + } + if stats.Rows() > 0 { + err = moerr.NewInternalErrorNoCtxf( + "live appendable object %s has non-empty row count %d", + stats.ObjectName().String(), + stats.Rows(), + ) + return + } + return prepared, true, nil +} + func (a *ApplyTableDataArg) createDatabase() (err error) { var database handle.Database diff --git a/pkg/vm/engine/tae/rpc/apply_table_data_test.go b/pkg/vm/engine/tae/rpc/apply_table_data_test.go new file mode 100644 index 0000000000000..5ce9ac5f601ee --- /dev/null +++ b/pkg/vm/engine/tae/rpc/apply_table_data_test.go @@ -0,0 +1,45 @@ +// Copyright 2026 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 rpc + +import ( + "testing" + + "github.com/matrixorigin/matrixone/pkg/objectio" + "github.com/stretchr/testify/require" +) + +func TestPrepareObjectStatsForApply(t *testing.T) { + id := objectio.NewObjectid() + flushed := objectio.NewObjectStatsWithObjectID(&id, true, false, false) + require.NoError(t, objectio.SetObjectStatsRowCnt(flushed, 2)) + + prepared, createLiveAppendable, err := prepareObjectStatsForApply(*flushed, true) + require.NoError(t, err) + require.False(t, createLiveAppendable) + require.False(t, prepared.GetAppendable()) + require.Equal(t, uint32(2), prepared.Rows()) + // Preparing the restored catalog state must not mutate the dump metadata. + require.True(t, flushed.GetAppendable()) + + _, _, err = prepareObjectStatsForApply(*flushed, false) + require.ErrorContains(t, err, "live appendable object") + + live := objectio.NewObjectStatsWithObjectID(&id, true, false, false) + prepared, createLiveAppendable, err = prepareObjectStatsForApply(*live, false) + require.NoError(t, err) + require.True(t, createLiveAppendable) + require.True(t, prepared.GetAppendable()) +} diff --git a/pkg/vm/engine/tae/rpc/dump_table.go b/pkg/vm/engine/tae/rpc/dump_table.go index 34d778347c1e0..37267e9c09172 100644 --- a/pkg/vm/engine/tae/rpc/dump_table.go +++ b/pkg/vm/engine/tae/rpc/dump_table.go @@ -282,11 +282,19 @@ func (c *DumpTableArg) Run() (err error) { return err } - p := &catalog.LoopProcessor{} - p.ObjectFn = c.onObject - p.TombstoneFn = c.onObject - if err = c.table.RecurLoop(p); err != nil { - return err + dataIt := c.table.MakeDataVisibleObjectIt(c.txn) + defer dataIt.Release() + for dataIt.Next() { + if err = c.onObject(dataIt.Item()); err != nil { + return err + } + } + tombstoneIt := c.table.MakeTombstoneVisibleObjectIt(c.txn) + defer tombstoneIt.Release() + for tombstoneIt.Next() { + if err = c.onObject(tombstoneIt.Item()); err != nil { + return err + } } if err := c.flush(DumpTableObjectList, c.objectListBatch); err != nil { return err @@ -512,11 +520,11 @@ func (c *DumpTableArg) collectObjectList(e *catalog.ObjectEntry) (isPersisted bo } else { deleteTS = e.DeletedAt } - if e.GetAppendable() && deleteTS.IsEmpty() { - isPersisted = false - } else { - isPersisted = true - } + // Appendable describes the object's original layout, not whether its data + // is still backed by a live memory node. A checkpoint-replayed appendable + // object can be forced persisted without a delete entry. An active drop is + // not sufficient here because its file-backed state is not committed yet. + isPersisted = !e.GetAppendable() || e.IsForcePNode() || e.HasDropCommitted() if err := vector.AppendFixed( c.objectListBatch.Vecs[ObjectListAttr_DeleteTS_Idx], deleteTS, false, c.mp, ); err != nil { diff --git a/pkg/vm/engine/tae/rpc/inspectMerge.go b/pkg/vm/engine/tae/rpc/inspectMerge.go index 908db0531d709..c2f155b41a26b 100644 --- a/pkg/vm/engine/tae/rpc/inspectMerge.go +++ b/pkg/vm/engine/tae/rpc/inspectMerge.go @@ -231,7 +231,10 @@ func (arg *mergeShowArg) Run() error { if arg.tbl != nil { target = catalog.ToMergeTable(arg.tbl) } - answer := arg.ctx.db.MergeScheduler.Query(target) + answer, err := arg.ctx.db.MergeScheduler.Query(context.Background(), target) + if err != nil { + return err + } out.WriteString(fmt.Sprintf( "auto merge for all: %v, msg queue len: %d\n", answer.GlobalAutoMergeOn, diff --git a/pkg/vm/engine/tae/rpc/inspect_test.go b/pkg/vm/engine/tae/rpc/inspect_test.go index 24078ff7355dd..f0511e97e6735 100644 --- a/pkg/vm/engine/tae/rpc/inspect_test.go +++ b/pkg/vm/engine/tae/rpc/inspect_test.go @@ -75,11 +75,7 @@ func TestMergeCommand(t *testing.T) { vector := containers.NewVector(types.T_varchar.ToType()) { id := objectio.NewObjectid() - stats := objectio.NewObjectStatsWithObjectID(&id, true, true, false) - vector.Append(stats.Marshal(), false) - - id = objectio.NewObjectid() - stats = objectio.NewObjectStatsWithObjectID(&id, false, true, false) + stats := objectio.NewObjectStatsWithObjectID(&id, false, true, false) zm := index.NewZM(types.T_int32, 0) v1 := int32(1) v2 := int32(2) diff --git a/pkg/vm/engine/tae/tables/aobj.go b/pkg/vm/engine/tae/tables/aobj.go index b46a1a9ed60c5..70401470db355 100644 --- a/pkg/vm/engine/tae/tables/aobj.go +++ b/pkg/vm/engine/tae/tables/aobj.go @@ -118,7 +118,6 @@ func (obj *aobject) PrepareCompact() bool { } return false } - // see more notes in flushtabletail.go obj.freezelock.Lock() obj.FreezeAppend() diff --git a/pkg/vm/engine/tae/tables/jobs/flushTableTail.go b/pkg/vm/engine/tae/tables/jobs/flushTableTail.go index 4a242494a91a9..b07ebf6b50bd5 100644 --- a/pkg/vm/engine/tae/tables/jobs/flushTableTail.go +++ b/pkg/vm/engine/tae/tables/jobs/flushTableTail.go @@ -506,6 +506,10 @@ func (task *flushTableTailTask) prepareAObjSortedData( if err != nil { return } + if bat == nil { + empty = true + return + } for i := range idxs { if vec := bat.Vecs[i]; vec == nil || vec.Length() == 0 { empty = true @@ -648,6 +652,7 @@ func (task *flushTableTailTask) mergeAObjs(ctx context.Context, isTombstone bool } if !isTombstone && task.doTransfer { mergesort.ReleaseTransferMaps(task.transMappings) + task.transMappings = nil } return nil } diff --git a/pkg/vm/engine/tae/tables/table_scan.go b/pkg/vm/engine/tae/tables/table_scan.go index 5b4465d448e0f..e78449f3cff7f 100644 --- a/pkg/vm/engine/tae/tables/table_scan.go +++ b/pkg/vm/engine/tae/tables/table_scan.go @@ -80,8 +80,10 @@ TombstoneRangeScanByObject scans the an object's tombstones committed in the ran Since the returned batch must have accruate ts for each row, we need collect the data from appendable objects. Targets: -1. CNCreated entries where start <= CreatedAt <= end -2. Appendable entries where x <= CreatedAt <= end, where x is the first appendable entry with CreatedAt < start + 1. CNCreated entries where start <= CreatedAt <= end + 2. Appendable entries whose catalog lifetime can overlap the range. All live + appendable entries created before end remain candidates because their rows + can commit out of object creation order. */ func TombstoneRangeScanByObject( ctx context.Context, @@ -94,15 +96,16 @@ func TombstoneRangeScanByObject( tableEntry.WaitTombstoneObjectCommitted(end) it := tableEntry.MakeTombstoneObjectIt() defer it.Release() - earlybreak := false + // CreatedAt orders catalog publication, not the commit timestamps of rows + // appended later. Concurrent flushes can populate multiple appendable + // tombstone objects and commit them out of creation order, so an older + // object can still contain deletes in [start, end]. Do not stop the scan + // solely because an object's catalog lifetime precedes start. for ok := it.Last(); ok; ok = it.Prev() { - if earlybreak { - break - } - tombstone := it.Item() - // we only check the created version of the object. - if tombstone.HasDropIntent() { + if tombstone.IsCEntry() && tombstone.HasDCounterpart() && tombstone.GetNextVersion().HasDropCommitted() { + // The dropped counterpart owns the persisted appendable tombstone data. + // Scanning both versions duplicates the same committed delete rows. continue } @@ -111,11 +114,21 @@ func TombstoneRangeScanByObject( // committing create object is excluded here continue } - // first committed appendable object with CreatedAt < start, stop at next round - if tombstone.CreatedAt.LT(&start) { - earlybreak = true + if tombstone.DeletedAt.Equal(&txnif.UncommitTS) { + // Its C counterpart remains visible until the drop commits. + continue + } + if tombstone.HasDropCommitted() { + deleteAt := tombstone.GetDeleteAt() + if tombstone.CreatedAt.GT(&end) || deleteAt.LT(&start) { + continue + } } } else { + // we only check the created version of the object. + if tombstone.HasDropIntent() { + continue + } if !tombstone.ObjectStats.GetCNCreated() { continue } diff --git a/pkg/vm/engine/tae/tables/table_scan_test.go b/pkg/vm/engine/tae/tables/table_scan_test.go new file mode 100644 index 0000000000000..ea9f6c0df98ae --- /dev/null +++ b/pkg/vm/engine/tae/tables/table_scan_test.go @@ -0,0 +1,92 @@ +// Copyright 2026 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 tables + +import ( + "context" + "testing" + + "github.com/matrixorigin/matrixone/pkg/common/mpool" + "github.com/matrixorigin/matrixone/pkg/container/types" + "github.com/matrixorigin/matrixone/pkg/objectio" + "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/catalog" + "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/containers" + "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/iface/data" + "github.com/stretchr/testify/require" +) + +type tombstoneScanCounter struct { + data.Object + scans int +} + +func (counter *tombstoneScanCounter) CollectObjectTombstoneInRange( + context.Context, + types.TS, + types.TS, + *types.Objectid, + **containers.Batch, + *mpool.MPool, + *containers.VectorPool, +) error { + counter.scans++ + return nil +} + +func TestTombstoneRangeScanKeepsOlderAppendableCandidates(t *testing.T) { + c := catalog.MockCatalog(nil) + defer c.Close() + db, err := c.CreateDBEntry("db", "", "", nil) + require.NoError(t, err) + table, err := db.CreateTableEntry(catalog.MockSchema(1, 0), nil, nil) + require.NoError(t, err) + + oldID := objectio.NewObjectid() + oldStats := objectio.NewObjectStatsWithObjectID(&oldID, true, false, false) + oldData := &tombstoneScanCounter{} + _, err = table.CreateCommittedObject( + types.BuildTS(1, 0), + &objectio.CreateObjOpt{Stats: oldStats, IsTombstone: true}, + func(*catalog.ObjectEntry) data.Object { return oldData }, + ) + require.NoError(t, err) + + newerID := objectio.NewObjectid() + newerStats := objectio.NewObjectStatsWithObjectID(&newerID, true, false, false) + newer, err := table.CreateCommittedObject( + types.BuildTS(2, 0), + &objectio.CreateObjOpt{Stats: newerStats, IsTombstone: true}, + nil, + ) + require.NoError(t, err) + catalog.MockDroppedObjectEntry2List(newer, types.BuildTS(3, 0)) + + targetID := objectio.NewObjectid() + bat, err := TombstoneRangeScanByObject( + context.Background(), + table, + targetID, + types.BuildTS(4, 0), + types.BuildTS(5, 0), + nil, + nil, + ) + require.NoError(t, err) + require.Nil(t, bat) + // Object creation order does not order later append commits. Even though + // the newer object's lifetime precedes start, the older appendable object + // can still contain rows committed in the requested range. + require.Equal(t, 1, oldData.scans) +} diff --git a/pkg/vm/engine/tae/tables/txnentries/flushTableTail.go b/pkg/vm/engine/tae/tables/txnentries/flushTableTail.go index 3f847b3893395..f6e7152fce5a4 100644 --- a/pkg/vm/engine/tae/tables/txnentries/flushTableTail.go +++ b/pkg/vm/engine/tae/tables/txnentries/flushTableTail.go @@ -197,6 +197,7 @@ func (entry *flushTableTailEntry) addTransferPages(ctx context.Context) error { func (entry *flushTableTailEntry) collectDelsAndTransfer( ctx context.Context, from, to types.TS, ) (transCnt int, err error) { + scanStart := from.Next() if len(entry.aobjHandles) == 0 { return } @@ -226,7 +227,7 @@ func (entry *flushTableTailEntry) collectDelsAndTransfer( ctx, entry.tableEntry, *obj.ID(), - from.Next(), // NOTE HERE + scanStart, // NOTE HERE to, common.MergeAllocator, entry.rt.VectorPool.Small, @@ -285,7 +286,12 @@ func (entry *flushTableTailEntry) PrepareCommit() error { return nil } ctx := context.Background() - trans, err := entry.collectDelsAndTransfer(ctx, entry.collectTs, entry.txn.GetPrepareTS().Prev()) + txnStart := entry.txn.GetStartTS() + flushScanStart := txnStart.Next() + commitScanStart := entry.collectTs.Next() + prepareTS := entry.txn.GetPrepareTS() + preparePrev := prepareTS.Prev() + trans, err := entry.collectDelsAndTransfer(ctx, entry.collectTs, preparePrev) if err != nil { return err } @@ -294,7 +300,14 @@ func (entry *flushTableTailEntry) PrepareCommit() error { logutil.Info( "[FLUSH-PREPARE-COMMIT]", zap.String("task", entry.taskName), - zap.String("commit-ts", entry.txn.GetPrepareTS().ToString()), + zap.String("commit-ts", prepareTS.ToString()), + zap.String("transfer-split-ts", entry.collectTs.ToString()), + zap.String("flush-range-from", txnStart.ToString()), + zap.String("flush-range-scan-start", flushScanStart.ToString()), + zap.String("flush-range-to", entry.collectTs.ToString()), + zap.String("commit-range-from", entry.collectTs.ToString()), + zap.String("commit-range-scan-start", commitScanStart.ToString()), + zap.String("commit-range-to", preparePrev.ToString()), zap.Int("ablks", aconflictCnt), zap.Int("transfer-rows", totalTrans), zap.Int("in-queue-transfers", trans), diff --git a/pkg/vm/engine/tae/tables/updates/append.go b/pkg/vm/engine/tae/tables/updates/append.go index 0c087ee3b85c5..6901fef455067 100644 --- a/pkg/vm/engine/tae/tables/updates/append.go +++ b/pkg/vm/engine/tae/tables/updates/append.go @@ -139,10 +139,11 @@ func (node *AppendNode) ApplyCommit(id string) error { } node.TxnMVCCNode.ApplyCommit(id) listener := node.mvcc.GetAppendListener() - if listener == nil { - return nil + var err error + if listener != nil { + err = listener(node) } - return listener(node) + return err } func (node *AppendNode) ApplyRollback() (err error) { diff --git a/pkg/vm/engine/tae/tables/updates/mvcc.go b/pkg/vm/engine/tae/tables/updates/mvcc.go index 3f0ecfd48c5b8..749f2a0810d2b 100644 --- a/pkg/vm/engine/tae/tables/updates/mvcc.go +++ b/pkg/vm/engine/tae/tables/updates/mvcc.go @@ -315,6 +315,9 @@ func (n *AppendMVCCHandle) PrepareCompact() bool { } func (n *AppendMVCCHandle) GetLatestAppendPrepareTSLocked() types.TS { + if n.appends == nil || n.appends.IsEmpty() { + return types.TS{} + } return n.appends.GetUpdateNodeLocked().Prepare } func (n *AppendMVCCHandle) GetMeta() *catalog.ObjectEntry { @@ -331,6 +334,9 @@ func (n *AppendMVCCHandle) allAppendsCommittedLocked() bool { meta.GetDeleteAt().ToString()) return false } + if n.appends.IsEmpty() { + return true + } return n.appends.IsCommitted() } diff --git a/pkg/vm/engine/tae/txn/txnbase/handle.go b/pkg/vm/engine/tae/txn/txnbase/handle.go index d0df9f0de63cd..8292f3d138831 100644 --- a/pkg/vm/engine/tae/txn/txnbase/handle.go +++ b/pkg/vm/engine/tae/txn/txnbase/handle.go @@ -82,6 +82,9 @@ func (rel *TxnRelation) GetObject(id *types.Objectid, isTombstone bool) (obj han } func (rel *TxnRelation) SoftDeleteObject(id *types.Objectid, isTombstone bool) (err error) { return } func (rel *TxnRelation) CreateObject(bool) (obj handle.Object, err error) { return } +func (rel *TxnRelation) CreateObjectWithOpt(bool, *objectio.CreateObjOpt) (obj handle.Object, err error) { + return +} func (rel *TxnRelation) CreateNonAppendableObject(bool, *objectio.CreateObjOpt) (obj handle.Object, err error) { return } diff --git a/pkg/vm/engine/tae/txn/txnbase/store.go b/pkg/vm/engine/tae/txn/txnbase/store.go index 8635f38bf261b..8db2b4fb02748 100644 --- a/pkg/vm/engine/tae/txn/txnbase/store.go +++ b/pkg/vm/engine/tae/txn/txnbase/store.go @@ -104,6 +104,9 @@ func (store *NoopTxnStore) GetObject(id *common.ID, isTombstone bool) (obj handl func (store *NoopTxnStore) CreateObject(dbId, tid uint64, isTombstone bool) (obj handle.Object, err error) { return } +func (store *NoopTxnStore) CreateObjectWithOpt(dbId, tid uint64, _ bool, _ *objectio.CreateObjOpt) (obj handle.Object, err error) { + return +} func (store *NoopTxnStore) CreateNonAppendableObject(dbId, tid uint64, _ bool, _ *objectio.CreateObjOpt) (obj handle.Object, err error) { return } diff --git a/pkg/vm/engine/tae/txn/txnbase/txnmgr.go b/pkg/vm/engine/tae/txn/txnbase/txnmgr.go index 447658811b8e4..7ef4fbb88e05b 100644 --- a/pkg/vm/engine/tae/txn/txnbase/txnmgr.go +++ b/pkg/vm/engine/tae/txn/txnbase/txnmgr.go @@ -123,6 +123,51 @@ func (bl *batchTxnCommitListener) OnEndPrepareWAL(txn txnif.AsyncTxn) { type TxnStoreFactory = func() txnif.TxnStore type TxnFactory = func(*TxnManager, txnif.TxnStore, []byte, types.TS, types.TS) txnif.AsyncTxn +type txnWaiter struct { + mu sync.Mutex + count int + emptyCh chan struct{} +} + +func (w *txnWaiter) Add() { + w.mu.Lock() + defer w.mu.Unlock() + if w.count == 0 { + w.emptyCh = make(chan struct{}) + } + w.count++ +} + +func (w *txnWaiter) Done() { + w.mu.Lock() + defer w.mu.Unlock() + if w.count <= 0 { + panic("txn waiter: negative transaction count") + } + w.count-- + if w.count == 0 { + close(w.emptyCh) + w.emptyCh = nil + } +} + +func (w *txnWaiter) Wait(ctx context.Context) error { + w.mu.Lock() + if w.count == 0 { + w.mu.Unlock() + return nil + } + emptyCh := w.emptyCh + w.mu.Unlock() + + select { + case <-ctx.Done(): + return ctx.Err() + case <-emptyCh: + return nil + } +} + type TxnManager struct { sm.ClosedState preWalQueue sm.Queue @@ -142,8 +187,10 @@ type TxnManager struct { // store all txns store *sync.Map - // wg is used to wait all txns to be done - wg sync.WaitGroup + // waiter is used to wait all txns to be done. Unlike sync.WaitGroup, + // it supports cancelling a wait and starting a later transaction + // generation without leaving a blocked waiter behind. + waiter txnWaiter // TxnSkipFlag to skip some txn type // 0: skip nothing @@ -182,7 +229,6 @@ func NewTxnManager( CommitListener: newBatchCommitListener(), } mgr.txns.store = new(sync.Map) - mgr.txns.wg = sync.WaitGroup{} for _, opt := range opts { opt(mgr) } @@ -203,9 +249,29 @@ func (mgr *TxnManager) initMaxCommittedTS() { } func (mgr *TxnManager) TryUpdateMaxCommittedTS(ts types.TS) { - if ts.GT(&MinCommittedTS) { - mgr.MaxCommittedTS.CompareAndSwap(mgr.MaxCommittedTS.Load(), &ts) + for old := mgr.MaxCommittedTS.Load(); ts.GT(old); old = mgr.MaxCommittedTS.Load() { + if mgr.MaxCommittedTS.CompareAndSwap(old, &ts) { + return + } + } +} + +// AllocateAndPublishCommitTS serializes timestamp allocation with publishing +// the state committed at that timestamp. The publisher must make the state +// visible before returning so a later transaction timestamp cannot pass state +// that has not been published yet. +func (mgr *TxnManager) AllocateAndPublishCommitTS( + publish func(types.TS) error, +) (ts types.TS, err error) { + mgr.ts.mu.Lock() + defer mgr.ts.mu.Unlock() + + ts = mgr.ts.allocator.Alloc() + if err = publish(ts); err != nil { + return } + mgr.TryUpdateMaxCommittedTS(ts) + return } // Now gets a timestamp under the protect from a inner lock. The lock makes @@ -324,17 +390,7 @@ func (mgr *TxnManager) StartTxnWithStartTSAndSnapshotTS( } func (mgr *TxnManager) WaitEmpty(ctx context.Context) (err error) { - c := make(chan struct{}) - go func() { - mgr.txns.wg.Wait() - close(c) - }() - select { - case <-ctx.Done(): - return ctx.Err() - case <-c: - return - } + return mgr.txns.waiter.Wait(ctx) } func (mgr *TxnManager) loadTxn( @@ -350,7 +406,7 @@ func (mgr *TxnManager) loadAndDeleteTxn( id string, ) (txnif.AsyncTxn, bool) { if res, ok := mgr.txns.store.LoadAndDelete(id); ok { - mgr.txns.wg.Done() + mgr.txns.waiter.Done() return res.(txnif.AsyncTxn), true } return nil, false @@ -363,11 +419,11 @@ func (mgr *TxnManager) loadAndDeleteTxn( func (mgr *TxnManager) storeTxn( newTxn txnif.AsyncTxn, flag TxnFlag, ) (offline bool) { - mgr.txns.wg.Add(1) + mgr.txns.waiter.Add() skipFlags := TxnSkipFlag(mgr.txns.skipFlags.Load()) if skipFlags.Skip(flag) { - mgr.txns.wg.Done() + mgr.txns.waiter.Done() offline = true return } @@ -381,13 +437,19 @@ func (mgr *TxnManager) storeTxn( func (mgr *TxnManager) loadOrStoreTxn( newTxn txnif.AsyncTxn, flag TxnFlag, ) (retTxn txnif.AsyncTxn, loaded bool, offline bool) { - mgr.txns.wg.Add(1) + mgr.txns.waiter.Add() skipFlags := TxnSkipFlag(mgr.txns.skipFlags.Load()) if skipFlags.Skip(flag) { - mgr.txns.wg.Done() - offline = true + mgr.txns.waiter.Done() + if actual, ok := mgr.txns.store.Load(newTxn.GetID()); ok { + retTxn = actual.(txnif.AsyncTxn) + loaded = true + offline = retTxn.GetStore().IsOffline() + return + } retTxn = newTxn + offline = true return } @@ -395,7 +457,7 @@ func (mgr *TxnManager) loadOrStoreTxn( newTxn.GetID(), newTxn, ) if loaded { - mgr.txns.wg.Done() + mgr.txns.waiter.Done() retTxn = actual.(txnif.AsyncTxn) offline = retTxn.GetStore().IsOffline() } else { diff --git a/pkg/vm/engine/tae/txn/txnbase/txnmgr_test.go b/pkg/vm/engine/tae/txn/txnbase/txnmgr_test.go index 59af0513985cd..c254c486a95b6 100644 --- a/pkg/vm/engine/tae/txn/txnbase/txnmgr_test.go +++ b/pkg/vm/engine/tae/txn/txnbase/txnmgr_test.go @@ -4,7 +4,7 @@ // 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 +// 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, @@ -16,6 +16,7 @@ package txnbase import ( "context" + "errors" "sync" "testing" "time" @@ -25,6 +26,194 @@ import ( "github.com/stretchr/testify/require" ) +func TestTxnWaiterCancelAndReuse(t *testing.T) { + var waiter txnWaiter + waiter.Add() + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + require.ErrorIs(t, waiter.Wait(ctx), context.Canceled) + + firstWait := make(chan error, 1) + go func() { + firstWait <- waiter.Wait(context.Background()) + }() + waiter.Done() + require.NoError(t, <-firstWait) + + // A canceled wait must not leave a sync.WaitGroup-style waiter that makes + // the next transaction generation unsafe to start. + waiter.Add() + secondWait := make(chan error, 1) + go func() { + secondWait <- waiter.Wait(context.Background()) + }() + select { + case <-secondWait: + t.Fatal("new transaction generation reported empty before completion") + case <-time.After(20 * time.Millisecond): + } + waiter.Done() + require.NoError(t, <-secondWait) +} + +func TestLoadOrStoreTxnBalancesWaiterOnHit(t *testing.T) { + mgr := &TxnManager{} + mgr.txns.store = new(sync.Map) + startTS := types.BuildTS(1, 0) + id := []byte("same-txn") + + first := NewTxn(mgr, new(NoopTxnStore), id, startTS, types.TS{}) + stored, loaded, offline := mgr.loadOrStoreTxn(first, TxnFlag_Normal) + require.False(t, loaded) + require.False(t, offline) + require.Same(t, first, stored) + + duplicate := NewTxn(mgr, new(NoopTxnStore), id, startTS, types.TS{}) + stored, loaded, offline = mgr.loadOrStoreTxn(duplicate, TxnFlag_Normal) + require.True(t, loaded) + require.False(t, offline) + require.Same(t, first, stored) + + _, ok := mgr.loadAndDeleteTxn(first.GetID()) + require.True(t, ok) + require.NoError(t, mgr.WaitEmpty(context.Background())) +} + +func TestLoadOrStoreTxnRejectsNewTxnInReadonlyMode(t *testing.T) { + newManager := func() *TxnManager { + mgr := &TxnManager{} + mgr.txns.store = new(sync.Map) + return mgr + } + newTxn := func(mgr *TxnManager, id []byte) *Txn { + return NewTxn(mgr, new(NoopTxnStore), id, types.BuildTS(1, 0), types.TS{}) + } + + t.Run("new ID is not managed", func(t *testing.T) { + mgr := newManager() + WithTxnSkipFlag(TxnFlag_Normal)(mgr) + id := []byte("readonly-new-txn") + + txn, loaded, offline := mgr.loadOrStoreTxn(newTxn(mgr, id), TxnFlag_Normal) + require.False(t, loaded) + require.True(t, offline) + require.Equal(t, id, []byte(txn.GetID())) + _, ok := mgr.loadTxn(txn.GetID()) + require.False(t, ok) + + ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + require.NoError(t, mgr.WaitEmpty(ctx)) + }) + + t.Run("managed ID is reused", func(t *testing.T) { + mgr := newManager() + id := []byte("readonly-managed-txn") + managed := newTxn(mgr, id) + stored, loaded, offline := mgr.loadOrStoreTxn(managed, TxnFlag_Normal) + require.False(t, loaded) + require.False(t, offline) + require.Same(t, managed, stored) + + WithTxnSkipFlag(TxnFlag_Normal)(mgr) + stored, loaded, offline = mgr.loadOrStoreTxn(newTxn(mgr, id), TxnFlag_Normal) + require.True(t, loaded) + require.False(t, offline) + require.Same(t, managed, stored) + + _, ok := mgr.loadAndDeleteTxn(managed.GetID()) + require.True(t, ok) + require.NoError(t, mgr.WaitEmpty(context.Background())) + }) +} + +func TestTryUpdateMaxCommittedTSNeverMovesBackward(t *testing.T) { + mgr := &TxnManager{} + mgr.initMaxCommittedTS() + + newer := types.BuildTS(2, 0) + older := types.BuildTS(1, 0) + mgr.TryUpdateMaxCommittedTS(newer) + mgr.TryUpdateMaxCommittedTS(older) + + require.Equal(t, newer, *mgr.MaxCommittedTS.Load()) +} + +func TestTryUpdateMaxCommittedTSConcurrent(t *testing.T) { + mgr := &TxnManager{} + mgr.initMaxCommittedTS() + + const updates = 100 + var wg sync.WaitGroup + for i := 1; i <= updates; i++ { + ts := types.BuildTS(int64(i), 0) + wg.Add(1) + go func() { + defer wg.Done() + mgr.TryUpdateMaxCommittedTS(ts) + }() + } + wg.Wait() + + require.Equal(t, types.BuildTS(updates, 0), *mgr.MaxCommittedTS.Load()) +} + +func TestAllocateAndPublishCommitTSSerializesPublication(t *testing.T) { + mgr := NewTxnManager(nil, nil, types.NewMockHLCClock(1)) + defer mgr.workers.Release() + + firstStarted := make(chan types.TS, 1) + releaseFirst := make(chan struct{}) + firstDone := make(chan error, 1) + go func() { + _, err := mgr.AllocateAndPublishCommitTS(func(ts types.TS) error { + firstStarted <- ts + <-releaseFirst + return nil + }) + firstDone <- err + }() + + firstTS := <-firstStarted + require.True(t, mgr.MaxCommittedTS.Load().LT(&firstTS)) + + secondStarted := make(chan types.TS, 1) + secondDone := make(chan error, 1) + go func() { + _, err := mgr.AllocateAndPublishCommitTS(func(ts types.TS) error { + secondStarted <- ts + return nil + }) + secondDone <- err + }() + + select { + case <-secondStarted: + t.Fatal("later timestamp allocated before earlier state was published") + case <-time.After(20 * time.Millisecond): + } + + close(releaseFirst) + require.NoError(t, <-firstDone) + secondTS := <-secondStarted + require.True(t, secondTS.GT(&firstTS)) + require.NoError(t, <-secondDone) + require.Equal(t, secondTS, *mgr.MaxCommittedTS.Load()) +} + +func TestAllocateAndPublishCommitTSErrorDoesNotPublish(t *testing.T) { + mgr := NewTxnManager(nil, nil, types.NewMockHLCClock(1)) + defer mgr.workers.Release() + + publishErr := errors.New("publish failed") + ts, err := mgr.AllocateAndPublishCommitTS(func(types.TS) error { + return publishErr + }) + require.ErrorIs(t, err, publishErr) + require.True(t, mgr.MaxCommittedTS.Load().LT(&ts)) +} + func newTxnManagerForLifecycleTest() *TxnManager { mgr := &TxnManager{} mgr.txns.store = new(sync.Map) @@ -42,7 +231,7 @@ func waitTxnManagerEmpty(t *testing.T, mgr *TxnManager) { require.NoError(t, mgr.WaitEmpty(ctx)) } -func TestLoadOrStoreTxnBalancesWaitGroupWhenLoaded(t *testing.T) { +func TestLoadOrStoreTxnBalancesLifecycleWaiterWhenLoaded(t *testing.T) { mgr := newTxnManagerForLifecycleTest() first, loaded, offline := mgr.loadOrStoreTxn( newTxnForLifecycleTest(mgr, "txn"), TxnFlag_Normal) diff --git a/pkg/vm/engine/tae/txn/txnimpl/base_table.go b/pkg/vm/engine/tae/txn/txnimpl/base_table.go index c78e0c3e6dc14..7e747f7d871f7 100644 --- a/pkg/vm/engine/tae/txn/txnimpl/base_table.go +++ b/pkg/vm/engine/tae/txn/txnimpl/base_table.go @@ -176,6 +176,9 @@ func (tbl *baseTable) getRowsByPK(ctx context.Context, pks containers.Vector) (r } for it.Next() { obj := it.Item() + if isEmptyDroppedAppendableObject(obj) { + continue + } objData := obj.GetObjectData() if objData == nil { continue @@ -233,6 +236,9 @@ func (tbl *baseTable) incrementalGetRowsByPK(ctx context.Context, pks containers break } obj := objIt.Item() + if isEmptyDroppedAppendableObject(obj) { + continue + } if obj.CreatedAt.GT(&to) { continue @@ -284,6 +290,29 @@ func (tbl *baseTable) incrementalGetRowsByPK(ctx context.Context, pks containers return } +func isEmptyDroppedAppendableObject(obj *catalog.ObjectEntry) bool { + stats := obj.GetObjectStats() + if !obj.IsAppendable() || stats.Rows() != 0 || stats.BlkCnt() != 0 { + return false + } + dropCommitted := obj.HasDropCommitted() + if !dropCommitted && obj.IsCEntry() && obj.HasDCounterpart() { + dropCommitted = obj.GetNextVersion().HasDropCommitted() + } + if !dropCommitted { + return false + } + objData := obj.GetObjectData() + if objData == nil { + return false + } + rows, err := objData.Rows() + if err != nil || rows != 0 { + return false + } + return true +} + func (tbl *baseTable) CleanUp() { if tbl.tableSpace != nil { tbl.tableSpace.CloseAppends() diff --git a/pkg/vm/engine/tae/txn/txnimpl/relation.go b/pkg/vm/engine/tae/txn/txnimpl/relation.go index b358377d215cf..236b28278462f 100644 --- a/pkg/vm/engine/tae/txn/txnimpl/relation.go +++ b/pkg/vm/engine/tae/txn/txnimpl/relation.go @@ -189,6 +189,13 @@ func (h *txnRelation) CreateObject(isTombstone bool) (obj handle.Object, err err return h.Txn.GetStore().CreateObject(h.table.entry.GetDB().ID, h.table.entry.GetID(), isTombstone) } +func (h *txnRelation) CreateObjectWithOpt(isTombstone bool, opt *objectio.CreateObjOpt) (obj handle.Object, err error) { + if err = validateCreateObjectOpt(opt); err != nil { + return + } + return h.Txn.GetStore().CreateObjectWithOpt(h.table.entry.GetDB().ID, h.table.entry.GetID(), isTombstone, opt) +} + func (h *txnRelation) CreateNonAppendableObject(isTombstone bool, opt *objectio.CreateObjOpt) (obj handle.Object, err error) { if opt == nil { noid := objectio.NewObjectid() diff --git a/pkg/vm/engine/tae/txn/txnimpl/replaystore.go b/pkg/vm/engine/tae/txn/txnimpl/replaystore.go index e94079d60db5b..644a9c1d97076 100644 --- a/pkg/vm/engine/tae/txn/txnimpl/replaystore.go +++ b/pkg/vm/engine/tae/txn/txnimpl/replaystore.go @@ -20,6 +20,7 @@ import ( "github.com/matrixorigin/matrixone/pkg/common/moerr" "github.com/matrixorigin/matrixone/pkg/container/types" "github.com/matrixorigin/matrixone/pkg/logutil" + "github.com/matrixorigin/matrixone/pkg/objectio" "github.com/matrixorigin/matrixone/pkg/util/fault" "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/catalog" "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/common" @@ -31,6 +32,10 @@ import ( var ErrDebugReplay = moerr.NewInternalErrorNoCtx("debug") +type replayAObjectCreateObserver interface { + RecordReplayAObjectCreate(id *common.ID, isTombstone bool, ts types.TS) +} + type replayTxnStore struct { txnbase.NoopTxnStore Cmd *txnbase.TxnCmd @@ -163,17 +168,7 @@ func (store *replayTxnStore) replayAppendData(cmd *AppendCmd, observer wal.Repla _, sarg, _ := fault.TriggerFault("replay debug log") for _, info := range cmd.Infos { id := info.GetDest() - database, err := store.catalog.GetDatabaseByID(id.DbID) - if sarg != "" { - err = ErrDebugReplay - } - if err != nil { - logutil.Infof("cmd %v\ncatalog: %v", cmd.String(), store.catalog.SimplePPString(3)) - if err != ErrDebugReplay { - panic(err) - } - } - blk, err := database.GetObjectEntryByID(id, cmd.IsTombstone) + blk, _, err := store.ensureReplayAObject(id, cmd.IsTombstone, cmd.Ts, observer) if sarg != "" { err = ErrDebugReplay } @@ -203,17 +198,7 @@ func (store *replayTxnStore) replayAppendData(cmd *AppendCmd, observer wal.Repla for _, info := range cmd.Infos { id := info.GetDest() - database, err := store.catalog.GetDatabaseByID(id.DbID) - if sarg != "" { - err = ErrDebugReplay - } - if err != nil { - logutil.Infof("cmd %v\ncatalog: %v", cmd.String(), store.catalog.SimplePPString(3)) - if err != ErrDebugReplay { - panic(err) - } - } - blk, err := database.GetObjectEntryByID(id, cmd.IsTombstone) + blk, _, err := store.ensureReplayAObject(id, cmd.IsTombstone, cmd.Ts, observer) if sarg != "" { err = ErrDebugReplay } @@ -254,18 +239,8 @@ func (store *replayTxnStore) replayDataCmds(cmd *updates.UpdateCmd, observer wal func (store *replayTxnStore) replayAppend(cmd *updates.UpdateCmd, observer wal.ReplayObserver) { appendNode := cmd.GetAppendNode() id := appendNode.GetID() - database, err := store.catalog.GetDatabaseByID(id.DbID) _, sarg, _ := fault.TriggerFault("replay debug log") - if sarg != "" { - err = ErrDebugReplay - } - if err != nil { - logutil.Infof("cmd %v\ncatalog: %v", cmd.String(), store.catalog.SimplePPString(3)) - if err != ErrDebugReplay { - panic(err) - } - } - obj, err := database.GetObjectEntryByID(id, cmd.GetAppendNode().IsTombstone()) + obj, _, err := store.ensureReplayAObject(id, appendNode.IsTombstone(), replayAppendNodeCreateTS(appendNode), observer) if sarg != "" { err = ErrDebugReplay } @@ -288,3 +263,54 @@ func (store *replayTxnStore) replayAppend(cmd *updates.UpdateCmd, observer wal.R } } } + +func replayAppendNodeCreateTS(node *updates.AppendNode) types.TS { + if txn := node.GetTxn(); txn != nil { + return txn.GetCommitTS() + } + if ts := node.GetCommitTS(); !ts.IsEmpty() { + return ts + } + return node.GetPrepare() +} + +func (store *replayTxnStore) ensureReplayAObject( + id *common.ID, + isTombstone bool, + createHint types.TS, + observer wal.ReplayObserver, +) (obj *catalog.ObjectEntry, created bool, err error) { + database, err := store.catalog.GetDatabaseByID(id.DbID) + if err != nil { + return + } + obj, err = database.GetObjectEntryByID(id, isTombstone) + if err == nil { + return + } + if !moerr.IsMoErrCode(err, moerr.OkExpectedEOB) { + return + } + table, err := database.GetTableEntryByID(id.TableID) + if err != nil { + return + } + stats := objectio.NewObjectStatsWithObjectID(id.ObjectID(), true, isTombstone, false) + var factory catalog.ObjectDataFactory + if store.catalog.DataFactory != nil { + factory = store.catalog.DataFactory.MakeObjectFactory() + } + obj, err = table.CreateCommittedObject( + createHint, + &objectio.CreateObjOpt{Stats: stats, IsTombstone: isTombstone}, + factory, + ) + if err != nil { + return + } + created = true + if recorder, ok := observer.(replayAObjectCreateObserver); ok { + recorder.RecordReplayAObjectCreate(id, isTombstone, createHint) + } + return +} diff --git a/pkg/vm/engine/tae/txn/txnimpl/store.go b/pkg/vm/engine/tae/txn/txnimpl/store.go index 597f9f818bbb1..c1a5a45bba557 100644 --- a/pkg/vm/engine/tae/txn/txnimpl/store.go +++ b/pkg/vm/engine/tae/txn/txnimpl/store.go @@ -674,6 +674,17 @@ func (store *txnStore) CreateObject(dbId, tid uint64, isTombstone bool) (obj han return db.CreateObject(tid, isTombstone) } +func (store *txnStore) CreateObjectWithOpt(dbId, tid uint64, isTombstone bool, opt *objectio.CreateObjOpt) (obj handle.Object, err error) { + if err = store.WantWrite("CreateObjectWithOpt"); err != nil { + return + } + var db *txnDB + if db, err = store.getOrSetDB(dbId); err != nil { + return + } + return db.CreateObjectWithOpt(tid, opt, isTombstone) +} + func (store *txnStore) CreateNonAppendableObject(dbId, tid uint64, isTombstone bool, opt *objectio.CreateObjOpt) (obj handle.Object, err error) { if err = store.WantWrite("CreateNonAppendableObject"); err != nil { return diff --git a/pkg/vm/engine/tae/txn/txnimpl/table.go b/pkg/vm/engine/tae/txn/txnimpl/table.go index 69370fe9eef63..5004f3f41a12e 100644 --- a/pkg/vm/engine/tae/txn/txnimpl/table.go +++ b/pkg/vm/engine/tae/txn/txnimpl/table.go @@ -48,6 +48,7 @@ import ( "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/index/indexwrapper" "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/logstore/wal" "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/model" + "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/txn/txnbase" "go.uber.org/zap" ) @@ -797,15 +798,43 @@ func (tbl *txnTable) CreateObject(isTombstone bool) (obj handle.Object, err erro sorted, false, ) - return tbl.createObject( + return tbl.createCommittedAppendableObject( &objectio.CreateObjOpt{Stats: stats, IsTombstone: isTombstone}, ) } +func (tbl *txnTable) CreateObjectWithOpt(opts *objectio.CreateObjOpt) (obj handle.Object, err error) { + return tbl.createCommittedAppendableObject(opts) +} + func (tbl *txnTable) CreateNonAppendableObject(opts *objectio.CreateObjOpt) (obj handle.Object, err error) { return tbl.createObject(opts) } +func (tbl *txnTable) createCommittedAppendableObject(opts *objectio.CreateObjOpt) (obj handle.Object, err error) { + if !opts.Stats.GetAppendable() { + return nil, moerr.NewInternalErrorNoCtx("CreateObject outside txn only supports appendable object") + } + baseTxn, ok := tbl.store.txn.GetBase().(*txnbase.Txn) + if !ok || baseTxn.Mgr == nil { + return nil, moerr.NewInternalErrorNoCtx("missing txn manager for appendable object create") + } + var factory catalog.ObjectDataFactory + if tbl.store.catalog.DataFactory != nil { + factory = tbl.store.catalog.DataFactory.MakeObjectFactory() + } + var meta *catalog.ObjectEntry + _, err = baseTxn.Mgr.AllocateAndPublishCommitTS(func(createTS types.TS) error { + meta, err = tbl.entry.CreateCommittedObject(createTS, opts, factory) + return err + }) + if err != nil { + return + } + obj = newObject(tbl, meta) + return +} + func (tbl *txnTable) createObject(opts *objectio.CreateObjOpt) (obj handle.Object, err error) { var factory catalog.ObjectDataFactory if tbl.store.catalog.DataFactory != nil { @@ -1540,6 +1569,9 @@ func (tbl *txnTable) PrepareCommit() (err error) { if tbl.txnEntries.IsDeleted(idx) { continue } + if isAppendableObjectCreateEntry(node) { + panic("appendable object create entry must not be committed through txn") + } if err = node.PrepareCommit(); err != nil { if moerr.IsMoErrCode(err, moerr.ErrTxnNotFound) { var buf bytes.Buffer @@ -1568,6 +1600,9 @@ func (tbl *txnTable) PrepareCommit() (err error) { if tbl.txnEntries.IsDeleted(idx) { continue } + if isAppendableObjectCreateEntry(tbl.txnEntries.entries[idx]) { + panic("appendable object create entry must not be committed through txn") + } if err = tbl.txnEntries.entries[idx].PrepareCommit(); err != nil { break } @@ -1586,6 +1621,9 @@ func (tbl *txnTable) ApplyCommit() (err error) { if tbl.txnEntries.IsDeleted(idx) { continue } + if isAppendableObjectCreateEntry(node) { + panic("appendable object create entry must not be committed through txn") + } if err = node.ApplyCommit(tbl.store.txn.GetID()); err != nil { if moerr.IsMoErrCode(err, moerr.ErrTxnNotFound) { var buf bytes.Buffer @@ -1626,6 +1664,11 @@ func (tbl *txnTable) ApplyCommit() (err error) { return } +func isAppendableObjectCreateEntry(entry txnif.TxnEntry) bool { + obj, ok := entry.(*catalog.ObjectEntry) + return ok && obj.ObjectState == catalog.ObjectState_Create_Active && obj.IsAppendable() +} + func (tbl *txnTable) ApplyRollback() (err error) { csn := tbl.csnStart for idx, node := range tbl.txnEntries.entries { diff --git a/pkg/vm/engine/tae/txn/txnimpl/txn_test.go b/pkg/vm/engine/tae/txn/txnimpl/txn_test.go index 8a0c763f2dcb9..3a6db3b92512d 100644 --- a/pkg/vm/engine/tae/txn/txnimpl/txn_test.go +++ b/pkg/vm/engine/tae/txn/txnimpl/txn_test.go @@ -34,10 +34,12 @@ import ( "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/iface/txnif" "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/logstore/wal" "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/tables" + "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/tables/updates" "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/testutils" "github.com/matrixorigin/matrixone/pkg/vm/engine/tae/txn/txnbase" "github.com/panjf2000/ants/v2" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" ) const ( @@ -908,7 +910,7 @@ func TestObject1(t *testing.T) { t.Log(iobj.String()) cnt++ } - assert.Equal(t, 2, cnt) + assert.Equal(t, 1, cnt) txn3, _ := mgr.StartTxn(nil) db, _ = txn3.GetDatabase(name) @@ -920,7 +922,7 @@ func TestObject1(t *testing.T) { t.Log(iobj.String()) cnt++ } - assert.Equal(t, 1, cnt) + assert.Equal(t, 2, cnt) err = txn2.Commit(context.Background()) assert.Nil(t, err) @@ -932,7 +934,7 @@ func TestObject1(t *testing.T) { t.Log(iobj.String()) cnt++ } - assert.Equal(t, 1, cnt) + assert.Equal(t, 2, cnt) } func TestObject2(t *testing.T) { @@ -955,18 +957,95 @@ func TestObject2(t *testing.T) { assert.Nil(t, err) } + err := txn1.Commit(context.Background()) + assert.Nil(t, err) + + txn2, _ := mgr.StartTxn(nil) + db, _ = txn2.GetDatabase("db") + rel, _ = db.GetRelationByName(schema.Name) it := rel.MakeObjectIt(false) cnt := 0 for it.Next() { cnt++ - // iobj := it.GetObject() } assert.Equal(t, objCnt, cnt) - // err := txn1.Commit() - // assert.Nil(t, err) + assert.Nil(t, txn2.Commit(context.Background())) t.Log(c.SimplePPString(common.PPL1)) } +func TestIsEmptyDroppedAppendableObject(t *testing.T) { + defer testutils.AfterTest(t)() + testutils.EnsureNoLeak(t) + + ctx := context.Background() + dir := testutils.InitTestEnv(ModuleName, t) + c, mgr, driver := initTestContext(ctx, t, dir) + defer driver.Close() + defer c.Close() + defer mgr.Stop() + + schema := catalog.MockSchema(1, 0) + txn, err := mgr.StartTxn(nil) + require.NoError(t, err) + db, err := txn.CreateDatabase("db", "", "") + require.NoError(t, err) + _, err = db.CreateRelation(schema) + require.NoError(t, err) + require.NoError(t, txn.Commit(ctx)) + + bat := catalog.MockBatch(schema, 1) + defer bat.Close() + txn, err = mgr.StartTxn(nil) + require.NoError(t, err) + db, err = txn.GetDatabase("db") + require.NoError(t, err) + rel, err := db.GetRelationByName(schema.Name) + require.NoError(t, err) + require.NoError(t, rel.Append(ctx, bat)) + require.NoError(t, txn.Commit(ctx)) + + txn, err = mgr.StartTxn(nil) + require.NoError(t, err) + db, err = txn.GetDatabase("db") + require.NoError(t, err) + rel, err = db.GetRelationByName(schema.Name) + require.NoError(t, err) + it := rel.MakeObjectIt(false) + require.True(t, it.Next()) + dataObjectID := *it.GetObject().GetID() + it.Close() + emptyObject, err := rel.CreateObject(false) + require.NoError(t, err) + emptyObjectID := *emptyObject.GetID() + require.NoError(t, emptyObject.Close()) + require.NoError(t, txn.Commit(ctx)) + + txn, err = mgr.StartTxn(nil) + require.NoError(t, err) + db, err = txn.GetDatabase("db") + require.NoError(t, err) + rel, err = db.GetRelationByName(schema.Name) + require.NoError(t, err) + require.NoError(t, rel.SoftDeleteObject(&dataObjectID, false)) + require.NoError(t, rel.SoftDeleteObject(&emptyObjectID, false)) + require.NoError(t, txn.Commit(ctx)) + + table := rel.GetMeta().(*catalog.TableEntry) + dataObject, err := table.GetObjectByID(&dataObjectID, false) + require.NoError(t, err) + require.True(t, dataObject.HasDropCommitted()) + require.Zero(t, dataObject.GetObjectStats().Rows()) + rows, err := dataObject.GetObjectData().Rows() + require.NoError(t, err) + require.Equal(t, 1, rows) + require.False(t, isEmptyDroppedAppendableObject(dataObject)) + + emptyObjectMeta, err := table.GetObjectByID(&emptyObjectID, false) + require.NoError(t, err) + require.True(t, emptyObjectMeta.HasDropCommitted()) + require.True(t, isEmptyDroppedAppendableObject(emptyObjectMeta)) +} + func TestDedup1(t *testing.T) { defer testutils.AfterTest(t)() ctx := context.Background() @@ -1051,3 +1130,236 @@ func TestDedup1(t *testing.T) { } t.Log(c.SimplePPString(common.PPL1)) } +func TestCreateAppendableObjectWithOptions(t *testing.T) { + defer testutils.AfterTest(t)() + testutils.EnsureNoLeak(t) + + ctx := context.Background() + dir := testutils.InitTestEnv(ModuleName, t) + c, mgr, driver := initTestContext(ctx, t, dir) + defer driver.Close() + defer c.Close() + defer mgr.Stop() + + schema := catalog.MockSchema(1, 0) + txn, err := mgr.StartTxn(nil) + assert.NoError(t, err) + db, err := txn.CreateDatabase("db", "", "") + assert.NoError(t, err) + _, err = db.CreateRelation(schema) + assert.NoError(t, err) + assert.NoError(t, txn.Commit(ctx)) + + txn, err = mgr.StartTxn(nil) + assert.NoError(t, err) + db, err = txn.GetDatabase("db") + assert.NoError(t, err) + rel, err := db.GetRelationByName(schema.Name) + assert.NoError(t, err) + + id := objectio.NewObjectid() + stats := objectio.NewObjectStatsWithObjectID(&id, true, false, false) + obj, err := rel.CreateObjectWithOpt(false, &objectio.CreateObjOpt{Stats: stats}) + assert.NoError(t, err) + assert.True(t, obj.GetMeta().(*catalog.ObjectEntry).IsAppendable()) + assert.Equal(t, id, *obj.GetID()) + + store := txn.GetStore().(*txnStore) + txnDB, err := store.getOrSetDB(db.GetID()) + assert.NoError(t, err) + txnTable, err := txnDB.getOrSetTable(rel.ID()) + assert.NoError(t, err) + assert.Zero(t, txnTable.txnEntries.Len()) + assert.NoError(t, txn.Commit(ctx)) + + txn, err = mgr.StartTxn(nil) + assert.NoError(t, err) + db, err = txn.GetDatabase("db") + assert.NoError(t, err) + rel, err = db.GetRelationByName(schema.Name) + assert.NoError(t, err) + _, err = rel.GetObject(&id, false) + assert.NoError(t, err) + assert.NoError(t, txn.Commit(ctx)) +} + +func TestCreateAppendableObjectWithOptionsRejectsInvalidOptions(t *testing.T) { + defer testutils.AfterTest(t)() + testutils.EnsureNoLeak(t) + + ctx := context.Background() + dir := testutils.InitTestEnv(ModuleName, t) + c, mgr, driver := initTestContext(ctx, t, dir) + defer driver.Close() + defer c.Close() + defer mgr.Stop() + + txn, err := mgr.StartTxn(nil) + assert.NoError(t, err) + db, err := txn.CreateDatabase("db", "", "") + assert.NoError(t, err) + rel, err := db.CreateRelation(catalog.MockSchema(1, 0)) + assert.NoError(t, err) + + _, err = rel.CreateObjectWithOpt(false, nil) + assert.True(t, moerr.IsMoErrCode(err, moerr.ErrInvalidInput), err) + + _, err = rel.CreateObjectWithOpt(false, &objectio.CreateObjOpt{}) + assert.True(t, moerr.IsMoErrCode(err, moerr.ErrInvalidInput), err) + + id := objectio.NewObjectid() + stats := objectio.NewObjectStatsWithObjectID(&id, false, false, false) + _, err = rel.CreateObjectWithOpt(false, &objectio.CreateObjOpt{Stats: stats}) + assert.Error(t, err) + assert.Contains(t, err.Error(), "only supports appendable object") + assert.NoError(t, txn.Rollback(ctx)) +} + +func TestCreateAppendableObjectWithOptionsErrors(t *testing.T) { + defer testutils.AfterTest(t)() + testutils.EnsureNoLeak(t) + + ctx := context.Background() + dir := testutils.InitTestEnv(ModuleName, t) + c, mgr, driver := initTestContext(ctx, t, dir) + defer driver.Close() + defer c.Close() + defer mgr.Stop() + + txn, err := mgr.StartTxn(nil) + assert.NoError(t, err) + db, err := txn.CreateDatabase("db", "", "") + assert.NoError(t, err) + rel, err := db.CreateRelation(catalog.MockSchema(1, 0)) + assert.NoError(t, err) + + store := txn.GetStore().(*txnStore) + newOpt := func() *objectio.CreateObjOpt { + id := objectio.NewObjectid() + return &objectio.CreateObjOpt{ + Stats: objectio.NewObjectStatsWithObjectID(&id, true, false, false), + } + } + + _, err = store.CreateObjectWithOpt(db.GetID()+1, rel.ID(), false, newOpt()) + assert.Error(t, err) + + txnDB, err := store.getOrSetDB(db.GetID()) + assert.NoError(t, err) + _, err = store.CreateObjectWithOpt(db.GetID(), rel.ID(), false, nil) + assert.True(t, moerr.IsMoErrCode(err, moerr.ErrInvalidInput), err) + _, err = txnDB.CreateObjectWithOpt(rel.ID(), nil, false) + assert.True(t, moerr.IsMoErrCode(err, moerr.ErrInvalidInput), err) + _, err = txnDB.CreateObjectWithOpt(rel.ID()+1, newOpt(), false) + assert.Error(t, err) + + store.isOffline = true + _, err = store.CreateObjectWithOpt(db.GetID(), rel.ID(), false, newOpt()) + assert.Error(t, err) + _, err = txnDB.CreateObjectWithOpt(rel.ID(), newOpt(), false) + assert.Error(t, err) + store.isOffline = false + assert.NoError(t, txn.Rollback(ctx)) +} + +type replayAObjectCreateRecorder struct { + created []replayAObjectCreateRecord +} + +type replayAObjectCreateRecord struct { + id common.ID + isTombstone bool + ts types.TS +} + +func (*replayAObjectCreateRecorder) OnTimeStamp(types.TS) {} + +func (r *replayAObjectCreateRecorder) RecordReplayAObjectCreate( + id *common.ID, + isTombstone bool, + ts types.TS, +) { + r.created = append(r.created, replayAObjectCreateRecord{*id, isTombstone, ts}) +} + +func TestEnsureReplayAObject(t *testing.T) { + defer testutils.AfterTest(t)() + testutils.EnsureNoLeak(t) + + ctx := context.Background() + dir := testutils.InitTestEnv(ModuleName, t) + c, mgr, driver := initTestContext(ctx, t, dir) + defer driver.Close() + defer c.Close() + defer mgr.Stop() + + txn, err := mgr.StartTxn(nil) + assert.NoError(t, err) + db, err := txn.CreateDatabase("db", "", "") + assert.NoError(t, err) + rel, err := db.CreateRelation(catalog.MockSchema(1, 0)) + assert.NoError(t, err) + assert.NoError(t, txn.Commit(ctx)) + + database, err := c.GetDatabaseByID(db.GetID()) + assert.NoError(t, err) + table, err := database.GetTableEntryByID(rel.ID()) + assert.NoError(t, err) + id := table.AsCommonID() + objectID := objectio.NewObjectid() + id.SetObjectID(&objectID) + + createTS := mgr.Now() + recorder := new(replayAObjectCreateRecorder) + store := &replayTxnStore{catalog: c} + obj, created, err := store.ensureReplayAObject(id, false, createTS, recorder) + assert.NoError(t, err) + assert.True(t, created) + assert.True(t, obj.IsAppendable()) + assert.Equal(t, createTS, obj.GetCreatedAt()) + assert.Len(t, recorder.created, 1) + assert.Equal(t, *id, recorder.created[0].id) + assert.Equal(t, createTS, recorder.created[0].ts) + + obj, created, err = store.ensureReplayAObject(id, false, mgr.Now(), recorder) + assert.NoError(t, err) + assert.False(t, created) + assert.True(t, obj.IsAppendable()) + assert.Len(t, recorder.created, 1) + + tombstoneID := table.AsCommonID() + objectID = objectio.NewObjectid() + tombstoneID.SetObjectID(&objectID) + obj, created, err = store.ensureReplayAObject(tombstoneID, true, mgr.Now(), nil) + assert.NoError(t, err) + assert.True(t, created) + assert.True(t, obj.IsTombstone) + + missingDB := *id + missingDB.DbID++ + _, _, err = store.ensureReplayAObject(&missingDB, false, mgr.Now(), recorder) + assert.Error(t, err) + + missingTable := *id + missingTable.TableID++ + _, _, err = store.ensureReplayAObject(&missingTable, false, mgr.Now(), recorder) + assert.Error(t, err) +} + +func TestReplayAppendNodeCreateTS(t *testing.T) { + prepareTS := types.BuildTS(10, 0) + commitTS := types.BuildTS(20, 0) + node := updates.NewEmptyAppendNode() + node.TxnMVCCNode.Prepare = prepareTS + node.TxnMVCCNode.End = commitTS + assert.Equal(t, commitTS, replayAppendNodeCreateTS(node)) + + txn := txnbase.NewTxn(nil, &txnbase.NoopTxnStore{}, []byte("txn"), prepareTS, prepareTS) + assert.NoError(t, txn.SetCommitTS(types.BuildTS(30, 0))) + node.TxnMVCCNode.Txn = txn + assert.Equal(t, types.BuildTS(30, 0), replayAppendNodeCreateTS(node)) + + node.TxnMVCCNode.Txn = nil + node.TxnMVCCNode.End = types.TS{} + assert.Equal(t, prepareTS, replayAppendNodeCreateTS(node)) +} diff --git a/pkg/vm/engine/tae/txn/txnimpl/txndb.go b/pkg/vm/engine/tae/txn/txnimpl/txndb.go index e986f3d8e321a..d78dd18c9ff9c 100644 --- a/pkg/vm/engine/tae/txn/txnimpl/txndb.go +++ b/pkg/vm/engine/tae/txn/txnimpl/txndb.go @@ -344,6 +344,32 @@ func (db *txnDB) CreateObject(tid uint64, isTombstone bool) (obj handle.Object, } return table.CreateObject(isTombstone) } + +func validateCreateObjectOpt(opt *objectio.CreateObjOpt) error { + if opt == nil { + return moerr.NewInvalidInputNoCtx("CreateObjectWithOpt requires non-nil options") + } + if opt.Stats == nil { + return moerr.NewInvalidInputNoCtx("CreateObjectWithOpt requires object stats") + } + return nil +} + +func (db *txnDB) CreateObjectWithOpt(tid uint64, opt *objectio.CreateObjOpt, isTombstone bool) (obj handle.Object, err error) { + if err = validateCreateObjectOpt(opt); err != nil { + return + } + if err = db.store.WantWrite("CreateObjectWithOpt"); err != nil { + return + } + var table *txnTable + if table, err = db.getOrSetTable(tid); err != nil { + return + } + opt.WithIsTombstone(isTombstone) + return table.CreateObjectWithOpt(opt) +} + func (db *txnDB) CreateNonAppendableObject(tid uint64, opt *objectio.CreateObjOpt, isTombstone bool) (obj handle.Object, err error) { if err = db.store.WantWrite("CreateNonAppendableObject"); err != nil { return diff --git a/pkg/vm/engine/test/partition_state_test.go b/pkg/vm/engine/test/partition_state_test.go index cd8cff067d9d2..74f878a3614d4 100644 --- a/pkg/vm/engine/test/partition_state_test.go +++ b/pkg/vm/engine/test/partition_state_test.go @@ -143,8 +143,10 @@ func Test_Append(t *testing.T) { } objIt.Close() - assert.Equal(t, expectObjCnt, objCnt) - assert.Equal(t, expectBlkCnt, blkCnt) + // Appendable objects are created committed outside the creating + // transaction, so the explicitly created empty object is visible here. + assert.Equal(t, expectObjCnt+1, objCnt) + assert.Equal(t, expectBlkCnt+1, blkCnt) } t.Log(taeEngine.GetDB().Catalog.SimplePPString(common.PPL1))