diff --git a/blobstore/proxy/cacher/cacher.go b/blobstore/proxy/cacher/cacher.go index afd04bb4d..4e358b8fd 100644 --- a/blobstore/proxy/cacher/cacher.go +++ b/blobstore/proxy/cacher/cacher.go @@ -243,6 +243,8 @@ func (c *cacher) getCachedValue(span trace.Span, id any, key string, expiration return value } reporter("memcache", "expired") + span.Debugf("memory cache value has expired:%s id:%v", key, id) + return nil } } reporter("memcache", "miss") diff --git a/blobstore/proxy/cacher/volume.go b/blobstore/proxy/cacher/volume.go index 211adf756..43b0e68ce 100644 --- a/blobstore/proxy/cacher/volume.go +++ b/blobstore/proxy/cacher/volume.go @@ -66,7 +66,7 @@ func (c *cacher) GetVolume(ctx context.Context, args *proxy.CacheVolumeArgs) (*c return &vol.VolumeInfo, nil } - if args.Version > 0 && args.Version != uint64(vol.RouteVersion) { + if args.Version > 0 && args.Version <= uint64(vol.RouteVersion) { span.Infof("request to flush, but version mismatch request(%d) != cache(%d)", args.Version, vol.RouteVersion) return &vol.VolumeInfo, nil @@ -97,23 +97,23 @@ func (c *cacher) GetVolume(ctx context.Context, args *proxy.CacheVolumeArgs) (*c } c.volumeReport("clustermgr", "hit") - var result *clustermgr.VolumeInfo - if err := c.withVolumeLock(vid, func() error { - if cached := c.getVolume(span, vid); cached != nil && cached.RouteVersion > volume.RouteVersion { - span.Infof("skip storing stale cm volume vid:%d since cached route version %d > %d", - vid, cached.RouteVersion, volume.RouteVersion) - result = &cached.VolumeInfo + go func() { + _span, _ctx := trace.StartSpanFromContextWithTraceID(context.Background(), "", span.TraceID()) + defer _span.Finish() + err := c.withVolumeLock(vid, func() error { + if cached := c.getVolume(_span, vid); cached != nil && cached.RouteVersion > volume.RouteVersion { + _span.Infof("skip storing stale cm volume vid:%d ", vid) + return nil + } + vol := c.newExpiryVolume(*volume) + c.storeVolume(_ctx, vid, vol) return nil + }) + if err != nil { + _span.Warnf("write volume info to kv failed: %s", err) } - - vol := c.newExpiryVolume(*volume) - c.storeVolume(ctx, vid, vol) - result = &vol.VolumeInfo - return nil - }); err != nil { - return nil, err - } - return result, nil + }() + return volume, nil } func (c *cacher) getVolume(span trace.Span, vid proto.Vid) *expiryVolume { @@ -193,7 +193,6 @@ func (c *cacher) syncVolumeRoutes(ctx context.Context) error { if err := c.applyVolumeRouteItem(ctx, item); err != nil { return err } - span.Debugf("applied volume route item: %+v", item) } if err := c.persistVolRouteVersion(ctx, ret.RouteVersion); err != nil { @@ -208,18 +207,23 @@ func (c *cacher) applyVolumeRouteItem(ctx context.Context, item clustermgr.Volum if item.Item == nil { return errors.New("volume route item missing payload") } + span := trace.SpanFromContextSafe(ctx) switch item.Type { case proto.RouteItemTypeAddVolume: payload := clustermgr.RouteItemAddVolume{} if err := payload.Unmarshal(item.Item.Value); err != nil { return err } + payload.VolumeInfoBase.RouteVersion = item.RouteVersion + span.Debugf("apply volume route item: %+v", payload) return c.applyVolumeAdd(ctx, &payload) case proto.RouteItemTypeUpdateVolume: payload := clustermgr.RouteItemUpdateVolume{} if err := payload.Unmarshal(item.Item.Value); err != nil { return err } + payload.VolumeInfoBase.RouteVersion = item.RouteVersion + span.Debugf("apply volume route item: %+v", payload) return c.applyVolumeUpdate(ctx, &payload) default: return fmt.Errorf("unknown volume route item type: %d", item.Type) @@ -276,7 +280,7 @@ func (c *cacher) newExpiryVolume(info clustermgr.VolumeInfo) *expiryVolume { func (c *cacher) storeVolume(ctx context.Context, vid proto.Vid, vol *expiryVolume) { c.volumeCache.Set(vid, vol) - go c.writeVolumeToDisk(ctx, vid, vol) + c.writeVolumeToDisk(ctx, vid, vol) } func (c *cacher) writeVolumeToDisk(ctx context.Context, vid proto.Vid, vol *expiryVolume) { diff --git a/blobstore/proxy/cacher/volume_test.go b/blobstore/proxy/cacher/volume_test.go index b924f1fb9..e8cab87ca 100644 --- a/blobstore/proxy/cacher/volume_test.go +++ b/blobstore/proxy/cacher/volume_test.go @@ -44,24 +44,51 @@ import ( func TestProxyCacherVolumeUpdate(t *testing.T) { c, cmCli, clean := newCacher(t, 2, nil) defer clean() + cc := c.(*cacher) + ctx := context.Background() cmCli.EXPECT().GetVolumeInfo(A, A).Return(&clustermgr.VolumeInfo{}, nil).Times(4) - for range [100]struct{}{} { - _, err := c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 1}) + span := trace.SpanFromContextSafe(ctx) + waitCached := func(vid proto.Vid) { + require.Eventually(t, func() bool { + return cc.getVolume(span, vid) != nil + }, time.Second, 5*time.Millisecond) + } + + // storeVolume writes disk synchronously, but withVolumeLock runs in a goroutine. + // After the first GetVolume triggers a CM fetch, wait for the cache to be written + // before subsequent requests can hit the cache. + _, err := c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 1}) + require.NoError(t, err) + waitCached(1) + for range [99]struct{}{} { + _, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 1}) require.NoError(t, err) } - for range [100]struct{}{} { - _, err := c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 2}) + + _, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 2}) + require.NoError(t, err) + waitCached(2) + for range [99]struct{}{} { + _, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 2}) require.NoError(t, err) } time.Sleep(time.Second * 4) // expired - for range [100]struct{}{} { - _, err := c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 1}) + + _, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 1}) + require.NoError(t, err) + waitCached(1) + for range [99]struct{}{} { + _, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 1}) require.NoError(t, err) } - for range [100]struct{}{} { - _, err := c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 2}) + + _, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 2}) + require.NoError(t, err) + waitCached(2) + for range [99]struct{}{} { + _, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 2}) require.NoError(t, err) } } @@ -69,35 +96,61 @@ func TestProxyCacherVolumeUpdate(t *testing.T) { func TestProxyCacherVolumeFlush(t *testing.T) { c, cmCli, clean := newCacher(t, 0, nil) defer clean() + cc := c.(*cacher) + ctx := context.Background() + span := trace.SpanFromContextSafe(ctx) volume := new(clustermgr.VolumeInfo) volume.Units = []clustermgr.Unit{{Vuid: 1234}, {Vuid: 5678}} version := proto.RouteVersion(12345678) - volume.RouteVersion = proto.RouteVersion(version) + volume.RouteVersion = version + + waitCached := func(vid proto.Vid) { + require.Eventually(t, func() bool { + return cc.getVolume(span, vid) != nil + }, time.Second, 5*time.Millisecond) + } + + // vid=1 first request hits CM; wait for cache write before subsequent requests cmCli.EXPECT().GetVolumeInfo(A, A).Return(volume, nil).Times(1) - for range [100]struct{}{} { - vol, err := c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 1}) + vol, err := c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 1}) + require.NoError(t, err) + require.Equal(t, version, vol.RouteVersion) + waitCached(1) + for range [99]struct{}{} { + vol, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 1}) require.NoError(t, err) require.Equal(t, version, vol.RouteVersion) } + // Flush=true with Version=0: version check is skipped, every request goes to CM cmCli.EXPECT().GetVolumeInfo(A, A).Return(volume, nil).Times(100) for range [100]struct{}{} { - vol, err := c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 1, Flush: true}) + vol, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 1, Flush: true}) require.NoError(t, err) require.Equal(t, version, vol.RouteVersion) } + // vid=3 first request has no cache, hits CM; Version=0x01 < cache version, subsequent requests skip CM cmCli.EXPECT().GetVolumeInfo(A, A).Return(volume, nil).Times(1) - for range [100]struct{}{} { - vol, err := c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 3, Flush: true, Version: 0x01}) + vol, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 3, Flush: true, Version: 0x01}) + require.NoError(t, err) + require.Equal(t, version, vol.RouteVersion) + waitCached(3) + for range [99]struct{}{} { + vol, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 3, Flush: true, Version: 0x01}) require.NoError(t, err) require.Equal(t, version, vol.RouteVersion) } + // vid=4 first request has no cache, hits CM; Version=version-1 < cache version, subsequent requests skip CM cmCli.EXPECT().GetVolumeInfo(A, A).Return(volume, nil).Times(1) - for range [100]struct{}{} { - vol, err := c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 4, Flush: true, Version: 0x9d31f755}) + vol, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 4, Flush: true, Version: uint64(version) - 1}) + require.NoError(t, err) + require.Equal(t, version, vol.RouteVersion) + waitCached(4) + for range [99]struct{}{} { + vol, err = c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 4, Flush: true, Version: uint64(version) - 1}) require.NoError(t, err) require.Equal(t, version, vol.RouteVersion) } @@ -126,11 +179,12 @@ func TestProxyCacherVolumeFlushSkipStaleCMVolume(t *testing.T) { } cmCli.EXPECT().GetVolumeInfo(A, A).Return(older, nil).Times(1) + // GetVolume now returns CM data directly (older); it no longer synchronously returns the newer cached value vol, err := c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 1, Flush: true}) require.NoError(t, err) - require.Equal(t, newer.RouteVersion, vol.RouteVersion) - require.Equal(t, newer.Units, vol.Units) + require.Equal(t, older.RouteVersion, vol.RouteVersion) + // goroutine skips write because cached.RouteVersion(11) > cm.RouteVersion(10); cache retains newer span := trace.SpanFromContextSafe(context.Background()) cached := cc.getVolume(span, proto.Vid(1)) require.NotNil(t, cached) @@ -212,6 +266,28 @@ func TestProxyCacherVolumeCacheMiss(t *testing.T) { } } +func TestProxyCacherVolumeMemExpiredSkipsDiskv(t *testing.T) { + c, cmCli, clean := newCacher(t, 0, nil) + defer clean() + cc := c.(*cacher) + ctx := context.Background() + + info := clustermgr.VolumeInfo{ + VolumeInfoBase: clustermgr.VolumeInfoBase{Vid: 1, RouteVersion: proto.RouteVersion(10)}, + } + vol := cc.newExpiryVolume(info) + cc.storeVolume(ctx, 1, vol) + <-cc.syncChan + + // Expire the in-memory entry without touching diskv, so diskv still holds valid data. + // This isolates the "mem expired -> skip diskv -> go to CM" path. + vol.expiration = time.Nanosecond + + cmCli.EXPECT().GetVolumeInfo(A, A).Return(&info, nil).Times(1) + _, err := c.GetVolume(ctx, &proxy.CacheVolumeArgs{Vid: 1}) + require.NoError(t, err) +} + func BenchmarkProxyMemoryHit(b *testing.B) { c, cmCli, clean := newCacher(b, 0, nil) defer clean() @@ -282,16 +358,18 @@ func TestProxyCacheCloser(t *testing.T) { func TestVolumeRouteSyncAddVolume(t *testing.T) { cmCli := mocks.NewMockClientAPI(C(t)) + // CM fills VolumeInfoBase.RouteVersion with the latest overall version (20), + // but item.RouteVersion is the version of this specific route record (10). + // The proxy must use item.RouteVersion (10) for the cached volume, not 20. payload := &clustermgr.RouteItemAddVolume{ - Vid: 1, - RouteVersion: proto.RouteVersion(10), + Vid: 1, Units: []clustermgr.VolumeUnitInfoBase{ {Vuid: proto.Vuid(1), DiskID: proto.DiskID(1), Host: "host-1"}, }, VolumeInfoBase: clustermgr.VolumeInfoBasePB{ Vid: 1, CodeMode: codemode.EC6P6, - RouteVersion: proto.RouteVersion(10), + RouteVersion: proto.RouteVersion(20), }, } cmCli.EXPECT().GetVolumeRoutes(A, A).Return(&clustermgr.GetVolumeRoutesRet{ @@ -316,16 +394,137 @@ func TestVolumeRouteSyncAddVolume(t *testing.T) { span := trace.SpanFromContextSafe(context.Background()) vol := cc.getVolume(span, proto.Vid(1)) require.NotNil(t, vol) + // must be item.RouteVersion (10), not payload.VolumeInfoBase.RouteVersion (20) require.Equal(t, proto.RouteVersion(10), vol.VolumeInfo.RouteVersion) require.Len(t, vol.VolumeInfo.Units, 1) require.Equal(t, proto.DiskID(1), vol.VolumeInfo.Units[0].DiskID) } +func TestProxyCacherVolumeFlushVersionComparison(t *testing.T) { + c, cmCli, clean := newCacher(t, 0, nil) + defer clean() + cc := c.(*cacher) + + cached := clustermgr.VolumeInfo{ + VolumeInfoBase: clustermgr.VolumeInfoBase{Vid: 1, RouteVersion: proto.RouteVersion(10)}, + Units: []clustermgr.Unit{{Vuid: proto.Vuid(1), DiskID: proto.DiskID(1)}}, + } + cc.storeVolume(context.Background(), 1, cc.newExpiryVolume(cached)) + // drain the syncChan signal from setup storeVolume + select { + case <-cc.syncChan: + default: + } + + newer := &clustermgr.VolumeInfo{ + VolumeInfoBase: clustermgr.VolumeInfoBase{Vid: 1, RouteVersion: proto.RouteVersion(11)}, + Units: []clustermgr.Unit{{Vuid: proto.Vuid(1), DiskID: proto.DiskID(2)}}, + } + + // cache ver=10 < access.Version=11: proxy is stale, should trigger a CM fetch + cmCli.EXPECT().GetVolumeInfo(A, A).Return(newer, nil).Times(1) + vol, err := c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 1, Flush: true, Version: 11}) + require.NoError(t, err) + require.Equal(t, proto.RouteVersion(11), vol.RouteVersion) + + // wait for the async goroutine to write ver=11 into cache + <-cc.syncChan + + span := trace.SpanFromContextSafe(context.Background()) + + // cache ver=11 == access.Version=11: skip CM + vol, err = c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 1, Flush: true, Version: 11}) + require.NoError(t, err) + require.Equal(t, proto.RouteVersion(11), vol.RouteVersion) + require.Equal(t, proto.RouteVersion(11), cc.getVolume(span, 1).RouteVersion) + + // cache ver=11 > access.Version=10: proxy is ahead, skip CM + vol, err = c.GetVolume(context.Background(), &proxy.CacheVolumeArgs{Vid: 1, Flush: true, Version: 10}) + require.NoError(t, err) + require.Equal(t, proto.RouteVersion(11), vol.RouteVersion) +} + +func TestVolumeRouteSyncUpdateMultipleUnits(t *testing.T) { + cmCli := mocks.NewMockClientAPI(C(t)) + + vid := proto.Vid(1) + // CM fills VolumeInfoBase.RouteVersion with the latest overall version (12) for all items, + // but each item carries its own item.RouteVersion (10, 11, 12 respectively). + addPayload := &clustermgr.RouteItemAddVolume{ + Vid: vid, + Units: []clustermgr.VolumeUnitInfoBase{ + {Vuid: proto.EncodeVuid(proto.EncodeVuidPrefix(vid, 0), 1), DiskID: proto.DiskID(1), Host: "host-1"}, + {Vuid: proto.EncodeVuid(proto.EncodeVuidPrefix(vid, 1), 1), DiskID: proto.DiskID(2), Host: "host-2"}, + }, + VolumeInfoBase: clustermgr.VolumeInfoBasePB{Vid: vid, RouteVersion: proto.RouteVersion(12)}, + } + // CM behavior: two UpdateVolume items for the same volume, both with payload.VolumeInfoBase.RouteVersion=12 (current latest). + // Correct idempotency must rely on item.RouteVersion (11, 12), not the version inside the payload. + update0 := &clustermgr.RouteItemUpdateVolume{ + Vid: vid, + Unit: clustermgr.VolumeUnitInfoBase{Vuid: proto.EncodeVuid(proto.EncodeVuidPrefix(vid, 0), 2), DiskID: proto.DiskID(3), Host: "host-3"}, + VolumeInfoBase: clustermgr.VolumeInfoBasePB{Vid: vid, RouteVersion: proto.RouteVersion(12)}, + } + update1 := &clustermgr.RouteItemUpdateVolume{ + Vid: vid, + Unit: clustermgr.VolumeUnitInfoBase{Vuid: proto.EncodeVuid(proto.EncodeVuidPrefix(vid, 1), 2), DiskID: proto.DiskID(4), Host: "host-4"}, + VolumeInfoBase: clustermgr.VolumeInfoBasePB{Vid: vid, RouteVersion: proto.RouteVersion(12)}, + } + + gomock.InOrder( + // New calls syncVolumeRoutes once on init when version=0 + cmCli.EXPECT().GetVolumeRoutes(A, A).Return(&clustermgr.GetVolumeRoutesRet{}, nil), + cmCli.EXPECT().GetVolumeRoutes(A, A).Return(&clustermgr.GetVolumeRoutesRet{ + RouteVersion: proto.RouteVersion(10), + Items: []clustermgr.VolumeRouteItem{{ + RouteVersion: proto.RouteVersion(10), + Type: proto.RouteItemTypeAddVolume, + Item: mustMarshalAny(t, addPayload), + }}, + }, nil), + cmCli.EXPECT().GetVolumeRoutes(A, A).Return(&clustermgr.GetVolumeRoutesRet{ + RouteVersion: proto.RouteVersion(12), + Items: []clustermgr.VolumeRouteItem{ + {RouteVersion: proto.RouteVersion(11), Type: proto.RouteItemTypeUpdateVolume, Item: mustMarshalAny(t, update0)}, + {RouteVersion: proto.RouteVersion(12), Type: proto.RouteItemTypeUpdateVolume, Item: mustMarshalAny(t, update1)}, + }, + }, nil), + ) + + c, _, clean := newCacher(t, 0, cmCli) + defer clean() + cc := c.(*cacher) + + require.NoError(t, cc.syncVolumeRoutes(context.Background())) + require.NoError(t, cc.syncVolumeRoutes(context.Background())) + require.Equal(t, proto.RouteVersion(12), cc.getVolRouteVersion()) + + span := trace.SpanFromContextSafe(context.Background()) + vol := cc.getVolume(span, vid) + require.NotNil(t, vol) + require.Equal(t, proto.RouteVersion(12), vol.VolumeInfo.RouteVersion) + require.Len(t, vol.VolumeInfo.Units, 2) + + unitByIndex := func(units []clustermgr.Unit, idx uint8) *clustermgr.Unit { + for i := range units { + if units[i].Vuid.Index() == idx { + return &units[i] + } + } + return nil + } + u0 := unitByIndex(vol.VolumeInfo.Units, 0) + require.NotNil(t, u0) + require.Equal(t, proto.DiskID(3), u0.DiskID, "unit-0 should be updated to disk3") + u1 := unitByIndex(vol.VolumeInfo.Units, 1) + require.NotNil(t, u1) + require.Equal(t, proto.DiskID(4), u1.DiskID, "unit-1 should be updated to disk4") +} + func TestVolumeRouteSyncUpdateUnit(t *testing.T) { cmCli := mocks.NewMockClientAPI(C(t)) addPayload := &clustermgr.RouteItemAddVolume{ - Vid: 1, - RouteVersion: proto.RouteVersion(10), + Vid: 1, Units: []clustermgr.VolumeUnitInfoBase{ {Vuid: proto.Vuid(1), DiskID: proto.DiskID(1), Host: "host-1"}, {Vuid: proto.Vuid(2), DiskID: proto.DiskID(2), Host: "host-2"}, @@ -336,9 +535,11 @@ func TestVolumeRouteSyncUpdateUnit(t *testing.T) { RouteVersion: proto.RouteVersion(10), }, } + // CM fills VolumeInfoBase.RouteVersion with the latest overall version (20), + // but item.RouteVersion is the version of this specific route record (11). + // The proxy must use item.RouteVersion (11) for the cached volume, not 20. updatePayload := &clustermgr.RouteItemUpdateVolume{ - Vid: 1, - RouteVersion: proto.RouteVersion(11), + Vid: 1, Unit: clustermgr.VolumeUnitInfoBase{ Vuid: proto.EncodeVuid(proto.VuidPrefix(2), 1), DiskID: proto.DiskID(3), @@ -347,7 +548,7 @@ func TestVolumeRouteSyncUpdateUnit(t *testing.T) { VolumeInfoBase: clustermgr.VolumeInfoBasePB{ Vid: 1, CodeMode: codemode.EC6P6, - RouteVersion: proto.RouteVersion(11), + RouteVersion: proto.RouteVersion(20), }, } @@ -382,6 +583,7 @@ func TestVolumeRouteSyncUpdateUnit(t *testing.T) { span := trace.SpanFromContextSafe(context.Background()) vol := cc.getVolume(span, proto.Vid(1)) require.NotNil(t, vol) + // must be item.RouteVersion (11), not payload.VolumeInfoBase.RouteVersion (20) require.Equal(t, proto.RouteVersion(11), vol.VolumeInfo.RouteVersion) require.Len(t, vol.VolumeInfo.Units, 2) require.Equal(t, proto.DiskID(3), vol.VolumeInfo.Units[0].DiskID)