From 68f94501b5a358447a114a6919efa4a0e1353bbc Mon Sep 17 00:00:00 2001 From: 2pk03 Date: Tue, 7 Jul 2026 13:18:43 +0200 Subject: [PATCH] fix(metadata): hold persistMu across topic create/delete mutations The snapshot watcher could refresh in-memory metadata from etcd between in-memory DeleteTopic/CreateTopic and persistSnapshot, causing rollback paths (e.g. restore failure cleanup) to leave stale topics visible under -race. Serialize the full mutation+persist critical section via persistSnapshotLocked. --- pkg/metadata/etcd_store.go | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/pkg/metadata/etcd_store.go b/pkg/metadata/etcd_store.go index 7c30d71f..fa3abec7 100644 --- a/pkg/metadata/etcd_store.go +++ b/pkg/metadata/etcd_store.go @@ -428,11 +428,14 @@ func (s *EtcdStore) CreatePartitions(ctx context.Context, topic string, partitio // CreateTopic currently updates only the in-memory snapshot; the operator is still responsible // for reconciling durable topic configuration into etcd/S3. func (s *EtcdStore) CreateTopic(ctx context.Context, spec TopicSpec) (*protocol.MetadataTopic, error) { + s.persistMu.Lock() + defer s.persistMu.Unlock() + topic, err := s.metadata.CreateTopic(ctx, spec) if err != nil { return nil, err } - if err := s.persistSnapshot(ctx); err != nil { + if err := s.persistSnapshotLocked(ctx); err != nil { return nil, err } return topic, nil @@ -460,6 +463,9 @@ func (s *EtcdStore) partitionExists(ctx context.Context, topic string, partition // DeleteTopic updates the local snapshot so admin APIs behave consistently. func (s *EtcdStore) DeleteTopic(ctx context.Context, name string) error { + s.persistMu.Lock() + defer s.persistMu.Unlock() + metaCtx, cancel := context.WithTimeout(ctx, 3*time.Second) defer cancel() state, err := s.metadata.Metadata(metaCtx, []string{name}) @@ -485,10 +491,7 @@ func (s *EtcdStore) DeleteTopic(ctx context.Context, name string) error { if err := s.deleteConsumerOffsets(ctx, name); err != nil { return err } - if err := s.persistSnapshot(ctx); err != nil { - return err - } - return nil + return s.persistSnapshotLocked(ctx) } func (s *EtcdStore) startWatchers() { @@ -552,7 +555,10 @@ func snapshotKey() string { func (s *EtcdStore) persistSnapshot(ctx context.Context) error { s.persistMu.Lock() defer s.persistMu.Unlock() + return s.persistSnapshotLocked(ctx) +} +func (s *EtcdStore) persistSnapshotLocked(ctx context.Context) error { state, err := s.metadata.Metadata(context.Background(), nil) if err != nil { return err