From 4d2e26c6d63e92328c99115c96ff4accc4637c73 Mon Sep 17 00:00:00 2001 From: XuPeng-SH Date: Fri, 11 Sep 2026 10:43:14 +0800 Subject: [PATCH 1/2] fix(embed): join startup workers and preserve service failures --- pkg/embed/cluster.go | 38 ++++++++++-------- pkg/embed/cluster_test.go | 83 +++++++++++++++++++++++++++++++++++---- 2 files changed, 96 insertions(+), 25 deletions(-) diff --git a/pkg/embed/cluster.go b/pkg/embed/cluster.go index 08400a9cf4892..bba24956cc804 100644 --- a/pkg/embed/cluster.go +++ b/pkg/embed/cluster.go @@ -201,40 +201,44 @@ func (c *cluster) Start() (err error) { } func (c *cluster) doStartLocked(from int) error { + services := c.services[from:] + // Each service owns one slot. Join only after every started CN has returned, + // so rollback cannot race startup and diagnostics retain all failures in + // configuration order regardless of completion order or concrete error type. + startErrors := make([]error, len(services)) var wg sync.WaitGroup - var startErr atomic.Value - for _, s := range c.services[from:] { + for i, s := range services { if s.serviceType != metadata.ServiceType_CN { if err := c.startServiceLocked(s); err != nil { - return err + startErrors[i] = err + break } continue } wg.Add(1) - go func(s *operator) { + go func(i int, s *operator) { defer wg.Done() - if err := c.startServiceLocked(s); err != nil { - // Only the first error is captured; concurrent failures - // from other services are discarded since knowing that - // any service failed is sufficient to abort startup. - startErr.CompareAndSwap(nil, err) - } - }(s) + startErrors[i] = c.startServiceLocked(s) + }(i, s) } wg.Wait() - if v := startErr.Load(); v != nil { - return v.(error) - } - return nil + return errors.Join(startErrors...) } func (c *cluster) startServiceLocked(op *operator) error { + var err error if c.startFn != nil { - return c.startFn(op) + err = c.startFn(op) + } else { + err = op.Start() + } + if err != nil { + return fmt.Errorf("embedded cluster %d start %s service %q: %w", + c.id, op.serviceType, op.sid, err) } - return op.Start() + return nil } func (c *cluster) Close() error { diff --git a/pkg/embed/cluster_test.go b/pkg/embed/cluster_test.go index 935b842875fca..d135dc7aecaa5 100644 --- a/pkg/embed/cluster_test.go +++ b/pkg/embed/cluster_test.go @@ -660,9 +660,7 @@ func TestCreateDB(t *testing.T) { // doStartLocked that are not reached by normal cluster startup tests. func TestDoStartLockedErrorPaths(t *testing.T) { t.Run("non-CN service error returns immediately", func(t *testing.T) { - // A non-CN operator whose state is already 'started' will return - // an error from Start(), exercising the direct-return path at - // cluster.go line 119-121. + // A non-CN operator rejects duplicate startup before allocating resources. op := &operator{ serviceType: metadata.ServiceType_LOG, state: started, // forces Start() to return error @@ -675,10 +673,8 @@ func TestDoStartLockedErrorPaths(t *testing.T) { assert.Contains(t, err.Error(), "already started") }) - t.Run("CN service error captured via atomic.Value", func(t *testing.T) { - // A CN operator whose state is already 'started' will return an - // error from Start(), exercising the goroutine error-capture path - // at cluster.go lines 128-133 and the error-return at 138-140. + t.Run("CN service error preserves cause", func(t *testing.T) { + // The concurrent startup path preserves the same rejection. op := &operator{ serviceType: metadata.ServiceType_CN, state: started, @@ -692,7 +688,7 @@ func TestDoStartLockedErrorPaths(t *testing.T) { }) t.Run("Start propagates doStartLocked error", func(t *testing.T) { - // Exercises the error propagation in Start() at line 107-109. + // Start preserves service errors through rollback. op := &operator{ serviceType: metadata.ServiceType_LOG, state: started, @@ -733,6 +729,77 @@ func TestDoStartLockedErrorPaths(t *testing.T) { }) } +// These startup regressions use the existing startFn seam: no real services, +// ports, cluster admission, or wall-clock sleeps are needed. +func TestDoStartLockedConcurrentErrors(t *testing.T) { + first := errors.New("connection refused") + second := &os.PathError{Op: "open", Path: "catalog", Err: os.ErrPermission} + c := &cluster{ + id: 42, + services: []*operator{ + {serviceType: metadata.ServiceType_CN, sid: "cn-a"}, + {serviceType: metadata.ServiceType_CN, sid: "cn-b"}, + }, + startFn: func(op *operator) error { + if op.sid == "cn-a" { + return first + } + return second + }, + } + err := c.doStartLocked(0) + require.ErrorIs(t, err, first) + require.ErrorIs(t, err, second) + var pathErr *os.PathError + require.ErrorAs(t, err, &pathErr) + require.Same(t, second, pathErr) + require.EqualError(t, err, + "embedded cluster 42 start CN service \"cn-a\": connection refused\n"+ + "embedded cluster 42 start CN service \"cn-b\": open catalog: permission denied") + + // Startup results belong to this invocation, including incremental CN starts. + c.startFn = func(op *operator) error { + assert.Equal(t, "cn-b", op.sid) + return nil + } + require.NoError(t, c.doStartLocked(1)) +} + +func TestDoStartLockedJoinsCNBeforeReturningInfrastructureFailure(t *testing.T) { + cnEntered := make(chan struct{}) + infrastructureFailed := make(chan struct{}) + cnErr := errors.New("CN startup failed") + logErr := errors.New("log startup failed") + var cnReturned atomic.Bool + c := &cluster{ + services: []*operator{ + {serviceType: metadata.ServiceType_CN, sid: "cn"}, + {serviceType: metadata.ServiceType_LOG, sid: "log"}, + {serviceType: metadata.ServiceType_TN, sid: "not-started"}, + }, + startFn: func(op *operator) error { + switch op.sid { + case "cn": + close(cnEntered) + <-infrastructureFailed + defer cnReturned.Store(true) + return cnErr + case "log": + <-cnEntered + close(infrastructureFailed) + return logErr + default: + t.Error("started a service after infrastructure failure") + return nil + } + }, + } + err := c.doStartLocked(0) + require.True(t, cnReturned.Load(), "all startup workers must finish before caller can roll back") + require.ErrorIs(t, err, cnErr) + require.ErrorIs(t, err, logErr) +} + func TestClusterStartRollbackClosesPartiallyStartedServices(t *testing.T) { startErr := errors.New("TN wait for HAKeeper timed out") portLease, err := acquireClusterPortLease() From 12e751377d039a493a9e0be134c076f1988f1bf7 Mon Sep 17 00:00:00 2001 From: XuPeng-SH Date: Fri, 11 Sep 2026 11:04:34 +0800 Subject: [PATCH 2/2] fix: satisfy embed SCA --- pkg/embed/cluster.go | 7 +++++-- pkg/embed/cluster_test.go | 6 ++++-- 2 files changed, 9 insertions(+), 4 deletions(-) diff --git a/pkg/embed/cluster.go b/pkg/embed/cluster.go index bba24956cc804..085bafc00f7b0 100644 --- a/pkg/embed/cluster.go +++ b/pkg/embed/cluster.go @@ -235,8 +235,11 @@ func (c *cluster) startServiceLocked(op *operator) error { err = op.Start() } if err != nil { - return fmt.Errorf("embedded cluster %d start %s service %q: %w", - c.id, op.serviceType, op.sid, err) + return errors.Join( + moerr.NewInternalErrorNoCtxf("embedded cluster %d start %s service %q", + c.id, op.serviceType, op.sid), + err, + ) } return nil } diff --git a/pkg/embed/cluster_test.go b/pkg/embed/cluster_test.go index d135dc7aecaa5..cbb9abcd16649 100644 --- a/pkg/embed/cluster_test.go +++ b/pkg/embed/cluster_test.go @@ -754,8 +754,10 @@ func TestDoStartLockedConcurrentErrors(t *testing.T) { require.ErrorAs(t, err, &pathErr) require.Same(t, second, pathErr) require.EqualError(t, err, - "embedded cluster 42 start CN service \"cn-a\": connection refused\n"+ - "embedded cluster 42 start CN service \"cn-b\": open catalog: permission denied") + "internal error: embedded cluster 42 start CN service \"cn-a\"\n"+ + "connection refused\n"+ + "internal error: embedded cluster 42 start CN service \"cn-b\"\n"+ + "open catalog: permission denied") // Startup results belong to this invocation, including incremental CN starts. c.startFn = func(op *operator) error {