Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
a98450b
fix(disttae): retain unpublished S3 cleanup ownership
XuPeng-SH Sep 23, 2026
751d0a1
fix(multi-update): release completed spill writers promptly
XuPeng-SH Sep 23, 2026
adab276
fix(multi-update): hand off spill object ownership safely
XuPeng-SH Sep 23, 2026
bd8d7c5
fix: preserve remote SQL errors after S3 cleanup
XuPeng-SH Sep 23, 2026
f72048b
fix: retain S3 cleanup ownership through remote handoff
XuPeng-SH Sep 23, 2026
18c7963
test: cover S3 cleanup ownership retries
XuPeng-SH Sep 24, 2026
81e69aa
fix: retry remote S3 cleanup after stream exit
XuPeng-SH Sep 24, 2026
5f39dff
fix: retain local S3 cleanup through rollback
XuPeng-SH Sep 24, 2026
0c3837e
fix: preserve S3 cleanup retries through CN shutdown
XuPeng-SH Sep 24, 2026
de4c539
fix: retain prepared cleanup ownership and share teardown budget
XuPeng-SH Sep 24, 2026
6648d47
docs: propose bounded unpublished S3 cleanup admission
XuPeng-SH Sep 24, 2026
fcfff3a
docs: reconcile bounded S3 cleanup design with lightweight retry shells
XuPeng-SH Sep 24, 2026
c7473f8
docs: record S3 cleanup telemetry and rollout alert thresholds
XuPeng-SH Sep 24, 2026
94f8f10
fix: bound unpublished CN S3 cleanup before upload
XuPeng-SH Sep 24, 2026
fd08bc9
test(compile): register S3 admission service for remote takeover
XuPeng-SH Sep 24, 2026
d615974
build(proto): regenerate pipeline bindings after rebase
XuPeng-SH Sep 24, 2026
4932683
fix: make unpublished S3 cleanup progress monotonic
XuPeng-SH Sep 25, 2026
4be21e9
test: synchronize immediate unpublished S3 retries
XuPeng-SH Sep 25, 2026
c7aa7b3
fix: keep partitioned multi-update writers on coordinator
XuPeng-SH Sep 25, 2026
d8352ec
Merge branch 'main' into codex/issue-29257-s3-orphan
XuPeng-SH Sep 25, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
364 changes: 364 additions & 0 deletions docs/design/unpublished_s3_object_ownership.md

Large diffs are not rendered by default.

29 changes: 21 additions & 8 deletions pkg/cnservice/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -544,6 +544,8 @@ func (s *service) closeService() error {
// otherwise retirement can wait for the work we have not stopped yet.
s.closeSiriusRuntime,
s.closeMongoDBRuntime,
// All transaction producers must stop before cleanup admission closes.
func() error { return s.colexecServer.CloseUnpublishedS3Cleanup(context.Background()) },
)
if s.closeErr != nil {
return
Expand Down Expand Up @@ -860,14 +862,7 @@ func (s *service) handleRequest(
}
}

// start a goroutine to handle one received message.
owned = false
cancelOwned = false
go func() {
defer release()
if value.Cancel != nil {
defer value.Cancel()
}
invoke := func() {
s.pipelines.counter.Add(1)
defer s.pipelines.counter.Add(-1)

Expand All @@ -885,6 +880,24 @@ func (s *service) handleRequest(
s._txnClient,
s.aicm,
s.acquireMessage)
}
// The connection read loop calls handleRequest in wire order. Ownership
// receipts are per batch, so ACKs must not be processed out of order by
// separate handler goroutines.
if msg.GetCmd() == pipeline.Method_PipelineBatchAck {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ACKs now execute on the connection read loop. Correct for ordering, but any blocking in the ACK path now stalls every stream on this connection — worth a comment pinning that invariant (lock-only, no channel sends / waits).

invoke()
return nil
}

// start a goroutine to handle one received message.
owned = false
cancelOwned = false
go func() {
defer release()
if value.Cancel != nil {
defer value.Cancel()
}
invoke()
}()
return nil
}
Expand Down
86 changes: 86 additions & 0 deletions pkg/cnservice/server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ import (
"github.com/matrixorigin/matrixone/pkg/pb/timestamp"
"github.com/matrixorigin/matrixone/pkg/queryservice"
qclient "github.com/matrixorigin/matrixone/pkg/queryservice/client"
"github.com/matrixorigin/matrixone/pkg/sql/colexec"
"github.com/matrixorigin/matrixone/pkg/sql/compile"
"github.com/matrixorigin/matrixone/pkg/txn/client"
"github.com/matrixorigin/matrixone/pkg/txn/trace"
Expand Down Expand Up @@ -325,6 +326,42 @@ func TestServiceCloseWithdrawalErrorIsLocallyComplete(t *testing.T) {
}
}

func TestServiceCloseCompletesAfterTransientS3CleanupFailure(t *testing.T) {
moruntime.RunTest(t.Name(), func(rt moruntime.Runtime) {
withdrawalErr := errors.New("withdrawal failed")
hc := &failingWithdrawalHeartbeatClient{testHAKClient: &testHAKClient{}, err: withdrawalErr}
mc := clusterservice.NewMOCluster(t.Name(), hc, time.Hour)
t.Cleanup(mc.Close)
ctrl := gomock.NewController(t)
ls := mock_lock.NewMockLockService(ctrl)
ls.EXPECT().Close().Return(nil).Times(2)
executor := colexec.NewServer(t.Name())
var attempts atomic.Int32
cleanup := func(context.Context) error {
if attempts.Add(1) == 1 {
return errors.New("temporary Delete failure")
}
return nil
}
sv := &service{
cfg: &Config{UUID: t.Name()}, logger: zap.NewNop(), config: util.NewConfigData(nil),
stopper: stopper.NewStopper(t.Name()), bootstrapService: &testBootService{},
mo: closeErrorMOServer{}, _hakeeperClient: hc, moCluster: mc,
server: closeOnlyRPCServer{}, lockService: ls, colexecServer: executor,
incrservice: closeOnlyIncrService{onClose: func() {
require.NoError(t, executor.RetryUnpublishedS3Cleanup(cleanup))
}},
viewMetadataAdmissionGeneration: 1,
}
// Remote withdrawal is diagnostic; it must not hide whether local
// teardown advanced past the recovered S3 cleanup and closed its tail.
require.ErrorIs(t, sv.Close(), withdrawalErr)
require.True(t, sv.CloseComplete())
require.Equal(t, int32(2), attempts.Load())
require.Equal(t, 1, hc.closed)
})
}

func TestMakeRSSCacheEvictorEvictsMemoryCacheOnly(t *testing.T) {
oldMemoryEvictor := evictMemoryCachesToCapacityPercent
defer func() {
Expand Down Expand Up @@ -1541,6 +1578,55 @@ func TestPipelineAdmissionRejectCancelsRequestOnce(t *testing.T) {
require.Equal(t, int32(1), cancelCount.Load())
}

func TestPipelineBatchAckRunsBeforeIngressReturns(t *testing.T) {
started := make(chan struct{})
releaseHandler := make(chan struct{})
var releaseOnce sync.Once
t.Cleanup(func() { releaseOnce.Do(func() { close(releaseHandler) }) })
var cancelCount atomic.Int32
s := &service{cfg: &Config{}}
s.requestHandler = func(
_ context.Context,
_ string,
_ morpc.Message,
_ morpc.ClientSession,
_ engine.Engine,
_ fileservice.FileService,
_ lockservice.LockService,
_ qclient.QueryClient,
_ logservice.CNHAKeeperClient,
_ udf.Service,
_ client.TxnClient,
_ *defines.AutoIncrCacheManager,
_ func() morpc.Message,
) error {
close(started)
<-releaseHandler
return nil
}
returned := make(chan error, 1)
go func() {
returned <- s.handleRequest(context.Background(), morpc.RPCMessage{
Message: &pipeline.Message{Cmd: pipeline.Method_PipelineBatchAck, Sid: pipeline.Status_Last},
Cancel: func() { cancelCount.Add(1) },
}, 0, nil)
}()
select {
case <-started:
case <-time.After(time.Second):
t.Fatal("ACK handler did not start")
}
select {
case <-returned:
t.Fatal("ingress returned before ACK processing completed")
case <-time.After(50 * time.Millisecond):
}
releaseOnce.Do(func() { close(releaseHandler) })
require.NoError(t, <-returned)
require.Equal(t, int32(1), cancelCount.Load())
require.Zero(t, s.pipelines.counter.Load())
}

func TestHandleRequestPropagatesConfiguredRPCMaxMessageSize(t *testing.T) {
const configuredLimit = 32 * 1024
s := &service{cfg: &Config{UUID: t.Name()}}
Expand Down
2 changes: 1 addition & 1 deletion pkg/objectio/ioutil/sink_pool.go
Original file line number Diff line number Diff line change
Expand Up @@ -194,7 +194,7 @@ func (p *SinkPool) runSyncWorker() {
syncStart := time.Now()
stats, err := job.fSinker.Sync(r.ctx)
atomic.AddInt64(&r.syncNs, int64(time.Since(syncStart)))
if err != nil {
if err != nil && !objectSyncWasNotStarted(err) {
if name := activeFileSinkerObjectName(job.fSinker); name != "" {
r.mu.Lock()
r.unpublished = append(r.unpublished, name)
Expand Down
62 changes: 62 additions & 0 deletions pkg/objectio/ioutil/sink_pool_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ package ioutil

import (
"context"
"errors"
"fmt"
"sync"
"sync/atomic"
Expand Down Expand Up @@ -151,6 +152,67 @@ func mockFactory(sinkErr, syncErr error) FileSinkerFactory {
}
}

func TestAdmittedFileSinkerReservesBeforeSync(t *testing.T) {
base := &mockFileSinker{activeName: "object", syncStart: make(chan struct{})}
admitted := &admittedFileSinker{
FileSinker: base,
reserve: func(name string) error {
require.Equal(t, "object", name)
return errors.New("cleanup capacity exhausted")
},
}
_, err := admitted.Sync(context.Background())
require.ErrorContains(t, err, "cleanup capacity exhausted")
select {
case <-base.syncStart:
t.Fatal("Sync ran despite failed pre-upload admission")
default:
}
}

func TestSinkerRejectedAdmissionDoesNotCreateCleanupDebt(t *testing.T) {
for _, pipeline := range []bool{false, true} {
t.Run(fmt.Sprintf("pipeline=%t", pipeline), func(t *testing.T) {
proc := testutil.NewProc(t)
fs, err := fileservice.NewMemoryFS("shared", fileservice.DisabledCacheConfig, nil)
require.NoError(t, err)
attrs, typs, _ := mockSchema(3, -1)
base := &mockFileSinker{activeName: "not-written", syncStart: make(chan struct{})}
admissionErr := errors.New("cleanup capacity exhausted")
opts := []SinkerOption{
WithMemorySizeThreshold(1),
WithObjectSyncAdmission(func(name string) error {
if name != "not-written" {
return fmt.Errorf("unexpected object name %q", name)
}
return admissionErr
}, func(string) {}),
}
if pipeline {
opts = append(opts, WithPipelineFlush())
}
sinker := NewSinker(-1, attrs, typs,
func(*mpool.MPool, fileservice.FileService) FileSinker { return base },
proc.Mp(), fs, opts...)
bat := containers.MockBatch(typs, 8192, -1, nil)
err = sinker.Write(context.Background(), containers.ToCNBatch(bat))
if err == nil {
err = sinker.Sync(context.Background())
}
require.ErrorIs(t, err, admissionErr)
files, cleanupErr := sinker.DeletePersisted(context.Background())
require.NoError(t, cleanupErr)
require.Empty(t, files, "a rejected upload must not enter cleanup debt")
select {
case <-base.syncStart:
t.Fatal("file Sync started despite rejected admission")
default:
}
_ = sinker.Close()
})
}
}

// ---------- SinkPool.Submit tests ----------

func TestSinkPool_SubmitSuccess(t *testing.T) {
Expand Down
89 changes: 85 additions & 4 deletions pkg/objectio/ioutil/sinker.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ import (
"github.com/matrixorigin/matrixone/pkg/logutil"
"github.com/matrixorigin/matrixone/pkg/objectio"
"github.com/matrixorigin/matrixone/pkg/objectio/mergeutil"
metricv2 "github.com/matrixorigin/matrixone/pkg/util/metric/v2"
"github.com/matrixorigin/matrixone/pkg/vm/engine/tae/containers"
)

Expand Down Expand Up @@ -107,6 +108,47 @@ func WithChunkedColumnPolicy(policy objectio.ChunkedColumnPolicy) SinkerOption {
}
}

// WithObjectSyncAdmission reserves cleanup capacity before each file Sync.
// Release runs only after this sinker successfully deletes an object.
func WithObjectSyncAdmission(reserve func(string) error, release func(string)) SinkerOption {
return func(sinker *Sinker) {
sinker.config.reserveObject = reserve
sinker.config.releaseObject = release
}
}

type admittedFileSinker struct {
FileSinker
reserve func(string) error
}

// objectSyncAdmissionError means no object write began. The ordinary Sync
// error path must retain the name because persistence may be ambiguous.
type objectSyncAdmissionError struct{ cause error }

func (e *objectSyncAdmissionError) Error() string { return e.cause.Error() }
func (e *objectSyncAdmissionError) Unwrap() error { return e.cause }

func objectSyncWasNotStarted(err error) bool {
var admissionErr *objectSyncAdmissionError
return errors.As(err, &admissionErr)
}

func (s *admittedFileSinker) ActiveObjectName() string {
return activeFileSinkerObjectName(s.FileSinker)
}

func (s *admittedFileSinker) Sync(ctx context.Context) (*objectio.ObjectStats, error) {
name := s.ActiveObjectName()
if name == "" {
return nil, &objectSyncAdmissionError{moerr.NewInvalidStateNoCtx("file sinker has no object name before Sync")}
}
if err := s.reserve(name); err != nil {
return nil, &objectSyncAdmissionError{err}
}
return s.FileSinker.Sync(ctx)
}

type FileSinker interface {
Sink(context.Context, *batch.Batch) error
Sync(context.Context) (*objectio.ObjectStats, error)
Expand Down Expand Up @@ -353,6 +395,12 @@ func NewSinker(
return fileSinker
}
}
if reserve := sinker.config.reserveObject; reserve != nil {
factory := sinker.fSinker.factory
sinker.fSinker.factory = func(mp *mpool.MPool, fs fileservice.FileService) FileSinker {
return &admittedFileSinker{FileSinker: factory(mp, fs), reserve: reserve}
}
}

sinker.fillDefaults()
return sinker
Expand Down Expand Up @@ -406,6 +454,8 @@ type Sinker struct {
tailSizeCap int
offHeap bool
chunkedColumnPolicy objectio.ChunkedColumnPolicy
reserveObject func(string) error
releaseObject func(string)
}
fSinker struct {
executor FileSinker
Expand Down Expand Up @@ -493,8 +543,33 @@ func DeleteUnpublishedObjects(
fs fileservice.FileService,
files ...string,
) (int, error) {
unique, _, err := deleteUnpublishedObjectBatches(ctx, fs, files)
return len(unique), err
}

// DeleteUnpublishedObjectsWithProgress reports only names in fully confirmed
// Delete batches. A failed batch can have partial effects, so it and every
// unattempted batch remain owned. The remaining slice does not retain the
// backing array of completed names after the caller stores it for retry.
func DeleteUnpublishedObjectsWithProgress(
ctx context.Context,
fs fileservice.FileService,
files ...string,
) (completed, remaining []string, err error) {
unique, completedCount, err := deleteUnpublishedObjectBatches(ctx, fs, files)
if err != nil {
return unique[:completedCount], append([]string(nil), unique[completedCount:]...), err
}
return unique, nil, nil
}

func deleteUnpublishedObjectBatches(
ctx context.Context,
fs fileservice.FileService,
files []string,
) (unique []string, completed int, err error) {
seen := make(map[string]struct{}, len(files))
unique := make([]string, 0, len(files))
unique = make([]string, 0, len(files))
for _, file := range files {
if file == "" {
continue
Expand All @@ -513,14 +588,15 @@ func DeleteUnpublishedObjects(
err := fs.Delete(deleteCtx, unique[start:end]...)
cancel()
if err != nil && !moerr.IsMoErrCode(err, moerr.ErrFileNotFound) {
return len(unique), errors.Join(
metricv2.UnpublishedS3DeleteFailuresCounter.Inc()
return unique, start, errors.Join(
moerr.NewInternalErrorf(
ctx, "delete unpublished objects [%d:%d]", start, end),
err,
)
}
}
return len(unique), nil
return unique, len(unique), nil
}

// DeletePersisted deletes every object that this sinker has persisted, or may
Expand Down Expand Up @@ -557,6 +633,11 @@ func (sinker *Sinker) DeletePersisted(ctx context.Context) ([]string, error) {
if err != nil {
return files, err
}
if sinker.config.releaseObject != nil {
for _, name := range files {
sinker.config.releaseObject(name)
}
}

sinker.staged.persisted = sinker.staged.persisted[:0]
sinker.result.persisted = sinker.result.persisted[:0]
Expand Down Expand Up @@ -795,7 +876,7 @@ func (sinker *Sinker) syncFileSinker(ctx context.Context, fSinker FileSinker) er
stats, err := fSinker.Sync(ctx)
atomic.AddInt64(&sinker.timing.syncNs, int64(time.Since(syncStart)))
if err != nil {
if name := activeFileSinkerObjectName(fSinker); name != "" {
if name := activeFileSinkerObjectName(fSinker); name != "" && !objectSyncWasNotStarted(err) {
sinker.staged.unpublished = append(sinker.staged.unpublished, name)
}
return err
Expand Down
Loading
Loading