mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
refactor(blobnode): optimize io context cancel, and it has become simpler
with: #1000146037 Signed-off-by: mawei029 <mawei2@oppo.com>
This commit is contained in:
parent
390b3554b4
commit
80b069d508
44
blobstore/blobnode/base/io_ctx.go
Normal file
44
blobstore/blobnode/base/io_ctx.go
Normal file
@ -0,0 +1,44 @@
|
||||
package base
|
||||
|
||||
import (
|
||||
"context"
|
||||
)
|
||||
|
||||
type CtxWriterAt interface {
|
||||
WriteAtCtx(ctx context.Context, p []byte, off int64) (n int, err error)
|
||||
}
|
||||
|
||||
// WriterWithCtx : wrap blobfile.go WriteAtCtx. After io dequeue, it is not carried out when ctx cancel
|
||||
type WriterWithCtx struct {
|
||||
Offset int64
|
||||
Wt CtxWriterAt
|
||||
Ctx context.Context
|
||||
}
|
||||
|
||||
func (w *WriterWithCtx) Write(val []byte) (n int, err error) {
|
||||
n, err = w.Wt.WriteAtCtx(w.Ctx, val, w.Offset)
|
||||
w.Offset += int64(n)
|
||||
return
|
||||
}
|
||||
|
||||
type CtxReaderAt interface {
|
||||
ReadAtCtx(ctx context.Context, p []byte, off int64) (n int, err error)
|
||||
}
|
||||
|
||||
// ReaderWithCtx : wrap blobfile.go ReadAtCtx. After io dequeue, it is not carried out when ctx cancel
|
||||
type ReaderWithCtx struct {
|
||||
offset int64
|
||||
Rd CtxReaderAt
|
||||
Ctx context.Context
|
||||
}
|
||||
|
||||
func (r *ReaderWithCtx) ReadAt(val []byte, off int64) (n int, err error) {
|
||||
n, err = r.Rd.ReadAtCtx(r.Ctx, val, off)
|
||||
return
|
||||
}
|
||||
|
||||
func (r *ReaderWithCtx) Read(val []byte) (n int, err error) {
|
||||
n, err = r.Rd.ReadAtCtx(r.Ctx, val, r.offset)
|
||||
r.offset += int64(n)
|
||||
return
|
||||
}
|
||||
151
blobstore/blobnode/base/io_ctx_test.go
Normal file
151
blobstore/blobnode/base/io_ctx_test.go
Normal file
@ -0,0 +1,151 @@
|
||||
package base
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"io"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// mockWriterAt implements io.WriterAt and CtxWriterAt for testing.
|
||||
type mockWriterAt struct {
|
||||
WriteAtCtxFunc func(ctx context.Context, p []byte, off int64) (n int, err error)
|
||||
}
|
||||
|
||||
func (m *mockWriterAt) WriteAtCtx(ctx context.Context, p []byte, off int64) (n int, err error) {
|
||||
if m.WriteAtCtxFunc != nil {
|
||||
return m.WriteAtCtxFunc(ctx, p, off)
|
||||
}
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
func newMockWriterAt(w io.Writer) *mockWriterAt {
|
||||
return &mockWriterAt{
|
||||
WriteAtCtxFunc: func(ctx context.Context, p []byte, off int64) (n int, err error) {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
n, err = 0, ctx.Err()
|
||||
default:
|
||||
n, err = w.Write(p) // len(p), nil
|
||||
}
|
||||
return
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// mockReaderAt implements io.ReaderAt and CtxReader for testing.
|
||||
type mockReaderAt struct {
|
||||
ReadAtCtxFunc func(ctx context.Context, p []byte, off int64) (n int, err error)
|
||||
}
|
||||
|
||||
func (m *mockReaderAt) ReadAtCtx(ctx context.Context, p []byte, off int64) (n int, err error) {
|
||||
if m.ReadAtCtxFunc != nil {
|
||||
return m.ReadAtCtxFunc(ctx, p, off)
|
||||
}
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
func newMockReaderAt(src io.ReaderAt) *mockReaderAt {
|
||||
return &mockReaderAt{
|
||||
ReadAtCtxFunc: func(ctx context.Context, p []byte, off int64) (n int, err error) {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
n, err = 0, ctx.Err()
|
||||
default:
|
||||
n, err = src.ReadAt(p, off)
|
||||
}
|
||||
return
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriterWithCtx_Write(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
|
||||
str := "hello"
|
||||
buf := bytes.NewBuffer(nil)
|
||||
mockWt := newMockWriterAt(buf)
|
||||
|
||||
writer := &WriterWithCtx{
|
||||
Offset: 0,
|
||||
Wt: mockWt,
|
||||
Ctx: ctx,
|
||||
}
|
||||
|
||||
// Test normal write
|
||||
n, err := writer.Write([]byte(str))
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 5, n)
|
||||
require.Equal(t, int64(5), writer.Offset)
|
||||
|
||||
// Test context cancellation
|
||||
cancel()
|
||||
n, err = writer.Write([]byte("world"))
|
||||
require.Error(t, err)
|
||||
require.Equal(t, context.Canceled, err)
|
||||
require.Equal(t, 0, n) // n should be 0 as no bytes were written after the context was canceled
|
||||
require.Equal(t, int64(5), writer.Offset) // Offset should remain unchanged
|
||||
}
|
||||
|
||||
func TestReaderWithCtx_ReadAt(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
|
||||
data := []byte("hello")
|
||||
buf := bytes.NewReader(data)
|
||||
mockRd := newMockReaderAt(buf)
|
||||
|
||||
reader := &ReaderWithCtx{
|
||||
Rd: mockRd,
|
||||
Ctx: ctx,
|
||||
}
|
||||
|
||||
// Test normal read at
|
||||
p := make([]byte, 5)
|
||||
n, err := reader.ReadAt(p, 0)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 5, n)
|
||||
require.Equal(t, string(data), string(p))
|
||||
|
||||
// Test context cancellation
|
||||
cancel()
|
||||
p = make([]byte, 5)
|
||||
n, err = reader.ReadAt(p, 0)
|
||||
require.Error(t, err)
|
||||
require.Equal(t, context.Canceled, err)
|
||||
require.Equal(t, 0, n) // n should be 0 as no bytes were read after the context was canceled
|
||||
}
|
||||
|
||||
func TestReaderWithCtx_Read(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
|
||||
data := []byte("hello")
|
||||
buf := bytes.NewReader(data)
|
||||
mockRd := newMockReaderAt(buf)
|
||||
|
||||
reader := &ReaderWithCtx{
|
||||
Rd: mockRd,
|
||||
Ctx: ctx,
|
||||
}
|
||||
|
||||
// Test normal read at
|
||||
p1 := make([]byte, 3)
|
||||
n, err := reader.Read(p1)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 3, n)
|
||||
p2 := make([]byte, 2)
|
||||
n, err = reader.Read(p2)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 2, n)
|
||||
p := append(p1, p2...)
|
||||
require.Equal(t, string(data), string(p))
|
||||
|
||||
// Test context cancellation
|
||||
cancel()
|
||||
p = make([]byte, 5)
|
||||
n, err = reader.Read(p)
|
||||
require.Error(t, err)
|
||||
require.Equal(t, context.Canceled, err)
|
||||
require.Equal(t, 0, n) // n should be 0 as no bytes were read after the context was canceled
|
||||
}
|
||||
@ -22,7 +22,6 @@ import (
|
||||
"golang.org/x/time/rate"
|
||||
|
||||
"github.com/cubefs/cubefs/blobstore/common/errors"
|
||||
"github.com/cubefs/cubefs/blobstore/common/iostat"
|
||||
"github.com/cubefs/cubefs/blobstore/common/trace"
|
||||
)
|
||||
|
||||
@ -33,7 +32,6 @@ type rateLimiter struct {
|
||||
reader io.Reader
|
||||
writer io.Writer
|
||||
writerAt io.WriterAt
|
||||
wAtCtx iostat.WriterAtCtx
|
||||
ctx context.Context
|
||||
bpsLimiter *rate.Limiter
|
||||
}
|
||||
@ -90,14 +88,6 @@ func (l *rateLimiter) WriteAt(p []byte, off int64) (n int, err error) {
|
||||
return l.writerAt.WriteAt(p, off)
|
||||
}
|
||||
|
||||
func (l *rateLimiter) WriteAtCtx(ctx context.Context, p []byte, off int64) (n int, err error) {
|
||||
n, err = l.wAtCtx.WriteAtCtx(ctx, p, off)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
return n, l.doWithLimit(n)
|
||||
}
|
||||
|
||||
func (l *rateLimiter) doWithLimit(n int) (err error) {
|
||||
return l.doWithSingleLimit(l.bpsLimiter, n)
|
||||
}
|
||||
|
||||
@ -31,7 +31,7 @@ import (
|
||||
|
||||
type Qos interface {
|
||||
ReaderAt(context.Context, bnapi.IOType, io.ReaderAt) io.ReaderAt
|
||||
WriterAt(context.Context, bnapi.IOType, iostat.WriterAtCtx) iostat.WriterAtCtx
|
||||
WriterAt(context.Context, bnapi.IOType, io.WriterAt) io.WriterAt
|
||||
Writer(context.Context, bnapi.IOType, io.Writer) io.Writer
|
||||
Reader(context.Context, bnapi.IOType, io.Reader) io.Reader
|
||||
TryAllow(rwType IOTypeRW) bool
|
||||
@ -111,16 +111,16 @@ func (qos *IoQueueQos) ReaderAt(ctx context.Context, ioType bnapi.IOType, reader
|
||||
return r
|
||||
}
|
||||
|
||||
func (qos *IoQueueQos) WriterAt(ctx context.Context, ioType bnapi.IOType, writer iostat.WriterAtCtx) iostat.WriterAtCtx {
|
||||
func (qos *IoQueueQos) WriterAt(ctx context.Context, ioType bnapi.IOType, writer io.WriterAt) io.WriterAt {
|
||||
w := writer
|
||||
if ios := qos.getIostat(ioType); ios != nil {
|
||||
w = ios.WriterAtCtx(writer)
|
||||
w = ios.WriterAt(writer)
|
||||
}
|
||||
|
||||
if lmt := qos.getBpsLimiter(LimitType(ioType)); lmt != nil {
|
||||
w = &rateLimiter{
|
||||
ctx: ctx,
|
||||
wAtCtx: w,
|
||||
writerAt: w,
|
||||
bpsLimiter: lmt,
|
||||
}
|
||||
}
|
||||
|
||||
@ -4,7 +4,6 @@ import (
|
||||
"context"
|
||||
"io"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/dustin/go-humanize"
|
||||
@ -16,11 +15,20 @@ import (
|
||||
)
|
||||
|
||||
type mockWriteCtx struct {
|
||||
w io.WriterAt
|
||||
w io.Writer
|
||||
wa io.WriterAt
|
||||
}
|
||||
|
||||
func (q mockWriteCtx) Write(p []byte) (n int, err error) {
|
||||
return q.w.Write(p)
|
||||
}
|
||||
|
||||
func (q mockWriteCtx) WriteCtx(ctx context.Context, p []byte) (n int, err error) {
|
||||
return q.w.Write(p)
|
||||
}
|
||||
|
||||
func (q mockWriteCtx) WriteAtCtx(ctx context.Context, p []byte, off int64) (n int, err error) {
|
||||
return q.w.WriteAt(p, off)
|
||||
return q.wa.WriteAt(p, off)
|
||||
}
|
||||
|
||||
func TestNewQosManager(t *testing.T) {
|
||||
@ -95,19 +103,21 @@ func TestNewQosManager(t *testing.T) {
|
||||
|
||||
{
|
||||
// write
|
||||
mockW := &mockWriteCtx{w: f}
|
||||
|
||||
ok = q.TryAcquireIO(ctx, 1, IOTypeWrite)
|
||||
require.True(t, ok)
|
||||
defer q.ReleaseIO(1, IOTypeWrite)
|
||||
reader := strings.NewReader(ss)
|
||||
writer := qos.Writer(ctx, bnapi.NormalIO, f)
|
||||
n, err := io.Copy(writer, reader)
|
||||
require.Equal(t, int64(len(ss)), n)
|
||||
writer := qos.Writer(ctx, bnapi.NormalIO, mockW)
|
||||
data := []byte(ss)
|
||||
n, err := writer.Write(data)
|
||||
require.Equal(t, len(ss), n)
|
||||
require.NoError(t, err)
|
||||
|
||||
reader2 := strings.NewReader(ss)
|
||||
writer = qos.Writer(ctx, bnapi.BackgroundIO, f)
|
||||
n, err = io.Copy(writer, reader2)
|
||||
require.Equal(t, int64(len(ss)), n)
|
||||
data2 := []byte(ss)
|
||||
writer = qos.Writer(ctx, bnapi.BackgroundIO, mockW)
|
||||
n, err = writer.Write(data2)
|
||||
require.Equal(t, len(ss), n)
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
@ -153,11 +163,10 @@ func TestNewQosManager(t *testing.T) {
|
||||
fi, err := f.Stat()
|
||||
require.NoError(t, err)
|
||||
oldSize := fi.Size()
|
||||
mockW := &mockWriteCtx{w: f}
|
||||
|
||||
wt := qos.WriterAt(ctx, bnapi.NormalIO, mockW)
|
||||
wt := qos.WriterAt(ctx, bnapi.NormalIO, f)
|
||||
data := []byte("hello")
|
||||
_, err = wt.WriteAtCtx(ctx, data, oldSize)
|
||||
_, err = wt.WriteAt(data, oldSize)
|
||||
require.NoError(t, err)
|
||||
f.Sync()
|
||||
fi, err = f.Stat()
|
||||
@ -165,7 +174,7 @@ func TestNewQosManager(t *testing.T) {
|
||||
require.Equal(t, oldSize+int64(len(data)), fi.Size())
|
||||
|
||||
require.Panics(t, func() {
|
||||
wt = qos.WriterAt(ctx, bnapi.IOTypeMax, mockW)
|
||||
wt = qos.WriterAt(ctx, bnapi.IOTypeMax, f)
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@ -40,8 +40,8 @@ type RawFile interface {
|
||||
|
||||
type BlobFile interface {
|
||||
RawFile
|
||||
ReadAtCtx(ctx context.Context, b []byte, off int64) (n int, err error)
|
||||
WriteAtCtx(ctx context.Context, b []byte, off int64) (n int, err error)
|
||||
base.CtxReaderAt
|
||||
base.CtxWriterAt
|
||||
Allocate(off int64, size int64) (err error)
|
||||
Discard(off int64, size int64) (err error)
|
||||
SysStat() (sysstat syscall.Stat_t, err error)
|
||||
@ -89,16 +89,19 @@ func (ef *blobFile) WriteAt(b []byte, off int64) (n int, err error) {
|
||||
}
|
||||
|
||||
func (ef *blobFile) ReadAtCtx(ctx context.Context, b []byte, off int64) (n int, err error) {
|
||||
// If io ctx has been cancelled(ctx.Err() is not nil), we hope to return the error immediately
|
||||
if ctx.Err() != nil {
|
||||
return 0, ctx.Err()
|
||||
}
|
||||
|
||||
task := taskpool.IoPoolTaskArgs{
|
||||
BucketId: ef.chunk,
|
||||
Tm: time.Now(),
|
||||
Ctx: ctx,
|
||||
TaskFn: func() {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
if ctx.Err() != nil {
|
||||
n, err = 0, ctx.Err()
|
||||
return
|
||||
default:
|
||||
}
|
||||
n, err = ef.file.ReadAt(b, off)
|
||||
},
|
||||
@ -110,16 +113,18 @@ func (ef *blobFile) ReadAtCtx(ctx context.Context, b []byte, off int64) (n int,
|
||||
}
|
||||
|
||||
func (ef *blobFile) WriteAtCtx(ctx context.Context, b []byte, off int64) (n int, err error) {
|
||||
if ctx.Err() != nil {
|
||||
return 0, ctx.Err()
|
||||
}
|
||||
|
||||
task := taskpool.IoPoolTaskArgs{
|
||||
BucketId: ef.chunk,
|
||||
Tm: time.Now(),
|
||||
Ctx: ctx,
|
||||
TaskFn: func() {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
if ctx.Err() != nil {
|
||||
n, err = 0, ctx.Err()
|
||||
return
|
||||
default:
|
||||
}
|
||||
n, err = ef.file.WriteAt(b, off)
|
||||
},
|
||||
|
||||
@ -140,3 +140,79 @@ func TestBlobFile_Op(t *testing.T) {
|
||||
// phy allocate == 0
|
||||
require.Equal(t, int(stat.Blocks), 0)
|
||||
}
|
||||
|
||||
func TestBlobFile_doTaskFnCtxCancel(t *testing.T) {
|
||||
testDir, err := os.MkdirTemp(os.TempDir(), "BlobFileTaskCancel")
|
||||
require.NoError(t, err)
|
||||
defer os.RemoveAll(testDir)
|
||||
|
||||
posixfilepath := filepath.Join(testDir, "PoxsixFile")
|
||||
log.Info(posixfilepath)
|
||||
|
||||
temppath := filepath.Join(posixfilepath, "xxxtemp")
|
||||
f, err := OpenFile(temppath, false)
|
||||
require.Error(t, err)
|
||||
require.Nil(t, f)
|
||||
|
||||
f, err = OpenFile(posixfilepath, true)
|
||||
require.NoError(t, err)
|
||||
|
||||
require.NotNil(t, f)
|
||||
|
||||
// create
|
||||
syncWorker := mergetask.NewMergeTask(-1, func(interface{}) error { return nil })
|
||||
|
||||
ctr := gomock.NewController(t)
|
||||
ioPool := mocks.NewMockIoPool(ctr)
|
||||
ioPools := map[qos.IOTypeRW]taskpool.IoPool{
|
||||
qos.IOTypeRead: ioPool,
|
||||
qos.IOTypeWrite: ioPool,
|
||||
qos.IOTypeDel: ioPool,
|
||||
}
|
||||
|
||||
ef := blobFile{f, 1, syncWorker, nil, ioPools}
|
||||
fd := ef.Fd()
|
||||
require.NotNil(t, fd)
|
||||
|
||||
info, err := ef.Stat()
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, info)
|
||||
|
||||
data := []byte("test data")
|
||||
|
||||
// WriteAtCtx
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Do(func(args taskpool.IoPoolTaskArgs) {
|
||||
args.TaskFn()
|
||||
})
|
||||
n, err := ef.WriteAtCtx(ctx, data, 0)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, len(data), n)
|
||||
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Do(func(args taskpool.IoPoolTaskArgs) {
|
||||
cancel()
|
||||
args.TaskFn()
|
||||
})
|
||||
n, err = ef.WriteAtCtx(ctx, data, 0)
|
||||
require.ErrorIs(t, context.Canceled, err)
|
||||
require.Equal(t, 0, n)
|
||||
|
||||
// ReadAtCtx
|
||||
ctx, cancel = context.WithCancel(context.Background())
|
||||
buf := make([]byte, len(data))
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Do(func(args taskpool.IoPoolTaskArgs) {
|
||||
args.TaskFn()
|
||||
})
|
||||
n, err = ef.ReadAtCtx(ctx, buf, 0)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, len(data), n)
|
||||
require.Equal(t, data, buf)
|
||||
|
||||
ioPool.EXPECT().Submit(gomock.Any()).Do(func(args taskpool.IoPoolTaskArgs) {
|
||||
cancel()
|
||||
args.TaskFn()
|
||||
})
|
||||
n, err = ef.ReadAtCtx(ctx, buf, 0)
|
||||
require.ErrorIs(t, context.Canceled, err)
|
||||
require.Equal(t, 0, n)
|
||||
}
|
||||
|
||||
@ -17,6 +17,7 @@ package chunk
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
@ -545,6 +546,10 @@ func (cs *chunk) rangeRead(ctx context.Context, stg core.Storage, s *core.Shard,
|
||||
span.AppendTrackLogWithDuration("net.w", tw.Duration(), err)
|
||||
span.AppendTrackLogWithDuration("dat.r", tr.Duration(), err)
|
||||
if err != nil {
|
||||
// prevent 5xx error code
|
||||
if errors.Is(err, context.Canceled) {
|
||||
err = bloberr.ErrIOCtxCancel
|
||||
}
|
||||
return n, err
|
||||
}
|
||||
|
||||
|
||||
@ -333,7 +333,7 @@ func (cd *datafile) Write(ctx context.Context, shard *core.Shard) (err error) {
|
||||
defer recycle()
|
||||
|
||||
// prepare reader and writer
|
||||
w := &bncomm.Writer{WriterAt: cd.ef, Offset: pos}
|
||||
w := &bncomm.WriterWithCtx{Offset: pos, Wt: cd.ef, Ctx: ctx}
|
||||
twRaw := bncomm.NewTimeWriter(w)
|
||||
|
||||
qosw := cd.qosWriter(ctx, twRaw)
|
||||
@ -364,9 +364,12 @@ func (cd *datafile) Write(ctx context.Context, shard *core.Shard) (err error) {
|
||||
buf := buffer[core.HeaderSize : len(buffer)-core.FooterSize]
|
||||
n, err := encoder.Read(buf)
|
||||
if err != nil {
|
||||
// prevent 5xx error code
|
||||
if _, is := err.(crc32block.ReaderError); is {
|
||||
span.Warnf("write shard:%+v -> %s", shard, err.Error())
|
||||
err = bloberr.ErrReaderError
|
||||
} else if errors.Is(err, context.Canceled) {
|
||||
err = bloberr.ErrIOCtxCancel
|
||||
}
|
||||
return err
|
||||
}
|
||||
@ -396,6 +399,9 @@ func (cd *datafile) Write(ctx context.Context, shard *core.Shard) (err error) {
|
||||
// write header+data+footer; header+data, data..., data+footer
|
||||
n, err = tw.Write(buf)
|
||||
if err != nil {
|
||||
if errors.Is(err, context.Canceled) {
|
||||
err = bloberr.ErrIOCtxCancel
|
||||
}
|
||||
return err
|
||||
}
|
||||
if n != len(buf) {
|
||||
@ -423,7 +429,8 @@ func (cd *datafile) Read(ctx context.Context, shard *core.Shard, from, to uint32
|
||||
pos := shard.Offset + core.GetShardHeaderSize()
|
||||
|
||||
// new reader
|
||||
iosr := cd.qosReaderAt(ctx, cd.ef)
|
||||
ra := &bncomm.ReaderWithCtx{Rd: cd.ef, Ctx: ctx}
|
||||
iosr := cd.qosReaderAt(ctx, ra)
|
||||
|
||||
// new buffer
|
||||
buffer := bytespool.Alloc(core.CrcBlockUnitSize)
|
||||
|
||||
@ -31,11 +31,13 @@ import (
|
||||
|
||||
bnapi "github.com/cubefs/cubefs/blobstore/api/blobnode"
|
||||
"github.com/cubefs/cubefs/blobstore/api/clustermgr"
|
||||
"github.com/cubefs/cubefs/blobstore/blobnode/base"
|
||||
"github.com/cubefs/cubefs/blobstore/blobnode/base/qos"
|
||||
"github.com/cubefs/cubefs/blobstore/blobnode/core"
|
||||
"github.com/cubefs/cubefs/blobstore/common/crc32block"
|
||||
bloberr "github.com/cubefs/cubefs/blobstore/common/errors"
|
||||
"github.com/cubefs/cubefs/blobstore/common/proto"
|
||||
bnmock "github.com/cubefs/cubefs/blobstore/testing/mockblobnode"
|
||||
"github.com/cubefs/cubefs/blobstore/testing/mocks"
|
||||
_ "github.com/cubefs/cubefs/blobstore/testing/nolog"
|
||||
"github.com/cubefs/cubefs/blobstore/util/bytespool"
|
||||
@ -947,3 +949,115 @@ func TestChunkHeader(t *testing.T) {
|
||||
s := chunkHeader.String()
|
||||
require.NotNil(t, s)
|
||||
}
|
||||
|
||||
func TestChunkData_WriteReadCancel(t *testing.T) {
|
||||
testDir, err := os.MkdirTemp(os.TempDir(), defaultDiskTestDir+"WriteCancel")
|
||||
require.NoError(t, err)
|
||||
defer os.RemoveAll(testDir)
|
||||
|
||||
ctx := context.Background()
|
||||
chunkname := clustermgr.NewChunkID(0).String()
|
||||
chunkname = filepath.Join(testDir, chunkname)
|
||||
log.Info(chunkname)
|
||||
|
||||
diskConfig := &core.Config{
|
||||
BaseConfig: core.BaseConfig{Path: testDir},
|
||||
RuntimeConfig: core.RuntimeConfig{
|
||||
BlockBufferSize: 64 * 1024,
|
||||
},
|
||||
}
|
||||
|
||||
ioPools := newIoPoolMock(t)
|
||||
ioQos, _ := qos.NewIoQueueQos(qos.Config{ReadQueueDepth: 2, WriteQueueDepth: 2, WriteChanQueCnt: 2})
|
||||
defer ioQos.Close()
|
||||
cd, err := NewChunkData(ctx, core.VuidMeta{}, chunkname, diskConfig, true, ioQos, ioPools)
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, cd)
|
||||
defer cd.Close()
|
||||
|
||||
// mock
|
||||
backup := cd.ef
|
||||
ctr := gomock.NewController(t)
|
||||
cd.ef = bnmock.NewMockBlobFile(ctr)
|
||||
a := gomock.Any()
|
||||
|
||||
log.Infof("chunkdata: \n%s", cd)
|
||||
require.Equal(t, int32(cd.wOff), int32(4096))
|
||||
sharddata := []byte("test data")
|
||||
|
||||
// build shard data
|
||||
shard := &core.Shard{
|
||||
Bid: 5,
|
||||
Vuid: 10,
|
||||
Flag: bnapi.ShardStatusNormal,
|
||||
Size: uint32(len(sharddata)),
|
||||
Body: bytes.NewBuffer(sharddata),
|
||||
}
|
||||
|
||||
// write ok, size 9.
|
||||
cd.ef.(*bnmock.MockBlobFile).EXPECT().WriteAtCtx(a, a, a).DoAndReturn(func(ctx context.Context, b []byte, off int64) (n int, err error) {
|
||||
return len(b), nil
|
||||
})
|
||||
err = cd.Write(ctx, shard)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int32(shard.Offset), int32(4096))
|
||||
require.Equal(t, int32(cd.wOff), int32(8192))
|
||||
|
||||
// fail, ctx cancel, before enqueue
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
shard2 := &core.Shard{
|
||||
Bid: 6,
|
||||
Vuid: 10,
|
||||
Flag: bnapi.ShardStatusNormal,
|
||||
Size: uint32(len(sharddata)),
|
||||
Body: bytes.NewBuffer(sharddata),
|
||||
}
|
||||
|
||||
cancel()
|
||||
cd.ef.(*bnmock.MockBlobFile).EXPECT().WriteAtCtx(a, a, a).DoAndReturn(func(ctx context.Context, b []byte, off int64) (n int, err error) {
|
||||
return 0, context.Canceled
|
||||
})
|
||||
err = cd.Write(ctx, shard2)
|
||||
require.NotNil(t, err)
|
||||
|
||||
// fail, ctx cancel, after dequeue
|
||||
ctx, cancel = context.WithCancel(context.Background())
|
||||
shard2.Body = bytes.NewBuffer(sharddata)
|
||||
|
||||
cd.ef.(*bnmock.MockBlobFile).EXPECT().WriteAtCtx(ctx, a, a).DoAndReturn(func(ctx context.Context, b []byte, off int64) (n int, err error) {
|
||||
cancel()
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
n, err = 0, ctx.Err()
|
||||
return
|
||||
default:
|
||||
}
|
||||
return len(b), nil
|
||||
})
|
||||
err = cd.Write(ctx, shard2)
|
||||
require.NotNil(t, err)
|
||||
require.ErrorIs(t, err, bloberr.ErrIOCtxCancel)
|
||||
|
||||
// read fail, cancel
|
||||
ctx, cancel = context.WithCancel(context.Background())
|
||||
readBuf := bytes.NewBuffer(nil)
|
||||
shard.Writer = readBuf
|
||||
|
||||
rc, err := cd.Read(ctx, shard, 0, shard.Size)
|
||||
require.NoError(t, err)
|
||||
|
||||
tw := base.NewTimeWriter(shard.Writer)
|
||||
tr := base.NewTimeReader(rc)
|
||||
|
||||
cancel()
|
||||
cd.ef.(*bnmock.MockBlobFile).EXPECT().ReadAtCtx(a, a, a).DoAndReturn(func(ctx context.Context, b []byte, off int64) (n int, err error) {
|
||||
return 0, context.Canceled
|
||||
}).Times(1)
|
||||
n, err := io.CopyN(tw, tr, int64(len(sharddata)))
|
||||
require.ErrorIs(t, err, context.Canceled)
|
||||
require.Equal(t, int64(0), n)
|
||||
|
||||
// resume
|
||||
cd.ef = backup
|
||||
}
|
||||
|
||||
@ -58,6 +58,7 @@ const (
|
||||
CodeRequestLimited = 673
|
||||
CodeUnsupportedTaskCodeMode = 674
|
||||
CodePutShardTimeout = 675
|
||||
CodeIOCtxCancel = 676
|
||||
)
|
||||
|
||||
var (
|
||||
@ -102,6 +103,7 @@ var (
|
||||
ErrRequestLimited = Error(CodeRequestLimited)
|
||||
ErrUnsupportedTaskCodeMode = Error(CodeUnsupportedTaskCodeMode)
|
||||
ErrPutShardTimeout = Error(CodePutShardTimeout)
|
||||
ErrIOCtxCancel = Error(CodeIOCtxCancel)
|
||||
)
|
||||
|
||||
var ErrShardMayBeLost = errors.New("shard may be lost")
|
||||
|
||||
@ -156,6 +156,7 @@ var errCodeMap = map[int]string{
|
||||
|
||||
CodeUnsupportedTaskCodeMode: "unsupported task codemode",
|
||||
CodePutShardTimeout: "put shard timeout",
|
||||
CodeIOCtxCancel: "io context cancel",
|
||||
|
||||
CodeShardNodeNotLeader: "shardnode:not leader",
|
||||
CodeShardRangeMismatch: "shardnode:range mismatch",
|
||||
|
||||
@ -18,7 +18,6 @@
|
||||
package iostat
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
@ -67,14 +66,6 @@ type StatData struct {
|
||||
Await int64 // nanoseconds
|
||||
}
|
||||
|
||||
type ReaderAtCtx interface {
|
||||
ReadAtCtx(ctx context.Context, p []byte, off int64) (n int, err error)
|
||||
}
|
||||
|
||||
type WriterAtCtx interface {
|
||||
WriteAtCtx(ctx context.Context, p []byte, off int64) (n int, err error)
|
||||
}
|
||||
|
||||
type Iostat interface {
|
||||
Get(out *Stat)
|
||||
Begin(size uint64)
|
||||
@ -90,8 +81,6 @@ type StatMgrAPI interface {
|
||||
WriterAt(underlying io.WriterAt) io.WriterAt
|
||||
Reader(underlying io.Reader) io.Reader
|
||||
ReaderAt(underlying io.ReaderAt) io.ReaderAt
|
||||
ReaderAtCtx(underlying ReaderAtCtx) ReaderAtCtx
|
||||
WriterAtCtx(underlying WriterAtCtx) WriterAtCtx
|
||||
}
|
||||
|
||||
func (sm *StatMgr) ReadBegin(size uint64) {
|
||||
|
||||
@ -15,7 +15,6 @@
|
||||
package iostat
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"time"
|
||||
)
|
||||
@ -35,11 +34,6 @@ type iostatReaderAt struct {
|
||||
ios *StatMgr
|
||||
}
|
||||
|
||||
type iostatReaderAtCtx struct {
|
||||
underlying ReaderAtCtx
|
||||
ios *StatMgr
|
||||
}
|
||||
|
||||
func (ior *iostatReaderAt) ReadAt(p []byte, off int64) (n int, err error) {
|
||||
ior.ios.ReadBegin(uint64(len(p)))
|
||||
defer ior.ios.ReadEnd(time.Now())
|
||||
@ -69,14 +63,6 @@ func (ior *iostatReadCloser) Close() error {
|
||||
return ior.underlying.Close()
|
||||
}
|
||||
|
||||
func (ior *iostatReaderAtCtx) ReadAtCtx(ctx context.Context, p []byte, off int64) (n int, err error) {
|
||||
ior.ios.ReadBegin(uint64(len(p)))
|
||||
defer ior.ios.ReadEnd(time.Now())
|
||||
|
||||
n, err = ior.underlying.ReadAtCtx(ctx, p, off)
|
||||
return
|
||||
}
|
||||
|
||||
func (sm *StatMgr) Reader(underlying io.Reader) io.Reader {
|
||||
return &iostatReader{
|
||||
underlying: underlying,
|
||||
@ -91,13 +77,6 @@ func (sm *StatMgr) ReaderAt(underlying io.ReaderAt) io.ReaderAt {
|
||||
}
|
||||
}
|
||||
|
||||
func (sm *StatMgr) ReaderAtCtx(underlying ReaderAtCtx) ReaderAtCtx {
|
||||
return &iostatReaderAtCtx{
|
||||
underlying: underlying,
|
||||
ios: sm,
|
||||
}
|
||||
}
|
||||
|
||||
func (sm *StatMgr) ReaderCloser(underlying io.ReadCloser) io.ReadCloser {
|
||||
return &iostatReadCloser{
|
||||
underlying: underlying,
|
||||
|
||||
@ -15,7 +15,6 @@
|
||||
package iostat
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"time"
|
||||
)
|
||||
@ -35,11 +34,6 @@ type iostatWriterAt struct {
|
||||
ios *StatMgr
|
||||
}
|
||||
|
||||
type iostatWriterAtCtx struct {
|
||||
underlying WriterAtCtx
|
||||
ios *StatMgr
|
||||
}
|
||||
|
||||
func (iow *iostatWriter) Write(p []byte) (written int, err error) {
|
||||
iow.ios.WriteBegin(uint64(len(p)))
|
||||
defer iow.ios.WriteEnd(time.Now())
|
||||
@ -71,15 +65,6 @@ func (iow *iostatWriteCloser) Close() error {
|
||||
return iow.underlying.Close()
|
||||
}
|
||||
|
||||
func (iow *iostatWriterAtCtx) WriteAtCtx(ctx context.Context, p []byte, off int64) (n int, err error) {
|
||||
iow.ios.WriteBegin(uint64(len(p)))
|
||||
defer iow.ios.WriteEnd(time.Now())
|
||||
|
||||
n, err = iow.underlying.WriteAtCtx(ctx, p, off)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (sm *StatMgr) Writer(underlying io.Writer) io.Writer {
|
||||
return &iostatWriter{
|
||||
underlying: underlying,
|
||||
@ -94,13 +79,6 @@ func (sm *StatMgr) WriterAt(underlying io.WriterAt) io.WriterAt {
|
||||
}
|
||||
}
|
||||
|
||||
func (sm *StatMgr) WriterAtCtx(underlying WriterAtCtx) WriterAtCtx {
|
||||
return &iostatWriterAtCtx{
|
||||
underlying: underlying,
|
||||
ios: sm,
|
||||
}
|
||||
}
|
||||
|
||||
func (sm *StatMgr) WriteCloser(underlying io.WriteCloser) io.WriteCloser {
|
||||
return &iostatWriteCloser{
|
||||
underlying: underlying,
|
||||
|
||||
211
blobstore/testing/mockblobnode/blobfile.go
Normal file
211
blobstore/testing/mockblobnode/blobfile.go
Normal file
@ -0,0 +1,211 @@
|
||||
// Code generated by MockGen. DO NOT EDIT.
|
||||
// Source: github.com/cubefs/cubefs/blobstore/blobnode/core (interfaces: BlobFile)
|
||||
|
||||
// Package mock is a generated GoMock package.
|
||||
package mock
|
||||
|
||||
import (
|
||||
context "context"
|
||||
fs "io/fs"
|
||||
reflect "reflect"
|
||||
syscall "syscall"
|
||||
|
||||
gomock "github.com/golang/mock/gomock"
|
||||
)
|
||||
|
||||
// MockBlobFile is a mock of BlobFile interface.
|
||||
type MockBlobFile struct {
|
||||
ctrl *gomock.Controller
|
||||
recorder *MockBlobFileMockRecorder
|
||||
}
|
||||
|
||||
// MockBlobFileMockRecorder is the mock recorder for MockBlobFile.
|
||||
type MockBlobFileMockRecorder struct {
|
||||
mock *MockBlobFile
|
||||
}
|
||||
|
||||
// NewMockBlobFile creates a new mock instance.
|
||||
func NewMockBlobFile(ctrl *gomock.Controller) *MockBlobFile {
|
||||
mock := &MockBlobFile{ctrl: ctrl}
|
||||
mock.recorder = &MockBlobFileMockRecorder{mock}
|
||||
return mock
|
||||
}
|
||||
|
||||
// EXPECT returns an object that allows the caller to indicate expected use.
|
||||
func (m *MockBlobFile) EXPECT() *MockBlobFileMockRecorder {
|
||||
return m.recorder
|
||||
}
|
||||
|
||||
// Allocate mocks base method.
|
||||
func (m *MockBlobFile) Allocate(arg0, arg1 int64) error {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "Allocate", arg0, arg1)
|
||||
ret0, _ := ret[0].(error)
|
||||
return ret0
|
||||
}
|
||||
|
||||
// Allocate indicates an expected call of Allocate.
|
||||
func (mr *MockBlobFileMockRecorder) Allocate(arg0, arg1 interface{}) *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Allocate", reflect.TypeOf((*MockBlobFile)(nil).Allocate), arg0, arg1)
|
||||
}
|
||||
|
||||
// Close mocks base method.
|
||||
func (m *MockBlobFile) Close() error {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "Close")
|
||||
ret0, _ := ret[0].(error)
|
||||
return ret0
|
||||
}
|
||||
|
||||
// Close indicates an expected call of Close.
|
||||
func (mr *MockBlobFileMockRecorder) Close() *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Close", reflect.TypeOf((*MockBlobFile)(nil).Close))
|
||||
}
|
||||
|
||||
// Discard mocks base method.
|
||||
func (m *MockBlobFile) Discard(arg0, arg1 int64) error {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "Discard", arg0, arg1)
|
||||
ret0, _ := ret[0].(error)
|
||||
return ret0
|
||||
}
|
||||
|
||||
// Discard indicates an expected call of Discard.
|
||||
func (mr *MockBlobFileMockRecorder) Discard(arg0, arg1 interface{}) *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Discard", reflect.TypeOf((*MockBlobFile)(nil).Discard), arg0, arg1)
|
||||
}
|
||||
|
||||
// Fd mocks base method.
|
||||
func (m *MockBlobFile) Fd() uintptr {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "Fd")
|
||||
ret0, _ := ret[0].(uintptr)
|
||||
return ret0
|
||||
}
|
||||
|
||||
// Fd indicates an expected call of Fd.
|
||||
func (mr *MockBlobFileMockRecorder) Fd() *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Fd", reflect.TypeOf((*MockBlobFile)(nil).Fd))
|
||||
}
|
||||
|
||||
// Name mocks base method.
|
||||
func (m *MockBlobFile) Name() string {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "Name")
|
||||
ret0, _ := ret[0].(string)
|
||||
return ret0
|
||||
}
|
||||
|
||||
// Name indicates an expected call of Name.
|
||||
func (mr *MockBlobFileMockRecorder) Name() *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Name", reflect.TypeOf((*MockBlobFile)(nil).Name))
|
||||
}
|
||||
|
||||
// ReadAt mocks base method.
|
||||
func (m *MockBlobFile) ReadAt(arg0 []byte, arg1 int64) (int, error) {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "ReadAt", arg0, arg1)
|
||||
ret0, _ := ret[0].(int)
|
||||
ret1, _ := ret[1].(error)
|
||||
return ret0, ret1
|
||||
}
|
||||
|
||||
// ReadAt indicates an expected call of ReadAt.
|
||||
func (mr *MockBlobFileMockRecorder) ReadAt(arg0, arg1 interface{}) *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ReadAt", reflect.TypeOf((*MockBlobFile)(nil).ReadAt), arg0, arg1)
|
||||
}
|
||||
|
||||
// ReadAtCtx mocks base method.
|
||||
func (m *MockBlobFile) ReadAtCtx(arg0 context.Context, arg1 []byte, arg2 int64) (int, error) {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "ReadAtCtx", arg0, arg1, arg2)
|
||||
ret0, _ := ret[0].(int)
|
||||
ret1, _ := ret[1].(error)
|
||||
return ret0, ret1
|
||||
}
|
||||
|
||||
// ReadAtCtx indicates an expected call of ReadAtCtx.
|
||||
func (mr *MockBlobFileMockRecorder) ReadAtCtx(arg0, arg1, arg2 interface{}) *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "ReadAtCtx", reflect.TypeOf((*MockBlobFile)(nil).ReadAtCtx), arg0, arg1, arg2)
|
||||
}
|
||||
|
||||
// Stat mocks base method.
|
||||
func (m *MockBlobFile) Stat() (fs.FileInfo, error) {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "Stat")
|
||||
ret0, _ := ret[0].(fs.FileInfo)
|
||||
ret1, _ := ret[1].(error)
|
||||
return ret0, ret1
|
||||
}
|
||||
|
||||
// Stat indicates an expected call of Stat.
|
||||
func (mr *MockBlobFileMockRecorder) Stat() *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Stat", reflect.TypeOf((*MockBlobFile)(nil).Stat))
|
||||
}
|
||||
|
||||
// Sync mocks base method.
|
||||
func (m *MockBlobFile) Sync() error {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "Sync")
|
||||
ret0, _ := ret[0].(error)
|
||||
return ret0
|
||||
}
|
||||
|
||||
// Sync indicates an expected call of Sync.
|
||||
func (mr *MockBlobFileMockRecorder) Sync() *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "Sync", reflect.TypeOf((*MockBlobFile)(nil).Sync))
|
||||
}
|
||||
|
||||
// SysStat mocks base method.
|
||||
func (m *MockBlobFile) SysStat() (syscall.Stat_t, error) {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "SysStat")
|
||||
ret0, _ := ret[0].(syscall.Stat_t)
|
||||
ret1, _ := ret[1].(error)
|
||||
return ret0, ret1
|
||||
}
|
||||
|
||||
// SysStat indicates an expected call of SysStat.
|
||||
func (mr *MockBlobFileMockRecorder) SysStat() *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "SysStat", reflect.TypeOf((*MockBlobFile)(nil).SysStat))
|
||||
}
|
||||
|
||||
// WriteAt mocks base method.
|
||||
func (m *MockBlobFile) WriteAt(arg0 []byte, arg1 int64) (int, error) {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "WriteAt", arg0, arg1)
|
||||
ret0, _ := ret[0].(int)
|
||||
ret1, _ := ret[1].(error)
|
||||
return ret0, ret1
|
||||
}
|
||||
|
||||
// WriteAt indicates an expected call of WriteAt.
|
||||
func (mr *MockBlobFileMockRecorder) WriteAt(arg0, arg1 interface{}) *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "WriteAt", reflect.TypeOf((*MockBlobFile)(nil).WriteAt), arg0, arg1)
|
||||
}
|
||||
|
||||
// WriteAtCtx mocks base method.
|
||||
func (m *MockBlobFile) WriteAtCtx(arg0 context.Context, arg1 []byte, arg2 int64) (int, error) {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "WriteAtCtx", arg0, arg1, arg2)
|
||||
ret0, _ := ret[0].(int)
|
||||
ret1, _ := ret[1].(error)
|
||||
return ret0, ret1
|
||||
}
|
||||
|
||||
// WriteAtCtx indicates an expected call of WriteAtCtx.
|
||||
func (mr *MockBlobFileMockRecorder) WriteAtCtx(arg0, arg1, arg2 interface{}) *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "WriteAtCtx", reflect.TypeOf((*MockBlobFile)(nil).WriteAtCtx), arg0, arg1, arg2)
|
||||
}
|
||||
18
blobstore/testing/mockblobnode/mock.go
Normal file
18
blobstore/testing/mockblobnode/mock.go
Normal file
@ -0,0 +1,18 @@
|
||||
// Copyright 2022 The CubeFS Authors.
|
||||
//
|
||||
// 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 mock
|
||||
|
||||
// github.com/cubefs/cubefs/blobstore/blobnode/... module blobnode interfaces
|
||||
//go:generate mockgen -destination=./blobfile.go -package=mock -mock_names BlobFile=MockBlobFile github.com/cubefs/cubefs/blobstore/blobnode/core BlobFile
|
||||
Loading…
Reference in New Issue
Block a user